Outbox Publisher
PostgreSQL Outbox 的 Claim、Lease、投递尝试、消费回执、隔离与恢复契约
首期 Outbox Publisher 是OceanWay内部模块共用的可靠投递实现,并按Producer事实库部署独立Worker Role:Core Admission/Billing经济边界Event投递给Operations Projector;Execution Eligibility与Gateway Availability投递给Metering Input Processor;Execution Closure及Billing Finalization Reconciliation Resolution投递给Billing Finalization Coordinator;Metering Settlement Input投递给Billing Catch-up Consumer。每个实例只从自己的PostgreSQL事实边界读取已提交事件、调用冻结的受控Consumer并记录可重试结果;不执行领域命令,也不把投递成功等同于业务成功。各链复用同一Contracts状态机和Golden,不共享Lease、游标或数据库事务。
为什么首期不上外部 Broker
当前只有四个边界清晰、用途固定的内部Consumer:Operations Projector、Metering Input Processor、Billing Finalization Coordinator与Billing Settlement Catch-up。各Owner PostgreSQL都能在本地领域事务内原子写自身Outbox;跨Owner只使用受控内部Transport和at-least-once重投,不伪造跨库事务。使用同一Contracts实现的独立Worker Role即可验证事件版本、幂等、跨库重投、重建和运维语义,同时避免提前引入第二套消息集群、权限、保留和故障恢复系统。Consumer提交后响应丢失时,源Dispatcher重投相同Event并由Owner Receipt幂等收敛。
这不是永久禁止 Broker。出现以下经过度量的条件时应另立 ADR:
- 多个独立团队或故障域需要各自消费和回放;
- PostgreSQL Claim/Retention 已影响在线事务目标;
- 跨区域复制、长保留或消费吞吐超过容量验证边界;
- 需要 Broker 原生分区、订阅隔离或流式生态,并且收益超过运维成本。
引入 Broker 时不改变 Canonical Event Envelope;Outbox Dispatcher 改为向 Broker 发布,Consumer Receipt 和 Projection 幂等语义继续保留。
存储边界
每个Producer事实库中的Event Payload在受治理保留期内只以不可变形式保存在自己的 outbox_events 中;离开热表只能经过受审计 Retention/Archive Procedure,绝不原地改写。不同事实库不得共享表、Lease或published_at,但必须实现同一严格逻辑结构:
outbox_events
event_id / event_type / schema_version
delivery_set_version / delivery_set_definition_schema_version
delivery_set_definition_digest_algorithm_version / delivery_set_definition_digest
aggregate_id / aggregate_revision
tenant and correlation envelope
event_payload / event_envelope_digest_algorithm_version / event_envelope_sha256
occurred_at / created_at
delivery_set_versions
delivery_set_version / delivery_set_definition_schema_version
delivery_set_definition_digest_algorithm_version / delivery_set_definition_digest
status / activated_at / source_release
delivery_set_members
delivery_set_version / event_type / schema_version
destination / consumer_name / mandatory
expected_receipt_type / expected_receipt_schema_version
projection_target =
{ state=not_applicable }
| { state=required; projection_name; projection_version }
outbox_deliveries
event_id / delivery_set_version / destination / consumer_name
expected_receipt_type / expected_receipt_schema_version
required_projection_name? / required_projection_version?
mandatory
status = pending | leased | acknowledged | quarantined
available_at
lease_owner / lease_token / lease_expires_at
current_attempt_id
attempt_count
last_attempt_at / delivered_at
last_error_code / last_error_ref
ack_consumer_name? / ack_receipt_type? / ack_receipt_key?
ack_event_envelope_digest_algorithm_version? / ack_event_envelope_sha256?
outbox_delivery_attempt_starts
attempt_id / destination / event_id
lease_token
started_at
outbox_delivery_attempt_finishes
attempt_id # 唯一外键到 Start
finished_at
outcome / retry_disposition
normalized_error_code / sanitized_error_ref
outbox_quarantine
destination / event_id
reason / event_envelope_digest_algorithm_version / event_envelope_sha256
quarantined_at / acknowledged_at?
case_id?具体表名可在迁移中调整,但必须保留以下原则:
- Event Envelope/Payload 不因重试、投递或人工处置而更新。
- 每个 Destination 拥有独立 Delivery State;不可变 Delivery Set Member 冻结
destination → consumerName映射,未来增加 Consumer 不复用一个全局“已发布”布尔值。 - Attempt 生命周期由不可变 Start 与至多一个 Finish 两类事实表达;不能原地补写终态,也不能只覆盖最后一次错误。
- 隔离是明确状态,不是把失败事件标记为成功或删除。
- Producer Append 在同一业务事务内冻结
deliverySetVersion + Definition Schema Version + Digest Algorithm Version + Digest,并按该不可变定义精确创建所有 Mandatory Delivery Row。 - Dispatcher 永远 Claim 显式未完成 Delivery Row,不扫描“可能尚未创建 Delivery 的事件”。
Delivery State是封闭四值:Append创建pending;合法Claim原子迁移为leased;retryable只是Attempt Outcome,完成当前Attempt后把Delivery迁回pending并更新配置化availableAt;成功Ack迁移为终态acknowledged;明确不可处理迁移为终态quarantined。不存在retryable|delivered|dead_letter等额外持久状态,任何未知值拒绝。只有pending或Lease已过期且完成恢复手续的leased可重新Claim。
Append 原子性与 Delivery Set
每次领域事务 Append Event 时必须在同一 PostgreSQL 事务中完成:
- 从 PostgreSQL 中读取已激活且不可变的 Delivery Set Version/Member;部署配置只能选择已登记版本,不能成为数据库外的集合事实源。
- 把精确
deliverySetVersion + deliverySetDefinitionSchemaVersion + deliverySetDefinitionDigestAlgorithmVersion + deliverySetDefinitionDigest固定到 Event Row。 - 写入不可变 Envelope、Payload,以及覆盖二者的 Contracts 固定
eventEnvelopeDigestAlgorithmVersion + eventEnvelopeSha256。 - 为该版本列出的每个 Mandatory
destination + consumerName创建精确 Delivery Row。 - 由 Deferred Constraint 或等价的事务末 anti-join 验证 Event 对应的 Mandatory Delivery 集合既无缺失也无多余后再提交业务事务。
Delivery Set Definition 本身是严格、封闭的内容寻址对象,Contracts 必须能从下列结构直接生成运行时 Schema、TypeScript 和 JSON Schema:
DeliverySetDefinition@N = {
deliverySetVersion
deliverySetDefinitionSchemaVersion
deliverySetDefinitionDigestAlgorithmVersion = jcs-sha256-v1
deliverySetDefinitionDigest
applicableEventCount
applicableEvents[] = sorted {
eventType
eventSchemaVersion
}
memberCount
members[] = sorted {
eventType
eventSchemaVersion
destination
consumerName
expectedReceiptType
expectedReceiptSchemaVersion
projectionTarget =
{ state=not_applicable }
| { state=required; projectionName; projectionVersion }
mandatory
}
}applicableEvents[] 按 (eventType, eventSchemaVersion) 排序,members[] 按 (eventType, eventSchemaVersion, destination, consumerName) 排序;两个数组都拒绝重复,Count 必须与实际项数一致。每个 Member 的 Event 键必须存在于 applicableEvents[],且同一 (eventType, eventSchemaVersion, destination) 只能映射一个 Consumer;mandatory=true 的精确子集就是 Producer 必须原子创建的 Delivery Row 集合。destination=core_operations必须且只能使用expectedReceiptType=OperationsProjectionAppliedReceipt和projectionTarget.state=required,冻结本Definition可接受的精确projectionName + projectionVersion;其他三个Destination必须分别绑定自己的注册Receipt Type/Schema且projectionTarget.state=not_applicable。Receipt返回的Projection坐标绝不能充当自己的Expected值。摘要 Candidate 精确覆盖 deliverySetDefinitionSchemaVersion + deliverySetDefinitionDigestAlgorithmVersion + deliverySetVersion + applicableEventCount + applicableEvents[] + memberCount + members[],包括每个Member的Receipt Contract与Projection Target,只排除 deliverySetDefinitionDigest 自身;Registry 的 Status、Activated Time 与 Source Release 是独立存储元数据,不属于 Definition Candidate。未知字段、未知 Schema/算法、数组未规范排序、Count 不等、孤立 Member、同 Destination重绑、Receipt/Projection分支错配、裸 Digest 或同 Version 异 Candidate/Digest 均拒绝。Contracts 必须发布原始 JSON、Canonical Bytes 与 Digest Golden;字段或排序语义变化升级 Definition Schema Version,摘要算法变化升级 Digest Algorithm Version。Registry 只能以完整四元组 insert-or-compare并激活不可变定义;Event和Delivery必须逐项冻结或以不可变外键解析Member的Expected Receipt/Projection Target,Ack验证与批准归档Source View必须解析同一Definition四元组和Member,不能用响应、自报Receipt、部署配置或当前Registry补算法/版本/Projection。
首个允许完整准入、Billing经济边界、Execution/Gateway→Metering输入和Billing收敛唤醒的Delivery Set Release,applicableEvents[] 至少必须精确包含下列十三个Pair;这里的Event Type与Schema Version都是精确字面量,不允许前缀、通配、别名或“latest”:
{ eventType=run.created@2.0; eventSchemaVersion=2.0 }
{ eventType=wallet.reserved@2.0; eventSchemaVersion=2.0 }
{ eventType=billing.finalization.committed@1; eventSchemaVersion=1 }
{ eventType=billing.settlement-transition.committed@1; eventSchemaVersion=1 }
{ eventType=billing.late-settlement-exposure.case-requested@1; eventSchemaVersion=1 }
{ eventType=billing.case-resolution.applied@1; eventSchemaVersion=1 }
{ eventType=billing.late-settlement-command-result.committed@1; eventSchemaVersion=1 }
{ eventType=execution.attempt-eligibility.finalized@1; eventSchemaVersion=1 }
{ eventType=execution.run-finalization.closed@1; eventSchemaVersion=1 }
{ eventType=billing.finalization-reconciliation.case-requested@1; eventSchemaVersion=1 }
{ eventType=billing.finalization-reconciliation-resolution.applied@1; eventSchemaVersion=1 }
{ eventType=gateway.evidence-availability.current-changed@1; eventSchemaVersion=1 }
{ eventType=metering.settlement-input.current-changed@1; eventSchemaVersion=1 }Run/Wallet v2及五个Billing经济边界Pair各自只有{destination=core_operations; consumerName=operations_projector; expectedReceipt=OperationsProjectionAppliedReceipt@1; projectionTarget=required(精确Active Name/Version); mandatory=true};billing.finalization-reconciliation.case-requested@1也只有该Operations Member;execution.run-finalization.closed@1只有{destination=billing_finalization_coordinator; consumerName=billing_finalization_coordinator_consumer; mandatory=true};billing.finalization-reconciliation-resolution.applied@1则同时拥有该Coordinator Member和core_operations → operations_projector两条Mandatory Member。Execution Eligibility与Gateway Availability两个Pair各自只有{destination=metering_input_processor; consumerName=metering_input_processor_consumer; mandatory=true};Settlement Input Pair只有{destination=billing_settlement_catch_up; consumerName=billing_settlement_catch_up_consumer; mandatory=true}。因此该最小集合是十三个Applicable Pair、十四个Mandatory Member;Count与Digest必须覆盖完整排序数组和每个Member的Expected Receipt/Projection Target。Projection切换必须在同一受控Cutover锁中先证明新版本Ready,再原子切换Active Pointer并激活引用该精确版本的新Delivery Set Definition;新Event只能冻结新Definition,旧Pending Event仍按自己的旧Definition Target收敛,Shadow或其他版本Receipt永远不能满足任一者。首期前四类Owner间触发Event不伪造Operations语义投影;两类Finalization Reconciliation Event的Operations投影则分别是预留Identity创建Case和基于Billing Fact终结Case,不得从它们推断Ledger或改写Billing Fence。任一Destination的Ack都不能满足另一Delivery。该Release还可包含Contracts已经支持的其他精确Event。任何Core Admission、Execution、Gateway、Billing或Metering Writer在启用某个新Event前,都必须在自己的事实库中证明当前激活Definition含精确Pair和全部对应Mandatory Member;Definition不适用、Event Type/Schema别名、Destination/Consumer/Expected Receipt/Projection Target错配或Member缺失时领域事务直接失败,不能提交后等待Registry升级。Consumer只按冻结Definition中的精确Pair解码,未知Event进入Quarantine而不是宽松映射成相似名称。
Billing边界事务必须先分配本事务内全部Event ID并冻结Causation,再对每个Event独立计算Canonical Envelope Digest、冻结同一时点激活的Delivery Set完整四元组并创建其全部Mandatory Row。Finalization事务至少写 billing.finalization.committed@1;每个成功Transition写 billing.settlement-transition.committed@1,request_open同事务再写Case Requested,resolve同事务再写Case Resolution Applied;Late Command成功同事务写Command Result Committed和Case Resolution Applied。每个Event拥有自己的Event Row、Digest和Delivery Rows,不允许一个“主事件”的Delivery代替同事务其他Event。Owner Fact/Result、Ledger/Release/Reservation或Dimension State、Exposure/Applied Fact、Audit、所有Event与其Delivery Set/Rows任一写入或Deferred集合校验失败,整个Billing事务回滚。
Execution追加唯一ExecutionEligibilityFact@1的事务必须同时写execution.attempt-eligibility.finalized@1、完整Event摘要、冻结Definition四元组与metering_input_processor Mandatory Row;嵌套Not-dispatched/Failure Fact、Event或Delivery任一失败全部回滚。Gateway每次为attempt-bound Usage/Cost链推进Current Availability时同样必须把Evidence(若有)、Availability Snapshot/Pointer、gateway.evidence-availability.current-changed@1和同Destination Delivery原子提交;Pre-binding Snapshot不适用且禁止产生该Event。Metering Input Consumer在自己的Inbox事务验证Event与Owner四元组,单调推进MeteringAttemptInputFence相应Watermark、insert-or-compare Work并写严格Receipt后才Ack。Eligibility与Gateway Usage可以任意先到并等待另一侧;Gateway Cost独立调度,新Current Revision重新置Pending。任何同步响应、Poll或best-effort Observation都不能满足Mandatory Delivery。
Execution关闭Attempt集合的CAS事务必须原子追加RunAttemptSetManifest@1 + RunExecutionFinalizationFact@1 + execution.run-finalization.closed@1,计算完整Event摘要、冻结Definition四元组并创建唯一billing_finalization_coordinator Mandatory Row。Event自包含Tenant/Account/Reservation/Run、Run Admission Manifest四元组、Closure Fact四元组/State Version、Attempt Set Manifest四元组与Count,Aggregate固定Run/Revision 1;任一Owner/Event/Delivery失败都回滚关闭CAS,同Run第二个Closure Event拒绝。Billing Coordinator Consumer在单一Inbox事务验证Event、推进Finalization Fence/Closure Watermark、冻结Trigger Snapshot、insert-or-compare Work和写版本化Receipt后才Ack;Execution在关闭提交后、同步调用Billing或排队前崩溃仍由该Delivery恢复。
Billing Finalization首次进入reconciliation_required前,必须先通过Operations受控Primitive预留Run级FinalizationReconciliationCaseIdentityReservation@1。Identity严格绑定规范化Tenant/Workspace/Project、Account/Reservation/Run/Fence、Case/Generation、确定性Case Request Operation和Predecessor;其摘要、Reserve/Read Audience/Scope不得绑定Triggering Work或billingFinalizationOperationId。Billing随后的单一Finalization事务同时提交Decision、绑定预留Identity的Fence/Work Finish/Audit、billing.finalization-reconciliation.case-requested@1完整Event摘要、冻结Definition四元组及唯一core_operations → operations_projector Mandatory Row;任一失败全部回滚。Event自包含Tenant、Account/Reservation/Run/Fence、Case/Generation/Expected Resolution Revision/Operation、Identity Reservation、Decision、Triggering Work四元组与triggeringBillingFinalizationWorkGeneration及Reason,Operations以Event完整四元组和全部Expected身份调用readBillingFinalizationDecisionForOperations,逐项验证最终winning Operation/Work。首代或直接前代已由Billing Fact终结时按预分配Ref创建Case;后继Request先到时则在同一Projector事务insert-or-compare严格PendingFinalizationCaseSuccessor@1与本Event的Operations Receipt并Ack,完整冻结后继Identity、Predecessor Identity、Case Requested Event Envelope、两Operation、Decision、Work/Generation与Reason,禁止创建第二Open Case或丢弃已Ack后继。前代Resolution/Request到达后在一个事务消费最大连续已验证Pending链,原子终结前代、物化后继并冻结消费所依据的Resolution Event/Fact。reserve(G1) → G2 supersedes → G1 CAS失败 → G2必须复用同Identity;不能因预留绑定旧Work产生孤儿。
后续受控Resolution命令事务必须原子消费Grant JTI、追加FinalizationReconciliationResolutionApplied@1、Audit及billing.finalization-reconciliation-resolution.applied@1,计算完整Event摘要、冻结Definition四元组并创建billing_finalization_coordinator + core_operations两条Mandatory Row。Event自包含Tenant/Account/Reservation/Run/Fence、Case Ref/Generation/Resolution Revision/Operation、Identity Reservation与Resolution Fact完整四元组及固定action=reevaluate_current_owner_state;Aggregate由Fence ID+代次确定性派生,Revision等于Resolution Revision。Fact绑定的非空investigationEvidenceReferences[]只是Contracts Registry内排序去重的授权调查引用;Billing只校验Request/Candidate/Decision/Grant/Fact间集合等值和Schema注册,不解引用它来证明Owner已修复。命令处理必须先按Operation调用readFinalizationReconciliationResolutionAppliedByOperation并完整比较Actor、Request/Candidate/Decision/Grant/JTI,found返回首次Fact后才跳过当前Grant检查;Operations消费Event则使用不同Audience的readFinalizationReconciliationResolutionAppliedForOperations。Fact、Grant消费、Audit、Event或任一Delivery失败全部回滚。Coordinator只用既有严格Owner Reads重新评估,Operations只根据Billing Fact终结对应Case;两者都不能携带/推断金额、Ledger、Reservation终态、替代Snapshot或忽略冲突标志。Owner状态仍冲突时重评创建新Case代次。
当前ADR-030只读阶段不启用上述命令。完整Release定义可以登记这两类Event,但在任何Producer可能进入reconciliation_required前,必须把ADR-031的严格Request → Preview/Approval/JIT → Candidate → Operations Decision → Signed Grant → Billing Fact/Event命令链与Producer同批启用;否则Producer保持禁用,不能产生无出口Fence。Grant的boundWorkload是实际呈递Domain调用、并等于Caller JWT sub的Admin Command Gateway Workload Principal,不是Billing接收方Worker。
Metering每次追加SettlementInputSnapshot@1并CAS推进Current Pointer时,必须在同一Metering事务写唯一metering.settlement-input.current-changed@1、完整Canonical Envelope Digest、冻结Definition四元组以及唯一billing_settlement_catch_up Mandatory Row。Event的Aggregate是Contracts确定性Settlement Dimension ID,Revision等于Snapshot State Version,Payload携带Current Snapshot与直接前驱完整五元组;Snapshot、Pointer、Event、Delivery任一失败全部回滚。Metering Dispatcher在事务外调用Billing Consumer;Billing只在自己的Inbox事务耐久推进Watermark:Reservation尚未成功Finalization时推进Admission阶段已存在的Finalization Fence、冻结Trigger Snapshot并产生/合并Finalization Work,成功Finalization后重新打开Dimension Catch-up Fence并insert-or-compare Catch-up Work,最后写版本化Receipt并返回Ack。极端乱序令Fence尚不可绑定时可以先耐久保存Inbox Watermark,但必须返回retryable且不写成功Receipt/Ack;不能因尚无Dimension State而丢弃,也不能永久停在awaiting_finalization。Billing提交后响应丢失只会导致同一Event重投,不能要求跨库两阶段提交。
Case Requested、Resolution Applied和Late Result都可能乱序投递,但Delivery Set不因此增加隐式先后依赖。每个Consumer以Event Receipt幂等,并靠Event自包含的Owner Fact/Result、Case Identity Reservation、Exposure、Aggregate Revision与Causation收敛;所有相关经济Event都冻结完整Billing Account/Reservation/Run/Step/Attempt/Dimension,Consumer不得从Attempt或相邻Event反查Step。pending_resolution只允许在Owner Fact与中间Exposure链全部可读且验证完成后,与Event Receipt同一Operations事务提交;暂不可读或链不完整必须返回retryable,不写Pending/Outcome/Receipt,让原Delivery保持non-terminal按Lease重试。Finalization后继Case Request若直接前代尚未Owner终结,则以严格pending_case_successors_vN行与Case Request Event Receipt同事务Ack;前代Resolution或缺失的前代Request到达后原子消费该行并物化后继,C2 Request → C1 Resolution → C1 Request与并发反序必须和Shadow Rebuild一致,不能只Ack Event却没有可唤醒载体。Transition的首个Causation指向Finalization,后续严格指向同Dimension直接前一个推进State Revision的Transition或Late Result Event,并验证前驱toRevision等于当前fromRevision。Transition Event还必须携带Owner Fact内稳定逻辑Operation与按(logical Operation, Reservation, Run Step, Attempt, Dimension, From State Revision, Basis完整摘要)确定性派生的Validation Operation,Consumer逐项验证且不得把旧Receipt对应的Validation Operation跨Revision复用。多Event事务崩溃、同事务只缺一个Event/Delivery Row、别名Event、缺失/跨Step、错Causation、Owner暂不可读且不得提交Receipt、F→T(open)→Late Result→T(correction)乱序/缺边、F→旧Receipt→并发前驱→新From State/Basis/Validation Operation重基线、重复Event ID、相同Owner Fact产生两个Event、Definition切换并发和响应丢失重放都必须有真实PostgreSQL负向/竞态测试。Settlement Input链另必须覆盖同Snapshot缺Event/Delivery、Event早于或晚于Finalization、Billing Receipt提交后Ack响应丢失、旧Revision晚到、同Revision异Digest、Snapshot连续推进与S1 Receipt → S2 Event先到且Ack → F(S1) → 无S3,并证明Watermark、Fence、Work最终让Billing追平S2。
Delivery Set 后续新增 Destination 只影响新 Append,不能把历史事件静默解释为当年也必须投递。首次迁移必须创建一个显式、不可变且带完整 Definition 四元组的 Bootstrap Delivery Set Version,并在受控迁移中为全部既有 Event 固定该四元组与精确 Delivery Row。对每个可按其原 Schema 严格解析的历史 Event,迁移必须使用 Contracts jcs-sha256-v1 Golden 实现对完整 Canonical Envelope(含 Payload)计算并逐条 insert-or-compare eventEnvelopeDigestAlgorithmVersion + eventEnvelopeSha256,记录分类、Count 与集合摘要,验证后才将两列设为 NOT NULL;禁止使用 PostgreSQL json/jsonb::text、语言运行时默认序列化或伪造零摘要。不支持 Schema、未知 Definition Schema/算法、Definition/Event 四元组错配、摘要冲突或损坏 Envelope/Payload 进入可见 Quarantine,不得生成虚假 Digest,也不能留下无 Delivery Set 的悬空 Event。Bootstrap 范围、数量、分类、Definition 四元组、摘要与 Golden 运行结果进入迁移证据。
不采用全局顺序游标或单值进度标记。Dispatcher 的工作集合始终由显式status=pending,或完成Lease过期恢复手续后可重新领取的status=leased,结合availableAt条件查询得到;acknowledged|quarantined均为终态。
现有 published_at 只能从 NULL 变为首次满足“冻结 Delivery Set 中全部 Mandatory Delivery 已由匹配 Receipt 确认”的时间戳。Dispatcher 不拥有该列的直接 UPDATE 权限,只能调用数据库受控函数。函数返回三态:当前 Ack/Receipt/Member/Fencing 不变量错误时 invalid 并抛错回滚;当前 Ack 合法但其他 Mandatory Delivery 仍未完成时返回 not_complete、保持 published_at=NULL,本 Delivery 的成功确认事务正常提交;当前 Ack 使全集完成时返回 complete 并以 CAS 将 NULL 写为首次时间戳。函数在同一事务中验证 destination → consumerName、Ack Receipt 坐标与 Event ID、Canonical Event Digest Algorithm Version/SHA-256;数据库 Trigger 再禁止回空、改时或二次覆盖。published_at 是兼容摘要,不是多 Destination 的权威状态,也不能解释为未来 Consumer、外部 Broker 或客户 Webhook 都已完成。
Event Envelope、Payload、Canonical Event Digest Algorithm Version/SHA-256、deliverySetVersion、Delivery Set Definition Schema/Digest Algorithm Version/Digest 和创建时间通过 Trigger 与 Column Grant 双重保护为不可变;Dispatcher Role 只能操作 Delivery/Attempt/Quarantine,并执行上述受控函数。
Claim 与 Lease
Dispatcher 循环使用短数据库事务 Claim 一组当前可用 Delivery:
BEGIN
SELECT eligible deliveries
ORDER BY available_at, created_at, event_id
FOR UPDATE SKIP LOCKED
LIMIT configured_batch_size
UPDATE selected rows
SET lease_owner = current_worker,
lease_token = new_fencing_token,
lease_expires_at = configured_deadline,
status = 'leased',
current_attempt_id = new_attempt_id,
attempt_count = attempt_count + 1
INSERT immutable delivery_attempt_started(new_attempt_id, destination, event_id,
new_fencing_token, started_at)
COMMIT只有 Start 与 Lease/Current Attempt 在同一 Claim 事务提交后,Worker 才能在事务之外执行 Schema 验证、完整性校验和投递。这样慢 Consumer 不会长期持有数据库锁,同时“可能已经调用 Consumer”的每个窗口都有稳定 Attempt ID。
所有 Ack、Retry、Lease Renew、Quarantine 和正常 Finish 都必须带 destination + eventId + currentAttemptId + leaseToken 条件。正常确认事务只能为当前 Start 追加唯一 Finish;Lease 过期后,恢复者通过受控 Primitive 验证旧 Token 已不再有效、尚无 Finish,再为旧 Attempt 追加 outcome=lease_expired_unknown,随后新 Claim 使用新 Attempt ID/Token。旧 Worker 即使稍后返回,也不能追加第二个 Finish或覆盖新 Owner 状态。这是 Fencing,不依赖进程时间判断“谁更新得晚”。
批次、Lease 期限、续租、并发、数据库语句超时和 Worker 停机排空都来自部署配置与容量验证。若单次投递可能超过 Lease,优先让目标操作具备幂等 Ack 与可查询状态;不能靠不断增加固定 Lease 掩盖阻塞。
投递与 at-least-once
首期 Canonical Destination 是 core_operations | metering_input_processor | billing_finalization_coordinator | billing_settlement_catch_up 四个精确字面量,分别只能映射 operations_projector | metering_input_processor_consumer | billing_finalization_coordinator_consumer | billing_settlement_catch_up_consumer。投递结果只有三类:
- acknowledged:Consumer 已持久接受并返回稳定 Ack,Dispatcher 用当前 Lease Token 标记完成。
- retryable:临时依赖故障、超时或容量保护;记录 Attempt,把Delivery迁回
pending并按配置化退避重新设置availableAt。 - non_retryable:未知 Schema、Payload 摘要冲突、Envelope/Payload 不变量破坏或明确不可处理版本;进入 Quarantine。
Active Consumer Ack使用封闭联合,不能以自由字符串或HTTP 2xx代替Receipt:
ActiveConsumerAck@1 = {
destination / consumerName / eventId
eventEnvelopeDigestAlgorithmVersion / eventEnvelopeSha256
appliedReceipt =
{ kind=operations_projection;
operationsProjectionAppliedReceiptRef;
operationsProjectionAppliedReceiptSchemaVersion;
operationsProjectionAppliedReceiptDigestAlgorithmVersion;
operationsProjectionAppliedReceiptDigest }
| { kind=metering_input_processor;
meteringInputProcessorConsumerReceiptRef;
meteringInputProcessorConsumerReceiptSchemaVersion;
meteringInputProcessorConsumerReceiptDigestAlgorithmVersion;
meteringInputProcessorConsumerReceiptDigest }
| { kind=billing_finalization_coordinator;
billingFinalizationCoordinatorConsumerReceiptRef;
billingFinalizationCoordinatorConsumerReceiptSchemaVersion;
billingFinalizationCoordinatorConsumerReceiptDigestAlgorithmVersion;
billingFinalizationCoordinatorConsumerReceiptDigest }
| { kind=billing_settlement_catch_up;
billingSettlementCatchUpConsumerReceiptRef;
billingSettlementCatchUpConsumerReceiptSchemaVersion;
billingSettlementCatchUpConsumerReceiptDigestAlgorithmVersion;
billingSettlementCatchUpConsumerReceiptDigest }
}Operations Owner的Applied Receipt同样是可跨库验证的内容寻址对象,不是只能靠源库外键猜测的坐标:
OperationsProjectionAppliedReceipt@1 = {
operationsProjectionAppliedReceiptRef
operationsProjectionAppliedReceiptSchemaVersion
operationsProjectionAppliedReceiptDigestAlgorithmVersion = jcs-sha256-v1
operationsProjectionAppliedReceiptDigest
destination=core_operations
consumerName=operations_projector
projectionName / projectionVersion
eventId / eventType / eventSchemaVersion
eventEnvelopeDigestAlgorithmVersion / eventEnvelopeSha256
appliedAt
}Receipt Candidate覆盖除Repository-owned Ref、appliedAt与摘要自身外的全部字段,(consumerName,projectionName,projectionVersion,eventId)唯一且insert-or-compare。readOperationsProjectionAppliedReceipt(Receipt完整四元组, expected destination/consumer/projection/event ID/Type/Schema/Envelope摘要)中的Expected Projection必须来自Event冻结Definition的匹配Member,而非请求中Receipt正文;Audience/Scope同时绑定Definition四元组、Member坐标和原Producer Dispatcher Workload,严格返回found{receipt=OperationsProjectionAppliedReceipt@1} | not_found | conflicting,Operations与调用方都重算摘要并逐项比较Definition Target,不提供列表、latest或Admin读取。Core本库Producer可另用外键做完整性加固,Billing等跨库Producer必须以冻结Member加受信响应或该定向Read校验同一Receipt四元组;不得伪造本库外键,也不得让Shadow/错误Projection的自洽Receipt完成Delivery。
Metering Owner的版本化Receipt是严格、低敏、内容寻址对象:
MeteringInputProcessorConsumerReceipt@1 = {
meteringInputProcessorConsumerReceiptRef
meteringInputProcessorConsumerReceiptSchemaVersion
meteringInputProcessorConsumerReceiptDigestAlgorithmVersion = jcs-sha256-v1
meteringInputProcessorConsumerReceiptDigest
destination=metering_input_processor
consumerName=metering_input_processor_consumer
eventId
eventEnvelopeDigestAlgorithmVersion / eventEnvelopeSha256
source =
{ kind=execution_eligibility;
eventType=execution.attempt-eligibility.finalized@1; eventSchemaVersion=1;
tenantKind / tenantId / workspaceId / projectId?;
billingReservationId / billingAccountId;
runId / runStepId / executionAttemptId;
attemptExecutionManifestRef / attemptExecutionManifestSchemaVersion;
attemptExecutionManifestDigestAlgorithmVersion / attemptExecutionManifestDigest;
executionEligibilityFactRef / executionEligibilityFactSchemaVersion;
executionEligibilityFactDigestAlgorithmVersion / executionEligibilityFactDigest }
| { kind=gateway_availability;
eventType=gateway.evidence-availability.current-changed@1; eventSchemaVersion=1;
executionAttemptId / gatewayDeploymentId;
attemptExecutionManifestRef / attemptExecutionManifestSchemaVersion;
attemptExecutionManifestDigestAlgorithmVersion / attemptExecutionManifestDigest;
providerEvidenceKind=provider_usage|provider_cost; evidenceDimensionKey;
currentAvailabilityRef / currentAvailabilitySchemaVersion;
currentAvailabilityDigestAlgorithmVersion / currentAvailabilityDigest;
currentAvailabilityStateVersion }
outcome = advanced | already_observed
lane =
{ kind=eligibility_and_settlement_input;
eligibilityFenceRevision / latestEligibilityWorkGeneration;
eligibilityTriggerDigestAlgorithmVersion / eligibilityTriggerDigest;
coordinationState =
{ state=waiting_for_execution|waiting_for_gateway_usage;
activeWork={ state=none };
lastFinishedWork=
{ state=none }
| { state=available;
workGeneration / meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest } }
| { state=pending;
activeWork={ state=pending; workGeneration; meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest } }
| { state=applied; activeWork={ state=none }; appliedWorkGeneration; meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest;
settlementInputSnapshotCount;
settlementInputSnapshots[] = sorted {
runStepId / executionAttemptId / chargeDimensionKey;
ref / schemaVersion / digestAlgorithmVersion / digest / stateVersion } }
| { state=blocked_on_owner_conflict;
activeWork={ state=pending; workGeneration; meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest };
meteringInputWorkAttemptBlockedRef / meteringInputWorkAttemptBlockedSchemaVersion;
meteringInputWorkAttemptBlockedDigestAlgorithmVersion / meteringInputWorkAttemptBlockedDigest } }
| { kind=provider_cost;
gatewayDeploymentId / sourceEvidenceDimensionKey;
providerCostFenceRevision / latestProviderCostWorkGeneration;
providerCostTriggerDigestAlgorithmVersion / providerCostTriggerDigest;
coordinationState =
{ state=pending;
activeWork={ state=pending; workGeneration; meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest } }
| { state=availability_processed; processingOutcome=available_applied; activeWork={ state=none };
workGeneration / meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
providerCostAvailabilityProcessingRef / providerCostAvailabilityProcessingSchemaVersion;
providerCostAvailabilityProcessingDigestAlgorithmVersion / providerCostAvailabilityProcessingDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest }
| { state=waiting_for_evidence; processingOutcome=waiting_for_evidence; activeWork={ state=none };
workGeneration / meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
providerCostAvailabilityProcessingRef / providerCostAvailabilityProcessingSchemaVersion;
providerCostAvailabilityProcessingDigestAlgorithmVersion / providerCostAvailabilityProcessingDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest }
| { state=no_cost_evidence; processingOutcome=no_cost_evidence; activeWork={ state=none };
workGeneration / meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
providerCostAvailabilityProcessingRef / providerCostAvailabilityProcessingSchemaVersion;
providerCostAvailabilityProcessingDigestAlgorithmVersion / providerCostAvailabilityProcessingDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest }
| { state=blocked_on_processing_conflict;
activeWork={ state=pending; workGeneration; meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest };
meteringInputWorkAttemptBlockedRef / meteringInputWorkAttemptBlockedSchemaVersion;
meteringInputWorkAttemptBlockedDigestAlgorithmVersion / meteringInputWorkAttemptBlockedDigest }
| { state=reconciliation_required; processingOutcome=reconciliation_required;
availabilityState=conflicting; stateReasonCode;
activeWork={ state=none };
workGeneration / meteringInputWorkOperationId;
meteringInputWorkRef / meteringInputWorkSchemaVersion;
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest;
providerCostAvailabilityProcessingRef / providerCostAvailabilityProcessingSchemaVersion;
providerCostAvailabilityProcessingDigestAlgorithmVersion / providerCostAvailabilityProcessingDigest;
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion;
meteringInputWorkAttemptFinishedDigestAlgorithmVersion / meteringInputWorkAttemptFinishedDigest } }
committedAt
}Work Finish必须能显式收敛被新Generation取代的Started Attempt:
MeteringInputWorkAttemptFinished@1 = {
meteringInputWorkAttemptFinishedRef / meteringInputWorkAttemptFinishedSchemaVersion
meteringInputWorkAttemptFinishedDigestAlgorithmVersion=jcs-sha256-v1
meteringInputWorkAttemptFinishedDigest
meteringInputWorkRef / meteringInputWorkSchemaVersion
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest
meteringInputWorkOperationId
executionAttemptId / workGeneration
lane = MeteringInputWork@1.lane
laneTriggerDigestAlgorithmVersion / laneTriggerDigest
attemptNumber / fencingToken
outcome =
{ kind=eligibility_and_settlement_input;
state=waiting_for_counterpart; waitingFor=execution_eligibility|gateway_usage }
| { kind=eligibility_and_settlement_input; state=applied;
settlementInputSnapshotCount;
settlementInputSnapshots[] = sorted {
runStepId / executionAttemptId / chargeDimensionKey;
ref / schemaVersion / digestAlgorithmVersion / digest / stateVersion } }
| { kind=provider_cost; state=availability_processed;
providerCostAvailabilityProcessingRef / providerCostAvailabilityProcessingSchemaVersion;
providerCostAvailabilityProcessingDigestAlgorithmVersion / providerCostAvailabilityProcessingDigest }
| { kind=eligibility_and_settlement_input|provider_cost; state=superseded;
supersededByWorkGeneration }
finishedAt
}新Current Event在旧Work已经Started时,创建新Generation的Inbox事务必须基于旧Work、Started Attempt与旧Fencing Token insert-or-compare唯一state=superseded Finish;若该事务未代写,旧Worker晚到只能在确认Lane Current已指向supersededByWorkGeneration后追加同一正文。两条路径竞争必须返回同一Fact,旧Token不能写任何业务Fact或推进新Lane。
Worker的非终态阻塞记录使用同一个封闭DTO;它是Attempt审计,不是Work Finish或Operations Case:
MeteringInputWorkAttemptBlocked@1 = {
meteringInputWorkAttemptBlockedRef / meteringInputWorkAttemptBlockedSchemaVersion
meteringInputWorkAttemptBlockedDigestAlgorithmVersion=jcs-sha256-v1
meteringInputWorkAttemptBlockedDigest
meteringInputWorkRef / meteringInputWorkSchemaVersion
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest
meteringInputWorkOperationId
executionAttemptId / workGeneration
lane = MeteringInputWork@1.lane
laneTriggerDigestAlgorithmVersion / laneTriggerDigest
attemptNumber / fencingToken
blockingCause =
{ kind=execution_eligibility_owner_conflict;
executionEligibilityConflictFactRef / executionEligibilityConflictFactSchemaVersion;
executionEligibilityConflictFactDigestAlgorithmVersion / executionEligibilityConflictFactDigest }
| { kind=gateway_usage_owner_conflict;
gatewayEvidenceConflictFactRef / gatewayEvidenceConflictFactSchemaVersion;
gatewayEvidenceConflictFactDigestAlgorithmVersion / gatewayEvidenceConflictFactDigest }
| { kind=gateway_binding_owner_conflict;
gatewayBindingConflictFactRef / gatewayBindingConflictFactSchemaVersion;
gatewayBindingConflictFactDigestAlgorithmVersion / gatewayBindingConflictFactDigest }
| { kind=eligibility_processing_conflict;
eligibilityProcessingConflictEvidenceRef / eligibilityProcessingConflictEvidenceSchemaVersion;
eligibilityProcessingConflictEvidenceDigestAlgorithmVersion / eligibilityProcessingConflictEvidenceDigest }
| { kind=provider_cost_processing_conflict;
providerCostProcessingConflictEvidenceRef / providerCostProcessingConflictEvidenceSchemaVersion;
providerCostProcessingConflictEvidenceDigestAlgorithmVersion / providerCostProcessingConflictEvidenceDigest }
blockedAt
}
EligibilityProcessingConflictEvidence@1 = {
eligibilityProcessingConflictEvidenceRef / eligibilityProcessingConflictEvidenceSchemaVersion
eligibilityProcessingConflictEvidenceDigestAlgorithmVersion=jcs-sha256-v1
eligibilityProcessingConflictEvidenceDigest
runId / runStepId / executionAttemptId
attemptExecutionManifestRef / attemptExecutionManifestSchemaVersion
attemptExecutionManifestDigestAlgorithmVersion / attemptExecutionManifestDigest
executionEligibilityFactRef / executionEligibilityFactSchemaVersion
executionEligibilityFactDigestAlgorithmVersion / executionEligibilityFactDigest
meteringInputWorkRef / meteringInputWorkSchemaVersion
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest
meteringInputWorkOperationId / workGeneration
laneTriggerDigestAlgorithmVersion / laneTriggerDigest
validationStage = eligibility_basis | availability_validation_receipt |
eligibility_decision | meter_event | settlement_input
conflictingCandidateDigestAlgorithmVersion / conflictingCandidateDigest
reasonCode
observedAt
}
ProviderCostProcessingConflictEvidence@1 = {
providerCostProcessingConflictEvidenceRef / providerCostProcessingConflictEvidenceSchemaVersion
providerCostProcessingConflictEvidenceDigestAlgorithmVersion=jcs-sha256-v1
providerCostProcessingConflictEvidenceDigest
meteringInputWorkRef / meteringInputWorkSchemaVersion
meteringInputWorkDigestAlgorithmVersion / meteringInputWorkDigest
meteringInputWorkOperationId / workGeneration
runId / runStepId / executionAttemptId / gatewayDeploymentId / sourceEvidenceDimensionKey
attemptExecutionManifestRef / attemptExecutionManifestSchemaVersion
attemptExecutionManifestDigestAlgorithmVersion / attemptExecutionManifestDigest
attemptRouteBindingRef / attemptRouteBindingSchemaVersion
attemptRouteBindingDigestAlgorithmVersion / attemptRouteBindingDigest
gatewayRouteSnapshotRef / gatewayRouteSnapshotSchemaVersion
gatewayRouteSnapshotDigestAlgorithmVersion / gatewayRouteSnapshotDigest
costAvailabilityRef / costAvailabilitySchemaVersion
costAvailabilityDigestAlgorithmVersion / costAvailabilityDigest / costAvailabilityStateVersion
cause =
{ kind=gateway_owner_conflict;
gatewayEvidenceConflictFactRef / gatewayEvidenceConflictFactSchemaVersion;
gatewayEvidenceConflictFactDigestAlgorithmVersion / gatewayEvidenceConflictFactDigest }
| { kind=gateway_binding_owner_conflict;
gatewayBindingConflictFactRef / gatewayBindingConflictFactSchemaVersion;
gatewayBindingConflictFactDigestAlgorithmVersion / gatewayBindingConflictFactDigest }
| { kind=metering_deterministic_validation;
validationStage=source_schema|source_identity|provider_cost_mapping|
provider_cost_normalization|availability_validation_receipt|
route_validation_receipt|provider_cost_repository;
conflictingCandidateDigestAlgorithmVersion / conflictingCandidateDigest;
reasonCode }
detectedAt
}三个对象的摘要分别覆盖除Repository-owned Ref、各自时间和摘要自身外的全部字段。Blocked提交与Fence的blocked_on_owner_conflict|blocked_on_processing_conflict、保持activeWork.state=pending及Lease释放同一Metering事务完成;Work/Operation/Generation保持不变,之后同Work取得新Lease重试。blockingCause与Lane严格细化:Eligibility Lane只允许Execution/Gateway Evidence/Gateway Binding直接Owner分支或Eligibility Processing;Provider Cost Lane只允许外层provider_cost_processing_conflict,其Evidence再区分Gateway Evidence、Gateway Binding与Metering本地确定性校验。Provider Cost本地校验分支必须冻结冲突Candidate带算法摘要;Deployment/Source Dimension不可裁剪。禁止夹带Finished、Processing Fact、Case Ref或把Blocked当作成功Ack结果。
Billing Catch-up使用独立但同语义的严格阻塞事实:
BillingSettlementCatchUpWorkAttemptBlocked@1 = {
billingSettlementCatchUpWorkAttemptBlockedRef / billingSettlementCatchUpWorkAttemptBlockedSchemaVersion
billingSettlementCatchUpWorkAttemptBlockedDigestAlgorithmVersion=jcs-sha256-v1
billingSettlementCatchUpWorkAttemptBlockedDigest
billingSettlementCatchUpWorkRef / billingSettlementCatchUpWorkSchemaVersion
billingSettlementCatchUpWorkDigestAlgorithmVersion / billingSettlementCatchUpWorkDigest
billingCatchUpValidationOperationId
billingReservationId / billingAccountId / runId / runStepId
executionAttemptId / chargeDimensionKey / workGeneration
attemptNumber / fencingToken
blockingCause =
{ kind=settlement_input_owner_conflict;
settlementInputConflictFactRef / settlementInputConflictFactSchemaVersion;
settlementInputConflictFactDigestAlgorithmVersion / settlementInputConflictFactDigest }
blockedAt
}该Fact摘要同样排除Repository-owned Ref、blockedAt和摘要自身;它与Catch-up Fence的coordinationState=blocked_on_owner_conflict + activeWork.state=pending、同一Work/Operation/Generation及Lease释放原子提交,绝不等同于BillingSettlementCatchUpWorkAttemptFinished@1。
Receipt Candidate覆盖除Repository-owned Ref、committedAt和摘要自身外的完整严格联合;(consumerName,eventId)终身唯一。Execution Source只允许lane.kind=eligibility_and_settlement_input,其Run/Step/Attempt、Manifest与Eligibility Fact必须和Event逐项相等,并使用Event-bound readAttemptExecutionManifest(purpose=metering_execution_eligibility);Gateway Usage Source同样只进入Eligibility Lane,Gateway Cost Source无论Availability为available|pending|not_reported|unavailable|conflicting都只进入与gatewayDeploymentId + sourceEvidenceDimensionKey相等的Provider Cost Lane。两个Gateway分支读取Execution Owner的Attempt Manifest时使用purpose=metering_gateway_availability;读取Gateway Owner的Availability/Evidence/Route正文则必须使用绑定已接受Current-changed Event完整身份、Envelope摘要、Aggregate/Revision与Payload身份的metering_source_event_bootstrap,Bootstrap禁止预先要求尚未形成的Eligibility Basis、Provider Cost Candidate或Validation Operation。不存在metering_fact_construction通用Purpose。
Receipt必须冻结同一Inbox事务提交后的本Lane Fence Revision + latest Work Generation + Trigger摘要 + coordinationState,不能假定每个Event都会产生Pending Work。coordinationState=pending必须且只能绑定本Lane最新Generation的Active Work完整四元组与确定性Operation;Eligibility的waiting_for_execution|waiting_for_gateway_usage只允许activeWork=none并保留上一完成Work的Generation、Operation、Work与Finish完整四元组或none,applied必须绑定完成Generation、Operation、Work/Finish四元组及其严格排序/Count相等的Settlement Snapshot集。任一Owner Read或本地确定性校验冲突使用Eligibility的blocked_on_owner_conflict或Provider Cost的blocked_on_processing_conflict,绑定相同activeWork.state=pending与MeteringInputWorkAttemptBlocked@1完整四元组,绝不绑定Case或Finished。Provider Cost的三个普通完成分支与仅Gateway Availability.state=conflicting允许的reconciliation_required必须同时绑定完成Generation、Operation、Work、Processing Fact与Finish完整四元组;本地确定性冲突只允许Blocked分支且禁止Processing Fact/Finish/Lane终态CAS。Receipt自身因此可验证Receipt → Work → Blocked或Receipt → Work → Finish/Processing,不依赖不存在的跨库隐式外键。新Current Event必须以新Generation重新进入Pending,旧Work/Blocked/Finish/Processing Fact不得完成新Head。Eligibility与Provider Cost两种Revision/Trigger/Generation禁止混用,Ref/Schema相同但Work/Blocked/Finish/Processing摘要不同同样冲突。Gateway-bound Eligibility可以等待Gateway Usage,Gateway Usage可以等待Execution Eligibility,Not-dispatched不得等待Gateway,Provider Cost没有Counterpart等待分支。Event同摘要重投返回首次Receipt;outcome=already_observed也必须冻结首次处理时的严格Lane State,不能用后来的Current状态改写。同Revision异摘要、跨Manifest/Attempt/Kind/Dimension或Source/Lane/State联合混合均拒绝。Receipt只证明输入与当时Lane协调结果已耐久接受,不能代替Eligibility Decision、MeterEvent、ProviderCostFact、ProviderCostAvailabilityProcessingFact或Settlement Input。
Billing Owner的两类版本化Receipt都是严格、低敏、内容寻址对象:
BillingFinalizationCoordinatorConsumerReceipt@1 = {
billingFinalizationCoordinatorConsumerReceiptRef
billingFinalizationCoordinatorConsumerReceiptSchemaVersion
billingFinalizationCoordinatorConsumerReceiptDigestAlgorithmVersion = jcs-sha256-v1
billingFinalizationCoordinatorConsumerReceiptDigest
destination=billing_finalization_coordinator
consumerName=billing_finalization_coordinator_consumer
eventId
eventEnvelopeDigestAlgorithmVersion / eventEnvelopeSha256
source =
{ kind=execution_closure;
eventType=execution.run-finalization.closed@1; eventSchemaVersion=1;
tenantKind / tenantId / workspaceId / projectId?;
billingReservationId / billingAccountId / runId / closureOperationId;
runAdmissionManifestRef / runAdmissionManifestSchemaVersion;
runAdmissionManifestDigestAlgorithmVersion / runAdmissionManifestDigest;
runExecutionFinalizationRef / runExecutionFinalizationSchemaVersion;
runExecutionFinalizationDigestAlgorithmVersion / runExecutionFinalizationDigest;
runExecutionFinalizationStateVersion;
attemptSetManifestRef / attemptSetManifestSchemaVersion;
attemptSetManifestDigestAlgorithmVersion / attemptSetManifestDigest;
attemptCount }
| { kind=finalization_reconciliation_resolution;
eventType=billing.finalization-reconciliation-resolution.applied@1; eventSchemaVersion=1;
tenantKind / tenantId / workspaceId / projectId?;
billingReservationId / billingAccountId / runId / billingFinalizationFenceId;
reconciliationCaseRef / reconciliationGeneration / resolutionRevision=1 / resolutionOperationId;
finalizationCaseIdentityReservationRef / finalizationCaseIdentityReservationSchemaVersion;
finalizationCaseIdentityReservationDigestAlgorithmVersion / finalizationCaseIdentityReservationDigest;
finalizationReconciliationResolutionAppliedRef / finalizationReconciliationResolutionAppliedSchemaVersion;
finalizationReconciliationResolutionAppliedDigestAlgorithmVersion / finalizationReconciliationResolutionAppliedDigest;
actorWorkforcePrincipalId;
action=reevaluate_current_owner_state }
outcome =
{ state=work_scheduled;
finalizationFenceRevision / workGeneration;
billingFinalizationTriggerSnapshotRef / billingFinalizationTriggerSnapshotSchemaVersion;
billingFinalizationTriggerSnapshotDigestAlgorithmVersion / billingFinalizationTriggerSnapshotDigest;
billingFinalizationWorkRef / billingFinalizationWorkSchemaVersion;
billingFinalizationWorkDigestAlgorithmVersion / billingFinalizationWorkDigest }
| { state=already_superseded; finalizationFenceRevision }
committedAt
}
BillingSettlementCatchUpConsumerReceipt@1 = {
billingSettlementCatchUpConsumerReceiptRef
billingSettlementCatchUpConsumerReceiptSchemaVersion
billingSettlementCatchUpConsumerReceiptDigestAlgorithmVersion = jcs-sha256-v1
billingSettlementCatchUpConsumerReceiptDigest
destination=billing_settlement_catch_up
consumerName=billing_settlement_catch_up_consumer
eventId / eventType=metering.settlement-input.current-changed@1 / eventSchemaVersion=1
eventEnvelopeDigestAlgorithmVersion / eventEnvelopeSha256
tenantKind / tenantId / workspaceId / projectId?
billingReservationId / billingAccountId / runId / runStepId
executionAttemptId / chargeDimensionKey
settlementEconomicIdentity = SettlementEconomicIdentity@1
eventSettlementInput = {
settlementInputSnapshotRef / settlementInputSnapshotSchemaVersion;
settlementInputSnapshotDigestAlgorithmVersion / settlementInputSnapshotDigest;
settlementInputStateVersion
}
outcome =
{ kind=watermark_advanced; inboxRevision;
durability=
{ state=finalization_work_scheduled;
finalizationFenceRevision; workGeneration;
billingFinalizationTriggerSnapshotRef / billingFinalizationTriggerSnapshotSchemaVersion;
billingFinalizationTriggerSnapshotDigestAlgorithmVersion / billingFinalizationTriggerSnapshotDigest;
billingFinalizationWorkRef / billingFinalizationWorkSchemaVersion;
billingFinalizationWorkDigestAlgorithmVersion / billingFinalizationWorkDigest }
| { state=settlement_catch_up_scheduled;
settlementCatchUpFenceRevision; workGeneration;
billingSettlementCatchUpWorkRef / billingSettlementCatchUpWorkSchemaVersion;
billingSettlementCatchUpWorkDigestAlgorithmVersion / billingSettlementCatchUpWorkDigest;
billingCatchUpValidationOperationId }
| { state=held_for_resolution; finalizationFenceRevision;
reconciliationCaseRef / reconciliationGeneration / expectedResolutionRevision=1 } }
| { kind=already_observed; inboxRevision;
highestObservedSettlementInput={
settlementInputSnapshotRef / settlementInputSnapshotSchemaVersion;
settlementInputSnapshotDigestAlgorithmVersion / settlementInputSnapshotDigest;
settlementInputStateVersion } }
committedAt
}两类Receipt的摘要Candidate都覆盖除Repository-owned Ref、committedAt和摘要自身外的全部字段;各自(consumerName,eventId)终身唯一。同Event同Envelope摘要的重投返回首次完整Receipt;同Event异摘要、跨身份重绑或未知Schema/算法不得写Receipt。Coordinator Receipt的source=execution_closure必须且只能搭配outcome=work_scheduled,Closure/Manifest/Count、Fence Revision、Trigger Snapshot四元组、Billing Finalization Work完整四元组及Generation必须与首次Billing Inbox事务逐项相等。Finalization进入Reconciliation前唯一Closure Event与该Receipt必须已应用;之后相同Event重投只返回首次Receipt,不会产生新Outcome、Watermark或Work。source=finalization_reconciliation_resolution必须逐项等于Resolution Event中的Identity Reservation/Fact、Case/Generation/Revision/Operation与Fence身份,当前期望Revision搭配work_scheduled并要求Trigger的Resolution Watermark及Work等值;只有同一Resolution已经成功创建其重评Work,或该Work的重评已推进/被该Resolution吸收,重投才允许already_superseded并冻结当时Fence Revision。普通Closure/Settlement Watermark推进不得命中该分支。Source/Outcome分支混合、Closure返回already_superseded、同Ref/Schema但Identity/Fact/Work摘要不同、同Revision异Event摘要都冲突。
Settlement Receipt的finalization_work_scheduled只允许Reservation尚未成功终局且Fence不在Reconciliation,并要求Finalization Fence/Trigger/Billing Finalization Work四元组完整;settlement_catch_up_scheduled只允许成功终局,且其Fence Revision、Work完整四元组、Generation与billingCatchUpValidationOperationId必须和同一Inbox事务中的BillingSettlementCatchUpFence@1.coordinationState.activeWork及不可变BillingSettlementCatchUpWork@1逐项相等。Fence正在Finalization reconciliation_required时,新Current Settlement Input必须先单调推进Watermark,再返回绑定Finalization Fence Revision、Case/Generation/Expected Resolution Revision的held_for_resolution;该分支严格禁止Trigger/Work/Generation/Catch-up Operation,仅当后续同代Current Resolution Applied重开Finalization时由新Trigger一次性吸收已更新Watermark。分支混合、insert-or-compare冲突或同Revision异Snapshot摘要都不得Ack。already_observed只允许Event Revision低于当前合法Watermark,且highestObservedSettlementInput.stateVersion >= eventSettlementInput.stateVersion;若该Event已有Receipt,必须返回原Receipt而不是基于新Watermark改写结果。两类Receipt都只证明触发已耐久接受,不能被Ledger、Settlement Consumption、Finalization或Transition当作Current Validation Receipt。
跨库时没有伪造的Receipt外键。Operations/Metering/Billing Consumer通过受信内部Transport返回完整Receipt及其Ack;Gateway/Execution/Metering/Billing Dispatcher严格解码、重算所选Receipt摘要并比较Event、Destination/Consumer和完整业务身份。若需要对已返回四元组复核,只能使用精确Audience/Scope的readOperationsProjectionAppliedReceipt、readMeteringInputProcessorConsumerReceipt(Receipt完整四元组, expected Source Event/Envelope/Manifest/Fact或Availability身份)、readBillingFinalizationCoordinatorConsumerReceipt(Receipt完整四元组, expected Source严格联合:Closure Event/Run/Manifest或Resolution Event/Fence/Case/Generation/Identity/Fact身份)或readBillingSettlementCatchUpConsumerReceipt(Receipt完整四元组, expected Event ID/Type/Schema/Envelope摘要, expected Tenant/Account/Reservation/Run/Step/Attempt/Dimension/Snapshot五元组),分别返回found{receipt=对应严格DTO} | not_found | conflicting,不提供列表、latest或Admin读取;响应丢失时直接重投原Event也必须返回同一Receipt。
Dispatcher只能在当前 destination + eventId + currentAttemptId + leaseToken 仍有效、Consumer及Expected Receipt/Projection与冻结Member匹配、Ack算法版本/摘要与Event逐项相等,并且所选Receipt联合可验证时接受Ack。四个分支都要求Owner Receipt完整四元组、正文摘要与受信响应/定向Read匹配;Operations分支额外强制Receipt projectionName/projectionVersion等于冻结Member的projectionTarget,不能把Receipt自报值回填成Expected。只有Producer与Operations Receipt同库时才可额外要求本库外键,跨库分支不得把外键当作必要前置。其他Consumer/Projection、Shadow Rebuild Receipt或另一联合分支不能满足当前Delivery。确认事务必须原子完成:为当前Start追加唯一成功Finish、把Delivery置为terminal acknowledged并冻结Ack Receipt类型/完整四元组、调用published_at受控函数。Retry/Quarantine事务同样先追加Finish,再按同一Fencing条件迁移Delivery State。not_complete是合法结果并提交当前Delivery;只有Ack/Fencing/Receipt/Member不变量错误才使确认事务回滚。最后一个Mandatory Delivery与published_at CAS同事务,进程崩溃后依靠既有Receipt幂等重投,不会留下“全集已完成但 published_at 永久为空”的裂缝。
以下崩溃窗口会造成重投,是设计内行为:
Consumer 已提交 Receipt + Projection
→ Ack 返回前 Dispatcher 崩溃
→ 恢复者为旧 Attempt 追加 lease_expired_unknown Finish
→ Lease 过期后同一 eventId 再投
→ Consumer 发现 Receipt,返回相同 Ack,不重复应用数据库测试必须覆盖 Consumer 已被调用后进程崩溃、调用前崩溃、恢复者与旧 Worker 竞态、旧 Lease 晚到、重复 Finish 和重复 Ack;每个 Start 最终至多一个 Finish,任何 stale token 都不能改变 Delivery,新 Attempt 仍通过 Consumer Receipt 幂等收敛。
因此不能宣称网络层 exactly-once。经济事实、外部副作用或资产登记的消费者除 Event Receipt 外,还必须使用自己的业务幂等键;“已收到事件”不能代替领域不变量。
Consumer 原子提交
Metering Input Consumer必须严格区分AttemptEligibilityFinalized@1 | GatewayEvidenceAvailabilityCurrentChanged@1,验证Canonical Event摘要、Owner/Manifest四元组与各分支Refinement,再在一个Metering Inbox事务中锁定相同Attempt Input Fence、单调推进对应Watermark、产生/合并不可变Durable Work并写MeteringInputProcessorConsumerReceipt@1。Eligibility先到可等待Gateway Usage,Gateway Usage先到可等待Eligibility,Cost独立处理;等待仍必须是可恢复Work而不是成功后丢弃。临时Owner Read/锁不可用返回retryable且不写Receipt;同Revision异摘要或严格身份冲突返回non_retryable。相同Event重投取回首次Receipt,新Gateway Revision重新置Fence Pending,旧Receipt不能吞掉新Current。
Gateway Event的Source读取与Current验证必须分阶段。Worker先用Event-boundmetering_source_event_bootstrap读取Event明确引用的Availability、Available Evidence与Route Cost Binding;Token绑定Event ID/Type/Schema/Envelope摘要、Aggregate/Revision和Payload中的Attempt/Deployment/Binding/Route/Kind/Dimension/Snapshot身份,不携带尚未从正文构造的Basis、Candidate或Validation Operation。只有严格解码并重算正文后,Worker才构造Eligibility Basis或Provider Cost Candidate、确定性派生对应Validation Operation,并用operation-bound验证Purpose取得Current Receipt;Bootstrap响应本身不能完成Work、写MeterEvent/ProviderCostFact或充当Receipt。
Provider Cost Work对每个Current Availability状态使用严格路径,而不是把缺证据解释成零成本:available成功固定available_applied并要求同一Metering事务提交ProviderCostFact及Availability/Route两份验证Receipt;pending固定waiting_for_evidence;not_reported|unavailable固定no_cost_evidence且禁止ProviderCostFact。只有Gateway Availability正文已经是conflicting时,才写ProviderCostAvailabilityProcessingFact@1.outcome={state=reconciliation_required; availabilityState=conflicting; stateReasonCode},并将Processing Fact、Work Finish与Lane CAS原子提交;Processing Fact完整四元组是Metering-owned耐久锚点,后继Gateway Current Event以新Generation恢复。
available路径中已证明的Schema/Identity/Mapping/Normalization/Receipt/Repository冲突写内容寻址ProviderCostProcessingConflictEvidence@1 + MeteringInputWorkAttemptBlocked@1,将Fence置为blocked_on_processing_conflict,保留原Work/Operation/Generation及activeWork.state=pending并释放Lease;禁止写Processing Fact、ProviderCostFact、Work Finish、Lane终态CAS或Case。Owner修复后同Work重新Lease即可重试,不需要不存在的reopen Event。Owner暂不可读、网络/锁/容量故障及其他可重试错误连上述领域Conflict/Blocked也不写,只返回retryable。Receipt必须区分Blocked与Gateway-conflicting终态,拒绝未知Stage、Conflict Evidence错摘要、Blocked夹带Finish/Processing、终态缺Processing/Finish或任何reconciliationCaseRef。当前Delivery Set没有Provider Cost Processing/Conflict Event:首期只能从Metering定向诊断并发送明确不完整的best-effort告警,不能创建Operations Reconciliation Case或声称已有可靠Admin队列;可靠投影须先新增Canonical Event、core_operations Mandatory Delivery、Owner Read与严格Projection Contract。
Billing Finalization Coordinator Consumer必须严格区分RunFinalizationClosed@1 | FinalizationReconciliationResolutionAppliedEvent@1。Closure分支验证完整Envelope/Payload、Closure Fact/Attempt Set/Run Manifest身份,再在Billing Inbox事务锁定Finalization Fence、推进Closure Watermark、冻结Trigger Snapshot、insert-or-compare新Work并写work_scheduled Receipt。Admission原子创建的首代Work在该Closure Watermark与Coordinator Receipt尚未提交时只能追加closure_event_not_applied阻塞Finish:不得调用Execution Owner Read,不得形成Finalization Input/Decision或写Ledger;Closure Event Consumer提交Receipt后用新Generation唤醒。最终成功的BillingFinalizationCommitted@1必须绑定提交时Fence Revision、这条Closure Event身份/Envelope摘要及Coordinator Receipt完整四元组并逐项复验。Execution Precondition/Finalization Owner Read的conflicting只返回RunExecutionFinalizationConflictFact@1完整四元组;Billing可把它纳入正式Finalization Decision,并且只有该Worker按Run级Case Identity协议创建Operations Case,Read端不得返回Case Ref。Finalization进入Reconciliation前该唯一Event/Receipt必须已提交;之后重投只返回首次Receipt,不再改写Fence或创建Work。
Resolution分支必须验证Event中Identity Reservation与Resolution Fact四元组、Case/Generation/Revision/Operation及当前Fence。仅当该Resolution是当前代期望输入时,才能在同一Inbox事务推进Resolution Watermark、把Case期间累积的Closure/Settlement Watermark一次性冻结进新Trigger/Work、退出Case协调态并写work_scheduled Receipt;只有同一Resolution已应用并创建重评Work,或该Work的重评已推进/被该次Resolution吸收时,旧投递才写already_superseded且不创建第二Work。普通Closure/Settlement输入在Case持有态只能推进Watermark,不能改变Resolution的Current/Stale判断。两类Event在Fence异常不可绑定、Owner Fact暂不可读或Current/Stale无法判定时必须返回retryable且不写成功Receipt;重复Event不能增加Work Generation。
Billing Catch-up Consumer必须先严格验证SettlementInputCurrentChanged@1完整Envelope/Payload、Canonical Event摘要、Aggregate/Revision、Snapshot五元组和直接前驱,再在一个Billing Inbox事务中锁定同维度Watermark并单调处理watermark_advanced | already_observed。Reservation尚未成功Finalization且Fence不在Case中时,必须推进Finalization Fence、冻结Trigger Snapshot并产生/合并Finalization Work;Fence处于正式Finalization reconciliation_required时只更新Watermark并写held_for_resolution Receipt,不创建Work或退出Case;成功Finalization后必须重新打开Dimension Catch-up Fence并insert-or-compare Catch-up Work。所选Watermark/Fence/Work/持有结果与BillingSettlementCatchUpConsumerReceipt@1最后在同一Inbox事务提交。极端乱序令Fence不可绑定时只保存Inbox Watermark并返回retryable,不得写成功Receipt或Ack。Event自身严格冲突返回non_retryable并让源Delivery进入Quarantine;这与Catch-up Worker随后遇到Owner确定性冲突不同,后者写SettlementInputConflictFact@1 + BillingSettlementCatchUpWorkAttemptBlocked@1,令Fence=blocked_on_owner_conflict且保留同一activeWork.state=pending,禁止Finished/Fence终态/Case。Consumer提交后响应丢失、Dispatcher确认前崩溃或Lease换主都依靠(consumerName,eventId)取回首次Receipt,不能重复创建Work或改写首次Outcome。
Catch-up Consumer Receipt与billing_catch_up_fence Current Validation Receipt是两个不同Owner、Purpose和时点的对象:前者只证明Event已转为耐久Watermark/Work,后者由Catch-up Worker按确定性billingCatchUpValidationOperationId取得并只允许推进Fence。两者都不能直接写Ledger;只有随后独立的post_finalization_transition验证、Consumption与Billing事务可以改变账本。
Catch-up Worker的transition_applied | already_evaluated也不是caught_up:它必须在同一Billing事务把Fence lastEvaluatedInput重基线到结果指针,递增successorWorkGeneration,按新Basis确定性派生successorBillingCatchUpValidationOperationId,并insert-or-compare新pending_verification Work完整四元组;Finished Fact反向冻结这些字段。只有该后继Work的专用billing_catch_up_fence Current Receipt在Fence Revision和Watermark仍不变时才能进入caught_up。Transition与Successor Work中任一写入失败必须全部回滚,响应丢失从旧Operation恢复同一结果/后继Work,不得清空Active Work等待不存在的新Event。
每个Operations Projector以 consumerName + projectionName + projectionVersion + eventId 建立 Receipt,并在一个数据库事务内完成:
- 验证 Event Schema、Envelope/Payload 一致性,以及覆盖完整 Canonical Event 的
eventEnvelopeDigestAlgorithmVersion + eventEnvelopeSha256。 - 锁定或创建对应 Projection Aggregate。
- 检查单 Aggregate Revision;处理重复、连续版本或显式缺口。
- 写入 Projection Mutation 与 Exact Identifier Index。
- 写入包含 Canonical Event Digest Algorithm Version/SHA-256 的 Consumer Receipt。
- 写入该次 Projection Build 的运行摘要与最后成功应用时间;该摘要不用于推断此前记录已经全部处理。
任一步失败全部回滚。已存在相同 Receipt 且算法版本/摘要一致时返回幂等 Ack;同一 eventId 的算法版本或摘要不同立即隔离。遇到 Aggregate Revision 缺口时记录待补 Gap,但不写 applied receipt;Consumer 返回可恢复等待,等缺失 Revision 到达后再应用。只有经证据确认不可恢复时才按治理策略隔离,不能靠最大重试次数把未知事实静默变成成功。
Ordering 与完整性证明
- 不承诺全局顺序;Dispatcher 的扫描顺序只是公平与可预测性策略。
- 单 Aggregate 以
aggregateRevision判断重复和缺口。 - 跨 Aggregate 链通过
causationId、operationId和correlationId关联;例如 Wallet Event 可能先于 Run Event 被投影。 - 不采用全局源序号、已扫描最大值、连续位置或任何单值进度标记证明“此前全部完成”。
complete_for_supported_scope必须由 Core Operations Verifier/Projector Role 在一个数据库 Consistent Snapshot 内证明:当前 Projector 支持范围与冻结 Delivery Set 中没有显式未完成 Mandatory Delivery;适用 Event 对当前 Projection Version 的 Event Receipt Anti-Join 为空;已接受 Observation 对当前版本的 Observation Receipt Anti-Join 为空;不存在未解决 Aggregate Revision Gap、影响范围的 Quarantine/Unsupported Schema,且 Semantic Check 通过。Observation 完整性只针对 Core Accepted Source,不代表 Edge 请求全集。- 该Operations完整性快照只覆盖其Source Query Definition明确纳入的Core Event/Accepted Observation和
core_operationsDelivery,不伪称横跨Metering事实库证明Billing已追平。Settlement Input链由Metering侧billing_settlement_catch_up非终态Delivery/Quarantine、Billing Watermark/Fence/Work与Catch-up Receipt各自证明;若未来要在统一Operations Query展示,必须新增版本化跨服务Owner Read与Definition,不能让Projector读取远端表或把新Event当成五类Billing经济边界之一。 - 上述状态只证明
completenessAsOf.snapshotId/evaluatedAt的完整 Proof Manifest,不自动代表稍后的查询时点,也不能直接当作 Cutover Readiness。只读 Query 必须用相同 Query/Semantic Definition,在一个 Consistent Snapshot 内将当前双 Source、Mandatory Delivery/Receipt、Gap、Quarantine、Unsupported Schema 与 Semantic Check 全部输入和 Snapshot Manifest 做精确比较。全部不变才返回verificationState=verified + snapshotCurrency=verified_current;存在完成且定义可比的 Snapshot、Guard 明确确认变化时返回verified + stale,确定 Gap 对应partial_source_gap,无确定 Gap 对应unknown,且 Stale 禁止返回 Complete。没有完成 Snapshot、定义不可比或 Guard 无法执行/比较时返回unverified + completenessAsOf=null + integrity=null + snapshotCurrency=unknown + projectionCompleteness=unknown + verificationReason。Cutover 在冻结 Candidate 前先暂停并排空 Shadow Claim;短 Consistent Snapshot 事务原子冻结 Candidate 精确成员/定义/Expected Receipt Coordinates并绑定新的 Shadow Claim Fencing Token,恢复后 Shadow 只对该边界 Catch-up。随后停 Active Claim/排空,在最终事务复核 Candidate、生成新的不可变 Readiness Proof 并同事务切换。边界后 Pending 不进入 Readiness,旧 Shadow Token/边界外 Mutation 被拒绝,恢复后由新 Active 追平。 - Operation Summary 可以暂时只见 Run;Wallet Event 尚未到达时 Reservation 字段保持缺失并显示来源延迟。Wallet-first 只能进入
pending_unlinked_reservations,在 Run Event 完成因果与等值校验前不得创建 Operation Summary 或生成伪状态。
Core Operations Dispatcher自身的Delivery、Attempt、Receipt、Build、Gap与Quarantine属于operations_runtime;Metering Dispatcher的Delivery/Attempt/Quarantine属于Metering事实边界,Billing Consumer Receipt、Watermark/Fence/Work属于Billing事实边界,三者不得互相复制为Owner状态。Core Operations Query只读取本地Runtime并与投影组合;禁止为了显示Publisher状态再向同一个Outbox递归追加Publisher Event,也禁止直连Metering/Billing数据库伪造全链完整性。
重试、隔离与人工处置
Retry Policy 由错误分类驱动:
- 数据库暂时不可用、Consumer 容量保护和可恢复网络失败可以重试。
- 严格 Schema 失败、未知版本、摘要冲突和不变量破坏不可盲目重试。
- 最大尝试、退避和隔离条件由配置、SLO 和容量验证确定,不在代码里散落固定数字。
- Quarantine 必须产生 Metrics、结构化低敏日志和可选 Reconciliation/Incident Link。
- 首期 Admin 只读展示,不提供“跳过事件”“强制成功”“修改 Payload”或任意重放按钮。
受控人工恢复必须在后续 Admin Command 契约中显式设计:绑定 Workforce、JIT、Case、影响预览、幂等键和 Audit;恢复只能重新投递原 Event 或部署修复后的 Projector,不能改写历史事件。
保留与重建
首期保留已投递 Outbox,以支持 Projection 重建和证据复核。具体热数据保留、归档、压缩与删除窗口由治理策略、恢复目标、数据分类和容量证据决定。
删除已投递事件前必须同时证明:
- Source Snapshot 只是成员/摘要/完整性证明,不是恢复快照,不能据此删除唯一 Event Payload;
- 全量 Projection State Snapshot 是派生状态,不能替代 Canonical Source,也不能独自授权删除;
- Canonical Event Envelope/Payload、
eventEnvelopeDigestAlgorithmVersion + eventEnvelopeSha256、Delivery Set Version/Definition 已进入完整、不可变且可读取的恢复归档,并且该归档继续作为 Rebuild/Verifier 批准 Source View 的组成部分;逻辑sourceKind + sourceId + sourceDigest(kind + algorithmVersion + sha256)不因冷热迁移消失; - Schema 与 Consumer 运行时仍可获得;
- 在隔离恢复环境完成归档校验与重放/恢复演练,并把归档摘要、运行时版本与演练证据写入受审计恢复基线;
- Audit、账务、事故和法律保留要求已满足;
- 删除过程不影响未完成 Delivery、Quarantine 或 Shadow Rebuild。
Applied Receipt 的保留必须绑定 Source 与 Projection 生命周期:Active/Shadow Projection Receipt、仍被 Delivery terminal Ack 引用的 Receipt,以及对应 Source 仍可投递/重建时的 Receipt 均不得删除。Projection Version 已退休且回滚窗口结束、Active Pointer/Consumer Target/Delivery/Quarantine/Legal Hold 均无引用,并且上述恢复基线验证通过后,Retention Procedure 才能清理该退休版本的 Receipt 与派生 Projection;Canonical Event Source 仍保留在批准归档 Source View。测试必须拒绝以 Projection State Snapshot 替代 Source、错/缺归档仍声明 Complete/Cutover、先删 Receipt、仅删 Receipt、删除 Active/Shadow Receipt和保留可重放 Source 却删除幂等 Receipt。
运维指标
至少暴露以下低基数 Metrics:
- Pending Delivery 数量与最旧年龄;
- Claim、Ack、Retry、Lease Expiry 与 Quarantine 速率;
- 按 Destination、Event Type、Schema Version 和规范化错误码聚合的失败;
- Consumer 应用延迟、Receipt Anti-Join 未应用数量、Revision Gap 与 Schema Reject;
- Worker 版本、部署区域和健康状态。
标签不得包含 Organization、User、Request、Trace、Run、Reservation、Event 或 Provider Task ID。精确 ID 只能进入授权查询索引和脱敏日志字段,不能进入时序标签。