From 23b74ed55416ba6cce1bfab9fe3ba5ebcfa0f488 Mon Sep 17 00:00:00 2001 From: windyboy Date: Mon, 21 Sep 2026 14:18:40 +0800 Subject: [PATCH] =?UTF-8?q?fix(processing):=20Wave1=20=E2=80=94=20IgnoreRu?= =?UTF-8?q?les/ROUT=E6=88=AA=E6=96=AD/Q9=E6=96=87=E6=A1=A3/schd=20FLTR=20?= =?UTF-8?q?=E6=95=B0=E7=BB=84=EF=BC=88ACM2-94/88/81/96=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit REGN/RSTA/EROR 退出忽略清单;ROUT/ERUT 升序留 4 且不落 SCAT/SCDT;C-1/Q9 统一库方清除;flushSchd 聚合成无 key 的 FLTR 数组。 Co-authored-by: Cursor --- docs/implementation.md | 6 +- docs/reference.md | 2 +- docs/requirements.md | 4 +- docs/specification.md | 4 +- .../omms/msgexchange/codec/JacksonXmlCodec.kt | 2 + .../omms/msgexchange/delivery/Dispatcher.kt | 45 +++++---- .../omms/msgexchange/domain/DecodedMessage.kt | 8 +- .../domain/flight/FlightStateEngine.kt | 26 ++++- .../infra/kafka/KafkaDeliveryPort.kt | 5 +- .../persistence/jdbc/JdbcPgRepositories.kt | 28 +++++- .../msgexchange/infra/stub/StubAdapters.kt | 6 +- .../msgexchange/processing/IgnoreRules.kt | 10 +- .../gzzn/omms/msgexchange/processing/Pump.kt | 5 + .../delivery/DispatcherTickTest.kt | 33 ++++--- .../domain/flight/FlightStateEngineTest.kt | 31 ++++++ .../infra/health/HealthIndicatorsTest.kt | 4 +- .../jdbc/JdbcFlightStateRoutePgTest.kt | 97 +++++++++++++++++++ .../processing/IgnoreBranchTest.kt | 23 +++-- 18 files changed, 268 insertions(+), 71 deletions(-) create mode 100644 src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt diff --git a/docs/implementation.md b/docs/implementation.md index fd1425d..cef8316 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -272,7 +272,7 @@ PENDING → SENT → DONE | 对象 | 终局判据 | 归档目标 | 清除证据 | 执行方 | 偏差 | |---|---|---|---|---|---| -| 共享信箱 `CMINMSGS` 原文 | 写回完成(`BACKFILL_AT` 非空) | — | 待确认(`Q9`) | 库方 | 清除协议未确认(`C-1`) | +| 共享信箱 `CMINMSGS` 原文 | 写回完成(`BACKFILL_AT` 非空) | — | 库方清除协议(保留期不早于写回完成) | 库方 | 本系统不删(`C-1`;`Q9` 按此闭合) | | `FLIGHT_SCHD` + 资源明细 | `US-14` AC2 的已结束条件 | 历史存储 | 历史写入确认 + 版本复查 | 我们 | — | | 航班历史存储 | 保留期 | — | — | 我们 | `G-FLIGHT-HIST-RETENTION` | | `SCHD_SNAP_LOG` | 保留期 | 无(本地可重建) | 无 | 我们 | — | @@ -332,7 +332,7 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 | `CHOT` | `CSNO`、`CHTM`、`CHID`、`CHST` | 99 | `SIS:3.22` | | `DELY` | `CODE`、`STRT`、`DURA`、文本 | — | [XSD](legacy/unisysaodbsis.xsd)「FLOP 元素」 | | `ABTM` | `ASNO`、`ABDG`、`ABOP`、`AOTM` | 99 | [XSD](legacy/unisysaodbsis.xsd)「FLOP 元素」 | -| `ROUT` / `ERUT` | `RTNO`、`APCD`、`SCAT`、`SCDT` | SIS 报文容量 6 / 7;本地各自最多保留 4 条(`US-05` AC3) | `SIS:3.40` | +| `ROUT` / `ERUT` | `RTNO`、`APCD` | SIS 报文容量 6 / 7;本地各自按序号升序最多保留 4 条,不落 `SCAT`/`SCDT`(`US-05` AC3) | `SIS:3.40` | | `SRVT` | `OPER`、`SRTC`、`SRQT`、`SRST`、`SRET`、`SRPR`、`SANR`、`SARR` | 无界 | `G-SRVT-VIPF` | | `VIPF` | `OPER`、`VPCD`、`VFES`、`VIPT/OPER`、`VIPT/VSCD`、`VIPT/VTQY`、`VIPT/VTST`、`VIPT/VTET` | 无界 | `G-SRVT-VIPF` | @@ -386,7 +386,7 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 | `SIS:3.37` | `HNAG` | `FHAG`/`PHAG`/`MHAG` 标量 | `FHAG` 空 = 删除该代理;`MHAG` 可缺席 | | `SIS:3.38` | `PSDT` | `PSDT` 集合 | 空 `` = 现无机位分配;AODB 实际会发,照常接收处理(`US-05` AC3) | | `SIS:3.39` | `RENO` | `RENO` 标量 | 空 = 清除注册号 | -| `SIS:3.40` | `ROUT` | `ROUT` 集合 | `SCAT`/`SCDT` 分别对起点/终点缺席。SIS 报文容量 6/7;本地 `ROUT`/`ERUT` **各自最多保留 4 条**(`US-05` AC3),不保存 `SCAT`/`SCDT` | +| `SIS:3.40` | `ROUT` | `ROUT` 集合 | `SCAT`/`SCDT` 不落当前态。SIS 报文容量 6/7;本地 `ROUT`/`ERUT` **各自按序号升序最多保留 4 条**(`US-05` AC3) | | `SIS:3.41` | `TAOP` | `TAOP`/`TAFL`/`TAID` 标量 | 任一为空 = 该到达航班的经停连接断开 | | `SIS:3.42` | `TRML` | `TRML` 标量 | 空 = 删除航站楼 | | `SIS:3.43` | `VIPP` | `VIPP`/`VIPR` 标量 | 空 = 删除;SIS 另要求 RMS 忽略 `VIPP`(忽略事件还是忽略字段,SIS 未写明,须以真实报文确认,见 `G-FLOP-SEMANTICS`) | diff --git a/docs/reference.md b/docs/reference.md index 21a3279..9f883c7 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -111,7 +111,7 @@ | `msgx.pipeline.job.last_sweep_selected` | 上次检查找到多少条待写回信箱的记录 | 持续达到每批上限:待办太多 | | `msgx.pipeline.codec.srvt_seen.total` | 收到含 `SRVT` 段的消息数 | 非零:核对 `G-SRVT-VIPF` | | `msgx.pipeline.codec.vipf_seen.total` | 收到含 `VIPF` 段的消息数 | 非零:核对 `G-SRVT-VIPF` | -| `msgx.pipeline.processing.ignored.total` | 被 `IgnoreRules` 跳过的消息数 | 清单含 `REGN`、`RSTA`、`EROR`,与 `US-13`、`US-09` 冲突 | +| `msgx.pipeline.processing.ignored.total` | 被 `IgnoreRules` 跳过的消息数 | 现行清单仅 `LDM-*`;`REGN`/`RSTA`/`EROR` 已改走 US-13/US-09 | | `msgx.pipeline.delivery.send_failures.total{target}` | 该投递目标累计发送失败次数(进程内,重启归零) | 持续增长且 `dead` 非零:投递链路故障(`OPS-2`) | | `msgx.pipeline.delivery.dead{target}` | 该投递目标当前死信(`DEAD`)行数 | 非零即告警:死信需人工处置(`OPS-2`;状态语义见 [implementation.md](implementation.md)「状态与错误分类」) | | `msgx.pipeline.watermark.lag` | 信箱最新编号比已扫描编号大多少 | 改用处理时间筛选后删除(`G-SCAN-PREDICATE`) | diff --git a/docs/requirements.md b/docs/requirements.md index 867f05c..0c9f347 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -11,7 +11,7 @@ **非目标**: - 航班当前态权威只在自有 PG,开发测试用 PostgreSQL;生产物理选型待 `Q14`,验证通过前不承诺 Oracle 兼容;Redis 仅作查询投影,不作权威或处理状态。 -- 共享 MySQL 只做读写消息和写回处理标记,不改表结构、不建表;已回填的入站行超过保留期后由本系统清理(specification.md 的 `C-1`)。 +- 共享 MySQL 只做读写消息和写回处理标记,不改表结构、不建表;入站原文保留与清除由库方负责(specification.md 的 `C-1`)。 - 本消息网关只有一个实例。 - 对外投递只承诺至少一次;同一航班(`FLID`)内保序,不同航班之间不承诺顺序。 - 已结束的航班写入 Elasticsearch 历史库后从实时数据删除。 @@ -53,7 +53,7 @@ **验收标准** 1. 一次只处理一条消息,取编号最小的未完成消息;处理中的消息不让后面的越过。 -2. 报文不合法:进死信。报文合法但本系统不支持该类型:跳过留档,按已处理写回标记。原始报文留在信箱,已回填的行由本系统按保留期清理(`C-1`)。 +2. 报文不合法:进死信。报文合法但本系统不支持该类型:跳过留档,按已处理写回标记。原始报文留在信箱,由库方按保留期清除(`C-1`)。 3. 处理或提交失败:失败的事务回滚,消息保持未完成,下一轮自动重新处理。 4. 发 Kafka 和回填信箱在处理完成之后做;下游失败不影响消息处理结果。 5. 错误必须记录到对应消息的处理记录上,不能被外层吞掉。 diff --git a/docs/specification.md b/docs/specification.md index 6fbb5f8..72a152f 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -41,7 +41,7 @@ ### 2.1 共享信箱(库方) -- **C-1** 本系统不删除入站行(`CMINMSGS`):处理完成后只写回处理标记,原文保留与清除由库方负责,保留期不早于该行的写回完成时刻。`(待确认 Q9)` +- **C-1** 本系统不删除入站行(`CMINMSGS`):处理完成后只写回处理标记,原文保留与清除由库方负责,保留期不早于该行的写回完成时刻。 - **C-2** 共享 MySQL 不改表结构;本系统只读写 `CMINMSGS` 与 `COUTMSGS`。 ### 2.2 上游(AODB / SIS) @@ -121,7 +121,7 @@ | Q6 | 待对方 | 网页客户端怎么读 Redis 快照 | 见 `C-11`;`MAFL` 是否提供未定 | | Q7 | 已定 | 信箱编号只增不减、不重用 | `INV-1`、`CLM-2`、`US-01` AC3 | | Q8 | 已定 | 入站处理时间列名 `CMINMSGS_DATE_PROCESSED` | 术语「处理标记」、`US-10` | -| Q9 | 待对方 | 入站行删除条件与保留期 | 本系统只写回处理标记、不删除;清除协议与保留期待库方确认(`C-1`) | +| Q9 | 已定 | 入站行删除条件与保留期 | 本系统只写回处理标记、不删除;清除与保留期由库方负责(`C-1`) | | Q10 | 已定 | `SEQN` 重置与身份规则 | `C-3` | | Q11 | 已定 | 日计划没带字段是否删除 | `C-6` | | Q12 | 已定 | 主航班与共享航班删除联动 | `US-06` AC2 | 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 bcfb631..1fb9ddc 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodec.kt @@ -68,6 +68,7 @@ class JacksonXmlCodec : XmlCodec { MsgKind.Fdel -> msg.flop?.let(SisWireMapper::flopPayload) // 参考数据类别报文的正文结构随 13 类各异(ACM2-93 实装 ReferenceDataProcessor 时绑定) is MsgKind.RefData -> null + MsgKind.Eror -> null is MsgKind.Unsupported -> null } @@ -106,6 +107,7 @@ class JacksonXmlCodec : XmlCodec { else -> MsgKind.Unsupported("FLOP-$styp") } in REF_DATA_TYPES -> MsgKind.RefData(type) + "EROR" -> MsgKind.Eror else -> MsgKind.Unsupported("$type-$styp") } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt index 239c66d..ab8d209 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt @@ -18,8 +18,10 @@ interface DeliveryPort { /** 发一条变化通知,消息 key 是 FLID(航班实例 ID)。 */ fun sendKafka(topic: String, key: String, payloadJson: String) - /** 发一条完整状态,消息 key 是 FLID。适配层负责把 target 映射成 topic:KAFKA:msg 对应 "msg",KAFKA:schd 对应 "schd"。 */ - fun sendKafkaSchd(topic: String, key: String, payloadJson: String) + /** + * 发一整批 schd 状态:`payloadJson` 是 `SCHD.FLTR` JSON 数组(`C-9`),**不设 message key**。 + */ + fun sendKafkaSchd(topic: String, payloadJson: String) /** 发一条删除通知:key 是 FLID、value 为空;下游按"整态里这个键没了"理解成删除。 */ fun sendKafkaNull(topic: String, key: String) @@ -33,8 +35,8 @@ interface DeliveryPort { * * KAFKA_MSG 按 `EVENT_ID` 顺序发,保序以 FLID 为单位(D2):某航班的队头失败或退避未到期 * 只暂停该 FLID,其他航班照常推进。KAFKA_SCHD 走 flushSchd:outbox 每个 FLID 只保留 - * 一行(写入侧单行 upsert),发出后按读取时刻的代次做条件确认;删除通知发 value 为空的 tombstone。 - * 两个主题之间不保证先后顺序。 + * 一行(写入侧单行 upsert),聚合成一条无 key 的 `SCHD.FLTR` JSON 数组发出(`C-9`); + * 删航班只走 msg。两个主题之间不保证先后顺序。 * * 失败处理:某 FLID 队首的重试时间没到就跳过该 FLID;发送失败重试次数加一并推后退避, * 次数用尽转 DEAD 当死信。见 docs/implementation.md「Kafka 与读取」。 @@ -46,6 +48,7 @@ class Dispatcher( private val props: PipelineProps, private val scheduler: FailureScheduler, private val counters: com.gzzn.omms.msgexchange.infra.metrics.PipelineCounters, + private val mapper: com.fasterxml.jackson.databind.ObjectMapper, ) { private val log = org.slf4j.LoggerFactory.getLogger(Dispatcher::class.java) @@ -152,7 +155,10 @@ class Dispatcher( } } - /** 发 KAFKA_SCHD:每个 FLID 只有一行待发事件,删除通知发空 value,成功后按代次条件确认。 */ + /** + * 发 `KAFKA:schd`:本 tick 窗口内待发 UPSERT 聚成一条 `SCHD.FLTR` JSON 数组、不设 key(`C-9`)。 + * 空窗口不发。schd 侧 TOMBSTONE 不发送(删航班只走 msg)。 + */ internal fun flushSchd() { val batch = try { msgEvents.mergePendingSchd(scheduler.now(), props.schd.flushLimit) @@ -160,25 +166,28 @@ class Dispatcher( log.warn("mergePendingSchd failed: {}", e.message) return } - if (batch.isEmpty()) { + val upserts = batch.filter { it.eventType == EventType.UPSERT } + if (upserts.isEmpty()) { lastFlush = scheduler.now() return } - val failures = mutableListOf>() - for (e in batch) { + val fltrs: List = upserts.map { e -> try { - when (e.eventType) { - EventType.TOMBSTONE -> port.sendKafkaNull("schd", e.partitionKey) - EventType.UPSERT -> port.sendKafkaSchd("schd", e.partitionKey, e.payloadJson) - } - // 条件确认:读取时刻的代次(EVENT_ID + STATE_VERSION)被新写入覆盖时不标记,留待下一轮重发。 - e.eventId?.let { msgEvents.markSentIfVersion(it, e.stateVersion, scheduler.now()) } - } catch (ex: Exception) { - // 首次失败就带上真实原因;不能退回字面量把根因抹掉 - failures += e to (ex.message ?: ex.javaClass.simpleName) + mapper.readValue(e.payloadJson, Any::class.java) + } catch (_: Exception) { + e.payloadJson } } - failures.forEach { (e, error) -> retryOrDead(e, error) } + val arrayJson = mapper.writeValueAsString(fltrs) + try { + port.sendKafkaSchd("schd", arrayJson) + for (e in upserts) { + e.eventId?.let { msgEvents.markSentIfVersion(it, e.stateVersion, scheduler.now()) } + } + } catch (ex: Exception) { + val error = ex.message ?: ex.javaClass.simpleName + upserts.forEach { e -> retryOrDead(e, error) } + } lastFlush = scheduler.now() } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt index 40f573e..2bbe7cc 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt @@ -14,7 +14,7 @@ data class MetaFields( val dttm: Long, ) -/** 报文的种类:日计划、运行动态、FDEL(航班终止)、静态参考数据,或者还没支持的类型。分派时用穷尽 when 保证不漏分支。 */ +/** 报文的种类:日计划、运行动态、FDEL(航班终止)、静态参考数据、错误回报,或者还没支持的类型。分派时用穷尽 when 保证不漏分支。 */ sealed interface MsgKind { data class Schd(val subtype: SchdSubtype) : MsgKind data class Flop(val subtype: String) : MsgKind // 运行动态,subtype 是报文子类型;FDEL 单独成类,不走这里 @@ -26,6 +26,11 @@ sealed interface MsgKind { */ data class RefData(val type: String) : MsgKind + /** + * AODB 错误回报(US-09 AC3):须匹配出站请求并标 FAILED;协调器落地前保持可重试失败并告警。 + */ + data object Eror : MsgKind + /** 已识别但还没有对应处理器的类型:跳过留档记 `SKIPPED(unsupported)`,不作为可重试失败占住队头。 */ data class Unsupported(val tag: String) : MsgKind @@ -45,6 +50,7 @@ data class DecodedMessage( is MsgKind.Flop -> "FLOP-${k.subtype}" MsgKind.Fdel -> "FDEL" is MsgKind.RefData -> "REF-${k.type}" + MsgKind.Eror -> "EROR" is MsgKind.Unsupported -> k.tag } } 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 b7c8d71..12a7976 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 @@ -22,6 +22,28 @@ object FlightStateEngine { "GTDT", "CKDT", "CLDT", "PSDT", "CHDT", "DELY", "ABTM", "CHOT", "ROUT", "ERUT", ) + private val ROUTE_COLLECTION_KEYS = setOf("ROUT", "ERUT") + private const val ROUTE_POINT_LIMIT = 4 + + /** + * US-05 AC3:ROUT/ERUT 各自最多 4 条;按上游序号(RTNO / 读回的 SOURCE_SEQ)升序保留, + * 不写入当前态的 SCAT/SCDT。 + */ + fun normalizeRouteCollection(items: List>): List> = + items.withIndex() + .sortedWith(compareBy({ (_, row) -> routeSeqKey(row) }, { (idx, _) -> idx })) + .map { it.value } + .take(ROUTE_POINT_LIMIT) + .map { row -> row.filterKeys { it != "SCAT" && it != "SCDT" } } + + private fun routeSeqKey(row: Map): Int { + val raw = row["RTNO"] ?: row["SOURCE_SEQ"] + return raw?.trim()?.toIntOrNull() ?: Int.MAX_VALUE + } + + private fun normalizeCollection(key: String, items: List>): List> = + if (key in ROUTE_COLLECTION_KEYS) normalizeRouteCollection(items) else items + /** * 校验一包日计划:报头声明的记录数(RECS)要在 0–9999 之间、且等于实际收到的条数; * 每条记录的 FLID 必须是 1–12 位数字,快照内不能有重复 FLID,每条记录的运营日都要 @@ -99,7 +121,7 @@ object FlightStateEngine { putAll(current.collections) // 报文没带的集合保持原值 } record.collections.forEach { (key, items) -> - if (key in COLLECTION_KEYS) put(key, items) // 报文带了整个集合就整体替换 + if (key in COLLECTION_KEYS) put(key, normalizeCollection(key, items)) } } return FlightSnapshot( @@ -124,7 +146,7 @@ object FlightStateEngine { val collections = buildMap { putAll(current.collections) change.collections.forEach { (key, items) -> - if (key in COLLECTION_KEYS) put(key, items) + if (key in COLLECTION_KEYS) put(key, normalizeCollection(key, items)) } } return current.copy( diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/kafka/KafkaDeliveryPort.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/kafka/KafkaDeliveryPort.kt index 6826f98..d67ba32 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/kafka/KafkaDeliveryPort.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/kafka/KafkaDeliveryPort.kt @@ -29,8 +29,9 @@ class KafkaDeliveryPort( producer.send(ProducerRecord(topic, key, payloadJson)).get() } - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) { - producer.send(ProducerRecord(topic, key, payloadJson)).get() + override fun sendKafkaSchd(topic: String, payloadJson: String) { + // C-9:不设 message key + producer.send(ProducerRecord(topic, null, payloadJson)).get() } override fun sendKafkaNull(topic: String, key: String) { 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 19c0a74..d11c125 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 @@ -8,6 +8,7 @@ import com.gzzn.omms.msgexchange.domain.MsgEvent import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.ProcStatus import com.gzzn.omms.msgexchange.domain.Targets +import com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine import com.gzzn.omms.msgexchange.domain.flight.FlightMainRow import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot import com.gzzn.omms.msgexchange.domain.flight.FlightState @@ -679,7 +680,12 @@ class JdbcFlightStateRepository( val collections = linkedMapOf>>() COLLECTIONS.forEach { (key, spec) -> - collections[key] = loadDetails(flid, key, spec) + val loaded = loadDetails(flid, key, spec) + collections[key] = if (key == "ROUT" || key == "ERUT") { + FlightStateEngine.normalizeRouteCollection(loaded) + } else { + loaded + } } return FlightSnapshot( flid = main.flid, @@ -812,8 +818,20 @@ class JdbcFlightStateRepository( private fun replaceDetails(snapshot: FlightSnapshot, now: Instant) { val flid = snapshot.flid COLLECTIONS.forEach { (key, spec) -> - ds.update("DELETE FROM ${spec.table} WHERE flid = ?", { ps -> ps.setString(1, flid) }) - val items = snapshot.collections[key].orEmpty() + if (spec.routeKind != null) { + ds.update( + "DELETE FROM ${spec.table} WHERE flid = ? AND route_kind = ?", + { ps -> + ps.setString(1, flid) + ps.setString(2, spec.routeKind) + }, + ) + } else { + ds.update("DELETE FROM ${spec.table} WHERE flid = ?", { ps -> ps.setString(1, flid) }) + } + val items = snapshot.collections[key].orEmpty().let { raw -> + if (key == "ROUT" || key == "ERUT") FlightStateEngine.normalizeRouteCollection(raw) else raw + } items.forEachIndexed { ordinal, item -> val cols = mutableListOf("flid", "ordinal", "source_seq") val vals = mutableListOf(flid) @@ -915,8 +933,8 @@ class JdbcFlightStateRepository( "DELY" to DetailSpec("flight_delay", listOf("code", "strt", "dura", "remc"), null), "ABTM" to DetailSpec("flight_bridge_op", listOf("abdg", "abop", "aotm"), "ASNO"), "CHOT" to DetailSpec("flight_chock_op", listOf("chid", "chst", "chtm"), "CSNO"), - "ROUT" to DetailSpec("flight_route_point", listOf("apcd", "scat", "scdt"), "RTNO", routeKind = "ROUT"), - "ERUT" to DetailSpec("flight_route_point", listOf("apcd", "scat", "scdt"), "RTNO", routeKind = "ERUT"), + "ROUT" to DetailSpec("flight_route_point", listOf("apcd"), "RTNO", routeKind = "ROUT"), + "ERUT" to DetailSpec("flight_route_point", listOf("apcd"), "RTNO", routeKind = "ERUT"), ) } } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt index ec5ed5a..bde7ca3 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt @@ -13,7 +13,7 @@ import jakarta.inject.Singleton @Requires(property = "msgx.stubs", value = "true") @Singleton class StubDeliveryPort : DeliveryPort { - data class Sent(val topic: String, val key: String, val payload: String?) + data class Sent(val topic: String, val key: String?, val payload: String?) val sent = mutableListOf() val tombstones: List get() = sent.filter { it.payload == null } @@ -24,8 +24,8 @@ class StubDeliveryPort : DeliveryPort { sent += Sent(topic, key, payloadJson) } - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) { - sent += Sent(topic, key, payloadJson) + override fun sendKafkaSchd(topic: String, payloadJson: String) { + sent += Sent(topic, null, payloadJson) } override fun sendKafkaNull(topic: String, key: String) { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreRules.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreRules.kt index 45ebe5d..daf40ea 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreRules.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreRules.kt @@ -1,18 +1,14 @@ package com.gzzn.omms.msgexchange.processing /** - * US-04 忽略清单:基线规则为 `LDM-*`、`REGN-*`、`RSTA-*`、`EROR-*`。 + * 忽略清单:仅保留文档未要求处理的 `LDM-*`。 * - * 匹配按 `TYPE` 前缀做(`TYPE-*` 表示该 TYPE 下所有 STYP 均忽略); - * `MetaFields.type` 已由 codec 统一大写,这里直接用大写常量比较。 - * 返回命中的规则标签(如 `LDM-*`),未命中返回 null。 + * `REGN`/`RSTA` 走参考数据(US-13);`EROR` 走出站错误回报(US-09)——不得在此跳过。 + * 匹配按 `TYPE` 精确相等;返回命中规则标签,未命中返回 null。 */ object IgnoreRules { private val rules: List> = listOf( "LDM" to "LDM-*", - "REGN" to "REGN-*", - "RSTA" to "RSTA-*", - "EROR" to "EROR-*", ) fun match(type: String): String? { 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 95243bc..399e87b 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -251,6 +251,11 @@ class MessageProcessor( log.warn("refdata without handler -> FAILED(UNSUPPORTED) msgId={} type={}", head.msgId, kind.type) return procFailure.fail(head, ErrorClass.UNSUPPORTED, "refdata-pending:${kind.type}") } + MsgKind.Eror -> { + // US-09 AC3:须匹配出站请求;REQ_TRACK 协调器落地前(ACM2-92)先按可重试失败留队。 + log.warn("eror without handler -> FAILED(UNSUPPORTED) msgId={}", head.msgId) + return procFailure.fail(head, ErrorClass.UNSUPPORTED, "eror-pending") + } is MsgKind.Unsupported -> { // 合法但不支持的类型:跳过留档记 SKIPPED(US-03 AC2;markTerminal 同一条 UPDATE // 登记回填意图),不当可重试失败占住队头 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt index d02f67e..647c342 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt @@ -1,5 +1,6 @@ package com.gzzn.omms.msgexchange.delivery +import com.fasterxml.jackson.databind.ObjectMapper import com.gzzn.omms.msgexchange.MutableClock import com.gzzn.omms.msgexchange.config.PipelineProps import com.gzzn.omms.msgexchange.domain.ErrorClass @@ -21,20 +22,21 @@ import java.time.Instant /** * 守着投递的几条规矩: * - KAFKA_MSG 一条条按顺序发,而且不会顺手把 KAFKA_SCHD 的事件发掉; - * - KAFKA_SCHD 只有 flushSchd 一个出口,同一个 FLID(航班实例 ID)只发版本号最新的那条; - * - 删除通知发成 value 为空的 tombstone 消息; + * - KAFKA_SCHD 只有 flushSchd 一个出口:本 tick 待发 UPSERT 聚成一条无 key 的 FLTR 数组(C-9); + * - schd 不发 tombstone(删航班只走 msg); * - 发送失败按退避重试,次数用尽转 DEAD 当死信。 */ class DispatcherTickTest { private val clock = MutableClock(MutableClock.BASE) + private val mapper = ObjectMapper() private fun dispatcher( repo: MsgEventRepository, port: DeliveryPort, p: PipelineProps = PipelineProps(), counters: PipelineCounters = PipelineCounters(), - ): Dispatcher = Dispatcher(repo, port, p, FailureScheduler(p, clock), counters) + ): Dispatcher = Dispatcher(repo, port, p, FailureScheduler(p, clock), counters, mapper) private fun ev(id: Long, target: String, key: String, payload: String, version: Long = 0) = MsgEvent(eventId = id, target = target, partitionKey = key, stateVersion = version, payloadJson = payload, createdAt = MutableClock.BASE) @@ -70,7 +72,7 @@ class DispatcherTickTest { } @Test - fun `flushSchd sends one row per flight and marks it sent`() { + fun `flushSchd sends one FLTR array without key and marks rows sent`() { val repo = StubMsgEvents() val port = StubDeliveryPort() repo.insertAll( @@ -84,8 +86,12 @@ class DispatcherTickTest { d.flushSchd() val schd = port.sent.filter { it.topic == "schd" } - assertEquals(2, schd.size) - assertEquals(1, schd.count { it.key == "F1" && it.payload!!.contains("new") && !it.payload.contains("old") }) + assertEquals(1, schd.size) + assertNull(schd.single().key) + val payload = schd.single().payload!! + assertTrue(payload.startsWith("[")) + assertTrue(payload.contains("new") && !payload.contains("old")) + assertTrue(payload.contains("f2")) assertEquals(0, repo.rows.values.count { it.state.name == "PENDING" }) } @@ -104,7 +110,7 @@ class DispatcherTickTest { } @Test - fun `tombstone is delivered as null value message`() { + fun `schd tombstone is not sent on schd topic`() { val repo = StubMsgEvents() val port = StubDeliveryPort() repo.insertAll( @@ -119,9 +125,8 @@ class DispatcherTickTest { val d = dispatcher(repo, port) d.flushSchd() - assertEquals(1, port.tombstones.size) - assertEquals("F1", port.tombstones.single().key) // value 为空表示删掉这个 key 的旧值 - assertNull(port.tombstones.single().payload) + assertEquals(0, port.sent.size) + assertEquals(0, port.tombstones.size) } @Test @@ -325,7 +330,7 @@ private class FailingMsgKeyPort(private val failingKey: String) : DeliveryPort { sent.add(StubDeliveryPort.Sent(topic, key, payloadJson)) } - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) = Unit + override fun sendKafkaSchd(topic: String, payloadJson: String) = Unit override fun sendKafkaNull(topic: String, key: String) = Unit } @@ -333,7 +338,7 @@ private class FailingMsgKeyPort(private val failingKey: String) : DeliveryPort { private class FailingSchdPort : DeliveryPort { var calls = 0 override fun sendKafka(topic: String, key: String, payloadJson: String) = Unit - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) { + override fun sendKafkaSchd(topic: String, payloadJson: String) { calls++ throw IllegalStateException("broker-down") } @@ -346,10 +351,10 @@ private class FailingSchdPort : DeliveryPort { /** 在发送过程中给同一 FLID 写入新代次,用来验证「条件确认」不会把新内容标记成已发。 */ private class ReentrantSchdPort(private val repo: StubMsgEvents) : DeliveryPort { override fun sendKafka(topic: String, key: String, payloadJson: String) = Unit - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) { + override fun sendKafkaSchd(topic: String, payloadJson: String) { repo.insertAll( listOf( - MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = key, stateVersion = 99, payloadJson = """{"v":"new"}""", createdAt = MutableClock.BASE), + MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = "F1", stateVersion = 99, payloadJson = """{"v":"new"}""", createdAt = MutableClock.BASE), ), ) } 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 b71c4b1..5c14edc 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 @@ -169,6 +169,37 @@ class FlightStateEngineTest { assertTrue(empty.flags.contains(SnapshotFlag.EMPTY)) } + @Test + fun `route collections keep lowest four by rtno ascending and drop scat scdt`() { + val sixRout = (1..6).map { n -> + mapOf("RTNO" to n.toString(), "APCD" to "AP$n", "SCAT" to "15DEC031000", "SCDT" to "15DEC031100") + } + val merged = FlightStateEngine.snapshotState( + null, + ScheduleRecord("121", emptyMap(), mapOf("ROUT" to sixRout)), + operationDay = day, + keepDeleted = false, + ) + val rout = merged.collections["ROUT"]!! + assertEquals(4, rout.size) + assertEquals(listOf("1", "2", "3", "4"), rout.map { it["RTNO"] }) + assertEquals(listOf("AP1", "AP2", "AP3", "AP4"), rout.map { it["APCD"] }) + assertTrue(rout.all { !it.containsKey("SCAT") && !it.containsKey("SCDT") }) + } + + @Test + fun `rout and erut route caps are independent`() { + val five = (1..5).map { n -> mapOf("RTNO" to n.toString(), "APCD" to "X$n") } + val merged = FlightStateEngine.snapshotState( + null, + ScheduleRecord("121", emptyMap(), mapOf("ROUT" to five, "ERUT" to five)), + operationDay = day, + keepDeleted = false, + ) + assertEquals(4, merged.collections["ROUT"]!!.size) + assertEquals(4, merged.collections["ERUT"]!!.size) + } + @Test fun `validation rejects non numeric or oversized flid`() { val bad = FlightStateEngine.validateMessage( diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt index 5004a99..3708056 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt @@ -19,14 +19,14 @@ class HealthIndicatorsTest { private class FakePort(private val pingResult: Boolean) : DeliveryPort { override fun sendKafka(topic: String, key: String, payloadJson: String) = Unit - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) = Unit + override fun sendKafkaSchd(topic: String, payloadJson: String) = Unit override fun sendKafkaNull(topic: String, key: String) = Unit override fun ping(): Boolean = pingResult } private class ThrowingPort : DeliveryPort { override fun sendKafka(topic: String, key: String, payloadJson: String) = Unit - override fun sendKafkaSchd(topic: String, key: String, payloadJson: String) = Unit + override fun sendKafkaSchd(topic: String, payloadJson: String) = Unit override fun sendKafkaNull(topic: String, key: String) = Unit override fun ping(): Boolean = throw RuntimeException("metadata fetch failed") } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt new file mode 100644 index 0000000..00d63d9 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcFlightStateRoutePgTest.kt @@ -0,0 +1,97 @@ +package com.gzzn.omms.msgexchange.infra.persistence.jdbc + +import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot +import com.gzzn.omms.msgexchange.domain.flight.FlightState +import com.gzzn.omms.msgexchange.support.PgTestSupport +import com.zaxxer.hikari.HikariConfig +import com.zaxxer.hikari.HikariDataSource +import org.flywaydb.core.Flyway +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Assumptions.assumeTrue +import org.junit.jupiter.api.Test +import java.sql.DriverManager +import java.time.Clock +import java.time.Instant +import java.time.LocalDate +import java.time.ZoneOffset +import java.util.UUID + +/** US-05 AC3:ROUT/ERUT 各自最多 4 条、不存 SCAT/SCDT,且两类路线同 FLID 并存。 */ +class JdbcFlightStateRoutePgTest { + + private val t0: Instant = Instant.parse("2026-09-12T00:00:00Z") + + @Test + fun `persist keeps four rout and four erut without scat scdt`() { + assumeTrue(PgTestSupport.canConnect(), PgTestSupport.skipMessage()) + val ds = dataSource() + val clock = Clock.fixed(t0, ZoneOffset.UTC) + val repo = JdbcFlightStateRepository(ds, clock) + val flid = "RT-" + UUID.randomUUID().toString().take(8) + val rout = (1..6).map { n -> + mapOf("RTNO" to n.toString(), "APCD" to "R$n", "SCAT" to "15DEC031000", "SCDT" to "15DEC031100") + } + val erut = (1..5).map { n -> mapOf("RTNO" to n.toString(), "APCD" to "E$n", "SCAT" to "x", "SCDT" to "y") } + repo.persistFullState( + FlightSnapshot( + flid, LocalDate.of(2026, 9, 12), FlightState.ACTIVE, 1, + mapOf("SODT" to "12Sep261200"), + mapOf("ROUT" to rout, "ERUT" to erut), + ), + msgId = 1, + now = t0, + ) + + DriverManager.getConnection(PgTestSupport.jdbcUrl, PgTestSupport.user, PgTestSupport.password).use { conn -> + conn.prepareStatement( + "SELECT route_kind, ordinal, source_seq, apcd, scat, scdt FROM flight_route_point WHERE flid = ? ORDER BY route_kind, ordinal", + ).use { ps -> + ps.setString(1, flid) + ps.executeQuery().use { rs -> + val rows = mutableListOf>() + while (rs.next()) { + rows.add( + listOf( + rs.getString("route_kind"), + rs.getInt("ordinal").toString(), + rs.getString("source_seq"), + rs.getString("apcd"), + rs.getString("scat"), + rs.getString("scdt"), + ), + ) + } + assertEquals(8, rows.size) + assertEquals(4, rows.count { it[0] == "ROUT" }) + assertEquals(4, rows.count { it[0] == "ERUT" }) + assertTrue(rows.all { it[4] == null && it[5] == null }) + assertEquals(listOf("1", "2", "3", "4"), rows.filter { it[0] == "ROUT" }.map { it[2] }) + assertEquals(listOf("1", "2", "3", "4"), rows.filter { it[0] == "ERUT" }.map { it[2] }) + } + } + } + + val loaded = repo.loadFullSnapshot(flid)!! + assertFalse(loaded.collections["ROUT"]!!.any { it.containsKey("SCAT") || it.containsKey("SCDT") }) + assertEquals(4, loaded.collections["ROUT"]!!.size) + assertEquals(4, loaded.collections["ERUT"]!!.size) + } + + private fun dataSource(): javax.sql.DataSource { + Flyway.configure() + .dataSource(PgTestSupport.jdbcUrl, PgTestSupport.user, PgTestSupport.password) + .locations("classpath:db/migration") + .load() + .migrate() + return HikariDataSource( + HikariConfig().apply { + jdbcUrl = PgTestSupport.jdbcUrl + username = PgTestSupport.user + password = PgTestSupport.password + maximumPoolSize = 2 + }, + ) + } +} diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt index 2f5e285..6938785 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/IgnoreBranchTest.kt @@ -5,6 +5,7 @@ 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.DecodedMessage +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.ProcState @@ -105,23 +106,25 @@ class IgnoreBranchTest { } @Test - fun `EROR subtype hits ignore rule`() { + fun `EROR routes to eror path not ignore list`() { val proc = StubProcState() proc.insertIfAbsent(msgId, null) val inbox = StubInbox() inbox.raws[msgId] = "" - val msg = decoded("EROR", "GEN") + val msg = decoded("EROR", "GEN").copy(kind = MsgKind.Eror) val p = processor(proc = proc, inbox = inbox, codec = codecReturning(msg)) p.processOne(head()) - assertEquals(ProcStatus.SKIPPED, proc.find(msgId)!!.state) - assertEquals("ignored:EROR-*", proc.find(msgId)!!.lastError) + val row = proc.find(msgId)!! + assertEquals(ProcStatus.FAILED, row.state) + assertEquals(ErrorClass.UNSUPPORTED, row.errorClass) + assertEquals("eror-pending", row.lastError) } @Test - fun `REGN and RSTA both hit ignore rules`() { - for ((type, rule) in listOf("REGN" to "REGN-*", "RSTA" to "RSTA-*")) { + fun `REGN and RSTA route to refdata path not ignore list`() { + for (type in listOf("REGN", "RSTA")) { val proc = StubProcState() val id = msgId + type.hashCode().toLong() proc.insertIfAbsent(id, null) @@ -129,15 +132,17 @@ class IgnoreBranchTest { inbox.raws[id] = "" val msg = DecodedMessage( meta = MetaFields("AODB", type, "X", 300L, 1L), - kind = MsgKind.Unsupported("$type-X"), + kind = MsgKind.RefData(type), rawXml = "", ) val p = processor(proc = proc, inbox = inbox, codec = codecReturning(msg)) p.processOne(ProcState(id, ProcStatus.PENDING, updatedAt = Instant.EPOCH)) - assertEquals(ProcStatus.SKIPPED, proc.find(id)!!.state) - assertEquals("ignored:$rule", proc.find(id)!!.lastError) + val row = proc.find(id)!! + assertEquals(ProcStatus.FAILED, row.state) + assertEquals(ErrorClass.UNSUPPORTED, row.errorClass) + assertEquals("refdata-pending:$type", row.lastError) } }