From 99a0f5f738c280f30b1082a55f2d5c2cf4890030 Mon Sep 17 00:00:00 2001 From: windyboy Date: Wed, 9 Sep 2026 22:17:04 +0800 Subject: [PATCH] =?UTF-8?q?refactor(flight-state):=20SCHD=20=E6=97=A5?= =?UTF-8?q?=E8=AE=A1=E5=88=92=E6=94=B6=E6=95=9B=E4=B8=BA=E7=BC=BA=E5=A4=B1?= =?UTF-8?q?=E4=BF=9D=E7=95=99=E5=90=88=E5=B9=B6=E8=AF=AD=E4=B9=89=E5=B9=B6?= =?UTF-8?q?=E5=AF=B9=E9=BD=90=E7=8E=B0=E8=A1=8C=E8=AE=BE=E8=AE=A1=E6=96=87?= =?UTF-8?q?=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 按现行 docs/flight-state.md(111 行版,§1-§7)全面对齐 domain 包及其消费方, 删除整套指向已退役长版文档(§5.x-§10)的引用与死代码: - 语义:FlightStateEngine.snapshotState 由"完整快照整体替换"改为 §3.1 合并语义 ——出现 Set/Replace、缺失保留、标量空串显式清空;DELETED 不被日计划恢复(§3.3)。 - 校验:validateMessage 移除从未接线的报文覆盖范围(scope)参数,只保留 §4 步骤 2 的声明数量/航班标识/运营日推导校验;ScheduleBody 删除 scopeStart/scopeEnd。 - 删除死代码:domain/Decision.kt、FlightModel 的 SnapshotPatch/UpsertOutcome、 ProcState.isTerminal/archivable、SnapshotFlag.SEQN_REGRESSION、 ScheduleRecord.seqn(及 wire FlightRecordXml.SEQN);codec 移除未消费 FFID。 - 删除未接线且引用已退役列(fday/last_message_id)的 SqlDialect 方言脚手架, oracle11g README 改为按 V1 现列重建的口径。 - 注释/测试:domain、processing、infra 仓储与 jobs/delivery、配置类及对应测试的 KDoc 章节引用全部对齐现行 flight-state.md/design.md;FlightStateEngineTest 重写为合并语义(62/62 通过)。 V1__flight_state_baseline.sql 保留原样(内容注释仍带旧章节号,改动会破坏 已应用迁移的 Flyway checksum,待重建基线或 V2 净迁移时收敛)。 --- docs/design.md | 2 +- .../omms/msgexchange/codec/SisMessageBody.kt | 5 - .../omms/msgexchange/codec/SisWireMapper.kt | 4 +- .../omms/msgexchange/config/HistoryProps.kt | 8 +- .../msgexchange/config/OperationDayProps.kt | 5 +- .../omms/msgexchange/delivery/Dispatcher.kt | 10 +- .../gzzn/omms/msgexchange/domain/Decision.kt | 31 ------- .../omms/msgexchange/domain/DecodedMessage.kt | 7 +- .../gzzn/omms/msgexchange/domain/MsgEvent.kt | 13 +-- .../omms/msgexchange/domain/OperationDay.kt | 9 +- .../gzzn/omms/msgexchange/domain/ProcState.kt | 18 ++-- .../omms/msgexchange/domain/SnapshotLog.kt | 15 ++- .../gzzn/omms/msgexchange/domain/Targets.kt | 6 +- .../msgexchange/domain/flight/FlightModel.kt | 52 +++++------ .../domain/flight/FlightStateEngine.kt | 74 ++++++++------- .../infra/persistence/Repositories.kt | 51 ++++++----- .../infra/persistence/dialect/SqlDialect.kt | 72 --------------- .../persistence/jdbc/JdbcPgRepositories.kt | 24 ++--- .../msgexchange/infra/stub/StubAdapters.kt | 2 +- .../infra/stub/StubRepositories.kt | 6 +- .../omms/msgexchange/jobs/HistorySweepJob.kt | 16 ++-- .../gzzn/omms/msgexchange/jobs/JobRunner.kt | 7 +- .../processing/DynamicProcessors.kt | 46 +++++----- .../gzzn/omms/msgexchange/processing/Pump.kt | 15 +-- .../processing/ScheduleProcessor.kt | 36 ++++---- .../db/migration/oracle11g/README.md | 14 +-- .../omms/msgexchange/PipelineSmokeTest.kt | 2 +- .../msgexchange/codec/JacksonXmlCodecTest.kt | 4 +- .../delivery/DispatcherTickTest.kt | 4 +- .../msgexchange/domain/OperationDayTest.kt | 2 +- .../domain/flight/FlightStateEngineTest.kt | 91 ++++++++++++++----- .../persistence/jdbc/FlywayMigrationTest.kt | 2 +- .../msgexchange/jobs/HistorySweepJobTest.kt | 10 +- .../processing/FdelAndAdftProcessorTest.kt | 10 +- .../processing/ScheduleProcessorTest.kt | 17 ++-- 35 files changed, 320 insertions(+), 370 deletions(-) delete mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt delete mode 100644 src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/dialect/SqlDialect.kt diff --git a/docs/design.md b/docs/design.md index c6b528b..dd415af 100644 --- a/docs/design.md +++ b/docs/design.md @@ -110,7 +110,7 @@ SNDR | TYPE | STYP | SEQN `SCHD-DNLD` 与 `SCHD-RESP` 共用 `ScheduleProcessor.applyScheduleRecords`: 1. **重放判定**:`PROC_STATE` 已存在成功终态 → 幂等成功,仅追加留痕,不重复写入。 -2. **整包校验**:声明记录数、记录范围与运营日归属等校验失败 → 整包 `DEAD(PROTOCOL)`,不写半包,既有状态保持不变。 +2. **整包校验**:声明记录数、航班标识与运营日推导等校验失败 → 整包 `DEAD(PROTOCOL)`,不写半包,既有状态保持不变。 3. **事务写入**:锁内按 `FLID` 点查归属日,发现同一航班跨运营日即整包回滚并 `DEAD(PROTOCOL)`;通过后合并写主表与资源明细。报文未携带的航班不因本次日计划报文被删除。 4. **提交结果**:同一事务保存 `KAFKA:schd` / `KAFKA:msg` 事件、置消息 `SUCCEEDED` 并预登记回填待办;事务提交后执行信箱回填,留痕在事务外追加。 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 1d9dd60..15e0c94 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisMessageBody.kt @@ -6,7 +6,6 @@ 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 java.time.LocalDate /** Stable payload boundary used by processing. It is deliberately not an XML model. */ data class FlopPayload( @@ -18,8 +17,6 @@ data class FlopPayload( data class ScheduleBody( val recsDeclared: Int, val records: List, - val scopeStart: LocalDate? = null, - val scopeEnd: LocalDate? = null, ) /** @@ -54,9 +51,7 @@ data class SchdXml( @JsonIgnoreProperties(ignoreUnknown = true) data class FlightRecordXml( @param:JacksonXmlProperty(localName = "FLID") val flid: String? = null, - @param:JacksonXmlProperty(localName = "FFID") val ffid: String? = null, @param:JacksonXmlProperty(localName = "FDEL") val fdel: FdelXml? = null, - @param:JacksonXmlProperty(localName = "SEQN") val seqn: Long? = null, @param:JacksonXmlProperty(localName = "ALCD") val alcd: String? = null, @param:JacksonXmlProperty(localName = "ALSC") val alsc: String? = null, @param:JacksonXmlProperty(localName = "FLNO") val flno: String? = null, 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 45ed807..0c33f70 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/codec/SisWireMapper.kt @@ -16,12 +16,12 @@ internal object SisWireMapper { private fun scheduleRecord(xml: FlightRecordXml): ScheduleRecord? = xml.flid?.trim()?.takeIf(String::isNotEmpty)?.let { - ScheduleRecord(it, xml.scalars(), xml.collections(), xml.seqn ?: 0L) + ScheduleRecord(it, xml.scalars(), xml.collections()) } private fun FlightRecordXml.scalars(): Map = linkedMapOf().apply { listOf( - "FFID" to ffid, "ALCD" to alcd, "ALSC" to alsc, "FLNO" to flno, "MVIN" to mvin, "SODT" to sodt, + "ALCD" to alcd, "ALSC" to alsc, "FLNO" to flno, "MVIN" to mvin, "SODT" to sodt, "FLTY" to flty, "FLIN" to flin, "ACFT" to acft, "RENO" to reno, "TAOP" to taop, "TAFL" to tafl, "TAID" to taid, "TRML" to trml, "MAXP" to maxp, "CSOP" to csop, "CSFT" to csft, "MAID" to maid, "ESTT" to estt, "ACTT" to actt, "STND" to stnd, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/HistoryProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/HistoryProps.kt index 8e34d70..e4c39ae 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/HistoryProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/HistoryProps.kt @@ -2,13 +2,13 @@ package com.gzzn.omms.msgexchange.config import io.micronaut.context.annotation.ConfigurationProperties -/** 生命周期与保留期(docs/flight-state.md §8)。窗口按机场时区计算。 */ +/** 生命周期与保留期(docs/flight-state.md §6;design.md §6.2)。窗口按机场时区计算。 */ @ConfigurationProperties("msgx.history") class HistoryProps { /** 已取消(CNCL 非空)超过 N 小时。 */ var cancelledHours: Long = 48 - /** 到港/离港终态(NAAT/NEAT,含义待术语表确认 §10)超过 N 小时。 */ + /** 到港/离港终态(NAAT/NEAT,含义待术语表确认 §6 开放项)超过 N 小时。 */ var terminalHours: Long = 48 /** STATE = DELETED 超过 N 小时。 */ @@ -17,9 +17,9 @@ class HistoryProps { /** 无终态字段:最后有效更新超过兜底期限(默认 7 天)。 */ var idleHours: Long = 24 * 7 - /** 留痕 SCHD_SNAP_LOG 保留天数(§8.3,按 (SCOPE_END, RECV_AT) 清理)。 */ + /** 留痕 SCHD_SNAP_LOG 保留天数(design.md §6.2,按 (SCOPE_END, RECV_AT) 清理)。 */ var snapLogRetentionDays: Long = 90 - /** 是否接通历史存储——未接通时 HISTORY_SWEEP 必须删 0 条(§8.2)。 */ + /** 是否接通历史存储——未接通时清理必须删 0 条(flight-state.md §6 红线)。 */ var historyStoreEnabled: Boolean = false } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/OperationDayProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/OperationDayProps.kt index 69e887a..b822e6e 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/OperationDayProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/OperationDayProps.kt @@ -3,8 +3,9 @@ package com.gzzn.omms.msgexchange.config import io.micronaut.context.annotation.ConfigurationProperties /** - * 运营日计算(docs/flight-state.md §3.5):由 SODT(ddMMMyyHHmm)与机场时区计算; - * 切日边界业务配置——不得假设等于接收日期或自然日零点(默认 0 点为占位,待业务确认 §10)。 + * 运营日计算(docs/flight-state.md §2.1):由 SODT(ddMMMyyHHmm)与机场时区按切日边界推导; + * OPERATION_DAY 不是接收日/落库日,一经确定不可变。切日边界业务配置, + * 不得假设等于自然日零点(默认 0 点为占位,业务口径待确认)。 */ @ConfigurationProperties("msgx.operation-day") class OperationDayProps { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt index 0787088..9454238 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/delivery/Dispatcher.kt @@ -23,7 +23,7 @@ interface DeliveryPort { */ fun sendKafkaSchd(topic: String, key: String, payloadJson: String) - /** TOMBSTONE:key=FLID、value=null——整态键缺失表示删除旧值(§7.3)。 */ + /** TOMBSTONE:key=FLID、value=null——整态键缺失表示删除旧值(flight-state.md §5)。 */ fun sendKafkaNull(topic: String, key: String) /** 连通性探测(健康检查用);默认 true,真实 Kafka 实装时覆写为 producer metadata 校验。 */ @@ -31,9 +31,9 @@ interface DeliveryPort { } /** - * 投递调度(docs/flight-state.md §7.3):逐条 KAFKA_MSG 严格 FIFO; + * 投递调度(docs/flight-state.md §5 + design.md §5.1/§5.2):逐条 KAFKA_MSG 严格 FIFO; * KAFKA_SCHD 走 flushSchd 批量——同一 FLID 未发事件按最新 STATE_VERSION 合并输出, - * TOMBSTONE 发 null 值消息。两主题间不保证顺序(§7.3)。 + * TOMBSTONE 发 null 值消息。两主题间不保证顺序(§5)。 * 批量闭环:队首退避未到期不 claim;发送失败整批 attempts+1 退避,达上限整批 DEAD/DLQ。 */ @Singleton @@ -91,7 +91,7 @@ class Dispatcher( } } - /** flushSchd:同 FLID 未发事件按最新 STATE_VERSION 合并(§7.3);TOMBSTONE 发 null。 */ + /** flushSchd:同 FLID 未发事件按最新 STATE_VERSION 合并(§5);TOMBSTONE 发 null。 */ internal fun flushSchd() { val batch = try { msgEvents.mergePendingSchd(props.schd.flushLimit) @@ -116,7 +116,7 @@ class Dispatcher( } val sentIds = batch.mapNotNull { it.eventId }.toSet() - failures.mapNotNull { it.eventId }.toSet() if (sentIds.isNotEmpty()) msgEvents.markAllSent(sentIds.toList()) - // 被最新版本合并压掉的未发事件同样关闭(§7.3:同 FLID 只按最新 STATE_VERSION 输出一次) + // 被最新版本合并压掉的未发事件同样关闭(§5:同 FLID 只按最新 STATE_VERSION 输出一次) val sentVersions = batch.associate { it.partitionKey to it.stateVersion } runCatching { while (true) { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt deleted file mode 100644 index 311580b..0000000 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Decision.kt +++ /dev/null @@ -1,31 +0,0 @@ -package com.gzzn.omms.msgexchange.domain - -import com.gzzn.omms.msgexchange.domain.flight.MergeChange -import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord -import java.time.LocalDate - -/** - * 处理决策(docs/flight-state.md)——解码与校验进,落库计划出,纯数据不触碰 DB/Kafka。 - * 事务边界、PIPELINE_LOCK、事件登记由处理层执行(§5.1/§6)。 - */ -sealed interface Decision { - /** - * SCHD DNLD/RESP:已过 §5.2 报文完整性校验的记录集。 - * §5.3 归属校验与 §5.4 upsert 在事务内执行(需读既有 OPERATION_DAY)。 - */ - data class Schedule( - val records: List, - /** 报文覆盖运营日范围(单日快照时两者相等;无法确定时为 null → 整包拒绝)。 */ - val scopeStart: LocalDate?, - val scopeEnd: LocalDate?, - ) : Decision - - /** FLOP 增量合并(§6.1)/ ADFT(§2.1 语义待确认,按 MergeChange 承载)。 */ - data class Dynamic(val change: MergeChange) : Decision - - /** FDEL 标记删除(§6.2)。 */ - data class Delete(val flid: String) : Decision - - /** 幂等无变化(重复/迟到等,记成功但不推进任何状态)。 */ - data object NoOp : Decision -} diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt index 167a898..28c2d4c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/DecodedMessage.kt @@ -1,7 +1,8 @@ package com.gzzn.omms.msgexchange.domain /** - * 解码后的入站报文(统一内部模型)。META 字段实名 SNDR/SEQN/DTTM(legacy META.java,I3 幂等键来源)。 + * 解码后的入站报文(docs/design.md §2.2 统一载荷)。META 字段实名 + * SNDR/SEQN/DTTM,其中 SNDR|TYPE|STYP|SEQN 是业务幂等键来源(Identity.of)。 */ data class MetaFields( val sndr: String, @@ -11,13 +12,13 @@ data class MetaFields( val dttm: Long, ) -/** 消息分派(sealed + 穷尽 when;FDEL 一等公民——终止航班实例 §6.2)。 */ +/** 消息分派(design.md §2.2:一等分派键,穷尽 when;FDEL 一等公民——终止航班实例)。 */ sealed interface MsgKind { data class Schd(val subtype: SchdSubtype) : MsgKind data class Flop(val subtype: String) : MsgKind // 运行动态 STYP(FDEL 除外) data object Fdel : MsgKind - /** 未支持类型(§9:FAILED(UNSUPPORTED),达阈值转 DEAD)。 */ + /** 未支持类型(design.md §2.3:FAILED(UNSUPPORTED),达阈值转 DEAD)。 */ data class Unsupported(val tag: String) : MsgKind enum class SchdSubtype { RESP, DNLD, ADFT } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/MsgEvent.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/MsgEvent.kt index 484e10c..52a2bb5 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/MsgEvent.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/MsgEvent.kt @@ -1,22 +1,23 @@ package com.gzzn.omms.msgexchange.domain -/** MSG_EVENT 事件形态(§7.3)。 */ +/** MSG_EVENT 事件形态(design.md §2.1:UPSERT 整态/通知,TOMBSTONE 删除)。 */ enum class EventType { UPSERT, TOMBSTONE } -/** MSG_EVENT 投递状态(outbox 状态机)。 */ +/** MSG_EVENT 投递状态机(design.md §2.3:PENDING → SENT;失败退避重试;耗尽转 DEAD 保留作 DLQ)。 */ enum class EventStatus { PENDING, SENT, DEAD } /** - * MSG_EVENT(outbox,§3.2):状态、变更、删除通知。 - * KAFKA_SCHD 整态 + KAFKA_MSG 变化通知;TOMBSTONE 仅在 ACTIVE→DELETED 时 - * 与删除同事务登记(§7.3),投递失败持续重试。 + * MSG_EVENT(outbox,design.md §2.1):状态、变更、删除通知。 + * KAFKA_SCHD 发整态、KAFKA_MSG 只通知变化(docs/flight-state.md §5,两主题不承诺顺序); + * TOMBSTONE 仅在 ACTIVE→DELETED(§3.3)或生命周期清理前补发(§6),与删除同事务登记。 + * stateVersion:发布时的航班版本;KAFKA_SCHD 聚合按 FLID 取最新(§5)。 */ data class MsgEvent( val eventId: Long? = null, val target: String, val partitionKey: String, // 恒为 FLID val eventType: EventType = EventType.UPSERT, - val stateVersion: Long = 0, // 发布时航班版本;Dispatcher 合并同 FLID 未发事件取最新(§7.3) + val stateVersion: Long = 0, val payloadJson: String, val state: EventStatus = EventStatus.PENDING, val attempts: Int = 0, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/OperationDay.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/OperationDay.kt index 301c7af..7836366 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/OperationDay.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/OperationDay.kt @@ -5,15 +5,16 @@ import java.time.LocalDateTime import java.time.ZoneId /** - * 运营日计算(docs/flight-state.md §3.5):由计划运行时间字段 SODT(ddMMMyyHHmm) - * 与机场时区计算;切日边界业务配置——默认 0 点为占位,待业务确认(§10)。 + * 运营日计算(docs/flight-state.md §2.1):由计划运行时间字段 SODT(ddMMMyyHHmm) + * 与机场时区按切日边界推导;OPERATION_DAY 不是消息接收日或落库日,一经确定不可变。 + * 切日边界由 msgx.operation-day.cutoff-hour 配置(默认 0 = 自然日零点,业务口径待确认)。 */ class OperationDayCalculator( zone: ZoneId, cutoffHour: Int, ) { private val zone: ZoneId = zone - /** 切日边界:SODT 本地时刻早于该小时的归属前一运营日(0–23,越界按 0)。 */ + /** 切日边界:SODT 本地时刻早于该小时的归属前一运营日(0–23,越界收敛到边界值)。 */ private val cutoffHour: Int = cutoffHour.coerceIn(0, 23) /** SODT → 运营日;输入 null/空/非法返回 null(不抛异常,调用方按"运营日不可计算"处理)。 */ @@ -28,7 +29,7 @@ class OperationDayCalculator( /** * SIS SODT 线格式:ddMMMyyHHmm(如 15DEC031723),月份英文三字母、大小写不敏感。 * 两位年显式按 2000 基准展开(java.time 的 yy reduced-value 解析跨实现不一致, - * 显式展开保证 AODB 侧年份窗口唯一口径;基准年待真实报文验收确认 §10)。 + * 显式展开保证 AODB 侧年份窗口唯一口径;基准年待真实报文验收确认)。 */ private val SODT_REGEX = Regex("(\\d{1,2})([A-Za-z]{3})(\\d{2})(\\d{2})(\\d{2})") diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt index 07805d1..6953365 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt @@ -1,27 +1,23 @@ package com.gzzn.omms.msgexchange.domain /** - * ACMA-8 数据模型 / PROC_STATE 状态机(docs/flight-state.md §3.2:每消息一行, - * 处理状态与重试结果;兼作快照重放判定 §5.1)。 + * PROC_STATE 处理伴生状态(docs/design.md §2.1/§2.3):每消息一行, + * MSG_ID = 信箱 ID 主键防重复入队;IDENTITY_KEY 唯一约束防业务重复; + * 处理状态机与错误分类同设计文档 §2.3。SUCCEEDED 终态兼作 SCHD 快照 + * 重放判定(docs/flight-state.md §4 步骤 / design.md §4.1)。 */ enum class ProcStatus { PENDING, FAILED, SUCCEEDED, SKIPPED, DEAD } -/** §9 错误分类:PROTOCOL = 整包拒绝(归属日不符等),不重试交人工。 */ +/** 错误分类(design.md §2.3):MALFORMED/PROTOCOL 直接 DEAD 不重试;其余退避重试。 */ enum class ErrorClass { MALFORMED, PROTOCOL, CODEC_ERROR, EXHAUSTED, INFRA, UNSUPPORTED } data class ProcState( val msgId: Long, val state: ProcStatus, - val identityKey: String? = null, // SNDR|TYPE|STYP|SEQN;decode 后首次绑定,FAILED 重试不重绑 + val identityKey: String? = null, // SNDR|TYPE|STYP|SEQN(design.md §2.2);decode 后首次绑定,FAILED 重试不重绑 val attempts: Int = 0, val nextAttemptAt: java.time.Instant? = null, val errorClass: ErrorClass? = null, val lastError: String? = null, val updatedAt: java.time.Instant = java.time.Instant.now(), -) { - val isTerminal: Boolean - get() = state == ProcStatus.SUCCEEDED || state == ProcStatus.SKIPPED || state == ProcStatus.DEAD - - /** 终态皆可归档;PENDING/FAILED 不迁。 */ - val archivable: Boolean get() = isTerminal -} +) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/SnapshotLog.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/SnapshotLog.kt index 45dc91f..089a761 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/SnapshotLog.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/SnapshotLog.kt @@ -4,13 +4,18 @@ import java.time.Instant import java.time.LocalDate /** - * SCHD 快照留痕模型(docs/flight-state.md §5.5): - * RESULT 与 FLAGS 分列(可「成功且告警」);一行 = 一次尝试,重放也记; - * 留痕不参与决策;写失败只记指标;保留 90 天,按 (SCOPE_END, RECV_AT) 清理。 + * SCHD 快照留痕 SCHD_SNAP_LOG(docs/flight-state.md §2 权威模型表; + * 清理规则见 design.md §6.2):只追加、可重建、不参与状态决策,写失败只记指标; + * 一行 = 一次尝试,重放也记;保留 90 天,按 (SCOPE_END, RECV_AT) 清理。 */ enum class SnapshotResult { COMMITTED, REPLAY_SKIPPED, ROLLED_BACK } -enum class SnapshotFlag { EMPTY, RECS_DROP, SEQN_REGRESSION, DAY_MISMATCH, SCHD_REVIVE_CONFLICT } +/** + * 留痕告警 flags:EMPTY(空快照合法)、RECS_DROP(声明数量不符)、 + * DAY_MISMATCH(记录运营日不可计算,整包拒绝)、 + * SCHD_REVIVE_CONFLICT(日计划命中 DELETED 航班,保持 DELETED 不恢复)。 + */ +enum class SnapshotFlag { EMPTY, RECS_DROP, DAY_MISMATCH, SCHD_REVIVE_CONFLICT } data class SnapshotLogEntry( val msgId: Long, @@ -22,5 +27,5 @@ data class SnapshotLogEntry( val durationMs: Long, val result: SnapshotResult, val flags: Set = emptySet(), - val archiveKey: String? = null, // 证据层引用(尚未交付 §3.2) + val archiveKey: String? = null, // 证据层引用(尚未交付 §2) ) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Targets.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Targets.kt index b1fe80f..2ddd75f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Targets.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/Targets.kt @@ -1,8 +1,10 @@ package com.gzzn.omms.msgexchange.domain /** - * MSG_EVENT 投递目标(§7.3:KAFKA_SCHD 发整态、KAFKA_MSG 只通知变化; - * 两主题间不保证顺序)。ES 投影属阶段 B,暂不登记目标。 + * MSG_EVENT 投递目标(docs/flight-state.md §5:KAFKA_SCHD 按 FLID 发完整状态投影、 + * KAFKA_MSG 只通知变化;两主题间不保证顺序;删除用 tombstone)。 + * 常量值即 MSG_EVENT.TARGET 落库值(design.md §2.1);topic 名由投递适配层映射 + * (KAFKA:msg → "msg",KAFKA:schd → "schd")。ES 投影属阶段 B,暂不登记目标。 */ object Targets { const val KAFKA_MSG = "KAFKA:msg" diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt index 82a47a1..6d0b954 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/flight/FlightModel.kt @@ -4,14 +4,15 @@ import java.time.Instant import java.time.LocalDate /** - * 航班实例当前态模型(docs/flight-state.md §3)。 + * 航班实例当前态模型(docs/flight-state.md §2/§3)。 * - * 身份:FLID 唯一关联键;OPERATION_DAY 一经确定不可变(§3.5/§5.3); - * STATE 仅 ACTIVE/DELETED(§3.1,无 ARCHIVED——物理清除只发生在历史归档成功之后 §8.2)。 + * 身份:FLID 是唯一关联键(§2.1);OPERATION_DAY 一经确定不可变(§2.1), + * 未被日计划收录前可为 NULL。STATE 仅 ACTIVE/DELETED(§3.3,无 ARCHIVED—— + * 物理清除只发生在历史归档成功之后 §6)。 */ enum class FlightState { ACTIVE, DELETED } -/** FLIGHT_SCHD 主行的身份与追踪字段(不含标量载荷)。 */ +/** FLIGHT_SCHD 主行的身份与追踪字段(§2 权威模型;不含标量载荷)。 */ data class FlightMainRow( val flid: String, val operationDay: LocalDate?, @@ -22,35 +23,19 @@ data class FlightMainRow( ) /** - * SCHD DNLD/RESP 单条 FLTR 记录(解码产物,§5 入口)。 - * scalars/collections 只含报文中出现的字段——完整快照语义下 - * 出现 = Set/Replace,未出现 = 清除/替换空集(§5.4)。 + * SCHD DNLD/RESP 单条 FLTR 记录(解码产物,§3.1 日计划入口)。 + * scalars/collections 只含报文中出现的字段:出现 = Set/Replace(标量空值 = 显式清空), + * 未出现 = 按合并语义保留本地值(§3.1/§7 不变量);集合按完整合并结果写入(§2.2)。 */ data class ScheduleRecord( val flid: String, val scalars: Map, val collections: Map>> = emptyMap(), - val seqn: Long = 0, ) /** - * 完整快照落库载荷(§5.1 步骤 6):按记录整体替换映射内字段, - * 仓储层负责明细先删后插与 `OPERATION_DAY` 不可变条件更新(§7.4)。 - */ -data class SnapshotPatch( - val flid: String, - val operationDay: LocalDate, - val scalars: Map, - val collections: Map>>, -) - -/** 快照 upsert 结果(§5.1 步骤 6:DELETED 航班保持 DELETED 并告警,不恢复)。 */ -enum class UpsertOutcome { INSERTED, UPDATED, REVIVE_CONFLICT } - -/** - * FLOP/ADFT 增量载荷(§6.1/§2.1)。 - * scalars 出现 = Set,缺失 = 保留本地值;collections 出现 = Replace - * (FLOP 集合语义未定案前沿用 Replace,§10 当前偏差明示)。 + * FLOP/ADFT 增量载荷(§3.2/§3.3):出现字段/集合更新,缺失保留本地值; + * ADFT 缺失字段语义待上游确认前按保守 Set-only 处理(§3.3),不沿用全量替换。 */ data class MergeChange( val flid: String, @@ -59,8 +44,9 @@ data class MergeChange( ) /** - * 完整当前态(§3.3):主行 + 全部明细 = 完整当前态; - * 读取须在一致性读事务中,且展示层过滤 STATE = ACTIVE(§9)。 + * 完整当前态(§2.2/§3):主行 + 全部明细 = 完整当前态。 + * 写入前在内存生成完整新状态再落库(引擎产物),读取须在主表与全部明细的 + * 一致性读边界内进行(§5),展示层过滤 STATE = ACTIVE。 */ data class FlightSnapshot( val flid: String, @@ -71,19 +57,23 @@ data class FlightSnapshot( val collections: Map>>, ) -/** §8.1 历史判定窗口(按机场时区计算;窗口值业务配置)。 */ +/** §6 生命周期判定窗口(按机场时区计算;窗口值业务配置)。 */ data class HistoryRules( val cancelledHours: Long = 48, // CNCL 非空超过 N 小时 - val terminalHours: Long = 48, // NAAT/NEAT 终态超过 N 小时(字段含义待术语表确认 §10) + val terminalHours: Long = 48, // NAAT/NEAT 终态超过 N 小时(字段含义待术语表确认 §6) val deletedHours: Long = 48, // STATE = DELETED 超过 N 小时 val idleHours: Long = 24 * 7, // 无终态字段:最后更新超过兜底期限 ) -/** §8.1 命中历史判定的航班(归档 → 物理清除候选)。 */ +/** + * §6 命中生命周期判定的航班(先归档 → 后物理清除;历史存储失败必须删 0 行)。 + * wasNeverFdel:是否未经 FDEL 而被生命周期清除——清除前须补发一次删除事件(§3.3/§5)。 + * 当前为推断口径(无持久化"曾 FDEL"痕迹):state=ACTIVE 视为未经 FDEL, + * FDEL→ADFT 重激活后的航班会被误判并重复补发,收敛见 §6 开放项。 + */ data class HistoryCandidate( val flid: String, val state: FlightState, val stateVersion: Long, - /** 是否未经 FDEL 而被生命周期清除——清除前须补发一次删除事件(§7.3/§8.2)。 */ val wasNeverFdel: Boolean, ) 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 6ec0fa4..ad8d634 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 @@ -5,33 +5,40 @@ import com.gzzn.omms.msgexchange.domain.SnapshotFlag import java.time.LocalDate /** - * 航班状态引擎(docs/flight-state.md §3.3/§5.2/§5.4/§6.1)——纯函数, - * 内存生成完整新状态再落库;不触碰 DB/Kafka。 + * 航班状态引擎(docs/flight-state.md §2–§4)——纯函数,内存推导完整新状态, + * 不触碰 DB/Kafka: + * - [validateMessage]:日计划整包校验(§4 步骤 2:声明数量、航班标识、运营日推导); + * - [snapshotState]:DNLD/RESP 日计划合并(§3.1)——出现 Set/Replace、缺失保留、 + * 标量显式空值清空;DELETED 航班保持 DELETED(§3.3,不承担恢复语义); + * - [mergedState]:FLOP/ADFT 增量合并(§3.2/§3.3)——只修改报文表达的字段/集合。 */ object FlightStateEngine { - /** 集合键白名单:10 类集合 ↔ 9 张明细表(ROUT/ERUT 共用 FLIGHT_ROUTE_POINT,§3.2)。 */ + /** 集合键白名单:10 类集合 ↔ 9 张明细表(8 张资源表 + FLIGHT_ROUTE_POINT, + * ROUT/ERUT 共用路线表,flight-state.md §2.2)。 */ val COLLECTION_KEYS: Set = setOf( "GTDT", "CKDT", "CLDT", "PSDT", "CHDT", "DELY", "ABTM", "CHOT", "ROUT", "ERUT", ) /** - * §5.2 报文完整性五项校验:RECS 0–9999 且等于实收数;每条含合法数字型 FLID; - * 快照内不重复;每条记录运营日可计算且在报文覆盖范围内。任一失败整包不落地。 + * 日计划整包校验(§4 步骤 2 / §3.1):RECS 0–9999 且等于实收数;每条含合法 + * 数字型 FLID(SIS Number(1-12));快照内不重复;每条记录运营日可计算。 + * 任一失败整包不落地(DEAD(PROTOCOL))。 */ fun validateMessage( recsDeclared: Int, records: List, - scopeStart: LocalDate?, - scopeEnd: LocalDate?, opDay: OperationDayCalculator, ): SnapshotValidation { - val flags = linkedSetOf() - if (recsDeclared !in 0..9999) return SnapshotValidation.Invalid("RECS out of range: $recsDeclared", setOf(SnapshotFlag.RECS_DROP)) - if (records.size != recsDeclared) return SnapshotValidation.Invalid( - "RECS ($recsDeclared) != received FLTR count (${records.size})", - setOf(SnapshotFlag.RECS_DROP), - ) + if (recsDeclared !in 0..9999) { + return SnapshotValidation.Invalid("RECS out of range: $recsDeclared", setOf(SnapshotFlag.RECS_DROP)) + } + if (records.size != recsDeclared) { + return SnapshotValidation.Invalid( + "RECS ($recsDeclared) != received FLTR count (${records.size})", + setOf(SnapshotFlag.RECS_DROP), + ) + } if (records.isEmpty()) return SnapshotValidation.Ok(emptyMap(), setOf(SnapshotFlag.EMPTY)) val perRecordDay = linkedMapOf() @@ -45,24 +52,22 @@ object FlightStateEngine { } val day = opDay.compute(record.scalars["SODT"]) if (day == null) { - return SnapshotValidation.Invalid("operation day not computable, flid=${record.flid}", setOf(SnapshotFlag.DAY_MISMATCH)) - } - val inScope = (scopeStart == null || !day.isBefore(scopeStart)) && - (scopeEnd == null || !day.isAfter(scopeEnd)) - if (!inScope) { return SnapshotValidation.Invalid( - "operation day $day outside coverage [$scopeStart, $scopeEnd], flid=${record.flid}", + "operation day not computable, flid=${record.flid}", setOf(SnapshotFlag.DAY_MISMATCH), ) } perRecordDay[record.flid] = day } - return SnapshotValidation.Ok(perRecordDay, flags) + return SnapshotValidation.Ok(perRecordDay, emptySet()) } /** - * 完整快照语义(§5.4 DNLD/RESP):标量出现 Set、缺失 Clear;集合出现 Replace、缺失 Replace 空集。 - * `keepDeleted = true` 时 STATE 保持 DELETED(§5.1 步骤 6:普通 SCHD 不恢复)。 + * 日计划合并(§3.1):标量出现 Set(空值 = 显式清空)、缺失保留本地值; + * 集合出现 Replace、缺失保留。新航班(current == null)无本地值可保留, + * 集合按全键输出空集(§2.2 完整合并结果写入口径)。 + * `keepDeleted = true` 时 STATE 保持 DELETED(§3.3:日计划不承担恢复, + * 恢复入口只有 ADFT);冲突告警由调用方按 SCHD_REVIVE_CONFLICT 记录。 */ fun snapshotState( current: FlightSnapshot?, @@ -75,28 +80,33 @@ object FlightStateEngine { keepDeleted -> FlightState.DELETED else -> current.state } - // §5.4 完整快照:标量出现 Set、缺失 Clear;集合出现 Replace、缺失 Replace 空集。 - // 完整替换 = 新状态只由记录决定,不从 current 继承任何字段。 - val scalars: Map = record.scalars + val scalars = buildMap { + current?.scalars?.let(::putAll) // 缺失 = 保留(§3.1) + putAll(record.scalars) // 出现 = Set;空串 = 显式清空 + } val collections = buildMap { - COLLECTION_KEYS.forEach { key -> put(key, emptyList()) } // 缺失 = Replace 空集(§5.4) + if (current == null) { + COLLECTION_KEYS.forEach { key -> put(key, emptyList()) } // 新航班全键输出 + } else { + putAll(current.collections) // 缺失 = 保留(§3.1) + } record.collections.forEach { (key, items) -> - if (key in COLLECTION_KEYS) put(key, items) + if (key in COLLECTION_KEYS) put(key, items) // 出现 = Replace(§3.1) } } return FlightSnapshot( flid = record.flid, operationDay = operationDay, state = state, - stateVersion = (current?.stateVersion ?: 0L) + 1, + stateVersion = (current?.stateVersion ?: 0L) + 1, // 每次成功写入 +1(§3.1/§4) scalars = scalars, collections = collections, ) } /** - * FLOP 增量合并(§6.1):标量出现覆盖、缺失保留;集合出现 Replace、缺失保留 - * ——FLOP 集合语义未定案,Replace 为当前偏差明示沿用(§10)。 + * FLOP/ADFT 增量合并(§3.2/§3.3):标量出现覆盖、缺失保留;集合出现 Replace、 + * 缺失保留——只修改报文表达的字段/集合,其余航班状态保持不变;不推导运营日。 */ fun mergedState(current: FlightSnapshot, change: MergeChange): FlightSnapshot { val scalars = buildMap { @@ -116,11 +126,11 @@ object FlightStateEngine { ) } - /** FLID:数字型,SIS §3.16.2 Number(1-12)。 */ + /** FLID:数字型(SIS Number(1-12))。 */ private val FLID_REGEX = Regex("\\d{1,12}") } -/** §5.2 校验结果:Ok 携带每条记录的归属运营日与观测 flags;Invalid 整包拒绝。 */ +/** 整包校验结果:Ok 携带每条记录的归属运营日与观测 flags;Invalid 整包拒绝。 */ sealed interface SnapshotValidation { data class Ok( val perRecordDay: Map, 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 14ba3e1..6243294 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 @@ -14,22 +14,22 @@ import java.time.LocalDate import java.time.ZoneId // ===================================================================== -// 仓储契约(docs/flight-state.md §3.2 表职责)。 +// 仓储契约(docs/flight-state.md §2 权威模型/§3 写入语义,design.md §2.1)。 // 决策层写路径约定:所有 FLIGHT_SCHD 及明细写操作必须发生在 -// 「持有 PIPELINE_LOCK 的同一事务」内(§1:锁只串行化 DB 事务)。 +// 「持有 PIPELINE_LOCK 的同一事务」内(§1/§4:锁只串行化 DB 事务)。 // ===================================================================== -/** 自有 PG 单事务原子保障(§1 目标 3:状态、事件、处理终态同事务提交)。 */ +/** 自有 PG 单事务原子保障(flight-state.md §1/§4:状态、事件、处理终态同事务提交)。 */ interface PipelineTransactionManager { fun inTransaction(block: () -> T): T } -/** 单行锁(§3.2 PIPELINE_LOCK):事务内第一步 SELECT ... FOR UPDATE,串行化状态写事务。 */ +/** 单行锁 PIPELINE_LOCK:事务内第一步 SELECT ... FOR UPDATE,串行化状态写事务(§1/§4)。 */ interface PipelineLockRepository { fun lock() } -/** PROC_STATE:每消息一行;兼作快照重放判定(§5.1 步骤 2:MSG_ID 已有成功终态 → 重放)。 */ +/** PROC_STATE:每消息一行(design.md §2.1);SUCCEEDED 终态兼作日计划重放判定。 */ interface ProcStateRepository { fun insert(msgId: Long, state: ProcStatus = ProcStatus.PENDING) @@ -61,7 +61,7 @@ interface ProcStateRepository { fun requeueByErrorClasses(errorClasses: List): Int } -/** MSG_EVENT outbox(§7.3)。KAFKA_SCHD 合并同 FLID 未发事件按最新 STATE_VERSION 输出。 */ +/** MSG_EVENT outbox(flight-state.md §5)。KAFKA_SCHD 合并同 FLID 未发事件按最新 STATE_VERSION 输出。 */ interface MsgEventRepository { fun insertAll(events: List): List @@ -69,7 +69,7 @@ interface MsgEventRepository { fun claimBatch(target: String, limit: Int): List - /** 同 FLID 未发 KAFKA_SCHD 事件合并:每 FLID 取最新 STATE_VERSION 一条(§7.3)。 */ + /** 同 FLID 未发 KAFKA_SCHD 事件合并:每 FLID 取最新 STATE_VERSION 一条(§5)。 */ fun mergePendingSchd(limit: Int): List fun markSent(eventId: Long) @@ -81,62 +81,63 @@ interface MsgEventRepository { fun markDead(eventId: Long, errorClass: ErrorClass, lastError: String, attempts: Int? = null) } -/** 完整态落库结果:DAY_GUARD_VIOLATION = OPERATION_DAY 不可变条件更新未命中(§7.4)。 */ +/** 完整态落库结果:DAY_GUARD_VIOLATION = OPERATION_DAY 不可变条件更新未命中(§2.1)。 */ enum class PersistOutcome { INSERTED, UPDATED, DAY_GUARD_VIOLATION } /** - * 航班当前态权威(§3.2 FLIGHT_SCHD + 9 张明细表)。 + * 航班当前态权威(§2 FLIGHT_SCHD + 9 张明细表)。 * 唯一写路径 = persistFullState:写入前在内存生成完整新状态(引擎产物)再落库, - * 明细集合按组先删后插(§3.3)。 + * 明细集合按组先删后插(§2.2 完整合并结果为准)。 */ interface FlightStateRepository { - /** 主行点查(身份/版本判定;批量用于快照归属校验 §5.3)。 */ + /** 主行点查(身份/版本判定;批量用于日计划归属校验 §2.1)。 */ fun findMainRow(flid: String): FlightMainRow? fun findMainRows(flids: Collection): Map - /** 完整当前态:主行 + 全部明细(一致性读边界由调用方事务保证,§9)。 */ + /** 完整当前态:主行 + 全部明细(一致性读边界由调用方事务保证,§5)。 */ fun loadFullSnapshot(flid: String): FlightSnapshot? /** - * 完整当前态落库(§3.3):主行 upsert + 全部明细按组先删后插; + * 完整当前态落库(§3):主行 upsert + 全部明细按组先删后插; * STATE_VERSION 以 snapshot.stateVersion 落库。 - * SCHD(整体替换/清除)、FLOP/ADFT(合并后全量写)共用此唯一写路径。 - * 条件更新带 `WHERE operation_day IS NULL OR operation_day = :day`(§7.4 不可变强化)。 + * SCHD(§3.1)、FLOP/ADFT(§3.2/§3.3)合并后的全量结果共用此唯一写路径。 + * 条件更新带 `WHERE operation_day IS NULL OR operation_day = :day`(§2.1 不可变强化)。 */ fun persistFullState(snapshot: FlightSnapshot, msgId: Long, now: Instant): PersistOutcome /** - * FDEL(§6.2):ACTIVE → 置 DELETED、推进版本、明细保留、返回 true(发布删除事件); + * FDEL(§3.3):ACTIVE → 置 DELETED、推进版本、明细保留、返回 true(发布删除事件); * 已 DELETED 或不存在 → 返回 false(幂等成功,不推进版本不重复发布)。 */ fun markDeleted(flid: String, msgId: Long, now: Instant): Boolean - /** ADFT 生命周期重激活(§6.3):DELETED → ACTIVE,推进版本;非 DELETED 返回 false。 */ + /** ADFT 生命周期重激活(§3.3):DELETED → ACTIVE,推进版本;非 DELETED 返回 false。 */ fun revive(flid: String, msgId: Long, now: Instant): Boolean - /** §8.1:按窗口规则选出历史候选(含 DELETED;候选时间按机场时区折算)。 */ + /** §6:按保留期与终态/静默判据选出历史候选(含 DELETED;窗口按机场时区折算)。 */ fun findHistoryCandidates(rules: HistoryRules, zone: ZoneId, now: Instant): List /** - * §8.2 步骤 3:物理删除主行与明细(仅历史存储确认成功后调用; + * §6:物理删除主行与明细(仅历史存储确认成功后调用; * 历史存储未接通时调用方必须传空集合——删 0 条)。 */ fun purgeArchived(flids: Collection): Int - /** §8.3 观测:OPERATION_DAY 仍为 NULL 的航班数(只增不删,终止规则未定 §10)。 */ + /** §2.1 观测:OPERATION_DAY 仍为 NULL 的航班数(只增不删,终止规则未定 §6)。 */ fun countOperationDayNull(): Int } -/** §5.5 SCHD_SNAP_LOG:事务外追加留痕,不参与决策;一行 = 一次尝试(重放也记)。 */ +/** SCHD_SNAP_LOG(design.md §6.2):事务外追加留痕,不参与决策;一行 = 一次尝试(重放也记)。 */ interface SnapshotLogRepository { fun append(entry: SnapshotLogEntry) } /** - * 请求状态机(§7.1):只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型) - * 匹配最新一条 PENDING;同类请求只留一条有效,新请求置旧为 EXPIRED。 + * 请求状态机 REQ_TRACK(design.md §4.2 目标机制,尚无运行时协调器): + * 只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型)匹配最新一条 PENDING; + * 同类请求只留一条有效,新请求置旧为 EXPIRED。 */ interface ReqTrackRepository { enum class ReqState { PENDING, SENT, DONE, EXPIRED } @@ -156,7 +157,7 @@ interface ReqTrackRepository { fun findLatest(reqType: String, operationDay: LocalDate, sender: String, states: List): Req? - /** RESP 完成请求(§7.1):匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */ + /** RESP 完成请求(design.md §4.2):匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */ fun completeLatest(reqType: String, operationDay: LocalDate, sender: String): Boolean fun linkCoutmsgs(reqId: Long, coutmsgsId: Long) @@ -167,7 +168,7 @@ interface ReqTrackRepository { } /** - * 共享信箱回填补偿待办(§7.2):业务事务内预登记(消除崩溃窗口 §10 偏差), + * 共享信箱回填补偿待办(design.md §3.3/§6.1):业务事务内预登记, * 提交后由 BackfillSweepJob 重试;回填失败不得把 SUCCEEDED 改回 FAILED。 */ interface BackfillTodoRepository { diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/dialect/SqlDialect.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/dialect/SqlDialect.kt deleted file mode 100644 index ad912ea..0000000 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/dialect/SqlDialect.kt +++ /dev/null @@ -1,72 +0,0 @@ -package com.gzzn.omms.msgexchange.infra.persistence.dialect - -/** - * v2 §6(ACM2-29 P3-B):数据库方言接缝。 - * - * 原则(v2 §1/§6):同一业务模型支持 PostgreSQL 与 Oracle 11g,方言只存在于基础设施适配层, - * 业务层与仓储语义不按数据库类型分支。 - * - * 接缝范围(随仓储 SQL 演进扩充): - * - 航班主行快照 upsert:PG `INSERT .. ON CONFLICT` vs 11g `MERGE INTO`(11g 无 ON CONFLICT) - * - 增量路径主行存在性保障:同上 - * - * 激活门控(本接缝为编译级交付,**未经目标库验证**): - * Oracle 11.2 补丁级别 × JDK 25 × ojdbc 驱动 × Flyway(db/migration/oracle11g location) - * × 连接池组合必须在现场 11g 实测通过后,才允许把生产方言切到 [Oracle11gDialect]; - * MERGE 绑定顺序适配、CLOB(SRVT/VIPF/MAFL_TEXT)与空串=NULL 语义回归一并纳入激活清单。 - * 验收证据要求见 docs/flight-state.md §8。 - */ -interface SqlDialect { - /** 航班主行快照 upsert:按 SCALAR_WRITE_COLUMNS 生成整体替换 SQL。 */ - fun flightSnapshotUpsertSql(columns: List): String - - /** 增量路径主行存在性保障(新插 FDAY=NULL,已有行仅推进 updated_at)。 */ - fun flightRowEnsureSql(): String -} - -/** PostgreSQL 方言:当前生产路径。 */ -object PostgreSqlDialect : SqlDialect { - - override fun flightSnapshotUpsertSql(columns: List): String = buildString { - append("INSERT INTO flight_schd (flid, fday, ") - append(columns.joinToString(", ")) - append(", last_message_id, state_version, created_at, updated_at) VALUES (?, ?, ") - append(columns.joinToString(", ") { "?" }) - append(", ?, ?, ?, ?) ON CONFLICT (flid) DO UPDATE SET fday = EXCLUDED.fday, ") - append(columns.joinToString(", ") { "$it = EXCLUDED.$it" }) - append(", last_message_id = EXCLUDED.last_message_id, state_version = EXCLUDED.state_version, updated_at = EXCLUDED.updated_at") - } - - override fun flightRowEnsureSql(): String = - "INSERT INTO flight_schd (flid, created_at, updated_at) VALUES (?, ?, ?) " + - "ON CONFLICT (flid) DO UPDATE SET updated_at = EXCLUDED.updated_at" -} - -/** - * Oracle 11g 方言:`MERGE INTO` 单语句 upsert(11g 无 `ON CONFLICT`)。 - * - * **编译级交付,未在目标 11g 上验收**(见 [SqlDialect] 激活门控)。 - * 绑定顺序与 PG 不同(USING 子句的 flid 绑定 + MATCHED/NOT MATCHED 两份列绑定), - * 激活时需同步实现 11g 绑定适配,禁止直接复用 PG 绑定序列。 - */ -object Oracle11gDialect : SqlDialect { - - override fun flightSnapshotUpsertSql(columns: List): String = buildString { - append("MERGE INTO flight_schd t ") - append("USING (SELECT ? AS FLID FROM dual) s ON (t.flid = s.FLID) ") - append("WHEN MATCHED THEN UPDATE SET t.fday = ?, ") - append(columns.joinToString(", ") { "t.$it = ?" }) - append(", t.last_message_id = ?, t.state_version = ?, t.updated_at = ? ") - append("WHEN NOT MATCHED THEN INSERT (flid, fday, ") - append(columns.joinToString(", ")) - append(", last_message_id, state_version, created_at, updated_at) VALUES (") - append((listOf("?", "?") + columns.map { "?" } + listOf("?", "?", "?")).joinToString(", ")) - append(")") - } - - override fun flightRowEnsureSql(): String = - "MERGE INTO flight_schd t " + - "USING (SELECT ? AS FLID FROM dual) s ON (t.flid = s.FLID) " + - "WHEN MATCHED THEN UPDATE SET t.updated_at = ? " + - "WHEN NOT MATCHED THEN INSERT (flid, created_at, updated_at) VALUES (?, ?, ?)" -} 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 611ca93..f1e39eb 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 @@ -36,9 +36,9 @@ import java.util.Locale import javax.sql.DataSource // ===================================================================== -// 自有 PostgreSQL 仓储实现(docs/flight-state.md §3.2 表职责;schema 见 +// 自有 PostgreSQL 仓储实现(docs/flight-state.md §2 权威模型/§3 写入语义;schema 见 // db/migration/V1__flight_state_baseline.sql)。全部写路径约定在 -// withTransaction + PIPELINE_LOCK 内调用(§1 目标 3/§3.2)。 +// withTransaction + PIPELINE_LOCK 内调用(§1/§4)。 // ===================================================================== @Singleton @@ -56,7 +56,7 @@ class JdbcPipelineTransactionManager( class JdbcPipelineLockRepository( private val ds: DataSource, ) : PipelineLockRepository { - /** 事务内第一步:单行 FOR UPDATE 串行化状态写事务(§3.2 PIPELINE_LOCK)。 */ + /** 事务内第一步:单行 FOR UPDATE 串行化状态写事务(§1/§4 PIPELINE_LOCK)。 */ override fun lock() { ds.queryOne("SELECT lock_id FROM pipeline_lock WHERE lock_id = 1 FOR UPDATE", {}) { 1 } ?: error("PIPELINE_LOCK row missing") } @@ -211,7 +211,7 @@ class JdbcMsgEventRepository( ::mapEvent, ) - /** §7.3:同 FLID 未发事件按最新 STATE_VERSION 合并;PG DISTINCT ON 方言(注释明示)。 */ + /** §5:同 FLID 未发事件按最新 STATE_VERSION 合并;PG DISTINCT ON 方言(注释明示)。 */ override fun mergePendingSchd(limit: Int): List = ds.query( """ @@ -348,7 +348,7 @@ class JdbcFlightStateRepository( }, ) if (updated == 0) { - // §7.4 不可变强化:条件更新未命中 = 归属日冲突 + // §2.1 不可变强化:条件更新未命中 = 归属日冲突 return if (existed != null) PersistOutcome.DAY_GUARD_VIOLATION else PersistOutcome.UPDATED } replaceDetails(snapshot, now) @@ -370,7 +370,7 @@ class JdbcFlightStateRepository( ) == 1 override fun findHistoryCandidates(rules: HistoryRules, zone: ZoneId, now: Instant): List { - // §8.1 四条判定在应用层执行(SIS 时间串解析无法下推 SQL); + // §6 四条判定在应用层执行(SIS 时间串解析无法下推 SQL); // 不按 updated_at 粗筛——取消/终态时间可能早于最近一次更新,粗筛会漏删。 val rows = ds.query( "SELECT flid, state, state_version, operation_day, cncl, naat, neat, updated_at FROM flight_schd", @@ -399,7 +399,7 @@ class JdbcFlightStateRepository( r.updatedAt < now.minus(Duration.ofHours(rules.idleHours)) -> true else -> false } - // NAAT/NEAT 业务含义待术语表确认(§10):解析失败视为不命中,不误删 + // NAAT/NEAT 业务含义待术语表确认(§6 开放项):解析失败视为不命中,不误删 if (hit) HistoryCandidate(r.flid, r.state, r.version, wasNeverFdel = r.state == FlightState.ACTIVE) else null } } @@ -456,7 +456,7 @@ class JdbcFlightStateRepository( items.forEachIndexed { ordinal, item -> val cols = mutableListOf("flid", "ordinal", "source_seq") val vals = mutableListOf(flid) - vals.add(ordinal + 1) // ORDINAL 保留输入顺序(§3.2) + vals.add(ordinal + 1) // ORDINAL 保留输入顺序(§2.2) vals.add(item[spec.seqAttr]) spec.columns.forEach { col -> cols.add(col) @@ -488,7 +488,7 @@ class JdbcFlightStateRepository( private fun loadDetails(flid: String, key: String, spec: DetailSpec): List> { val sql = if (spec.routeKind != null) { - // ROUT/ERUT 共表:按 ROUTE_KIND 过滤(§3.2) + // ROUT/ERUT 共表:按 ROUTE_KIND 过滤(§2.2) "SELECT * FROM ${spec.table} WHERE flid = ? AND route_kind = ? ORDER BY ordinal ASC" } else { "SELECT * FROM ${spec.table} WHERE flid = ? ORDER BY ordinal ASC" @@ -544,7 +544,7 @@ class JdbcFlightStateRepository( "flight_delay", "flight_bridge_op", "flight_chock_op", "flight_route_point", ) - /** 10 类集合 ↔ 明细表/列映射(§3.2;列名与基线一致)。 */ + /** 10 类集合 ↔ 明细表/列映射(§2.2;列名与基线一致)。 */ internal val COLLECTIONS: Map = mapOf( "GTDT" to DetailSpec("flight_gate", listOf("gate", "pgot", "pgct", "gotm", "gctm", "gtyp"), "GTNO"), "CKDT" to DetailSpec("flight_checkin", listOf("chkc", "ccls", "pcot", "pcct", "cotm", "cctm", "ctyp"), "CKNO"), @@ -566,7 +566,7 @@ class JdbcFlightStateRepository( class JdbcSnapshotLogRepository( private val ds: DataSource, ) : SnapshotLogRepository { - /** §5.5 只追加;写失败由调用方捕获记指标(append 自身不抛出)。 */ + /** design.md §6.2:只追加;写失败由调用方捕获记指标(append 自身不抛出)。 */ override fun append(entry: SnapshotLogEntry) { ds.update( """ @@ -628,7 +628,7 @@ class JdbcReqTrackRepository( ) } - /** §7.1 RESP 完成请求:匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */ + /** design.md §4.2 RESP 完成请求:匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */ override fun completeLatest(reqType: String, operationDay: LocalDate, sender: String): Boolean = ds.update( """ diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt index 89aa0f5..ea3ccc2 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubAdapters.kt @@ -6,7 +6,7 @@ import jakarta.inject.Singleton /** * stub 适配层——DeliveryPort 内存实现,仅在 msgx.stubs=true 时生效。 - * 记录 (topic, key, payload):payload = null 表示 TOMBSTONE(§7.3 键缺失=删除旧值)。 + * 记录 (topic, key, payload):payload = null 表示 TOMBSTONE(flight-state.md §5 键缺失=删除旧值)。 */ @Requires(property = "msgx.stubs", value = "true") @Singleton 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 f73ab68..eed10d5 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 @@ -132,7 +132,7 @@ class StubMsgEvents : MsgEventRepository { override fun claimBatch(target: String, limit: Int): List = rows.values.filter { it.target == target && it.state == EventStatus.PENDING }.sortedBy { it.eventId!! }.take(limit) - /** §7.3:同 FLID 取最新 STATE_VERSION,按事件序输出。 */ + /** §5:同 FLID 取最新 STATE_VERSION,按事件序输出。 */ override fun mergePendingSchd(limit: Int): List = rows.values .filter { it.target == "KAFKA:schd" && it.state == EventStatus.PENDING } @@ -178,7 +178,7 @@ class StubFlightState : FlightStateRepository { override fun loadFullSnapshot(flid: String): FlightSnapshot? = snapshots[flid] - /** §7.4 不可变条件:已有非空 OPERATION_DAY 且与新值不同 → DAY_GUARD_VIOLATION。 */ + /** §2.1 不可变条件:已有非空 OPERATION_DAY 且与新值不同 → DAY_GUARD_VIOLATION。 */ override fun persistFullState(snapshot: FlightSnapshot, msgId: Long, now: Instant): PersistOutcome { val existing = mains[snapshot.flid] if (existing?.operationDay != null && existing.operationDay != snapshot.operationDay) { @@ -260,7 +260,7 @@ class StubReqTrack : ReqTrackRepository { fun clear() = rows.clear() override fun insert(reqType: String, operationDay: LocalDate, sender: String): Long { - // §7.1:同类(类型+运营日+发送方)旧有效请求先置 EXPIRED + // design.md §4.2:同类(类型+运营日+发送方)旧有效请求先置 EXPIRED rows.values.filter { it.reqType == reqType && it.operationDay == operationDay && it.sender == sender && (it.state == ReqTrackRepository.ReqState.PENDING || it.state == ReqTrackRepository.ReqState.SENT) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt index e2eefe7..0dd81e3 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt @@ -14,13 +14,13 @@ import java.time.Instant import java.time.ZoneId /** - * 历史归档与物理清除(docs/flight-state.md §8.2,顺序不可颠倒): - * 1. HISTORY_SWEEP 选出满足 §8.1 的航班(含 DELETED); + * 历史归档与物理清除(docs/flight-state.md §6 生命周期,顺序不可颠倒;design.md §6.2): + * 1. 选出满足保留期 + 终态/静默判据的航班(含 DELETED); * 2. 写入历史存储; * 3. 历史存储返回成功的 FLID 集合 → 物理删除主行与明细(归档结果只记录在历史存储,当前态无 ARCHIVED 态); - * 4. 未经 FDEL 的航班在清除前补发一次删除事件(§7.3);其余不发; + * 4. 未经 FDEL 的航班在清除前补发一次删除事件(§3.3/§5);其余不发; * 5. 失败或不明确的保留重试。 - * 前提红线:历史存储未接通时必须删 0 条(§8.2)。 + * 前提红线:历史存储未接通时必须删 0 条(§6)。 */ @Singleton class HistorySweepJob( @@ -29,7 +29,7 @@ class HistorySweepJob( private val props: HistoryProps, /** 历史存储端口:返回归档成功的 FLID 集合;未接通时不注入(null)。 */ private val historyStore: HistoryStore? = null, - /** §8.3 留痕清理端口(90 天,按 (SCOPE_END, RECV_AT));未接通时不注入。 */ + /** design.md §6.2 留痕清理端口(90 天,按 (SCOPE_END, RECV_AT));未接通时不注入。 */ private val snapLogPurge: SnapshotLogPurge? = null, ) { /** 历史存储端口(由部署侧适配实现;脚手架默认未接通)。 */ @@ -46,12 +46,12 @@ class HistorySweepJob( fun run(now: Instant = Instant.now()): SweepOutcome { if (!props.historyStoreEnabled || historyStore == null) { - // 红线:历史存储未接通必须删 0 条;绝不允许先删当前态再补历史(§8.2) + // 红线:历史存储未接通必须删 0 条;绝不允许先删当前态再补历史(§6) return SweepOutcome(selected = 0, archived = 0, purged = 0) } val rules = HistoryRules(props.cancelledHours, props.terminalHours, props.deletedHours, props.idleHours) - val zone = ZoneId.of("Asia/Shanghai") // §8.1:窗口按机场时区计算 + val zone = ZoneId.of("Asia/Shanghai") // §6:窗口按机场时区计算 val candidates = flightState.findHistoryCandidates(rules, zone, now) if (candidates.isEmpty()) return SweepOutcome(0, 0, 0) @@ -59,7 +59,7 @@ class HistorySweepJob( if (archivedFlids.isEmpty()) return SweepOutcome(candidates.size, archived = 0, purged = 0) val toPurge = candidates.filter { it.flid in archivedFlids } - // §8.2 步骤 4:未经 FDEL、由生命周期直接清除的航班,清除前补发一次删除事件 + // §6/§3.3:未经 FDEL、由生命周期直接清除的航班,清除前补发一次删除事件 val preDelete = toPurge.filter { it.wasNeverFdel } if (preDelete.isNotEmpty()) { msgEvents.insertAll(preDelete.map { tombstone(it) }) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt index 695fac8..34fdc8c 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt @@ -9,10 +9,9 @@ import java.time.LocalDate import java.time.ZoneId /** - * 维护作业调度(docs/flight-state.md §8):单 daemon 线程,独立于主泵—— - * 旧「作业不插队 PUMP_JOB 队列」机制随审计口径移除(§3.2 表清单无 PUMP_JOB); - * 历史归档/留痕清理均为内部清理路径,不参与 FIFO 消息序。 - * 触发:回填补偿 30s 固定间隔;历史归档/留痕清理每日 03:30(机场时区)后首个 tick。 + * 维护作业调度(docs/design.md §6.1/§6.2):单 daemon 线程,独立于主泵—— + * 作业不参与消息 FIFO,也不使到期消息饥饿(PUMP_JOB 队列机制已随审计口径移除)。 + * 触发:回填补偿扫描 30s 固定间隔;历史归档/留痕清理每日机场时区 03:30 后首个 tick。 */ @Singleton class JobRunner( 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 2eb3faf..dc79a22 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt @@ -10,8 +10,10 @@ import com.gzzn.omms.msgexchange.domain.OperationDayCalculator import com.gzzn.omms.msgexchange.domain.ProcState import com.gzzn.omms.msgexchange.domain.Targets import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot +import com.gzzn.omms.msgexchange.domain.flight.FlightState import com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine import com.gzzn.omms.msgexchange.domain.flight.MergeChange +import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord import com.gzzn.omms.msgexchange.infra.persistence.BackfillTodoRepository import com.gzzn.omms.msgexchange.infra.persistence.FlightStateRepository import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository @@ -22,8 +24,9 @@ import java.time.Instant import java.time.ZoneId /** - * FLOP(§6.1):读取完整当前态 → 合并变化 → 保留运营日 → STATE_VERSION+1 → - * 同事务登记 KAFKA_MSG / KAFKA_SCHD 与处理终态。 + * FLOP(docs/flight-state.md §3.2 动态运行事件):读取完整当前态 → 合并变化 → + * 保留运营日 → STATE_VERSION+1 → 同事务登记 KAFKA_MSG / KAFKA_SCHD 与处理终态。 + * 未知/迟到航班按幂等成功处理,不创建实例(创建入口只有 SCHD/ADFT)。 */ @Singleton class FlopProcessor( @@ -49,8 +52,8 @@ class FlopProcessor( } /** - * FDEL(§6.2):ACTIVE → 置 DELETED、推进版本、明细保留、发布 tombstone; - * 已 DELETED / 不存在 → 幂等成功,不推进版本、不重复发布。 + * FDEL(docs/flight-state.md §3.3 删除):ACTIVE → 置 DELETED、推进版本、明细保留、 + * 与删除同事务登记 tombstone(§5);已 DELETED / 不存在 → 幂等成功,不推进版本、不重复发布。 */ @Singleton class FdelProcessor( @@ -66,7 +69,7 @@ class FdelProcessor( val deleted = flightState.markDeleted(payload.flid, msgId = head.msgId, now = Instant.now()) if (deleted) { val current = flightState.loadFullSnapshot(payload.flid) - // tombstone 仅在 ACTIVE→DELETED 时登记(§7.3),与删除同事务 + // tombstone 仅在 ACTIVE→DELETED 时登记(§3.3/§5),与删除同事务 msgEvents.insertAll( listOf( MsgEvent( @@ -94,14 +97,14 @@ class FdelProcessor( ) preRegisterBackfill(head, msg, backfillTodo) } - ApplyResult.Succeeded // 未命中 = 迟到/重复,幂等成功(§6.2 步骤 3/4) + ApplyResult.Succeeded // 未命中 = 迟到/重复,幂等成功(§3.3) } } /** - * ADFT(§6.3 + §2.1):字段缺失语义待确认——确认前按保守 Set-only 处理 - * (出现字段覆盖,缺失不 Clear,不沿用 FLOP 全量合并规则)。 - * FLID 已存在且 DELETED → 生命周期重激活;不存在 → 新实例建立(含运营日计算 §3.5)。 + * ADFT(docs/flight-state.md §3.3):字段缺失语义待上游确认——确认前按保守 Set-only + * 处理(出现字段覆盖、缺失不清空)。FLID 已存在且 DELETED → 生命周期重激活;不存在 → + * 新实例建立(含 SODT 时直接计算 OPERATION_DAY,§2.1;否则保留 NULL 待日计划收录)。 */ @Singleton class AdftProcessor( @@ -118,12 +121,12 @@ class AdftProcessor( cutoffHour = operationDayProps.cutoffHour, ) - fun apply(head: ProcState, msg: DecodedMessage, record: com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord): ApplyResult = + fun apply(head: ProcState, msg: DecodedMessage, record: ScheduleRecord): ApplyResult = txManager.inTransaction { lock.lock() val main = flightState.findMainRow(record.flid) - if (main != null && main.state == com.gzzn.omms.msgexchange.domain.flight.FlightState.DELETED) { - // §6.3 重激活:DELETED → ACTIVE,推进版本,登记状态事件 + if (main != null && main.state == FlightState.DELETED) { + // §3.3 重激活:DELETED → ACTIVE,推进版本,登记状态事件 if (flightState.revive(record.flid, msgId = head.msgId, now = Instant.now())) { val current = flightState.loadFullSnapshot(record.flid) if (current != null) { @@ -138,15 +141,15 @@ class AdftProcessor( val current = flightState.loadFullSnapshot(record.flid) val next: FlightSnapshot = if (current == null) { - // 新实例建立:ADFT 含 SODT 时直接计算运营日(§2.1),不可算则置 null 待快照收录 + // 新实例建立:ADFT 含 SODT 时直接计算运营日(§2.1),不可算则置 null 待日计划收录 val day = opDay.compute(record.scalars["SODT"]) FlightSnapshot( flid = record.flid, - operationDay = day, // 待确认项 §2.1:不可算时不得默认写接收日 - state = com.gzzn.omms.msgexchange.domain.flight.FlightState.ACTIVE, + operationDay = day, // §2.1:不可算时不得默认写接收日 + state = FlightState.ACTIVE, stateVersion = 1L, scalars = record.scalars, - collections = com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.COLLECTION_KEYS.associateWith { key -> + collections = FlightStateEngine.COLLECTION_KEYS.associateWith { key -> record.collections[key] ?: emptyList() }, ) @@ -162,8 +165,8 @@ class AdftProcessor( ApplyResult.Succeeded } - /** §2.1 保守语义:仅出现字段 Set;集合出现 Replace、缺失保留。 */ - private fun setOnly(record: com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord) = MergeChange( + /** §3.3 保守语义:仅出现字段 Set;集合出现 Replace、缺失保留。 */ + private fun setOnly(record: ScheduleRecord) = MergeChange( flid = record.flid, scalars = record.scalars, collections = record.collections, @@ -174,10 +177,10 @@ class AdftProcessor( // 共享小工具(处理器层私有约定) // ===================================================================== -/** 航班不存在/迟到:幂等成功(§9 队头不阻塞;不创建实例——创建入口只有 SCHD/ADFT)。 */ +/** 航班不存在/迟到:幂等成功(§3.2/§3.3;不阻塞队头,不创建实例——创建入口只有 SCHD/ADFT)。 */ private fun idempotentAbsent(head: ProcState, msg: DecodedMessage): ApplyResult = ApplyResult.Succeeded -/** KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(§7.3)。 */ +/** KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1/§5)。 */ internal fun eventsFor(next: FlightSnapshot, mapper: ObjectMapper): List { val payload = linkedMapOf( "flid" to next.flid, @@ -201,7 +204,7 @@ internal fun eventsFor(next: FlightSnapshot, mapper: ObjectMapper): List { val body = decoded.body as? ScheduleBody @@ -194,7 +195,7 @@ class MessageProcessor( flopProcessor.apply(head, decoded, payload) } is MsgKind.Unsupported -> { - // §9:未支持类型 → FAILED(UNSUPPORTED) 退避重试,达阈值转 DEAD;绝不写终态 + // design.md §2.3:未支持类型 → FAILED(UNSUPPORTED) 退避重试,达阈值转 DEAD;绝不写终态 log.warn("unsupported type -> FAILED(UNSUPPORTED) msgId={} tag={}", head.msgId, kind.tag) procFailure.fail(head, ErrorClass.UNSUPPORTED, "no-handler:${kind.tag}") return @@ -205,7 +206,7 @@ class MessageProcessor( is ApplyResult.Succeeded, ApplyResult.ReplaySkipped -> procState.update(head.msgId, ProcStatus.SUCCEEDED) is ApplyResult.DeadProtocol -> { - // §9:整包拒绝 DEAD(PROTOCOL),立即释放队头,交人工确认 + // design.md §2.3:整包拒绝 DEAD(PROTOCOL),立即释放队头,交人工确认 log.error("DEAD(PROTOCOL) msgId={} reason={} flags={}", head.msgId, result.reason, result.flags) procState.update(head.msgId, ProcStatus.DEAD, errorClass = ErrorClass.PROTOCOL, lastError = result.reason.take(1000)) compensateBackfill(head, decoded) // 拒绝包同样要回填信箱,防止反复轮询 @@ -217,7 +218,7 @@ class MessageProcessor( log.info("SUCCEEDED msgId={} kind={}", head.msgId, decoded.typeTag) } - /** §7.2:提交后回填共享信箱;失败不得把 SUCCEEDED 改回 FAILED,待办已事务内预登记。 */ + /** design.md §3.3/§6.1:提交后回填共享信箱;失败不得把 SUCCEEDED 改回 FAILED,待办已事务内预登记。 */ private fun backfill(head: ProcState, decoded: DecodedMessage) { try { inbox.backfillOnSuccess(head.msgId, decoded.meta.sndr, decoded.meta.type, decoded.meta.styp, decoded.meta.seqn) 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 3c1b911..d4e6ac6 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt @@ -6,7 +6,6 @@ import com.gzzn.omms.msgexchange.config.OperationDayProps import com.gzzn.omms.msgexchange.domain.DecodedMessage 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.SnapshotFlag import com.gzzn.omms.msgexchange.domain.SnapshotLogEntry import com.gzzn.omms.msgexchange.domain.SnapshotResult @@ -30,24 +29,25 @@ import java.time.Instant import java.time.LocalDate import java.time.ZoneId -/** 处理器执行结果——终态迁移由 MessageProcessor 统一落库。 */ +/** 处理器执行结果——终态迁移由 MessageProcessor 统一落库(docs/design.md §2.3)。 */ sealed interface ApplyResult { /** 业务成功(含幂等成功)。 */ data object Succeeded : ApplyResult - /** §5.1 步骤 2:MSG_ID 已有成功终态 → 重放,直接记幂等成功。 */ + /** MSG_ID 已有成功终态 → 重放,直接记幂等成功(docs/flight-state.md §4 重放判定)。 */ data object ReplaySkipped : ApplyResult - /** §9:整包拒绝 DEAD(PROTOCOL),不重试,交人工确认。 */ + /** 整包拒绝 DEAD(PROTOCOL)(flight-state.md §3.1/§4),不重试,交人工确认。 */ data class DeadProtocol(val reason: String, val flags: Set = emptySet()) : ApplyResult } -/** §5.3 第四行:归属日不符 = 串日/错发/污染,整包拒绝。 */ +/** 归属日冲突(flight-state.md §2.1:OPERATION_DAY 一经确定不可变)= 串日/错发/污染,整包拒绝。 */ class ProtocolViolation(message: String) : RuntimeException(message) /** - * SCHD 快照主链路(docs/flight-state.md §5.1 applyScheduleRecords,同一事务): - * 对单日快照与滚动窗口统一适用,不做名单层面的处理。 + * SCHD 日计划主链路(docs/flight-state.md §3.1/§4):对 DNLD 与 RESP 统一适用。 + * 同一事务内:锁 → 归属日校验 → 逐条合并写完整当前态 → 登记事件与回填待办; + * 整包校验失败或运营日冲突整包不落地。 */ @Singleton class ScheduleProcessor( @@ -70,18 +70,16 @@ class ScheduleProcessor( val body = msg.body as? ScheduleBody ?: return ApplyResult.DeadProtocol("missing-schd-body") val started = System.nanoTime() - // ② 重放判定:MSG_ID 已有成功终态 → 幂等成功(§5.1 步骤 2,重放也记留痕) + // 重放判定:MSG_ID 已有成功终态 → 幂等成功(flight-state.md §4),重放也记留痕 if (procState.findSuccessTerminal(head.msgId)) { logSnapshot(head, body, SnapshotResult.REPLAY_SKIPPED, upserted = 0, flags = emptySet(), started) return ApplyResult.ReplaySkipped } - // ③ 报文完整性(§5.2 五项):任一失败整包不落地 → DEAD(PROTOCOL) + // 整包校验(§4 步骤 2 / design.md §4.1):任一失败整包不落地 → DEAD(PROTOCOL) val validation = FlightStateEngine.validateMessage( recsDeclared = body.recsDeclared, records = body.records, - scopeStart = body.scopeStart, - scopeEnd = body.scopeEnd, opDay = opDay, ) if (validation is SnapshotValidation.Invalid) { @@ -91,18 +89,18 @@ class ScheduleProcessor( val ok = validation as SnapshotValidation.Ok if (ok.perRecordDay.isEmpty()) { - // 空快照:合法但无写入(§5.5 EMPTY),仍算成功终态 + // 空快照:合法但无写入,仍算成功终态 logSnapshot(head, body, SnapshotResult.COMMITTED, upserted = 0, setOf(SnapshotFlag.EMPTY), started) return ApplyResult.Succeeded } val flags = linkedSetOf() return try { - // ①④⑤⑥⑦ 同一事务:锁 → 归属校验 → upsert → 版本/事件/待办预登记/终态 + // 同一事务(flight-state.md §4):锁 → 归属校验 → 逐条合并写 → 事件/待办预登记 val upserted = txManager.inTransaction { lock.lock() - // ⑤ 归属校验(§5.3):OPERATION_DAY 不可变,批量点查避免逐航班往返 + // 归属校验(§2.1):OPERATION_DAY 不可变,批量点查避免逐航班往返 val mains = flightState.findMainRows(ok.perRecordDay.keys) ok.perRecordDay.forEach { (flid, day) -> val existing = mains[flid] ?: return@forEach @@ -119,7 +117,7 @@ class ScheduleProcessor( val record = body.records.first { it.flid == flid } val existingMain = mains[flid] val keepDeleted = existingMain?.state == FlightState.DELETED - if (keepDeleted) flags.add(SnapshotFlag.SCHD_REVIVE_CONFLICT) // §5.1 步骤 6:不恢复 + if (keepDeleted) flags.add(SnapshotFlag.SCHD_REVIVE_CONFLICT) // §3.3:日计划不复活 DELETED val current = if (existingMain != null) flightState.loadFullSnapshot(flid) else null val next = FlightStateEngine.snapshotState( current = current, @@ -135,7 +133,7 @@ class ScheduleProcessor( events += snapshotEvents(next) } if (events.isNotEmpty()) msgEvents.insertAll(events) - // §7.2 目标形态:业务事务内预登记回填待办(提交后由 MessageProcessor 回填并删待办) + // 回填待办与业务终态同事务预登记(design.md §3.3/§6.1;提交后由 MessageProcessor 回填并删待办) backfillTodo?.record( BackfillTodoRepository.BackfillTask( msgId = head.msgId, sndr = msg.meta.sndr, type = msg.meta.type, @@ -153,9 +151,8 @@ class ScheduleProcessor( } } - /** 状态事件:整态出站(§7.3 KAFKA_SCHD)+ 变化通知(KAFKA_MSG)。 */ + /** 状态事件:KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1 登记、§5 语义)。 */ private fun snapshotEvents(next: FlightSnapshot): List { - // KAFKA_SCHD 整态:缺失集合输出空集 → 消费者删除旧值(§5.4/§7.3) val payload = linkedMapOf( "flid" to next.flid, "stateVersion" to next.stateVersion, @@ -174,7 +171,8 @@ class ScheduleProcessor( ) } - /** §5.5 留痕:事务外追加,失败只记 error 不阻塞;scope 按记录归属运营日推导。 */ + /** 留痕 SCHD_SNAP_LOG(design.md §6.2):事务外追加,失败只记 error 不阻塞; + * scope 为报文各记录归属运营日的最小/最大(单日快照两者相等)。 */ private fun logSnapshot( head: ProcState, body: ScheduleBody, diff --git a/src/main/resources/db/migration/oracle11g/README.md b/src/main/resources/db/migration/oracle11g/README.md index 8d8b7eb..a9b386c 100644 --- a/src/main/resources/db/migration/oracle11g/README.md +++ b/src/main/resources/db/migration/oracle11g/README.md @@ -7,17 +7,19 @@ Flyway 配置**;PG 路径使用 `classpath:db/migration`,两者互不混用 1. 现场 11.2 补丁级别、数据库字符集、DBA 权限清单拿到,且可提供可测试的目标库。 2. JDK 25 × ojdbc 驱动(具体版本)× Flyway Oracle 支持 × 连接池组合在目标库实测通过 - ——不能以"PG 通过"代替 Oracle 验收(docs/flight-state.md)。 -3. `SqlDialect` 切到 `Oracle11gDialect` 前,MERGE 绑定顺序适配完成并通过 - `FlightSchdJdbcPgTest` 同等粒度的 11g 集成测试。 + ——不能以"PG 通过"代替 Oracle 验收(Oracle 11g 适配为开放项,docs/flight-state.md §6)。 +3. 11g 的 upsert(MERGE INTO)与绑定顺序适配完成并通过与 PG 同粒度的集成测试后, + 才允许把仓储 SQL 切到 11g 方言。先前编译级交付的 `SqlDialect` 方言接缝已随 + V1 基线列名更替(fday/last_message_id → operation_day/last_msg_id)退役删除; + 接入时必须按 V1 基线现有列名重建,禁止直接复用旧接缝。 ## 计划内容 - V1 起:`V1__flight_state_baseline.sql` 的 11g 等价 DDL—— PIPELINE_LOCK/PROC_STATE/MSG_EVENT/REQ_TRACK/BACKFILL_TODO/FLIGHT_SCHD + 9 张明细表/SCHD_SNAP_LOG; `NUMBER`/序列替代 `BIGSERIAL`、`TIMESTAMP WITH TIME ZONE`、`VARCHAR2` BYTE/CHAR 语义钉死。 -- JdbcFlightStateRepository 的 `INSERT ... ON CONFLICT` 改 MERGE(OPERATION_DAY 不可变条件 §7.4)。 -- 空串按 NULL 的语义回归:命令层 presence 信息不得被 11g 空串语义吞掉 - (docs/flight-state.md §5/§10)。 +- `INSERT ... ON CONFLICT` 改 11g MERGE(OPERATION_DAY 不可变条件,flight-state.md §2.1)。 +- 空串按 NULL 的语义回归:显式清空的 presence 信息不得被 11g 空串语义吞掉 + (flight-state.md §3.1 显式清空口径)。 版本号与 PG location 各自独立推进,禁止复用版本号语义。 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt index b46e3d6..7324eff 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt @@ -58,7 +58,7 @@ class PipelineSmokeTest { } companion object { - /** 合法 META,但 TYPE 未在支持范围 → FAILED(UNSUPPORTED)(§9) */ + /** 合法 META,但 TYPE 未在支持范围 → FAILED(UNSUPPORTED)(design.md §2.3) */ val UNSUPPORTED_XML = """ AODB120260908120000XYZQFOO diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt index eeb4099..386a61e 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/codec/JacksonXmlCodecTest.kt @@ -3,6 +3,7 @@ package com.gzzn.omms.msgexchange.codec import com.gzzn.omms.msgexchange.domain.MsgKind import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Assertions.assertInstanceOf +import org.junit.jupiter.api.Assertions.assertNull import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test @@ -57,7 +58,8 @@ class JacksonXmlCodecTest { val body = msg.body as FlopPayload assertEquals("121112312", body.flid) - assertEquals("CA-CA101-A-12DEC261345-D", body.scalars["FFID"]) + // FFID(友好航班号)是 wire 必填字段,但当前状态模型不消费它,映射层应忽略。 + assertNull(body.scalars["FFID"]) assertEquals(3, body.collections["GTDT"]!!.size) assertEquals("G28", body.collections["GTDT"]!![0]["GATE"]) assertEquals("1", body.collections["GTDT"]!![0]["GTNO"]) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt index ebcd6e7..f556abb 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/delivery/DispatcherTickTest.kt @@ -16,7 +16,7 @@ import org.junit.jupiter.api.Test import java.time.Instant /** - * 投递调度(docs/flight-state.md §7.3): + * 投递调度(docs/flight-state.md §5 + design.md §5.1/§5.2): * KAFKA_MSG 逐条 FIFO;KAFKA_SCHD 唯一出口 flushSchd——同 FLID 未发事件按最新 * STATE_VERSION 合并;TOMBSTONE 发 null 值消息;失败退避重试、达上限 DEAD。 */ @@ -84,7 +84,7 @@ class DispatcherTickTest { d.flushSchd() assertEquals(1, port.tombstones.size) - assertEquals("F1", port.tombstones.single().key) // 整态键缺失表示删除旧值(§7.3) + assertEquals("F1", port.tombstones.single().key) // §5:整态键缺失表示删除旧值 assertNull(port.tombstones.single().payload) } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/domain/OperationDayTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/domain/OperationDayTest.kt index 729c264..7d135df 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/domain/OperationDayTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/domain/OperationDayTest.kt @@ -8,7 +8,7 @@ import java.time.LocalDate import java.time.ZoneId /** - * 运营日计算(docs/flight-state.md §3.5):SODT(ddMMMyyHHmm)+ 机场时区 + 业务切日边界。 + * 运营日计算(docs/flight-state.md §2.1):SODT(ddMMMyyHHmm)+ 机场时区 + 业务切日边界。 */ class OperationDayTest { 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 480a27f..41b1da9 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 @@ -10,9 +10,10 @@ import java.time.LocalDate import java.time.ZoneId /** - * 航班状态引擎(docs/flight-state.md §3.3/§5.2/§5.4/§6.1)不变量: - * 快照整体替换(缺失=清除)、DELETED 不被 SCHD 恢复、FLOP 增量保留缺失字段、 - * §5.2 五项校验整包拒绝语义。 + * 航班状态引擎(docs/flight-state.md §2–§4)不变量: + * 日计划合并(§3.1)——出现 Set/Replace、缺失保留、标量空值显式清空; + * DELETED 不被日计划恢复(§3.3);FLOP/ADFT 增量合并(§3.2/§3.3); + * 整包校验(§4 步骤 2:声明数量、航班标识、运营日推导)失败整包拒绝。 */ class FlightStateEngineTest { @@ -20,7 +21,7 @@ class FlightStateEngineTest { private val day = LocalDate.of(2026, 12, 15) @Test - fun `snapshot replaces scalars and clears absent ones`() { + fun `day plan merges scalars overwriting overlaps retaining absent and clearing on explicit empty`() { val current = FlightSnapshot( "121", day, FlightState.ACTIVE, 5, scalars = mapOf("FLNO" to "CA001", "REMC" to "old-note"), @@ -28,25 +29,66 @@ class FlightStateEngineTest { ) val next = FlightStateEngine.snapshotState( current, - ScheduleRecord("121", scalars = mapOf("FLNO" to "CA002")), + ScheduleRecord("121", scalars = mapOf("FLNO" to "CA002", "CNCL" to "")), operationDay = day, keepDeleted = false, ) - assertEquals("CA002", next.scalars["FLNO"]) - assertFalse(next.scalars.containsKey("REMC")) // 缺失 = Clear(§5.4) - assertEquals(emptyList>(), next.collections["GTDT"]) // 缺失 = Replace 空集 - assertEquals(6, next.stateVersion) // 每次成功写入 +1(§4) + assertEquals("CA002", next.scalars["FLNO"]) // 重叠字段被覆盖(§3.1) + assertEquals("old-note", next.scalars["REMC"]) // 缺失标量保留(§3.1/§7) + assertEquals("", next.scalars["CNCL"]) // 空串 = 显式清空(落库置 NULL) + assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"]) // 缺失集合保留 + assertEquals(6, next.stateVersion) // 每次成功写入 +1 assertEquals(FlightState.ACTIVE, next.state) } @Test - fun `snapshot keeps DELETED state and flags revive conflict by caller`() { + fun `day plan replaces collections that are present in the record`() { + val current = FlightSnapshot( + "121", day, FlightState.ACTIVE, 3, + scalars = mapOf("FLNO" to "CA001"), + collections = mapOf("GTDT" to listOf(mapOf("GATE" to "G1"))), + ) + val next = FlightStateEngine.snapshotState( + current, + ScheduleRecord( + "121", + scalars = mapOf("FLNO" to "CA002"), + collections = mapOf("GTDT" to listOf(mapOf("GATE" to "G2"), mapOf("GATE" to "G3"))), + ), + operationDay = day, + keepDeleted = false, + ) + + assertEquals(listOf(mapOf("GATE" to "G2"), mapOf("GATE" to "G3")), next.collections["GTDT"]) + assertEquals(4, next.stateVersion) + } + + @Test + fun `new flight snapshot carries only its record scalars and full collection key set`() { + val next = FlightStateEngine.snapshotState( + null, + ScheduleRecord("121", mapOf("FLNO" to "CA001"), mapOf("GTDT" to listOf(mapOf("GATE" to "G1")))), + operationDay = day, + keepDeleted = false, + ) + + assertEquals(FlightState.ACTIVE, next.state) + assertEquals(1L, next.stateVersion) + assertEquals("CA001", next.scalars["FLNO"]) + assertFalse(next.scalars.containsKey("REMC")) // 无历史可保留 + assertEquals(10, next.collections.size) // 全键输出(§2.2) + assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"]) + assertEquals(emptyList>(), next.collections["DELY"]) + } + + @Test + fun `day plan keeps DELETED state for revive conflict handling`() { val current = FlightSnapshot("121", day, FlightState.DELETED, 3, mapOf(), emptyMap()) val next = FlightStateEngine.snapshotState( current, ScheduleRecord("121", mapOf("FLNO" to "CA001")), day, keepDeleted = true, ) - assertEquals(FlightState.DELETED, next.state) // §5.1 步骤 6:普通 SCHD 不恢复 + assertEquals(FlightState.DELETED, next.state) // §3.3:日计划不恢复 DELETED assertEquals(4, next.stateVersion) } @@ -63,51 +105,54 @@ class FlightStateEngineTest { ) assertEquals("15DEC261900", next.scalars["ESTT"]) - assertEquals("CA001", next.scalars["FLNO"]) // 缺失 = 保留(§6.1) + assertEquals("CA001", next.scalars["FLNO"]) // 缺失 = 保留(§3.2) assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"]) assertEquals(2, next.stateVersion) } @Test - fun `validation rejects RECS mismatch and duplicate flid and day outside scope`() { + fun `validation accepts valid snapshot and rejects RECS mismatch`() { val ok = FlightStateEngine.validateMessage( - 1, listOf(ScheduleRecord("121", mapOf("SODT" to "15DEC261723"))), day, day, opDay, + 1, listOf(ScheduleRecord("121", mapOf("SODT" to "15DEC261723"))), opDay, ) assertTrue(ok is SnapshotValidation.Ok) val recsMismatch = FlightStateEngine.validateMessage( - 2, listOf(ScheduleRecord("121", mapOf("SODT" to "15DEC261723"))), day, day, opDay, + 2, listOf(ScheduleRecord("121", mapOf("SODT" to "15DEC261723"))), opDay, ) as SnapshotValidation.Invalid assertTrue(recsMismatch.flags.contains(SnapshotFlag.RECS_DROP)) + } + @Test + fun `validation rejects duplicate flid and non computable operation day`() { val duplicate = FlightStateEngine.validateMessage( 2, listOf( ScheduleRecord("121", mapOf("SODT" to "15DEC261723")), ScheduleRecord("121", mapOf("SODT" to "15DEC261723")), ), - day, day, opDay, + opDay, ) assertTrue(duplicate is SnapshotValidation.Invalid) - val outOfScope = FlightStateEngine.validateMessage( - 1, listOf(ScheduleRecord("121", mapOf("SODT" to "20DEC261723"))), day, day, opDay, + val noSodt = FlightStateEngine.validateMessage( + 1, listOf(ScheduleRecord("121", emptyMap())), opDay, ) as SnapshotValidation.Invalid - assertTrue(outOfScope.flags.contains(SnapshotFlag.DAY_MISMATCH)) // §5.2 覆盖范围校验 + assertTrue(noSodt.flags.contains(SnapshotFlag.DAY_MISMATCH)) // 运营日不可计算 → 整包拒绝 - val empty = FlightStateEngine.validateMessage(0, emptyList(), day, day, opDay) as SnapshotValidation.Ok + val empty = FlightStateEngine.validateMessage(0, emptyList(), opDay) as SnapshotValidation.Ok assertTrue(empty.flags.contains(SnapshotFlag.EMPTY)) } @Test fun `validation rejects non numeric or oversized flid`() { val bad = FlightStateEngine.validateMessage( - 1, listOf(ScheduleRecord("ABC", mapOf("SODT" to "15DEC261723"))), day, day, opDay, + 1, listOf(ScheduleRecord("ABC", mapOf("SODT" to "15DEC261723"))), opDay, ) assertTrue(bad is SnapshotValidation.Invalid) val tooLong = FlightStateEngine.validateMessage( - 1, listOf(ScheduleRecord("1234567890123", mapOf("SODT" to "15DEC261723"))), day, day, opDay, + 1, listOf(ScheduleRecord("1234567890123", mapOf("SODT" to "15DEC261723"))), opDay, ) - assertTrue(tooLong is SnapshotValidation.Invalid) // SIS §3.16.2 Number(1-12) + assertTrue(tooLong is SnapshotValidation.Invalid) // FLID 数字型 Number(1-12) } } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt index b825006..de564b7 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/FlywayMigrationTest.kt @@ -9,7 +9,7 @@ import org.junit.jupiter.api.Test import java.sql.DriverManager /** - * Flyway 迁移引擎端到端验证(docs/flight-state.md §3.2 表职责基线): + * Flyway 迁移引擎端到端验证(docs/flight-state.md §2 权威模型 + design.md §2.1): * 1. V1__flight_state_baseline.sql 在真实 PostgreSQL 上自动迁移成功; * 2. flyway_schema_history 落库且 success = true; * 3. 决策层/管道层/留痕层全表就绪;PIPELINE_LOCK 单行种子就位。 diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt index 35a2c92..a7f11ee 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt @@ -12,7 +12,7 @@ import org.junit.jupiter.api.Test import java.time.Instant /** - * 历史归档与物理清除(docs/flight-state.md §8.2,顺序不可颠倒): + * 历史归档与物理清除(docs/flight-state.md §6,顺序不可颠倒;design.md §6.2): * 历史存储未接通必须删 0 条;先归档确认再物理清除; * 未经 FDEL 的航班在清除前补发删除事件。 */ @@ -53,7 +53,7 @@ class HistorySweepJobTest { val outcome = job.run(now) - assertEquals(0, outcome.purged) // §8.2 红线:先删历史再补当前态绝不允许 + assertEquals(0, outcome.purged) // §6 红线:历史存储未接通必须删 0 条 assertTrue(f.findMainRow("F1") != null) assertEquals(0, events.rows.size) } @@ -78,14 +78,14 @@ class HistorySweepJobTest { assertEquals(2, outcome.selected) assertEquals(1, outcome.archived) assertEquals(1, outcome.purged) - assertEquals(null, f.findMainRow("F1")) // §8.2 步骤 3:确认成功 → 物理删除 + assertEquals(null, f.findMainRow("F1")) // §6:归档确认成功 → 物理删除 assertTrue(f.findMainRow("F2") != null) // 失败或不明确的保留重试(步骤 5) assertEquals(0, events.rows.size) // DELETED 航班清除不再发业务删除事件 } @Test fun `never-fdel lifecycle purge emits tombstone before deletion`() { - val f = seededFlight("F3", deleted = false, idleDays = 30) // ACTIVE 且超兜底窗(§8.1) + val f = seededFlight("F3", deleted = false, idleDays = 30) // ACTIVE 且超兜底窗(§6 静默判据) val events = StubMsgEvents() val store = RecordingHistoryStore() val job = HistorySweepJob( @@ -94,7 +94,7 @@ class HistorySweepJobTest { job.run(now) - // §7.3/§8.2 步骤 4:未经 FDEL 的航班,清除前补发一次删除事件 + // §3.3/§5/§6:未经 FDEL、由生命周期直接清除的航班,清除前补发一次删除事件 val tombstone = events.rows.values.single { it.eventType == EventType.TOMBSTONE } assertEquals("F3", tombstone.partitionKey) assertTrue(tombstone.payloadJson.contains("\"deleted\":true")) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt index 53b95b5..4a85f14 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt @@ -24,9 +24,9 @@ import org.junit.jupiter.api.Test import java.time.LocalDate /** - * FDEL(§6.2)+ ADFT(§6.3/§2.1)不变量: + * FDEL(docs/flight-state.md §3.3)+ ADFT(§3.3/§2.1)不变量: * ACTIVE→DELETED 发布一次 tombstone 且明细保留;重复/迟到 FDEL 幂等不推进版本; - * ADFT 重激活恢复 ACTIVE;新实例建立;缺失标量不误删(§2.1 保守语义)。 + * ADFT 重激活恢复 ACTIVE;新实例建立;缺失标量不误删(§3.3 保守 Set-only)。 */ class FdelAndAdftProcessorTest { @@ -84,7 +84,7 @@ class FdelAndAdftProcessorTest { val version = f.findMainRow("121")!!.stateVersion val result = proc.apply(head(), msg(), FlopPayload("121")) - assertEquals(ApplyResult.Succeeded, result) // §6.2 步骤 3:幂等成功 + assertEquals(ApplyResult.Succeeded, result) // §3.3:重复 FDEL 幂等成功 assertEquals(version, f.findMainRow("121")!!.stateVersion) assertEquals(1, events.rows.values.count { it.eventType == EventType.TOMBSTONE }) } @@ -97,7 +97,7 @@ class FdelAndAdftProcessorTest { val result = FdelProcessor(StubPipelineTx(), StubPipelineLock(), f, events, null, ObjectMapper()) .apply(head(), msg(), FlopPayload("999")) - assertEquals(ApplyResult.Succeeded, result) // §6.2 步骤 4:迟到不报错 + assertEquals(ApplyResult.Succeeded, result) // §3.3:迟到/不存在幂等成功 assertEquals(0, events.rows.size) } @@ -118,7 +118,7 @@ class FdelAndAdftProcessorTest { assertEquals(ApplyResult.Succeeded, result) val main = f.findMainRow("121")!! - assertEquals(FlightState.ACTIVE, main.state) // §6.3 生命周期重激活 + assertEquals(FlightState.ACTIVE, main.state) // §3.3 生命周期重激活 assertEquals(LocalDate.of(2026, 12, 15), main.operationDay) assertEquals("CA002", f.loadFullSnapshot("121")!!.scalars["FLNO"]) } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt index 2414317..fb65e02 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt @@ -32,8 +32,9 @@ import org.junit.jupiter.api.Test import java.time.LocalDate /** - * §5.1 applyScheduleRecords 不变量:重放幂等、完整性校验整包拒绝 DEAD(PROTOCOL)、 - * 归属日冲突不落地、DELETED 不恢复且记冲突告警、成功路径版本推进+事件+留痕+待办预登记。 + * 日计划主链路(docs/flight-state.md §3.1/§4)不变量:重放幂等、整包校验失败 + * DEAD(PROTOCOL)、归属日冲突不落地、DELETED 不恢复且记冲突告警、 + * 成功路径版本推进 + 事件 + 留痕 + 待办预登记。 */ class ScheduleProcessorTest { @@ -107,7 +108,7 @@ class ScheduleProcessorTest { assertEquals(1, log.entries.size) assertEquals(SnapshotResult.COMMITTED, log.entries.single().result) assertEquals(1, log.entries.single().upserted) - assertEquals(1, todo.count()) // §7.2 事务内预登记 + assertEquals(1, todo.count()) // design.md §3.3/§6.1 事务内预登记 } @Test @@ -120,7 +121,7 @@ class ScheduleProcessorTest { val result = processor(proc = proc, flights = flights, log = log) .applyScheduleRecords(head(), message(makeBody("121" to "15DEC261723"))) - assertEquals(ApplyResult.ReplaySkipped, result) // §5.1 步骤 2 + assertEquals(ApplyResult.ReplaySkipped, result) // §4 重放判定:幂等成功 assertNull(flights.findMainRow("121")) assertEquals(SnapshotResult.REPLAY_SKIPPED, log.entries.single().result) } @@ -132,7 +133,7 @@ class ScheduleProcessorTest { val schdBody = ScheduleBody(recsDeclared = 2, records = makeBody("121" to "15DEC261723").records) val result = processor(flights = flights, log = log).applyScheduleRecords(head(), message(schdBody)) as ApplyResult.DeadProtocol - assertTrue(result.flags.contains(SnapshotFlag.RECS_DROP)) // §5.2 RECS = 实收数 + assertTrue(result.flags.contains(SnapshotFlag.RECS_DROP)) // §4:声明数量不符整包拒绝 assertNull(flights.findMainRow("121")) assertEquals(SnapshotResult.ROLLED_BACK, log.entries.single().result) } @@ -140,7 +141,7 @@ class ScheduleProcessorTest { @Test fun `same flid across operation days is rejected whole-batch without writes`() { val flights = StubFlightState() - // 预置 121 已归属 12-15;报文声称 12-16 → §5.3 第四行 + // 预置 121 已归属 12-15;报文声称 12-16 → §2.1 归属冲突 flights.persistFullState( com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot( "121", LocalDate.of(2026, 12, 15), FlightState.ACTIVE, 1, mapOf("SODT" to "15DEC261723"), emptyMap(), @@ -178,7 +179,7 @@ class ScheduleProcessorTest { ) assertEquals(ApplyResult.Succeeded, result) - assertEquals(FlightState.DELETED, flights.findMainRow("121")!!.state) // §5.1 步骤 6 不恢复 + assertEquals(FlightState.DELETED, flights.findMainRow("121")!!.state) // §3.3:日计划不恢复 assertEquals(8L, flights.findMainRow("121")!!.stateVersion) assertEquals(SnapshotResult.COMMITTED, log.entries.single().result) assertTrue(log.entries.single().flags.contains(SnapshotFlag.SCHD_REVIVE_CONFLICT)) @@ -193,7 +194,7 @@ class ScheduleProcessorTest { proc.insert(msgId) val p = processor(proc = proc, flights = flights, log = log, todo = todo, tx = TxRunner { true }) - // §9:数据库/内部故障 → 异常上抛,MessageProcessor 记 FAILED(INFRA) 退避;不写任何终态 + // design.md §2.3:数据库/内部故障 → 异常上抛,MessageProcessor 记 FAILED(INFRA) 退避;不写任何终态 org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException::class.java) { p.applyScheduleRecords(head(), message(makeBody("121" to "15DEC261723"))) }