fix(delivery,health): U14 deleteOf 结构化 JSON + U12 健康检查真实 ping(ACM2-10 复审核验项②③)

- U14(复审②):deleteOf 字符串模板在 partitionKey=null/非空时均产出非法 JSON,
  改 Jackson 结构化序列化(null→null 字面量,非空→带引号转义);
  DispatcherTickTest 新增 phase B 两样例(readTree 合法性 + refs 精确断言)。
  注:该分支仅 Phase B 且 ES 投递成功后可达,当前属潜伏缺陷。
- U12(复审③):Redis/Kafka 健康指示器由仅查 bean 存在升级为真实 ping
  (false/异常→DOWN;缺 bean 仍 DOWN)。FlightRedisClient/DeliveryPort 增 ping()
  (stub 返回 true;DeliveryPort 默认实现,真实 Kafka 实装时覆写为 metadata 校验)。
  ping 判定拆为可单测纯函数,HealthIndicatorsTest 覆盖 UP/DOWN/异常/缺 bean 四态。

验证:./gradlew test 37 个测试全绿(本地,未推送)。
This commit is contained in:
windyboy
2026-09-06 22:12:23 +08:00
parent 6de3fa6c9c
commit 6c887adf56
7 changed files with 121 additions and 18 deletions
@@ -1,5 +1,6 @@
package com.gzzn.omms.msgexchange.nextgen.delivery
import com.fasterxml.jackson.databind.ObjectMapper
import com.gzzn.omms.msgexchange.nextgen.config.PipelineProps
import com.gzzn.omms.msgexchange.nextgen.domain.ErrorClass
import com.gzzn.omms.msgexchange.nextgen.domain.EventStatus
@@ -21,6 +22,9 @@ interface DeliveryPort {
/** 阶段 B:Redis 投影写(仅阶段 BI5Delivery 阶段 A 不写 Redis)。 */
fun projectRedis(payloadJson: String)
/** 连通性探测(健康检查用);默认 true,真实 Kafka 实装时覆写为 producer metadata 校验。 */
fun ping(): Boolean = true
}
/**
@@ -115,7 +119,14 @@ class Dispatcher(
else -> error("unknown target $target")
}
private fun deleteOf(e: MsgEvent) = """{"op":"delete","refs":${e.partitionKey ?: ""}}"""
private val jsonMapper = ObjectMapper()
/**
* U14:删除事件 wire JSON 结构化序列化——refs 可空,产出必须恒为合法 JSON
* null → null 字面量;非空 → 带引号并转义)。禁止字符串模板拼接。
*/
private fun deleteOf(e: MsgEvent): String =
jsonMapper.writeValueAsString(linkedMapOf("op" to "delete", "refs" to e.partitionKey))
/**
* 流程 3 flushSchd:批上限 + 每 FLID 最新一态聚合(topic "schd"wire=FLTR JSON 数组)。
@@ -12,8 +12,9 @@ import org.reactivestreams.Publisher
/**
* U12(R05):阶段 A 关键依赖的自定义健康指示器——
* RedisflightInfo 权威存储)与 Kafka(投递端口)。依赖经 BeanProvider 可选解析:
* 缺 bean(如未用 stub 也未实装)时指示 DOWN 而非启动失败
* RedisflightInfo 权威存储)与 Kafka(投递端口)。经 BeanProvider 可选解析:
* 缺 bean(如未用 stub 也未实装)时指示 DOWN 而非启动失败
* UP 判据为真实 pingfalse/异常 → DOWN),而非仅 bean 存在(复审 P1 修正)。
*/
@Singleton
class FlightRedisHealthIndicator(
@@ -21,8 +22,7 @@ class FlightRedisHealthIndicator(
) : HealthIndicator {
override fun getResult(): Publisher<HealthResult> =
Publishers.just(resultOf("redis-flight-store", "flight store", redis.isPresent,
redis.isPresent.takeIf { it }?.let { redis.get()::class.java.simpleName }))
Publishers.just(redisHealth(if (redis.isPresent) redis.get() else null))
}
@Singleton
@@ -31,18 +31,25 @@ class KafkaDeliveryHealthIndicator(
) : HealthIndicator {
override fun getResult(): Publisher<HealthResult> =
Publishers.just(resultOf("kafka-delivery", "delivery port", port.isPresent,
port.isPresent.takeIf { it }?.let { port.get()::class.java.simpleName }))
Publishers.just(kafkaHealth(if (port.isPresent) port.get() else null))
}
private fun resultOf(name: String, what: String, up: Boolean, impl: String?): HealthResult {
val builder = HealthResult.builder(name)
.status(if (up) HealthStatus.UP else HealthStatus.DOWN)
.details(
mapOf(
"message" to if (up) "$what present" else "$what bean missing (stub off, impl pending)",
"implementation" to (impl ?: "none"),
),
)
return builder.build()
/** ping 判定独立成纯函数便于单测:client 为 null = bean 缺失;ping false/异常 = DOWN。 */
internal fun redisHealth(client: FlightRedisClient?): HealthResult =
healthOf("redis-flight-store", "flight store", client?.let {
try { it.ping() } catch (e: Exception) { false }
})
internal fun kafkaHealth(port: DeliveryPort?): HealthResult =
healthOf("kafka-delivery", "delivery port", port?.let {
try { it.ping() } catch (e: Exception) { false }
})
private fun healthOf(name: String, what: String, pingOk: Boolean?): HealthResult {
val (status, message) = when (pingOk) {
null -> HealthStatus.DOWN to "$what bean missing (stub off, impl pending)"
false -> HealthStatus.DOWN to "$what ping failed"
true -> HealthStatus.UP to "$what ping ok"
}
return HealthResult.builder(name).status(status).details(mapOf("message" to message)).build()
}
@@ -21,6 +21,9 @@ interface FlightRedisClient {
/** 阶段 A 权威读(处理决策 loadState、3:30 清场 findAll)。 */
fun hgetAllFlightInfo(): Map<String, String>
/** 连通性探测(健康检查用;实现必须为快速调用,失败返回 false 而非抛出穿出)。 */
fun ping(): Boolean
}
// TODO(阶段1后续): 基于 micronaut-redis-lettuce 的实装(脚本自 classpath 装载并缓存 SHA)。
@@ -49,6 +49,8 @@ class StubRedis : FlightRedisClient {
}
override fun hgetAllFlightInfo(): Map<String, String> = hash.toMap()
override fun ping(): Boolean = true
}
@Requires(property = "msgx.stubs", value = "true")
@@ -56,7 +56,9 @@ class DispatcherTickTest {
attempts = attempts ?: events[i].attempts)
}
override fun insertSync(events: List<MsgEvent>) = Unit
override fun insertSync(events: List<MsgEvent>) { syncInserted += events }
val syncInserted = mutableListOf<MsgEvent>()
}
private class FakePort : DeliveryPort {
@@ -181,4 +183,24 @@ class DispatcherTickTest {
assertEquals(1, repo.events.first { it.eventId == 1L }.attempts) // schd 进退避
assertEquals(EventStatus.PENDING, repo.events.first { it.eventId == 1L }.state)
}
@Test
fun `delete events always carry legal JSON refs - null and non-null partitionKey`() {
val props = PipelineProps().apply { phase = PipelineProps.Phase.B }
val repo = FakeRepo()
val port = FakePort()
val d = dispatcher(repo, port, props)
repo.enqueue(ev(1, Targets.ES_FLIGHT_HTS, "F1", """{"h":1}"""))
repo.enqueue(ev(2, Targets.ES_FLIGHT_HTS, null, """{"h":2}"""))
d.tick(); d.tick() // 每 target 每 tick 仅出队一条:两条 ES 事件需两个 tick
val deletes = repo.syncInserted.filter { it.target == Targets.REDIS_FLIGHT_INFO }
assertEquals(2, deletes.size)
val mapper = com.fasterxml.jackson.databind.ObjectMapper()
val parsed = deletes.map { mapper.readTree(it.payloadJson) } // readTree 即合法性断言
assertEquals(listOf("delete", "delete"), parsed.map { it.get("op").asText() })
assertEquals("F1", parsed[0].get("refs").asText()) // 非空 → 带引号字符串
assertTrue(parsed[1].get("refs").isNull) // null → null 字面量
}
}
@@ -0,0 +1,57 @@
package com.gzzn.omms.msgexchange.nextgen.infra.health
import com.gzzn.omms.msgexchange.nextgen.delivery.DeliveryPort
import com.gzzn.omms.msgexchange.nextgen.infra.redis.FlightRedisClient
import com.gzzn.omms.msgexchange.nextgen.infra.redis.RedisScript
import io.micronaut.health.HealthStatus
import org.junit.jupiter.api.Assertions.assertEquals
import org.junit.jupiter.api.Test
/**
* U12:健康指示器 UP 判据为真实 ping(复审 P1 修正)——
* bean 缺失 / ping false / ping 抛异常 → DOWN;仅 ping 成功 → UP。
*/
class HealthIndicatorsTest {
private class FakeRedis(private val pingResult: Boolean) : FlightRedisClient {
override fun eval(script: RedisScript, setPairs: List<Pair<String, String>>, delFields: List<String>) = Unit
override fun hgetAllFlightInfo(): Map<String, String> = emptyMap()
override fun ping(): Boolean = pingResult
}
private class ThrowingRedis : FlightRedisClient {
override fun eval(script: RedisScript, setPairs: List<Pair<String, String>>, delFields: List<String>) = Unit
override fun hgetAllFlightInfo(): Map<String, String> = emptyMap()
override fun ping(): Boolean = throw RuntimeException("connection refused")
}
private class FakePort(private val pingResult: Boolean) : DeliveryPort {
override fun sendKafka(topic: String, payloadJson: String) = Unit
override fun indexFlightHts(payloadJson: String) = Unit
override fun projectRedis(payloadJson: String) = Unit
override fun ping(): Boolean = pingResult
}
private class ThrowingPort : DeliveryPort {
override fun sendKafka(topic: String, payloadJson: String) = Unit
override fun indexFlightHts(payloadJson: String) = Unit
override fun projectRedis(payloadJson: String) = Unit
override fun ping(): Boolean = throw RuntimeException("metadata fetch failed")
}
@Test
fun `redis up only when ping succeeds`() {
assertEquals(HealthStatus.UP, redisHealth(FakeRedis(true)).status)
assertEquals(HealthStatus.DOWN, redisHealth(FakeRedis(false)).status)
assertEquals(HealthStatus.DOWN, redisHealth(ThrowingRedis()).status)
assertEquals(HealthStatus.DOWN, redisHealth(null).status) // bean 缺失
}
@Test
fun `kafka up only when ping succeeds`() {
assertEquals(HealthStatus.UP, kafkaHealth(FakePort(true)).status)
assertEquals(HealthStatus.DOWN, kafkaHealth(FakePort(false)).status)
assertEquals(HealthStatus.DOWN, kafkaHealth(ThrowingPort()).status)
assertEquals(HealthStatus.DOWN, kafkaHealth(null).status) // bean 缺失
}
}
@@ -138,6 +138,7 @@ class MessageProcessorTest {
private object FakeRedis : FlightRedisClient {
override fun eval(script: RedisScript, setPairs: List<Pair<String, String>>, delFields: List<String>) = Unit
override fun hgetAllFlightInfo(): Map<String, String> = emptyMap()
override fun ping(): Boolean = true
}
private class FakeRefData : RefDataRepository {