diff --git a/examples/workflows/role_pipeline.toml b/examples/workflows/role_pipeline.toml index f4e41883..ce134bf6 100644 --- a/examples/workflows/role_pipeline.toml +++ b/examples/workflows/role_pipeline.toml @@ -1,13 +1,20 @@ -# Role pipeline: analyst → architect → brief_echo → planner → implementer ⇄ reviewer +# Role pipeline: analyst → architect → decomposer → ⟲(claim → implement → review)⟲ → done # -# Each stage produces a typed artifact that the next stage `needs`, so work flows forward -# without a human relaying notes. The decision journal (pinned into every stage's context) -# carries steering/approvals/verdicts across the whole run, so the reviewer sees the same -# ground truth as the planner. +# The architect's design is broken into a task graph by the decomposer, then a single loop +# drains those tasks one at a time. Each implementer entry deterministically CLAIMS the next +# ready task (claim_task = true) — the kernel picks it, not the LLM — and injects that task's +# acceptance criteria + blockers as L0 context. While a task is claimed, file writes are +# auto-scoped to its affected_paths; the implementer widens scope only via propose_scope +# (operator-approved). The loop is gated by the tasks_ready predicate, so it exits when the +# board is drained. # -# The implementer⇄reviewer loop is gated by review_report.verdict: -# approved → done -# changes_requested → back to implementer (capped by implementer.max_retries, then escalates) +# Loop control depends on transition ORDERING. The resolver evaluates a stage's outgoing +# transitions sorted by transition id and takes the first whose condition holds, so the +# reviewer edges are named review-1/2/3 to force this precedence: +# 1. review-1-changes (verdict ≠ approved) → implementer (fix the current task) +# 2. review-2-more (tasks_ready) → implementer (approved: claim next) +# 3. review-3-done (always_true, fallthrough) → done (approved + board empty) +# This avoids needing composite (all_of/not) conditions in the flat TOML schema. # # Requires these artifact kinds in ~/.config/correx/config.toml (schemas under docs/schemas/): # [[artifacts]] @@ -15,11 +22,7 @@ # [[artifacts]] # id = "design"; schema_path = "schemas/design.json"; llm_emitted = true # [[artifacts]] -# id = "impl_plan"; schema_path = "schemas/impl_plan.json"; llm_emitted = true -# [[artifacts]] # id = "review_report"; schema_path = "schemas/review_report.json"; llm_emitted = true -# [[artifacts]] -# id = "brief_echo"; schema_path = "schemas/brief_echo.json"; llm_emitted = true # # Prompt files (prompts/*.md, relative to this workflow) must exist for a real run. @@ -47,61 +50,42 @@ produces = [{ name = "design", kind = "design" }] token_budget = 16384 max_retries = 2 -# 3a. Before planning: echo the analyst brief back as a structured artifact. -# The brief_echo gate diffs the restatement against the original analysis; -# on divergence (dropped requirements or hallucinated files) it fails the -# stage (retryable) so the pipeline cannot reach plan generation on a misread brief. +# 3. Break the design into a task graph via task_decompose. require_task_decompose is a +# deterministic kernel gate: stage_complete is blocked until task_decompose (or task_create +# for a trivial single-task request) returns. Each child carries acceptance_criteria and +# affected_paths — those drive the per-task L0 context and the per-task write scope below. [[stages]] -id = "brief_echo" -prompt = "prompts/brief_echo.md" -needs = ["analysis"] -produces = [{ name = "brief_echo", kind = "brief_echo" }] -brief_echo = true -token_budget = 8192 -max_retries = 2 - -# 3b. Break the design into ordered, verifiable steps. -[[stages]] -id = "planner" -prompt = "prompts/planner.md" +id = "decomposer" +prompt = "prompts/task_planner.md" needs = ["design"] -produces = [{ name = "impl_plan", kind = "impl_plan" }] +allowed_tools = ["file_read", "ShellTool", "task_search", "task_context", "task_decompose", "task_create"] +require_task_decompose = true token_budget = 16384 max_retries = 2 -# 4. Implement the plan. Writes files (jailed to the workspace). The loop target — -# max_retries here caps how many review→implement refinement rounds are allowed. -# Optional: a `writes` manifest (workspace-relative globs) hard-bounds where this -# stage may write — a FILE_WRITE outside it is blocked as scope creep. Left open -# here because the targets are task-specific; a task-scoped workflow would set e.g. -# writes = ["core/sessions/**", "testing/sessions/**"] -# Optional: `static_analysis` runs deterministic tools (compiler/detekt/formatters) -# against the patch BEFORE the reviewer (role-reliability §5). A non-clean command -# fails this stage retryably with its output fed back verbatim, so only static-clean -# code reaches the reviewer — the reviewer never spends inference on what tools catch -# for free. Commands are workspace-specific, so left commented; a Kotlin/Gradle repo -# would set e.g. -# static_analysis = ["./gradlew compileKotlin -q", "./gradlew detekt -q"] +# 4. Implement one task. claim_task = true: on entry the kernel claims the next ready task +# (or re-surfaces the one already claimed on a review bounce) and injects its acceptance +# criteria as L0. Writes are auto-jailed to the claimed task's affected_paths; propose_scope +# is the operator-approved valve to widen that scope mid-task. [[stages]] id = "implementer" prompt = "prompts/implementer.md" -needs = ["impl_plan"] produces = [{ name = "patch", kind = "file_written" }] +claim_task = true allowed_tools = [ - "file_read", "file_write", "file_edit", "ShellTool", + "file_read", "file_write", "file_edit", "ShellTool", "propose_scope", "task_create", "task_update", "task_context", "task_search", ] -# static_analysis = ["./gradlew compileKotlin -q", "./gradlew detekt -q"] token_budget = 32768 max_retries = 3 -# 5. Review the patch against the plan AND the analyst's acceptance criteria. The reviewer -# needs `analysis` so it judges the diff against concrete, pre-stated criteria (§5 narrow -# question) rather than whole files against taste. +# 5. Review the patch against the claimed task's acceptance criteria and the analyst's brief. +# On approval the reviewer marks the task done via task_update — that drops it from the +# active claim so the next implementer entry claims a fresh task and the loop advances. [[stages]] id = "reviewer" prompt = "prompts/reviewer.md" -needs = ["patch", "impl_plan", "analysis"] +needs = ["patch", "analysis"] produces = [{ name = "review_report", kind = "review_report" }] allowed_tools = ["task_context", "task_update"] token_budget = 32768 @@ -117,25 +101,18 @@ condition_type = "artifact_validated" condition_artifact_id = "analysis" [[transitions]] -id = "architect-to-brief-echo" +id = "architect-to-decomposer" from = "architect" -to = "brief_echo" +to = "decomposer" condition_type = "artifact_validated" condition_artifact_id = "design" +# Only proceed to implementation once the decomposer actually produced ready tasks. [[transitions]] -id = "brief-echo-to-planner" -from = "brief_echo" -to = "planner" -condition_type = "artifact_validated" -condition_artifact_id = "brief_echo" - -[[transitions]] -id = "planner-to-implementer" -from = "planner" +id = "decomposer-to-implementer" +from = "decomposer" to = "implementer" -condition_type = "artifact_validated" -condition_artifact_id = "impl_plan" +condition_type = "tasks_ready" [[transitions]] id = "implementer-to-reviewer" @@ -144,20 +121,11 @@ to = "reviewer" condition_type = "artifact_validated" condition_artifact_id = "patch" -# --- verdict-gated loop exit / re-entry --- +# --- verdict + task-board gated loop (id order = evaluation precedence, see header) --- +# 1. Changes requested → keep refining the current (still-claimed) task. [[transitions]] -id = "review-approved" -from = "reviewer" -to = "done" -condition_type = "artifact_field_equals" -condition_artifact_id = "review_report" -condition_field = "verdict" -condition_value = "approved" -condition_operator = "eq" - -[[transitions]] -id = "review-changes-requested" +id = "review-1-changes" from = "reviewer" to = "implementer" condition_type = "artifact_field_equals" @@ -165,3 +133,17 @@ condition_artifact_id = "review_report" condition_field = "verdict" condition_value = "approved" condition_operator = "neq" + +# 2. Approved and tasks remain → claim the next one. +[[transitions]] +id = "review-2-more" +from = "reviewer" +to = "implementer" +condition_type = "tasks_ready" + +# 3. Approved and the board is drained → done (fallthrough). +[[transitions]] +id = "review-3-done" +from = "reviewer" +to = "done" +condition_type = "always_true" diff --git a/infrastructure/workflow/src/main/kotlin/com/correx/infrastructure/workflow/TomlWorkflowLoader.kt b/infrastructure/workflow/src/main/kotlin/com/correx/infrastructure/workflow/TomlWorkflowLoader.kt index e3ac9e73..21c94f6b 100644 --- a/infrastructure/workflow/src/main/kotlin/com/correx/infrastructure/workflow/TomlWorkflowLoader.kt +++ b/infrastructure/workflow/src/main/kotlin/com/correx/infrastructure/workflow/TomlWorkflowLoader.kt @@ -50,6 +50,7 @@ private data class StageSection( @param:JsonProperty("ground_references") val groundReferences: Boolean = false, @param:JsonProperty("brief_echo") val briefEcho: Boolean = false, @param:JsonProperty("require_task_decompose") val requireTaskDecompose: Boolean = false, + @param:JsonProperty("claim_task") val claimTask: Boolean = false, ) // condition fields flattened into the transition row @@ -119,15 +120,7 @@ class TomlWorkflowLoader( maxTokens = s.tokenBudget, ), maxRetries = s.maxRetries, - metadata = buildMap { - s.prompt?.let { put("prompt", resolvePromptPath(it, workflowDir)) } - s.systemPrompt?.let { put("systemPrompt", resolvePromptPath(it, workflowDir)) } - if (s.requiresApproval) put("requiresApproval", "true") - if (s.injectArtifactKinds) put("injectArtifactKinds", "true") - if (s.groundReferences) put("groundReferences", "true") - if (s.briefEcho) put("briefEcho", "true") - if (s.requireTaskDecompose) put("requireTaskDecompose", "true") - }, + metadata = s.toMetadata(workflowDir), ) } @@ -164,6 +157,17 @@ class TomlWorkflowLoader( return WorkflowGraph(id = id, stages = stageMap, transitions = transitionSet, start = startId) } + private fun StageSection.toMetadata(workflowDir: Path?): Map = buildMap { + prompt?.let { put("prompt", resolvePromptPath(it, workflowDir)) } + systemPrompt?.let { put("systemPrompt", resolvePromptPath(it, workflowDir)) } + if (requiresApproval) put("requiresApproval", "true") + if (injectArtifactKinds) put("injectArtifactKinds", "true") + if (groundReferences) put("groundReferences", "true") + if (briefEcho) put("briefEcho", "true") + if (requireTaskDecompose) put("requireTaskDecompose", "true") + if (claimTask) put("claimTask", "true") + } + @Suppress("ThrowsCount") private fun validate( start: String, diff --git a/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/RolePipelineWorkflowTest.kt b/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/RolePipelineWorkflowTest.kt index 52d58de9..fbba416a 100644 --- a/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/RolePipelineWorkflowTest.kt +++ b/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/RolePipelineWorkflowTest.kt @@ -7,6 +7,7 @@ import com.correx.core.artifacts.kind.JsonSchemaProperty import com.correx.core.events.types.StageId import com.correx.core.transitions.conditions.ArtifactFieldEquals import com.correx.core.transitions.conditions.FieldOperator +import com.correx.core.transitions.conditions.TasksReady import org.junit.jupiter.api.Assumptions.assumeTrue import org.junit.jupiter.api.Test import java.nio.file.Path @@ -29,8 +30,6 @@ class RolePipelineWorkflowTest { private val registry = DefaultArtifactKindRegistry().apply { register(llmKind("analysis")) register(llmKind("design")) - register(llmKind("brief_echo")) - register(llmKind("impl_plan")) register(llmKind("review_report")) } @@ -46,49 +45,67 @@ class RolePipelineWorkflowTest { return null } - @Test - fun `shipped role_pipeline toml parses with the full role graph`() { - val path = repoFile("examples/workflows/role_pipeline.toml") - assumeTrue(path != null, "role_pipeline.toml not found from ${System.getProperty("user.dir")}") + private fun loadGraph() = + repoFile("examples/workflows/role_pipeline.toml")?.let { TomlWorkflowLoader(registry).load(it) } - val graph = TomlWorkflowLoader(registry).load(path!!) + @Test + fun `shipped role_pipeline toml parses with the tasked execution graph`() { + val graph = loadGraph() + assumeTrue(graph != null, "role_pipeline.toml not found from ${System.getProperty("user.dir")}") + graph!! assertEquals( - setOf("analyst", "architect", "brief_echo", "planner", "implementer", "reviewer"), + setOf("analyst", "architect", "decomposer", "implementer", "reviewer"), graph.stages.keys.map { it.value }.toSet(), ) assertEquals(7, graph.transitions.size) - // Forward chaining via needs. - assertTrue(graph.stages[StageId("reviewer")]!!.needs.map { it.value }.containsAll(listOf("patch", "impl_plan"))) + // The decomposer is the mandatory task-graph gate. + assertEquals("true", graph.stages[StageId("decomposer")]!!.metadata["requireTaskDecompose"]) + // The implementer claims a task on entry. + assertEquals("true", graph.stages[StageId("implementer")]!!.metadata["claimTask"]) + } - // Verdict-gated loop: approved → done, changes → implementer. - val approved = graph.transitions.first { it.to == StageId("done") }.condition as ArtifactFieldEquals - assertEquals(FieldOperator.EQ, approved.operator) - assertEquals("verdict", approved.field) - val loop = graph.transitions - .first { it.from == StageId("reviewer") && it.to == StageId("implementer") } - .condition as ArtifactFieldEquals - assertEquals(FieldOperator.NEQ, loop.operator) + @Test + fun `loop is gated by verdict and the task board`() { + val graph = loadGraph() + assumeTrue(graph != null, "role_pipeline.toml not found") + graph!! + + // Decomposer only advances once tasks exist. + val toImpl = graph.transitions + .first { it.from == StageId("decomposer") && it.to == StageId("implementer") } + assertEquals(TasksReady, toImpl.condition) + + // Reviewer edges, in id (= evaluation) order: changes → impl, tasks_ready → impl, else → done. + val reviewerEdges = graph.transitions + .filter { it.from == StageId("reviewer") } + .sortedBy { it.id.value } + assertEquals( + listOf("review-1-changes", "review-2-more", "review-3-done"), + reviewerEdges.map { it.id.value }, + ) + // 1. changes requested (verdict != approved) loops back to the implementer. + val changes = reviewerEdges[0] + assertEquals(StageId("implementer"), changes.to) + assertEquals(FieldOperator.NEQ, (changes.condition as ArtifactFieldEquals).operator) + // 2. approved + tasks remain → claim the next. + assertEquals(StageId("implementer"), reviewerEdges[1].to) + assertEquals(TasksReady, reviewerEdges[1].condition) + // 3. approved + board drained → done (fallthrough). + assertEquals(StageId("done"), reviewerEdges[2].to) } @Test fun `task tools are granted to the right stages`() { - val path = repoFile("examples/workflows/role_pipeline.toml") - assumeTrue(path != null, "role_pipeline.toml not found from ${System.getProperty("user.dir")}") - val graph = TomlWorkflowLoader(registry).load(path!!) + val graph = loadGraph() + assumeTrue(graph != null, "role_pipeline.toml not found") + graph!! - // analyst opens + searches (read-only on files); creation is approval-gated (T2). - assertEquals( - setOf("file_read", "ShellTool", "task_search", "task_context", "task_create"), - graph.stages[StageId("analyst")]!!.allowedTools, - ) - // implementer owns the lifecycle: claim/submit + the full file set. - assertEquals( - setOf("file_read", "file_write", "file_edit", "ShellTool") + - setOf("task_create", "task_update", "task_context", "task_search"), - graph.stages[StageId("implementer")]!!.allowedTools, - ) + // implementer owns the lifecycle: claim/submit + the full file set + the scope valve. + val implTools = graph.stages[StageId("implementer")]!!.allowedTools + assertTrue("propose_scope" in implTools, "implementer must be able to propose scope widening") + assertTrue(implTools.containsAll(setOf("file_write", "file_edit", "task_update"))) // reviewer reads the task and completes it on an approved verdict. assertEquals( setOf("task_context", "task_update"),