From 6c887adf5693f63d9d72ce87c1136ed84888e38d Mon Sep 17 00:00:00 2001 From: windyboy Date: Sun, 6 Sep 2026 22:12:23 +0800 Subject: [PATCH] =?UTF-8?q?fix(delivery,health):=20U14=20deleteOf=20?= =?UTF-8?q?=E7=BB=93=E6=9E=84=E5=8C=96=20JSON=20+=20U12=20=E5=81=A5?= =?UTF-8?q?=E5=BA=B7=E6=A3=80=E6=9F=A5=E7=9C=9F=E5=AE=9E=20ping=EF=BC=88AC?= =?UTF-8?q?M2-10=20=E5=A4=8D=E5=AE=A1=E6=A0=B8=E9=AA=8C=E9=A1=B9=E2=91=A1?= =?UTF-8?q?=E2=91=A2=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 个测试全绿(本地,未推送)。 --- .../nextgen/delivery/Dispatcher.kt | 13 ++++- .../nextgen/infra/health/HealthIndicators.kt | 39 +++++++------ .../nextgen/infra/redis/RedisScripts.kt | 3 + .../nextgen/infra/stub/StubAdapters.kt | 2 + .../nextgen/delivery/DispatcherTickTest.kt | 24 +++++++- .../infra/health/HealthIndicatorsTest.kt | 57 +++++++++++++++++++ .../processing/MessageProcessorTest.kt | 1 + 7 files changed, 121 insertions(+), 18 deletions(-) create mode 100644 src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicatorsTest.kt diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/Dispatcher.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/Dispatcher.kt index f5e89fb..b1ac518 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/Dispatcher.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/Dispatcher.kt @@ -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 投影写(仅阶段 B;I5:Delivery 阶段 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 数组)。 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicators.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicators.kt index 387999e..ce160c2 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicators.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicators.kt @@ -12,8 +12,9 @@ import org.reactivestreams.Publisher /** * U12(R05):阶段 A 关键依赖的自定义健康指示器—— - * Redis(flightInfo 权威存储)与 Kafka(投递端口)。依赖经 BeanProvider 可选解析: - * 缺 bean(如未用 stub 也未实装)时指示 DOWN 而非启动失败。 + * Redis(flightInfo 权威存储)与 Kafka(投递端口)。经 BeanProvider 可选解析: + * 缺 bean(如未用 stub 也未实装)时指示 DOWN 而非启动失败; + * UP 判据为真实 ping(false/异常 → DOWN),而非仅 bean 存在(复审 P1 修正)。 */ @Singleton class FlightRedisHealthIndicator( @@ -21,8 +22,7 @@ class FlightRedisHealthIndicator( ) : HealthIndicator { override fun getResult(): Publisher = - 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 = - 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() } diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/redis/RedisScripts.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/redis/RedisScripts.kt index 4d516d2..17c0e77 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/redis/RedisScripts.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/redis/RedisScripts.kt @@ -21,6 +21,9 @@ interface FlightRedisClient { /** 阶段 A 权威读(处理决策 loadState、3:30 清场 findAll)。 */ fun hgetAllFlightInfo(): Map + + /** 连通性探测(健康检查用;实现必须为快速调用,失败返回 false 而非抛出穿出)。 */ + fun ping(): Boolean } // TODO(阶段1后续): 基于 micronaut-redis-lettuce 的实装(脚本自 classpath 装载并缓存 SHA)。 diff --git a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubAdapters.kt b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubAdapters.kt index 7598342..98e872b 100644 --- a/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubAdapters.kt +++ b/src/main/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/stub/StubAdapters.kt @@ -49,6 +49,8 @@ class StubRedis : FlightRedisClient { } override fun hgetAllFlightInfo(): Map = hash.toMap() + + override fun ping(): Boolean = true } @Requires(property = "msgx.stubs", value = "true") diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/DispatcherTickTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/DispatcherTickTest.kt index e5efc64..2b91c3e 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/DispatcherTickTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/delivery/DispatcherTickTest.kt @@ -56,7 +56,9 @@ class DispatcherTickTest { attempts = attempts ?: events[i].attempts) } - override fun insertSync(events: List) = Unit + override fun insertSync(events: List) { syncInserted += events } + + val syncInserted = mutableListOf() } 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 字面量 + } } diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicatorsTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicatorsTest.kt new file mode 100644 index 0000000..9b0a1bd --- /dev/null +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/infra/health/HealthIndicatorsTest.kt @@ -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>, delFields: List) = Unit + override fun hgetAllFlightInfo(): Map = emptyMap() + override fun ping(): Boolean = pingResult + } + + private class ThrowingRedis : FlightRedisClient { + override fun eval(script: RedisScript, setPairs: List>, delFields: List) = Unit + override fun hgetAllFlightInfo(): Map = 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 缺失 + } +} diff --git a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt index 6df2f64..a324203 100644 --- a/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt +++ b/src/test/kotlin/com/gzzn/omms/msgexchange/nextgen/processing/MessageProcessorTest.kt @@ -138,6 +138,7 @@ class MessageProcessorTest { private object FakeRedis : FlightRedisClient { override fun eval(script: RedisScript, setPairs: List>, delFields: List) = Unit override fun hgetAllFlightInfo(): Map = emptyMap() + override fun ping(): Boolean = true } private class FakeRefData : RefDataRepository {