From 1fc2b2c8958d80281ed20389cb6b308352378cfb Mon Sep 17 00:00:00 2001 From: windyboy Date: Tue, 8 Sep 2026 18:50:02 +0800 Subject: [PATCH] =?UTF-8?q?fix(persistence):=20=E5=B7=AE=E5=88=A0=E5=85=88?= =?UTF-8?q?=E6=8C=89=20FDAY=20=E5=9C=88=E5=AE=9A=E5=88=A0=E9=99=A4?= =?UTF-8?q?=E6=88=90=E5=91=98=EF=BC=8C=E4=BF=AE=E5=A4=8D=E4=BF=9D=E7=95=99?= =?UTF-8?q?=E8=88=AA=E7=8F=AD=E6=98=8E=E7=BB=86=E8=AF=AF=E5=88=A0=20(ACM2-?= =?UTF-8?q?30)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../domain/flight/FlightStateEngine.kt | 2 +- .../persistence/jdbc/FlightDetailTables.kt | 21 +-------- .../persistence/jdbc/JdbcPgRepositories.kt | 33 +++++--------- .../gzzn/omms/msgexchange/processing/Pump.kt | 2 +- .../msgexchange/processing/SnapshotFlow.kt | 2 +- .../processing/handlers/GtdtHandler.kt | 2 +- .../domain/flight/FlightStateEngineTest.kt | 10 ++--- .../persistence/jdbc/FlightSchdJdbcPgTest.kt | 45 ++++++++++++++----- .../omms/msgexchange/support/SeedHelpers.kt | 4 +- 9 files changed, 56 insertions(+), 65 deletions(-) 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 9623be8..e60c790 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 @@ -39,7 +39,7 @@ object FlightStateEngine { // 注意:DELY 无协议序号属性(DLNO 非法)→ 不支持 Apply;清除走空数组 Replace(空集) /** DNLD/FLOP 字段集 → 命令(出现即 Set/Replace;未出现即 Unchanged;序号 0 条目 = 显式清除标记)。 */ - fun commandsFromFields(flid: String, fields: FlightFields, snapshotReplace: Boolean): FlightFieldCommands { + fun commandsFromFields(flid: String, fields: FlightFields): FlightFieldCommands { val scalars = linkedMapOf() val collections = linkedMapOf() fields.forEach { (key, value) -> diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightDetailTables.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightDetailTables.kt index 6e5f472..0b60fb6 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightDetailTables.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightDetailTables.kt @@ -1,20 +1,14 @@ package com.gzzn.omms.msgexchange.infra.persistence.jdbc import com.gzzn.omms.msgexchange.domain.flight.FlightNextState -import java.sql.Timestamp import java.time.Instant import javax.sql.DataSource /** * v2 明细表读写(flight-state-design-v2 §3.2)。 - * 与宽表双写;读路径优先明细表(完整顺序/源序号保真)。 + * 主表保存标量,明细表保存重复集合;不存在旧槽位读写回退。 */ internal object FlightDetailTables { - /** 由明细表承载的集合键(v2 读写的权威来源)。 */ - val DETAIL_COLLECTION_KEYS: Set = setOf( - "GTDT", "CKDT", "CLDT", "PSDT", "CHDT", "DELY", "ABTM", "CHOT", "ROUT", "ERUT", - ) - private data class TableSpec( val table: String, val collectionKey: String, @@ -165,17 +159,4 @@ internal object FlightDetailTables { return out } - fun hasAnyDetailRows(ds: DataSource, flid: String): Boolean { - for (spec in SPECS) { - val count = ds.queryOne( - "SELECT 1 FROM ${spec.table} WHERE flid = ? LIMIT 1", - { ps -> ps.setString(1, flid) }, - ) { 1 } - if (count != null) return true - } - return ds.queryOne( - "SELECT 1 FROM flight_route_point WHERE flid = ? LIMIT 1", - { ps -> ps.setString(1, flid) }, - ) { 1 } != null - } } 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 9880376..a818262 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 @@ -396,26 +396,10 @@ class JdbcFlightSchdRepository( companion object { /** SCHD.FLTR 标量列 + legacy 派生列(与 V1.1.0 迁移一致,全部可空 VARCHAR)。 */ - private val SCALAR_COLUMNS: List = listOf( - "ALCD", "ALSC", "FLNO", "MVIN", "SODT", "FLTY", "FLIN", "ACFT", "RENO", - "TAOP", "TAFL", "TAID", "TRML", "MAXP", "CSOP", "CSFT", "MAID", - "ESTT", "ACTT", "STND", "PHAG", "CNCL", "REMC", "BOTM", "LACL", "FINT", - "APPT", "EGSR", "EGST", "FHAG", "MHAG", "VIPP", "VIPR", "LBNO", "LBWT", - "PAXC", "EXSC", "EXSR", "FTSS", "PEDT", "NEAT", "PADT", "NAAT", - "ABDG", "LPSDT", "ABN", - ) - - /** 异常明细前缀标量列(FDIV/FRET/FLAB:规范 1:0..1 单值异常)。 */ - private val EXCEPTION_COLUMNS: List = listOf( - "FDIV_DDES", "FDIV_DDIR", "FDIV_REMC", "FRET_REID", "FRET_RSN", - "FLAB_ARES", "FLAB_RSN", - ) - - /** 无界集合紧凑 JSON 数组字串列(键名去 _TEXT 后缀即 legacy 视图键,直存直取)。 */ - private val TEXT_COLUMNS: List = listOf("SRVT_TEXT", "VIPF_TEXT", "MAFL_TEXT") + private val SCALAR_COLUMNS: List = FlightSchdReadAssembler.SCALAR_COLUMNS /** v2 主表写入列:标量 + 异常 + 无界文本;集合由明细表承载(P4 起槽位/里程碑/航路列已退场)。 */ - private val SCALAR_WRITE_COLUMNS: List = SCALAR_COLUMNS + EXCEPTION_COLUMNS + TEXT_COLUMNS + private val SCALAR_WRITE_COLUMNS: List = FlightSchdReadAssembler.ALL_COLUMNS /** 库列名全集(读侧 SELECT 稳定顺序;P4 退场后 = 无损承载列,与 V1.4.0 后表结构一致)。 */ private val ALL_COLUMNS: List = FlightSchdReadAssembler.ALL_COLUMNS @@ -441,7 +425,7 @@ class JdbcFlightSchdRepository( return out } - /** FDIV/FRET/FLAB:1:0..1 单值异常对象 → 前缀标量列(自由文本取首个命中的文本键)。 */ + /** FDIV/FRET/FLAB:1:0..1 单值异常对象 → 前缀标量列(自由文本取首个命中的文本键)。 */ private fun flattenExceptions(fields: FlightFields, out: MutableMap) { fun flatten(key: String, mappings: Map, textTarget: String) { val raw = fields[key] ?: return @@ -462,11 +446,19 @@ class JdbcFlightSchdRepository( override fun deleteDiffByDay(day: String, delFlids: Collection): Int { if (delFlids.isEmpty()) return 0 - FlightDetailTables.deleteForFlids(ds, delFlids) var totalDeleted = 0 val sqlDate = toSqlDate(day) for (chunk in delFlids.chunked(200)) { val placeholders = chunk.joinToString(",") { "?" } + // 与主表使用相同日期域;跨日迁移和 FDAY=NULL 的航班必须连同明细保留。 + val ownedFlids = ds.query( + "SELECT flid FROM flight_schd WHERE fday = ? AND flid IN ($placeholders)", + { ps -> + ps.setDate(1, sqlDate) + chunk.forEachIndexed { i, flid -> ps.setString(i + 2, flid) } + }, + ) { rs -> rs.getString("flid") } + FlightDetailTables.deleteForFlids(ds, ownedFlids) val sql = "DELETE FROM flight_schd WHERE fday = ? AND flid IN ($placeholders)" totalDeleted += ds.update(sql) { ps -> ps.setDate(1, sqlDate) @@ -874,4 +866,3 @@ class JdbcReqTrackRepository( } } } - 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 5dd557d..70eca56 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -232,7 +232,7 @@ class MessageProcessor( if (decision.flightChanges.isNotEmpty()) { val nextStates = decision.flightChanges.map { change -> val current = flightSchd.findNextStateByFlid(change.flid) - val commands = FlightStateEngine.commandsFromFields(change.flid, change.fields, snapshotReplace = false) + val commands = FlightStateEngine.commandsFromFields(change.flid, change.fields) FlightStateEngine.apply(current, commands, messageId, bumpVersion = true) } flightSchd.persistNextStates(null, nextStates, snapshotReplace = false) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt index 20325d4..e366a7e 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/SnapshotFlow.kt @@ -78,7 +78,7 @@ class SnapshotFlow( // 锁内读取当前航班态,计算 nextState(单写者互斥下读到的一定是已提交最新态) val nextStates = normalized.map { (flid, fields) -> val current = flightSchd.findNextStateByFlid(flid) - val commands = FlightStateEngine.commandsFromFields(flid, fields, snapshotReplace = true) + val commands = FlightStateEngine.commandsFromFields(flid, fields) FlightStateEngine.apply(current, commands, messageId, bumpVersion = true) } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/handlers/GtdtHandler.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/handlers/GtdtHandler.kt index 707ad15..a0d099f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/handlers/GtdtHandler.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/handlers/GtdtHandler.kt @@ -32,7 +32,7 @@ class GtdtHandler( fields["GTDT"] = mapper.writeValueAsString(gtdtItems) val current = flightView[body.flid]?.let { FlightStateEngine.fromFlightFields(body.flid, it) } - val commands = FlightStateEngine.commandsFromFields(body.flid, fields, snapshotReplace = false) + val commands = FlightStateEngine.commandsFromFields(body.flid, fields) val preview = FlightStateEngine.apply(current, commands, messageId = "", bumpVersion = false) val payloadJson = FlightFieldsJson.toJson(preview.toFlightFields(mapper)) 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 1df0187..fa49691 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 @@ -112,7 +112,6 @@ class FlightStateEngineTest { val commands = FlightStateEngine.commandsFromFields( "F1", mapOf("GTDT" to """[{"GTNO":"0"}]"""), - snapshotReplace = false, ) assertEquals(CollectionCommand.Clear, commands.collections["GTDT"]) } @@ -123,7 +122,6 @@ class FlightStateEngineTest { FlightStateEngine.commandsFromFields( "F1", mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"0"}]"""), - snapshotReplace = false, ) }.exceptionOrNull() assertTrue(error is IllegalArgumentException) @@ -138,7 +136,7 @@ class FlightStateEngineTest { 1L, "m1", ) - val commands = FlightStateEngine.commandsFromFields("F1", mapOf("GTDT" to "[]"), snapshotReplace = false) + val commands = FlightStateEngine.commandsFromFields("F1", mapOf("GTDT" to "[]")) val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true) assertEquals(emptyList>(), next.collections["GTDT"]) } @@ -175,7 +173,7 @@ class FlightStateEngineTest { 3L, "m1", ) - val commands = FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to "null"), snapshotReplace = false) + val commands = FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to "null")) assertEquals(ScalarCommand.Clear, commands.scalars["FRET"]) val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true) @@ -186,7 +184,7 @@ class FlightStateEngineTest { // 非空载荷仍为 Set,且清空后的下一跳 Set 恢复键(clearedKeys 不跨消息携带) val restored = FlightStateEngine.apply( next, - FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to """{"REID":"R2"}"""), snapshotReplace = false), + FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to """{"REID":"R2"}""")), "m3", bumpVersion = true, ) @@ -197,7 +195,7 @@ class FlightStateEngineTest { @Test fun `exception clear accepts null literal empty object and empty string`() { for (payload in listOf("null", "{}", "")) { - val commands = FlightStateEngine.commandsFromFields("F1", mapOf("FDIV" to payload), snapshotReplace = false) + val commands = FlightStateEngine.commandsFromFields("F1", mapOf("FDIV" to payload)) assertEquals(ScalarCommand.Clear, commands.scalars["FDIV"], "payload=[$payload] must be Clear") } } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt index 2b27c4f..1169fd2 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlightSchdJdbcPgTest.kt @@ -89,12 +89,36 @@ class FlightSchdJdbcPgTest { ) = flights.map { (flid, fields) -> com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply( null, - com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields(flid, fields, snapshotReplace = true), + com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields(flid, fields), messageId, bumpVersion = false, ).copy(stateVersion = version) } + @Test + fun `day diff preserves all details of migrated and unowned flights`() { + val day = "2026-09-07" + val otherDay = "2026-09-08" + val fields = mapOf("GTDT" to """[{"GTNO":"3","GATE":"G23"}]""") + val removed = "TEST_DIFF_REMOVED" + val moved = "TEST_DIFF_MOVED" + val unowned = "TEST_DIFF_ADFT" + repo.persistNextStates(day, snapshotStates(day, listOf(removed to fields)), snapshotReplace = true) + repo.persistNextStates(otherDay, snapshotStates(otherDay, listOf(moved to fields)), snapshotReplace = true) + repo.persistNextStates(null, snapshotStates(day, listOf(unowned to fields)), snapshotReplace = false) + val beforeMoved = repo.findByFlid(moved) + val beforeUnowned = repo.findByFlid(unowned) + + JdbcPipelineTransactionManager(ds).inTransaction { + assertEquals(1, repo.deleteDiffByDay(day, listOf(removed, moved, unowned))) + } + + assertNull(repo.findByFlid(removed)) + assertTrue(FlightDetailTables.loadCollections(ds, removed).isEmpty()) + assertEquals(beforeMoved, repo.findByFlid(moved)) + assertEquals(beforeUnowned, repo.findByFlid(unowned)) + } + @Test fun `PG dialect - snapshot batch upsert, find, and domain diff delete`() { val day = "2026-09-07" @@ -142,7 +166,7 @@ class FlightSchdJdbcPgTest { com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply( null, com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields( - flid, mapOf("FLID" to flid, "REMC" to remc), snapshotReplace = false, + flid, mapOf("FLID" to flid, "REMC" to remc), ), "msg-inc-01", bumpVersion = true, @@ -183,7 +207,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( null, - engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"2","GATE":"A2"}]"""), snapshotReplace = true), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"2","GATE":"A2"}]""")), "msg-res-1", bumpVersion = true, ), ), @@ -206,7 +230,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( repo.findNextStateByFlid(flid), - engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"3","GATE":"B1"}]"""), snapshotReplace = false), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"3","GATE":"B1"}]""")), "msg-res-2", bumpVersion = true, ), ), @@ -223,7 +247,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( repo.findNextStateByFlid(flid), - engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"0"}]"""), snapshotReplace = false), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"0"}]""")), "msg-res-3", bumpVersion = true, ), ), @@ -251,7 +275,6 @@ class FlightSchdJdbcPgTest { "ABTM" to """[{"ASNO":"1","ABDG":"B01","ABOP":"A","AOTM":"07SEP261700"}]""", "FDIV" to """{"DDES":"PEK","DDIR":"TO","REMC":"weather"}""", ), - snapshotReplace = true, ), "msg-scalar-1", bumpVersion = true, ), @@ -304,7 +327,6 @@ class FlightSchdJdbcPgTest { "DELY" to "[]", "ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""", ), - snapshotReplace = false, ), "msg-scalar-2", bumpVersion = true, ), @@ -638,7 +660,6 @@ class FlightSchdJdbcPgTest { "CLDT" to """[{"CLNO":"1","BELT":"B01"}]""", "DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200"}]""", ), - snapshotReplace = true, ), "msg-display-1", bumpVersion = true, ), @@ -677,7 +698,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( repo.findNextStateByFlid(flid), - engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"5","GATE":"Z1"}]"""), snapshotReplace = false), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"5","GATE":"Z1"}]""")), "msg-display-2", bumpVersion = true, ), ), @@ -702,7 +723,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( null, - engine.commandsFromFields(flid, mapOf("FRET" to """{"REID":"CA108","RSN":"diverted"}"""), snapshotReplace = true), + engine.commandsFromFields(flid, mapOf("FRET" to """{"REID":"CA108","RSN":"diverted"}""")), "msg-exc-1", bumpVersion = true, ), ), @@ -720,7 +741,7 @@ class FlightSchdJdbcPgTest { listOf( engine.apply( repo.findNextStateByFlid(flid), - engine.commandsFromFields(flid, mapOf("FRET" to "null"), snapshotReplace = false), + engine.commandsFromFields(flid, mapOf("FRET" to "null")), "msg-exc-2", bumpVersion = true, ), ), @@ -755,7 +776,7 @@ class FlightSchdJdbcPgTest { "FDIV" to """{"DDES":"PEK","DDIR":"TO"}""", "SRVT" to """[{"SRTC":"WHEEL"}]""", ) - val state = engine.apply(null, engine.commandsFromFields(flid, fields, snapshotReplace = true), "msg-matrix-1", bumpVersion = true) + val state = engine.apply(null, engine.commandsFromFields(flid, fields), "msg-matrix-1", bumpVersion = true) repo.persistNextStates(day, listOf(state), snapshotReplace = true) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/support/SeedHelpers.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/support/SeedHelpers.kt index 8d81b36..6d780a8 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/support/SeedHelpers.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/support/SeedHelpers.kt @@ -18,7 +18,7 @@ fun seedSnapshot( flights.map { (flid, fields) -> FlightStateEngine.apply( null, - FlightStateEngine.commandsFromFields(flid, fields, snapshotReplace = true), + FlightStateEngine.commandsFromFields(flid, fields), "seed-$day", bumpVersion = false, ).copy(stateVersion = 1L) @@ -34,7 +34,7 @@ fun seedIncremental(repo: FlightSchdRepository, changes: List, mes changes.map { change -> FlightStateEngine.apply( repo.findNextStateByFlid(change.flid), - FlightStateEngine.commandsFromFields(change.flid, change.fields, snapshotReplace = false), + FlightStateEngine.commandsFromFields(change.flid, change.fields), messageId, bumpVersion = true, )