fix(delivery): 删除通知只登记 KAFKA:msg,schd 不再发 tombstone(ACM2-87)

- FdelProcessor 移除 schd TOMBSTONE,仅保留 msg 删除通知(C-9)
- HistorySweepJob 补发的清退通知由 schd TOMBSTONE 改为 msg 删除通知
- msg value 形态沿用现状(deleted:true JSON),待 Q5 定稿不改成 null tombstone
- 同步修订 FdelAndAdftProcessorTest / HistorySweepJobTest / HistorySweepPurgePgTest 断言
This commit is contained in:
windyboy
2026-09-21 12:31:02 +08:00
parent df087c7933
commit a1c56f546f
5 changed files with 27 additions and 38 deletions
@@ -68,7 +68,7 @@ class HistorySweepJob(
val toPurge = candidates.filter { it.flid in archivedFlids } val toPurge = candidates.filter { it.flid in archivedFlids }
// 归档与删除之间主泵可能已写入同一 FLID:进锁事务后按 (FLID, STATE_VERSION) 复查, // 归档与删除之间主泵可能已写入同一 FLID:进锁事务后按 (FLID, STATE_VERSION) 复查,
// 仍合格才登记 tombstone 并删除;事件与删除同事务提交或回滚(D1、US-14 AC4)。 // 仍合格才登记删除通知并删除;事件与删除同事务提交或回滚(D1、US-14 AC4)。
var purged = 0 var purged = 0
txManager.inTransaction { txManager.inTransaction {
lock.lock() lock.lock()
@@ -78,7 +78,7 @@ class HistorySweepJob(
// 从没收到过 FDEL、被这里直接清掉的航班,删除前补发一次通知,否则下游一直以为它还在 // 从没收到过 FDEL、被这里直接清掉的航班,删除前补发一次通知,否则下游一直以为它还在
val preDelete = rechecked.filter { it.wasNeverFdel } val preDelete = rechecked.filter { it.wasNeverFdel }
if (preDelete.isNotEmpty()) { if (preDelete.isNotEmpty()) {
msgEvents.insertAll(preDelete.map { tombstone(it, now) }) msgEvents.insertAll(preDelete.map { deletionNotice(it, now) })
} }
purged = flightState.purgeArchived(rechecked) purged = flightState.purgeArchived(rechecked)
} }
@@ -86,10 +86,10 @@ class HistorySweepJob(
return SweepOutcome(candidates.size, archivedFlids.size, purged, snapLogPurged) return SweepOutcome(candidates.size, archivedFlids.size, purged, snapLogPurged)
} }
private fun tombstone(candidate: HistoryCandidate, createdAt: Instant) = MsgEvent( /** C-9:删除通知只走 KAFKA:msgvalue 形态("deleted":true 的 JSON)沿用现状,待 Q5 定稿。 */
target = Targets.KAFKA_SCHD, private fun deletionNotice(candidate: HistoryCandidate, createdAt: Instant) = MsgEvent(
target = Targets.KAFKA_MSG,
partitionKey = candidate.flid, partitionKey = candidate.flid,
eventType = EventType.TOMBSTONE,
stateVersion = candidate.stateVersion, stateVersion = candidate.stateVersion,
payloadJson = """{"flid":"${candidate.flid}","stateVersion":${candidate.stateVersion},"deleted":true}""", payloadJson = """{"flid":"${candidate.flid}","stateVersion":${candidate.stateVersion},"deleted":true}""",
createdAt = createdAt, createdAt = createdAt,
@@ -86,20 +86,8 @@ class FdelProcessor(
val createdAt = clock.instant() val createdAt = clock.instant()
msgEvents.insertAll( msgEvents.insertAll(
listOf( listOf(
MsgEvent( // C-9:删除通知只走 KAFKA:msgschd 不再发 tombstone
target = Targets.KAFKA_SCHD, // value 形态("deleted":true 的 JSON)沿用现状,待 Q5 定稿
partitionKey = payload.flid,
eventType = EventType.TOMBSTONE,
stateVersion = current?.stateVersion ?: 0L,
payloadJson = mapper.writeValueAsString(
mapOf(
"flid" to payload.flid,
"stateVersion" to (current?.stateVersion ?: 0L),
"deleted" to true,
),
),
createdAt = createdAt,
),
MsgEvent( MsgEvent(
target = Targets.KAFKA_MSG, target = Targets.KAFKA_MSG,
partitionKey = payload.flid, partitionKey = payload.flid,
@@ -2,7 +2,7 @@ package com.gzzn.omms.msgexchange.jobs
import com.gzzn.omms.msgexchange.config.HistoryProps import com.gzzn.omms.msgexchange.config.HistoryProps
import com.gzzn.omms.msgexchange.config.OperationDayProps import com.gzzn.omms.msgexchange.config.OperationDayProps
import com.gzzn.omms.msgexchange.domain.EventType import com.gzzn.omms.msgexchange.domain.Targets
import com.gzzn.omms.msgexchange.domain.flight.FlightState import com.gzzn.omms.msgexchange.domain.flight.FlightState
import com.gzzn.omms.msgexchange.domain.flight.HistoryCandidate import com.gzzn.omms.msgexchange.domain.flight.HistoryCandidate
import com.gzzn.omms.msgexchange.infra.stub.StubFlightState import com.gzzn.omms.msgexchange.infra.stub.StubFlightState
@@ -88,7 +88,7 @@ class HistorySweepJobTest {
} }
@Test @Test
fun `never-fdel lifecycle purge emits tombstone before deletion`() { fun `never-fdel lifecycle purge emits msg deletion notice before deletion`() {
val f = seededFlight("F3", deleted = false, idleDays = 30) // 还在用,但已经静默超过兜底期限,命中清理条件 val f = seededFlight("F3", deleted = false, idleDays = 30) // 还在用,但已经静默超过兜底期限,命中清理条件
val events = StubMsgEvents() val events = StubMsgEvents()
val store = RecordingHistoryStore() val store = RecordingHistoryStore()
@@ -98,15 +98,15 @@ class HistorySweepJobTest {
job.run(now) job.run(now)
// 从没收到过 FDEL、被生命周期直接清掉的航班,删除前补发一次删除通知 // 从没收到过 FDEL、被生命周期直接清掉的航班,删除前补发一次 KAFKA:msg 删除通知(C-9
val tombstone = events.rows.values.single { it.eventType == EventType.TOMBSTONE } val notice = events.rows.values.single { it.target == Targets.KAFKA_MSG }
assertEquals("F3", tombstone.partitionKey) assertEquals("F3", notice.partitionKey)
assertTrue(tombstone.payloadJson.contains("\"deleted\":true")) assertTrue(notice.payloadJson.contains("\"deleted\":true"))
assertEquals(null, f.findMainRow("F3")) assertEquals(null, f.findMainRow("F3"))
} }
@Test @Test
fun `a flight updated after archival is not purged and emits no tombstone`() { fun `a flight updated after archival is not purged and emits no deletion notice`() {
val f = seededFlight("F4", deleted = false, idleDays = 30) val f = seededFlight("F4", deleted = false, idleDays = 30)
val events = StubMsgEvents() val events = StubMsgEvents()
val store = object : HistorySweepJob.HistoryStore { val store = object : HistorySweepJob.HistoryStore {
@@ -126,7 +126,7 @@ class HistorySweepJobTest {
assertEquals(1, outcome.archived) assertEquals(1, outcome.archived)
assertEquals(0, outcome.purged) // 归档后版本变了:本次不删 assertEquals(0, outcome.purged) // 归档后版本变了:本次不删
assertTrue(f.findMainRow("F4") != null) // 在用航班不许被清掉 assertTrue(f.findMainRow("F4") != null) // 在用航班不许被清掉
assertEquals(0, events.rows.size) // 也不登记 tombstone assertEquals(0, events.rows.size) // 也不登记删除通知
} }
@Test @Test
@@ -27,14 +27,14 @@ import java.util.UUID
/** /**
* 在真实 PostgreSQL 上验证历史清退的原子性:归档确认后若在删当前态之前失败, * 在真实 PostgreSQL 上验证历史清退的原子性:归档确认后若在删当前态之前失败,
* 同事务登记的 tombstone 必须一起回滚(`D1`)。 * 同事务登记的删除通知必须一起回滚(`D1`)。
*/ */
class HistorySweepPurgePgTest { class HistorySweepPurgePgTest {
private val t0: Instant = Instant.parse("2026-09-12T00:00:00Z") private val t0: Instant = Instant.parse("2026-09-12T00:00:00Z")
@Test @Test
fun `a failure after the tombstone insert rolls back the whole purge transaction`() { fun `a failure after the deletion-notice insert rolls back the whole purge transaction`() {
val ds = dataSource() val ds = dataSource()
val clock = Clock.fixed(t0, ZoneOffset.UTC) val clock = Clock.fixed(t0, ZoneOffset.UTC)
val flights = JdbcFlightStateRepository(ds, clock) val flights = JdbcFlightStateRepository(ds, clock)
@@ -46,7 +46,7 @@ class HistorySweepPurgePgTest {
now = t0.minusSeconds(30L * 86400), // 静默超过 idle-hours,命中生命周期清退 now = t0.minusSeconds(30L * 86400), // 静默超过 idle-hours,命中生命周期清退
) )
// 删除阶段注入失败:tombstone 已登记,删除抛错 → 整个事务必须回滚。 // 删除阶段注入失败:删除通知已登记,删除抛错 → 整个事务必须回滚。
val failing = object : FlightStateRepository by flights { val failing = object : FlightStateRepository by flights {
override fun purgeArchived(candidates: List<HistoryCandidate>): Int = override fun purgeArchived(candidates: List<HistoryCandidate>): Int =
throw IllegalStateException("purge-failed") throw IllegalStateException("purge-failed")
@@ -61,8 +61,8 @@ class HistorySweepPurgePgTest {
assertNotNull(flights.findMainRow(flid), "当前态不得被删") assertNotNull(flights.findMainRow(flid), "当前态不得被删")
assertTrue( assertTrue(
events.mergePendingSchd(t0, 1000).none { it.partitionKey == flid }, events.claimBatch(com.gzzn.omms.msgexchange.domain.Targets.KAFKA_MSG, 1000).none { it.partitionKey == flid },
"tombstone 必须随事务回滚,不留半状态", "删除通知必须随事务回滚,不留半状态",
) )
} }
@@ -57,7 +57,7 @@ class FdelAndAdftProcessorTest {
} }
@Test @Test
fun `active flight deletion publishes single tombstone and keeps details`() { fun `active flight deletion publishes single msg deletion notice and keeps details`() {
val f = flights() val f = flights()
val events = StubMsgEvents() val events = StubMsgEvents()
val proc = StubProcState() val proc = StubProcState()
@@ -69,10 +69,11 @@ class FdelAndAdftProcessorTest {
assertEquals(ApplyResult.Succeeded, result) assertEquals(ApplyResult.Succeeded, result)
assertEquals(FlightState.DELETED, f.findMainRow("121")!!.state) assertEquals(FlightState.DELETED, f.findMainRow("121")!!.state)
assertEquals(5L, f.findMainRow("121")!!.stateVersion) assertEquals(5L, f.findMainRow("121")!!.stateVersion)
val tombstones = events.rows.values.filter { it.eventType == EventType.TOMBSTONE } // C-9:删除通知只登记一条 KAFKA:msgschd 无 tombstone
assertEquals(1, tombstones.size) assertEquals(1, events.rows.size)
assertEquals(Targets.KAFKA_SCHD, tombstones.single().target) val notice = events.rows.values.single()
assertTrue(tombstones.single().payloadJson.contains("\"deleted\":true")) assertEquals(Targets.KAFKA_MSG, notice.target)
assertTrue(notice.payloadJson.contains("\"deleted\":true"))
// 终态和回填待办由处理器在自己的事务里写好 // 终态和回填待办由处理器在自己的事务里写好
val row = proc.find(msgId)!! val row = proc.find(msgId)!!
assertEquals(ProcStatus.SUCCEEDED, row.state) assertEquals(ProcStatus.SUCCEEDED, row.state)
@@ -91,7 +92,7 @@ class FdelAndAdftProcessorTest {
assertEquals(ApplyResult.Succeeded, result) // 重复的删除报文算成功,但不再动数据 assertEquals(ApplyResult.Succeeded, result) // 重复的删除报文算成功,但不再动数据
assertEquals(version, f.findMainRow("121")!!.stateVersion) assertEquals(version, f.findMainRow("121")!!.stateVersion)
assertEquals(1, events.rows.values.count { it.eventType == EventType.TOMBSTONE }) assertEquals(1, events.rows.values.count { it.target == Targets.KAFKA_MSG })
} }
@Test @Test