diff --git a/docs/README.md b/docs/README.md index d9ca2e1..4be4c62 100644 --- a/docs/README.md +++ b/docs/README.md @@ -1,6 +1,6 @@ # 设计文档入口 -当前实现核对日期:2026-09-08。文档中的目标能力不等于已实现;测试通过不等于现场已发布。 +当前实现核对日期:2026-09-09(与代码基线 d53a0a1 一致)。文档中的目标能力不等于已实现;测试通过不等于现场已发布。 | 文档 | 唯一职责 | |---|---| diff --git a/docs/architecture.md b/docs/architecture.md index 154274b..6c6554e 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -29,7 +29,7 @@ CIIMS / AODB 等上游 ┌──────────────── msgexchange-v2(单实例)────────────────┐ │ ingress:发现报文 → PostgreSQL 持久化入队 │ │ │ │ -│ processing:取 FIFO 队头 → 解析 / 去重 → Handler 决策 │ +│ processing:取 FIFO 队头 → 解析 / 去重 → 处理器决策 │ │ └─ PG 单事务:航班变更 + 终态 + 待发事件│ │ │ │ jobs:独立维护线程(回填补偿 / 历史归档 / 留痕清理) │ @@ -52,9 +52,9 @@ CIIMS / AODB 等上游 |---|---| | `ingress` | 轮询信箱、持久化入队、补偿重扫及兼容 HTTP 写入;不解析业务报文。 | | `codec` | XML 解码,区分非法报文与可修复的解码失败。 | -| `processing` | FIFO 调度、业务身份绑定与去重、处理器决策(SCHD/FLOP/FDEL/ADFT)、事务提交。处理器只产出决策与落库计划,不直接触碰 Kafka。 | +| `processing` | FIFO 调度、业务身份绑定与去重、处理器决策与落库(SCHD/FLOP/FDEL/ADFT):处理器在锁事务内完成状态写入、事件与回填待办登记,不直接触碰 Kafka。 | | `delivery` | 消费待发事件,负责按目标保序、`schd` 聚合、投递和失败重试。 | -| `jobs` | 回填补偿扫描、历史归档与物理清除(§8.2)、留痕保留期清理;独立线程执行,不参与 FIFO。 | +| `jobs` | 回填补偿扫描、航班历史清理与留痕保留期清理;独立 job 线程执行(调度见 design.md §6.1,红线见 flight-state.md §6),不参与 FIFO。 | | `domain` / `config` | 领域状态、事件和决策模型,以及运行参数。 | | `infra` | 仓储(JDBC/stub)、外部适配器、重试、健康检查与日志;通过接口隔离基础设施。 | @@ -78,8 +78,8 @@ CIIMS / AODB 等上游 - **消息严格 FIFO**:队头失败并退避时,后续消息仍不能越过它。只有队头完成或按失败策略进入终态后,队列才继续推进。收报重扫和水位设计必须防止较小 ID 漏入队而被后续消息越过。 - **动态状态单写者**:`FLIGHT_SCHD` 及明细表只由主泵单线程写入。事务内第一步对 `PIPELINE_LOCK` 单行 `SELECT ... FOR UPDATE` 互斥;该方案在 PG/Oracle 11g 均无需数据库扩展或额外 DBA 特权。不能通过增加实例或处理线程直接扩容。 - **身份去重**:同一业务身份只能绑定一条有效处理记录,重复报文不应再次产生业务副作用。具体身份组成和重放规则见设计文档。 -- **快照可恢复**:快照以 `PROC_STATE` 成功终态判定重放(§5.1);每次成功写入推进 `STATE_VERSION`;`OPERATION_DAY` 一经确定不可变(§5.3)。 -- **物理清除只发生在历史归档**:删除一律先标记(FDEL)或由生命周期清除;历史存储未接通时必须删 0 条(§8.2)。 +- **快照可恢复**:快照以 `PROC_STATE` 成功终态判定重放;每次成功写入推进 `STATE_VERSION`;`OPERATION_DAY` 一经确定不可变(见 flight-state.md §2.1)。 +- **物理清除只发生在历史归档**:删除一律先标记(FDEL)或由生命周期清除;历史存储未接通时必须删 0 条(见 flight-state.md §6)。 这些约束优先于吞吐量优化。单写者降低了并发复杂度,代价是队头阻塞和吞吐上限;如需并行化,必须先重新定义顺序与状态归属,不能只调整线程数。 @@ -89,7 +89,7 @@ CIIMS / AODB 等上游 |---|---|---| | 自有 PostgreSQL | 单行锁 `PIPELINE_LOCK`、处理状态 `PROC_STATE`、待发事件 `MSG_EVENT`、请求跟踪 `REQ_TRACK`、回填待办 `BACKFILL_TODO`、航班当前态 `FLIGHT_SCHD` + 9 张明细表、留痕 `SCHD_SNAP_LOG` | 本系统唯一业务数据库。消息处理、状态推进与待发事件在单事务内原子提交;本地事务只在此库。 | | 共享 MySQL | `CMINMSGS` 入站信箱、`COUTMSGS` 出站信箱 | 外部系统所有。仅执行约定的信箱读写和处理标记回填,不建表、不迁移 schema、不写历史表。兼容 HTTP 入口可按既有契约写入入站信箱。 | -**不使用跨库事务。** PG 事务只能保证“处理结果与待发事件一起提交”,不能覆盖 Redis 更新、MySQL 回填或 Kafka 发送。跨存储依靠幂等、重试和持久化补偿恢复: +**不使用跨库事务。** PG 事务只能保证“处理结果与待发事件一起提交”,不能覆盖 MySQL 回填或 Kafka 发送等外部副作用。跨存储依靠幂等、重试和持久化补偿恢复: | 中断位置 | 恢复要求 | |---|---| @@ -100,24 +100,16 @@ CIIMS / AODB 等上游 对外投递按**至少一次**设计,不承诺端到端恰好一次。Kafka 生产者幂等不能消除应用重启或 outbox 重发带来的所有重复。 -## 7. 关键决策索引 +## 7. 关键决策 -保留 D1–D12 编号,便于设计文档和工程历史引用;以下是决策摘要,而非完成清单。 +仅保留仍具约束价值、且无法从正文(§4–§6、design.md)直接推出的决策,按 D1–D4 连续编号供正文与 design.md 引用;其余曾编号条目(严格 FIFO、stub 门控、本地事务、UNSUPPORTED 处理等)已在正文以约束形式表达,不再重复列表。状态只反映是否已落地,不代表决策被撤销。 -| 编号 | 决策及理由 | -|---|---| -| D1 | 业务报文严格 FIFO,优先保护航班状态的时序正确性;维护作业独立线程执行,不参与消息序。 | -| D2 | 阶段 B 暂缓。历史写入成功后才可生成删除事件;顺序调用本身不保证原子性,恢复与去重方案需在启用前补齐。 | -| D3 | `schd` 只从 `flushSchd` 聚合发送,减少已被覆盖的中间状态通知。 | -| D4 | 未实现的报文类型按 `UNSUPPORTED` 可恢复失败处理,不当作非法报文直接丢弃;补齐能力后按重放规则恢复。 | -| D5 | 在持有具体消息或批次上下文的位置记录失败和退避;不吞掉线程中断或 JVM 严重错误。 | -| D6 | 基础设施通过接口注入,时间通过 `Clock` 注入,便于确定性测试顺序、重试与超时。 | -| D7 | 内存 stub 仅显式开启时装配,生产禁止使用,避免把未持久化的数据误当作已落库。 | -| D8 | 使用编译期依赖注入,并以启动冒烟测试验证关键 Bean 装配。 | -| D9 | 自有 PG 内完成本地事务,共享 MySQL 仅作信箱;跨存储采用补偿,不使用 XA。 | -| D10 | 动态状态单写者,生产只允许一个活动实例;多实例必须先具备可靠的排他保护。 | -| D11 | Kafka 生产要求 `acks=all`、`enable.idempotence=true`、`max.in.flight=1`,切流前验证 Broker 兼容性;不允许通过关闭幂等来满足生产接入。 | -| D12 | 仅将自有库终态记录归档到 `PROC_STATE_HST`,不侵入共享库的表结构或保留策略。 | +| 编号 | 决策及理由 | 当前状态 | +|---|---|---| +| D1 | 航班清场只在历史写入成功后进行,未接通时删 0 条;未经 FDEL 的清场须先补发删除事件。ES 历史投影(阶段 B)暂缓。 | 红线已实现于 `HistorySweepJob`;恢复/去重方案未闭合 | +| D2 | 动态状态单写者,生产只允许一个活动实例;多实例必须先具备可靠的排他保护。 | 事务行锁已实现;实例级排他未完成 | +| D3 | Kafka 生产要求 `acks=all`、`enable.idempotence=true`、`max.in.flight=1`;不允许通过关闭幂等来满足生产接入。 | 约束未强制:默认 in-flight=5,且可用环境变量覆盖 | +| D4 | 自有库终态记录只归档到 `PROC_STATE_HST`,不侵入共享库的表结构或保留策略。 | 目标表未建,尚无归档作业 | ## 8. 部署、切换与运维 @@ -137,14 +129,14 @@ CIIMS / AODB 等上游 ## 9. 当前实现与上线门槛 -当前已有管道骨架、重试机制、部分 JDBC 适配和开发环境 stub 冒烟能力,**不能据此认定生产链路已闭环**。默认配置关闭管道自动启动及真实数据库/信箱适配。 +当前已实现收报入队与判重、严格 FIFO 主泵、SCHD/FLOP/FDEL/ADFT 处理器与 PG 单事务写入、outbox 与 `schd` 聚合投递骨架、回填补偿与航班历史清理脚手架。**这些只证明机制可用,不证明生产链路已闭环**:默认配置不自动启动管道,真实数据库、信箱与出站适配必须显式开启(`msgx.pipeline.autostart`、`mailbox.shared-mysql.enabled`、真实 `DeliveryPort` 适配)。 -上线前至少需要完成并验证: +上线前必须完成并验证: -- 真实 PG 事务、信箱水位与补扫、回填补偿、出站信箱,以及所需业务 Handler。 -- `FLIGHT_SCHD` 事务原子性、`STATE_VERSION` 推进与 `OPERATION_DAY` 不可变校验、故障中断回滚恢复。 -- FIFO、身份去重、FDEL/ADFT 生命周期、历史归档顺序和投递故障下的回归测试。 -- 生产启动校验、单实例排他保护、影子隔离和 Kafka 配置约束;当前配置仍允许 Kafka 参数覆盖,且默认 in-flight 值与 D11 要求不同。 -- 死信告警、人工重放、端到端追踪、积压指标及安全边界。 +- 真实 PG + 共享 MySQL 信箱的端到端处理、补偿与投递,以及出站信箱适配;未闭合的缺口清单见 design.md §10。 +- 航班状态不变量与恢复证据:`STATE_VERSION` 推进、`OPERATION_DAY` 不可变、故障中断回滚(见 flight-state.md §7)。 +- FIFO 越序、身份去重、FDEL/ADFT、清场顺序与投递故障的回归测试(见 design.md §8)。 +- 单实例排他保护与启动校验、影子隔离、Kafka 生产配置约束——当前配置允许环境变量覆盖 `acks`/幂等/in-flight,且默认 in-flight 值与 D3 不同,切流前必须按 D3 收敛。 +- 死信与一致性异常的告警、可执行的人工重放流程、端到端追踪、积压指标与安全边界。 -具体实现差异见 design.md §10 与 ACM2-29 核查报告,工作由 Plane 跟踪。航班状态的有效规则统一见 flight-state.md;历史阶段口径不覆盖当前规则。 +验收与进度由 Plane 跟踪,缺口逐项见 design.md §10 与 user-stories.md;航班状态规则统一以 flight-state.md 为准。 diff --git a/docs/design.md b/docs/design.md index 70824d0..c6b528b 100644 --- a/docs/design.md +++ b/docs/design.md @@ -4,9 +4,9 @@ 本文说明模块如何协作、状态如何流转,以及失败后如何恢复。系统范围、存储归属和部署约束见 [architecture.md](architecture.md),不在这里重复。 -阶段 A 采纳 ACM2-28 选项 C 定案:运营航班权威状态落自有 PostgreSQL(`FLIGHT_SCHD` 及资源明细表),Redis 彻底退出动态权威与全部写路径。阶段 B 的历史投影和清场暂不启用。 +运营航班权威状态只落自有 PostgreSQL(`FLIGHT_SCHD` 及资源明细表),Redis 已彻底退出动态权威与全部写路径;ES 历史投影属暂缓范围,不参与当前设计。 -本文的流程是验收目标,当前差异见 §10。航班状态设计统一由 +本文描述处理机制与流程;尚未交付的能力在本文件中明确标注,并以 §10 差异为准。航班状态规则统一由 [航班状态设计](flight-state.md) 维护。 ## 2. 数据与领域模型 @@ -17,21 +17,19 @@ | 记录 | 用途 | 关键约束 | |---|---|---| -| `PROC_STATE` | 入站消息的处理状态、身份、重试次数和错误原因 | `CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头。 | -| `MSG_EVENT` | 等待投递的事件(outbox) | `EVENT_ID` 决定投递顺序;`TARGET` 区分目标;`PARTITION_KEY` 在 `schd` 中为 `FLID`。 | -| `PUMP_JOB` | 持久化维护作业 | 状态为 `QUEUED / RUNNING / DONE / FAILED`;不与业务消息共用排序序号。 | -| `REQ_TRACK` | 上游请求及应答关联 | 保存请求类型、参数、出站信箱 ID、发送和完成时间;同类只允许一个开放请求。 | -| `REF_MASTER` | 静态参考数据 | `(RTYPE, RKEY)` 唯一,`SOURCE` 记录数据来源。 | -| `FLIGHT_SCHD` | 航班标量及单值异常、3 类文本载荷 | `FLID` 主键;运营日、版本和消息身份用于追踪。10 类重复集合存于 9 张明细表,见航班状态设计。 | -| `BACKFILL_TODO` | 共享信箱回填重试 | 当前失败后落账,终态与待办同事务持久化尚未完成。 | -| `PROC_STATE_HST` | 终态处理记录的归档目标 | 属于目标设计,当前迁移尚未建表;不得改写为共享库历史表。 | +| `PROC_STATE` | 入站消息的处理状态、身份、重试次数和错误原因 | `MSG_ID = CMINMSGS_ID` 主键防止重复入队;`IDENTITY_KEY` 唯一约束防止业务重复;按最小未完成消息 ID 取队头。 | +| `MSG_EVENT` | 等待投递的事件(outbox) | `EVENT_ID` 决定投递顺序;`TARGET` 区分 `KAFKA:msg` / `KAFKA:schd`;`PARTITION_KEY` 恒为 `FLID`;`EVENT_TYPE` 区分 UPSERT 与 TOMBSTONE。 | +| `REQ_TRACK` | 上游请求及应答关联 | 保存请求类型、覆盖运营日、发送方、出站信箱 ID 与发送/完成时间;同类只允许一个开放请求。登记、超时与应答匹配尚未实现(§10)。 | +| `REF_MASTER` | 静态参考数据(目标表) | `(RTYPE, RKEY)` 唯一;尚未建表,客户端与刷新流程见 user-stories.md US-13/US-14。 | +| `FLIGHT_SCHD` | 航班标量及单值异常字段 | `FLID` 主键;`OPERATION_DAY` 一经确定不可变;版本与最近消息 ID 用于追踪。变长资源集合存于 9 张明细表,规则见 [flight-state.md](flight-state.md) §2,不在此重复。 | +| `BACKFILL_TODO` | 共享信箱回填重试 | 回填意图与业务终态同事务预登记,提交后回填成功即删除;待办再次落账失败的崩溃窗口仍在(§10)。 | +| `PROC_STATE_HST` | 终态处理记录的归档目标 | 尚未建表;不得改写为共享库历史表。 | -字段与索引定义以 `src/main/resources/db/migration/` 为准(含 `V1.1.0__flight_schd.sql`)。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理。 +字段与索引定义以 `src/main/resources/db/migration/V1__flight_state_baseline.sql` 为准;Oracle 11g 的迁移形态见 `src/main/resources/db/migration/oracle11g/`(占位,未接入任何 Flyway 配置)。报文原文仍从共享信箱读取,因此必须协调原文保留期,不能在消息尚需处理或重放时提前清理。 -Redis 已彻底退出动态权威与写路径;`FLIGHT_SCHD` 及资源明细表在自有 PG 中由主泵单线程独占写入。 ### 2.2 消息、身份与决策 -`XmlCodec` 将 XML 解码为 `DecodedMessage`,包含 `SNDR / TYPE / STYP / SEQN / DTTM` 元数据和业务载荷。`MsgKind` 区分 `SCHD` 与 `FLOP` 子类型;Handler 查找与日志类型标识使用同一套映射。 +`XmlCodec`(实装 `JacksonXmlCodec`)将 XML 解码为 `DecodedMessage`,包含 `SNDR / TYPE / STYP / SEQN / DTTM` 元数据、`MsgKind` 与业务载荷。解码失败区分 `MALFORMED`(报文非法,不重试)与可随 codec 修复的编码错误。`MsgKind` 为一等分派键:`Schd(RESP/DNLD/ADFT)`、`Flop`、`Fdel`、`Unsupported`。 业务身份统一由 `Identity.of` 生成: @@ -41,35 +39,29 @@ SNDR | TYPE | STYP | SEQN 接收时只按信箱 ID 去重;解码后才首次绑定业务身份。重试保留原有绑定,不能把自己判为重复消息。身份被另一条记录占用时,当前消息转为 `SKIPPED`,记录 `duplicate-of:`。是否加入日期边界取决于上游序号重置规则,默认关闭;上线后不能随意更换身份算法。 -Handler 是纯函数: - -```text -Handler.decide(flightView, message) → Decision -Decision = 航班变更 + msg 通知 + schd 状态 + 出站意图 + 静态数据变更 -``` - -Handler 不写 Kafka 或数据库。主泵负责应用决策;各类变更在自有 PG 单事务内提交。 +分派与落库由 `MessageProcessor` 统一执行:按 `MsgKind` 把已绑定身份的队头消息交给对应处理器(DNLD/RESP → `ScheduleProcessor`,ADFT → `AdftProcessor`,FLOP → `FlopProcessor`,FDEL → `FdelProcessor`,其余 → `FAILED(UNSUPPORTED)`)。处理器在 `PIPELINE_LOCK` 事务内读取当前完整态、计算下一态并落库、登记事件与回填待办,不直接触碰 Kafka;提交与终态迁移仍由主泵边界负责(见 flight-state.md §4)。 ### 2.3 状态与错误分类 ```text 处理:PENDING / FAILED → SUCCEEDED(成功) - → SKIPPED(忽略或重复) - → FAILED(等待重试) - → DEAD(非法报文或重试耗尽) + → SKIPPED(业务重复;忽略/无匹配分支见 US-04/US-06,尚未实现) + → FAILED(等待退避重试) + → DEAD(MALFORMED / PROTOCOL / EXHAUSTED,均需人工处置) 投递:PENDING → SENT → PENDING(退避后重试) - → DEAD(重试耗尽) + → DEAD(重试耗尽,保留记录作 DLQ) ``` `SUCCEEDED / SKIPPED / DEAD` 是处理终态,不再阻塞后续消息;`FAILED` 不是终态,仍占据队头。`DEAD` 表示需要处置,不等于业务成功。 | 错误类别 | 处理方式 | |---|---| -| `MALFORMED` | 报文非法,直接 `DEAD`,不在原记录重放白名单内。 | +| `MALFORMED` | 报文非法或原文缺失,直接 `DEAD`,不在重放白名单内。 | +| `PROTOCOL` | 整包协议拒绝(运营日冲突、声明数量不符、缺载荷等),立即 `DEAD`,整包不落地、不重试。 | | `CODEC_ERROR` | 解码能力问题,退避重试;修复后允许重放。 | -| `UNSUPPORTED` | Handler 或快照能力未实现,按可恢复失败处理,不直接当作非法报文;仍受重试上限约束。 | +| `UNSUPPORTED` | 处理器或快照能力未实现,按可恢复失败处理,不直接当作非法报文;仍受重试上限约束。 | | `INFRA` | 基础设施或执行异常,退避重试。 | | `EXHAUSTED` | 重试耗尽或滞留超时,转 `DEAD`,人工复核后允许重放。 | @@ -77,101 +69,111 @@ Handler 不写 Kafka 或数据库。主泵负责应用决策;各类变更在 ### 3.1 收报 -`InboxPoller` 默认每秒读取未处理信箱记录,在 PG 建立 `PENDING`,不解析业务载荷。PG 插入必须按信箱 ID 幂等,失败由后续扫描补建。 +`InboxPoller` 默认每秒按 ID 升序、有限批次(`claim-batch`,默认 50)读取 `DATE_PROCESSED IS NULL` 的信箱记录,经 `InboxEnqueue` 在自有 PG 建立 `PENDING`;重复扫描幂等,入队失败留待下一轮。收报层不解析业务载荷,也不回填已处理标记。 -水位优化分为两条路径:快路径读取水位之后的新记录,补偿路径重扫遗漏的未处理记录。只有本批 PG 入队全部确认后才能推进水位。**水位不是已处理标记,也不能单独证明较小 ID 已收齐**;迟提交和补扫场景的顺序保证需要在启用前验证。 +快路径水位与受控补扫是目标形态,当前仍按 `afterId=0` 每轮全量扫描(§10)。**水位不是已处理标记,也不能单独证明较小 ID 已收齐**;迟提交、ID 空洞与补扫场景的顺序保证在启用前必须验证。 兼容 HTTP 入口执行“写入共享信箱 → PG 入队”。两步不在同一事务中:信箱成功而 PG 失败时,原文不能丢失,由轮询补建;客户端失败重试可能再次写信箱,业务身份去重仍然必需。 ### 3.2 主泵调度 -每次 `Pump.tick`: +每次 `Pump.tick` 只围绕最小未完成消息(`PENDING` 与 `FAILED` 都占队头): -1. 读取最小未完成消息,必须包含 `FAILED`,不能只查当前可执行的记录。 -2. 无消息,或队头仍在退避窗口内时,允许执行一个维护作业;消息已可执行时优先处理消息。 -3. 队头达到重试或滞留上限时转 `DEAD(EXHAUSTED)`;未到重试时间则等待,不领取后续消息。 -4. 其余情况调用 `MessageProcessor.processOne`。 +1. 无队头:按轮询间隔休眠。 +2. 队头 `FAILED` 且未到 `next_attempt_at`:未超限则等到可重试时刻;已达重试上限或超过队头滞留时限(`head-deadline`,默认 10 分钟)则转 `DEAD(EXHAUSTED)`。 +3. 队头可执行:调用 `MessageProcessor.processOne`,失败迁移在该边界内完成。 -作业执行时长和饥饿边界需要限制,不能用长期作业阻塞已到期消息。滞留超时应基于稳定的起始时刻,不能用每次失败都会刷新的 `updatedAt` 代替。 +维护作业由独立 job 线程调度(§6.1),不占用消息循环;作业有界且不使到期消息无限饥饿。滞留判据当前仍使用 `updatedAt` 与直取系统时间,未基于稳定起始时刻(§10)。 ### 3.3 单条处理 ```text -读取原文 → 解码 → 忽略规则 → 首次绑定身份 - ├─ RESP / DNLD:快照流程 - └─ 其他:Handler 决策 - ↓ - PG 单事务:FLIGHT_SCHD 增量更新 + 待发事件 + 处理终态 - ↓ - 提交后补偿回填信箱 +读取原文 → 解码(MALFORMED → DEAD;编码错误 → FAILED 退避) + → 首次绑定身份(冲突 → SKIPPED,记 duplicate-of) + → 按 MsgKind 分派处理器 + DNLD / RESP → ScheduleProcessor(快照事务) + ADFT / FLOP / FDEL → Adft / Flop / FdelProcessor(单航班事务) + Unsupported → FAILED(UNSUPPORTED) + PG 单事务:锁 + 航班变更 + 待发事件 + 终态 + 回填待办预登记 + → 提交后回填信箱;失败由补偿待办重试 ``` -- 原文缺失当前归为 `MALFORMED`;读取异常不能伪装成“缺失”,应进入基础设施重试。 -- 忽略规则覆盖约定的 `LDM / REGN / RSTA / EROR`,转 `SKIPPED` 并审计;不能产生业务副作用。 -- 航班变更(`FLIGHT_SCHD`)与待发事件、处理终态必须在同一 PG 事务原子提交。跨存储双写窗口已根除。 -- 所有终态都需要回填信箱,包括成功、忽略、重复和死信;非终态禁止回填。回填必须在 PG 提交后执行,并有持久化补偿、退避与告警;影子环境禁写。 +- 原文缺失归为 `MALFORMED`;读取异常不能伪装成“缺失”,应进入基础设施重试。 +- 忽略规则(`LDM / REGN / RSTA / EROR` → `SKIPPED`)尚未实现(§10);不能因类型未覆盖就把合法忽略报文当非法报文处理。 +- 航班变更、待发事件、处理终态与回填待办在同一 PG 事务原子提交;跨存储双写窗口已根除。 +- 终态回填发生在提交后:成功、业务重复与协议拒绝包持有可用 META 并回填;缺 META 或解码失败的死信无法回填,外部处理方式待确认(US-09/Q7)。`PENDING / FAILED` 禁止回填;影子环境禁写。 ## 4. 日计划快照与请求匹配 ### 4.1 快照发布 -`SCHD-RESP` 和 `SCHD-DNLD` 都进入 `ScheduleProcessor`: +`SCHD-DNLD` 与 `SCHD-RESP` 共用 `ScheduleProcessor.applyScheduleRecords`: -1. **暂存校验**:流式解析后完成整包校验和航班规范化;失败前不修改权威状态。暂存数据可在崩溃后从原文重建。 -2. **应答守卫**:`RESP` 必须匹配开放的 `RQFD` 请求;无匹配、已过期或报文时间早于发送时间时,转 `SKIPPED` 并审计,不更新快照。 -3. **SQL 原子写入**:在自有 PG 单事务内写入校验通过的 `FLIGHT_SCHD` 航班状态及资源明细;报文未携带的航班不因本次日计划报文被删除。 -4. **提交结果**:在同一 PG 事务中保存 `MSG_EVENT`、将消息置为 `SUCCEEDED`,并将匹配 `RESP` 的请求置为 `DONE`;事务提交后执行信箱回填。 +1. **重放判定**:`PROC_STATE` 已存在成功终态 → 幂等成功,仅追加留痕,不重复写入。 +2. **整包校验**:声明记录数、记录范围与运营日归属等校验失败 → 整包 `DEAD(PROTOCOL)`,不写半包,既有状态保持不变。 +3. **事务写入**:锁内按 `FLID` 点查归属日,发现同一航班跨运营日即整包回滚并 `DEAD(PROTOCOL)`;通过后合并写主表与资源明细。报文未携带的航班不因本次日计划报文被删除。 +4. **提交结果**:同一事务保存 `KAFKA:schd` / `KAFKA:msg` 事件、置消息 `SUCCEEDED` 并预登记回填待办;事务提交后执行信箱回填,留痕在事务外追加。 单事务保证未提交变更整体回滚。消息重放由 `PROC_STATE` 的消息 ID 与业务身份控制;版本号不能单独证明消息身份。 +**应答守卫仍是缺口**:`RESP` 应匹配开放 `RQFD` 请求(无匹配、过期或报文早于发送时间则不更新快照);当前 RESP 与 DNLD 无差别进入快照写入(§4.2/§10),不能视为 RESP 匹配闭环。 + ### 4.2 上游请求与静态数据 -`RequestCoordinator` 管理请求生命周期: +`REQ_TRACK` 表与仓储已存在(状态 `PENDING / SENT / DONE / EXPIRED`),但没有运行时协调器:出站 `COUTMSGS` 适配、请求编码、超时与应答匹配均未实现(`XmlCodec.encodeRqrd` 只是占位)。本节是目标机制,不是现状。 + +请求生命周期目标: ```text REGISTERED → SENT → WAITING → DONE └──→ EXPIRED ``` -注册同类新请求前使旧开放请求过期。只有 `COUTMSGS` 写入确认后才标记 `SENT` 并关联出站记录;写信箱成功但本地未确认的情况需要补偿与去重,不能无条件重新发送。 +- 注册同类新请求前使旧开放请求过期;只有 `COUTMSGS` 写入确认后才标记 `SENT` 并关联出站记录;写信箱成功但本地未确认的情况需要补偿与去重,不能无条件重新发送。 +- 应答优先按已确认的回显字段精确匹配;回显契约未确认时的降级匹配(同类开放请求且 `DTTM ≥ sentAt`)存在跨代误配风险,必须明确接受并审计,不能宣称精确关联。比较前统一时区和时间单位。 +- 参考应答写入自有 `REF_MASTER`(尚未建表),日计划应答走快照流程;请求完成必须在相应数据处理成功之后,超时和迟到应答不能修改已关闭请求对应的状态。 -应答优先按已确认的回显字段精确匹配。回显契约未确认时,按同类开放请求和 `DTTM ≥ sentAt` 判断的降级方式存在跨代误配风险,必须明确接受并审计,不能宣称精确关联。比较前统一时区和时间单位。 - -静态应答写入自有 PG 的 `REF_MASTER`,日计划应答走快照流程。请求完成必须在相应数据处理成功之后;超时和迟到应答不能修改已关闭请求对应的状态。 +请求与参考数据的交付范围见 user-stories.md US-08/US-13/US-14 与 §10。 ## 5. 事件投递 ### 5.1 普通事件 -`Dispatcher` 按 `TARGET` 读取最小未发送 `EVENT_ID`。队头退避未到期时,该目标停止推进;发送确认后才标记 `SENT`,失败记录次数和下次执行时间。所有外部调用需要有界超时,避免阻塞整个投递线程。 +`Dispatcher` 按 `TARGET` 读取最小未发送 `EVENT_ID`(当前逐条投递只处理 `KAFKA:msg`)。队头退避未到期时,该目标停止推进;发送确认后才标记 `SENT`,失败记录次数并按退避推后,达到上限转 `DEAD`(记录保留作 DLQ)。所有外部调用需要有界超时,避免阻塞整个投递线程。 -投递是至少一次:下游已接收但本地未标记成功时可能重发。Kafka 生产约束沿用架构决策 D11,但生产者幂等不替代应用层事件去重;跨重启的事件身份和下游去重契约仍需落实。共享出站信箱也必须单独解决重复写入,不能假设 Kafka 的保证适用于 MySQL。 +投递是至少一次:下游已接收但本地未标记成功时可能重发。Kafka 生产约束沿用架构决策 D3(architecture.md §7),但生产者幂等不替代应用层事件去重;跨重启的事件身份和下游去重契约仍需落实。共享出站信箱也必须单独解决重复写入,不能假设 Kafka 的保证适用于 MySQL。真实 Kafka 适配器尚未实现(`DeliveryPort` 仅有 stub),投递闭环须先交付适配与配置强制校验。 ### 5.2 `schd` 聚合 -`KAFKA_SCHD` 不进入逐条投递循环,只由 `flushSchd` 发送: +`KAFKA:schd` 不进入逐条投递循环,只由 `flushSchd` 发送: -1. 队头可执行后,按 `EVENT_ID` 顺序领取有界批次。 -2. 按 `FLID` 分组,保留批次内最大 `EVENT_ID` 对应的状态,组成 FLTR JSON 数组发送。 -3. 成功后将本批被代表的事件一起标记完成,推进 `lastFlush`;失败则整批增加次数并退避,达到上限整批转 `DEAD`。 +1. 到期领取批次:`mergePendingSchd` 按 `FLID` 合并未发事件,每个 `FLID` 只保留最新 `STATE_VERSION`,批次大小受 `flush-limit` 约束。 +2. 逐条发送:UPSERT 发送该 `FLID` 的最新整态(key = `FLID`);TOMBSTONE 发送 null 值删除通知。 +3. 成功后把本批被代表的事件(含被最新版本合并压掉的旧事件)一起标记完成,推进 `lastFlush`;失败的单条增加次数并退避,达到上限转 `DEAD`。 -默认聚合周期 3 秒、批上限 500。它提供最新状态通知,不保留每次中间变化;批次不能绕过尚在退避的队头。 +默认聚合周期 3 秒、批上限 500。它提供最新状态通知,不保留每次中间变化;`KAFKA:msg` 与 `KAFKA:schd` 之间不承诺顺序。 ## 6. 失败恢复与维护作业 ### 6.1 失败、重试与重放 -`ProcFailure` 和 `FailureScheduler` 统一处理侧的失败落账;投递侧按单条或聚合批次执行同样的次数与退避规则。默认最多 5 次,退避档位为 1、2、4、8、16 秒,单档封顶 60 秒。 +`ProcFailure` 与 `FailureScheduler` 统一处理侧失败落账,投递侧(`Dispatcher`)按同一套次数与退避规则迁移事件。默认最多 5 次(attempts ≥ 5 判耗尽),退避档位 1、2、4、8、16 秒、单档封顶 60 秒;时间经可注入 `Clock` 判定。 -失败必须在持有具体消息或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不捕获 JVM `Error` 作为普通业务失败。所有时间判断通过注入的 `Clock` 完成。 +失败必须在持有具体消息、事件或批次的位置记录,外层循环只做兜底日志和等待,不重复增加次数。线程中断应恢复中断标记并向上传递;不把 JVM `Error` 当普通业务失败捕获。 -`ReplayService` 只允许 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED` 从 `FAILED / DEAD` 回到 `PENDING`,重置次数和下次执行时间,保留身份与错误审计。旧消息进入终态后,后续消息可能已经执行;因此**重新入队不等于恢复历史顺序**,人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。 +`ReplayService` 只允许 `CODEC_ERROR / UNSUPPORTED / INFRA / EXHAUSTED` 从 `FAILED / DEAD` 回到 `PENDING`,重置次数与下次执行时间,保留身份与错误审计。它按错误类整批重放,尚无按记录预检、操作审计与管理入口(US-10)。旧消息进入终态后后续消息可能已执行,**重新入队不等于恢复历史顺序**;人工重放前必须评估状态覆盖和版本保护,不能直接批量重放到生产。 -### 6.2 归档与阶段 B 作业 +**维护作业**:`JobRunner` 用独立 daemon 线程每 30 秒触发 `BackfillSweepJob`(回填补偿扫描,指数退避 30 秒起步、封顶 15 分钟),每天机场时区 03:30 后触发一次 `HistorySweepJob`(§6.2)。作业不再经 `PUMP_JOB` 队列插队,不参与消息 FIFO,也不使到期消息饥饿。 -`ARCHIVE` 仅归档自有库的终态记录,默认按接收时间保留 1 天,保留期可配置为 1~7 天。归档表、接收时间依据和关联事件处理尚需落地;不得归档未完成记录,也不能因移走身份记录而意外失去业务去重能力。共享信箱保留策略由库所有方管理。 +### 6.2 历史清理与归档 -`HISTORY_SWEEP` 与 `PROJECTION_REBUILD` 暂缓。未来清场必须只删除已确认成功写入历史存储的集合;历史写入未接通时默认删除零条。历史写入与删除事件入队之间仍需恢复方案,顺序调用不构成原子提交。 +**航班历史清理**(`HistorySweepJob`,每天 03:30 触发):按 `HistoryProps` 的保留期与终态/静默判据选出候选(含 `DELETED`),先写历史存储,成功后物理删除主行与明细;历史存储未接通或 `history-store-enabled=false` 时删除 0 条。未经 FDEL、由生命周期直接清除的航班,清除前补发一次删除事件。语义与红线见 flight-state.md §6,不在此重复。 + +**留痕清理**:`SCHD_SNAP_LOG` 保留 90 天,在历史清理窗口内按 `(SCOPE_END, RECV_AT)` 删除。 + +**处理终态归档**:`PROC_STATE_HST` 仍是目标表(user-stories.md US-11),尚未建表;不得归档 `PENDING / FAILED`,也不能因移走身份记录而失去业务去重能力。ES 历史投影(阶段 B)不启用。 + +共享信箱保留策略由库所有方管理;历史写入与删除事件入队之间仍需恢复方案,顺序调用不构成原子提交。 ## 7. 接口与运行配置 @@ -179,11 +181,11 @@ REGISTERED → SENT → WAITING → DONE 运行配置以 `application.yml`、`application-dev.yml` 和 `.env.example` 为准,设计上重点区分: -- `pipeline.autostart` 与 `msgx.stubs`:分别控制管道启动和内存适配器;生产禁止 stub,默认不自动启动。 -- `pipeline.poll-interval / max-attempts / head-deadline`:控制轮询、重试上限和队头滞留;不能改变 FIFO。 +- `pipeline.autostart` 与 `msgx.stubs`:分别控制管道启动与内存适配器;生产禁止 stub,默认不自动启动。 +- `pipeline.poll-interval / claim-batch / max-attempts / backoff / head-deadline`:控制轮询节奏、批次、重试上限、退避与队头滞留;这些参数不能改变 FIFO。 - `schd.flush-period / flush-limit`:控制状态通知的聚合延迟与批量大小。 - `identity.include-day-boundary`:影响去重语义,不能作为普通调优项切换。 -- `phase`:阶段 A 是当前范围,不应把切为 B 当成已具备历史投影能力。 +- `mailbox.shared-mysql.enabled` 与 `history.history-store-enabled`:分别门控真实信箱与历史存储接线,默认关闭。 日志关联消息 ID、事件 ID 和批次;失败记录错误分类、次数、下次执行时间。健康检查反映依赖实际可用性;队头滞留、积压、死信和补偿失败需要指标及告警。日志出口故障不得阻塞业务线程。 @@ -197,7 +199,8 @@ REGISTERED → SENT → WAITING → DONE | 队头失败、退避及作业竞争 | 消息不越队;到期后恢复;作业不使消息无限饥饿。 | | 同身份多条记录、失败后重试、归档后重复 | 只产生一次有效业务处理,不把自身重试判为重复。 | | PG 事务失败、快照重复或迟到 | 整体回滚重试、不重复推进版本、不回退状态、不误删增量航班。 | -| PG 提交失败、信箱回填失败 | 事件与处理结果一起回滚;已提交结果只补偿回填。 | +| 整包协议拒绝(运营日冲突、声明数不符) | 整包不落地、整体回滚,既有状态与版本不变,消息终态为 `DEAD(PROTOCOL)`。 | +| PG 提交失败、信箱回填失败 | 事件与处理结果一起回滚;已提交结果只由回填待办补偿,不重放业务。 | | 投递确认丢失、批次失败、次数耗尽 | 允许可识别的重发、保持目标顺序、整批退避并保留死信。 | | 请求超时、无匹配 RESP、时间单位不一致 | 不误用迟到应答,不提前完成请求。 | | stub 误配置、重复实例、停机中断 | 生产拒绝不安全启动,工作线程能正确退出。 | @@ -210,22 +213,23 @@ REGISTERED → SENT → WAITING → DONE | 关注点 | 主要入口 | |---|---| -| 收报与兼容接口 | `ingress/InboxPoller.kt`、`InboxService.kt`、`InboxController.kt` | -| 调度与处理 | `processing/Pump.kt`(含 `MessageProcessor`)、`Handler.kt`、`Identity.kt` | -| 日计划与请求 | `processing/ScheduleProcessor.kt`、`reference/RequestCoordinator.kt` | -| 投递与作业 | `delivery/Dispatcher.kt`、`SchdAggregation.kt`、`jobs/JobExecutor.kt` | -| 持久化与恢复 | `infra/persistence/`、`infra/retry/`、`jobs/BackfillSweepJob.kt` | -| 启停与配置 | `PipelineLifecycle.kt`、`config/PipelineProps.kt` | +| 收报与兼容接口 | `ingress/InboxPoller.kt`、`InboxEnqueue.kt`、`InboxService.kt`、`InboxController.kt` | +| 解码 | `codec/JacksonXmlCodec.kt`、`SisWireMapper.kt`、`SisMessageBody.kt` | +| 调度与处理 | `processing/Pump.kt`(含 `MessageProcessor`)、`DynamicProcessors.kt`(FLOP/FDEL/ADFT)、`Identity.kt` | +| 日计划 | `processing/ScheduleProcessor.kt`;请求协调尚无实现(`REQ_TRACK` 见 `infra/persistence/`) | +| 投递与作业 | `delivery/Dispatcher.kt`、`jobs/JobRunner.kt`、`BackfillSweepJob.kt`、`HistorySweepJob.kt` | +| 持久化与恢复 | `infra/persistence/`、`infra/retry/`(`ProcFailure` / `ReplayService` / `FailureScheduler`) | +| 启停与配置 | `PipelineLifecycle.kt`、`config/PipelineProps.kt`、`config/HistoryProps.kt` | ## 10. 当前实现差异 以下缺口直接影响上述设计是否成立,不能以类或接口已存在作为完成依据: -- **事务与外部副作用**:状态、outbox 与终态同 PG 事务已实现;提交后尝试回填,失败落 BACKFILL_TODO。提交后崩溃及待办落账再次失败仍有恢复缺口,不能认定补偿闭环。 -- **收报与调度**:轮询仍从 `afterId=0` 扫描,持久水位与补扫策略未完成;作业以“队头为 FAILED”近似窗口,未区分是否已到期。主泵直取系统时间,滞留判据使用更新时刻,不能保证设计要求的超时升级。 -- **快照与业务能力**:快照事务内核已有测试,但生产 staging parser 未接通;Handler 注册检查在 DNLD 路由之前,生产只有 GTDT Handler。RESP 匹配和其他 Handler 待完成;10000 条保护不等于整包协议校验完成。 -- **航班读写**:唯一写入口和明细权威读已落地;ROUT/ERUT 主键冲突、空值/未知属性保真、事件重复计算和逐航班多次查询仍需修正。`/all/flights` 尚未实现。 -- **请求、静态数据与归档**:请求在真实出站前就标记发送,时间匹配与审计仍需修正;静态数据存在进程内过渡实现,归档表与关联保留策略尚未落地。 -- **生产与运维**:已有事务行锁,但没有整个实例的排他/FIFO 协调;启动校验、影子隔离、Kafka 强制配置、告警与指标尚未闭环。对拍比较工具存在不等于现场对拍已完成。Oracle 尚需完整适配。 +- **事务与外部副作用**:状态、事件、终态与回填待办同 PG 事务已实现,回填失败落 `BACKFILL_TODO` 并由 `BackfillSweepJob` 到期重试。剩余缺口在“提交后回填前崩溃”与“待办落账再次失败”两个窗口,不能认定补偿最终必达。 +- **收报与调度**:轮询仍从 `afterId=0` 每轮全量扫描,持久水位与受控补扫未完成;较小 ID 迟提交与空洞场景的顺序保证未验证。滞留判据使用 `updatedAt` 与直取系统时间,未基于稳定起始时刻。 +- **快照与业务能力**:DNLD/RESP/ADFT 与 FLOP/FDEL 处理器、整包校验与跨运营日整包拒绝均已接入;但 RESP 应答守卫与出站请求未实现,忽略规则(US-04)未实现,29 类 FLOP 与参考应答的逐类矩阵未补全,ADFT 缺失字段与 `FLID` 重用语义待上游确认。 +- **航班读写**:唯一写入口与权威读已落地;ROUT/ERUT 联合主键、空值/未知属性保真、事件在事务内只算一次、逐航班多次查询仍待修正(见 flight-state.md §6)。`/all/flights` 尚未实现。 +- **请求、参考数据与归档**:`REQ_TRACK` 表与仓储已建,但无运行时协调与 `COUTMSGS` 出站适配;`REF_MASTER` 未建表;`PROC_STATE_HST` 未建表。航班历史清理脚手架已实现,历史存储未接通时删 0 条。 +- **生产与运维**:已有事务行锁,但没有整个实例的排他/FIFO 协调;真实 Kafka 与信箱出站适配未交付(仅 stub),配置允许环境变量覆盖 `acks`/幂等/in-flight;启动校验、影子隔离、告警与指标未闭环。对拍比较工具存在不等于现场对拍已完成。Oracle 11g 方言与迁移未接入验证。 -进度与验收项见 Plane ACM2-10 实施计划(U01–U30),业务契约与待确认事项见 [user-stories.md](user-stories.md)。本文件不维护工单流水账、测试数量或历史方案全文。 +业务契约与待确认事项见 [user-stories.md](user-stories.md),进度由 Plane 跟踪;本文件不维护工单流水账、测试数量或历史方案全文。 diff --git a/docs/user-stories.md b/docs/user-stories.md index 2e66157..1506ba6 100644 --- a/docs/user-stories.md +++ b/docs/user-stories.md @@ -2,7 +2,7 @@ ## 1. 如何使用本文 -本文是把现有脚手架补成可用系统的实施入口:**故事定义要交付什么,验收标准定义怎样证明完成,代码落点说明从哪里改起**。保留 US-01~US-15、OPS-1~OPS-4 编号,便于关联已有任务和测试。 +本文定义阶段 A 的实施范围与验收口径:**故事定义要交付什么,验收标准定义怎样证明完成,代码落点说明从哪里改起**。保留 US-01~US-15、OPS-1~OPS-4 编号,便于关联已有任务和测试。 - 系统边界见 [architecture.md](architecture.md),模块流程见 [design.md](design.md)。本文不重复设计全文,也不以工单状态代替代码验收。 - “当前基础”来自本轮代码核对,只表示有接口或部分实现,不表示故事完成。`KEEP` 是保留业务兼容,`FIX` 是明确修正旧缺陷,`DEFERRED` 不进入阶段 A。 @@ -26,7 +26,7 @@ | S5:运维恢复 | US-10、US-11,完成 OPS-1~OPS-3 | 安全重放、归档后去重、故障告警、单写者保护、影子禁写均有验证证据。 | | S6:上线验证 | OPS-4,复核所有进入切流范围的故事 | 对拍、故障演练、配置和恢复 Runbook 验收后切流;US-15 不作为门槛。 | -**依赖口径**:US-03 是基础管道,不依赖具体业务 Handler;US-09 依赖其终态提交接口,US-04 复用 US-09。US-08 的请求登记与匹配基础不依赖 US-06;US-06 消费该基础,二者共同完成 RESP 集成验收,不形成开发依赖环。US-05 只有 PSDT 子范围依赖 US-14,不应阻塞其余 Handler。 +**依赖口径**:US-03 是基础管道,不依赖具体业务处理器;US-09 依赖其终态提交接口,US-04 复用 US-09。US-08 的请求登记与匹配基础不依赖 US-06;US-06 消费该基础,二者共同完成 RESP 集成验收,不形成开发依赖环。US-05 只有 PSDT 子范围依赖 US-14,不应阻塞其余处理器。 ## 3. 阶段 A 用户故事 @@ -42,7 +42,7 @@ 4. PG 不可用或批次中途失败时不改信箱标记;恢复后补建遗漏,记录失败次数与扫描进度。 5. 较小 ID 迟提交、ID 有空洞、兼容入口先入队较大 ID 时,必须遵守经 Q2 确认的发现与顺序协议;不能用“最终会重扫”冒充严格 FIFO。 -**当前基础与落点**:`ingress/InboxPoller.kt`、`InboxEnqueue.kt`、`infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt` 已有轮询和判重;固定 `afterId=0`,水位与补偿未实现。扩展 `InboxPollerTest`,补真实 PG/MySQL 中断恢复测试。 +**当前基础与落点**:`ingress/InboxPoller.kt`、`InboxEnqueue.kt`、`infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt` 已有轮询与判重;固定 `afterId=0` 全量扫描,水位与补扫策略未实现(每轮顺带对账回填待办,见 US-09)。扩展 `InboxPollerTest`,补真实 PG/MySQL 中断恢复测试。 **前置**:共享库读契约;Q2 决定严格顺序的端到端验收。 @@ -72,7 +72,7 @@ **目标**:报文失败和重试不造成航班状态倒序,也不重复产生副作用。 -**实施拆分**:调度与时钟 → 安全解码及路由 → 身份绑定 → 状态应用与 PG 提交。先用假 Handler 验证管道,不等 US-05 全部实现。 +**实施拆分**:调度与时钟 → 安全解码及路由 → 身份绑定 → 状态应用与 PG 提交。先用假处理器验证管道,不等 US-05 全部实现。 **验收标准** @@ -80,12 +80,12 @@ 2. 安全解码 XML,至少覆盖 META、SCHD、FLOP、参考应答与忽略类路由;合法但能力未支持是 `UNSUPPORTED`,不能一律归为非法报文。保留原文以支持诊断和回放。 3. 解码后首次绑定 `SNDR|TYPE|STYP|SEQN`;冲突转 `SKIPPED` 并记录原 ID;自身重试保留绑定。生产 `include-day-boundary=false`,更改算法须另行评审上游序号规则。 4. `MALFORMED` 直接 `DEAD`;`CODEC_ERROR / UNSUPPORTED / INFRA` 按次数和退避处理,耗尽转 `DEAD(EXHAUSTED)`。不能无限重试未实现类型,也不能立即当非法报文丢弃。 -5. HOL deadline 使用稳定的 `PROC_STATE.CREATED_AT`,不使用每次重试刷新的 `updatedAt`;所有调度判断注入 `Clock`。默认 5 次重试、10 分钟滞留限制;积压与人工重放的 deadline 边界按 Q6 验证。 +5. HOL deadline 必须基于稳定起始时刻(必要时在 `PROC_STATE` 增加入队时间列),不能用每次重试刷新的 `updatedAt` 代替;调度判断注入 `Clock`。默认 5 次重试、10 分钟滞留限制;积压与人工重放的 deadline 边界按 Q6 验证。 6. 主泵在同一 PG 事务提交航班主表/明细、事件与处理结果;终态回填意图通过 US-09 同事务保存。任一步失败整体回滚;提交后只重试外部回填,不重复生成业务事件。 -7. Handler 只返回 `Decision`,副作用由管道执行;失败只在持有消息上下文的边界落账,中断向上传递,不作为普通失败吞掉。 +7. 处理器只读取当前完整态、计算下一态,并把状态写入、事件与终态提交收敛在同一事务边界内,不直接触碰 Kafka;失败只在持有消息上下文的边界落账,中断向上传递,不作为普通失败吞掉。 8. 权威存储不可用或未完成恢复时停止业务处理;不能把“整个状态丢失”误判为“单航班不存在”而批量成功结束增量报文。 -**当前基础与落点**:`processing/Pump.kt`、`Identity.kt`、`codec/JacksonXmlCodec.kt`、`infra/retry/`、`JdbcPgRepositories.kt`。已有 PG 共享事务、身份/重试及 GTDT XML 链路;完整 SCHD 解码/路由、时钟与 deadline、同事务回填意图仍未完成。 +**当前基础与落点**:`processing/Pump.kt`(含 `MessageProcessor`)、`DynamicProcessors.kt`、`Identity.kt`、`codec/JacksonXmlCodec.kt`、`infra/retry/`。严格 FIFO 主泵、SCHD(DNLD/RESP/ADFT)/FLOP/FDEL 处理器、PG 单事务(含回填待办预登记)、身份绑定与重试已实现。尚未完成:忽略规则分支(US-04)、基于稳定起始时刻的滞留判据与可注入时钟、以及逐类矩阵与 golden 样例(US-05)。 **前置**:US-01;Q1 已定单库方向,Q6 决定 deadline 边界。数据库迁移只落自有库。 @@ -95,11 +95,11 @@ **验收标准** -1. 解码 META 后、身份绑定及业务 Handler 查找前,大小写不敏感匹配 `TYPE-STYP` 或 `TYPE-*`;基线为 `LDM-* / REGN-* / RSTA-* / EROR-*`,不混用 `ERROR`。 +1. 解码 META 后、身份绑定及处理器分派前,大小写不敏感匹配 `TYPE-STYP` 或 `TYPE-*`;基线为 `LDM-* / REGN-* / RSTA-* / EROR-*`,不混用 `ERROR`。 2. 命中后转 `SKIPPED`,记录 `ignored:` 和计数;不更新航班、不创建业务通知。 3. 通过 US-09 保存回填意图;命中、未命中、大小写和重扫均有测试。合法忽略报文不应因 `MsgKind` 尚不能表达它而先解码失败。 -**当前基础与落点**:`MessageProcessor` 尚无忽略分支,`DecodedMessage/MsgKind` 主要覆盖 SCHD/FLOP;扩展解码与路由,不把忽略逻辑散落到 Handler。 +**当前基础与落点**:`MessageProcessor` 尚无忽略分支;`MsgKind` 已有 SCHD/FLOP/FDEL/Unsupported 分派。在解码后的路由边界补忽略匹配,不把规则散落到各处理器。 **前置**:US-03 解码/终态接口、US-09。 @@ -109,14 +109,14 @@ **验收标准** -1. `SCHD-ADFT` 与 29 个 FLOP 子类型逐项列入覆盖矩阵,每项有纯函数 Handler;未知类型可恢复失败。RESP/DNLD 不计入这批 Handler,走 US-06。 +1. `SCHD-ADFT` 与 29 个 FLOP 子类型逐项列入覆盖矩阵,每项有对应的处理器规则与回归测试;未知类型可恢复失败。RESP/DNLD 不计入这批处理器,走 US-06。 2. 每类固定“输入与前态 → 后态 → msg → schd → 终态”五面样例;区分字段缺失、显式清空、重复报文和主/共享航班。清单和 golden 样例按 Q8 补齐,不以“已写 29 个类”替代验收。 -3. 对按 KEEP 规则需忽略的不存在航班,以 `SUCCEEDED` 无副作用结束,并由 US-09 回填;ADFT 建航班等行为按各类型矩阵执行。在 ACM2-28 运营航班落自有 PG 后,Redis 全损导致报文持续终结与白名单无法找回的 F1 损坏路径已被根除,PG 作为唯一权威重启即恢复。 +3. 对按 KEEP 规则需忽略的不存在航班,以 `SUCCEEDED` 无副作用结束,并由 US-09 回填;ADFT 建航班等行为按各类型矩阵执行。航班当前态以自有 PG 为唯一权威,重启即恢复,不存在 Redis 全损后白名单无法找回的损坏路径。 4. 共享航班通常更新并通知主航班,不直接发共享通知。FDEL 删除共享航班时更新主航班 MAFL 并通知;删除主航班时删除主航班及其子共享关联并发删除通知;目标不存在幂等成功。 5. ADFT/FDEL 使用值相等比较;主/共享关系一次原子变更,不出现主已删、子残留等半状态。 -6. PSDT 通过 US-14 的只读映射计算 `abdg`,Handler 不直接调用 admin-api。 +6. PSDT 通过 US-14 的只读映射计算 `abdg`,处理器不直接调用 admin-api。 -**当前基础与落点**:`processing/Handler.kt`、`domain/Decision.kt`、`infra/persistence/FlightSchdRepository.kt` 为主要落点;变更直接入自有 PG 事务 2。按类型新增 Handler 与同包测试。 +**当前基础与落点**:落点已变为 `processing/DynamicProcessors.kt`(`FlopProcessor`/`FdelProcessor`/`AdftProcessor`)与 `domain/flight/FlightStateEngine`、`FlightStateRepository`。FDEL/ADFT 与通用 FLOP 处理器已接入(含 DELETED 幂等与重激活、tombstone 同事务登记);29 类逐类语义矩阵与 golden 样例仍未补全,不能因处理器存在就视为覆盖完成。 **前置**:US-03;PSDT 另依赖 US-14;Q1、Q8。 @@ -128,7 +128,7 @@ |---|---|---|---| | `SCHD-DNLD` | ScheduleProcessor | 不更新请求 | 不要求开放请求 | | `SCHD-RESP` | 匹配守卫后进入 ScheduleProcessor | 成功提交时匹配 RQFD → DONE | SKIPPED、审计,禁止更新快照 | -| `SCHD-ADFT` | US-05 增量 Handler | 不更新请求 | 不适用 | +| `SCHD-ADFT` | US-05 增量处理器 | 不更新请求 | 不适用 | **验收标准** @@ -137,7 +137,7 @@ 3. 在自有 PG 单事务内,批处理写入已校验的 `FLIGHT_SCHD` 航班状态与资源明细;本次日计划中未出现的航班不因此被删除。 4. 在同一 PG 事务中提交 `FLIGHT_SCHD` 变更、`MSG_EVENT` 待发通知与 `PROC_STATE(SUCCEEDED)`;匹配 RESP 同事务完成请求并置 `DONE`;事务提交后执行信箱回填。 5. 相同报文重放不二次写入或重复发事件;单事务崩溃整体回滚,重放幂等。 -**当前基础与落点**:主处理只分流 DNLD;需补 RESP 路由、`REQ_TRACK` 应答关联字段和事务测试。 +**当前基础与落点**:DNLD 与 RESP 已共同路由到 `processing/ScheduleProcessor.applyScheduleRecords`,整包校验、归属日冲突整包拒绝与单事务写入已实现;`REQ_TRACK` 表与 `ReqTrackRepository` 已建。仍需补 RESP 应答守卫(开放 RQFD 匹配、时间比对)与请求完成关联逻辑。 **前置**:US-03、US-08 请求登记/匹配基础;Q1、Q5。 @@ -150,11 +150,11 @@ 1. `KAFKA:msg` 按目标内 `EVENT_ID` 顺序发送,确认后才标 `SENT`;队头退避时不跳过,发送有超时上限。 2. `KAFKA:schd` 只通过 `flushSchd` 聚合,默认 3 秒/500 条;同一 FLID 取批内最新状态,成功确认覆盖对应原事件,失败保持批次可恢复并退避,耗尽可见为 `DEAD`。 3. 外部接收成功、本地确认失败或进程重启后允许重发;事件标识跨重发稳定,消费者有去重约定,不宣称端到端恰好一次。 -4. msg 的目标分区键为 `SNDR`。schd 的聚合粒度与分区键先按 Q4 定案,明确哪些顺序只在单分区成立;不能把多航班数组声称为“每个 FLID 都是该 Kafka record 的 key”。 +4. 当前 `KAFKA:msg` 与 `KAFKA:schd` 的分区键均为 `FLID`,schd 逐 `FLID` 发送最新状态,不再是 legacy 的多航班数组。`msg` 是否需按 `SNDR` 分区、发送粒度与去重标识的放置以 Q4 定案为准;定案前不宣称单分区之外的顺序保证。 5. 生产强制 `acks=all`、`enable.idempotence=true`、`max.in.flight.requests.per.connection=1`;Broker 支持幂等生产协议并完成实际验证,不允许非幂等降级通过验收。 6. 普通/聚合发送失败、确认丢失、批次标记中断和目标阻塞均有测试;DEAD 保留记录并告警。 -**当前基础与落点**:`delivery/Dispatcher.kt`、`SchdAggregation.kt` 已有聚合和重试测试;`DeliveryPort.sendKafka` 仅接 topic/payload,无 key 或事件标识参数,真实生产适配与强制配置校验需补齐。 +**当前基础与落点**:`delivery/Dispatcher.kt` 已实现逐条 `KAFKA:msg` 与 `flushSchd` 聚合(按 `FLID` 合并最新 `STATE_VERSION`、TOMBSTONE 发 null)以及退避/DEAD 迁移;`DeliveryPort` 接口含 topic/key/事件类型参数。生产 Kafka 适配器未交付(仅 stub),强制配置校验与真实 Broker 验证需补齐。 **前置**:US-03 事件提交;Q4、现网 Broker 验证。wire 不兼容的标识字段不能直接加到现役载荷。 @@ -173,7 +173,7 @@ 5. 迟到、无匹配或已关闭请求的应答不得更新数据,转 `SKIPPED` 并审计。参考应答成功写入 REF_MASTER 后,与请求完成、处理终态和事件在 PG 边界内保持所需原子性。 6. `POST /schd/sync` 复用请求入口,采用 24 小时制和非空/区间校验;响应明确已登记还是已落信,不承诺计划已更新。 -**当前基础与落点**:`reference/RequestCoordinator.kt`、`ReqTrackRepository` 和 JDBC 表已存在;当前真实出站未实现就标发送,回显匹配和超时未闭环。需补 COUTMSGS 适配器、编码、并发约束、应答路由与故障测试。 +**当前基础与落点**:`ReqTrackRepository` 与 JDBC 实现、`REQ_TRACK` 表已存在;运行时协调器、`COUTMSGS` 出站适配与请求编码未实现(`XmlCodec.encodeRqrd` 仅占位),超时与应答匹配未闭环。需补协调器、出站适配、并发约束、应答路由与故障测试。 **前置**:US-01、US-03;Q5、Q8、出站信箱去重契约。请求基础不依赖 US-06。 @@ -189,7 +189,7 @@ 4. 重复补偿效果幂等,保留稳定的完成时间与审计;重放后的新处理结果不能被旧回填任务覆盖。非法报文缺 META 时也有明确回填方式。 5. 影子模式禁写,双跑仅一个系统持有标记写权;暴露 PG 终态、回填状态、积压、最老年龄与持续失败告警。 -**当前基础与落点**:`backfillOnSuccess` 当前仅同步写 `PROCESSED`,无持久意图。补自有 PG 意图模型/迁移、事务提交入口和独立补偿执行器;不在共享库新增补偿表。 +**当前基础与落点**:回填意图已在业务事务内预登记到自有 `BACKFILL_TODO`;提交后同步回填(写 `DATE_PROCESSED` 与 `PROCESSED`),失败由 `BackfillSweepJob` 到期重试(每 30 秒、指数退避)。剩余缺口:缺 META 或解码失败的死信回填方式,以及“待办落账后再失败 / 提交后崩溃”两个窗口的补偿闭环。 **前置**:US-03 终态接口;Q7、共享库更新权限。测试覆盖四类终态、事务回滚、重复补偿和重放竞争。 @@ -205,7 +205,7 @@ 4. DEAD 之后可能已有新状态,必须预检版本与覆盖风险;不安全时拒绝直接重放,改用经批准的隔离重建或恢复流程,禁止无保护的全量 `replayAll` 生产入口。 5. 操作有认证、授权、范围限制与审计;死信、持续补偿失败、队列年龄越界有告警和处理 Runbook。 -**当前基础与落点**:`infra/retry/ReplayService.kt` 仅按错误类批量重排并返回数量;需补按记录选择、预检、操作审计和管理入口,扩展 `ReplayServiceTest`。 +**当前基础与落点**:`infra/retry/ReplayService.kt` 已按错误类批量把 FAILED/DEAD 置回 `PENDING`(重置次数与下次执行时间,保留身份与错误审计)并返回数量;需补按记录选择、版本/覆盖预检、操作审计和管理入口,扩展 `ReplayServiceTest`。 **前置**:US-03、US-09 的恢复状态;Q6、OPS-1/OPS-2 的安全与可观测基础。 @@ -220,7 +220,7 @@ 3. 迁移与删除在自有库事务内完成,重复执行幂等;失败保留源记录并报告计数。归档后同信箱 ID/业务身份再次到达,仍能按约定去重。 4. 不写共享 MySQL `CMINMSGS_HST`,不清理外部信箱;原文可用性与重放保留期由 Q7/Q8 关联确认。 -**当前基础与落点**:`jobs/JobExecutor.kt`、`PumpJobRepository` 为入口;当前迁移无归档表。先确定去重记录保留与关联策略,再补迁移及归档中断测试。 +**当前基础与落点**:`PROC_STATE_HST` 未建表,也没有归档处理记录的作业;先确定去重记录保留与关联策略,再补迁移与归档中断测试。航班历史清理(`HistorySweepJob`,属 US-15 红线范围)与本文档处理记录归档不是同一件事,不能混为一谈。 **前置**:US-03、US-09;US-07 提供事件终态规则,US-10 提供恢复保留要求。不依赖 US-15。 @@ -230,11 +230,11 @@ **验收标准** -1. 保留 `GET /all/flights`,直接从自有 PostgreSQL `FLIGHT_SCHD` 查询(ACM2-28 定案:读源由 Redis 转为 PG),过滤 `MAID != NULL` 的共享航班;不改写业务状态。 +1. 保留 `GET /all/flights`,直接从自有 PostgreSQL `FLIGHT_SCHD` 查询,过滤 `MAID != NULL` 的共享航班;不改写业务状态。 2. 固定响应样例、空结果、排序、大小限制及一致性时点。现役未分页时不能无声改为只返回第一页;分页或响应结构变更按 Q3 决定。 3. 依赖异常不能伪装为空数组成功;影子只读影子状态,入口有约定的访问控制、限流与审计。 -**当前基础与落点**:`InboxController.kt` 中仅有该接口 TODO;新增查询控制器与只读服务,复用权威读取端口,不能从空占位仓储返回成功。新增接口/状态不可用测试。 +**当前基础与落点**:`InboxController.kt` 仅有该接口的 TODO(阶段 2);权威读取端口为 `FlightStateRepository` 的全量快照读取。新增查询控制器与只读服务,不能从空占位仓储返回成功,并补接口/状态不可用测试。 **前置**:Q1、Q3;所查询的 US-05/US-06 状态发布能力。 @@ -249,9 +249,9 @@ 3. 一类失败不破坏其他类或该类旧版本;同类由 AODB/admin-api 都提供时明确覆盖优先级,全量刷新时明确已删除项的处理,不能仅靠 SOURCE 日志解决冲突。 4. 影子默认不主动刷新生产数据;需要参考样本时显式导入隔离副本。 -**当前基础与落点**:`reference/ReferenceService.kt` 是空刷新入口,`JdbcStaticRefRepository` 有逐条 upsert;补客户端、类型清单、批次发布/事务和失败保旧测试。 +**当前基础与落点**:`reference/` 包与 `REF_MASTER` 表均未落地(无迁移、无客户端、无刷新服务)。按 Q8 清单从零补建:先固定 21 类端点与字段契约,再补客户端、批次发布/事务与失败保旧测试。 -**前置**:admin-api 访问契约、Q8。可独立于消息 Handler 开发。 +**前置**:admin-api 访问契约、Q8。可独立于消息处理器开发。 ### US-14 提供机位与登机桥映射 @@ -262,9 +262,9 @@ 1. 保留 `ORMS_STAND / ORMS_STAND_AIRBRIDGE` 两类,与 US-13 的 21 类分开统计;适配器拉取、完整校验后原子发布只读缓存。 2. 近机位产生登机桥值,远机位或清空机位时 `abdg` 为空;一机位多桥、缺失映射与共享航班规则用 golden 固定。 3. admin-api 不可用时使用最后可用版本;无可用版本或映射不完整时明确失败,不用空映射冒充正常清空,也不发布半批。 -4. Handler 输入包含所需只读参考视图,不允许其直接 HTTP 或写缓存。 +4. 处理器输入包含所需只读参考视图,不允许其直接 HTTP 或写缓存。 -**当前基础与落点**:在 `reference/` 与 `infra/` 增加映射服务/适配器,必要时扩展 Handler 输入上下文;PSDT 测试使用固定映射,无需在线 admin-api。 +**当前基础与落点**:在 `reference/` 与 `infra/` 增加映射服务/适配器,必要时扩展处理器输入上下文;PSDT 测试使用固定映射,无需在线 admin-api。 **前置**:机位/桥数据契约及 Q8;不要求 US-13 全部完成。 @@ -272,9 +272,9 @@ ### US-15 历史航班清场(DEFERRED,阶段 B) -历史存储确认成功后,才允许主泵删除对应实时航班;逐条隔离坏数据,不能删除写历史失败的集合。五条判史规则、业务时区 `Asia/Shanghai`、历史写入与删除事件之间的恢复协议需在启用前完成 golden 对拍。 +历史存储确认成功后,才允许删除对应实时航班;逐条隔离坏数据,不能删除写历史失败的集合。判史规则与保留期(`HistoryProps`)、业务时区 `Asia/Shanghai`、历史写入与删除事件之间的恢复协议需在启用前完成 golden 对拍。 -阶段 A 不依赖 ES,不接通 `HISTORY_SWEEP / PROJECTION_REBUILD` 占位能力;没有历史写入能力时默认删除零条。不要因代码已有 B 枚举就启用它。 +阶段 A 不依赖 ES,不启用 `PROJECTION_REBUILD`。`HistorySweepJob` 已作为每日清理脚手架接入,但历史存储未接通(或 `history-store-enabled=false`)时删除 0 条;脚手架存在不等于清场已交付,红线见 flight-state.md §6。 ## 5. 运行与切流验收 @@ -306,7 +306,7 @@ | Q1 权威存储(方向已定) | 当前 PG 单库权威,主表 + 无损明细;现场供库目标为 Oracle 11g。 | 单库决策已采纳,Oracle 完整适配与部署验收仍待交付;不得退回 Redis 双写。 | | Q2 入队顺序 | 水位+补扫无法自动保证较小 ID 迟提交不越序;当前有限批扫描也可能被未回填记录挡住。 | 库方 ID/提交顺序约束,或明确的发现完整性与暂停/恢复协议;晚提交、空洞、兼容入口与重扫联合测试。不能凭空假定 ID 连续。 | | Q3 HTTP 契约 | 目标 ResponseDto 与现有 text/plain ID 不同;请求媒体类型目标已列出,错误码、状态码、查询格式等仍需对拍。 | 每个保留接口的真实请求/响应样例、错误表和契约测试;10MB 的字节口径、字符集及兼容变更说明一起固定。 | -| Q4 Kafka wire | schd 当前多 FLID 聚为一个数组/一条 record,与“record 按 FLID 分区”冲突;msg 的 SNDR 尚未接线。 | 下游确认发送粒度、key、去重标识放置、分区内顺序及批次确认策略;同步 Dispatcher/设计,不擅自把现役数组改成逐航班消息。 | +| Q4 Kafka wire | 当前 msg/schd 均按 `FLID` 逐航班发送(schd 由 `flushSchd` 合并为每个 FLID 的最新状态),与 legacy“多 FLID 数组”形态不同;msg 是否需按 `SNDR` 分区尚未接线。 | 下游确认发送粒度、key、去重标识放置、分区内顺序及批次确认策略;未定案前不改动现役消费契约。 | | Q5 请求匹配 | 目标优先 SEQN 回显,但回显是否可靠需确认;DTTM 降级存在跨代误匹配,尤其旧应答到达新请求期间。 | 15 类请求/响应样例、回显字段与时间格式;降级风险是否接受及拒绝条件。默认超时仍为 RQFD 60 秒、RQRD 30 秒。 | | Q6 deadline 与重放 | 入队时间作为稳定锚点会把长期排队消息计入滞留;历史重放保留 CREATED_AT 后可能立即过期。 | 明确首次处理/排队过期策略与“本次恢复尝试”计时方式,保留原始时间审计;测试积压恢复和旧 DEAD 重放,不用刷新 UPDATED_AT 绕过超时。 | | Q7 信箱外部契约 | 原目标为成功→SUCCESS、忽略→SKIPPED、重复→DUPLICATE、死信→DEAD;当前 JDBC 写 PROCESSED。新 STATUS 值尚不能假定库方支持。 | 库方认可的状态值、META/缺失字段、权限、原文保留期、出站去重与处理时间语义;若只允许 legacy 集合,显式映射内部原因,不新增外部枚举。 |