MessageStore 与 MessageCommitter 契约¶
依赖:shared.md。MessageStore 是提交事实源;MessageCommitter 是唯一能把成功提交映射为 SEND_ACK 的边界。
pub trait MessageStore: Send + Sync {
async fn reserve(
&self,
request: &AuthenticatedSend,
delivery: FrozenDeliveryInput,
candidate: ReserveCandidate,
) -> Result<ReserveResult, MessageStoreError>;
async fn put_record(
&self,
intent: &CommitIntent,
) -> Result<(), MessageStoreError>;
async fn mark_outbox_published(
&self,
intent: &CommitIntent,
publication: OutboxPublication,
) -> Result<(), MessageStoreError>;
async fn mark_published_and_committed(
&self,
intent: &CommitIntent,
publication: OutboxPublication,
) -> Result<CommittedMessage, MessageStoreError>;
async fn mark_committed(
&self,
intent: &CommitIntent,
) -> Result<CommittedMessage, MessageStoreError>;
async fn load_commit(
&self,
key: &CommitKey,
conversation_id: ConversationId,
) -> Result<Option<LoadedCommit>, MessageStoreError>;
async fn load_commit_for_request(
&self,
request: &AuthenticatedSend,
) -> Result<Option<LoadedCommit>, MessageStoreError>;
async fn scan_recoverable_commits(
&self,
older_than: TimestampMillis,
max_items: u32,
) -> Result<Vec<CommitIntent>, MessageStoreError>;
async fn history_page(
&self,
query: HistoryQuery,
) -> Result<HistoryPage, MessageStoreError>;
async fn readiness(&self) -> Result<(), MessageStoreError>;
}
pub trait MessageCommitter: Send + Sync {
async fn commit(
&self,
request: &AuthenticatedSend,
delivery: FrozenDeliveryInput,
candidate: ReserveCandidate,
) -> Result<CommittedMessage, CommitError>;
async fn resume(&self, key: &CommitKey) -> Result<CommittedMessage, CommitError>;
async fn resume_existing(
&self,
request: &AuthenticatedSend,
) -> Result<Option<CommittedMessage>, CommitError>;
async fn recover_due(
&self,
older_than: TimestampMillis,
max_items: u32,
) -> Result<u32, CommitError>;
async fn readiness(&self) -> Result<(), CommitError>;
}
前置条件¶
commit的AuthenticatedSend已携带网关认证身份;写入者仍须完成发送者成员准入、内容与配额检查。无效ClientMessageId必须返回CommitError::InvalidClientMessageId。reserve之前固定FrozenDeliveryInput与完整 canonicalMessageRecord;不得用恢复时的当前成员关系或客户端再次上传的 payload 替换它。load_commit_for_request必须比较持久的规范请求身份;相同(tenant,sender,cmid)但 device、conversation、content、mention 或 reply 任一不同均返回冲突,绝不能回放旧 ACK。- 已提交的消息只承诺回放
SEND_ACK所需的坐标:reserve命中返回ReserveResult::Committed,load_commit*返回LoadedCommit::Committed。后端可在进入 COMMITTED 时丢弃记录本体(此时它已在 会话历史与 Outbox 中);Redis 把提交 hash 压缩为身份摘要 + 坐标 + publication(2026-09-28, 100 B 正文每条 2,640 B → 552 B 且保持 listpack),升级前的完整格式照常读取。持久请求身份可存 域分隔摘要(AuthenticatedSend::reserve_identity_digest,头QIMD);读到升级前的完整身份(头QIMC)时逐字节比较。 put_record只能写CommitIntent.record;history_page必须先对query.actor执行会话成员授权。readiness必须真实访问生产依赖;参数为零的本地早退、lazy client/producer 构造成功都不是就绪证据。
后置条件¶
- 同一
CommitKey的首次和重复reserve返回相同坐标,且重复请求不消耗新的消息或会话序号。 reserve成功持久化的 intent 单独包含恢复后续步骤所需的全部消息事实;客户端永不重试也不能形成不可恢复 RESERVED。put_record与mark_outbox_published可安全重试并收敛为同一值;mark_committed只有在两者均已有可恢复证据时才能返回CommittedMessage。mark_published_and_committed接受 durableOutboxPublication,并在成功时返回该冻结 intent 的CommittedMessage。Memory 与 Redis 必须将Recorded + publication → Committed作为一个后端原子迁移:同值重试回放原 ACK,异 publication、缺少 record 或错误前态均不得留下部分写入;成功后必须移除恢复候选并登记应有的 history expiry。- 后端没有此原子能力时,默认实现依次调用
mark_outbox_published与mark_committed。若第二步失败,必须保留Outboxed及 publication 作为可恢复事实,随后由resume的Outboxed分支调用mark_committed。Scylla 当前采用该默认两步条件写,不宣称单次 LWT 或跨表原子提交。 MessageCommitter在Recorded分支取得 Outbox durable ack 并经过 M-COMMIT-04 故障点后调用mark_published_and_committed;已持久化的Outboxed仍只调用mark_committed。调用响应丢失或并发推进时,必须重读事实源并只向前收敛。commit或resume命中中间态必须从持久状态继续;命中Committed只回放原CommittedMessage。CommitRecovery必须调用同一resume语义。- 成功
CommittedMessage是且仅是生成SEND_ACK的依据;它不表示邮箱物化、推送或已读。 - 历史页保留
MessageType、custom_type、正文引用、方向、锚点、字节上限和可用边界;历史空洞不能被解释为丢失。
错误¶
- 存储边界只暴露
MessageStoreError,不得透出具体后端条件写、键、表或脚本错误。 - 提交边界只暴露
CommitError;Recoverable { key }表示可由同一幂等键继续,绝不能伪造成功 ACK。