mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-11 03:40:05 +00:00
Merge remote-tracking branch 'origin/main' into pr-599
# Conflicts: # docs/public/agents/outputs.mdx # docs/public/reference/dot-language.mdx
This commit is contained in:
commit
7bdee5b494
24 changed files with 1618 additions and 220 deletions
|
|
@ -12,7 +12,7 @@ digraph ImplementPlan {
|
|||
implement [label="Implement", prompt="Read the plan file referenced in the goal and implement every step. Make all the code changes described in the plan. Use red/green TDD.", model="openai/gpt-5.6-sol", provider="openrouter", reasoning_effort="xhigh"]
|
||||
simplify_fable [label="Simplify (Claude Fable 5)", prompt="@prompts/simplify.md", model="anthropic/claude-fable-5", provider="openrouter", reasoning_effort="xhigh"]
|
||||
simplify_sol [label="Simplify (GPT-5.6 Sol)", prompt="@prompts/simplify.md", model="openai/gpt-5.6-sol", provider="openrouter", reasoning_effort="max"]
|
||||
verify [label="Verify", shape=parallelogram, script="git fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && { command -v rg >/dev/null 2>&1 || { echo 'rg is required for verify'; exit 127; }; } && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\bActorRef\b|\bActorKind\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\s*==\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1", goal_gate=true, retry_target="fixup"]
|
||||
verify [label="Verify", shape=parallelogram, script="git fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && { command -v rg >/dev/null 2>&1 || { echo 'rg is required for verify'; exit 127; }; } && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\bActorRef\b|\bActorKind\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\s*==\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1", timeout="20m", goal_gate=true, retry_target="fixup"]
|
||||
fixup [label="Fixup", prompt="The verify step failed. Read the build output from context and fix all format, clippy, Rust test, docs, TypeScript typecheck/test, and build failures.", model="anthropic/claude-fable-5", provider="openrouter", reasoning_effort="xhigh", max_visits=3]
|
||||
|
||||
start -> toolchain
|
||||
|
|
|
|||
103
Cargo.lock
generated
103
Cargo.lock
generated
|
|
@ -2239,7 +2239,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-acp"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"agent-client-protocol",
|
||||
"agent-client-protocol-tokio",
|
||||
|
|
@ -2258,7 +2258,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-agent"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -2300,7 +2300,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-api"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"fabro-automation",
|
||||
|
|
@ -2323,7 +2323,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-auth"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -2348,7 +2348,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-automation"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -2367,11 +2367,11 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-build-support"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
|
||||
[[package]]
|
||||
name = "fabro-checkpoint"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"fabro-config",
|
||||
|
|
@ -2387,7 +2387,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-cli"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"assert_cmd",
|
||||
|
|
@ -2489,7 +2489,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-client"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"bytes",
|
||||
|
|
@ -2518,7 +2518,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-config"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -2547,7 +2547,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-core"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"fabro-types",
|
||||
|
|
@ -2562,7 +2562,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-db"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -2574,7 +2574,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-dev"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"assert_cmd",
|
||||
|
|
@ -2593,7 +2593,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-dump"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"bytes",
|
||||
|
|
@ -2607,7 +2607,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-environment"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -2629,7 +2629,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-github"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"base64",
|
||||
|
|
@ -2651,7 +2651,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-graphviz"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"fabro-types",
|
||||
|
|
@ -2659,13 +2659,14 @@ dependencies = [
|
|||
"nom",
|
||||
"regex",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"strum 0.28.0",
|
||||
"thiserror 2.0.18",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fabro-hooks"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"fabro-agent",
|
||||
|
|
@ -2688,7 +2689,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-http"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"fabro-static",
|
||||
"http 1.4.0",
|
||||
|
|
@ -2698,7 +2699,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-install"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"base64",
|
||||
|
|
@ -2717,7 +2718,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-interview"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"dialoguer",
|
||||
|
|
@ -2732,7 +2733,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-llm"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -2773,7 +2774,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-macros"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"clap",
|
||||
"fabro-options-metadata",
|
||||
|
|
@ -2784,7 +2785,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-manifest"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"fabro-api",
|
||||
|
|
@ -2802,7 +2803,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-mcp"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum",
|
||||
|
|
@ -2822,7 +2823,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-mcp-server"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -2849,7 +2850,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-mcp-store"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"fabro-db",
|
||||
|
|
@ -2867,7 +2868,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-model"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"fabro-static",
|
||||
"http 1.4.0",
|
||||
|
|
@ -2883,7 +2884,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-oauth"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"axum",
|
||||
|
|
@ -2905,7 +2906,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-options-metadata"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
|
|
@ -2913,7 +2914,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-proc"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"cc",
|
||||
"libc",
|
||||
|
|
@ -2922,7 +2923,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-redact"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"aho-corasick",
|
||||
"ref-cast",
|
||||
|
|
@ -2938,7 +2939,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-sandbox"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -2983,7 +2984,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-server"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -3075,7 +3076,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-slack"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"fabro-http",
|
||||
"fabro-interview",
|
||||
|
|
@ -3097,18 +3098,18 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-spa"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"rust-embed",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fabro-static"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
|
||||
[[package]]
|
||||
name = "fabro-store"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
|
|
@ -3138,7 +3139,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-telemetry"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"base64",
|
||||
|
|
@ -3164,7 +3165,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-template"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"fabro-types",
|
||||
|
|
@ -3178,7 +3179,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-test"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"assert_cmd",
|
||||
|
|
@ -3203,7 +3204,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-tool"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -3224,7 +3225,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-tracker"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-trait",
|
||||
|
|
@ -3238,7 +3239,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-types"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"clap",
|
||||
|
|
@ -3260,7 +3261,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-util"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"console 0.15.11",
|
||||
|
|
@ -3281,7 +3282,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-validate"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"fabro-acp",
|
||||
"fabro-graphviz",
|
||||
|
|
@ -3294,7 +3295,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-variable"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -3311,7 +3312,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-vault"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"chrono",
|
||||
|
|
@ -3330,7 +3331,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "fabro-workflow"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"assert_cmd",
|
||||
|
|
@ -8494,7 +8495,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "twin-github"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"base64",
|
||||
|
|
@ -8513,7 +8514,7 @@ dependencies = [
|
|||
|
||||
[[package]]
|
||||
name = "twin-openai"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"async-stream",
|
||||
|
|
|
|||
|
|
@ -5,7 +5,7 @@ resolver = "2"
|
|||
|
||||
[workspace.package]
|
||||
edition = "2021"
|
||||
version = "0.303.0-nightly.3"
|
||||
version = "0.303.0-nightly.4"
|
||||
license = "MIT"
|
||||
|
||||
[workspace.dependencies]
|
||||
|
|
|
|||
|
|
@ -97,6 +97,8 @@ Agent nodes can provide routing directives through fallback files. Fabro checks
|
|||
|
||||
This fallback chain applies to normal routing extraction and to `output_schema="routing"`. For validated routing, Fabro only advances to the next source when the current source has no JSON object or no object with recognized routing fields. If the current source contains malformed routing JSON or valid JSON with wrong routing field types, validation fails and Fabro starts the repair loop instead.
|
||||
|
||||
The last-file fallback only reads `.json` and `.md` files (case-insensitive), and the routing JSON must be the final JSON object in the file, with only whitespace after it. Fabro ignores other file types and routing JSON followed by any other content. These restrictions do not apply to the dedicated `status.json` fallback.
|
||||
|
||||
Prompt nodes do not use file fallbacks; they validate or extract routing directives from the response text only. Command nodes likewise have no file fallback and validate only their merged stdout and stderr.
|
||||
|
||||
If no source provides routing directives, the transition falls through to condition matching, unconditional edges, or weight-based tiebreaking as described in [Transitions](/workflows/transitions).
|
||||
|
|
|
|||
|
|
@ -139,6 +139,16 @@ Fidelity can be set at three levels. The first match wins:
|
|||
|
||||
If none of these are set, fidelity defaults to `compact`.
|
||||
|
||||
### Parallel branch fidelity
|
||||
|
||||
The first node in each parallel branch uses this precedence:
|
||||
|
||||
1. `fidelity` on the fork-to-branch edge
|
||||
2. `fidelity` on the branch node
|
||||
3. Otherwise, inherit the fork's preamble unchanged
|
||||
|
||||
Fabro renders any branch-specific preambles before fan-out from the fork's context snapshot, then places them into the isolated branch contexts. An explicit branch-level `full` degrades to `summary:high` because concurrent branches cannot share conversation sessions. `thread_id` on a branch node or fork-to-branch edge is inert.
|
||||
|
||||
### Full fidelity and threads
|
||||
|
||||
`full` fidelity is typically used with `thread_id` to create a shared conversation across multiple nodes. Nodes with the same `thread_id` share a single LLM session, preserving full context continuity:
|
||||
|
|
|
|||
|
|
@ -201,8 +201,8 @@ Start nodes can also be identified by ID (`start` or `Start`). Exit nodes can be
|
|||
| `prompt` | String | Task instructions for the LLM. Supports file references with `@path/to/file.md` |
|
||||
| `reasoning_effort` | String | `low`, `medium`, or `high` (default: `high`) |
|
||||
| `max_tokens` | Integer | Maximum output tokens |
|
||||
| `fidelity` | String | How much prior context is passed: `compact`, `full`, `summary:high`, `summary:medium`, `summary:low`, `truncate` |
|
||||
| `thread_id` | String | Groups nodes into a shared conversation thread |
|
||||
| `fidelity` | String | How much prior context is passed: `compact`, `full`, `summary:high`, `summary:medium`, `summary:low`, `truncate`. On a node entered directly from a parallel fork, this is overridden by the fork-to-branch edge; explicit `full` degrades to `summary:high`. |
|
||||
| `thread_id` | String | Groups nodes into a shared conversation thread. Inert when the node is entered directly from a parallel fork. |
|
||||
| `model` | String | Explicit model ID (overrides stylesheet) |
|
||||
| `provider` | String | Explicit provider name (overrides stylesheet). Auto-inferred from the model catalog when omitted. |
|
||||
| `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. |
|
||||
|
|
@ -235,7 +235,7 @@ audit [
|
|||
- 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.
|
||||
- `output_retries` defaults to `2` and controls only these corrective structured-output turns. Negative values are treated as `0`. It is not the same as `max_retries` and does not consume workflow retry attempts.
|
||||
- Custom schema output is stored in context at `output.{node_id}`. Routing schema output updates routing fields and any `context_updates`.
|
||||
- Agent routing fallbacks still apply to `output_schema="routing"`: response text first, then `status.json`, then the last file touched by the agent. Custom schemas and prompt nodes validate response text only.
|
||||
- Agent routing fallbacks still apply to `output_schema="routing"`: response text first, then `status.json`, then the last file touched by the agent. The last-file fallback only accepts `.json` and `.md` files (case-insensitive) whose final JSON object contains the routing directive; only whitespace may follow it. Custom schemas and prompt nodes validate response text only.
|
||||
- Command nodes validate merged stdout and stderr only after the script exits with code `0`, using the same object selection as agents and prompts: custom schemas validate the last JSON object, and `routing` validates the last JSON object containing a recognized routing field. Print the intended JSON object last. A validation error is a deterministic, non-retryable failure with no repair turn or `status.json` fallback. `output_retries`, `retry_policy`, and `max_retries` do not retry it. Nonzero exits retain normal command failure behavior without schema validation.
|
||||
- For commands, custom schema output is stored at `output.{node_id}`; edge conditions cannot traverse into its fields. The `routing` schema applies routing fields and merges `context_updates` into flat context keys, which conditions can read (for example, `context.kept_count`).
|
||||
- `backend="acp"` with `output_schema` is unsupported in this release.
|
||||
|
|
@ -255,6 +255,8 @@ audit [
|
|||
| `join_policy` | String | When the merge can proceed: `wait_all` (default), `first_success` |
|
||||
| `max_parallel` | Integer | Maximum concurrent branches (default: 4) |
|
||||
|
||||
For the first node in each branch, `fidelity` resolves from the fork-to-branch edge, then the branch node; without either, the fork preamble is inherited unchanged. Branch-specific preambles are rendered before fan-out from the fork's context snapshot. Concurrent branches cannot share sessions, so explicit branch `full` becomes `summary:high`, and branch-level `thread_id` is inert.
|
||||
|
||||
### Wait nodes
|
||||
|
||||
| Attribute | Type | Description |
|
||||
|
|
@ -285,8 +287,8 @@ audit [
|
|||
| `label` | String | Display text; also used for human gate option matching |
|
||||
| `condition` | String | Boolean expression for conditional routing (see below) |
|
||||
| `weight` | Integer | Priority for tiebreaking (higher wins, default: 0) |
|
||||
| `fidelity` | String | Override fidelity level for this transition |
|
||||
| `thread_id` | String | Override thread ID for this transition |
|
||||
| `fidelity` | String | Override fidelity level for this transition. On a fork-to-branch edge, takes precedence over the branch node; explicit `full` degrades to `summary:high`. |
|
||||
| `thread_id` | String | Override thread ID for this transition. Inert on fork-to-branch edges. |
|
||||
| `loop_restart` | Boolean | Restart the workflow from this edge's target when taken: stage history and retry counts clear and the context resets to empty (visit counts are kept). Failed outcomes may only take it for `transient_infra` failures — see [Failures](/execution/failures#loop-restart-edges) |
|
||||
| `freeform` | Boolean | When `true` on a human-gate edge, accept free-text input instead of fixed choices |
|
||||
|
||||
|
|
|
|||
|
|
@ -170,6 +170,8 @@ fork -> quality
|
|||
| `join_policy` | When the merge can proceed (see table below) |
|
||||
| `max_parallel` | Maximum concurrent branches (default: 4) |
|
||||
|
||||
For each branch's first node, fidelity resolves from the fork-to-branch edge, then the branch node; otherwise it inherits the fork preamble unchanged. Fabro renders branch-specific preambles before fan-out from the fork snapshot. Branch-level `full` degrades to `summary:high` because concurrent branches cannot share sessions, and `thread_id` on a branch node or fork-to-branch edge is inert.
|
||||
|
||||
**Join policies:**
|
||||
|
||||
| Policy | Behavior |
|
||||
|
|
|
|||
|
|
@ -338,10 +338,16 @@ fn keep_existing_settings_persists_secrets_without_rewriting_server_target() {
|
|||
let mut context = test_context!();
|
||||
let storage_dir = context.temp_dir.join("install-storage");
|
||||
context.manage_storage_dir(&storage_dir);
|
||||
let existing_web_url = unused_loopback_web_url();
|
||||
let requested_web_url = unused_loopback_web_url();
|
||||
write_http_install_settings(&context, &storage_dir, &existing_web_url, "keep-me");
|
||||
login_with_storage_dev_token(&context, &storage_dir, &existing_web_url);
|
||||
let socket_path = context.temp_dir.join("install-storage.sock");
|
||||
let existing_web_url = "https://existing.example.test".to_string();
|
||||
let requested_web_url = "https://requested.example.test".to_string();
|
||||
write_unix_install_settings(
|
||||
&context,
|
||||
&storage_dir,
|
||||
&socket_path,
|
||||
&existing_web_url,
|
||||
"keep-me",
|
||||
);
|
||||
|
||||
let path = fake_gh_path(&context, "ghp_keep_existing");
|
||||
let output = context
|
||||
|
|
@ -382,15 +388,19 @@ fn keep_existing_settings_persists_secrets_without_rewriting_server_target() {
|
|||
.and_then(toml::Value::as_str),
|
||||
Some(existing_web_url.as_str())
|
||||
);
|
||||
let target = parsed
|
||||
.get("cli")
|
||||
.and_then(toml::Value::as_table)
|
||||
.and_then(|cli| cli.get("target"))
|
||||
.and_then(toml::Value::as_table)
|
||||
.expect("cli.target should remain configured");
|
||||
assert_eq!(
|
||||
parsed
|
||||
.get("cli")
|
||||
.and_then(toml::Value::as_table)
|
||||
.and_then(|cli| cli.get("target"))
|
||||
.and_then(toml::Value::as_table)
|
||||
.and_then(|target| target.get("url"))
|
||||
.and_then(toml::Value::as_str),
|
||||
Some(existing_web_url.as_str())
|
||||
target.get("type").and_then(toml::Value::as_str),
|
||||
Some("unix")
|
||||
);
|
||||
assert_eq!(
|
||||
target.get("path").and_then(toml::Value::as_str),
|
||||
socket_path.to_str()
|
||||
);
|
||||
assert!(
|
||||
!settings.contains(&requested_web_url),
|
||||
|
|
@ -874,6 +884,53 @@ mode = "{}"
|
|||
);
|
||||
}
|
||||
|
||||
fn write_unix_install_settings(
|
||||
context: &fabro_test::TestContext,
|
||||
storage_dir: &std::path::Path,
|
||||
socket_path: &std::path::Path,
|
||||
web_url: &str,
|
||||
metadata_mode: &str,
|
||||
) {
|
||||
write_raw_home_settings(
|
||||
context,
|
||||
&format!(
|
||||
r#"
|
||||
_version = 1
|
||||
|
||||
[server.storage]
|
||||
root = "{}"
|
||||
|
||||
[server.api]
|
||||
url = "{}/api/v1"
|
||||
|
||||
[server.web]
|
||||
enabled = true
|
||||
url = "{}"
|
||||
|
||||
[server.auth]
|
||||
methods = ["dev-token"]
|
||||
|
||||
[server.listen]
|
||||
type = "unix"
|
||||
path = "{}"
|
||||
|
||||
[cli.target]
|
||||
type = "unix"
|
||||
path = "{}"
|
||||
|
||||
[project.metadata]
|
||||
mode = "{}"
|
||||
"#,
|
||||
storage_dir.display(),
|
||||
web_url,
|
||||
web_url,
|
||||
socket_path.display(),
|
||||
socket_path.display(),
|
||||
metadata_mode
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
fn login_with_storage_dev_token(
|
||||
context: &fabro_test::TestContext,
|
||||
storage_dir: &std::path::Path,
|
||||
|
|
|
|||
|
|
@ -21,3 +21,6 @@ regex = { workspace = true }
|
|||
serde = { workspace = true }
|
||||
strum.workspace = true
|
||||
thiserror = { workspace = true }
|
||||
|
||||
[dev-dependencies]
|
||||
serde_json = { workspace = true }
|
||||
|
|
|
|||
|
|
@ -1,9 +1,24 @@
|
|||
use serde::{Deserialize, Serialize};
|
||||
use strum::{Display, EnumString, VariantArray};
|
||||
|
||||
/// Fidelity mode controlling how much prior context is provided to LLM
|
||||
/// sessions.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Display, EnumString, VariantArray)]
|
||||
#[derive(
|
||||
Debug,
|
||||
Clone,
|
||||
Copy,
|
||||
Default,
|
||||
PartialEq,
|
||||
Eq,
|
||||
Hash,
|
||||
Display,
|
||||
EnumString,
|
||||
VariantArray,
|
||||
Serialize,
|
||||
Deserialize,
|
||||
)]
|
||||
#[strum(serialize_all = "lowercase")]
|
||||
#[serde(rename_all = "lowercase")]
|
||||
pub enum Fidelity {
|
||||
/// Complete context, no summarization — sessions share a thread.
|
||||
Full,
|
||||
|
|
@ -14,12 +29,15 @@ pub enum Fidelity {
|
|||
Compact,
|
||||
/// Brief textual summary (~600 token target).
|
||||
#[strum(serialize = "summary:low")]
|
||||
#[serde(rename = "summary:low")]
|
||||
SummaryLow,
|
||||
/// Moderate textual summary (~1500 token target).
|
||||
#[strum(serialize = "summary:medium")]
|
||||
#[serde(rename = "summary:medium")]
|
||||
SummaryMedium,
|
||||
/// Detailed per-stage Markdown report.
|
||||
#[strum(serialize = "summary:high")]
|
||||
#[serde(rename = "summary:high")]
|
||||
SummaryHigh,
|
||||
}
|
||||
|
||||
|
|
@ -73,4 +91,14 @@ mod tests {
|
|||
fn fidelity_unknown_mode_errors() {
|
||||
assert!("bogus".parse::<Fidelity>().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn fidelity_serde_matches_strum_display() {
|
||||
for mode in Fidelity::variants() {
|
||||
let json = serde_json::to_value(mode).unwrap();
|
||||
assert_eq!(json, serde_json::Value::String(mode.to_string()));
|
||||
let parsed: Fidelity = serde_json::from_value(json).unwrap();
|
||||
assert_eq!(parsed, *mode);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -497,8 +497,15 @@ fn with_session_lock<T>(root: &Path, f: impl FnOnce() -> T) -> T {
|
|||
// cleanup_session_root can remove_dir_all between the two calls.
|
||||
let deadline = std::time::Instant::now() + SESSION_LOCK_TIMEOUT;
|
||||
let lock_file = loop {
|
||||
std::fs::create_dir_all(root)
|
||||
.unwrap_or_else(|err| panic!("failed to create {}: {err}", root.display()));
|
||||
if let Err(err) = std::fs::create_dir_all(root) {
|
||||
assert!(
|
||||
std::time::Instant::now() < deadline,
|
||||
"failed to create {} after retries: {err}",
|
||||
root.display()
|
||||
);
|
||||
std::thread::sleep(Duration::from_millis(10));
|
||||
continue;
|
||||
}
|
||||
ensure_parent_dir(&lock_path);
|
||||
match File::create(&lock_path) {
|
||||
Ok(f) => break f,
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ mod inert_attribute;
|
|||
mod model_support;
|
||||
mod node_model_known;
|
||||
mod orphan_custom_outcome;
|
||||
mod parallel_branch;
|
||||
mod parallel_branch_inert_attribute;
|
||||
mod prompt_on_llm_nodes;
|
||||
mod random_selection_no_conditions;
|
||||
|
|
|
|||
59
lib/crates/fabro-validate/src/rules/parallel_branch.rs
Normal file
59
lib/crates/fabro-validate/src/rules/parallel_branch.rs
Normal file
|
|
@ -0,0 +1,59 @@
|
|||
use std::collections::BTreeSet;
|
||||
|
||||
use fabro_graphviz::graph::{Edge, Graph};
|
||||
|
||||
pub(super) struct ParallelBranches<'a> {
|
||||
graph: &'a Graph,
|
||||
fork_ids: BTreeSet<&'a str>,
|
||||
}
|
||||
|
||||
impl<'a> ParallelBranches<'a> {
|
||||
pub(super) fn new(graph: &'a Graph) -> Self {
|
||||
let fork_ids = graph
|
||||
.nodes
|
||||
.values()
|
||||
.filter(|node| node.handler_type() == Some("parallel"))
|
||||
.map(|node| node.id.as_str())
|
||||
.collect();
|
||||
Self { graph, fork_ids }
|
||||
}
|
||||
|
||||
pub(super) fn is_empty(&self) -> bool {
|
||||
self.fork_ids.is_empty()
|
||||
}
|
||||
|
||||
pub(super) fn is_fork_edge(&self, edge: &Edge) -> bool {
|
||||
self.fork_ids.contains(edge.from.as_str())
|
||||
}
|
||||
|
||||
pub(super) fn branch_targets(&self) -> BTreeSet<&str> {
|
||||
self.graph
|
||||
.edges
|
||||
.iter()
|
||||
.filter(|edge| self.is_fork_edge(edge))
|
||||
.map(|edge| edge.to.as_str())
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// True when every incoming edge of `node_id` comes from a parallel fork
|
||||
/// (and there is at least one). Such a node only ever runs as a branch.
|
||||
pub(super) fn is_branch_only_node(&self, node_id: &str) -> bool {
|
||||
let incoming = self.graph.incoming_edges(node_id);
|
||||
!incoming.is_empty() && incoming.iter().all(|edge| self.is_fork_edge(edge))
|
||||
}
|
||||
|
||||
/// The sorted, deduplicated fork parents of a branch-only node, or `None`
|
||||
/// when the node has a non-fork entry path (or no entry at all).
|
||||
pub(super) fn branch_only_parents(&self, node_id: &str) -> Option<Vec<String>> {
|
||||
if !self.is_branch_only_node(node_id) {
|
||||
return None;
|
||||
}
|
||||
let parents: BTreeSet<&str> = self
|
||||
.graph
|
||||
.incoming_edges(node_id)
|
||||
.into_iter()
|
||||
.map(|edge| edge.from.as_str())
|
||||
.collect();
|
||||
Some(parents.into_iter().map(String::from).collect())
|
||||
}
|
||||
}
|
||||
|
|
@ -1,18 +1,20 @@
|
|||
use std::collections::BTreeSet;
|
||||
|
||||
use fabro_graphviz::graph::Graph;
|
||||
|
||||
use super::parallel_branch::ParallelBranches;
|
||||
use crate::{Diagnostic, LintRule, Severity};
|
||||
|
||||
pub(super) fn rule() -> Box<dyn LintRule> {
|
||||
Box::new(Rule)
|
||||
}
|
||||
|
||||
/// Attributes that parallel branch execution does not resolve. Branch nodes
|
||||
/// are dispatched with a snapshot of the context taken when the parallel node
|
||||
/// started, so per-branch `fidelity` never changes what a branch sees, and
|
||||
/// per-branch `thread_id` never replaces the thread inherited in that snapshot.
|
||||
const BRANCH_IGNORED_ATTRS: &[&str] = &["fidelity", "thread_id"];
|
||||
/// Attributes that parallel branch execution does not resolve. Only
|
||||
/// `thread_id` is inert on branches (concurrent branches cannot share an LLM
|
||||
/// session); per-branch `fidelity` is honored via pre-rendered preambles.
|
||||
const BRANCH_IGNORED_ATTRS: &[&str] = &["thread_id"];
|
||||
|
||||
const FULL_FIDELITY_MESSAGE: &str = "Parallel branches run at most at summary:high; full is degraded at runtime because branches cannot share a session";
|
||||
|
||||
const THREAD_ID_FIX: &str = "Remove 'thread_id': parallel branches inherit the thread resolved when the parallel node started";
|
||||
|
||||
struct Rule;
|
||||
|
||||
|
|
@ -24,25 +26,31 @@ fn quoted_list(ids: &[String]) -> String {
|
|||
.join(", ")
|
||||
}
|
||||
|
||||
fn fix_message(attr: &str, parallel_ids: &[String]) -> String {
|
||||
match attr {
|
||||
"fidelity" => {
|
||||
if parallel_ids.len() == 1 {
|
||||
format!(
|
||||
"Set fidelity on the parallel node {} (or its incoming edge) to control what every branch sees",
|
||||
quoted_list(parallel_ids),
|
||||
)
|
||||
} else {
|
||||
format!(
|
||||
"Set fidelity on the parallel nodes {} (or their incoming edges) to control what every branch sees",
|
||||
quoted_list(parallel_ids),
|
||||
)
|
||||
}
|
||||
}
|
||||
"thread_id" => format!(
|
||||
"Remove '{attr}': parallel branches inherit the thread resolved when the parallel node started"
|
||||
),
|
||||
_ => format!("Remove '{attr}'"),
|
||||
fn full_fidelity_fix(parallel_ids: &[String]) -> String {
|
||||
let parent = if parallel_ids.len() == 1 {
|
||||
format!("parallel node {}", quoted_list(parallel_ids))
|
||||
} else {
|
||||
format!("parallel nodes {}", quoted_list(parallel_ids))
|
||||
};
|
||||
format!(
|
||||
"Use fidelity=\"summary:high\" or another lower mode on this branch; to reuse a full session before fan-out, set fidelity=\"full\" on {parent} or its incoming edge"
|
||||
)
|
||||
}
|
||||
|
||||
fn full_fidelity_diagnostic(
|
||||
rule_name: &str,
|
||||
node_id: Option<String>,
|
||||
edge: Option<(String, String)>,
|
||||
parallel_ids: &[String],
|
||||
) -> Diagnostic {
|
||||
Diagnostic {
|
||||
rule: rule_name.to_string(),
|
||||
severity: Severity::Warning,
|
||||
message: FULL_FIDELITY_MESSAGE.to_string(),
|
||||
node_id,
|
||||
edge,
|
||||
fix: Some(full_fidelity_fix(parallel_ids)),
|
||||
..Diagnostic::default()
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -52,24 +60,25 @@ impl LintRule for Rule {
|
|||
}
|
||||
|
||||
fn apply(&self, graph: &Graph) -> Vec<Diagnostic> {
|
||||
let parallel_ids: BTreeSet<&str> = graph
|
||||
.nodes
|
||||
.values()
|
||||
.filter(|n| n.handler_type() == Some("parallel"))
|
||||
.map(|n| n.id.as_str())
|
||||
.collect();
|
||||
if parallel_ids.is_empty() {
|
||||
let branches = ParallelBranches::new(graph);
|
||||
if branches.is_empty() {
|
||||
return Vec::new();
|
||||
}
|
||||
|
||||
let mut diagnostics = Vec::new();
|
||||
|
||||
// Branch edges (parallel node -> branch target) carrying an attribute
|
||||
// that branch dispatch never reads.
|
||||
for edge in &graph.edges {
|
||||
if !parallel_ids.contains(edge.from.as_str()) {
|
||||
if !branches.is_fork_edge(edge) {
|
||||
continue;
|
||||
}
|
||||
if edge.fidelity() == Some("full") {
|
||||
diagnostics.push(full_fidelity_diagnostic(
|
||||
self.name(),
|
||||
None,
|
||||
Some((edge.from.clone(), edge.to.clone())),
|
||||
std::slice::from_ref(&edge.from),
|
||||
));
|
||||
}
|
||||
for attr in BRANCH_IGNORED_ATTRS {
|
||||
if !edge.attrs.contains_key(*attr) {
|
||||
continue;
|
||||
|
|
@ -83,42 +92,29 @@ impl LintRule for Rule {
|
|||
),
|
||||
node_id: None,
|
||||
edge: Some((edge.from.clone(), edge.to.clone())),
|
||||
fix: Some(fix_message(attr, std::slice::from_ref(&edge.from))),
|
||||
fix: Some(THREAD_ID_FIX.to_string()),
|
||||
..Diagnostic::default()
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Branch target nodes carrying such an attribute — but only when every
|
||||
// incoming edge comes from a parallel node. A node that is also
|
||||
// reachable through a normal edge resolves the attribute on that path,
|
||||
// so it is not inert there.
|
||||
let branch_targets: BTreeSet<&str> = graph
|
||||
.edges
|
||||
.iter()
|
||||
.filter(|e| parallel_ids.contains(e.from.as_str()))
|
||||
.map(|e| e.to.as_str())
|
||||
.collect();
|
||||
for target in branch_targets {
|
||||
let only_branch_entries = graph
|
||||
.edges
|
||||
.iter()
|
||||
.filter(|e| e.to == target)
|
||||
.all(|e| parallel_ids.contains(e.from.as_str()));
|
||||
if !only_branch_entries {
|
||||
// A node with any normal incoming path still resolves its attributes on
|
||||
// that path, so branch-only diagnostics do not apply to it.
|
||||
for target in branches.branch_targets() {
|
||||
let Some(parents) = branches.branch_only_parents(target) else {
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let Some(node) = graph.nodes.get(target) else {
|
||||
continue;
|
||||
};
|
||||
let parents: Vec<String> = graph
|
||||
.edges
|
||||
.iter()
|
||||
.filter(|e| e.to == target && parallel_ids.contains(e.from.as_str()))
|
||||
.map(|e| e.from.clone())
|
||||
.collect::<BTreeSet<_>>()
|
||||
.into_iter()
|
||||
.collect();
|
||||
if node.fidelity() == Some("full") {
|
||||
diagnostics.push(full_fidelity_diagnostic(
|
||||
self.name(),
|
||||
Some(node.id.clone()),
|
||||
None,
|
||||
&parents,
|
||||
));
|
||||
}
|
||||
for attr in BRANCH_IGNORED_ATTRS {
|
||||
if !node.attrs.contains_key(*attr) {
|
||||
continue;
|
||||
|
|
@ -133,7 +129,7 @@ impl LintRule for Rule {
|
|||
),
|
||||
node_id: Some(node.id.clone()),
|
||||
edge: None,
|
||||
fix: Some(fix_message(attr, &parents)),
|
||||
fix: Some(THREAD_ID_FIX.to_string()),
|
||||
..Diagnostic::default()
|
||||
});
|
||||
}
|
||||
|
|
@ -181,7 +177,7 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn warns_on_fidelity_on_branch_node() {
|
||||
fn accepts_non_full_fidelity_on_branch_node() {
|
||||
let mut g = parallel_graph();
|
||||
g.nodes
|
||||
.get_mut("branch_a")
|
||||
|
|
@ -191,14 +187,73 @@ mod tests {
|
|||
"fidelity".to_string(),
|
||||
AttrValue::String("truncate".to_string()),
|
||||
);
|
||||
|
||||
assert!(Rule.apply(&g).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn warns_when_full_fidelity_on_branch_node_degrades() {
|
||||
let mut g = parallel_graph();
|
||||
g.nodes
|
||||
.get_mut("branch_a")
|
||||
.expect("graph has branch_a")
|
||||
.attrs
|
||||
.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String("full".to_string()),
|
||||
);
|
||||
|
||||
let d = Rule.apply(&g);
|
||||
|
||||
assert_eq!(d.len(), 1);
|
||||
assert_eq!(d[0].severity, Severity::Warning);
|
||||
assert_eq!(d[0].node_id.as_deref(), Some("branch_a"));
|
||||
assert!(d[0].message.contains("'fidelity'"));
|
||||
assert!(d[0].message.contains("full"));
|
||||
assert!(d[0].message.contains("summary:high"));
|
||||
assert!(d[0].fix.as_deref().is_some_and(|f| f.contains("'fork'")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn accepts_every_non_full_fidelity_on_branch_edges() {
|
||||
for fidelity in [
|
||||
"truncate",
|
||||
"compact",
|
||||
"summary:low",
|
||||
"summary:medium",
|
||||
"summary:high",
|
||||
] {
|
||||
let mut g = parallel_graph();
|
||||
g.edges[1].attrs.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String(fidelity.to_string()),
|
||||
);
|
||||
|
||||
assert!(
|
||||
Rule.apply(&g).is_empty(),
|
||||
"{fidelity} should be accepted on a branch edge"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn warns_when_full_fidelity_on_branch_edge_degrades() {
|
||||
let mut g = parallel_graph();
|
||||
g.edges[1].attrs.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String("full".to_string()),
|
||||
);
|
||||
|
||||
let d = Rule.apply(&g);
|
||||
|
||||
assert_eq!(d.len(), 1);
|
||||
assert_eq!(
|
||||
d[0].edge,
|
||||
Some(("fork".to_string(), "branch_a".to_string()))
|
||||
);
|
||||
assert!(d[0].message.contains("full"));
|
||||
assert!(d[0].message.contains("summary:high"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn warns_on_thread_id_on_branch_edge() {
|
||||
let mut g = parallel_graph();
|
||||
|
|
@ -221,6 +276,31 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn warns_on_thread_id_on_branch_only_node() {
|
||||
let mut g = parallel_graph();
|
||||
g.nodes
|
||||
.get_mut("branch_a")
|
||||
.expect("graph has branch_a")
|
||||
.attrs
|
||||
.insert(
|
||||
"thread_id".to_string(),
|
||||
AttrValue::String("impl".to_string()),
|
||||
);
|
||||
|
||||
let d = Rule.apply(&g);
|
||||
|
||||
assert_eq!(d.len(), 1);
|
||||
assert_eq!(d[0].node_id.as_deref(), Some("branch_a"));
|
||||
assert!(d[0].message.contains("'thread_id'"));
|
||||
assert_eq!(
|
||||
d[0].fix.as_deref(),
|
||||
Some(
|
||||
"Remove 'thread_id': parallel branches inherit the thread resolved when the parallel node started"
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn accepts_fidelity_on_the_parallel_node_itself() {
|
||||
let mut g = parallel_graph();
|
||||
|
|
@ -265,11 +345,10 @@ mod tests {
|
|||
.attrs
|
||||
.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String("truncate".to_string()),
|
||||
AttrValue::String("full".to_string()),
|
||||
);
|
||||
let d = Rule.apply(&g);
|
||||
assert_eq!(d.len(), 1);
|
||||
assert!(d[0].message.contains("'fork', 'fork2'"));
|
||||
let fix = d[0].fix.as_deref().expect("diagnostic has a fix");
|
||||
assert!(fix.contains("'fork', 'fork2'"));
|
||||
assert!(fix.contains("parallel nodes"));
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
use fabro_graphviz::graph::Graph;
|
||||
|
||||
use super::parallel_branch::ParallelBranches;
|
||||
use crate::{Diagnostic, LintRule, Severity};
|
||||
|
||||
pub(super) fn rule() -> Box<dyn LintRule> {
|
||||
|
|
@ -20,9 +21,16 @@ impl LintRule for Rule {
|
|||
fn apply(&self, graph: &Graph) -> Vec<Diagnostic> {
|
||||
let mut diagnostics = Vec::new();
|
||||
let graph_default_full = graph.default_fidelity() == Some("full");
|
||||
let branches = ParallelBranches::new(graph);
|
||||
|
||||
// thread_id is inert on parallel branches, where
|
||||
// parallel_branch_inert_attribute already says "remove thread_id" —
|
||||
// advising fidelity="full" there would contradict it.
|
||||
for node in graph.nodes.values() {
|
||||
if node.thread_id().is_some() && node.fidelity() != Some("full") && !graph_default_full
|
||||
if node.thread_id().is_some()
|
||||
&& !branches.is_branch_only_node(&node.id)
|
||||
&& node.fidelity() != Some("full")
|
||||
&& !graph_default_full
|
||||
{
|
||||
diagnostics.push(Diagnostic {
|
||||
rule: self.name().to_string(),
|
||||
|
|
@ -41,7 +49,7 @@ impl LintRule for Rule {
|
|||
}
|
||||
|
||||
for edge in &graph.edges {
|
||||
if edge.thread_id().is_some() {
|
||||
if edge.thread_id().is_some() && !branches.is_fork_edge(edge) {
|
||||
let edge_full = edge.fidelity() == Some("full");
|
||||
let target_full =
|
||||
graph.nodes.get(&edge.to).and_then(|n| n.fidelity()) == Some("full");
|
||||
|
|
@ -82,12 +90,29 @@ impl LintRule for Rule {
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use fabro_graphviz::graph::{AttrValue, Edge, Node};
|
||||
use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node};
|
||||
|
||||
use super::Rule;
|
||||
use crate::rules::test_support::minimal_graph;
|
||||
use crate::{LintRule, Severity};
|
||||
|
||||
fn parallel_graph() -> Graph {
|
||||
let mut g = minimal_graph();
|
||||
let mut fork = Node::new("fork");
|
||||
fork.attrs.insert(
|
||||
"shape".to_string(),
|
||||
AttrValue::String("component".to_string()),
|
||||
);
|
||||
g.nodes.insert("fork".to_string(), fork);
|
||||
g.nodes.insert("branch".to_string(), Node::new("branch"));
|
||||
g.edges = vec![
|
||||
Edge::new("start", "fork"),
|
||||
Edge::new("fork", "branch"),
|
||||
Edge::new("branch", "exit"),
|
||||
];
|
||||
g
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn thread_id_requires_fidelity_full_node_warns() {
|
||||
let mut g = minimal_graph();
|
||||
|
|
@ -206,6 +231,51 @@ mod tests {
|
|||
assert!(d.is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn skips_thread_id_on_parallel_branch_edge() {
|
||||
let mut g = parallel_graph();
|
||||
g.edges[1].attrs.insert(
|
||||
"thread_id".to_string(),
|
||||
AttrValue::String("branch-thread".to_string()),
|
||||
);
|
||||
|
||||
assert!(Rule.apply(&g).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn skips_thread_id_on_branch_only_node() {
|
||||
let mut g = parallel_graph();
|
||||
g.nodes
|
||||
.get_mut("branch")
|
||||
.expect("graph has branch")
|
||||
.attrs
|
||||
.insert(
|
||||
"thread_id".to_string(),
|
||||
AttrValue::String("branch-thread".to_string()),
|
||||
);
|
||||
|
||||
assert!(Rule.apply(&g).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn checks_thread_id_on_branch_node_with_normal_entry() {
|
||||
let mut g = parallel_graph();
|
||||
g.edges.push(Edge::new("start", "branch"));
|
||||
g.nodes
|
||||
.get_mut("branch")
|
||||
.expect("graph has branch")
|
||||
.attrs
|
||||
.insert(
|
||||
"thread_id".to_string(),
|
||||
AttrValue::String("shared-thread".to_string()),
|
||||
);
|
||||
|
||||
let d = Rule.apply(&g);
|
||||
|
||||
assert_eq!(d.len(), 1);
|
||||
assert_eq!(d[0].node_id.as_deref(), Some("branch"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn thread_id_requires_fidelity_full_graph_warns() {
|
||||
let mut g = minimal_graph();
|
||||
|
|
|
|||
|
|
@ -78,11 +78,18 @@ pub fn format_artifact_reference(path: &str) -> String {
|
|||
|
||||
pub fn durable_context_snapshot(context: &Context) -> HashMap<String, Value> {
|
||||
let mut snapshot = context.snapshot();
|
||||
snapshot.remove(context::keys::CURRENT_PREAMBLE);
|
||||
strip_transient_keys(&mut snapshot);
|
||||
normalize_durable_updates(&mut snapshot);
|
||||
snapshot
|
||||
}
|
||||
|
||||
/// Remove runtime-only keys that must never reach durable storage.
|
||||
pub(crate) fn strip_transient_keys(values: &mut HashMap<String, Value>) {
|
||||
for key in context::keys::TRANSIENT_CONTEXT_KEYS {
|
||||
values.remove(*key);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn normalize_durable_updates(updates: &mut HashMap<String, Value>) {
|
||||
for value in updates.values_mut() {
|
||||
normalize_durable_value(value);
|
||||
|
|
@ -96,9 +103,7 @@ pub fn normalize_durable_outcomes(node_outcomes: &mut HashMap<String, Outcome>)
|
|||
}
|
||||
|
||||
pub fn normalize_checkpoint_for_resume(checkpoint: &mut Checkpoint) {
|
||||
checkpoint
|
||||
.context_values
|
||||
.remove(context::keys::CURRENT_PREAMBLE);
|
||||
strip_transient_keys(&mut checkpoint.context_values);
|
||||
normalize_durable_updates(&mut checkpoint.context_values);
|
||||
normalize_durable_outcomes(&mut checkpoint.node_outcomes);
|
||||
}
|
||||
|
|
@ -515,6 +520,59 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn durable_context_snapshot_drops_parallel_branch_preambles() {
|
||||
let context = Context::new();
|
||||
context.set(
|
||||
context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::json!({"branch-a": "runtime only"}),
|
||||
);
|
||||
context.set("response.work", serde_json::json!("durable"));
|
||||
|
||||
let snapshot = durable_context_snapshot(&context);
|
||||
|
||||
assert!(!snapshot.contains_key(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES));
|
||||
assert_eq!(
|
||||
snapshot.get("response.work"),
|
||||
Some(&serde_json::json!("durable"))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_checkpoint_for_resume_drops_parallel_branch_preambles() {
|
||||
let mut checkpoint = crate::records::Checkpoint {
|
||||
timestamp: chrono::Utc::now(),
|
||||
current_node: "work".to_string(),
|
||||
completed_nodes: vec!["work".to_string()],
|
||||
node_retries: HashMap::new(),
|
||||
context_values: HashMap::from([
|
||||
(
|
||||
context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES.to_string(),
|
||||
serde_json::json!({"branch-a": "runtime only"}),
|
||||
),
|
||||
("response.work".to_string(), serde_json::json!("durable")),
|
||||
]),
|
||||
node_outcomes: HashMap::new(),
|
||||
next_node_id: Some("exit".to_string()),
|
||||
git_commit_sha: None,
|
||||
loop_failure_signatures: HashMap::new(),
|
||||
restart_failure_signatures: HashMap::new(),
|
||||
node_visits: HashMap::new(),
|
||||
};
|
||||
|
||||
normalize_checkpoint_for_resume(&mut checkpoint);
|
||||
|
||||
assert!(
|
||||
!checkpoint
|
||||
.context_values
|
||||
.contains_key(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES)
|
||||
);
|
||||
assert_eq!(
|
||||
checkpoint.context_values.get("response.work"),
|
||||
Some(&serde_json::json!("durable"))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn normalize_checkpoint_for_resume_converts_managed_blob_file_refs_and_drops_preamble() {
|
||||
let blob_id = fabro_types::RunBlobId::new(b"managed");
|
||||
|
|
|
|||
|
|
@ -25,6 +25,10 @@ pub mod keys {
|
|||
pub const INTERNAL_PARENT_PREAMBLE: &str = "internal.parent_preamble";
|
||||
pub const INTERNAL_PARALLEL_GROUP_ID: &str = "internal.parallel_group_id";
|
||||
pub const INTERNAL_PARALLEL_BRANCH_ID: &str = "internal.parallel_branch_id";
|
||||
/// Stash of pre-rendered per-branch preambles for a parallel node; see
|
||||
/// [`super::ParallelBranchPreamble`] for the entry shape and the
|
||||
/// producer/consumer contract.
|
||||
pub const INTERNAL_PARALLEL_BRANCH_PREAMBLES: &str = "internal.parallel_branch_preambles";
|
||||
|
||||
// --- current.* keys ---
|
||||
pub const CURRENT_PREAMBLE: &str = "current.preamble";
|
||||
|
|
@ -44,6 +48,10 @@ pub mod keys {
|
|||
pub const PARALLEL_FAN_IN_BEST_OUTCOME: &str = "parallel.fan_in.best_outcome";
|
||||
pub const PARALLEL_FAN_IN_BEST_HEAD_SHA: &str = "parallel.fan_in.best_head_sha";
|
||||
|
||||
/// Runtime-only keys stripped from durable context projections.
|
||||
pub(crate) const TRANSIENT_CONTEXT_KEYS: &[&str] =
|
||||
&[CURRENT_PREAMBLE, INTERNAL_PARALLEL_BRANCH_PREAMBLES];
|
||||
|
||||
// --- Prefix constants (for filtering and dynamic keys) ---
|
||||
pub const GRAPH_PREFIX: &str = "graph.";
|
||||
pub const INTERNAL_PREFIX: &str = "internal.";
|
||||
|
|
@ -135,9 +143,21 @@ pub mod keys {
|
|||
pub use fabro_core::Context;
|
||||
use fabro_graphviz::Fidelity;
|
||||
use fabro_types::{ParallelBranchId, StageId};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::event::StageScope;
|
||||
|
||||
/// One entry of the [`keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES`] stash.
|
||||
///
|
||||
/// The stash is a JSON array indexed by the parallel node's outgoing-edge
|
||||
/// order. `null` entries mean the branch inherits the fork's preamble.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
#[serde(deny_unknown_fields)]
|
||||
pub(crate) struct ParallelBranchPreamble {
|
||||
pub(crate) fidelity: Fidelity,
|
||||
pub(crate) preamble: String,
|
||||
}
|
||||
|
||||
/// Domain-specific typed accessors for workflow context values.
|
||||
pub trait WorkflowContext {
|
||||
fn fidelity(&self) -> Fidelity;
|
||||
|
|
|
|||
|
|
@ -19,6 +19,8 @@ use crate::event::{Emitter, Event, StageScope};
|
|||
use crate::interview_runtime::WorkflowAgentQuestionRuntime;
|
||||
use crate::outcome::{BilledModelUsage, Outcome, OutcomeExt};
|
||||
|
||||
const LAST_FILE_ROUTING_EXTENSIONS: &[&str] = &["json", "md"];
|
||||
|
||||
/// Result from a `CodergenBackend` invocation.
|
||||
#[allow(
|
||||
clippy::large_enum_variant,
|
||||
|
|
@ -169,8 +171,8 @@ pub(crate) async fn validate_agent_output_sources(
|
|||
}
|
||||
|
||||
if let Some(path) = last_file_touched {
|
||||
if let Some(contents) = read_sandbox_file(sandbox, path).await {
|
||||
return structured_output::validate_response_text(schema, &contents);
|
||||
if let Some(routing_json) = read_last_file_routing_json(sandbox, path).await {
|
||||
return structured_output::validate_response_text(schema, &routing_json);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -181,6 +183,22 @@ async fn read_sandbox_file(sandbox: &Arc<dyn Sandbox>, path: &str) -> Option<Str
|
|||
sandbox.read_file_text(path).await.ok()
|
||||
}
|
||||
|
||||
/// Extract the terminal JSON object from the last-touched file when it has an
|
||||
/// eligible extension. Does not check that the object contains routing fields;
|
||||
/// callers validate that.
|
||||
async fn read_last_file_routing_json(sandbox: &Arc<dyn Sandbox>, path: &str) -> Option<String> {
|
||||
let extension = Path::new(path).extension()?.to_str()?;
|
||||
if !LAST_FILE_ROUTING_EXTENSIONS
|
||||
.iter()
|
||||
.any(|allowed| extension.eq_ignore_ascii_case(allowed))
|
||||
{
|
||||
return None;
|
||||
}
|
||||
|
||||
let contents = read_sandbox_file(sandbox, path).await?;
|
||||
structured_output::terminal_json_object(&contents).map(str::to_owned)
|
||||
}
|
||||
|
||||
/// Truncate a string to at most `max_chars` characters (char-boundary safe).
|
||||
pub(crate) fn truncate(s: &str, max_chars: usize) -> &str {
|
||||
if s.len() <= max_chars {
|
||||
|
|
@ -382,7 +400,7 @@ impl Handler for AgentHandler {
|
|||
} else {
|
||||
// 7b. Parse routing directives from response text, falling back to
|
||||
// status.json written by the agent into the sandbox CWD, then to
|
||||
// the last file the agent wrote.
|
||||
// a terminal JSON object in an eligible last-written file.
|
||||
let found_in_response = extract_status_fields(&response_text, &mut outcome);
|
||||
if !found_in_response {
|
||||
let mut found_in_status_json = false;
|
||||
|
|
@ -393,9 +411,10 @@ impl Handler for AgentHandler {
|
|||
}
|
||||
if !found_in_status_json {
|
||||
if let Some(ref path) = last_file_touched {
|
||||
if let Some(contents) = read_sandbox_file(&services.run.sandbox, path).await
|
||||
if let Some(routing_json) =
|
||||
read_last_file_routing_json(&services.run.sandbox, path).await
|
||||
{
|
||||
extract_status_fields(&contents, &mut outcome);
|
||||
extract_status_fields(&routing_json, &mut outcome);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -506,6 +525,67 @@ mod tests {
|
|||
context
|
||||
}
|
||||
|
||||
struct LastFileBackend {
|
||||
path: String,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl CodergenBackend for LastFileBackend {
|
||||
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
|
||||
Ok(CodergenResult::Text {
|
||||
text: "Done writing results.".to_string(),
|
||||
usage: None,
|
||||
files_touched: vec![self.path.clone()],
|
||||
last_file_touched: Some(self.path.clone()),
|
||||
timing: StageTiming::default(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn sandbox_with_file(path: &str, contents: &str) -> (TempDir, Arc<dyn Sandbox>) {
|
||||
let sandbox_dir = TempDir::new().unwrap();
|
||||
std::fs::write(sandbox_dir.path().join(path), contents).unwrap();
|
||||
let sandbox: Arc<dyn Sandbox> = Arc::new(fabro_agent::LocalSandbox::new(
|
||||
sandbox_dir.path().to_path_buf(),
|
||||
));
|
||||
(sandbox_dir, sandbox)
|
||||
}
|
||||
|
||||
async fn execute_with_last_file(path: &str, contents: &str) -> Outcome {
|
||||
let (_sandbox_dir, sandbox) = sandbox_with_file(path, contents);
|
||||
|
||||
let handler = AgentHandler::new(Some(Box::new(LastFileBackend {
|
||||
path: path.to_string(),
|
||||
})));
|
||||
let node = Node::new("step");
|
||||
let context = test_context();
|
||||
let graph = Graph::new("test");
|
||||
let tmp = TempDir::new().unwrap();
|
||||
|
||||
let mut services = EngineServices::test_default();
|
||||
services.run = services.run.with_sandbox(sandbox);
|
||||
|
||||
handler
|
||||
.execute(&node, &context, &graph, tmp.path(), &services)
|
||||
.await
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn validate_routing_with_last_file(
|
||||
path: &str,
|
||||
contents: &str,
|
||||
) -> Result<ValidatedStructuredOutput, StructuredOutputError> {
|
||||
let (_sandbox_dir, sandbox) = sandbox_with_file(path, contents);
|
||||
|
||||
validate_agent_output_sources(
|
||||
&OutputSchemaKind::Routing,
|
||||
"Done writing results.",
|
||||
&sandbox,
|
||||
Some(path),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_simulate() {
|
||||
let handler = AgentHandler::new(None);
|
||||
|
|
@ -731,49 +811,13 @@ mod tests {
|
|||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_extracts_status_from_last_file_touched() {
|
||||
struct LastFileBackend;
|
||||
|
||||
#[async_trait]
|
||||
impl CodergenBackend for LastFileBackend {
|
||||
async fn run(&self, _request: CodergenRunRequest<'_>) -> Result<CodergenResult, Error> {
|
||||
Ok(CodergenResult::Text {
|
||||
text: "Done writing results.".to_string(),
|
||||
usage: None,
|
||||
files_touched: vec!["results.md".to_string()],
|
||||
last_file_touched: Some("results.md".to_string()),
|
||||
timing: StageTiming::default(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
let sandbox_dir = TempDir::new().unwrap();
|
||||
// Write status fields into the file the agent "touched" — no status.json
|
||||
std::fs::write(
|
||||
sandbox_dir.path().join("results.md"),
|
||||
let outcome = execute_with_last_file(
|
||||
"results.md",
|
||||
r#"# Results
|
||||
{"context_updates": {"verified": "true"}}
|
||||
"#,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let handler = AgentHandler::new(Some(Box::new(LastFileBackend)));
|
||||
let node = Node::new("step");
|
||||
let context = test_context();
|
||||
let graph = Graph::new("test");
|
||||
let tmp = TempDir::new().unwrap();
|
||||
|
||||
let mut services = EngineServices::test_default();
|
||||
services.run =
|
||||
services
|
||||
.run
|
||||
.with_sandbox(std::sync::Arc::new(fabro_agent::LocalSandbox::new(
|
||||
sandbox_dir.path().to_path_buf(),
|
||||
)));
|
||||
|
||||
let outcome = handler
|
||||
.execute(&node, &context, &graph, tmp.path(), &services)
|
||||
.await
|
||||
.unwrap();
|
||||
.await;
|
||||
|
||||
assert_eq!(outcome.status, crate::outcome::StageOutcome::Succeeded);
|
||||
assert_eq!(
|
||||
|
|
@ -782,6 +826,33 @@ mod tests {
|
|||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_ignores_nonterminal_status_in_last_markdown_file() {
|
||||
let outcome = execute_with_last_file(
|
||||
"results.md",
|
||||
r#"{"outcome":"failed","failure_reason":"tests failed"}
|
||||
|
||||
All checks passed.
|
||||
"#,
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(outcome.status, crate::outcome::StageOutcome::Succeeded);
|
||||
assert!(outcome.failure.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_ignores_terminal_status_in_disallowed_last_file() {
|
||||
let outcome = execute_with_last_file(
|
||||
"command.rs",
|
||||
r#"{"outcome":"failed","failure_reason":"tests failed"}"#,
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(outcome.status, crate::outcome::StageOutcome::Succeeded);
|
||||
assert!(outcome.failure.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_output_schema_routing_uses_status_json_fallback_when_response_has_no_json()
|
||||
{
|
||||
|
|
@ -819,6 +890,51 @@ mod tests {
|
|||
assert_eq!(outcome.preferred_label.as_deref(), Some("review"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn validated_routing_accepts_terminal_status_in_json_file_case_insensitively() {
|
||||
let validated = validate_routing_with_last_file(
|
||||
"results.JSON",
|
||||
"# Results\n\n{\"preferred_next_label\":\"review\"}\n",
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
validated.value,
|
||||
serde_json::json!({"preferred_next_label": "review"}),
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn validated_routing_ignores_nonterminal_status_in_last_markdown_file() {
|
||||
let error = validate_routing_with_last_file(
|
||||
"results.md",
|
||||
"{\"outcome\":\"failed\",\"failure_reason\":\"tests failed\"}\nAll checks passed.",
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
assert_eq!(
|
||||
error.kind(),
|
||||
structured_output::StructuredOutputErrorKind::NoJsonObject,
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn validated_routing_ignores_terminal_status_in_disallowed_last_file() {
|
||||
let error = validate_routing_with_last_file(
|
||||
"command.rs",
|
||||
r#"{"outcome":"failed","failure_reason":"tests failed"}"#,
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
assert_eq!(
|
||||
error.kind(),
|
||||
structured_output::StructuredOutputErrorKind::NoJsonObject,
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn codergen_handler_output_schema_routing_rejects_malformed_response_before_status_json_fallback()
|
||||
{
|
||||
|
|
|
|||
|
|
@ -24,7 +24,7 @@ use fabro_model::{AgentProfileKind, Catalog, FallbackTarget, ModelRef, ProviderI
|
|||
use fabro_types::settings::run::RunModelControls;
|
||||
use fabro_types::{PermissionLevel, RunId, SessionCapability, StageId, StageTiming};
|
||||
use serde::de::DeserializeOwned;
|
||||
use tokio::sync::Mutex as TokioMutex;
|
||||
use tokio::sync::{Mutex as TokioMutex, mpsc};
|
||||
use tokio::task::JoinHandle;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
|
|
@ -532,16 +532,50 @@ fn emit_agent_tools_available(
|
|||
/// Spawn a task that subscribes to session events and:
|
||||
/// 1. Tracks file changes (write_file/edit_file tool calls) into shared state.
|
||||
/// 2. Forwards non-streaming agent events to the pipeline emitter.
|
||||
///
|
||||
/// The returned handle exposes a per-input barrier. A successful
|
||||
/// `process_input_with_runtime` emits `ProcessingEnd` after all events for
|
||||
/// that input, so waiting for the barrier keeps terminal stage events from
|
||||
/// overtaking queued agent events.
|
||||
struct EventForwarder {
|
||||
processing_end_rx: mpsc::UnboundedReceiver<()>,
|
||||
task: JoinHandle<()>,
|
||||
}
|
||||
|
||||
impl EventForwarder {
|
||||
async fn wait_for_processing_end(&mut self) {
|
||||
if self.processing_end_rx.recv().await.is_none() {
|
||||
tracing::warn!("Agent event forwarder stopped before processing input events");
|
||||
}
|
||||
}
|
||||
|
||||
fn abort(&self) {
|
||||
self.task.abort();
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for EventForwarder {
|
||||
fn drop(&mut self) {
|
||||
self.task.abort();
|
||||
}
|
||||
}
|
||||
|
||||
fn spawn_event_forwarder(
|
||||
session: &Session,
|
||||
node_id: String,
|
||||
scope: StageScope,
|
||||
emitter: Arc<Emitter>,
|
||||
file_tracking: Arc<Mutex<FileTracking>>,
|
||||
) {
|
||||
) -> EventForwarder {
|
||||
let mut rx = session.subscribe();
|
||||
tokio::spawn(async move {
|
||||
let root_session_id = session.id().to_string();
|
||||
let (processing_end_tx, processing_end_rx) = mpsc::unbounded_channel();
|
||||
let task = tokio::spawn(async move {
|
||||
while let Ok(event) = rx.recv().await {
|
||||
let is_root_processing_end = event.session_id == root_session_id
|
||||
&& event.parent_session_id.is_none()
|
||||
&& matches!(&event.event, AgentEvent::ProcessingEnd);
|
||||
|
||||
// Reset watchdog on every event, including streaming deltas
|
||||
emitter.touch();
|
||||
|
||||
|
|
@ -573,8 +607,17 @@ fn spawn_event_forwarder(
|
|||
&scope,
|
||||
);
|
||||
}
|
||||
|
||||
if is_root_processing_end {
|
||||
let _ = processing_end_tx.send(());
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
EventForwarder {
|
||||
processing_end_rx,
|
||||
task,
|
||||
}
|
||||
}
|
||||
|
||||
/// LLM backend that delegates to an `agent` Session per invocation.
|
||||
|
|
@ -1213,7 +1256,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
let stage_scope = StageScope::for_handler(context, &node.id);
|
||||
|
||||
// Subscribe to session events: forward to pipeline emitter + track files.
|
||||
spawn_event_forwarder(
|
||||
let mut event_forwarder = spawn_event_forwarder(
|
||||
&session,
|
||||
node.id.clone(),
|
||||
stage_scope.clone(),
|
||||
|
|
@ -1284,6 +1327,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
inference_duration = inference_duration.saturating_add(timing.inference);
|
||||
tool_duration = tool_duration.saturating_add(timing.tool);
|
||||
if process_result.is_ok() {
|
||||
event_forwarder.wait_for_processing_end().await;
|
||||
total_usage += session.last_input_usage();
|
||||
UsdMicros::accumulate(&mut total_cost, session.last_input_cost());
|
||||
}
|
||||
|
|
@ -1315,6 +1359,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
let mut succeeded = false;
|
||||
|
||||
bridge.abort();
|
||||
event_forwarder.abort();
|
||||
discard_session(&mut session, &mut lease, emitter);
|
||||
|
||||
for (index, target) in self.fallback_chain.iter().enumerate() {
|
||||
|
|
@ -1368,7 +1413,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
bridge.replace(cancel_token.clone(), &session);
|
||||
|
||||
// Re-subscribe to forward events + track files from the new session
|
||||
spawn_event_forwarder(
|
||||
event_forwarder = spawn_event_forwarder(
|
||||
&session,
|
||||
node.id.clone(),
|
||||
stage_scope.clone(),
|
||||
|
|
@ -1420,6 +1465,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
tool_duration = tool_duration.saturating_add(timing.tool);
|
||||
match process_result {
|
||||
Ok(()) => {
|
||||
event_forwarder.wait_for_processing_end().await;
|
||||
total_usage += session.last_input_usage();
|
||||
UsdMicros::accumulate(&mut total_cost, session.last_input_cost());
|
||||
succeeded = true;
|
||||
|
|
@ -1492,6 +1538,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
tool_duration = tool_duration.saturating_add(timing.tool);
|
||||
match repair_result {
|
||||
Ok(()) => {
|
||||
event_forwarder.wait_for_processing_end().await;
|
||||
total_usage += session.last_input_usage();
|
||||
UsdMicros::accumulate(&mut total_cost, session.last_input_cost());
|
||||
repair_attempts += 1;
|
||||
|
|
@ -1534,6 +1581,7 @@ impl CodergenBackend for AgentApiBackend {
|
|||
|
||||
// Collect files_touched from the shared tracking state.
|
||||
let (files_touched, last_file_touched) = file_tracking_snapshot(&file_tracking);
|
||||
drop(event_forwarder);
|
||||
|
||||
if let Some(lease) = lease.take() {
|
||||
lease.release();
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ use fabro_types::{ParallelBranchId, RunId, StageId};
|
|||
use tokio::sync::Semaphore;
|
||||
|
||||
use super::{EngineServices, Handler};
|
||||
use crate::context::{Context, WorkflowContext, keys};
|
||||
use crate::context::{Context, ParallelBranchPreamble, WorkflowContext, keys};
|
||||
use crate::error::Error;
|
||||
use crate::event::{Event, RunNoticeCode, RunNoticeLevel, StageScope};
|
||||
use crate::git::sanitize_ref_component;
|
||||
|
|
@ -56,6 +56,31 @@ struct BranchResult {
|
|||
worktree_path: Option<PathBuf>,
|
||||
}
|
||||
|
||||
/// Parse the per-branch preamble stash produced by `FidelityLifecycle`.
|
||||
///
|
||||
/// Outer `None` means the stash is absent, malformed, or has the wrong branch
|
||||
/// count — every branch then inherits the fork context (legacy behavior).
|
||||
/// Inner `None` means that single branch inherits.
|
||||
fn parse_branch_preambles(
|
||||
value: Option<serde_json::Value>,
|
||||
branch_count: usize,
|
||||
) -> Option<Vec<Option<ParallelBranchPreamble>>> {
|
||||
let serde_json::Value::Array(entries) = value? else {
|
||||
return None;
|
||||
};
|
||||
if entries.len() != branch_count {
|
||||
return None;
|
||||
}
|
||||
|
||||
entries
|
||||
.into_iter()
|
||||
.map(|entry| match entry {
|
||||
serde_json::Value::Null => Some(None),
|
||||
entry => serde_json::from_value(entry).ok().map(Some),
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl Handler for ParallelHandler {
|
||||
async fn simulate(
|
||||
|
|
@ -220,6 +245,17 @@ impl Handler for ParallelHandler {
|
|||
None
|
||||
};
|
||||
|
||||
let branch_preambles = parse_branch_preambles(
|
||||
context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
branches.len(),
|
||||
);
|
||||
// Clear the stash before forking so branch contexts never carry the
|
||||
// outer array — a nested parallel branch target must not misread it as
|
||||
// its own. The write-back diff also clears it on the run state.
|
||||
context.set(
|
||||
keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::Value::Null,
|
||||
);
|
||||
let mut branch_setups: Vec<BranchSetup> = Vec::new();
|
||||
for (branch_index, edge) in branches.iter().enumerate() {
|
||||
let target_id = edge.to.clone();
|
||||
|
|
@ -236,6 +272,20 @@ impl Handler for ParallelHandler {
|
|||
keys::INTERNAL_PARALLEL_BRANCH_ID,
|
||||
serde_json::Value::String(parallel_branch_id.to_string()),
|
||||
);
|
||||
if let Some(entry) = branch_preambles
|
||||
.as_ref()
|
||||
.and_then(|entries| entries.get(branch_index))
|
||||
.and_then(Option::as_ref)
|
||||
{
|
||||
branch_context.set(
|
||||
keys::CURRENT_PREAMBLE,
|
||||
serde_json::Value::String(entry.preamble.clone()),
|
||||
);
|
||||
branch_context.set(
|
||||
keys::INTERNAL_FIDELITY,
|
||||
serde_json::Value::String(entry.fidelity.to_string()),
|
||||
);
|
||||
}
|
||||
|
||||
let (branch_sandbox, worktree_path): (Arc<dyn Sandbox>, Option<PathBuf>) = if let (
|
||||
Some(ref gs),
|
||||
|
|
@ -693,7 +743,7 @@ fn parallel_branch_commit_cmd(
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
|
||||
use fabro_graphviz::graph::{AttrValue, Edge};
|
||||
|
|
@ -756,6 +806,185 @@ mod tests {
|
|||
context
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
struct BranchContextCapture {
|
||||
node_id: String,
|
||||
preamble: String,
|
||||
fidelity: String,
|
||||
stash: Option<serde_json::Value>,
|
||||
}
|
||||
|
||||
struct BranchContextRecordingHandler {
|
||||
captures: Arc<Mutex<Vec<BranchContextCapture>>>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl Handler for BranchContextRecordingHandler {
|
||||
async fn execute(
|
||||
&self,
|
||||
node: &Node,
|
||||
context: &Context,
|
||||
_graph: &Graph,
|
||||
_run_dir: &Path,
|
||||
_services: &EngineServices,
|
||||
) -> Result<Outcome, Error> {
|
||||
self.captures.lock().unwrap().push(BranchContextCapture {
|
||||
node_id: node.id.clone(),
|
||||
preamble: context.preamble(),
|
||||
fidelity: context.fidelity().to_string(),
|
||||
stash: context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
});
|
||||
Ok(Outcome::success())
|
||||
}
|
||||
}
|
||||
|
||||
async fn execute_with_branch_stash(
|
||||
stash: Option<serde_json::Value>,
|
||||
duplicate_target: bool,
|
||||
) -> (Context, Vec<BranchContextCapture>) {
|
||||
let captures = Arc::new(Mutex::new(Vec::new()));
|
||||
let recorder = BranchContextRecordingHandler {
|
||||
captures: Arc::clone(&captures),
|
||||
};
|
||||
let mut registry = super::super::HandlerRegistry::new(Box::new(recorder));
|
||||
registry.register(
|
||||
"record",
|
||||
Box::new(BranchContextRecordingHandler {
|
||||
captures: Arc::clone(&captures),
|
||||
}),
|
||||
);
|
||||
let mut services = EngineServices::test_default();
|
||||
services.registry = Arc::new(registry);
|
||||
|
||||
let mut node = Node::new("par");
|
||||
node.attrs.insert(
|
||||
"shape".to_string(),
|
||||
AttrValue::String("component".to_string()),
|
||||
);
|
||||
let mut branch_a = Node::new("branch_a");
|
||||
branch_a
|
||||
.attrs
|
||||
.insert("type".to_string(), AttrValue::String("record".to_string()));
|
||||
let mut branch_b = Node::new("branch_b");
|
||||
branch_b
|
||||
.attrs
|
||||
.insert("type".to_string(), AttrValue::String("record".to_string()));
|
||||
|
||||
let mut graph = Graph::new("test");
|
||||
graph.nodes.insert(node.id.clone(), node.clone());
|
||||
graph.nodes.insert(branch_a.id.clone(), branch_a);
|
||||
graph.nodes.insert(branch_b.id.clone(), branch_b);
|
||||
graph.edges.push(Edge::new("par", "branch_a"));
|
||||
graph.edges.push(Edge::new(
|
||||
"par",
|
||||
if duplicate_target {
|
||||
"branch_a"
|
||||
} else {
|
||||
"branch_b"
|
||||
},
|
||||
));
|
||||
|
||||
let context = test_context();
|
||||
context.set(keys::CURRENT_PREAMBLE, serde_json::json!("fork preamble"));
|
||||
context.set(keys::INTERNAL_FIDELITY, serde_json::json!("compact"));
|
||||
if let Some(stash) = stash {
|
||||
context.set(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES, stash);
|
||||
}
|
||||
|
||||
let run_dir = tempfile::tempdir().unwrap();
|
||||
ParallelHandler
|
||||
.execute(&node, &context, &graph, run_dir.path(), &services)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let captures = captures.lock().unwrap().clone();
|
||||
(context, captures)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_handler_applies_indexed_branch_preambles_and_clears_stash() {
|
||||
let stash = serde_json::json!([
|
||||
{"fidelity": "truncate", "preamble": "branch zero"},
|
||||
{"fidelity": "summary:high", "preamble": "branch one"}
|
||||
]);
|
||||
|
||||
let (context, mut captures) = execute_with_branch_stash(Some(stash), false).await;
|
||||
captures.sort_by(|left, right| left.node_id.cmp(&right.node_id));
|
||||
|
||||
assert_eq!(captures.len(), 2);
|
||||
assert_eq!(captures[0].node_id, "branch_a");
|
||||
assert_eq!(captures[0].preamble, "branch zero");
|
||||
assert_eq!(captures[0].fidelity, "truncate");
|
||||
assert_eq!(captures[0].stash, Some(serde_json::Value::Null));
|
||||
assert_eq!(captures[1].node_id, "branch_b");
|
||||
assert_eq!(captures[1].preamble, "branch one");
|
||||
assert_eq!(captures[1].fidelity, "summary:high");
|
||||
assert_eq!(captures[1].stash, Some(serde_json::Value::Null));
|
||||
assert_eq!(
|
||||
context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
Some(serde_json::Value::Null)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_handler_uses_edge_index_for_duplicate_targets() {
|
||||
let stash = serde_json::json!([
|
||||
{"fidelity": "truncate", "preamble": "first edge"},
|
||||
{"fidelity": "summary:low", "preamble": "second edge"}
|
||||
]);
|
||||
|
||||
let (_context, captures) = execute_with_branch_stash(Some(stash), true).await;
|
||||
let observed = captures
|
||||
.iter()
|
||||
.map(|capture| (capture.preamble.as_str(), capture.fidelity.as_str()))
|
||||
.collect::<std::collections::HashSet<_>>();
|
||||
|
||||
assert_eq!(observed.len(), 2);
|
||||
assert!(observed.contains(&("first edge", "truncate")));
|
||||
assert!(observed.contains(&("second edge", "summary:low")));
|
||||
assert!(
|
||||
captures
|
||||
.iter()
|
||||
.all(|capture| capture.stash == Some(serde_json::Value::Null))
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_handler_legacy_stashes_inherit_fork_context() {
|
||||
for stash in [
|
||||
None,
|
||||
Some(serde_json::Value::Null),
|
||||
Some(serde_json::json!({
|
||||
"fidelity": "truncate",
|
||||
"preamble": "not an array"
|
||||
})),
|
||||
Some(serde_json::json!([
|
||||
{"fidelity": "truncate", "preamble": "wrong length"}
|
||||
])),
|
||||
Some(serde_json::json!([
|
||||
{"fidelity": "truncate"},
|
||||
null
|
||||
])),
|
||||
Some(serde_json::json!([
|
||||
{"fidelity": "not-a-fidelity", "preamble": "malformed fidelity"},
|
||||
null
|
||||
])),
|
||||
] {
|
||||
let (context, captures) = execute_with_branch_stash(stash, false).await;
|
||||
|
||||
assert_eq!(captures.len(), 2);
|
||||
assert!(captures.iter().all(|capture| {
|
||||
capture.preamble == "fork preamble"
|
||||
&& capture.fidelity == "compact"
|
||||
&& capture.stash == Some(serde_json::Value::Null)
|
||||
}));
|
||||
assert_eq!(
|
||||
context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
Some(serde_json::Value::Null)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_handler_no_branches() {
|
||||
let services = make_services();
|
||||
|
|
|
|||
|
|
@ -263,6 +263,15 @@ fn find_json_objects(text: &str) -> Vec<&str> {
|
|||
results
|
||||
}
|
||||
|
||||
/// Return the outermost balanced JSON object that ends the text, ignoring
|
||||
/// trailing whitespace.
|
||||
pub(crate) fn terminal_json_object(text: &str) -> Option<&str> {
|
||||
let trimmed = text.trim_end();
|
||||
find_json_objects(trimmed)
|
||||
.into_iter()
|
||||
.find(|candidate| trimmed.ends_with(candidate))
|
||||
}
|
||||
|
||||
pub(crate) fn extract_status_fields(text: &str, outcome: &mut Outcome) -> bool {
|
||||
let candidates = find_json_objects(text);
|
||||
|
||||
|
|
@ -492,6 +501,31 @@ mod tests {
|
|||
assert!(error.messages()[0].contains("recognized routing field"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn terminal_json_object_accepts_final_object_after_prose() {
|
||||
let object =
|
||||
terminal_json_object("# Results\n\n{\"context_updates\":{\"verified\":true}}\n\n");
|
||||
|
||||
assert_eq!(object, Some(r#"{"context_updates":{"verified":true}}"#),);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn terminal_json_object_returns_outermost_nested_object() {
|
||||
let object = terminal_json_object(r#"Results: {"context_updates":{"verified":true}}"#);
|
||||
|
||||
assert_eq!(object, Some(r#"{"context_updates":{"verified":true}}"#),);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn terminal_json_object_rejects_object_followed_by_content() {
|
||||
assert_eq!(
|
||||
terminal_json_object(
|
||||
"{\"outcome\":\"failed\",\"failure_reason\":\"tests failed\"}\nMore details",
|
||||
),
|
||||
None,
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn routing_json_with_wrong_field_type_is_invalid() {
|
||||
let error =
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ use fabro_types::{Principal, RunId, StageTiming};
|
|||
|
||||
use super::circuit_breaker::CircuitBreakerLifecycle;
|
||||
use super::git::GitCheckpointResult;
|
||||
use crate::context::WorkflowContext;
|
||||
use crate::context::{Context, WorkflowContext};
|
||||
use crate::event::{Emitter, Event, StageScope};
|
||||
use crate::graph::{WorkflowGraph, WorkflowNode};
|
||||
use crate::outcome::{BilledModelUsage, FailureCategory, FailureDetail, Outcome, StageOutcome};
|
||||
|
|
@ -92,6 +92,16 @@ fn response_from_outcome(node_id: &str, outcome: &Outcome) -> Option<String> {
|
|||
.and_then(|value| value.as_str().map(ToOwned::to_owned))
|
||||
}
|
||||
|
||||
/// Context values for `StageCompleted` events. Unlike
|
||||
/// `artifact::strip_transient_keys`, this keeps `CURRENT_PREAMBLE` — stage
|
||||
/// events have always included the active preamble — and drops only the
|
||||
/// parallel stash, which can embed every branch's rendered preamble.
|
||||
fn stage_context_values(workflow_context: &Context) -> Option<BTreeMap<String, serde_json::Value>> {
|
||||
let mut snapshot = workflow_context.snapshot();
|
||||
snapshot.remove(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES);
|
||||
(!snapshot.is_empty()).then(|| snapshot.into_iter().collect())
|
||||
}
|
||||
|
||||
pub(super) fn stage_visit(state: &WfRunState, node_id: &str) -> u32 {
|
||||
let visits = state.node_visits.get(node_id).copied().unwrap_or(1);
|
||||
u32::try_from(visits).unwrap_or(u32::MAX)
|
||||
|
|
@ -318,11 +328,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
|
|||
.collect::<BTreeMap<_, _>>()
|
||||
}),
|
||||
jump_to_node: outcome.jump_to_node.clone(),
|
||||
context_values: {
|
||||
let snapshot = state.context.snapshot();
|
||||
(!snapshot.is_empty())
|
||||
.then(|| snapshot.into_iter().collect::<BTreeMap<_, _>>())
|
||||
},
|
||||
context_values: stage_context_values(&state.context),
|
||||
node_visits: (!state.node_visits.is_empty()).then(|| {
|
||||
state
|
||||
.node_visits
|
||||
|
|
@ -446,3 +452,26 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
|
|||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn stage_context_values_drops_parallel_branch_preambles() {
|
||||
let workflow_context = Context::new();
|
||||
workflow_context.set(
|
||||
context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::json!([{"fidelity": "summary:high", "preamble": "runtime only"}]),
|
||||
);
|
||||
workflow_context.set("response.work", serde_json::json!("durable"));
|
||||
|
||||
let values = stage_context_values(&workflow_context).expect("snapshot should not be empty");
|
||||
|
||||
assert!(!values.contains_key(context::keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES));
|
||||
assert_eq!(
|
||||
values.get("response.work"),
|
||||
Some(&serde_json::json!("durable"))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,3 +1,4 @@
|
|||
use std::collections::HashMap;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
|
|
@ -10,10 +11,10 @@ use fabro_core::state::ExecutionState;
|
|||
use fabro_graphviz::graph::types::{Edge as GvEdge, Graph as GvGraph, Node as GvNode};
|
||||
|
||||
use crate::artifact;
|
||||
use crate::context::keys;
|
||||
use crate::context::{Context, ParallelBranchPreamble, keys};
|
||||
use crate::graph::{WorkflowGraph, WorkflowNode};
|
||||
use crate::handler::llm::preamble::build_preamble;
|
||||
use crate::outcome::BilledModelUsage;
|
||||
use crate::outcome::{BilledModelUsage, Outcome};
|
||||
use crate::runtime_store::RunStoreHandle;
|
||||
|
||||
type WfRunState = ExecutionState<Option<BilledModelUsage>>;
|
||||
|
|
@ -61,6 +62,65 @@ impl FidelityLifecycle {
|
|||
"fidelity mutex should not be poisoned: no code panics while holding this lock",
|
||||
) = flag;
|
||||
}
|
||||
|
||||
/// Render the per-branch preamble stash for a parallel node, indexed by
|
||||
/// outgoing-edge order (the same order `ParallelHandler` fans out in).
|
||||
/// `Null` entries inherit the fork's preamble.
|
||||
fn build_parallel_branch_preambles(
|
||||
&self,
|
||||
node_id: &str,
|
||||
fork_fidelity: keys::Fidelity,
|
||||
resolved_context: &Context,
|
||||
resolved_outcomes: &HashMap<String, Outcome>,
|
||||
completed_nodes: &[String],
|
||||
) -> Vec<serde_json::Value> {
|
||||
let edges = self.graph.outgoing_edges(node_id);
|
||||
let mut preambles: Vec<serde_json::Value> = Vec::with_capacity(edges.len());
|
||||
let mut rendered: HashMap<keys::Fidelity, usize> = HashMap::new();
|
||||
|
||||
for (branch_index, edge) in edges.into_iter().enumerate() {
|
||||
let Some(target_node) = self.graph.nodes.get(&edge.to) else {
|
||||
preambles.push(serde_json::Value::Null);
|
||||
continue;
|
||||
};
|
||||
let resolution = resolve_parallel_branch_fidelity(edge, target_node, fork_fidelity);
|
||||
if resolution.requested == Some(keys::Fidelity::Full) {
|
||||
tracing::warn!(
|
||||
parallel_node = %node_id,
|
||||
branch = %edge.to,
|
||||
branch_index,
|
||||
effective_fidelity = %keys::Fidelity::Full.degraded(),
|
||||
"Parallel branch fidelity degraded from full"
|
||||
);
|
||||
}
|
||||
let Some(branch_fidelity) = resolution.effective else {
|
||||
preambles.push(serde_json::Value::Null);
|
||||
continue;
|
||||
};
|
||||
if let Some(&rendered_index) = rendered.get(&branch_fidelity) {
|
||||
preambles.push(preambles[rendered_index].clone());
|
||||
continue;
|
||||
}
|
||||
|
||||
let entry = ParallelBranchPreamble {
|
||||
fidelity: branch_fidelity,
|
||||
preamble: build_preamble(
|
||||
branch_fidelity,
|
||||
resolved_context,
|
||||
&self.graph,
|
||||
completed_nodes,
|
||||
resolved_outcomes,
|
||||
),
|
||||
};
|
||||
rendered.insert(branch_fidelity, preambles.len());
|
||||
preambles.push(
|
||||
serde_json::to_value(entry)
|
||||
.expect("ParallelBranchPreamble serialization cannot fail"),
|
||||
);
|
||||
}
|
||||
|
||||
preambles
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
|
|
@ -78,6 +138,11 @@ impl RunLifecycle<WorkflowGraph> for FidelityLifecycle {
|
|||
node: &WorkflowNode,
|
||||
state: &WfRunState,
|
||||
) -> CoreResult<WfNodeDecision> {
|
||||
state.context.set(
|
||||
keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::Value::Null,
|
||||
);
|
||||
|
||||
let incoming = self
|
||||
.incoming_edge_data
|
||||
.lock()
|
||||
|
|
@ -138,7 +203,23 @@ impl RunLifecycle<WorkflowGraph> for FidelityLifecycle {
|
|||
.context
|
||||
.set(keys::CURRENT_PREAMBLE, serde_json::json!(preamble));
|
||||
|
||||
// 5. Thread ID resolution via resolve_thread_id: edge → node → graph default →
|
||||
// 5. Parallel nodes: pre-render per-branch preambles into the stash that
|
||||
// ParallelHandler consumes at fan-out.
|
||||
if gv_node.handler_type() == Some("parallel") {
|
||||
let branch_preambles = self.build_parallel_branch_preambles(
|
||||
node.id(),
|
||||
fidelity,
|
||||
&resolved_context,
|
||||
&resolved_outcomes,
|
||||
&state.completed_nodes,
|
||||
);
|
||||
state.context.set(
|
||||
keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::Value::Array(branch_preambles),
|
||||
);
|
||||
}
|
||||
|
||||
// 6. Thread ID resolution via resolve_thread_id: edge → node → graph default →
|
||||
// class → previous
|
||||
let thread_id = resolve_thread_id(
|
||||
incoming_edge_ref,
|
||||
|
|
@ -147,13 +228,13 @@ impl RunLifecycle<WorkflowGraph> for FidelityLifecycle {
|
|||
state.previous_node_id.as_deref(),
|
||||
);
|
||||
|
||||
// 6. Set thread.{tid}.current_node
|
||||
// 7. Set thread.{tid}.current_node
|
||||
if let Some(ref tid) = thread_id {
|
||||
let key = keys::thread_current_node_key(tid);
|
||||
state.context.set(key, serde_json::json!(node.id()));
|
||||
}
|
||||
|
||||
// 7. Set INTERNAL_THREAD_ID (or null)
|
||||
// 8. Set INTERNAL_THREAD_ID (or null)
|
||||
match thread_id {
|
||||
Some(tid) => {
|
||||
state
|
||||
|
|
@ -167,7 +248,7 @@ impl RunLifecycle<WorkflowGraph> for FidelityLifecycle {
|
|||
}
|
||||
}
|
||||
|
||||
// 8. Set INTERNAL_NODE_VISIT_COUNT and CURRENT_NODE
|
||||
// 9. Set INTERNAL_NODE_VISIT_COUNT and CURRENT_NODE
|
||||
let visits = state.node_visits.get(node.id()).copied().unwrap_or(1);
|
||||
state
|
||||
.context
|
||||
|
|
@ -198,6 +279,53 @@ impl RunLifecycle<WorkflowGraph> for FidelityLifecycle {
|
|||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
struct ParallelBranchFidelityResolution {
|
||||
/// The explicit fidelity requested on the edge or node, pre-degradation.
|
||||
requested: Option<keys::Fidelity>,
|
||||
/// The fidelity to render an entry for; `None` inherits the fork preamble.
|
||||
effective: Option<keys::Fidelity>,
|
||||
}
|
||||
|
||||
/// Resolve explicit branch fidelity with edge-over-node precedence.
|
||||
///
|
||||
/// Branches with no explicit fidelity inherit the parallel node's preamble.
|
||||
/// Explicit full fidelity is degraded because concurrent branches cannot share
|
||||
/// an LLM session. An effective fidelity equal to the parallel node also
|
||||
/// inherits, avoiding a redundant preamble render.
|
||||
fn resolve_parallel_branch_fidelity(
|
||||
edge: &GvEdge,
|
||||
target_node: &GvNode,
|
||||
parallel_fidelity: keys::Fidelity,
|
||||
) -> ParallelBranchFidelityResolution {
|
||||
let requested = explicit_fidelity(Some(edge), target_node).map(|(fidelity, _)| fidelity);
|
||||
let effective = requested
|
||||
.map(keys::Fidelity::degraded)
|
||||
.filter(|fidelity| *fidelity != parallel_fidelity);
|
||||
|
||||
ParallelBranchFidelityResolution {
|
||||
requested,
|
||||
effective,
|
||||
}
|
||||
}
|
||||
|
||||
/// Explicit fidelity from the incoming edge attribute, else the node
|
||||
/// attribute, with the winning source labeled for logging.
|
||||
fn explicit_fidelity(
|
||||
incoming_edge: Option<&GvEdge>,
|
||||
node: &GvNode,
|
||||
) -> Option<(keys::Fidelity, &'static str)> {
|
||||
incoming_edge
|
||||
.and_then(|e| e.fidelity())
|
||||
.and_then(|s| s.parse().ok())
|
||||
.map(|f| (f, "edge"))
|
||||
.or_else(|| {
|
||||
node.fidelity()
|
||||
.and_then(|s| s.parse().ok())
|
||||
.map(|f| (f, "node"))
|
||||
})
|
||||
}
|
||||
|
||||
/// Resolve the context fidelity for a node, following the precedence:
|
||||
/// 1. Incoming edge `fidelity` attribute
|
||||
/// 2. Target node `fidelity` attribute
|
||||
|
|
@ -208,13 +336,8 @@ fn resolve_fidelity(
|
|||
node: &GvNode,
|
||||
graph: &GvGraph,
|
||||
) -> keys::Fidelity {
|
||||
let (resolved, source) = if let Some(f) = incoming_edge
|
||||
.and_then(|e| e.fidelity())
|
||||
.and_then(|s| s.parse().ok())
|
||||
{
|
||||
(f, "edge")
|
||||
} else if let Some(f) = node.fidelity().and_then(|s| s.parse().ok()) {
|
||||
(f, "node")
|
||||
let (resolved, source) = if let Some((f, source)) = explicit_fidelity(incoming_edge, node) {
|
||||
(f, source)
|
||||
} else if let Some(f) = graph.default_fidelity().and_then(|s| s.parse().ok()) {
|
||||
(f, "graph")
|
||||
} else {
|
||||
|
|
@ -263,11 +386,214 @@ fn resolve_thread_id(
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::Path;
|
||||
use std::time::Duration;
|
||||
|
||||
use fabro_core::graph::Graph as CoreGraph;
|
||||
use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node};
|
||||
use fabro_store::Database;
|
||||
use fabro_types::fixtures;
|
||||
use object_store::memory::InMemory;
|
||||
|
||||
use super::*;
|
||||
use crate::context::WorkflowContext;
|
||||
use crate::context::keys::Fidelity;
|
||||
|
||||
fn str_attr(value: &str) -> AttrValue {
|
||||
AttrValue::String(value.to_string())
|
||||
}
|
||||
|
||||
fn parallel_workflow_graph(
|
||||
fork_fidelity: Option<&str>,
|
||||
branch_a_fidelity: Option<&str>,
|
||||
) -> WorkflowGraph {
|
||||
let mut graph = Graph::new("parallel-fidelity");
|
||||
let mut start = Node::new("start");
|
||||
start
|
||||
.attrs
|
||||
.insert("shape".to_string(), str_attr("Mdiamond"));
|
||||
let mut fork = Node::new("fork");
|
||||
fork.attrs
|
||||
.insert("shape".to_string(), str_attr("component"));
|
||||
if let Some(fidelity) = fork_fidelity {
|
||||
fork.attrs
|
||||
.insert("fidelity".to_string(), str_attr(fidelity));
|
||||
}
|
||||
let mut branch_a = Node::new("branch_a");
|
||||
if let Some(fidelity) = branch_a_fidelity {
|
||||
branch_a
|
||||
.attrs
|
||||
.insert("fidelity".to_string(), str_attr(fidelity));
|
||||
}
|
||||
let branch_b = Node::new("branch_b");
|
||||
let mut work = Node::new("work");
|
||||
work.attrs.insert("shape".to_string(), str_attr("box"));
|
||||
|
||||
graph.nodes.insert(start.id.clone(), start);
|
||||
graph.nodes.insert(fork.id.clone(), fork);
|
||||
graph.nodes.insert(branch_a.id.clone(), branch_a);
|
||||
graph.nodes.insert(branch_b.id.clone(), branch_b);
|
||||
graph.nodes.insert(work.id.clone(), work);
|
||||
graph.edges.push(Edge::new("start", "fork"));
|
||||
graph.edges.push(Edge::new("fork", "branch_a"));
|
||||
graph.edges.push(Edge::new("fork", "branch_b"));
|
||||
|
||||
WorkflowGraph(Arc::new(graph))
|
||||
}
|
||||
|
||||
async fn test_lifecycle(graph: &WorkflowGraph, run_dir: &Path) -> FidelityLifecycle {
|
||||
let store = Arc::new(Database::new(
|
||||
Arc::new(InMemory::new()),
|
||||
"",
|
||||
Duration::from_millis(1),
|
||||
None,
|
||||
));
|
||||
let run_store = store.create_run(&fixtures::RUN_1).await.unwrap();
|
||||
let sandbox: Arc<dyn Sandbox> =
|
||||
Arc::new(fabro_agent::LocalSandbox::new(run_dir.to_path_buf()));
|
||||
FidelityLifecycle::new(
|
||||
graph.0.clone(),
|
||||
sandbox,
|
||||
RunStoreHandle::local(run_store),
|
||||
run_dir.to_path_buf(),
|
||||
)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_branch_fidelity_edge_overrides_node() {
|
||||
let mut node = Node::new("branch");
|
||||
node.attrs
|
||||
.insert("fidelity".to_string(), str_attr("compact"));
|
||||
let mut edge = Edge::new("fork", "branch");
|
||||
edge.attrs
|
||||
.insert("fidelity".to_string(), str_attr("truncate"));
|
||||
|
||||
let resolved = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh);
|
||||
|
||||
assert_eq!(resolved.requested, Some(Fidelity::Truncate));
|
||||
assert_eq!(resolved.effective, Some(Fidelity::Truncate));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_branch_fidelity_without_attribute_inherits() {
|
||||
let node = Node::new("branch");
|
||||
let edge = Edge::new("fork", "branch");
|
||||
|
||||
let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Compact);
|
||||
|
||||
assert_eq!(resolution.requested, None);
|
||||
assert_eq!(resolution.effective, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_branch_full_fidelity_degrades_to_summary_high() {
|
||||
let mut node = Node::new("branch");
|
||||
node.attrs.insert("fidelity".to_string(), str_attr("full"));
|
||||
let edge = Edge::new("fork", "branch");
|
||||
|
||||
let resolved = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Compact);
|
||||
|
||||
assert_eq!(resolved.requested, Some(Fidelity::Full));
|
||||
assert_eq!(resolved.effective, Some(Fidelity::SummaryHigh));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parallel_branch_fidelity_equal_to_fork_inherits() {
|
||||
let mut node = Node::new("branch");
|
||||
node.attrs
|
||||
.insert("fidelity".to_string(), str_attr("summary:high"));
|
||||
let edge = Edge::new("fork", "branch");
|
||||
|
||||
let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh);
|
||||
|
||||
assert_eq!(resolution.requested, Some(Fidelity::SummaryHigh));
|
||||
assert_eq!(resolution.effective, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn explicit_full_branch_equal_to_degraded_fork_inherits() {
|
||||
let mut node = Node::new("branch");
|
||||
node.attrs.insert("fidelity".to_string(), str_attr("full"));
|
||||
let edge = Edge::new("fork", "branch");
|
||||
|
||||
let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::SummaryHigh);
|
||||
|
||||
assert_eq!(resolution.requested, Some(Fidelity::Full));
|
||||
assert_eq!(resolution.effective, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn full_fork_without_branch_fidelity_does_not_create_entry() {
|
||||
let node = Node::new("branch");
|
||||
let edge = Edge::new("fork", "branch");
|
||||
|
||||
let resolution = resolve_parallel_branch_fidelity(&edge, &node, Fidelity::Full);
|
||||
|
||||
assert_eq!(resolution.requested, None);
|
||||
assert_eq!(resolution.effective, None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_before_node_rebuilds_branch_preamble_stash() {
|
||||
let graph = parallel_workflow_graph(None, Some("truncate"));
|
||||
let run_dir = tempfile::tempdir().unwrap();
|
||||
let lifecycle = test_lifecycle(&graph, run_dir.path()).await;
|
||||
let state: WfRunState = ExecutionState::new(&graph).unwrap();
|
||||
let fork = graph.get_node("fork").unwrap();
|
||||
|
||||
lifecycle.before_node(&fork, &state).await.unwrap();
|
||||
state.context.set(
|
||||
keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES,
|
||||
serde_json::json!(["stale", "entries", "must disappear"]),
|
||||
);
|
||||
lifecycle.before_node(&fork, &state).await.unwrap();
|
||||
|
||||
let stash = state
|
||||
.context
|
||||
.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES)
|
||||
.expect("parallel stash should be set");
|
||||
let entries = stash.as_array().expect("parallel stash should be an array");
|
||||
assert_eq!(entries.len(), 2);
|
||||
assert!(entries[0].is_object());
|
||||
assert!(entries[1].is_null());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn non_parallel_before_node_overwrites_branch_preamble_stash_with_null() {
|
||||
let graph = parallel_workflow_graph(None, Some("truncate"));
|
||||
let run_dir = tempfile::tempdir().unwrap();
|
||||
let lifecycle = test_lifecycle(&graph, run_dir.path()).await;
|
||||
let state: WfRunState = ExecutionState::new(&graph).unwrap();
|
||||
let fork = graph.get_node("fork").unwrap();
|
||||
let work = graph.get_node("work").unwrap();
|
||||
|
||||
lifecycle.before_node(&fork, &state).await.unwrap();
|
||||
lifecycle.before_node(&work, &state).await.unwrap();
|
||||
|
||||
assert_eq!(
|
||||
state.context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
Some(serde_json::Value::Null)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn resumed_full_fork_degrades_without_rendering_fallback_branches() {
|
||||
let graph = parallel_workflow_graph(Some("full"), None);
|
||||
let run_dir = tempfile::tempdir().unwrap();
|
||||
let lifecycle = test_lifecycle(&graph, run_dir.path()).await;
|
||||
lifecycle.set_degrade_fidelity_on_resume(true);
|
||||
let state: WfRunState = ExecutionState::new(&graph).unwrap();
|
||||
let fork = graph.get_node("fork").unwrap();
|
||||
|
||||
lifecycle.before_node(&fork, &state).await.unwrap();
|
||||
|
||||
assert_eq!(state.context.fidelity(), Fidelity::SummaryHigh);
|
||||
assert_eq!(
|
||||
state.context.get(keys::INTERNAL_PARALLEL_BRANCH_PREAMBLES),
|
||||
Some(serde_json::json!([null, null]))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn fidelity_defaults_to_compact() {
|
||||
let node = Node::new("work");
|
||||
|
|
|
|||
|
|
@ -2395,7 +2395,7 @@ reasoning = false
|
|||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn workflow_persists_authoritative_openrouter_cost_for_agent_stage() {
|
||||
use fabro_auth::EnvCredentialSource;
|
||||
use fabro_workflow::steering_hub::SteeringHub;
|
||||
|
|
@ -2474,8 +2474,18 @@ base_url = "{}"
|
|||
registry.register("start", Box::new(StartHandler));
|
||||
registry.register("exit", Box::new(ExitHandler));
|
||||
|
||||
let events = Arc::new(std::sync::Mutex::new(Vec::new()));
|
||||
let events_for_listener = Arc::clone(&events);
|
||||
let emitter = Arc::new(Emitter::default());
|
||||
emitter.on_event(move |event| {
|
||||
if event.event_name() == "agent.message" {
|
||||
std::thread::sleep(Duration::from_millis(50));
|
||||
}
|
||||
events_for_listener.lock().unwrap().push(event.clone());
|
||||
});
|
||||
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env());
|
||||
let engine = WorkflowRunner::new(registry, emitter, local_env());
|
||||
let run_options = RunOptions {
|
||||
settings: WorkflowSettings::default(),
|
||||
run_dir: dir.path().to_path_buf(),
|
||||
|
|
@ -2508,6 +2518,22 @@ base_url = "{}"
|
|||
Some(AUTHORITATIVE_COST_USD_MICROS),
|
||||
"provider-reported usage.cost should override the catalog estimate"
|
||||
);
|
||||
|
||||
let events = events.lock().unwrap();
|
||||
let agent_message = events
|
||||
.iter()
|
||||
.position(|event| event.event_name() == "agent.message")
|
||||
.expect("agent message should be emitted");
|
||||
let stage_completed = events
|
||||
.iter()
|
||||
.position(|event| {
|
||||
event.event_name() == "stage.completed" && event.node_id.as_deref() == Some("work")
|
||||
})
|
||||
.expect("work stage completion should be emitted");
|
||||
assert!(
|
||||
agent_message < stage_completed,
|
||||
"agent messages must be forwarded before terminal stage events"
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
|
@ -5014,6 +5040,27 @@ struct FidelityCapturingHandler {
|
|||
captures: FidelityCaptures,
|
||||
}
|
||||
|
||||
struct ParallelFidelitySeedHandler;
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Handler for ParallelFidelitySeedHandler {
|
||||
async fn execute(
|
||||
&self,
|
||||
_node: &Node,
|
||||
_context: &Context,
|
||||
_graph: &Graph,
|
||||
_run_dir: &Path,
|
||||
_services: &fabro_workflow::handler::EngineServices,
|
||||
) -> Result<Outcome, Error> {
|
||||
let mut outcome = Outcome::success();
|
||||
outcome.context_updates.insert(
|
||||
"parallel_fidelity_marker".to_string(),
|
||||
serde_json::json!("marker visible to inherited preambles"),
|
||||
);
|
||||
Ok(outcome)
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Handler for FidelityCapturingHandler {
|
||||
async fn execute(
|
||||
|
|
@ -9411,6 +9458,176 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String {
|
|||
.expect("report prompt should exist")
|
||||
}
|
||||
|
||||
async fn run_parallel_fidelity_capture(
|
||||
fork_fidelity: Option<&str>,
|
||||
branch_node_fidelity: Option<&str>,
|
||||
branch_edge_fidelity: Option<&str>,
|
||||
) -> FidelityCaptures {
|
||||
use fabro_workflow::handler::fan_in::FanInHandler;
|
||||
use fabro_workflow::handler::parallel::ParallelHandler;
|
||||
|
||||
let mut graph = make_graph_with_start_exit("ParallelFidelityTest");
|
||||
graph.attrs.insert(
|
||||
"goal".to_string(),
|
||||
AttrValue::String("Verify parallel branch context".to_string()),
|
||||
);
|
||||
|
||||
let mut seed = Node::new("seed");
|
||||
seed.attrs.insert(
|
||||
"type".to_string(),
|
||||
AttrValue::String("parallel_fidelity_seed".to_string()),
|
||||
);
|
||||
let mut fork = Node::new("fork");
|
||||
fork.attrs.insert(
|
||||
"shape".to_string(),
|
||||
AttrValue::String("component".to_string()),
|
||||
);
|
||||
if let Some(fidelity) = fork_fidelity {
|
||||
fork.attrs.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String(fidelity.to_string()),
|
||||
);
|
||||
}
|
||||
let mut branch_a = Node::new("branch_a");
|
||||
branch_a.attrs.insert(
|
||||
"type".to_string(),
|
||||
AttrValue::String("fidelity_capture".to_string()),
|
||||
);
|
||||
if let Some(fidelity) = branch_node_fidelity {
|
||||
branch_a.attrs.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String(fidelity.to_string()),
|
||||
);
|
||||
}
|
||||
let mut branch_b = Node::new("branch_b");
|
||||
branch_b.attrs.insert(
|
||||
"type".to_string(),
|
||||
AttrValue::String("fidelity_capture".to_string()),
|
||||
);
|
||||
let mut fan_in = Node::new("fan_in");
|
||||
fan_in.attrs.insert(
|
||||
"shape".to_string(),
|
||||
AttrValue::String("tripleoctagon".to_string()),
|
||||
);
|
||||
|
||||
graph.nodes.insert(seed.id.clone(), seed);
|
||||
graph.nodes.insert(fork.id.clone(), fork);
|
||||
graph.nodes.insert(branch_a.id.clone(), branch_a);
|
||||
graph.nodes.insert(branch_b.id.clone(), branch_b);
|
||||
graph.nodes.insert(fan_in.id.clone(), fan_in);
|
||||
graph.edges.push(Edge::new("start", "seed"));
|
||||
graph.edges.push(Edge::new("seed", "fork"));
|
||||
let mut branch_a_edge = Edge::new("fork", "branch_a");
|
||||
if let Some(fidelity) = branch_edge_fidelity {
|
||||
branch_a_edge.attrs.insert(
|
||||
"fidelity".to_string(),
|
||||
AttrValue::String(fidelity.to_string()),
|
||||
);
|
||||
}
|
||||
graph.edges.push(branch_a_edge);
|
||||
graph.edges.push(Edge::new("fork", "branch_b"));
|
||||
graph.edges.push(Edge::new("branch_a", "fan_in"));
|
||||
graph.edges.push(Edge::new("branch_b", "fan_in"));
|
||||
graph.edges.push(Edge::new("fan_in", "exit"));
|
||||
|
||||
let captures = FidelityCaptures::new();
|
||||
let mut registry = HandlerRegistry::new(Box::new(StartHandler));
|
||||
registry.register("start", Box::new(StartHandler));
|
||||
registry.register("exit", Box::new(ExitHandler));
|
||||
registry.register("parallel", Box::new(ParallelHandler));
|
||||
registry.register(
|
||||
"parallel.fan_in",
|
||||
Box::new(FanInHandler::new(Some(Box::new(MockCodergenBackend)))),
|
||||
);
|
||||
registry.register(
|
||||
"parallel_fidelity_seed",
|
||||
Box::new(ParallelFidelitySeedHandler),
|
||||
);
|
||||
registry.register(
|
||||
"fidelity_capture",
|
||||
Box::new(FidelityCapturingHandler {
|
||||
captures: captures.clone(),
|
||||
}),
|
||||
);
|
||||
|
||||
let dir = tempfile::tempdir().expect("parallel fidelity run directory should be created");
|
||||
let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env());
|
||||
let run_options = RunOptions {
|
||||
settings: WorkflowSettings::default(),
|
||||
run_dir: dir.path().to_path_buf(),
|
||||
cancel_token: CancellationToken::new(),
|
||||
run_id: test_run_id("parallel-fidelity"),
|
||||
labels: std::collections::HashMap::new(),
|
||||
workflow_slug: None,
|
||||
github_app: None,
|
||||
base_branch: None,
|
||||
display_base_sha: None,
|
||||
pre_run_git: None,
|
||||
fork_source_ref: None,
|
||||
git: None,
|
||||
};
|
||||
let (outcome, _state) = engine
|
||||
.run_with_state(&graph, &run_options)
|
||||
.await
|
||||
.expect("parallel fidelity workflow should succeed");
|
||||
assert_eq!(outcome.status, StageOutcome::Succeeded);
|
||||
captures
|
||||
}
|
||||
|
||||
fn captured_fidelity_preamble(captures: &FidelityCaptures, node_id: &str) -> (String, String) {
|
||||
let fidelity = captures
|
||||
.fidelities
|
||||
.lock()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.find(|(captured_node_id, _)| captured_node_id == node_id)
|
||||
.map(|(_, fidelity)| fidelity.clone())
|
||||
.expect("branch fidelity should be captured");
|
||||
let preamble = captures
|
||||
.preambles
|
||||
.lock()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.find(|(captured_node_id, _)| captured_node_id == node_id)
|
||||
.map(|(_, preamble)| preamble.clone())
|
||||
.expect("branch preamble should be captured");
|
||||
(fidelity, preamble)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_branches_get_per_branch_preambles_by_fidelity() {
|
||||
let captures = run_parallel_fidelity_capture(None, Some("truncate"), None).await;
|
||||
|
||||
let (branch_a_fidelity, branch_a_preamble) = captured_fidelity_preamble(&captures, "branch_a");
|
||||
let (branch_b_fidelity, branch_b_preamble) = captured_fidelity_preamble(&captures, "branch_b");
|
||||
|
||||
assert_eq!(branch_a_fidelity, "truncate");
|
||||
assert!(!branch_a_preamble.contains("parallel_fidelity_marker"));
|
||||
assert_eq!(branch_b_fidelity, "compact");
|
||||
assert!(branch_b_preamble.contains("parallel_fidelity_marker"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_fork_fidelity_still_applies_to_all_branches() {
|
||||
let captures = run_parallel_fidelity_capture(Some("truncate"), None, None).await;
|
||||
|
||||
for branch_id in ["branch_a", "branch_b"] {
|
||||
let (fidelity, preamble) = captured_fidelity_preamble(&captures, branch_id);
|
||||
assert_eq!(fidelity, "truncate");
|
||||
assert!(!preamble.contains("parallel_fidelity_marker"));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn parallel_branch_edge_fidelity_overrides_node_fidelity() {
|
||||
let captures =
|
||||
run_parallel_fidelity_capture(None, Some("summary:high"), Some("truncate")).await;
|
||||
|
||||
let (fidelity, preamble) = captured_fidelity_preamble(&captures, "branch_a");
|
||||
assert_eq!(fidelity, "truncate");
|
||||
assert!(!preamble.contains("parallel_fidelity_marker"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn fidelity_prompt_compact() {
|
||||
let prompt = run_fidelity_prompt_pipeline("compact").await;
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue