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
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.persistence.RefDataRepository
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.redis.FlightRedisClient
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.redis.RedisScript
|
|
|
|
|
import com.gzzn.omms.msgexchange.infra.retry.ProcFailure
|
2026-09-06 16:13:50 +08:00
|
|
|
import jakarta.inject.Singleton
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* ACMA-8 流程 4:日计划快照(generation,主泵内执行)。
|
|
|
|
|
* staging(内存瞬态,崩溃从 raw 整包重放)→ Lua 原子“覆盖+按代差删”(I4/I5)
|
2026-09-06 18:19:04 +08:00
|
|
|
* → putGenIfVersion CAS + SUCCEEDED(重放幂等,版本不二次自增——恢复协议细节见 U09)。
|
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,
|
|
|
|
|
private val refData: RefDataRepository,
|
|
|
|
|
private val redis: FlightRedisClient,
|
2026-09-06 19:08:02 +08:00
|
|
|
private val procFailure: ProcFailure,
|
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-06 16:13:50 +08:00
|
|
|
// 1) staging:流式解析 + 整包校验(TODO(阶段2): 流式 codec;千级 FLTR 为 MB 级,内存瞬态)
|
|
|
|
|
val staged = StageResult.stagingOf(msg) // 骨架:TODO 解析 FLTR 集与重组(KEEP 现役 MAFL/登机桥规则)
|
|
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
val normalized = (staged as StageResult.Ok).flights // (flid, payloadJson)
|
|
|
|
|
|
|
|
|
|
// 2) 发布:Lua 原子覆盖 + 按代差删(删除集 = gen.flids − 新代,ADFT 自动存活)
|
|
|
|
|
val day = staged.day
|
|
|
|
|
val gen = refData.getGen(day)
|
|
|
|
|
val newFlids = normalized.map { it.first }
|
|
|
|
|
val delFields = gen?.flids?.minus(newFlids.toSet()) ?: emptyList()
|
|
|
|
|
redis.eval(RedisScript.SNAPSHOT_REPLACE, setPairs = normalized, delFields = delFields)
|
|
|
|
|
|
|
|
|
|
// 3) 事务:putGen CAS + SUCCEEDED(实装后同 @Transactional;CAS 失败=并发,串行泵下不应发生→告警)
|
|
|
|
|
val expected = gen?.version ?: 0L
|
|
|
|
|
if (!refData.putGenIfVersion(day, expected, RefDataRepository.GenMeta(newFlids, expected + 1))) {
|
|
|
|
|
// 重放路径:version 已是目标值 → no-op 视为成功
|
|
|
|
|
val again = refData.getGen(day)
|
|
|
|
|
if (again == null || again.version != expected + 1) {
|
2026-09-06 21:02:27 +08:00
|
|
|
log.warn("gen CAS conflict -> FAILED(INFRA) id={}", head.cminmsgsId)
|
2026-09-06 19:08:02 +08:00
|
|
|
procFailure.fail(head, ErrorClass.INFRA, "gen-cas-conflict") // N06/N28:带退避,禁止紧循环
|
2026-09-06 16:13:50 +08:00
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-09-06 18:19:04 +08:00
|
|
|
procState.update(head.cminmsgsId, ProcStatus.SUCCEEDED)
|
2026-09-06 21:02:27 +08:00
|
|
|
log.info("snapshot SUCCEEDED id={} day={} flights={}", head.cminmsgsId, day, normalized.size)
|
2026-09-06 18:19:04 +08:00
|
|
|
}
|
|
|
|
|
|
2026-09-06 16:13:50 +08:00
|
|
|
/** staging 结果(骨架)。 */
|
|
|
|
|
sealed interface StageResult {
|
|
|
|
|
data class Ok(val day: String, val flights: List<Pair<String, String>>) : StageResult
|
|
|
|
|
data class Invalid(val reason: String) : StageResult
|
|
|
|
|
|
|
|
|
|
companion object {
|
|
|
|
|
fun stagingOf(msg: DecodedMessage): StageResult =
|
|
|
|
|
Invalid("staging-not-implemented(${msg.typeTag})") // TODO(阶段2)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|