Files
msgexchange-v2/src/test/kotlin/com/gzzn/omms/msgexchange/processing/BackfillServiceTest.kt
T
windyboyandCursor dc1c49af6e docs(acm2-75): 精简不变量与声明边界,去掉无依据条目
审改 specification 管道/航班域/投影与 CLM:作废与需求重复或依据不足的 INV/CLM,白话重写保留条款,并同步架构、实现与引用注释。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-16 17:26:49 +08:00

384 lines
17 KiB
Kotlin
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package com.gzzn.omms.msgexchange.processing
import com.gzzn.omms.msgexchange.config.MailboxProps
import com.gzzn.omms.msgexchange.config.PipelineProps
import com.gzzn.omms.msgexchange.domain.ErrorClass
import com.gzzn.omms.msgexchange.domain.ProcStatus
import com.gzzn.omms.msgexchange.infra.persistence.CminmsgInboxRepository
import com.gzzn.omms.msgexchange.infra.persistence.MailboxRow
import com.gzzn.omms.msgexchange.infra.persistence.MailboxMarkResult
import com.gzzn.omms.msgexchange.infra.stub.StubInbox
import com.gzzn.omms.msgexchange.infra.stub.StubProcState
import com.gzzn.omms.msgexchange.infra.retry.ReplayService
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Assertions.assertFalse
import org.junit.jupiter.api.Assertions.assertNotNull
import org.junit.jupiter.api.Assertions.assertNull
import org.junit.jupiter.api.Assertions.assertTrue
import org.junit.jupiter.api.Assertions.assertThrows
import org.junit.jupiter.api.Test
import java.time.Clock
import java.time.Duration
import java.time.Instant
import java.time.ZoneOffset
import java.util.concurrent.CountDownLatch
import java.util.concurrent.Callable
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit
import java.util.concurrent.TimeoutException
/**
* 回填环节的规矩:
* - 处理完马上写标记,写不进去就按退避重试,重启后接着重试(状态都在数据库里);
* - 等得太久的消息无视退避强制补写,保证标记最终一定会写上;
* - 只写还没有标记的行,已有的值不覆盖;
* - 还没处理完的消息(PENDING / FAILED)不许写标记;
* - 回填失败只是回填的事,不会把处理结果改回去。
*/
class BackfillServiceTest {
private val t0: Instant = Instant.parse("2026-09-08T03:00:00Z")
private val props = PipelineProps()
/** 可以人为制造故障的信箱,用来验证回填失败时怎么处理。 */
private class FakeMailbox(var fail: Boolean = false, var missing: Boolean = false) : CminmsgInboxRepository {
val marked = linkedSetOf<Long>()
override fun insertRaw(rawXml: String): Long = 1L
override fun rawOf(msgId: Long): String? = null
override fun receivedAtOf(msgId: Long): Instant? = null
override fun readRange(fromExclusive: Long, limit: Int): List<MailboxRow> = emptyList()
override fun maxId(): Long? = null
override fun minId(): Long? = null
override fun markProcessedIfUnmarked(msgId: Long, value: String): MailboxMarkResult {
if (fail) throw IllegalStateException("mysql-down")
if (missing) return MailboxMarkResult.MISSING
return if (marked.add(msgId)) MailboxMarkResult.MARKED else MailboxMarkResult.ALREADY_MARKED
}
}
private fun service(proc: StubProcState, mailbox: CminmsgInboxRepository, now: Instant = t0) =
BackfillService(proc, mailbox, MailboxProps(), props, Clock.fixed(now, ZoneOffset.UTC), MessageLifecycleGate())
/** 终态 + 回填意图(固定时刻,避免依赖真实时钟)。 */
private fun succeeded(proc: StubProcState, id: Long) {
proc.insertIfAbsent(id, t0)
proc.markTerminal(id, ProcStatus.SUCCEEDED, now = t0)
}
@Test
fun `terminal message is marked immediately and the intent is cleared`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val id = inbox.insertRaw("<MSG/>")
succeeded(proc, id)
service(proc, inbox).attempt(id)
assertEquals("PROCESSED", inbox.markOf(id))
val row = proc.find(id)!!
assertNotNull(row.backfillAt)
assertNull(row.backfillNextAt)
assertNull(row.backfillError)
}
@Test
fun `marking is monotonic - an existing value is never overwritten`() {
val inbox = StubInbox().apply { clear() }
val id = inbox.insertRaw("<MSG/>")
assertEquals(MailboxMarkResult.MARKED, inbox.markProcessedIfUnmarked(id, "PROCESSED"))
assertEquals(MailboxMarkResult.ALREADY_MARKED, inbox.markProcessedIfUnmarked(id, "OTHER"))
assertEquals("PROCESSED", inbox.markOf(id))
}
@Test
fun `already marked rows are idempotent and never recorded as failures`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val id = inbox.insertRaw("<MSG/>")
succeeded(proc, id)
val backfill = service(proc, inbox)
backfill.attempt(id)
val firstCompletion = proc.find(id)!!.backfillAt
backfill.attempt(id) // 重复补写没有副作用
assertNull(proc.find(id)!!.backfillError)
assertEquals(0, proc.find(id)!!.backfillAttempts)
assertEquals(firstCompletion, proc.find(id)!!.backfillAt)
}
@Test
fun `failed attempt records backoff and never touches the terminal state`() {
val proc = StubProcState()
val mailbox = FakeMailbox(fail = true)
succeeded(proc, 901L)
service(proc, mailbox).attempt(901L)
val row = proc.find(901L)!!
assertEquals(ProcStatus.SUCCEEDED, row.state) // 回填失败不会把处理结果改回去
assertEquals(1, row.backfillAttempts)
assertEquals("mysql-down", row.backfillError)
assertEquals(t0.plus(Duration.ofSeconds(30)), row.backfillNextAt)
assertNull(row.backfillAt)
}
/**
* 评审修正 R3:信箱行不存在是**确定性结论**,不再当作"可重试失败"——
* 重试不会改变结果,只会每 30 秒重试一次并永久占满扫描批次。
* 但仍必须区分"停止自动重试"与"标记已确认"backfillAt 保持为空。
*/
@Test
fun `missing mailbox row is abandoned instead of retried forever`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val id = inbox.insertRaw("<MSG/>")
succeeded(proc, id)
inbox.removeRow(id)
service(proc, inbox).attempt(id)
val row = proc.find(id)!!
assertNull(row.backfillAt) // 放弃 ≠ 已确认
assertEquals(BackfillService.ABANDON_MISSING_ROW, row.backfillAbandonedReason)
assertNotNull(row.backfillAbandonedAt)
assertNull(row.backfillNextAt) // 不再排下一次重试
assertNull(row.backfillError) // 不再是"失败重试",而是放弃
assertEquals(0, service(proc, inbox).sweep(t0)) // 也不再被扫描
}
@Test
fun `sweep retries due rows and completes once the mailbox recovers`() {
val proc = StubProcState()
val mailbox = FakeMailbox(fail = true)
succeeded(proc, 901L)
val backfill = service(proc, mailbox)
assertEquals(1, backfill.sweep(t0)) // 到期 → 失败 → 退避
assertEquals(0, backfill.sweep(t0.plusSeconds(29))) // 未到期
assertEquals(1, backfill.sweep(t0.plusSeconds(30))) // 到期再试 → 仍失败(重启后同样收敛)
mailbox.fail = false
assertEquals(1, backfill.sweep(t0.plusSeconds(90)))
assertTrue(901L in mailbox.marked)
assertNotNull(proc.find(901L)!!.backfillAt)
}
/** 等得太久的消息不再等退避、直接补写:保证标记最终一定会写上。 */
@Test
fun `overdue rows bypass the retry backoff`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val id = inbox.insertRaw("<MSG/>")
val old = t0.minus(props.pipeline.overdueBackfill).minusSeconds(60)
proc.insertIfAbsent(id, old) // 接收时间早于 R
proc.markTerminal(id, ProcStatus.SUCCEEDED, now = t0)
proc.recordBackfillFailure(id, "mysql-down", 5, t0.plus(Duration.ofMinutes(15)), t0) // 退避推到很远之后
assertEquals(1, service(proc, inbox).sweep(t0))
assertNotNull(proc.find(id)!!.backfillAt)
}
/** 还没处理完的消息(PENDING / FAILED)永远不打标,等再久也不行。 */
@Test
fun `mid states are never marked even when far past the deadline`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val pending = inbox.insertRaw("<MSG/>")
val failed = inbox.insertRaw("<MSG/>")
val old = t0.minus(props.pipeline.overdueBackfill).minusSeconds(60)
proc.insertIfAbsent(pending, old)
proc.insertIfAbsent(failed, old)
proc.update(failed, ProcStatus.FAILED, errorClass = ErrorClass.INFRA)
assertEquals(0, service(proc, inbox).sweep(t0))
assertFalse(inbox.isMarked(pending))
assertFalse(inbox.isMarked(failed))
}
@Test
fun `immediate attempt also refuses pending and failed messages`() {
val proc = StubProcState()
val inbox = StubInbox().apply { clear() }
val pending = inbox.insertRaw("<MSG/>")
val failed = inbox.insertRaw("<MSG/>")
proc.insertIfAbsent(pending, t0)
proc.insertIfAbsent(failed, t0)
proc.update(failed, ProcStatus.FAILED, errorClass = ErrorClass.INFRA)
service(proc, inbox).attempt(pending)
service(proc, inbox).attempt(failed)
assertFalse(inbox.isMarked(pending))
assertFalse(inbox.isMarked(failed))
assertNull(proc.find(pending)!!.backfillAt)
assertNull(proc.find(failed)!!.backfillAt)
}
@Test
fun `backoff delay doubles per attempt and caps at fifteen minutes`() {
val svc = BackfillService(StubProcState(), StubInbox(), MailboxProps(), props, Clock.fixed(t0, ZoneOffset.UTC), MessageLifecycleGate())
assertEquals(Duration.ofSeconds(30), svc.backoffDelayFor(1))
assertEquals(Duration.ofMinutes(2), svc.backoffDelayFor(3))
assertEquals(Duration.ofMinutes(15), svc.backoffDelayFor(10))
assertEquals(Duration.ofMinutes(15), svc.backoffDelayFor(50))
}
@Test
fun `replay waits for an in-flight backfill before reopening the message`() {
val entered = CountDownLatch(1)
val release = CountDownLatch(1)
val replayStarted = CountDownLatch(1)
val mailbox = object : CminmsgInboxRepository {
override fun insertRaw(rawXml: String) = 1L
override fun rawOf(msgId: Long): String? = "<MSG/>"
override fun receivedAtOf(msgId: Long) = t0
override fun readRange(fromExclusive: Long, limit: Int) = emptyList<MailboxRow>()
override fun maxId(): Long? = 1L
override fun minId(): Long? = 1L
override fun markProcessedIfUnmarked(msgId: Long, value: String): MailboxMarkResult {
entered.countDown()
release.await()
return MailboxMarkResult.MARKED
}
}
val proc = StubProcState().apply {
insertIfAbsent(1L, t0)
markTerminal(1L, ProcStatus.DEAD, errorClass = ErrorClass.EXHAUSTED, now = t0)
}
val gate = MessageLifecycleGate()
val backfill = BackfillService(proc, mailbox, MailboxProps(), props, Clock.fixed(t0, ZoneOffset.UTC), gate)
val replay = ReplayService(proc, gate)
val pool = Executors.newFixedThreadPool(2)
try {
val backfillTask = pool.submit { backfill.attempt(1L) }
assertTrue(entered.await(1, TimeUnit.SECONDS))
val replayTask = pool.submit(Callable {
replayStarted.countDown()
replay.replay(listOf(ErrorClass.EXHAUSTED))
})
assertTrue(replayStarted.await(1, TimeUnit.SECONDS))
assertThrows(TimeoutException::class.java) { replayTask.get(100, TimeUnit.MILLISECONDS) }
release.countDown()
backfillTask.get(1, TimeUnit.SECONDS)
assertEquals(1, replayTask.get(1, TimeUnit.SECONDS))
assertEquals(ProcStatus.PENDING, proc.find(1L)!!.state)
} finally {
release.countDown()
pool.shutdownNow()
}
}
// ------------------------------------------------------------------
// 回填闭环(评审修正 R3 + 饥饿回归)
// ------------------------------------------------------------------
@Test
fun `missing mailbox row is abandoned immediately and is never treated as marked`() {
val proc = StubProcState()
val mailbox = FakeMailbox(missing = true)
succeeded(proc, 911L)
service(proc, mailbox).attempt(911L)
val row = proc.find(911L)!!
assertEquals(BackfillService.ABANDON_MISSING_ROW, row.backfillAbandonedReason)
assertNotNull(row.backfillAbandonedAt)
// 放弃 ≠ 标记已确认:清除前提(边界内全部行已打标)因此仍然不成立。
assertNull(row.backfillAt)
assertNull(row.backfillNextAt)
// 已放弃的行不再进入扫描:不会每 30 秒无限重试。
assertEquals(0, service(proc, mailbox).sweep(t0))
// 但它**不是**永久失去补偿:人工恢复后可以重新排队。
assertTrue(service(proc, mailbox).reopen(911L))
assertNull(proc.find(911L)!!.backfillAbandonedAt)
assertNotNull(proc.find(911L)!!.backfillNextAt)
}
/**
* 门禁裁决 D1:暂时性故障**不按次数放弃**。即使超过 `backfill-max-attempts` 这个
* 告警阈值,也继续退避重试——一次小时级的共享库故障不该把待回填行成批判死。
*/
@Test
fun `transient failures keep retrying past the warning threshold`() {
props.pipeline.backfillMaxAttempts = 3
val proc = StubProcState()
val mailbox = FakeMailbox(fail = true)
succeeded(proc, 912L)
val svc = service(proc, mailbox)
repeat(5) { svc.attempt(912L) }
val row = proc.find(912L)!!
assertEquals(5, row.backfillAttempts)
assertNull(row.backfillAbandonedAt) // 次数不是放弃判据
assertNull(row.backfillAt)
assertNotNull(row.backfillNextAt) // 仍在退避重试
}
/** 到 `R` 仍未打标(overdue)时,暂时性故障才停止自动重试并进放弃清单。 */
@Test
fun `overdue transient failure is abandoned with the deadline reason and stays recoverable`() {
val proc = StubProcState()
val mailbox = FakeMailbox(fail = true)
succeeded(proc, 913L)
val svc = service(proc, mailbox)
svc.attempt(913L, overdue = true)
val row = proc.find(913L)!!
assertEquals(BackfillService.ABANDON_TRANSIENT_DEADLINE, row.backfillAbandonedReason)
assertNotNull(row.backfillAbandonedAt)
assertNull(row.backfillAt) // 放弃 ≠ 标记已确认
assertTrue(svc.reopen(913L)) // 人工恢复入口存在且有效
assertNull(proc.find(913L)!!.backfillAbandonedAt)
}
@Test
fun `permanently failing oldest rows do not starve later rows (fair rotation)`() {
props.pipeline.backfillBatch = 3
val proc = StubProcState()
val overdueBefore = t0.plusSeconds(3600)
// 1..3:模拟"最旧且已尝试多次"的永久失败行,保持 duenextAt <= now
(1L..3L).forEach { id ->
proc.insertIfAbsent(id, t0)
proc.markTerminal(id, ProcStatus.SUCCEEDED, now = t0)
proc.recordBackfillFailure(id, "still-failing", attempts = 50, nextAttemptAt = t0, now = t0)
}
// 200:新行,从未尝试
proc.insertIfAbsent(200L, t0)
proc.markTerminal(200L, ProcStatus.SUCCEEDED, now = t0)
val due = proc.findBackfillDue(t0, overdueBefore, limit = 3)
assertEquals(3, due.size)
// 公平轮转:尝试次数少的先被扫描,因此新行不会被最旧的一批永久失败行饿死。
assertEquals(200L, due.first().msgId)
assertTrue(due.any { it.msgId == 200L })
}
/**
* 超期判据是**本地入队时间** `ENQUEUED_AT`,与 `RECEIVED_AT` 无关。
* 上游没给接收时间(NULL)不再让 `R` 兜底失效——这正是引入 `ENQUEUED_AT` 要消除的窗口。
*/
@Test
fun `null received time no longer disables the overdue shortcut`() {
val proc = StubProcState()
val old = t0.minus(props.pipeline.overdueBackfill).minusSeconds(60)
proc.insertIfAbsent(921L, null, enqueuedAt = old) // 上游未提供接收时间
proc.markTerminal(921L, ProcStatus.SUCCEEDED, now = t0)
proc.recordBackfillFailure(921L, "mysql-down", attempts = 1, nextAttemptAt = t0.plusSeconds(3600), now = t0)
val due = proc.findBackfillDue(now = t0, overdueBefore = t0.minus(props.pipeline.overdueBackfill), limit = 10)
assertEquals(listOf(921L), due.map { it.msgId }) // 退避未到期,但入队时间已超 R → 仍被扫描
assertTrue(due.single().overdue)
}
}