From 44fbac0d8b3049188273a98ff83545e19cf6308a Mon Sep 17 00:00:00 2001 From: windyboy Date: Tue, 8 Sep 2026 12:11:57 +0800 Subject: [PATCH] =?UTF-8?q?feat(persistence):=20FLIGHT=5FSCHD=5FDISPLAY=20?= =?UTF-8?q?=E5=85=BC=E5=AE=B9=E8=A7=86=E5=9B=BE=20+=20=E7=94=9F=E4=BA=A7?= =?UTF-8?q?=E8=AF=BB=E5=88=87=E6=98=8E=E7=BB=86=E6=9D=83=E5=A8=81=20(ACM2-?= =?UTF-8?q?29=20P2-2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - V1.3.1 迁移:flight_schd_display 视图投影首个/第二个资源与总数(gate1/gate2/ gate_total、chkc1/chkc2/checkin_total、belt1/belt_total、dely_code/strt), 完整明细仍在各明细表可查;PG LATERAL 方言,11g 随 P3-B 单独提供 - 生产读路径 mapFlightRow 切 assembleDetailOnly:集合键仅由明细表重建, 不再回退 V1.2.0 槽位/里程碑列(过渡双读仅保留在 FS7 对拍工具内) - FlightStateEngine 正确性修复:序号属性 "0" 条目现在正确翻译为显式 Clear (原实现 items 非空走 Replace,会存出 GTNO=0 脏明细行);混入常规条目的 0 标记 fail fast 拒绝猜测(v2 §4 非法结构显式失败) - JDBC 用例全部切 v2 写路径(persistNextStates),删除 legacy 写路径断言; 新增:显示视图投影/替换同步用例、明细级联删除断言、引擎 0 标记语义单测 验证:MSGX_PG_PORT=5433 真实 PG ./gradlew test --rerun-tasks 97 用例 0 失败 0 跳过 --- .../domain/flight/FlightStateEngine.kt | 25 +- .../persistence/jdbc/JdbcPgRepositories.kt | 4 +- .../migration/V1.3.1__flight_schd_display.sql | 57 ++++ .../domain/flight/FlightStateEngineTest.kt | 36 ++ .../persistence/jdbc/FlightSchdJdbcPgTest.kt | 318 ++++++++++++++---- 5 files changed, 348 insertions(+), 92 deletions(-) create mode 100644 src/main/resources/db/migration/V1.3.1__flight_schd_display.sql 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 6c6a552..04ac662 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 @@ -27,7 +27,7 @@ object FlightStateEngine { "ERUT" to "RTNO", ) - /** DNLD/FLOP 字段集 → 命令(出现即 Set/Replace;未出现即 Unchanged)。 */ + /** DNLD/FLOP 字段集 → 命令(出现即 Set/Replace;未出现即 Unchanged;序号 0 条目 = 显式清除标记)。 */ fun commandsFromFields(flid: String, fields: FlightFields, snapshotReplace: Boolean): FlightFieldCommands { val scalars = linkedMapOf() val collections = linkedMapOf() @@ -35,15 +35,17 @@ object FlightStateEngine { if (key == "FLID") return@forEach when { key in COLLECTION_KEYS -> { - val items = parseCollection(value) val seqAttr = SEQ_ATTR[key] - if (items.isEmpty() && seqAttr != null && fieldsContainsClearOnly(value)) { - collections[key] = CollectionCommand.Clear - } else { - collections[key] = CollectionCommand.Replace(items) + val parsed = parseCollection(value) + val markers = parsed.filter { seqAttr != null && it[seqAttr] == "0" } + collections[key] = when { + markers.isEmpty() -> CollectionCommand.Replace(parsed) + markers.size == parsed.size -> CollectionCommand.Clear + else -> throw IllegalArgumentException( + "$key mixes explicit clear marker (seq=0) with regular items; refusing to guess", + ) } } - snapshotReplace -> scalars[key] = ScalarCommand.Set(value) else -> scalars[key] = ScalarCommand.Set(value) } } @@ -116,15 +118,6 @@ object FlightStateEngine { } } - private fun fieldsContainsClearOnly(raw: String): Boolean { - val node = runCatching { mapper.readTree(raw) }.getOrNull() ?: return false - if (node.isArray && node.size() == 1) { - val seq = node[0].fields().asSequence().firstOrNull { it.key.endsWith("NO") }?.value?.asText() - return seq == "0" - } - return false - } - /** 当前库态 → FlightNextState(供增量合并)。 */ fun fromFlightFields(flid: String, fields: FlightFields, stateVersion: Long = 0L, lastMessageId: String = ""): FlightNextState { val scalars = linkedMapOf() 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 5fd4fb8..52914c6 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 @@ -657,7 +657,7 @@ class JdbcFlightSchdRepository( return totalDeleted } - /** 行 → 字段视图:标量列直读,集合键由明细表或槽位列重建(FS7 双读对拍见 FlightSchdReadAssembler)。 */ + /** 行 → 字段视图:标量列直读,集合键由明细表重建(v2 §8:完整明细权威读,P2-2 起不再回退槽位列)。 */ private fun mapFlightRow(rs: java.sql.ResultSet): Pair { val flid = rs.getString("flid") val row = linkedMapOf() @@ -665,7 +665,7 @@ class JdbcFlightSchdRepository( row[column] = rs.getString(column) } val detailCollections = FlightDetailTables.loadCollections(ds, flid) - return flid to FlightSchdReadAssembler.assembleMerged(flid, row, detailCollections) + return flid to FlightSchdReadAssembler.assembleDetailOnly(flid, row, detailCollections) } override fun findByFlid(flid: String): FlightFields? = diff --git a/src/main/resources/db/migration/V1.3.1__flight_schd_display.sql b/src/main/resources/db/migration/V1.3.1__flight_schd_display.sql new file mode 100644 index 0000000..f383c5f --- /dev/null +++ b/src/main/resources/db/migration/V1.3.1__flight_schd_display.sql @@ -0,0 +1,57 @@ +-- flight-state-design-v2 §8 (ACM2-29 P2-2): compatibility read view. +-- 提供「首个/第二个资源 + 总数」的运营查询投影;完整明细仍在各明细表可查。 +-- PG 专属语法(LATERAL);Oracle 11g 方言随 P3-B 单独提供 location。 + +CREATE OR REPLACE VIEW flight_schd_display AS +SELECT + f.flid, + f.fday, + f.flno, + f.mvin, + f.sodt, + f.stnd, + f.abdg, + f.state_version, + f.last_message_id, + + g1.gate AS gate1, + g2.gate AS gate2, + COALESCE(g_total.cnt, 0) AS gate_total, + + c1.chkc AS chkc1, + c2.chkc AS chkc2, + COALESCE(c_total.cnt, 0) AS checkin_total, + + b1.belt AS belt1, + COALESCE(b_total.cnt, 0) AS belt_total, + + dly.code AS dely_code, + dly.strt AS dely_strt +FROM flight_schd f +LEFT JOIN LATERAL ( + SELECT gate FROM flight_gate WHERE flid = f.flid ORDER BY ordinal LIMIT 1 OFFSET 0 +) g1 ON TRUE +LEFT JOIN LATERAL ( + SELECT gate FROM flight_gate WHERE flid = f.flid ORDER BY ordinal LIMIT 1 OFFSET 1 +) g2 ON TRUE +LEFT JOIN LATERAL ( + SELECT count(*) AS cnt FROM flight_gate WHERE flid = f.flid +) g_total ON TRUE +LEFT JOIN LATERAL ( + SELECT chkc FROM flight_checkin WHERE flid = f.flid ORDER BY ordinal LIMIT 1 OFFSET 0 +) c1 ON TRUE +LEFT JOIN LATERAL ( + SELECT chkc FROM flight_checkin WHERE flid = f.flid ORDER BY ordinal LIMIT 1 OFFSET 1 +) c2 ON TRUE +LEFT JOIN LATERAL ( + SELECT count(*) AS cnt FROM flight_checkin WHERE flid = f.flid +) c_total ON TRUE +LEFT JOIN LATERAL ( + SELECT belt FROM flight_belt WHERE flid = f.flid ORDER BY ordinal LIMIT 1 OFFSET 0 +) b1 ON TRUE +LEFT JOIN LATERAL ( + SELECT count(*) AS cnt FROM flight_belt WHERE flid = f.flid +) b_total ON TRUE +LEFT JOIN LATERAL ( + SELECT code, strt FROM flight_delay WHERE flid = f.flid ORDER BY ordinal LIMIT 1 +) dly ON TRUE; 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 020e7ff..f0743da 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 @@ -106,4 +106,40 @@ class FlightStateEngineTest { ) assertEquals("NEW", next.collections["GTDT"]!![0]["GATE"]) } + + @Test + fun `zero sequence marker alone translates to explicit clear`() { + val commands = FlightStateEngine.commandsFromFields( + "F1", + mapOf("GTDT" to """[{"GTNO":"0"}]"""), + snapshotReplace = false, + ) + assertEquals(CollectionCommand.Clear, commands.collections["GTDT"]) + } + + @Test + fun `zero sequence marker mixed with regular items fails fast instead of guessing`() { + val error = runCatching { + FlightStateEngine.commandsFromFields( + "F1", + mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"0"}]"""), + snapshotReplace = false, + ) + }.exceptionOrNull() + assertTrue(error is IllegalArgumentException) + } + + @Test + fun `empty array replaces collection with empty set keeping key semantics`() { + val current = FlightNextState( + "F1", + emptyMap(), + mapOf("GTDT" to listOf(mapOf("GTNO" to "1", "GATE" to "A1"))), + 1L, + "m1", + ) + val commands = FlightStateEngine.commandsFromFields("F1", mapOf("GTDT" to "[]"), snapshotReplace = false) + val next = FlightStateEngine.apply(current, commands, "m2", bumpVersion = true) + assertEquals(emptyList>(), next.collections["GTDT"]) + } } 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 55aaaf0..4bdaf45 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 @@ -35,6 +35,7 @@ class FlightSchdJdbcPgTest { private lateinit var ds: HikariDataSource private lateinit var repo: JdbcFlightSchdRepository + private val mapper = com.fasterxml.jackson.databind.ObjectMapper() @BeforeEach fun setUp() { @@ -66,21 +67,42 @@ class FlightSchdJdbcPgTest { private fun cleanup() { try { + val testFlids = ds.query( + "SELECT flid FROM flight_schd WHERE flid LIKE 'TEST_%'", + {}, + ) { rs -> rs.getString("flid") } + if (testFlids.isNotEmpty()) { + com.gzzn.omms.msgexchange.infra.persistence.jdbc.FlightDetailTables.deleteForFlids(ds, testFlids) + } ds.update("DELETE FROM flight_schd WHERE flid LIKE 'TEST_%'") {} ds.update("DELETE FROM schd_gen WHERE fday >= '2026-09-01'") {} } catch (_: Exception) { } } + /** v2 写路径便捷构造:字段集 → 快照 nextState(版本顺序推进)。 */ + private fun snapshotStates( + day: String, + flights: List>>, + version: Long = 1L, + messageId: String = "msg-test", + ) = 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), + messageId, + bumpVersion = false, + ).copy(stateVersion = version) + } + @Test fun `PG dialect - snapshot batch upsert, find, and domain diff delete`() { val day = "2026-09-07" - // 构造 250 条记录跨越 batchSize 200 门限 + // 构造 250 条记录跨越 JDBC batch 200 门限 val flights = (1..250).map { i -> "TEST_FL_$i" to mapOf("FLID" to "TEST_FL_$i", "FLNO" to "CA$i") } - - repo.upsertSnapshotBatch(day, flights) + repo.persistNextStates(day, snapshotStates(day, flights), snapshotReplace = true) // 验证全量写入成功 val byDay = repo.findByDay(day) @@ -107,14 +129,27 @@ class FlightSchdJdbcPgTest { fun `PG dialect - incremental upsert preserves existing FDAY and sets NULL for new`() { val day = "2026-09-07" // 1. 快照写入 REG_01,FDAY 为 2026-09-07 - repo.upsertSnapshotBatch(day, listOf("TEST_REG_01" to mapOf("FLID" to "TEST_REG_01", "REMC" to "1"))) + repo.persistNextStates( + day, + snapshotStates(day, listOf("TEST_REG_01" to mapOf("FLID" to "TEST_REG_01", "REMC" to "1"))), + snapshotReplace = true, + ) // 2. 增量更新 REG_01(已有行)与 ADFT_01(新行) - val changes = listOf( - FlightChange("TEST_REG_01", mapOf("FLID" to "TEST_REG_01", "REMC" to "2")), - FlightChange("TEST_ADFT_01", mapOf("FLID" to "TEST_ADFT_01", "REMC" to "1")), + repo.persistNextStates( + null, + listOf("TEST_REG_01" to "2", "TEST_ADFT_01" to "1").map { (flid, remc) -> + 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, + ), + "msg-inc-01", + bumpVersion = true, + ) + }, + snapshotReplace = false, ) - repo.upsertIncremental(changes) // 3. 验证 REG_01 内容更新但 FDAY 依然保留为 2026-09-07 val reg01 = ds.queryOne("SELECT fday, remc FROM flight_schd WHERE flid = 'TEST_REG_01'", {}) { rs -> @@ -138,101 +173,150 @@ class FlightSchdJdbcPgTest { } @Test - fun `resource slots replace the whole collection and zero sequence clears it`() { + fun `resource collections replace as a whole and zero sequence clears them losslessly`() { val day = "2026-09-07" - val mapper = com.fasterxml.jackson.databind.ObjectMapper() - repo.upsertSnapshotBatch( + val flid = "TEST_RESOURCE" + val engine = com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine + // 快照:两个登机门 → 明细表两行(保序保 GTNO) + repo.persistNextStates( day, listOf( - "TEST_RESOURCE" to mapOf( - "GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"2","GATE":"A2"}]""", + engine.apply( + null, + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"1","GATE":"A1"},{"GTNO":"2","GATE":"A2"}]"""), snapshotReplace = true), + "msg-res-1", bumpVersion = true, ), ), + snapshotReplace = true, ) - // 零子表:两个登机门落平铺槽位列 - val slots = ds.queryOne( - "SELECT gate1, pgot1, gate2, pgot2 FROM flight_schd WHERE flid = 'TEST_RESOURCE'", - {}, - ) { rs -> listOf(rs.getString("gate1"), rs.getString("pgot1"), rs.getString("gate2"), rs.getString("pgot2")) } - assertEquals(listOf("A1", null, "A2", null), slots) - // 读侧视图重建 legacy 同构 GTDT 数组(槽位序号合成 GTNO) - val view = repo.findByFlid("TEST_RESOURCE")!!["GTDT"]?.let { mapper.readTree(it) } + val gates = ds.query( + "SELECT ordinal, source_seq, gate FROM flight_gate WHERE flid = ? ORDER BY ordinal", + { ps -> ps.setString(1, flid) }, + ) { rs -> Triple(rs.getInt("ordinal"), rs.getString("source_seq"), rs.getString("gate")) } + assertEquals(listOf(Triple(1, "1", "A1"), Triple(2, "2", "A2")), gates) + // 读侧明细权威视图:两个登机门完整读回 + val view = repo.findByFlid(flid)!!["GTDT"]?.let { mapper.readTree(it) } assertEquals(2, view!!.size()) assertEquals("A2", view[1]["GATE"].asText()) assertEquals("2", view[1]["GTNO"].asText()) - // 增量 = 单资源集合级全量快照替换:1 个登机门 → 槽位 1 覆盖、槽位 2 清空 - repo.upsertIncremental( - listOf(FlightChange("TEST_RESOURCE", mapOf("GTDT" to """[{"GTNO":"3","GATE":"B1"}]"""))), + // 增量 = 集合级全量替换:1 个登机门(非连续 GTNO=3 保真) + repo.persistNextStates( + null, + listOf( + engine.apply( + repo.findNextStateByFlid(flid), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"3","GATE":"B1"}]"""), snapshotReplace = false), + "msg-res-2", bumpVersion = true, + ), + ), + snapshotReplace = false, ) - val replaced = mapper.readTree(repo.findByFlid("TEST_RESOURCE")!!["GTDT"]) + val replaced = mapper.readTree(repo.findByFlid(flid)!!["GTDT"]) assertEquals(1, replaced.size()) assertEquals("B1", replaced[0]["GATE"].asText()) + assertEquals("3", replaced[0]["GTNO"].asText()) - // 序号属性 "0" = 显式删除标记 → 全槽位清空,视图无 GTDT 键 - repo.upsertIncremental( - listOf(FlightChange("TEST_RESOURCE", mapOf("GTDT" to """[{"GTNO":"0"}]"""))), + // 序号属性 "0" = 显式删除标记 → 集合清除,明细行删除,视图无 GTDT 键 + repo.persistNextStates( + null, + listOf( + engine.apply( + repo.findNextStateByFlid(flid), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"0"}]"""), snapshotReplace = false), + "msg-res-3", bumpVersion = true, + ), + ), + snapshotReplace = false, ) - assertNull(repo.findByFlid("TEST_RESOURCE")!!["GTDT"]) - val cleared = ds.queryOne( - "SELECT gate1, gate2 FROM flight_schd WHERE flid = 'TEST_RESOURCE'", - {}, - ) { rs -> listOf(rs.getString("gate1"), rs.getString("gate2")) } - assertEquals(listOf(null, null), cleared) + assertNull(repo.findByFlid(flid)!!["GTDT"]) + assertEquals(0, ds.query("SELECT COUNT(*) FROM flight_gate WHERE flid = ?", { ps -> ps.setString(1, flid) }) { rs -> rs.getInt(1) }.first()) + assertLegacyCollectionColumnsAllNull(flid) } @Test - fun `milestone delay route and exception collections flatten and rebuild the legacy view`() { + fun `milestone delay route and exception collections persist losslessly in v2 storage`() { val day = "2026-09-07" - val mapper = com.fasterxml.jackson.databind.ObjectMapper() - repo.upsertSnapshotBatch( + val flid = "TEST_SCALAR" + val engine = com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine + repo.persistNextStates( day, listOf( - "TEST_SCALAR" to mapOf( - "DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200","REMC":"Flight Delayed"}]""", - "ROUT" to """[{"RTNO":"1","APCD":"ORD","SCAT":"07SEP261125","SCDT":"07SEP261315"},{"RTNO":"2","APCD":"MSP"}]""", - "ABTM" to """[{"ASNO":"1","ABOP":"A","AOTM":"07SEP261700"}]""", - "FDIV" to """{"DDES":"PEK","DDIR":"TO","REMC":"weather"}""", + engine.apply( + null, + engine.commandsFromFields( + flid, + mapOf( + "DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200","REMC":"Flight Delayed"}]""", + "ROUT" to """[{"RTNO":"1","APCD":"ORD","SCAT":"07SEP261125","SCDT":"07SEP261315"},{"RTNO":"2","APCD":"MSP"}]""", + "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, ), ), + snapshotReplace = true, ) - // 库内为零子表平铺列:延误覆盖、紧凑航路字串、里程碑时刻、异常前缀标量 - val stored = ds.queryOne( - "SELECT dely_code, dely_strt, rout_path, abtm_a, abtm_d, fdiv_ddes, fdiv_remc FROM flight_schd WHERE flid = 'TEST_SCALAR'", - {}, - ) { rs -> - listOf( - rs.getString("dely_code"), rs.getString("dely_strt"), rs.getString("rout_path"), - rs.getString("abtm_a"), rs.getString("abtm_d"), rs.getString("fdiv_ddes"), rs.getString("fdiv_remc"), - ) - } + // 库内明细表保序保源序号:延误、航路(route_kind 区分)、靠桥操作 + val delay = ds.queryOne( + "SELECT source_seq, code, strt, remc FROM flight_delay WHERE flid = ?", + { ps -> ps.setString(1, flid) }, + ) { rs -> listOf(rs.getString("source_seq"), rs.getString("code"), rs.getString("strt"), rs.getString("remc")) } + assertEquals(listOf(null, "YY", "07SEP261605", "Flight Delayed"), delay) + val routes = ds.query( + "SELECT route_kind, source_seq, apcd, scdt FROM flight_route_point WHERE flid = ? ORDER BY ordinal", + { ps -> ps.setString(1, flid) }, + ) { rs -> listOf(rs.getString("route_kind"), rs.getString("source_seq"), rs.getString("apcd"), rs.getString("scdt")) } assertEquals( - listOf("YY", "07SEP261605", "ORD/07SEP261125/07SEP261315,MSP//", "07SEP261700", null, "PEK", "weather"), - stored, + listOf( + listOf("ROUT", "1", "ORD", "07SEP261315"), + listOf("ROUT", "2", "MSP", null), + ), + routes, ) - // 读侧视图重建 legacy 同构集合键 - val fields = repo.findByFlid("TEST_SCALAR")!! + val bridge = ds.queryOne( + "SELECT source_seq, abdg, abop, aotm FROM flight_bridge_op WHERE flid = ?", + { ps -> ps.setString(1, flid) }, + ) { rs -> listOf(rs.getString("source_seq"), rs.getString("abdg"), rs.getString("abop"), rs.getString("aotm")) } + assertEquals(listOf("1", "B01", "A", "07SEP261700"), bridge) + // 异常对象 = 主表前缀标量列;读侧视图重建 FDIV 键 + val fdiv = ds.queryOne( + "SELECT fdiv_ddes, fdiv_ddir, fdiv_remc FROM flight_schd WHERE flid = ?", + { ps -> ps.setString(1, flid) }, + ) { rs -> listOf(rs.getString("fdiv_ddes"), rs.getString("fdiv_ddir"), rs.getString("fdiv_remc")) } + assertEquals(listOf("PEK", "TO", "weather"), fdiv) + val fields = repo.findByFlid(flid)!! assertEquals("YY", mapper.readTree(fields["DELY"])[0]["CODE"].asText()) - val route = mapper.readTree(fields["ROUT"]) - assertEquals(2, route.size()) - assertEquals("MSP", route[1]["APCD"].asText()) - assertFalse(route[1].has("SCAT")) + assertEquals(2, mapper.readTree(fields["ROUT"]).size()) assertEquals("A", mapper.readTree(fields["ABTM"])[0]["ABOP"].asText()) assertEquals("PEK", mapper.readTree(fields["FDIV"])["DDES"].asText()) - // 增量:空延误数组 = 覆盖清空;航路整集合替换 - repo.upsertIncremental( + // 增量:空延误数组 = 集合替换为空 → 明细行删除;航路整集合替换 + repo.persistNextStates( + null, listOf( - FlightChange("TEST_SCALAR", mapOf("DELY" to "[]")), - FlightChange("TEST_SCALAR", mapOf("ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""")), + engine.apply( + repo.findNextStateByFlid(flid), + engine.commandsFromFields( + flid, + mapOf( + "DELY" to "[]", + "ROUT" to """[{"RTNO":"1","APCD":"CTU"}]""", + ), + snapshotReplace = false, + ), + "msg-scalar-2", bumpVersion = true, + ), ), + snapshotReplace = false, ) - val updated = ds.queryOne( - "SELECT dely_code, rout_path FROM flight_schd WHERE flid = 'TEST_SCALAR'", - {}, - ) { rs -> listOf(rs.getString("dely_code"), rs.getString("rout_path")) } - assertEquals(listOf(null, "CTU//"), updated) - assertNull(repo.findByFlid("TEST_SCALAR")!!["DELY"]) + assertEquals(0, ds.query("SELECT COUNT(*) FROM flight_delay WHERE flid = ?", { ps -> ps.setString(1, flid) }) { rs -> rs.getInt(1) }.first()) + assertNull(repo.findByFlid(flid)!!["DELY"]) + val replacedRoute = mapper.readTree(repo.findByFlid(flid)!!["ROUT"]) + assertEquals(1, replacedRoute.size()) + assertEquals("CTU", replacedRoute[0]["APCD"].asText()) } @Test @@ -300,7 +384,11 @@ class FlightSchdJdbcPgTest { val day = "2026-09-07" try { ds.withTransaction { - repo.upsertSnapshotBatch(day, listOf("TEST_TX_01" to mapOf("REMC" to "tx"))) + repo.persistNextStates( + day, + snapshotStates(day, listOf("TEST_TX_01" to mapOf("REMC" to "tx"))), + snapshotReplace = true, + ) repo.putGenIfVersion(day, 0L, FlightSchdRepository.GenMeta(day, 1L, setOf("TEST_TX_01"))) // 模拟事务内部抛出异常 throw IllegalStateException("forced-abort-for-rollback-test") @@ -318,13 +406,19 @@ class FlightSchdJdbcPgTest { fun `PG dialect - ES history sweep deleteByFlids idempotent batch execution`() { val day = "2026-09-07" val flids = (1..50).map { "TEST_SWEEP_$it" } - repo.upsertSnapshotBatch(day, flids.map { it to mapOf("FLID" to it) }) + repo.persistNextStates( + day, + snapshotStates(day, flids.map { it to mapOf("FLID" to it, "GTDT" to """[{"GTNO":"1","GATE":"A"}]""") }), + snapshotReplace = true, + ) assertEquals(50, repo.findByFlids(flids).size) + assertEquals(50, ds.query("SELECT COUNT(*) FROM flight_gate WHERE flid LIKE 'TEST_SWEEP_%'", {}) { rs -> rs.getInt(1) }.first()) - // 第一次分批删除 + // 第一次分批删除:主表与明细级联删除 val deleted1 = repo.deleteByFlids(flids.toSet()) assertEquals(50, deleted1) assertEquals(0, repo.findByFlids(flids).size) + assertEquals(0, ds.query("SELECT COUNT(*) FROM flight_gate WHERE flid LIKE 'TEST_SWEEP_%'", {}) { rs -> rs.getInt(1) }.first()) // 重放删除:安全返回 0,无副作用 val deleted2 = repo.deleteByFlids(flids.toSet()) @@ -340,7 +434,12 @@ class FlightSchdJdbcPgTest { val day = "2026-09-07" val fixedInstant = Instant.parse("2026-09-07T08:30:15.123456Z") - repo.upsertSnapshotBatch(day, listOf("TEST_TZ_01" to mapOf("FLID" to "TEST_TZ_01")), fixedInstant) + repo.persistNextStates( + day, + snapshotStates(day, listOf("TEST_TZ_01" to mapOf("FLID" to "TEST_TZ_01"))), + snapshotReplace = true, + now = fixedInstant, + ) repo.putGenIfVersion(day, 0L, FlightSchdRepository.GenMeta(day, 1L, setOf("TEST_TZ_01")), fixedInstant) // 东京时区下回读 @@ -523,6 +622,77 @@ class FlightSchdJdbcPgTest { assertEquals("m2", gen.lastMessageId) } + @Test + fun `PG dialect - flight_schd_display view projects first second resources and totals from detail tables`() { + val day = "2026-09-07" + val flid = "TEST_DISPLAY" + val engine = com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine + repo.persistNextStates( + day, + listOf( + engine.apply( + null, + engine.commandsFromFields( + flid, + mapOf( + "FLNO" to "CA display", + "GTDT" to """[{"GTNO":"1","GATE":"G28"},{"GTNO":"2","GATE":"G33"},{"GTNO":"3","GATE":"G23"}]""", + "CKDT" to """[{"CKNO":"1","CHKC":"01","CCLS":"Y"},{"CKNO":"2","CHKC":"17","CCLS":"F"}]""", + "CLDT" to """[{"CLNO":"1","BELT":"B01"}]""", + "DELY" to """[{"CODE":"YY","STRT":"07SEP261605","DURA":"0200"}]""", + ), + snapshotReplace = true, + ), + "msg-display-1", bumpVersion = true, + ), + ), + snapshotReplace = true, + ) + + val row = ds.queryOne( + """ + SELECT flno, gate1, gate2, gate_total, chkc1, chkc2, checkin_total, belt1, belt_total, dely_code, dely_strt + FROM flight_schd_display WHERE flid = ? + """.trimIndent(), + { ps -> ps.setString(1, flid) }, + ) { rs -> + listOf( + rs.getString("flno"), rs.getString("gate1"), rs.getString("gate2"), rs.getInt("gate_total"), + rs.getString("chkc1"), rs.getString("chkc2"), rs.getInt("checkin_total"), + rs.getString("belt1"), rs.getInt("belt_total"), + rs.getString("dely_code"), rs.getString("dely_strt"), + ) + } + // 视图提供首个/第二个资源 + 总数;完整明细仍在明细表可查(三门不丢失) + assertEquals( + listOf( + "CA display", "G28", "G33", 3, + "01", "17", 2, + "B01", 1, + "YY", "07SEP261605", + ), + row, + ) + + // 集合替换后视图随明细表同步(3 门 → 1 门) + repo.persistNextStates( + null, + listOf( + engine.apply( + repo.findNextStateByFlid(flid), + engine.commandsFromFields(flid, mapOf("GTDT" to """[{"GTNO":"5","GATE":"Z1"}]"""), snapshotReplace = false), + "msg-display-2", bumpVersion = true, + ), + ), + snapshotReplace = false, + ) + val after = ds.queryOne( + "SELECT gate1, gate2, gate_total FROM flight_schd_display WHERE flid = ?", + { ps -> ps.setString(1, flid) }, + ) { rs -> Triple(rs.getString("gate1"), rs.getString("gate2"), rs.getInt("gate_total")) } + assertEquals(Triple("Z1", null, 1), after) + } + private fun assertLegacyCollectionColumnsAllNull(flid: String) { val columnList = com.gzzn.omms.msgexchange.support.LEGACY_COLLECTION_STORAGE_COLUMNS.joinToString(", ") ds.queryOne(