From d475feb7907ed2d288b08c0e718dc3f189877ec6 Mon Sep 17 00:00:00 2001 From: windyboy Date: Thu, 10 Sep 2026 11:05:55 +0800 Subject: [PATCH] =?UTF-8?q?docs(kdoc):=20=E9=87=8D=E5=86=99=E4=BF=A1?= =?UTF-8?q?=E7=AE=B1=E7=94=9F=E5=91=BD=E5=91=A8=E6=9C=9F=E7=9B=B8=E5=85=B3?= =?UTF-8?q?=20KDoc=EF=BC=8C=E6=94=B9=E4=B8=BA=E7=9B=B4=E7=99=BD=E8=AF=B4?= =?UTF-8?q?=E6=98=8E=E5=B9=B6=E5=8E=BB=E6=8E=89=E5=86=85=E9=83=A8=E7=BC=96?= =?UTF-8?q?=E5=8F=B7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 原有注释大量引用 §5.1/§5.2/Q2/US-01 这类文档编号和内部简称,跳过了"这段代码在做什么、 为什么这么做",没有读过设计文档的人基本读不懂。本次统一改成先讲清这件事本身、 再说为什么要这样做,编号只在末尾留一处指路。 覆盖本次改动涉及的 24 个 Kotlin 文件(生产 16 个 + 测试 8 个): - 领域与端口:ProcState(补齐全量字段说明与状态/错误分类逐项注释)、 ProcStateRepository / InboxCursorRepository / CminmsgInboxRepository 及 MailboxRow / BackfillDue / Backlog; - 收报:InboxPoller(把"水位连续、遇缺口停下、缺口老化"用大白话讲透)、 InboxService、JdbcCminmsgInboxRepository; - 处理:Pump / MessageProcessor、ScheduleProcessor、DynamicProcessors、 ProcFailure、JdbcProcStateRepository 与游标实现; - 回填与观测:BackfillService、InboxLifecycleHealthIndicator、JobRunner; - 配置与 stub:PipelineProps(三个新增参数说清取值理由)、MailboxProps、 StubRepositories; - 测试:8 个测试类改为"这些用例在守哪几条规矩",并保留 H2 不覆盖 ON CONFLICT 的说明。 术语统一按第一次出现就地解释:水位、处理标记、回填、死信、队头、终态。 纯注释改动;除拆分枚举时按仓库风格补的两个行尾逗号外无代码变更 (已用剥离注释后比对 HEAD 的方式逐文件核对)。测试仍为 78 passed / 1 skipped。 --- .../omms/msgexchange/config/MailboxProps.kt | 5 +- .../omms/msgexchange/config/PipelineProps.kt | 24 +++-- .../gzzn/omms/msgexchange/domain/ProcState.kt | 66 ++++++++++--- .../health/InboxLifecycleHealthIndicator.kt | 12 ++- .../infra/persistence/Repositories.kt | 93 +++++++++++++------ .../jdbc/JdbcCminmsgInboxRepository.kt | 16 ++-- .../persistence/jdbc/JdbcPgRepositories.kt | 12 +-- .../msgexchange/infra/retry/ProcFailure.kt | 12 ++- .../infra/stub/StubRepositories.kt | 20 ++-- .../omms/msgexchange/ingress/InboxPoller.kt | 35 ++++--- .../omms/msgexchange/ingress/InboxService.kt | 7 +- .../gzzn/omms/msgexchange/jobs/JobRunner.kt | 11 ++- .../msgexchange/processing/BackfillService.kt | 33 ++++--- .../processing/DynamicProcessors.kt | 24 +++-- .../gzzn/omms/msgexchange/processing/Pump.kt | 31 ++++--- .../processing/ScheduleProcessor.kt | 39 +++++--- .../omms/msgexchange/PipelineSmokeTest.kt | 15 +-- .../infra/health/HealthIndicatorsTest.kt | 6 +- .../persistence/jdbc/FlywayMigrationTest.kt | 11 +-- .../jdbc/InboxLifecycleJdbcSqlTest.kt | 14 +-- .../msgexchange/ingress/InboxPollerTest.kt | 11 ++- .../processing/BackfillServiceTest.kt | 9 +- .../processing/FdelAndAdftProcessorTest.kt | 6 +- .../processing/ScheduleProcessorTest.kt | 6 +- 24 files changed, 332 insertions(+), 186 deletions(-) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/MailboxProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/MailboxProps.kt index 2aa1184..4bf79b5 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/MailboxProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/MailboxProps.kt @@ -2,11 +2,12 @@ package com.gzzn.omms.msgexchange.config import io.micronaut.context.annotation.ConfigurationProperties +/** 共享信箱的连接信息与写入约定,对应配置里的 `mailbox.*`。 */ @ConfigurationProperties("mailbox") class MailboxProps { /** - * 处理标记写入值(message-lifecycle §5.2/Q7):仅限库方认可的 legacy 值集; - * 值集与写权限书面确认前保持 legacy 现役值。 + * 回填处理标记时写进 `CMINMSGS_STATUS` 的值。 + * 库里认可哪些取值还没和库方确认,暂时沿用 legacy 现役的 `PROCESSED`。 */ var processedValue: String = "PROCESSED" diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt index ec11ed1..bad45dc 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/config/PipelineProps.kt @@ -4,9 +4,10 @@ import io.micronaut.context.annotation.ConfigurationProperties import java.time.Duration /** - * ACMA-8 参数表(v4)初值;阶段 0 现网基线校准。 - * U03(N02):Micronaut 要求嵌套配置类同样标注 @ConfigurationProperties,否则 - * msgx.pipeline/schd/identity.* 全部静默回落 Kotlin 默认值。 + * 管道运行参数,对应配置文件里的 `msgx.*`。 + * + * 注意:嵌套的配置类也必须标 `@ConfigurationProperties`,否则 Micronaut 不会绑定 + * 里面的键,配置会静默失效、悄悄用回代码里的默认值。 */ @ConfigurationProperties("msgx") class PipelineProps { @@ -26,19 +27,24 @@ class PipelineProps { var headDeadline: Duration = Duration.ofMinutes(10) // 最坏 HOL 上界(毒丸升级) /** - * message-lifecycle §5.1 空洞老化:水位 W+1 处的空洞持续超过该时延即判定为永久并放行。 - * 取值口径 = 库方承诺的最大提交时延(Q2);太小会把迟到消息判成永久空洞(FIFO 越序风险)。 + * 缺口等待时长:水位后面缺了一个 ID 时,等这么久还没出现就认定它永远不会来了, + * 跳过缺口继续推进水位。 + * + * 取值应该等于库方承诺的"上游提交到消息可见的最长时间"。设太小,可能把一条 + * 迟到的消息误判成永久缺失,导致它排到后面的消息之后;设太大,收报会在缺口上白等。 */ var maxCommitDelay: Duration = Duration.ofMinutes(5) /** - * message-lifecycle §5.2 超期补写期限 R:终态后仍无处理标记的行到达该期限即强制补写, - * 保证库方清除前提「边界内无未标记行」在有限时间内成立。 - * R ≥ 人工重放期限 + 人工处置期限(Q6);确认前不得下调。 + * 超期补写期限:一条消息处理完之后,如果过了这么久还是没能把处理标记写回信箱 + * (比如回填一直失败),就直接强制补写一次,不再等退避。 + * + * 这是保证库方能清理信箱的兜底期限,必须覆盖人工重放所需的保留期, + * 确认之前不要调小,否则还在重放窗口内的消息会先被库方清掉。 */ var overdueBackfill: Duration = Duration.ofDays(30) - /** 回填扫描单批条数。 */ + /** 每次回填扫描最多处理多少条。 */ var backfillBatch: Int = 100 /** U07:启动即拉起 Pump/Dispatcher 循环(默认关——需要真实仓储或 msgx.stubs=true 才可安全开启)。 */ 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 f0ad5ba..93b8f64 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/domain/ProcState.kt @@ -3,32 +3,72 @@ package com.gzzn.omms.msgexchange.domain import java.time.Instant /** - * 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)。 + * 一条入站消息的处理记录,一行对应共享信箱 CMINMSGS 里的一条报文。 * - * 回填事实(message-lifecycle.md §5.2/§11)与处理事实同体同行:终态与回填意图 - * 由同一条 UPDATE 落下,因此不存在"业务已提交、回填待办未记"的崩溃窗口。 + * 这张表同时承担四件事: + * 1. 防重复入队——主键是信箱 ID,同一条报文只会有一行; + * 2. 防业务重复——IDENTITY_KEY 唯一,同一条业务报文只处理一次; + * 3. 记录处理进度——状态、重试次数、下次重试时间和失败原因; + * 4. 记录回填进度——处理完要把"已处理"标记写回共享信箱,写成功之前一直留着待办意图。 + * + * 第 4 条和终态写在同一条 UPDATE 里,所以不会出现"业务处理完了,却没人记得去回填"。 */ -enum class ProcStatus { PENDING, FAILED, SUCCEEDED, SKIPPED, DEAD } +enum class ProcStatus { + /** 已入队,等待处理。 */ + PENDING, -/** 错误分类(design.md §2.3):MALFORMED/PROTOCOL 直接 DEAD 不重试;其余退避重试。 */ -enum class ErrorClass { MALFORMED, PROTOCOL, CODEC_ERROR, EXHAUSTED, INFRA, UNSUPPORTED } + /** 处理失败,等退避时间到了再重试。 */ + FAILED, + + /** 处理成功。 */ + SUCCEEDED, + + /** 判定为业务重复,跳过不处理。 */ + SKIPPED, + + /** 处理失败且不再重试,等人工处置。 */ + DEAD, +} + +/** 失败原因分类,决定失败后是重试还是直接进死信。 */ +enum class ErrorClass { + /** 报文本身不合法,重试也没用,直接进死信。 */ + MALFORMED, + + /** 整包被拒绝(运营日冲突、声明条数不符等),整包不落地,直接进死信。 */ + PROTOCOL, + + /** 解码逻辑的问题;修好 codec 之后可以重放。 */ + CODEC_ERROR, + + /** 重试次数用尽或队头滞留超时;人工复核后可以重放。 */ + EXHAUSTED, + + /** 数据库、网络等基础设施抖动,重试通常就能过。 */ + INFRA, + + /** 报文类型还没有对应处理器;属于能力未实现,先退避重试等补齐。 */ + UNSUPPORTED, +} data class ProcState( + /** 信箱 CMINMSGS_ID,也是本表主键。 */ val msgId: Long, val state: ProcStatus, - val identityKey: String? = null, // SNDR|TYPE|STYP|SEQN(design.md §2.2);decode 后首次绑定,FAILED 重试不重绑 + /** 业务身份 SNDR|TYPE|STYP|SEQN;解码成功后绑定一次,重试不会重绑。 */ + val identityKey: String? = null, + /** 处理失败次数,用来算退避档位和判断是否已到上限。 */ val attempts: Int = 0, + /** FAILED 状态下,下次可以重试的时刻。 */ val nextAttemptAt: Instant? = null, val errorClass: ErrorClass? = null, + /** 最近一次失败的原因(截断后落库,供排查)。 */ val lastError: String? = null, - /** 信箱 CMINMSGS_DATE_RECEIVED:§5.2 超期补写的 R 判据与 OPS-2「最老未处理信龄」锚点。 */ + /** 信箱里的接收时间:用来判断"超期仍未回填",也是最老未处理信龄的计算依据。 */ val receivedAt: Instant? = null, - /** 非空 = 已确认信箱行持有处理标记(回填完成)。 */ + /** 非空表示已确认信箱行带上了处理标记。 */ val backfillAt: Instant? = null, - /** 非空 = 待回填;终态事务内登记为 now,失败按退避推后。 */ + /** 非空表示还欠一次回填:写终态时置为当前时间,失败后退避推后。 */ val backfillNextAt: Instant? = null, val backfillAttempts: Int = 0, val backfillError: String? = null, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/health/InboxLifecycleHealthIndicator.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/health/InboxLifecycleHealthIndicator.kt index bf499cc..01b07bb 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/health/InboxLifecycleHealthIndicator.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/health/InboxLifecycleHealthIndicator.kt @@ -14,11 +14,14 @@ import java.time.Duration import java.time.Instant /** - * 信箱生命周期观测(docs/message-lifecycle.md §5.3 验收 / user-stories.md OPS-2): - * 输出剩余积压、最老未处理信龄、未回填终态条数与水位滞后,供积压消化期间持续观察。 + * 在 /health 里输出收报与回填的当前情况,用来观察积压消化得怎么样: + * - `backlog`:还没处理完的消息条数; + * - `oldestUnprocessedSeconds`:最老一条未处理消息从收到到现在过了多久; + * - `unmarkedTerminal`:已经处理完、但还没把标记写回信箱的条数(回填跟不上时这个数会涨); + * - `watermark` / `watermarkLag`:收报读到哪个 ID 了、落后信箱最新 ID 多少。 * - * 端口缺省(未接通共享信箱或自有 PG)时报告未绑定而不判 DOWN——可用性由各依赖自身的 - * 健康指示器承担,本指示器只反映生命周期状态;端口查询失败判 DOWN。 + * 相关工作没接上时(比如没连共享信箱)只提示"未绑定",不判 DOWN——依赖本身是否可用 + * 由各自的健康指示器回答,这里只报告业务状态。只有查询出错才判 DOWN。 */ @Singleton class InboxLifecycleHealthIndicator( @@ -39,6 +42,7 @@ class InboxLifecycleHealthIndicator( ) } +/** 组装上面那几个指标;单独抽出来是为了能在测试里直接调用。 */ internal fun lifecycleHealth( procState: ProcStateRepository?, cursor: InboxCursorRepository?, 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 118511f..78b2b71 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 @@ -30,30 +30,36 @@ interface PipelineLockRepository { } /** - * PROC_STATE:每消息一行(design.md §2.1);SUCCEEDED 终态兼作日计划重放判定。 - * 回填事实(BACKFILL_AT/NEXT_AT/ATTEMPTS/ERROR)与处理事实同行,取代独立的回填待办表。 + * PROC_STATE 表的读写入口:一条消息一行的处理记录。 + * + * 状态分两类:PENDING / FAILED 是还在处理中,SUCCEEDED / SKIPPED / DEAD 是终态。 + * 只有终态才允许往信箱回填处理标记。回填进度就记在同一行上,不需要另一张待办表。 */ interface ProcStateRepository { /** - * 入队:MSG_ID 主键幂等(重复扫描与 compat 入口并发都不会重复建行)。 - * @param receivedAt 信箱 DATE_RECEIVED,用于 §5.2 超期判据与信龄观测。 - * @return true = 本次实际新建 + * 入队:把信箱里发现的消息登记成 PENDING。 + * + * 幂等:同一个消息 ID 重复登记既不报错、也不会建第二行(收报重扫和兼容入口并发调用都安全)。 + * + * @param receivedAt 信箱里的接收时间,用于判断超期未回填和统计最老信龄 + * @return true 表示这次真的新建了一行 */ fun insertIfAbsent(msgId: Long, receivedAt: Instant?): Boolean fun find(msgId: Long): ProcState? + /** 这条消息是否已经处理成功过(判断重放用,避免把同一份快照重复应用)。 */ fun findSuccessTerminal(msgId: Long): Boolean - /** 严格 FIFO 队头(最小未完成 MSG_ID)。 */ + /** 当前队头:还没处理完的消息里 ID 最小的那条。 */ fun headUnfinished(): ProcState? - /** identity 首次绑定;返回 false = 另一条消息已持有该键。 */ + /** 给消息绑定业务身份;返回 false 表示这个身份已经被另一条消息占了(业务重复)。 */ fun tryBindIdentity(msgId: Long, identityKey: String): Boolean fun ownerOfIdentity(identityKey: String): Long? - /** 非终态迁移(PENDING / FAILED 及退避),不触碰回填列。 */ + /** 改写 PENDING / FAILED 这类非终态(含重试次数与下次重试时间),不动回填字段。 */ fun update( msgId: Long, state: ProcStatus, @@ -64,8 +70,11 @@ interface ProcStateRepository { ) /** - * 终态 + 回填意图同一条 UPDATE(message-lifecycle §2/§4),由处理器在自己的业务 - * 事务内调用:航班变更、事件、终态、回填意图同提交同回滚。 + * 写终态,同时记下"这条消息还欠一次回填"。 + * + * 两件事是同一条 UPDATE,必须由处理器在自己的业务事务里调用:业务改动、待发事件、 + * 终态和回填意图一起提交或一起回滚。这样即使进程恰好在这里挂掉,也不会留下 + * "业务已经改了、却没人记得回填信箱"的记录。 */ fun markTerminal( msgId: Long, @@ -76,29 +85,36 @@ interface ProcStateRepository { now: Instant = Instant.now(), ) - /** 信箱行已确认持有处理标记。 */ + /** 回填成功:记下完成时间,清掉待办。 */ fun markBackfilled(msgId: Long, now: Instant = Instant.now()) - /** 回填失败:次数 +1、按退避推后、留错误;终态不得回改(§11)。 */ + /** 回填失败:次数 +1、按退避推后、记下原因。处理终态不受影响,不会被改回去。 */ fun recordBackfillFailure(msgId: Long, error: String?, attempts: Int, nextAttemptAt: Instant, now: Instant) /** - * 待回填待办(message-lifecycle §3/§5.2):终态 + 未确认标记, - * 且(已到期 或 接收时间已达超期期限 overdueBefore)。 + * 找出现在该回填的记录:已经到终态、还没确认回填,并且退避时间已到。 + * + * [overdueBefore] 是兜底:消息接收时间早于它的(已经等了很久)无视退避直接补写。 + * 没有这条兜底,退避一直失败的话这些行就永远打不上标记,库方也没法清理信箱。 */ fun findBackfillDue(now: Instant, overdueBefore: Instant, limit: Int): List - /** 显式重放入口:仅把给定 errorClass 集合中的行从 FAILED/DEAD 置回 PENDING(ATTEMPTS=0)。 */ + /** 人工重放:把指定错误类别的 FAILED / DEAD 记录改回 PENDING,重试次数清零。 */ fun requeueByErrorClasses(errorClasses: List): Int - /** OPS-2 观测口径(message-lifecycle §5.3 验收):积压、最老信龄锚点、未回填终态。 */ + /** 积压观测:还没处理完的条数、最老一条的接收时间、处理完但还没回填的条数。 */ fun backlog(): Backlog } -/** 某条待回填记录(扫描输入)。 */ +/** 扫描到的待回填记录。 */ data class BackfillDue(val msgId: Long, val attempts: Int) -/** 处理侧积压快照。 */ +/** + * 处理侧积压快照。 + * @param unfinished 还没处理完的消息条数 + * @param oldestReceivedAt 其中最早一条的接收时间(据此算信龄) + * @param unmarkedTerminal 已经处理完、但还没把标记写回信箱的条数 + */ data class Backlog(val unfinished: Int, val oldestReceivedAt: Instant?, val unmarkedTerminal: Int) /** MSG_EVENT outbox(flight-state.md §5)。KAFKA_SCHD 合并同 FLID 未发事件按最新 STATE_VERSION 输出。 */ @@ -208,13 +224,20 @@ interface ReqTrackRepository { } /** - * 消费水位 W(message-lifecycle §5.1):单行游标,只随新 ID 成功入队推进、遇空洞即停, - * 与入队在同一 PG 事务提交——中断后 W 未前进,重扫即补建(§4 第一行)。 + * 收报进度(水位 W):记下"信箱里到哪个 ID 为止已经全部读进自有库",全表只有一行。 * - * `holeSince` 记录 W+1 处空洞首次被观测到的时刻:超过最大提交时延(Q2 承诺)即判定为 - * 永久空洞并放行,否则水位会永久停摆于一次自增回滚留下的空位,后续 ID 再无入队机会。 + * 光记一个数字不够,还要记住缺口是什么时候出现的:如果 W 后面缺了一个 ID,就先停在 + * 缺口前面等(可能是上游还没提交完,随时会补上)。等的时间超过最大提交时延,就改判为 + * 永久缺失、跳过去继续推进——否则一次自增回滚留下的空位就能让水位永远卡住, + * 它后面的消息再也进不了队。 + * + * 水位推进与入队在同一个 PG 事务里提交:中途崩溃时水位没动,重启后重扫一遍即可补齐。 */ interface InboxCursorRepository { + /** + * @param committedUpTo 水位 W + * @param holeSince W 后面那个缺口最早被发现的时刻;当前没有缺口时为 null + */ data class Cursor(val committedUpTo: Long = 0L, val holeSince: Instant? = null) fun load(): Cursor @@ -223,24 +246,36 @@ interface InboxCursorRepository { } /** - * 共享 MySQL 信箱 CMINMSGS 访问(他人系统库,本系统不建表,只做 DML)。 - * 三个事实互不替代(message-lifecycle §5.1/§11):**发现**按 ID 区间读、**水位**只表示 - * 读取进度、**处理标记**只用于回填与库方清除——标记不得作为扫描谓词。 + * 共享 MySQL 信箱 CMINMSGS 的读写入口。这个库是别人的,本系统只做约定的读写, + * 不建表、不改结构。 + * + * 这里把三件事分得很清楚,谁也不代替谁: + * - **发现**:按 ID 区间读有哪些新消息([readRange]); + * - **进度**:读到哪儿了记在自有库的水位里(见 [InboxCursorRepository]); + * - **标记**:处理完了把"已处理"写回信箱([markProcessedIfUnmarked])。 + * + * 特别是发现,不能拿"有没有处理标记"当筛选条件:处理完但还没回填的行,以及永远不会 + * 回填的死信,会一直占着每一批的名额,攒够一批之后新消息就再也读不到了。 */ interface CminmsgInboxRepository { + /** 兼容入口往信箱写一条报文,返回新的信箱 ID。 */ fun insertRaw(rawXml: String): Long + /** 读某条消息的原文;返回 null 表示读不到(行已被清除,或原文本身为空)。 */ fun rawOf(msgId: Long): String? - /** 按 ID 区间升序有界读取(`ID > fromExclusive`),不以处理标记为谓词。 */ + /** 按 ID 升序读一批 `ID > fromExclusive` 的行,不带别的过滤条件。 */ fun readRange(fromExclusive: Long, limit: Int): List - /** 信箱当前最大 ID(空表 null);仅用于观测水位滞后。 */ + /** 信箱当前最大 ID,空表返回 null;只用来观测收报落后了多少。 */ fun maxId(): Long? - /** 只把空标写为已处理(§11 单调);返回 true = 本次实际写入,重复执行无副作用。 */ + /** + * 把处理标记写回信箱,并且**只写还是空标记的行**:库里已有值时不覆盖、不回退, + * 重复调用没有副作用。返回 true 表示这次真的写进去了。 + */ fun markProcessedIfUnmarked(msgId: Long, value: String): Boolean } -/** 信箱行读取结果(发现阶段只需要身份与接收时间,原文按需再取)。 */ +/** 从信箱读到的一行:ID 加接收时间。原文按需再取,扫描时不读大字段。 */ data class MailboxRow(val msgId: Long, val receivedAt: Instant?) diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt index 27b81e8..00d15ee 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcCminmsgInboxRepository.kt @@ -8,11 +8,13 @@ import jakarta.inject.Singleton import javax.sql.DataSource /** - * 共享 MySQL CMINMSGS 信箱适配(仅 DML,不建表;message-lifecycle §5.1/§6/§11)。 - * 列名与 legacy `entity/Cminmsg.java` 一致。 + * 共享 MySQL CMINMSGS 信箱的 JDBC 实现。列名沿用 legacy 的 `CMINMSGS_*`。 + * 只做增删改查,不建表、不改结构——这个库属于别的系统。 * - * 发现按 ID 区间读(`ID > ?`),**不以 `DATE_PROCESSED` 为扫描谓词**:已入队但尚未 - * 回填的行否则会永久占据批次,正是 §5.1 与 US-01 条目 3 要求排除的场景。 + * 两处刻意为之: + * - 取新消息只看 ID(`ID > ?`),不看 `DATE_PROCESSED`。用处理标记当条件的话, + * 处理完但还没回填的行会长期占住每批名额,死信攒够一批就再也发现不了新消息。 + * - 写回处理标记带 `DATE_PROCESSED IS NULL` 条件,一行只会被标记一次,不会覆盖已有值。 */ @Singleton @Requires(property = "msgx.stubs", notEquals = "true") @@ -62,9 +64,9 @@ class JdbcCminmsgInboxRepository( } /** - * §11 单调:`DATE_PROCESSED IS NULL` 守卫保证只把空标写为已处理,已有值不回撤、 - * 不覆盖;影响 0 行 = 已被其他路径标记,调用方按幂等成功处理。 - * 写入值(DATE_PROCESSED 时间语义与 STATUS 值集)以库方契约为准(Q7)。 + * 只更新还是空标记的行,所以重复调用不会覆盖库里已有的值; + * 影响 0 行说明已经被标记过了,调用方按"成功"处理即可。 + * 具体写什么值、时间怎么解释,以与库方约定为准。 */ override fun markProcessedIfUnmarked(msgId: Long, value: String): Boolean = ds.update( 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 e9672c9..4d8fc11 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 @@ -70,7 +70,7 @@ class JdbcPipelineLockRepository( class JdbcProcStateRepository( private val ds: DataSource, ) : ProcStateRepository { - /** MSG_ID 主键幂等入队:重复扫描与 compat 入口并发都不重复建行(§5.1)。 */ + /** 入队(幂等):主键冲突时什么都不做,所以重复扫描和兼容入口并发调用都安全。 */ override fun insertIfAbsent(msgId: Long, receivedAt: Instant?): Boolean = ds.update( "INSERT INTO proc_state (msg_id, state, received_at, updated_at) VALUES (?, 'PENDING', ?, ?) " + @@ -144,7 +144,7 @@ class JdbcProcStateRepository( ) } - /** 终态与回填意图同一条 UPDATE:业务事务内调用即原子提交(message-lifecycle §2/§4)。 */ + /** 写终态并同时登记回填待办,一条 SQL 搞定;由处理器在业务事务里调用。 */ override fun markTerminal( msgId: Long, state: ProcStatus, @@ -199,8 +199,8 @@ class JdbcProcStateRepository( } /** - * 待回填:终态 + 未确认标记,且已到期或已达 §5.2 超期期限(R 覆盖退避,保证 - * 有限时间内必然补写,否则库方清除的前提"边界内无未标记行"无法成立)。 + * 到了该回填的时候:终态 + 还没有标记 + (退避到期 或 收信时间已经很久)。 + * 后面这个"很久"是兜底,保证标记最终一定会补上,库方才能按标记清理信箱。 */ override fun findBackfillDue(now: Instant, overdueBefore: Instant, limit: Int): List = ds.query( @@ -231,7 +231,7 @@ class JdbcProcStateRepository( ) } - /** OPS-2(message-lifecycle §5.3):积压条数、最老未处理接收时刻、未回填终态条数。 */ + /** 积压观测:未处理条数、最老一条的接收时间、处理完但未回填的条数。 */ override fun backlog(): Backlog = ds.queryOne( """ @@ -273,7 +273,7 @@ class JdbcProcStateRepository( } } -/** 消费水位单行游标(message-lifecycle §5.1);与入队同事务写入。 */ +/** 收报水位游标(单行)的 JDBC 实现;水位推进与入队在同一个事务里提交。 */ @Singleton @Requires(property = "datasources.default.enabled", value = "true") @Requires(missingProperty = "msgx.stubs") diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/retry/ProcFailure.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/retry/ProcFailure.kt index 17642f0..c2e093f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/infra/retry/ProcFailure.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/infra/retry/ProcFailure.kt @@ -7,16 +7,20 @@ import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository import jakarta.inject.Singleton /** - * ProcState 侧统一失败迁移(U08/U10):处理/快照路径共用—— - * attempts+1 后若 exhausted → DEAD(EXHAUSTED)(终态 + 回填意图,errorClass 规范化,原因保留在 lastError); - * 否则 FAILED + attempts + nextAttemptAt(退避)(可重放)。任何“失败”都不得在无退避下直接终态化。 + * 处理失败时统一改状态,主泵和快照路径共用。 + * + * 规则很简单:失败次数 +1 之后 + * - 还没到上限:改成 FAILED,并按退避表算好下次重试时间(这种记录以后可以重放); + * - 已经到上限:改成 DEAD(EXHAUSTED) 终态,等人工复核,不再自动重试。 + * + * 失败一定先退避、再重试,不允许一次失败就直接判死。 */ @Singleton class ProcFailure( private val procState: ProcStateRepository, val scheduler: FailureScheduler, ) { - /** @return 是否已达终态(DEAD 才是终态;FAILED 仍可重放/重试) */ + /** @return 是否已经落到终态:DEAD 是终态,FAILED 还会再试 */ fun fail(head: ProcState, ec: ErrorClass, reason: String): Boolean { val attempts = head.attempts + 1 if (scheduler.exhausted(attempts)) { 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 ec46669..21a6d6e 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 @@ -33,8 +33,10 @@ import java.time.ZoneId import java.util.concurrent.atomic.AtomicLong /** - * stub 仓储(msgx.stubs=true 时装配)——内存实现全部契约,测试与无库环境用。 - * 事务管理器直接执行 block(无嵌套语义);锁为 no-op(单线程测试前提)。 + * 内存版仓储,`msgx.stubs=true` 时装配,供测试和无外部依赖的环境使用。 + * 行为对齐 JDBC 实现(例如入队幂等、写终态时同时登记回填待办)。 + * + * 事务管理器直接执行代码块、不做回滚;锁是空操作——都建立在单线程测试的前提上。 */ @Singleton @Requires(property = "msgx.stubs", value = "true") @@ -50,6 +52,7 @@ class StubPipelineLock : PipelineLockRepository { @Singleton @Requires(property = "msgx.stubs", value = "true") +/** 内存版 PROC_STATE:入队幂等,写终态时一并写下回填待办。 */ class StubProcState : ProcStateRepository { val rows = linkedMapOf() val bound = linkedMapOf() @@ -100,7 +103,7 @@ class StubProcState : ProcStateRepository { ) } - /** 终态 + 回填意图同一动作(模拟单条 UPDATE 的原子性)。 */ + /** 写终态并按同一次操作登记回填待办,模拟 JDBC 实现里"一条 SQL 写两件事"的原子性。 */ override fun markTerminal( msgId: Long, state: ProcStatus, @@ -366,6 +369,7 @@ class StubReqTrack : ReqTrackRepository { @Singleton @Requires(property = "msgx.stubs", value = "true") +/** 内存版共享信箱:可以模拟上游写入、库方清除,并记录哪些行被打上了处理标记。 */ class StubInbox : CminmsgInboxRepository { val raws = linkedMapOf() private val received = linkedMapOf() @@ -390,7 +394,7 @@ class StubInbox : CminmsgInboxRepository { override fun maxId(): Long? = raws.keys.maxOrNull() - /** §11 单调:只把空标写为已处理;已有值不回撤、不覆盖。 */ + /** 只写还没有标记的行;已经标记过就返回 false,不覆盖已有值。 */ override fun markProcessedIfUnmarked(msgId: Long, value: String): Boolean { if (!raws.containsKey(msgId) || marks.containsKey(msgId)) return false marks[msgId] = value @@ -401,16 +405,16 @@ class StubInbox : CminmsgInboxRepository { fun isMarked(msgId: Long): Boolean = marks.containsKey(msgId) - /** 测试辅助:模拟上游外部写入共享信箱(不经本系统)。 */ + /** 测试辅助:模拟上游直接往信箱写报文(不经过本系统)。 */ fun simulateExternalWrite(rawXml: String): Long = insertRaw(rawXml) - /** 测试辅助:模拟库方清除(原文不可读),用于 §9 原文缺失与空洞场景。 */ + /** 测试辅助:模拟库方清除这条行,用来构造"原文读不到"和 ID 缺口。 */ fun removeRow(msgId: Long) { raws.remove(msgId); received.remove(msgId); marks.remove(msgId) } } -/** 消费水位游标(stub)。 */ +/** 内存版收报水位游标。 */ @Singleton @Requires(property = "msgx.stubs", value = "true") class StubInboxCursor : InboxCursorRepository { @@ -427,6 +431,6 @@ class StubInboxCursor : InboxCursorRepository { } } -/** 终态判定(§2 状态总纲)。 */ +/** 判断是不是终态:处理已经结束、不会再重试的状态。 */ private fun ProcStatus.isTerminal(): Boolean = this == ProcStatus.SUCCEEDED || this == ProcStatus.SKIPPED || this == ProcStatus.DEAD diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPoller.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPoller.kt index a41b9b8..b4ea41f 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPoller.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPoller.kt @@ -11,18 +11,22 @@ import java.time.Duration import java.time.Instant /** - * 收报(docs/message-lifecycle.md §5.1):共享 MySQL 信箱按 **ID 区间**升序有界读取, - * 在自有 PG 建 `PROC_STATE(PENDING)`,并把消费水位推进到连续上界;入队与水位推进在同一 - * PG 事务内提交,中断即重扫补建(§4 第一行)。收报层不解析业务载荷、不写处理标记。 + * 收报:轮询共享信箱,把新消息登记到自有 PG 的 PROC_STATE,等着主泵处理。 * - * 三条红线: - * - **扫描谓词不含处理标记**:标记只用于回填与库方清除,不参与消息发现;否则已入队而未 - * 回填的行会永久占据批次,死信累积到批大小时收报整体停摆; - * - **水位遇空洞即停**:不越过空洞入队(越过后较小 ID 迟到即 FIFO 越序,architecture §5); - * - **空洞老化**:超过最大提交时延(Q2 承诺)的空洞判定为永久并放行,否则水位会永久停摆于 - * 一次自增回滚留下的空位。 + * 每一轮只做三件事: + * 1. 从水位 W 之后按 ID 升序读一批行(只看 ID,不看处理标记); + * 2. 把读到的行登记成 PENDING,并把水位推到"连续"的位置; + * 3. 登记和水位推进写在同一个事务里——中途崩溃时水位没动,下一轮重扫即可补齐。 * - * 空转代价为每轮一次区间 SELECT。 + * "连续"是这里唯一需要理解的规则。如果 W 后面缺了一个 ID,说明可能有 ID 更小的消息 + * 还没提交上来。此时先停在缺口前,不把缺口后面的消息放进队列:否则那条迟到的消息 + * 会排到它们后面,破坏"先来先处理"的约定,同一航班的报文可能被乱序应用。 + * + * 缺口等超过 `msgx.pipeline.max-commit-delay` 仍未出现,就认定它永远不会来了 + * (典型情况是自增回滚留下的空位),跳过它继续推进——否则水位会卡在第一个空位上 + * 再也不动。 + * + * 这一层不解析报文,也不写信箱处理标记:标记由处理完成后的 BackfillService 负责补。 */ @Singleton class InboxPoller( @@ -37,14 +41,14 @@ class InboxPoller( @Volatile private var running = false - /** @return 本轮新建的入队条数 */ + /** @return 本轮新登记的消息条数(已登记过的行不计入,也不影响水位推进)。 */ fun pollOnce(now: Instant = Instant.now()): Int { val batch = props.pipeline.claimBatch.coerceAtLeast(1) val watermark = cursor.load() val rows = mailbox.readRange(watermark.committedUpTo, batch) if (rows.isEmpty()) return 0 - // 连续上界;读取区间内出现空洞时,只推进到连续部分,空洞之后的行暂不入队(防较小 ID 迟到被越过) + // 找到连续部分的末尾;如果这批里出现了缺口,缺口后面的行这一轮先不入队 val contiguous = contiguousUpTo(watermark.committedUpTo, rows) ?: watermark.committedUpTo var committedTo = contiguous var holeSince: Instant? = null @@ -53,7 +57,7 @@ class InboxPoller( if (Duration.between(since, now) < props.pipeline.maxCommitDelay) { holeSince = since } else { - // 空洞老化:超过最大提交时延仍缺席即判永久(Q2),放行水位,否则永久停摆 + // 缺口等太久了:当成永久缺失跳过,让水位继续往前走 committedTo = rows.first { it.msgId > contiguous }.msgId - 1 log.warn("hole after W={} aged out, watermark advanced to {}", contiguous, committedTo) } @@ -95,7 +99,10 @@ class InboxPoller( running = false } - /** 连续上界:从 W+1 起 ID 逐 1 相邻的最后一个;首个空位即停(rows 为升序且覆盖该区间)。 */ + /** + * 从 W+1 开始数,返回 ID 逐 1 相连的最后一个 ID;遇到第一个缺号就停。 + * 返回 null 表示 W+1 本身就不存在。 + */ private fun contiguousUpTo(from: Long, rows: List): Long? { var expected = from + 1 var last: Long? = null diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxService.kt index c033e95..cfdf8b8 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/ingress/InboxService.kt @@ -6,8 +6,11 @@ import jakarta.inject.Singleton import java.time.Instant /** - * ACMA-8 流程 1 · compat 写路径:HTTP 落信 + PG 入队(非生产主拓扑;主路径=InboxPoller 轮询)。 - * 两步不在同一事务:信箱成功而 PG 失败时原文不丢失,由轮询按 ID 区间补建(design.md §3.1)。 + * 兼容 HTTP 入口:把一条报文写进共享信箱,再在自有 PG 里入队,供联调和影子对拍使用。 + * 生产的主路径是 InboxPoller 轮询,不是这里。 + * + * 两步不在同一个事务里(跨库没有事务)。如果信箱写成功、PG 入队失败,原文仍然在信箱里, + * 收报轮询会按 ID 把它补进来,所以不会丢消息。 */ @Singleton class InboxService( 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 da71eef..52858ce 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/JobRunner.kt @@ -10,10 +10,13 @@ import java.time.LocalDate import java.time.ZoneId /** - * 维护作业调度(docs/design.md §6.1/§6.2):单 daemon 线程,独立于主泵—— - * 作业不参与消息 FIFO,也不使到期消息饥饿(PUMP_JOB 队列机制已随审计口径移除)。 - * 触发:回填补写扫描 30s 固定间隔;历史归档/留痕清理每日机场时区 03:30 后首个 tick。 - * 回填意图由处理器在终态事务内登记(message-lifecycle §4),本线程只负责到期重试。 + * 维护作业线程:定时做那些不需要跟消息一起排队的事。 + * + * - 每 30 秒扫一次还欠回填的记录,把处理标记补写回共享信箱; + * - 每天机场时间 3:30 之后跑一次航班历史归档与留痕清理。 + * + * 作业跑在自己的线程上,不占用消息处理循环,也不会让到期的消息饿死在这里。 + * 回填意图在处理完成时就已经写进数据库了,本线程只负责到期重试。 */ @Singleton class JobRunner( diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/BackfillService.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/BackfillService.kt index 5f4587c..4ba121b 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/BackfillService.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/BackfillService.kt @@ -10,15 +10,20 @@ import java.time.Duration import java.time.Instant /** - * 信箱处理标记回填(docs/message-lifecycle.md §3/§4/§5.2)。 + * 把"已处理"标记写回共享信箱。 * - * 回填意图(`BACKFILL_NEXT_AT`)由处理器在终态事务内登记,与业务写入同提交同回滚, - * 因此不存在"业务已提交、待办未记"的窗口;本服务只做两件事: - * 1. [attempt]:终态提交后立即尝试一次(低延迟,失败静默留给扫描); - * 2. [sweep]:到期或已达超期期限 R 的记录批量补写(§5.2),指数退避 30s 起步、封顶 15 分钟。 + * 一条消息处理完要做两件事:记下终态、把标记写回信箱。第一件在业务事务里完成, + * 同时留下"还欠一次回填"的意图(PROC_STATE 上的 BACKFILL_NEXT_AT); + * 这个类负责第二件: * - * 回填失败绝不重放业务变更,也绝不回改终态(§11);写入侧只把空标写为已处理, - * 重复执行无副作用。 + * - [attempt]:处理刚结束时马上试一次,让标记尽快落到信箱。失败也不影响处理结果, + * 留给扫描重试即可。 + * - [sweep]:定时把还欠回填的记录挑出来重试。失败就按 30 秒起步、最长 15 分钟的 + * 退避往后推;如果一条消息从收到现在已经超过超期期限,则无视退避强制补写—— + * 否则退避可能一直失败下去,这些行永远打不上标记,库方就没法清理信箱。 + * + * 两条底线:回填失败不会把终态改回去,也不会重新执行业务逻辑;写标记只写还是空标记的 + * 行,重复执行没有副作用。 */ @Singleton class BackfillService( @@ -40,7 +45,10 @@ class BackfillService( } } - /** 单条最佳努力回填;失败只登记退避(异常不外抛,不阻塞提交后的处理路径)。 */ + /** + * 处理完立刻试一次。失败只记一笔退避信息就返回,不抛异常—— + * 调用方是主泵的处理路径,不能被回填问题拖住。 + */ fun attempt(msgId: Long, now: Instant = clock.instant()) { record(msgId, attempts = 0, now = now)?.let { log.warn("backfill failed msgId={} error={} (sweep will retry)", msgId, it) @@ -48,8 +56,9 @@ class BackfillService( } /** - * 批量补写(JobRunner 每 30s 触发;重启即继续,不依赖内存状态)。 - * @return 本批检查条数 + * 批量补写,由 JobRunner 每 30 秒调用一次。待办状态都在数据库里, + * 进程重启后接着跑,不需要额外恢复步骤。 + * @return 本批处理的条数 */ fun sweep(now: Instant = clock.instant()): Int { val due = procState.findBackfillDue(now, now.minus(props.pipeline.overdueBackfill), props.pipeline.backfillBatch) @@ -57,10 +66,10 @@ class BackfillService( return due.size } - /** @return 失败原因;null = 已确认标记(含"已被其他路径标记"的幂等成功) */ + /** 回填一条。@return 失败原因;返回 null 表示标记已确认(包括"别的路径已经标过了")。 */ private fun record(msgId: Long, attempts: Int, now: Instant): String? = try { - // 影响 0 行 = 已有标记;按幂等成功处理(§11 标记单调:不回撤、不覆盖) + // 没有真正写进去说明库里已经有标记了,同样算成功(不覆盖已有值) mailbox.markProcessedIfUnmarked(msgId, mailboxProps.processedValue) procState.markBackfilled(msgId, now) null 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 768d2fa..1ed87bc 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt @@ -25,9 +25,11 @@ import java.time.Instant import java.time.ZoneId /** - * FLOP(docs/flight-state.md §3.2 动态运行事件):读取完整当前态 → 合并变化 → - * 保留运营日 → STATE_VERSION+1 → 同事务登记 KAFKA_MSG / KAFKA_SCHD 与处理终态。 - * 未知/迟到航班按幂等成功处理,不创建实例(创建入口只有 SCHD/ADFT)。 + * 处理 FLOP 动态运行事件:读出航班当前态,把报文里的变化合并进去,版本号加一, + * 并登记状态事件。 + * + * 航班不存在时按"迟到的消息"处理——直接算处理成功,不会顺手建一个新航班 + * (新建只发生在 SCHD 和 ADFT 里)。 */ @Singleton class FlopProcessor( @@ -57,8 +59,10 @@ class FlopProcessor( } /** - * FDEL(docs/flight-state.md §3.3 删除):ACTIVE → 置 DELETED、推进版本、明细保留、 - * 与删除同事务登记 tombstone(§5);已 DELETED / 不存在 → 幂等成功,不推进版本、不重复发布。 + * 处理 FDEL 删除报文:把在用的航班标记为 DELETED,版本号加一,明细数据保留, + * 并且只发布一次删除事件。 + * + * 已经删除过、或航班本来就不存在时,算处理成功但不再动版本、不重复发事件。 */ @Singleton class FdelProcessor( @@ -107,9 +111,11 @@ class FdelProcessor( } /** - * ADFT(docs/flight-state.md §3.3):字段缺失语义待上游确认——确认前按保守 Set-only - * 处理(出现字段覆盖、缺失不清空)。FLID 已存在且 DELETED → 生命周期重激活;不存在 → - * 新实例建立(含 SODT 时直接计算 OPERATION_DAY,§2.1;否则保留 NULL 待日计划收录)。 + * 处理 ADFT 异常航班报文:已删除的航班重新激活,不存在的航班新建。 + * + * 字段按"出现才覆盖"处理:报文里带来的字段写进去,没带到的字段保持原值、不清空。 + * 这是上游缺失字段的语义确认之前的保守做法。新航班如果带了计划时间就直接算出运营日, + * 否则先留空,等日计划报文来收录。 */ @Singleton class AdftProcessor( @@ -182,7 +188,7 @@ class AdftProcessor( // 共享小工具(处理器层私有约定) // ===================================================================== -/** KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1/§5)。 */ +/** 航班状态变化后要发的两类事件:KAFKA_SCHD 发完整状态,KAFKA_MSG 发一条变更通知。 */ internal fun eventsFor(next: FlightSnapshot, mapper: ObjectMapper): List { val payload = linkedMapOf( "flid" to next.flid, 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 8adf937..84cac69 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt @@ -19,11 +19,16 @@ import java.time.Duration import java.time.Instant /** - * 处理主泵(docs/flight-state.md §1 严格有序 + §4 处理事务与失败规则): - * 单活动主泵严格 FIFO + HOL 阻塞 + 队头滞留转 DEAD。 + * 处理主泵:一个线程按消息 ID 从小到大一条条处理,保证先来的先处理。 * - * 终态与业务写入在处理器事务内原子提交;本泵只负责调度、边界化失败迁移与 - * 提交后的最佳努力回填(message-lifecycle §3/§4)。 + * 每次 tick 只看当前最小的未完成消息("队头"): + * - 没有待处理消息就睡一个轮询间隔; + * - 队头失败了还在退避期,就等到能重试的时刻;如果重试次数用尽或滞留太久, + * 直接转死信,不放任它一直堵着; + * - 其余情况交给 [MessageProcessor] 处理。 + * + * 一次只处理一条是刻意的。后面的消息不能越过卡住的队头,否则同一条航班的报文 + * 可能被乱序应用,几十秒后才到的旧报文会把新状态覆盖回去。 */ @Singleton class Pump( @@ -38,7 +43,7 @@ class Pump( @Volatile private var running = true - /** 优雅停机:loop 在当前 tick 收尾后退出;线程中断由 PipelineLifecycle 负责。 */ + /** 请求停机:当前 tick 跑完就退出。线程中断由 PipelineLifecycle 负责。 */ fun stop() { running = false } @@ -51,7 +56,7 @@ class Pump( Thread.currentThread().interrupt() return } catch (e: Exception) { - // 最后防线:失败状态迁移已在 processOne 边界内完成;致命 Error 不捕获 + // 兜底:单条消息的失败状态已在 processOne 内记录,这里只避免线程退出 sleepQuietly(props.pipeline.pollInterval) } } @@ -74,7 +79,7 @@ class Pump( } else { sleepQuietly(Duration.between(clock.instant(), head.nextAttemptAt)) } - // PENDING、或 FAILED 退避已到期:交处理入口(内部有边界化失败迁移与 attempts 守卫) + // 其余情况(新消息,或退避到期的重试)交给处理入口 else -> processor.processOne(head) } } @@ -89,9 +94,13 @@ class Pump( } /** - * processOne:解码 → 绑定 → 处理器(事务内决策 + 落库 + 终态 + 回填意图)→ 提交后最佳努力回填。 - * 边界化失败迁移(ProcFailure):任何意外异常归于本条 head,FAILED(INFRA)+退避,不穿出杀泵; - * MALFORMED / PROTOCOL 直接 DEAD 不重试(docs/design.md §2.3 错误分类)。 + * 处理一条消息:读原文 → 解码 → 绑定业务身份 → 分派给对应处理器 → 提交后回填标记。 + * + * 业务数据和终态由各处理器在自己的事务里写入。终态一旦落下(成功、跳过或死信), + * 这里马上试一次把处理标记写回信箱;写不进去也没关系,回填扫描会按退避继续重试。 + * + * 任何意外异常都算在当前这条消息头上(记 FAILED(INFRA) 后重试),不会把主泵线程带崩。 + * 报文非法和整包协议拒绝不重试,直接进死信等人工处置。 */ @Singleton class MessageProcessor( @@ -126,7 +135,7 @@ class MessageProcessor( } } - /** @return 是否已达终态(终态才允许回填信箱标记) */ + /** @return 这条消息是否已经落到终态(只有终态才允许回填信箱标记) */ private fun processInternal(head: ProcState): Boolean { // 守卫:手工/遗留 FAILED 行若 attempts 已达上限,直接终态(防止退避到期后无限重试) if (head.state == ProcStatus.FAILED && procFailure.scheduler.exhausted(head.attempts)) { 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 89ae6fa..553814b 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessor.kt @@ -29,25 +29,32 @@ import java.time.Instant import java.time.LocalDate import java.time.ZoneId -/** 处理器执行结果——终态与回填意图由处理器在自己的业务事务内落库(message-lifecycle §2/§4)。 */ +/** + * 处理器告诉调用方这一条消息处理成了什么。 + * 无论哪种结果,终态和回填待办都由处理器在自己的事务里写好。 + */ sealed interface ApplyResult { - /** 业务成功(含幂等成功)。 */ + /** 业务处理成功(包含"重复写入但结果一致"这种幂等成功)。 */ data object Succeeded : ApplyResult - /** MSG_ID 已有成功终态 → 重放,直接记幂等成功(docs/flight-state.md §4 重放判定)。 */ + /** 这条消息之前已经成功处理过,这次只是重放,不重复写数据。 */ data object ReplaySkipped : ApplyResult - /** 整包拒绝 DEAD(PROTOCOL)(flight-state.md §3.1/§4),不重试,交人工确认。 */ + /** 整包被拒绝(如运营日冲突、声明条数不符):整包不落地、不重试,交人工确认。 */ data class DeadProtocol(val reason: String, val flags: Set = emptySet()) : ApplyResult } -/** 归属日冲突(flight-state.md §2.1:OPERATION_DAY 一经确定不可变)= 串日/错发/污染,整包拒绝。 */ +/** 同一航班的运营日对不上(一个航班只能属于一个运营日):可能是串日或错发,整包拒绝。 */ class ProtocolViolation(message: String) : RuntimeException(message) /** - * SCHD 日计划主链路(docs/flight-state.md §3.1/§4):对 DNLD 与 RESP 统一适用。 - * 同一事务内:锁 → 归属日校验 → 逐条合并写完整当前态 → 登记事件与回填待办; - * 整包校验失败或运营日冲突整包不落地。 + * 处理 SCHD 日计划报文(DNLD 和 RESP 走同一条路):把报文里的航班记录合并进航班当前态。 + * + * 顺序是:整包校验 → 加锁 → 核对每条航班的运营日 → 逐条合并写入 → 登记待发事件 → + * 写终态和回填待办。除了校验,后面所有步骤都在同一个事务里,任何一步失败整包回滚, + * 不会留下写了一半的数据。 + * + * 报文里没提到的航班不会被删除——日计划只负责写它带来的那部分。 */ @Singleton class ScheduleProcessor( @@ -75,7 +82,7 @@ class ScheduleProcessor( return ApplyResult.ReplaySkipped } - // 整包校验(§4 步骤 2 / design.md §4.1):任一失败整包不落地 → DEAD(PROTOCOL) + // 整包校验:任何一项不通过就整包拒绝,不写半份数据 val validation = FlightStateEngine.validateMessage( recsDeclared = body.recsDeclared, records = body.records, @@ -88,18 +95,18 @@ class ScheduleProcessor( val ok = validation as SnapshotValidation.Ok if (ok.perRecordDay.isEmpty()) { - // 空快照:合法但无写入,仍算成功终态 + // 报文合法但没有记录:不需要写数据,照样算处理成功 logSnapshot(head, body, SnapshotResult.COMMITTED, upserted = 0, setOf(SnapshotFlag.EMPTY), started) return ApplyResult.Succeeded } val flags = linkedSetOf() return try { - // 同一事务(flight-state.md §4):锁 → 归属校验 → 逐条合并写 → 事件/待办预登记 + // 以下都在同一个事务里:加锁 → 校验运营日 → 逐条合并写入 → 登记事件与回填待办 val upserted = txManager.inTransaction { lock.lock() - // 归属校验(§2.1):OPERATION_DAY 不可变,批量点查避免逐航班往返 + // 一次批量查出这些航班现有的运营日,逐个比对(避免逐条查询) val mains = flightState.findMainRows(ok.perRecordDay.keys) ok.perRecordDay.forEach { (flid, day) -> val existing = mains[flid] ?: return@forEach @@ -132,7 +139,7 @@ class ScheduleProcessor( events += snapshotEvents(next) } if (events.isNotEmpty()) msgEvents.insertAll(events) - // 终态与回填意图同一事务(message-lifecycle §2/§4):业务写入、事件、终态、回填意图同提交同回滚 + // 终态与回填待办跟业务数据同事务提交:要么全成,要么全回滚 procState.markTerminal(head.msgId, ProcStatus.SUCCEEDED) written } @@ -164,8 +171,10 @@ class ScheduleProcessor( ) } - /** 留痕 SCHD_SNAP_LOG(design.md §6.2):事务外追加,失败只记 error 不阻塞; - * scope 为报文各记录归属运营日的最小/最大(单日快照两者相等)。 */ + /** + * 写一条处理留痕(仅供排查和统计,不参与业务判断)。放在事务外做, + * 写失败也只记日志,不会连累处理结果。scope 是这份报文覆盖的运营日范围。 + */ private fun logSnapshot( head: ProcState, body: ScheduleBody, diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt index 8dd3afb..e60c1b0 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/PipelineSmokeTest.kt @@ -19,11 +19,12 @@ import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test /** - * U07 端到端(stub 装配,等价 /beans 核验):无 MySQL/Redis/Kafka 下, - * ①核心 bean 装配齐全;②compat HTTP 写→PG 入队→主泵领取→解码未实装→FAILED(CODEC_ERROR)+退避; - * ③U11 重放把 FAILED 拉回 PENDING;④Dispatcher flushSchd 聚合发出 schd。 - * (生产主路径 JDBC 轮询 InboxPoller;compat HTTP 写路径见 InboxService)。 - * 后台循环关闭(autostart=false),按需手动 tick,避免测试泄漏线程。 + * 不接任何外部中间件,用内存实现把整条链跑通: + * 核心 bean 能装配、HTTP 写入到进队再到主泵处理、人工重放能把失败消息拉回队列、 + * 投递按航班聚合并发出去。另外验证死信不会卡住后续消息的发现。 + * + * 生产主路径是 InboxPoller 轮询,HTTP 只是兼容入口。后台循环默认关闭, + * 测试里手动 tick,避免留下多余线程。 */ @MicronautTest class PipelineSmokeTest { @@ -102,8 +103,8 @@ class PipelineSmokeTest { } /** - * 回归(message-lifecycle §5.2/§5.3):死信到达终态时同时登记回填意图并立即回填—— - * 否则永不回填的行会永久占据发现窗口,累积到批大小后收报整体停摆。 + * 回归用例:死信在进入终态时同时记下回填待办,并马上把标记写回信箱。 + * 少了这一步,这些行永远占着每批的名额,攒够一批就再也发现不了新消息了。 */ @Test fun `dead letter reaches terminal state, gets marked and cannot block later discovery`() { diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt index 44146f9..d4bf3f7 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/health/HealthIndicatorsTest.kt @@ -12,8 +12,8 @@ import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Test /** - * U12:健康指示器 UP 判据为真实 ping(复审 P1 修正)—— - * bean 缺失 / ping false / ping 抛异常 → DOWN;仅 ping 成功 → UP。 + * 健康指示器的判定规则:投递端口要靠真实 ping 通过才算 UP,端口缺失、ping 返回 false + * 或抛异常都算 DOWN;收报与回填的状态指标按实际数据计算,端口没接上只提示未绑定。 */ class HealthIndicatorsTest { @@ -39,7 +39,7 @@ class HealthIndicatorsTest { assertEquals(HealthStatus.DOWN, kafkaHealth(null).status) // bean 缺失 } - /** OPS-2(message-lifecycle §5.3):积压、最老未处理信龄、未回填终态与水位滞后。 */ + /** 检查 /health 里那几个积压指标算得对不对。 */ @Test fun `inbox lifecycle reports backlog, oldest age, unmarked terminals and watermark lag`() { val proc = StubProcState() 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 c77cd95..f78550f 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,12 +9,11 @@ import org.junit.jupiter.api.Test import java.sql.DriverManager /** - * Flyway 迁移引擎端到端验证(docs/flight-state.md §2 权威模型 + design.md §2.1 + - * message-lifecycle.md §5.1/§5.2): - * 1. V1 基线 + V2 信箱生命周期迁移在真实 PostgreSQL 上自动成功; - * 2. flyway_schema_history 落库且 success = true; - * 3. 决策层/管道层/留痕层全表就绪;PIPELINE_LOCK 与 INBOX_CURSOR 单行种子就位; - * 4. V2 收敛结果成立:回填事实并入 PROC_STATE,BACKFILL_TODO 下线。 + * 在真实 PostgreSQL 上跑一遍迁移,确认结果符合预期: + * V1 基线加 V2 信箱生命周期都能成功执行、迁移记录显示成功、该建的表和单行种子 + * (PIPELINE_LOCK、INBOX_CURSOR)都在,回填相关字段进了 PROC_STATE、BACKFILL_TODO 已下线。 + * + * 没有可用的 PostgreSQL 时跳过(不假装通过)。 */ class FlywayMigrationTest { 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 0e1bfa9..c1e1e3a 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 @@ -19,14 +19,14 @@ import java.util.UUID import javax.sql.DataSource /** - * 信箱生命周期 SQL 语义(message-lifecycle §4/§5.1/§5.2/§11)在真实 JDBC 上的验证。 - * 以 H2 的 PostgreSQL 兼容模式承载与 V1+V2 等价的表结构——不依赖 docker/外接库, - * 覆盖 PG 侧(终态与回填意图同体、到期/超期筛选、积压观测、水位游标)与 MySQL 侧 - * (区间发现、标记单调)的实际语句行为。 + * 用真实 JDBC 跑一遍生命周期相关的 SQL,确认语句行为符合预期(不只是接口签名对)。 * - * 说明:H2 不支持 `INSERT ... ON CONFLICT DO NOTHING`, - * [JdbcProcStateRepository.insertIfAbsent] 的入队幂等由 PG 语义与 - * `InboxPollerTest`(重复轮询不再入队)分别保证。 + * 库用 H2 的 PostgreSQL 兼容模式,表结构照抄 V1 + V2,因此不需要 docker 或外接数据库 + * 就能跑。覆盖两边的真实语句:自有 PG 侧(写终态时一并写下回填待办、挑选待回填记录、 + * 积压统计、水位游标读写)和共享信箱侧(按 ID 区间读、标记只写一次)。 + * + * 有一处覆盖不到:H2 不支持 `INSERT ... ON CONFLICT DO NOTHING`,所以入队幂等没在这里验证, + * 由 `InboxPollerTest`(重复轮询不再登记)和 PostgreSQL 本身的语义来保证。 */ class InboxLifecycleJdbcSqlTest { 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 1a8193d..d46bd19 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/ingress/InboxPollerTest.kt @@ -16,11 +16,12 @@ import org.junit.jupiter.api.Test import java.time.Instant /** - * 收报发现权不变量(docs/message-lifecycle.md §5.1/§5.3 + architecture.md §5 严格 FIFO): - * - 扫描按 ID 区间,**不受处理标记影响**:终态而未回填的行不得占据批次、不得阻断新信发现; - * - 水位只随成功入队推进,且与入队同事务(中断后由重扫补建); - * - 遇空洞即停(较小 ID 未入队时不得被后续消息越过);空洞老化后放行(水位不得永久停摆); - * - 收报层不写处理标记。 + * 收报环节最要紧的几条规矩: + * - 取新消息只看 ID,不看处理标记。处理完却没能回填的行(尤其是永远不回填的死信) + * 不允许占住批次,也不允许挡住后面的新消息——这是曾经的线上隐患; + * - 水位只在成功登记后才推进,而且和登记写在同一个事务里,中断后重扫就能补齐; + * - 遇到 ID 缺口先停下来(可能有更小的消息还没到),缺口等太久则跳过(否则水位永远卡住); + * - 这一层不碰信箱的处理标记,标记留给回填环节写。 */ class InboxPollerTest { 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 6b14e87..11e1de5 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt @@ -20,9 +20,12 @@ import java.time.Instant import java.time.ZoneOffset /** - * 回填通道不变量(docs/message-lifecycle.md §3/§4/§5.2/§11): - * 终态即刻补写、失败退避重试、超期期限 R 覆盖退避、标记单调只写空标、 - * 中间态不适用、回填失败不回改终态、重启后扫描不依赖内存状态。 + * 回填环节的规矩: + * - 处理完马上写标记,写不进去就按退避重试,重启后接着重试(状态都在数据库里); + * - 等得太久的消息无视退避强制补写,保证标记最终一定会写上; + * - 只写还没有标记的行,已有的值不覆盖; + * - 还没处理完的消息(PENDING / FAILED)不许写标记; + * - 回填失败只是回填的事,不会把处理结果改回去。 */ class BackfillServiceTest { 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 164368d..ece8bbd 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt @@ -25,9 +25,9 @@ import org.junit.jupiter.api.Test import java.time.LocalDate /** - * FDEL(docs/flight-state.md §3.3)+ ADFT(§3.3/§2.1)不变量: - * ACTIVE→DELETED 发布一次 tombstone 且明细保留;重复/迟到 FDEL 幂等不推进版本; - * ADFT 重激活恢复 ACTIVE;新实例建立;缺失标量不误删(§3.3 保守 Set-only)。 + * 删除与异常航班处理的规矩:删除在用航班时只发一次删除事件、明细保留; + * 重复或迟到的删除报文算成功但不推进版本;ADFT 能把删除过的航班重新激活、 + * 也能新建航班;报文没带的字段不会被清空。 */ class FdelAndAdftProcessorTest { 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 e4a0827..d2d35aa 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/ScheduleProcessorTest.kt @@ -31,9 +31,9 @@ import org.junit.jupiter.api.Test import java.time.LocalDate /** - * 日计划主链路(docs/flight-state.md §3.1/§4)不变量:重放幂等、整包校验失败 - * DEAD(PROTOCOL)、归属日冲突不落地、DELETED 不恢复且记冲突告警、 - * 成功路径版本推进 + 事件 + 留痕 + 待办预登记。 + * 日计划处理的几条规矩:重复消息不重复写数据;整包校验不过就整包不落地; + * 航班运营日对不上时整包拒绝;已删除的航班不因为日计划又活过来; + * 正常路径要推进版本、发事件、留痕,并把终态和回填待办一起写在同一事务里。 */ class ScheduleProcessorTest {