From 17a4bfe91b60340b7c72632c7742f765949cfdb0 Mon Sep 17 00:00:00 2001 From: windyboy Date: Sun, 13 Sep 2026 09:21:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(codec):=20=E8=A7=82=E6=B5=8B=E6=9C=AA?= =?UTF-8?q?=E8=90=BD=E5=BA=93=E7=9A=84=20SRVT/VIPF=20=E9=9B=86=E5=90=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit SRVT/VIPF 只保留在 wire/domain 并计数告警,不落明细表、不参与合并: 出现事实不再被静默丢弃,为 Q13 定案提供真实流量证据([G-SRVT-VIPF])。 - wire DTO:SRVT/SERVICEDATA、VIPF/VIPDATA 与嵌套 VIPT(OPER 为元素属性) - 出现即留键:缺席与"出现但为空"不再等价;已落库 10 类集合语义不变 - 新增 msgx.pipeline.codec.srvt_seen.total / vipf_seen.total(reference.md 已登记) - 一并提交此前的 MAFL 文档改动(INV-21/INV-22、flight-state §2.3、G-MAFL 措辞) --- docs/flight-state.md | 16 +++- docs/invariants.md | 8 +- docs/reference.md | 4 +- .../omms/msgexchange/codec/SisMessageBody.kt | 71 ++++++++++++++++++ .../omms/msgexchange/codec/SisWireMapper.kt | 48 +++++++++--- .../infra/metrics/PipelineCounters.kt | 21 +++++- .../infra/metrics/PipelineMetrics.kt | 12 +++ .../gzzn/omms/msgexchange/processing/Pump.kt | 11 +++ .../msgexchange/codec/JacksonXmlCodecTest.kt | 74 +++++++++++++++++++ .../domain/flight/FlightStateEngineTest.kt | 24 ++++++ .../infra/metrics/PipelineMetricsTest.kt | 12 +++ 11 files changed, 283 insertions(+), 18 deletions(-) diff --git a/docs/flight-state.md b/docs/flight-state.md index df0f94c..3eb8fcc 100644 --- a/docs/flight-state.md +++ b/docs/flight-state.md @@ -33,7 +33,7 @@ ### 2.2 字段与集合 -标量与异常对象前缀字段存于主表。协议中的 `SRVT`、`VIPF` 是无界集合,目标形态必须按集合完整保存到专用明细表示;当前 wire mapper 与持久化尚未实现 `[G-SRVT-VIPF]`。`MAFL` 不是 SIS/XML 入站字段,而是由共享航班的 `MAID`、`FLID`、`FLNO` 生成的主航班派生投影;当前尚未实现 `[G-MAFL]`。 +标量与异常对象前缀字段存于主表。协议中的 `SRVT`、`VIPF` 是无界集合,目标形态必须按集合完整保存到专用明细表示;当前只在 wire/domain 保留其出现事实与原始内容,专用明细、合并与投递尚未实现 `[G-SRVT-VIPF]`。`MAFL` 不是 SIS/XML 入站字段,而是由共享航班的 `MAID`、`FLID`、`FLNO` 生成的主航班派生投影;当前尚未实现 `[G-MAFL]`。 - `ORDINAL` 是持久化顺序,从 1 开始;`SOURCE_SEQ` 是上游序号,允许为空或重复。 - 相同资源号不代表同一条分配,禁止按资源号去重。 @@ -41,6 +41,17 @@ - ROUT 与 ERUT 是两类独立集合,不能因相同序号覆盖彼此。 - 主/共享关系以主表的 `MAID` 为事实来源:`MAID` 是共享航班指向主航班 `FLID` 的引用(非共享航班为 `NULL`);`MAFL` 只在读取和事件投影时从子航班事实派生,不按入站标量解析或保存。 +### 2.3 主/共享投影(`MAFL`) + +`MAFL` 是主航班的派生集合,元素为子航班的 `FLID` 与 `FLNO`;内容与变更传播分别由 `INV-21`、`INV-22` 保证。 + +- 子航班集合 = `STATE = ACTIVE` 且 `MAID = 主航班 FLID` 的 `FLIGHT_SCHD` 行;已 FDEL 的子航班(`STATE = DELETED`)自然退出投影,不需要改写主航班行。 +- 只有 `MAID` 为空的主航班携带 `MAFL`;共享航班只携带自身 `MAID`、`CSOP`、`CSFT`,不携带 `MAFL`,避免下游双向合并。 +- 投影按 `FLID` 升序,与到达顺序及 `FLNO` 变更无关:同一 `STATE_VERSION` 的投影逐字节稳定,重发与消费端比对才有意义。 +- `MAID = FLID` 的自引用行不进入任何 `MAFL`;`MAID` 指向不存在主航班的悬挂引用不阻断该子航班自身处理,只是不产生投影。 +- 子航班集合变化(新增、删除、`MAID` 迁移)必须让涉及的主航班在同一事务内推进 `STATE_VERSION` 并登记主航班事件(`KAFKA:msg` + `KAFKA:schd`);否则整态投影的只进不退写入会丢弃它(`design.md`「`schd` 聚合」)。共享航班自身不单独发通知。 +- 派生主航班投影与产生它的状态写入必须同一事务或一致读快照;按 `MAID` 取子航班要求该列有索引(`INV-17`)。 + ## 3. 合并与写入语义 领域决策逻辑(如 `FlightStateEngine` 及各类 Handler 规则)保持纯粹:它根据当前完整态和已解码报文,返回下一完整态与待发事件,不执行数据库或 Kafka I/O。`ScheduleProcessor` / `FlopProcessor` / `FdelProcessor` / `AdftProcessor` 是事务协调器,负责在统一事务边界内调用决策逻辑并持久化结果。 @@ -71,7 +82,7 @@ FDEL 是业务删除入口:仅在 `ACTIVE → DELETED` 时推进版本、保 ADFT 的字段缺失语义尚待上游确认。在确认前采用保守的 Set-only 规则:出现字段可更新,缺失字段不清空;不得把它当成日计划或动态全量替换。新建 ADFT 若带可解析的 `SODT`,按同一运营日规则计算 `OPERATION_DAY`;否则保留为 `NULL`。 -主/共享航班级联:删除共享航班时更新主航班 `MAFL` 并向主航班通知;删除主航班时级联删除其子共享关联并发出删除通知;主/共享关系必须一次原子变更,不出现主已删、子残留的半状态。共享航班增量通常只更新并通知主航班,不直接发共享通知。这些语义同样约束 FDEL 之外的生命周期清理。主/共享关联的增删按 `FLID` 做值比较,不使用引用比较。 +主/共享航班级联:删除共享航班时重算主航班 `MAFL`(见「主/共享投影」)并向主航班通知;删除主航班时级联删除其子共享关联并发出删除通知;主/共享关系必须一次原子变更,不出现主已删、子残留的半状态。共享航班增量通常只更新并通知主航班,不直接发共享通知。这些语义同样约束 FDEL 之外的生命周期清理。主/共享关联的增删按 `FLID` 做值比较,不使用引用比较。 SIS 规定删除主航班时必须先删子共享航班、再删主航班,顺序不符时 RMS 应向 AODB 回发 EROR(`SIS_AODB_RMS-V0.1.md` §1.6.1-1.d,事件定义见 SIS §4.8)。本文的原子级联不发该回报,两者取舍见 Q14。 @@ -106,6 +117,7 @@ SIS 规定删除主航班时必须先删子共享航班、再删主航班,顺 - ADFT 缺失字段和 `FLID` 重用的上游语义(`FLID` 复用见 `Q16`); - 日计划缺失可选字段的删除语义(与 SIS 的冲突见本文件「合并与写入语义」,`Q13`); - 主/共享删除顺序与 EROR 回报义务(Q14); +- `MAFL` 投影是否出现在查询视图,随 `Q3`;字段名与空集合表示随 `Q4`; - 未归属 `OPERATION_DAY` 航班的终止与保留策略; - ROUT/ERUT 联合主键迁移; - 批量一致性读取,以及动态事件仅在事务内计算一次的收敛。 diff --git a/docs/invariants.md b/docs/invariants.md index ae6e6ec..ecc4bf9 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -48,6 +48,8 @@ - **INV-18** 航班表的写者集合是「主泵处理器」与「历史清理」;两者必须互斥(同一 `PIPELINE_LOCK`,或清理在同一事务内复查判据后再删除),不得出现清理删除与处理器更新同一 `FLID` 的竞态。 - **INV-19** 整包校验失败或运营日冲突时整包不落地,既有状态与版本保持不变。 - **INV-20** 处理器幂等:同一消息重复执行只产生一次业务效果。身份唯一只防「重复记录」,不防「重新执行」;29 类 FLOP 幂等矩阵补全前,本条**不可声明**。`[G-FLOP-IDEMPOTENT]` +- **INV-21** `MAFL` 是派生投影:内容恒等于「`STATE = ACTIVE` 且 `MAID = 主航班 FLID`」的子航班集合(元素 `FLID` + `FLNO`,按 `FLID` 升序),不落库、不从入站解析;自引用与悬挂引用不入投影。 +- **INV-22** 子航班集合变化必须使涉及的主航班在同一事务内推进 `STATE_VERSION` 并登记主航班事件;投影只进不退,版本不推进即被下游丢弃。 ## 3. 声明边界 @@ -92,6 +94,8 @@ | INV-17 | 业务型终态四件套同事务 | 待核对;真实 PG 用例待补(ACM2-39) | | INV-18 | 清理与处理并发 | `HistorySweepJobTest`(归档后被主泵更新的航班不删除、不发 tombstone)+ `HistorySweepPurgePgTest`(删除阶段失败时 tombstone 与删除整体回滚) | | INV-20 / CLM-3 | 重放同一条消息 | 缺口:29 类 FLOP 幂等矩阵未补全 | +| INV-21 | `MAFL` 投影与 `ACTIVE` 子航班集合一致(子航班删除后退出、自引用与悬挂引用不入、顺序确定) | 缺口:投影未实现(`[G-MAFL]`) | +| INV-22 | 子航班新增、删除、`MAID` 迁移时主航班版本与事件 | 缺口:主/共享级联未实现(`[G-MAFL]`) | | CLM-4 | 放弃行与清除前提 | 断言放弃行不写标记、不被当作已打标(关联 ACM2-36) | | CLM-9 | 回填/积压完成时限 | 指标已就位:`msgx.pipeline.job.heartbeat_age_seconds` / `ticks.total` / `failures.total` / `last_sweep_selected` 与 `msgx.pipeline.backfill.oldest_unmarked_seconds`(关联 ACM2-38);实际延迟仍需现场数据,CLM-9 不可声明 | | — | 请求超时、无匹配 RESP、时间单位不一致 | 不误用迟到应答、不提前完成请求 | @@ -112,8 +116,8 @@ | `G-BACKFILL-BACKOFF` | 回填独立退避键(`backfill-backoff-ms` / `-cap-ms`)未实现,当前为代码内硬编码(取值见 reference) | 回填重试节奏 | | `G-KAFKA-D3` | `kafka.producers.default.max-in-flight` 与 D3 要求的 1 不一致(取值见 reference) | 投递幂等前提 | | `G-REPLAY-CHANNEL` | 「打标即清除」语义下的独立原文保留通道未设计 | CLM-5 | -| `G-MAFL` | 主航班 `MAFL` 派生投影及主/共享原子级联未实现;`MAFL` 不是 SIS/XML 入站字段 | 航班完整态;删除与重建 | -| `G-SRVT-VIPF` | SIS/XML 的 `SRVT`、`VIPF` 无界集合尚未映射到 wire/domain/持久化明细 | 航班完整态;无损字段保存 | +| `G-MAFL` | 主航班 `MAFL` 派生投影及主/共享原子级联未实现(规则见 `INV-21`/`INV-22`);`MAFL` 不是 SIS/XML 入站字段 | 航班完整态;删除与重建 | +| `G-SRVT-VIPF` | SIS/XML 的 `SRVT`、`VIPF` 无界集合尚未映射到持久化明细;wire/domain 只保留出现事实与原始内容,不参与合并与投递(清空语义见 `Q13`) | 航班完整态;无损字段保存 | | `G-COMPAT-HTTP` | compat 入口仍未实现 Q3 定案后的 ResponseDto、媒体类型、字符集、失败响应与请求体上限 | `C-28`;US-02 | | `G-REQ-OPEN-UNIQUE` | `REQ_TRACK` 尚无约束开放态 `(REQ_TYPE, OPERATION_DAY, SENDER)` 唯一性的部分索引 | US-08;`G-REQ-TRACK` | diff --git a/docs/reference.md b/docs/reference.md index 61ae39b..5a62537 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -83,8 +83,10 @@ | `msgx.pipeline.job.ticks.total` | 作业 tick 完成次数 | 不增长 → 作业停摆 | | `msgx.pipeline.job.failures.total` | 作业 tick 抛错次数 | 增长 → 扫描/历史作业异常 | | `msgx.pipeline.job.last_sweep_selected` | 上一轮回填扫描选中的待办条数(扫描积压) | 持续顶到批次上限 → 扫描吃不消 | +| `msgx.pipeline.codec.srvt_seen.total` | 入站记录中出现 `SRVT` 段的条数(尚未落明细表,`[G-SRVT-VIPF]`) | > 0 → 真实流量确有该段,按真实报文定案 `Q13` | +| `msgx.pipeline.codec.vipf_seen.total` | 入站记录中出现 `VIPF` 段的条数(尚未落明细表,`[G-SRVT-VIPF]`) | 同上 | -取数规则:统一走 `BacklogSnapshotProvider`(`PARAM:msgx.health.backlog-cache-ttl-ms`),`/health` 与 `/metrics` 共用同一快照——`backlog()` 是 `PROC_STATE` 的全表聚合,不能被高频抓取打穿;**无法取数上报 `NaN`,无可比记录的年龄/滞后类仪表上报 `-1`,都不伪造 0**。日志出口故障不得阻塞业务线程。 +取数规则:统一走 `BacklogSnapshotProvider`(`PARAM:msgx.health.backlog-cache-ttl-ms`),`/health` 与 `/metrics` 共用同一快照——`backlog()` 是 `PROC_STATE` 的全表聚合,不能被高频抓取打穿;**无法取数上报 `NaN`,无可比记录的年龄/滞后类仪表上报 `-1`,都不伪造 0**。日志出口故障不得阻塞业务线程。进程内计数类指标不经快照,重启归零。 作业健康:回填的唯一驱动是扫描作业,因此作业存活必须独立可观测——`msgx.pipeline.job.heartbeat_age_seconds` / `ticks.total` / `failures.total` 是作业心跳,`msgx.pipeline.job.last_sweep_selected` 是扫描积压,`msgx.pipeline.backfill.oldest_unmarked_seconds` 是实际回填延迟;作业线程停摆由 `/health` 的作业指示器判 `DOWN`(心跳超过 `3 × 扫描周期`,周期是代码常量)。 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 15e0c94..71dde51 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt @@ -19,6 +19,32 @@ data class ScheduleBody( val records: List, ) +/** + * 只保留在 wire/domain、尚未映射到持久化明细的集合键(`[G-SRVT-VIPF]`)。 + * + * 它们不参与合并、不进快照、不落库;保留的目的是不让入站事实在解码层被静默抹平,并让真实 + * 流量里的出现情况可观测——决定"缺席是否等于删除"的 `Q13` 需要真实报文才能定案。 + */ +private val UNPERSISTED_COLLECTION_KEYS: Set = setOf("SRVT", "VIPF") + +/** + * 一条解码载荷里各未落库集合命中的**记录数**(`[G-SRVT-VIPF]`)。 + * + * 按记录计(一条记录带该段即算 1),不按段内元素计;只返回命中项。载荷类型不在 + * [ScheduleBody] / [FlopPayload] 之内时返回空——调用方不得据此改变处理结果。 + */ +fun unpersistedCollectionHits(body: Any?): Map { + val perRecord = when (body) { + is ScheduleBody -> body.records.map { it.collections } + is FlopPayload -> listOf(body.collections) + else -> return emptyMap() + } + if (perRecord.isEmpty()) return emptyMap() + return UNPERSISTED_COLLECTION_KEYS + .mapNotNull { key -> perRecord.count { key in it }.takeIf { it > 0 }?.let { key to it } } + .toMap() +} + /** * Wire DTOs for SIS in docs/legacy/SIS_AODB_RMS-V0.1.md. * SIS permits a subsystem to ignore standard fields it does not use; those fields are @@ -105,6 +131,10 @@ data class FlightRecordXml( @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "CHOT") val chot: List = emptyList(), @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ROUT") val rout: List = emptyList(), @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ERUT") val erut: List = emptyList(), + // SRVT/VIPF 用可空表达"段是否出现":null = 未出现;出现即为列表(空元素得到一行空行)。 + // 两者都不落明细表、不参与合并,清空语义待 `Q13`(`[G-SRVT-VIPF]`)。 + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "SRVT") val srvt: List? = null, + @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "VIPF") val vipf: List? = null, ) @JsonIgnoreProperties(ignoreUnknown = true) @@ -198,3 +228,44 @@ data class RoutXml( @param:JacksonXmlProperty(localName = "SCAT") val scat: String? = null, @param:JacksonXmlProperty(localName = "SCDT") val scdt: String? = null, ) + +/** + * SIS `SRVT`/`SERVICEDATA`(SCHD)与 `OPT_SERVICEDATA`(FLOP):一次服务明细,可重复出现。 + * + * 字段形态以 `docs/legacy/unisysaodbsis.xsd` 为准。`OPER` 是元素属性,SIS 正文未记载其值域; + * `SANR`/`SARR` 只出现在 FLOP 的 `OPT_SERVICEDATA`。协议未给序号属性,故没有源序号。 + */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class SrvtXml( + @param:JacksonXmlProperty(isAttribute = true, localName = "OPER") val oper: String? = null, + @param:JacksonXmlProperty(localName = "SRTC") val srtc: String? = null, + @param:JacksonXmlProperty(localName = "SRQT") val srqt: String? = null, + @param:JacksonXmlProperty(localName = "SRST") val srst: String? = null, + @param:JacksonXmlProperty(localName = "SRET") val sret: String? = null, + @param:JacksonXmlProperty(localName = "SRPR") val srpr: String? = null, + @param:JacksonXmlProperty(localName = "SANR") val sanr: String? = null, + @param:JacksonXmlProperty(localName = "SARR") val sarr: String? = null, +) + +/** + * SIS `VIPF`/`VIPDATA`(SCHD)与 `OPT_VIPDATA`(FLOP):一位 VIP,可重复出现。 + * + * `VIPT` 按 XSD 至多一个;SIS 正文"每个 VIP 重复"的说法与之冲突,形态以 XSD 为准。 + */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class VipfXml( + @param:JacksonXmlProperty(isAttribute = true, localName = "OPER") val oper: String? = null, + @param:JacksonXmlProperty(localName = "VPCD") val vpcd: String? = null, + @param:JacksonXmlProperty(localName = "VFES") val vfes: String? = null, + @param:JacksonXmlProperty(localName = "VIPT") val vipt: ViptXml? = null, +) + +/** SIS `VIPT`/`VIPTXNDATA`:VIP 关联的服务交易明细,嵌套在 [VipfXml] 内。 */ +@JsonIgnoreProperties(ignoreUnknown = true) +data class ViptXml( + @param:JacksonXmlProperty(isAttribute = true, localName = "OPER") val oper: String? = null, + @param:JacksonXmlProperty(localName = "VSCD") val vscd: String? = null, + @param:JacksonXmlProperty(localName = "VTQY") val vtqy: String? = null, + @param:JacksonXmlProperty(localName = "VTST") val vtst: String? = null, + @param:JacksonXmlProperty(localName = "VTET") val vtet: String? = null, +) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt index 0c33f70..ac189b1 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt @@ -33,18 +33,42 @@ internal object SisWireMapper { ).forEach { (key, value) -> value?.let { put(key, it.trim()) } } } - private fun FlightRecordXml.collections(): Map>> = linkedMapOf( - "GTDT" to gtdt.map { it.toMap("GTNO" to it.gtno, "GATE" to it.gate, "PGOT" to it.pgot, "PGCT" to it.pgct, "GOTM" to it.gotm, "GCTM" to it.gctm, "GTYP" to it.gtyp) }, - "CKDT" to ckdt.map { it.toMap("CKNO" to it.ckno, "CHKC" to it.chkc, "CCLS" to it.ccls, "PCOT" to it.pcot, "PCCT" to it.pcct, "COTM" to it.cotm, "CCTM" to it.cctm, "CTYP" to it.ctyp) }, - "CLDT" to cldt.map { it.toMap("CLNO" to it.clno, "BELT" to it.belt, "BCLS" to it.bcls, "PCOT" to it.pcot, "PCCT" to it.pcct, "FBAG" to it.fbag, "LBAG" to it.lbag, "BTYP" to it.btyp) }, - "PSDT" to psdt.map { it.toMap("PSNO" to it.psno, "PSST" to it.psst, "STST" to it.stst, "STET" to it.stet) }, - "CHDT" to chdt.map { it.toMap("CHNO" to it.chno, "CHUT" to it.chut, "CHCLS" to it.chcls, "PCBT" to it.pcbt, "PCET" to it.pcet, "CBTM" to it.cbtm, "CETM" to it.cetm, "CHTYP" to it.chtyp) }, - "DELY" to dely.map { it.toMap("CODE" to it.code, "STRT" to it.strt, "DURA" to it.dura, "REMC" to it.text) }, - "ABTM" to abtm.map { it.toMap("ASNO" to it.asno, "ABDG" to it.abdg, "ABOP" to it.abop, "AOTM" to it.aotm) }, - "CHOT" to chot.map { it.toMap("CSNO" to it.csno, "CHID" to it.chid, "CHST" to it.chst, "CHTM" to it.chtm) }, - "ROUT" to rout.map { it.toMap("RTNO" to it.rtno, "APCD" to it.apcd, "SCAT" to it.scat, "SCDT" to it.scdt) }, - "ERUT" to erut.map { it.toMap("RTNO" to it.rtno, "APCD" to it.apcd, "SCAT" to it.scat, "SCDT" to it.scdt) }, - ).filterValues { it.isNotEmpty() } + /** + * 已落库的 10 类集合:缺席(字段默认空列表)不产生键,出现但为空的元素得到一行空行 `[{}]`; + * `filterValues` 去掉的正是"缺席",避免整包凭空清空本地明细(合并语义见 flight-state.md §3.1)。 + * + * `SRVT`/`VIPF` 尚未有明细表(`[G-SRVT-VIPF]`):只用"键是否存在"表达段是否出现,保留原始 + * 内容与顺序,不参与合并、不判断清空语义(`Q13`)——出现(哪怕为空)与缺席不再被抹平。 + */ + private fun FlightRecordXml.collections(): Map>> { + val mapped = linkedMapOf( + "GTDT" to gtdt.map { it.toMap("GTNO" to it.gtno, "GATE" to it.gate, "PGOT" to it.pgot, "PGCT" to it.pgct, "GOTM" to it.gotm, "GCTM" to it.gctm, "GTYP" to it.gtyp) }, + "CKDT" to ckdt.map { it.toMap("CKNO" to it.ckno, "CHKC" to it.chkc, "CCLS" to it.ccls, "PCOT" to it.pcot, "PCCT" to it.pcct, "COTM" to it.cotm, "CCTM" to it.cctm, "CTYP" to it.ctyp) }, + "CLDT" to cldt.map { it.toMap("CLNO" to it.clno, "BELT" to it.belt, "BCLS" to it.bcls, "PCOT" to it.pcot, "PCCT" to it.pcct, "FBAG" to it.fbag, "LBAG" to it.lbag, "BTYP" to it.btyp) }, + "PSDT" to psdt.map { it.toMap("PSNO" to it.psno, "PSST" to it.psst, "STST" to it.stst, "STET" to it.stet) }, + "CHDT" to chdt.map { it.toMap("CHNO" to it.chno, "CHUT" to it.chut, "CHCLS" to it.chcls, "PCBT" to it.pcbt, "PCET" to it.pcet, "CBTM" to it.cbtm, "CETM" to it.cetm, "CHTYP" to it.chtyp) }, + "DELY" to dely.map { it.toMap("CODE" to it.code, "STRT" to it.strt, "DURA" to it.dura, "REMC" to it.text) }, + "ABTM" to abtm.map { it.toMap("ASNO" to it.asno, "ABDG" to it.abdg, "ABOP" to it.abop, "AOTM" to it.aotm) }, + "CHOT" to chot.map { it.toMap("CSNO" to it.csno, "CHID" to it.chid, "CHST" to it.chst, "CHTM" to it.chtm) }, + "ROUT" to rout.map { it.toMap("RTNO" to it.rtno, "APCD" to it.apcd, "SCAT" to it.scat, "SCDT" to it.scdt) }, + "ERUT" to erut.map { it.toMap("RTNO" to it.rtno, "APCD" to it.apcd, "SCAT" to it.scat, "SCDT" to it.scdt) }, + ).filterValues { it.isNotEmpty() }.toMutableMap() + srvt?.let { mapped["SRVT"] = it.map(::srvtRow) } + vipf?.let { mapped["VIPF"] = it.map(::vipfRow) } + return mapped + } + + /** `SRVT` 一行:`OPER` 取自属性;`VIPT_*` 前缀供嵌套字段使用,避免与外层 `OPER` 撞键。 */ + private fun srvtRow(x: SrvtXml): Map = x.toMap( + "OPER" to x.oper, "SRTC" to x.srtc, "SRQT" to x.srqt, "SRST" to x.srst, + "SRET" to x.sret, "SRPR" to x.srpr, "SANR" to x.sanr, "SARR" to x.sarr, + ) + + private fun vipfRow(x: VipfXml): Map = x.toMap( + "OPER" to x.oper, "VPCD" to x.vpcd, "VFES" to x.vfes, + "VIPT_OPER" to x.vipt?.oper, "VIPT_VSCD" to x.vipt?.vscd, "VIPT_VTQY" to x.vipt?.vtqy, + "VIPT_VTST" to x.vipt?.vtst, "VIPT_VTET" to x.vipt?.vtet, + ) private fun Any.toMap(vararg values: Pair): Map = values.mapNotNull { (key, value) -> value?.trim()?.let { key to it } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineCounters.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineCounters.kt index 4633df5..6010f9c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineCounters.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineCounters.kt @@ -1,6 +1,7 @@ package com.gzzn.omms.msgexchange.infra.metrics import jakarta.inject.Singleton +import java.util.concurrent.atomic.AtomicLong /** * 管道运行期的**进程内**计数。 @@ -11,4 +12,22 @@ import jakarta.inject.Singleton * 注意:计数在重启后归零。需要跨重启的累计值应由指标后端聚合,不在这里做持久化。 */ @Singleton -class PipelineCounters +class PipelineCounters { + private val srvtSeen = AtomicLong(0) + private val vipfSeen = AtomicLong(0) + + /** + * 入站记录里出现 `SRVT`/`VIPF` 段的条数(`[G-SRVT-VIPF]`)。 + * + * 这两个集合目前只保留在 wire/domain,不落明细表、不参与合并;计数是"真实报文有没 + * 有在用"的唯一取证渠道(清空语义 `Q13` 需要真实样例才能定案)。> 0 表示确有流量携带该段。 + */ + fun unpersistedCollectionSeenAdd(hits: Map) { + hits["SRVT"]?.let { srvtSeen.addAndGet(it.toLong()) } + hits["VIPF"]?.let { vipfSeen.addAndGet(it.toLong()) } + } + + fun srvtSeenCount(): Long = srvtSeen.get() + + fun vipfSeenCount(): Long = vipfSeen.get() +} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetrics.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetrics.kt index 368dbaa..35d0261 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetrics.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetrics.kt @@ -26,6 +26,8 @@ import java.time.Duration * - `msgx.pipeline.job.last_failure_age_seconds`:距最近一次作业 tick 失败的秒数(从未失败为 -1) * - `msgx.pipeline.job.ticks.total` / `msgx.pipeline.job.failures.total`:作业 tick 完成/抛错次数 * - `msgx.pipeline.job.last_sweep_selected`:上一轮回填扫描选中的待办条数(扫描积压) + * - `msgx.pipeline.codec.srvt_seen.total` / `msgx.pipeline.codec.vipf_seen.total`:入站记录里出现 + * `SRVT`/`VIPF` 段的条数(尚未落明细表,`[G-SRVT-VIPF]`;> 0 表示真实流量确有该段) * * 取数统一走 [BacklogSnapshotProvider](30 秒 TTL),因此指标抓取不会打穿数据库。 * 无法取数时以 `NaN` 上报(Micrometer 的惯例表示"本次无值"),而不是伪造 0。 @@ -45,6 +47,7 @@ class PipelineMetrics( private val mailbox: BeanProvider, private val activity: JobActivity, private val clock: Clock, + private val counters: PipelineCounters, ) { @PostConstruct @@ -89,6 +92,15 @@ class PipelineMetrics( Gauge.builder("msgx.pipeline.job.last_sweep_selected", activity) { it.snapshot().lastSweepSelected.toDouble() } .strongReference(true) .register(registry) + + // SRVT/VIPF 尚未落明细表([G-SRVT-VIPF]):计数替代静默丢弃,为 Q13 提供真实流量证据。 + Gauge.builder("msgx.pipeline.codec.srvt_seen.total", counters) { it.srvtSeenCount().toDouble() } + .strongReference(true) + .register(registry) + + Gauge.builder("msgx.pipeline.codec.vipf_seen.total", counters) { it.vipfSeenCount().toDouble() } + .strongReference(true) + .register(registry) } private fun backlogGauge(name: String, value: (com.gzzn.omms.msgexchange.infra.persistence.Backlog) -> Double) { 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 ad6a5b5..7056af6 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -3,6 +3,7 @@ package com.gzzn.omms.msgexchange.processing import com.gzzn.omms.msgexchange.codec.FlopPayload import com.gzzn.omms.msgexchange.codec.ScheduleBody import com.gzzn.omms.msgexchange.codec.XmlCodec +import com.gzzn.omms.msgexchange.codec.unpersistedCollectionHits import com.gzzn.omms.msgexchange.config.OperationDayProps import com.gzzn.omms.msgexchange.config.PipelineProps import com.gzzn.omms.msgexchange.domain.DecodedMessage @@ -11,6 +12,7 @@ import com.gzzn.omms.msgexchange.domain.MsgKind import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.infra.log.TraceLog +import com.gzzn.omms.msgexchange.infra.metrics.PipelineCounters import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository import com.gzzn.omms.msgexchange.infra.persistence.InboxCursorRepository import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository @@ -150,6 +152,7 @@ class MessageProcessor( private val props: PipelineProps, private val clock: Clock, private val operationDayProps: OperationDayProps, + private val counters: PipelineCounters, ) { private val log = org.slf4j.LoggerFactory.getLogger(MessageProcessor::class.java) @@ -187,6 +190,14 @@ class MessageProcessor( } } + // [G-SRVT-VIPF]:SRVT/VIPF 段只保留在解码载荷里,尚未落明细表(清空语义待 Q13)。 + // 计数 + 告警替代此前的静默丢弃;出现即证明真实报文携带该段,可作为定案依据。 + val unpersisted = unpersistedCollectionHits(decoded.body) + if (unpersisted.isNotEmpty()) { + counters.unpersistedCollectionSeenAdd(unpersisted) + log.warn("unpersisted collection(s) {} present msgId={} [G-SRVT-VIPF]", unpersisted, head.msgId) + } + // I3:identity 仅首次绑定(head.identityKey == null);FAILED 重试不重绑 if (head.identityKey == null) { val identity = Identity.of( 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 259187b..f8767b2 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt @@ -2,6 +2,7 @@ package com.gzzn.omms.msgexchange.codec import com.gzzn.omms.msgexchange.domain.MsgKind import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse import org.junit.jupiter.api.Assertions.assertInstanceOf import org.junit.jupiter.api.Assertions.assertNull import org.junit.jupiter.api.Assertions.assertTrue @@ -87,6 +88,79 @@ class JacksonXmlCodecTest { assertEquals(listOf(mapOf("GTNO" to "0")), body.collections["GTDT"]) } + @Test + fun `SRVT and VIPF sections keep order attributes and nested VIPT`() { + val raw = """ + + AODB720021010090311SCHDDNLD + + 1 + + 121112312 + CACA10115DEC261723 + SR0129001 + SR02 + V0013 + VS01415DEC262000 + + + + + """.trimIndent() + + val result = codec.decode(raw) as DecodeResult.Ok + val record = (result.message.body as ScheduleBody).records.single() + + assertEquals(2, record.collections["SRVT"]!!.size) + assertEquals("ADD", record.collections["SRVT"]!![0]["OPER"]) // OPER 是元素属性 + assertEquals("SR01", record.collections["SRVT"]!![0]["SRTC"]) + assertEquals("2", record.collections["SRVT"]!![0]["SRQT"]) + assertEquals("9001", record.collections["SRVT"]!![0]["SRPR"]) + assertEquals("DEL", record.collections["SRVT"]!![1]["OPER"]) // 保留报文先后顺序 + assertEquals("SR02", record.collections["SRVT"]!![1]["SRTC"]) + + val vip = record.collections["VIPF"]!!.single() + assertEquals("V001", vip["VPCD"]) + assertEquals("3", vip["VFES"]) + assertEquals("UPD", vip["VIPT_OPER"]) + assertEquals("VS01", vip["VIPT_VSCD"]) + assertEquals("4", vip["VIPT_VTQY"]) + assertEquals("15DEC262000", vip["VIPT_VTET"]) + } + + @Test + fun `absent and empty SRVT VIPF stay distinguishable while mapped collections keep dropping empties`() { + val withEmpty = """ + + AODB820021010090311SCHDDNLD + 1121 + + """.trimIndent() + val present = ((codec.decode(withEmpty) as DecodeResult.Ok).message.body as ScheduleBody) + .records.single().collections + + // 出现但为空段保留键,不再与"未出现"混为一谈 [G-SRVT-VIPF] + assertTrue(present.containsKey("SRVT")) + assertTrue(present.containsKey("VIPF")) + assertTrue(present.getValue("SRVT").all { it.isEmpty() }) + assertTrue(present.getValue("VIPF").all { it.isEmpty() }) + // 已落库的 10 类集合行为未变:空元素 → 一行空行;缺席 → 没有这个键 + assertEquals(listOf(emptyMap()), present["CLDT"]) + + val absent = """ + + AODB920021010090311SCHDDNLD + 1121 + + """.trimIndent() + val absentKeys = ((codec.decode(absent) as DecodeResult.Ok).message.body as ScheduleBody) + .records.single().collections + + assertFalse(absentKeys.containsKey("SRVT")) + assertFalse(absentKeys.containsKey("VIPF")) + assertFalse(absentKeys.containsKey("CLDT")) + } + @Test fun `malformed xml returns MALFORMED not CODEC_ERROR`() { val err = codec.decode("not-xml") as DecodeResult.Err diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngineTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngineTest.kt index 368c663..b71c4b1 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngineTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightStateEngineTest.kt @@ -111,6 +111,30 @@ class FlightStateEngineTest { assertEquals(2, next.stateVersion) } + @Test + fun `unpersisted collections are carried but never merged until their detail tables exist`() { + val current = FlightSnapshot( + "121", day, FlightState.ACTIVE, 1, + scalars = mapOf("FLNO" to "CA001"), + collections = mapOf("GTDT" to listOf(mapOf("GATE" to "G1"))), + ) + val next = FlightStateEngine.snapshotState( + current, + ScheduleRecord( + "121", + scalars = mapOf("FLNO" to "CA002"), + collections = mapOf("SRVT" to listOf(mapOf("SRTC" to "SR01")), "VIPF" to emptyList()), + ), + operationDay = day, + keepDeleted = false, + ) + + // [G-SRVT-VIPF]:wire/domain 保留这两个集合,但合并层既不落库,也不因"出现"清空任何明细 + assertFalse(next.collections.containsKey("SRVT")) + assertFalse(next.collections.containsKey("VIPF")) + assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"]) + } + @Test fun `validation accepts valid snapshot and rejects RECS mismatch`() { val ok = FlightStateEngine.validateMessage( diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetricsTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetricsTest.kt index ddfaac6..a0b23a1 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetricsTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/metrics/PipelineMetricsTest.kt @@ -62,6 +62,18 @@ class PipelineMetricsTest { org.junit.jupiter.api.Assertions.assertTrue(gauge("msgx.pipeline.job.heartbeat_age_seconds") >= 0.0) } + @Test + fun `unpersisted SRVT VIPF counters are registered and reflect decoded hits`() { + val counters = ctx.getBean(PipelineCounters::class.java) + val srvtBefore = gauge("msgx.pipeline.codec.srvt_seen.total") + val vipfBefore = gauge("msgx.pipeline.codec.vipf_seen.total") + + counters.unpersistedCollectionSeenAdd(mapOf("SRVT" to 2, "VIPF" to 1)) + + assertEquals(srvtBefore + 2.0, gauge("msgx.pipeline.codec.srvt_seen.total"), 0.001) + assertEquals(vipfBefore + 1.0, gauge("msgx.pipeline.codec.vipf_seen.total"), 0.001) + } + private fun gauge(name: String): Double = requireNotNull(registry.find(name).gauge()) { "gauge not registered: $name" }.value() }