From 15c035182f15e9a7b81ff8614af3d0d326c24aff Mon Sep 17 00:00:00 2001 From: windyboy Date: Sat, 12 Sep 2026 20:45:41 +0800 Subject: [PATCH] test(delivery): verify schd upsert on real postgres MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增 JdbcMsgEventUpsertPgTest(Testcontainers 本地 PG):单行只进不退、旧版本被拒、同版本 tombstone 胜出且不可复活、条件确认旧代次 0 行/当前代次 1 行、DEAD 行被更高代次替换后单行 PENDING 且错误清空。此前该 SQL 无任何执行覆盖。 --- .../jdbc/JdbcMsgEventUpsertPgTest.kt | 96 +++++++++++++++++++ 1 file changed, 96 insertions(+) create mode 100644 src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcMsgEventUpsertPgTest.kt diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcMsgEventUpsertPgTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcMsgEventUpsertPgTest.kt new file mode 100644 index 0000000..fdde8b5 --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/infra/persistence/jdbc/JdbcMsgEventUpsertPgTest.kt @@ -0,0 +1,96 @@ +package com.gzzn.omms.msgexchange.infra.persistence.jdbc + +import com.gzzn.omms.msgexchange.domain.ErrorClass +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.support.PgTestSupport +import com.zaxxer.hikari.HikariConfig +import com.zaxxer.hikari.HikariDataSource +import org.flywaydb.core.Flyway +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertNull +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Assumptions.assumeTrue +import org.junit.jupiter.api.Test +import java.time.Clock +import java.time.Instant +import java.util.UUID + +/** + * 在真实 PostgreSQL 上验证 `KAFKA:schd` 的单行 upsert 与条件确认 SQL 行为 + * (Testcontainers 本地容器;拿不到 PG 时跳过,不假装通过)。 + */ +class JdbcMsgEventUpsertPgTest { + + private val now: Instant = Instant.parse("2026-09-12T00:00:00Z") + + @Test + fun `schd upsert keeps one row, guards versions and confirms by generation`() { + val repo = JdbcMsgEventRepository(dataSource(), Clock.systemUTC()) + val key = "PG-" + UUID.randomUUID().toString().take(8) + + val firstId = repo.insertAll(listOf(schd(key, 1, """{"v":1}"""))).single() + repo.insertAll(listOf(schd(key, 2, """{"v":2}"""))) + assertEquals(emptyList(), repo.insertAll(listOf(schd(key, 1, """{"v":"late-old"}""")))) // 旧版本被拒 + + val latest = pendingSchd(repo, key).single() + assertEquals(2, latest.stateVersion) + assertTrue(latest.payloadJson.contains("\"v\":2")) + + // 同版本 tombstone 胜出;后续同版本 UPSERT 不得复活 + val tombstoneId = repo.insertAll(listOf(schd(key, 2, """{"flid":"x","deleted":true}""", EventType.TOMBSTONE))).single() + assertEquals(EventType.TOMBSTONE, pendingSchd(repo, key).single().eventType) + repo.insertAll(listOf(schd(key, 2, """{"v":"again"}"""))) + assertEquals(EventType.TOMBSTONE, pendingSchd(repo, key).single().eventType) + + // 条件确认:旧代次影响 0 行,当前代次标记成功 + assertEquals(0, repo.markSentIfVersion(firstId, 1)) + assertEquals(1, repo.markSentIfVersion(tombstoneId, 2)) + assertEquals(0, pendingSchd(repo, key).size) + } + + @Test + fun `a higher generation replaces a dead schd row and clears the old error`() { + val repo = JdbcMsgEventRepository(dataSource(), Clock.systemUTC()) + val key = "PG-" + UUID.randomUUID().toString().take(8) + + val deadId = repo.insertAll(listOf(schd(key, 1, """{"v":1}"""))).single() + repo.markDead(deadId, ErrorClass.EXHAUSTED, "broker-down", attempts = 5) + + val newId = repo.insertAll(listOf(schd(key, 2, """{"v":2}"""))).single() + + assertTrue(newId != deadId, "接受新代次必须替换行主键(可变写代次)") + val row = pendingSchd(repo, key).single() // 仍是一行 + assertEquals(2, row.stateVersion) + assertEquals(0, row.attempts) + assertNull(row.errorClass) + assertNull(row.lastError) + } + + private fun pendingSchd(repo: JdbcMsgEventRepository, key: String): List = + repo.mergePendingSchd(now, 1000).filter { it.partitionKey == key } + + private fun schd(key: String, version: Long, payload: String, type: EventType = EventType.UPSERT) = + MsgEvent( + target = Targets.KAFKA_SCHD, partitionKey = key, eventType = type, + stateVersion = version, payloadJson = payload, createdAt = now, + ) + + private fun dataSource(): javax.sql.DataSource { + assumeTrue(PgTestSupport.canConnect(), PgTestSupport.skipMessage()) + Flyway.configure() + .dataSource(PgTestSupport.jdbcUrl, PgTestSupport.user, PgTestSupport.password) + .locations("classpath:db/migration") + .load() + .migrate() + return HikariDataSource( + HikariConfig().apply { + jdbcUrl = PgTestSupport.jdbcUrl + username = PgTestSupport.user + password = PgTestSupport.password + maximumPoolSize = 2 + }, + ) + } +}