From bcc20509e0fc73d0e4b35216ea6e190b8728fa2e Mon Sep 17 00:00:00 2001 From: kami Date: Thu, 4 Jun 2026 00:48:07 +0400 Subject: [PATCH] feat(kernel): inject decision journal into stage context + wire repository --- apps/server/build.gradle | 1 + .../kotlin/com/correx/apps/server/Main.kt | 2 + .../server/lifecycle/LifecycleTestSupport.kt | 2 + .../apps/server/ws/WorkspaceHandshakeTest.kt | 1 + core/kernel/build.gradle | 1 + .../DefaultSessionOrchestrator.kt | 4 +- .../orchestration/ReplayOrchestrator.kt | 6 +++ .../orchestration/SessionOrchestrator.kt | 20 +++++++- infrastructure/build.gradle | 1 + .../infrastructure/InfrastructureModule.kt | 11 ++++ testing/integration/build.gradle | 1 + .../SessionOrchestratorIntegrationTest.kt | 16 ++++++ .../src/test/kotlin/ToolCallGateTest.kt | 7 +++ .../kotlin/WorkspaceScopedToolRegistryTest.kt | 7 +++ testing/replay/build.gradle | 1 + .../test/kotlin/DecisionJournalReplayTest.kt | 50 +++++++++++++++++++ 16 files changed, 129 insertions(+), 2 deletions(-) create mode 100644 testing/replay/src/test/kotlin/DecisionJournalReplayTest.kt diff --git a/apps/server/build.gradle b/apps/server/build.gradle index cce43481..9071c0c7 100644 --- a/apps/server/build.gradle +++ b/apps/server/build.gradle @@ -21,6 +21,7 @@ dependencies { implementation project(':core:approvals') implementation project(':core:sessions') implementation project(':core:kernel') + implementation project(':core:journal') implementation project(':core:inference') implementation project(':core:transitions') implementation project(':core:context') diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt index 1140358e..5b0127d7 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt @@ -230,12 +230,14 @@ fun main() { workspacePolicy = workspacePolicy, workspaceToolRegistryProvider = wsToolRegistryProvider, ) + val decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore) val orchestrator = DefaultSessionOrchestrator( repositories = repositories, engines = engines, retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = artifactStore, tokenizer = firstProvider.tokenizer, + decisionJournalRepository = decisionJournalRepository, ) val defaultOrchestrationConfig = OrchestrationConfig( sandboxRoot = sandboxRoot, diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/lifecycle/LifecycleTestSupport.kt b/apps/server/src/test/kotlin/com/correx/apps/server/lifecycle/LifecycleTestSupport.kt index f567d527..9418f022 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/lifecycle/LifecycleTestSupport.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/lifecycle/LifecycleTestSupport.kt @@ -106,12 +106,14 @@ fun buildTestServerModule( val repositories = buildRepositories(eventStore) + val decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore) val orchestrator = DefaultSessionOrchestrator( repositories = repositories, engines = engines, retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = artifactStore, tokenizer = firstProvider.tokenizer, + decisionJournalRepository = decisionJournalRepository, ) val routerFacade = InfrastructureModule.createRouterFacade( diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/ws/WorkspaceHandshakeTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/ws/WorkspaceHandshakeTest.kt index 00468c29..398341a5 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/ws/WorkspaceHandshakeTest.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/ws/WorkspaceHandshakeTest.kt @@ -183,6 +183,7 @@ class WorkspaceHandshakeTest { retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = noopArtifactStore, tokenizer = provider.tokenizer, + decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore), ) val routerFacade = InfrastructureModule.createRouterFacade( diff --git a/core/kernel/build.gradle b/core/kernel/build.gradle index 6ecf2db5..d22f293a 100644 --- a/core/kernel/build.gradle +++ b/core/kernel/build.gradle @@ -17,6 +17,7 @@ dependencies { implementation project(':core:artifacts-store') implementation project(':core:risk') implementation project(':core:toolintent') + implementation(project(":core:journal")) implementation "org.slf4j:slf4j-api:2.0.16" } tasks.named("koverVerify").configure { enabled = false } diff --git a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/DefaultSessionOrchestrator.kt b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/DefaultSessionOrchestrator.kt index 83c7f9cb..36d59b86 100644 --- a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/DefaultSessionOrchestrator.kt +++ b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/DefaultSessionOrchestrator.kt @@ -4,6 +4,7 @@ import com.correx.core.approvals.ApprovalOutcome import com.correx.core.approvals.ApprovalStatus import com.correx.core.approvals.model.ApprovalDecision import com.correx.core.inference.Tokenizer +import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.artifacts.ArtifactState import com.correx.core.artifactstore.ArtifactStore import com.correx.core.events.events.ApprovalDecisionResolvedEvent @@ -40,7 +41,8 @@ class DefaultSessionOrchestrator( private val retryCoordinator: RetryCoordinator, artifactStore: ArtifactStore, tokenizer: Tokenizer? = null, -) : SessionOrchestrator(repositories, engines, artifactStore), ApprovalGateway { + decisionJournalRepository: DefaultDecisionJournalRepository, +) : SessionOrchestrator(repositories, engines, artifactStore, decisionJournalRepository), ApprovalGateway { override val tokenizer: Tokenizer? = tokenizer override val cancellations: ConcurrentHashMap = ConcurrentHashMap() diff --git a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/ReplayOrchestrator.kt b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/ReplayOrchestrator.kt index 279751cb..3e1eb012 100644 --- a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/ReplayOrchestrator.kt +++ b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/ReplayOrchestrator.kt @@ -2,6 +2,9 @@ package com.correx.core.kernel.orchestration import com.correx.core.approvals.domain.NoOpApprovalEngine import com.correx.core.artifactstore.ArtifactStore +import com.correx.core.journal.DefaultDecisionJournalRepository +import com.correx.core.journal.model.DecisionJournalState +import com.correx.core.sessions.projections.replay.EventReplayer import com.correx.core.context.model.ContextPack import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatingEvent @@ -58,6 +61,9 @@ class ReplayOrchestrator( repositories, engines.copy(approvalEngine = NoOpApprovalEngine(), riskAssessor = NoOpRiskAssessor()), artifactStore, + DefaultDecisionJournalRepository(object : EventReplayer { + override fun rebuild(sessionId: SessionId) = DecisionJournalState() + }), ) { private val replayProvider = ReplayInferenceProvider(repositories.eventStore) diff --git a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/SessionOrchestrator.kt b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/SessionOrchestrator.kt index 1a495270..181dfd64 100644 --- a/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/SessionOrchestrator.kt +++ b/core/kernel/src/main/kotlin/com/correx/core/kernel/orchestration/SessionOrchestrator.kt @@ -107,6 +107,8 @@ import kotlinx.datetime.Clock import kotlinx.serialization.encodeToString import kotlinx.serialization.json.Json import kotlinx.serialization.json.JsonObject +import com.correx.core.journal.DecisionJournalRenderer +import com.correx.core.journal.DefaultDecisionJournalRepository import org.slf4j.LoggerFactory import java.util.* import java.util.concurrent.* @@ -129,6 +131,8 @@ abstract class SessionOrchestrator( repositories: OrchestratorRepositories, engines: OrchestratorEngines, private val artifactStore: ArtifactStore, + private val decisionJournalRepository: DefaultDecisionJournalRepository, + private val decisionJournalRenderer: DecisionJournalRenderer = DecisionJournalRenderer(), ) { private val log = LoggerFactory.getLogger(this::class.java) private val eventStore: EventStore = repositories.eventStore @@ -285,7 +289,21 @@ abstract class SessionOrchestrator( // from currentContext.layers (a Map grouped by layer) would scramble turn order // across rounds; instead we grow our own ordered list and let the builder restamp // ordinals from it each round. - var accumulatedEntries = systemPrompt + schemaEntries + promptEntries + steeringEntries + val journalText = decisionJournalRenderer.render( + decisionJournalRepository.getJournal(sessionId), + ) + val journalEntries = if (journalText.isBlank()) emptyList() else listOf( + ContextEntry( + id = ContextEntryId(UUID.randomUUID().toString()), + layer = ContextLayer.L0, + content = journalText, + sourceType = "decisionJournal", + sourceId = "decision-journal", + tokenEstimate = journalText.length / 4, + role = EntryRole.SYSTEM, + ), + ) + var accumulatedEntries = systemPrompt + journalEntries + schemaEntries + promptEntries + steeringEntries val contextPack = contextPackBuilder.build( id = ContextPackId(UUID.randomUUID().toString()), sessionId = sessionId, diff --git a/infrastructure/build.gradle b/infrastructure/build.gradle index 28e68ac6..77e47c5a 100644 --- a/infrastructure/build.gradle +++ b/infrastructure/build.gradle @@ -22,6 +22,7 @@ dependencies { implementation project(":infrastructure:workflow") implementation project(":core:artifacts") implementation project(":core:artifacts-store") + implementation(project(":core:journal")) implementation project(":infrastructure:artifacts-cas") implementation "io.ktor:ktor-client-core:$ktor_version" implementation "io.ktor:ktor-client-cio:$ktor_version" diff --git a/infrastructure/src/main/kotlin/com/correx/infrastructure/InfrastructureModule.kt b/infrastructure/src/main/kotlin/com/correx/infrastructure/InfrastructureModule.kt index 125cdfbb..30e22581 100644 --- a/infrastructure/src/main/kotlin/com/correx/infrastructure/InfrastructureModule.kt +++ b/infrastructure/src/main/kotlin/com/correx/infrastructure/InfrastructureModule.kt @@ -3,6 +3,9 @@ package com.correx.infrastructure import com.correx.core.approvals.ApprovalProjector import com.correx.core.approvals.DefaultApprovalReducer import com.correx.core.approvals.DefaultApprovalRepository +import com.correx.core.journal.DecisionJournalProjector +import com.correx.core.journal.DefaultDecisionJournalReducer +import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.artifacts.DefaultArtifactReducer import com.correx.core.artifacts.repository.ArtifactRepository import com.correx.core.artifactstore.ArtifactStore @@ -177,6 +180,14 @@ object InfrastructureModule { ), ) + fun createDecisionJournalRepository(eventStore: EventStore): DefaultDecisionJournalRepository = + DefaultDecisionJournalRepository( + DefaultEventReplayer( + eventStore, + DecisionJournalProjector(DefaultDecisionJournalReducer()), + ), + ) + fun createWorkflowLoader(extraKinds: List = emptyList()): WorkflowLoader { val registry = DefaultArtifactKindRegistry() extraKinds.forEach { registry.register(it) } diff --git a/testing/integration/build.gradle b/testing/integration/build.gradle index f5e9b574..acb1787c 100644 --- a/testing/integration/build.gradle +++ b/testing/integration/build.gradle @@ -13,6 +13,7 @@ dependencies { testImplementation(project(":core:context")) testImplementation(project(":core:inference")) testImplementation(project(":core:kernel")) + testImplementation(project(":core:journal")) testImplementation(project(":core:risk")) testImplementation(project(":core:tools")) testImplementation(project(":core:toolintent")) diff --git a/testing/integration/src/test/kotlin/SessionOrchestratorIntegrationTest.kt b/testing/integration/src/test/kotlin/SessionOrchestratorIntegrationTest.kt index 7dcc2289..82fddf84 100644 --- a/testing/integration/src/test/kotlin/SessionOrchestratorIntegrationTest.kt +++ b/testing/integration/src/test/kotlin/SessionOrchestratorIntegrationTest.kt @@ -52,6 +52,9 @@ import com.correx.core.inference.Tokenizer import com.correx.core.inference.ToolCallFunction import com.correx.core.inference.ToolCallRequest import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer +import com.correx.core.journal.DecisionJournalProjector +import com.correx.core.journal.DefaultDecisionJournalReducer +import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.kernel.orchestration.DefaultSessionOrchestrator import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationProjector @@ -124,6 +127,10 @@ class SessionOrchestratorIntegrationTest { DefaultEventReplayer(eventStore, ApprovalProjector(DefaultApprovalReducer())), ) + private val decisionJournalRepository = DefaultDecisionJournalRepository( + DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())), + ) + private val repositories = OrchestratorRepositories( eventStore = eventStore, inferenceRepository = inferenceRepository, @@ -147,6 +154,7 @@ class SessionOrchestratorIntegrationTest { engines = engines, retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) @Test @@ -194,6 +202,7 @@ class SessionOrchestratorIntegrationTest { ), retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) failingOrchestrator.run(sessionId, graph, config) @@ -221,6 +230,7 @@ class SessionOrchestratorIntegrationTest { engines = engines.copy(validationPipeline = approvingPipeline), retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) val runJob = launch { orchestrator.run(sessionId, graph, config) } @@ -269,6 +279,7 @@ class SessionOrchestratorIntegrationTest { engines = engines.copy(validationPipeline = approvingPipeline), retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) val runJob = launch { approvalOrchestrator.run(sessionId, graph, config) } @@ -322,6 +333,7 @@ class SessionOrchestratorIntegrationTest { engines = engines, retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) val sessionId = SessionId("s6") @@ -364,6 +376,7 @@ class SessionOrchestratorIntegrationTest { engines = engines.copy(contextPackBuilder = droppingBuilder, promptResolver = PromptResolver { it }), retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) val sessionId = SessionId("s-trunc") val config = OrchestrationConfig(retryPolicy = RetryPolicy(maxAttempts = 1, backoffMs = 0)) @@ -386,6 +399,7 @@ class SessionOrchestratorIntegrationTest { engines = engines, retryCoordinator = retryCoordinator, artifactStore = recordingStore, + decisionJournalRepository = decisionJournalRepository, ) recordingOrchestrator.run(sessionId, graph, config) @@ -480,6 +494,7 @@ class SessionOrchestratorIntegrationTest { engines = engines, retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) livenessOrchestrator.run(sessionId, livenessGraph, config) @@ -596,6 +611,7 @@ class SessionOrchestratorIntegrationTest { ), retryCoordinator = retryCoordinator, artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) livenessOrchestrator.run(sessionId, livenessGraph, config) diff --git a/testing/integration/src/test/kotlin/ToolCallGateTest.kt b/testing/integration/src/test/kotlin/ToolCallGateTest.kt index b5b4f2c3..ae7f6e72 100644 --- a/testing/integration/src/test/kotlin/ToolCallGateTest.kt +++ b/testing/integration/src/test/kotlin/ToolCallGateTest.kt @@ -38,6 +38,9 @@ import com.correx.core.inference.ToolCallFunction import com.correx.core.inference.ToolCallRequest import com.correx.core.inference.TokenUsage import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer +import com.correx.core.journal.DecisionJournalProjector +import com.correx.core.journal.DefaultDecisionJournalReducer +import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.kernel.orchestration.DefaultSessionOrchestrator import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationProjector @@ -220,11 +223,15 @@ class ToolCallGateTest { workspacePolicy = policy, ) + val decisionJournalRepository = DefaultDecisionJournalRepository( + DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())), + ) val orchestrator = DefaultSessionOrchestrator( repositories = repositories, engines = engines, retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) return Triple(orchestrator, eventStore, provider) diff --git a/testing/integration/src/test/kotlin/WorkspaceScopedToolRegistryTest.kt b/testing/integration/src/test/kotlin/WorkspaceScopedToolRegistryTest.kt index cb46b94f..496bbe0e 100644 --- a/testing/integration/src/test/kotlin/WorkspaceScopedToolRegistryTest.kt +++ b/testing/integration/src/test/kotlin/WorkspaceScopedToolRegistryTest.kt @@ -31,6 +31,9 @@ import com.correx.core.inference.ToolCallFunction import com.correx.core.inference.ToolCallRequest import com.correx.core.inference.TokenUsage import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer +import com.correx.core.journal.DecisionJournalProjector +import com.correx.core.journal.DefaultDecisionJournalReducer +import com.correx.core.journal.DefaultDecisionJournalRepository import com.correx.core.kernel.orchestration.DefaultSessionOrchestrator import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationProjector @@ -210,11 +213,15 @@ class WorkspaceScopedToolRegistryTest { workspaceToolRegistryProvider = workspaceProvider, ) + val decisionJournalRepository = DefaultDecisionJournalRepository( + DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())), + ) return DefaultSessionOrchestrator( repositories = repositories, engines = engines, retryCoordinator = DefaultRetryCoordinator(eventStore), artifactStore = artifactStore, + decisionJournalRepository = decisionJournalRepository, ) to eventStore } diff --git a/testing/replay/build.gradle b/testing/replay/build.gradle index 42d0ac05..b2179194 100644 --- a/testing/replay/build.gradle +++ b/testing/replay/build.gradle @@ -6,6 +6,7 @@ plugins { dependencies { testImplementation(project(":core:events")) + testImplementation(project(":core:journal")) testImplementation(project(":core:sessions")) testImplementation(project(":core:transitions")) testImplementation(project(":core:validation")) diff --git a/testing/replay/src/test/kotlin/DecisionJournalReplayTest.kt b/testing/replay/src/test/kotlin/DecisionJournalReplayTest.kt new file mode 100644 index 00000000..3a2f82e3 --- /dev/null +++ b/testing/replay/src/test/kotlin/DecisionJournalReplayTest.kt @@ -0,0 +1,50 @@ +import com.correx.core.events.events.EventMetadata +import com.correx.core.events.events.NewEvent +import com.correx.core.events.events.SteeringNoteAddedEvent +import com.correx.core.events.events.TransitionExecutedEvent +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.events.types.TransitionId +import com.correx.core.journal.DecisionJournalProjector +import com.correx.core.journal.DecisionJournalRenderer +import com.correx.core.journal.DefaultDecisionJournalReducer +import com.correx.core.journal.DefaultDecisionJournalRepository +import com.correx.core.sessions.projections.replay.DefaultEventReplayer +import com.correx.infrastructure.persistence.InMemoryEventStore +import kotlinx.coroutines.runBlocking +import kotlinx.datetime.Clock +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class DecisionJournalReplayTest { + + private fun meta(sessionId: SessionId, id: String) = + EventMetadata(EventId(id), sessionId, Clock.System.now(), 1, null, null) + + @Test + fun `journal replay is deterministic`(): Unit = runBlocking { + val sessionId = SessionId("s") + val store = InMemoryEventStore() + store.append(NewEvent(meta(sessionId, "steer"), SteeringNoteAddedEvent(sessionId, "use jwt"))) + store.append( + NewEvent( + meta(sessionId, "trans"), + TransitionExecutedEvent(sessionId, StageId("a"), StageId("b"), TransitionId("t")), + ), + ) + + val repo = DefaultDecisionJournalRepository( + DefaultEventReplayer(store, DecisionJournalProjector(DefaultDecisionJournalReducer())), + ) + val renderer = DecisionJournalRenderer() + + val a = renderer.render(repo.getJournal(sessionId)) + val b = renderer.render(repo.getJournal(sessionId)) + + assertEquals(a, b) + assertTrue(a.contains("use jwt")) + assertTrue(a.contains("a → b")) + } +}