feat(router,server): wire workflow inventory and project profile into router context per turn
This commit is contained in:
@@ -292,6 +292,9 @@ fun main() {
|
|||||||
// persisted vectors survive restart but the metadata map does not (no-op for
|
// persisted vectors survive restart but the metadata map does not (no-op for
|
||||||
// non-rehydratable backends like in_memory).
|
// non-rehydratable backends like in_memory).
|
||||||
runBlocking { L3MetadataRehydrator(eventStore).rehydrate(l3MemoryStore) }
|
runBlocking { L3MetadataRehydrator(eventStore).rehydrate(l3MemoryStore) }
|
||||||
|
val workflowRegistry = FileSystemWorkflowRegistry(
|
||||||
|
InfrastructureModule.createWorkflowLoader(configArtifactKindsEarly),
|
||||||
|
)
|
||||||
// Builds the router facade from a config snapshot, mapping the [router] block onto the domain
|
// Builds the router facade from a config snapshot, mapping the [router] block onto the domain
|
||||||
// RouterConfig. Reused by ConfigService's rebuild hook so router knob edits apply live (the
|
// RouterConfig. Reused by ConfigService's rebuild hook so router knob edits apply live (the
|
||||||
// embedder and L3 store are construction-time and reused; the facade is what gets swapped).
|
// embedder and L3 store are construction-time and reused; the facade is what gets swapped).
|
||||||
@@ -320,6 +323,16 @@ fun main() {
|
|||||||
tokenizer = firstProvider.tokenizer,
|
tokenizer = firstProvider.tokenizer,
|
||||||
embedder = embedder,
|
embedder = embedder,
|
||||||
l3MemoryStore = l3MemoryStore,
|
l3MemoryStore = l3MemoryStore,
|
||||||
|
workflowSummaryProvider = {
|
||||||
|
workflowRegistry.listAll().map { s ->
|
||||||
|
com.correx.core.router.model.WorkflowSummary(s.workflowId, s.description, s.stageIds)
|
||||||
|
}
|
||||||
|
},
|
||||||
|
sessionProfileProvider = { sid ->
|
||||||
|
runCatching { repositories.sessionRepository.getSession(sid).state.boundProjectProfile }
|
||||||
|
.getOrNull()
|
||||||
|
?.let { renderProjectProfileText(it) }
|
||||||
|
},
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
val routerFacade = buildRouterFacade(correxConfig)
|
val routerFacade = buildRouterFacade(correxConfig)
|
||||||
@@ -383,9 +396,7 @@ fun main() {
|
|||||||
eventStore = eventStore,
|
eventStore = eventStore,
|
||||||
artifactStore = artifactStore,
|
artifactStore = artifactStore,
|
||||||
sessionRepository = repositories.sessionRepository,
|
sessionRepository = repositories.sessionRepository,
|
||||||
workflowRegistry = FileSystemWorkflowRegistry(
|
workflowRegistry = workflowRegistry,
|
||||||
InfrastructureModule.createWorkflowLoader(configArtifactKindsEarly),
|
|
||||||
),
|
|
||||||
providerRegistry = infraRegistry.asServerRegistry(),
|
providerRegistry = infraRegistry.asServerRegistry(),
|
||||||
defaultOrchestrationConfig = defaultOrchestrationConfig,
|
defaultOrchestrationConfig = defaultOrchestrationConfig,
|
||||||
routerFacade = routerFacade,
|
routerFacade = routerFacade,
|
||||||
@@ -418,6 +429,22 @@ fun main() {
|
|||||||
embeddedServer(Netty, port = 8080) { configureServer(module) }.start(wait = true)
|
embeddedServer(Netty, port = 8080) { configureServer(module) }.start(wait = true)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fun renderProjectProfileText(profile: com.correx.core.sessions.BoundProjectProfile?): String? {
|
||||||
|
if (profile == null) return null
|
||||||
|
return buildString {
|
||||||
|
append("## Project profile\n")
|
||||||
|
if (profile.about.isNotBlank()) append(profile.about).append("\n")
|
||||||
|
if (profile.conventions.isNotEmpty()) {
|
||||||
|
append("### Conventions\n")
|
||||||
|
profile.conventions.forEach { append("- ").append(it).append("\n") }
|
||||||
|
}
|
||||||
|
if (profile.commands.isNotEmpty()) {
|
||||||
|
append("### Commands\n")
|
||||||
|
profile.commands.forEach { (name, cmd) -> append(name).append(": ").append(cmd).append("\n") }
|
||||||
|
}
|
||||||
|
}.trimEnd().takeIf { it.isNotBlank() }
|
||||||
|
}
|
||||||
|
|
||||||
private fun logModelInfo(
|
private fun logModelInfo(
|
||||||
modelId: String,
|
modelId: String,
|
||||||
infraRegistry: DefaultProviderRegistry,
|
infraRegistry: DefaultProviderRegistry,
|
||||||
|
|||||||
@@ -26,6 +26,7 @@ import com.correx.core.router.l3.L3Query
|
|||||||
import com.correx.core.router.model.NarrationTrigger
|
import com.correx.core.router.model.NarrationTrigger
|
||||||
import com.correx.core.router.model.RouterConfig
|
import com.correx.core.router.model.RouterConfig
|
||||||
import com.correx.core.router.model.RouterResponse
|
import com.correx.core.router.model.RouterResponse
|
||||||
|
import com.correx.core.router.model.WorkflowSummary
|
||||||
import kotlinx.coroutines.CancellationException
|
import kotlinx.coroutines.CancellationException
|
||||||
import kotlinx.datetime.Clock
|
import kotlinx.datetime.Clock
|
||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
@@ -52,6 +53,8 @@ class DefaultRouterFacade(
|
|||||||
private val validateSteering: (suspend (String) -> String?)? = null,
|
private val validateSteering: (suspend (String) -> String?)? = null,
|
||||||
private val embedder: Embedder,
|
private val embedder: Embedder,
|
||||||
private val l3MemoryStore: L3MemoryStore,
|
private val l3MemoryStore: L3MemoryStore,
|
||||||
|
private val workflowSummaryProvider: () -> List<WorkflowSummary> = { emptyList() },
|
||||||
|
private val sessionProfileProvider: suspend (SessionId) -> String? = { null },
|
||||||
) : RouterFacade {
|
) : RouterFacade {
|
||||||
|
|
||||||
override suspend fun onUserInput(
|
override suspend fun onUserInput(
|
||||||
@@ -115,7 +118,7 @@ class DefaultRouterFacade(
|
|||||||
// Rebuild state so lastRetrievedMemory reflects this turn
|
// Rebuild state so lastRetrievedMemory reflects this turn
|
||||||
state = routerRepository.getRouterState(sessionId)
|
state = routerRepository.getRouterState(sessionId)
|
||||||
|
|
||||||
val contextPack = routerContextBuilder.build(state, config.tokenBudget)
|
val contextPack = routerContextBuilder.build(state, config.tokenBudget, workflowSummaryProvider(), sessionProfileProvider(sessionId))
|
||||||
|
|
||||||
if (contextPack.compressionMetadata.entriesDropped > 0) {
|
if (contextPack.compressionMetadata.entriesDropped > 0) {
|
||||||
val truncNow = Clock.System.now()
|
val truncNow = Clock.System.now()
|
||||||
|
|||||||
@@ -31,6 +31,8 @@ import com.correx.core.router.RouterProjector
|
|||||||
import com.correx.core.router.l3.InMemoryL3MemoryStore
|
import com.correx.core.router.l3.InMemoryL3MemoryStore
|
||||||
import com.correx.core.router.l3.L3MemoryStore
|
import com.correx.core.router.l3.L3MemoryStore
|
||||||
import com.correx.core.router.model.RouterConfig
|
import com.correx.core.router.model.RouterConfig
|
||||||
|
import com.correx.core.router.model.WorkflowSummary
|
||||||
|
import com.correx.core.events.types.SessionId
|
||||||
import com.correx.core.sessions.projections.replay.DefaultEventReplayer
|
import com.correx.core.sessions.projections.replay.DefaultEventReplayer
|
||||||
import com.correx.core.tools.contract.ToolExecutor
|
import com.correx.core.tools.contract.ToolExecutor
|
||||||
import com.correx.core.tools.registry.ToolRegistry
|
import com.correx.core.tools.registry.ToolRegistry
|
||||||
@@ -291,12 +293,14 @@ object InfrastructureModule {
|
|||||||
tokenizer: Tokenizer? = null,
|
tokenizer: Tokenizer? = null,
|
||||||
embedder: Embedder = createNoopEmbedder(),
|
embedder: Embedder = createNoopEmbedder(),
|
||||||
l3MemoryStore: L3MemoryStore = createInMemoryL3MemoryStore(),
|
l3MemoryStore: L3MemoryStore = createInMemoryL3MemoryStore(),
|
||||||
|
workflowSummaryProvider: () -> List<WorkflowSummary> = { emptyList() },
|
||||||
|
sessionProfileProvider: suspend (SessionId) -> String? = { null },
|
||||||
): RouterFacade {
|
): RouterFacade {
|
||||||
val reducer = DefaultRouterReducer()
|
val reducer = DefaultRouterReducer()
|
||||||
val projector = RouterProjector(reducer)
|
val projector = RouterProjector(reducer)
|
||||||
val replayer = DefaultEventReplayer(eventStore, projector)
|
val replayer = DefaultEventReplayer(eventStore, projector)
|
||||||
val repository = DefaultRouterRepository(replayer)
|
val repository = DefaultRouterRepository(replayer)
|
||||||
val contextBuilder = DefaultRouterContextBuilder(config, tokenizer)
|
val contextBuilder = DefaultRouterContextBuilder(config, tokenizer)
|
||||||
return DefaultRouterFacade(repository, contextBuilder, inferenceRouter, eventStore, config, null, embedder, l3MemoryStore)
|
return DefaultRouterFacade(repository, contextBuilder, inferenceRouter, eventStore, config, null, embedder, l3MemoryStore, workflowSummaryProvider, sessionProfileProvider)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1274,6 +1274,77 @@ class RouterFacadeTest {
|
|||||||
assertEquals(0, threeArgCallCount, "onUserInput must not call 3-arg route()")
|
assertEquals(0, threeArgCallCount, "onUserInput must not call 3-arg route()")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --------------------------------------------------------------------------
|
||||||
|
// A4 — workflowSummaryProvider and sessionProfileProvider forwarding
|
||||||
|
// --------------------------------------------------------------------------
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `workflowSummaryProvider is called and forwarded to context builder on each turn`(): Unit = runBlocking {
|
||||||
|
val captured = mutableListOf<List<WorkflowSummary>>()
|
||||||
|
val capturingBuilder = object : RouterContextBuilder {
|
||||||
|
override suspend fun build(
|
||||||
|
state: RouterState,
|
||||||
|
budget: TokenBudget,
|
||||||
|
availableWorkflows: List<WorkflowSummary>,
|
||||||
|
projectProfileText: String?,
|
||||||
|
): ContextPack {
|
||||||
|
captured.add(availableWorkflows)
|
||||||
|
return emptyContextPack()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
val workflow = WorkflowSummary("wf", "desc", listOf("s1"))
|
||||||
|
val facade = facadeWith(contextBuilder = capturingBuilder, workflowSummaryProvider = { listOf(workflow) })
|
||||||
|
facade.onUserInput(SessionId("s"), "hello")
|
||||||
|
assertEquals(listOf(listOf(workflow)), captured)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `sessionProfileProvider is called and forwarded to context builder`(): Unit = runBlocking {
|
||||||
|
val capturedProfiles = mutableListOf<String?>()
|
||||||
|
val capturingBuilder = object : RouterContextBuilder {
|
||||||
|
override suspend fun build(
|
||||||
|
state: RouterState,
|
||||||
|
budget: TokenBudget,
|
||||||
|
availableWorkflows: List<WorkflowSummary>,
|
||||||
|
projectProfileText: String?,
|
||||||
|
): ContextPack {
|
||||||
|
capturedProfiles.add(projectProfileText)
|
||||||
|
return emptyContextPack()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
val facade = facadeWith(
|
||||||
|
contextBuilder = capturingBuilder,
|
||||||
|
sessionProfileProvider = { "profile text" },
|
||||||
|
)
|
||||||
|
facade.onUserInput(SessionId("s"), "hello")
|
||||||
|
assertEquals(listOf("profile text"), capturedProfiles)
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun facadeWith(
|
||||||
|
contextBuilder: RouterContextBuilder = object : RouterContextBuilder {
|
||||||
|
override suspend fun build(
|
||||||
|
state: RouterState,
|
||||||
|
budget: TokenBudget,
|
||||||
|
availableWorkflows: List<WorkflowSummary>,
|
||||||
|
projectProfileText: String?,
|
||||||
|
): ContextPack = emptyContextPack()
|
||||||
|
},
|
||||||
|
workflowSummaryProvider: () -> List<WorkflowSummary> = { emptyList() },
|
||||||
|
sessionProfileProvider: suspend (SessionId) -> String? = { null },
|
||||||
|
): RouterFacade = DefaultRouterFacade(
|
||||||
|
routerRepository = object : RouterRepository {
|
||||||
|
override suspend fun getRouterState(sessionId: SessionId): RouterState = RouterState(sessionId = sessionId)
|
||||||
|
},
|
||||||
|
routerContextBuilder = contextBuilder,
|
||||||
|
inferenceRouter = mockInferenceRouter("inference response"),
|
||||||
|
eventStore = mockEventStore(),
|
||||||
|
config = RouterConfig(tokenBudget = TokenBudget(limit = 5000)),
|
||||||
|
embedder = NoopEmbedder(dimension = 8),
|
||||||
|
l3MemoryStore = InMemoryL3MemoryStore(),
|
||||||
|
workflowSummaryProvider = workflowSummaryProvider,
|
||||||
|
sessionProfileProvider = sessionProfileProvider,
|
||||||
|
)
|
||||||
|
|
||||||
private fun emptyContextPack(): ContextPack = ContextPack(
|
private fun emptyContextPack(): ContextPack = ContextPack(
|
||||||
id = ContextPackId("empty"),
|
id = ContextPackId("empty"),
|
||||||
sessionId = SessionId("unknown"),
|
sessionId = SessionId("unknown"),
|
||||||
|
|||||||
Reference in New Issue
Block a user