test(delivery): verify schd upsert on real postgres

新增 JdbcMsgEventUpsertPgTest(Testcontainers 本地 PG):单行只进不退、旧版本被拒、同版本 tombstone 胜出且不可复活、条件确认旧代次 0 行/当前代次 1 行、DEAD 行被更高代次替换后单行 PENDING 且错误清空。此前该 SQL 无任何执行覆盖。
This commit is contained in:
windyboy
2026-09-12 20:45:41 +08:00
parent 5beb67dadb
commit 15c035182f
@@ -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<Long>(), 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<MsgEvent> =
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
},
)
}
}