From 6e7d819034364e6f87589fdde3e9f9447fbb83e6 Mon Sep 17 00:00:00 2001 From: windyboy Date: Thu, 10 Sep 2026 11:09:19 +0800 Subject: [PATCH] =?UTF-8?q?docs(kdoc):=20V2=20=E8=BF=81=E7=A7=BB=E4=B8=8E?= =?UTF-8?q?=E5=A4=84=E7=90=86/=E6=8C=81=E4=B9=85=E5=8C=96=E5=B1=82?= =?UTF-8?q?=E5=89=A9=E4=BD=99=E6=B3=A8=E9=87=8A=E6=94=B9=E4=B8=BA=E7=9B=B4?= =?UTF-8?q?=E7=99=BD=E8=AF=B4=E6=98=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 接上一提交,把本次工作范围内还没改到的注释补齐:V2 迁移脚本的注释原文几乎全是 段落编号引用,处理层与持久化层还留着一批 `(§3.3/§5)` 形式的行内注释。 V2__inbox_lifecycle.sql:文件头改为先讲清这次迁移解决的三个问题(记住收报读到哪里、 回填记录并进 PROC_STATE 而不是单开一张表、记下收信时间做什么用),再落到每个字段; 字段级注释补上"这个值非空代表什么"。SQL 语句一字未动(已用去注释后比对确认)。 processing:Pump 的分派与终态注释、ScheduleProcessor 的重放判定与运营日冲突、 DynamicProcessors 的删除通知与重新激活,都改成说明白"这一步在做什么、为什么这么做"。 infra/persistence:Repositories.kt 的仓储约定、航班状态读写、待发事件与请求跟踪接口, JdbcPgRepositories 的锁、整态合并、历史清理判据、明细表映射,StubRepositories 的 对应实现,一律先说清用途再谈规则。 测试:8 个测试类里的行内引用改为说明这条断言在守什么。 注释里保留的文档指向只在需要延伸阅读时出现,不再作为解释本身。 全部为注释改动,测试仍为 78 passed / 1 skipped。 --- .../infra/persistence/Repositories.kt | 76 ++++++++++++------- .../persistence/jdbc/JdbcPgRepositories.kt | 30 ++++---- .../infra/stub/StubRepositories.kt | 6 +- .../processing/DynamicProcessors.kt | 10 +-- .../gzzn/omms/msgexchange/processing/Pump.kt | 10 +-- .../processing/ScheduleProcessor.kt | 6 +- .../db/migration/V2__inbox_lifecycle.sql | 66 ++++++++++------ .../omms/msgexchange/PipelineSmokeTest.kt | 2 +- .../persistence/jdbc/FlywayMigrationTest.kt | 4 +- .../jdbc/InboxLifecycleJdbcSqlTest.kt | 12 +-- .../infra/retry/ReplayServiceTest.kt | 4 +- .../msgexchange/ingress/InboxPollerTest.kt | 6 +- .../processing/BackfillServiceTest.kt | 12 +-- .../processing/FdelAndAdftProcessorTest.kt | 12 +-- .../processing/ScheduleProcessorTest.kt | 12 +-- 15 files changed, 154 insertions(+), 114 deletions(-) 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 78b2b71..c78344f 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,17 +14,21 @@ import java.time.LocalDate import java.time.ZoneId // ===================================================================== -// 仓储契约(docs/flight-state.md §2 权威模型/§3 写入语义,design.md §2.1)。 -// 决策层写路径约定:所有 FLIGHT_SCHD 及明细写操作必须发生在 -// 「持有 PIPELINE_LOCK 的同一事务」内(§1/§4:锁只串行化 DB 事务)。 +// 仓储接口定义。 +// +// 有一条约定贯穿航班相关的写操作:改航班状态(FLIGHT_SCHD 和明细表)必须发生在 +// 同一个事务里,并且进事务后第一件事就是拿 PIPELINE_LOCK 这把单行锁。锁只负责让 +// 写事务排队,不负责选主或故障切换。 +// +// 航班模型与合并规则的完整说明见 docs/flight-state.md。 // ===================================================================== -/** 自有 PG 单事务原子保障(flight-state.md §1/§4:状态、事件、处理终态同事务提交)。 */ +/** 把一段代码包进一个数据库事务:航班状态、待发事件、处理终态要么一起提交,要么一起回滚。 */ interface PipelineTransactionManager { fun inTransaction(block: () -> T): T } -/** 单行锁 PIPELINE_LOCK:事务内第一步 SELECT ... FOR UPDATE,串行化状态写事务(§1/§4)。 */ +/** 单行锁:进事务后先 `SELECT ... FOR UPDATE`,让并发的状态写事务排队执行。 */ interface PipelineLockRepository { fun lock() } @@ -117,7 +121,12 @@ data class BackfillDue(val msgId: Long, val attempts: Int) */ data class Backlog(val unfinished: Int, val oldestReceivedAt: Instant?, val unmarkedTerminal: Int) -/** MSG_EVENT outbox(flight-state.md §5)。KAFKA_SCHD 合并同 FLID 未发事件按最新 STATE_VERSION 输出。 */ +/** + * 待发事件(outbox)表:业务提交时把要发的事件一起写进来,投递线程再从这张表往外发, + * 这样发送失败不会影响业务事务。 + * + * 其中 KAFKA_SCHD 是"最新状态"通知:同一条航班的多个未发事件只留最新的那个版本。 + */ interface MsgEventRepository { fun insertAll(events: List): List @@ -125,7 +134,7 @@ interface MsgEventRepository { fun claimBatch(target: String, limit: Int): List - /** 同 FLID 未发 KAFKA_SCHD 事件合并:每 FLID 取最新 STATE_VERSION 一条(§5)。 */ + /** 挑出待发的整态事件,同一条航班只取版本最高的那一条(中间的版本不用发)。 */ fun mergePendingSchd(limit: Int): List fun markSent(eventId: Long) @@ -137,63 +146,74 @@ interface MsgEventRepository { fun markDead(eventId: Long, errorClass: ErrorClass, lastError: String, attempts: Int? = null) } -/** 完整态落库结果:DAY_GUARD_VIOLATION = OPERATION_DAY 不可变条件更新未命中(§2.1)。 */ +/** + * 落库结果。 + * [DAY_GUARD_VIOLATION] 表示这次写入被"运营日不可变"的条件挡住:库里已有运营日, + * 和这次的报文对不上,说明这条报文可能串了日子。 + */ enum class PersistOutcome { INSERTED, UPDATED, DAY_GUARD_VIOLATION } /** - * 航班当前态权威(§2 FLIGHT_SCHD + 9 张明细表)。 - * 唯一写路径 = persistFullState:写入前在内存生成完整新状态(引擎产物)再落库, - * 明细集合按组先删后插(§2.2 完整合并结果为准)。 + * 航班当前态的读写入口:主表 FLIGHT_SCHD 加 9 张明细表。 + * + * 写只有一个入口 [persistFullState]:先在内存里算出合并后的完整状态,再整体落库。 + * 明细数据按组先删后插,所以入库的永远是"合并后的完整结果",不依赖增量更新。 */ interface FlightStateRepository { - /** 主行点查(身份/版本判定;批量用于日计划归属校验 §2.1)。 */ + /** 查航班主行,用于判断它是否存在、属于哪个运营日、当前版本是多少。 */ fun findMainRow(flid: String): FlightMainRow? fun findMainRows(flids: Collection): Map - /** 完整当前态:主行 + 全部明细(一致性读边界由调用方事务保证,§5)。 */ + /** 读一条航班的完整当前态(主行加全部明细)。要读得一致,得由调用方放在事务里读。 */ fun loadFullSnapshot(flid: String): FlightSnapshot? /** - * 完整当前态落库(§3):主行 upsert + 全部明细按组先删后插; - * STATE_VERSION 以 snapshot.stateVersion 落库。 - * SCHD(§3.1)、FLOP/ADFT(§3.2/§3.3)合并后的全量结果共用此唯一写路径。 - * 条件更新带 `WHERE operation_day IS NULL OR operation_day = :day`(§2.1 不可变强化)。 + * 把合并后的完整状态写库:主行做 upsert,明细按组先删后插, + * 版本号用 [snapshot] 里算好的那个。 + * + * 日计划、动态事件、异常航班的合并结果都走这一个入口。 + * 写入带运营日条件(`WHERE operation_day IS NULL OR operation_day = :day`): + * 一旦航班已经归属某个运营日,别的日子的报文就写不进来。 */ fun persistFullState(snapshot: FlightSnapshot, msgId: Long, now: Instant): PersistOutcome /** - * FDEL(§3.3):ACTIVE → 置 DELETED、推进版本、明细保留、返回 true(发布删除事件); - * 已 DELETED 或不存在 → 返回 false(幂等成功,不推进版本不重复发布)。 + * 删除航班:把在用的航班标成 DELETED、版本号加一,明细数据保留不删。 + * + * 返回 true 表示这次真的删了(调用方据此发一次删除通知);返回 false 表示航班已经 + * 是删除状态或根本不存在,按成功处理,不推进版本也不重复发通知。 */ fun markDeleted(flid: String, msgId: Long, now: Instant): Boolean - /** ADFT 生命周期重激活(§3.3):DELETED → ACTIVE,推进版本;非 DELETED 返回 false。 */ + /** 把已删除的航班重新激活(ADFT 报文触发):DELETED → ACTIVE、版本号加一;不是删除状态则返回 false。 */ fun revive(flid: String, msgId: Long, now: Instant): Boolean - /** §6:按保留期与终态/静默判据选出历史候选(含 DELETED;窗口按机场时区折算)。 */ + /** 按保留期挑出可以归档清理的航班(含已删除的)。时间窗口按机场时区计算。 */ fun findHistoryCandidates(rules: HistoryRules, zone: ZoneId, now: Instant): List /** - * §6:物理删除主行与明细(仅历史存储确认成功后调用; - * 历史存储未接通时调用方必须传空集合——删 0 条)。 + * 物理删除主行与明细。**只能在历史存储确认归档成功之后调用**; + * 历史存储没接通时,调用方必须传空集合,也就是一条都不删。 */ fun purgeArchived(flids: Collection): Int - /** §2.1 观测:OPERATION_DAY 仍为 NULL 的航班数(只增不删,终止规则未定 §6)。 */ + /** 运营日还没定下来的航班有多少条(观测用;这些航班暂时不参与清理)。 */ fun countOperationDayNull(): Int } -/** SCHD_SNAP_LOG(design.md §6.2):事务外追加留痕,不参与决策;一行 = 一次尝试(重放也记)。 */ +/** 日计划处理留痕:只追加、不参与业务判断,一次处理(含重放)记一行,写失败不影响业务。 */ interface SnapshotLogRepository { fun append(entry: SnapshotLogEntry) } /** - * 请求状态机 REQ_TRACK(design.md §4.2 目标机制,尚无运行时协调器): - * 只有 RESP 完成 RQFD 请求,按(运营日、发送方、请求类型)匹配最新一条 PENDING; - * 同类请求只留一条有效,新请求置旧为 EXPIRED。 + * 上游请求的跟踪表:记录我们发出去的请求、以及对方回来的应答。 + * + * 目前只有表和读写方法,**还没有运行时协调器**(出站写信箱、超时、应答匹配都没实现)。 + * 设计意图是:RESP 报文按(运营日、发送方、请求类型)匹配最近一条待应答的请求; + * 同一类请求同时只保留一条有效,发新请求时把旧的置为已过期。 */ interface ReqTrackRepository { enum class ReqState { PENDING, SENT, DONE, EXPIRED } 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 4d8fc11..e314c83 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 @@ -38,9 +38,10 @@ import java.util.Locale import javax.sql.DataSource // ===================================================================== -// 自有 PostgreSQL 仓储实现(docs/flight-state.md §2 权威模型/§3 写入语义;schema 见 -// db/migration/V1__flight_state_baseline.sql)。全部写路径约定在 -// withTransaction + PIPELINE_LOCK 内调用(§1/§4)。 +// 自有 PostgreSQL 的仓储实现,表结构见 db/migration/V1__flight_state_baseline.sql。 +// +// 航班相关的写操作都要求调用方先开事务、再拿 PIPELINE_LOCK 单行锁(见 Repositories.kt +// 顶部说明)。航班模型与合并规则见 docs/flight-state.md。 // ===================================================================== @Singleton @@ -58,7 +59,7 @@ class JdbcPipelineTransactionManager( class JdbcPipelineLockRepository( private val ds: DataSource, ) : PipelineLockRepository { - /** 事务内第一步:单行 FOR UPDATE 串行化状态写事务(§1/§4 PIPELINE_LOCK)。 */ + /** 进事务后的第一步:对锁行 `FOR UPDATE`,让并发的状态写事务排队。 */ override fun lock() { ds.queryOne("SELECT lock_id FROM pipeline_lock WHERE lock_id = 1 FOR UPDATE", {}) { 1 } ?: error("PIPELINE_LOCK row missing") } @@ -340,7 +341,7 @@ class JdbcMsgEventRepository( ::mapEvent, ) - /** §5:同 FLID 未发事件按最新 STATE_VERSION 合并;PG DISTINCT ON 方言(注释明示)。 */ + /** 同一条航班只取版本最高的待发整态事件;用了 PostgreSQL 的 DISTINCT ON 语法。 */ override fun mergePendingSchd(limit: Int): List = ds.query( """ @@ -477,7 +478,7 @@ class JdbcFlightStateRepository( }, ) if (updated == 0) { - // §2.1 不可变强化:条件更新未命中 = 归属日冲突 + // 一行都没更新到,说明被运营日条件挡住了:这条航班已经属于别的运营日 return if (existed != null) PersistOutcome.DAY_GUARD_VIOLATION else PersistOutcome.UPDATED } replaceDetails(snapshot, now) @@ -499,7 +500,7 @@ class JdbcFlightStateRepository( ) == 1 override fun findHistoryCandidates(rules: HistoryRules, zone: ZoneId, now: Instant): List { - // §6 四条判定在应用层执行(SIS 时间串解析无法下推 SQL); + // 四条清理判据都在应用层算:里面的时间字段是 SIS 格式的字符串,没法交给 SQL 比较; // 不按 updated_at 粗筛——取消/终态时间可能早于最近一次更新,粗筛会漏删。 val rows = ds.query( "SELECT flid, state, state_version, operation_day, cncl, naat, neat, updated_at FROM flight_schd", @@ -528,7 +529,7 @@ class JdbcFlightStateRepository( r.updatedAt < now.minus(Duration.ofHours(rules.idleHours)) -> true else -> false } - // NAAT/NEAT 业务含义待术语表确认(§6 开放项):解析失败视为不命中,不误删 + // NAAT/NEAT 的业务含义还没确认,所以解析不出来时一律当成"不满足条件",宁可漏删 if (hit) HistoryCandidate(r.flid, r.state, r.version, wasNeverFdel = r.state == FlightState.ACTIVE) else null } } @@ -585,7 +586,7 @@ class JdbcFlightStateRepository( items.forEachIndexed { ordinal, item -> val cols = mutableListOf("flid", "ordinal", "source_seq") val vals = mutableListOf(flid) - vals.add(ordinal + 1) // ORDINAL 保留输入顺序(§2.2) + vals.add(ordinal + 1) // 序号按报文里的先后顺序编,读回来顺序不变 vals.add(item[spec.seqAttr]) spec.columns.forEach { col -> cols.add(col) @@ -617,7 +618,7 @@ class JdbcFlightStateRepository( private fun loadDetails(flid: String, key: String, spec: DetailSpec): List> { val sql = if (spec.routeKind != null) { - // ROUT/ERUT 共表:按 ROUTE_KIND 过滤(§2.2) + // ROUT 和 ERUT 存在同一张表里,靠 ROUTE_KIND 区分,读的时候要带上这个条件 "SELECT * FROM ${spec.table} WHERE flid = ? AND route_kind = ? ORDER BY ordinal ASC" } else { "SELECT * FROM ${spec.table} WHERE flid = ? ORDER BY ordinal ASC" @@ -673,7 +674,7 @@ class JdbcFlightStateRepository( "flight_delay", "flight_bridge_op", "flight_chock_op", "flight_route_point", ) - /** 10 类集合 ↔ 明细表/列映射(§2.2;列名与基线一致)。 */ + /** 报文里的 10 类集合分别存到哪张表、哪些列(列名与基线脚本一致)。 */ 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"), @@ -695,7 +696,7 @@ class JdbcFlightStateRepository( class JdbcSnapshotLogRepository( private val ds: DataSource, ) : SnapshotLogRepository { - /** design.md §6.2:只追加;写失败由调用方捕获记指标(append 自身不抛出)。 */ + /** 只追加写;这里不抛异常,写失败由调用方捕获并记为指标。 */ override fun append(entry: SnapshotLogEntry) { ds.update( """ @@ -757,7 +758,10 @@ class JdbcReqTrackRepository( ) } - /** design.md §4.2 RESP 完成请求:匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */ + /** + * 应答报文到达时,把对应请求标记为已完成:找最近一条待应答的请求(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/StubRepositories.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/stub/StubRepositories.kt index 21a6d6e..7cc5c0a 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 @@ -199,7 +199,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) - /** §5:同 FLID 取最新 STATE_VERSION,按事件序输出。 */ + /** 同一条航班只取版本最高的待发整态事件,按事件先后输出。 */ override fun mergePendingSchd(limit: Int): List = rows.values .filter { it.target == "KAFKA:schd" && it.state == EventStatus.PENDING } @@ -245,7 +245,7 @@ class StubFlightState : FlightStateRepository { override fun loadFullSnapshot(flid: String): FlightSnapshot? = snapshots[flid] - /** §2.1 不可变条件:已有非空 OPERATION_DAY 且与新值不同 → DAY_GUARD_VIOLATION。 */ + /** 运营日不可变:库里已有运营日且与新值不同时,返回 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) { @@ -327,7 +327,7 @@ class StubReqTrack : ReqTrackRepository { fun clear() = rows.clear() override fun insert(reqType: String, operationDay: LocalDate, sender: String): Long { - // 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/processing/DynamicProcessors.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt index 1ed87bc..43a9217 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt @@ -78,7 +78,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 时登记(§3.3/§5),与删除同事务 + // 只有"在用 → 删除"这一步才发删除通知,而且和状态变更写在同一个事务里 msgEvents.insertAll( listOf( MsgEvent( @@ -105,7 +105,7 @@ class FdelProcessor( ), ) } - procState.markTerminal(head.msgId, ProcStatus.SUCCEEDED) // 未命中 = 迟到/重复,幂等成功(§3.3) + procState.markTerminal(head.msgId, ProcStatus.SUCCEEDED) // 没删到东西说明是迟到或重复报文,照样算成功 ApplyResult.Succeeded } } @@ -137,7 +137,7 @@ class AdftProcessor( lock.lock() val main = flightState.findMainRow(record.flid) if (main != null && main.state == FlightState.DELETED) { - // §3.3 重激活:DELETED → ACTIVE,推进版本,登记状态事件 + // 已删除的航班重新激活:状态改回 ACTIVE、版本号加一,并登记状态事件 if (flightState.revive(record.flid, msgId = head.msgId, now = Instant.now())) { val current = flightState.loadFullSnapshot(record.flid) if (current != null) { @@ -152,11 +152,11 @@ class AdftProcessor( val current = flightState.loadFullSnapshot(record.flid) val next: FlightSnapshot = if (current == null) { - // 新实例建立:ADFT 含 SODT 时直接计算运营日(§2.1),不可算则置 null 待日计划收录 + // 新建航班:带了计划时间就算出运营日,算不出来先留空,等日计划报文来收录 val day = opDay.compute(record.scalars["SODT"]) FlightSnapshot( flid = record.flid, - operationDay = day, // §2.1:不可算时不得默认写接收日 + operationDay = day, // 算不出来就留空,不能默认拿收报当天顶上 state = FlightState.ACTIVE, stateVersion = 1L, scalars = record.scalars, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt index 84cac69..8c15c71 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -129,8 +129,8 @@ class MessageProcessor( log.warn("processOne unexpected failure msgId={} ec=INFRA msg={}", head.msgId, e.message ?: e.javaClass.simpleName) procFailure.fail(head, ErrorClass.INFRA, e.message ?: e.javaClass.simpleName) } - // 终态已提交:立即尝试一次回填;失败留待回填扫描按退避重试(意图已在终态事务内登记)。 - // 中间态(PENDING/FAILED)不适用(message-lifecycle §5.2)。 + // 终态已经写好了,马上试一次回填;写不进去也没关系,回填扫描会按退避继续重试 + // (待办在写终态时就一起登记了)。还没处理完的消息不打标记。 if (terminal) backfill.attempt(head.msgId) } } @@ -178,7 +178,7 @@ class MessageProcessor( } } - // 处理器分派:SCHD 日计划主链路(§3.1)/ FLOP / FDEL / ADFT;缺载荷按 MALFORMED 终态 + // 按报文类型分派:日计划走 SCHD,其余走 FLOP / FDEL / ADFT;报文缺载荷直接判为非法报文的死信 val result: ApplyResult = when (val kind = decoded.kind) { is MsgKind.Schd -> { val body = decoded.body as? ScheduleBody @@ -204,14 +204,14 @@ class MessageProcessor( flopProcessor.apply(head, decoded, payload) } is MsgKind.Unsupported -> { - // design.md §2.3:未支持类型 → FAILED(UNSUPPORTED) 退避重试,达阈值转 DEAD;绝不写终态 + // 还没有对应处理器的报文类型:先按可重试的失败处理,等能力补齐,不直接判死 log.warn("unsupported type -> FAILED(UNSUPPORTED) msgId={} tag={}", head.msgId, kind.tag) return procFailure.fail(head, ErrorClass.UNSUPPORTED, "no-handler:${kind.tag}") } } if (result is ApplyResult.DeadProtocol) { - // design.md §2.3:整包拒绝 DEAD(PROTOCOL),立即释放队头,交人工确认 + // 整包被拒绝:不重试、立刻放掉队头,等人工确认 log.error("DEAD(PROTOCOL) msgId={} reason={} flags={}", head.msgId, result.reason, result.flags) procState.markTerminal( head.msgId, ProcStatus.DEAD, 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 553814b..4374353 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt @@ -76,7 +76,7 @@ class ScheduleProcessor( val body = msg.body as? ScheduleBody ?: return ApplyResult.DeadProtocol("missing-schd-body") val started = System.nanoTime() - // 重放判定: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 @@ -123,7 +123,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) // §3.3:日计划不复活 DELETED + if (keepDeleted) flags.add(SnapshotFlag.SCHD_REVIVE_CONFLICT) // 日计划不会让已删除的航班复活 val current = if (existingMain != null) flightState.loadFullSnapshot(flid) else null val next = FlightStateEngine.snapshotState( current = current, @@ -151,7 +151,7 @@ class ScheduleProcessor( } } - /** 状态事件:KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1 登记、§5 语义)。 */ + /** 每次航班状态变化登记两个事件:KAFKA_SCHD 发完整状态,KAFKA_MSG 发一条变更通知。 */ private fun snapshotEvents(next: FlightSnapshot): List { val payload = linkedMapOf( "flid" to next.flid, diff --git a/src/main/resources/db/migration/V2__inbox_lifecycle.sql b/src/main/resources/db/migration/V2__inbox_lifecycle.sql index 888016f..e7f86c8 100644 --- a/src/main/resources/db/migration/V2__inbox_lifecycle.sql +++ b/src/main/resources/db/migration/V2__inbox_lifecycle.sql @@ -1,44 +1,60 @@ -- ===================================================================== --- V2:信箱生命周期边界层(权威设计:docs/message-lifecycle.md §2/§4/§5.1/§5.2) +-- V2:收报进度与回填记录的存放方式调整 -- --------------------------------------------------------------------- --- 三处契约级修正,均只动自有 PG(共享 MySQL 不建表、不改结构,红线不破): --- 1. 消费水位 W 落库(INBOX_CURSOR):只随新 ID 成功入队推进、遇空洞即停, --- 与入队同事务——中断后 W 未前进,重扫即补建(§4 第一行); --- HOLE_SINCE 支撑"空洞老化",避免一次自增回滚空位永久停摆水位(§5.1)。 --- 2. 回填事实并入 PROC_STATE:终态与回填意图同一条 UPDATE 落下, --- 消除"业务已提交、回填待办未记"与"待办二次落账失败"两个崩溃窗口(§4); --- BACKFILL_TODO 随之下线(它承载的 sndr/type/styp/seqn 从未参与回填, --- 共享信箱回填只需要 MSG_ID——§5.2 补写同样如此)。 --- 3. 超期补写判据 RECEIVED_AT(§5.2 的 R)与 OPS-2「最老未处理信龄」锚点。 +-- 这次迁移解决三个问题,全部只动自有 PostgreSQL。共享 MySQL 那边不建表、 +-- 不改结构,边界没有变化。 +-- +-- 1) 收报读到哪儿了要能记住(新增 INBOX_CURSOR) +-- 原先每一轮都从 0 开始扫信箱,靠"有没有处理标记"判断该不该取。这个做法有问题: +-- 处理完但还没把标记写回信箱的行(以及永远不会回填的死信)会一直占着每批的 +-- 名额,攒够一批之后新消息就再也读不到了。 +-- 改成记住"读到哪个 ID 了"(水位),每轮只往后读,并且登记和水位推进放在同一个 +-- 事务里——中途崩溃时水位没动,重启后重扫一遍就补齐了。 +-- 水位还有个附带规则:如果后面的 ID 缺号,先停下来等(可能是上游还没提交完), +-- 等太久就认定它不会来了、跳过去继续,否则水位会卡在第一个空位上再也不动。 +-- +-- 2) 回填记录并进 PROC_STATE,不再单独建表 +-- 处理完有两件事要做:记下终态、把处理标记写回信箱。原先第二步靠一张独立的 +-- 待办表,两张表两次写入,中间崩溃就会出现"业务处理完了却没人记得回填"。 +-- 现在这两件事是同一条 UPDATE:在一个事务里提交,要么都成、要么都不做。 +-- BACKFILL_TODO 因此下线——它保存的 sndr/type/styp/seqn 在回填时从来没用上, +-- 回填只需要一个消息 ID。 +-- +-- 3) 记下收信时间(新增 RECEIVED_AT) +-- 有两个用途:判断一条消息等了多久还是没能回填(超期就强制补写), +-- 以及统计"最老一条未处理消息收了多久",供运维观察积压。 -- ===================================================================== --- ① 消费水位:单行游标(CURSOR_ID 恒为 1) +-- ① 收报水位:全表只有一行 CREATE TABLE INBOX_CURSOR ( - CURSOR_ID INT NOT NULL PRIMARY KEY, -- 恒为 1 - COMMITTED_UP_TO BIGINT NOT NULL, -- 水位 W:连续上界,无空洞 - HOLE_SINCE TIMESTAMP(6) WITH TIME ZONE, -- W+1 处空洞首次观测时刻;无空洞为 NULL + CURSOR_ID INT NOT NULL PRIMARY KEY, -- 固定为 1 + COMMITTED_UP_TO BIGINT NOT NULL, -- 水位:到哪个 ID 为止已经全部读进自有库 + HOLE_SINCE TIMESTAMP(6) WITH TIME ZONE, -- 后面那个缺号最早是什么时候发现的;不缺号时为 NULL UPDATED_AT TIMESTAMP(6) WITH TIME ZONE NOT NULL ); --- 初值 0:首轮从最小 ID 起重扫全部信箱行,入队幂等(ON CONFLICT DO NOTHING), --- 上线即自愈既有"已入队未回填"造成的窗口污染。 + +-- 初值 0:上线后第一轮会把信箱里所有行重新读一遍,重复登记不会建出第二行, +-- 所以顺手把历史上"已入队但没回填"造成的漏读一并补上。 INSERT INTO INBOX_CURSOR (CURSOR_ID, COMMITTED_UP_TO, HOLE_SINCE, UPDATED_AT) VALUES (1, 0, NULL, now()); --- ② PROC_STATE:接收时间 + 回填事实(替代 BACKFILL_TODO) -ALTER TABLE PROC_STATE ADD COLUMN RECEIVED_AT TIMESTAMP(6) WITH TIME ZONE; -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_AT TIMESTAMP(6) WITH TIME ZONE; -- 非空 = 已确认信箱行持有标记 -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_NEXT_AT TIMESTAMP(6) WITH TIME ZONE; -- 非空 = 待回填(终态同事务登记) +-- ② PROC_STATE:加上收信时间与回填进度(取代 BACKFILL_TODO) +ALTER TABLE PROC_STATE ADD COLUMN RECEIVED_AT TIMESTAMP(6) WITH TIME ZONE; -- 信箱里的接收时间 +ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_AT TIMESTAMP(6) WITH TIME ZONE; -- 有值 = 已确认信箱行带上了标记 +ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_NEXT_AT TIMESTAMP(6) WITH TIME ZONE; -- 有值 = 还欠一次回填,写终态时一起写下 ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ATTEMPTS INT NOT NULL DEFAULT 0; -ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ERROR VARCHAR(512); +ALTER TABLE PROC_STATE ADD COLUMN BACKFILL_ERROR VARCHAR(512); -- 回填失败的原因,供排查 --- 存量行接收时间:无法跨库回读 DATE_RECEIVED,以入队时间为下界(§5.2 判据只会偏晚,不会提前补写) +-- 存量数据补收信时间:跨库读不到信箱里的 DATE_RECEIVED,只能拿入队时间当兜底。 +-- 这个值只会偏晚,宁可晚一点补写标记,也不会提前把还会重放的消息标掉。 UPDATE PROC_STATE SET RECEIVED_AT = UPDATED_AT WHERE RECEIVED_AT IS NULL; --- 存量终态行补登记回填意图:立即交给回填扫描(已持有标记的行由空标守卫返回 0 行,幂等) +-- 存量里已经处理完的行补一条回填待办,交给回填扫描处理。 +-- 其中已经打过标记的行会被"只写空标记"的条件挡住,重复执行没有副作用。 UPDATE PROC_STATE SET BACKFILL_NEXT_AT = now() WHERE STATE IN ('SUCCEEDED', 'SKIPPED', 'DEAD') AND BACKFILL_AT IS NULL; --- 回填扫描索引:仅覆盖未确认标记的行 +-- 回填扫描用的索引,只覆盖还没确认回填的行 CREATE INDEX idx_proc_backfill_due ON PROC_STATE (BACKFILL_NEXT_AT) WHERE BACKFILL_AT IS NULL; --- ③ BACKFILL_TODO 下线(回填意图已并入 PROC_STATE,机制收敛为一套;索引随表一并释放) +-- ③ BACKFILL_TODO 下线:回填进度已经并进 PROC_STATE,机制只剩一套,索引随表一起释放 DROP TABLE BACKFILL_TODO; diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt index e60c1b0..847bcb0 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt @@ -60,7 +60,7 @@ class PipelineSmokeTest { } companion object { - /** 合法 META,但 TYPE 未在支持范围 → FAILED(UNSUPPORTED)(design.md §2.3) */ + /** 报文头合法,但类型没有对应处理器,会走"暂不支持、先重试"这条路 */ val UNSUPPORTED_XML = """ AODB120260908120000XYZQFOO 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 f78550f..26dc397 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 @@ -77,7 +77,7 @@ class FlywayMigrationTest { assertEquals(setOf("flid", "operation_day", "state", "state_version", "last_msg_id"), cols) } - // 回填事实并入 PROC_STATE(§5.2 判据 + 回填待办) + // 回填相关的列都落在 PROC_STATE 上(收信时间用来判断超期,其余记录回填进度) stmt.executeQuery( "SELECT column_name FROM information_schema.columns WHERE table_name = 'proc_state' " + "AND column_name IN ('received_at', 'backfill_at', 'backfill_next_at', " + @@ -91,7 +91,7 @@ class FlywayMigrationTest { ) } - // PIPELINE_LOCK 与 INBOX_CURSOR 单行种子(§3.2 / §5.1) + // 两张单行表(管道锁、收报水位)的种子数据都要在 stmt.executeQuery("SELECT count(*) FROM pipeline_lock WHERE lock_id = 1").use { rs -> assertTrue(rs.next()) assertEquals(1, rs.getInt(1)) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/InboxLifecycleJdbcSqlTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/InboxLifecycleJdbcSqlTest.kt index c1e1e3a..14095ec 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/InboxLifecycleJdbcSqlTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/InboxLifecycleJdbcSqlTest.kt @@ -70,7 +70,7 @@ class InboxLifecycleJdbcSqlTest { val row = proc.find(11L)!! assertEquals(ProcStatus.DEAD, row.state) assertEquals(ErrorClass.MALFORMED, row.errorClass) - assertEquals(t0, row.backfillNextAt) // 终态与回填意图同一条 UPDATE(§2/§4) + assertEquals(t0, row.backfillNextAt) // 写终态的同时就登记了回填待办(同一条 UPDATE) assertNull(row.backfillAt) assertEquals(t0, row.receivedAt) } @@ -83,7 +83,7 @@ class InboxLifecycleJdbcSqlTest { proc.recordBackfillFailure(11L, "mysql-down", attempts = 2, nextAttemptAt = t0.plusSeconds(120), now = t0) val row = proc.find(11L)!! - assertEquals(ProcStatus.SUCCEEDED, row.state) // §11 终态不可逆 + assertEquals(ProcStatus.SUCCEEDED, row.state) // 回填失败不改处理结果 assertEquals(2, row.backfillAttempts) assertEquals(t0.plusSeconds(120), row.backfillNextAt) assertEquals("mysql-down", row.backfillError) @@ -98,11 +98,11 @@ class InboxLifecycleJdbcSqlTest { seed(2L, t0) proc.markTerminal(2L, ProcStatus.SUCCEEDED, now = t0) proc.recordBackfillFailure(2L, "mysql-down", 1, t0.plus(Duration.ofMinutes(15)), t0) - // ③ 未到期但接收时间已达超期期限 R(§5.2 覆盖退避) + // ③ 退避还没到,但收信时间已经很久了:应当无视退避直接补写 seed(3L, overdue) proc.markTerminal(3L, ProcStatus.DEAD, errorClass = ErrorClass.MALFORMED, now = t0) proc.recordBackfillFailure(3L, "mysql-down", 1, t0.plus(Duration.ofMinutes(15)), t0) - // ④ 中间态:永不补写(§5.2) + // ④ 还没处理完:不补写 seed(4L, overdue) // ⑤ 已确认标记 seed(5L, t0) @@ -124,7 +124,7 @@ class InboxLifecycleJdbcSqlTest { val backlog = proc.backlog() assertEquals(1, backlog.unfinished) - assertEquals(overdue, backlog.oldestReceivedAt) // OPS-2 最老未处理信龄锚点 + assertEquals(overdue, backlog.oldestReceivedAt) // 最老一条未处理消息的接收时间 assertEquals(1, backlog.unmarkedTerminal) } @@ -150,7 +150,7 @@ class InboxLifecycleJdbcSqlTest { assertEquals(third, mailbox.maxId()) assertTrue(mailbox.markProcessedIfUnmarked(second, "PROCESSED")) - assertFalse(mailbox.markProcessedIfUnmarked(second, "OTHER")) // §11 只把空标写为已处理 + assertFalse(mailbox.markProcessedIfUnmarked(second, "OTHER")) // 已有标记,不覆盖 assertEquals("PROCESSED", statusOf(second)) // 已标记行仍出现在区间读结果中(发现与标记彻底解耦) assertEquals(listOf(first, second, third), mailbox.readRange(0L, 50).map { it.msgId }) diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/retry/ReplayServiceTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/retry/ReplayServiceTest.kt index 4182ca5..d45f62a 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/retry/ReplayServiceTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/retry/ReplayServiceTest.kt @@ -10,8 +10,8 @@ import org.junit.jupiter.api.Test import java.time.Instant /** - * U11:显式重放入口只允许可恢复错误从 FAILED/DEAD 返回 PENDING(attempts 清零、立即重试); - * MALFORMED(报文非法)永不被重放。 + * 人工重放的规矩:只有"可以重放"的失败原因才能把消息从 FAILED / DEAD 拉回队列 + * (重试次数清零,马上再试一次);报文本身非法的记录永远不允许重放。 */ class ReplayServiceTest { diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt index d46bd19..f8d8d2f 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt @@ -52,14 +52,14 @@ class InboxPollerTest { assertEquals(ProcStatus.PENDING, proc.find(second)!!.state) assertEquals(second, cursor.cursor.committedUpTo) assertNull(cursor.cursor.holeSince) - // §5.3:消化阶段只写自有 PG,不触碰信箱标记 + // 收报只写自有库,不碰信箱的处理标记 assertFalse(inbox.isMarked(first)) assertEquals(0, poller.pollOnce(t0)) // 重复扫描幂等 } /** - * 回归(US-01 条目 3 / §5.3「积压挡批」):终态且永不回填的行(解码失败死信等)曾占满 - * 有限批次使收报整体停摆——发现必须与处理标记彻底解耦。 + * 回归用例:处理完却永远不会回填的行(典型是解码失败的死信)曾经占满每一批的名额, + * 导致收报整体停摆。取新消息这件事必须和"有没有处理标记"彻底分开。 */ @Test fun `terminal rows without a mailbox mark do not block discovery of later messages`() { diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt index 11e1de5..f19b14b 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt @@ -32,7 +32,7 @@ class BackfillServiceTest { private val t0: Instant = Instant.parse("2026-09-08T03:00:00Z") private val props = PipelineProps() - /** 可注入故障的信箱:验证回填失败路径(§3 失败处理)。 */ + /** 可以人为制造故障的信箱,用来验证回填失败时怎么处理。 */ private class FakeMailbox(var fail: Boolean = false) : CminmsgInboxRepository { val marked = linkedSetOf() override fun insertRaw(rawXml: String): Long = 1L @@ -76,7 +76,7 @@ class BackfillServiceTest { val id = inbox.insertRaw("") assertTrue(inbox.markProcessedIfUnmarked(id, "PROCESSED")) - assertFalse(inbox.markProcessedIfUnmarked(id, "OTHER")) // §11 只把空标写为已处理 + assertFalse(inbox.markProcessedIfUnmarked(id, "OTHER")) // 已经有标记了,不再写第二次 assertEquals("PROCESSED", inbox.markOf(id)) } @@ -89,7 +89,7 @@ class BackfillServiceTest { val backfill = service(proc, inbox) backfill.attempt(id) - backfill.attempt(id) // 重复执行无副作用(§5.2 补写只针对空标记) + backfill.attempt(id) // 重复补写没有副作用 assertNull(proc.find(id)!!.backfillError) assertEquals(0, proc.find(id)!!.backfillAttempts) @@ -104,7 +104,7 @@ class BackfillServiceTest { service(proc, mailbox).attempt(901L) val row = proc.find(901L)!! - assertEquals(ProcStatus.SUCCEEDED, row.state) // §11 终态不可逆 + assertEquals(ProcStatus.SUCCEEDED, row.state) // 回填失败不会把处理结果改回去 assertEquals(1, row.backfillAttempts) assertEquals("mysql-down", row.backfillError) assertEquals(t0.plus(Duration.ofSeconds(30)), row.backfillNextAt) @@ -129,7 +129,7 @@ class BackfillServiceTest { assertNotNull(proc.find(901L)!!.backfillAt) } - /** §5.2:超期期限 R 覆盖退避,保证库方清除前提「边界内无未标记行」在有限时间内成立。 */ + /** 等得太久的消息不再等退避、直接补写:保证标记最终一定会写上。 */ @Test fun `overdue rows bypass the retry backoff`() { val proc = StubProcState() @@ -145,7 +145,7 @@ class BackfillServiceTest { assertNotNull(proc.find(id)!!.backfillAt) } - /** §5.2:中间态(PENDING / FAILED)不适用超期补写——处理未完成时不打标。 */ + /** 还没处理完的消息(PENDING / FAILED)永远不打标,等再久也不行。 */ @Test fun `mid states are never marked even when far past the deadline`() { val proc = StubProcState() 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 ece8bbd..4a0f60b 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt @@ -73,7 +73,7 @@ class FdelAndAdftProcessorTest { assertEquals(1, tombstones.size) assertEquals(Targets.KAFKA_SCHD, tombstones.single().target) assertTrue(tombstones.single().payloadJson.contains("\"deleted\":true")) - // 终态 + 回填意图由处理器在自己的事务内落库(message-lifecycle §2/§4) + // 终态和回填待办由处理器在自己的事务里写好 val row = proc.find(msgId)!! assertEquals(ProcStatus.SUCCEEDED, row.state) assertNotNull(row.backfillNextAt) @@ -89,7 +89,7 @@ class FdelAndAdftProcessorTest { val version = f.findMainRow("121")!!.stateVersion val result = proc.apply(head(), msg(), FlopPayload("121")) - assertEquals(ApplyResult.Succeeded, result) // §3.3:重复 FDEL 幂等成功 + assertEquals(ApplyResult.Succeeded, result) // 重复的删除报文算成功,但不再动数据 assertEquals(version, f.findMainRow("121")!!.stateVersion) assertEquals(1, events.rows.values.count { it.eventType == EventType.TOMBSTONE }) } @@ -102,7 +102,7 @@ class FdelAndAdftProcessorTest { val result = FdelProcessor(StubPipelineTx(), StubPipelineLock(), f, events, StubProcState(), ObjectMapper()) .apply(head(), msg(), FlopPayload("999")) - assertEquals(ApplyResult.Succeeded, result) // §3.3:迟到/不存在幂等成功 + assertEquals(ApplyResult.Succeeded, result) // 航班不存在(迟到或多余)也算成功 assertEquals(0, events.rows.size) } @@ -123,7 +123,7 @@ class FdelAndAdftProcessorTest { assertEquals(ApplyResult.Succeeded, result) val main = f.findMainRow("121")!! - assertEquals(FlightState.ACTIVE, main.state) // §3.3 生命周期重激活 + assertEquals(FlightState.ACTIVE, main.state) // 已删除的航班被重新激活 assertEquals(LocalDate.of(2026, 12, 15), main.operationDay) assertEquals("CA002", f.loadFullSnapshot("121")!!.scalars["FLNO"]) } @@ -145,7 +145,7 @@ class FdelAndAdftProcessorTest { assertEquals(ApplyResult.Succeeded, result) val main = f.findMainRow("555")!! assertEquals(FlightState.ACTIVE, main.state) - assertEquals(LocalDate.of(2026, 12, 20), main.operationDay) // §2.1:含 SODT 直接计算 + assertEquals(LocalDate.of(2026, 12, 20), main.operationDay) // 报文带了计划时间,运营日直接算出来 } @Test @@ -167,7 +167,7 @@ class FdelAndAdftProcessorTest { val scalars = f.loadFullSnapshot("555")!!.scalars assertEquals("XX200", scalars["FLNO"]) - assertEquals("keep-me", scalars["REMC"]) // §2.1 保守 Set-only:缺失不 Clear + assertEquals("keep-me", scalars["REMC"]) // 报文没带的字段保持原值,不清空 } @Suppress("unused") 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 d2d35aa..6425d71 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt @@ -106,7 +106,7 @@ class ScheduleProcessorTest { assertEquals(1, log.entries.size) assertEquals(SnapshotResult.COMMITTED, log.entries.single().result) assertEquals(1, log.entries.single().upserted) - // message-lifecycle §2/§4:终态与回填意图随业务写入在同一事务提交(不依赖主泵补写) + // 终态和回填待办跟业务数据一起提交,不靠主泵事后再补 val row = proc.find(msgId)!! assertEquals(ProcStatus.SUCCEEDED, row.state) assertNotNull(row.backfillNextAt) @@ -123,7 +123,7 @@ class ScheduleProcessorTest { val result = processor(proc = proc, flights = flights, log = log) .applyScheduleRecords(head(), message(makeBody("121" to "15DEC261723"))) - assertEquals(ApplyResult.ReplaySkipped, result) // §4 重放判定:幂等成功 + assertEquals(ApplyResult.ReplaySkipped, result) // 已经成功处理过,重放不重复写 assertNull(flights.findMainRow("121")) assertEquals(SnapshotResult.REPLAY_SKIPPED, log.entries.single().result) } @@ -135,7 +135,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)) // §4:声明数量不符整包拒绝 + assertTrue(result.flags.contains(SnapshotFlag.RECS_DROP)) // 声明条数和实收条数不符,整包拒绝 assertNull(flights.findMainRow("121")) assertEquals(SnapshotResult.ROLLED_BACK, log.entries.single().result) } @@ -143,7 +143,7 @@ class ScheduleProcessorTest { @Test fun `same flid across operation days is rejected whole-batch without writes`() { val flights = StubFlightState() - // 预置 121 已归属 12-15;报文声称 12-16 → §2.1 归属冲突 + // 库里 121 已经属于 12-15,报文却说是 12-16:运营日对不上 flights.persistFullState( com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot( "121", LocalDate.of(2026, 12, 15), FlightState.ACTIVE, 1, mapOf("SODT" to "15DEC261723"), emptyMap(), @@ -181,7 +181,7 @@ class ScheduleProcessorTest { ) assertEquals(ApplyResult.Succeeded, result) - assertEquals(FlightState.DELETED, flights.findMainRow("121")!!.state) // §3.3:日计划不恢复 + assertEquals(FlightState.DELETED, flights.findMainRow("121")!!.state) // 日计划不会让已删除的航班复活 assertEquals(8L, flights.findMainRow("121")!!.stateVersion) assertEquals(SnapshotResult.COMMITTED, log.entries.single().result) assertTrue(log.entries.single().flags.contains(SnapshotFlag.SCHD_REVIVE_CONFLICT)) @@ -195,7 +195,7 @@ class ScheduleProcessorTest { proc.insertIfAbsent(msgId, null) val p = processor(proc = proc, flights = flights, log = log, tx = TxRunner { true }) - // design.md §2.3:数据库/内部故障 → 异常上抛,MessageProcessor 记 FAILED(INFRA) 退避;不写任何终态 + // 数据库故障时异常直接上抛,由 MessageProcessor 记成可重试的失败;这里不留下任何终态 org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException::class.java) { p.applyScheduleRecords(head(), message(makeBody("121" to "15DEC261723"))) }