diff --git a/run.json b/run.json
index a6e078a7b..0f844e0cc 100644
--- a/run.json
+++ b/run.json
@@ -492,7 +492,7 @@
"kind": "running"
},
"status_updated_at": "2026-05-23T19:56:01.493793Z",
- "last_event_at": "2026-05-23T20:26:00.219582Z",
+ "last_event_at": "2026-05-23T20:42:40.964151Z",
"pending_control": null,
"checkpoints": [
{
@@ -740,9 +740,9 @@
}
},
{
- "seq": 0,
+ "seq": 736,
"checkpoint": {
- "timestamp": "2026-05-23T20:26:00.320170Z",
+ "timestamp": "2026-05-23T20:26:04.061726Z",
"current_node": "implement",
"completed_nodes": [
"start",
@@ -753,31 +753,157 @@
],
"node_retries": {},
"context_values": {
- "internal.thread_id": "preflight_lint",
- "failure_class": "",
- "graph.goal": "# Output Schema Validation Implementation Plan\n\n> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.\n\n**Goal:** Add `output_schema` validation for agent and prompt nodes, with context-preserving repair turns when structured output does not validate.\n\n**Architecture:** Introduce a small structured-output layer in `fabro-workflow` that resolves node-level schema declarations, extracts JSON output, validates it, and produces either routing side effects or a parsed custom output context update. Agent and prompt execution must perform schema repair inside the active LLM conversation instead of using the workflow executor retry path.\n\n**Tech Stack:** Rust, Graphviz workflow attrs, `serde_json`, workspace `jsonschema`, existing `fabro-llm::ResponseFormat`, Fabro agent sessions, `cargo nextest`.\n\n---\n\n## Public Interface\n\nWorkflow authors can opt in on agent and prompt nodes:\n\n```dot\nreview [\n shape=tab,\n output_schema=\"routing\",\n output_retries=2\n]\n\naudit [\n shape=tab,\n output_schema=\"@schemas/audit-result.schema.json\",\n output_retries=2\n]\n```\n\n- `output_schema=\"routing\"` uses Fabro's built-in routing directive schema.\n- `output_schema=\"@path/to/schema.json\"` loads a JSON Schema file through existing workflow file-reference rules.\n- `output_retries` controls corrective turns inside the same node execution. Default: `2`. `0` means validate once and fail without a repair turn.\n- Schema failures are terminal node failures after `output_retries` is exhausted. They are not `retry_requested` outcomes and do not consume `max_retries`.\n- `backend=\"acp\"` with `output_schema` is unsupported in v1 and returns a clear validation error.\n\n## Implementation Tasks\n\n### Task 1: Node Attributes And File Reference Resolution\n\n**Files:**\n- Modify: `lib/crates/fabro-types/src/graph.rs`\n- Modify: `lib/crates/fabro-workflow/src/static_reference.rs`\n- Modify: `lib/crates/fabro-workflow/src/transforms/file_inlining.rs`\n- Test: existing unit tests in those files\n\n- [ ] Add `Node::output_schema(&self) -> Option<&str>` next to other agent/prompt attrs.\n- [ ] Add `Node::output_retries(&self) -> i64` returning `self.int_attr(\"output_retries\").unwrap_or(2).max(0)`.\n- [ ] Teach static reference validation that node attr `output_schema` values starting with `@` are file inline references.\n- [ ] Extend file inlining so `output_schema=\"@schemas/foo.json\"` is replaced with the schema file contents before execution, while `output_schema=\"routing\"` stays unchanged.\n- [ ] Add tests for absent attrs, default retries, zero retries, file inlining, and unresolved schema reference diagnostics.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-types -p fabro-workflow graph:: file_inlining static_reference\n```\n\nExpected: targeted tests pass.\n\n### Task 2: Structured Output Module\n\n**Files:**\n- Create: `lib/crates/fabro-workflow/src/handler/structured_output.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/mod.rs`\n- Modify: `lib/crates/fabro-workflow/Cargo.toml`\n- Test: unit tests in `structured_output.rs`\n\n- [ ] Add `jsonschema.workspace = true` to `fabro-workflow` dependencies.\n- [ ] Define `OutputSchemaKind` with `Routing` and `JsonSchema { schema: serde_json::Value }`.\n- [ ] Parse `node.output_schema()` into `None`, `Routing`, or custom JSON Schema. Treat literal `routing` as the only built-in keyword.\n- [ ] Add a built-in routing schema requiring an object with at least one recognized field: `preferred_next_label`, `outcome`, `failure_reason`, `suggested_next_ids`, or `context_updates`.\n- [ ] Reuse balanced-object scanning semantics for response text: validate the last JSON object that is relevant to the selected schema.\n- [ ] Return a structured validation result containing the parsed JSON object, concise error messages, and enough information to build a repair prompt.\n- [ ] Add tests for valid routing JSON, missing routing fields, wrong routing field types, valid custom schema, invalid custom schema, invalid JSON, and no JSON object.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow structured_output\n```\n\nExpected: structured-output unit tests pass.\n\n### Task 3: Routing Extraction Compatibility\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/agent.rs`\n- Test: existing agent handler unit tests\n\n- [ ] Keep the loose default unchanged when `output_schema` is absent.\n- [ ] Move current `STATUS_FIELDS`, balanced JSON scanning, and routing-field application behind reusable functions in `structured_output.rs` or call the new module from `agent.rs`.\n- [ ] For `output_schema=\"routing\"`, require schema-valid routing JSON and surface validation failures for repair instead of silently ignoring bad candidates.\n- [ ] Preserve existing routing fallback priority for agent nodes: response text first, then `status.json`, then last file touched.\n- [ ] Keep prompt-node routing behavior response-only unless later tasks explicitly add prompt `status.json` support.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::agent\n```\n\nExpected: existing loose routing tests still pass, plus new strict routing tests pass.\n\n### Task 4: Prompt Node Same-Context Repair\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/api.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/prompt.rs`\n- Test: prompt/API backend tests in those files\n\n- [ ] In `AgentApiBackend::one_shot`, keep `messages` mutable across attempts.\n- [ ] When a prompt node has a custom JSON Schema, set `response_format=JsonSchema` on the initial and repair LLM requests. For `routing`, use `JsonObject` or no provider-native schema if provider behavior would conflict with Fabro's routing extraction.\n- [ ] After each LLM response, validate according to `output_schema`.\n- [ ] On validation failure with repair attempts remaining, append `Message::assistant(response.text())`, then append a corrective `Message::user(repair_message)`, and call `client.complete` again with the same messages.\n- [ ] On success, return the validated response text and aggregate usage across all attempts.\n- [ ] On exhaustion, return a terminal failed outcome with failure reason `output schema validation failed after N repair attempt(s)`.\n- [ ] Update `PromptHandler` so validated custom output is added to `context_updates[\"output.{node_id}\"]`; routing output still updates outcome routing fields.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::prompt handler::llm::api\n```\n\nExpected: prompt repair keeps previous assistant output in the message list and succeeds after a corrective response.\n\n### Task 5: Agent Node Same-Session Repair\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/api.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/agent.rs`\n- Test: agent/API backend tests in those files\n\n- [ ] In `AgentApiBackend::run`, validate the final assistant response before releasing, closing, or caching the session.\n- [ ] On validation failure with repair attempts remaining, call `session.process_input(repair_message)` on the same `Session`.\n- [ ] Recompute the final assistant response after each repair turn from `session.history()`.\n- [ ] Aggregate usage across all new assistant turns, including repair turns, without double-counting reused session history.\n- [ ] Do not set provider-native `response_format` for agent sessions in v1, because agent sessions may need normal tool-use messages before final output.\n- [ ] Return terminal failure after exhaustion; do not return a retryable backend error and do not request workflow node retry.\n- [ ] Update `AgentHandler` to apply validated routing/custom output to the final `Outcome`.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::agent handler::llm::api\n```\n\nExpected: agent repair sends a second `process_input` to the same session and final validated output drives outcome/context updates.\n\n### Task 6: ACP Guardrail\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/acp.rs`\n- Test: ACP backend tests in that file\n\n- [ ] At the start of `AgentAcpBackend::run`, reject nodes where `node.output_schema().is_some()`.\n- [ ] Use a clear error message: `output_schema is not supported with backend=\"acp\" in this release`.\n- [ ] Add a test proving the ACP backend does not launch a process when `output_schema` is present.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::llm::acp\n```\n\nExpected: ACP guardrail test passes.\n\n### Task 7: Docs\n\n**Files:**\n- Modify: `docs/public/agents/outputs.mdx`\n- Modify: `docs/public/reference/dot-language.mdx`\n\n- [ ] Document `output_schema=\"routing\"` and `output_schema=\"@schema.json\"` under routing/structured outputs.\n- [ ] Document same-context repair behavior explicitly: Fabro sends validation feedback to the same agent/prompt context before failing.\n- [ ] Document `output_retries`, default `2`, and distinction from `max_retries`.\n- [ ] Document v1 scope: agent/prompt nodes only; ACP unsupported; custom schema output stored at `output.{node_id}`.\n\nRun:\n\n```bash\nrg -n \"output_schema|output_retries|output\\\\.\" docs/public/agents/outputs.mdx docs/public/reference/dot-language.mdx\n```\n\nExpected: docs mention the new attrs and storage behavior.\n\n### Task 8: Full Verification\n\n**Files:**\n- No new files beyond prior tasks\n\n- [ ] Run focused workflow tests:\n\n```bash\ncargo nextest run -p fabro-workflow\n```\n\n- [ ] Run formatting check:\n\n```bash\ncargo +nightly-2026-04-14 fmt --check --all\n```\n\n- [ ] Run clippy:\n\n```bash\ncargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings\n```\n\n- [ ] If snapshots change, inspect before accepting:\n\n```bash\ncargo insta pending-snapshots\n```\n\nOnly run `cargo insta accept` after verifying every pending snapshot is expected.\n\n## Acceptance Criteria\n\n- Existing workflows without `output_schema` behave exactly as before.\n- `output_schema=\"routing\"` prevents malformed/missing routing JSON from silently falling through to normal edge selection.\n- Invalid structured output results in a corrective LLM turn in the same context window.\n- Prompt repair preserves previous assistant output in the message list.\n- Agent repair preserves the same live session and does not re-run the node from scratch.\n- Exhausted output repair attempts produce a clear terminal failure.\n- Custom schema output is available to downstream nodes at `output.{node_id}`.\n- Docs clearly distinguish `output_retries` from `max_retries`.\n\n## Assumptions\n\n- `output_retries=2` is the default.\n- Custom schema validation targets the final JSON object in the response text.\n- `status.json` fallback remains routing-specific.\n- Provider-native response schema is used for prompt nodes only where it is safe.\n- ACP support can be added later after there is a guaranteed context-preserving repair mechanism.\n",
- "failure_signature": "",
- "internal.retry_count.toolchain": 0,
- "internal.retry_count.implement": 0,
- "internal.run_id": "01KSB6GTZ00T5V6BNMXN3SPKZF",
"internal.retry_count.start": 0,
- "graph.rankdir": "LR",
- "current_node": "implement",
- "thread.start.current_node": "toolchain",
- "thread.toolchain.current_node": "preflight_compile",
- "internal.node_visit_count": 1,
- "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126",
"internal.retry_count.preflight_compile": 0,
"thread.preflight_compile.current_node": "preflight_lint",
- "internal.retry_count.preflight_lint": 0,
+ "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126",
+ "thread.start.current_node": "toolchain",
+ "failure_class": "",
+ "current_node": "implement",
+ "failure_signature": "",
+ "graph.model_stylesheet": "\n * { model: claude-opus-4-7; }\n ",
+ "internal.node_visit_count": 1,
+ "internal.retry_count.implement": 0,
+ "internal.retry_count.toolchain": 0,
+ "last_response": "Implemented output schema validation end-to-end.\n\nSummary:\n- Added `output_schema` / `output_retries` node attrs.\n- Added `output_schema` static/file reference handling and inlining.\n- Added `handler:",
"thread.preflight_lint.current_node": "implement",
+ "internal.thread_id": "preflight_lint",
+ "response.implement": "Implemented output schema validation end-to-end.\n\nSummary:\n- Added `output_schema` / `output_retries` node attrs.\n- Added `output_schema` static/file reference handling and inlining.\n- Added `handler::structured_output` for:\n - routing/custom schema parsing\n - balanced JSON extraction\n - JSON Schema validation\n - repair prompts\n - routing/custom context application\n- Integrated strict routing/custom output into agent and prompt handlers.\n- Added same-context repair for:\n - prompt API calls via mutable message history\n - agent API sessions via the same live `Session`\n- Added ACP guardrail:\n - `output_schema is not supported with backend=\"acp\" in this release`\n - test verifies no process launch\n- Added custom output storage at `output.{node_id}`.\n- Updated docs for routing/custom schemas, `output_retries`, repair behavior, ACP limitation.\n\nValidation run:\n- `cargo nextest run -p fabro-types -p fabro-workflow graph:: file_inlining static_reference` ✅\n- `cargo nextest run -p fabro-workflow` ✅\n- `cargo +nightly-2026-04-14 fmt --check --all` ✅\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` ✅\n- `cargo insta pending-snapshots` unavailable in this environment (`cargo-insta` not installed); checked for `**/*.snap.new` and found none.",
+ "internal.work_dir": "/home/daytona/workspace/fabro",
+ "internal.retry_count.preflight_lint": 0,
+ "graph.rankdir": "LR",
+ "internal.run_id": "01KSB6GTZ00T5V6BNMXN3SPKZF",
+ "outcome": "succeeded",
+ "internal.fidelity": "compact",
+ "thread.toolchain.current_node": "preflight_compile",
+ "last_stage": "implement",
+ "graph.goal": "# Output Schema Validation Implementation Plan\n\n> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.\n\n**Goal:** Add `output_schema` validation for agent and prompt nodes, with context-preserving repair turns when structured output does not validate.\n\n**Architecture:** Introduce a small structured-output layer in `fabro-workflow` that resolves node-level schema declarations, extracts JSON output, validates it, and produces either routing side effects or a parsed custom output context update. Agent and prompt execution must perform schema repair inside the active LLM conversation instead of using the workflow executor retry path.\n\n**Tech Stack:** Rust, Graphviz workflow attrs, `serde_json`, workspace `jsonschema`, existing `fabro-llm::ResponseFormat`, Fabro agent sessions, `cargo nextest`.\n\n---\n\n## Public Interface\n\nWorkflow authors can opt in on agent and prompt nodes:\n\n```dot\nreview [\n shape=tab,\n output_schema=\"routing\",\n output_retries=2\n]\n\naudit [\n shape=tab,\n output_schema=\"@schemas/audit-result.schema.json\",\n output_retries=2\n]\n```\n\n- `output_schema=\"routing\"` uses Fabro's built-in routing directive schema.\n- `output_schema=\"@path/to/schema.json\"` loads a JSON Schema file through existing workflow file-reference rules.\n- `output_retries` controls corrective turns inside the same node execution. Default: `2`. `0` means validate once and fail without a repair turn.\n- Schema failures are terminal node failures after `output_retries` is exhausted. They are not `retry_requested` outcomes and do not consume `max_retries`.\n- `backend=\"acp\"` with `output_schema` is unsupported in v1 and returns a clear validation error.\n\n## Implementation Tasks\n\n### Task 1: Node Attributes And File Reference Resolution\n\n**Files:**\n- Modify: `lib/crates/fabro-types/src/graph.rs`\n- Modify: `lib/crates/fabro-workflow/src/static_reference.rs`\n- Modify: `lib/crates/fabro-workflow/src/transforms/file_inlining.rs`\n- Test: existing unit tests in those files\n\n- [ ] Add `Node::output_schema(&self) -> Option<&str>` next to other agent/prompt attrs.\n- [ ] Add `Node::output_retries(&self) -> i64` returning `self.int_attr(\"output_retries\").unwrap_or(2).max(0)`.\n- [ ] Teach static reference validation that node attr `output_schema` values starting with `@` are file inline references.\n- [ ] Extend file inlining so `output_schema=\"@schemas/foo.json\"` is replaced with the schema file contents before execution, while `output_schema=\"routing\"` stays unchanged.\n- [ ] Add tests for absent attrs, default retries, zero retries, file inlining, and unresolved schema reference diagnostics.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-types -p fabro-workflow graph:: file_inlining static_reference\n```\n\nExpected: targeted tests pass.\n\n### Task 2: Structured Output Module\n\n**Files:**\n- Create: `lib/crates/fabro-workflow/src/handler/structured_output.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/mod.rs`\n- Modify: `lib/crates/fabro-workflow/Cargo.toml`\n- Test: unit tests in `structured_output.rs`\n\n- [ ] Add `jsonschema.workspace = true` to `fabro-workflow` dependencies.\n- [ ] Define `OutputSchemaKind` with `Routing` and `JsonSchema { schema: serde_json::Value }`.\n- [ ] Parse `node.output_schema()` into `None`, `Routing`, or custom JSON Schema. Treat literal `routing` as the only built-in keyword.\n- [ ] Add a built-in routing schema requiring an object with at least one recognized field: `preferred_next_label`, `outcome`, `failure_reason`, `suggested_next_ids`, or `context_updates`.\n- [ ] Reuse balanced-object scanning semantics for response text: validate the last JSON object that is relevant to the selected schema.\n- [ ] Return a structured validation result containing the parsed JSON object, concise error messages, and enough information to build a repair prompt.\n- [ ] Add tests for valid routing JSON, missing routing fields, wrong routing field types, valid custom schema, invalid custom schema, invalid JSON, and no JSON object.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow structured_output\n```\n\nExpected: structured-output unit tests pass.\n\n### Task 3: Routing Extraction Compatibility\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/agent.rs`\n- Test: existing agent handler unit tests\n\n- [ ] Keep the loose default unchanged when `output_schema` is absent.\n- [ ] Move current `STATUS_FIELDS`, balanced JSON scanning, and routing-field application behind reusable functions in `structured_output.rs` or call the new module from `agent.rs`.\n- [ ] For `output_schema=\"routing\"`, require schema-valid routing JSON and surface validation failures for repair instead of silently ignoring bad candidates.\n- [ ] Preserve existing routing fallback priority for agent nodes: response text first, then `status.json`, then last file touched.\n- [ ] Keep prompt-node routing behavior response-only unless later tasks explicitly add prompt `status.json` support.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::agent\n```\n\nExpected: existing loose routing tests still pass, plus new strict routing tests pass.\n\n### Task 4: Prompt Node Same-Context Repair\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/api.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/prompt.rs`\n- Test: prompt/API backend tests in those files\n\n- [ ] In `AgentApiBackend::one_shot`, keep `messages` mutable across attempts.\n- [ ] When a prompt node has a custom JSON Schema, set `response_format=JsonSchema` on the initial and repair LLM requests. For `routing`, use `JsonObject` or no provider-native schema if provider behavior would conflict with Fabro's routing extraction.\n- [ ] After each LLM response, validate according to `output_schema`.\n- [ ] On validation failure with repair attempts remaining, append `Message::assistant(response.text())`, then append a corrective `Message::user(repair_message)`, and call `client.complete` again with the same messages.\n- [ ] On success, return the validated response text and aggregate usage across all attempts.\n- [ ] On exhaustion, return a terminal failed outcome with failure reason `output schema validation failed after N repair attempt(s)`.\n- [ ] Update `PromptHandler` so validated custom output is added to `context_updates[\"output.{node_id}\"]`; routing output still updates outcome routing fields.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::prompt handler::llm::api\n```\n\nExpected: prompt repair keeps previous assistant output in the message list and succeeds after a corrective response.\n\n### Task 5: Agent Node Same-Session Repair\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/api.rs`\n- Modify: `lib/crates/fabro-workflow/src/handler/agent.rs`\n- Test: agent/API backend tests in those files\n\n- [ ] In `AgentApiBackend::run`, validate the final assistant response before releasing, closing, or caching the session.\n- [ ] On validation failure with repair attempts remaining, call `session.process_input(repair_message)` on the same `Session`.\n- [ ] Recompute the final assistant response after each repair turn from `session.history()`.\n- [ ] Aggregate usage across all new assistant turns, including repair turns, without double-counting reused session history.\n- [ ] Do not set provider-native `response_format` for agent sessions in v1, because agent sessions may need normal tool-use messages before final output.\n- [ ] Return terminal failure after exhaustion; do not return a retryable backend error and do not request workflow node retry.\n- [ ] Update `AgentHandler` to apply validated routing/custom output to the final `Outcome`.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::agent handler::llm::api\n```\n\nExpected: agent repair sends a second `process_input` to the same session and final validated output drives outcome/context updates.\n\n### Task 6: ACP Guardrail\n\n**Files:**\n- Modify: `lib/crates/fabro-workflow/src/handler/llm/acp.rs`\n- Test: ACP backend tests in that file\n\n- [ ] At the start of `AgentAcpBackend::run`, reject nodes where `node.output_schema().is_some()`.\n- [ ] Use a clear error message: `output_schema is not supported with backend=\"acp\" in this release`.\n- [ ] Add a test proving the ACP backend does not launch a process when `output_schema` is present.\n\nRun:\n\n```bash\ncargo nextest run -p fabro-workflow handler::llm::acp\n```\n\nExpected: ACP guardrail test passes.\n\n### Task 7: Docs\n\n**Files:**\n- Modify: `docs/public/agents/outputs.mdx`\n- Modify: `docs/public/reference/dot-language.mdx`\n\n- [ ] Document `output_schema=\"routing\"` and `output_schema=\"@schema.json\"` under routing/structured outputs.\n- [ ] Document same-context repair behavior explicitly: Fabro sends validation feedback to the same agent/prompt context before failing.\n- [ ] Document `output_retries`, default `2`, and distinction from `max_retries`.\n- [ ] Document v1 scope: agent/prompt nodes only; ACP unsupported; custom schema output stored at `output.{node_id}`.\n\nRun:\n\n```bash\nrg -n \"output_schema|output_retries|output\\\\.\" docs/public/agents/outputs.mdx docs/public/reference/dot-language.mdx\n```\n\nExpected: docs mention the new attrs and storage behavior.\n\n### Task 8: Full Verification\n\n**Files:**\n- No new files beyond prior tasks\n\n- [ ] Run focused workflow tests:\n\n```bash\ncargo nextest run -p fabro-workflow\n```\n\n- [ ] Run formatting check:\n\n```bash\ncargo +nightly-2026-04-14 fmt --check --all\n```\n\n- [ ] Run clippy:\n\n```bash\ncargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings\n```\n\n- [ ] If snapshots change, inspect before accepting:\n\n```bash\ncargo insta pending-snapshots\n```\n\nOnly run `cargo insta accept` after verifying every pending snapshot is expected.\n\n## Acceptance Criteria\n\n- Existing workflows without `output_schema` behave exactly as before.\n- `output_schema=\"routing\"` prevents malformed/missing routing JSON from silently falling through to normal edge selection.\n- Invalid structured output results in a corrective LLM turn in the same context window.\n- Prompt repair preserves previous assistant output in the message list.\n- Agent repair preserves the same live session and does not re-run the node from scratch.\n- Exhausted output repair attempts produce a clear terminal failure.\n- Custom schema output is available to downstream nodes at `output.{node_id}`.\n- Docs clearly distinguish `output_retries` from `max_retries`.\n\n## Assumptions\n\n- `output_retries=2` is the default.\n- Custom schema validation targets the final JSON object in the response text.\n- `status.json` fallback remains routing-specific.\n- Provider-native response schema is used for prompt nodes only where it is safe.\n- ACP support can be added later after there is a guaranteed context-preserving repair mechanism.\n"
+ },
+ "node_outcomes": {
+ "toolchain": {
+ "status": "succeeded",
+ "context_updates": {
+ "command.output": "blob://sha256/fc14b2ba2d770e5cd3169df7a29525c962adfc4cfa3097b9098c63ebd61a748c"
+ },
+ "notes": "Script completed: command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1",
+ "usage": null
+ },
+ "implement": {
+ "status": "succeeded",
+ "context_updates": {
+ "last_response": "Implemented output schema validation end-to-end.\n\nSummary:\n- Added `output_schema` / `output_retries` node attrs.\n- Added `output_schema` static/file reference handling and inlining.\n- Added `handler:",
+ "response.implement": "Implemented output schema validation end-to-end.\n\nSummary:\n- Added `output_schema` / `output_retries` node attrs.\n- Added `output_schema` static/file reference handling and inlining.\n- Added `handler::structured_output` for:\n - routing/custom schema parsing\n - balanced JSON extraction\n - JSON Schema validation\n - repair prompts\n - routing/custom context application\n- Integrated strict routing/custom output into agent and prompt handlers.\n- Added same-context repair for:\n - prompt API calls via mutable message history\n - agent API sessions via the same live `Session`\n- Added ACP guardrail:\n - `output_schema is not supported with backend=\"acp\" in this release`\n - test verifies no process launch\n- Added custom output storage at `output.{node_id}`.\n- Updated docs for routing/custom schemas, `output_retries`, repair behavior, ACP limitation.\n\nValidation run:\n- `cargo nextest run -p fabro-types -p fabro-workflow graph:: file_inlining static_reference` ✅\n- `cargo nextest run -p fabro-workflow` ✅\n- `cargo +nightly-2026-04-14 fmt --check --all` ✅\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` ✅\n- `cargo insta pending-snapshots` unavailable in this environment (`cargo-insta` not installed); checked for `**/*.snap.new` and found none.",
+ "last_stage": "implement"
+ },
+ "notes": "Stage completed: implement",
+ "usage": {
+ "input": {
+ "usage": {
+ "model": {
+ "provider": "openai",
+ "model_id": "gpt-5.5"
+ },
+ "tokens": {
+ "input_tokens": 329828,
+ "output_tokens": 33087,
+ "reasoning_tokens": 21712,
+ "cache_read_tokens": 29338112,
+ "cache_write_tokens": 0
+ }
+ },
+ "facts": {
+ "algorithm": "openai"
+ }
+ },
+ "total_usd_micros": 17962166
+ },
+ "files_touched": [
+ "/home/daytona/workspace/fabro/lib/crates/fabro-workflow/src/handler/structured_output.rs"
+ ]
+ },
+ "start": {
+ "status": "succeeded",
+ "usage": null
+ },
+ "preflight_compile": {
+ "status": "succeeded",
+ "context_updates": {
+ "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126"
+ },
+ "notes": "Script completed: cargo check -q --workspace 2>&1",
+ "usage": null
+ },
+ "preflight_lint": {
+ "status": "succeeded",
+ "context_updates": {
+ "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126"
+ },
+ "notes": "Script completed: cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1",
+ "usage": null
+ }
+ },
+ "next_node_id": "simplify_opus",
+ "git_commit_sha": "43b8310e6051d435c815d935333e893833fb72a9",
+ "node_visits": {
+ "preflight_lint": 1,
+ "preflight_compile": 1,
+ "implement": 1,
+ "start": 1,
+ "toolchain": 1
+ }
+ },
+ "diff": {
+ "patch": "diff --git a/Cargo.lock b/Cargo.lock\nindex c50b6c665..28b5f81e8 100644\n--- a/Cargo.lock\n+++ b/Cargo.lock\n@@ -2603,6 +2603,7 @@ dependencies = [\n \"git2\",\n \"hex\",\n \"httpmock\",\n+ \"jsonschema\",\n \"md5\",\n \"miette\",\n \"mime_guess\",\ndiff --git a/docs/public/agents/outputs.mdx b/docs/public/agents/outputs.mdx\nindex 1407f4db6..7467f050b 100644\n--- a/docs/public/agents/outputs.mdx\n+++ b/docs/public/agents/outputs.mdx\n@@ -64,6 +64,23 @@ The test coverage is below the threshold.\n \n JSON objects without recognized fields are ignored.\n \n+### Validated routing output\n+\n+Set `output_schema=\"routing\"` on an agent or prompt node to require Fabro's built-in routing directive schema:\n+\n+```dot\n+review [\n+ shape=tab,\n+ prompt=\"Review the implementation and return routing JSON.\",\n+ output_schema=\"routing\",\n+ output_retries=2\n+]\n+```\n+\n+With `output_schema=\"routing\"`, the routing JSON must be an object with at least one recognized routing field (`preferred_next_label`, `outcome`, `failure_reason`, `suggested_next_ids`, or `context_updates`) and those fields must have the expected types. Malformed routing JSON fails validation instead of being silently ignored.\n+\n+Fabro repairs invalid structured output inside the same LLM context before failing the node. For prompt nodes, Fabro appends the invalid assistant response and a corrective user message to the same message list. For agent nodes using the API backend, Fabro sends the corrective message to the same live agent session. `output_retries` controls these repair turns and defaults to `2`; `output_retries=0` validates once and fails without a repair turn. These repair turns are separate from workflow `max_retries` and do not consume node retry attempts.\n+\n ### Fallback: status.json file\n \n If no routing directives are found in the response text, Fabro checks whether the agent wrote a `status.json` file into the sandbox working directory. If the file exists, Fabro extracts routing directives from it using the same logic. This is useful for agents that write structured output to files rather than including JSON in their response text.\n@@ -90,6 +107,35 @@ review -> fix [label=\"Fix\"]\n review -> approve [label=\"Approve\"]\n ```\n \n+## Custom structured outputs\n+\n+Agent and prompt nodes can also validate their final JSON object against a JSON Schema file:\n+\n+```dot\n+audit [\n+ shape=tab,\n+ prompt=\"Audit the change and return JSON that matches the schema.\",\n+ output_schema=\"@schemas/audit-result.schema.json\",\n+ output_retries=2\n+]\n+```\n+\n+`output_schema=\"@path/to/schema.json\"` uses the same workflow file-reference rules as prompt files: the schema is loaded relative to the workflow file and inlined before execution. The final JSON object in the LLM response is validated with `jsonschema`.\n+\n+When custom schema validation succeeds, Fabro stores the parsed JSON value in context at:\n+\n+| Key | Value |\n+|---|---|\n+| `output.{node_id}` | The parsed JSON object that matched the custom schema |\n+\n+For example, node `audit` writes its parsed custom output to `output.audit`. Fabro still stores the raw response text at `response.audit`.\n+\n+If custom schema validation fails, Fabro sends concise validation feedback to the same prompt conversation or agent session and asks for corrected JSON. After `output_retries` repair turns are exhausted, the node fails terminally with `output schema validation failed after N repair attempt(s)`.\n+\n+\n+Structured output validation currently applies to agent and prompt nodes. `backend=\"acp\"` does not support `output_schema` in this release. Custom schemas update `output.{node_id}`; routing schemas update routing fields and `context_updates` instead.\n+\n+\n ## Output logging\n \n Fabro writes several files per stage to `stages/{rank:03}-{node_id}@{visit}/` in metadata snapshots and `fabro dump` output:\ndiff --git a/docs/public/reference/dot-language.mdx b/docs/public/reference/dot-language.mdx\nindex 1c1656057..74a366632 100644\n--- a/docs/public/reference/dot-language.mdx\n+++ b/docs/public/reference/dot-language.mdx\n@@ -206,10 +206,37 @@ Start nodes can also be identified by ID (`start` or `Start`). Exit nodes can be\n | `model` | String | Explicit model ID (overrides stylesheet) |\n | `provider` | String | Explicit provider name (overrides stylesheet). Auto-inferred from the model catalog when omitted. |\n | `project_memory` | Boolean | When `true` (default), prompt nodes discover and include project docs (`AGENTS.md`, `CLAUDE.md`, etc.) as a system prompt. Set to `false` to disable. |\n+| `output_schema` | String | Optional structured output validation. Use `routing` for Fabro's built-in routing directive schema, or `@path/to/schema.json` for a JSON Schema file. Supported on agent and prompt nodes. |\n+| `output_retries` | Integer | Corrective structured-output turns inside the same prompt conversation or agent session. Default `2`; `0` validates once and fails without repair. Separate from `max_retries`. |\n | `backend` | String | Agent execution backend: `api` (default) or `acp`. `api` runs Fabro's tool loop through provider APIs; `acp` runs an Agent Client Protocol stdio agent inside the active sandbox. Prompt nodes are API-only. See [Agents — Backends](/core-concepts/agents#backends). |\n | `acp.command` | String | Shell command for nodes with `backend=\"acp\"`. Mutually exclusive with `acp.config`. The value is always parsed as a command string, not JSON. |\n | `acp.config` | String | JSON stdio ACP config for nodes with `backend=\"acp\"`. Mutually exclusive with `acp.command`. |\n \n+#### Structured output validation\n+\n+`output_schema` opts an agent or prompt node into strict JSON validation:\n+\n+```dot\n+review [\n+ shape=tab,\n+ output_schema=\"routing\",\n+ output_retries=2\n+]\n+\n+audit [\n+ shape=tab,\n+ output_schema=\"@schemas/audit-result.schema.json\",\n+ output_retries=2\n+]\n+```\n+\n+- `output_schema=\"routing\"` requires a JSON object with at least one recognized routing field: `preferred_next_label`, `outcome`, `failure_reason`, `suggested_next_ids`, or `context_updates`.\n+- `output_schema=\"@schemas/audit-result.schema.json\"` loads a JSON Schema file using workflow file-reference rules and validates the final JSON object in the response text.\n+- On validation failure, Fabro sends validation feedback to the same active context before failing: prompt nodes keep the prior assistant response in the message list, and API-backed agent nodes repair in the same live session.\n+- `output_retries` defaults to `2` and controls only these corrective structured-output turns. It is not the same as `max_retries` and does not consume workflow retry attempts.\n+- Custom schema output is stored in context at `output.{node_id}`. Routing schema output updates routing fields and any `context_updates`.\n+- `backend=\"acp\"` with `output_schema` is unsupported in this release.\n+\n ### Command nodes\n \n | Attribute | Type | Description |\n@@ -331,15 +358,16 @@ gate -> v2 [condition=\"context.version matches ^v2\\\\.\"]\n gate -> slow_path\n ```\n \n-## Prompt file references\n+## Prompt and schema file references\n \n-Instead of inlining long prompts, reference an external file:\n+Instead of inlining long prompts or JSON Schemas, reference an external file:\n \n ```dot\n simplify [label=\"Simplify\", prompt=\"@prompts/simplify.md\"]\n+audit [shape=tab, output_schema=\"@schemas/audit-result.schema.json\"]\n ```\n \n-The `@` prefix tells Fabro to load the prompt from a file path relative to the workflow file. Paths support `~` (home directory) and `..` (parent directory):\n+The `@` prefix tells Fabro to load the referenced file relative to the workflow file. Paths support `~` (home directory) and `..` (parent directory):\n \n ```dot\n shared [prompt=\"@~/shared-prompts/review.md\"]\ndiff --git a/lib/crates/fabro-types/src/graph.rs b/lib/crates/fabro-types/src/graph.rs\nindex 27344f8c9..7ef9bad99 100644\n--- a/lib/crates/fabro-types/src/graph.rs\n+++ b/lib/crates/fabro-types/src/graph.rs\n@@ -168,6 +168,16 @@ impl Node {\n self.str_attr(\"prompt\")\n }\n \n+ #[must_use]\n+ pub fn output_schema(&self) -> Option<&str> {\n+ self.str_attr(\"output_schema\")\n+ }\n+\n+ #[must_use]\n+ pub fn output_retries(&self) -> i64 {\n+ self.int_attr(\"output_retries\").unwrap_or(2).max(0)\n+ }\n+\n #[must_use]\n pub fn max_retries(&self) -> Option {\n self.int_attr(\"max_retries\")\n@@ -576,6 +586,8 @@ mod tests {\n assert_eq!(node.shape(), \"box\");\n assert_eq!(node.node_type(), None);\n assert_eq!(node.prompt(), None);\n+ assert_eq!(node.output_schema(), None);\n+ assert_eq!(node.output_retries(), 2);\n assert_eq!(node.max_retries(), None);\n assert!(!node.goal_gate());\n assert_eq!(node.retry_target(), None);\n@@ -602,6 +614,31 @@ mod tests {\n assert!(!node.project_memory());\n }\n \n+ #[test]\n+ fn node_output_retries_defaults_and_clamps_to_zero() {\n+ let mut node = Node::new(\"x\");\n+ assert_eq!(node.output_retries(), 2);\n+\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(0));\n+ assert_eq!(node.output_retries(), 0);\n+\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(-3));\n+ assert_eq!(node.output_retries(), 0);\n+ }\n+\n+ #[test]\n+ fn node_output_schema_returns_string_attr() {\n+ let mut node = Node::new(\"x\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\"routing\".to_string()),\n+ );\n+\n+ assert_eq!(node.output_schema(), Some(\"routing\"));\n+ }\n+\n #[test]\n fn node_with_attrs() {\n let mut node = Node::new(\"plan\");\ndiff --git a/lib/crates/fabro-workflow/Cargo.toml b/lib/crates/fabro-workflow/Cargo.toml\nindex d5025f1bd..6021bca4c 100644\n--- a/lib/crates/fabro-workflow/Cargo.toml\n+++ b/lib/crates/fabro-workflow/Cargo.toml\n@@ -46,6 +46,7 @@ fabro-http.workspace = true\n thiserror.workspace = true\n serde.workspace = true\n serde_json.workspace = true\n+jsonschema.workspace = true\n tokio.workspace = true\n bytes.workspace = true\n object_store.workspace = true\ndiff --git a/lib/crates/fabro-workflow/src/error.rs b/lib/crates/fabro-workflow/src/error.rs\nindex 6fcd565e3..4373a16db 100644\n--- a/lib/crates/fabro-workflow/src/error.rs\n+++ b/lib/crates/fabro-workflow/src/error.rs\n@@ -304,6 +304,9 @@ pub enum Error {\n #[error(\"Unsupported operation: {0}\")]\n Unsupported(String),\n \n+ #[error(\"{0}\")]\n+ OutputSchemaValidation(String),\n+\n #[error(\"Pipeline cancelled\")]\n Cancelled,\n }\n@@ -428,8 +431,8 @@ impl Error {\n /// Retryable: Handler (transient handler failures), Engine (could be\n /// transient), Io (network/disk issues are often transient),\n /// Llm (delegates to SdkError). Terminal: Parse, Validation,\n- /// Stylesheet (configuration errors), Checkpoint (storage\n- /// integrity), Cancelled (explicit cancellation).\n+ /// OutputSchemaValidation, Stylesheet (configuration errors), Checkpoint\n+ /// (storage integrity), Cancelled (explicit cancellation).\n #[must_use]\n pub fn is_retryable(&self) -> bool {\n match self {\n@@ -444,6 +447,7 @@ impl Error {\n | Self::Precondition(_)\n | Self::RunNotFound(_)\n | Self::Unsupported(_)\n+ | Self::OutputSchemaValidation(_)\n | Self::Cancelled => false,\n }\n }\n@@ -461,7 +465,8 @@ impl Error {\n | Self::Template { .. }\n | Self::Stylesheet(_)\n | Self::Checkpoint(_)\n- | Self::Unsupported(_) => FailureCategory::Deterministic,\n+ | Self::Unsupported(_)\n+ | Self::OutputSchemaValidation(_) => FailureCategory::Deterministic,\n Self::Precondition(_) | Self::RunNotFound(_) => FailureCategory::Structural,\n Self::Handler { failure_class, .. } | Self::Engine { failure_class, .. } => {\n *failure_class\ndiff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs\nindex 0788ee1e5..0f2118df8 100644\n--- a/lib/crates/fabro-workflow/src/handler/agent.rs\n+++ b/lib/crates/fabro-workflow/src/handler/agent.rs\n@@ -2,20 +2,21 @@ use std::path::Path;\n use std::sync::Arc;\n \n use async_trait::async_trait;\n-use fabro_agent::Sandbox;\n+use fabro_agent::{Sandbox, shell_quote};\n use fabro_graphviz::graph::{Graph, Node};\n use fabro_types::{RunId, StageModelUsage};\n use tokio_util::sync::CancellationToken;\n \n use super::llm::api::EffectiveRequestControls;\n+use super::structured_output::{\n+ self, OutputSchemaKind, StructuredOutputError, ValidatedStructuredOutput,\n+};\n use super::{EngineServices, Handler, NodeTimeoutPolicy};\n use crate::context::{Context, WorkflowContext, keys};\n use crate::error::Error;\n use crate::event::{Emitter, Event, StageScope};\n use crate::interview_runtime::WorkflowAgentQuestionRuntime;\n-use crate::outcome::{\n- BilledModelUsage, FailureCategory, FailureDetail, Outcome, OutcomeExt, StageOutcome,\n-};\n+use crate::outcome::{BilledModelUsage, Outcome, OutcomeExt};\n \n /// Result from a `CodergenBackend` invocation.\n pub enum CodergenResult {\n@@ -132,110 +133,65 @@ impl AgentHandler {\n }\n }\n \n-/// Status fields that indicate a JSON object contains routing directives.\n-const STATUS_FIELDS: &[&str] = &[\n- \"preferred_next_label\",\n- \"outcome\",\n- \"failure_reason\",\n- \"suggested_next_ids\",\n- \"context_updates\",\n-];\n-\n-/// Find all balanced `{...}` JSON object substrings in the text.\n-fn find_json_objects(text: &str) -> Vec<&str> {\n- let mut results = Vec::new();\n- let bytes = text.as_bytes();\n- let mut i = 0;\n- while i < bytes.len() {\n- if bytes[i] == b'{' {\n- let start = i;\n- let mut depth = 0;\n- let mut in_string = false;\n- let mut escape = false;\n- let mut j = i;\n- while j < bytes.len() {\n- let c = bytes[j];\n- if escape {\n- escape = false;\n- } else if c == b'\\\\' && in_string {\n- escape = true;\n- } else if c == b'\"' {\n- in_string = !in_string;\n- } else if !in_string {\n- if c == b'{' {\n- depth += 1;\n- } else if c == b'}' {\n- depth -= 1;\n- if depth == 0 {\n- results.push(&text[start..=j]);\n- break;\n- }\n- }\n- }\n- j += 1;\n- }\n- }\n- i += 1;\n- }\n- results\n-}\n-\n /// Extract routing directives from LLM response text.\n ///\n /// Searches for the last JSON object in the response that contains at least\n /// one status field (`preferred_next_label`, `outcome`, `suggested_next_ids`,\n /// `context_updates`). Merges extracted fields into the outcome.\n pub(crate) fn extract_status_fields(text: &str, outcome: &mut Outcome) -> bool {\n- let candidates = find_json_objects(text);\n-\n- let parsed = candidates.iter().rev().find_map(|candidate| {\n- let value: serde_json::Value = serde_json::from_str(candidate).ok()?;\n- if let Some(obj) = value.as_object() {\n- if STATUS_FIELDS.iter().any(|f| obj.contains_key(*f)) {\n- return Some(value);\n- }\n- }\n- None\n- });\n-\n- let Some(value) = parsed else { return false };\n- let Some(obj) = value.as_object() else {\n- return false;\n- };\n+ structured_output::extract_status_fields_loose(text, outcome)\n+}\n \n- if let Some(label) = obj.get(\"preferred_next_label\").and_then(|v| v.as_str()) {\n- outcome.preferred_label = Some(label.to_string());\n+pub(crate) async fn validate_agent_output_sources(\n+ schema: &OutputSchemaKind,\n+ response_text: &str,\n+ sandbox: &Arc,\n+ last_file_touched: Option<&str>,\n+) -> Result {\n+ if !matches!(schema, OutputSchemaKind::Routing) {\n+ return structured_output::validate_response_text(schema, response_text);\n }\n \n- if let Some(ids) = obj.get(\"suggested_next_ids\").and_then(|v| v.as_array()) {\n- let string_ids: Vec = ids\n- .iter()\n- .filter_map(|v| v.as_str().map(String::from))\n- .collect();\n- if !string_ids.is_empty() {\n- outcome.suggested_next_ids = string_ids;\n- }\n+ match structured_output::validate_response_text(schema, response_text) {\n+ Ok(validated) => return Ok(validated),\n+ Err(error) if error.allows_routing_fallback() => {}\n+ Err(error) => return Err(error),\n }\n \n- if let Some(status_str) = obj.get(\"outcome\").and_then(|v| v.as_str()) {\n- if let Ok(status) = status_str.parse::() {\n- outcome.status = status;\n- if outcome.status.is_failure() {\n- if let Some(reason) = obj.get(\"failure_reason\").and_then(|v| v.as_str()) {\n- outcome.failure =\n- Some(FailureDetail::new(reason, FailureCategory::Deterministic));\n- }\n+ let mut fallback_error = None;\n+ if let Some(status_json) = read_sandbox_file(sandbox, \"status.json\").await {\n+ match structured_output::validate_response_text(schema, &status_json) {\n+ Ok(validated) => return Ok(validated),\n+ Err(error) if error.allows_routing_fallback() => {\n+ fallback_error = Some(error);\n }\n+ Err(error) => return Err(error),\n }\n }\n \n- if let Some(updates) = obj.get(\"context_updates\").and_then(|v| v.as_object()) {\n- for (key, val) in updates {\n- outcome.context_updates.insert(key.clone(), val.clone());\n+ if let Some(path) = last_file_touched {\n+ if let Some(contents) = read_sandbox_file(sandbox, path).await {\n+ return structured_output::validate_response_text(schema, &contents);\n }\n }\n \n- true\n+ Err(fallback_error.unwrap_or_else(|| {\n+ structured_output::validate_response_text(schema, response_text)\n+ .expect_err(\"response text should have failed routing validation\")\n+ }))\n+}\n+\n+async fn read_sandbox_file(sandbox: &Arc, path: &str) -> Option {\n+ let cmd = format!(\"cat {}\", shell_quote(path));\n+ let result = sandbox\n+ .exec_command(&cmd, 5_000, None, None, None)\n+ .await\n+ .ok()?;\n+ if result.is_success() {\n+ Some(result.stdout)\n+ } else {\n+ None\n+ }\n }\n \n /// Truncate a string to at most `max_chars` characters (char-boundary safe).\n@@ -416,34 +372,46 @@ impl Handler for AgentHandler {\n serde_json::json!(&response_text),\n );\n \n- // 7b. Parse routing directives from response text, falling back to\n- // status.json written by the agent into the sandbox CWD, then to\n- // the last file the agent wrote.\n- let found_in_response = extract_status_fields(&response_text, &mut outcome);\n- if !found_in_response {\n- let mut found_in_status_json = false;\n- if let Ok(result) = services\n- .run\n- .sandbox\n- .exec_command(\"cat status.json\", 5_000, None, None, None)\n- .await\n+ if let Some(schema) = structured_output::parse_node_output_schema(node)? {\n+ match validate_agent_output_sources(\n+ &schema,\n+ &response_text,\n+ &services.run.sandbox,\n+ last_file_touched.as_deref(),\n+ )\n+ .await\n {\n- if result.is_success() {\n- found_in_status_json = extract_status_fields(&result.stdout, &mut outcome);\n+ Ok(validated) => {\n+ structured_output::apply_validated_output(\n+ node,\n+ &schema,\n+ &validated,\n+ &mut outcome,\n+ );\n+ }\n+ Err(_) => {\n+ return Ok(structured_output::exhausted_failure_outcome(\n+ node.output_retries(),\n+ ));\n }\n }\n- if !found_in_status_json {\n- if let Some(ref path) = last_file_touched {\n- let quoted = shlex::try_quote(path).unwrap_or_else(|_| path.into());\n- let cmd = format!(\"cat {quoted}\");\n- if let Ok(result) = services\n- .run\n- .sandbox\n- .exec_command(&cmd, 5_000, None, None, None)\n- .await\n- {\n- if result.is_success() {\n- extract_status_fields(&result.stdout, &mut outcome);\n+ } else {\n+ // 7b. Parse routing directives from response text, falling back to\n+ // status.json written by the agent into the sandbox CWD, then to\n+ // the last file the agent wrote.\n+ let found_in_response = extract_status_fields(&response_text, &mut outcome);\n+ if !found_in_response {\n+ let mut found_in_status_json = false;\n+ if let Some(status_json) =\n+ read_sandbox_file(&services.run.sandbox, \"status.json\").await\n+ {\n+ found_in_status_json = extract_status_fields(&status_json, &mut outcome);\n+ }\n+ if !found_in_status_json {\n+ if let Some(ref path) = last_file_touched {\n+ if let Some(contents) = read_sandbox_file(&services.run.sandbox, path).await\n+ {\n+ extract_status_fields(&contents, &mut outcome);\n }\n }\n }\n@@ -795,6 +763,142 @@ mod tests {\n );\n }\n \n+ #[tokio::test]\n+ async fn codergen_handler_output_schema_routing_uses_status_json_fallback_when_response_has_no_json()\n+ {\n+ let sandbox_dir = TempDir::new().unwrap();\n+ std::fs::write(\n+ sandbox_dir.path().join(\"status.json\"),\n+ r#\"{\"preferred_next_label\": \"review\"}\"#,\n+ )\n+ .unwrap();\n+\n+ let handler = AgentHandler::new(None);\n+ let mut node = Node::new(\"step\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\"routing\".to_string()),\n+ );\n+ let context = test_context();\n+ let graph = Graph::new(\"test\");\n+ let tmp = TempDir::new().unwrap();\n+\n+ let mut services = EngineServices::test_default();\n+ services.run =\n+ services\n+ .run\n+ .with_sandbox(std::sync::Arc::new(fabro_agent::LocalSandbox::new(\n+ sandbox_dir.path().to_path_buf(),\n+ )));\n+\n+ let outcome = handler\n+ .execute(&node, &context, &graph, tmp.path(), &services)\n+ .await\n+ .unwrap();\n+\n+ assert_eq!(outcome.status, crate::outcome::StageOutcome::Succeeded);\n+ assert_eq!(outcome.preferred_label.as_deref(), Some(\"review\"));\n+ }\n+\n+ #[tokio::test]\n+ async fn codergen_handler_output_schema_routing_rejects_malformed_response_before_status_json_fallback()\n+ {\n+ struct BadRoutingBackend;\n+\n+ #[async_trait]\n+ impl CodergenBackend for BadRoutingBackend {\n+ async fn run(&self, _request: CodergenRunRequest<'_>) -> Result {\n+ Ok(CodergenResult::Text {\n+ text: r#\"{\"suggested_next_ids\": [1]}\"#.to_string(),\n+ usage: None,\n+ files_touched: Vec::new(),\n+ last_file_touched: None,\n+ })\n+ }\n+ }\n+\n+ let sandbox_dir = TempDir::new().unwrap();\n+ std::fs::write(\n+ sandbox_dir.path().join(\"status.json\"),\n+ r#\"{\"preferred_next_label\": \"should_not_use\"}\"#,\n+ )\n+ .unwrap();\n+\n+ let handler = AgentHandler::new(Some(Box::new(BadRoutingBackend)));\n+ let mut node = Node::new(\"step\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\"routing\".to_string()),\n+ );\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(0));\n+ let context = test_context();\n+ let graph = Graph::new(\"test\");\n+ let tmp = TempDir::new().unwrap();\n+\n+ let mut services = EngineServices::test_default();\n+ services.run =\n+ services\n+ .run\n+ .with_sandbox(std::sync::Arc::new(fabro_agent::LocalSandbox::new(\n+ sandbox_dir.path().to_path_buf(),\n+ )));\n+\n+ let outcome = handler\n+ .execute(&node, &context, &graph, tmp.path(), &services)\n+ .await\n+ .unwrap();\n+\n+ assert_eq!(outcome.status, crate::outcome::StageOutcome::Failed {\n+ retry_requested: false,\n+ });\n+ assert_eq!(\n+ outcome.failure_reason(),\n+ Some(\"output schema validation failed after 0 repair attempt(s)\")\n+ );\n+ assert!(outcome.preferred_label.is_none());\n+ }\n+\n+ #[tokio::test]\n+ async fn codergen_handler_custom_output_schema_updates_output_context_key() {\n+ struct CustomOutputBackend;\n+\n+ #[async_trait]\n+ impl CodergenBackend for CustomOutputBackend {\n+ async fn run(&self, _request: CodergenRunRequest<'_>) -> Result {\n+ Ok(CodergenResult::Text {\n+ text: r#\"{\"passed\": true}\"#.to_string(),\n+ usage: None,\n+ files_touched: Vec::new(),\n+ last_file_touched: None,\n+ })\n+ }\n+ }\n+\n+ let handler = AgentHandler::new(Some(Box::new(CustomOutputBackend)));\n+ let mut node = Node::new(\"audit\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\n+ r#\"{\"type\":\"object\",\"required\":[\"passed\"],\"properties\":{\"passed\":{\"type\":\"boolean\"}}}\"#\n+ .to_string(),\n+ ),\n+ );\n+ let context = test_context();\n+ let graph = Graph::new(\"test\");\n+ let tmp = TempDir::new().unwrap();\n+\n+ let outcome = handler\n+ .execute(&node, &context, &graph, tmp.path(), &make_services())\n+ .await\n+ .unwrap();\n+\n+ assert_eq!(\n+ outcome.context_updates.get(\"output.audit\"),\n+ Some(&serde_json::json!({\"passed\": true})),\n+ );\n+ }\n+\n #[tokio::test]\n async fn codergen_handler_projects_provider_used_from_agent_session_events() {\n struct ProviderEventBackend;\ndiff --git a/lib/crates/fabro-workflow/src/handler/llm/acp.rs b/lib/crates/fabro-workflow/src/handler/llm/acp.rs\nindex 827f09e8a..3318c64ef 100644\n--- a/lib/crates/fabro-workflow/src/handler/llm/acp.rs\n+++ b/lib/crates/fabro-workflow/src/handler/llm/acp.rs\n@@ -325,6 +325,11 @@ impl Default for AgentAcpBackend {\n #[async_trait]\n impl CodergenBackend for AgentAcpBackend {\n async fn run(&self, request: CodergenRunRequest<'_>) -> Result {\n+ if request.node.output_schema().is_some() {\n+ return Err(Error::Validation(\n+ \"output_schema is not supported with backend=\\\"acp\\\" in this release\".to_string(),\n+ ));\n+ }\n let stage_scope = StageScope::for_handler(request.context, &request.node.id);\n self.run_turn(\n request.node,\n@@ -476,6 +481,56 @@ mod tests {\n assert_eq!(files_touched, vec![\"hello.txt\"]);\n }\n \n+ #[tokio::test]\n+ async fn acp_backend_rejects_output_schema_without_launching_process() {\n+ let tempdir = tempfile::tempdir().unwrap();\n+ let launched_path = tempdir.path().join(\"launched\");\n+\n+ let mut node = Node::new(\"work\");\n+ node.attrs\n+ .insert(\"backend\".to_string(), AttrValue::String(\"acp\".to_string()));\n+ node.attrs.insert(\n+ \"acp.command\".to_string(),\n+ AttrValue::String(\"sh -c 'touch launched'\".to_string()),\n+ );\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\"routing\".to_string()),\n+ );\n+\n+ let backend = AgentAcpBackend::new();\n+ let sandbox: Arc = Arc::new(LocalSandbox::new(tempdir.path().to_path_buf()));\n+ let emitter = Arc::new(Emitter::default());\n+ let context = Context::new();\n+ let result = backend\n+ .run(CodergenRunRequest {\n+ node: &node,\n+ prompt: \"write hello\",\n+ context: &context,\n+ thread_id: None,\n+ emitter: &emitter,\n+ sandbox: &sandbox,\n+ tool_hooks: None,\n+ cancel_token: CancellationToken::new(),\n+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),\n+ })\n+ .await;\n+\n+ let Err(error) = result else {\n+ panic!(\"expected output_schema guardrail error\");\n+ };\n+ assert!(\n+ error\n+ .to_string()\n+ .contains(\"output_schema is not supported with backend=\\\"acp\\\" in this release\"),\n+ \"unexpected error: {error}\",\n+ );\n+ assert!(\n+ !launched_path.exists(),\n+ \"ACP process should not launch when output_schema is present\",\n+ );\n+ }\n+\n #[tokio::test]\n async fn acp_backend_accepts_steer_and_incorporates_followup_result() {\n let tempdir = tempfile::tempdir().unwrap();\ndiff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs\nindex 8da02640d..22d08dc4b 100644\n--- a/lib/crates/fabro-workflow/src/handler/llm/api.rs\n+++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs\n@@ -13,7 +13,8 @@ use fabro_auth::{CredentialSource, EnvCredentialSource};\n use fabro_graphviz::graph::{AttrValue, Node};\n use fabro_llm::client::Client;\n use fabro_llm::types::{\n- Message, ReasoningEffort, Request, Speed, TokenCounts, ToolDefinition as LlmToolDefinition,\n+ Message, ReasoningEffort, Request, Response, Speed, TokenCounts,\n+ ToolDefinition as LlmToolDefinition,\n };\n use fabro_mcp::config::McpServerSettings;\n #[cfg(test)]\n@@ -26,7 +27,11 @@ use tokio::sync::Mutex as TokioMutex;\n use tokio::task::JoinHandle;\n use tokio_util::sync::CancellationToken;\n \n-use super::super::agent::{CodergenBackend, CodergenResult, CodergenRunRequest, OneShotRequest};\n+use super::super::agent::{\n+ CodergenBackend, CodergenResult, CodergenRunRequest, OneShotRequest,\n+ validate_agent_output_sources,\n+};\n+use super::super::structured_output;\n use super::activation_lease::{ActivationLease, ActivationLeaseOptions};\n use super::routing;\n use super::routing::ProviderContext;\n@@ -460,6 +465,32 @@ fn track_file_event(event: &AgentEvent, state: &mut FileTracking) {\n }\n }\n \n+fn file_tracking_snapshot(\n+ file_tracking: &Arc>,\n+) -> (Vec, Option) {\n+ let state = file_tracking.lock().unwrap();\n+ let mut files: Vec = state.touched.iter().cloned().collect();\n+ files.sort();\n+ (files, state.last.clone())\n+}\n+\n+fn last_assistant_response(session: &Session) -> String {\n+ session\n+ .history()\n+ .turns()\n+ .iter()\n+ .rev()\n+ .find_map(|turn| {\n+ if let AgentMessage::Assistant { content, .. } = turn {\n+ if !content.is_empty() {\n+ return Some(content.clone());\n+ }\n+ }\n+ None\n+ })\n+ .unwrap_or_default()\n+}\n+\n /// Spawn a task that subscribes to session events and:\n /// 1. Tracks file changes (write_file/edit_file tool calls) into shared state.\n /// 2. Forwards non-streaming agent events to the pipeline emitter.\n@@ -521,6 +552,13 @@ pub struct AgentApiBackend {\n fabro_run_tools: Option,\n }\n \n+struct OneShotCompletion {\n+ response: Response,\n+ actual_model: String,\n+ actual_provider: String,\n+ actual_speed: Option,\n+}\n+\n impl AgentApiBackend {\n #[must_use]\n pub fn new(\n@@ -812,71 +850,18 @@ impl AgentApiBackend {\n }\n }\n }\n-}\n-\n-#[async_trait]\n-impl CodergenBackend for AgentApiBackend {\n- async fn shutdown(&self, emitter: &Arc) {\n- self.shutdown_cached_sessions(emitter);\n- }\n-\n- fn effective_request_controls(&self, node: &Node) -> Result {\n- self.resolve_effective_request_controls(node)\n- }\n-\n- async fn one_shot(&self, request: OneShotRequest<'_>) -> Result {\n- let node = request.node;\n- let prompt = request.prompt;\n- let system_prompt = request.system_prompt;\n- let emitter = request.emitter;\n- let stage_scope = request.stage_scope;\n-\n- let client = Client::from_source(self.source.as_ref(), Arc::clone(&self.catalog))\n- .await\n- .map_err(|e| Error::handler_with_source(\"Failed to create LLM client\", e))?;\n-\n- let model = node.model().unwrap_or(&self.model);\n- let provider = self.resolve_provider_context(model, node.provider())?;\n- let provider_id = provider.provider_id.to_string();\n- let controls = self.resolve_effective_request_controls(node)?;\n-\n- let max_tokens = node\n- .max_tokens()\n- .or_else(|| self.catalog.get(model).and_then(|m| m.limits.max_output));\n-\n- let mut messages = Vec::new();\n- if let Some(sys) = system_prompt {\n- messages.push(Message::system(sys));\n- }\n- messages.push(Message::user(prompt));\n-\n- let request = Request {\n- model: model.to_string(),\n- messages,\n- provider: Some(provider_id),\n- reasoning_effort: controls.reasoning_effort,\n- speed: controls.speed,\n- tools: None,\n- tool_choice: None,\n- response_format: None,\n- temperature: None,\n- top_p: None,\n- max_tokens,\n- stop_sequences: None,\n- metadata: None,\n- provider_options: None,\n- };\n-\n- // Build per-request fallback chain: if the node overrides the provider,\n- // no failover is available; otherwise use the backend's.\n- let fallback_chain: &[FallbackTarget] = if node.provider().is_some() {\n- &[]\n- } else {\n- &self.fallback_chain\n- };\n-\n- let result = client.complete(&request).await;\n \n+ async fn complete_one_shot_request(\n+ &self,\n+ client: &Client,\n+ node: &Node,\n+ emitter: &Arc,\n+ stage_scope: &StageScope,\n+ request: &Request,\n+ controls: EffectiveRequestControls,\n+ fallback_chain: &[FallbackTarget],\n+ ) -> Result {\n+ let result = client.complete(request).await;\n let default_provider = self.provider_id.to_string();\n \n let (response, actual_model, actual_provider, actual_speed) = match result {\n@@ -916,7 +901,7 @@ impl CodergenBackend for AgentApiBackend {\n let max_tokens = node.max_tokens().or_else(|| {\n self.catalog\n .get(&target.model)\n- .and_then(|m| m.limits.max_output)\n+ .and_then(|model| model.limits.max_output)\n });\n \n let fallback_request = Request {\n@@ -930,12 +915,12 @@ impl CodergenBackend for AgentApiBackend {\n \n match client.complete(&fallback_request).await {\n Ok(resp) => {\n- found = Some((\n- resp,\n- target.model.clone(),\n- target.provider.clone(),\n- controls.speed,\n- ));\n+ found = Some(OneShotCompletion {\n+ response: resp,\n+ actual_model: target.model.clone(),\n+ actual_provider: target.provider.clone(),\n+ actual_speed: controls.speed,\n+ });\n break;\n }\n Err(err) if err.failover_eligible() => {\n@@ -946,30 +931,144 @@ impl CodergenBackend for AgentApiBackend {\n }\n \n match found {\n- Some(triple) => triple,\n+ Some(completion) => return Ok(completion),\n None => return Err(Error::Llm(last_err)),\n }\n }\n Err(sdk_err) => return Err(Error::Llm(sdk_err)),\n };\n \n- let stage_usage = billed_model_usage_from_llm(\n- self.catalog.as_ref(),\n- &ModelRef {\n- provider: ProviderId::from(actual_provider),\n- model_id: actual_model,\n- speed: actual_speed,\n- },\n- &response.usage,\n- )?;\n-\n- Ok(CodergenResult::Text {\n- text: response.text(),\n- usage: Some(stage_usage),\n- files_touched: Vec::new(),\n- last_file_touched: None,\n+ Ok(OneShotCompletion {\n+ response,\n+ actual_model,\n+ actual_provider,\n+ actual_speed,\n })\n }\n+}\n+\n+#[async_trait]\n+impl CodergenBackend for AgentApiBackend {\n+ async fn shutdown(&self, emitter: &Arc) {\n+ self.shutdown_cached_sessions(emitter);\n+ }\n+\n+ fn effective_request_controls(&self, node: &Node) -> Result {\n+ self.resolve_effective_request_controls(node)\n+ }\n+\n+ async fn one_shot(&self, request: OneShotRequest<'_>) -> Result {\n+ let node = request.node;\n+ let prompt = request.prompt;\n+ let system_prompt = request.system_prompt;\n+ let emitter = request.emitter;\n+ let stage_scope = request.stage_scope;\n+\n+ let client = Client::from_source(self.source.as_ref(), Arc::clone(&self.catalog))\n+ .await\n+ .map_err(|e| Error::handler_with_source(\"Failed to create LLM client\", e))?;\n+\n+ let model = node.model().unwrap_or(&self.model);\n+ let provider = self.resolve_provider_context(model, node.provider())?;\n+ let provider_id = provider.provider_id.to_string();\n+ let controls = self.resolve_effective_request_controls(node)?;\n+\n+ let max_tokens = node\n+ .max_tokens()\n+ .or_else(|| self.catalog.get(model).and_then(|m| m.limits.max_output));\n+\n+ let mut messages = Vec::new();\n+ if let Some(sys) = system_prompt {\n+ messages.push(Message::system(sys));\n+ }\n+ messages.push(Message::user(prompt));\n+\n+ // Build per-request fallback chain: if the node overrides the provider,\n+ // no failover is available; otherwise use the backend's.\n+ let fallback_chain: &[FallbackTarget] = if node.provider().is_some() {\n+ &[]\n+ } else {\n+ &self.fallback_chain\n+ };\n+\n+ let output_schema = structured_output::parse_node_output_schema(node)?;\n+ let response_format = output_schema\n+ .as_ref()\n+ .map(structured_output::prompt_response_format);\n+ let mut repair_attempts = 0_i64;\n+ let mut total_usage = TokenCounts::default();\n+\n+ loop {\n+ let request = Request {\n+ model: model.to_string(),\n+ messages: messages.clone(),\n+ provider: Some(provider_id.clone()),\n+ reasoning_effort: controls.reasoning_effort,\n+ speed: controls.speed,\n+ tools: None,\n+ tool_choice: None,\n+ response_format: response_format.clone(),\n+ temperature: None,\n+ top_p: None,\n+ max_tokens,\n+ stop_sequences: None,\n+ metadata: None,\n+ provider_options: None,\n+ };\n+\n+ let completion = self\n+ .complete_one_shot_request(\n+ &client,\n+ node,\n+ emitter,\n+ stage_scope,\n+ &request,\n+ controls,\n+ fallback_chain,\n+ )\n+ .await?;\n+ total_usage += completion.response.usage.clone();\n+ let response_text = completion.response.text();\n+\n+ let validation_error = if let Some(schema) = &output_schema {\n+ match structured_output::validate_response_text(schema, &response_text) {\n+ Ok(_) => None,\n+ Err(error) => Some((schema, error)),\n+ }\n+ } else {\n+ None\n+ };\n+\n+ if let Some((schema, error)) = validation_error {\n+ if repair_attempts >= node.output_retries() {\n+ return Err(Error::OutputSchemaValidation(\n+ structured_output::exhausted_failure_reason(node.output_retries()),\n+ ));\n+ }\n+ messages.push(Message::assistant(response_text));\n+ messages.push(Message::user(error.repair_message(schema)));\n+ repair_attempts += 1;\n+ continue;\n+ }\n+\n+ let stage_usage = billed_model_usage_from_llm(\n+ self.catalog.as_ref(),\n+ &ModelRef {\n+ provider: ProviderId::from(completion.actual_provider),\n+ model_id: completion.actual_model,\n+ speed: completion.actual_speed,\n+ },\n+ &total_usage,\n+ )?;\n+\n+ return Ok(CodergenResult::Text {\n+ text: response_text,\n+ usage: Some(stage_usage),\n+ files_touched: Vec::new(),\n+ last_file_touched: None,\n+ });\n+ }\n+ }\n \n async fn run(&self, request: CodergenRunRequest<'_>) -> Result {\n let node = request.node;\n@@ -981,6 +1080,7 @@ impl CodergenBackend for AgentApiBackend {\n let tool_hooks = request.tool_hooks;\n let cancel_token = request.cancel_token;\n let agent_tool_runtime = request.agent_tool_runtime;\n+ let output_schema = structured_output::parse_node_output_schema(node)?;\n \n let fidelity = context.fidelity();\n let reuse_key = if fidelity == Fidelity::Full {\n@@ -1258,8 +1358,59 @@ impl CodergenBackend for AgentApiBackend {\n return Err(err);\n }\n \n+ let mut response = last_assistant_response(&session);\n+ if let Some(schema) = &output_schema {\n+ let mut repair_attempts = 0_i64;\n+ loop {\n+ let (_, last_file_touched) = file_tracking_snapshot(&file_tracking);\n+ match validate_agent_output_sources(\n+ schema,\n+ &response,\n+ sandbox,\n+ last_file_touched.as_deref(),\n+ )\n+ .await\n+ {\n+ Ok(_) => break,\n+ Err(error) => {\n+ if repair_attempts >= node.output_retries() {\n+ bridge.abort();\n+ discard_session(&mut session, &mut lease, emitter);\n+ return Err(Error::OutputSchemaValidation(\n+ structured_output::exhausted_failure_reason(node.output_retries()),\n+ ));\n+ }\n+ let repair_message = error.repair_message(schema);\n+ match session.process_input(&repair_message).await {\n+ Ok(()) => {\n+ repair_attempts += 1;\n+ response = last_assistant_response(&session);\n+ }\n+ Err(err) => match classify_agent_error(err, false) {\n+ AgentApiErrorDisposition::Cancelled => {\n+ bridge.abort();\n+ discard_session(&mut session, &mut lease, emitter);\n+ return Err(Error::Cancelled);\n+ }\n+ AgentApiErrorDisposition::Terminal(err) => {\n+ bridge.abort();\n+ discard_session(&mut session, &mut lease, emitter);\n+ return Err(err);\n+ }\n+ AgentApiErrorDisposition::FailoverEligible(sdk_err) => {\n+ bridge.abort();\n+ discard_session(&mut session, &mut lease, emitter);\n+ return Err(Error::Llm(sdk_err));\n+ }\n+ },\n+ }\n+ }\n+ }\n+ }\n+ }\n+\n // Aggregate token usage only from new turns (prevents double-counting on\n- // reuse).\n+ // reuse), including any output-schema repair turns.\n let mut total_usage = TokenCounts::default();\n for turn in &session.history().turns()[turns_before..] {\n if let AgentMessage::Assistant { usage, .. } = turn {\n@@ -1278,29 +1429,8 @@ impl CodergenBackend for AgentApiBackend {\n &total_usage,\n )?;\n \n- // Extract last assistant response from the session history.\n- let response = session\n- .history()\n- .turns()\n- .iter()\n- .rev()\n- .find_map(|turn| {\n- if let AgentMessage::Assistant { content, .. } = turn {\n- if !content.is_empty() {\n- return Some(content.clone());\n- }\n- }\n- None\n- })\n- .unwrap_or_default();\n-\n // Collect files_touched from the shared tracking state.\n- let (files_touched, last_file_touched) = {\n- let s = file_tracking.lock().unwrap();\n- let mut v: Vec = s.touched.iter().cloned().collect();\n- v.sort();\n- (v, s.last.clone())\n- };\n+ let (files_touched, last_file_touched) = file_tracking_snapshot(&file_tracking);\n \n if let Some(lease) = lease.take() {\n lease.release();\n@@ -1376,10 +1506,13 @@ mod tests {\n };\n use fabro_vault::{SecretType, Vault};\n use futures::stream;\n+ use httpmock::Method::POST;\n+ use httpmock::MockServer;\n use tokio::sync::RwLock as AsyncRwLock;\n use tokio_util::sync::CancellationToken;\n \n use super::*;\n+ use crate::context::Context;\n use crate::services::FabroRunToolServices;\n \n struct ShutdownTestProfile {\n@@ -1447,6 +1580,109 @@ mod tests {\n }\n }\n \n+ fn mock_llm_catalog(server: &MockServer) -> Arc {\n+ let settings: LlmCatalogSettings = toml::from_str(&format!(\n+ r#\"\n+[providers.mock]\n+adapter = \"openai_compatible\"\n+agent_profile = \"openai\"\n+base_url = \"{}\"\n+\n+[providers.mock.auth]\n+credentials = [\"env:MOCK_API_KEY\"]\n+\n+[models.mock-model]\n+provider = \"mock\"\n+display_name = \"Mock Model\"\n+family = \"mock\"\n+default = true\n+\n+[models.mock-model.limits]\n+context_window = 8192\n+max_output = 1024\n+\n+[models.mock-model.features]\n+tools = true\n+vision = false\n+reasoning = false\n+\"#,\n+ server.base_url()\n+ ))\n+ .unwrap();\n+ Arc::new(Catalog::from_builtin_with_overrides(&settings).unwrap())\n+ }\n+\n+ fn mock_api_backend(server: &MockServer) -> AgentApiBackend {\n+ let source = EnvCredentialSource::with_env_lookup(Arc::new(|name| {\n+ if name == \"MOCK_API_KEY\" {\n+ Some(\"sk-test\".to_string())\n+ } else {\n+ None\n+ }\n+ }));\n+ AgentApiBackend::new_with_catalog(\n+ \"mock-model\".to_string(),\n+ ProviderId::from(\"mock\"),\n+ Vec::new(),\n+ Arc::new(source),\n+ SteeringHub::for_tests(),\n+ mock_llm_catalog(server),\n+ )\n+ }\n+\n+ fn chat_completion_response(\n+ text: &str,\n+ input_tokens: i64,\n+ output_tokens: i64,\n+ ) -> serde_json::Value {\n+ serde_json::json!({\n+ \"id\": uuid::Uuid::new_v4().to_string(),\n+ \"model\": \"mock-model\",\n+ \"choices\": [{\n+ \"message\": {\n+ \"content\": text\n+ },\n+ \"finish_reason\": \"stop\"\n+ }],\n+ \"usage\": {\n+ \"prompt_tokens\": input_tokens,\n+ \"completion_tokens\": output_tokens,\n+ \"total_tokens\": input_tokens + output_tokens\n+ }\n+ })\n+ }\n+\n+ fn chat_completion_stream(text: &str, input_tokens: i64, output_tokens: i64) -> String {\n+ let text_chunk = serde_json::json!({\n+ \"id\": uuid::Uuid::new_v4().to_string(),\n+ \"model\": \"mock-model\",\n+ \"choices\": [{\n+ \"delta\": {\n+ \"content\": text\n+ },\n+ \"finish_reason\": null\n+ }]\n+ });\n+ let usage_chunk = serde_json::json!({\n+ \"id\": uuid::Uuid::new_v4().to_string(),\n+ \"model\": \"mock-model\",\n+ \"choices\": [],\n+ \"usage\": {\n+ \"prompt_tokens\": input_tokens,\n+ \"completion_tokens\": output_tokens,\n+ \"total_tokens\": input_tokens + output_tokens\n+ }\n+ });\n+ format!(\"data: {text_chunk}\\n\\ndata: {usage_chunk}\\n\\ndata: [DONE]\\n\\n\")\n+ }\n+\n+ fn custom_output_schema_attr() -> AttrValue {\n+ AttrValue::String(\n+ r#\"{\"type\":\"object\",\"required\":[\"passed\"],\"properties\":{\"passed\":{\"type\":\"boolean\"}}}\"#\n+ .to_string(),\n+ )\n+ }\n+\n #[test]\n fn agent_backend_stores_config() {\n let backend = AgentApiBackend::new_from_env(\n@@ -2212,6 +2448,127 @@ reasoning = false\n assert_eq!(client.provider_names(), vec![\"anthropic\"]);\n }\n \n+ #[tokio::test]\n+ async fn one_shot_repairs_custom_output_schema_with_previous_assistant_message() {\n+ let server = MockServer::start();\n+ let first = server.mock(|when, then| {\n+ when.method(POST)\n+ .path(\"/chat/completions\")\n+ .body_includes(r#\"\"type\":\"json_schema\"\"#)\n+ .body_excludes(r#\"\"role\":\"assistant\"\"#);\n+ then.status(200)\n+ .header(\"content-type\", \"application/json\")\n+ .json_body(chat_completion_response(\"not json\", 10, 1));\n+ });\n+ let repair = server.mock(|when, then| {\n+ when.method(POST)\n+ .path(\"/chat/completions\")\n+ .body_includes(r#\"\"type\":\"json_schema\"\"#)\n+ .body_includes(r#\"\"role\":\"assistant\"\"#)\n+ .body_includes(\"not json\")\n+ .body_includes(\"output_schema\");\n+ then.status(200)\n+ .header(\"content-type\", \"application/json\")\n+ .json_body(chat_completion_response(r#\"{\"passed\":true}\"#, 11, 2));\n+ });\n+ let backend = mock_api_backend(&server);\n+ let mut node = Node::new(\"audit\");\n+ node.attrs\n+ .insert(\"output_schema\".to_string(), custom_output_schema_attr());\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(1));\n+ let context = Context::new();\n+ let stage_scope = StageScope::for_handler(&context, &node.id);\n+ let emitter = Arc::new(Emitter::new(fabro_types::RunId::new()));\n+ let workspace = tempfile::tempdir().unwrap();\n+ let sandbox: Arc =\n+ Arc::new(LocalSandbox::new(workspace.path().to_path_buf()));\n+\n+ let result = backend\n+ .one_shot(OneShotRequest {\n+ node: &node,\n+ prompt: \"Audit the result\",\n+ system_prompt: None,\n+ emitter: &emitter,\n+ stage_scope: &stage_scope,\n+ sandbox: &sandbox,\n+ cancel_token: CancellationToken::new(),\n+ })\n+ .await\n+ .unwrap();\n+\n+ first.assert_calls(1);\n+ repair.assert_calls(1);\n+ let CodergenResult::Text { text, usage, .. } = result else {\n+ panic!(\"one_shot should return text\");\n+ };\n+ assert_eq!(text, r#\"{\"passed\":true}\"#);\n+ let usage = usage.expect(\"usage should be aggregated\");\n+ assert_eq!(usage.tokens().input_tokens, 21);\n+ assert_eq!(usage.tokens().output_tokens, 3);\n+ }\n+\n+ #[tokio::test]\n+ async fn agent_run_repairs_custom_output_schema_in_same_session() {\n+ let server = MockServer::start();\n+ let first = server.mock(|when, then| {\n+ when.method(POST)\n+ .path(\"/chat/completions\")\n+ .body_includes(r#\"\"stream\":true\"#)\n+ .body_excludes(r#\"\"role\":\"assistant\"\"#);\n+ then.status(200)\n+ .header(\"content-type\", \"text/event-stream\")\n+ .body(chat_completion_stream(\"not json\", 20, 3));\n+ });\n+ let repair = server.mock(|when, then| {\n+ when.method(POST)\n+ .path(\"/chat/completions\")\n+ .body_includes(r#\"\"stream\":true\"#)\n+ .body_includes(r#\"\"role\":\"assistant\"\"#)\n+ .body_includes(\"not json\")\n+ .body_includes(\"output_schema\");\n+ then.status(200)\n+ .header(\"content-type\", \"text/event-stream\")\n+ .body(chat_completion_stream(r#\"{\"passed\":true}\"#, 21, 4));\n+ });\n+ let backend = mock_api_backend(&server);\n+ let mut node = Node::new(\"audit\");\n+ node.attrs\n+ .insert(\"output_schema\".to_string(), custom_output_schema_attr());\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(1));\n+ let context = Context::new();\n+ let emitter = Arc::new(Emitter::new(fabro_types::RunId::new()));\n+ let workspace = tempfile::tempdir().unwrap();\n+ let sandbox: Arc =\n+ Arc::new(LocalSandbox::new(workspace.path().to_path_buf()));\n+\n+ let result = backend\n+ .run(CodergenRunRequest {\n+ node: &node,\n+ prompt: \"Audit the result\",\n+ context: &context,\n+ thread_id: None,\n+ emitter: &emitter,\n+ sandbox: &sandbox,\n+ tool_hooks: None,\n+ cancel_token: CancellationToken::new(),\n+ agent_tool_runtime: fabro_agent::AgentToolRuntime::default(),\n+ })\n+ .await\n+ .unwrap();\n+\n+ first.assert_calls(1);\n+ repair.assert_calls(1);\n+ let CodergenResult::Text { text, usage, .. } = result else {\n+ panic!(\"run should return text\");\n+ };\n+ assert_eq!(text, r#\"{\"passed\":true}\"#);\n+ let usage = usage.expect(\"usage should be aggregated\");\n+ assert_eq!(usage.tokens().input_tokens, 41);\n+ assert_eq!(usage.tokens().output_tokens, 7);\n+ }\n+\n #[tokio::test]\n async fn api_backend_shutdown_closes_cached_sessions_once() {\n let backend = AgentApiBackend::new_from_env(\ndiff --git a/lib/crates/fabro-workflow/src/handler/mod.rs b/lib/crates/fabro-workflow/src/handler/mod.rs\nindex 675f24b39..bc94af3b8 100644\n--- a/lib/crates/fabro-workflow/src/handler/mod.rs\n+++ b/lib/crates/fabro-workflow/src/handler/mod.rs\n@@ -9,6 +9,7 @@ pub mod manager_loop;\n pub mod parallel;\n pub mod prompt;\n pub mod start;\n+pub mod structured_output;\n pub mod wait;\n \n use std::any::Any;\ndiff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs\nindex 6dbc0576d..27b9fa524 100644\n--- a/lib/crates/fabro-workflow/src/handler/prompt.rs\n+++ b/lib/crates/fabro-workflow/src/handler/prompt.rs\n@@ -10,7 +10,7 @@ use super::agent::{\n truncate,\n };\n use super::llm::routing;\n-use super::{EngineServices, Handler};\n+use super::{EngineServices, Handler, structured_output};\n use crate::context::{Context, WorkflowContext, keys};\n use crate::error::Error;\n use crate::event::{Emitter, Event};\n@@ -191,7 +191,25 @@ impl Handler for PromptHandler {\n serde_json::json!(&response_text),\n );\n \n- extract_status_fields(&response_text, &mut outcome);\n+ if let Some(schema) = structured_output::parse_node_output_schema(node)? {\n+ match structured_output::validate_response_text(&schema, &response_text) {\n+ Ok(validated) => {\n+ structured_output::apply_validated_output(\n+ node,\n+ &schema,\n+ &validated,\n+ &mut outcome,\n+ );\n+ }\n+ Err(_) => {\n+ return Ok(structured_output::exhausted_failure_outcome(\n+ node.output_retries(),\n+ ));\n+ }\n+ }\n+ } else {\n+ extract_status_fields(&response_text, &mut outcome);\n+ }\n outcome.usage = stage_usage;\n outcome.files_touched = backend_files_touched;\n \n@@ -214,6 +232,7 @@ mod tests {\n use super::*;\n use crate::event::Emitter;\n use crate::handler::agent::CodergenRunRequest;\n+ use crate::outcome::OutcomeExt;\n \n fn make_services() -> EngineServices {\n EngineServices::test_default()\n@@ -368,6 +387,102 @@ mod tests {\n );\n }\n \n+ #[tokio::test]\n+ async fn prompt_handler_custom_output_schema_updates_output_context_key() {\n+ struct CustomOutputBackend;\n+\n+ #[async_trait]\n+ impl CodergenBackend for CustomOutputBackend {\n+ async fn run(&self, _request: CodergenRunRequest<'_>) -> Result {\n+ panic!(\"run() should not be called for prompt handler\");\n+ }\n+\n+ async fn one_shot(\n+ &self,\n+ _request: OneShotRequest<'_>,\n+ ) -> Result {\n+ Ok(CodergenResult::Text {\n+ text: r#\"{\"passed\": true}\"#.to_string(),\n+ usage: None,\n+ files_touched: Vec::new(),\n+ last_file_touched: None,\n+ })\n+ }\n+ }\n+\n+ let handler = PromptHandler::new(Some(Box::new(CustomOutputBackend)));\n+ let mut node = Node::new(\"audit\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\n+ r#\"{\"type\":\"object\",\"required\":[\"passed\"],\"properties\":{\"passed\":{\"type\":\"boolean\"}}}\"#\n+ .to_string(),\n+ ),\n+ );\n+ let context = Context::new();\n+ let graph = Graph::new(\"test\");\n+ let tmp = TempDir::new().unwrap();\n+\n+ let outcome = handler\n+ .execute(&node, &context, &graph, tmp.path(), &make_services())\n+ .await\n+ .unwrap();\n+\n+ assert_eq!(\n+ outcome.context_updates.get(\"output.audit\"),\n+ Some(&serde_json::json!({\"passed\": true})),\n+ );\n+ }\n+\n+ #[tokio::test]\n+ async fn prompt_handler_routing_output_schema_requires_valid_routing_json() {\n+ struct BadRoutingBackend;\n+\n+ #[async_trait]\n+ impl CodergenBackend for BadRoutingBackend {\n+ async fn run(&self, _request: CodergenRunRequest<'_>) -> Result {\n+ panic!(\"run() should not be called for prompt handler\");\n+ }\n+\n+ async fn one_shot(\n+ &self,\n+ _request: OneShotRequest<'_>,\n+ ) -> Result {\n+ Ok(CodergenResult::Text {\n+ text: r#\"{\"outcome\": 123}\"#.to_string(),\n+ usage: None,\n+ files_touched: Vec::new(),\n+ last_file_touched: None,\n+ })\n+ }\n+ }\n+\n+ let handler = PromptHandler::new(Some(Box::new(BadRoutingBackend)));\n+ let mut node = Node::new(\"route\");\n+ node.attrs.insert(\n+ \"output_schema\".to_string(),\n+ AttrValue::String(\"routing\".to_string()),\n+ );\n+ node.attrs\n+ .insert(\"output_retries\".to_string(), AttrValue::Integer(0));\n+ let context = Context::new();\n+ let graph = Graph::new(\"test\");\n+ let tmp = TempDir::new().unwrap();\n+\n+ let outcome = handler\n+ .execute(&node, &context, &graph, tmp.path(), &make_services())\n+ .await\n+ .unwrap();\n+\n+ assert_eq!(outcome.status, crate::outcome::StageOutcome::Failed {\n+ retry_requested: false,\n+ });\n+ assert_eq!(\n+ outcome.failure_reason(),\n+ Some(\"output schema validation failed after 0 repair attempt(s)\")\n+ );\n+ }\n+\n #[tokio::test]\n async fn prompt_handler_projects_provider_used_from_prompt_events() {\n struct ProviderOneShotBackend;\ndiff --git a/lib/crates/fabro-workflow/src/handler/structured_output.rs b/lib/crates/fabro-workflow/src/handler/structured_output.rs\nnew file mode 100644\nindex 000000000..d98a42afb\n--- /dev/null\n+++ b/lib/crates/fabro-workflow/src/handler/structured_output.rs\n@@ -0,0 +1,608 @@\n+use fabro_graphviz::graph::Node;\n+use fabro_llm::types::{ResponseFormat, ResponseFormatType};\n+use serde_json::Value;\n+\n+use crate::error::Error;\n+use crate::outcome::{FailureCategory, FailureDetail, Outcome, StageOutcome};\n+\n+pub(crate) const ROUTING_STATUS_FIELDS: &[&str] = &[\n+ \"preferred_next_label\",\n+ \"outcome\",\n+ \"failure_reason\",\n+ \"suggested_next_ids\",\n+ \"context_updates\",\n+];\n+\n+#[derive(Debug, Clone, PartialEq)]\n+pub(crate) enum OutputSchemaKind {\n+ Routing,\n+ JsonSchema { schema: Value },\n+}\n+\n+#[derive(Debug, Clone, Copy, PartialEq, Eq)]\n+pub(crate) enum StructuredOutputErrorKind {\n+ NoJsonObject,\n+ NoRelevantJsonObject,\n+ InvalidJson,\n+ SchemaValidation,\n+}\n+\n+#[derive(Debug, Clone, PartialEq, Eq)]\n+pub(crate) struct StructuredOutputError {\n+ kind: StructuredOutputErrorKind,\n+ messages: Vec,\n+}\n+\n+impl StructuredOutputError {\n+ fn new(kind: StructuredOutputErrorKind, message: impl Into) -> Self {\n+ Self {\n+ kind,\n+ messages: vec![message.into()],\n+ }\n+ }\n+\n+ fn validation(messages: Vec) -> Self {\n+ Self {\n+ kind: StructuredOutputErrorKind::SchemaValidation,\n+ messages,\n+ }\n+ }\n+\n+ #[cfg(test)]\n+ #[must_use]\n+ pub(crate) fn kind(&self) -> StructuredOutputErrorKind {\n+ self.kind\n+ }\n+\n+ #[cfg(test)]\n+ #[must_use]\n+ pub(crate) fn messages(&self) -> &[String] {\n+ &self.messages\n+ }\n+\n+ #[must_use]\n+ pub(crate) fn allows_routing_fallback(&self) -> bool {\n+ matches!(\n+ self.kind,\n+ StructuredOutputErrorKind::NoJsonObject\n+ | StructuredOutputErrorKind::NoRelevantJsonObject\n+ )\n+ }\n+\n+ #[must_use]\n+ pub(crate) fn repair_message(&self, schema: &OutputSchemaKind) -> String {\n+ let expectation = match schema {\n+ OutputSchemaKind::Routing => format!(\n+ \"Return a single JSON object with at least one routing field: {}.\",\n+ ROUTING_STATUS_FIELDS.join(\", \")\n+ ),\n+ OutputSchemaKind::JsonSchema { .. } => {\n+ \"Return a single JSON object that satisfies the configured JSON Schema.\".to_string()\n+ }\n+ };\n+ let errors = self\n+ .messages\n+ .iter()\n+ .map(|message| format!(\"- {message}\"))\n+ .collect::>()\n+ .join(\"\\n\");\n+ format!(\n+ \"Your previous response did not satisfy the node's output_schema.\\n\\n\\\n+ Validation errors:\\n{errors}\\n\\n\\\n+ {expectation}\\n\\\n+ Do not include Markdown fences or explanatory prose; reply only with the corrected JSON object.\"\n+ )\n+ }\n+}\n+\n+#[derive(Debug, Clone, PartialEq)]\n+pub(crate) struct ValidatedStructuredOutput {\n+ pub(crate) value: Value,\n+}\n+\n+#[must_use]\n+pub(crate) fn output_key(node_id: &str) -> String {\n+ format!(\"output.{node_id}\")\n+}\n+\n+#[must_use]\n+pub(crate) fn exhausted_failure_reason(repair_attempts: i64) -> String {\n+ format!(\"output schema validation failed after {repair_attempts} repair attempt(s)\")\n+}\n+\n+#[must_use]\n+pub(crate) fn exhausted_failure_outcome(repair_attempts: i64) -> Outcome {\n+ Outcome {\n+ status: StageOutcome::Failed {\n+ retry_requested: false,\n+ },\n+ failure: Some(FailureDetail::new(\n+ exhausted_failure_reason(repair_attempts),\n+ FailureCategory::Deterministic,\n+ )),\n+ ..Outcome::default()\n+ }\n+}\n+\n+pub(crate) fn parse_node_output_schema(node: &Node) -> Result