diff --git a/docs/contracts/interface-contract.md b/docs/contracts/interface-contract.md index 0cc3b95..7d1d354 100644 --- a/docs/contracts/interface-contract.md +++ b/docs/contracts/interface-contract.md @@ -11,7 +11,8 @@ | 接口 | 请求契约 | 成功响应契约 | 失败契约 | 尚需确定 | |---|---|---|---|---| | `POST /cminmsgs/send` | 请求体是 XML 原文;接受 `text/xml`、`application/xml`、`text/plain`,默认 UTF-8;仅限内网,网络层限制来源。 | 报文写入 `CMINMSGS` 后返回信箱编号;写入的报文与上游投递走同一条处理路径、效果一致;该响应只证明已落信,不证明业务处理或下游投递(`US-02`)。 | 空报文、超大小上限、非法 XML 不落信并返回错误;写信失败不返回编号;XML 解析禁用外部实体和外部资源访问。 | 大小上限、请求编码与 `Content-Type` 的精确处理规则、HTTP 状态码、成功/失败响应体字段及样例(`Q15`)。 | -| `POST /schd/sync` | 触发一次 `RQFD` 日计划请求;请求体为空、请求全量(`C-4`);登记与落信规则见「在途与作废」。 | HTTP 200,返回请求编号;只表示已登记,不代表已落信、更不代表 AODB 已收到。 | 已有未结案请求时 HTTP 409,响应体 `open-request-exists`。 | — | +| `POST /schd/sync` | 触发一次 `RQFD` 日计划请求;请求体是网页选定的时间条件,写入规则见 `C-4`;登记与落信规则见「在途与作废」。 | HTTP 200,返回请求编号;只表示已登记,不代表已落信、更不代表 AODB 已收到。 | 已有未结案请求时 HTTP 409,响应体 `open-request-exists`。 | — | +| `POST /refdata/sync` | 触发一次 `RQRD` 参考数据请求;请求体携带网页选定的类别码 `STYP`;`STYP=RSTA` 时须带资源类型 `RTYP`,其余类别不得带;`STYP` 取值见 `C-4`,`RTYP` 取值见 implementation.md「静态参考数据」;登记与落信规则见「在途与作废」。 | HTTP 200,返回请求编号;只表示已登记,不代表已落信、更不代表 AODB 已收到。 | 非法 `STYP` 时 HTTP 400,响应体 `unknown-styp`;`RSTA` 缺 `RTYP` 时 HTTP 400,响应体 `missing-rtyp`;非 `RSTA` 带 `RTYP` 时 HTTP 400,响应体 `unexpected-rtyp`;已有未结案请求时 HTTP 409,响应体 `open-request-exists`。 | — | | `GET /all/flights` | 无已定义的请求字段;从 Redis 投影读取当前全部动态航班,不含共享航班,与网页客户端同源(`US-12`)。 | HTTP 200;响应体是裸 JSON 数组(不套旧 `ResponseDto`),元素为日计划 `SCHD.FLTR` 转成的 JSON,与 `KAFKA:schd` 数组元素同形(`C-9`、`C-11`);共享航班(`MAID` 非空)不出现在数组里;不分页。 | HTTP 503;JSON 对象 `{"error":"FLIGHT_PROJECTION_UNAVAILABLE","reason":"<细节>"}`;Redis 或投影读失败时不得返回 200 空数组(`INV-11`)。 | — | 旧系统线索(来源:旧项目用户故事「HTTP 接口清单」「日计划请求」): @@ -21,7 +22,7 @@ - 旧 `GET /all/flights` 用 `ResponseDto`,`body` 是非共享航班的 `SCHD.FLTR` 列表(对象即日计划 XML 解码后的 `FLTR`,再序列化为 JSON);新版成功体直接返回该列表对应的裸 JSON 数组(`C-11`),元素形状不变。 - `POST /schd/sync` 的请求体是 `{startDate, endDate}`,时间格式 `yyyy-MM-dd hh:mm` 用无 AM/PM 的 12 小时制,解析结果写进 `RQFD` 的 `STDB`/`STDE`(`ddMMMyyHHmm`,大写)。 -旧系统的 `{startDate, endDate}` 与 `RQFD` 的 `STDB`/`STDE` 均不沿用:新版日计划是 AODB 当前时刻的完整航班列表(`US-07`),出站编码已定(`C-4`),不带日期筛选;`POST /schd/sync` 请求与响应口径见 `Q16`。旧系统响应体是否能作为 `POST /cminmsgs/send` 的定稿样例,见 `Q15`。 +旧系统只把 `startDate`、`endDate` 写成 `STDB`、`STDE`,时间格式是无 AM/PM 的 12 小时制,新版不沿用这套入参。新版由网页选择时间条件,四个筛选都可传,格式与含义见 `C-4`;`POST /schd/sync` 请求与响应口径见 `Q16`。旧系统响应体是否能作为 `POST /cminmsgs/send` 的定稿样例,见 `Q15`。 ## 入站报文与请求应答 @@ -40,7 +41,7 @@ AODB 经 CIIMS adapter 把 XML 报文写入 `CMINMSGS`,格式以架构指定 | 主题 | 已确定的消息语义 | 尚需确定 | |---|---|---| -| `msg` | 单条航班变更通知;Kafka value 是整条 `MSG` 的 JSON(`META` 加对应业务体,空字段不输出),变更与删除由 `META` 的类型与子类型区分(`Q5`);航班动态与删除处理完成后投递;发送失败自动重试,一直失败的记录保留可查并告警(`US-08`);同一 `FLID` 的变更保序,对外按至少一次投递(`CLM-3`,单分区)。删除通知的来源有三处:`FDEL` 删除(`US-06`)、日计划快照覆盖范围内缺席删除(`US-07`)、历史清理在物理删除前必要时登记(架构 `D1`)。 | 编码方式(JSON 之外的压缩或封装是否引入)。 | +| `msg` | 单条航班变更通知;Kafka value 是整条 `MSG` 的 JSON(`META` 加对应业务体,空字段不输出),变更与删除由 `META` 的类型与子类型区分(`Q5`);航班动态与删除处理完成后投递;发送失败自动重试,一直失败的记录保留可查并告警(`US-08`);同一 `FLID` 的变更保序,对外按至少一次投递(`CLM-3`,单分区)。删除通知的来源有三处:`FDEL` 删除(`US-06`)、日计划完整名单覆盖范围内缺席删除(`US-07`)、历史清理在物理删除前必要时登记(架构 `D1`)。 | 编码方式(JSON 之外的压缩或封装是否引入)。 | | `schd` | 定时批量发送最新航班状态;两次 tick 之间积累的航班组成 `SCHD.FLTR` 数组 JSON,整批作为单条 record 发出,条数不设上限(沿用旧系统,`Q1`);空字段不输出,字段与类型见 [XSD](../legacy/unisysaodbsis.xsd) 的 `FLTR`;没有变化不发;删除航班不进入本主题,由 `msg` 发一条删除通知;发送失败自动重试,对外按至少一次投递(`US-08`)。 | — | 旧系统线索(来源:旧项目用户故事「前端通知」「动态类(FLOP-*)处理」): @@ -80,7 +81,7 @@ AODB 经 CIIMS adapter 把 XML 报文写入 `CMINMSGS`,格式以架构指定 旧系统线索(来源:旧项目用户故事「Redis key 汇总」「术语与数据语义」「动态航班转历史」): - 投影是 hash `flightInfo`:field 为 `FLID`,value 为完整 `SCHD.FLTR` 对象的带类型 JSON,不设过期;Redis 里没有名为 `schd` 的 key,`schd` 只是 Kafka 主题。新版 value 仍是 `FLTR` 转 JSON,但不使用旧系统的 Jackson 默认类型标注。 -- 写入路径:日计划下载(`DNLD`/`RESP`)整体写入当天航班,单条变更(`ADFT`、`FLOP`)只写对应的一条,转历史时按 `FLID` 逐条移除。整体写入不删除本次映射中缺席的航班,与 `US-07` AC5 相反,新版按 `US-07` AC5 在覆盖范围内刷新。 +- 写入路径:日计划下载(`DNLD`/`RESP`)写入报文里的航班,单条变更(`ADFT`、`FLOP`)只写对应的一条,转历史时按 `FLID` 逐条移除。旧系统整体写入不删除本次映射中缺席的航班。新版按 `US-07` AC5:完整名单在覆盖范围内移除缺席航班;带了时间条件的 `RESP` 只替换回信里的航班,不因缺席移除其他航班。 - 写入前生成主航班的共享航班列表 `MAFL`,共享航班不单列、随主航班下发(`Q6`);旧系统在日计划下载时按 `MAID` 把共享航班挂进主航班 `MAFL`,登机桥字段 `abdg` 本版不提供——旧系统拼它的数据源是机位与登机桥映射缓存,已列入需求「范围与非目标」不交付。 ### 自有 PostgreSQL:内部存储与 admin-api 只读 diff --git a/docs/implementation.md b/docs/implementation.md index e6b13dc..c9cfa30 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -195,7 +195,7 @@ LIMIT PARAM:msgx.pipeline.backfill-batch 1. 已成功 → 幂等,只追加留痕。 2. 整包校验失败 → `DEAD(PROTOCOL)`,不写半包(`INV-4`)。 -3. 分批写:跨运营日整包失败;**覆盖范围内**快照缺席航班标删、发删除事件、删 Redis,范围外的不受影响(`INV-7`)。每批同事务写变更 + 事件(`INV-3`)。 +3. 分批写:跨运营日整包失败。报文里的航班整份替换。完整名单才在覆盖范围内把缺席航班标删、发删除事件、删 Redis;带了时间条件的 `RESP` 不因缺席删除(`INV-7`)。每批同事务写变更 + 事件(`INV-3`)。 4. 整包成功 → `SUCCEEDED` + 回填意图;留痕在事务外。 字段语义见「航班域」。`RESP` 须匹配开放请求,否则不更新。 @@ -374,7 +374,7 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 ### 12.1 SCHD -整包校验通过后分批写入;**覆盖范围内**快照缺席标删,范围外的不受影响(`INV-7`)。未带字段清空(`C-6`)。成功航班推进 `STATE_VERSION` 并写 `schd`+`msg` 事件。重复由 `PROC_STATE` 控制;校验失败整包不写(`INV-4`)。 +整包校验通过后分批写入。报文里的航班整份替换,未带字段清空(`C-6`)。完整名单才在覆盖范围内把缺席航班标删;带了时间条件的 `RESP` 不因缺席删除(`INV-7`)。成功航班推进 `STATE_VERSION` 并写 `schd`+`msg` 事件。重复由 `PROC_STATE` 控制;校验失败整包不写(`INV-4`)。 ### 12.2 动态运行事件(FLOP) @@ -416,7 +416,7 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 ### 12.3 删除与重建 -FDEL:`ACTIVE→DELETED`,只登记 `KAFKA:msg` 变更通知(`C-9`;不再写 `schd` tombstone)。物理删除仅历史清理成功后(`US-14`、`D1`)。覆盖范围内快照缺席也标删(`INV-7`)。 +FDEL:`ACTIVE→DELETED`,只登记 `KAFKA:msg` 变更通知(`C-9`;不再写 `schd` tombstone)。物理删除仅历史清理成功后(`US-14`、`D1`)。完整名单在覆盖范围内的缺席也标删(`INV-7`)。 ADFT:Set-only(`US-04` AC2),未带字段不清。有 `SODT` 则算 `OPERATION_DAY`。 diff --git a/docs/requirements.md b/docs/requirements.md index b1c7488..f989c03 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -90,7 +90,7 @@ ### US-07 导入日计划(DNLD / RESP) -**目标**:日计划是 AODB 在某个时间范围内的完整航班列表:AODB 主动下发(DNLD)或本系统请求后应答(RESP),收到后分批同步本地数据;整包成功时本地航班当前态与快照一致——请求日计划就是主动与 AODB 全量同步一次。覆盖范围由报文自身给出:`DNLD` 覆盖下发时刻起的约 48 小时,`RESP` 覆盖请求的日期区间(`SIS:3.16`)。 +**目标**:日计划报文里的每一班,用这份记录整份替换本地同一班,报文没带的字段清掉。机场主动下发(`DNLD`)是下发时刻起约 48 小时的完整名单。本系统请求后的应答(`RESP`)是否为完整名单,取决于请求有没有带时间条件(`C-4`)。四个条件都没传时,应答是当天全部航班,与机场主动下发的日计划同一范围。带了任一时间条件时,应答只包含这一部分航班,不是完整名单。 | 报文 | 说明 | |---|---| @@ -100,10 +100,10 @@ **验收标准** 1. 报文整体校验(声明的航班数、航班标识、覆盖范围等)通过才处理;校验失败整包拒绝,本地数据不变。 -2. 报文里的航班逐条写入或更新;**覆盖范围内**快照里没有的航班,在本地标记已删除,并登记待发删除消息。覆盖范围外的航班不受本报文影响:前一日延误航班不在 `DNLD` 窗口内,不因缺席被判为已删除。 -3. 以 AODB 下发的数据为准:日计划里某航班没携带的字段,视为 AODB 已删除该值,本地同步清掉。 +2. 报文里的航班逐条整份替换本地同一班。只有完整名单才把范围内有、报文里没有的航班标为已删除,并登记待发删除消息。完整名单是 `DNLD`,以及没带时间条件的 `RESP`。带了时间条件的 `RESP` 不删除报文里没有的航班。完整名单范围外的航班不受本报文影响:前一日延误航班不在 `DNLD` 窗口内,不因缺席被判为已删除。 +3. 以 AODB 下发的数据为准:出现在报文里的航班,没携带的字段视为 AODB 已删除该值,本地同步清掉。 4. 航班量大,分批写入数据库,每批一个事务;处理失败不标记已处理,下轮整包重新处理。 -5. 按快照结果刷新 Redis:报文里的航班写入,**覆盖范围内**缺席的航班移除,覆盖范围外的投影保留。 +5. 按快照结果刷新 Redis:报文里的航班写入。完整名单在覆盖范围内移除缺席航班,范围外的投影保留。带了时间条件的 `RESP` 不因缺席移除其他航班。 ### US-08 通知网页客户端(Kafka) @@ -117,7 +117,7 @@ ### US-09 向 AODB 请求数据 -**目标**:本系统可以主动向 AODB 要数据:14 类参考数据(`RQRD`)+ 1 类日计划(`RQFD`),子类型以消息接口规范为准。参考数据请求由人工发起,日计划请求由 `POST /schd/sync` 触发。 +**目标**:本系统可以主动向 AODB 要数据:14 类参考数据(`RQRD`)+ 1 类日计划(`RQFD`),子类型以消息接口规范为准。参考数据请求由人工发起。日计划请求由网页选定时间条件后,经 `POST /schd/sync` 发出(`C-4`、`C-8`)。 **验收标准** diff --git a/docs/specification.md b/docs/specification.md index 85325f4..fcfd522 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -47,7 +47,7 @@ ### 2.2 上游(AODB / SIS) - **C-3** 一条报文的身份 = `SNDR` + `TYPE` + `STYP` + `SEQN`。`SEQN` 自增,极少在消息服务器重启时重置;重置后不与历史冲突。 -- **C-4** 出站请求写入 `COUTMSGS`;本系统只保证写入信箱,不保证 AODB 收到。编码:`SNDR=OMMS`;`SEQN` 本系统生成;`DTTM` 北京时间 `YYYYMMDDHHMMSS`;`RQRD` 共 14 种子类型(以 SIS 为准);`RQFD` 的 `STYP=NONE`,全量不带 `STDB`/`STDE` 筛选。季度计划不属本系统出站范围(`Q25`)。 +- **C-4** 出站请求写入 `COUTMSGS`;本系统只保证写入信箱,不保证 AODB 收到。编码:`SNDR=OMMS`;`SEQN` 本系统生成;`DTTM` 北京时间 `YYYYMMDDHHMMSS`;`RQRD` 共 14 种子类型(以 SIS 为准);`RQFD` 的 `STYP=NONE`。日计划请求的时间条件由网页传入,本系统不补、不改。`STDB`、`STDE` 按计划到港或计划离港时间筛,`ETDB`、`ETDE` 按预计到港或预计离港时间筛;带 `B` 的是大于等于该时刻,带 `E` 的是小于等于该时刻;多个条件同时成立。时刻格式为 `DDMONYYHHMM`(见 [SIS](legacy/SIS_AODB_RMS-V0.1.md)「RMS 日航班计划请求事件」)。没传的条件不写入报文。四个都不传时不带筛选,AODB 返回当天全部记录,与它主动下发的日计划同一范围,这份回信是完整名单。带了任一时间条件时,回信只是筛选出来的一部分,不是完整名单。回信里的每一班整份替换本地同一班(`C-6`、`US-07` AC3)。只有完整名单才删除范围内缺席的航班;带了时间条件时不删除回信里没有的航班(`US-07` AC2)。季度计划不属本系统出站范围(`Q25`)。 - **C-5** 删航班只打删除标记;主航班与共享航班各自独立标记。 - **C-6** 日计划快照里没带的字段,视为 AODB 已删掉该值,本地也清掉(`US-07` AC3)。 @@ -55,7 +55,8 @@ - **C-7** `POST /cminmsgs/send`:XML 写进入站信箱,与 adapter 走同一套处理。成功返回消息编号(只表示已写入、尚未处理);失败不返回编号。接受 `text/xml`、`application/xml`、`text/plain`(UTF-8);拒收空报文、超长、非法 XML;解析禁止访问外部资源。内网访问由网络层控制;仅供内部联调。 - 待确认:大小上限、HTTP 状态码与响应体 → `Q15`。 -- **C-8** `POST /schd/sync`:登记一次 `RQFD` 日计划请求。请求体为空、请求全量(`C-4`);成功返回请求编号,已有未结案请求时返回 `409`(`Q16`)。同类型若还有未处理完的请求,新请求先不写信箱,等旧的处理完再写 `COUTMSGS`(`US-09` AC1)。 +- **C-8** `POST /schd/sync`:登记一次 `RQFD` 日计划请求。请求体是网页选定的时间条件,写入规则见 `C-4`;成功返回请求编号,已有未结案请求时返回 `409`(`Q16`)。同类型若还有未处理完的请求,新请求先不写信箱,等旧的处理完再写 `COUTMSGS`(`US-09` AC1)。 +- **C-12** `POST /refdata/sync`:登记一次 `RQRD` 参考数据请求。请求体携带网页选定的类别码 `STYP`(取值见 `C-4`,以 SIS 为准);`STYP=RSTA` 时须带资源类型 `RTYP`,其余类别不得带;成功返回请求编号,已有未结案请求时返回 `409`(`Q17`)。登记后等在途请求结案再写 `COUTMSGS`(`US-09` AC1)。 ### 2.4 下游(admin-api、网页客户端) @@ -84,7 +85,7 @@ - **INV-5** 航班当前数据以自有 PG 为准。 - **INV-6** `FLID` 全局唯一。 -- **INV-7** 日计划:**覆盖范围内**快照里没有的航班打删除标记、记删除事件、从 Redis 删掉;覆盖范围外的航班不受本报文影响;快照没带的字段本地清掉(`US-07` AC2/AC3)。 +- **INV-7** 日计划报文里的航班整份替换,没带的字段清掉(`US-07` AC3、`C-6`)。完整名单(`DNLD`,以及没带时间条件的 `RESP`)在覆盖范围内把快照里没有的航班打删除标记、记删除事件、从 Redis 删掉;范围外的航班不受本报文影响。带了时间条件的 `RESP` 不因缺席删除(`US-07` AC2)。 - **INV-8** 已打删除标记的航班必须从 Redis 删掉;删掉才算这条消息处理完成(`US-06` AC1)。 - **INV-9** 日计划可以分批写 PG,但整份 PG 写完且 Redis 按快照刷完才算完成;失败则全部重来(`US-07` AC4/AC5)。 @@ -126,8 +127,8 @@ | Q13 | 已定 | `RQRD` / `RQFD` 编码字段 | `C-4` | | Q14 | 本系统 | 生产用 PostgreSQL 还是 Oracle 11g | Oracle 未验证前不作支持承诺 | | Q15 | 本系统 | `POST /cminmsgs/send` HTTP 约定 | 大小上限、状态码、最终响应体;见 `C-7`。旧系统成功响应是 `ResponseDto`(`err_code=1`,`body` 为信箱编号),且未配置大小上限 | -| Q16 | 已定 | `POST /schd/sync` 请求与响应 | 请求体为空、请求全量(`C-4`);成功返回请求编号,开放请求未结案返回 `409`(`C-8`) | -| Q17 | 本系统 | 人工发 `RQRD` 的入口 | `US-09` 要求能人工发,HTTP 清单里没有;旧系统也没有 | +| Q16 | 已定 | `POST /schd/sync` 请求与响应 | 请求体携带网页选定的时间条件(`C-4`);成功返回请求编号,开放请求未结案返回 `409`(`C-8`) | +| Q17 | 已定 | 人工发 `RQRD` 的入口 | `POST /refdata/sync` 登记一次请求;请求体携带类别码 `STYP`(`RSTA` 时加带 `RTYP`);成功返回请求编号,开放请求未结案返回 `409`(`C-12`) | | Q18 | 已定 | 回退时正在处理的消息怎么办 | 未写回处理标记的消息仍算未处理,由旧系统继续;旧系统停机只等在途任务跑完(`OPS-4`) | | Q19 | — | (未分配) | — | | Q20 | — | (未分配) | — | @@ -151,8 +152,7 @@ | `G-FLOP-UNMAPPED` | [XSD](legacy/unisysaodbsis.xsd) FLOP 字段映射不全 | `US-05` | | `G-MAFL` | 主航班共享列表未做 | `US-06` AC2 | | `G-REF-DATA` | admin-api 生产侧只读接入与联调验收未闭合 | `US-13`;本网关落库与 `ReferenceDataProcessor` 已做(ACM2-93) | -| `G-REQ-OPEN-UNIQUE` | 新请求等待在途请求结案的登记模型未闭合;当前开放态唯一约束会拒绝新登记 | `US-09` | -| `G-REQ-TRACK` | `RQFD` 跟踪已落地;14 类 `RQRD` 人工登记、子类型作废与等待投递未做 | `US-09` | +| `G-REQ-TRACK` | `RQRD` 人工登记已落地(`C-12`);子类型作废与等待投递未做 | `US-09` | | `G-REQ-TRACK-RETENTION` | `REQ_TRACK` 已结案保留期未定(`Q24`) | `US-09` | | `G-SRVT-VIPF` | `SRVT`、`VIPF` 缺席是否清除待 `Q2`;段出现时已落 `FLIGHT_SRVT`/`FLIGHT_VIPF` | `US-05` | @@ -165,6 +165,7 @@ | C-5、`US-06` AC2 | `US-06` AC2 | 删共享联动主航班;删主级联删共享 | | C-7 | `US-02` AC1~AC5 | 三种 Content-Type、UTF-8;空/超长/非法 XML 不写入;禁外部资源;成功编号只表示已写入;与 adapter 同路径建记录 | | C-8、CLM-4 | `US-09` AC1~AC3 | 出站写入 `COUTMSGS`;HTTP 触发先登记,旧的处理完才写入信箱 | +| C-12、CLM-4 | `US-09` AC1~AC3 | `RQRD` 按类别登记并写入 `COUTMSGS`;非法类别或开放槽占用时拒绝 | | C-9、CLM-3 | `US-08` AC1~AC3 | `msg` 单条、`schd` 批量;失败可重试;`msg` 同 `FLID` 保序;同条可能发多次 | | C-10 | `US-13` AC1~AC4 | 全量替换、增删改逐条;空值表示无值不是删除;单类校验失败只停该类 | | INV-1(建立处理记录) | `US-01` AC1/AC2/AC4 | 重扫与重启后记录数不变、行不丢;同一编号重复出现时不新增记录 | @@ -173,7 +174,7 @@ | INV-4 | `US-07` AC1 | 校验失败后 PG 航班数据不变 | | INV-5 | 架构「数据归属与一致性」 | 航班数据只写入自有 PG | | INV-6 | implementation.md「数据模型」 | `FLID` 唯一 | -| INV-7 | `US-07` AC2/AC3 | 覆盖范围内快照缺席航班标删并从 Redis 删,范围外不动;未带字段清空 | +| INV-7 | `US-07` AC2/AC3 | 报文里的航班整份替换、未带字段清空;只有完整名单才在覆盖范围内把缺席航班标删并从 Redis 删 | | INV-8 | `US-06` AC1 | 标删后从 Redis 删;失败下轮重做 | | INV-9 | `US-07` AC4/AC5 | 分批失败全部重来;PG 整份写完后再刷 Redis | | INV-11 | `US-12` AC1/AC2 | 返回全部非共享航班且与 Redis 一致;Redis 故障返回错误 | diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt index 716ea30..c1ede92 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt @@ -10,6 +10,7 @@ import com.gzzn.omms.msgexchange.domain.ErrorClass import com.gzzn.omms.msgexchange.domain.MetaFields import com.gzzn.omms.msgexchange.domain.MsgKind import com.gzzn.omms.msgexchange.domain.OutboundRequestKeys +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter import jakarta.inject.Singleton import javax.xml.stream.XMLInputFactory @@ -90,10 +91,18 @@ class JacksonXmlCodec : XmlCodec { } override fun encodeOutboundRqfd(seqn: Long, dttm: Long): String = + encodeOutboundRqfd(seqn, dttm, RqfdTimeFilter()) + + override fun encodeOutboundRqfd(seqn: Long, dttm: Long, filter: RqfdTimeFilter): String = writeOutbound( SisOutboundMessageXml( meta = outboundMeta("RQFD", "NONE", seqn, dttm), - rqfd = RqfdXml(), + rqfd = RqfdXml( + stdb = filter.stdb, + stde = filter.stde, + etdb = filter.etdb, + etde = filter.etde, + ), ), ) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt index 6921508..5651ba2 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt @@ -1,6 +1,7 @@ package com.gzzn.omms.msgexchange.codec import com.fasterxml.jackson.annotation.JsonIgnoreProperties +import com.fasterxml.jackson.annotation.JsonInclude import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlElementWrapper import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlProperty import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlRootElement @@ -83,9 +84,15 @@ data class ErorXml( @param:JacksonXmlProperty(localName = "ETEX") val etex: String? = null, ) -/** 出站 RQFD 空体占位(C-4 全量不带 STDB/STDE)。 */ +/** 出站 RQFD。没传的时间条件保持 null,序列化时不写出(C-4)。 */ +@JsonInclude(JsonInclude.Include.NON_NULL) @JsonIgnoreProperties(ignoreUnknown = true) -class RqfdXml +data class RqfdXml( + @param:JacksonXmlProperty(localName = "STDB") val stdb: String? = null, + @param:JacksonXmlProperty(localName = "STDE") val stde: String? = null, + @param:JacksonXmlProperty(localName = "ETDB") val etdb: String? = null, + @param:JacksonXmlProperty(localName = "ETDE") val etde: String? = null, +) @JsonIgnoreProperties(ignoreUnknown = true) data class RqrdXml( diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/XmlCodec.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/XmlCodec.kt index e55fc9a..7ff010a 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/XmlCodec.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/XmlCodec.kt @@ -2,6 +2,7 @@ package com.gzzn.omms.msgexchange.codec import com.gzzn.omms.msgexchange.domain.DecodedMessage import com.gzzn.omms.msgexchange.domain.ErrorClass +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter /** * 解码失败的结果:错误分类加上一句原因。 @@ -22,9 +23,13 @@ sealed interface DecodeResult { interface XmlCodec { fun decode(rawXml: String): DecodeResult - /** C-4:RQFD STYP=NONE,空 RQFD 体。 */ + /** C-4:RQFD STYP=NONE。不带时间条件时体为空。 */ fun encodeOutboundRqfd(seqn: Long, dttm: Long): String + /** C-4:只写入网页传来的时间条件,没传的标签不出现。 */ + fun encodeOutboundRqfd(seqn: Long, dttm: Long, filter: RqfdTimeFilter): String = + encodeOutboundRqfd(seqn, dttm) + /** C-4:RQRD,STYP 为参考数据子类型;RSTA 时可带 RTYP。 */ fun encodeOutboundRqrd(styp: String, seqn: Long, dttm: Long, rtyp: String? = null): String } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqfdTimeFilter.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqfdTimeFilter.kt new file mode 100644 index 0000000..ae60ec2 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqfdTimeFilter.kt @@ -0,0 +1,51 @@ +package com.gzzn.omms.msgexchange.domain + +/** + * 日计划请求的时间条件(`C-4`)。网页传入,本系统不补、不改。 + * 空白视为没传。非空但不符 `DDMONYYHHMM` 则整单拒绝。 + */ +data class RqfdTimeFilter( + val stdb: String? = null, + val stde: String? = null, + val etdb: String? = null, + val etde: String? = null, +) { + /** 四个条件都没传:回信是当天完整名单,才删除缺席航班(`US-07` AC2)。 */ + fun isCompleteDay(): Boolean = stdb == null && stde == null && etdb == null && etde == null + + sealed interface Parse { + data class Ok(val filter: RqfdTimeFilter) : Parse + data object Invalid : Parse + } + + companion object { + private val TOKEN = Regex( + """^(0[1-9]|[12]\d|3[01])(JAN|FEB|MAR|APR|MAY|JUN|JUL|AUG|SEP|OCT|NOV|DEC)\d{2}([01]\d|2[0-3])[0-5]\d$""", + ) + + fun parse(stdb: String?, stde: String?, etdb: String?, etde: String?): Parse { + val fields = listOf(stdb, stde, etdb, etde).map { one(it) } + if (fields.any { it is Field.Bad }) return Parse.Invalid + return Parse.Ok( + RqfdTimeFilter( + stdb = (fields[0] as Field.Value).text, + stde = (fields[1] as Field.Value).text, + etdb = (fields[2] as Field.Value).text, + etde = (fields[3] as Field.Value).text, + ), + ) + } + + private fun one(raw: String?): Field { + val text = raw?.trim().orEmpty() + if (text.isEmpty()) return Field.Value(null) + if (!TOKEN.matches(text)) return Field.Bad + return Field.Value(text) + } + } + + private sealed interface Field { + data class Value(val text: String?) : Field + data object Bad : Field + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqrdTarget.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqrdTarget.kt new file mode 100644 index 0000000..6693cbc --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/RqrdTarget.kt @@ -0,0 +1,41 @@ +package com.gzzn.omms.msgexchange.domain + +/** + * RQRD 参考数据请求的目标类别(`C-12`)。 + * 类别码以 SIS 为准(`SIS:4.6.2`);`RSTA` 须带资源类型,其余类别不得带。 + */ +data class RqrdTarget( + val styp: String, + val rtyp: String? = null, +) { + sealed interface Parse { + data class Ok(val target: RqrdTarget) : Parse + data object UnknownStyp : Parse + data object MissingRtyp : Parse + data object UnexpectedRtyp : Parse + } + + companion object { + /** SIS `4.6.2` 的 14 种 RQRD 子类型。 */ + val STYPS = setOf( + "COUL", "ARPT", "AIRL", "AIRC", "REGN", "ORGN", "FLTL", + "TLST", "SLST", "CLST", "GLST", "BLST", "CHLT", "RSTA", + ) + + /** 资源状态类型,见 implementation.md「静态参考数据」。 */ + val RTYPS = setOf("BELT", "CNTR", "GATE", "STND") + + fun parse(styp: String?, rtyp: String?): Parse { + val s = styp?.trim().orEmpty().uppercase() + if (s !in STYPS) return Parse.UnknownStyp + val r = rtyp?.trim().orEmpty().uppercase().ifEmpty { null } + if (s == "RSTA") { + if (r == null) return Parse.MissingRtyp + if (r !in RTYPS) return Parse.UnknownStyp + return Parse.Ok(RqrdTarget(s, r)) + } + if (r != null) return Parse.UnexpectedRtyp + return Parse.Ok(RqrdTarget(s)) + } + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt index d50b857..757119a 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt @@ -26,9 +26,8 @@ data class FlightMainRow( /** * 日计划报文(SCHD DNLD/RESP)里的一条 FLTR 记录,是解码后的产物。 * - * scalars 和 collections 只装报文里真正出现过的字段:出现就覆盖本地值(标量给空串表示 - * 显式清空),没出现就保留库里已有的值;集合一旦出现就按合并后的完整结果整体覆盖写入。 - * 合并规则见 docs/implementation.md「SCHD 日计划」。 + * scalars 和 collections 只装报文里真正出现过的字段。写入当前态时这一班整份替换, + * 没出现的字段清掉(`C-6`)。规则见 docs/implementation.md「SCHD」。 */ data class ScheduleRecord( val flid: String, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngine.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngine.kt index 4ad1a30..82e2323 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngine.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngine.kt @@ -91,8 +91,8 @@ object FlightStateEngine { /** * 把一条日计划记录写成新的当前态:**以 AODB 下发的这份快照为准**——报文带的字段写进去, - * 没带的字段和集合一律清掉(`C-6`、`US-07` AC3)。日计划是覆盖范围内的完整列表, - * 不是增量,所以没有"没出现就保留"这回事(那是 FLOP/ADFT 的语义,见 [mergedState])。 + * 没带的字段和集合一律清掉(`C-6`、`US-07` AC3)。这是这一班的整份替换,不是增量 + * (增量是 FLOP/ADFT,见 [mergedState])。报文里没有的其他航班是否删除见 `INV-7`。 * * keepDeleted = true 时即使收到日计划也保持 DELETED:日计划不能把删掉的航班救回来, * 唯一的恢复入口是 ADFT。调用方负责记一条 SCHD_REVIVE_CONFLICT 告警。 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt index 701ab9a..9e53af4 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt @@ -1,6 +1,8 @@ package com.gzzn.omms.msgexchange.infra.persistence import com.gzzn.omms.msgexchange.domain.ErrorClass +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter +import com.gzzn.omms.msgexchange.domain.RqrdTarget import com.gzzn.omms.msgexchange.domain.MsgEvent import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus @@ -321,7 +323,6 @@ interface SnapshotLogRepository { /** * 上游请求的跟踪表:记录我们发出去的请求、以及对方回来的应答。 * - * 目前只有表和读写方法,**还没有运行时协调器**(出站写信箱、超时、应答匹配都没实现)。 * 设计意图是:RESP 报文按(运营日、发送方、请求类型)匹配最近一条待应答的请求; * 同一类请求同时只保留一条有效,发新请求时把旧的置为已过期。 */ @@ -339,10 +340,25 @@ interface ReqTrackRepository { val writeUncertain: Boolean = false, val sentAt: Instant? = null, val createdAt: Instant? = null, - ) + val stdb: String? = null, + val stde: String? = null, + val etdb: String? = null, + val etde: String? = null, + val styp: String? = null, + val rtyp: String? = null, + ) { + fun rqfdFilter() = RqfdTimeFilter(stdb, stde, etdb, etde) + fun rqrdTarget() = if (styp != null) RqrdTarget(styp, rtyp) else null + } - /** 登记新请求:同类(类型+运营日+发送方)旧有效请求先置 EXPIRED。 */ - fun insert(reqType: String, operationDay: LocalDate, sender: String): Long + /** 登记新请求:同类(类型+运营日+发送方)旧有效请求先置 EXPIRED。日计划时间条件原样留下(C-4)。 */ + fun insert( + reqType: String, + operationDay: LocalDate, + sender: String, + filter: RqfdTimeFilter = RqfdTimeFilter(), + target: RqrdTarget? = null, + ): Long fun findLatest(reqType: String, operationDay: LocalDate, sender: String, states: List): Req? diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt index 9aef4e6..f28d6a9 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcPgRepositories.kt @@ -1091,7 +1091,13 @@ class JdbcReqTrackRepository( private val ds: DataSource, private val clock: Clock, ) : ReqTrackRepository { - override fun insert(reqType: String, operationDay: LocalDate, sender: String): Long { + override fun insert( + reqType: String, + operationDay: LocalDate, + sender: String, + filter: com.gzzn.omms.msgexchange.domain.RqfdTimeFilter, + target: com.gzzn.omms.msgexchange.domain.RqrdTarget?, + ): Long { ds.update( """ UPDATE req_track SET state = 'EXPIRED' @@ -1103,12 +1109,19 @@ class JdbcReqTrackRepository( ps.setString(3, sender) } return ds.updateReturningLong( - "INSERT INTO req_track (req_type, operation_day, sender, state, created_at) VALUES (?, ?, ?, 'PENDING', ?) RETURNING req_id", + "INSERT INTO req_track (req_type, operation_day, sender, state, created_at, stdb, stde, etdb, etde, styp, rtyp) " + + "VALUES (?, ?, ?, 'PENDING', ?, ?, ?, ?, ?, ?, ?) RETURNING req_id", { ps -> ps.setString(1, reqType) ps.setDate(2, java.sql.Date.valueOf(operationDay)) ps.setString(3, sender) ps.setTimestamp(4, clock.instant().toSqlTimestamp()) + ps.setString(5, filter.stdb) + ps.setString(6, filter.stde) + ps.setString(7, filter.etdb) + ps.setString(8, filter.etde) + ps.setString(9, target?.styp) + ps.setString(10, target?.rtyp) }, ) } @@ -1245,5 +1258,11 @@ class JdbcReqTrackRepository( writeUncertain = rs.getBoolean("write_uncertain"), sentAt = rs.getInstant("sent_at"), createdAt = rs.getInstant("created_at"), + stdb = rs.getString("stdb"), + stde = rs.getString("stde"), + etdb = rs.getString("etdb"), + etde = rs.getString("etde"), + styp = rs.getString("styp"), + rtyp = rs.getString("rtyp"), ) } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt index 9914c52..c510dc8 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt @@ -496,7 +496,13 @@ class StubReqTrack(private val clock: Clock = Clock.systemUTC()) : ReqTrackRepos fun clear() = rows.clear() - override fun insert(reqType: String, operationDay: LocalDate, sender: String): Long { + override fun insert( + reqType: String, + operationDay: LocalDate, + sender: String, + filter: com.gzzn.omms.msgexchange.domain.RqfdTimeFilter, + target: com.gzzn.omms.msgexchange.domain.RqrdTarget?, + ): Long { // 同一类请求(类型 + 运营日 + 发送方)只保留一条有效,旧的先置为已过期 rows.values.filter { it.reqType == reqType && it.operationDay == operationDay && it.sender == sender && @@ -506,6 +512,12 @@ class StubReqTrack(private val clock: Clock = Clock.systemUTC()) : ReqTrackRepos rows[id] = ReqTrackRepository.Req( id, reqType, operationDay, sender, ReqTrackRepository.ReqState.PENDING, createdAt = clock.instant(), + stdb = filter.stdb, + stde = filter.stde, + etdb = filter.etdb, + etde = filter.etde, + styp = target?.styp, + rtyp = target?.rtyp, ) return id } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxController.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxController.kt index 3bb0e88..5b2c9b5 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxController.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxController.kt @@ -1,8 +1,12 @@ package com.gzzn.omms.msgexchange.ingress +import com.fasterxml.jackson.annotation.JsonIgnoreProperties +import com.fasterxml.jackson.annotation.JsonProperty import io.micronaut.http.HttpResponse +import io.micronaut.http.HttpStatus import io.micronaut.http.MediaType import io.micronaut.http.annotation.Body +import io.micronaut.http.annotation.Consumes import io.micronaut.http.annotation.Controller import io.micronaut.http.annotation.Post import io.micronaut.http.annotation.Produces @@ -16,6 +20,7 @@ import io.micronaut.http.annotation.Produces class InboxController( private val inbox: InboxService, private val schdSync: SchdSyncService, + private val refdataSync: RefdataSyncService, ) { @Post("/cminmsgs/send") @@ -25,14 +30,51 @@ class InboxController( return HttpResponse.ok(receipt.msgId.toString()) // TODO: 先返回消息 ID 文本,等跟现役响应体逐字对拍过再定稿 } - /** C-8:登记 RQFD 出站请求;开放槽占用时 409,不重复登记。 */ + /** C-8:登记 RQFD。请求体是网页选定的时间条件;开放槽占用时 409。 */ @Post("/schd/sync") + @Consumes(MediaType.APPLICATION_JSON) @Produces(MediaType.TEXT_PLAIN) - fun schdSync(): HttpResponse = - when (val outcome = schdSync.trigger()) { + fun schdSync(@Body body: SchdSyncBody?): HttpResponse = + when (val outcome = schdSync.trigger(body?.stdb, body?.stde, body?.etdb, body?.etde)) { is SchdSyncService.Outcome.Registered -> HttpResponse.ok(outcome.reqId.toString()) SchdSyncService.Outcome.OpenExists -> - HttpResponse.status(io.micronaut.http.HttpStatus.CONFLICT).body("open-request-exists") + HttpResponse.status(HttpStatus.CONFLICT).body("open-request-exists") + SchdSyncService.Outcome.InvalidTime -> + HttpResponse.status(HttpStatus.BAD_REQUEST).body("invalid-time-filter") + } + + /** C-12:登记 RQRD。请求体是网页选定的类别码;开放槽占用时 409。 */ + @Post("/refdata/sync") + @Consumes(MediaType.APPLICATION_JSON) + @Produces(MediaType.TEXT_PLAIN) + fun refdataSync(@Body body: RefdataSyncBody?): HttpResponse = + when (val outcome = refdataSync.trigger(body?.styp, body?.rtyp)) { + is RefdataSyncService.Outcome.Registered -> + HttpResponse.ok(outcome.reqId.toString()) + RefdataSyncService.Outcome.OpenExists -> + HttpResponse.status(HttpStatus.CONFLICT).body("open-request-exists") + RefdataSyncService.Outcome.UnknownStyp -> + HttpResponse.status(HttpStatus.BAD_REQUEST).body("unknown-styp") + RefdataSyncService.Outcome.MissingRtyp -> + HttpResponse.status(HttpStatus.BAD_REQUEST).body("missing-rtyp") + RefdataSyncService.Outcome.UnexpectedRtyp -> + HttpResponse.status(HttpStatus.BAD_REQUEST).body("unexpected-rtyp") } } + +/** 网页提交的参考数据类别。`RSTA` 须带 `RTYP`,其余类别不得带(C-12)。 */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class RefdataSyncBody( + @param:JsonProperty("STYP") val styp: String? = null, + @param:JsonProperty("RTYP") val rtyp: String? = null, +) + +/** 网页提交的日计划时间条件。字段可缺,缺了表示不按该项筛选(C-4)。 */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class SchdSyncBody( + @param:JsonProperty("STDB") val stdb: String? = null, + @param:JsonProperty("STDE") val stde: String? = null, + @param:JsonProperty("ETDB") val etdb: String? = null, + @param:JsonProperty("ETDE") val etde: String? = null, +) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/RefdataSyncService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/RefdataSyncService.kt new file mode 100644 index 0000000..7ce1ea8 --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/RefdataSyncService.kt @@ -0,0 +1,32 @@ +package com.gzzn.omms.msgexchange.ingress + +import com.gzzn.omms.msgexchange.domain.RqrdTarget +import com.gzzn.omms.msgexchange.processing.OutboundRequestService +import jakarta.inject.Singleton + +/** + * C-12:登记 RQRD。请求体是网页选定的类别码;RSTA 须带资源类型。 + */ +@Singleton +class RefdataSyncService(private val outbound: OutboundRequestService) { + sealed interface Outcome { + data class Registered(val reqId: Long) : Outcome + data object OpenExists : Outcome + data object UnknownStyp : Outcome + data object MissingRtyp : Outcome + data object UnexpectedRtyp : Outcome + } + + fun trigger(styp: String?, rtyp: String?): Outcome { + val target = when (val parsed = RqrdTarget.parse(styp, rtyp)) { + is RqrdTarget.Parse.Ok -> parsed.target + RqrdTarget.Parse.UnknownStyp -> return Outcome.UnknownStyp + RqrdTarget.Parse.MissingRtyp -> return Outcome.MissingRtyp + RqrdTarget.Parse.UnexpectedRtyp -> return Outcome.UnexpectedRtyp + } + return when (val r = outbound.registerRqrdSync(target)) { + is OutboundRequestService.RegisterOutcome.Registered -> Outcome.Registered(r.reqId) + OutboundRequestService.RegisterOutcome.OpenExists -> Outcome.OpenExists + } + } +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/SchdSyncService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/SchdSyncService.kt index e323c0d..71fb629 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/SchdSyncService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/SchdSyncService.kt @@ -1,5 +1,6 @@ package com.gzzn.omms.msgexchange.ingress +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter import com.gzzn.omms.msgexchange.processing.OutboundRequestService import jakarta.inject.Singleton @@ -8,11 +9,17 @@ class SchdSyncService(private val outbound: OutboundRequestService) { sealed interface Outcome { data class Registered(val reqId: Long) : Outcome data object OpenExists : Outcome + data object InvalidTime : Outcome } - fun trigger(): Outcome = - when (val r = outbound.registerRqfdSync()) { + fun trigger(stdb: String? = null, stde: String? = null, etdb: String? = null, etde: String? = null): Outcome { + val filter = when (val parsed = RqfdTimeFilter.parse(stdb, stde, etdb, etde)) { + RqfdTimeFilter.Parse.Invalid -> return Outcome.InvalidTime + is RqfdTimeFilter.Parse.Ok -> parsed.filter + } + return when (val r = outbound.registerRqfdSync(filter)) { is OutboundRequestService.RegisterOutcome.Registered -> Outcome.Registered(r.reqId) OutboundRequestService.RegisterOutcome.OpenExists -> Outcome.OpenExists } + } } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt index c0961f3..e2a616d 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestService.kt @@ -4,6 +4,8 @@ import com.gzzn.omms.msgexchange.codec.XmlCodec import com.gzzn.omms.msgexchange.config.OperationDayProps import com.gzzn.omms.msgexchange.config.PipelineProps import com.gzzn.omms.msgexchange.domain.OutboundRequestKeys +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter +import com.gzzn.omms.msgexchange.domain.RqrdTarget import com.gzzn.omms.msgexchange.infra.persistence.CoutmsgOutboxRepository import com.gzzn.omms.msgexchange.infra.persistence.ReqTrackRepository import jakarta.inject.Singleton @@ -50,20 +52,40 @@ class OutboundRequestService( listOf(ReqTrackRepository.ReqState.PENDING, ReqTrackRepository.ReqState.SENT), ) != null - fun registerRqfdSync(): RegisterOutcome { + fun registerRqfdSync(filter: RqfdTimeFilter = RqfdTimeFilter()): RegisterOutcome { val day = currentOperationDay() if (hasOpenRequest(OutboundRequestKeys.RQFD_REQ_TYPE, day)) return RegisterOutcome.OpenExists - val reqId = reqTrack.insert(OutboundRequestKeys.RQFD_REQ_TYPE, day, OutboundRequestKeys.SENDER) + val reqId = reqTrack.insert(OutboundRequestKeys.RQFD_REQ_TYPE, day, OutboundRequestKeys.SENDER, filter) dispatchPending(limit = 1) return RegisterOutcome.Registered(reqId) } - fun hasOpenSentRqfd(operationDay: LocalDate = currentOperationDay()): Boolean = + /** C-12:登记一次 RQRD。开放槽占用时 OpenExists;落信沿用排队规则。 */ + fun registerRqrdSync(target: RqrdTarget): RegisterOutcome { + val day = currentOperationDay() + if (hasOpenRequest(OutboundRequestKeys.RQRD_REQ_TYPE, day)) return RegisterOutcome.OpenExists + val reqId = reqTrack.insert(OutboundRequestKeys.RQRD_REQ_TYPE, day, OutboundRequestKeys.SENDER, target = target) + dispatchPending(limit = 1) + return RegisterOutcome.Registered(reqId) + } + + fun openSentRqfd(operationDay: LocalDate = currentOperationDay()): ReqTrackRepository.Req? = reqTrack.findLatest( OutboundRequestKeys.RQFD_REQ_TYPE, operationDay, OutboundRequestKeys.SENDER, listOf(ReqTrackRepository.ReqState.SENT), + ) + + fun hasOpenSentRqfd(operationDay: LocalDate = currentOperationDay()): Boolean = + openSentRqfd(operationDay) != null + + fun hasOpenSentRqrd(operationDay: LocalDate = currentOperationDay()): Boolean = + reqTrack.findLatest( + OutboundRequestKeys.RQRD_REQ_TYPE, + operationDay, + OutboundRequestKeys.SENDER, + listOf(ReqTrackRepository.ReqState.SENT), ) != null fun completeRqfdResponse(operationDay: LocalDate = currentOperationDay()): Boolean = @@ -102,8 +124,13 @@ class OutboundRequestService( val seqn = outbox.nextOutboundSeqn() val dttm = beijingMetaDttm() val xml = when (req.reqType) { - OutboundRequestKeys.RQFD_REQ_TYPE -> codec.encodeOutboundRqfd(seqn, dttm) - OutboundRequestKeys.RQRD_REQ_TYPE -> codec.encodeOutboundRqrd("AIRL", seqn, dttm) + OutboundRequestKeys.RQFD_REQ_TYPE -> codec.encodeOutboundRqfd(seqn, dttm, req.rqfdFilter()) + OutboundRequestKeys.RQRD_REQ_TYPE -> { + val styp = req.styp ?: return false.also { + log.warn("RQRD row without STYP reqId={}", req.reqId) + } + codec.encodeOutboundRqrd(styp, seqn, dttm, req.rtyp) + } else -> { log.warn("unknown req_type for dispatch reqId={} type={}", req.reqId, req.reqType) return false diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt index 48d1ca0..7b6bea5 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -203,9 +203,10 @@ class MessageProcessor( if (body == null) return deadMalformed(head, "missing-schd-body") when (kind.subtype) { MsgKind.SchdSubtype.DNLD -> - scheduleProcessor.applyScheduleRecords(head, decoded) + scheduleProcessor.applyScheduleRecords(head, decoded, deleteAbsent = true) MsgKind.SchdSubtype.RESP -> { - if (!outbound.hasOpenSentRqfd()) { + val open = outbound.openSentRqfd() + if (open == null) { log.info("SCHD-RESP without open RQFD -> SKIPPED msgId={}", head.msgId) procState.markTerminal( head.msgId, ProcStatus.SKIPPED, @@ -214,7 +215,11 @@ class MessageProcessor( ) return } - val respResult = scheduleProcessor.applyScheduleRecords(head, decoded) + val respResult = scheduleProcessor.applyScheduleRecords( + head, + decoded, + deleteAbsent = open.rqfdFilter().isCompleteDay(), + ) if (respResult is ApplyResult.Succeeded || respResult is ApplyResult.ReplaySkipped) { outbound.completeRqfdResponse() } @@ -240,6 +245,15 @@ class MessageProcessor( is MsgKind.RefData -> { val body = decoded.body as? com.gzzn.omms.msgexchange.domain.ref.RefDataBody // validated below if (body == null) return deadMalformed(head, "missing-refdata-body") + if (body.styp.equals("RESP", ignoreCase = true) && !outbound.hasOpenSentRqrd()) { + log.info("REF-RESP without open RQRD -> SKIPPED msgId={}", head.msgId) + procState.markTerminal( + head.msgId, ProcStatus.SKIPPED, + lastError = "resp-guard:no-open-req", + now = clock.instant(), + ) + return + } referenceDataProcessor.apply(head, decoded, kind.type) } MsgKind.Eror -> { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt index 8e329bd..947679f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt @@ -50,7 +50,8 @@ class ProtocolViolation(message: String) : RuntimeException(message) * 校验或运营日核对不过就整包拒绝、一条都不落(`INV-4`);写入阶段按批分事务, * 整份写完且投影刷完才算处理完成,中途失败下轮整包重来(`INV-9`)。 * - * 日计划是覆盖范围内的完整列表:范围内缺席的航班要标删(`INV-7`),范围外的不受影响。 + * 报文里的航班整份替换。只有完整名单才删除覆盖范围内的缺席航班(`INV-7`)。 + * 带了时间条件的应答不是完整名单,不删除回信里没有的航班。 */ @Singleton class ScheduleProcessor( @@ -70,7 +71,7 @@ class ScheduleProcessor( cutoffHour = operationDayProps.cutoffHour, ) - fun applyScheduleRecords(head: ProcState, msg: DecodedMessage): ApplyResult { + fun applyScheduleRecords(head: ProcState, msg: DecodedMessage, deleteAbsent: Boolean = true): ApplyResult { val body = msg.body as? ScheduleBody ?: return ApplyResult.DeadProtocol("missing-schd-body") val started = System.nanoTime() @@ -103,7 +104,7 @@ class ScheduleProcessor( val batchSize = props.schd.snapshotBatch return try { // 分批写:每批一个事务(`US-07` AC4),整份写完才刷投影、才记终态(`INV-9`)。 - // 覆盖范围内缺席的航班在最后清扫(`INV-7`)。 + // 完整名单才在最后清扫覆盖范围内的缺席航班(`INV-7`)。 val upserted = commit.commitBatched(head) { batches -> // 一次批量查出这些航班现有的运营日,逐个比对(避免逐条查询)。 // 这一比对必须在任何一批写入之前跑完:运营日冲突要整包拒绝、一条都不落(`INV-4`)。 @@ -150,14 +151,16 @@ class ScheduleProcessor( } } - sweepAbsent( - head = head, - batches = batches, - coverage = ok.perRecordDay.values.toSet(), - present = ok.perRecordDay.keys, - batchSize = batchSize, - flags = flags, - ) + if (deleteAbsent) { + sweepAbsent( + head = head, + batches = batches, + coverage = ok.perRecordDay.values.toSet(), + present = ok.perRecordDay.keys, + batchSize = batchSize, + flags = flags, + ) + } written } logSnapshot(head, body, SnapshotResult.COMMITTED, upserted, flags, started) @@ -169,7 +172,7 @@ class ScheduleProcessor( } /** - * 覆盖范围内缺席的航班:标删、登记删除通知、从投影里删掉(`INV-7`、`US-07` AC2/AC5)。 + * 完整名单覆盖范围内缺席的航班:标删、登记删除通知、从投影里删掉(`INV-7`、`US-07` AC2/AC5)。 * * [coverage] 是这份报文覆盖的运营日,取自报文自身——每条记录的 `SODT` 推出的运营日 * (整包校验已保证每条都算得出)。覆盖范围外的航班一条都不碰:前一日延误的航班不在 diff --git a/src/main/resources/db/migration/V1__flight_state_baseline.sql b/src/main/resources/db/migration/V1__flight_state_baseline.sql index d8a099d..1628ecc 100644 --- a/src/main/resources/db/migration/V1__flight_state_baseline.sql +++ b/src/main/resources/db/migration/V1__flight_state_baseline.sql @@ -287,14 +287,13 @@ CREATE UNIQUE INDEX uq_schd_event WHERE TARGET = 'KAFKA:schd'; -- ⑥ 请求状态机:只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型)匹配最新一条 --- 开放请求。同类请求只留一条有效,新请求置旧请求为 EXPIRED。登记、超时与应答匹配 --- 尚未实现(G-REQ-TRACK、G-REQ-OPEN-UNIQUE)。 +-- 开放请求。同类请求只留一条有效,新请求置旧请求为 EXPIRED。 CREATE TABLE REQ_TRACK ( REQ_ID BIGSERIAL PRIMARY KEY, REQ_TYPE VARCHAR(20) NOT NULL, OPERATION_DAY DATE NOT NULL, -- 请求覆盖运营日 SENDER VARCHAR(64) NOT NULL, -- 请求发送方(匹配键之一) - STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DONE/EXPIRED + STATE VARCHAR(16) NOT NULL, -- PENDING/SENT/DONE/EXPIRED/FAILED COUTMSGS_ID BIGINT, SENT_AT TIMESTAMP(6) WITH TIME ZONE, COMPLETED_AT TIMESTAMP(6) WITH TIME ZONE, diff --git a/src/main/resources/db/migration/V7__req_track_rqfd_window.sql b/src/main/resources/db/migration/V7__req_track_rqfd_window.sql new file mode 100644 index 0000000..bdabe73 --- /dev/null +++ b/src/main/resources/db/migration/V7__req_track_rqfd_window.sql @@ -0,0 +1,6 @@ +-- 日计划请求的时间条件:网页传入后原样留到落信,空表示没传(C-4)。 +ALTER TABLE req_track + ADD COLUMN IF NOT EXISTS stdb VARCHAR(11), + ADD COLUMN IF NOT EXISTS stde VARCHAR(11), + ADD COLUMN IF NOT EXISTS etdb VARCHAR(11), + ADD COLUMN IF NOT EXISTS etde VARCHAR(11); diff --git a/src/main/resources/db/migration/V8__req_track_rqrd_target.sql b/src/main/resources/db/migration/V8__req_track_rqrd_target.sql new file mode 100644 index 0000000..7acc72c --- /dev/null +++ b/src/main/resources/db/migration/V8__req_track_rqrd_target.sql @@ -0,0 +1,5 @@ +-- RQRD 参考数据请求的子类型:网页选定后原样留到落信(C-12)。 +-- RQFD 行两列为空。开放槽仍按 REQ_TYPE 单槽,不按 STYP 分槽。 +ALTER TABLE req_track + ADD COLUMN IF NOT EXISTS styp VARCHAR(4), + ADD COLUMN IF NOT EXISTS rtyp VARCHAR(4); diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt index c6fe0f5..5a5f4da 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt @@ -1,6 +1,7 @@ package com.gzzn.omms.msgexchange.codec import com.gzzn.omms.msgexchange.domain.MsgKind +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot import com.gzzn.omms.msgexchange.domain.flight.FlightState import com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine @@ -17,12 +18,27 @@ class JacksonXmlCodecTest { private val codec = JacksonXmlCodec() @Test - fun `encode outbound RQFD uses OMMS meta and empty body per C-4`() { + fun `encode outbound RQFD uses OMMS meta and omits time filters that were not sent`() { val xml = codec.encodeOutboundRqfd(1243L, 20021010090311L) assertTrue(xml.contains("OMMS")) assertTrue(xml.contains("RQFD")) assertTrue(xml.contains("NONE")) assertTrue(xml.contains("12JAN041730")) + assertTrue(xml.contains("13JAN042359")) + assertFalse(xml.contains("STDE")) + assertFalse(xml.contains("ETDB")) } @Test diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt index d1b8c6c..3b0b68a 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt @@ -11,7 +11,7 @@ import java.sql.DriverManager /** * 在真实 PostgreSQL 上跑一遍迁移链,确认结果符合预期:`V1__flight_state_baseline.sql` - * 基线加 V2~V6 增量迁移 + * 基线加 V2~V8 增量迁移 * 依序执行成功,该建的表和 PIPELINE_LOCK 单行种子都在,回填事实落在基线里,`INBOX_CURSOR`、 * `BACKFILL_TODO`、`idx_evt_flid`、`PROC_STATE` 的处理开始时间列都不复存在;FLIGHT_CHUTE * 的类字段列已由 V2 更名为 CCLS/CTYP(SIS 口径)。 @@ -59,6 +59,8 @@ class FlywayMigrationTest { "4" to "V4__req_track_outbound_seqn.sql", "5" to "V5__unmapped_field_srvt_vipf.sql", "6" to "V6__basicdata_ref_data.sql", + "7" to "V7__req_track_rqfd_window.sql", + "8" to "V8__req_track_rqrd_target.sql", ), records.map { it.first to it.second }, ) @@ -133,6 +135,14 @@ class FlywayMigrationTest { assertTrue(indexDef.contains("UNIQUE"), "uq_req_open 必须是唯一索引:$indexDef") assertTrue(indexDef.contains("WHERE"), "uq_req_open 必须是仅约束开放态的部分索引:$indexDef") } + stmt.executeQuery( + "SELECT column_name FROM information_schema.columns WHERE table_name = 'req_track' " + + "AND column_name IN ('stdb', 'stde', 'etdb', 'etde', 'styp', 'rtyp')", + ).use { rs -> + val cols = mutableSetOf() + while (rs.next()) cols.add(rs.getString("column_name")) + assertEquals(setOf("stdb", "stde", "etdb", "etde", "styp", "rtyp"), cols) + } stmt.executeUpdate( "INSERT INTO req_track (req_type, operation_day, sender, state, created_at) " + "VALUES ('RQFD-NONE', DATE '2026-09-12', 'RMS', 'PENDING', now())", diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestTest.kt index 856a5b8..dc8aece 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/OutboundRequestTest.kt @@ -5,6 +5,8 @@ import com.gzzn.omms.msgexchange.codec.JacksonXmlCodec import com.gzzn.omms.msgexchange.config.OperationDayProps import com.gzzn.omms.msgexchange.config.PipelineProps import com.gzzn.omms.msgexchange.domain.OutboundRequestKeys +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter +import com.gzzn.omms.msgexchange.ingress.SchdSyncService import com.gzzn.omms.msgexchange.infra.persistence.CoutmsgOutboxRepository import com.gzzn.omms.msgexchange.infra.persistence.ReqTrackRepository import com.gzzn.omms.msgexchange.infra.stub.StubCoutmsgOutbox @@ -45,6 +47,28 @@ class OutboundRequestTest { assertTrue(outbox.messages.values.first().contains("RQFD")) } + @Test + fun `schd sync writes the page time filter into COUTMSGS and leaves out the rest`() { + val req = StubReqTrack(clock) + val outbox = StubCoutmsgOutbox() + val svc = service(req, outbox) + + svc.registerRqfdSync(RqfdTimeFilter(stdb = "12JAN041730")) + val xml = outbox.messages.values.first() + assertTrue(xml.contains("12JAN041730")) + assertFalse(xml.contains("STDE")) + assertFalse(xml.contains("ETDB")) + assertFalse(xml.contains("ETDE")) + } + + @Test + fun `invalid time filter is rejected before registration`() { + val req = StubReqTrack(clock) + val sync = SchdSyncService(service(req, StubCoutmsgOutbox())) + assertEquals(SchdSyncService.Outcome.InvalidTime, sync.trigger(stdb = "today")) + assertTrue(req.rows.isEmpty()) + } + @Test fun `open slot rejects duplicate schd sync registration`() { val req = StubReqTrack(clock) @@ -108,4 +132,71 @@ class OutboundRequestTest { assertTrue(svc.failFromEror(42L, "RQFD", "NONE")) assertEquals(ReqTrackRepository.ReqState.FAILED, req.rows[reqId]!!.state) } + + @Test + fun `refdata sync registers AIRL and dispatches RQRD into COUTMSGS`() { + val req = StubReqTrack(clock) + val outbox = StubCoutmsgOutbox() + val sync = com.gzzn.omms.msgexchange.ingress.RefdataSyncService(service(req, outbox)) + + val outcome = sync.trigger("AIRL", null) + assertTrue(outcome is com.gzzn.omms.msgexchange.ingress.RefdataSyncService.Outcome.Registered) + val xml = outbox.messages.values.first() + assertTrue(xml.contains("RQRD")) + assertTrue(xml.contains("AIRL")) + } + + @Test + fun `refdata sync writes RTYP only for RSTA`() { + val req = StubReqTrack(clock) + val outbox = StubCoutmsgOutbox() + val svc = service(req, outbox) + + svc.registerRqrdSync(com.gzzn.omms.msgexchange.domain.RqrdTarget("RSTA", "GATE")) + val xml = outbox.messages.values.first() + assertTrue(xml.contains("RSTA")) + assertTrue(xml.contains("GATE")) + } + + @Test + fun `refdata sync rejects unknown styp and rtyp misuse before registration`() { + val req = StubReqTrack(clock) + val sync = com.gzzn.omms.msgexchange.ingress.RefdataSyncService(service(req, StubCoutmsgOutbox())) + + assertEquals( + com.gzzn.omms.msgexchange.ingress.RefdataSyncService.Outcome.UnknownStyp, + sync.trigger("NOPE", null), + ) + assertEquals( + com.gzzn.omms.msgexchange.ingress.RefdataSyncService.Outcome.MissingRtyp, + sync.trigger("RSTA", null), + ) + assertEquals( + com.gzzn.omms.msgexchange.ingress.RefdataSyncService.Outcome.UnexpectedRtyp, + sync.trigger("AIRL", "GATE"), + ) + assertTrue(req.rows.isEmpty()) + } + + @Test + fun `open slot rejects duplicate refdata sync registration`() { + val req = StubReqTrack(clock) + val sync = com.gzzn.omms.msgexchange.ingress.RefdataSyncService(service(req, StubCoutmsgOutbox())) + + sync.trigger("AIRL", null) + assertEquals( + com.gzzn.omms.msgexchange.ingress.RefdataSyncService.Outcome.OpenExists, + sync.trigger("AIRL", null), + ) + } + + @Test + fun `RQRD and RQFD use separate open slots`() { + val req = StubReqTrack(clock) + val rqrd = com.gzzn.omms.msgexchange.ingress.RefdataSyncService(service(req, StubCoutmsgOutbox())) + val rqfd = SchdSyncService(service(req, StubCoutmsgOutbox())) + + rqrd.trigger("AIRL", null) + assertTrue(rqfd.trigger() is SchdSyncService.Outcome.Registered) + } } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt index f32690b..afd4423 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/RespGuardTest.kt @@ -9,6 +9,10 @@ import com.gzzn.omms.msgexchange.domain.DecodedMessage import com.gzzn.omms.msgexchange.domain.MetaFields import com.gzzn.omms.msgexchange.domain.MsgKind import com.gzzn.omms.msgexchange.domain.OutboundRequestKeys +import com.gzzn.omms.msgexchange.domain.RqfdTimeFilter +import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot +import com.gzzn.omms.msgexchange.domain.flight.FlightState +import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.infra.metrics.PipelineCounters @@ -25,6 +29,7 @@ import com.gzzn.omms.msgexchange.infra.stub.StubReqTrack import com.gzzn.omms.msgexchange.infra.stub.StubSnapshotLog import com.gzzn.omms.msgexchange.MutableClock import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test import java.time.Instant import java.time.LocalDate @@ -49,10 +54,10 @@ class RespGuardTest { inbox: StubInbox, req: StubReqTrack, decoded: DecodedMessage, + flights: StubFlightState = StubFlightState(), ): MessageProcessor { val props = PipelineProps() val opDay = OperationDayProps().apply { zone = "Asia/Shanghai"; cutoffHour = 0 } - val flights = StubFlightState() val events = StubMsgEvents() val commit = testCommit(procState = proc, msgEvents = events, projection = StubFlightProjectionPort(), clock = clock) return MessageProcessor( @@ -94,6 +99,41 @@ class RespGuardTest { assertEquals("resp-guard:no-open-req", proc.find(msgId)!!.lastError) } + @Test + fun `filtered SCHD-RESP replaces flights in the reply and does not delete the rest`() { + val req = StubReqTrack(clock) + val reqId = req.insert( + OutboundRequestKeys.RQFD_REQ_TYPE, + day, + OutboundRequestKeys.SENDER, + RqfdTimeFilter(stdb = "15DEC261700"), + ) + req.linkCoutmsgs(reqId, 1L, 1L) + req.markSent(reqId, clock.instant()) + + val flights = StubFlightState() + val flightDay = LocalDate.of(2026, 12, 15) + listOf("121", "122").forEach { flid -> + flights.persistFullState( + FlightSnapshot(flid, flightDay, FlightState.ACTIVE, 1, mapOf("SODT" to "15DEC261723"), emptyMap()), + msgId = 1, + now = Instant.EPOCH, + ) + } + val body = ScheduleBody(1, listOf(ScheduleRecord("121", mapOf("SODT" to "15DEC261800")))) + val proc = StubProcState() + proc.insertIfAbsent(msgId, null) + val inbox = StubInbox() + inbox.raws[msgId] = "" + val p = processor(proc, inbox, req, schdMsg(MsgKind.SchdSubtype.RESP, body), flights) + + p.processOne(ProcState(msgId, ProcStatus.PENDING, updatedAt = Instant.EPOCH)) + + assertEquals(ProcStatus.SUCCEEDED, proc.find(msgId)!!.state) + assertEquals(FlightState.ACTIVE, flights.findMainRow("122")!!.state) + assertTrue(flights.findMainRow("121")!!.stateVersion > 1L) + } + @Test fun `SCHD-DNLD does not require open request`() { val proc = StubProcState() @@ -125,4 +165,23 @@ class RespGuardTest { assertEquals(ProcStatus.SKIPPED, proc.find(msgId)!!.state) } + + @Test + fun `REF-RESP without open RQRD is SKIPPED`() { + val proc = StubProcState() + proc.insertIfAbsent(msgId, null) + val inbox = StubInbox() + inbox.raws[msgId] = "" + val decoded = DecodedMessage( + meta = MetaFields("AODB", "AIRL", "RESP", 1L, 1L), + kind = MsgKind.RefData("AIRL"), + rawXml = "", + body = com.gzzn.omms.msgexchange.domain.ref.RefDataBody(category = "AIRL", styp = "RESP", records = emptyList()), + ) + val p = processor(proc, inbox, StubReqTrack(clock), decoded) + p.processOne(ProcState(msgId, ProcStatus.PENDING, updatedAt = Instant.EPOCH)) + + assertEquals(ProcStatus.SKIPPED, proc.find(msgId)!!.state) + assertEquals("resp-guard:no-open-req", proc.find(msgId)!!.lastError) + } } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt index 271cd72..66b3322 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt @@ -233,6 +233,34 @@ class ScheduleProcessorTest { assertTrue(log.entries.single().flags.contains(SnapshotFlag.SCHD_ABSENT_DELETED)) } + /** 带了时间条件的应答不是完整名单:回信里的航班整份替换,其余航班不删(`US-07` AC2)。 */ + @Test + fun `partial reply replaces flights in the message and leaves the others`() { + val flights = StubFlightState() + val events = StubMsgEvents() + val projection = StubFlightProjectionPort() + seed(flights, "122", LocalDate.of(2026, 12, 15)) + projection.snapshots["122"] = """{"flid":"122"}""" + flights.persistFullState( + com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot( + "121", LocalDate.of(2026, 12, 15), FlightState.ACTIVE, 1, + mapOf("SODT" to "15DEC261723", "REMC" to "old-note"), emptyMap(), + ), + msgId = 1, + now = java.time.Instant.now(), + ) + + val result = processor(flights = flights, events = events, projection = projection) + .applyScheduleRecords(head(), message(makeBody("121" to "15DEC261800")), deleteAbsent = false) + + assertEquals(ApplyResult.Succeeded, result) + assertEquals(FlightState.ACTIVE, flights.findMainRow("122")!!.state) + assertEquals(1L, flights.findMainRow("122")!!.stateVersion) + assertTrue(projection.snapshots.containsKey("122")) + assertFalse(events.rows.values.any { it.partitionKey == "122" }) + assertFalse(flights.loadFullSnapshot("121")!!.scalars.containsKey("REMC")) + } + /** 还没被日计划收录的航班(运营日为空)不属于任何覆盖范围,缺席清扫碰不到它。 */ @Test fun `flights without an operation day are out of every coverage range`() {