docs(kdoc): 其余模块注释统一为直白说明

把上一提交没覆盖到的模块也改完,至此仓库内不再有"只写编号、不写说明"的注释。

涉及 33 个文件(生产 24 + 测试 9),包括航班状态模型与合并引擎、事件与留痕模型、
运营日计算、投递器、历史归档作业、编解码与身份、重试与重放、健康指示器、
生命周期装配、以及各自的测试类。

改法与前面一致:先说这段代码做什么,再说为什么这么做,术语第一次出现就地解释
(FLID、OPERATION_DAY、航班版本号、删除通知、待发事件表等)。工单编号
(U07/N02/ACM2-10 之类)已全部移除;段落编号只剩 10 处,全部带文件名位于句尾,
用作延伸阅读,例如"见 docs/flight-state.md §5"。

保留了本来就自解释的行内注释、PgTestSupport 顶部的环境变量默认值表格,
以及 InboxController 里描述待补接口的 TODO。

校验:逐文件做"去注释后比对"(块注释 + 行注释清除、空白归一),33 个文件代码
零差异;另做逐行非注释代码比对,同样零差异。clean test 为 78 passed / 1 skipped。
This commit is contained in:
windyboy
2026-09-10 11:14:11 +08:00
parent 6e7d819034
commit 1183c7522b
33 changed files with 283 additions and 210 deletions
@@ -11,9 +11,11 @@ import jakarta.annotation.PreDestroy
import jakarta.inject.Singleton
/**
* U07(T03+N29):管道生命周期装配——服务启动后拉起 inbox 轮询、主泵、投递三条专用单线程
* 停机时 requestStop + interrupt + join。仅当 `msgx.pipeline.autostart=true` 时装配
* (默认关:需要真实仓储或 msgx.stubs=true 才安全开启)。
* 管道的启动与停机:服务起来后拉起 inbox 轮询、主泵、投递三条循环,每条各占一个专用平台线程
* 停机时先让它们各自收尾,再中断线程解除 sleep 并等待退出。
*
* 只有配置 msgx.pipeline.autostart=true 时才装配,默认不开:得有真实仓储或者 msgx.stubs=true
* 才能安全自启。
*/
@Requires(property = "msgx.pipeline.autostart", value = "true")
@Singleton
@@ -53,7 +55,7 @@ class PipelineLifecycle(
pump.stop()
dispatcher.stop()
jobRunner.stop()
threads.forEach { it.interrupt() } // 解除 Thread.sleep 阻塞,加速退出
threads.forEach { it.interrupt() } // 打断 sleep,让循环立刻醒过来退出,而不是干等一轮间隔
threads.forEach { runCatching { it.join(3000) } }
log.info("pipeline loops stopped")
}
@@ -4,7 +4,8 @@ import com.gzzn.omms.msgexchange.domain.DecodedMessage
import com.gzzn.omms.msgexchange.domain.ErrorClass
/**
* ACMA-8 流程 2:解码失败分类(矩阵/状态机:MALFORMED 不重试;CODEC_ERROR 可一键重放)
* 解码失败的结果:错误分类加上一句原因
* MALFORMED 是报文本身不合法、重试也不会过;CODEC_ERROR 是解码逻辑的问题,修好后可以重放。
*/
data class DecodeFailure(val errorClass: ErrorClass, val detail: String)
@@ -14,8 +15,9 @@ sealed interface DecodeResult {
}
/**
* XML codec:使用 SIS 文档对应的 Jackson XML 注解 DTO,再显式转换为领域载荷;
* XMLInputFactory 禁用 DTD/外部实体,原文仍由入站层保留。
* XML 解码入口:先用标注了 Jackson XML 注解 DTO 接住报文体,再显式转成领域模型。
*
* 解析时关掉 DTD 和外部实体,避免 XML 外部实体注入;报文原文由入站层另行留存,这里不管。
*/
interface XmlCodec {
fun decode(rawXml: String): DecodeResult
@@ -2,24 +2,24 @@ package com.gzzn.omms.msgexchange.config
import io.micronaut.context.annotation.ConfigurationProperties
/** 生命周期与保留期(docs/flight-state.md §6design.md §6.2)。窗口按机场时区算。 */
/** 航班归档删除的时间窗口和留痕保留期,对应配置里的 `msgx.history.*`窗口按机场时区算。 */
@ConfigurationProperties("msgx.history")
class HistoryProps {
/** 取消(CNCL 非空)超过 N 小时。 */
/** 取消时间字段CNCL非空,且已经过去这么多小时。 */
var cancelledHours: Long = 48
/** 到/离终态(NAAT/NEAT,含义待术语表确认 §6 开放项)超过 N 小时。 */
/** 到/离终态字段NAAT/NEAT)非空且已过去这么多小时;这两个字段的确切含义还要跟业务确认。 */
var terminalHours: Long = 48
/** STATE = DELETED 超过 N 小时。 */
/** 航班已标记删除(STATE = DELETED)并且过了这么多小时。 */
var deletedHours: Long = 48
/** 无终态字段:最后有效更新超过兜底期限(默认 7 天)。 */
/** 兜底期限:航班既没取消也没到终态时,最后一次有效更新超过这么多小时就算静默,可以清理。 */
var idleHours: Long = 24 * 7
/** 留痕 SCHD_SNAP_LOG 保留天数(design.md §6.2,按 (SCOPE_END, RECV_AT) 清理)。 */
/** 日计划留痕 SCHD_SNAP_LOG 保留多少天,按(快照覆盖截止日, 接收时间)删旧行。 */
var snapLogRetentionDays: Long = 90
/** 是否接通历史存储——未接通时清理必须删 0 条(flight-state.md §6 红线)。 */
/** 有没有历史存储可用。没接通时清理必须一条都不删,这是硬性前提。 */
var historyStoreEnabled: Boolean = false
}
@@ -3,15 +3,17 @@ package com.gzzn.omms.msgexchange.config
import io.micronaut.context.annotation.ConfigurationProperties
/**
* 运营日计算(docs/flight-state.md §2.1):由 SODTddMMMyyHHmm)与机场时区按切日边界推导;
* OPERATION_DAY 不是接收日/落库日,一经确定不可变。切日边界业务配置
* 不得假设等于自然日零点(默认 0 点为占位,业务口径待确认)
* 算航班运营保障日(OPERATION_DAY)要用的两个参数,对应配置里的 `msgx.operation-day.*`
* 机场时区和切日边界。运营日由计划运行时间 SODT 加时区推出来,不是接收日也不是落库日
* 算出来就固定不变
*
* 不能假设运营日等于自然日,民航的一天从几点开始由业务定。默认 0 点只是占位,待确认。
*/
@ConfigurationProperties("msgx.operation-day")
class OperationDayProps {
/** 机场时区(IANA。 */
/** 机场所在时区,用 IANA 名称,比如 Asia/Shanghai。 */
var zone: String = "Asia/Shanghai"
/** 切日边界:SODT 本地时刻早于该小时的归属前一运营日023。 */
/** 切日边界小时SODT 本地时刻早于这个点就归属前一运营日。取值 023。 */
var cutoffHour: Int = 0
}
@@ -47,7 +47,10 @@ class PipelineProps {
/** 每次回填扫描最多处理多少条。 */
var backfillBatch: Int = 100
/** U07:启动即拉起 Pump/Dispatcher 循环(默认关——需要真实仓储或 msgx.stubs=true 才可安全开启)。 */
/**
* 服务启动后是否自动拉起收报、主泵、投递三个循环。
* 默认关闭:只有接了真实仓储、或者明确用内存 stub 跑的时候才安全。
*/
var autostart: Boolean = false
/** N28attempt ≤ 0(如 FAILED 未递增 attempts 的行)不得抛异常,取下界=首档退避。 */
@@ -12,29 +12,30 @@ import jakarta.inject.Singleton
import java.time.Duration
import java.time.Instant
/** 对外投递端口(Kafka 同步确认,at-least-once;阶段 B 追加 ES 投影写入)。 */
/** 对外投递的出口:同步等 Kafka 确认,语义是至少发一次(可能重复,但不会丢)。以后 ES 投影也从这里加。 */
interface DeliveryPort {
/** KAFKA_MSG 变化通知(key=FLID)。 */
/** 发一条变化通知,消息 key 是 FLID(航班实例 ID)。 */
fun sendKafka(topic: String, key: String, payloadJson: String)
/**
* KAFKA_SCHD 整态(key=FLID)。KAFKA_MSG 与 KAFKA_SCHD 映射同一 topic 语义由适配层定;
* target→topicKAFKA:msg→"msg"KAFKA:schd→"schd"。
*/
/** 发一条完整状态,消息 key 是 FLID。适配层负责把 target 映射成 topicKAFKA:msg 对应 "msg"KAFKA:schd 对应 "schd"。 */
fun sendKafkaSchd(topic: String, key: String, payloadJson: String)
/** TOMBSTONEkey=FLID、value=null——整态键缺失表示删除旧值(flight-state.md §5。 */
/** 发一条删除通知keyFLID、value 为空;下游按"整态里这个键没了"理解成删除。 */
fun sendKafkaNull(topic: String, key: String)
/** 连通性探测(健康检查用);默认 true,真实 Kafka 实装时覆写为 producer metadata 校验。 */
/** 给健康检查用的连通性探测;默认返回 true,真实 Kafka 实现要覆写成向 broker 拉一次 metadata 来判断。 */
fun ping(): Boolean = true
}
/**
* 投递调度docs/flight-state.md §5 + design.md §5.1/§5.2):逐条 KAFKA_MSG 严格 FIFO
* KAFKA_SCHD 走 flushSchd 批量——同一 FLID 未发事件按最新 STATE_VERSION 合并输出,
* TOMBSTONE 发 null 值消息。两主题间不保证顺序(§5)。
* 批量闭环:队首退避未到期不 claim;发送失败整批 attempts+1 退避,达上限整批 DEAD/DLQ
* 投递调度:把 outbox(待发事件表)里的事件发给下游。
*
* KAFKA_MSG 一条一条按登记顺序发,不插队。KAFKA_SCHD 走 flushSchd 批量发:同一个 FLID 攒了
* 多条未发事件时只发版本号最新的那条,旧的自然作废;删除通知发 value 为空的 tombstone
* 两个主题之间不保证先后顺序。
*
* 失败处理:队首的重试时间没到就不取;一批里有发送失败,整批重试次数加一并推后退避,
* 次数用尽整批转 DEAD 当死信。见 docs/flight-state.md §5。
*/
@Singleton
class Dispatcher(
@@ -50,7 +51,7 @@ class Dispatcher(
private var lastFlush: Instant? = null
/** 优雅停机:loop 收尾后退出;线程中断由 PipelineLifecycle 负责。 */
/** 请求停机:置位后 loop 走完当前一轮就退出;真正中断线程由 PipelineLifecycle 负责。 */
fun stop() {
running = false
}
@@ -91,7 +92,7 @@ class Dispatcher(
}
}
/** flushSchd:同 FLID 未发事件按最新 STATE_VERSION 合并(§5);TOMBSTONE 发 null。 */
/** 批量发 KAFKA_SCHD:每个 FLID 只发版本号最新的那条未发事件,删除通知发空 value。 */
internal fun flushSchd() {
val batch = try {
msgEvents.mergePendingSchd(props.schd.flushLimit)
@@ -116,7 +117,7 @@ class Dispatcher(
}
val sentIds = batch.mapNotNull { it.eventId }.toSet() - failures.mapNotNull { it.eventId }.toSet()
if (sentIds.isNotEmpty()) msgEvents.markAllSent(sentIds.toList())
// 被新版本合并压掉的未发事件同样关闭(§5:同 FLID 只按最新 STATE_VERSION 输出一次
// 被新版本压掉的旧事件也要标成已发,否则它们会一直留在队里:同一个 FLID 只按最新版本输出一次
val sentVersions = batch.associate { it.partitionKey to it.stateVersion }
runCatching {
while (true) {
@@ -130,7 +131,7 @@ class Dispatcher(
lastFlush = scheduler.now()
}
/** 条事件失败迁移:attempts+1;达上限 DEAD(EXHAUSTED)DLQattempts 落库审计),否则退避重试。 */
/** 条事件失败之后怎么走:重试次数加一,到上限就标成 DEAD(EXHAUSTED) 留作死信(次数落库便于追查),否则退避推到下次再发。 */
private fun retryOrDead(e: MsgEvent, lastError: String) {
val eventId = e.eventId ?: return
val attempts = e.attempts + 1
@@ -1,8 +1,10 @@
package com.gzzn.omms.msgexchange.domain
/**
* 解码后的入站报文(docs/design.md §2.2 统一载荷)。META 字段实名
* SNDR/SEQN/DTTM,其中 SNDR|TYPE|STYP|SEQN 是业务幂等键来源(Identity.of
* 一条入站报文的元信息(报文头 META 段)。SNDR 是发送方,TYPE/STYP 是报文类型和子类型,
* SEQN 是发送方流水号,DTTM 是发送时间
*
* SNDR|TYPE|STYP|SEQN 拼起来就是业务幂等键,用来判断同一条报文是否已经处理过。
*/
data class MetaFields(
val sndr: String,
@@ -12,13 +14,13 @@ data class MetaFields(
val dttm: Long,
)
/** 消息分派(design.md §2.2:一等分派键,穷尽 when;FDEL 一等公民——终止航班实例)。 */
/** 报文的种类:日计划、运行动态、FDEL(航班终止),或者还没支持的类型。分派时用穷尽 when 保证不漏分支。 */
sealed interface MsgKind {
data class Schd(val subtype: SchdSubtype) : MsgKind
data class Flop(val subtype: String) : MsgKind // 运行动态 STYPFDEL 除外)
data class Flop(val subtype: String) : MsgKind // 运行动态subtype 是报文子类型;FDEL 单独成类,不走这里
data object Fdel : MsgKind
/** 未支持类型(design.md §2.3FAILED(UNSUPPORTED),达阈值转 DEAD。 */
/** 还没有对应处理器的类型:先按 UNSUPPORTED 记失败并退避重试,重试次数用尽转 DEAD 等人工处置。 */
data class Unsupported(val tag: String) : MsgKind
enum class SchdSubtype { RESP, DNLD, ADFT }
@@ -28,7 +30,7 @@ data class DecodedMessage(
val meta: MetaFields,
val kind: MsgKind,
val rawXml: String,
/** 各类型载荷(vendor POJO 或 Kotlin 模型,阶段 1 后续统一。 */
/** 解码出来的载荷对象:目前有的类型是厂商 DTO、有的是本地领域模型,之后会统一。 */
val body: Any? = null,
) {
val typeTag: String
@@ -1,21 +1,25 @@
package com.gzzn.omms.msgexchange.domain
/** MSG_EVENT 事件形态(design.md §2.1UPSERT 整态/通知,TOMBSTONE 删除。 */
/** 事件的两种形态:UPSERT 是新增或更新(可能是整态也可能是变化通知TOMBSTONE 删除。 */
enum class EventType { UPSERT, TOMBSTONE }
/** MSG_EVENT 投递状态机(design.md §2.3PENDING → SENT失败退避重试;耗尽转 DEAD 保留作 DLQ。 */
/** 投递状态PENDING 待发 → SENT 已发出;发送失败退避重试,重试次数用尽转 DEAD,留在表里当死信队列。 */
enum class EventStatus { PENDING, SENT, DEAD }
/**
* MSG_EVENToutboxdesign.md §2.1):状态变更、删除通知
* KAFKA_SCHD 发整态、KAFKA_MSG 只通知变化(docs/flight-state.md §5,两主题不承诺顺序);
* TOMBSTONE 仅在 ACTIVE→DELETED(§3.3)或生命周期清理前补发(§6),与删除同事务登记。
* stateVersion:发布时的航班版本;KAFKA_SCHD 聚合按 FLID 取最新(§5)。
* MSG_EVENT:一张待发事件表(outbox),记录航班状态变更要对外发什么
*
* 写业务数据时在同一个事务里往这里插一行,投递线程随后按行发出,这样业务提交和"该发的事件"
* 不会脱节。KAFKA_SCHD 发完整状态,KAFKA_MSG 只发"这个航班变了"的通知;两个主题之间不保证
* 先后顺序。TOMBSTONE 只在两种情况下登记:航班从在用变成删除,或者被生命周期清理前补发
* 一次删除通知。见 docs/flight-state.md §5。
*
* stateVersion 是发布时的航班版本号;同一 FLID 攒了多条待发事件时,只发版本号最新的那条。
*/
data class MsgEvent(
val eventId: Long? = null,
val target: String,
val partitionKey: String, // 恒为 FLID
val partitionKey: String, // 分区键,恒为 FLID(航班实例 ID
val eventType: EventType = EventType.UPSERT,
val stateVersion: Long = 0,
val payloadJson: String,
@@ -5,19 +5,21 @@ import java.time.LocalDateTime
import java.time.ZoneId
/**
* 运营日计算(docs/flight-state.md §2.1):计划运行时间字段 SODTddMMMyyHHmm
* 与机场时区按切日边界推导;OPERATION_DAY 不是消息接收日或落库日,一经确定不变。
* 切日边界由 msgx.operation-day.cutoff-hour 配置(默认 0 = 自然日零点,业务口径待确认)。
* 算航班的运营保障日(OPERATION_DAY):计划运行时间 SODT 和机场时区推出来。它不是消息
* 接收日、也不是落库日,算出来之后就固定不变。
*
* 民航的一天不一定从零点开始,切日边界由 msgx.operation-day.cutoff-hour 配。默认 0 点只是
* 占位,真实口径还要业务确认。见 docs/flight-state.md §2.1。
*/
class OperationDayCalculator(
zone: ZoneId,
cutoffHour: Int,
) {
private val zone: ZoneId = zone
/** 切日边界:SODT 本地时刻早于该小时的归属前一运营日023越界收敛到边界值。 */
/** 切日边界小时SODT 本地时刻早于这个点就归属前一运营日。取值 023超出范围会被夹到边界值。 */
private val cutoffHour: Int = cutoffHour.coerceIn(0, 23)
/** SODT 运营日;输入 null/空/非法返回 null不抛异常,调用方按"运营日不可计算"处理。 */
/** SODT 换算成运营日;入参为空、全空白或格式不对时返回 null不抛异常,调用方按"运营日算不出来"处理。 */
fun compute(sodt: String?): LocalDate? {
if (sodt.isNullOrBlank()) return null
val local = parseSodt(sodt.trim()) ?: return null
@@ -27,9 +29,10 @@ class OperationDayCalculator(
companion object {
/**
* SIS SODT 线格式:ddMMMyyHHmm如 15DEC031723,月份英文三字母、大小写不敏感
* 两位年显式按 2000 基准展开(java.time 的 yy reduced-value 解析跨实现不一致,
* 显式展开保证 AODB 侧年份窗口唯一口径;基准年待真实报文验收确认)。
* SODT 的报文格式:ddMMMyyHHmm,比如 15DEC031723,月份三字母英文缩写、大小写都认
*
* 两位年份显式按 2000 年展开,不交给 java.time 的 yy 解析——不同实现对两位年的处理
* 不一致,自己展开才能保证各方对年份窗口的理解相同。基准年等真实报文到了再确认。
*/
private val SODT_REGEX = Regex("(\\d{1,2})([A-Za-z]{3})(\\d{2})(\\d{2})(\\d{2})")
@@ -4,16 +4,20 @@ import java.time.Instant
import java.time.LocalDate
/**
* SCHD 快照留痕 SCHD_SNAP_LOGdocs/flight-state.md §2 权威模型表;
* 清理规则见 design.md §6.2):只追加、可重建、不参与状态决策,写失败只记指标;
* 一行 = 一次尝试,重放也记;保留 90 天,按 (SCOPE_END, RECV_AT) 清理。
* SCHD 快照留痕 SCHD_SNAP_LOG,记录每一包日计划被处理的经过。
*
* 只追加不修改;写失败不影响主流程,只记一条指标。数据丢了可以靠重放重建,它也不参与任何
* 状态决策。一行对应一次处理尝试,被重放的包会再记一行。默认保留 90 天,按(快照覆盖截止日,
* 接收时间)清理。见 docs/flight-state.md §2、docs/design.md §6.2。
*/
enum class SnapshotResult { COMMITTED, REPLAY_SKIPPED, ROLLED_BACK }
/**
* 留痕告警 flagsEMPTY(空快照合法)、RECS_DROP(声明数量不符)、
* DAY_MISMATCH(记录运营日不可计算,整包拒绝)、
* SCHD_REVIVE_CONFLICT(日计划命中 DELETED 航班,保持 DELETED 不恢复)。
* 处理日计划时记下的告警标记:
* EMPTY——空快照,合法但不是常态;
* RECS_DROP——报头声明的条数和实际收到的对不上;
* DAY_MISMATCH——有记录算不出运营日,整包拒收;
* SCHD_REVIVE_CONFLICT——日计划打到了已删除的航班上,航班保持删除不恢复。
*/
enum class SnapshotFlag { EMPTY, RECS_DROP, DAY_MISMATCH, SCHD_REVIVE_CONFLICT }
@@ -27,5 +31,5 @@ data class SnapshotLogEntry(
val durationMs: Long,
val result: SnapshotResult,
val flags: Set<SnapshotFlag> = emptySet(),
val archiveKey: String? = null, // 证据层引用(尚未交付 §2
val archiveKey: String? = null, // 指向归档原文的引用;归档侧还没交付,暂为 null
)
@@ -1,10 +1,11 @@
package com.gzzn.omms.msgexchange.domain
/**
* MSG_EVENT 投递目标(docs/flight-state.md §5KAFKA_SCHD 按 FLID 发完整状态投影、
* KAFKA_MSG 只通知变化;两主题间不保证顺序删除用 tombstone
* 常量值即 MSG_EVENT.TARGET 落库值(design.md §2.1);topic 名由投递适配层映射
* KAFKA:msg → "msg"KAFKA:schd → "schd")。ES 投影属阶段 B,暂不登记目标。
* 事件往哪儿投。KAFKA_SCHD 按 FLID(航班实例 ID)发航班的完整状态,KAFKA_MSG 只发"变化了"
* 的通知;两主题间不保证顺序删除用 value 为空的 tombstone 消息表示
*
* 常量值就是 MSG_EVENT.TARGET 落库的值,真正的 topic 名由投递适配层映射
* KAFKA:msg → "msg"KAFKA:schd → "schd")。ES 投影还没做,先不登记目标。
*/
object Targets {
const val KAFKA_MSG = "KAFKA:msg"
@@ -4,15 +4,16 @@ import java.time.Instant
import java.time.LocalDate
/**
* 航班实例当前态模型docs/flight-state.md §2/§3
* 航班实例当前态模型,规则见 docs/flight-state.md §2/§3。
*
* 身份:FLID唯一关联键(§2.1);OPERATION_DAY 一经确定不可变(§2.1),
* 未被日计划收录前可为 NULL。STATE 仅 ACTIVE/DELETED(§3.3,无 ARCHIVED——
* 物理清除只发生在历史归档成功之后 §6)。
* FLID(航班实例 ID)是航班唯一关联键,同一航班的所有报文都靠它对上号。OPERATION_DAY
* (运营保障日)在航班首次入库时确定,之后不允许再改;日计划还没收录它之前可以是 NULL。
* STATE 只有 ACTIVE(在用)和 DELETED(已删除)两种,没有"已归档"这种态——航班退出当前态
* 的唯一方式是被物理删除,而删除必须发生在历史归档成功之后。
*/
enum class FlightState { ACTIVE, DELETED }
/** FLIGHT_SCHD 主行的身份追踪字段(§2 权威模型;不含标量载荷)。 */
/** FLIGHT_SCHD 主行:航班的身份追踪字段,标量载荷存在别的字段里。 */
data class FlightMainRow(
val flid: String,
val operationDay: LocalDate?,
@@ -23,9 +24,11 @@ data class FlightMainRow(
)
/**
* SCHD DNLD/RESP条 FLTR 记录(解码产物,§3.1 日计划入口)
* scalars/collections 只含报文中出现的字段:出现 = Set/Replace(标量空值 = 显式清空),
* 未出现 = 按合并语义保留本地值(§3.1/§7 不变量);集合按完整合并结果写入(§2.2)。
* 日计划报文(SCHD DNLD/RESP)里的一条 FLTR 记录,是解码后的产物
*
* scalars 和 collections 只装报文里真正出现过的字段:出现就覆盖本地值(标量给空串表示
* 显式清空),没出现就保留库里已有的值;集合一旦出现就按合并后的完整结果整体覆盖写入。
* 合并规则见 docs/flight-state.md §3.1。
*/
data class ScheduleRecord(
val flid: String,
@@ -34,8 +37,10 @@ data class ScheduleRecord(
)
/**
* FLOP/ADFT 增量载荷(§3.2/§3.3):出现字段/集合更新,缺失保留本地值;
* ADFT 缺失字段语义待上游确认前按保守 Set-only 处理(§3.3),不沿用全量替换。
* FLOP/ADFT 报文的增量载荷:只带这次要改的字段集合,没出现的字段保留库里已有的值。
*
* ADFT 里缺失字段到底算清空还是算保留,上游还没给准话,所以暂时保守处理成"只设不改",
* 不按整包替换来理解。见 docs/flight-state.md §3.2/§3.3。
*/
data class MergeChange(
val flid: String,
@@ -44,9 +49,11 @@ data class MergeChange(
)
/**
* 完整当前态(§2.2/§3):主行 + 全部明细 = 完整当前态
* 写入前在内存生成完整新状态再落库(引擎产物),读取须在主表与全部明细的
* 一致性读边界内进行(§5),展示层过滤 STATE = ACTIVE
* 一个航班的完整当前态:主行字段加上它的全部明细
*
* 状态引擎先在内存里推出完整的新状态再落库,所以这个对象既是引擎的输出,也是落库的输入
* 读的时候要在主表和全部明细的一致读边界内取,否则会读到半新半旧;对外展示时只挑
* STATE = ACTIVE 的航班。
*/
data class FlightSnapshot(
val flid: String,
@@ -57,19 +64,22 @@ data class FlightSnapshot(
val collections: Map<String, List<Map<String, String>>>,
)
/** §6 生命周期判定窗口(按机场时区算;窗口值业务配置。 */
/** 判定一个航班可以归档删除的四条时间窗口,单位是小时,按机场所在时区算;具体取值走业务配置。 */
data class HistoryRules(
val cancelledHours: Long = 48, // CNCL 非空过 N 小时
val terminalHours: Long = 48, // NAAT/NEAT 终态超过 N 小时(字段含义待术语表确认 §6
val deletedHours: Long = 48, // STATE = DELETED过 N 小时
val idleHours: Long = 24 * 7, // 无终态字段:最后更新超过兜底期限
val cancelledHours: Long = 48, // 取消时间字段(CNCL非空且已过 N 小时
val terminalHours: Long = 48, // 到达/离开终态字段(NAAT/NEAT)非空且已过 N 小时
val deletedHours: Long = 48, // 已标记删除(STATE = DELETED)且已过 N 小时
val idleHours: Long = 24 * 7, // 既没有终态也没有取消:最后更新距今超过这个兜底期限
)
/**
* §6 命中生命周期判定的航班(先归档 → 后物理清除;历史存储失败必须删 0 行)。
* wasNeverFdel:是否未经 FDEL 而被生命周期清除——清除前须补发一次删除事件(§3.3/§5)
* 当前为推断口径(无持久化"曾 FDEL"痕迹):state=ACTIVE 视为未经 FDEL
* FDEL→ADFT 重激活后的航班会被误判并重复补发,收敛见 §6 开放项。
* 生命周期判定选中、准备归档并物理删除的航班。顺序固定是"先归档、再删除",
* 归档没确认成功就一行都不能删
*
* wasNeverFdel 表示这个航班从没收到过 FDEL(航班终止报文),是被生命周期直接清掉的,
* 所以清除前要补发一次删除通知,否则下游不知道它已经没了。目前只能靠推断:没有地方记录
* "曾经收过 FDEL",于是把 state = ACTIVE 当成没收到过。副作用是 FDEL 之后又被 ADFT
* 重新激活的航班会被误判、重复补发,解决办法待定(见 docs/flight-state.md §6 开放项)。
*/
data class HistoryCandidate(
val flid: String,
@@ -5,25 +5,30 @@ import com.gzzn.omms.msgexchange.domain.SnapshotFlag
import java.time.LocalDate
/**
* 航班状态引擎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)——只修改报文表达的字段/集合。
* 航班状态引擎:把报文内容和库里已有的状态合成出新状态,全部在内存里算,不碰数据库和
* Kafka,方便单独测。三个入口
* - [validateMessage]:整包校验日计划,检查声明条数、航班标识、运营日能不能算出来
* - [snapshotState]把日计划(DNLD/RESP)合并进当前态——报文里出现的字段覆盖本地值,
* 没出现的保留,标量给空串表示清空;已删除的航班保持删除,日计划救不回来
* - [mergedState]FLOP/ADFT 增量合并进当前态,只改报文明确表达的字段集合。
*
* 合并规则见 docs/flight-state.md §3.1/§3.2。
*/
object FlightStateEngine {
/** 集合键白名单10 集合 9 张明细表8 张资源表 + FLIGHT_ROUTE_POINT
* ROUT/ERUT 共用路线表,flight-state.md §2.2。 */
/** 允许落库的集合键:10 集合对应 9 张明细表——8 张资源表,加上存经停点的
* FLIGHT_ROUTE_POINT;其中 ROUT 和 ERUT 共用同一张路线表。 */
val COLLECTION_KEYS: Set<String> = setOf(
"GTDT", "CKDT", "CLDT", "PSDT", "CHDT", "DELY", "ABTM", "CHOT", "ROUT", "ERUT",
)
/**
* 日计划整包校验(§4 步骤 2 / §3.1):RECS 09999 且等于实收数;每条含合法
* 数字型 FLIDSIS Number(1-12));快照内不重复;每条记录运营日可计算。
* 任一失败整包不落地(DEAD(PROTOCOL)
* 校验一包日计划:报头声明的记录数(RECS)要在 09999 之间、且等于实际收到的条数;
* 每条记录的 FLID 必须是 1–12 位数字,快照内不能有重复 FLID,每条记录运营日都要
* 能算出来
*
* 任何一条不合格就整包拒收、一条都不入库,记为 DEAD(PROTOCOL) 等人工处置——日计划是
* 整包语义,落一半会让当前态自相矛盾。
*/
fun validateMessage(
recsDeclared: Int,
@@ -63,11 +68,14 @@ object FlightStateEngine {
}
/**
* 日计划合并(§3.1):标量出现 Set(空值 = 显式清空)、缺失保留本地值;
* 集合出现 Replace、缺失保留。新航班(current == null)无本地值可保留,
* 集合按全键输出空集(§2.2 完整合并结果写入口径)。
* `keepDeleted = true` 时 STATE 保持 DELETED(§3.3:日计划不承担恢复,
* 恢复入口只有 ADFT);冲突告警由调用方按 SCHD_REVIVE_CONFLICT 记录
* 把一条日计划记录合并进当前态:报文里出现的标量覆盖本地值(给空串表示显式清空),
* 没出现的保留;集合一旦出现就整体替换,没出现就保留。
*
* 新航班(current == null)没有本地值可保留,所以集合按全部键输出空集,保证落库时
* 每张明细表都有确定的行集
*
* keepDeleted = true 时即使收到日计划也保持 DELETED:日计划不能把删掉的航班救回来,
* 唯一的恢复入口是 ADFT。调用方负责记一条 SCHD_REVIVE_CONFLICT 告警。
*/
fun snapshotState(
current: FlightSnapshot?,
@@ -81,32 +89,32 @@ object FlightStateEngine {
else -> current.state
}
val scalars = buildMap {
current?.scalars?.let(::putAll) // 缺失 = 保留(§3.1
current?.scalars?.let(::putAll) // 报文没带的标量保持原值
putAll(record.scalars) // 出现 = Set;空串 = 显式清空
}
val collections = buildMap {
if (current == null) {
COLLECTION_KEYS.forEach { key -> put(key, emptyList()) } // 新航班全键输出
} else {
putAll(current.collections) // 缺失 = 保留(§3.1
putAll(current.collections) // 报文没带的集合保持原值
}
record.collections.forEach { (key, items) ->
if (key in COLLECTION_KEYS) put(key, items) // 出现 = Replace(§3.1
if (key in COLLECTION_KEYS) put(key, items) // 报文带了整个集合就整体替换
}
}
return FlightSnapshot(
flid = record.flid,
operationDay = operationDay,
state = state,
stateVersion = (current?.stateVersion ?: 0L) + 1, // 每次成功写入 +1(§3.1/§4
stateVersion = (current?.stateVersion ?: 0L) + 1, // 每次成功写入版本号加一
scalars = scalars,
collections = collections,
)
}
/**
* FLOP/ADFT 增量合并(§3.2/§3.3):标量出现覆盖、缺失保留;集合出现 Replace、
* 缺失保留——只修改报文表达的字段/集合,其余航班状态保持不变;不推导运营日
* FLOP/ADFT 增量合并进当前态:出现的标量覆盖、出现的集合整体替换,没出现的都保留。
* 改报文明确表达的字段集合,不重新推导运营日,也不动航班的在用/删除状态
*/
fun mergedState(current: FlightSnapshot, change: MergeChange): FlightSnapshot {
val scalars = buildMap {
@@ -126,11 +134,11 @@ object FlightStateEngine {
)
}
/** FLID:数字型(SIS Number(1-12)。 */
/** FLID 必须是 112 位纯数字。 */
private val FLID_REGEX = Regex("\\d{1,12}")
}
/** 整包校验结果:Ok 带每条记录的归属运营日与观测 flagsInvalid 整包拒。 */
/** 整包校验结果:Ok 带每条记录的运营日和观察到的告警标记,Invalid 表示整包拒收、带上原因。 */
sealed interface SnapshotValidation {
data class Ok(
val perRecordDay: Map<String, LocalDate>,
@@ -10,11 +10,13 @@ import jakarta.inject.Singleton
import org.reactivestreams.Publisher
/**
* U12R05):自定义健康指示器——
* Kafka(投递端口)。经 BeanProvider 可选解析:
* 缺 bean(如未用 stub 也未实装)时指示 DOWN 而非启动失败;
* UP 判据为真实 pingfalse/异常 → DOWN),而非仅 bean 存在(复审 P1 修正)。
* ACM2-28Redis 退出阶段 A 权威与写路径,redis-flight-store 指示器移除)
* 把 Kafka 投递端口的状态报给 /health。
*
* 用 BeanProvider 可选注入:没配 stub、也还没有真实实现时,端口 bean 根本不存在,这时报 DOWN
* 而不是让服务起不来。判 UP 的依据是真的调一次 ping——返回 false 或抛异常都算 DOWN,
* 不是"bean 在就算好"
*
* Redis 已经不再参与航班状态的写入和判定,对应的健康指示器也一并删掉了。
*/
@Singleton
class KafkaDeliveryHealthIndicator(
@@ -3,8 +3,8 @@ package com.gzzn.omms.msgexchange.infra.log
import org.slf4j.MDC
/**
* U12R05):处理路径入口写入 MDC traceId=cminmsgsId/eventId),
* 使一条消息全链路日志可串(logback %X{traceId} + logstash includeMdcKeyName
* 处理一条消息时把它的 ID(信箱 ID 或事件 ID)放进 MDC traceId,用 withTrace 把整段处理
* 逻辑包起来;日志配置里带上 %X{traceId},同一条消息的日志就能串到一起
*/
object TraceLog {
fun <T> withTrace(id: Any, body: () -> T): T {
@@ -233,7 +233,10 @@ interface ReqTrackRepository {
fun findLatest(reqType: String, operationDay: LocalDate, sender: String, states: List<ReqState>): Req?
/** RESP 完成请求(design.md §4.2):匹配最新一条 PENDING/SENT;无匹配返回 false(迟到不报错)。 */
/**
* 应答报文到达时把对应请求标记为已完成:匹配最近一条待应答的请求(PENDING 或 SENT)。
* 找不到返回 false——迟到或多余的应答不算错误。
*/
fun completeLatest(reqType: String, operationDay: LocalDate, sender: String): Boolean
fun linkCoutmsgs(reqId: Long, coutmsgsId: Long)
@@ -7,28 +7,29 @@ import java.time.Clock
import java.time.Instant
/**
* 统一重试策略ACM2-10 U08 MessageProcessor / SnapshotFlowProcState
* DispatcherMsgEvent 共用attempts 递增后按 backoff 表给 nextAttemptAt
* exhausted 判定与两侧同源maxAttempts时间一律经可注入 Clock测试用固定钟避免脆弱睡眠
* 重试节奏的唯一出处处理失败的入站消息和待发事件都用它算下次重试时间保证两边口径一致
*
* 重试次数加一之后按退避表算出下次可以重试的时刻"次数是否已经用尽"也在这里判断
* 上限取 msgx.pipeline.max-attempts时间一律走注入的 Clock测试可以塞固定时钟不用 sleep
*/
@Singleton
class FailureScheduler(
private val props: PipelineProps,
private val clock: Clock,
) {
/** 便捷构造:默认系统时钟(生产路径)。 */
/** 生产环境用这个构造:时钟取系统 UTC 时间。 */
constructor(props: PipelineProps) : this(props, Clock.systemUTC())
fun now(): Instant = clock.instant()
fun exhausted(attempts: Int): Boolean = attempts >= props.pipeline.maxAttempts
/** attempts 指递增后的值;N28attempt ≤ 0 由 backoffFor 兜底为首档。 */
/** 传入的 attempts 是已经加一之后的值。若传 0 或负数(比如 FAILED 行还没记过次数),退避表会兜底给第一档。 */
fun nextAttemptAt(attemptsAfterIncrement: Int): Instant =
now().plusMillis(props.pipeline.backoffFor(attemptsAfterIncrement))
}
/** 提供可注入 Clockjava.time.Clock);测试可用 Clock.fixed(...) 或自定义可变钟覆盖。 */
/** Clock 注册成可注入的 bean;测试里可以换成 Clock.fixed(...) 或自己写的可推进时钟。 */
@Factory
class TimeFactory {
@Singleton
@@ -5,19 +5,21 @@ import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository
import jakarta.inject.Singleton
/**
* U11 显式重放入口只允许可恢复的错误类 FAILED/DEAD PENDING主泵重
* 不可恢复类MALFORMED报文非法重放必再失败与未知类一律不在白名单内
* 人工重放入口把失败的消息 FAILED/DEAD PENDING主泵重新领走处理
*
* 只有"再试一次有可能成功"的错误类才放行报文本身不合法的MALFORMED重放多少次都一样
* 所以不在白名单里没列出的错误类也一律不放行
*/
@Singleton
class ReplayService(
private val procState: ProcStateRepository,
) {
private val log = org.slf4j.LoggerFactory.getLogger(ReplayService::class.java)
/** 可恢复错误类:codec 修复可重放 / 未实装补齐可重放 / 基础设施抖动可重放 / 重试耗尽人工复核可重放。 */
/** 可以重放的错误类:解码逻辑修好后能过、处理器补齐后能过、基础设施抖动已恢复、以及重试耗尽人工复核认为还能再试的。 */
val replayableErrorClasses: Set<ErrorClass> =
setOf(ErrorClass.CODEC_ERROR, ErrorClass.UNSUPPORTED, ErrorClass.INFRA, ErrorClass.EXHAUSTED)
/** 只放白名单内的类;请求 MALFORMED 等非法类时静默忽略该类。 */
/** 只放白名单内的错误类;请求里带了 MALFORMED 这类不可重放的,直接忽略不报错。 */
fun replay(requested: Collection<ErrorClass>): Int {
val allowed = requested.filter { it in replayableErrorClasses }
if (allowed.isEmpty()) {
@@ -29,6 +31,6 @@ class ReplayService(
return n
}
/** 默认入口:重放全部可恢复类。 */
/** 什么都不指定时走这个:把白名单里的错误类全部重放一遍。 */
fun replayAll(): Int = replay(replayableErrorClasses)
}
@@ -5,8 +5,10 @@ import io.micronaut.context.annotation.Requires
import jakarta.inject.Singleton
/**
* stub 适配层DeliveryPort 内存实现仅在 msgx.stubs=true 生效
* 记录 (topic, key, payload)payload = null 表示 TOMBSTONEflight-state.md §5 键缺失=删除旧值
* 投递端口的内存实现只有配置 msgx.stubs=true 才装配本地开发和测试用
*
* 每次发送都往 sent 里记一条topickeypayloadpayload null 的那条就是删除通知
* tombstone下游按"这个键没了"理解成删除 docs/flight-state.md §5
*/
@Requires(property = "msgx.stubs", value = "true")
@Singleton
@@ -8,8 +8,9 @@ import io.micronaut.http.annotation.Post
import io.micronaut.http.annotation.Produces
/**
* ACMA-8 契约冻结compat HTTP 路径 `POST /cminmsgs/send`手工/对拍
* 生产主路径 = JDBC 轮询共享 CMINMSGSInboxPollerU05非本 Controller
* 兼容用的 HTTP 入口 `POST /cminmsgs/send`手工投递和跟现役系统逐字对拍用
*
* 生产上收报走的是轮询共享信箱表 CMINMSGSInboxPoller不是这个 Controller
*/
@Controller
class InboxController(private val inbox: InboxService) {
@@ -18,8 +19,8 @@ class InboxController(private val inbox: InboxService) {
@Produces(MediaType.TEXT_PLAIN)
fun send(@Body rawXml: String): HttpResponse<String> {
val receipt = inbox.accept(rawXml)
return HttpResponse.ok(receipt.msgId.toString()) // TODO: 现役响应体逐字对拍后固化
return HttpResponse.ok(receipt.msgId.toString()) // TODO: 先返回消息 ID 文本,等跟现役响应体逐字对拍过再定稿
}
// TODO(阶段2): /schd/sync、/all/flights、/kafka/topics/{name}/msgs 按契约冻结清单补齐
// TODO(阶段2): 按契约清单补齐 /schd/sync、/all/flights、/kafka/topics/{name}/msgs 这几个接口
}
@@ -14,13 +14,14 @@ import java.time.Instant
* 收报轮询共享信箱把新消息登记到自有 PG PROC_STATE等着主泵处理
*
* 每一轮只做三件事
* 1. 水位 W 之后按 ID 升序读一批行只看 ID不看处理标记
* 2. 把读到的行登记成 PENDING并把水位推到"连续"的位置
* 1. "水位"之后按 ID 升序读一批行只看 ID不看处理标记
* 水位记的是"信箱里到哪个 ID 为止已经全部读进自有库"存在自己的库里重启不丢
* 2. 把读到的行登记成 PENDING并把水位往前推
* 3. 登记和水位推进写在同一个事务里中途崩溃时水位没动下一轮重扫即可补齐
*
* "连续"是这里唯一需要理解的规则如果 W 后面缺了一个 ID说明可能有 ID 更小的消息
* 还没提交上来此时先停在缺口前不把缺口后面的消息放进队列否则那条迟到的消息
* 会排到它们后面破坏"先来先处理"的约定同一航班的报文可能被乱序应用
* 水位一次只推到"连续"的位置是这里唯一需要理解的规则如果下一个 ID 缺号
* 说明可能有 ID 更小的消息还没提交上来此时先停在缺口前不把缺口后面的消息放进队列
* 否则那条迟到的消息会排到它们后面破坏"先来先处理"的约定同一航班的报文可能被乱序应用
*
* 缺口等超过 `msgx.pipeline.max-commit-delay` 仍未出现就认定它永远不会来了
* 典型情况是自增回滚留下的空位跳过它继续推进否则水位会卡在第一个空位上
@@ -14,27 +14,30 @@ import java.time.Instant
import java.time.ZoneId
/**
* 历史归档与物理清除docs/flight-state.md §6 生命周期顺序不颠倒design.md §6.2
* 1. 选出满足保留期 + 终态/静默判据的航班 DELETED
* 2. 写入历史存储
* 3. 历史存储返回成功的 FLID 集合 物理删除主行明细归档结果只记录在历史存储当前态无 ARCHIVED
* 4. 未经 FDEL 的航班在清除前补发一次删除事件§3.3/§5其余不发
* 5. 失败或不明确的保留重试
* 前提红线历史存储未接通时必须删 0 §6
* 历史归档与物理清除把到期航班搬进历史存储再从当前态删掉五步顺序不颠倒
* 1. 按保留期和终态/静默条件挑出候选航班已删除的也算
* 2. 交给历史存储归档
* 3. 只对历史存储确认归档成功的 FLID航班实例 ID物理删除主行明细当前态不设
* "已归档"状态归档结果只留在历史存储里
* 4. 从没收到过 FDEL航班终止报文被这里直接清掉的航班清除前补发一次删除通知其余的不用发
* 5. 归档失败或者结果说不清的原样留着下次再来
*
* 红线历史存储没接通时必须一条都不删先删当前态事后再补历史是不允许的
* docs/flight-state.md §6
*/
@Singleton
class HistorySweepJob(
private val flightState: FlightStateRepository,
private val msgEvents: MsgEventRepository,
private val props: HistoryProps,
/** 历史存储端口:返回归档成功的 FLID 集合;未接通时不注入(null。 */
/** 历史存储端口:返回确认归档成功的 FLID 集合;部署侧没接通时是 null。 */
private val historyStore: HistoryStore? = null,
/** design.md §6.2 留痕清理端口(90 天,按 (SCOPE_END, RECV_AT));未接通时不注入。 */
/** 留痕清理端口:删掉超过保留期(默认 90 天)的日计划留痕行;没接通时是 null。 */
private val snapLogPurge: SnapshotLogPurge? = null,
) {
/** 历史存储端口(由部署侧适配实现;脚手架默认接通。 */
/** 归档到历史存储的接口,由部署侧适配具体存储;默认接通。 */
fun interface HistoryStore {
/** 归档候选航班返回归档确认成功的 FLID 集合(部分成功允许)。 */
/** 归档这些候选航班返回其中确认成功的 FLID;允许部分成功,没在返回值里的留着重试。 */
fun archive(candidates: List<HistoryCandidate>): Set<String>
}
@@ -46,12 +49,12 @@ class HistorySweepJob(
fun run(now: Instant = Instant.now()): SweepOutcome {
if (!props.historyStoreEnabled || historyStore == null) {
// 红线:历史存储接通必须删 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") // §6窗口按机场时区
val zone = ZoneId.of("Asia/Shanghai") // 保留期窗口按机场时区算,不用 UTC
val candidates = flightState.findHistoryCandidates(rules, zone, now)
if (candidates.isEmpty()) return SweepOutcome(0, 0, 0)
@@ -59,7 +62,7 @@ class HistorySweepJob(
if (archivedFlids.isEmpty()) return SweepOutcome(candidates.size, archived = 0, purged = 0)
val toPurge = candidates.filter { it.flid in archivedFlids }
// §6/§3.3:未经 FDEL、由生命周期直接清的航班,除前补发一次删除事件
// 从没收到过 FDEL、被这里直接清的航班,除前补发一次通知,否则下游一直以为它还在
val preDelete = toPurge.filter { it.wasNeverFdel }
if (preDelete.isNotEmpty()) {
msgEvents.insertAll(preDelete.map { tombstone(it) })
@@ -176,7 +176,10 @@ class AdftProcessor(
ApplyResult.Succeeded
}
/** §3.3 保守语义:仅出现字段 Set;集合出现 Replace、缺失保留。 */
/**
* 上游对"字段缺失"的含义还没确认这里取保守做法
* 报文里出现的字段覆盖本地值没出现的字段保持原值不动
*/
private fun setOnly(record: ScheduleRecord) = MergeChange(
flid = record.flid,
scalars = record.scalars,
@@ -5,9 +5,12 @@ import com.gzzn.omms.msgexchange.domain.DecodedMessage
import java.time.LocalDate
/**
* ACMA-8 I3幂等键 = SNDR|TYPE|STYP|SEQN
* 计算集中此处唯一入口是否含日边界可配置且默认关闭CONFIRM 矩阵 #11
* SEQN 重置作用域确认前不改语义上线后不改幂等键
* 幂等键的唯一算法发送方报文类型子类型流水号四段用竖线拼起来用来判断同一条业务
* 报文是不是已经处理过
*
* 也可以再拼上日期做"按天去重" msgx.identity.include-day-boundary 控制默认关闭
* 发送方的流水号到底按什么范围重置还没确认语义定下来之前不动幂等键上线后就不能再改
* 改了会让已有的去重记录全部对不上
*/
object Identity {
fun of(
@@ -8,9 +8,11 @@ import java.sql.Connection
import javax.sql.DataSource
/**
* U03N02/R09/R01a 验收的启动级版本基础设施配置在真实上下文启动时绑定生效
* 注入 DataSource 并建立真实连接H2 内存application-test.yml证明 datasources.default.*
* 键位与驱动解析正确而非看似配置实则未生效Pump 等业务 bean 懒加载不依赖实仓储
* 守着启动期配置真的生效在完整应用上下文里注入 DataSource 并真连一次H2 内存库配在
* application-test.yml证明 datasources.default.* 的键位和驱动都能解析而不是
* "看着配了、其实没生效"
*
* 业务 bean比如主泵是懒加载的所以这里不需要真实仓储
*/
@MicronautTest
class InfraBindingStartupTest {
@@ -33,7 +35,7 @@ class InfraBindingStartupTest {
@Test
fun `registration disabled by test profile`() {
assertNotNull(props)
// eureka 注册在 application-test.yml 经 msgx.register-eureka=false 关闭——启动级验证绑定可达
// 测试配置里关掉了 eureka 注册(msgx.register-eureka=false),这里验证上下文照样能起来
assertNotNull(dataSource)
}
}
@@ -11,10 +11,11 @@ import org.junit.jupiter.api.TestInstance
import java.time.Duration
/**
* U03N02验收msgx.* 嵌套键改值回读离线化
* 三个嵌套类Pipeline/Schd/Identity须标注 @ConfigurationProperties 才会被绑定
* 未标注时下列覆盖值会静默回落 Kotlin 默认值本测试即红仅注入配置 bean不触碰仓储/DB
* Micronaut Test 5.x 移除 @Property 注解改以 TestPropertyProvider 注入测试属性 PER_CLASS
* 守着配置绑定真的生效msgx.* 下面每个嵌套配置类都要标 @ConfigurationProperties否则键不会
* 被绑定测试里给的覆盖值会静默失效用回代码里的默认值这个测试就会红
*
* 只注入配置 bean不碰仓储和数据库另外注意 Micronaut Test 5.x 去掉了 @Property 注解
* 改由 TestPropertyProvider 提供测试属性所以类上要加 PER_CLASS
*/
@MicronautTest
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
@@ -41,6 +42,6 @@ class PipelinePropsBindingTest : TestPropertyProvider {
@Test
fun `top-level scalars still bind`() {
assertEquals("msgexchangeapi", props.serviceName)
assertFalse(props.registerEureka) // application-test.yml false
assertFalse(props.registerEureka) // application-test.yml 里配成 false
}
}
@@ -4,7 +4,8 @@ import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Test
/**
* U08/N28backoffFor 必须对 attempt 0FAILED 行未递增 attempts给出首档退避而非抛异常
* 守着退避表的规矩按档位递增超出表长就封顶传入 0 或负数比如还没记过重试次数的
* FAILED 要返回第一档不能抛异常
*/
class PipelinePropsTest {
@@ -16,9 +16,11 @@ import org.junit.jupiter.api.Test
import java.time.Instant
/**
* 投递调度docs/flight-state.md §5 + design.md §5.1/§5.2
* KAFKA_MSG 逐条 FIFOKAFKA_SCHD 唯一出口 flushSchd FLID 未发事件按最新
* STATE_VERSION 合并TOMBSTONE null 值消息失败退避重试达上限 DEAD
* 守着投递的几条规矩
* - KAFKA_MSG 一条条按顺序发而且不会顺手把 KAFKA_SCHD 的事件发掉
* - KAFKA_SCHD 只有 flushSchd 一个出口同一个 FLID航班实例 ID只发版本号最新的那条
* - 删除通知发成 value 为空的 tombstone 消息
* - 发送失败按退避重试次数用尽转 DEAD 当死信
*/
class DispatcherTickTest {
@@ -45,7 +47,7 @@ class DispatcherTickTest {
d.tick()
assertEquals(1, port.sent.count { it.topic == "msg" })
assertEquals(0, port.sent.count { it.topic == "schd" }) // N03schd 唯一出口 flushSchd
assertEquals(0, port.sent.count { it.topic == "schd" }) // KAFKA_SCHD 只能由 flushSchd 发,tick 不许碰
}
@Test
@@ -84,7 +86,7 @@ class DispatcherTickTest {
d.flushSchd()
assertEquals(1, port.tombstones.size)
assertEquals("F1", port.tombstones.single().key) // §5:整态键缺失表示删除旧值
assertEquals("F1", port.tombstones.single().key) // value 为空表示删掉这个 key 的旧值
assertNull(port.tombstones.single().payload)
}
@@ -117,7 +119,7 @@ class DispatcherTickTest {
}
}
/** schd 发送失败的端口(退避/终态闭环验证用)。 */
/** KAFKA_SCHD 发送永远失败的投递端口,用来验证退避重试和最终转死信的闭环。 */
private class FailingSchdPort : DeliveryPort {
var calls = 0
override fun sendKafka(topic: String, key: String, payloadJson: String) = Unit
@@ -8,7 +8,8 @@ import java.time.LocalDate
import java.time.ZoneId
/**
* 运营日计算docs/flight-state.md §2.1SODTddMMMyyHHmm+ 机场时区 + 业务切日边界
* 守着运营日的算法用计划运行时间 SODTddMMMyyHHmm加机场时区算本地时刻早于切日边界的
* 归到前一个运营日越界的切日边界收敛到 23SODT 缺失或格式不对就返回 null
*/
class OperationDayTest {
@@ -10,10 +10,11 @@ import java.time.LocalDate
import java.time.ZoneId
/**
* 航班状态引擎docs/flight-state.md §2§4不变量
* 日计划合并§3.1出现 Set/Replace缺失保留标量空值显式清空
* DELETED 不被日计划恢复§3.3FLOP/ADFT 增量合并§3.2/§3.3
* 整包校验§4 步骤 2声明数量航班标识运营日推导失败整包拒绝
* 守着航班状态合并的几条规矩
* - 日计划合并报文里出现的字段覆盖本地值没出现的保留标量给空串等于显式清空
* - 日计划不能把已删除的航班恢复成在用
* - FLOP/ADFT 只改报文表达的字段和集合
* - 整包校验条数FLID 格式运营日能不能算出来有一条不过就整包拒收
*/
class FlightStateEngineTest {
@@ -34,11 +35,11 @@ class FlightStateEngineTest {
keepDeleted = false,
)
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("CA002", next.scalars["FLNO"]) // 报文里出现的字段覆盖本地值
assertEquals("old-note", next.scalars["REMC"]) // 报文没提的标量保留原值
assertEquals("", next.scalars["CNCL"]) // 空串表示显式清空落库时置成 NULL
assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"]) // 报文没提的集合保留原值
assertEquals(6, next.stateVersion) // 每成功合并一次,版本号加一
assertEquals(FlightState.ACTIVE, next.state)
}
@@ -76,8 +77,8 @@ class FlightStateEngineTest {
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
assertFalse(next.scalars.containsKey("REMC")) // 新航班没有历史可保留
assertEquals(10, next.collections.size) // 新航班按全部集合键输出,没内容的给空集
assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"])
assertEquals(emptyList<Map<String, String>>(), next.collections["DELY"])
}
@@ -88,7 +89,7 @@ class FlightStateEngineTest {
val next = FlightStateEngine.snapshotState(
current, ScheduleRecord("121", mapOf("FLNO" to "CA001")), day, keepDeleted = true,
)
assertEquals(FlightState.DELETED, next.state) // §3.3:日计划不恢复 DELETED
assertEquals(FlightState.DELETED, next.state) // 日计划不能把已删除的航班恢复成在用
assertEquals(4, next.stateVersion)
}
@@ -105,7 +106,7 @@ class FlightStateEngineTest {
)
assertEquals("15DEC261900", next.scalars["ESTT"])
assertEquals("CA001", next.scalars["FLNO"]) // 缺失 = 保留(§3.2
assertEquals("CA001", next.scalars["FLNO"]) // 增量里没提的字段保留原值
assertEquals(listOf(mapOf("GATE" to "G1")), next.collections["GTDT"])
assertEquals(2, next.stateVersion)
}
@@ -153,6 +154,6 @@ class FlightStateEngineTest {
val tooLong = FlightStateEngine.validateMessage(
1, listOf(ScheduleRecord("1234567890123", mapOf("SODT" to "15DEC261723"))), opDay,
)
assertTrue(tooLong is SnapshotValidation.Invalid) // FLID 数字型 Number(1-12)
assertTrue(tooLong is SnapshotValidation.Invalid) // FLID 只能是 112 位纯数字
}
}
@@ -12,9 +12,10 @@ import org.junit.jupiter.api.Test
import java.time.Instant
/**
* 历史归档与物理清除docs/flight-state.md §6顺序不可颠倒design.md §6.2
* 历史存储接通必须删 0 先归档确认再物理清除
* 未经 FDEL 的航班在清除前补发删除事件
* 守着归档清理的三条规矩
* - 历史存储接通就一条都不删
* - 先拿到归档确认再物理删除没确认的留着下次重试
* - 从没收到过 FDEL航班终止报文就被清掉的航班删除前要补发一次删除通知
*/
class HistorySweepJobTest {
@@ -53,7 +54,7 @@ class HistorySweepJobTest {
val outcome = job.run(now)
assertEquals(0, outcome.purged) // §6 红线:历史存储接通必须删 0 条
assertEquals(0, outcome.purged) // 历史存储接通,一条都不许删
assertTrue(f.findMainRow("F1") != null)
assertEquals(0, events.rows.size)
}
@@ -78,14 +79,14 @@ class HistorySweepJobTest {
assertEquals(2, outcome.selected)
assertEquals(1, outcome.archived)
assertEquals(1, outcome.purged)
assertEquals(null, f.findMainRow("F1")) // §6归档确认成功物理删除
assertTrue(f.findMainRow("F2") != null) // 失败或不明确的保留重试(步骤 5
assertEquals(0, events.rows.size) // DELETED 航班清除不再发业务删除事件
assertEquals(null, f.findMainRow("F1")) // 归档确认成功的才物理删除
assertTrue(f.findMainRow("F2") != null) // 归档没确认的留着下次再试
assertEquals(0, events.rows.size) // 早就标记删除的航班,清理时不用再发删除通知
}
@Test
fun `never-fdel lifecycle purge emits tombstone before deletion`() {
val f = seededFlight("F3", deleted = false, idleDays = 30) // ACTIVE 且超兜底窗(§6 静默判据)
val f = seededFlight("F3", deleted = false, idleDays = 30) // 还在用,但已经静默超过兜底期限,命中清理条件
val events = StubMsgEvents()
val store = RecordingHistoryStore()
val job = HistorySweepJob(
@@ -94,7 +95,7 @@ class HistorySweepJobTest {
job.run(now)
// §3.3/§5/§6:未经 FDEL、生命周期直接清的航班,除前补发一次删除事件
// 从没收到过 FDEL、生命周期直接清的航班,除前补发一次删除通知
val tombstone = events.rows.values.single { it.eventType == EventType.TOMBSTONE }
assertEquals("F3", tombstone.partitionKey)
assertTrue(tombstone.payloadJson.contains("\"deleted\":true"))
@@ -9,7 +9,8 @@ import java.time.LocalDate
import kotlin.test.assertEquals
/**
* ACMA-8 I3identity = SNDR|TYPE|STYP|SEQN日边界含否集中可配CONFIRM 矩阵 #11默认关
* 守着幂等键的规矩发送方类型子类型流水号四段用竖线拼起来只有配置打开时才追加
* 日期段用来按天去重默认不追加
*/
class IdentityTest {
@@ -3,13 +3,13 @@ package com.gzzn.omms.msgexchange.support
import org.testcontainers.containers.PostgreSQLContainer
/**
* JDBC 集成测试 PostgreSQL 连接解析ACM2-29 P3-A
*
* 优先级
* 1. 显式环境变量 `MSGX_PG_URL` / `MSGX_PG_HOST` / `MSGX_PG_PORT` / `MSGX_PG_NAME`外接库/CI 固定库
* 2. Testcontainers 自动起 `postgres:17-alpine` 隔离容器未配置环境变量且 docker 可用时
* JVM 单例首个用例启动后全程复用容器随 JVM 退出由 Ryuk 清理
* 3. 都不可用 canConnect()=false用例 assumeTrue 跳过绝不误报通过
* JDBC 集成测试找一个能用的 PostgreSQL按下面的顺序挑
* 1. 环境变量显式指定MSGX_PG_URL / MSGX_PG_HOST / MSGX_PG_PORT / MSGX_PG_NAME
* 适合外接库或 CI 上的固定库
* 2. 没配环境变量 docker 可用时 Testcontainers 起一个 postgres:17-alpine 容器
* JVM 内只起一个全程复用进程退出后由 Ryuk 清掉
* 3. 两条路都走不通时 canConnect() 返回 false用例里用 assumeTrue 跳过
* 绝不因为连不上就当成通过
*
* | 变量 | 默认 |
* |---|---|