1a9e1019ec
DefaultSessionReducer skipped WorkflowCompletedEvent to avoid terminalizing on freestyle planning's mid-session completion, but that left execution/normal completions stuck ACTIVE forever — headless watchers and approval loops never saw a terminal state. Add workflowId to WorkflowCompletedEvent (defaulted for replay compat), pass graph.id at all completeWorkflow sites, and terminalize to COMPLETED unless it's the freestyle_planning handoff. Correct stale tests that pinned the ACTIVE-forever behavior.
65 lines
2.5 KiB
Kotlin
65 lines
2.5 KiB
Kotlin
import com.correx.core.events.events.EventMetadata
|
|
import com.correx.core.events.events.NewEvent
|
|
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
|
|
import com.correx.core.sessions.projections.replay.DefaultEventReplayer
|
|
import com.correx.infrastructure.persistence.InMemoryEventStore
|
|
import kotlinx.datetime.Clock
|
|
import org.junit.jupiter.api.Assertions.assertEquals
|
|
import org.junit.jupiter.api.Test
|
|
|
|
class SessionReplayTest {
|
|
|
|
private val store = InMemoryEventStore()
|
|
private val projector = SessionProjector(DefaultSessionReducer())
|
|
private val replayer = DefaultEventReplayer(store, projector)
|
|
|
|
@Test
|
|
fun `rebuild session from event stream`(): Unit = kotlinx.coroutines.runBlocking {
|
|
val sessionId = SessionId("s1")
|
|
|
|
val metadataToPayload = mapOf(
|
|
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 OrchestrationPausedEvent(
|
|
sessionId,
|
|
stageId = StageId("stage-1"),
|
|
reason = "APPROVAL_PENDING",
|
|
),
|
|
EventMetadata(
|
|
EventId("resumed"), sessionId, Clock.System.now(), 1, null, null,
|
|
) to OrchestrationResumedEvent(
|
|
sessionId,
|
|
stageId = StageId("stage-1"),
|
|
),
|
|
EventMetadata(
|
|
EventId("completed"),
|
|
sessionId,
|
|
Clock.System.now(),
|
|
1,
|
|
null,
|
|
null,
|
|
) to WorkflowCompletedEvent(sessionId, terminalStageId = StageId("stage-1"), totalStages = 1),
|
|
)
|
|
|
|
metadataToPayload.map { (meta, payload) ->
|
|
store.append(NewEvent(meta, payload))
|
|
}
|
|
|
|
val state = replayer.rebuild(sessionId)
|
|
|
|
assertEquals(SessionStatus.COMPLETED, state.status)
|
|
}
|
|
}
|