import com.correx.core.approvals.ApprovalProjector import com.correx.core.approvals.DefaultApprovalReducer import com.correx.core.approvals.DefaultApprovalRepository import com.correx.core.approvals.domain.DefaultApprovalEngine import com.correx.core.artifacts.DefaultArtifactReducer import com.correx.core.events.events.RefinementIterationEvent import com.correx.core.events.events.WorkflowCompletedEvent import com.correx.core.events.events.WorkflowFailedEvent import com.correx.core.events.execution.RetryPolicy import com.correx.core.events.types.SessionId import com.correx.core.events.types.StageId import com.correx.core.events.types.TransitionId import com.correx.core.inference.InferenceRepository import com.correx.core.inference.InferenceState import com.correx.core.journal.DecisionJournalProjector import com.correx.core.journal.DefaultDecisionJournalReducer import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer import com.correx.core.kernel.orchestration.DefaultSessionOrchestrator import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationProjector import com.correx.core.kernel.orchestration.OrchestrationRepository import com.correx.core.kernel.orchestration.OrchestratorEngines import com.correx.core.kernel.orchestration.OrchestratorRepositories import com.correx.core.kernel.retry.DefaultRetryCoordinator import com.correx.core.risk.DefaultRiskAssessor import com.correx.core.sessions.DefaultSessionRepository import com.correx.core.sessions.projections.replay.DefaultEventReplayer import com.correx.core.sessions.projections.replay.EventReplayer import com.correx.core.transitions.graph.StageConfig import com.correx.core.transitions.graph.TransitionEdge import com.correx.core.transitions.graph.WorkflowGraph import com.correx.core.validation.pipeline.ValidationPipeline import com.correx.infrastructure.persistence.InMemoryEventStore import com.correx.infrastructure.persistence.artifact.LiveArtifactRepository import com.correx.testing.contracts.fixtures.artifactstore.NoopArtifactStore import com.correx.testing.fixtures.InferenceFixtures import com.correx.testing.fixtures.context.ContextFixtures import com.correx.testing.fixtures.cyclePolicyMissingValidator import com.correx.testing.fixtures.transitions.TransitionFixtures import com.correx.testing.kernel.MockSessionEventReplayer import kotlinx.coroutines.runBlocking import org.junit.jupiter.api.Assertions.assertEquals import org.junit.jupiter.api.Assertions.assertNotNull import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test class RefinementLoopTest { private val eventStore = InMemoryEventStore() private val sessionRepository = DefaultSessionRepository(MockSessionEventReplayer()) private val orchestrationRepository = OrchestrationRepository( DefaultEventReplayer(eventStore, OrchestrationProjector(DefaultOrchestrationReducer())), ) private val inferenceRepository = InferenceRepository( object : EventReplayer { override fun rebuild(sessionId: SessionId) = InferenceState() }, ) private val approvalRepository = DefaultApprovalRepository( DefaultEventReplayer(eventStore, ApprovalProjector(DefaultApprovalReducer())), ) private val decisionJournalRepository = DefaultDecisionJournalRepository( DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())), ) private val repositories = OrchestratorRepositories( eventStore = eventStore, inferenceRepository = inferenceRepository, orchestrationRepository = orchestrationRepository, sessionRepository = sessionRepository, artifactRepository = LiveArtifactRepository(eventStore, DefaultArtifactReducer()), approvalRepository = approvalRepository, ) private val engines = OrchestratorEngines( transitionResolver = TransitionFixtures.simpleResolver(), contextPackBuilder = ContextFixtures.simpleBuilder(), inferenceRouter = InferenceFixtures.fixedRouter(), validationPipeline = ValidationPipeline(validators = listOf(cyclePolicyMissingValidator())), approvalEngine = DefaultApprovalEngine(), riskAssessor = DefaultRiskAssessor(), ) private val orchestrator = DefaultSessionOrchestrator( repositories = repositories, engines = engines, retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = NoopArtifactStore(), decisionJournalRepository = decisionJournalRepository, ) // A→B and B→A, both unconditional: an unbounded loop absent the runtime guard. // B declares maxRetries=2 → the "A->B" cycle is capped at 2 iterations. private fun loopGraph() = WorkflowGraph( id = "refine-loop", stages = mapOf( StageId("A") to StageConfig(), StageId("B") to StageConfig(maxRetries = 2), ), transitions = setOf( TransitionEdge(TransitionId("t-ab"), StageId("A"), StageId("B"), condition = { true }), TransitionEdge(TransitionId("t-ba"), StageId("B"), StageId("A"), condition = { true }), ), start = StageId("A"), ) @Test fun `guard terminates an unbounded loop at maxIterations`(): Unit = runBlocking { val sessionId = SessionId("loop-1") orchestrator.run( sessionId, loopGraph(), OrchestrationConfig(retryPolicy = RetryPolicy(maxAttempts = 1, backoffMs = 0)), ) val events = eventStore.read(sessionId) val iterations = events.mapNotNull { it.payload as? RefinementIterationEvent } assertTrue(iterations.isNotEmpty(), "expected RefinementIterationEvent on back-edges") val failed = events.mapNotNull { it.payload as? WorkflowFailedEvent } assertTrue(failed.isNotEmpty(), "guard must escalate to a terminal failure") assertTrue( failed.first().reason.contains("refinement loop"), "unexpected failure reason: ${failed.first().reason}", ) // The "A->B" cycle is capped at 2 → it fires the terminal failure on iteration 3. val ab = iterations.filter { it.cycleKey == "A->B" } assertEquals(3, ab.maxOf { it.iteration }) } @Test fun `forward-only workflow emits no refinement iteration`(): Unit = runBlocking { val sessionId = SessionId("forward-1") val graph = WorkflowGraph( id = "forward", stages = mapOf(StageId("A") to StageConfig()), transitions = setOf( TransitionEdge(TransitionId("t1"), StageId("A"), StageId("done"), condition = { true }), ), start = StageId("A"), ) orchestrator.run( sessionId, graph, OrchestrationConfig(retryPolicy = RetryPolicy(maxAttempts = 1, backoffMs = 0)), ) val events = eventStore.read(sessionId) assertNotNull(events.find { it.payload is WorkflowCompletedEvent }) assertTrue( events.none { it.payload is RefinementIterationEvent }, "forward edges must not trigger the refinement guard", ) } }