Files
msgexchange-v2/src/main/kotlin/com/gzzn/omms/msgexchange/processing/Pump.kt
T

255 lines
12 KiB
Kotlin
Raw Normal View History

package com.gzzn.omms.msgexchange.processing
import com.gzzn.omms.msgexchange.codec.FlopPayload
import com.gzzn.omms.msgexchange.codec.ScheduleBody
import com.gzzn.omms.msgexchange.codec.XmlCodec
import com.gzzn.omms.msgexchange.config.OperationDayProps
import com.gzzn.omms.msgexchange.config.PipelineProps
import com.gzzn.omms.msgexchange.domain.DecodedMessage
import com.gzzn.omms.msgexchange.domain.ErrorClass
import com.gzzn.omms.msgexchange.domain.MsgKind
import com.gzzn.omms.msgexchange.domain.ProcState
import com.gzzn.omms.msgexchange.domain.ProcStatus
import com.gzzn.omms.msgexchange.infra.log.TraceLog
import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository
import com.gzzn.omms.msgexchange.infra.persistence.InboxCursorRepository
import com.gzzn.omms.msgexchange.infra.persistence.ProcStateRepository
import com.gzzn.omms.msgexchange.infra.retry.ProcFailure
import jakarta.inject.Singleton
import java.time.Clock
import java.time.Duration
import java.time.Instant
import java.time.LocalDate
import java.util.concurrent.atomic.AtomicLong
/**
* 处理主泵:一个线程按消息 ID 从小到大一条条处理,保证先来的先处理。
*
* 每次 tick 只看当前最小的未完成消息("队头"):
* - 没有待处理消息就睡一个轮询间隔;
* - 队头失败了还在退避期,就等到能重试的时刻;重试次数用尽才转死信,不放任它一直堵着;
* - 其余情况交给 [MessageProcessor] 处理。
*
* 一次只处理一条是刻意的。后面的消息不能越过卡住的队头,否则同一条航班的报文
* 可能被乱序应用,几十秒后才到的旧报文会把新状态覆盖回去。
*/
@Singleton
class Pump(
private val procState: ProcStateRepository,
private val cursor: InboxCursorRepository,
private val processor: MessageProcessor,
private val props: PipelineProps,
private val clock: Clock,
) {
private val log = org.slf4j.LoggerFactory.getLogger(Pump::class.java)
/** 连续失败计数:仅用于日志/排障,不代表业务状态。 */
private val tickFailures = AtomicLong(0)
@Volatile
private var running = true
/** 请求停机:当前 tick 跑完就退出。线程中断由 PipelineLifecycle 负责。 */
fun stop() {
running = false
}
fun loop() {
while (running) {
try {
tick()
tickFailures.set(0)
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
return
} catch (e: Exception) {
// 兜底:单条消息的失败状态已在 processOne 内记录;但 tick 整体失败(DB 断连、
// 锁超时)在别处没有痕迹,必须留日志,否则排障无据可查。
val failures = tickFailures.incrementAndGet()
log.error("pump tick failed (failure #{})", failures, e)
sleepQuietly(props.pipeline.pollInterval)
}
}
}
internal fun tick() {
val head = procState.headUnfinished()
if (head == null) {
sleepQuietly(props.pipeline.pollInterval)
return
}
// 只领取"已被水位覆盖"的队头(`msgId <= W`)。
//
// 水位以内的行都是收报按 ID 顺序发现并登记的;水位之外的行只可能来自兼容入口
// 直接写 PROC_STATE(它不参与水位)。若允许领取,它就会越过那些尚未入队的较小 ID,
// 破坏 FIFO(不变量"只领取已发现的行"`invariants.md` INV-4)。这种行在空洞补齐、`W` 追平之后自然可领取。
val watermark = cursor.load().committedUpTo
if (head.msgId > watermark) {
warnBeyondWatermark(head.msgId, watermark)
sleepQuietly(props.pipeline.pollInterval)
return
}
val now = clock.instant()
when {
head.state == ProcStatus.FAILED && head.attempts >= props.pipeline.maxAttempts -> {
log.error("poison -> DEAD msgId={} attempts={} lastError={}", head.msgId, head.attempts, head.lastError)
// markTerminal 在同一条 UPDATE 里登记回填意图;回填由扫描补写,不在这里做跨库写。
procState.markTerminal(
head.msgId, ProcStatus.DEAD,
errorClass = ErrorClass.EXHAUSTED,
lastError = head.lastError ?: "attempts-exhausted",
attempts = head.attempts,
now = now,
)
}
head.state == ProcStatus.FAILED && (head.nextAttemptAt ?: Instant.EPOCH) > now ->
sleepQuietly(Duration.between(now, head.nextAttemptAt))
// 其余情况(新消息,或退避到期的重试)交给处理入口
else -> processor.processOne(head)
}
}
/** 上一次"队头在水位之外"告警时的水位值:只在它变化时告警,避免每秒刷屏。 */
private val warnedWatermark = AtomicLong(Long.MIN_VALUE)
private fun warnBeyondWatermark(msgId: Long, watermark: Long) {
if (warnedWatermark.getAndSet(watermark) != watermark) {
log.warn(
"head msgId={} is beyond watermark W={}; waiting for discovery " +
"(row injected by the compat entry point?)",
msgId, watermark,
)
}
}
private fun sleepQuietly(d: Duration) {
if (!d.isNegative && !d.isZero) Thread.sleep(d.toMillis().coerceAtLeast(1))
}
}
/**
* 处理一条消息:读原文 → 解码 → 绑定业务身份 → 分派给对应处理器。
*
* 业务数据、终态与回填意图都由各处理器在自己的事务里写入(终态与回填意图是同一条 UPDATE)。
* **这里不做信箱回填**:主泵是 FIFO 关键路径,跨库写会把它绑在共享 MySQL 的可用性上。
* 回填由 `JobRunner` 定时的 `BackfillService.sweep` 驱动(调度周期不等于完成时限)。
*
* 任何意外异常都算在当前这条消息头上(记 FAILED(INFRA) 后重试),不会把主泵线程带崩。
* 报文非法和整包协议拒绝不重试,直接进死信等人工处置。
*/
@Singleton
class MessageProcessor(
private val inbox: CminmsgInboxRepository,
private val procState: ProcStateRepository,
private val codec: XmlCodec,
private val scheduleProcessor: ScheduleProcessor,
private val flopProcessor: FlopProcessor,
private val fdelProcessor: FdelProcessor,
private val adftProcessor: AdftProcessor,
private val procFailure: ProcFailure,
private val props: PipelineProps,
private val clock: Clock,
private val operationDayProps: OperationDayProps,
) {
private val log = org.slf4j.LoggerFactory.getLogger(MessageProcessor::class.java)
fun processOne(head: ProcState) {
TraceLog.withTrace(head.msgId) {
try {
processInternal(head)
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
throw e
} catch (e: Exception) {
// 边界化:异常归于本条 head,写 FAILED(INFRA)/DEAD,而不是穿出杀 pump
log.warn("processOne unexpected failure msgId={} ec=INFRA msg={}", head.msgId, e.message ?: e.javaClass.simpleName)
procFailure.fail(head, ErrorClass.INFRA, e.message ?: e.javaClass.simpleName)
}
}
}
private fun processInternal(head: ProcState) {
val raw = inbox.rawOf(head.msgId)
if (raw == null) {
log.error("raw missing -> DEAD(MALFORMED) msgId={}", head.msgId)
return deadMalformed(head, "raw-missing")
}
val decoded = when (val r = codec.decode(raw)) {
is com.gzzn.omms.msgexchange.codec.DecodeResult.Ok -> r.message
is com.gzzn.omms.msgexchange.codec.DecodeResult.Err -> {
// MALFORMED(报文非法)→ DEAD 不重试;CODEC_ERROR(可随 codec 修复重放)→ FAILED 退避
if (r.failure.errorClass == ErrorClass.MALFORMED) {
log.error("decode MALFORMED -> DEAD msgId={} detail={}", head.msgId, r.failure.detail)
return deadMalformed(head, r.failure.detail)
}
log.warn("decode {} -> FAILED msgId={} detail={}", r.failure.errorClass, head.msgId, r.failure.detail)
return procFailure.fail(head, r.failure.errorClass, r.failure.detail)
}
}
// I3identity 仅首次绑定(head.identityKey == null);FAILED 重试不重绑
if (head.identityKey == null) {
val identity = Identity.of(
decoded, props.identity, LocalDate.now(clock.withZone(operationDayProps.zoneId())),
)
if (!procState.tryBindIdentity(head.msgId, identity)) {
val owner = procState.ownerOfIdentity(identity) ?: -1L
log.info("duplicate-of:{} -> SKIPPED msgId={}", owner, head.msgId)
procState.markTerminal(head.msgId, ProcStatus.SKIPPED, lastError = "duplicate-of:$owner", now = clock.instant())
return
}
}
// 按报文类型分派:日计划走 SCHD,其余走 FLOP / FDEL / ADFT;报文缺载荷直接判为非法报文的死信
val result: ApplyResult = when (val kind = decoded.kind) {
is MsgKind.Schd -> {
val body = decoded.body as? ScheduleBody
if (body == null) return deadMalformed(head, "missing-schd-body")
when (kind.subtype) {
MsgKind.SchdSubtype.DNLD, MsgKind.SchdSubtype.RESP ->
scheduleProcessor.applyScheduleRecords(head, decoded)
MsgKind.SchdSubtype.ADFT -> {
val record = body.records.singleOrNull()
if (record == null) return deadMalformed(head, "adft-needs-single-fltr")
adftProcessor.apply(head, decoded, record)
}
}
}
MsgKind.Fdel -> {
val payload = decoded.body as? FlopPayload
if (payload == null) return deadMalformed(head, "missing-fdel-flid")
fdelProcessor.apply(head, decoded, payload)
}
is MsgKind.Flop -> {
val payload = decoded.body as? FlopPayload
if (payload == null) return deadMalformed(head, "missing-flop-body")
flopProcessor.apply(head, decoded, payload)
}
is MsgKind.Unsupported -> {
// 还没有对应处理器的报文类型:先按可重试的失败处理,等能力补齐,不直接判死
log.warn("unsupported type -> FAILED(UNSUPPORTED) msgId={} tag={}", head.msgId, kind.tag)
return procFailure.fail(head, ErrorClass.UNSUPPORTED, "no-handler:${kind.tag}")
}
}
if (result is ApplyResult.DeadProtocol) {
// 整包被拒绝:不重试、立刻放掉队头,等人工确认
log.error("DEAD(PROTOCOL) msgId={} reason={} flags={}", head.msgId, result.reason, result.flags)
procState.markTerminal(
head.msgId, ProcStatus.DEAD,
errorClass = ErrorClass.PROTOCOL,
lastError = result.reason.take(1000),
now = clock.instant(),
)
return
}
// Succeeded / ReplaySkippedSUCCEEDED 终态与回填意图已由处理器在自己的事务内落库
log.info("SUCCEEDED msgId={} kind={}", head.msgId, decoded.typeTag)
}
private fun deadMalformed(head: ProcState, detail: String) {
procState.markTerminal(head.msgId, ProcStatus.DEAD, errorClass = ErrorClass.MALFORMED, lastError = detail, now = clock.instant())
}
}