fix(processing): Wave1 — IgnoreRules/ROUT截断/Q9文档/schd FLTR 数组(ACM2-94/88/81/96)

REGN/RSTA/EROR 退出忽略清单;ROUT/ERUT 升序留 4 且不落 SCAT/SCDT;C-1/Q9 统一库方清除;flushSchd 聚合成无 key 的 FLTR 数组。

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
windyboy
2026-09-21 14:18:40 +08:00
co-authored by Cursor
parent d298fcc48f
commit 23b74ed554
18 changed files with 268 additions and 71 deletions
+3 -3
View File
@@ -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` 集合 | 空 `<PSDT PSNO="0">` = 现无机位分配;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` |
+1 -1
View File
@@ -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` |
+2 -2
View File
@@ -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. 错误必须记录到对应消息的处理记录上,不能被外层吞掉。
+2 -2
View File
@@ -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 |
@@ -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")
}
@@ -18,8 +18,10 @@ interface DeliveryPort {
/** 发一条变化通知,消息 key 是 FLID(航班实例 ID)。 */
fun sendKafka(topic: String, key: String, payloadJson: String)
/** 发一条完整状态,消息 key 是 FLID。适配层负责把 target 映射成 topicKAFKA: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 走 flushSchdoutbox 每个 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<Pair<MsgEvent, String>>()
for (e in batch) {
val fltrs: List<Any> = 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()
}
@@ -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
}
}
@@ -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 AC3ROUT/ERUT 各自最多 4 条;按上游序号(RTNO / 读回的 SOURCE_SEQ)升序保留,
* 不写入当前态的 SCAT/SCDT。
*/
fun normalizeRouteCollection(items: List<Map<String, String>>): List<Map<String, String>> =
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<String, String>): Int {
val raw = row["RTNO"] ?: row["SOURCE_SEQ"]
return raw?.trim()?.toIntOrNull() ?: Int.MAX_VALUE
}
private fun normalizeCollection(key: String, items: List<Map<String, String>>): List<Map<String, String>> =
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(
@@ -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<String, String>(topic, null, payloadJson)).get()
}
override fun sendKafkaNull(topic: String, key: String) {
@@ -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<String, List<Map<String, String>>>()
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<Any?>(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"),
)
}
}
@@ -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<Sent>()
val tombstones: List<Sent> 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) {
@@ -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<Pair<String, String>> = listOf(
"LDM" to "LDM-*",
"REGN" to "REGN-*",
"RSTA" to "RSTA-*",
"EROR" to "EROR-*",
)
fun match(type: String): String? {
@@ -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 -> {
// 合法但不支持的类型:跳过留档记 SKIPPEDUS-03 AC2markTerminal 同一条 UPDATE
// 登记回填意图),不当可重试失败占住队头
@@ -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),
),
)
}
@@ -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(
@@ -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")
}
@@ -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 AC3ROUT/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<List<String?>>()
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
},
)
}
}
@@ -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] = "<RAW/>"
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] = "<RAW/>"
val msg = DecodedMessage(
meta = MetaFields("AODB", type, "X", 300L, 1L),
kind = MsgKind.Unsupported("$type-X"),
kind = MsgKind.RefData(type),
rawXml = "<MSG/>",
)
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)
}
}