package com.gzzn.omms.msgexchange.delivery import com.gzzn.omms.msgexchange.config.PipelineProps import com.gzzn.omms.msgexchange.domain.ErrorClass import com.gzzn.omms.msgexchange.domain.EventStatus import com.gzzn.omms.msgexchange.domain.EventType import com.gzzn.omms.msgexchange.domain.MsgEvent import com.gzzn.omms.msgexchange.domain.Targets import com.gzzn.omms.msgexchange.infra.persistence.MsgEventRepository import com.gzzn.omms.msgexchange.infra.retry.FailureScheduler import jakarta.inject.Singleton import java.time.Duration import java.time.Instant import java.util.concurrent.atomic.AtomicLong /** 对外投递的出口:同步等 Kafka 确认,语义是至少发一次(可能重复,但不会丢)。以后 ES 投影也从这里加。 */ interface DeliveryPort { /** 发一条变化通知,消息 key 是 FLID(航班实例 ID)。 */ fun sendKafka(topic: String, key: String, payloadJson: String) /** 发一条完整状态,消息 key 是 FLID。适配层负责把 target 映射成 topic:KAFKA:msg 对应 "msg",KAFKA:schd 对应 "schd"。 */ fun sendKafkaSchd(topic: String, key: String, payloadJson: String) /** 发一条删除通知:key 是 FLID、value 为空;下游按"整态里这个键没了"理解成删除。 */ fun sendKafkaNull(topic: String, key: String) /** 给健康检查用的连通性探测;默认返回 true,真实 Kafka 实现要覆写成向 broker 拉一次 metadata 来判断。 */ fun ping(): Boolean = true } /** * 投递调度:把 outbox(待发事件表)里的事件发给下游。 * * KAFKA_MSG 一条一条按登记顺序发,不插队。KAFKA_SCHD 走 flushSchd:outbox 每个 FLID 只保留 * 一行(写入侧单行 upsert),发出后按读取时刻的代次做条件确认;删除通知发 value 为空的 tombstone。 * 两个主题之间不保证先后顺序。 * * 失败处理:队首的重试时间没到就不取;一批里有发送失败,整批重试次数加一并推后退避, * 次数用尽整批转 DEAD 当死信。见 docs/flight-state.md §5。 */ @Singleton class Dispatcher( private val msgEvents: MsgEventRepository, private val port: DeliveryPort, private val props: PipelineProps, private val scheduler: FailureScheduler, ) { private val log = org.slf4j.LoggerFactory.getLogger(Dispatcher::class.java) /** 连续失败计数:仅用于日志/排障。 */ private val tickFailures = AtomicLong(0) @Volatile private var running = true private var lastFlush: Instant? = null /** 请求停机:置位后 loop 走完当前一轮就退出;真正中断线程由 PipelineLifecycle 负责。 */ fun stop() { running = false } fun loop() { while (running) { try { tick() tickFailures.set(0) } catch (e: InterruptedException) { Thread.currentThread().interrupt() return } catch (e: Exception) { // tick 整体失败(DB 断连等)必须在别处留痕,否则只剩"投递不动"没有原因。 val failures = tickFailures.incrementAndGet() log.error("dispatcher tick failed (failure #{})", failures, e) sleepQuietly(props.pipeline.pollInterval) } } } internal fun tick() { val sent = drainKafkaMsg() if (flushDue()) flushSchd() // 有活就不睡:pollInterval 只用于"本轮无事可做或队头在退避",不能当投递节流用。 if (sent == 0 && running) sleepQuietly(props.pipeline.pollInterval) } /** * 按 `EVENT_ID` 顺序批量投递 `KAFKA:msg`。 * * 保序规则不变:**同一目标内不越序**——队头失败或退避未到期时立即停止本轮, * 不跳过它去投后面的。批量只用来省掉"每条一次 DB 往返 + 一轮一次 sleep"。 * * @return 本轮成功发出的条数 */ private fun drainKafkaMsg(): Int { val batchSize = props.pipeline.deliveryBatch.coerceAtLeast(1) val maxRounds = props.pipeline.deliveryDrainRounds.coerceAtLeast(1) var sent = 0 var round = 0 while (round < maxRounds) { round++ val batch = try { msgEvents.claimBatch(Targets.KAFKA_MSG, batchSize) } catch (e: Exception) { log.error("claimBatch(KAFKA_MSG) failed", e) return sent } if (batch.isEmpty()) return sent for (e in batch) { val next = e.nextAttemptAt if (next != null && next > scheduler.now()) return sent // 队头退避未到期:停止推进 if (!deliver(Targets.KAFKA_MSG, e)) return sent // 失败已记账:停止以保序 sent++ } if (batch.size < batchSize) return sent } log.warn("kafka msg drain hit round cap ({} rounds, sent={}); yielding to schd flush", maxRounds, sent) return sent } private fun flushDue(): Boolean = lastFlush?.let { Duration.between(it, scheduler.now()) >= props.schd.flushPeriod } ?: true /** @return 是否已确认发出(false 表示已按退避/死信记账,调用方应停止本轮以保序) */ private fun deliver(target: String, e: MsgEvent): Boolean { val eventId = e.eventId ?: return false return try { when (target) { Targets.KAFKA_MSG -> port.sendKafka("msg", e.partitionKey, e.payloadJson) // 未知目标不能在"未发送"的情况下被标记为已发;记账后停止本轮。 else -> error("unknown delivery target: $target") } msgEvents.markSent(eventId) true } catch (ex: Exception) { retryOrDead(e, ex.message ?: ex.javaClass.simpleName) false } } /** 发 KAFKA_SCHD:每个 FLID 只有一行待发事件,删除通知发空 value,成功后按代次条件确认。 */ internal fun flushSchd() { val batch = try { msgEvents.mergePendingSchd(scheduler.now(), props.schd.flushLimit) } catch (e: Exception) { log.warn("mergePendingSchd failed: {}", e.message) return } if (batch.isEmpty()) { lastFlush = scheduler.now() return } val failures = mutableListOf() for (e in batch) { try { when (e.eventType) { EventType.TOMBSTONE -> port.sendKafkaNull("schd", e.partitionKey) EventType.UPSERT -> port.sendKafkaSchd("schd", e.partitionKey, e.payloadJson) } // 条件确认:读取时刻的代次(EVENT_ID + STATE_VERSION)被新写入覆盖时不标记,留待下一轮重发。 e.eventId?.let { msgEvents.markSentIfVersion(it, e.stateVersion) } } catch (ex: Exception) { failures.add(e) } } failures.forEach { retryOrDead(it, it.lastError ?: "send-failed") } lastFlush = scheduler.now() } /** 一条事件发失败之后怎么走:重试次数加一,到上限就标成 DEAD(EXHAUSTED) 留作死信(次数落库便于追查),否则按退避推到下次再发。 */ private fun retryOrDead(e: MsgEvent, lastError: String) { val eventId = e.eventId ?: return val attempts = e.attempts + 1 if (scheduler.exhausted(attempts)) { msgEvents.markDead(eventId, ErrorClass.EXHAUSTED, lastError, attempts) } else { msgEvents.scheduleRetry(eventId, scheduler.nextAttemptAt(attempts), attempts) } } private fun sleepQuietly(d: Duration) { if (!d.isNegative && !d.isZero) Thread.sleep(d.toMillis().coerceAtLeast(1)) } }