fix(persistence): 差删先按 FDAY 圈定删除成员,修复保留航班明细误删 (ACM2-30)

This commit is contained in:
windyboy
2026-09-08 18:50:02 +08:00
parent 7c2d22e8f1
commit 1fc2b2c895
9 changed files with 56 additions and 65 deletions
@@ -39,7 +39,7 @@ object FlightStateEngine {
// 注意:DELY 无协议序号属性(DLNO 非法)→ 不支持 Apply;清除走空数组 Replace(空集) // 注意:DELY 无协议序号属性(DLNO 非法)→ 不支持 Apply;清除走空数组 Replace(空集)
/** DNLD/FLOP 字段集 → 命令(出现即 Set/Replace;未出现即 Unchanged;序号 0 条目 = 显式清除标记)。 */ /** 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<String, ScalarCommand>() val scalars = linkedMapOf<String, ScalarCommand>()
val collections = linkedMapOf<String, CollectionCommand>() val collections = linkedMapOf<String, CollectionCommand>()
fields.forEach { (key, value) -> fields.forEach { (key, value) ->
@@ -1,20 +1,14 @@
package com.gzzn.omms.msgexchange.infra.persistence.jdbc package com.gzzn.omms.msgexchange.infra.persistence.jdbc
import com.gzzn.omms.msgexchange.domain.flight.FlightNextState import com.gzzn.omms.msgexchange.domain.flight.FlightNextState
import java.sql.Timestamp
import java.time.Instant import java.time.Instant
import javax.sql.DataSource import javax.sql.DataSource
/** /**
* v2 明细表读写(flight-state-design-v2 §3.2)。 * v2 明细表读写(flight-state-design-v2 §3.2)。
* 与宽表双写;读路径优先明细表(完整顺序/源序号保真) * 主表保存标量,明细表保存重复集合;不存在旧槽位读写回退
*/ */
internal object FlightDetailTables { internal object FlightDetailTables {
/** 由明细表承载的集合键(v2 读写的权威来源)。 */
val DETAIL_COLLECTION_KEYS: Set<String> = setOf(
"GTDT", "CKDT", "CLDT", "PSDT", "CHDT", "DELY", "ABTM", "CHOT", "ROUT", "ERUT",
)
private data class TableSpec( private data class TableSpec(
val table: String, val table: String,
val collectionKey: String, val collectionKey: String,
@@ -165,17 +159,4 @@ internal object FlightDetailTables {
return out 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
}
} }
@@ -396,26 +396,10 @@ class JdbcFlightSchdRepository(
companion object { companion object {
/** SCHD.FLTR 标量列 + legacy 派生列(与 V1.1.0 迁移一致,全部可空 VARCHAR)。 */ /** SCHD.FLTR 标量列 + legacy 派生列(与 V1.1.0 迁移一致,全部可空 VARCHAR)。 */
private val SCALAR_COLUMNS: List<String> = listOf( private val SCALAR_COLUMNS: List<String> = FlightSchdReadAssembler.SCALAR_COLUMNS
"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<String> = listOf(
"FDIV_DDES", "FDIV_DDIR", "FDIV_REMC", "FRET_REID", "FRET_RSN",
"FLAB_ARES", "FLAB_RSN",
)
/** 无界集合紧凑 JSON 数组字串列(键名去 _TEXT 后缀即 legacy 视图键,直存直取)。 */
private val TEXT_COLUMNS: List<String> = listOf("SRVT_TEXT", "VIPF_TEXT", "MAFL_TEXT")
/** v2 主表写入列:标量 + 异常 + 无界文本;集合由明细表承载(P4 起槽位/里程碑/航路列已退场)。 */ /** v2 主表写入列:标量 + 异常 + 无界文本;集合由明细表承载(P4 起槽位/里程碑/航路列已退场)。 */
private val SCALAR_WRITE_COLUMNS: List<String> = SCALAR_COLUMNS + EXCEPTION_COLUMNS + TEXT_COLUMNS private val SCALAR_WRITE_COLUMNS: List<String> = FlightSchdReadAssembler.ALL_COLUMNS
/** 库列名全集(读侧 SELECT 稳定顺序;P4 退场后 = 无损承载列,与 V1.4.0 后表结构一致)。 */ /** 库列名全集(读侧 SELECT 稳定顺序;P4 退场后 = 无损承载列,与 V1.4.0 后表结构一致)。 */
private val ALL_COLUMNS: List<String> = FlightSchdReadAssembler.ALL_COLUMNS private val ALL_COLUMNS: List<String> = FlightSchdReadAssembler.ALL_COLUMNS
@@ -441,7 +425,7 @@ class JdbcFlightSchdRepository(
return out return out
} }
/** FDIV/FRET/FLAB1:0..1 单值异常对象 → 前缀标量列(自由文本取首个命中的文本键)。 */ /** FDIV/FRET/FLAB1:0..1 单值异常对象 → 前缀标量列(自由文本取首个命中的文本键)。 */
private fun flattenExceptions(fields: FlightFields, out: MutableMap<String, String?>) { private fun flattenExceptions(fields: FlightFields, out: MutableMap<String, String?>) {
fun flatten(key: String, mappings: Map<String, String>, textTarget: String) { fun flatten(key: String, mappings: Map<String, String>, textTarget: String) {
val raw = fields[key] ?: return val raw = fields[key] ?: return
@@ -462,11 +446,19 @@ class JdbcFlightSchdRepository(
override fun deleteDiffByDay(day: String, delFlids: Collection<String>): Int { override fun deleteDiffByDay(day: String, delFlids: Collection<String>): Int {
if (delFlids.isEmpty()) return 0 if (delFlids.isEmpty()) return 0
FlightDetailTables.deleteForFlids(ds, delFlids)
var totalDeleted = 0 var totalDeleted = 0
val sqlDate = toSqlDate(day) val sqlDate = toSqlDate(day)
for (chunk in delFlids.chunked(200)) { for (chunk in delFlids.chunked(200)) {
val placeholders = chunk.joinToString(",") { "?" } 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)" val sql = "DELETE FROM flight_schd WHERE fday = ? AND flid IN ($placeholders)"
totalDeleted += ds.update(sql) { ps -> totalDeleted += ds.update(sql) { ps ->
ps.setDate(1, sqlDate) ps.setDate(1, sqlDate)
@@ -874,4 +866,3 @@ class JdbcReqTrackRepository(
} }
} }
} }
@@ -232,7 +232,7 @@ class MessageProcessor(
if (decision.flightChanges.isNotEmpty()) { if (decision.flightChanges.isNotEmpty()) {
val nextStates = decision.flightChanges.map { change -> val nextStates = decision.flightChanges.map { change ->
val current = flightSchd.findNextStateByFlid(change.flid) 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) FlightStateEngine.apply(current, commands, messageId, bumpVersion = true)
} }
flightSchd.persistNextStates(null, nextStates, snapshotReplace = false) flightSchd.persistNextStates(null, nextStates, snapshotReplace = false)
@@ -78,7 +78,7 @@ class SnapshotFlow(
// 锁内读取当前航班态,计算 nextState(单写者互斥下读到的一定是已提交最新态) // 锁内读取当前航班态,计算 nextState(单写者互斥下读到的一定是已提交最新态)
val nextStates = normalized.map { (flid, fields) -> val nextStates = normalized.map { (flid, fields) ->
val current = flightSchd.findNextStateByFlid(flid) 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) FlightStateEngine.apply(current, commands, messageId, bumpVersion = true)
} }
@@ -32,7 +32,7 @@ class GtdtHandler(
fields["GTDT"] = mapper.writeValueAsString(gtdtItems) fields["GTDT"] = mapper.writeValueAsString(gtdtItems)
val current = flightView[body.flid]?.let { FlightStateEngine.fromFlightFields(body.flid, it) } 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 preview = FlightStateEngine.apply(current, commands, messageId = "", bumpVersion = false)
val payloadJson = FlightFieldsJson.toJson(preview.toFlightFields(mapper)) val payloadJson = FlightFieldsJson.toJson(preview.toFlightFields(mapper))
@@ -112,7 +112,6 @@ class FlightStateEngineTest {
val commands = FlightStateEngine.commandsFromFields( val commands = FlightStateEngine.commandsFromFields(
"F1", "F1",
mapOf("GTDT" to """[{"GTNO":"0"}]"""), mapOf("GTDT" to """[{"GTNO":"0"}]"""),
snapshotReplace = false,
) )
assertEquals(CollectionCommand.Clear, commands.collections["GTDT"]) assertEquals(CollectionCommand.Clear, commands.collections["GTDT"])
} }
@@ -123,7 +122,6 @@ class FlightStateEngineTest {
FlightStateEngine.commandsFromFields( FlightStateEngine.commandsFromFields(
"F1", "F1",
mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"0"}]"""), mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"0"}]"""),
snapshotReplace = false,
) )
}.exceptionOrNull() }.exceptionOrNull()
assertTrue(error is IllegalArgumentException) assertTrue(error is IllegalArgumentException)
@@ -138,7 +136,7 @@ class FlightStateEngineTest {
1L, 1L,
"m1", "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) val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true)
assertEquals(emptyList<Map<String, String>>(), next.collections["GTDT"]) assertEquals(emptyList<Map<String, String>>(), next.collections["GTDT"])
} }
@@ -175,7 +173,7 @@ class FlightStateEngineTest {
3L, 3L,
"m1", "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"]) assertEquals(ScalarCommand.Clear, commands.scalars["FRET"])
val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true) val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true)
@@ -186,7 +184,7 @@ class FlightStateEngineTest {
// 非空载荷仍为 Set,且清空后的下一跳 Set 恢复键(clearedKeys 不跨消息携带) // 非空载荷仍为 Set,且清空后的下一跳 Set 恢复键(clearedKeys 不跨消息携带)
val restored = FlightStateEngine.apply( val restored = FlightStateEngine.apply(
next, next,
FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to """{"REID":"R2"}"""), snapshotReplace = false), FlightStateEngine.commandsFromFields("F1", mapOf("FRET" to """{"REID":"R2"}""")),
"m3", "m3",
bumpVersion = true, bumpVersion = true,
) )
@@ -197,7 +195,7 @@ class FlightStateEngineTest {
@Test @Test
fun `exception clear accepts null literal empty object and empty string`() { fun `exception clear accepts null literal empty object and empty string`() {
for (payload in listOf("null", "{}", "")) { 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") assertEquals(ScalarCommand.Clear, commands.scalars["FDIV"], "payload=[$payload] must be Clear")
} }
} }
@@ -89,12 +89,36 @@ class FlightSchdJdbcPgTest {
) = flights.map { (flid, fields) -> ) = flights.map { (flid, fields) ->
com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply( com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply(
null, null,
com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields(flid, fields, snapshotReplace = true), com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields(flid, fields),
messageId, messageId,
bumpVersion = false, bumpVersion = false,
).copy(stateVersion = version) ).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 @Test
fun `PG dialect - snapshot batch upsert, find, and domain diff delete`() { fun `PG dialect - snapshot batch upsert, find, and domain diff delete`() {
val day = "2026-09-07" val day = "2026-09-07"
@@ -142,7 +166,7 @@ class FlightSchdJdbcPgTest {
com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply( com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.apply(
null, null,
com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.commandsFromFields( 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", "msg-inc-01",
bumpVersion = true, bumpVersion = true,
@@ -183,7 +207,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
null, 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, "msg-res-1", bumpVersion = true,
), ),
), ),
@@ -206,7 +230,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
repo.findNextStateByFlid(flid), 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, "msg-res-2", bumpVersion = true,
), ),
), ),
@@ -223,7 +247,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
repo.findNextStateByFlid(flid), 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, "msg-res-3", bumpVersion = true,
), ),
), ),
@@ -251,7 +275,6 @@ class FlightSchdJdbcPgTest {
"ABTM" to """[{"ASNO":"1","ABDG":"B01","ABOP":"A","AOTM":"07SEP261700"}]""", "ABTM" to """[{"ASNO":"1","ABDG":"B01","ABOP":"A","AOTM":"07SEP261700"}]""",
"FDIV" to """{"DDES":"PEK","DDIR":"TO","REMC":"weather"}""", "FDIV" to """{"DDES":"PEK","DDIR":"TO","REMC":"weather"}""",
), ),
snapshotReplace = true,
), ),
"msg-scalar-1", bumpVersion = true, "msg-scalar-1", bumpVersion = true,
), ),
@@ -304,7 +327,6 @@ class FlightSchdJdbcPgTest {
"DELY" to "[]", "DELY" to "[]",
"ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""", "ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""",
), ),
snapshotReplace = false,
), ),
"msg-scalar-2", bumpVersion = true, "msg-scalar-2", bumpVersion = true,
), ),
@@ -638,7 +660,6 @@ class FlightSchdJdbcPgTest {
"CLDT" to """[{"CLNO":"1","BELT":"B01"}]""", "CLDT" to """[{"CLNO":"1","BELT":"B01"}]""",
"DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200"}]""", "DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200"}]""",
), ),
snapshotReplace = true,
), ),
"msg-display-1", bumpVersion = true, "msg-display-1", bumpVersion = true,
), ),
@@ -677,7 +698,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
repo.findNextStateByFlid(flid), 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, "msg-display-2", bumpVersion = true,
), ),
), ),
@@ -702,7 +723,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
null, 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, "msg-exc-1", bumpVersion = true,
), ),
), ),
@@ -720,7 +741,7 @@ class FlightSchdJdbcPgTest {
listOf( listOf(
engine.apply( engine.apply(
repo.findNextStateByFlid(flid), repo.findNextStateByFlid(flid),
engine.commandsFromFields(flid, mapOf("FRET" to "null"), snapshotReplace = false), engine.commandsFromFields(flid, mapOf("FRET" to "null")),
"msg-exc-2", bumpVersion = true, "msg-exc-2", bumpVersion = true,
), ),
), ),
@@ -755,7 +776,7 @@ class FlightSchdJdbcPgTest {
"FDIV" to """{"DDES":"PEK","DDIR":"TO"}""", "FDIV" to """{"DDES":"PEK","DDIR":"TO"}""",
"SRVT" to """[{"SRTC":"WHEEL"}]""", "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) repo.persistNextStates(day, listOf(state), snapshotReplace = true)
@@ -18,7 +18,7 @@ fun seedSnapshot(
flights.map { (flid, fields) -> flights.map { (flid, fields) ->
FlightStateEngine.apply( FlightStateEngine.apply(
null, null,
FlightStateEngine.commandsFromFields(flid, fields, snapshotReplace = true), FlightStateEngine.commandsFromFields(flid, fields),
"seed-$day", "seed-$day",
bumpVersion = false, bumpVersion = false,
).copy(stateVersion = 1L) ).copy(stateVersion = 1L)
@@ -34,7 +34,7 @@ fun seedIncremental(repo: FlightSchdRepository, changes: List<FlightChange>, mes
changes.map { change -> changes.map { change ->
FlightStateEngine.apply( FlightStateEngine.apply(
repo.findNextStateByFlid(change.flid), repo.findNextStateByFlid(change.flid),
FlightStateEngine.commandsFromFields(change.flid, change.fields, snapshotReplace = false), FlightStateEngine.commandsFromFields(change.flid, change.fields),
messageId, messageId,
bumpVersion = true, bumpVersion = true,
) )