diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt index 37599c2..5d68d2e 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJob.kt @@ -68,7 +68,7 @@ class HistorySweepJob( val toPurge = candidates.filter { it.flid in archivedFlids } // 归档与删除之间主泵可能已写入同一 FLID:进锁事务后按 (FLID, STATE_VERSION) 复查, - // 仍合格才登记 tombstone 并删除;事件与删除同事务提交或回滚(D1、US-14 AC4)。 + // 仍合格才登记删除通知并删除;事件与删除同事务提交或回滚(D1、US-14 AC4)。 var purged = 0 txManager.inTransaction { lock.lock() @@ -78,7 +78,7 @@ class HistorySweepJob( // 从没收到过 FDEL、被这里直接清掉的航班,删除前补发一次通知,否则下游一直以为它还在 val preDelete = rechecked.filter { it.wasNeverFdel } if (preDelete.isNotEmpty()) { - msgEvents.insertAll(preDelete.map { tombstone(it, now) }) + msgEvents.insertAll(preDelete.map { deletionNotice(it, now) }) } purged = flightState.purgeArchived(rechecked) } @@ -86,10 +86,10 @@ class HistorySweepJob( return SweepOutcome(candidates.size, archivedFlids.size, purged, snapLogPurged) } - private fun tombstone(candidate: HistoryCandidate, createdAt: Instant) = MsgEvent( - target = Targets.KAFKA_SCHD, + /** C-9:删除通知只走 KAFKA:msg;value 形态("deleted":true 的 JSON)沿用现状,待 Q5 定稿。 */ + private fun deletionNotice(candidate: HistoryCandidate, createdAt: Instant) = MsgEvent( + target = Targets.KAFKA_MSG, partitionKey = candidate.flid, - eventType = EventType.TOMBSTONE, stateVersion = candidate.stateVersion, payloadJson = """{"flid":"${candidate.flid}","stateVersion":${candidate.stateVersion},"deleted":true}""", createdAt = createdAt, diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt index a4f80cf..abde0d3 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/processing/DynamicProcessors.kt @@ -86,20 +86,8 @@ class FdelProcessor( val createdAt = clock.instant() msgEvents.insertAll( listOf( - MsgEvent( - target = Targets.KAFKA_SCHD, - 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, - ), + // C-9:删除通知只走 KAFKA:msg,schd 不再发 tombstone; + // value 形态("deleted":true 的 JSON)沿用现状,待 Q5 定稿 MsgEvent( target = Targets.KAFKA_MSG, partitionKey = payload.flid, diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt index 53dd4ea..f1f13c4 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepJobTest.kt @@ -2,7 +2,7 @@ package com.gzzn.omms.msgexchange.jobs import com.gzzn.omms.msgexchange.config.HistoryProps 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.HistoryCandidate import com.gzzn.omms.msgexchange.infra.stub.StubFlightState @@ -88,7 +88,7 @@ class HistorySweepJobTest { } @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 events = StubMsgEvents() val store = RecordingHistoryStore() @@ -98,15 +98,15 @@ class HistorySweepJobTest { job.run(now) - // 从没收到过 FDEL、被生命周期直接清掉的航班,删除前要补发一次删除通知 - val tombstone = events.rows.values.single { it.eventType == EventType.TOMBSTONE } - assertEquals("F3", tombstone.partitionKey) - assertTrue(tombstone.payloadJson.contains("\"deleted\":true")) + // 从没收到过 FDEL、被生命周期直接清掉的航班,删除前补发一次 KAFKA:msg 删除通知(C-9) + val notice = events.rows.values.single { it.target == Targets.KAFKA_MSG } + assertEquals("F3", notice.partitionKey) + assertTrue(notice.payloadJson.contains("\"deleted\":true")) assertEquals(null, f.findMainRow("F3")) } @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 events = StubMsgEvents() val store = object : HistorySweepJob.HistoryStore { @@ -126,7 +126,7 @@ class HistorySweepJobTest { assertEquals(1, outcome.archived) assertEquals(0, outcome.purged) // 归档后版本变了:本次不删 assertTrue(f.findMainRow("F4") != null) // 在用航班不许被清掉 - assertEquals(0, events.rows.size) // 也不登记 tombstone + assertEquals(0, events.rows.size) // 也不登记删除通知 } @Test diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepPurgePgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepPurgePgTest.kt index bc569aa..46a6007 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepPurgePgTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/jobs/HistorySweepPurgePgTest.kt @@ -27,14 +27,14 @@ import java.util.UUID /** * 在真实 PostgreSQL 上验证历史清退的原子性:归档确认后若在删当前态之前失败, - * 同事务登记的 tombstone 必须一起回滚(`D1`)。 + * 同事务登记的删除通知必须一起回滚(`D1`)。 */ class HistorySweepPurgePgTest { private val t0: Instant = Instant.parse("2026-09-12T00:00:00Z") @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 clock = Clock.fixed(t0, ZoneOffset.UTC) val flights = JdbcFlightStateRepository(ds, clock) @@ -46,7 +46,7 @@ class HistorySweepPurgePgTest { now = t0.minusSeconds(30L * 86400), // 静默超过 idle-hours,命中生命周期清退 ) - // 删除阶段注入失败:tombstone 已登记,删除抛错 → 整个事务必须回滚。 + // 删除阶段注入失败:删除通知已登记,删除抛错 → 整个事务必须回滚。 val failing = object : FlightStateRepository by flights { override fun purgeArchived(candidates: List): Int = throw IllegalStateException("purge-failed") @@ -61,8 +61,8 @@ class HistorySweepPurgePgTest { assertNotNull(flights.findMainRow(flid), "当前态不得被删") assertTrue( - events.mergePendingSchd(t0, 1000).none { it.partitionKey == flid }, - "tombstone 必须随事务回滚,不留半状态", + events.claimBatch(com.gzzn.omms.msgexchange.domain.Targets.KAFKA_MSG, 1000).none { it.partitionKey == flid }, + "删除通知必须随事务回滚,不留半状态", ) } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt index 5cd4ab1..423763b 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/processing/FdelAndAdftProcessorTest.kt @@ -57,7 +57,7 @@ class FdelAndAdftProcessorTest { } @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 events = StubMsgEvents() val proc = StubProcState() @@ -69,10 +69,11 @@ class FdelAndAdftProcessorTest { assertEquals(ApplyResult.Succeeded, result) assertEquals(FlightState.DELETED, f.findMainRow("121")!!.state) assertEquals(5L, f.findMainRow("121")!!.stateVersion) - val tombstones = events.rows.values.filter { it.eventType == EventType.TOMBSTONE } - assertEquals(1, tombstones.size) - assertEquals(Targets.KAFKA_SCHD, tombstones.single().target) - assertTrue(tombstones.single().payloadJson.contains("\"deleted\":true")) + // C-9:删除通知只登记一条 KAFKA:msg,schd 无 tombstone + assertEquals(1, events.rows.size) + val notice = events.rows.values.single() + assertEquals(Targets.KAFKA_MSG, notice.target) + assertTrue(notice.payloadJson.contains("\"deleted\":true")) // 终态和回填待办由处理器在自己的事务里写好 val row = proc.find(msgId)!! assertEquals(ProcStatus.SUCCEEDED, row.state) @@ -91,7 +92,7 @@ class FdelAndAdftProcessorTest { assertEquals(ApplyResult.Succeeded, result) // 重复的删除报文算成功,但不再动数据 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