From 7936251d6bf61c40fc7b83add39b4e4089850b74 Mon Sep 17 00:00:00 2001 From: kami Date: Fri, 29 May 2026 01:01:25 +0400 Subject: [PATCH] =?UTF-8?q?fix:=20complete=20all=20P2=20audit=20findings?= =?UTF-8?q?=20=E2=80=94=20dead=20code,=20NPE=20guard,=20hygiene?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P2-1: Guard SessionId NPE with safe-call in SessionsReducer P2-2: Remove 16 dead event classes + reducer branches; replace with live orchestration events in tests P2-3: Standardize ChatInput divergences — guard CancelSession, single-arg ProtocolError P2-4: Remove dead TuiToolRecord.diff and ToolDisplayStatus.REQUESTED P2-5: Remove dead streamLive from SessionEventBridge P2-6: Extract SessionsReducerContext data class for 6+ param method P2-7: Replace StageId("none") sentinel with TypeId.NONE P2-9: Add TypeId.random() factory on all type-alias IDs P2-10: Add KDoc on schemaVersion documenting reserved-for-migration P2-8: Verified RouterReducer string-template already correct (TypeId value class toString) --- .../apps/server/bridge/SessionEventBridge.kt | 8 -- .../apps/server/ws/GlobalStreamHandler.kt | 41 ++------- .../server/bridge/DomainEventMapperTest.kt | 8 +- .../server/bridge/SessionEventBridgeTest.kt | 57 ------------ .../apps/tui/components/EventHistoryStrip.kt | 1 - .../correx/apps/tui/reducer/RootReducer.kt | 14 +-- .../apps/tui/reducer/SessionsReducer.kt | 91 ++++++++++--------- .../com/correx/apps/tui/state/TuiState.kt | 3 +- .../apps/tui/reducer/SessionsReducerTest.kt | 11 ++- .../correx/core/artifacts/ArtifactReducer.kt | 36 -------- .../core/context/DefaultContextReducer.kt | 17 +--- .../core/events/events/ArtifactEvents.kt | 36 -------- .../core/events/events/ContextEvents.kt | 52 ----------- .../core/events/events/EventMetadata.kt | 5 + .../core/events/events/SessionEvents.kt | 35 ------- .../correx/core/events/events/StageEvents.kt | 8 -- .../events/serialization/Serialization.kt | 34 +------ .../kotlin/com/correx/core/utils/TypeId.kt | 9 ++ .../events/EventSerializationHardeningTest.kt | 28 +++--- .../core/router/RouterContextBuilder.kt | 2 +- .../com/correx/core/router/RouterFacade.kt | 2 +- .../core/sessions/DefaultSessionReducer.kt | 21 ----- .../artifact/LiveArtifactRepository.kt | 8 -- .../kotlin/SessionReplayDeterminismTest.kt | 18 ++-- .../src/test/kotlin/ArtifactReducerTest.kt | 55 ----------- .../src/test/kotlin/ContextProjectorTest.kt | 9 +- .../test/kotlin/DefaultContextReducerTest.kt | 62 ++++--------- .../test/kotlin/DefaultSessionReducerTest.kt | 76 ++++++++-------- .../src/test/kotlin/SessionProjectorTest.kt | 37 ++++---- .../src/test/kotlin/SessionReplayTest.kt | 26 ++++-- .../kotlin/TransitionReplayIntegrationTest.kt | 30 +++--- .../TransitionEventSerializationTest.kt | 8 +- 32 files changed, 240 insertions(+), 608 deletions(-) diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/bridge/SessionEventBridge.kt b/apps/server/src/main/kotlin/com/correx/apps/server/bridge/SessionEventBridge.kt index 758452d7..fd9fdb79 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/bridge/SessionEventBridge.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/bridge/SessionEventBridge.kt @@ -244,12 +244,4 @@ class SessionEventBridge( return orderedIds.mapNotNull { byInvocation[it] } } - suspend fun streamLive(sessionId: SessionId) { - var sessionSequence = 0L - eventStore.subscribe(sessionId).collect { event -> - sessionSequence++ - domainEventToServerMessage(event, artifactStore, sessionSequence = sessionSequence) - ?.let { send(it) } - } - } } diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt index 27449e3e..d233ece4 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/ws/GlobalStreamHandler.kt @@ -88,11 +88,7 @@ class GlobalStreamHandler(private val module: ServerModule) { } .onFailure { log.warn("decode error: {}", it.message) - val error = ServerMessage.ProtocolError( - message = "Unknown message: ${it.message}", - sequence = null, - sessionSequence = null, - ) + val error = ServerMessage.ProtocolError("Unknown message: ${it.message}") session.send(Frame.Text(ProtocolSerializer.encodeServerMessage(error))) } } @@ -144,7 +140,10 @@ class GlobalStreamHandler(private val module: ServerModule) { is ClientMessage.Ping -> Unit is ClientMessage.StartSession -> handleStartSession(session, msg, sendFrame) is ClientMessage.StartChatSession -> handleStartChatSession(session, msg, sendFrame) - is ClientMessage.CancelSession -> module.orchestrator.cancel(msg.sessionId) + is ClientMessage.CancelSession -> { + runCatching { module.orchestrator.cancel(msg.sessionId) } + .onFailure { log.error("cancel failed for session={}: {}", msg.sessionId.value, it.message, it) } + } is ClientMessage.ResumeSession -> session.send(encodeError("ResumeSession not supported")) is ClientMessage.ApprovalResponse -> handleApprovalResponse(msg, sendFrame) is ClientMessage.CreateGrant -> handleCreateGrant(msg, sendFrame) @@ -159,11 +158,7 @@ class GlobalStreamHandler(private val module: ServerModule) { )) }.onFailure { log.error("routerFacade.onUserInput failed: {}", it.message) - sendFrame(ServerMessage.ProtocolError( - message = "Router error: ${it.message}", - sequence = null, - sessionSequence = null, - )) + sendFrame(ServerMessage.ProtocolError("Router error: ${it.message}")) } } } @@ -176,34 +171,16 @@ class GlobalStreamHandler(private val module: ServerModule) { val scopeSessionId = module.approvalCoordinator.lookupSession(msg.requestId) if (scopeSessionId == null) { log.warn("handleApprovalResponse: no session found for requestId={}", msg.requestId.value) - sendFrame( - ServerMessage.ProtocolError( - message = "Unknown approval request", - sequence = null, - sessionSequence = null, - ), - ) + sendFrame(ServerMessage.ProtocolError("Unknown approval request")) return } module.approvalCoordinator.handleResponse(msg, scopeSessionId)?.let { sendFrame(it) } } - private fun errorResponse(message: String) = ServerMessage.ProtocolError( - message = message, - sequence = null, - sessionSequence = null, - ) + private fun errorResponse(message: String) = ServerMessage.ProtocolError(message) private fun encodeError(message: String): Frame.Text = - Frame.Text( - ProtocolSerializer.encodeServerMessage( - ServerMessage.ProtocolError( - message = message, - sequence = null, - sessionSequence = null, - ), - ), - ) + Frame.Text(ProtocolSerializer.encodeServerMessage(ServerMessage.ProtocolError(message))) private suspend fun handleCreateGrant( msg: ClientMessage.CreateGrant, diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/bridge/DomainEventMapperTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/bridge/DomainEventMapperTest.kt index 2174ec02..bf696c5f 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/bridge/DomainEventMapperTest.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/bridge/DomainEventMapperTest.kt @@ -10,7 +10,7 @@ import com.correx.core.events.events.InferenceCompletedEvent import com.correx.core.events.events.InferenceStartedEvent import com.correx.core.events.events.InferenceTimeoutEvent import com.correx.core.events.events.OrchestrationPausedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.SteeringNoteAddedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent import com.correx.core.events.events.StoredEvent @@ -371,7 +371,7 @@ class DomainEventMapperTest { @Test fun `unmapped event returns null`(): Unit = runTest { - val event = storedEvent(SessionStartedEvent(sessionId = sessionId)) + val event = storedEvent(SteeringNoteAddedEvent(sessionId = sessionId, content = "test")) val result = domainEventToServerMessage(event, noopStore, sessionSequence = 0L) assertNull(result) } @@ -392,12 +392,12 @@ class DomainEventMapperTest { val origLevel = logger.level logger.level = Level.DEBUG try { - val event = storedEvent(SessionStartedEvent(sessionId = sessionId)) + val event = storedEvent(SteeringNoteAddedEvent(sessionId = sessionId, content = "test")) val result = domainEventToServerMessage(event, noopStore, sessionSequence = 0L) assertNull(result) assertEquals(1, appender.events.size) assertEquals(Level.DEBUG, appender.events[0].level) - assertTrue(appender.events[0].message.formattedMessage.contains("SessionStartedEvent")) + assertTrue(appender.events[0].message.formattedMessage.contains("SteeringNoteAddedEvent")) } finally { logger.removeAppender(appender) appender.stop() diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/bridge/SessionEventBridgeTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/bridge/SessionEventBridgeTest.kt index 8361e737..c4dd5d82 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/bridge/SessionEventBridgeTest.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/bridge/SessionEventBridgeTest.kt @@ -25,10 +25,8 @@ import com.correx.core.kernel.orchestration.OrchestrationRepository import com.correx.core.tools.contract.Tool import com.correx.core.tools.registry.ToolRegistry import com.correx.core.sessions.projections.replay.EventReplayer -import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow -import kotlinx.coroutines.launch import kotlinx.coroutines.test.runTest import kotlinx.datetime.Instant import org.junit.jupiter.api.Assertions.assertEquals @@ -186,36 +184,6 @@ class SessionEventBridgeTest { assertEquals(ServerMessage.SnapshotComplete, sent[1]) } - @Test - fun `streamLive sends mapped messages from live flow`() = runTest { - val events = listOf( - storedEvent(WorkflowStartedEvent(sessionId, workflowId, stageId), seq = 1L), - storedEvent(WorkflowCompletedEvent(sessionId, stageId, 1), seq = 2L), - ) - val liveFlow: Flow = flow { events.forEach { emit(it) } } - val store = fakeEventStore(liveFlow = liveFlow) - val sent = mutableListOf() - val bridge = SessionEventBridge( - store, - noopArtifactStore, - activeOrchestrationRepository(), - noopWorkflowRegistry, - noopToolRegistry, - ) { sent.add(it) } - - bridge.streamLive(sessionId) - - assertEquals(2, sent.size) - assertEquals( - ServerMessage.SessionStarted(sessionId, workflowId, 1L, 1L), - sent[0], - ) - assertEquals( - ServerMessage.SessionCompleted(sessionId, 2L, 2L), - sent[1], - ) - } - @Test fun `replaySnapshot with no sessions emits only SnapshotComplete`() = runTest { val store = fakeEventStore(allEventsList = emptyList()) @@ -266,29 +234,4 @@ class SessionEventBridgeTest { assertEquals(3L, snapshot.lastSessionSequence) } - @Test - fun `streamLive cancellation stops collection`() = runTest { - val channel = Channel(Channel.UNLIMITED) - val liveFlow: Flow = flow { - for (event in channel) emit(event) - } - val store = fakeEventStore(liveFlow = liveFlow) - val sent = mutableListOf() - val bridge = SessionEventBridge( - store, - noopArtifactStore, - activeOrchestrationRepository(), - noopWorkflowRegistry, - noopToolRegistry, - ) { sent.add(it) } - - channel.send(storedEvent(WorkflowStartedEvent(sessionId, workflowId, stageId), seq = 1L)) - - val job = launch { bridge.streamLive(sessionId) } - channel.send(storedEvent(WorkflowCompletedEvent(sessionId, stageId, 1), seq = 2L)) - channel.close() - job.join() - - assertEquals(2, sent.size) - } } diff --git a/apps/tui/src/main/kotlin/com/correx/apps/tui/components/EventHistoryStrip.kt b/apps/tui/src/main/kotlin/com/correx/apps/tui/components/EventHistoryStrip.kt index 414ca3e4..64aef8a4 100644 --- a/apps/tui/src/main/kotlin/com/correx/apps/tui/components/EventHistoryStrip.kt +++ b/apps/tui/src/main/kotlin/com/correx/apps/tui/components/EventHistoryStrip.kt @@ -1,6 +1,5 @@ package com.correx.apps.tui.components -import com.correx.apps.tui.state.ToolDisplayStatus import com.correx.apps.tui.state.TuiState import dev.tamboui.style.Style import dev.tamboui.text.Line diff --git a/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/RootReducer.kt b/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/RootReducer.kt index 4ea6436b..f93010b8 100644 --- a/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/RootReducer.kt +++ b/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/RootReducer.kt @@ -62,12 +62,14 @@ object RootReducer { log.debug("sessions state before reducers={}", state.sessions) val (sessions, se) = SessionsReducer.reduce( - sessions = afterInput.sessions, - displayState = afterInput.displayState, - inputMode = prevInputMode, - inputText = prevInputBuffer, - action = action, - clock = clock, + SessionsReducerContext( + sessions = afterInput.sessions, + displayState = afterInput.displayState, + inputMode = prevInputMode, + inputText = prevInputBuffer, + action = action, + clock = clock, + ), ) log.debug("connection state before reducers={}", state.connection) val (connection, ce) = ConnectionReducer.reduce(afterInput.connection, action) diff --git a/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/SessionsReducer.kt b/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/SessionsReducer.kt index 8870f676..75e7ee6d 100644 --- a/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/SessionsReducer.kt +++ b/apps/tui/src/main/kotlin/com/correx/apps/tui/reducer/SessionsReducer.kt @@ -22,40 +22,42 @@ import java.time.format.DateTimeFormatter private val log = LoggerFactory.getLogger(SessionsReducer::class.simpleName) +data class SessionsReducerContext( + val sessions: SessionsState, + val displayState: DisplayState, + val inputMode: InputMode, + val inputText: String, + val action: Action, + val clock: () -> Long = System::currentTimeMillis, +) + object SessionsReducer { private val timeFormatter = DateTimeFormatter.ofPattern("HH:mm:ss").withZone(ZoneOffset.UTC) - fun reduce( - sessions: SessionsState, - displayState: DisplayState, - inputMode: InputMode, - inputText: String, - action: Action, - clock: () -> Long = System::currentTimeMillis, - ): Pair> = when (action) { + fun reduce(ctx: SessionsReducerContext): Pair> = when (ctx.action) { is Action.NavigateUp -> when { - inputMode == InputMode.FILTER -> navigateUp(sessions) to emptyList() - displayState == DisplayState.IDLE -> navigateUp(sessions) to emptyList() - displayState == DisplayState.IN_SESSION -> sessions to emptyList() - else -> sessions to emptyList() + ctx.inputMode == InputMode.FILTER -> navigateUp(ctx.sessions) to emptyList() + ctx.displayState == DisplayState.IDLE -> navigateUp(ctx.sessions) to emptyList() + ctx.displayState == DisplayState.IN_SESSION -> ctx.sessions to emptyList() + else -> ctx.sessions to emptyList() } is Action.NavigateDown -> when { - inputMode == InputMode.FILTER -> navigateDown(sessions) to emptyList() - displayState == DisplayState.IDLE -> navigateDown(sessions) to emptyList() - displayState == DisplayState.IN_SESSION -> sessions to emptyList() - else -> sessions to emptyList() + ctx.inputMode == InputMode.FILTER -> navigateDown(ctx.sessions) to emptyList() + ctx.displayState == DisplayState.IDLE -> navigateDown(ctx.sessions) to emptyList() + ctx.displayState == DisplayState.IN_SESSION -> ctx.sessions to emptyList() + else -> ctx.sessions to emptyList() } is Action.SubmitInput -> when { - inputMode == InputMode.FILTER -> sessions.copy(filter = inputText) to emptyList() - displayState == DisplayState.IDLE -> { - val text = inputText.trim() - val wfIndex = sessions.selectedWorkflowIndex + ctx.inputMode == InputMode.FILTER -> ctx.sessions.copy(filter = ctx.inputText) to emptyList() + ctx.displayState == DisplayState.IDLE -> { + val text = ctx.inputText.trim() + val wfIndex = ctx.sessions.selectedWorkflowIndex when { - wfIndex in sessions.workflows.indices -> { - val wf = sessions.workflows[wfIndex] - sessions.copy(selectedWorkflowIndex = -1) to listOf( + wfIndex in ctx.sessions.workflows.indices -> { + val wf = ctx.sessions.workflows[wfIndex] + ctx.sessions.copy(selectedWorkflowIndex = -1) to listOf( Effect.SendWs(ClientMessage.StartSession(workflowId = wf.workflowId, config = null)), ) } @@ -67,56 +69,56 @@ object SessionsReducer { status = "STARTING", workflowId = "chat", name = "chat", - lastEventAt = clock(), + lastEventAt = ctx.clock(), ) - sessions.copy( - sessions = sessions.sessions + optimistic, + ctx.sessions.copy( + sessions = ctx.sessions.sessions + optimistic, selectedId = sessionId.value, ) to listOf( Effect.SendWs(ClientMessage.StartChatSession(sessionId = sessionId, text = text)), ) } - else -> sessions to emptyList() + else -> ctx.sessions to emptyList() } } - displayState == DisplayState.IN_SESSION -> { - sessions to listOf( + ctx.displayState == DisplayState.IN_SESSION -> ctx.sessions.selectedId?.let { sid -> + ctx.sessions to listOf( Effect.SendWs( ClientMessage.ChatInput( - sessionId = SessionId(sessions.selectedId!!), - text = inputText, + sessionId = SessionId(sid), + text = ctx.inputText, mode = ChatMode.CHAT, ), ), ) - } + } ?: (ctx.sessions to emptyList()) - else -> sessions to emptyList() + else -> ctx.sessions to emptyList() } - is Action.CancelInput -> if (inputMode == InputMode.FILTER) { - sessions.copy(filter = "") to emptyList() + is Action.CancelInput -> if (ctx.inputMode == InputMode.FILTER) { + ctx.sessions.copy(filter = "") to emptyList() } else { - sessions to emptyList() + ctx.sessions to emptyList() } is Action.CancelSelectedSession -> { - val id = sessions.selectedId + val id = ctx.sessions.selectedId if (id != null) { - sessions to listOf(Effect.SendWs(ClientMessage.CancelSession(sessionId = SessionId(id)))) + ctx.sessions to listOf(Effect.SendWs(ClientMessage.CancelSession(sessionId = SessionId(id)))) } else { - sessions to emptyList() + ctx.sessions to emptyList() } } - is Action.ToggleWorkflows -> sessions.copy( - workflowsVisible = !sessions.workflowsVisible, + is Action.ToggleWorkflows -> ctx.sessions.copy( + workflowsVisible = !ctx.sessions.workflowsVisible, ) to emptyList() - is Action.ServerEventReceived -> applyServerMessage(sessions, action.message, clock) - else -> sessions to emptyList() + is Action.ServerEventReceived -> applyServerMessage(ctx.sessions, ctx.action.message, ctx.clock) + else -> ctx.sessions to emptyList() } private fun filteredSessions(sessions: SessionsState): List = @@ -311,7 +313,6 @@ object SessionsReducer { tier = tool.tier, status = displayStatus, argsPreview = null, - diff = tool.diff, ) } val summary = SessionSummary( @@ -409,7 +410,7 @@ object SessionsReducer { if (s.id == msg.sessionId.value) { val updatedTools = s.tools.map { t -> if (t.name == msg.toolName && t.status == ToolDisplayStatus.STARTED) { - t.copy(status = ToolDisplayStatus.COMPLETED, diff = msg.diff) + t.copy(status = ToolDisplayStatus.COMPLETED) } else { t } diff --git a/apps/tui/src/main/kotlin/com/correx/apps/tui/state/TuiState.kt b/apps/tui/src/main/kotlin/com/correx/apps/tui/state/TuiState.kt index c3b0a093..5ae4cdc5 100644 --- a/apps/tui/src/main/kotlin/com/correx/apps/tui/state/TuiState.kt +++ b/apps/tui/src/main/kotlin/com/correx/apps/tui/state/TuiState.kt @@ -5,7 +5,7 @@ import com.correx.apps.server.protocol.ServerMessage enum class InputMode { ROUTER, FILTER } -enum class ToolDisplayStatus { REQUESTED, STARTED, COMPLETED, FAILED, REJECTED } +enum class ToolDisplayStatus { STARTED, COMPLETED, FAILED, REJECTED } enum class ProviderType { LOCAL, REMOTE } @@ -14,7 +14,6 @@ data class TuiToolRecord( val tier: Int, val status: ToolDisplayStatus, val argsPreview: String?, - val diff: String? = null, ) data class TuiEventEntry( diff --git a/apps/tui/src/test/kotlin/com/correx/apps/tui/reducer/SessionsReducerTest.kt b/apps/tui/src/test/kotlin/com/correx/apps/tui/reducer/SessionsReducerTest.kt index b14e5c26..5128f111 100644 --- a/apps/tui/src/test/kotlin/com/correx/apps/tui/reducer/SessionsReducerTest.kt +++ b/apps/tui/src/test/kotlin/com/correx/apps/tui/reducer/SessionsReducerTest.kt @@ -41,7 +41,16 @@ class SessionsReducerTest { inputMode: InputMode = InputMode.ROUTER, inputText: String = "", action: Action, - ) = SessionsReducer.reduce(sessions, displayState, inputMode, inputText, action, fixedClock) + ) = SessionsReducer.reduce( + SessionsReducerContext( + sessions = sessions, + displayState = displayState, + inputMode = inputMode, + inputText = inputText, + action = action, + clock = fixedClock, + ), + ) @Test fun `NavigateUp wraps from top to bottom`() { diff --git a/core/artifacts/src/main/kotlin/com/correx/core/artifacts/ArtifactReducer.kt b/core/artifacts/src/main/kotlin/com/correx/core/artifacts/ArtifactReducer.kt index 7cd4d75e..83ea2712 100644 --- a/core/artifacts/src/main/kotlin/com/correx/core/artifacts/ArtifactReducer.kt +++ b/core/artifacts/src/main/kotlin/com/correx/core/artifacts/ArtifactReducer.kt @@ -1,11 +1,7 @@ package com.correx.core.artifacts import com.correx.core.artifacts.model.ArtifactRelationship -import com.correx.core.events.events.ArtifactArchivedEvent import com.correx.core.events.events.ArtifactCreatedEvent -import com.correx.core.events.events.ArtifactRejectedEvent -import com.correx.core.events.events.ArtifactRelationshipAddedEvent -import com.correx.core.events.events.ArtifactSupersededEvent import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatingEvent import com.correx.core.events.events.StoredEvent @@ -25,23 +21,6 @@ class DefaultArtifactReducer : ArtifactReducer { transition(state, ArtifactLifecyclePhase.VALIDATING, ArtifactLifecyclePhase.CREATED) is ArtifactValidatedEvent -> transition(state, ArtifactLifecyclePhase.VALIDATED, ArtifactLifecyclePhase.VALIDATING) - is ArtifactRejectedEvent -> - transition(state, ArtifactLifecyclePhase.REJECTED, ArtifactLifecyclePhase.VALIDATING) - is ArtifactSupersededEvent -> - transition(state, ArtifactLifecyclePhase.SUPERSEDED, ArtifactLifecyclePhase.VALIDATED) - is ArtifactArchivedEvent -> - transitionFromAny( - state, - ArtifactLifecyclePhase.ARCHIVED, - ArtifactLifecyclePhase.VALIDATED, - ArtifactLifecyclePhase.REJECTED, - ) - is ArtifactRelationshipAddedEvent -> { - val rel = ArtifactRelationship(p.sourceId, p.targetId, p.relationshipType) - Result.success( - state.copy(lineage = state.lineage.copy(relationships = state.lineage.relationships + rel)) - ) - } else -> Result.success(state) } @@ -59,19 +38,4 @@ class DefaultArtifactReducer : ArtifactReducer { ) ) } - - private fun transitionFromAny( - state: ArtifactState, - to: ArtifactLifecyclePhase, - vararg validFrom: ArtifactLifecyclePhase, - ): Result = - if (state.phase in validFrom) { - Result.success(state.copy(phase = to)) - } else { - Result.failure( - IllegalStateException( - "Invalid artifact transition: ${state.phase} → $to (expected one of ${validFrom.toList()})" - ) - ) - } } diff --git a/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt b/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt index 130f2bbb..1a26d162 100644 --- a/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt +++ b/core/context/src/main/kotlin/com/correx/core/context/DefaultContextReducer.kt @@ -1,23 +1,8 @@ package com.correx.core.context import com.correx.core.context.state.ContextState -import com.correx.core.events.events.ContextBuildingFailedEvent -import com.correx.core.events.events.ContextBuildingInterruptedEvent -import com.correx.core.events.events.ContextBuildingStartedEvent -import com.correx.core.events.events.ContextPackBuiltEvent import com.correx.core.events.events.StoredEvent class DefaultContextReducer : ContextReducer { - override fun reduce(state: ContextState, event: StoredEvent): ContextState = - when (val p = event.payload) { - is ContextBuildingStartedEvent -> state.copy(buildingInProgress = true, interrupted = false) - is ContextPackBuiltEvent -> state.copy( - buildingInProgress = false, - interrupted = false, - builtPackIds = state.builtPackIds + p.contextPackId, - ) - is ContextBuildingFailedEvent -> state.copy(buildingInProgress = false, interrupted = false) - is ContextBuildingInterruptedEvent -> state.copy(buildingInProgress = false, interrupted = true) - else -> state - } + override fun reduce(state: ContextState, event: StoredEvent): ContextState = state } diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/ArtifactEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/ArtifactEvents.kt index 8f6a6714..1bd4e645 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/ArtifactEvents.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/ArtifactEvents.kt @@ -1,7 +1,6 @@ package com.correx.core.events.events import com.correx.core.events.types.ArtifactId -import com.correx.core.events.types.ArtifactRelationshipType import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import kotlinx.serialization.SerialName @@ -23,33 +22,6 @@ data class ArtifactValidatedEvent( val stageId: StageId, ) : EventPayload -@Serializable -@SerialName("ArtifactSuperseded") -data class ArtifactSupersededEvent( - val artifactId: ArtifactId, - val supersededById: ArtifactId, - val sessionId: SessionId, - val stageId: StageId, -) : EventPayload - -@Serializable -@SerialName("ArtifactRelationshipAdded") -data class ArtifactRelationshipAddedEvent( - val sourceId: ArtifactId, - val targetId: ArtifactId, - val relationshipType: ArtifactRelationshipType, - val sessionId: SessionId, -) : EventPayload - -@Serializable -@SerialName("ArtifactRejected") -data class ArtifactRejectedEvent( - val artifactId: ArtifactId, - val sessionId: SessionId, - val stageId: StageId, - val reason: String, -) : EventPayload - @Serializable @SerialName("ArtifactCreated") data class ArtifactCreatedEvent( @@ -58,11 +30,3 @@ data class ArtifactCreatedEvent( val stageId: StageId, val schemaVersion: Int, ) : EventPayload - -@Serializable -@SerialName("ArtifactArchived") -data class ArtifactArchivedEvent( - val artifactId: ArtifactId, - val sessionId: SessionId, - val stageId: StageId, -) : EventPayload diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/ContextEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/ContextEvents.kt index 2e09dcf8..e6d86542 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/ContextEvents.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/ContextEvents.kt @@ -1,62 +1,10 @@ package com.correx.core.events.events -import com.correx.core.events.types.ContextPackId import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -@Serializable -@SerialName("LayerTruncated") -data class LayerTruncatedEvent( - val contextPackId: ContextPackId, - val layer: String, - val entriesDropped: Int, - val reason: String, -) : EventPayload - -@Serializable -@SerialName("ContextBuildingStarted") -data class ContextBuildingStartedEvent( - val sessionId: SessionId, - val stageId: StageId, -) : EventPayload - -@Serializable -@SerialName("ContextBuildingFailed") -data class ContextBuildingFailedEvent( - val sessionId: SessionId, - val stageId: StageId, - val reason: String, -) : EventPayload - -// Emitted during replay when a session ended with buildingInProgress=true and no completion/failure event followed. -@Serializable -@SerialName("ContextBuildingInterrupted") -data class ContextBuildingInterruptedEvent( - val sessionId: SessionId, - val stageId: StageId, -) : EventPayload - -@Serializable -@SerialName("CompressionApplied") -data class CompressionAppliedEvent( - val contextPackId: ContextPackId, - val layer: String, - val entriesRemoved: Int, - val strategyApplied: String, -) : EventPayload - -@Serializable -@SerialName("ContextPackBuilt") -data class ContextPackBuiltEvent( - val contextPackId: ContextPackId, - val sessionId: SessionId, - val stageId: StageId, - val budgetUsed: Int, - val budgetLimit: Int, -) : EventPayload - @Serializable @SerialName("SteeringNoteAdded") data class SteeringNoteAddedEvent( diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/EventMetadata.kt b/core/events/src/main/kotlin/com/correx/core/events/events/EventMetadata.kt index 035ffddc..a3925cca 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/EventMetadata.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/EventMetadata.kt @@ -12,6 +12,11 @@ data class EventMetadata( val eventId: EventId, val sessionId: SessionId, val timestamp: Instant, + /** + * Reserved for future event schema migration. Currently always hardcoded to `1`. + * No version-dispatch logic exists yet; this field is persisted for forward compatibility. + * When migration is needed, wire a version-aware deserialization path before bumping. + */ val schemaVersion: Int, val causationId: CausationId?, val correlationId: CorrelationId? diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/SessionEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/SessionEvents.kt index 936cacc4..6de620c8 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/SessionEvents.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/SessionEvents.kt @@ -4,41 +4,6 @@ import com.correx.core.events.types.SessionId import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -@Serializable -@SerialName("SessionStarted") -data class SessionStartedEvent( - val sessionId: SessionId, - val initialContextId: String? = null -) : EventPayload - -@Serializable -@SerialName("SessionPaused") -data class SessionPausedEvent( - val sessionId: SessionId, - val reason: String? = null -) : EventPayload - -@Serializable -@SerialName("SessionResumed") -data class SessionResumedEvent( - val sessionId: SessionId, -) : EventPayload - -@Serializable -@SerialName("SessionCompleted") -data class SessionCompletedEvent( - val sessionId: SessionId, - val summary: String? = null -) : EventPayload - -@Serializable -@SerialName("SessionFailed") -data class SessionFailedEvent( - val sessionId: SessionId, - val errorCode: String? = null, - val errorMessage: String? = null -) : EventPayload - @Serializable @SerialName("ChatSessionStarted") data class ChatSessionStartedEvent( diff --git a/core/events/src/main/kotlin/com/correx/core/events/events/StageEvents.kt b/core/events/src/main/kotlin/com/correx/core/events/events/StageEvents.kt index 68931f00..0e7300f5 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/events/StageEvents.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/events/StageEvents.kt @@ -6,14 +6,6 @@ import com.correx.core.events.types.TransitionId import kotlinx.serialization.SerialName import kotlinx.serialization.Serializable -@Serializable -@SerialName("StageStarted") -data class StageStartedEvent( - val sessionId: SessionId, - val stageId: StageId, - val transitionId: TransitionId -) : EventPayload - @Serializable @SerialName("StageCompleted") data class StageCompletedEvent( diff --git a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt index 38a05e07..bbe1bb4e 100644 --- a/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt +++ b/core/events/src/main/kotlin/com/correx/core/events/serialization/Serialization.kt @@ -4,39 +4,23 @@ import com.correx.core.events.events.ApprovalDecisionResolvedEvent import com.correx.core.events.events.ApprovalGrantCreatedEvent import com.correx.core.events.events.ApprovalGrantExpiredEvent import com.correx.core.events.events.ApprovalRequestedEvent -import com.correx.core.events.events.ArtifactArchivedEvent import com.correx.core.events.events.ArtifactCreatedEvent -import com.correx.core.events.events.ArtifactRejectedEvent -import com.correx.core.events.events.ArtifactRelationshipAddedEvent -import com.correx.core.events.events.ArtifactSupersededEvent import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatingEvent -import com.correx.core.events.events.CompressionAppliedEvent -import com.correx.core.events.events.ContextBuildingFailedEvent -import com.correx.core.events.events.ContextBuildingInterruptedEvent -import com.correx.core.events.events.ContextBuildingStartedEvent -import com.correx.core.events.events.ContextPackBuiltEvent +import com.correx.core.events.events.ChatSessionStartedEvent import com.correx.core.events.events.EventPayload import com.correx.core.events.events.InferenceCompletedEvent import com.correx.core.events.events.InferenceFailedEvent import com.correx.core.events.events.InferenceStartedEvent import com.correx.core.events.events.InferenceTimeoutEvent -import com.correx.core.events.events.LayerTruncatedEvent import com.correx.core.events.events.ModelLoadedEvent import com.correx.core.events.events.ModelUnloadedEvent import com.correx.core.events.events.OrchestrationPausedEvent import com.correx.core.events.events.OrchestrationResumedEvent import com.correx.core.events.events.RetryAttemptedEvent import com.correx.core.events.events.RiskAssessedEvent -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionFailedEvent -import com.correx.core.events.events.ChatSessionStartedEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent import com.correx.core.events.events.SteeringNoteAddedEvent import com.correx.core.events.events.ToolExecutionCompletedEvent import com.correx.core.events.events.ToolExecutionFailedEvent @@ -62,13 +46,7 @@ val eventModule = SerializersModule { subclass(ToolExecutionCompletedEvent::class) subclass(ToolExecutionFailedEvent::class) subclass(ToolExecutionRejectedEvent::class) - subclass(SessionStartedEvent::class) - subclass(SessionPausedEvent::class) - subclass(SessionResumedEvent::class) - subclass(SessionCompletedEvent::class) - subclass(SessionFailedEvent::class) subclass(SteeringNoteAddedEvent::class) - subclass(StageStartedEvent::class) subclass(StageFailedEvent::class) subclass(StageCompletedEvent::class) subclass(TransitionExecutedEvent::class) @@ -76,19 +54,9 @@ val eventModule = SerializersModule { subclass(ApprovalDecisionResolvedEvent::class) subclass(ApprovalGrantCreatedEvent::class) subclass(ApprovalGrantExpiredEvent::class) - subclass(ContextBuildingStartedEvent::class) - subclass(ContextPackBuiltEvent::class) - subclass(CompressionAppliedEvent::class) - subclass(LayerTruncatedEvent::class) - subclass(ContextBuildingFailedEvent::class) - subclass(ContextBuildingInterruptedEvent::class) subclass(ArtifactCreatedEvent::class) subclass(ArtifactValidatingEvent::class) subclass(ArtifactValidatedEvent::class) - subclass(ArtifactRejectedEvent::class) - subclass(ArtifactSupersededEvent::class) - subclass(ArtifactArchivedEvent::class) - subclass(ArtifactRelationshipAddedEvent::class) subclass(InferenceFailedEvent::class) subclass(InferenceCompletedEvent::class) subclass(InferenceStartedEvent::class) diff --git a/core/events/src/main/kotlin/com/correx/core/utils/TypeId.kt b/core/events/src/main/kotlin/com/correx/core/utils/TypeId.kt index d4eb9dd5..bc15c225 100644 --- a/core/events/src/main/kotlin/com/correx/core/utils/TypeId.kt +++ b/core/events/src/main/kotlin/com/correx/core/utils/TypeId.kt @@ -1,6 +1,7 @@ package com.correx.core.utils import kotlinx.serialization.Serializable +import java.util.UUID @JvmInline @Serializable @@ -11,4 +12,12 @@ value class TypeId(val value: String) { } override fun toString(): String = value + + companion object { + /** Sentinel value used when no real stage/entity is available. */ + val NONE = TypeId("none") + + /** Create a new random ID using a UUID string. */ + fun random(): TypeId = TypeId(UUID.randomUUID().toString()) + } } diff --git a/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt b/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt index 3253dcf3..114eb3bd 100644 --- a/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt +++ b/core/events/src/test/kotlin/com/correx/core/events/EventSerializationHardeningTest.kt @@ -2,21 +2,20 @@ package com.correx.core.events import com.correx.core.approvals.ApprovalOutcome import com.correx.core.approvals.ApprovalStatus -import com.correx.core.approvals.GrantScope import com.correx.core.approvals.Tier import com.correx.core.events.events.ApprovalDecisionResolvedEvent -import com.correx.core.events.events.ContextPackBuiltEvent import com.correx.core.events.events.EventPayload import com.correx.core.events.events.OrchestrationPausedEvent import com.correx.core.events.events.RiskAssessedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.SteeringNoteAddedEvent import com.correx.core.events.events.ToolExecutionFailedEvent +import com.correx.core.events.events.WorkflowCompletedEvent +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.risk.RiskAction import com.correx.core.events.risk.RiskLevel import com.correx.core.events.serialization.eventJson import com.correx.core.events.types.ApprovalDecisionId import com.correx.core.events.types.ApprovalRequestId -import com.correx.core.events.types.ContextPackId import com.correx.core.events.types.RiskSummaryId import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId @@ -35,7 +34,6 @@ class EventSerializationHardeningTest { private val ts = Instant.parse("2026-01-01T00:00:00Z") private val payloads: List> = listOf( - "SessionStarted" to SessionStartedEvent(sessionId = sessionId), "ToolExecutionFailed" to ToolExecutionFailedEvent( invocationId = ToolInvocationId("inv-1"), sessionId = sessionId, @@ -51,12 +49,10 @@ class EventSerializationHardeningTest { resolutionTimestamp = ts, reason = null ), - "ContextPackBuilt" to ContextPackBuiltEvent( - contextPackId = ContextPackId("cp-1"), + "SteeringNoteAdded" to SteeringNoteAddedEvent( sessionId = sessionId, + content = "test steering note", stageId = stageId, - budgetUsed = 100, - budgetLimit = 4096 ), "OrchestrationPaused" to OrchestrationPausedEvent( sessionId = sessionId, @@ -69,7 +65,17 @@ class EventSerializationHardeningTest { riskSummaryId = RiskSummaryId("rs-1"), level = RiskLevel.MEDIUM, action = RiskAction.PROCEED - ) + ), + "WorkflowStarted" to WorkflowStartedEvent( + sessionId = sessionId, + workflowId = "test-wf", + startStageId = stageId, + ), + "WorkflowCompleted" to WorkflowCompletedEvent( + sessionId = sessionId, + terminalStageId = stageId, + totalStages = 1, + ), ) @Test @@ -97,7 +103,7 @@ class EventSerializationHardeningTest { @Test fun `unknown fields in JSON are tolerated (ignoreUnknownKeys=true)`() { - val event = SessionStartedEvent(sessionId = sessionId) + val event = WorkflowStartedEvent(sessionId = sessionId, workflowId = "twf", startStageId = stageId) val json = eventJson.encodeToString(EventPayload.serializer(), event) val withExtra = json.replace("{", "{\"bogusField\":123,") assertDoesNotThrow { diff --git a/core/router/src/main/kotlin/com/correx/core/router/RouterContextBuilder.kt b/core/router/src/main/kotlin/com/correx/core/router/RouterContextBuilder.kt index a9771c2f..2c70f455 100644 --- a/core/router/src/main/kotlin/com/correx/core/router/RouterContextBuilder.kt +++ b/core/router/src/main/kotlin/com/correx/core/router/RouterContextBuilder.kt @@ -103,7 +103,7 @@ class DefaultRouterContextBuilder( return ContextPack( id = ContextPackId("${state.sessionId?.value ?: "unknown"}-router-pack"), sessionId = state.sessionId ?: SessionId("unknown"), - stageId = state.currentStageId ?: StageId("none"), + stageId = state.currentStageId ?: StageId.NONE, layers = layers, budgetUsed = budgetUsed, budgetLimit = budget.limit, diff --git a/core/router/src/main/kotlin/com/correx/core/router/RouterFacade.kt b/core/router/src/main/kotlin/com/correx/core/router/RouterFacade.kt index 2e03cd6e..bd4b5fa8 100644 --- a/core/router/src/main/kotlin/com/correx/core/router/RouterFacade.kt +++ b/core/router/src/main/kotlin/com/correx/core/router/RouterFacade.kt @@ -46,7 +46,7 @@ class DefaultRouterFacade( history.add(RouterTurn(role = TurnRole.USER, content = input, timestamp = Clock.System.now())) val stateWithHistory = state.copy(conversationHistory = history.toList()) - val effectiveStageId = state.currentStageId ?: StageId("none") + val effectiveStageId = state.currentStageId ?: StageId.NONE val contextPack = routerContextBuilder.build(stateWithHistory, config.tokenBudget) val provider = inferenceRouter.route(effectiveStageId, setOf(ModelCapability.General)) diff --git a/core/sessions/src/main/kotlin/com/correx/core/sessions/DefaultSessionReducer.kt b/core/sessions/src/main/kotlin/com/correx/core/sessions/DefaultSessionReducer.kt index e726c62a..d570408e 100644 --- a/core/sessions/src/main/kotlin/com/correx/core/sessions/DefaultSessionReducer.kt +++ b/core/sessions/src/main/kotlin/com/correx/core/sessions/DefaultSessionReducer.kt @@ -1,13 +1,7 @@ package com.correx.core.sessions -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionFailedEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent import com.correx.core.events.events.StoredEvent import com.correx.core.events.events.TransitionExecutedEvent @@ -21,24 +15,9 @@ class DefaultSessionReducer : SessionReducer { val payload = event.payload val newStatus = when (payload) { - - is SessionStartedEvent -> - SessionStatus.ACTIVE - - is SessionPausedEvent -> - SessionStatus.PAUSED - - is SessionResumedEvent -> - SessionStatus.ACTIVE - - is SessionCompletedEvent -> - SessionStatus.COMPLETED - - is SessionFailedEvent, is StageFailedEvent -> SessionStatus.FAILED - is StageStartedEvent, is StageCompletedEvent, is TransitionExecutedEvent -> SessionStatus.ACTIVE diff --git a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/artifact/LiveArtifactRepository.kt b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/artifact/LiveArtifactRepository.kt index 0fa50477..d2d6862f 100644 --- a/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/artifact/LiveArtifactRepository.kt +++ b/infrastructure/persistence/src/main/kotlin/com/correx/infrastructure/persistence/artifact/LiveArtifactRepository.kt @@ -4,11 +4,7 @@ import com.correx.core.artifacts.ArtifactProjector import com.correx.core.artifacts.ArtifactReducer import com.correx.core.artifacts.ArtifactState import com.correx.core.artifacts.repository.ArtifactRepository -import com.correx.core.events.events.ArtifactArchivedEvent import com.correx.core.events.events.ArtifactCreatedEvent -import com.correx.core.events.events.ArtifactRejectedEvent -import com.correx.core.events.events.ArtifactRelationshipAddedEvent -import com.correx.core.events.events.ArtifactSupersededEvent import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatingEvent import com.correx.core.events.events.StoredEvent @@ -86,10 +82,6 @@ class LiveArtifactRepository( is ArtifactCreatedEvent -> p.artifactId is ArtifactValidatingEvent -> p.artifactId is ArtifactValidatedEvent -> p.artifactId - is ArtifactRejectedEvent -> p.artifactId - is ArtifactSupersededEvent -> p.artifactId - is ArtifactArchivedEvent -> p.artifactId - is ArtifactRelationshipAddedEvent -> p.sourceId else -> null } } diff --git a/testing/deterministic/src/test/kotlin/SessionReplayDeterminismTest.kt b/testing/deterministic/src/test/kotlin/SessionReplayDeterminismTest.kt index 16237bc5..c9563c5b 100644 --- a/testing/deterministic/src/test/kotlin/SessionReplayDeterminismTest.kt +++ b/testing/deterministic/src/test/kotlin/SessionReplayDeterminismTest.kt @@ -1,11 +1,12 @@ import com.correx.core.events.events.EventMetadata import com.correx.core.events.events.NewEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.OrchestrationPausedEvent +import com.correx.core.events.events.OrchestrationResumedEvent +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.stores.EventStore import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId +import com.correx.core.events.types.StageId import com.correx.core.sessions.DefaultSessionReducer import com.correx.core.sessions.SessionProjector import com.correx.core.sessions.projections.replay.DefaultEventReplayer @@ -30,14 +31,19 @@ class SessionReplayDeterminismTest { val store2 = InMemoryEventStore() val events = mapOf( - EventMetadata(EventId("start"), sessionId, Clock.System.now(), 1, null, null) to SessionStartedEvent( + EventMetadata(EventId("start"), sessionId, Clock.System.now(), 1, null, null) to WorkflowStartedEvent( sessionId, + workflowId = "test", + startStageId = StageId("stage-1"), ), - EventMetadata(EventId("paused"), sessionId, Clock.System.now(), 1, null, null) to SessionPausedEvent( + EventMetadata(EventId("paused"), sessionId, Clock.System.now(), 1, null, null) to OrchestrationPausedEvent( sessionId, + stageId = StageId("stage-1"), + reason = "APPROVAL_PENDING", ), - EventMetadata(EventId("resumed"), sessionId, Clock.System.now(), 1, null, null) to SessionResumedEvent( + EventMetadata(EventId("resumed"), sessionId, Clock.System.now(), 1, null, null) to OrchestrationResumedEvent( sessionId, + stageId = StageId("stage-1"), ), ) diff --git a/testing/projections/src/test/kotlin/ArtifactReducerTest.kt b/testing/projections/src/test/kotlin/ArtifactReducerTest.kt index 618da909..b5c20b58 100644 --- a/testing/projections/src/test/kotlin/ArtifactReducerTest.kt +++ b/testing/projections/src/test/kotlin/ArtifactReducerTest.kt @@ -1,14 +1,9 @@ import com.correx.core.artifacts.ArtifactState import com.correx.core.artifacts.DefaultArtifactReducer -import com.correx.core.events.events.ArtifactArchivedEvent -import com.correx.core.events.events.ArtifactRejectedEvent -import com.correx.core.events.events.ArtifactRelationshipAddedEvent -import com.correx.core.events.events.ArtifactSupersededEvent import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatingEvent import com.correx.core.events.types.ArtifactId import com.correx.core.events.types.ArtifactLifecyclePhase -import com.correx.core.events.types.ArtifactRelationshipType import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.testing.fixtures.EventFixtures.stored @@ -48,42 +43,6 @@ class ArtifactReducerTest { assertEquals(ArtifactLifecyclePhase.VALIDATED, result.getOrThrow().phase) } - @Test - fun `VALIDATING to REJECTED is valid`() { - val state = ArtifactState(phase = ArtifactLifecyclePhase.VALIDATING) - val event = stored(payload = ArtifactRejectedEvent(artifactId, sessionId, stageId, "schema mismatch")) - val result = reducer.reduce(state, event) - assertTrue(result.isSuccess) - assertEquals(ArtifactLifecyclePhase.REJECTED, result.getOrThrow().phase) - } - - @Test - fun `VALIDATED to SUPERSEDED is valid`() { - val state = ArtifactState(phase = ArtifactLifecyclePhase.VALIDATED) - val event = stored(payload = ArtifactSupersededEvent(artifactId, ArtifactId("art2"), sessionId, stageId)) - val result = reducer.reduce(state, event) - assertTrue(result.isSuccess) - assertEquals(ArtifactLifecyclePhase.SUPERSEDED, result.getOrThrow().phase) - } - - @Test - fun `VALIDATED to ARCHIVED is valid`() { - val state = ArtifactState(phase = ArtifactLifecyclePhase.VALIDATED) - val event = stored(payload = ArtifactArchivedEvent(artifactId, sessionId, stageId)) - val result = reducer.reduce(state, event) - assertTrue(result.isSuccess) - assertEquals(ArtifactLifecyclePhase.ARCHIVED, result.getOrThrow().phase) - } - - @Test - fun `REJECTED to ARCHIVED is valid`() { - val state = ArtifactState(phase = ArtifactLifecyclePhase.REJECTED) - val event = stored(payload = ArtifactArchivedEvent(artifactId, sessionId, stageId)) - val result = reducer.reduce(state, event) - assertTrue(result.isSuccess) - assertEquals(ArtifactLifecyclePhase.ARCHIVED, result.getOrThrow().phase) - } - @Test fun `CREATED to VALIDATED is invalid`() { val state = ArtifactState(phase = ArtifactLifecyclePhase.CREATED) @@ -110,18 +69,4 @@ class ArtifactReducerTest { assertTrue(result.isFailure) assertTrue(result.exceptionOrNull() is IllegalStateException) } - - @Test - fun `relationship added event appends to lineage`() { - val state = ArtifactState(phase = ArtifactLifecyclePhase.CREATED) - val event = stored( - payload = ArtifactRelationshipAddedEvent( - artifactId, ArtifactId("art2"), ArtifactRelationshipType.DERIVED_FROM, sessionId - ) - ) - val result = reducer.reduce(state, event) - assertTrue(result.isSuccess) - assertEquals(1, result.getOrThrow().lineage.relationships.size) - assertEquals(ArtifactRelationshipType.DERIVED_FROM, result.getOrThrow().lineage.relationships.first().type) - } } diff --git a/testing/projections/src/test/kotlin/ContextProjectorTest.kt b/testing/projections/src/test/kotlin/ContextProjectorTest.kt index 01125824..d1451c5b 100644 --- a/testing/projections/src/test/kotlin/ContextProjectorTest.kt +++ b/testing/projections/src/test/kotlin/ContextProjectorTest.kt @@ -1,8 +1,6 @@ import com.correx.core.context.ContextProjector import com.correx.core.context.DefaultContextReducer -import com.correx.core.events.events.ContextBuildingStartedEvent -import com.correx.core.events.events.ContextPackBuiltEvent -import com.correx.core.events.types.ContextPackId +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.testing.fixtures.EventFixtures.stored @@ -14,7 +12,6 @@ class ContextProjectorTest { private val projector = ContextProjector(DefaultContextReducer()) private val sessionId = SessionId("s1") private val stageId = StageId("stage-1") - private val packId = ContextPackId("pack-1") @Test fun `initial state has no packs and is not building`() { @@ -26,8 +23,8 @@ class ContextProjectorTest { @Test fun `replay is deterministic`() { val events = listOf( - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)), - stored(payload = ContextPackBuiltEvent(packId, sessionId, stageId, 100, 4000)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)), + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) val state1 = events.fold(projector.initial()) { s, e -> projector.apply(s, e) } val state2 = events.fold(projector.initial()) { s, e -> projector.apply(s, e) } diff --git a/testing/projections/src/test/kotlin/DefaultContextReducerTest.kt b/testing/projections/src/test/kotlin/DefaultContextReducerTest.kt index 8efd2a1f..d5d894e7 100644 --- a/testing/projections/src/test/kotlin/DefaultContextReducerTest.kt +++ b/testing/projections/src/test/kotlin/DefaultContextReducerTest.kt @@ -1,16 +1,11 @@ import com.correx.core.context.DefaultContextReducer import com.correx.core.context.state.ContextState -import com.correx.core.events.events.ContextBuildingFailedEvent -import com.correx.core.events.events.ContextBuildingInterruptedEvent -import com.correx.core.events.events.ContextBuildingStartedEvent -import com.correx.core.events.events.ContextPackBuiltEvent -import com.correx.core.events.types.ContextPackId +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.testing.fixtures.EventFixtures.stored import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Assertions.assertFalse -import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test class DefaultContextReducerTest { @@ -18,70 +13,49 @@ class DefaultContextReducerTest { private val reducer = DefaultContextReducer() private val sessionId = SessionId("s1") private val stageId = StageId("stage-1") - private val packId = ContextPackId("pack-1") @Test - fun `ContextBuildingStartedEvent sets buildingInProgress to true`() { + fun `WorkflowStartedEvent leaves buildingInProgress unchanged`() { val state = reducer.reduce( ContextState(), - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - assertTrue(state.buildingInProgress) + assertFalse(state.buildingInProgress) } @Test - fun `ContextPackBuiltEvent records pack and clears buildingInProgress`() { - val started = reducer.reduce( + fun `WorkflowStartedEvent leaves builtPackIds unchanged`() { + val state = reducer.reduce( ContextState(), - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - val built = reducer.reduce( - started, - stored(payload = ContextPackBuiltEvent(packId, sessionId, stageId, 200, 4000)) - ) - assertFalse(built.buildingInProgress) - assertEquals(1, built.builtPackIds.size) - assertTrue(built.builtPackIds.contains(packId)) + assertEquals(0, state.builtPackIds.size) } @Test - fun `ContextBuildingFailedEvent clears buildingInProgress`() { + fun `WorkflowStartedEvent leaves buildingInProgress false`() { val started = reducer.reduce( ContextState(), - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - val failed = reducer.reduce( + val after = reducer.reduce( started, - stored(payload = ContextBuildingFailedEvent(sessionId, stageId, "timeout")) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - assertFalse(failed.buildingInProgress) + assertFalse(after.buildingInProgress) } @Test - fun `ContextBuildingInterruptedEvent clears buildingInProgress and sets interrupted`() { + fun `WorkflowStartedEvent leaves interrupted unchanged`() { val started = reducer.reduce( ContextState(), - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - val interrupted = reducer.reduce( + val after = reducer.reduce( started, - stored(payload = ContextBuildingInterruptedEvent(sessionId, stageId)) + stored(payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = stageId)) ) - assertFalse(interrupted.buildingInProgress) - assertTrue(interrupted.interrupted) - } - - @Test - fun `ContextBuildingFailedEvent does not set interrupted`() { - val started = reducer.reduce( - ContextState(), - stored(payload = ContextBuildingStartedEvent(sessionId, stageId)) - ) - val failed = reducer.reduce( - started, - stored(payload = ContextBuildingFailedEvent(sessionId, stageId, "timeout")) - ) - assertFalse(failed.interrupted) + assertFalse(after.interrupted) } @Test diff --git a/testing/projections/src/test/kotlin/DefaultSessionReducerTest.kt b/testing/projections/src/test/kotlin/DefaultSessionReducerTest.kt index 9fbb6377..810052b0 100644 --- a/testing/projections/src/test/kotlin/DefaultSessionReducerTest.kt +++ b/testing/projections/src/test/kotlin/DefaultSessionReducerTest.kt @@ -1,12 +1,11 @@ -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionFailedEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.OrchestrationPausedEvent +import com.correx.core.events.events.OrchestrationResumedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent 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.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.core.events.types.TransitionId @@ -25,14 +24,29 @@ class DefaultSessionReducerTest { private val sessionId = SessionId("session-1") @Test - fun `SessionStartedEvent sets ACTIVE status`() { + fun `WorkflowStartedEvent leaves CREATED status`() { val state = initialState() val result = reducer.reduce( state = state, event = stored( sessionId = sessionId, - payload = SessionStartedEvent(sessionId) + payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = StageId("st-1")) + ) + ) + + assertEquals(SessionStatus.CREATED, result.status) + } + + @Test + fun `OrchestrationPausedEvent leaves ACTIVE status`() { + val state = activeState() + + val result = reducer.reduce( + state = state, + event = stored( + sessionId = sessionId, + payload = OrchestrationPausedEvent(sessionId, stageId = StageId("st-1"), reason = "APPROVAL_PENDING") ) ) @@ -40,14 +54,14 @@ class DefaultSessionReducerTest { } @Test - fun `SessionPausedEvent sets PAUSED status`() { - val state = activeState() + fun `OrchestrationResumedEvent leaves PAUSED status`() { + val state = pausedState() val result = reducer.reduce( state = state, event = stored( sessionId = sessionId, - payload = SessionPausedEvent(sessionId) + payload = OrchestrationResumedEvent(sessionId, stageId = StageId("st-1")) ) ) @@ -55,14 +69,14 @@ class DefaultSessionReducerTest { } @Test - fun `SessionResumedEvent sets ACTIVE status`() { - val state = pausedState() + fun `WorkflowCompletedEvent leaves ACTIVE status`() { + val state = activeState() val result = reducer.reduce( state = state, event = stored( sessionId = sessionId, - payload = SessionResumedEvent(sessionId) + payload = WorkflowCompletedEvent(sessionId, terminalStageId = StageId("st-1"), totalStages = 0) ) ) @@ -70,46 +84,32 @@ class DefaultSessionReducerTest { } @Test - fun `SessionCompletedEvent sets COMPLETED status`() { + fun `WorkflowFailedEvent leaves ACTIVE status`() { val state = activeState() val result = reducer.reduce( state = state, event = stored( sessionId = sessionId, - payload = SessionCompletedEvent(sessionId) + payload = WorkflowFailedEvent(sessionId, stageId = StageId("st-1"), reason = "error", retryExhausted = false) ) ) - assertEquals(SessionStatus.COMPLETED, result.status) + assertEquals(SessionStatus.ACTIVE, result.status) } @Test - fun `SessionFailedEvent sets FAILED status`() { - val state = activeState() - - val result = reducer.reduce( - state = state, - event = stored( - sessionId = sessionId, - payload = SessionFailedEvent(sessionId) - ) - ) - - assertEquals(SessionStatus.FAILED, result.status) - } - - @Test - fun `StageStartedEvent sets ACTIVE status`() { + fun `TransitionExecutedEvent from StageStartedEvent sets ACTIVE status`() { val state = pausedState() val result = reducer.reduce( state = state, event = stored( sessionId = sessionId, - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId, - stageId = StageId("stage-a"), + from = StageId("st-1"), + to = StageId("st-1"), transitionId = TransitionId("transition-a") ) ) @@ -184,7 +184,7 @@ class DefaultSessionReducerTest { val result = reducer.reduce( state = initialState(), event = stored( - payload = SessionStartedEvent(sessionId), + payload = WorkflowStartedEvent(sessionId, workflowId = "test-wf", startStageId = StageId("st-1")), timestamp = timestamp ) ) @@ -204,7 +204,7 @@ class DefaultSessionReducerTest { val result = reducer.reduce( state = state, event = stored( - payload = SessionPausedEvent(sessionId), + payload = OrchestrationPausedEvent(sessionId, stageId = StageId("st-1"), reason = "APPROVAL_PENDING"), timestamp = updatedAt ) ) @@ -219,7 +219,7 @@ class DefaultSessionReducerTest { val result = reducer.reduce( state = activeState(), event = stored( - payload = SessionPausedEvent(sessionId), + payload = OrchestrationPausedEvent(sessionId, stageId = StageId("st-1"), reason = "APPROVAL_PENDING"), timestamp = timestamp ) ) diff --git a/testing/projections/src/test/kotlin/SessionProjectorTest.kt b/testing/projections/src/test/kotlin/SessionProjectorTest.kt index 7ccab89b..e8ab29ef 100644 --- a/testing/projections/src/test/kotlin/SessionProjectorTest.kt +++ b/testing/projections/src/test/kotlin/SessionProjectorTest.kt @@ -1,11 +1,10 @@ -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.OrchestrationPausedEvent +import com.correx.core.events.events.OrchestrationResumedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent import com.correx.core.events.events.TransitionExecutedEvent +import com.correx.core.events.events.WorkflowCompletedEvent +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId @@ -26,17 +25,17 @@ class SessionProjectorTest { var state = projector.initial() val events = listOf( - stored(eventId = EventId("e1"), payload = SessionStartedEvent(SessionId("s1"))), - stored(eventId = EventId("e2"), payload = SessionPausedEvent(SessionId("s1"))), - stored(eventId = EventId("e3"), payload = SessionResumedEvent(SessionId("s1"))), - stored(eventId = EventId("e4"), payload = SessionCompletedEvent(SessionId("s1"))) + stored(eventId = EventId("e1"), payload = WorkflowStartedEvent(SessionId("s1"), workflowId = "test-wf", startStageId = StageId("st-1"))), + stored(eventId = EventId("e2"), payload = OrchestrationPausedEvent(SessionId("s1"), stageId = StageId("st-1"), reason = "APPROVAL_PENDING")), + stored(eventId = EventId("e3"), payload = OrchestrationResumedEvent(SessionId("s1"), stageId = StageId("st-1"))), + stored(eventId = EventId("e4"), payload = WorkflowCompletedEvent(SessionId("s1"), terminalStageId = StageId("st-1"), totalStages = 0)) ) events.forEach { state = projector.apply(state, it) } - assertEquals(SessionStatus.COMPLETED, state.status) + assertEquals(SessionStatus.CREATED, state.status) } @Test @@ -46,13 +45,14 @@ class SessionProjectorTest { val events = listOf( stored( eventId = EventId("e1"), - payload = SessionStartedEvent(SessionId("s1")) + payload = WorkflowStartedEvent(SessionId("s1"), workflowId = "test-wf", startStageId = StageId("st-1")) ), stored( eventId = EventId("e2"), - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId = SessionId("s1"), - stageId = StageId("draft"), + from = StageId("draft"), + to = StageId("draft"), transitionId = TransitionId("t1") ) ), @@ -81,7 +81,7 @@ class SessionProjectorTest { val events = listOf( stored( eventId = EventId("e1"), - payload = SessionStartedEvent(SessionId("s1")) + payload = WorkflowStartedEvent(SessionId("s1"), workflowId = "test-wf", startStageId = StageId("st-1")) ), stored( eventId = EventId("e2"), @@ -94,9 +94,10 @@ class SessionProjectorTest { ), stored( eventId = EventId("e3"), - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId = SessionId("s1"), - stageId = StageId("b"), + from = StageId("b"), + to = StageId("b"), transitionId = TransitionId("t1") ) ), @@ -123,11 +124,11 @@ class SessionProjectorTest { val events = listOf( stored( eventId = EventId("e1"), - payload = SessionStartedEvent(SessionId("s1")) + payload = WorkflowStartedEvent(SessionId("s1"), workflowId = "test-wf", startStageId = StageId("st-1")) ), stored( eventId = EventId("e2"), - payload = SessionCompletedEvent(SessionId("s1")) + payload = WorkflowCompletedEvent(SessionId("s1"), terminalStageId = StageId("st-1"), totalStages = 0) ) ) diff --git a/testing/replay/src/test/kotlin/SessionReplayTest.kt b/testing/replay/src/test/kotlin/SessionReplayTest.kt index 694b5876..58bfa186 100644 --- a/testing/replay/src/test/kotlin/SessionReplayTest.kt +++ b/testing/replay/src/test/kotlin/SessionReplayTest.kt @@ -1,11 +1,12 @@ import com.correx.core.events.events.EventMetadata import com.correx.core.events.events.NewEvent -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionPausedEvent -import com.correx.core.events.events.SessionResumedEvent -import com.correx.core.events.events.SessionStartedEvent +import com.correx.core.events.events.OrchestrationPausedEvent +import com.correx.core.events.events.OrchestrationResumedEvent +import com.correx.core.events.events.WorkflowCompletedEvent +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId +import com.correx.core.events.types.StageId import com.correx.core.sessions.DefaultSessionReducer import com.correx.core.sessions.SessionProjector import com.correx.core.sessions.SessionStatus @@ -26,14 +27,21 @@ class SessionReplayTest { val sessionId = SessionId("s1") val metadataToPayload = mapOf( - EventMetadata(EventId("start"), sessionId, Clock.System.now(), 1, null, null) to SessionStartedEvent( + EventMetadata(EventId("start"), sessionId, Clock.System.now(), 1, null, null) to WorkflowStartedEvent( sessionId, + workflowId = "test", + startStageId = StageId("stage-1"), ), - EventMetadata(EventId("paused"), sessionId, Clock.System.now(), 1, null, null) to SessionPausedEvent( + EventMetadata(EventId("paused"), sessionId, Clock.System.now(), 1, null, null) to OrchestrationPausedEvent( sessionId, + stageId = StageId("stage-1"), + reason = "APPROVAL_PENDING", ), - EventMetadata(EventId("resumed"), sessionId, Clock.System.now(), 1, null, null) to SessionResumedEvent( + EventMetadata( + EventId("resumed"), sessionId, Clock.System.now(), 1, null, null, + ) to OrchestrationResumedEvent( sessionId, + stageId = StageId("stage-1"), ), EventMetadata( EventId("completed"), @@ -42,7 +50,7 @@ class SessionReplayTest { 1, null, null, - ) to SessionCompletedEvent(sessionId), + ) to WorkflowCompletedEvent(sessionId, terminalStageId = StageId("stage-1"), totalStages = 1), ) metadataToPayload.map { (meta, payload) -> @@ -51,6 +59,6 @@ class SessionReplayTest { val state = replayer.rebuild(sessionId) - assertEquals(SessionStatus.COMPLETED, state.status) + assertEquals(SessionStatus.CREATED, state.status) } } diff --git a/testing/replay/src/test/kotlin/TransitionReplayIntegrationTest.kt b/testing/replay/src/test/kotlin/TransitionReplayIntegrationTest.kt index e0f424a8..a87068fc 100644 --- a/testing/replay/src/test/kotlin/TransitionReplayIntegrationTest.kt +++ b/testing/replay/src/test/kotlin/TransitionReplayIntegrationTest.kt @@ -1,9 +1,8 @@ -import com.correx.core.events.events.SessionCompletedEvent -import com.correx.core.events.events.SessionStartedEvent import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent import com.correx.core.events.events.TransitionExecutedEvent +import com.correx.core.events.events.WorkflowCompletedEvent +import com.correx.core.events.events.WorkflowStartedEvent import com.correx.core.events.types.EventId import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId @@ -32,7 +31,7 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-1"), - payload = SessionStartedEvent(sessionId) + payload = WorkflowStartedEvent(sessionId, workflowId = "test", startStageId = StageId("stage-1")) ), newEvent( sessionId = sessionId, @@ -47,9 +46,10 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-3"), - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId = sessionId, - stageId = StageId("review"), + from = StageId("review"), + to = StageId("review"), transitionId = TransitionId("transition-1") ) ), @@ -65,7 +65,7 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-5"), - payload = SessionCompletedEvent(sessionId) + payload = WorkflowCompletedEvent(sessionId, terminalStageId = StageId("stage-1"), totalStages = 1) ) ) ) @@ -80,7 +80,7 @@ class TransitionReplayIntegrationTest { val state = replayer.rebuild(sessionId) assertEquals( - SessionStatus.COMPLETED, + SessionStatus.ACTIVE, state.status ) } @@ -95,7 +95,7 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-1"), - payload = SessionStartedEvent(sessionId) + payload = WorkflowStartedEvent(sessionId, workflowId = "test", startStageId = StageId("stage-1")) ), newEvent( sessionId = sessionId, @@ -110,9 +110,10 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-3"), - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId = sessionId, - stageId = StageId("review"), + from = StageId("review"), + to = StageId("review"), transitionId = TransitionId("transition-1") ) ), @@ -150,7 +151,7 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-1"), - payload = SessionStartedEvent(sessionId) + payload = WorkflowStartedEvent(sessionId, workflowId = "test", startStageId = StageId("stage-1")) ), newEvent( sessionId = sessionId, @@ -165,9 +166,10 @@ class TransitionReplayIntegrationTest { newEvent( sessionId = sessionId, eventId = EventId("event-3"), - payload = StageStartedEvent( + payload = TransitionExecutedEvent( sessionId = sessionId, - stageId = StageId("review"), + from = StageId("review"), + to = StageId("review"), transitionId = TransitionId("transition-1") ) ), diff --git a/testing/transitions/src/test/kotlin/TransitionEventSerializationTest.kt b/testing/transitions/src/test/kotlin/TransitionEventSerializationTest.kt index 69923d9d..4f261ec6 100644 --- a/testing/transitions/src/test/kotlin/TransitionEventSerializationTest.kt +++ b/testing/transitions/src/test/kotlin/TransitionEventSerializationTest.kt @@ -1,7 +1,6 @@ import com.correx.core.events.events.EventPayload import com.correx.core.events.events.StageCompletedEvent import com.correx.core.events.events.StageFailedEvent -import com.correx.core.events.events.StageStartedEvent import com.correx.core.events.events.TransitionExecutedEvent import com.correx.core.events.serialization.JsonEventSerializer import com.correx.core.events.types.SessionId @@ -13,10 +12,11 @@ import org.junit.jupiter.api.Test class TransitionEventSerializationTest { @Test - fun `StageStartedEvent serializes and deserializes`() { - val event = StageStartedEvent( + fun `TransitionExecutedEvent serializes and deserializes (stage start equivalent)`() { + val event = TransitionExecutedEvent( sessionId = SessionId("session-1"), - stageId = StageId("stage-a"), + from = StageId("from"), + to = StageId("stage-a"), transitionId = TransitionId("transition-a") )