diff --git a/docs/implementation.md b/docs/implementation.md index d780d71..fc2f848 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -28,6 +28,7 @@ msgexchange-v2 怎么处理报文:记录模型、状态机、事务边界、 | `REF_MASTER` | SIS 消息提供的静态参考数据与资源状态的逻辑视图(物理为独立数据表组) | `(RTYPE, RKEY)` 唯一;`RTYPE` 类别、合并语义与资源状态见「静态参考数据」;取数路径见 [requirements.md](requirements.md) `US-13`。 | | `FLIGHT_SCHD` | 航班标量及单值异常字段 | `FLID` 主键;运营日与版本、最近消息 ID 用于追踪。变长集合存于资源明细表与 `FLIGHT_ROUTE_POINT`,规则见「航班域」。 | | `SCHD_SNAP_LOG` | 日计划处理留痕 | 只追加、可重建,不参与状态决策;保留期见 [reference.md](reference.md)。 | +| `UNMAPPED_FIELD` | 未落入航班当前态的字段原值(`US-05` AC3) | 唯一键 `(MSG_ID, PATH, ORDINAL)`;不随 `PROC_STATE` 到期删除;路径与 `RAW_VALUE` 编码见「未映射字段」。 | 字段与索引以 `src/main/resources/db/migration/` 为准。原文读共享信箱;清理见 `C-1`。 @@ -302,7 +303,8 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 | 对象 | 职责 | |---|---| | `FLIGHT_SCHD` | 一行一个 `FLID`,保存标量字段、`STATE`、`STATE_VERSION`、`OPERATION_DAY`、最近消息 ID 和审计时间。 | -| 资源明细表 | 保存登机门、值机柜台、转盘、计划机位、滑槽、延误、靠撤桥、轮挡等变长集合;`SRVT`/`VIPF` 专用明细见 `G-SRVT-VIPF`。主键为 `(FLID, ORDINAL)`。 | +| 资源明细表 | 保存登机门、值机柜台、转盘、计划机位、滑槽、延误、靠撤桥、轮挡等变长集合;主键为 `(FLID, ORDINAL)`。 | +| `FLIGHT_SRVT` / `FLIGHT_VIPF` | 服务与 VIP 明细(`G-SRVT-VIPF`);主键 `(FLID, ORDINAL)`。段**出现**时按报文顺序整段替换;段缺席不删已有行(`Q2` 默认)。 | | `FLIGHT_ROUTE_POINT` | ROUT 与 ERUT 两类路线点,使用 `ROUTE_KIND` 区分;主键应包含该列,避免两类路线的序号冲突。 | | `PROC_STATE` | 信箱消息的处理终态、业务身份幂等记录,以及回填事实(`RECEIVED_AT` / `BACKFILL_*`)。 | | `MSG_EVENT` | 事务 outbox,承载整态快照与变更通知;删航班只走 `msg`(`C-9`)。 | @@ -341,7 +343,21 @@ PG 是航班数据源(`INV-5`);Redis 从 PG 同步,处理完成前写入 | `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` | -`SRVT`/`VIPF` 只存原文,不参与合并。`MAFL` 由子航班派生(`G-MAFL`)。 +`SRVT`/`VIPF` 落专用明细表,不进 Redis 快照与合并层(`G-SRVT-VIPF`)。`MAFL` 由子航班派生(`G-MAFL`)。 + +### 11.3.1 未映射字段(`UNMAPPED_FIELD`) + +现行白名单内的 FLOP 报文:wire 已解码、且不能写入 `FLIGHT_SCHD` 或资源明细的字段,按消息写入 `UNMAPPED_FIELD`(`G-FLOP-UNMAPPED`、`US-05` AC3)。非白名单 FLOP 子类型仍走 `SKIPPED`(`US-03` AC2),不在此表留证。 + +| 列 | 含义 | +|---|---| +| `MSG_ID` | 信箱消息 ID,与 `PROC_STATE` 对齐 | +| `FLID` | 报文内航班 ID(便于按航班查询) | +| `PATH` | 字段路径:标量元素名为 SIS 标签(如 `FDIV`);属性为 `TAG/@ATTR`(如 `FDIV/@DDES`);同父下重复兄弟以 `PATH` 相同、`ORDINAL` 递增区分 | +| `ORDINAL` | 从 1 起;单元素多属性(如 `FDIV/@DDES` 与 `FDIV/@DDIR`)各自 `ORDINAL=1`,因 `PATH` 不同而不冲突 | +| `RAW_VALUE` | 原值 UTF-8 字符串;显式空标签记空串,不省略行 | + +写入与航班变更同一 PG 事务:对该 `MSG_ID` 先删后插,重处理结果一致即幂等。清理保留期不含本表(`G-PROC-CLEANUP` 只清 `PROC_STATE`)。 - `ORDINAL` 是持久化顺序,从 1 开始;`SOURCE_SEQ` 是上游序号,允许为空或重复。 - 相同资源号不代表同一条分配,禁止按资源号去重。 diff --git a/docs/specification.md b/docs/specification.md index 907ba6f..8912fca 100644 --- a/docs/specification.md +++ b/docs/specification.md @@ -160,7 +160,7 @@ | `G-REQ-TRACK-RETENTION` | `REQ_TRACK` 已结案保留期未定(`Q24`) | `US-09` | | `G-RESP-GUARD` | `SCHD-RESP` 过期判断未做 | `US-07` | | `G-SCHD-SNAPSHOT` | 日计划快照删除、清空、分批未做 | `INV-7`、`INV-9` | -| `G-SRVT-VIPF` | `SRVT`、`VIPF` 明细未入库(`Q2`) | `US-05` | +| `G-SRVT-VIPF` | `SRVT`、`VIPF` 缺席是否清除待 `Q2`;段出现时已落 `FLIGHT_SRVT`/`FLIGHT_VIPF` | `US-05` | ## 8. 验证映射 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 a000a87..48dec9f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt @@ -6,12 +6,14 @@ import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlProperty import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlRootElement import com.fasterxml.jackson.dataformat.xml.annotation.JacksonXmlText import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord +import com.gzzn.omms.msgexchange.domain.flight.UnmappedFieldEntry /** Stable payload boundary used by processing. It is deliberately not an XML model. */ data class FlopPayload( val flid: String, val scalars: Map = emptyMap(), val collections: Map>> = emptyMap(), + val unmapped: List = emptyList(), ) data class ScheduleBody( @@ -20,10 +22,8 @@ data class ScheduleBody( ) /** - * 只保留在 wire/domain、尚未映射到持久化明细的集合键(`[G-SRVT-VIPF]`)。 - * - * 它们不参与合并、不进快照、不落库;保留的目的是不让入站事实在解码层被静默抹平,并让真实 - * 流量里的出现情况可观测——决定"缺席是否等于删除"的 `Q2` 需要真实报文才能定案。 + * `SRVT`/`VIPF` 集合键(`[G-SRVT-VIPF]`):用于入站出现次数指标;明细落库见 `FLIGHT_SRVT`/`FLIGHT_VIPF`, + * 不参与合并层与 Redis 快照。 */ private val UNPERSISTED_COLLECTION_KEYS: Set = setOf("SRVT", "VIPF") @@ -160,11 +160,30 @@ data class FlightRecordXml( @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 = 未出现;出现即为列表(空元素得到一行空行)。 - // 两者都不落明细表、不参与合并,清空语义待 `Q2`(`[G-SRVT-VIPF]`)。 + // 明细表见 `G-SRVT-VIPF`;缺席是否清除待 `Q2`,现行默认不删已有行。 @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "SRVT") val srvt: List? = null, @param:JacksonXmlElementWrapper(useWrapping = false) @param:JacksonXmlProperty(localName = "VIPF") val vipf: List? = null, + @param:JacksonXmlProperty(localName = "FDIV") val fdiv: FdivXml? = null, + @param:JacksonXmlProperty(localName = "FRET") val fret: FretXml? = null, ) +@JsonIgnoreProperties(ignoreUnknown = true) +data class FdivXml( + @param:JacksonXmlProperty(isAttribute = true, localName = "DDES") val ddes: String? = null, + @param:JacksonXmlProperty(isAttribute = true, localName = "DDIR") val ddir: String? = null, +) { + @field:JacksonXmlText + var text: String? = null +} + +@JsonIgnoreProperties(ignoreUnknown = true) +data class FretXml( + @param:JacksonXmlProperty(isAttribute = true, localName = "REID") val reid: String? = null, +) { + @field:JacksonXmlText + var text: String? = null +} + @JsonIgnoreProperties(ignoreUnknown = true) class FdelXml 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 771f009..a6a259c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt @@ -1,6 +1,7 @@ package com.gzzn.omms.msgexchange.codec import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord +import com.gzzn.omms.msgexchange.domain.flight.UnmappedFieldEntry /** Converts the annotated wire DTOs into the stable payload types used by processing. */ internal object SisWireMapper { @@ -11,7 +12,7 @@ internal object SisWireMapper { fun flopPayload(xml: FlightRecordXml): FlopPayload? = xml.flid?.trim()?.takeIf(String::isNotEmpty)?.let { - FlopPayload(it, xml.scalars(), xml.collections()) + FlopPayload(it, xml.scalars(), xml.collections(), xml.unmappedFields()) } private fun scheduleRecord(xml: FlightRecordXml): ScheduleRecord? = @@ -58,6 +59,25 @@ internal object SisWireMapper { return mapped } + /** 尚未落入当前态的 FLOP 元素(`G-FLOP-UNMAPPED`);路径编码见 implementation.md「未映射字段」。 */ + private fun FlightRecordXml.unmappedFields(): List = buildList { + fdiv?.let { addAll(it.toUnmapped()) } + fret?.let { addAll(it.toUnmapped()) } + } + + private fun FdivXml.toUnmapped(): List = buildList { + val ord = 1 + if (ddes != null) add(UnmappedFieldEntry("FDIV/@DDES", ord, ddes.trim())) + if (ddir != null) add(UnmappedFieldEntry("FDIV/@DDIR", ord, ddir.trim())) + if (text != null) add(UnmappedFieldEntry("FDIV", ord, text!!.trim())) + } + + private fun FretXml.toUnmapped(): List = buildList { + val ord = 1 + if (reid != null) add(UnmappedFieldEntry("FRET/@REID", ord, reid.trim())) + if (text != null) add(UnmappedFieldEntry("FRET", ord, text!!.trim())) + } + /** `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, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/UnmappedFieldEntry.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/UnmappedFieldEntry.kt new file mode 100644 index 0000000..69f035f --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/UnmappedFieldEntry.kt @@ -0,0 +1,8 @@ +package com.gzzn.omms.msgexchange.domain.flight + +/** 一条未映射字段记录(`US-05` AC3);`path`/`ordinal`/`raw` 编码见 implementation.md「未映射字段」。 */ +data class UnmappedFieldEntry( + val path: String, + val ordinal: Int, + val raw: String, +) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt index f27f0d4..701ab9a 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/Repositories.kt @@ -9,6 +9,7 @@ import com.gzzn.omms.msgexchange.domain.flight.FlightMainRow import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot import com.gzzn.omms.msgexchange.domain.flight.HistoryCandidate import com.gzzn.omms.msgexchange.domain.flight.HistoryRules +import com.gzzn.omms.msgexchange.domain.flight.UnmappedFieldEntry import java.time.Instant import java.time.LocalDate import java.time.ZoneId @@ -268,6 +269,25 @@ interface FlightStateRepository { */ fun persistFullState(snapshot: FlightSnapshot, msgId: Long, now: Instant): PersistOutcome + /** + * `SRVT`/`VIPF` 段**出现**时写入专用明细(`G-SRVT-VIPF`);键缺席时不删已有行(`Q2` 默认)。 + */ + fun persistSrvtVipfIfPresent( + flid: String, + collections: Map>>, + now: Instant, + ) + + /** + * 按消息重写未映射字段证据(`US-05` AC3):先删该 `MSG_ID` 再插入,与航班变更同事务,重处理幂等。 + */ + fun replaceUnmappedForMessage( + msgId: Long, + flid: String, + fields: List, + now: Instant, + ) + /** * 删除航班:把在用的航班标成 DELETED、版本号加一,明细数据保留不删。 * 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 29e91cd..54989ec 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 @@ -698,6 +698,40 @@ class JdbcFlightStateRepository( return if (existed == null) PersistOutcome.INSERTED else PersistOutcome.UPDATED } + override fun persistSrvtVipfIfPresent( + flid: String, + collections: Map>>, + now: Instant, + ) { + if ("SRVT" in collections) replaceSrvt(flid, collections.getValue("SRVT"), now) + if ("VIPF" in collections) replaceVipf(flid, collections.getValue("VIPF"), now) + } + + override fun replaceUnmappedForMessage( + msgId: Long, + flid: String, + fields: List, + now: Instant, + ) { + ds.update("DELETE FROM unmapped_field WHERE msg_id = ?", { ps -> ps.setLong(1, msgId) }) + fields.forEach { field -> + ds.update( + """ + INSERT INTO unmapped_field (msg_id, flid, path, ordinal, raw_value, created_at) + VALUES (?, ?, ?, ?, ?, ?) + """.trimIndent(), + { ps -> + ps.setLong(1, msgId) + ps.setString(2, flid) + ps.setString(3, field.path) + ps.setInt(4, field.ordinal) + ps.setString(5, field.raw) + ps.setTimestamp(6, now.toSqlTimestamp()) + }, + ) + } + } + override fun markDeleted(flid: String, msgId: Long, now: Instant): Boolean = ds.update( "UPDATE flight_schd SET state = 'DELETED', state_version = state_version + 1, last_msg_id = ?, updated_at = ? " + @@ -761,12 +795,99 @@ class JdbcFlightStateRepository( { ps -> ps.setString(1, candidate.flid) }, ) } + SEGMENT_TABLES.forEach { table -> + purged += ds.update( + "DELETE FROM $table WHERE flid = ?", + { ps -> ps.setString(1, candidate.flid) }, + ) + } } return purged } // ---- 内部实现 ---- + private fun replaceSrvt(flid: String, items: List>, now: Instant) { + ds.update("DELETE FROM flight_srvt WHERE flid = ?", { ps -> ps.setString(1, flid) }) + items.forEachIndexed { idx, item -> + insertSegmentRow( + "flight_srvt", + flid, + idx + 1, + listOf("oper", "srtc", "srqt", "srst", "sret", "srpr", "sanr", "sarr"), + item, + now, + ) + } + } + + private fun replaceVipf(flid: String, items: List>, now: Instant) { + ds.update("DELETE FROM flight_vipf WHERE flid = ?", { ps -> ps.setString(1, flid) }) + items.forEachIndexed { idx, item -> + insertSegmentRow( + "flight_vipf", + flid, + idx + 1, + listOf("oper", "vpcd", "vfes", "vipt_oper", "vipt_vscd", "vipt_vtqy", "vipt_vtst", "vipt_vtet"), + item, + mapOf( + "OPER" to "oper", + "VPCD" to "vpcd", + "VFES" to "vfes", + "VIPT_OPER" to "vipt_oper", + "VIPT_VSCD" to "vipt_vscd", + "VIPT_VTQY" to "vipt_vtqy", + "VIPT_VTST" to "vipt_vtst", + "VIPT_VTET" to "vipt_vtet", + ), + now, + ) + } + } + + private fun insertSegmentRow( + table: String, + flid: String, + ordinal: Int, + columns: List, + item: Map, + now: Instant, + ) = insertSegmentRow(table, flid, ordinal, columns, item, columns.associateWith { it.uppercase() }, now) + + private fun insertSegmentRow( + table: String, + flid: String, + ordinal: Int, + columns: List, + item: Map, + keyToCol: Map, + now: Instant, + ) { + val cols = mutableListOf("flid", "ordinal") + val vals = mutableListOf(flid, ordinal) + columns.forEach { col -> + cols.add(col) + val wireKey = keyToCol.entries.firstOrNull { it.value == col }?.key ?: col.uppercase() + vals.add(item[wireKey]?.takeIf { it.isNotEmpty() }) + } + cols.add("created_at"); vals.add(now) + cols.add("updated_at"); vals.add(now) + val placeholders = cols.indices.joinToString(",") { "?" } + ds.update( + "INSERT INTO $table (${cols.joinToString(", ")}) VALUES ($placeholders)", + { ps -> + vals.forEachIndexed { idx, v -> + when (v) { + null -> ps.setNull(idx + 1, java.sql.Types.VARCHAR) + is Int -> ps.setInt(idx + 1, v) + is Instant -> ps.setTimestamp(idx + 1, v.toSqlTimestamp()) + else -> ps.setString(idx + 1, v.toString()) + } + } + }, + ) + } + private fun bindMain( ps: java.sql.PreparedStatement, snapshot: FlightSnapshot, @@ -895,6 +1016,8 @@ class JdbcFlightStateRepository( "flight_delay", "flight_bridge_op", "flight_chock_op", "flight_route_point", ) + internal val SEGMENT_TABLES = listOf("flight_srvt", "flight_vipf") + /** 报文里的 10 类集合分别存到哪张表、哪些列(列名与基线脚本一致)。 */ internal val COLLECTIONS: Map = mapOf( "GTDT" to DetailSpec("flight_gate", listOf("gate", "pgot", "pgct", "gotm", "gctm", "gtyp"), "GTNO"), diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt index b211210..9914c52 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt @@ -356,9 +356,16 @@ class StubMsgEvents : MsgEventRepository { class StubFlightState : FlightStateRepository { val mains = linkedMapOf() val snapshots = linkedMapOf() // 内存全量态(主行+明细一体) + val srvtByFlid = linkedMapOf>>() + val vipfByFlid = linkedMapOf>>() + val unmappedByMsgId = linkedMapOf>() fun clear() { - mains.clear(); snapshots.clear() + mains.clear() + snapshots.clear() + srvtByFlid.clear() + vipfByFlid.clear() + unmappedByMsgId.clear() } override fun findMainRow(flid: String): FlightMainRow? = mains[flid] @@ -395,6 +402,24 @@ class StubFlightState : FlightStateRepository { return if (existing == null) PersistOutcome.INSERTED else PersistOutcome.UPDATED } + override fun persistSrvtVipfIfPresent( + flid: String, + collections: Map>>, + now: Instant, + ) { + if ("SRVT" in collections) srvtByFlid[flid] = collections.getValue("SRVT") + if ("VIPF" in collections) vipfByFlid[flid] = collections.getValue("VIPF") + } + + override fun replaceUnmappedForMessage( + msgId: Long, + flid: String, + fields: List, + now: Instant, + ) { + unmappedByMsgId[msgId] = fields + } + override fun markDeleted(flid: String, msgId: Long, now: Instant): Boolean { val main = mains[flid] ?: return false if (main.state != FlightState.ACTIVE) return false diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt index aa1476e..d4d32e4 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt @@ -45,6 +45,13 @@ class FlopProcessor( val change = MergeChange(flid = payload.flid, scalars = payload.scalars, collections = payload.collections) val next = FlightStateEngine.mergedState(current, change) flightState.persistFullState(next, msgId = head.msgId, now = clock.instant()) + flightState.persistAuxiliaryFlightFields( + head.msgId, + payload.flid, + payload.collections, + payload.unmapped, + clock.instant(), + ) msgEvents.insertAll(eventsFor(next, head.msgId, mapper, clock.instant())) DomainOutcome(ApplyResult.Succeeded, listOf(projectionOf(next, mapper))) } @@ -117,6 +124,7 @@ class AdftProcessor( if (current != null) { val next = FlightStateEngine.mergedState(current, setOnly(record)) flightState.persistFullState(next, msgId = head.msgId, now = clock.instant()) + flightState.persistSrvtVipfIfPresent(record.flid, record.collections, clock.instant()) msgEvents.insertAll(eventsFor(next, head.msgId, mapper, clock.instant())) revived = next } @@ -148,6 +156,7 @@ class AdftProcessor( FlightStateEngine.mergedState(current, setOnly(record)) } val outcome = flightState.persistFullState(next, msgId = head.msgId, now = clock.instant()) + flightState.persistSrvtVipfIfPresent(record.flid, record.collections, clock.instant()) if (outcome == PersistOutcome.DAY_GUARD_VIOLATION) { // 运营日冲突是**协议级**问题:与 SCHD 同分类 —— 整笔回滚、不重试、交人工。 // 若按 INFRA 抛出去,会被当成暂时性故障白白重试到耗尽,并给出误导的错误类别。 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/FlightSegments.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/FlightSegments.kt new file mode 100644 index 0000000..5effd1b --- /dev/null +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/FlightSegments.kt @@ -0,0 +1,17 @@ +package com.gzzn.omms.msgexchange.processing + +import com.gzzn.omms.msgexchange.domain.flight.UnmappedFieldEntry +import com.gzzn.omms.msgexchange.infra.persistence.FlightStateRepository +import java.time.Instant + +/** `SRVT`/`VIPF` 与未映射字段的落库辅助(与航班主事务同批调用)。 */ +internal fun FlightStateRepository.persistAuxiliaryFlightFields( + msgId: Long, + flid: String, + collections: Map>>, + unmapped: List, + now: Instant, +) { + persistSrvtVipfIfPresent(flid, collections, now) + replaceUnmappedForMessage(msgId, flid, unmapped, now) +} 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 9a45f0d..24f3114 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -167,12 +167,10 @@ class MessageProcessor( } } - // [G-SRVT-VIPF]:SRVT/VIPF 段只保留在解码载荷里,尚未落明细表(清空语义待 Q2)。 - // 计数 + 告警替代此前的静默丢弃;出现即证明真实报文携带该段,可作为定案依据。 + // [G-SRVT-VIPF]:段出现计数,供真实流量观测;明细落库见 `G-SRVT-VIPF`(缺席是否清除待 Q2)。 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 重试不重绑 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt index f016bbd..8e329bd 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt @@ -141,6 +141,7 @@ class ScheduleProcessor( throw ProtocolViolation("operation-day guard violated flid=$flid") else -> batchWritten++ } + flightState.persistSrvtVipfIfPresent(flid, record.collections, clock.instant()) events += eventsFor(next, head.msgId, mapper, clock.instant()) projections += projectionOf(next, mapper) } diff --git a/src/main/resources/db/migration/V5__unmapped_field_srvt_vipf.sql b/src/main/resources/db/migration/V5__unmapped_field_srvt_vipf.sql new file mode 100644 index 0000000..ff444d9 --- /dev/null +++ b/src/main/resources/db/migration/V5__unmapped_field_srvt_vipf.sql @@ -0,0 +1,47 @@ +-- UNMAPPED_FIELD(US-05 AC3):未落入航班当前态的字段原值,按消息留存,不随 PROC_STATE 清理。 +-- SRVT/VIPF 专用明细(G-SRVT-VIPF):缺席段不清明细(Q2 默认)。 + +CREATE TABLE UNMAPPED_FIELD ( + MSG_ID BIGINT NOT NULL, + FLID VARCHAR(32) NOT NULL, + PATH VARCHAR(256) NOT NULL, + ORDINAL INT NOT NULL, + RAW_VALUE VARCHAR(4000) NOT NULL, + CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + PRIMARY KEY (MSG_ID, PATH, ORDINAL) +); +CREATE INDEX idx_unmapped_field_flid ON UNMAPPED_FIELD (FLID); + +CREATE TABLE FLIGHT_SRVT ( + FLID VARCHAR(32) NOT NULL, + ORDINAL INT NOT NULL, + OPER VARCHAR(64), + SRTC VARCHAR(64), + SRQT VARCHAR(64), + SRST VARCHAR(64), + SRET VARCHAR(64), + SRPR VARCHAR(64), + SANR VARCHAR(64), + SARR VARCHAR(64), + CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + PRIMARY KEY (FLID, ORDINAL) +); +CREATE INDEX idx_flight_srvt_flid ON FLIGHT_SRVT (FLID); + +CREATE TABLE FLIGHT_VIPF ( + FLID VARCHAR(32) NOT NULL, + ORDINAL INT NOT NULL, + OPER VARCHAR(64), + VPCD VARCHAR(64), + VFES VARCHAR(64), + VIPT_OPER VARCHAR(64), + VIPT_VSCD VARCHAR(64), + VIPT_VTQY VARCHAR(64), + VIPT_VTST VARCHAR(64), + VIPT_VTET VARCHAR(64), + CREATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL, + PRIMARY KEY (FLID, ORDINAL) +); +CREATE INDEX idx_flight_vipf_flid ON FLIGHT_VIPF (FLID); diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/SrvtVipfUnmappedTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/SrvtVipfUnmappedTest.kt new file mode 100644 index 0000000..b4df5c5 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/SrvtVipfUnmappedTest.kt @@ -0,0 +1,148 @@ +package com.gzzn.omms.msgexchange.processing + +import com.fasterxml.jackson.databind.ObjectMapper +import com.gzzn.omms.msgexchange.codec.FlopPayload +import com.gzzn.omms.msgexchange.codec.JacksonXmlCodec +import com.gzzn.omms.msgexchange.domain.DecodedMessage +import com.gzzn.omms.msgexchange.domain.MetaFields +import com.gzzn.omms.msgexchange.domain.MsgKind +import com.gzzn.omms.msgexchange.domain.ProcState +import com.gzzn.omms.msgexchange.domain.ProcStatus +import com.gzzn.omms.msgexchange.domain.flight.FlightState +import com.gzzn.omms.msgexchange.domain.flight.UnmappedFieldEntry +import com.gzzn.omms.msgexchange.infra.stub.StubFlightProjectionPort +import com.gzzn.omms.msgexchange.infra.stub.StubFlightState +import com.gzzn.omms.msgexchange.infra.stub.StubMsgEvents +import com.gzzn.omms.msgexchange.infra.stub.StubPipelineLock +import com.gzzn.omms.msgexchange.infra.stub.StubPipelineTx +import com.gzzn.omms.msgexchange.infra.stub.StubProcState +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test +import java.time.LocalDate + +class SrvtVipfUnmappedTest { + + private val msgId = 42L + private val flid = "900001" + + private fun head() = ProcState(msgId, ProcStatus.PENDING, updatedAt = java.time.Instant.EPOCH) + + private fun flopProcessor(flights: StubFlightState): FlopProcessor { + val commit = FlightCommit( + StubPipelineTx(), + StubPipelineLock(), + StubProcState(), + StubMsgEvents(), + StubFlightProjectionPort(), + java.time.Clock.systemUTC(), + ) + return FlopProcessor(commit, flights, StubMsgEvents(), ObjectMapper(), java.time.Clock.systemUTC()) + } + + private fun seedFlight(flights: StubFlightState) { + flights.persistFullState( + com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot( + flid, + LocalDate.of(2026, 12, 15), + FlightState.ACTIVE, + 1, + mapOf("FLNO" to "CA101"), + emptyMap(), + ), + msgId = 1, + now = java.time.Instant.EPOCH, + ) + } + + @Test + fun `SRVT present writes detail absent keeps prior rows`() { + val flights = StubFlightState() + seedFlight(flights) + val proc = flopProcessor(flights) + val srvt = listOf(mapOf("OPER" to "ADD", "SRTC" to "SR01")) + proc.apply( + head(), + flopMsg("GTDT"), + FlopPayload(flid, collections = mapOf("SRVT" to srvt)), + ) + assertEquals(srvt, flights.srvtByFlid[flid]) + + proc.apply( + head(), + flopMsg("GTDT"), + FlopPayload(flid, collections = mapOf("GTDT" to listOf(mapOf("GTNO" to "1", "GATE" to "D01")))), + ) + assertEquals(srvt, flights.srvtByFlid[flid], "absent SRVT must not clear Q2 default") + } + + @Test + fun `VIPF present and reprocess replaces segment idempotently`() { + val flights = StubFlightState() + seedFlight(flights) + val proc = flopProcessor(flights) + val vipf = listOf(mapOf("VPCD" to "V001", "VFES" to "3")) + val payload = FlopPayload(flid, collections = mapOf("VIPF" to vipf)) + proc.apply(head(), flopMsg("GTDT"), payload) + proc.apply(head(), flopMsg("GTDT"), payload) + assertEquals(1, flights.vipfByFlid[flid]!!.size) + assertEquals("V001", flights.vipfByFlid[flid]!![0]["VPCD"]) + } + + @Test + fun `whitelisted FLOP with FDIV stores unmapped fields`() { + val codec = JacksonXmlCodec() + val raw = """ + AODB120021010090311 + FLOPGTDT + $flid + D08 + note + + """.trimIndent() + val decoded = (codec.decode(raw) as com.gzzn.omms.msgexchange.codec.DecodeResult.Ok).message + val body = decoded.body as FlopPayload + assertTrue(body.unmapped.any { it.path == "FDIV/@DDES" && it.raw == "CTU" }) + assertTrue(body.unmapped.any { it.path == "FDIV/@DDIR" && it.raw == "TO" }) + assertTrue(body.unmapped.any { it.path == "FDIV" && it.raw == "note" }) + + val flights = StubFlightState() + seedFlight(flights) + flopProcessor(flights).apply(head(), decoded, body) + assertEquals(body.unmapped, flights.unmappedByMsgId[msgId]) + } + + @Test + fun `reprocess rewrites same unmapped rows`() { + val flights = StubFlightState() + seedFlight(flights) + val proc = flopProcessor(flights) + val unmapped = listOf( + UnmappedFieldEntry("FRET/@REID", 1, "G"), + UnmappedFieldEntry("FRET", 1, ""), + ) + val payload = FlopPayload(flid, unmapped = unmapped) + proc.apply(head(), flopMsg("DELY"), payload) + proc.apply(head(), flopMsg("DELY"), payload) + assertEquals(unmapped, flights.unmappedByMsgId[msgId]) + } + + @Test + fun `FLOP FDIV subtype stays unsupported not unmapped`() { + val codec = JacksonXmlCodec() + val raw = """ + AODB220021010090311 + FLOPFDIV + $flid + """.trimIndent() + val kind = (codec.decode(raw) as com.gzzn.omms.msgexchange.codec.DecodeResult.Ok).message.kind + assertEquals(MsgKind.Unsupported("FLOP-FDIV"), kind) + } + + private fun flopMsg(styp: String) = DecodedMessage( + meta = MetaFields("AODB", "FLOP", styp, 1L, 1L), + kind = MsgKind.Flop(styp), + rawXml = "", + body = FlopPayload(flid), + ) +}