2026-09-07 15:11:33 +08:00
|
|
|
|
package com.gzzn.omms.msgexchange.processing
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-07 15:11:33 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.domain.DecodedMessage
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.ErrorClass
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.ProcState
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.ProcStatus
|
2026-09-07 16:12:00 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.domain.MsgEvent
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.MsgKind
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.domain.Targets
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository
|
2026-09-07 17:04:57 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.FlightFields
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.FlightFieldsJson
|
2026-09-07 16:12:00 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.FlightSchdRepository
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository
|
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.PipelineTransactionManager
|
2026-09-07 15:11:33 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository
|
2026-09-07 16:12:00 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.ReqTrackRepository
|
2026-09-07 15:11:33 +08:00
|
|
|
|
import com.gzzn.omms.msgexchange.infra.retry.ProcFailure
|
2026-09-06 16:13:50 +08:00
|
|
|
|
import jakarta.inject.Singleton
|
|
|
|
|
|
|
|
|
|
|
|
/**
|
2026-09-07 16:12:00 +08:00
|
|
|
|
* staging(内存瞬态,崩溃从 raw 整包重放)→ PG 单事务原子“覆盖+按代差删+SQL CAS 推进+事件入队+SUCCEEDED”(I4/I5,ACM2-28 定案)。
|
2026-09-06 19:08:02 +08:00
|
|
|
|
* U10/T07 占位安全化 + U08 统一失败迁移(ProcFailure):staging 未实装 → FAILED(UNSUPPORTED)+退避
|
|
|
|
|
|
* (可重放,绝不写终态);CAS 冲突 → FAILED(INFRA)+退避;达上限统一 DEAD(EXHAUSTED)。
|
2026-09-06 16:13:50 +08:00
|
|
|
|
*/
|
|
|
|
|
|
@Singleton
|
|
|
|
|
|
class SnapshotFlow(
|
|
|
|
|
|
private val procState: ProcStateRepository,
|
2026-09-07 16:12:00 +08:00
|
|
|
|
private val flightSchd: FlightSchdRepository,
|
|
|
|
|
|
private val msgEvents: MsgEventRepository,
|
|
|
|
|
|
private val reqTrack: ReqTrackRepository,
|
2026-09-06 19:08:02 +08:00
|
|
|
|
private val procFailure: ProcFailure,
|
2026-09-07 16:12:00 +08:00
|
|
|
|
private val txManager: PipelineTransactionManager,
|
|
|
|
|
|
private val inbox: CminmsgInboxRepository,
|
2026-09-06 16:13:50 +08:00
|
|
|
|
) {
|
2026-09-06 21:02:27 +08:00
|
|
|
|
private val log = org.slf4j.LoggerFactory.getLogger(SnapshotFlow::class.java)
|
|
|
|
|
|
|
2026-09-06 18:19:04 +08:00
|
|
|
|
fun publishSnapshot(head: ProcState, msg: DecodedMessage) {
|
2026-09-07 16:12:00 +08:00
|
|
|
|
// 1) staging:流式解析 + 整包校验(内存瞬态,崩溃从 raw 整包重放)
|
|
|
|
|
|
val staged = StageResult.stagingOf(msg)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
if (staged is StageResult.Invalid) {
|
2026-09-06 21:02:27 +08:00
|
|
|
|
log.warn("staging not implemented -> FAILED(UNSUPPORTED) id={} reason={}", head.cminmsgsId, staged.reason)
|
2026-09-06 19:08:02 +08:00
|
|
|
|
procFailure.fail(head, ErrorClass.UNSUPPORTED, staged.reason) // U10:未实装 → 可重放,非终态
|
2026-09-06 16:13:50 +08:00
|
|
|
|
return
|
|
|
|
|
|
}
|
2026-09-07 16:12:00 +08:00
|
|
|
|
val ok = staged as StageResult.Ok
|
2026-09-07 17:04:57 +08:00
|
|
|
|
val normalized = ok.flights // (flid, fields)
|
2026-09-07 16:12:00 +08:00
|
|
|
|
val day = ok.day
|
2026-09-06 16:13:50 +08:00
|
|
|
|
|
2026-09-07 16:12:00 +08:00
|
|
|
|
// 内存与超大包熔断防御(ACM2-28 评论 4 项 3)
|
|
|
|
|
|
if (normalized.size > MAX_FLIGHTS_PER_SNAPSHOT) {
|
|
|
|
|
|
log.error("snapshot flights exceed limit -> DEAD(MALFORMED) id={} size={}", head.cminmsgsId, normalized.size)
|
|
|
|
|
|
procState.update(head.cminmsgsId, ProcStatus.DEAD, errorClass = ErrorClass.MALFORMED,
|
|
|
|
|
|
lastError = "flights-exceed-limit:${normalized.size}")
|
|
|
|
|
|
return
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// 2) 准备代元数据与差集(删除集 = gen.flids − 新代,ADFT 自动存活)
|
|
|
|
|
|
val gen = flightSchd.getGen(day)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
val expected = gen?.version ?: 0L
|
2026-09-07 16:12:00 +08:00
|
|
|
|
val newFlids = normalized.map { it.first }.toSet()
|
|
|
|
|
|
val delFields = gen?.flids?.minus(newFlids) ?: emptySet()
|
|
|
|
|
|
val newVersion = expected + 1L
|
|
|
|
|
|
|
|
|
|
|
|
// 3) 自有 PG 单事务原子提交(ACM2-28 定案:覆盖新代 + 域内差删 + SQL CAS + 事件 + SUCCEEDED)
|
|
|
|
|
|
try {
|
|
|
|
|
|
txManager.inTransaction {
|
|
|
|
|
|
// 覆盖新代全量:强行声明 FDAY 归属(JDBC batch 批处理)
|
|
|
|
|
|
flightSchd.upsertSnapshotBatch(day, normalized)
|
|
|
|
|
|
|
|
|
|
|
|
// 按代差删域化:仅删除 FDAY = day 且在 delFields 中的记录(ADFT 与跨代已迁移行存活)
|
|
|
|
|
|
if (delFields.isNotEmpty()) {
|
|
|
|
|
|
flightSchd.deleteDiffByDay(day, delFields)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// SQL CAS 版本推进(防双写断言)
|
|
|
|
|
|
val casSuccess = flightSchd.putGenIfVersion(
|
|
|
|
|
|
day = day,
|
|
|
|
|
|
expected = expected,
|
|
|
|
|
|
newGen = FlightSchdRepository.GenMeta(fday = day, version = newVersion, flids = newFlids),
|
|
|
|
|
|
)
|
|
|
|
|
|
if (!casSuccess) {
|
|
|
|
|
|
// 重放路径:version 已是目标值 → no-op 视为成功,否则为 CAS 冲突
|
|
|
|
|
|
val again = flightSchd.getGen(day)
|
|
|
|
|
|
if (again == null || again.version != newVersion) {
|
|
|
|
|
|
throw CasConflictException("gen-cas-conflict: expected=$expected current=${again?.version}")
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// RESP 匹配 REQ_TRACK -> DONE
|
|
|
|
|
|
val schdKind = msg.kind as? MsgKind.Schd
|
|
|
|
|
|
if (schdKind?.subtype == MsgKind.SchdSubtype.RESP) {
|
|
|
|
|
|
reqTrack.findOpenByKind("SCHD")?.let { req ->
|
|
|
|
|
|
reqTrack.markDone(req.reqId)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-07 17:04:57 +08:00
|
|
|
|
// 构造并批量写入 schd 投递事件(载荷 = 字段集序列化,KAFKA_SCHD 线格式)
|
|
|
|
|
|
val events = normalized.map { (flid, fields) ->
|
|
|
|
|
|
MsgEvent(target = Targets.KAFKA_SCHD, partitionKey = flid, payloadJson = FlightFieldsJson.toJson(fields))
|
2026-09-07 16:12:00 +08:00
|
|
|
|
}
|
|
|
|
|
|
if (events.isNotEmpty()) {
|
|
|
|
|
|
msgEvents.insertAll(events)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// 终态置 SUCCEEDED
|
|
|
|
|
|
procState.update(head.cminmsgsId, ProcStatus.SUCCEEDED)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
2026-09-07 16:12:00 +08:00
|
|
|
|
|
|
|
|
|
|
// 提交后:backfill CMINMSGS(外部副作用,补偿保障)
|
|
|
|
|
|
inbox.backfillOnSuccess(head.cminmsgsId, msg.meta.sndr, msg.meta.type, msg.meta.styp, msg.meta.seqn)
|
|
|
|
|
|
log.info("snapshot SUCCEEDED id={} day={} flights={}", head.cminmsgsId, day, normalized.size)
|
|
|
|
|
|
} catch (e: CasConflictException) {
|
|
|
|
|
|
log.warn("gen CAS conflict -> FAILED(INFRA) id={} msg={}", head.cminmsgsId, e.message)
|
|
|
|
|
|
procFailure.fail(head, ErrorClass.INFRA, e.message ?: "gen-cas-conflict")
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
2026-09-06 18:19:04 +08:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-09-06 16:13:50 +08:00
|
|
|
|
/** staging 结果(骨架)。 */
|
|
|
|
|
|
sealed interface StageResult {
|
2026-09-07 17:04:57 +08:00
|
|
|
|
data class Ok(val day: String, val flights: List<Pair<String, FlightFields>>) : StageResult
|
2026-09-06 16:13:50 +08:00
|
|
|
|
data class Invalid(val reason: String) : StageResult
|
|
|
|
|
|
|
|
|
|
|
|
companion object {
|
2026-09-07 16:12:00 +08:00
|
|
|
|
var parser: ((DecodedMessage) -> StageResult)? = null
|
|
|
|
|
|
|
2026-09-06 16:13:50 +08:00
|
|
|
|
fun stagingOf(msg: DecodedMessage): StageResult =
|
2026-09-07 16:12:00 +08:00
|
|
|
|
parser?.invoke(msg) ?: Invalid("staging-not-implemented(${msg.typeTag})") // TODO(阶段2)
|
2026-09-06 16:13:50 +08:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-09-07 16:12:00 +08:00
|
|
|
|
class CasConflictException(message: String) : RuntimeException(message)
|
|
|
|
|
|
const val MAX_FLIGHTS_PER_SNAPSHOT = 10000
|