commit 1f3b09be58ad9a964a391519df243bac8e60a4aa Author: Fabro Date: Wed May 27 17:40:08 2026 -0400 init run ⚒️ Generated with [Fabro](https://fabro.sh) diff --git a/graph.fabro b/graph.fabro new file mode 100644 index 000000000..d4d99bf9d --- /dev/null +++ b/graph.fabro @@ -0,0 +1,35 @@ +digraph ImplementPlan { + graph [ + goal="Implement and simplify", + model_stylesheet=" + * { model: claude-opus-4-7; } + " + ] + rankdir=LR + + start [shape=Mdiamond, label="Start"] + exit [shape=Msquare, label="Exit"] + + toolchain [label="Toolchain", shape=parallelogram, script="command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1", max_retries=0] + preflight_compile [label="Preflight Compile", shape=parallelogram, script="cargo check -q --workspace 2>&1", max_retries=0] + preflight_lint [label="Preflight Lint", shape=parallelogram, script="cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", max_retries=0] + fix_lints [label="Fix Lints", prompt="The preflight lint step failed. Read the build output from context and fix all clippy lint warnings.", max_visits=3] + 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="gpt-55", reasoning_effort="xhigh"] + simplify_opus [label="Simplify (Opus)", prompt="@prompts/simplify.md"] + simplify_gpt [label="Simplify (GPT-55)", prompt="@prompts/simplify.md", model="gpt-55"] + 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"] + 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.", max_visits=3] + + start -> toolchain + toolchain -> preflight_compile [condition="outcome=succeeded"] + toolchain -> exit + preflight_compile -> preflight_lint [condition="outcome=succeeded"] + preflight_compile -> exit + preflight_lint -> implement [condition="outcome=succeeded"] + preflight_lint -> fix_lints + fix_lints -> preflight_lint + implement -> simplify_opus -> simplify_gpt -> verify + verify -> exit [condition="outcome=succeeded"] + verify -> fixup + fixup -> verify +} diff --git a/run.json b/run.json new file mode 100644 index 000000000..00eb98ea3 --- /dev/null +++ b/run.json @@ -0,0 +1,533 @@ +{ + "title": "Worker Control Bus Implementation Plan", + "spec": { + "run_id": "01KSNP2DXVS2TD1HEASQAFBGFK", + "settings": { + "project": { + "name": null, + "description": null, + "metadata": {} + }, + "workflow": { + "name": null, + "description": null, + "graph": "workflow.fabro", + "metadata": {} + }, + "run": { + "goal": { + "type": "inline", + "value": "# Worker Control Bus Implementation Plan\n\n> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.\n\n**Goal:** Replace the server-to-worker stdin JSONL control pipe with a backend-agnostic worker control bus, implemented now with a local in-memory bus and delivered to workers over a worker-initiated WebSocket.\n\n**Architecture:** API handlers publish `WorkerControlEnvelope` messages to a `WorkerControlBus`; the worker WebSocket route subscribes to that bus and forwards ordered delivery frames to the worker. Workers track the last fully applied delivery id and reconnect with `?after=` after any unexpected WebSocket close. The first backend is an in-process `LocalWorkerControlBus` for local and single-node deployments. A Redis Streams backend must fit behind the same trait later, but Redis is explicitly out of scope for this implementation plan.\n\n**Tech Stack:** Rust, Axum WebSockets, tokio-tungstenite, UnixStream, async-trait or boxed async traits, tokio channels/Notify, worker JWT auth, existing `WorkerControlEnvelope`.\n\n---\n\n## Key Decisions\n\n- WebSocket fully replaces stdin control. Do not keep stdin JSONL as a compatibility path.\n- The worker protocol stays identical across local, single-node, ECS, and later SaaS deployments.\n- The server-side delivery backend is the only thing that varies by deployment.\n- This plan implements only `LocalWorkerControlBus`.\n- This plan does not add Redis dependencies, Redis configuration, Redis tests, Redis health checks, or Redis runtime behavior.\n- Redis Streams are covered only as a future backend contract so the local design does not paint us into a corner.\n- There is no `Latest` cursor. The first worker connection starts at the beginning of the run's retained control stream; reconnects resume after the worker's last fully applied delivery id.\n- Every WebSocket text frame is a delivery frame with an id and envelope. The worker advances `last_applied_id` only after applying the envelope.\n- Workers reconnect forever while the local run is not terminal, using backoff from 100ms, doubled after each failure, capped at 5s.\n- The worker must complete its first control-stream connection before starting or resuming workflow execution. Temporary first-connect failures wait and retry; they do not start a control-disconnected run.\n- Invalid cursor means the bus can no longer prove replay correctness. The worker treats it as fatal control-channel loss and fails/aborts the run as infrastructure failure, not as user cancellation.\n- WebSocket liveness is handled at the WebSocket layer with explicit ping/pong and timeout logic. The bus does not know about heartbeats.\n- ECS task launch, ECS stop/reconciliation, Redis-backed multi-node delivery, and remote hard-kill behavior are follow-up work.\n\n## Redis Fit Later: Out of Scope Now\n\nRedis should later implement the same `WorkerControlBus` API introduced here.\n\n- `publish(run_id, envelope)` maps to `XADD fabro:run:{run_id}:control ...`.\n- First `subscribe(run_id, Start)` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control 0-0`.\n- Reconnect `subscribe(run_id, After(id))` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control {id}`.\n- Local message ids use an opaque string format such as `local:1`; Redis message ids can use Redis stream ids such as `1716810000000-0`.\n- The WebSocket route should not care whether the subscription source is local memory or Redis.\n- The worker should not care whether the frame came from a local bus or Redis.\n- Redis trimming/retention, consumer groups, per-tenant key naming, TLS/auth, reconnect-after-redeploy semantics, and SaaS config validation are not part of this plan.\n\n## Proposed File Structure\n\n- Create `lib/crates/fabro-server/src/worker_control/mod.rs`\n - Owns the server-side control bus abstraction and re-exports the local backend.\n- Create `lib/crates/fabro-server/src/worker_control/bus.rs`\n - Defines `WorkerControlBus`, `WorkerControlDelivery`, `WorkerControlMessageId`, `WorkerControlCursor`, and bus errors.\n- Create `lib/crates/fabro-server/src/worker_control/local.rs`\n - Implements `LocalWorkerControlBus` using process memory.\n- Create `lib/crates/fabro-server/src/server/handler/worker_control.rs`\n - Adds the worker-only WebSocket route.\n- Modify `lib/crates/fabro-server/src/server.rs`\n - Adds the bus to `AppState`, replaces subprocess `RunAnswerTransport` sends with bus publishes, removes stdin pumping.\n- Modify `lib/crates/fabro-server/src/server/handler/mod.rs`\n - Registers the worker control route.\n- Modify `lib/crates/fabro-server/src/server/handler/lifecycle.rs`\n - Sends pause/unpause/cancel controls through the transport/bus where appropriate.\n- Modify `lib/crates/fabro-cli/src/commands/run/runner.rs`\n - Replaces stdin reading with worker WebSocket client handling.\n- Modify `lib/crates/fabro-cli/Cargo.toml`\n - Adds `tokio-tungstenite` as a direct dependency if needed.\n- Modify `lib/crates/fabro-interview/src/control_protocol.rs`\n - Adds pause/unpause control messages and a transport delivery frame type shared by server and worker.\n\n## Task 1: Define the Control Bus Contract\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/worker_control/mod.rs`\n- Create: `lib/crates/fabro-server/src/worker_control/bus.rs`\n- Modify: `lib/crates/fabro-server/src/lib.rs`\n\n- [ ] Add a private `worker_control` module in `fabro-server`.\n- [ ] Define `WorkerControlMessageId` as an opaque cloneable id rather than a numeric type.\n- [ ] Define `WorkerControlCursor` with `Start` and `After(WorkerControlMessageId)` variants.\n- [ ] Define `WorkerControlDelivery { id: WorkerControlMessageId, envelope: WorkerControlEnvelope }`.\n- [ ] Define `WorkerControlBus` with async `publish(run_id, envelope)` and `subscribe(run_id, cursor)` methods.\n- [ ] Make `subscribe` return a stream-like receiver owned by the caller, so the WebSocket handler can forward messages without knowing the backend.\n- [ ] Define explicit bus errors for closed backend, unavailable backend, invalid cursor, and publish timeout.\n- [ ] Document in code comments that `Start` maps to Redis stream id `0-0` and `After(id)` maps to Redis `XREAD` after that id, but do not add Redis code.\n- [ ] Add unit tests for id equality/debug formatting and cursor parsing from the optional `after` query parameter.\n- [ ] Test that absent `after` parses as `WorkerControlCursor::Start`.\n- [ ] Test that present `after=local:42` parses as `WorkerControlCursor::After(...)`.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 2: Implement the Local In-Memory Bus\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/worker_control/local.rs`\n- Test: `lib/crates/fabro-server/src/worker_control/local.rs`\n\n- [ ] Implement `LocalWorkerControlBus` as `Arc>>`.\n- [ ] Store messages per run in insertion order with a monotonic local sequence id.\n- [ ] Wake active subscribers when `publish` appends a message.\n- [ ] Support `subscribe(run_id, Start)` for first worker startup; it must replay retained messages from the beginning of the run control stream.\n- [ ] Support `subscribe(run_id, After(id))` so reconnect uses the same API that later maps to Redis `XREAD`.\n- [ ] Allow `publish` before the worker subscribes; retained messages must be visible to the first `Start` subscriber.\n- [ ] Trim retained local messages to a bounded per-run size so a disconnected local worker cannot grow memory without bound. Use a named constant with initial value 1024 messages per run.\n- [ ] Return a clear `invalid cursor` error when a subscriber asks for an id that has been trimmed or belongs to a different local stream.\n- [ ] Add a cleanup method for terminal runs so completed/cancelled runs can release retained control messages.\n- [ ] Test that messages publish in order.\n- [ ] Test that an active subscriber receives a message published after subscription.\n- [ ] Test that messages published before subscription are replayed to a `Start` subscriber.\n- [ ] Test that `After(id)` receives only later messages.\n- [ ] Test that trimming bounds retained messages and reports an invalid old cursor.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 3: Add Control Bus to Server State\n\n**Files:**\n- Modify: `lib/crates/fabro-server/src/server.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Add `worker_control_bus: Arc` to `AppState`.\n- [ ] Construct `LocalWorkerControlBus` in normal server state initialization.\n- [ ] Add a test-only way to inject a fake or local bus without exposing test helpers to production builds.\n- [ ] Keep demo/in-process execution behavior unchanged unless it currently depends on subprocess control.\n- [ ] Add a state construction test proving the default bus is local and available.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 4: Extend the Control Protocol\n\n**Files:**\n- Modify: `lib/crates/fabro-interview/src/control_protocol.rs`\n- Test: `lib/crates/fabro-interview/src/control_protocol.rs`\n\n- [ ] Add `WorkerControlEnvelope::pause_run()` and `WorkerControlEnvelope::unpause_run()` constructors.\n- [ ] Add `WorkerControlMessage::RunPause` serialized as `\"run.pause\"`.\n- [ ] Add `WorkerControlMessage::RunUnpause` serialized as `\"run.unpause\"`.\n- [ ] Add `WorkerControlDeliveryFrame { id: String, envelope: WorkerControlEnvelope }` as the WebSocket text-frame payload shared by server and worker.\n- [ ] Add round-trip serde tests for both new messages.\n- [ ] Add round-trip serde tests for `WorkerControlDeliveryFrame`.\n- [ ] Run `cargo nextest run -p fabro-interview control_protocol`.\n\n## Task 5: Share Worker Message Handling\n\n**Files:**\n- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n\n- [ ] Split `apply_worker_control_line(...)` into parsing and `apply_worker_control_message(...)`.\n- [ ] Route WebSocket delivery frames through `apply_worker_control_message(...)`.\n- [ ] Route `run.pause` to `RunControlState::request_pause()`.\n- [ ] Route `run.unpause` to `RunControlState::request_unpause()`.\n- [ ] Add a small in-memory applied-id dedupe set in the worker control task; ignore duplicate delivery ids before applying envelopes.\n- [ ] Update `last_applied_id` only after `apply_worker_control_message(...)` returns.\n- [ ] Treat all current control messages as idempotent under delivery-id dedupe. `run.steer` must not be applied twice for the same delivery id.\n- [ ] Keep control stream close behavior explicit: an unexpected close triggers reconnect; a fatal invalid cursor interrupts pending interviews and fails/aborts the run as control-channel loss.\n- [ ] Update existing stdin-era tests to exercise the shared message handler directly.\n- [ ] Add tests for pause and unpause routing.\n- [ ] Add a test proving duplicate delivery ids are not applied twice.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n\n## Task 6: Add Worker WebSocket Client\n\n**Files:**\n- Modify: `lib/crates/fabro-cli/Cargo.toml`\n- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n\n- [ ] Add `tokio-tungstenite.workspace = true` as a direct `fabro-cli` dependency if the crate does not already have it.\n- [ ] Add a helper that builds the control-stream request for a `ServerTarget` and `RunId`.\n- [ ] For HTTP URLs, convert `http` to `ws` and `https` to `wss`.\n- [ ] For Unix socket paths, connect `tokio::net::UnixStream` and use `ws://fabro/api/v1/runs/{run_id}/worker/control-stream` for the handshake host/path.\n- [ ] Add the worker bearer token as an `Authorization: Bearer ...` request header.\n- [ ] On the first connection, omit the `after` query parameter so the server maps it to `WorkerControlCursor::Start`.\n- [ ] On reconnect, include `?after=` when `last_applied_id` is set.\n- [ ] Spawn a WebSocket control manager task in `execute(...)` after `ControlInterviewer`, `RunControlState`, `CancellationToken`, and `SteeringHub` are created, and before `operations::start` or `operations::resume`.\n- [ ] Gate `operations::start` and `operations::resume` on the first successful control-stream connection.\n- [ ] The control manager should keep reconnecting while the run is not locally terminal, with backoff starting at 100ms, doubling after each failure, and capped at 5s.\n- [ ] Deserialize each text frame into `WorkerControlDeliveryFrame`.\n- [ ] Apply each envelope through `apply_worker_control_message(...)`, then record the frame id as `last_applied_id`.\n- [ ] Respond to received WebSocket ping frames with pong frames.\n- [ ] Send worker-initiated ping frames every 15s.\n- [ ] Track pongs for worker-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness.\n- [ ] Treat normal close/error as reconnectable while the run is not terminal.\n- [ ] Treat HTTP 410 Gone or a WebSocket close reason of `invalid_cursor` as fatal control-channel loss.\n- [ ] On fatal control-channel loss, interrupt pending interviews and fail/abort the run with an infrastructure/control-channel error, not a user cancellation.\n- [ ] Wire fatal control-channel loss back into `execute(...)` so the worker returns an error instead of silently continuing workflow execution.\n- [ ] Add tests for URL/request construction for `http`, `https`, and Unix socket targets.\n- [ ] Add tests proving first connection has no `after` query and reconnect includes `after=`.\n- [ ] Add tests for reconnect backoff bounds.\n- [ ] Add tests for ping/pong timeout behavior using paused Tokio time.\n- [ ] Add a local Unix-socket WebSocket test proving the client can complete a handshake against an Axum route.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n\n## Task 7: Add Worker-Only Control Stream Route\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/server/handler/worker_control.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/mod.rs`\n- Modify: `lib/crates/fabro-server/src/principal_middleware.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Add a narrow helper or extractor that accepts only authenticated worker principals whose token run id matches the route run id.\n- [ ] Add `GET /runs/{id}/worker/control-stream` to real API routes only.\n- [ ] Reject missing runs, terminal runs, and archived runs before upgrading.\n- [ ] Reject user JWTs and cross-run worker JWTs.\n- [ ] Parse absent `after` into `WorkerControlCursor::Start`.\n- [ ] Parse present `after` into `WorkerControlCursor::After(id)`.\n- [ ] On upgrade, call `worker_control_bus.subscribe(run_id, cursor)`.\n- [ ] If `subscribe` returns invalid cursor before upgrade, reject with HTTP 410 Gone.\n- [ ] Serialize each `WorkerControlDelivery` to `WorkerControlDeliveryFrame` and send it as a WebSocket text frame.\n- [ ] Send server-initiated ping frames every 15s.\n- [ ] Respond to received WebSocket ping frames with pong frames.\n- [ ] Track pongs for server-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness.\n- [ ] On timeout or disconnect, drop the bus subscription so local resources are released.\n- [ ] Do not store the live WebSocket sender in `ManagedRun`; the bus is now the delivery boundary.\n- [ ] Add tests for auth rejection, successful `Start` subscription, successful `After(id)` subscription, frame delivery, invalid cursor rejection as 410 Gone, ping/pong timeout cleanup, and cross-run worker rejection.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 8: Replace Server-Side Stdin Transport with Bus Publishing\n\n**Files:**\n- Modify: `lib/crates/fabro-server/src/server.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/lifecycle.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/pair.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Replace `RunAnswerTransport::Subprocess { control_tx }` with a bus-backed subprocess/worker variant.\n- [ ] Ensure the bus-backed variant has enough context to publish messages for the correct `RunId`.\n- [ ] Delete `pump_worker_control_jsonl(...)`.\n- [ ] Change `worker_command(...)` so `__run-worker` uses `stdin(Stdio::null())` instead of `stdin(Stdio::piped())`.\n- [ ] Remove child-stdin extraction and the control pump task from `execute_run_subprocess(...)`.\n- [ ] Keep stderr capture and worker exit handling unchanged.\n- [ ] Update `RunAnswerTransport` methods so answer, cancel, steer, interrupt, pair start/message/end all publish the existing envelope to `WorkerControlBus`.\n- [ ] Add `pause_run()` and `unpause_run()` methods on `RunAnswerTransport`.\n- [ ] Update pause/unpause lifecycle handlers to send `run.pause` and `run.unpause` over the bus for running workers.\n- [ ] Keep process signals only for hard cleanup paths such as cancel fallback, shutdown, terminal delete, and force removal.\n- [ ] Update existing tests that assert subprocess transport enqueue behavior to assert bus publish behavior instead.\n- [ ] Add a `worker_command` test proving stdin is null/not piped and `FABRO_WORKER_TOKEN` still travels only through env.\n- [ ] Run `cargo nextest run -p fabro-server worker_command`.\n\n## Task 9: End-to-End Local Control Flow Regression\n\n**Files:**\n- Test: `lib/crates/fabro-cli/tests/it/cmd/runner.rs`\n- Test: `lib/crates/fabro-server/tests/it/scenario/lifecycle.rs`\n\n- [ ] Add a test run where the worker connects to the control WebSocket and receives a cancel request through `LocalWorkerControlBus`.\n- [ ] Add a test where the server publishes a control message before the worker connects and the worker receives it on first connection.\n- [ ] Add a reconnect test where the worker applies message A, reconnects with `after=`, and then receives only message B.\n- [ ] Add an invalid-cursor test proving the worker reports control-channel loss as infrastructure failure rather than user cancellation.\n- [ ] Add a human-interview test proving submitted answers reach the worker through the bus and WebSocket.\n- [ ] Add a steer or interrupt test proving live agent controls still reach the worker transport.\n- [ ] Add a local Unix-socket server test proving the default local server target works without stdin.\n- [ ] Run `cargo nextest run -p fabro-cli --test it runner`.\n- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`.\n\n## Task 10: Final Verification\n\n**Files:**\n- Modify only if failures expose necessary fixes.\n\n- [ ] Run `cargo nextest run -p fabro-interview control_protocol`.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n- [ ] Run `cargo nextest run -p fabro-server worker_command`.\n- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`.\n- [ ] Run `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`.\n- [ ] Confirm no code path still writes `WorkerControlEnvelope` to child stdin.\n- [ ] Confirm no Redis dependency, Redis config key, or Redis runtime path was added.\n- [ ] Confirm there is no `Latest` cursor or wait-for-subscriber behavior in the control bus.\n- [ ] Confirm WebSocket ping/pong handling is explicit on both worker and server.\n- [ ] Confirm `run.steer` and other controls are protected from duplicate delivery-id application.\n- [ ] Confirm `__run-worker` still scrubs `FABRO_WORKER_TOKEN` from process env before launching descendants.\n\n## Acceptance Criteria\n\n- All worker control traffic uses the worker control bus plus WebSocket last-mile transport.\n- Local Unix-socket server targets and remote HTTP(S) server targets both support worker control without Redis.\n- Local and single-node deployments require no external control-channel service.\n- The bus API can later be implemented by Redis Streams without changing API handlers or worker message handling.\n- First worker connection replays retained messages from the beginning of the run control stream; reconnect resumes after the last fully applied id.\n- Invalid cursor is the only fatal control-stream replay failure and is surfaced as infrastructure/control-channel failure, not user cancellation.\n- WebSocket liveness is explicit and backend-agnostic.\n- Existing run event/blob/artifact HTTP paths are unchanged.\n- Existing worker JWT scope rules remain authoritative.\n- Temporary WebSocket disconnects reconnect and replay through the bus; only unrecoverable replay loss reports worker-control-unavailable behavior.\n- Worker stdout/stderr behavior remains unchanged except that stdin is no longer a control channel.\n- Redis is clearly documented as future work and is not required by this plan.\n" + }, + "working_dir": null, + "metadata": {}, + "inputs": {}, + "model": { + "provider": "anthropic", + "name": "claude-sonnet-4-6", + "fallbacks": [], + "controls": { + "reasoning_effort": null, + "speed": null + } + }, + "git": { + "author": null + }, + "prepare": { + "commands": [], + "timeout_ms": 300000 + }, + "execution": { + "mode": "normal", + "approval": "prompt" + }, + "checkpoint": { + "exclude_globs": [], + "skip_git_hooks": false + }, + "clone": { + "enabled": true + }, + "run_branch": { + "enabled": true, + "push": true + }, + "meta_branch": { + "enabled": true, + "push": true + }, + "environment": { + "id": "fabro-dev", + "provider": "daytona", + "image": { + "docker": null, + "dockerfile": { + "type": "inline", + "value": "FROM ubuntu:24.04\n\nRUN apt-get update && apt-get install -y --no-install-recommends \\\n curl git ripgrep ca-certificates build-essential pkg-config libssl-dev unzip python3 \\\n xvfb xfce4 xfce4-terminal x11vnc novnc dbus-x11 \\\n libx11-6 libxrandr2 libxext6 libxrender1 libxfixes3 libxss1 libxtst6 libxi6 \\\n && rm -rf /var/lib/apt/lists/*\n\n# Install real Chromium (not the snap stub) via xtradeb PPA\nRUN apt-get update && apt-get install -y --no-install-recommends \\\n software-properties-common curl gnupg \\\n && add-apt-repository -y ppa:xtradeb/apps \\\n && apt-get update \\\n && apt-get install -y --no-install-recommends chromium \\\n && rm -rf /var/lib/apt/lists/*\n\n# Wrapper: Chromium needs --no-sandbox when running as root in a container,\n# and --disable-dev-shm-usage avoids crashes from small /dev/shm\nRUN printf '#!/bin/bash\\nexec /usr/bin/chromium --no-sandbox --disable-dev-shm-usage \"$@\"\\n' \\\n > /usr/local/bin/chromium-wrapper \\\n && chmod +x /usr/local/bin/chromium-wrapper\n\n# Make the wrapper the default in the system .desktop file and via alternatives\nRUN sed -i 's|^Exec=.*|Exec=/usr/local/bin/chromium-wrapper %U|' \\\n /usr/share/applications/chromium.desktop \\\n && update-alternatives --install /usr/bin/x-www-browser x-www-browser \\\n /usr/local/bin/chromium-wrapper 100\n\n# Tell XFCE's exo-open that Chromium is the WebBrowser helper (system-wide)\nRUN mkdir -p /etc/xdg/xfce4 /usr/share/xfce4/helpers \\\n && printf 'WebBrowser=custom-WebBrowser\\n' > /etc/xdg/xfce4/helpers.rc \\\n && printf '[Desktop Entry]\\n\\\nVersion=1.0\\n\\\nType=X-XFCE-Helper\\n\\\nName=Chromium\\n\\\nIcon=chromium\\n\\\nX-XFCE-Category=WebBrowser\\n\\\nX-XFCE-CommandsWithParameter=/usr/local/bin/chromium-wrapper \"%%s\"\\n\\\nX-XFCE-Commands=/usr/local/bin/chromium-wrapper\\n' \\\n > /usr/share/xfce4/helpers/custom-WebBrowser.desktop\n\n# GitHub CLI\nRUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg \\\n | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg \\\n && echo \"deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main\" \\\n | tee /etc/apt/sources.list.d/github-cli.list > /dev/null \\\n && apt-get update && apt-get install -y --no-install-recommends gh \\\n && rm -rf /var/lib/apt/lists/*\n\n# Rust\nRUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y\nENV PATH=\"/root/.cargo/bin:${PATH}\"\nRUN rustup toolchain install nightly-2026-04-14 --profile minimal --component clippy,rustfmt\nRUN cargo install cargo-nextest --locked\nENV CARGO_INCREMENTAL=0\n\n# Bun\nRUN curl -fsSL https://bun.sh/install | bash\nENV PATH=\"/root/.bun/bin:${PATH}\"\n\nWORKDIR /root\n" + } + }, + "resources": { + "cpu": 8, + "memory": "16GB", + "disk": "20GB" + }, + "network": { + "mode": "allow_all", + "allow": [] + }, + "lifecycle": { + "preserve": false, + "stop_on_terminal": true, + "auto_stop": "30m" + }, + "labels": { + "repo": "fabro-sh/fabro" + }, + "volumes": [], + "env": {} + }, + "notifications": { + "feed": { + "enabled": true, + "provider": "slack", + "events": [ + "run.started", + "run.completed", + "run.failed" + ], + "slack": { + "channel": "#feed-fabro" + } + } + }, + "interviews": { + "provider": null, + "slack": null + }, + "agent": { + "fabro_tools": false, + "permissions": null, + "mcps": {} + }, + "hooks": [], + "scm": { + "provider": null, + "owner": null, + "repository": null, + "github": null + }, + "pull_request": { + "enabled": true, + "draft": false, + "auto_merge": false, + "merge_strategy": "squash" + }, + "artifacts": { + "include": [] + }, + "integrations": { + "github": { + "permissions": {} + } + } + } + }, + "graph": { + "name": "ImplementPlan", + "nodes": { + "simplify_opus": { + "id": "simplify_opus", + "attrs": { + "label": { + "String": "Simplify (Opus)" + }, + "prompt": { + "String": "# Simplify: Code Review and Cleanup\n\nReview changes vs. origin for reuse, quality, and efficiency. Fix any issues found.\n\n## Phase 1: Identify Changes\n\nRun git diff (or git diff HEAD if there are staged changes) to see what changed. If there are no git changes, review the most recently modified files that the user mentioned or that you edited earlier in this conversation.\n\n## Phase 2: Launch Three Review Agents in Parallel\n\nUse the Agent tool to launch all three agents concurrently in a single message. Pass each agent the full diff so it has the complete context.\n\n### Agent 1: Code Reuse Review\n\nFor each change:\n\n1. Search for existing utilities and helpers that could replace newly written code. Use Grep to find similar patterns elsewhere in the codebase — common locations are utility directories, shared modules, and files adjacent to the changed ones.\n2. Flag any new function that duplicates existing functionality. Suggest the existing function to use instead.\n3. Flag any inline logic that could use an existing utility — hand-rolled string manipulation, manual path handling, custom environment checks, ad-hoc type guards, and similar patterns are common candidates.\n\nNote: This is a greenfield app, so focus on maximizing simplicity and don't worry about changing things to achieve it.\n\n### Agent 2: Code Quality Review\n\nReview the same changes for hacky patterns:\n\n1. Redundant state: state that duplicates existing state, cached values that could be derived, observers/effects that could be direct calls\n2. Parameter sprawl: adding new parameters to a function instead of generalizing or restructuring existing ones\n3. Copy-paste with slight variation: near-duplicate code blocks that should be unified with a shared abstraction\n4. Leaky abstractions: exposing internal details that should be encapsulated, or breaking existing abstraction boundaries\n5. Stringly-typed code: using raw strings where constants, enums (string unions), or branded types already exist in the codebase\n\nNote: This is a greenfield app, so be aggressive in optimizing quality.\n\n### Agent 3: Efficiency Review\n\nReview the same changes for efficiency:\n\n1. Unnecessary work: redundant computations, repeated file reads, duplicate network/API calls, N+1 patterns\n2. Missed concurrency: independent operations run sequentially when they could run in parallel\n3. Hot-path bloat: new blocking work added to startup or per-request/per-render hot paths\n4. Unnecessary existence checks: pre-checking file/resource existence before operating (TOCTOU anti-pattern) — operate directly and handle the error\n5. Memory: unbounded data structures, missing cleanup, event listener leaks\n6. Overly broad operations: reading entire files when only a portion is needed, loading all items when filtering for one\n\n## Phase 3: Fix Issues\n\nWait for all three agents to complete. Aggregate their findings and fix each issue directly. If a finding is a false positive or not worth addressing, note it and move on — do not argue with the finding, just skip it.\n\nWhen done, briefly summarize what was fixed (or confirm the code was already clean)." + }, + "provider": { + "String": "anthropic" + }, + "model": { + "String": "claude-opus-4-7" + } + } + }, + "implement": { + "id": "implement", + "attrs": { + "model": { + "String": "gpt-5.5" + }, + "prompt": { + "String": "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." + }, + "reasoning_effort": { + "String": "xhigh" + }, + "label": { + "String": "Implement" + }, + "provider": { + "String": "openai" + } + } + }, + "preflight_compile": { + "id": "preflight_compile", + "attrs": { + "model": { + "String": "claude-opus-4-7" + }, + "script": { + "String": "cargo check -q --workspace 2>&1" + }, + "provider": { + "String": "anthropic" + }, + "max_retries": { + "Integer": 0 + }, + "label": { + "String": "Preflight Compile" + }, + "shape": { + "String": "parallelogram" + } + } + }, + "verify": { + "id": "verify", + "attrs": { + "shape": { + "String": "parallelogram" + }, + "retry_target": { + "String": "fixup" + }, + "script": { + "String": "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" + }, + "model": { + "String": "claude-opus-4-7" + }, + "label": { + "String": "Verify" + }, + "goal_gate": { + "Boolean": true + }, + "provider": { + "String": "anthropic" + } + } + }, + "preflight_lint": { + "id": "preflight_lint", + "attrs": { + "shape": { + "String": "parallelogram" + }, + "model": { + "String": "claude-opus-4-7" + }, + "label": { + "String": "Preflight Lint" + }, + "provider": { + "String": "anthropic" + }, + "script": { + "String": "cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1" + }, + "max_retries": { + "Integer": 0 + } + } + }, + "exit": { + "id": "exit", + "attrs": { + "shape": { + "String": "Msquare" + }, + "label": { + "String": "Exit" + }, + "model": { + "String": "claude-opus-4-7" + }, + "provider": { + "String": "anthropic" + } + } + }, + "start": { + "id": "start", + "attrs": { + "model": { + "String": "claude-opus-4-7" + }, + "shape": { + "String": "Mdiamond" + }, + "label": { + "String": "Start" + }, + "provider": { + "String": "anthropic" + } + } + }, + "fix_lints": { + "id": "fix_lints", + "attrs": { + "model": { + "String": "claude-opus-4-7" + }, + "prompt": { + "String": "The preflight lint step failed. Read the build output from context and fix all clippy lint warnings." + }, + "provider": { + "String": "anthropic" + }, + "max_visits": { + "Integer": 3 + }, + "label": { + "String": "Fix Lints" + } + } + }, + "simplify_gpt": { + "id": "simplify_gpt", + "attrs": { + "model": { + "String": "gpt-5.5" + }, + "provider": { + "String": "openai" + }, + "prompt": { + "String": "# Simplify: Code Review and Cleanup\n\nReview changes vs. origin for reuse, quality, and efficiency. Fix any issues found.\n\n## Phase 1: Identify Changes\n\nRun git diff (or git diff HEAD if there are staged changes) to see what changed. If there are no git changes, review the most recently modified files that the user mentioned or that you edited earlier in this conversation.\n\n## Phase 2: Launch Three Review Agents in Parallel\n\nUse the Agent tool to launch all three agents concurrently in a single message. Pass each agent the full diff so it has the complete context.\n\n### Agent 1: Code Reuse Review\n\nFor each change:\n\n1. Search for existing utilities and helpers that could replace newly written code. Use Grep to find similar patterns elsewhere in the codebase — common locations are utility directories, shared modules, and files adjacent to the changed ones.\n2. Flag any new function that duplicates existing functionality. Suggest the existing function to use instead.\n3. Flag any inline logic that could use an existing utility — hand-rolled string manipulation, manual path handling, custom environment checks, ad-hoc type guards, and similar patterns are common candidates.\n\nNote: This is a greenfield app, so focus on maximizing simplicity and don't worry about changing things to achieve it.\n\n### Agent 2: Code Quality Review\n\nReview the same changes for hacky patterns:\n\n1. Redundant state: state that duplicates existing state, cached values that could be derived, observers/effects that could be direct calls\n2. Parameter sprawl: adding new parameters to a function instead of generalizing or restructuring existing ones\n3. Copy-paste with slight variation: near-duplicate code blocks that should be unified with a shared abstraction\n4. Leaky abstractions: exposing internal details that should be encapsulated, or breaking existing abstraction boundaries\n5. Stringly-typed code: using raw strings where constants, enums (string unions), or branded types already exist in the codebase\n\nNote: This is a greenfield app, so be aggressive in optimizing quality.\n\n### Agent 3: Efficiency Review\n\nReview the same changes for efficiency:\n\n1. Unnecessary work: redundant computations, repeated file reads, duplicate network/API calls, N+1 patterns\n2. Missed concurrency: independent operations run sequentially when they could run in parallel\n3. Hot-path bloat: new blocking work added to startup or per-request/per-render hot paths\n4. Unnecessary existence checks: pre-checking file/resource existence before operating (TOCTOU anti-pattern) — operate directly and handle the error\n5. Memory: unbounded data structures, missing cleanup, event listener leaks\n6. Overly broad operations: reading entire files when only a portion is needed, loading all items when filtering for one\n\n## Phase 3: Fix Issues\n\nWait for all three agents to complete. Aggregate their findings and fix each issue directly. If a finding is a false positive or not worth addressing, note it and move on — do not argue with the finding, just skip it.\n\nWhen done, briefly summarize what was fixed (or confirm the code was already clean)." + }, + "label": { + "String": "Simplify (GPT-55)" + } + } + }, + "toolchain": { + "id": "toolchain", + "attrs": { + "shape": { + "String": "parallelogram" + }, + "label": { + "String": "Toolchain" + }, + "script": { + "String": "command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1" + }, + "max_retries": { + "Integer": 0 + }, + "provider": { + "String": "anthropic" + }, + "model": { + "String": "claude-opus-4-7" + } + } + }, + "fixup": { + "id": "fixup", + "attrs": { + "prompt": { + "String": "The verify step failed. Read the build output from context and fix all format, clippy, Rust test, docs, TypeScript typecheck/test, and build failures." + }, + "max_visits": { + "Integer": 3 + }, + "provider": { + "String": "anthropic" + }, + "model": { + "String": "claude-opus-4-7" + }, + "label": { + "String": "Fixup" + } + } + } + }, + "edges": [ + { + "from": "start", + "to": "toolchain", + "attrs": {} + }, + { + "from": "toolchain", + "to": "preflight_compile", + "attrs": { + "condition": { + "String": "outcome=succeeded" + } + } + }, + { + "from": "toolchain", + "to": "exit", + "attrs": {} + }, + { + "from": "preflight_compile", + "to": "preflight_lint", + "attrs": { + "condition": { + "String": "outcome=succeeded" + } + } + }, + { + "from": "preflight_compile", + "to": "exit", + "attrs": {} + }, + { + "from": "preflight_lint", + "to": "implement", + "attrs": { + "condition": { + "String": "outcome=succeeded" + } + } + }, + { + "from": "preflight_lint", + "to": "fix_lints", + "attrs": {} + }, + { + "from": "fix_lints", + "to": "preflight_lint", + "attrs": {} + }, + { + "from": "implement", + "to": "simplify_opus", + "attrs": {} + }, + { + "from": "simplify_opus", + "to": "simplify_gpt", + "attrs": {} + }, + { + "from": "simplify_gpt", + "to": "verify", + "attrs": {} + }, + { + "from": "verify", + "to": "exit", + "attrs": { + "condition": { + "String": "outcome=succeeded" + } + } + }, + { + "from": "verify", + "to": "fixup", + "attrs": {} + }, + { + "from": "fixup", + "to": "verify", + "attrs": {} + } + ], + "attrs": { + "goal": { + "String": "# Worker Control Bus Implementation Plan\n\n> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.\n\n**Goal:** Replace the server-to-worker stdin JSONL control pipe with a backend-agnostic worker control bus, implemented now with a local in-memory bus and delivered to workers over a worker-initiated WebSocket.\n\n**Architecture:** API handlers publish `WorkerControlEnvelope` messages to a `WorkerControlBus`; the worker WebSocket route subscribes to that bus and forwards ordered delivery frames to the worker. Workers track the last fully applied delivery id and reconnect with `?after=` after any unexpected WebSocket close. The first backend is an in-process `LocalWorkerControlBus` for local and single-node deployments. A Redis Streams backend must fit behind the same trait later, but Redis is explicitly out of scope for this implementation plan.\n\n**Tech Stack:** Rust, Axum WebSockets, tokio-tungstenite, UnixStream, async-trait or boxed async traits, tokio channels/Notify, worker JWT auth, existing `WorkerControlEnvelope`.\n\n---\n\n## Key Decisions\n\n- WebSocket fully replaces stdin control. Do not keep stdin JSONL as a compatibility path.\n- The worker protocol stays identical across local, single-node, ECS, and later SaaS deployments.\n- The server-side delivery backend is the only thing that varies by deployment.\n- This plan implements only `LocalWorkerControlBus`.\n- This plan does not add Redis dependencies, Redis configuration, Redis tests, Redis health checks, or Redis runtime behavior.\n- Redis Streams are covered only as a future backend contract so the local design does not paint us into a corner.\n- There is no `Latest` cursor. The first worker connection starts at the beginning of the run's retained control stream; reconnects resume after the worker's last fully applied delivery id.\n- Every WebSocket text frame is a delivery frame with an id and envelope. The worker advances `last_applied_id` only after applying the envelope.\n- Workers reconnect forever while the local run is not terminal, using backoff from 100ms, doubled after each failure, capped at 5s.\n- The worker must complete its first control-stream connection before starting or resuming workflow execution. Temporary first-connect failures wait and retry; they do not start a control-disconnected run.\n- Invalid cursor means the bus can no longer prove replay correctness. The worker treats it as fatal control-channel loss and fails/aborts the run as infrastructure failure, not as user cancellation.\n- WebSocket liveness is handled at the WebSocket layer with explicit ping/pong and timeout logic. The bus does not know about heartbeats.\n- ECS task launch, ECS stop/reconciliation, Redis-backed multi-node delivery, and remote hard-kill behavior are follow-up work.\n\n## Redis Fit Later: Out of Scope Now\n\nRedis should later implement the same `WorkerControlBus` API introduced here.\n\n- `publish(run_id, envelope)` maps to `XADD fabro:run:{run_id}:control ...`.\n- First `subscribe(run_id, Start)` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control 0-0`.\n- Reconnect `subscribe(run_id, After(id))` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control {id}`.\n- Local message ids use an opaque string format such as `local:1`; Redis message ids can use Redis stream ids such as `1716810000000-0`.\n- The WebSocket route should not care whether the subscription source is local memory or Redis.\n- The worker should not care whether the frame came from a local bus or Redis.\n- Redis trimming/retention, consumer groups, per-tenant key naming, TLS/auth, reconnect-after-redeploy semantics, and SaaS config validation are not part of this plan.\n\n## Proposed File Structure\n\n- Create `lib/crates/fabro-server/src/worker_control/mod.rs`\n - Owns the server-side control bus abstraction and re-exports the local backend.\n- Create `lib/crates/fabro-server/src/worker_control/bus.rs`\n - Defines `WorkerControlBus`, `WorkerControlDelivery`, `WorkerControlMessageId`, `WorkerControlCursor`, and bus errors.\n- Create `lib/crates/fabro-server/src/worker_control/local.rs`\n - Implements `LocalWorkerControlBus` using process memory.\n- Create `lib/crates/fabro-server/src/server/handler/worker_control.rs`\n - Adds the worker-only WebSocket route.\n- Modify `lib/crates/fabro-server/src/server.rs`\n - Adds the bus to `AppState`, replaces subprocess `RunAnswerTransport` sends with bus publishes, removes stdin pumping.\n- Modify `lib/crates/fabro-server/src/server/handler/mod.rs`\n - Registers the worker control route.\n- Modify `lib/crates/fabro-server/src/server/handler/lifecycle.rs`\n - Sends pause/unpause/cancel controls through the transport/bus where appropriate.\n- Modify `lib/crates/fabro-cli/src/commands/run/runner.rs`\n - Replaces stdin reading with worker WebSocket client handling.\n- Modify `lib/crates/fabro-cli/Cargo.toml`\n - Adds `tokio-tungstenite` as a direct dependency if needed.\n- Modify `lib/crates/fabro-interview/src/control_protocol.rs`\n - Adds pause/unpause control messages and a transport delivery frame type shared by server and worker.\n\n## Task 1: Define the Control Bus Contract\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/worker_control/mod.rs`\n- Create: `lib/crates/fabro-server/src/worker_control/bus.rs`\n- Modify: `lib/crates/fabro-server/src/lib.rs`\n\n- [ ] Add a private `worker_control` module in `fabro-server`.\n- [ ] Define `WorkerControlMessageId` as an opaque cloneable id rather than a numeric type.\n- [ ] Define `WorkerControlCursor` with `Start` and `After(WorkerControlMessageId)` variants.\n- [ ] Define `WorkerControlDelivery { id: WorkerControlMessageId, envelope: WorkerControlEnvelope }`.\n- [ ] Define `WorkerControlBus` with async `publish(run_id, envelope)` and `subscribe(run_id, cursor)` methods.\n- [ ] Make `subscribe` return a stream-like receiver owned by the caller, so the WebSocket handler can forward messages without knowing the backend.\n- [ ] Define explicit bus errors for closed backend, unavailable backend, invalid cursor, and publish timeout.\n- [ ] Document in code comments that `Start` maps to Redis stream id `0-0` and `After(id)` maps to Redis `XREAD` after that id, but do not add Redis code.\n- [ ] Add unit tests for id equality/debug formatting and cursor parsing from the optional `after` query parameter.\n- [ ] Test that absent `after` parses as `WorkerControlCursor::Start`.\n- [ ] Test that present `after=local:42` parses as `WorkerControlCursor::After(...)`.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 2: Implement the Local In-Memory Bus\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/worker_control/local.rs`\n- Test: `lib/crates/fabro-server/src/worker_control/local.rs`\n\n- [ ] Implement `LocalWorkerControlBus` as `Arc>>`.\n- [ ] Store messages per run in insertion order with a monotonic local sequence id.\n- [ ] Wake active subscribers when `publish` appends a message.\n- [ ] Support `subscribe(run_id, Start)` for first worker startup; it must replay retained messages from the beginning of the run control stream.\n- [ ] Support `subscribe(run_id, After(id))` so reconnect uses the same API that later maps to Redis `XREAD`.\n- [ ] Allow `publish` before the worker subscribes; retained messages must be visible to the first `Start` subscriber.\n- [ ] Trim retained local messages to a bounded per-run size so a disconnected local worker cannot grow memory without bound. Use a named constant with initial value 1024 messages per run.\n- [ ] Return a clear `invalid cursor` error when a subscriber asks for an id that has been trimmed or belongs to a different local stream.\n- [ ] Add a cleanup method for terminal runs so completed/cancelled runs can release retained control messages.\n- [ ] Test that messages publish in order.\n- [ ] Test that an active subscriber receives a message published after subscription.\n- [ ] Test that messages published before subscription are replayed to a `Start` subscriber.\n- [ ] Test that `After(id)` receives only later messages.\n- [ ] Test that trimming bounds retained messages and reports an invalid old cursor.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 3: Add Control Bus to Server State\n\n**Files:**\n- Modify: `lib/crates/fabro-server/src/server.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Add `worker_control_bus: Arc` to `AppState`.\n- [ ] Construct `LocalWorkerControlBus` in normal server state initialization.\n- [ ] Add a test-only way to inject a fake or local bus without exposing test helpers to production builds.\n- [ ] Keep demo/in-process execution behavior unchanged unless it currently depends on subprocess control.\n- [ ] Add a state construction test proving the default bus is local and available.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 4: Extend the Control Protocol\n\n**Files:**\n- Modify: `lib/crates/fabro-interview/src/control_protocol.rs`\n- Test: `lib/crates/fabro-interview/src/control_protocol.rs`\n\n- [ ] Add `WorkerControlEnvelope::pause_run()` and `WorkerControlEnvelope::unpause_run()` constructors.\n- [ ] Add `WorkerControlMessage::RunPause` serialized as `\"run.pause\"`.\n- [ ] Add `WorkerControlMessage::RunUnpause` serialized as `\"run.unpause\"`.\n- [ ] Add `WorkerControlDeliveryFrame { id: String, envelope: WorkerControlEnvelope }` as the WebSocket text-frame payload shared by server and worker.\n- [ ] Add round-trip serde tests for both new messages.\n- [ ] Add round-trip serde tests for `WorkerControlDeliveryFrame`.\n- [ ] Run `cargo nextest run -p fabro-interview control_protocol`.\n\n## Task 5: Share Worker Message Handling\n\n**Files:**\n- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n\n- [ ] Split `apply_worker_control_line(...)` into parsing and `apply_worker_control_message(...)`.\n- [ ] Route WebSocket delivery frames through `apply_worker_control_message(...)`.\n- [ ] Route `run.pause` to `RunControlState::request_pause()`.\n- [ ] Route `run.unpause` to `RunControlState::request_unpause()`.\n- [ ] Add a small in-memory applied-id dedupe set in the worker control task; ignore duplicate delivery ids before applying envelopes.\n- [ ] Update `last_applied_id` only after `apply_worker_control_message(...)` returns.\n- [ ] Treat all current control messages as idempotent under delivery-id dedupe. `run.steer` must not be applied twice for the same delivery id.\n- [ ] Keep control stream close behavior explicit: an unexpected close triggers reconnect; a fatal invalid cursor interrupts pending interviews and fails/aborts the run as control-channel loss.\n- [ ] Update existing stdin-era tests to exercise the shared message handler directly.\n- [ ] Add tests for pause and unpause routing.\n- [ ] Add a test proving duplicate delivery ids are not applied twice.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n\n## Task 6: Add Worker WebSocket Client\n\n**Files:**\n- Modify: `lib/crates/fabro-cli/Cargo.toml`\n- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs`\n\n- [ ] Add `tokio-tungstenite.workspace = true` as a direct `fabro-cli` dependency if the crate does not already have it.\n- [ ] Add a helper that builds the control-stream request for a `ServerTarget` and `RunId`.\n- [ ] For HTTP URLs, convert `http` to `ws` and `https` to `wss`.\n- [ ] For Unix socket paths, connect `tokio::net::UnixStream` and use `ws://fabro/api/v1/runs/{run_id}/worker/control-stream` for the handshake host/path.\n- [ ] Add the worker bearer token as an `Authorization: Bearer ...` request header.\n- [ ] On the first connection, omit the `after` query parameter so the server maps it to `WorkerControlCursor::Start`.\n- [ ] On reconnect, include `?after=` when `last_applied_id` is set.\n- [ ] Spawn a WebSocket control manager task in `execute(...)` after `ControlInterviewer`, `RunControlState`, `CancellationToken`, and `SteeringHub` are created, and before `operations::start` or `operations::resume`.\n- [ ] Gate `operations::start` and `operations::resume` on the first successful control-stream connection.\n- [ ] The control manager should keep reconnecting while the run is not locally terminal, with backoff starting at 100ms, doubling after each failure, and capped at 5s.\n- [ ] Deserialize each text frame into `WorkerControlDeliveryFrame`.\n- [ ] Apply each envelope through `apply_worker_control_message(...)`, then record the frame id as `last_applied_id`.\n- [ ] Respond to received WebSocket ping frames with pong frames.\n- [ ] Send worker-initiated ping frames every 15s.\n- [ ] Track pongs for worker-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness.\n- [ ] Treat normal close/error as reconnectable while the run is not terminal.\n- [ ] Treat HTTP 410 Gone or a WebSocket close reason of `invalid_cursor` as fatal control-channel loss.\n- [ ] On fatal control-channel loss, interrupt pending interviews and fail/abort the run with an infrastructure/control-channel error, not a user cancellation.\n- [ ] Wire fatal control-channel loss back into `execute(...)` so the worker returns an error instead of silently continuing workflow execution.\n- [ ] Add tests for URL/request construction for `http`, `https`, and Unix socket targets.\n- [ ] Add tests proving first connection has no `after` query and reconnect includes `after=`.\n- [ ] Add tests for reconnect backoff bounds.\n- [ ] Add tests for ping/pong timeout behavior using paused Tokio time.\n- [ ] Add a local Unix-socket WebSocket test proving the client can complete a handshake against an Axum route.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n\n## Task 7: Add Worker-Only Control Stream Route\n\n**Files:**\n- Create: `lib/crates/fabro-server/src/server/handler/worker_control.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/mod.rs`\n- Modify: `lib/crates/fabro-server/src/principal_middleware.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Add a narrow helper or extractor that accepts only authenticated worker principals whose token run id matches the route run id.\n- [ ] Add `GET /runs/{id}/worker/control-stream` to real API routes only.\n- [ ] Reject missing runs, terminal runs, and archived runs before upgrading.\n- [ ] Reject user JWTs and cross-run worker JWTs.\n- [ ] Parse absent `after` into `WorkerControlCursor::Start`.\n- [ ] Parse present `after` into `WorkerControlCursor::After(id)`.\n- [ ] On upgrade, call `worker_control_bus.subscribe(run_id, cursor)`.\n- [ ] If `subscribe` returns invalid cursor before upgrade, reject with HTTP 410 Gone.\n- [ ] Serialize each `WorkerControlDelivery` to `WorkerControlDeliveryFrame` and send it as a WebSocket text frame.\n- [ ] Send server-initiated ping frames every 15s.\n- [ ] Respond to received WebSocket ping frames with pong frames.\n- [ ] Track pongs for server-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness.\n- [ ] On timeout or disconnect, drop the bus subscription so local resources are released.\n- [ ] Do not store the live WebSocket sender in `ManagedRun`; the bus is now the delivery boundary.\n- [ ] Add tests for auth rejection, successful `Start` subscription, successful `After(id)` subscription, frame delivery, invalid cursor rejection as 410 Gone, ping/pong timeout cleanup, and cross-run worker rejection.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n\n## Task 8: Replace Server-Side Stdin Transport with Bus Publishing\n\n**Files:**\n- Modify: `lib/crates/fabro-server/src/server.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/lifecycle.rs`\n- Modify: `lib/crates/fabro-server/src/server/handler/pair.rs`\n- Test: `lib/crates/fabro-server/src/server/tests.rs`\n\n- [ ] Replace `RunAnswerTransport::Subprocess { control_tx }` with a bus-backed subprocess/worker variant.\n- [ ] Ensure the bus-backed variant has enough context to publish messages for the correct `RunId`.\n- [ ] Delete `pump_worker_control_jsonl(...)`.\n- [ ] Change `worker_command(...)` so `__run-worker` uses `stdin(Stdio::null())` instead of `stdin(Stdio::piped())`.\n- [ ] Remove child-stdin extraction and the control pump task from `execute_run_subprocess(...)`.\n- [ ] Keep stderr capture and worker exit handling unchanged.\n- [ ] Update `RunAnswerTransport` methods so answer, cancel, steer, interrupt, pair start/message/end all publish the existing envelope to `WorkerControlBus`.\n- [ ] Add `pause_run()` and `unpause_run()` methods on `RunAnswerTransport`.\n- [ ] Update pause/unpause lifecycle handlers to send `run.pause` and `run.unpause` over the bus for running workers.\n- [ ] Keep process signals only for hard cleanup paths such as cancel fallback, shutdown, terminal delete, and force removal.\n- [ ] Update existing tests that assert subprocess transport enqueue behavior to assert bus publish behavior instead.\n- [ ] Add a `worker_command` test proving stdin is null/not piped and `FABRO_WORKER_TOKEN` still travels only through env.\n- [ ] Run `cargo nextest run -p fabro-server worker_command`.\n\n## Task 9: End-to-End Local Control Flow Regression\n\n**Files:**\n- Test: `lib/crates/fabro-cli/tests/it/cmd/runner.rs`\n- Test: `lib/crates/fabro-server/tests/it/scenario/lifecycle.rs`\n\n- [ ] Add a test run where the worker connects to the control WebSocket and receives a cancel request through `LocalWorkerControlBus`.\n- [ ] Add a test where the server publishes a control message before the worker connects and the worker receives it on first connection.\n- [ ] Add a reconnect test where the worker applies message A, reconnects with `after=`, and then receives only message B.\n- [ ] Add an invalid-cursor test proving the worker reports control-channel loss as infrastructure failure rather than user cancellation.\n- [ ] Add a human-interview test proving submitted answers reach the worker through the bus and WebSocket.\n- [ ] Add a steer or interrupt test proving live agent controls still reach the worker transport.\n- [ ] Add a local Unix-socket server test proving the default local server target works without stdin.\n- [ ] Run `cargo nextest run -p fabro-cli --test it runner`.\n- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`.\n\n## Task 10: Final Verification\n\n**Files:**\n- Modify only if failures expose necessary fixes.\n\n- [ ] Run `cargo nextest run -p fabro-interview control_protocol`.\n- [ ] Run `cargo nextest run -p fabro-server worker_control`.\n- [ ] Run `cargo nextest run -p fabro-cli runner`.\n- [ ] Run `cargo nextest run -p fabro-server worker_command`.\n- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`.\n- [ ] Run `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`.\n- [ ] Confirm no code path still writes `WorkerControlEnvelope` to child stdin.\n- [ ] Confirm no Redis dependency, Redis config key, or Redis runtime path was added.\n- [ ] Confirm there is no `Latest` cursor or wait-for-subscriber behavior in the control bus.\n- [ ] Confirm WebSocket ping/pong handling is explicit on both worker and server.\n- [ ] Confirm `run.steer` and other controls are protected from duplicate delivery-id application.\n- [ ] Confirm `__run-worker` still scrubs `FABRO_WORKER_TOKEN` from process env before launching descendants.\n\n## Acceptance Criteria\n\n- All worker control traffic uses the worker control bus plus WebSocket last-mile transport.\n- Local Unix-socket server targets and remote HTTP(S) server targets both support worker control without Redis.\n- Local and single-node deployments require no external control-channel service.\n- The bus API can later be implemented by Redis Streams without changing API handlers or worker message handling.\n- First worker connection replays retained messages from the beginning of the run control stream; reconnect resumes after the last fully applied id.\n- Invalid cursor is the only fatal control-stream replay failure and is surfaced as infrastructure/control-channel failure, not user cancellation.\n- WebSocket liveness is explicit and backend-agnostic.\n- Existing run event/blob/artifact HTTP paths are unchanged.\n- Existing worker JWT scope rules remain authoritative.\n- Temporary WebSocket disconnects reconnect and replay through the bus; only unrecoverable replay loss reports worker-control-unavailable behavior.\n- Worker stdout/stderr behavior remains unchanged except that stdin is no longer a control channel.\n- Redis is clearly documented as future work and is not required by this plan.\n" + }, + "rankdir": { + "String": "LR" + }, + "model_stylesheet": { + "String": "\n * { model: claude-opus-4-7; }\n " + } + } + }, + "graph_source": "digraph ImplementPlan {\n graph [\n goal=\"Implement and simplify\",\n model_stylesheet=\"\n * { model: claude-opus-4-7; }\n \"\n ]\n rankdir=LR\n\n start [shape=Mdiamond, label=\"Start\"]\n exit [shape=Msquare, label=\"Exit\"]\n\n toolchain [label=\"Toolchain\", shape=parallelogram, script=\"command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1\", max_retries=0]\n preflight_compile [label=\"Preflight Compile\", shape=parallelogram, script=\"cargo check -q --workspace 2>&1\", max_retries=0]\n preflight_lint [label=\"Preflight Lint\", shape=parallelogram, script=\"cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1\", max_retries=0]\n fix_lints [label=\"Fix Lints\", prompt=\"The preflight lint step failed. Read the build output from context and fix all clippy lint warnings.\", max_visits=3]\n 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=\"gpt-55\", reasoning_effort=\"xhigh\"]\n simplify_opus [label=\"Simplify (Opus)\", prompt=\"@prompts/simplify.md\"]\n simplify_gpt [label=\"Simplify (GPT-55)\", prompt=\"@prompts/simplify.md\", model=\"gpt-55\"]\n 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\"]\n 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.\", max_visits=3]\n\n start -> toolchain\n toolchain -> preflight_compile [condition=\"outcome=succeeded\"]\n toolchain -> exit\n preflight_compile -> preflight_lint [condition=\"outcome=succeeded\"]\n preflight_compile -> exit\n preflight_lint -> implement [condition=\"outcome=succeeded\"]\n preflight_lint -> fix_lints\n fix_lints -> preflight_lint\n implement -> simplify_opus -> simplify_gpt -> verify\n verify -> exit [condition=\"outcome=succeeded\"]\n verify -> fixup\n fixup -> verify\n}\n", + "workflow_slug": "implement-plan", + "source_directory": "/Users/bhelmkamp/p/fabro-sh/fabro", + "provenance": { + "server": { + "version": "0.246.0-nightly.0" + }, + "client": { + "user_agent": "fabro-cli/0.246.0-nightly.0", + "name": "fabro-cli", + "version": "0.246.0-nightly.0" + }, + "subject": { + "kind": "user", + "identity": { + "issuer": "https://github.com", + "subject": "19" + }, + "login": "brynary", + "auth_method": "github", + "avatar_url": "https://avatars.githubusercontent.com/u/19?v=4" + } + }, + "manifest_blob": "1da0686f5910966d75c4e28c00ba963e045a0adf66d7533b26a5b8a32035baae", + "definition_blob": "f13f100f4f1f386c63f22b375c331077f1c1b3bd2df0dd17388b26786f23b9a6", + "git": { + "origin_url": "https://github.com/fabro-sh/fabro", + "branch": "main", + "sha": "c767db897f14b4e1bf4c7572f83b164c5650ce58", + "dirty": "dirty", + "push_outcome": { + "type": "not_attempted" + } + } + }, + "web_url": "http://127.0.0.1:32276/runs/01KSNP2DXVS2TD1HEASQAFBGFK", + "start": null, + "status": { + "kind": "starting" + }, + "status_updated_at": "2026-05-27T21:39:53.971436Z", + "last_event_at": "2026-05-27T21:40:07.698938Z", + "pending_control": null, + "checkpoints": [], + "conclusion": null, + "sandbox": { + "kind": "ready", + "plan": { + "provider": "daytona" + }, + "instance": { + "provider": "daytona", + "snapshot": "fabro-fdb28dec-1233-892c-b9d7-9f88f8353e7a", + "runtime": { + "id": "fabro-01KSNP2DXVS2TD1HEASQAFBGFK", + "working_directory": "/home/daytona/workspace/fabro", + "repo_cloned": true, + "clone_origin_url": "https://github.com/fabro-sh/fabro", + "clone_branch": "main", + "workspace_root": "/home/daytona/workspace", + "repos_root": "/home/daytona/repos", + "primary_repo_path": "/home/daytona/repos/fabro-sh/fabro", + "primary_repo_link": "/home/daytona/workspace/fabro" + } + } + }, + "pull_request": null, + "superseded_by": null, + "pending_interviews": {}, + "stages": {} +} \ No newline at end of file