refactor(flight-state): SCHD 日计划收敛为缺失保留合并语义并对齐现行设计文档
按现行 docs/flight-state.md(111 行版,§1-§7)全面对齐 domain 包及其消费方, 删除整套指向已退役长版文档(§5.x-§10)的引用与死代码: - 语义:FlightStateEngine.snapshotState 由"完整快照整体替换"改为 §3.1 合并语义 ——出现 Set/Replace、缺失保留、标量空串显式清空;DELETED 不被日计划恢复(§3.3)。 - 校验:validateMessage 移除从未接线的报文覆盖范围(scope)参数,只保留 §4 步骤 2 的声明数量/航班标识/运营日推导校验;ScheduleBody 删除 scopeStart/scopeEnd。 - 删除死代码:domain/Decision.kt、FlightModel 的 SnapshotPatch/UpsertOutcome、 ProcState.isTerminal/archivable、SnapshotFlag.SEQN_REGRESSION、 ScheduleRecord.seqn(及 wire FlightRecordXml.SEQN);codec 移除未消费 FFID。 - 删除未接线且引用已退役列(fday/last_message_id)的 SqlDialect 方言脚手架, oracle11g README 改为按 V1 现列重建的口径。 - 注释/测试:domain、processing、infra 仓储与 jobs/delivery、配置类及对应测试的 KDoc 章节引用全部对齐现行 flight-state.md/design.md;FlightStateEngineTest 重写为合并语义(62/62 通过)。 V1__flight_state_baseline.sql 保留原样(内容注释仍带旧章节号,改动会破坏 已应用迁移的 Flyway checksum,待重建基线或 V2 净迁移时收敛)。
This commit is contained in:
@@ -10,8 +10,10 @@ import com.gzzn.omms.msgexchange.domain.OperationDayCalculator
|
||||
import com.gzzn.omms.msgexchange.domain.ProcState
|
||||
import com.gzzn.omms.msgexchange.domain.Targets
|
||||
import com.gzzn.omms.msgexchange.domain.flight.FlightSnapshot
|
||||
import com.gzzn.omms.msgexchange.domain.flight.FlightState
|
||||
import com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine
|
||||
import com.gzzn.omms.msgexchange.domain.flight.MergeChange
|
||||
import com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord
|
||||
import com.gzzn.omms.msgexchange.infra.persistence.BackfillTodoRepository
|
||||
import com.gzzn.omms.msgexchange.infra.persistence.FlightStateRepository
|
||||
import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository
|
||||
@@ -22,8 +24,9 @@ import java.time.Instant
|
||||
import java.time.ZoneId
|
||||
|
||||
/**
|
||||
* FLOP(§6.1):读取完整当前态 → 合并变化 → 保留运营日 → STATE_VERSION+1 →
|
||||
* 同事务登记 KAFKA_MSG / KAFKA_SCHD 与处理终态。
|
||||
* FLOP(docs/flight-state.md §3.2 动态运行事件):读取完整当前态 → 合并变化 →
|
||||
* 保留运营日 → STATE_VERSION+1 → 同事务登记 KAFKA_MSG / KAFKA_SCHD 与处理终态。
|
||||
* 未知/迟到航班按幂等成功处理,不创建实例(创建入口只有 SCHD/ADFT)。
|
||||
*/
|
||||
@Singleton
|
||||
class FlopProcessor(
|
||||
@@ -49,8 +52,8 @@ class FlopProcessor(
|
||||
}
|
||||
|
||||
/**
|
||||
* FDEL(§6.2):ACTIVE → 置 DELETED、推进版本、明细保留、发布 tombstone;
|
||||
* 已 DELETED / 不存在 → 幂等成功,不推进版本、不重复发布。
|
||||
* FDEL(docs/flight-state.md §3.3 删除):ACTIVE → 置 DELETED、推进版本、明细保留、
|
||||
* 与删除同事务登记 tombstone(§5);已 DELETED / 不存在 → 幂等成功,不推进版本、不重复发布。
|
||||
*/
|
||||
@Singleton
|
||||
class FdelProcessor(
|
||||
@@ -66,7 +69,7 @@ class FdelProcessor(
|
||||
val deleted = flightState.markDeleted(payload.flid, msgId = head.msgId, now = Instant.now())
|
||||
if (deleted) {
|
||||
val current = flightState.loadFullSnapshot(payload.flid)
|
||||
// tombstone 仅在 ACTIVE→DELETED 时登记(§7.3),与删除同事务
|
||||
// tombstone 仅在 ACTIVE→DELETED 时登记(§3.3/§5),与删除同事务
|
||||
msgEvents.insertAll(
|
||||
listOf(
|
||||
MsgEvent(
|
||||
@@ -94,14 +97,14 @@ class FdelProcessor(
|
||||
)
|
||||
preRegisterBackfill(head, msg, backfillTodo)
|
||||
}
|
||||
ApplyResult.Succeeded // 未命中 = 迟到/重复,幂等成功(§6.2 步骤 3/4)
|
||||
ApplyResult.Succeeded // 未命中 = 迟到/重复,幂等成功(§3.3)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* ADFT(§6.3 + §2.1):字段缺失语义待确认——确认前按保守 Set-only 处理
|
||||
* (出现字段覆盖,缺失不 Clear,不沿用 FLOP 全量合并规则)。
|
||||
* FLID 已存在且 DELETED → 生命周期重激活;不存在 → 新实例建立(含运营日计算 §3.5)。
|
||||
* ADFT(docs/flight-state.md §3.3):字段缺失语义待上游确认——确认前按保守 Set-only
|
||||
* 处理(出现字段覆盖、缺失不清空)。FLID 已存在且 DELETED → 生命周期重激活;不存在 →
|
||||
* 新实例建立(含 SODT 时直接计算 OPERATION_DAY,§2.1;否则保留 NULL 待日计划收录)。
|
||||
*/
|
||||
@Singleton
|
||||
class AdftProcessor(
|
||||
@@ -118,12 +121,12 @@ class AdftProcessor(
|
||||
cutoffHour = operationDayProps.cutoffHour,
|
||||
)
|
||||
|
||||
fun apply(head: ProcState, msg: DecodedMessage, record: com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord): ApplyResult =
|
||||
fun apply(head: ProcState, msg: DecodedMessage, record: ScheduleRecord): ApplyResult =
|
||||
txManager.inTransaction {
|
||||
lock.lock()
|
||||
val main = flightState.findMainRow(record.flid)
|
||||
if (main != null && main.state == com.gzzn.omms.msgexchange.domain.flight.FlightState.DELETED) {
|
||||
// §6.3 重激活:DELETED → ACTIVE,推进版本,登记状态事件
|
||||
if (main != null && main.state == FlightState.DELETED) {
|
||||
// §3.3 重激活:DELETED → ACTIVE,推进版本,登记状态事件
|
||||
if (flightState.revive(record.flid, msgId = head.msgId, now = Instant.now())) {
|
||||
val current = flightState.loadFullSnapshot(record.flid)
|
||||
if (current != null) {
|
||||
@@ -138,15 +141,15 @@ class AdftProcessor(
|
||||
|
||||
val current = flightState.loadFullSnapshot(record.flid)
|
||||
val next: FlightSnapshot = if (current == null) {
|
||||
// 新实例建立:ADFT 含 SODT 时直接计算运营日(§2.1),不可算则置 null 待快照收录
|
||||
// 新实例建立:ADFT 含 SODT 时直接计算运营日(§2.1),不可算则置 null 待日计划收录
|
||||
val day = opDay.compute(record.scalars["SODT"])
|
||||
FlightSnapshot(
|
||||
flid = record.flid,
|
||||
operationDay = day, // 待确认项 §2.1:不可算时不得默认写接收日
|
||||
state = com.gzzn.omms.msgexchange.domain.flight.FlightState.ACTIVE,
|
||||
operationDay = day, // §2.1:不可算时不得默认写接收日
|
||||
state = FlightState.ACTIVE,
|
||||
stateVersion = 1L,
|
||||
scalars = record.scalars,
|
||||
collections = com.gzzn.omms.msgexchange.domain.flight.FlightStateEngine.COLLECTION_KEYS.associateWith { key ->
|
||||
collections = FlightStateEngine.COLLECTION_KEYS.associateWith { key ->
|
||||
record.collections[key] ?: emptyList()
|
||||
},
|
||||
)
|
||||
@@ -162,8 +165,8 @@ class AdftProcessor(
|
||||
ApplyResult.Succeeded
|
||||
}
|
||||
|
||||
/** §2.1 保守语义:仅出现字段 Set;集合出现 Replace、缺失保留。 */
|
||||
private fun setOnly(record: com.gzzn.omms.msgexchange.domain.flight.ScheduleRecord) = MergeChange(
|
||||
/** §3.3 保守语义:仅出现字段 Set;集合出现 Replace、缺失保留。 */
|
||||
private fun setOnly(record: ScheduleRecord) = MergeChange(
|
||||
flid = record.flid,
|
||||
scalars = record.scalars,
|
||||
collections = record.collections,
|
||||
@@ -174,10 +177,10 @@ class AdftProcessor(
|
||||
// 共享小工具(处理器层私有约定)
|
||||
// =====================================================================
|
||||
|
||||
/** 航班不存在/迟到:幂等成功(§9 队头不阻塞;不创建实例——创建入口只有 SCHD/ADFT)。 */
|
||||
/** 航班不存在/迟到:幂等成功(§3.2/§3.3;不阻塞队头,不创建实例——创建入口只有 SCHD/ADFT)。 */
|
||||
private fun idempotentAbsent(head: ProcState, msg: DecodedMessage): ApplyResult = ApplyResult.Succeeded
|
||||
|
||||
/** KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(§7.3)。 */
|
||||
/** KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1/§5)。 */
|
||||
internal fun eventsFor(next: FlightSnapshot, mapper: ObjectMapper): List<MsgEvent> {
|
||||
val payload = linkedMapOf<String, Any>(
|
||||
"flid" to next.flid,
|
||||
@@ -201,7 +204,7 @@ internal fun eventsFor(next: FlightSnapshot, mapper: ObjectMapper): List<MsgEven
|
||||
)
|
||||
}
|
||||
|
||||
/** §7.2 业务事务内预登记回填待办(与终态同事务)。 */
|
||||
/** 回填待办与业务终态同事务预登记(design.md §3.3/§6.1;提交后回填并删待办)。 */
|
||||
internal fun preRegisterBackfill(head: ProcState, msg: DecodedMessage, backfillTodo: BackfillTodoRepository?) {
|
||||
backfillTodo?.record(
|
||||
BackfillTodoRepository.BackfillTask(
|
||||
@@ -211,4 +214,3 @@ internal fun preRegisterBackfill(head: ProcState, msg: DecodedMessage, backfillT
|
||||
null,
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -19,8 +19,9 @@ import java.time.Duration
|
||||
import java.time.Instant
|
||||
|
||||
/**
|
||||
* 处理主泵(docs/flight-state.md §1 目标 2):单活动主泵严格 FIFO + HOL + 毒丸 DEAD 升级;
|
||||
* 航班状态、事件、处理终态在处理器事务内原子提交(§1 目标 3)。
|
||||
* 处理主泵(docs/flight-state.md §1 严格有序 + §4 处理事务与失败规则):
|
||||
* 单活动主泵严格 FIFO + HOL 阻塞 + 队头滞留转 DEAD;航班状态、事件、处理终态
|
||||
* 与回填待办在处理器事务内原子提交。
|
||||
*/
|
||||
@Singleton
|
||||
class Pump(
|
||||
@@ -84,7 +85,7 @@ class Pump(
|
||||
/**
|
||||
* processOne:解码 → 绑定 → 处理器(事务内决策+落库)→ 终态迁移 → 回填。
|
||||
* 边界化失败迁移(ProcFailure):任何意外异常归于本条 head,FAILED(INFRA)+退避,不穿出杀泵;
|
||||
* MALFORMED / PROTOCOL 直接 DEAD 不重试(§9)。
|
||||
* MALFORMED / PROTOCOL 直接 DEAD 不重试(docs/design.md §2.3 错误分类)。
|
||||
*/
|
||||
@Singleton
|
||||
class MessageProcessor(
|
||||
@@ -156,7 +157,7 @@ class MessageProcessor(
|
||||
}
|
||||
}
|
||||
|
||||
// 处理器分派:SCHD 快照主链路(§5.1)/ FLOP / FDEL / ADFT;缺载荷按 MALFORMED 终态
|
||||
// 处理器分派:SCHD 日计划主链路(§3.1)/ FLOP / FDEL / ADFT;缺载荷按 MALFORMED 终态
|
||||
val result: ApplyResult = when (val kind = decoded.kind) {
|
||||
is MsgKind.Schd -> {
|
||||
val body = decoded.body as? ScheduleBody
|
||||
@@ -194,7 +195,7 @@ class MessageProcessor(
|
||||
flopProcessor.apply(head, decoded, payload)
|
||||
}
|
||||
is MsgKind.Unsupported -> {
|
||||
// §9:未支持类型 → FAILED(UNSUPPORTED) 退避重试,达阈值转 DEAD;绝不写终态
|
||||
// design.md §2.3:未支持类型 → FAILED(UNSUPPORTED) 退避重试,达阈值转 DEAD;绝不写终态
|
||||
log.warn("unsupported type -> FAILED(UNSUPPORTED) msgId={} tag={}", head.msgId, kind.tag)
|
||||
procFailure.fail(head, ErrorClass.UNSUPPORTED, "no-handler:${kind.tag}")
|
||||
return
|
||||
@@ -205,7 +206,7 @@ class MessageProcessor(
|
||||
is ApplyResult.Succeeded, ApplyResult.ReplaySkipped ->
|
||||
procState.update(head.msgId, ProcStatus.SUCCEEDED)
|
||||
is ApplyResult.DeadProtocol -> {
|
||||
// §9:整包拒绝 DEAD(PROTOCOL),立即释放队头,交人工确认
|
||||
// design.md §2.3:整包拒绝 DEAD(PROTOCOL),立即释放队头,交人工确认
|
||||
log.error("DEAD(PROTOCOL) msgId={} reason={} flags={}", head.msgId, result.reason, result.flags)
|
||||
procState.update(head.msgId, ProcStatus.DEAD, errorClass = ErrorClass.PROTOCOL, lastError = result.reason.take(1000))
|
||||
compensateBackfill(head, decoded) // 拒绝包同样要回填信箱,防止反复轮询
|
||||
@@ -217,7 +218,7 @@ class MessageProcessor(
|
||||
log.info("SUCCEEDED msgId={} kind={}", head.msgId, decoded.typeTag)
|
||||
}
|
||||
|
||||
/** §7.2:提交后回填共享信箱;失败不得把 SUCCEEDED 改回 FAILED,待办已事务内预登记。 */
|
||||
/** design.md §3.3/§6.1:提交后回填共享信箱;失败不得把 SUCCEEDED 改回 FAILED,待办已事务内预登记。 */
|
||||
private fun backfill(head: ProcState, decoded: DecodedMessage) {
|
||||
try {
|
||||
inbox.backfillOnSuccess(head.msgId, decoded.meta.sndr, decoded.meta.type, decoded.meta.styp, decoded.meta.seqn)
|
||||
|
||||
@@ -6,7 +6,6 @@ import com.gzzn.omms.msgexchange.config.OperationDayProps
|
||||
import com.gzzn.omms.msgexchange.domain.DecodedMessage
|
||||
import com.gzzn.omms.msgexchange.domain.MsgEvent
|
||||
import com.gzzn.omms.msgexchange.domain.ProcState
|
||||
import com.gzzn.omms.msgexchange.domain.ProcStatus
|
||||
import com.gzzn.omms.msgexchange.domain.SnapshotFlag
|
||||
import com.gzzn.omms.msgexchange.domain.SnapshotLogEntry
|
||||
import com.gzzn.omms.msgexchange.domain.SnapshotResult
|
||||
@@ -30,24 +29,25 @@ import java.time.Instant
|
||||
import java.time.LocalDate
|
||||
import java.time.ZoneId
|
||||
|
||||
/** 处理器执行结果——终态迁移由 MessageProcessor 统一落库。 */
|
||||
/** 处理器执行结果——终态迁移由 MessageProcessor 统一落库(docs/design.md §2.3)。 */
|
||||
sealed interface ApplyResult {
|
||||
/** 业务成功(含幂等成功)。 */
|
||||
data object Succeeded : ApplyResult
|
||||
|
||||
/** §5.1 步骤 2:MSG_ID 已有成功终态 → 重放,直接记幂等成功。 */
|
||||
/** MSG_ID 已有成功终态 → 重放,直接记幂等成功(docs/flight-state.md §4 重放判定)。 */
|
||||
data object ReplaySkipped : ApplyResult
|
||||
|
||||
/** §9:整包拒绝 DEAD(PROTOCOL),不重试,交人工确认。 */
|
||||
/** 整包拒绝 DEAD(PROTOCOL)(flight-state.md §3.1/§4),不重试,交人工确认。 */
|
||||
data class DeadProtocol(val reason: String, val flags: Set<SnapshotFlag> = emptySet()) : ApplyResult
|
||||
}
|
||||
|
||||
/** §5.3 第四行:归属日不符 = 串日/错发/污染,整包拒绝。 */
|
||||
/** 归属日冲突(flight-state.md §2.1:OPERATION_DAY 一经确定不可变)= 串日/错发/污染,整包拒绝。 */
|
||||
class ProtocolViolation(message: String) : RuntimeException(message)
|
||||
|
||||
/**
|
||||
* SCHD 快照主链路(docs/flight-state.md §5.1 applyScheduleRecords,同一事务):
|
||||
* 对单日快照与滚动窗口统一适用,不做名单层面的处理。
|
||||
* SCHD 日计划主链路(docs/flight-state.md §3.1/§4):对 DNLD 与 RESP 统一适用。
|
||||
* 同一事务内:锁 → 归属日校验 → 逐条合并写完整当前态 → 登记事件与回填待办;
|
||||
* 整包校验失败或运营日冲突整包不落地。
|
||||
*/
|
||||
@Singleton
|
||||
class ScheduleProcessor(
|
||||
@@ -70,18 +70,16 @@ class ScheduleProcessor(
|
||||
val body = msg.body as? ScheduleBody ?: return ApplyResult.DeadProtocol("missing-schd-body")
|
||||
val started = System.nanoTime()
|
||||
|
||||
// ② 重放判定:MSG_ID 已有成功终态 → 幂等成功(§5.1 步骤 2,重放也记留痕)
|
||||
// 重放判定:MSG_ID 已有成功终态 → 幂等成功(flight-state.md §4),重放也记留痕
|
||||
if (procState.findSuccessTerminal(head.msgId)) {
|
||||
logSnapshot(head, body, SnapshotResult.REPLAY_SKIPPED, upserted = 0, flags = emptySet(), started)
|
||||
return ApplyResult.ReplaySkipped
|
||||
}
|
||||
|
||||
// ③ 报文完整性(§5.2 五项):任一失败整包不落地 → DEAD(PROTOCOL)
|
||||
// 整包校验(§4 步骤 2 / design.md §4.1):任一失败整包不落地 → DEAD(PROTOCOL)
|
||||
val validation = FlightStateEngine.validateMessage(
|
||||
recsDeclared = body.recsDeclared,
|
||||
records = body.records,
|
||||
scopeStart = body.scopeStart,
|
||||
scopeEnd = body.scopeEnd,
|
||||
opDay = opDay,
|
||||
)
|
||||
if (validation is SnapshotValidation.Invalid) {
|
||||
@@ -91,18 +89,18 @@ class ScheduleProcessor(
|
||||
val ok = validation as SnapshotValidation.Ok
|
||||
|
||||
if (ok.perRecordDay.isEmpty()) {
|
||||
// 空快照:合法但无写入(§5.5 EMPTY),仍算成功终态
|
||||
// 空快照:合法但无写入,仍算成功终态
|
||||
logSnapshot(head, body, SnapshotResult.COMMITTED, upserted = 0, setOf(SnapshotFlag.EMPTY), started)
|
||||
return ApplyResult.Succeeded
|
||||
}
|
||||
|
||||
val flags = linkedSetOf<SnapshotFlag>()
|
||||
return try {
|
||||
// ①④⑤⑥⑦ 同一事务:锁 → 归属校验 → upsert → 版本/事件/待办预登记/终态
|
||||
// 同一事务(flight-state.md §4):锁 → 归属校验 → 逐条合并写 → 事件/待办预登记
|
||||
val upserted = txManager.inTransaction {
|
||||
lock.lock()
|
||||
|
||||
// ⑤ 归属校验(§5.3):OPERATION_DAY 不可变,批量点查避免逐航班往返
|
||||
// 归属校验(§2.1):OPERATION_DAY 不可变,批量点查避免逐航班往返
|
||||
val mains = flightState.findMainRows(ok.perRecordDay.keys)
|
||||
ok.perRecordDay.forEach { (flid, day) ->
|
||||
val existing = mains[flid] ?: return@forEach
|
||||
@@ -119,7 +117,7 @@ class ScheduleProcessor(
|
||||
val record = body.records.first { it.flid == flid }
|
||||
val existingMain = mains[flid]
|
||||
val keepDeleted = existingMain?.state == FlightState.DELETED
|
||||
if (keepDeleted) flags.add(SnapshotFlag.SCHD_REVIVE_CONFLICT) // §5.1 步骤 6:不恢复
|
||||
if (keepDeleted) flags.add(SnapshotFlag.SCHD_REVIVE_CONFLICT) // §3.3:日计划不复活 DELETED
|
||||
val current = if (existingMain != null) flightState.loadFullSnapshot(flid) else null
|
||||
val next = FlightStateEngine.snapshotState(
|
||||
current = current,
|
||||
@@ -135,7 +133,7 @@ class ScheduleProcessor(
|
||||
events += snapshotEvents(next)
|
||||
}
|
||||
if (events.isNotEmpty()) msgEvents.insertAll(events)
|
||||
// §7.2 目标形态:业务事务内预登记回填待办(提交后由 MessageProcessor 回填并删待办)
|
||||
// 回填待办与业务终态同事务预登记(design.md §3.3/§6.1;提交后由 MessageProcessor 回填并删待办)
|
||||
backfillTodo?.record(
|
||||
BackfillTodoRepository.BackfillTask(
|
||||
msgId = head.msgId, sndr = msg.meta.sndr, type = msg.meta.type,
|
||||
@@ -153,9 +151,8 @@ class ScheduleProcessor(
|
||||
}
|
||||
}
|
||||
|
||||
/** 状态事件:整态出站(§7.3 KAFKA_SCHD)+ 变化通知(KAFKA_MSG)。 */
|
||||
/** 状态事件:KAFKA_SCHD 整态 + KAFKA_MSG 变化通知(flight-state.md §3.1 登记、§5 语义)。 */
|
||||
private fun snapshotEvents(next: FlightSnapshot): List<MsgEvent> {
|
||||
// KAFKA_SCHD 整态:缺失集合输出空集 → 消费者删除旧值(§5.4/§7.3)
|
||||
val payload = linkedMapOf<String, Any>(
|
||||
"flid" to next.flid,
|
||||
"stateVersion" to next.stateVersion,
|
||||
@@ -174,7 +171,8 @@ class ScheduleProcessor(
|
||||
)
|
||||
}
|
||||
|
||||
/** §5.5 留痕:事务外追加,失败只记 error 不阻塞;scope 按记录归属运营日推导。 */
|
||||
/** 留痕 SCHD_SNAP_LOG(design.md §6.2):事务外追加,失败只记 error 不阻塞;
|
||||
* scope 为报文各记录归属运营日的最小/最大(单日快照两者相等)。 */
|
||||
private fun logSnapshot(
|
||||
head: ProcState,
|
||||
body: ScheduleBody,
|
||||
|
||||
Reference in New Issue
Block a user