diff --git a/AGENTS.md b/AGENTS.md index a496d70..128321d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -13,6 +13,8 @@ - 核心约束锚点:`specification.md` (C/PRE/INV/CLM/Q/G), `architecture.md` (D1–D4), `implementation.md` (机制与航班域), `reference.md` (PARAM)。 - **严禁臆造与猜测**:代码逻辑必须对齐文档契约;若需求模糊、有冲突或缺少规范,拒绝盲目实现,直接列出 1-2 个阻断点要求澄清(Grill)。 - **文档引用纪律**:交叉引用仅使用稳定 ID(如 `PRE-x`, `C-x`, `PARAM:`)或「文件名 + 小节名」指针,禁止使用章节号;参数值和默认值不重复书写,统一指向原处;文档绝不记录进度(进度走 Plane)。 +- **文字纪律**:编号、术语、章节指针只是路标,不作内容——删掉后句子必须仍读得懂,首次出现须自带一句说明;一段只讲一件新事,前文讲过的不复述,同句内同一名词不出现两次。 +- **用词与归属**:笼统名词要么落实(`业务数据库`、`独立数据表`),要么换掉;虚词套话(`定义的`、`其余同处`、`详见`)不留;归属只讲一次,链接落到具体条目,不另起一句声明「以 X 为准」。 # 项目核心约束 (com.gzzn.omms.msgexchange) diff --git a/docs/README.md b/docs/README.md index b133f4a..69080ea 100644 --- a/docs/README.md +++ b/docs/README.md @@ -48,7 +48,7 @@ - `C-x`、`INV-x` 用加粗定义行(`- **C-5** …`); - `CLM-n`、`OPS-n`、`Dn`、`PRE-n`、`Qn`、`G-NAME` 与 `PARAM:` 用注册表首列;首列必须是**单个裸 ID** (`` `ID` `` 或 `ID`)。成组登记(`` `a` / `b` ``)、带括注的首列与写成 `` `ID` `` 的引用行都不算定义。 - - 同一行登记多个 ID(如验证映射的 `INV-20 / CLM-3`)是引用行,不构成定义。 + - 同一行登记多个 ID(如验证映射的 `INV-20b / CLM-3`)是引用行,不构成定义。 其他位置一律是引用。 - 编号稳定:条款被取代时标 `[作废 by C-y]` 并保留原文;不静默改写,不重编号。 - **引用只用稳定 ID**,不用章节号:写 `INV-3`、`C-8`、`PARAM:msgx.pipeline.claim-batch`,或「见 implementation.md『收报』」这类文件名 + 小节名指针。章节号随增删章节腐烂,指针失效后必然被改写为复述。 diff --git a/docs/architecture.md b/docs/architecture.md index 151c647..8189240 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1,101 +1,126 @@ # 架构 +本文写架构约束与归属,不代表能力已实现(可声明性见 [specification.md](specification.md));上线与切流的操作步骤另立规程。历史报文契约以 [SIS 接口规范](legacy/SIS_AODB_RMS-V0.1.md) 和 [XSD](legacy/unisysaodbsis.xsd) 为兼容依据。 + ## 1. 系统定位与范围 -msgexchange-v2 是机场 OMMS 的上游报文处理中间件,用于替换旧版 `msgexchange-api`。 -它读取 CIIMS、AODB 等系统写入共享 MySQL 信箱的 XML 报文,按顺序更新航班动态与静态参考数据,再将结果提供给下游。静态参考数据由 SIS 消息获得,是本版本的主要新增能力之一。 +msgexchange-v2 是 OMMS H5 查询系统的消息网关,替换旧版 `msgexchange-api`,收取 CIIMS adapter 信箱中 AODB 下发的 XML 报文:航班动态写入业务数据库、写 Redis 投影供页面查询,并经 Kafka 通知运营航班显示界面;静态参考数据(机场、航空公司、机型、机位等 13 类基础数据与资源状态)写入独立数据表。 -本系统负责**收报、解析、状态更新和结果投递**,不生成上游业务报文;admin-api 是处理结果的下游数据库读取方,本网关不调用 admin-api。其余非目标见 [requirements.md](requirements.md)「范围与非目标」。 +功能需求是 [requirements.md](requirements.md) 的十四条用户故事(采集、处理、投递、查询、维护,`US-01`~`US-14`),运行需求是四条验收(单实例、可观测、测试隔离、切换回退,`OPS-1`~`OPS-4`)。不生成航班/业务数据类报文,不替代 CIIMS/AODB,不调用 admin-api;完整非目标见同文件「范围与非目标」。 -- **主要入口**:轮询共享 MySQL 的 `CMINMSGS`。 -- **兼容入口**:`POST /cminmsgs/send`,供现役兼容、手工工具和对拍使用;写入信箱后返回记录 ID,不是生产收报主路径。 -- **输出**:Kafka 的 `msg` / `schd` 消息、共享 MySQL 的 `COUTMSGS` 出站信箱、供 admin-api 只读的业务数据库结果,以及查询 HTTP 接口;不直接推送前端。 -- **当前范围(阶段 A)**:航班当前态落自有 PostgreSQL(`FLIGHT_SCHD` + 资源明细表 + `FLIGHT_ROUTE_POINT`,权威口径见 [implementation.md](implementation.md)「航班域」)。无 Redis 依赖;ES 历史投影属阶段 B。 +- **主要入口**:轮询共享 MySQL 入站表 `CMINMSGS`,只取处理标记为空的行,按编号升序、每批有上限(`US-01`;`C-30`)。 +- **兼容入口**:`POST /cminmsgs/send` 供联调工具把报文写进信箱,与上游投递走同一条处理路径;返回的编号只表示已进信箱,不代表已处理或下游已收到(`US-02`)。 +- **查询入口**:`GET /all/flights` 返回当前全部动态航班(不含共享航班),读 Redis,与网页客户端同源(`US-12`;`INV-24`)。 +- **出站**:只向 AODB 发参考数据和日计划两类请求,经共享 MySQL 出站表 `COUTMSGS`,由 CIIMS adapter 消费;只保证请求写入信箱,不保证 AODB 收到(`US-09`;`C-24`)。 +- **输出**:Kafka 主题 `msg` 发单条航班变更、`schd` 定时发批量最新状态;Redis 存航班投影;静态参考数据表供 admin-api 只读(`US-08`、`US-13`)。 +- **权威**:航班当前态的权威是自有 PostgreSQL(`FLIGHT_SCHD`、资源明细表、`FLIGHT_ROUTE_POINT`;字段含义见 [implementation.md](implementation.md)「航班域」);信箱、Redis、Kafka、展示视图都不是(`INV-11b`)。 +- **航班历史**:已结束航班先写入 Elasticsearch 历史库,成功后才从实时数据删除(`US-14`;`D1`)。 -现场供库时目标为 Oracle 11g,否则自建 PostgreSQL;Oracle 适配必须通过方言与集成验证后才能作为运行时选项。 +测试环境用 PostgreSQL;生产环境用 PostgreSQL 还是 Oracle 11g 尚未确定,Oracle 适配验证通过前不作支持承诺(生产库选型见 `Q1`)。 -本文只描述架构约束与归属,不代表能力已实现:机制见 [implementation.md](implementation.md),前提、不变量与可声明性见 [specification.md](specification.md),对外契约见同文件「契约」,需求见 [requirements.md](requirements.md),参数与代码入口见 [reference.md](reference.md);操作步骤在上线/切流前另立规程(设计阶段只保留前置条件与红线)。历史报文契约以 [SIS 接口规范](legacy/SIS_AODB_RMS-V0.1.md) 和 [XSD](legacy/unisysaodbsis.xsd) 为兼容依据。 +监控指标与告警口径见 [reference.md](reference.md)「指标与健康」。 ## 2. 总体架构 ```text -CIIMS / AODB 等上游 - │ 写入 XML +CIIMS adapter(把 AODB 下发的 XML 写入信箱) + │ 落信 ▼ -共享 MySQL:CMINMSGS - │ 轮询未处理记录 +共享 MySQL:CMINMSGS(入站信箱) + │ ▼ -┌──────────────── msgexchange-v2(单实例)────────────────┐ -│ ingress:发现报文 → PostgreSQL 持久化入队 │ -│ │ │ -│ processing:取 FIFO 队头 → 解析 / 去重 → 处理器决策 │ -│ └─ PG 单事务:领域变更 + 终态 + 待发事件│ -│ │ -│ jobs:独立维护线程(回填补偿 / 历史归档 / 留痕清理) │ -│ delivery:读取 PG 待发事件 → 投递 / 重试 │ -└─────────────────────────┬──────────────────────────────┘ - ├─ Kafka:msg / schd - ├─ 共享 MySQL:COUTMSGS - └─ 业务数据库 ──只读──▶ admin-api - -处理结果提交后,再回填 CMINMSGS 的处理标记;失败需补偿。 -查询接口与 admin-api 读取处理后的业务数据,不参与状态写入,也不向本网关提供参考数据。 +──────── msgexchange-v2(单活动实例,同进程四组线程) ──────── +ingress 收报 processing 处理(流程见第 4 节) +delivery 投递 jobs 作业:回填 / 历史清理 / 记录清理 +────────────────────────────────────────────── + │ + ├─▶ Kafka:msg / schd ── 运营航班显示界面 + ├─▶ Redis:航班查询投影 + ├─▶ 共享 MySQL:COUTMSGS ── 参考数据与日计划请求发 AODB + ├─▶ Elasticsearch:航班历史 + └─▶ 自有 PostgreSQL:静态参考数据表 ── admin-api 只读 ``` -收报、处理、投递与维护作业各使用独立线程,不占用 HTTP 事件循环。**航班当前态的写入只发生在持有 `PIPELINE_LOCK` 的事务内**,由主泵串行驱动。 +四组线程在同一进程、互不调用,协作只经自有 PG 的持久记录交接;HTTP 接口走事件循环,不占这四组线程。重启后各段从记录接着做,不依赖内存进度(`US-01` AC4、`US-03` AC3、`US-10` AC2)。 -技术栈:Kotlin + JDK 25、Micronaut 编译期依赖注入、JDBC 持久化;数据库变更由 Flyway 管理,只作用于自有 PostgreSQL。依赖版本以 `build.gradle.kts` 为准,不在本文重复维护。 +线程之间不加锁,靠幂等写入:收报按信箱编号只登记一次(`US-01` AC2)、回填只写空标记(`C-5`)、记录清理只删已回填且超过保留期的行(`INV-25`);撞上时后到的一次重试即可。唯一的例外是处理消息的主循环(主泵)与航班历史清理之间要加锁(`INV-18`)。 + +运行边界:同一时刻只允许一个实例处理消息(`OPS-1`);处理进度只有信箱处理标记一处,停旧启新时未处理的消息由旧系统继续(`OPS-4`);积压、处理失败、投递失败、回填失败各有指标与告警(`OPS-2`)。 + +技术栈:Kotlin + JDK 25、Micronaut 编译期依赖注入、JDBC 持久化;数据库变更由 Flyway 管理,只作用于自有 PostgreSQL。依赖版本见 `build.gradle.kts`。 ## 3. 模块职责 | 模块 | 职责与边界 | |---|---| -| `ingress` | 轮询信箱、持久化入队及兼容 HTTP 写入;不解析业务报文。 | -| `codec` | XML 解码,区分非法报文与可修复的解码失败。 | -| `processing` | FIFO 调度、业务身份绑定与去重、领域决策与落库(SCHD/FLOP/FDEL/ADFT/静态参考数据);决策、事务与回填的职责边界见 [implementation.md](implementation.md)「消息、身份与决策」与 `INV-17`。 | -| `delivery` | 消费待发事件,负责按目标保序、`schd` 聚合、投递和失败重试。 | -| `jobs` | 回填补偿扫描、航班历史清理与留痕保留期清理;独立 job 线程执行,不参与 FIFO(与主泵的互斥见 `INV-18`)。 | +| `ingress` | 轮询信箱、登记入队、兼容入口落信(`US-01`、`US-02`);不解析业务报文。 | +| `codec` | XML 解码,区分非法报文与可修复的解码失败(`US-03`)。 | +| `processing` | FIFO 调度、业务身份绑定与去重、领域决策与落库(`SCHD`/`FLOP`/`FDEL`/`ADFT`/静态参考数据),航班类写 Redis 投影;决策与回填机制见 [implementation.md](implementation.md)「消息、身份与决策」。 | +| `delivery` | 读待发事件投 Kafka:按 `FLID` 保序、`schd` 聚合、失败重试(`US-08`;`C-29`)。 | +| `jobs` | 回填扫描、航班历史清理,以及 `PROC_STATE`、`MSG_EVENT`、`SCHD_SNAP_LOG`、`REQ_TRACK` 的到期清理;单独线程、不进 FIFO(与主泵的互斥见 `INV-18`),保留期一律见 [reference.md](reference.md)。 | | `domain` / `config` | 领域状态、事件和决策模型,以及运行参数。 | -| `infra` | 仓储(JDBC/stub)、外部适配器、重试、健康检查与日志;通过接口隔离基础设施。 | +| `infra` | 仓储(JDBC/stub)、外部适配器(共享信箱、Kafka、Redis、航班历史存储、AODB 出站)、重试、健康检查与日志;对其他模块只暴露接口。 | 代码入口清单见 [reference.md](reference.md)「模块与代码入口」。 ## 4. 主流程 -`InboxPoller` 发现 → 自有 PG 入队(`INV-2`)→ 主泵按最小未完成 `MSG_ID` 取队头、解码并按业务身份去重 → 处理器在持 `PIPELINE_LOCK` 的同一事务内写航班变更、待发事件、处理终态与回填意图(`INV-17`)→ 提交后回填共享信箱处理标记 → `Dispatcher` 投递 `KAFKA:msg` 与 `KAFKA:schd`(至少一次,`INV-10`)。 +航班动态消息走满全链,其余类别只换其中几步: -机制细节各有归属:扫描谓词与水位、分派与事务边界、回填、保序与 `schd` 聚合见 [implementation.md](implementation.md)。 +1. **收报**:按「处理标记为空」发现信箱行,登记入队(`INV-2b`)。 +2. **主泵**:按最小未完成 `MSG_ID` 取队头,解码,按业务身份去重(`INV-3`、`INV-9`);非法或不支持的报文无副作用,终态留档(`US-03`)。 +3. **事务一**:持 `PIPELINE_LOCK`,领域变更与待发事件一起提交(`INV-17b`)。 +4. **投影**:写 Redis,写成功才算处理完成(`INV-23`);失败保持未完成、下轮重写投影,业务效果幂等(`US-03` AC3)。 +5. **事务二**:处理终态与回填意图一起提交。 +6. **回填**:作业把处理标记写回共享信箱(`US-10`);写不上的记录在案并告警。 +7. **投递**:读待发事件发 `KAFKA:msg` / `KAFKA:schd`,至少一次、`FLID` 内保序(`INV-10`、`C-29`);一直失败的记录保留可查并告警(`US-08`)。 + +其余类别与这条主干的差异: + +| 类别 | 与主干的差异 | 完成判据 | 失败时 | +|---|---|---|---| +| 日计划(`DNLD`/`RESP`) | 第 3 步改为每批一个事务,第 5 步在整包完成后 | 整包成功,含 Redis 刷新(`US-07` AC4/AC5) | 不标记已处理,下轮整包重处理(`US-07` AC4) | +| 静态参考数据 | 无第 4 步;一个事务完成落库、终态与回填意图 | 该类落库成功(`US-13` AC2) | 校验不过整类不动,其他类照常(`US-13` AC2) | +| 出站请求 | 不走收报队列:登记新请求并作废同类旧请求,写入 `COUTMSGS` | 请求已写入信箱(`US-09` AC1) | 写入失败下轮重试;交付承诺只到落信(`C-24`) | +| 航班历史清理(作业) | 不走消息队列:历史写入成功后物理删除 | 实时数据已删(`US-14` AC3) | 历史写不成功不删,下轮重来(`US-14` AC3) | + +出站的后半程:应答按报文类型匹配等待中的请求,发错或迟到的不更新数据、记录后跳过;超过时限未等到应答,标记超时;`EROR` 定位到本系统发出的请求,标记失败并告警(`US-09`)。 + +机制细节各有归属:扫描谓词、分派与事务边界、回填四结果、保序与 `schd` 聚合见 [implementation.md](implementation.md)。 ## 5. 必须保持的约束 -本节只列约束的**归属**;完整定义与验证映射见 [specification.md](specification.md),实现与演进不得违反: +以下约束不得违反;完整定义与验证方式见 [specification.md](specification.md): -- 消息严格 FIFO:`INV-3`、`INV-4`、`INV-5`(发现完整性依赖 `PRE-2`/`PRE-3`)。 -- 动态状态单写者与写者集合互斥:`D2`、`INV-18`。 -- 身份去重:`INV-9`。 -- 快照可恢复与运营日不可变:`INV-12`、`INV-13`。 -- 航班当前态的物理清除只发生在历史归档之后:`D1`。 +- 只跑一个实例,航班当前态只有一个写入方:信箱读取不加锁,主泵与历史清理互斥(`OPS-1`、`PRE-5`、`INV-18`)。 +- 按顺序处理、只处理一次:每次只取编号最小的未完成消息,重复扫描、失败重处理与兼容入口并发都只登记一次、生效一次(`INV-2b`、`INV-3`、`INV-9`);不丢消息依赖「编号即到达顺序」且编号不复用、不回退(`PRE-2`、`PRE-3`)。 +- Redis 写成功才算处理完成;查询接口与网页客户端读同一份 Redis,出问题时报错,不返回空列表假装正常(`INV-23`、`INV-24`;`US-05`、`US-06`、`US-12`)。 +- 日计划快照以 AODB 下发为准:快照里没有的航班删除,未携带的字段清除;增量报文不适用这条(`INV-14b`、`INV-15b`;`US-07`)。 +- 对外投递至少一次,同一航班(`FLID`)内保序,跨航班不承诺顺序(`INV-10`、`C-29`;`US-08`)。 +- 航班只在历史写入成功后删除,历史库还没接通时一条也不删(`D1`、`INV-28`;`US-14`)。 +- 静态参考数据一类校验失败只停这一类,其他类照常;空值是「当前没有值」,不是删除(`INV-26`、`INV-27`;`US-13`)。 +- 航班唯一,版本只进不退:`FLID` 唯一,已写入非空的运营日不可改;每次成功写入版本号加一,重复消息不重复加(`INV-12`、`INV-13`)。 -这些约束优先于吞吐量优化。单写者降低了并发复杂度,代价是队头阻塞和吞吐上限;如需并行化,必须先重新定义顺序与状态归属,不能只调整线程数。 +这些约束是拿速度换来的:一个写入方、一次只推进一条,前一条没处理完,后面都得等。要提速、要多实例,光加线程没有用——得先重新设计消息顺序和数据由谁写,多实例还得补上可靠的互斥保护。 ## 6. 数据归属与一致性 | 存储 | 承载内容 | 职责说明 | |---|---|---| -| 自有 PostgreSQL | 单行锁 `PIPELINE_LOCK`、处理状态与回填事实 `PROC_STATE`、消费水位 `INBOX_CURSOR`、待发事件 `MSG_EVENT`、请求跟踪 `REQ_TRACK`、航班当前态 `FLIGHT_SCHD` + 资源明细表 + `FLIGHT_ROUTE_POINT`、静态参考数据有效视图(逻辑模型 `REF_MASTER`)、留痕 `SCHD_SNAP_LOG` | 本系统唯一业务数据库。消息处理、状态推进、处理终态、回填意图与待发事件在单事务内原子提交;本地事务只在此库。记录级定义见 [implementation.md](implementation.md)「持久化记录」。 | -| 共享 MySQL | `CMINMSGS` 入站信箱、`COUTMSGS` 出站信箱 | 外部系统所有。本系统仅执行约定的信箱读写与处理标记回填,不建表、不迁移 schema、不写历史表;由库方按 `Q9` 执行的清除与历史归档见 [specification.md](specification.md)「契约」。 | +| 自有 PostgreSQL | 处理锁 `PIPELINE_LOCK`、消息处理状态与回填意图 `PROC_STATE`、待发事件 `MSG_EVENT`、出站请求跟踪 `REQ_TRACK`、航班当前态 `FLIGHT_SCHD`、资源明细表、`FLIGHT_ROUTE_POINT`、静态参考数据 `REF_MASTER`、日计划快照留痕 `SCHD_SNAP_LOG` | 本系统唯一的业务数据库,也是航班当前态的唯一权威(`INV-11b`);本地事务只发生在这里,事务怎么分段见第 4 节。记录级定义见 [implementation.md](implementation.md)「持久化记录」。 | +| Redis | 航班查询投影 | 只作查询,不是权威,也不存处理状态(`INV-11b`);只由本系统写入和移除(`INV-24`),内容来自 PG 当前态;`GET /all/flights` 与网页客户端读的就是它。 | +| 共享 MySQL | `CMINMSGS` 入站信箱、`COUTMSGS` 出站信箱 | 信箱归外部系统所有。本系统只读写消息、回写处理标记,不建表、不改表结构、不清数据、不写历史表(`C-14`);原文保留多久、何时清除由库方定(`C-5`~`C-12`、`Q7`、`Q9`)。出站请求写进去就算交付(`C-24`)。 | +| 航班历史存储(Elasticsearch) | 已结束航班的历史副本 | 已结束航班写入这里作历史副本;写入确认成功后才删实时数据,写不成一条也不删(`D1`、`INV-28`)。保留期与容量上限未定(`G-FLIGHT-HIST-RETENTION`)。 | -**不使用跨库事务。** PG 事务只能保证「处理结果与待发事件一起提交」(`INV-17`),不能覆盖 MySQL 回填或 Kafka 发送等外部副作用。跨存储依靠幂等、重试和持久化补偿恢复;各中断位置的判定与恢复动作见 [implementation.md](implementation.md)「中断恢复」。 +**PG 的事务只管自己库。** Redis 写没写成、信箱标记写没写上、Kafka 发没发出,PG 事务都管不着;这些步骤各自可重试,重做多少遍结果都一样,重启后从 PG 记录接着走。哪种中断该怎么接续,见 [implementation.md](implementation.md)「中断恢复」。 -对外投递的承诺边界见 `C-29`;Kafka 生产者幂等不能消除应用重启或 outbox 重发带来的所有重复。 +对外投递只承诺至少一次(`C-29`):应用重启、待发事件重发都可能让同一条消息多发一次,Kafka 的生产端幂等挡不住这种重复。 ## 7. 关键决策 -仅保留仍具约束价值、且无法从正文与 [implementation.md](implementation.md) 直接推出的决策,按 `D1`–`D4` 连续编号;其余曾编号条目(严格 FIFO、stub 门控、本地事务、UNSUPPORTED 处理等)已在正文以约束形式表达,不再重复列表。第三列只给证据与偏差指针;可声明性见 [specification.md](specification.md),决策不因状态变化而撤销。 +只列正文推不出来、仍有约束价值的决策,按 `D1`–`D2` 编号。第三列只给证据与偏差指针;可声明性见 [specification.md](specification.md)。决策不随实现状态增删。 | 编号 | 决策及理由 | 证据 / 偏差 | |---|---|---| -| D1 | 航班清场只在历史写入成功后进行,未接通时删 0 条;未经 FDEL 的清场须先补发删除事件。历史写入是需求内交付(`US-14`),红线见 `INV-28`。 | 红线见 [implementation.md](implementation.md)「生命周期」 | -| D2 | 动态状态单写者,生产只允许一个活动实例;多实例必须先具备可靠的排他保护。 | 事务行锁见 [implementation.md](implementation.md)「事务边界」;实例级排他属 `PRE-5`,可声明性见 `CLM-6` | -| D3 | Kafka 生产必须同时满足确认级别、幂等生产与单连接在途上限三项约束;不允许通过关闭幂等来满足生产接入。 | 取值见 [reference.md](reference.md) 参数表 | -| D4 | 自有库终态处理记录到期直接删除,不归档;不侵入共享库的表结构或保留策略。 | 清理口径见 [implementation.md](implementation.md)「生命周期与清除」与 `INV-25` | +| D1 | 删除实时数据前,下游必须已收到删除事件:FDEL 自带;历史清理先补发再删。历史写入是需求内交付(`US-14`)。 | 红线见 [implementation.md](implementation.md)「生命周期」与 `INV-28` | +| D2 | Kafka 生产端同时满足三项:确认级别、幂等、单连接在途条数上限;不许关幂等绕开这条限制。 | 取值见 [reference.md](reference.md) 参数表;投递口径见 `C-29`、`INV-10`(`US-08`) | diff --git a/docs/implementation.md b/docs/implementation.md index 39eff5c..f6aa2c9 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -2,7 +2,7 @@ 本文件是实现设计的唯一出处,分三章: -- **处理管道**章:记录模型、状态机、收报与水位、主泵与事务边界、回填、快照与请求、投递、失败恢复与维护作业; +- **处理管道**章:记录模型、状态机、收报与扫描谓词、主泵与事务边界、回填、快照与请求、投递、失败恢复与维护作业; - **航班域**章:航班当前态的权威模型、合并与写入语义、删除与重建; - **静态参考数据**章:SIS 消息中的主数据类别与编码、`REF_MASTER` 结构与合并语义、资源状态,以及 admin-api 的下游读取边界。 @@ -14,11 +14,10 @@ | 术语 | 语义 | |---|---| -| `W`(水位) | 信箱 ID 的连续上界:`(min, W]` 已全部读入自有 PG;只随新 ID 成功入队推进(永久空洞放行是唯一例外)。 | -| `holeSince` | `W+1` 处空洞最早被观测到的时刻;无空洞时为 NULL,跨重启保留。 | +| 扫描谓词 | 信箱读取条件「处理标记为空」(`INV-2b`);本系统不以 ID 区间或水位作为消费边界。 | | 队头 | 最小的未完成消息(`PENDING` 与 `FAILED` 都占位)。 | | 终态 | `SUCCEEDED` / `SKIPPED` / `DEAD`;到达后队列方可推进。 | -| 回填意图 | 「还欠一次信箱标记」的持久化事实,与终态同一条语句落库。 | +| 回填意图 | 「还欠一次信箱标记」的持久化事实,与终态同一条语句落库,且发生在 Redis 投影写成功之后(`INV-23`)。 | | `R`、`R_keep`、处理标记 | 定义见 [specification.md](specification.md)。 | ### 1.2 持久化记录 @@ -42,7 +41,7 @@ **身份绑定是独立的幂等单语句**(`WHERE IDENTITY_KEY IS NULL`),不参与业务事务。它的前提是「报文不可变」(`PRE-7`):同一身份的重发不会被比对内容,若上游改发正文会被判为重复并跳过(`Q15`)。 -分派与落库由 `MessageProcessor` 协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(SCHD-DNLD/RESP → `ScheduleProcessor`,ADFT → `AdftProcessor`,FLOP → `FlopProcessor`,FDEL → `FdelProcessor`,`SIS:3.1`~`SIS:3.14` 的静态参考数据消息 → `ReferenceDataProcessor`,其余 → `SKIPPED(unsupported)`)。这些处理器在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,再统一落库并登记回填意图;它们不直接触碰 Kafka。领域决策逻辑不执行 I/O。 +分派与落库由 `MessageProcessor` 协调:按 `MsgKind` 把已绑定身份的队头消息交给对应事务协调器(SCHD-DNLD/RESP → `ScheduleProcessor`,ADFT → `AdftProcessor`,FLOP → `FlopProcessor`,FDEL → `FdelProcessor`,`SIS:3.1`~`SIS:3.14` 的静态参考数据消息 → `ReferenceDataProcessor`,其余 → `SKIPPED(unsupported)`)。这些处理器在 `PIPELINE_LOCK` 事务内读取当前完整态,调用纯领域决策逻辑得到下一完整态与待发事件,提交该事务后写 Redis 投影,再在另一事务中登记处理终态与回填意图(`INV-17b`、`INV-23`);它们不直接触碰 Kafka。领域决策逻辑不执行 I/O。处理步骤的锁跨越 Redis 写,因此与航班历史清理的互斥覆盖整个步骤(`INV-18`)。 合法但本系统不支持的消息类型:跳过留档、按已处理写回标记(`US-03` AC2),不重试。`REGN` / `RSTA` 是静态参考数据消息,必须分派给 `US-13`,不得跳过。 @@ -122,7 +121,10 @@ processOne(head): Unsupported → SKIPPED(unsupported)(跳过留档,按已处理写回标记) 载荷缺失 → DEAD(MALFORMED) 整包协议拒绝 → DEAD(PROTOCOL),不落半包 - 5. 业务型成功:处理器在自己的事务内写航班变更 + 待发事件 + SUCCEEDED + 回填意图 + 5. 业务型成功,分三步(INV-17b、INV-23): + ① 事务提交:航班变更 + 待发事件 + ② 写 Redis 投影;失败 → 保持未完成,下轮重处理 + ③ 事务提交:SUCCEEDED + 回填意图 6. 结束:主泵不做回填;回填意图已随终态落库,由扫描补写信箱标记 ``` @@ -135,13 +137,15 @@ processOne(head): |---|---|---|---|---| | 收报入队(`insertIfAbsent`) | 是 | 否 | 否 | 同库事务 | | 身份首次绑定 | 否 | 否 | 否 | 单语句 + 唯一约束 | -| 业务型终态(处理器产出 `SUCCEEDED`) | 是 | 是 | 是 | 同库事务:航班变更 + 事件 + 终态 + 回填意图 | +| 业务型领域变更(航班变更 + 待发事件) | 是 | 是 | 是 | 同库事务(`INV-17b`) | +| Redis 投影写 | 否 | 是(处理步骤锁跨越本步) | 否 | 外部副作用,不在 PG 事务内;写成功是终态事务的前置(`INV-23`) | +| 业务型终态(`SUCCEEDED` + 回填意图) | 是 | 是 | 否 | 同库事务(`INV-17b`) | | 非业务型终态(`MALFORMED` / `PROTOCOL` / `SKIPPED` / `EXHAUSTED`) | 否 | 否 | 否 | 单语句(终态与回填意图同一条 UPDATE) | | 航班历史清理的物理删除 | 是 | 是(`INV-18`) | 是 | 同库事务:复查判据 + 历史写入成功后删除 | | 回填(信箱标记 + `BACKFILL_AT`) | 否 | 否 | 否 | 跨库两次单写;幂等可重跑 | | 人工重放(批量改回 `PENDING`) | 否 | 否 | 否 | 单语句批量;`MessageLifecycleGate` 与回填互斥 | -结论:「航班变更与处理终态同事务」只对业务型终态成立。`PIPELINE_LOCK` 的竞争写者是**航班历史清理**(`INV-18`),不是别的处理器线程;没有第二写者时该锁不产生额外串行度。 +结论:航班变更与处理终态**不在同一事务**——两者之间夹着 Redis 投影写;同一事务只保证「航班变更 + 事件」与「终态 + 回填意图」各自原子(`INV-17b`、`INV-23`)。处理步骤的 `PIPELINE_LOCK` 跨越 Redis 写,使 `INV-18` 的互斥覆盖整个步骤;该锁的竞争写者是**航班历史清理**,不是别的处理器线程;没有第二写者时该锁不产生额外串行度。 ### 5.4 历史积压 @@ -283,6 +287,7 @@ PENDING → SENT → DONE |---|---|---| | 已落信、未入队 | 信箱行处理标记为空且 PG 无记录 | 重扫补建登记记录 | | 事务执行中 | PG 无该消息终态 | 事务整体回滚,按 `PENDING` 重新处理 | +| 领域事务已提交、Redis 写失败或终态未提交 | 该消息无终态(`PENDING`),仍占队头 | 整条消息重处理:投影按当前完整态重写,领域变更依赖逐类幂等(`INV-20b`,`G-FLOP-IDEMPOTENT`),已提交结果不回滚(`INV-16`) | | 事务已提交、标记未写 | 终态行仍持有回填意图 | 仅补写标记;业务处理结果保持不变 | | 标记写入中途 | 标记仍为空 | 重新写入;重复写入同一值无副作用 | | 回填时信箱行已不存在 | 写入 0 行且信箱行不存在 | 立即放弃自动重试(`MISSING_ROW`)并告警;放弃不等于标记已确认,仍需人工对账 | @@ -402,7 +407,7 @@ PENDING → SENT → DONE - `ORDINAL` 是持久化顺序,从 1 开始;`SOURCE_SEQ` 是上游序号,允许为空或重复。 - 相同资源号不代表同一条分配,禁止按资源号去重。 -- 每次持久化完整航班状态时,明细表按该 `FLID` 先删后插,以完整合并结果为准(`INV-14`)。 +- 每次持久化完整航班状态时,明细表按该 `FLID` 先删后插,以完整合并结果为准(`INV-14b`)。 - ROUT 与 ERUT 是两类独立集合,不能因相同序号覆盖彼此。 - `CHDT` 的类字段固定为 `CCLS`/`CTYP`;当前 wire DTO 与持久化列误写成 `CHCLS`/`CHTYP`,见 `G-FLOP-UNMAPPED`。 - 主/共享关系以主表的 `MAID` 为事实来源:`MAID` 是共享航班指向主航班 `FLID` 的引用(非共享航班为 `NULL`);`MAFL` 只在读取和事件投影时从子航班事实派生,不按入站标量解析或保存。 @@ -416,7 +421,7 @@ PENDING → SENT → DONE - 同一 `STATE_VERSION` 的投影逐字节稳定(顺序见 `INV-21`),与到达顺序及 `FLNO` 变更无关;重发与消费端比对才有意义。 - `MAID = FLID` 的自引用行不进入任何 `MAFL`;`MAID` 指向不存在主航班的悬挂引用不阻断该子航班自身处理,只是不产生投影。 - 子航班集合变化的传播见 `INV-22`,事件类型为 `KAFKA:msg` + `KAFKA:schd`;否则整态投影的只进不退写入会丢弃它(见「`schd` 聚合」)。共享航班自身不单独发通知。 -- 派生主航班投影与产生它的状态写入必须同一事务或一致读快照;按 `MAID` 取子航班要求该列有索引(`INV-17`)。 +- 派生主航班投影与产生它的状态写入必须同一事务或一致读快照;按 `MAID` 取子航班要求该列有索引(`INV-17b`、`INV-22`)。 ## 12. 航班域:合并、删除与生命周期 @@ -432,7 +437,7 @@ SCHD DNLD/RESP 在整包校验通过后,分批将报文携带的航班写入 ### 12.2 动态运行事件(FLOP) -FLOP 只修改报文表达的字段或集合,其余状态保持不变;目标形态与合并规则见「字段与集合」。`STYP` 必须命中下表白名单,未知值按不支持类型跳过留档(`US-03` AC2),不得进入通用合并。已确认的动态更新在同一事务推进 `STATE_VERSION`、登记 `KAFKA:msg` 与 `KAFKA:schd`、提交处理终态与回填意图;同一消息重复处理不得重复产生业务效果(`INV-17b`、`INV-20b`)。 +FLOP 只修改报文表达的字段或集合,其余状态保持不变;目标形态与合并规则见「字段与集合」。`STYP` 必须命中下表白名单,未知值按不支持类型跳过留档(`US-03` AC2),不得进入通用合并。已确认的动态更新在同一事务推进 `STATE_VERSION` 并登记 `KAFKA:msg` 与 `KAFKA:schd`(`INV-17b`);处理终态与回填意图在 Redis 投影写成功后的另一事务中提交(`INV-23`)。同一消息重复处理不得重复产生业务效果(`INV-20b`)。 逐类语义以 SIS 的字段表、空标签规则与 Processing Exceptions 为准;下表每一行都必须有一条回归用例钉住「输入与前态 → 目标状态 → 终态与事件」。 diff --git a/docs/reference.md b/docs/reference.md index 1771735..bbc5eb7 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -12,19 +12,19 @@ | 参数 | 默认 | 单位 | 依据 | 说明 | |---|---|---|---|---| -| `msgx.pipeline.poll-interval` | `1s` | Duration | 现役 | 收报轮询节奏;决定队头与空洞检查频率 | +| `msgx.pipeline.poll-interval` | `1s` | Duration | 现役 | 收报轮询节奏;决定队头检查频率 | | `msgx.pipeline.claim-batch` | `50` | 条 | 假定 | 单轮领取上限;过大延长单轮事务 | | `msgx.pipeline.max-attempts` | `5` | 次 | 假定 | 处理与投递共用;达到即转 `DEAD(EXHAUSTED)` | | `msgx.pipeline.backoff-ms` | `[1000,2000,4000,8000]` | ms / 档 | 假定 | **档位数必须 = `max-attempts − 1`**,启动自检拦截错位 | | `msgx.pipeline.backoff-cap-ms` | `60000` | ms | 假定 | 单档封顶;默认表内无档触及 | -| `msgx.pipeline.max-commit-delay` | `5m` | Duration | **假定(无依据)** | 空洞老化阈值;由 `C-2` 决定,**不可由 SIS `Expiry` 推导**(`Q2`) | +| `msgx.pipeline.max-commit-delay` | `5m` | Duration | **假定(无依据)** | **待退役**:水位空洞老化阈值,随收报改谓词扫描(`G-SCAN-PREDICATE`)删除;依据 `C-2` 已作废 | | `msgx.pipeline.overdue-backfill` | `30d` | Duration | 契约(`R ≤ R_keep`) | 即 `R`:进入强补写窗口、**取消退避**的阈值;**不是兜底保证** | | `msgx.pipeline.backfill-batch` | `100` | 条 | 假定 | 回填扫描单批条数 | | `msgx.pipeline.backfill-scan-period` | `30s`(代码常量,无配置键) | Duration | 现役 | 回填扫描作业周期;批次积压与单行超时会延长实际标记延迟(`CLM-9`) | | `msgx.pipeline.backfill-max-attempts` | `100` | 次 | 假定 | 单行重试的**告警阈值**;放弃判据是 `R` 超期,不是次数 | | `msgx.pipeline.backfill-backoff-ms` | `30000` | ms | 假定 | 回填独立退避起步间隔;`BackfillService` 指数退避的首档 | | `msgx.pipeline.backfill-backoff-cap-ms` | `900000` | ms | 假定 | 回填退避封顶(15 分钟) | -| `msgx.pipeline.cutover-watermark` | 不设置 | `min\|zero\|max\|` | 一次性运维决策 | 显式播种水位;非法值由启动自检挡下;升级实例拒绝重新播种 | +| `msgx.pipeline.cutover-watermark` | 不设置 | `min\|zero\|max\|` | 一次性运维决策 | **待退役**:显式播种水位(`G-SCAN-PREDICATE`);非法值由启动自检挡下,升级实例拒绝重新播种 | | `msgx.pipeline.delivery-batch` | `200` | 条 | 假定 | `KAFKA:msg` 每轮每目标领取上限 | | `msgx.pipeline.delivery-drain-rounds` | `10` | 轮 | 假定 | 连取批数上限,让出循环跑 `schd` flush,防状态通知被积压饿死 | | `msgx.pipeline.autostart` | `false` | 布尔 | 安全默认 | 启动即拉起收报 / 主泵 / 投递循环;需真实仓储或 `msgx.stubs=true` | @@ -76,7 +76,7 @@ | `msgx.pipeline.backfill.unmarked_terminal` | 已终态但未打标的条数 | 不降 → 回填失败 | | `msgx.pipeline.backfill.abandoned` | 已放弃自动回填的条数 | **非 0 需人工对账** | | `msgx.pipeline.backfill.oldest_unmarked_seconds` | 最老待回填年龄 | 决定实际回填延迟 | -| `msgx.pipeline.watermark.lag` | 水位落后信箱最新 ID 的距离 | 增长 → 收报停滞 | +| `msgx.pipeline.watermark.lag` | 水位落后信箱最新 ID 的距离;**待退役**(`G-SCAN-PREDICATE`) | 增长 → 收报停滞 | | `msgx.pipeline.job.heartbeat_age_seconds` | 距上一次作业 tick 完成的秒数(未跑过为 -1) | 持续增长 → 作业线程卡死 | | `msgx.pipeline.job.last_failure_age_seconds` | 距最近一次作业 tick 失败的秒数(从未失败为 -1) | 配合 `failures.total` 增长判断扫描/历史作业异常 | | `msgx.pipeline.job.ticks.total` | 作业 tick 完成次数 | 不增长 → 作业停摆 | diff --git a/docs/requirements.md b/docs/requirements.md index 45e6509..27e6ea6 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -8,7 +8,7 @@ ## 1. 范围与非目标 -**系统定位**:OMMS H5 查询系统的消息网关。收取 CIIMS adapter 信箱中 AODB 下发的 XML 报文:航班动态写入数据库并同步写 Redis,运营航班的动态消息经 Kafka 发给运营航班显示界面实现同步;静态参考数据写入数据库,供 admin-api 只读。出站仅向 AODB 发参考数据类请求(经 `COUTMSGS`,消费方为 CIIMS adapter)。 +**系统定位**:OMMS H5 查询系统的消息网关。收取 CIIMS adapter 信箱中 AODB 下发的 XML 报文:航班动态写入数据库并同步写 Redis,运营航班的动态消息经 Kafka 发给运营航班显示界面实现同步;静态参考数据写入数据库,供 admin-api 只读。出站仅向 AODB 发参考数据类请求和日计划请求(经 `COUTMSGS`,消费方为 CIIMS adapter)。 **交付范围**:`US-01`~`US-14`、`OPS-1`~`OPS-4`。 diff --git a/docs/specification.md b/docs/specification.md index 175ad02..df2f758 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -115,7 +115,7 @@ - **INV-15b** 日计划快照以 AODB 下发为准:快照里没有的航班删除——标记已删除、登记删除事件、从 Redis 投影移除;快照里未携带的字段视为 AODB 已删除该值,本地同步清除(`US-07` AC2/AC3)。 - **INV-16** 外部副作用(回填、Kafka 投递、出站信箱)失败可重试,但不回滚已提交的本地业务结果。 - **INV-17** 状态变更、待发事件、处理终态与回填意图在同一 PG 事务内原子提交。**[作废 by INV-17b]** -- **INV-17b** 状态变更与待发事件在同一 PG 事务内提交;处理终态与回填意图在同一 PG 事务内提交,且该事务在 Redis 投影写成功之后(`INV-23`)。日计划例外:分批写入,每批一个事务(状态变更 + 待发事件),终态与回填意图在整包完成后同一事务提交;失败不标记已处理,下轮整包重新处理(`US-07` AC4)。 +- **INV-17b** 状态变更与待发事件在同一 PG 事务内提交;处理终态与回填意图在同一 PG 事务内提交,且该事务在 Redis 投影写成功之后(`INV-23`)。日计划例外:分批写入,每批一个事务(状态变更 + 待发事件),终态与回填意图在整包完成后同一事务提交;失败不标记已处理,下轮整包重新处理(`US-07` AC4)。静态参考数据不产生待发事件与投影写,落库、终态与回填意图在同一事务提交。 - **INV-18** 航班表的写者集合是「主泵处理器」与「历史清理」;两者必须互斥(同一 `PIPELINE_LOCK`,或清理在同一事务内复查判据后再删除),不得出现清理删除与处理器更新同一 `FLID` 的竞态。 - **INV-19** 整包校验失败或运营日冲突时整包不落地,既有状态与版本保持不变。 - **INV-20** 处理器幂等:同一消息重复执行只产生一次业务效果。身份唯一只防「重复记录」,不防「重新执行」;SIS 25 类与经 `Q8` 定案启用的 legacy 子类型完成逐类幂等矩阵前,本条**不可声明**(`G-FLOP-IDEMPOTENT`)。**[作废 by INV-20b]**