From 75703ec0c816d9a421aabdfebe79b7edd0680290 Mon Sep 17 00:00:00 2001 From: kami Date: Fri, 12 Jun 2026 19:23:05 +0400 Subject: [PATCH] refactor(server,cli): shared payload-discriminator helper, single-read digest reuse, DOT label escaping --- .../correx/apps/cli/commands/EventsCommand.kt | 6 ++++-- .../server/replay/ReplayInspectionService.kt | 18 +++++++++--------- .../correx/apps/server/routes/SessionRoutes.kt | 5 ++--- .../serialization/PayloadDiscriminator.kt | 12 ++++++++++++ 4 files changed, 27 insertions(+), 14 deletions(-) create mode 100644 apps/server/src/main/kotlin/com/correx/apps/server/serialization/PayloadDiscriminator.kt diff --git a/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/EventsCommand.kt b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/EventsCommand.kt index 8b8ead5b..d6ddee25 100644 --- a/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/EventsCommand.kt +++ b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/EventsCommand.kt @@ -44,14 +44,16 @@ fun renderEventRows(rows: List): String { return (listOf(header) + lines).joinToString("\n") } +private fun String.dotEscape(): String = replace("\\", "\\\\").replace("\"", "\\\"") + fun renderDot(rows: List): String = buildString { appendLine("digraph session {") rows.forEach { r -> - appendLine(""""${r.eventId}" [label="${r.type} #${r.sequence}"]""") + appendLine(""""${r.eventId.dotEscape()}" [label="${r.type.dotEscape()} #${r.sequence}"]""") } rows.forEach { r -> r.causationId?.let { causation -> - appendLine(""""${r.eventId}" -> "$causation"""") + appendLine(""""${r.eventId.dotEscape()}" -> "${causation.dotEscape()}"""") } } append("}") diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/replay/ReplayInspectionService.kt b/apps/server/src/main/kotlin/com/correx/apps/server/replay/ReplayInspectionService.kt index 28245e46..9dbc4d84 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/replay/ReplayInspectionService.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/replay/ReplayInspectionService.kt @@ -15,12 +15,12 @@ import com.correx.core.events.events.TransitionExecutedEvent import com.correx.core.events.events.WorkflowCompletedEvent import com.correx.core.events.events.WorkflowFailedEvent import com.correx.core.events.events.WorkflowStartedEvent +import com.correx.apps.server.serialization.payloadDiscriminator +import com.correx.core.events.events.StoredEvent import com.correx.core.events.serialization.eventJson import com.correx.core.events.stores.EventStore import com.correx.core.events.types.SessionId import kotlinx.serialization.Serializable -import kotlinx.serialization.json.jsonObject -import kotlinx.serialization.json.jsonPrimitive import java.security.MessageDigest @Serializable @@ -46,11 +46,12 @@ class ReplayInspectionService(private val eventStore: EventStore) { val events = eventStore.read(sessionId) val timeline = events.map { stored -> val element = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload) - val type = element.jsonObject["type"]?.jsonPrimitive?.content ?: "unknown" - toTimelineEntry(stored.sequence, stored.payload, type) + toTimelineEntry(stored.sequence, stored.payload, payloadDiscriminator(element)) } - val digest = computeDigest(sessionId) - val digest2 = computeDigest(sessionId) + val digest = computeDigest(events) + // Second digest from an independent store read — catches store-order or + // serialization nondeterminism; must NOT reuse the already-loaded list. + val digest2 = computeDigest(eventStore.read(sessionId)) return ReplayReport( sessionId = sessionId.value, eventCount = events.size.toLong(), @@ -79,12 +80,11 @@ class ReplayInspectionService(private val eventStore: EventStore) { else -> TimelineEntry(sequence, type, null, type) } - private fun computeDigest(sessionId: SessionId): String { - val events = eventStore.read(sessionId) + private fun computeDigest(events: List): String { val digest = MessageDigest.getInstance("SHA-256") for (event in events) { val element = eventJson.encodeToJsonElement(EventPayload.serializer(), event.payload) - val type = element.jsonObject["type"]?.jsonPrimitive?.content ?: "unknown" + val type = payloadDiscriminator(element) val payloadJsonString = eventJson.encodeToString(EventPayload.serializer(), event.payload) val line = "${event.sequence}|$type|$payloadJsonString" digest.update(line.toByteArray(Charsets.UTF_8)) diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt b/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt index c120b65b..b93ff63a 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/routes/SessionRoutes.kt @@ -3,6 +3,7 @@ package com.correx.apps.server.routes import com.correx.apps.server.ServerModule import com.correx.apps.server.protocol.SessionConfigDto import com.correx.apps.server.replay.ReplayInspectionService +import com.correx.apps.server.serialization.payloadDiscriminator import com.correx.apps.server.ws.SessionStreamHandler import com.correx.core.events.events.EventPayload import com.correx.core.events.serialization.eventJson @@ -18,8 +19,6 @@ import io.ktor.server.routing.route import io.ktor.server.websocket.webSocket import kotlinx.serialization.Serializable import kotlinx.serialization.json.JsonElement -import kotlinx.serialization.json.JsonObject -import kotlinx.serialization.json.jsonPrimitive import org.slf4j.LoggerFactory import java.util.* @@ -164,7 +163,7 @@ private fun Route.getEventsRoute(module: ServerModule) { .asSequence() .map { stored -> val payloadJson = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload) - val type = (payloadJson as? JsonObject)?.get("type")?.jsonPrimitive?.content ?: "unknown" + val type = payloadDiscriminator(payloadJson) EventRow( sequence = stored.sequence, eventId = stored.metadata.eventId.value, diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/serialization/PayloadDiscriminator.kt b/apps/server/src/main/kotlin/com/correx/apps/server/serialization/PayloadDiscriminator.kt new file mode 100644 index 00000000..06a958fe --- /dev/null +++ b/apps/server/src/main/kotlin/com/correx/apps/server/serialization/PayloadDiscriminator.kt @@ -0,0 +1,12 @@ +package com.correx.apps.server.serialization + +import kotlinx.serialization.json.JsonElement +import kotlinx.serialization.json.JsonObject +import kotlinx.serialization.json.jsonPrimitive + +const val UNKNOWN_PAYLOAD_TYPE = "unknown" + +// The polymorphic discriminator emitted by eventJson for every EventPayload — +// payloads always encode to a JsonObject carrying a "type" key. +fun payloadDiscriminator(element: JsonElement): String = + (element as? JsonObject)?.get("type")?.jsonPrimitive?.content ?: UNKNOWN_PAYLOAD_TYPE