refactor(server,cli): shared payload-discriminator helper, single-read digest reuse, DOT label escaping
This commit is contained in:
@@ -44,14 +44,16 @@ fun renderEventRows(rows: List<EventRow>): String {
|
|||||||
return (listOf(header) + lines).joinToString("\n")
|
return (listOf(header) + lines).joinToString("\n")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fun String.dotEscape(): String = replace("\\", "\\\\").replace("\"", "\\\"")
|
||||||
|
|
||||||
fun renderDot(rows: List<EventRow>): String = buildString {
|
fun renderDot(rows: List<EventRow>): String = buildString {
|
||||||
appendLine("digraph session {")
|
appendLine("digraph session {")
|
||||||
rows.forEach { r ->
|
rows.forEach { r ->
|
||||||
appendLine(""""${r.eventId}" [label="${r.type} #${r.sequence}"]""")
|
appendLine(""""${r.eventId.dotEscape()}" [label="${r.type.dotEscape()} #${r.sequence}"]""")
|
||||||
}
|
}
|
||||||
rows.forEach { r ->
|
rows.forEach { r ->
|
||||||
r.causationId?.let { causation ->
|
r.causationId?.let { causation ->
|
||||||
appendLine(""""${r.eventId}" -> "$causation"""")
|
appendLine(""""${r.eventId.dotEscape()}" -> "${causation.dotEscape()}"""")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
append("}")
|
append("}")
|
||||||
|
|||||||
+9
-9
@@ -15,12 +15,12 @@ import com.correx.core.events.events.TransitionExecutedEvent
|
|||||||
import com.correx.core.events.events.WorkflowCompletedEvent
|
import com.correx.core.events.events.WorkflowCompletedEvent
|
||||||
import com.correx.core.events.events.WorkflowFailedEvent
|
import com.correx.core.events.events.WorkflowFailedEvent
|
||||||
import com.correx.core.events.events.WorkflowStartedEvent
|
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.serialization.eventJson
|
||||||
import com.correx.core.events.stores.EventStore
|
import com.correx.core.events.stores.EventStore
|
||||||
import com.correx.core.events.types.SessionId
|
import com.correx.core.events.types.SessionId
|
||||||
import kotlinx.serialization.Serializable
|
import kotlinx.serialization.Serializable
|
||||||
import kotlinx.serialization.json.jsonObject
|
|
||||||
import kotlinx.serialization.json.jsonPrimitive
|
|
||||||
import java.security.MessageDigest
|
import java.security.MessageDigest
|
||||||
|
|
||||||
@Serializable
|
@Serializable
|
||||||
@@ -46,11 +46,12 @@ class ReplayInspectionService(private val eventStore: EventStore) {
|
|||||||
val events = eventStore.read(sessionId)
|
val events = eventStore.read(sessionId)
|
||||||
val timeline = events.map { stored ->
|
val timeline = events.map { stored ->
|
||||||
val element = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload)
|
val element = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload)
|
||||||
val type = element.jsonObject["type"]?.jsonPrimitive?.content ?: "unknown"
|
toTimelineEntry(stored.sequence, stored.payload, payloadDiscriminator(element))
|
||||||
toTimelineEntry(stored.sequence, stored.payload, type)
|
|
||||||
}
|
}
|
||||||
val digest = computeDigest(sessionId)
|
val digest = computeDigest(events)
|
||||||
val digest2 = computeDigest(sessionId)
|
// 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(
|
return ReplayReport(
|
||||||
sessionId = sessionId.value,
|
sessionId = sessionId.value,
|
||||||
eventCount = events.size.toLong(),
|
eventCount = events.size.toLong(),
|
||||||
@@ -79,12 +80,11 @@ class ReplayInspectionService(private val eventStore: EventStore) {
|
|||||||
else -> TimelineEntry(sequence, type, null, type)
|
else -> TimelineEntry(sequence, type, null, type)
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun computeDigest(sessionId: SessionId): String {
|
private fun computeDigest(events: List<StoredEvent>): String {
|
||||||
val events = eventStore.read(sessionId)
|
|
||||||
val digest = MessageDigest.getInstance("SHA-256")
|
val digest = MessageDigest.getInstance("SHA-256")
|
||||||
for (event in events) {
|
for (event in events) {
|
||||||
val element = eventJson.encodeToJsonElement(EventPayload.serializer(), event.payload)
|
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 payloadJsonString = eventJson.encodeToString(EventPayload.serializer(), event.payload)
|
||||||
val line = "${event.sequence}|$type|$payloadJsonString"
|
val line = "${event.sequence}|$type|$payloadJsonString"
|
||||||
digest.update(line.toByteArray(Charsets.UTF_8))
|
digest.update(line.toByteArray(Charsets.UTF_8))
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package com.correx.apps.server.routes
|
|||||||
import com.correx.apps.server.ServerModule
|
import com.correx.apps.server.ServerModule
|
||||||
import com.correx.apps.server.protocol.SessionConfigDto
|
import com.correx.apps.server.protocol.SessionConfigDto
|
||||||
import com.correx.apps.server.replay.ReplayInspectionService
|
import com.correx.apps.server.replay.ReplayInspectionService
|
||||||
|
import com.correx.apps.server.serialization.payloadDiscriminator
|
||||||
import com.correx.apps.server.ws.SessionStreamHandler
|
import com.correx.apps.server.ws.SessionStreamHandler
|
||||||
import com.correx.core.events.events.EventPayload
|
import com.correx.core.events.events.EventPayload
|
||||||
import com.correx.core.events.serialization.eventJson
|
import com.correx.core.events.serialization.eventJson
|
||||||
@@ -18,8 +19,6 @@ import io.ktor.server.routing.route
|
|||||||
import io.ktor.server.websocket.webSocket
|
import io.ktor.server.websocket.webSocket
|
||||||
import kotlinx.serialization.Serializable
|
import kotlinx.serialization.Serializable
|
||||||
import kotlinx.serialization.json.JsonElement
|
import kotlinx.serialization.json.JsonElement
|
||||||
import kotlinx.serialization.json.JsonObject
|
|
||||||
import kotlinx.serialization.json.jsonPrimitive
|
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import java.util.*
|
import java.util.*
|
||||||
|
|
||||||
@@ -164,7 +163,7 @@ private fun Route.getEventsRoute(module: ServerModule) {
|
|||||||
.asSequence()
|
.asSequence()
|
||||||
.map { stored ->
|
.map { stored ->
|
||||||
val payloadJson = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload)
|
val payloadJson = eventJson.encodeToJsonElement(EventPayload.serializer(), stored.payload)
|
||||||
val type = (payloadJson as? JsonObject)?.get("type")?.jsonPrimitive?.content ?: "unknown"
|
val type = payloadDiscriminator(payloadJson)
|
||||||
EventRow(
|
EventRow(
|
||||||
sequence = stored.sequence,
|
sequence = stored.sequence,
|
||||||
eventId = stored.metadata.eventId.value,
|
eventId = stored.metadata.eventId.value,
|
||||||
|
|||||||
+12
@@ -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
|
||||||
Reference in New Issue
Block a user