feat(kernel): inject decision journal into stage context + wire repository

This commit is contained in:
2026-06-04 00:48:07 +04:00
parent bd4dd91bf1
commit bcc20509e0
16 changed files with 129 additions and 2 deletions
+1
View File
@@ -21,6 +21,7 @@ dependencies {
implementation project(':core:approvals') implementation project(':core:approvals')
implementation project(':core:sessions') implementation project(':core:sessions')
implementation project(':core:kernel') implementation project(':core:kernel')
implementation project(':core:journal')
implementation project(':core:inference') implementation project(':core:inference')
implementation project(':core:transitions') implementation project(':core:transitions')
implementation project(':core:context') implementation project(':core:context')
@@ -230,12 +230,14 @@ fun main() {
workspacePolicy = workspacePolicy, workspacePolicy = workspacePolicy,
workspaceToolRegistryProvider = wsToolRegistryProvider, workspaceToolRegistryProvider = wsToolRegistryProvider,
) )
val decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore)
val orchestrator = DefaultSessionOrchestrator( val orchestrator = DefaultSessionOrchestrator(
repositories = repositories, repositories = repositories,
engines = engines, engines = engines,
retryCoordinator = DefaultRetryCoordinator(eventStore), retryCoordinator = DefaultRetryCoordinator(eventStore),
artifactStore = artifactStore, artifactStore = artifactStore,
tokenizer = firstProvider.tokenizer, tokenizer = firstProvider.tokenizer,
decisionJournalRepository = decisionJournalRepository,
) )
val defaultOrchestrationConfig = OrchestrationConfig( val defaultOrchestrationConfig = OrchestrationConfig(
sandboxRoot = sandboxRoot, sandboxRoot = sandboxRoot,
@@ -106,12 +106,14 @@ fun buildTestServerModule(
val repositories = buildRepositories(eventStore) val repositories = buildRepositories(eventStore)
val decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore)
val orchestrator = DefaultSessionOrchestrator( val orchestrator = DefaultSessionOrchestrator(
repositories = repositories, repositories = repositories,
engines = engines, engines = engines,
retryCoordinator = DefaultRetryCoordinator(eventStore), retryCoordinator = DefaultRetryCoordinator(eventStore),
artifactStore = artifactStore, artifactStore = artifactStore,
tokenizer = firstProvider.tokenizer, tokenizer = firstProvider.tokenizer,
decisionJournalRepository = decisionJournalRepository,
) )
val routerFacade = InfrastructureModule.createRouterFacade( val routerFacade = InfrastructureModule.createRouterFacade(
@@ -183,6 +183,7 @@ class WorkspaceHandshakeTest {
retryCoordinator = DefaultRetryCoordinator(eventStore), retryCoordinator = DefaultRetryCoordinator(eventStore),
artifactStore = noopArtifactStore, artifactStore = noopArtifactStore,
tokenizer = provider.tokenizer, tokenizer = provider.tokenizer,
decisionJournalRepository = InfrastructureModule.createDecisionJournalRepository(eventStore),
) )
val routerFacade = InfrastructureModule.createRouterFacade( val routerFacade = InfrastructureModule.createRouterFacade(
+1
View File
@@ -17,6 +17,7 @@ dependencies {
implementation project(':core:artifacts-store') implementation project(':core:artifacts-store')
implementation project(':core:risk') implementation project(':core:risk')
implementation project(':core:toolintent') implementation project(':core:toolintent')
implementation(project(":core:journal"))
implementation "org.slf4j:slf4j-api:2.0.16" implementation "org.slf4j:slf4j-api:2.0.16"
} }
tasks.named("koverVerify").configure { enabled = false } tasks.named("koverVerify").configure { enabled = false }
@@ -4,6 +4,7 @@ import com.correx.core.approvals.ApprovalOutcome
import com.correx.core.approvals.ApprovalStatus import com.correx.core.approvals.ApprovalStatus
import com.correx.core.approvals.model.ApprovalDecision import com.correx.core.approvals.model.ApprovalDecision
import com.correx.core.inference.Tokenizer import com.correx.core.inference.Tokenizer
import com.correx.core.journal.DefaultDecisionJournalRepository
import com.correx.core.artifacts.ArtifactState import com.correx.core.artifacts.ArtifactState
import com.correx.core.artifactstore.ArtifactStore import com.correx.core.artifactstore.ArtifactStore
import com.correx.core.events.events.ApprovalDecisionResolvedEvent import com.correx.core.events.events.ApprovalDecisionResolvedEvent
@@ -40,7 +41,8 @@ class DefaultSessionOrchestrator(
private val retryCoordinator: RetryCoordinator, private val retryCoordinator: RetryCoordinator,
artifactStore: ArtifactStore, artifactStore: ArtifactStore,
tokenizer: Tokenizer? = null, tokenizer: Tokenizer? = null,
) : SessionOrchestrator(repositories, engines, artifactStore), ApprovalGateway { decisionJournalRepository: DefaultDecisionJournalRepository,
) : SessionOrchestrator(repositories, engines, artifactStore, decisionJournalRepository), ApprovalGateway {
override val tokenizer: Tokenizer? = tokenizer override val tokenizer: Tokenizer? = tokenizer
override val cancellations: ConcurrentHashMap<SessionId, AtomicBoolean> = override val cancellations: ConcurrentHashMap<SessionId, AtomicBoolean> =
ConcurrentHashMap<SessionId, AtomicBoolean>() ConcurrentHashMap<SessionId, AtomicBoolean>()
@@ -2,6 +2,9 @@ package com.correx.core.kernel.orchestration
import com.correx.core.approvals.domain.NoOpApprovalEngine import com.correx.core.approvals.domain.NoOpApprovalEngine
import com.correx.core.artifactstore.ArtifactStore 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.context.model.ContextPack
import com.correx.core.events.events.ArtifactValidatedEvent import com.correx.core.events.events.ArtifactValidatedEvent
import com.correx.core.events.events.ArtifactValidatingEvent import com.correx.core.events.events.ArtifactValidatingEvent
@@ -58,6 +61,9 @@ class ReplayOrchestrator(
repositories, repositories,
engines.copy(approvalEngine = NoOpApprovalEngine(), riskAssessor = NoOpRiskAssessor()), engines.copy(approvalEngine = NoOpApprovalEngine(), riskAssessor = NoOpRiskAssessor()),
artifactStore, artifactStore,
DefaultDecisionJournalRepository(object : EventReplayer<DecisionJournalState> {
override fun rebuild(sessionId: SessionId) = DecisionJournalState()
}),
) { ) {
private val replayProvider = ReplayInferenceProvider(repositories.eventStore) private val replayProvider = ReplayInferenceProvider(repositories.eventStore)
@@ -107,6 +107,8 @@ import kotlinx.datetime.Clock
import kotlinx.serialization.encodeToString import kotlinx.serialization.encodeToString
import kotlinx.serialization.json.Json import kotlinx.serialization.json.Json
import kotlinx.serialization.json.JsonObject import kotlinx.serialization.json.JsonObject
import com.correx.core.journal.DecisionJournalRenderer
import com.correx.core.journal.DefaultDecisionJournalRepository
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import java.util.* import java.util.*
import java.util.concurrent.* import java.util.concurrent.*
@@ -129,6 +131,8 @@ abstract class SessionOrchestrator(
repositories: OrchestratorRepositories, repositories: OrchestratorRepositories,
engines: OrchestratorEngines, engines: OrchestratorEngines,
private val artifactStore: ArtifactStore, private val artifactStore: ArtifactStore,
private val decisionJournalRepository: DefaultDecisionJournalRepository,
private val decisionJournalRenderer: DecisionJournalRenderer = DecisionJournalRenderer(),
) { ) {
private val log = LoggerFactory.getLogger(this::class.java) private val log = LoggerFactory.getLogger(this::class.java)
private val eventStore: EventStore = repositories.eventStore 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 // 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 // across rounds; instead we grow our own ordered list and let the builder restamp
// ordinals from it each round. // 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( val contextPack = contextPackBuilder.build(
id = ContextPackId(UUID.randomUUID().toString()), id = ContextPackId(UUID.randomUUID().toString()),
sessionId = sessionId, sessionId = sessionId,
+1
View File
@@ -22,6 +22,7 @@ dependencies {
implementation project(":infrastructure:workflow") implementation project(":infrastructure:workflow")
implementation project(":core:artifacts") implementation project(":core:artifacts")
implementation project(":core:artifacts-store") implementation project(":core:artifacts-store")
implementation(project(":core:journal"))
implementation project(":infrastructure:artifacts-cas") implementation project(":infrastructure:artifacts-cas")
implementation "io.ktor:ktor-client-core:$ktor_version" implementation "io.ktor:ktor-client-core:$ktor_version"
implementation "io.ktor:ktor-client-cio:$ktor_version" implementation "io.ktor:ktor-client-cio:$ktor_version"
@@ -3,6 +3,9 @@ package com.correx.infrastructure
import com.correx.core.approvals.ApprovalProjector import com.correx.core.approvals.ApprovalProjector
import com.correx.core.approvals.DefaultApprovalReducer import com.correx.core.approvals.DefaultApprovalReducer
import com.correx.core.approvals.DefaultApprovalRepository 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.DefaultArtifactReducer
import com.correx.core.artifacts.repository.ArtifactRepository import com.correx.core.artifacts.repository.ArtifactRepository
import com.correx.core.artifactstore.ArtifactStore 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<ArtifactKind> = emptyList()): WorkflowLoader { fun createWorkflowLoader(extraKinds: List<ArtifactKind> = emptyList()): WorkflowLoader {
val registry = DefaultArtifactKindRegistry() val registry = DefaultArtifactKindRegistry()
extraKinds.forEach { registry.register(it) } extraKinds.forEach { registry.register(it) }
+1
View File
@@ -13,6 +13,7 @@ dependencies {
testImplementation(project(":core:context")) testImplementation(project(":core:context"))
testImplementation(project(":core:inference")) testImplementation(project(":core:inference"))
testImplementation(project(":core:kernel")) testImplementation(project(":core:kernel"))
testImplementation(project(":core:journal"))
testImplementation(project(":core:risk")) testImplementation(project(":core:risk"))
testImplementation(project(":core:tools")) testImplementation(project(":core:tools"))
testImplementation(project(":core:toolintent")) testImplementation(project(":core:toolintent"))
@@ -52,6 +52,9 @@ import com.correx.core.inference.Tokenizer
import com.correx.core.inference.ToolCallFunction import com.correx.core.inference.ToolCallFunction
import com.correx.core.inference.ToolCallRequest import com.correx.core.inference.ToolCallRequest
import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer 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.DefaultSessionOrchestrator
import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationConfig
import com.correx.core.kernel.orchestration.OrchestrationProjector import com.correx.core.kernel.orchestration.OrchestrationProjector
@@ -124,6 +127,10 @@ class SessionOrchestratorIntegrationTest {
DefaultEventReplayer(eventStore, ApprovalProjector(DefaultApprovalReducer())), DefaultEventReplayer(eventStore, ApprovalProjector(DefaultApprovalReducer())),
) )
private val decisionJournalRepository = DefaultDecisionJournalRepository(
DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())),
)
private val repositories = OrchestratorRepositories( private val repositories = OrchestratorRepositories(
eventStore = eventStore, eventStore = eventStore,
inferenceRepository = inferenceRepository, inferenceRepository = inferenceRepository,
@@ -147,6 +154,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines, engines = engines,
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
@Test @Test
@@ -194,6 +202,7 @@ class SessionOrchestratorIntegrationTest {
), ),
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
failingOrchestrator.run(sessionId, graph, config) failingOrchestrator.run(sessionId, graph, config)
@@ -221,6 +230,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines.copy(validationPipeline = approvingPipeline), engines = engines.copy(validationPipeline = approvingPipeline),
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
val runJob = launch { orchestrator.run(sessionId, graph, config) } val runJob = launch { orchestrator.run(sessionId, graph, config) }
@@ -269,6 +279,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines.copy(validationPipeline = approvingPipeline), engines = engines.copy(validationPipeline = approvingPipeline),
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
val runJob = launch { approvalOrchestrator.run(sessionId, graph, config) } val runJob = launch { approvalOrchestrator.run(sessionId, graph, config) }
@@ -322,6 +333,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines, engines = engines,
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
val sessionId = SessionId("s6") val sessionId = SessionId("s6")
@@ -364,6 +376,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines.copy(contextPackBuilder = droppingBuilder, promptResolver = PromptResolver { it }), engines = engines.copy(contextPackBuilder = droppingBuilder, promptResolver = PromptResolver { it }),
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
val sessionId = SessionId("s-trunc") val sessionId = SessionId("s-trunc")
val config = OrchestrationConfig(retryPolicy = RetryPolicy(maxAttempts = 1, backoffMs = 0)) val config = OrchestrationConfig(retryPolicy = RetryPolicy(maxAttempts = 1, backoffMs = 0))
@@ -386,6 +399,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines, engines = engines,
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = recordingStore, artifactStore = recordingStore,
decisionJournalRepository = decisionJournalRepository,
) )
recordingOrchestrator.run(sessionId, graph, config) recordingOrchestrator.run(sessionId, graph, config)
@@ -480,6 +494,7 @@ class SessionOrchestratorIntegrationTest {
engines = engines, engines = engines,
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
livenessOrchestrator.run(sessionId, livenessGraph, config) livenessOrchestrator.run(sessionId, livenessGraph, config)
@@ -596,6 +611,7 @@ class SessionOrchestratorIntegrationTest {
), ),
retryCoordinator = retryCoordinator, retryCoordinator = retryCoordinator,
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
livenessOrchestrator.run(sessionId, livenessGraph, config) livenessOrchestrator.run(sessionId, livenessGraph, config)
@@ -38,6 +38,9 @@ import com.correx.core.inference.ToolCallFunction
import com.correx.core.inference.ToolCallRequest import com.correx.core.inference.ToolCallRequest
import com.correx.core.inference.TokenUsage import com.correx.core.inference.TokenUsage
import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer 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.DefaultSessionOrchestrator
import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationConfig
import com.correx.core.kernel.orchestration.OrchestrationProjector import com.correx.core.kernel.orchestration.OrchestrationProjector
@@ -220,11 +223,15 @@ class ToolCallGateTest {
workspacePolicy = policy, workspacePolicy = policy,
) )
val decisionJournalRepository = DefaultDecisionJournalRepository(
DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())),
)
val orchestrator = DefaultSessionOrchestrator( val orchestrator = DefaultSessionOrchestrator(
repositories = repositories, repositories = repositories,
engines = engines, engines = engines,
retryCoordinator = DefaultRetryCoordinator(eventStore), retryCoordinator = DefaultRetryCoordinator(eventStore),
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) )
return Triple(orchestrator, eventStore, provider) return Triple(orchestrator, eventStore, provider)
@@ -31,6 +31,9 @@ import com.correx.core.inference.ToolCallFunction
import com.correx.core.inference.ToolCallRequest import com.correx.core.inference.ToolCallRequest
import com.correx.core.inference.TokenUsage import com.correx.core.inference.TokenUsage
import com.correx.core.kernel.orchestration.DefaultOrchestrationReducer 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.DefaultSessionOrchestrator
import com.correx.core.kernel.orchestration.OrchestrationConfig import com.correx.core.kernel.orchestration.OrchestrationConfig
import com.correx.core.kernel.orchestration.OrchestrationProjector import com.correx.core.kernel.orchestration.OrchestrationProjector
@@ -210,11 +213,15 @@ class WorkspaceScopedToolRegistryTest {
workspaceToolRegistryProvider = workspaceProvider, workspaceToolRegistryProvider = workspaceProvider,
) )
val decisionJournalRepository = DefaultDecisionJournalRepository(
DefaultEventReplayer(eventStore, DecisionJournalProjector(DefaultDecisionJournalReducer())),
)
return DefaultSessionOrchestrator( return DefaultSessionOrchestrator(
repositories = repositories, repositories = repositories,
engines = engines, engines = engines,
retryCoordinator = DefaultRetryCoordinator(eventStore), retryCoordinator = DefaultRetryCoordinator(eventStore),
artifactStore = artifactStore, artifactStore = artifactStore,
decisionJournalRepository = decisionJournalRepository,
) to eventStore ) to eventStore
} }
+1
View File
@@ -6,6 +6,7 @@ plugins {
dependencies { dependencies {
testImplementation(project(":core:events")) testImplementation(project(":core:events"))
testImplementation(project(":core:journal"))
testImplementation(project(":core:sessions")) testImplementation(project(":core:sessions"))
testImplementation(project(":core:transitions")) testImplementation(project(":core:transitions"))
testImplementation(project(":core:validation")) testImplementation(project(":core:validation"))
@@ -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"))
}
}