feat(codec): 观测未落库的 SRVT/VIPF 集合

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 措辞)
This commit is contained in:
windyboy
2026-09-13 09:21:51 +08:00
parent b8738554e6
commit 17a4bfe91b
11 changed files with 283 additions and 18 deletions
+14 -2
View File
@@ -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 联合主键迁移;
- 批量一致性读取,以及动态事件仅在事务内计算一次的收敛。
+6 -2
View File
@@ -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` |
+3 -1
View File
@@ -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 × 扫描周期`,周期是代码常量)。
@@ -19,6 +19,32 @@ data class ScheduleBody(
val records: List<ScheduleRecord>,
)
/**
* 只保留在 wire/domain、尚未映射到持久化明细的集合键(`[G-SRVT-VIPF]`)。
*
* 它们不参与合并、不进快照、不落库;保留的目的是不让入站事实在解码层被静默抹平,并让真实
* 流量里的出现情况可观测——决定"缺席是否等于删除"的 `Q13` 需要真实报文才能定案。
*/
private val UNPERSISTED_COLLECTION_KEYS: Set<String> = setOf("SRVT", "VIPF")
/**
* 一条解码载荷里各未落库集合命中的**记录数**`[G-SRVT-VIPF]`)。
*
* 按记录计(一条记录带该段即算 1),不按段内元素计;只返回命中项。载荷类型不在
* [ScheduleBody] / [FlopPayload] 之内时返回空——调用方不得据此改变处理结果。
*/
fun unpersistedCollectionHits(body: Any?): Map<String, Int> {
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<ChotXml> = emptyList(),
@param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ROUT") val rout: List<RoutXml> = emptyList(),
@param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "ERUT") val erut: List<RoutXml> = emptyList(),
// SRVT/VIPF 用可空表达"段是否出现":null = 未出现;出现即为列表(空元素得到一行空行)。
// 两者都不落明细表、不参与合并,清空语义待 `Q13``[G-SRVT-VIPF]`)。
@param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "SRVT") val srvt: List<SrvtXml>? = null,
@param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "VIPF") val vipf: List<VipfXml>? = 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,
)
@@ -33,18 +33,42 @@ internal object SisWireMapper {
).forEach { (key, value) -> value?.let { put(key, it.trim()) } }
}
private fun FlightRecordXml.collections(): Map<String, List<Map<String, String>>> = 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<String, List<Map<String, String>>> {
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<String, String> = 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<String, String> = 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<String, String?>): Map<String, String> = values.mapNotNull { (key, value) ->
value?.trim()?.let { key to it }
@@ -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<String, Int>) {
hits["SRVT"]?.let { srvtSeen.addAndGet(it.toLong()) }
hits["VIPF"]?.let { vipfSeen.addAndGet(it.toLong()) }
}
fun srvtSeenCount(): Long = srvtSeen.get()
fun vipfSeenCount(): Long = vipfSeen.get()
}
@@ -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<CminmsgInboxRepository>,
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) {
@@ -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)
}
// I3identity 仅首次绑定(head.identityKey == null);FAILED 重试不重绑
if (head.identityKey == null) {
val identity = Identity.of(
@@ -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 = """
<MSG>
<META><SNDR>AODB</SNDR><SEQN>7</SEQN><DTTM>20021010090311</DTTM><TYPE>SCHD</TYPE><STYP>DNLD</STYP></META>
<SCHD>
<RECS>1</RECS>
<FLTR>
<FLID>121112312</FLID>
<ALCD>CA</ALCD><FLNO>CA101</FLNO><SODT>15DEC261723</SODT>
<SRVT OPER="ADD"><SRTC>SR01</SRTC><SRQT>2</SRQT><SRPR>9001</SRPR></SRVT>
<SRVT OPER="DEL"><SRTC>SR02</SRTC></SRVT>
<VIPF OPER="ADD"><VPCD>V001</VPCD><VFES>3</VFES>
<VIPT OPER="UPD"><VSCD>VS01</VSCD><VTQY>4</VTQY><VTET>15DEC262000</VTET></VIPT>
</VIPF>
</FLTR>
</SCHD>
</MSG>
""".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 = """
<MSG>
<META><SNDR>AODB</SNDR><SEQN>8</SEQN><DTTM>20021010090311</DTTM><TYPE>SCHD</TYPE><STYP>DNLD</STYP></META>
<SCHD><RECS>1</RECS><FLTR><FLID>121</FLID><SRVT/><VIPF></VIPF><CLDT/></FLTR></SCHD>
</MSG>
""".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<String, String>()), present["CLDT"])
val absent = """
<MSG>
<META><SNDR>AODB</SNDR><SEQN>9</SEQN><DTTM>20021010090311</DTTM><TYPE>SCHD</TYPE><STYP>DNLD</STYP></META>
<SCHD><RECS>1</RECS><FLTR><FLID>121</FLID></FLTR></SCHD>
</MSG>
""".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
@@ -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(
@@ -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()
}