From a0f0f439ee6d9509234f1404d849a1698855131b Mon Sep 17 00:00:00 2001 From: Fabro Date: Wed, 27 May 2026 18:58:48 -0400 Subject: [PATCH] =?UTF-8?q?checkpoint=20=E2=9A=92=EF=B8=8F=20Generated=20w?= =?UTF-8?q?ith=20[Fabro](https://fabro.sh)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- run.json | 482 +++++++++++++++++- stages/004-preflight_lint@1/output.log | 1 + .../004-preflight_lint@1/script_timing.json | 8 + stages/004-preflight_lint@1/status.json | 6 + stages/005-implement@1/prompt.md | 297 +++++++++++ stages/005-implement@1/provider_used.json | 6 + stages/005-implement@1/response.md | 61 +++ 7 files changed, 852 insertions(+), 9 deletions(-) create mode 100644 stages/004-preflight_lint@1/output.log create mode 100644 stages/004-preflight_lint@1/script_timing.json create mode 100644 stages/004-preflight_lint@1/status.json create mode 100644 stages/005-implement@1/prompt.md create mode 100644 stages/005-implement@1/provider_used.json create mode 100644 stages/005-implement@1/response.md diff --git a/run.json b/run.json index 5707c88d6..4075df371 100644 --- a/run.json +++ b/run.json @@ -505,7 +505,7 @@ "kind": "running" }, "status_updated_at": "2026-05-27T21:40:08.061760Z", - "last_event_at": "2026-05-27T21:42:32.328542Z", + "last_event_at": "2026-05-27T22:58:48.037234Z", "pending_control": null, "checkpoints": [ { @@ -690,9 +690,9 @@ } }, { - "seq": 0, + "seq": 48, "checkpoint": { - "timestamp": "2026-05-27T21:44:55.791425Z", + "timestamp": "2026-05-27T21:44:59.873073Z", "current_node": "preflight_lint", "completed_nodes": [ "start", @@ -701,21 +701,126 @@ "preflight_lint" ], "node_retries": {}, + "context_values": { + "graph.model_stylesheet": "\n * { model: claude-opus-4-7; }\n ", + "internal.retry_count.preflight_lint": 0, + "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "graph.goal": "# 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", + "thread.start.current_node": "toolchain", + "failure_signature": "", + "internal.run_id": "01KSNP2DXVS2TD1HEASQAFBGFK", + "current_node": "preflight_lint", + "thread.preflight_compile.current_node": "preflight_lint", + "graph.rankdir": "LR", + "internal.retry_count.toolchain": 0, + "internal.work_dir": "/home/daytona/workspace/fabro", + "internal.retry_count.preflight_compile": 0, + "thread.toolchain.current_node": "preflight_compile", + "internal.node_visit_count": 1, + "internal.retry_count.start": 0, + "failure_class": "", + "internal.fidelity": "compact", + "internal.thread_id": "preflight_compile", + "outcome": "succeeded" + }, + "node_outcomes": { + "preflight_lint": { + "status": "succeeded", + "context_updates": { + "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126" + }, + "notes": "Script completed: cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", + "usage": null, + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 0, + "tool_time_ms": 143452, + "active_time_ms": 143452 + } + }, + "toolchain": { + "status": "succeeded", + "context_updates": { + "command.output": "blob://sha256/fc14b2ba2d770e5cd3169df7a29525c962adfc4cfa3097b9098c63ebd61a748c" + }, + "notes": "Script completed: command -v cargo >/dev/null || { curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y && sudo ln -sf $HOME/.cargo/bin/* /usr/local/bin/; }; cargo --version 2>&1", + "usage": null, + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 0, + "tool_time_ms": 1331, + "active_time_ms": 1331 + } + }, + "preflight_compile": { + "status": "succeeded", + "context_updates": { + "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126" + }, + "notes": "Script completed: cargo check -q --workspace 2>&1", + "usage": null, + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 0, + "tool_time_ms": 130559, + "active_time_ms": 130559 + } + }, + "start": { + "status": "succeeded", + "usage": null + } + }, + "next_node_id": "implement", + "git_commit_sha": "49f19797eb68d07d201ab25abe1a7812d521bc2b", + "node_visits": { + "toolchain": 1, + "preflight_compile": 1, + "preflight_lint": 1, + "start": 1 + } + }, + "diff": { + "summary": { + "files_changed": 0, + "additions": 0, + "deletions": 0 + } + } + }, + { + "seq": 0, + "checkpoint": { + "timestamp": "2026-05-27T22:58:48.113127Z", + "current_node": "implement", + "completed_nodes": [ + "start", + "toolchain", + "preflight_compile", + "preflight_lint", + "implement" + ], + "node_retries": {}, "context_values": { "thread.start.current_node": "toolchain", "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "internal.retry_count.implement": 0, + "last_stage": "implement", "internal.retry_count.start": 0, "thread.preflight_compile.current_node": "preflight_lint", "internal.retry_count.preflight_lint": 0, - "internal.thread_id": "preflight_compile", + "internal.thread_id": "preflight_lint", "internal.retry_count.preflight_compile": 0, + "thread.preflight_lint.current_node": "implement", "failure_signature": "", "thread.toolchain.current_node": "preflight_compile", - "current_node": "preflight_lint", + "current_node": "implement", "internal.node_visit_count": 1, + "response.implement": "Implemented the worker control bus plan end-to-end.\n\n### What changed\n\n- Added shared protocol support:\n - `run.pause` / `run.unpause`\n - `WorkerControlDeliveryFrame { id, envelope }`\n\n- Added server-side worker control bus:\n - `WorkerControlBus` abstraction\n - opaque `WorkerControlMessageId`\n - `WorkerControlCursor::{Start, After(...)}`\n - `LocalWorkerControlBus` with ordered replay, `after` resume, bounded retention, invalid-cursor detection, and cleanup\n\n- Added server WebSocket route:\n - `GET /api/v1/runs/{id}/worker/control-stream`\n - worker-JWT-only and run-scoped\n - rejects missing/terminal/archived runs before upgrade\n - invalid cursor returns HTTP 410 before upgrade\n - sends JSON delivery frames\n - explicit ping/pong liveness\n\n- Removed worker stdin JSONL runtime transport:\n - server now publishes all controls through the bus\n - `__run-worker` subprocess stdin is configured as null\n - no `control_tx` / stdin control pump remains\n\n- Added CLI worker WebSocket control manager:\n - connects before start/resume for non-terminal runs\n - supports HTTP → WS, HTTPS → WSS, and Unix socket targets\n - reconnects with exponential backoff capped at 5s\n - first connect omits `after`; reconnect uses last fully applied id\n - applies frames only once by delivery id\n - updates `last_applied_id` only after successful application\n - handles pause/unpause/cancel/answer/steer/interrupt/pair controls\n - explicit ping/pong liveness and timeout\n - HTTP 410 / `invalid_cursor` close is fatal control-channel loss\n\n- Added/updated tests across protocol, bus, route, CLI, server transport, and lifecycle/cancel flows.\n\n- Fixed a lifecycle regression where cancellation at a human gate could race with the detached completion guard and emit a duplicate terminal failure.\n\n### Verification run\n\nPassed:\n\n- `cargo nextest run -p fabro-interview control_protocol`\n- `cargo nextest run -p fabro-server worker_control`\n- `cargo nextest run -p fabro-cli runner`\n- `cargo nextest run -p fabro-server worker_command`\n- `cargo nextest run -p fabro-cli --test it runner`\n- `cargo nextest run -p fabro-server --features test-support --test it lifecycle`\n- `cargo +nightly-2026-04-14 fmt --check --all`\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`\n- `git diff --check`\n\nAlso confirmed by grep:\n\n- no worker-control stdin pump / `control_tx` remnants in server or CLI runtime code\n- no `Latest` worker-control cursor\n- no Redis dependency/config/runtime path added, only future-backend comments in the bus contract.", "graph.rankdir": "LR", "outcome": "succeeded", "failure_class": "", + "last_response": "Implemented the worker control bus plan end-to-end.\n\n### What changed\n\n- Added shared protocol support:\n - `run.pause` / `run.unpause`\n - `WorkerControlDeliveryFrame { id, envelope }`\n\n- Added serve", "internal.run_id": "01KSNP2DXVS2TD1HEASQAFBGFK", "graph.goal": "# 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", "internal.work_dir": "/home/daytona/workspace/fabro", @@ -769,13 +874,56 @@ "start": { "status": "succeeded", "usage": null + }, + "implement": { + "status": "succeeded", + "context_updates": { + "response.implement": "Implemented the worker control bus plan end-to-end.\n\n### What changed\n\n- Added shared protocol support:\n - `run.pause` / `run.unpause`\n - `WorkerControlDeliveryFrame { id, envelope }`\n\n- Added server-side worker control bus:\n - `WorkerControlBus` abstraction\n - opaque `WorkerControlMessageId`\n - `WorkerControlCursor::{Start, After(...)}`\n - `LocalWorkerControlBus` with ordered replay, `after` resume, bounded retention, invalid-cursor detection, and cleanup\n\n- Added server WebSocket route:\n - `GET /api/v1/runs/{id}/worker/control-stream`\n - worker-JWT-only and run-scoped\n - rejects missing/terminal/archived runs before upgrade\n - invalid cursor returns HTTP 410 before upgrade\n - sends JSON delivery frames\n - explicit ping/pong liveness\n\n- Removed worker stdin JSONL runtime transport:\n - server now publishes all controls through the bus\n - `__run-worker` subprocess stdin is configured as null\n - no `control_tx` / stdin control pump remains\n\n- Added CLI worker WebSocket control manager:\n - connects before start/resume for non-terminal runs\n - supports HTTP → WS, HTTPS → WSS, and Unix socket targets\n - reconnects with exponential backoff capped at 5s\n - first connect omits `after`; reconnect uses last fully applied id\n - applies frames only once by delivery id\n - updates `last_applied_id` only after successful application\n - handles pause/unpause/cancel/answer/steer/interrupt/pair controls\n - explicit ping/pong liveness and timeout\n - HTTP 410 / `invalid_cursor` close is fatal control-channel loss\n\n- Added/updated tests across protocol, bus, route, CLI, server transport, and lifecycle/cancel flows.\n\n- Fixed a lifecycle regression where cancellation at a human gate could race with the detached completion guard and emit a duplicate terminal failure.\n\n### Verification run\n\nPassed:\n\n- `cargo nextest run -p fabro-interview control_protocol`\n- `cargo nextest run -p fabro-server worker_control`\n- `cargo nextest run -p fabro-cli runner`\n- `cargo nextest run -p fabro-server worker_command`\n- `cargo nextest run -p fabro-cli --test it runner`\n- `cargo nextest run -p fabro-server --features test-support --test it lifecycle`\n- `cargo +nightly-2026-04-14 fmt --check --all`\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`\n- `git diff --check`\n\nAlso confirmed by grep:\n\n- no worker-control stdin pump / `control_tx` remnants in server or CLI runtime code\n- no `Latest` worker-control cursor\n- no Redis dependency/config/runtime path added, only future-backend comments in the bus contract.", + "last_response": "Implemented the worker control bus plan end-to-end.\n\n### What changed\n\n- Added shared protocol support:\n - `run.pause` / `run.unpause`\n - `WorkerControlDeliveryFrame { id, envelope }`\n\n- Added serve", + "last_stage": "implement" + }, + "notes": "Stage completed: implement", + "usage": { + "input": { + "usage": { + "model": { + "provider": "openai", + "model_id": "gpt-5.5" + }, + "tokens": { + "input_tokens": 3359993, + "output_tokens": 62719, + "reasoning_tokens": 31380, + "cache_read_tokens": 43824128, + "cache_write_tokens": 0 + } + }, + "facts": { + "algorithm": "openai" + } + }, + "total_usd_micros": 41534999 + }, + "files_touched": [ + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/server/handler/worker_control.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/worker_control/bus.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/worker_control/local.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/worker_control/mod.rs" + ], + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 2814420, + "tool_time_ms": 1504977, + "active_time_ms": 4319397 + } } }, - "next_node_id": "implement", + "next_node_id": "simplify_opus", "node_visits": { "start": 1, "toolchain": 1, "preflight_lint": 1, + "implement": 1, "preflight_compile": 1 } }, @@ -808,6 +956,301 @@ "superseded_by": null, "pending_interviews": {}, "stages": { + "implement@1": { + "first_event_seq": 51, + "prompt": null, + "response": null, + "completion": null, + "provider_used": { + "mode": "agent", + "provider": "openai", + "model": "gpt-5.5", + "reasoning_effort": "xhigh" + }, + "diff": null, + "script_invocation": null, + "script_timing": null, + "parallel_results": null, + "output": null, + "started_at": "2026-05-27T21:44:59.876364Z", + "handler": "agent", + "usage": { + "input_tokens": 3359993, + "output_tokens": 62719, + "total_tokens": 47278220, + "reasoning_tokens": 31380, + "cache_read_tokens": 43824128, + "cache_write_tokens": 0, + "total_usd_micros": 41534999 + }, + "model": { + "provider": "openai", + "model_id": "gpt-5.5" + }, + "todos": { + "kind": "openai_plan", + "list_id": "openai_plan:472a10ab-7423-4272-9055-eff0a440f541", + "items": [ + { + "id": "70445a19f2b871d2", + "status": "completed", + "order": 0, + "subject": "Define shared control protocol and server WorkerControlBus contract" + }, + { + "id": "13038614eb922561", + "status": "completed", + "order": 1, + "subject": "Implement LocalWorkerControlBus with replay, after-cursor, trimming, cleanup, and tests" + }, + { + "id": "9034f7f00b2a72ab", + "status": "completed", + "order": 2, + "subject": "Add bus to server state and replace subprocess stdin control with bus publishing" + }, + { + "id": "6d3315f94f20bed5", + "status": "completed", + "order": 3, + "subject": "Add worker-only WebSocket route with worker JWT scoping, replay cursor handling, frames, and ping/pong liveness" + }, + { + "id": "ff4ec9f7aa35b244", + "status": "completed", + "order": 4, + "subject": "Replace CLI worker stdin reader with WebSocket control manager, shared message application, reconnect/after handling, dedupe, liveness, and fatal invalid-cursor handling" + }, + { + "id": "55876ccaf4abf52a", + "status": "completed", + "order": 5, + "subject": "Add unit/integration regressions for route auth/replay, CLI request/reconnect/liveness/Unix handshake, server command stdin/token behavior, cancel/answer local flows" + }, + { + "id": "14008af6f3f59e11", + "status": "completed", + "order": 6, + "subject": "Run targeted tests, lifecycle integration tests, formatting, clippy, and cleanup greps" + } + ] + }, + "subagents": [ + { + "agent_id": "a7258f87", + "depth": 1, + "task": "Inspect fabro-server worker auth/principal middleware and route registration. Report how to add a worker-only WebSocket route requiring worker principal run_id to match path run id, and identify relevant tests/helpers in server/tests.rs. Do not modify files.", + "status": { + "kind": "closed" + } + }, + { + "agent_id": "537c61a2", + "depth": 1, + "task": "Inspect fabro-cli runner.rs around stdin WorkerControlEnvelope handling, execute lifecycle, ServerTarget URL handling, and tests. Report exact structures/functions to change for WebSocket control manager. Do not modify files.", + "status": { + "kind": "closed" + } + }, + { + "agent_id": "2703e424", + "depth": 1, + "task": "Inspect fabro-server server.rs RunAnswerTransport, execute_run_subprocess, worker_command, lifecycle/pair handlers, and relevant tests. Report minimal changes to replace mpsc stdin transport with a bus-backed transport. Do not modify files.", + "status": { + "kind": "closed" + } + }, + { + "agent_id": "55b0fae3", + "depth": 1, + "task": "Review the current worker control bus/WebSocket implementation and propose concrete missing tests/fixes for fabro-server route coverage. Focus on lib/crates/fabro-server/src/server/tests.rs helpers for authenticated requests/WebSocket clients and how to add tests for /api/v1/runs/{id}/worker/control-stream. Do not modify files; return concise implementation guidance with existing helper names/locations.", + "status": { + "kind": "completed", + "success": true, + "turns_used": 9 + } + }, + { + "agent_id": "3fc6d823", + "depth": 1, + "task": "Review fabro-cli worker WebSocket client implementation in lib/crates/fabro-cli/src/commands/run/runner.rs and identify compile/clippy risks and missing unit tests for Task 6 (request construction, after behavior, backoff, ping/pong timeout, Unix handshake). Do not modify files; return concise guidance and suggested minimal tests.", + "status": { + "kind": "completed", + "success": true, + "turns_used": 9 + } + } + ], + "permission_level": "full", + "agent_tools": [ + { + "name": "apply_patch", + "description": "Use the `apply_patch` tool to edit files. This is a FREEFORM tool, so do not wrap the patch in JSON.", + "source": { + "kind": "native" + }, + "category": "write", + "invoked": true + }, + { + "name": "close_agent", + "description": "Close a running subagent that is no longer needed.", + "source": { + "kind": "native" + }, + "category": "subagent", + "invoked": false + }, + { + "name": "glob", + "description": "Find files by file names using a glob pattern. Use path to choose the search root. Prefer this over shell find or ls when locating repository files.", + "source": { + "kind": "native" + }, + "category": "read", + "invoked": true + }, + { + "name": "grep", + "description": "Search file contents with a regex pattern. Use path to choose the search root, glob_filter to limit matching files, case_insensitive for case folding, and max_results to cap output.", + "source": { + "kind": "native" + }, + "category": "read", + "invoked": true + }, + { + "name": "read_file", + "description": "Read files before editing them. Returns line-numbered text and supports offset/limit for large files. Use this instead of shell cat, head, tail, or sed when inspecting repository files.", + "source": { + "kind": "native" + }, + "category": "read", + "invoked": true + }, + { + "name": "request_user_input", + "description": "Ask the human one or more questions and wait for their answers before continuing this stage.", + "source": { + "kind": "native" + }, + "category": "other", + "invoked": false + }, + { + "name": "send_input", + "description": "Send a follow-up message to a running subagent when new information or corrected instructions are needed.", + "source": { + "kind": "native" + }, + "category": "subagent", + "invoked": false + }, + { + "name": "shell", + "description": "Execute shell commands for terminal operations, package managers, tests and builds. Use dedicated tools for file reads, file edits, filename searches, and content searches. Provide timeout_ms for long-running commands.", + "source": { + "kind": "native" + }, + "category": "shell", + "invoked": true + }, + { + "name": "spawn_agent", + "description": "Spawn a subagent for independent work or context isolation. Use it for tasks that can proceed separately, and avoid duplicating the same work in the parent session.", + "source": { + "kind": "native" + }, + "category": "subagent", + "invoked": true + }, + { + "name": "update_plan", + "description": "Update the multi-step plan for the current task. Submit the entire plan; existing steps are reconciled by exact step text.", + "source": { + "kind": "native" + }, + "category": "other", + "invoked": true + }, + { + "name": "wait", + "description": "Wait for a subagent to complete, then use the result to synthesize the outcome for the user.", + "source": { + "kind": "native" + }, + "category": "subagent", + "invoked": true + }, + { + "name": "web_fetch", + "description": "Fetch content from a URL that starts with http:// or https://. Pass a prompt to extract specific information or summarize the page; omit prompt to return the page content.", + "source": { + "kind": "native" + }, + "category": "other", + "invoked": false + }, + { + "name": "web_search", + "description": "Search the web using Brave Search when current external information is needed. Returns result titles, URLs, and descriptions; use web_fetch for a specific URL.", + "source": { + "kind": "native" + }, + "category": "other", + "invoked": false + }, + { + "name": "write_file", + "description": "Create new files, or overwrite an existing file only when replacement is explicitly intended. Prefer edit_file for targeted changes to existing files because write_file overwrites the full file content.", + "source": { + "kind": "native" + }, + "category": "write", + "invoked": true + } + ], + "context_window": { + "provider": "openai", + "model": "gpt-5.5", + "context_window_tokens": 272000, + "input_tokens": 209809, + "usage_percent": 77.13566176470589, + "count_method": "response_usage_scaled_breakdown", + "staleness": "live", + "generated_at": "2026-05-27T22:58:48.035992Z", + "event_seq": 1281, + "breakdown": [ + { + "category": "system_prompt", + "tokens": 926, + "usage_percent": 0.34044117647058825 + }, + { + "category": "tools", + "tokens": 1318, + "usage_percent": 0.48455882352941176 + }, + { + "category": "memory", + "tokens": 3137, + "usage_percent": 1.1533088235294118 + }, + { + "category": "conversation", + "tokens": 204422, + "usage_percent": 75.15514705882353 + }, + { + "category": "other", + "tokens": 6, + "usage_percent": 0.0022058823529411764 + } + ], + "warnings": [] + }, + "state": "running" + }, "toolchain@1": { "first_event_seq": 21, "prompt": null, @@ -860,7 +1303,12 @@ "first_event_seq": 41, "prompt": null, "response": null, - "completion": null, + "completion": { + "outcome": "succeeded", + "notes": "Script completed: cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", + "failure_reason": null, + "timestamp": "2026-05-27T21:44:55.790026Z" + }, "provider_used": null, "diff": null, "script_invocation": { @@ -868,11 +1316,27 @@ "command": "exec 2>&1\ncargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", "language": "shell" }, - "script_timing": null, + "script_timing": { + "output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "exit_code": 0, + "duration_ms": 143452, + "termination": "exited", + "output_bytes": 0, + "live_streaming": false + }, "parallel_results": null, "output": null, + "output_bytes": 0, + "live_streaming": false, + "termination": "exited", "started_at": "2026-05-27T21:42:32.328033Z", "handler": "command", + "timing": { + "wall_time_ms": 143461, + "inference_time_ms": 0, + "tool_time_ms": 143452, + "active_time_ms": 143452 + }, "usage": { "input_tokens": 0, "output_tokens": 0, @@ -881,7 +1345,7 @@ "cache_read_tokens": 0, "cache_write_tokens": 0 }, - "state": "running" + "state": "succeeded" }, "start@1": { "first_event_seq": 17, diff --git a/stages/004-preflight_lint@1/output.log b/stages/004-preflight_lint@1/output.log new file mode 100644 index 000000000..d87ba9545 --- /dev/null +++ b/stages/004-preflight_lint@1/output.log @@ -0,0 +1 @@ +blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126 \ No newline at end of file diff --git a/stages/004-preflight_lint@1/script_timing.json b/stages/004-preflight_lint@1/script_timing.json new file mode 100644 index 000000000..84980af84 --- /dev/null +++ b/stages/004-preflight_lint@1/script_timing.json @@ -0,0 +1,8 @@ +{ + "output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "exit_code": 0, + "duration_ms": 143452, + "termination": "exited", + "output_bytes": 0, + "live_streaming": false +} \ No newline at end of file diff --git a/stages/004-preflight_lint@1/status.json b/stages/004-preflight_lint@1/status.json new file mode 100644 index 000000000..ac36c2518 --- /dev/null +++ b/stages/004-preflight_lint@1/status.json @@ -0,0 +1,6 @@ +{ + "outcome": "succeeded", + "notes": "Script completed: cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1", + "failure_reason": null, + "timestamp": "2026-05-27T21:44:55.790026Z" +} \ No newline at end of file diff --git a/stages/005-implement@1/prompt.md b/stages/005-implement@1/prompt.md new file mode 100644 index 000000000..3ff925182 --- /dev/null +++ b/stages/005-implement@1/prompt.md @@ -0,0 +1,297 @@ +Goal: # Worker Control Bus Implementation Plan + +> **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. + +**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. + +**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. + +**Tech Stack:** Rust, Axum WebSockets, tokio-tungstenite, UnixStream, async-trait or boxed async traits, tokio channels/Notify, worker JWT auth, existing `WorkerControlEnvelope`. + +--- + +## Key Decisions + +- WebSocket fully replaces stdin control. Do not keep stdin JSONL as a compatibility path. +- The worker protocol stays identical across local, single-node, ECS, and later SaaS deployments. +- The server-side delivery backend is the only thing that varies by deployment. +- This plan implements only `LocalWorkerControlBus`. +- This plan does not add Redis dependencies, Redis configuration, Redis tests, Redis health checks, or Redis runtime behavior. +- Redis Streams are covered only as a future backend contract so the local design does not paint us into a corner. +- 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. +- Every WebSocket text frame is a delivery frame with an id and envelope. The worker advances `last_applied_id` only after applying the envelope. +- Workers reconnect forever while the local run is not terminal, using backoff from 100ms, doubled after each failure, capped at 5s. +- 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. +- 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. +- WebSocket liveness is handled at the WebSocket layer with explicit ping/pong and timeout logic. The bus does not know about heartbeats. +- ECS task launch, ECS stop/reconciliation, Redis-backed multi-node delivery, and remote hard-kill behavior are follow-up work. + +## Redis Fit Later: Out of Scope Now + +Redis should later implement the same `WorkerControlBus` API introduced here. + +- `publish(run_id, envelope)` maps to `XADD fabro:run:{run_id}:control ...`. +- First `subscribe(run_id, Start)` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control 0-0`. +- Reconnect `subscribe(run_id, After(id))` maps to `XREAD BLOCK ... STREAMS fabro:run:{run_id}:control {id}`. +- Local message ids use an opaque string format such as `local:1`; Redis message ids can use Redis stream ids such as `1716810000000-0`. +- The WebSocket route should not care whether the subscription source is local memory or Redis. +- The worker should not care whether the frame came from a local bus or Redis. +- 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. + +## Proposed File Structure + +- Create `lib/crates/fabro-server/src/worker_control/mod.rs` + - Owns the server-side control bus abstraction and re-exports the local backend. +- Create `lib/crates/fabro-server/src/worker_control/bus.rs` + - Defines `WorkerControlBus`, `WorkerControlDelivery`, `WorkerControlMessageId`, `WorkerControlCursor`, and bus errors. +- Create `lib/crates/fabro-server/src/worker_control/local.rs` + - Implements `LocalWorkerControlBus` using process memory. +- Create `lib/crates/fabro-server/src/server/handler/worker_control.rs` + - Adds the worker-only WebSocket route. +- Modify `lib/crates/fabro-server/src/server.rs` + - Adds the bus to `AppState`, replaces subprocess `RunAnswerTransport` sends with bus publishes, removes stdin pumping. +- Modify `lib/crates/fabro-server/src/server/handler/mod.rs` + - Registers the worker control route. +- Modify `lib/crates/fabro-server/src/server/handler/lifecycle.rs` + - Sends pause/unpause/cancel controls through the transport/bus where appropriate. +- Modify `lib/crates/fabro-cli/src/commands/run/runner.rs` + - Replaces stdin reading with worker WebSocket client handling. +- Modify `lib/crates/fabro-cli/Cargo.toml` + - Adds `tokio-tungstenite` as a direct dependency if needed. +- Modify `lib/crates/fabro-interview/src/control_protocol.rs` + - Adds pause/unpause control messages and a transport delivery frame type shared by server and worker. + +## Task 1: Define the Control Bus Contract + +**Files:** +- Create: `lib/crates/fabro-server/src/worker_control/mod.rs` +- Create: `lib/crates/fabro-server/src/worker_control/bus.rs` +- Modify: `lib/crates/fabro-server/src/lib.rs` + +- [ ] Add a private `worker_control` module in `fabro-server`. +- [ ] Define `WorkerControlMessageId` as an opaque cloneable id rather than a numeric type. +- [ ] Define `WorkerControlCursor` with `Start` and `After(WorkerControlMessageId)` variants. +- [ ] Define `WorkerControlDelivery { id: WorkerControlMessageId, envelope: WorkerControlEnvelope }`. +- [ ] Define `WorkerControlBus` with async `publish(run_id, envelope)` and `subscribe(run_id, cursor)` methods. +- [ ] Make `subscribe` return a stream-like receiver owned by the caller, so the WebSocket handler can forward messages without knowing the backend. +- [ ] Define explicit bus errors for closed backend, unavailable backend, invalid cursor, and publish timeout. +- [ ] 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. +- [ ] Add unit tests for id equality/debug formatting and cursor parsing from the optional `after` query parameter. +- [ ] Test that absent `after` parses as `WorkerControlCursor::Start`. +- [ ] Test that present `after=local:42` parses as `WorkerControlCursor::After(...)`. +- [ ] Run `cargo nextest run -p fabro-server worker_control`. + +## Task 2: Implement the Local In-Memory Bus + +**Files:** +- Create: `lib/crates/fabro-server/src/worker_control/local.rs` +- Test: `lib/crates/fabro-server/src/worker_control/local.rs` + +- [ ] Implement `LocalWorkerControlBus` as `Arc>>`. +- [ ] Store messages per run in insertion order with a monotonic local sequence id. +- [ ] Wake active subscribers when `publish` appends a message. +- [ ] Support `subscribe(run_id, Start)` for first worker startup; it must replay retained messages from the beginning of the run control stream. +- [ ] Support `subscribe(run_id, After(id))` so reconnect uses the same API that later maps to Redis `XREAD`. +- [ ] Allow `publish` before the worker subscribes; retained messages must be visible to the first `Start` subscriber. +- [ ] 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. +- [ ] Return a clear `invalid cursor` error when a subscriber asks for an id that has been trimmed or belongs to a different local stream. +- [ ] Add a cleanup method for terminal runs so completed/cancelled runs can release retained control messages. +- [ ] Test that messages publish in order. +- [ ] Test that an active subscriber receives a message published after subscription. +- [ ] Test that messages published before subscription are replayed to a `Start` subscriber. +- [ ] Test that `After(id)` receives only later messages. +- [ ] Test that trimming bounds retained messages and reports an invalid old cursor. +- [ ] Run `cargo nextest run -p fabro-server worker_control`. + +## Task 3: Add Control Bus to Server State + +**Files:** +- Modify: `lib/crates/fabro-server/src/server.rs` +- Test: `lib/crates/fabro-server/src/server/tests.rs` + +- [ ] Add `worker_control_bus: Arc` to `AppState`. +- [ ] Construct `LocalWorkerControlBus` in normal server state initialization. +- [ ] Add a test-only way to inject a fake or local bus without exposing test helpers to production builds. +- [ ] Keep demo/in-process execution behavior unchanged unless it currently depends on subprocess control. +- [ ] Add a state construction test proving the default bus is local and available. +- [ ] Run `cargo nextest run -p fabro-server worker_control`. + +## Task 4: Extend the Control Protocol + +**Files:** +- Modify: `lib/crates/fabro-interview/src/control_protocol.rs` +- Test: `lib/crates/fabro-interview/src/control_protocol.rs` + +- [ ] Add `WorkerControlEnvelope::pause_run()` and `WorkerControlEnvelope::unpause_run()` constructors. +- [ ] Add `WorkerControlMessage::RunPause` serialized as `"run.pause"`. +- [ ] Add `WorkerControlMessage::RunUnpause` serialized as `"run.unpause"`. +- [ ] Add `WorkerControlDeliveryFrame { id: String, envelope: WorkerControlEnvelope }` as the WebSocket text-frame payload shared by server and worker. +- [ ] Add round-trip serde tests for both new messages. +- [ ] Add round-trip serde tests for `WorkerControlDeliveryFrame`. +- [ ] Run `cargo nextest run -p fabro-interview control_protocol`. + +## Task 5: Share Worker Message Handling + +**Files:** +- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs` +- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs` + +- [ ] Split `apply_worker_control_line(...)` into parsing and `apply_worker_control_message(...)`. +- [ ] Route WebSocket delivery frames through `apply_worker_control_message(...)`. +- [ ] Route `run.pause` to `RunControlState::request_pause()`. +- [ ] Route `run.unpause` to `RunControlState::request_unpause()`. +- [ ] Add a small in-memory applied-id dedupe set in the worker control task; ignore duplicate delivery ids before applying envelopes. +- [ ] Update `last_applied_id` only after `apply_worker_control_message(...)` returns. +- [ ] Treat all current control messages as idempotent under delivery-id dedupe. `run.steer` must not be applied twice for the same delivery id. +- [ ] 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. +- [ ] Update existing stdin-era tests to exercise the shared message handler directly. +- [ ] Add tests for pause and unpause routing. +- [ ] Add a test proving duplicate delivery ids are not applied twice. +- [ ] Run `cargo nextest run -p fabro-cli runner`. + +## Task 6: Add Worker WebSocket Client + +**Files:** +- Modify: `lib/crates/fabro-cli/Cargo.toml` +- Modify: `lib/crates/fabro-cli/src/commands/run/runner.rs` +- Test: `lib/crates/fabro-cli/src/commands/run/runner.rs` + +- [ ] Add `tokio-tungstenite.workspace = true` as a direct `fabro-cli` dependency if the crate does not already have it. +- [ ] Add a helper that builds the control-stream request for a `ServerTarget` and `RunId`. +- [ ] For HTTP URLs, convert `http` to `ws` and `https` to `wss`. +- [ ] 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. +- [ ] Add the worker bearer token as an `Authorization: Bearer ...` request header. +- [ ] On the first connection, omit the `after` query parameter so the server maps it to `WorkerControlCursor::Start`. +- [ ] On reconnect, include `?after=` when `last_applied_id` is set. +- [ ] Spawn a WebSocket control manager task in `execute(...)` after `ControlInterviewer`, `RunControlState`, `CancellationToken`, and `SteeringHub` are created, and before `operations::start` or `operations::resume`. +- [ ] Gate `operations::start` and `operations::resume` on the first successful control-stream connection. +- [ ] 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. +- [ ] Deserialize each text frame into `WorkerControlDeliveryFrame`. +- [ ] Apply each envelope through `apply_worker_control_message(...)`, then record the frame id as `last_applied_id`. +- [ ] Respond to received WebSocket ping frames with pong frames. +- [ ] Send worker-initiated ping frames every 15s. +- [ ] Track pongs for worker-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness. +- [ ] Treat normal close/error as reconnectable while the run is not terminal. +- [ ] Treat HTTP 410 Gone or a WebSocket close reason of `invalid_cursor` as fatal control-channel loss. +- [ ] On fatal control-channel loss, interrupt pending interviews and fail/abort the run with an infrastructure/control-channel error, not a user cancellation. +- [ ] Wire fatal control-channel loss back into `execute(...)` so the worker returns an error instead of silently continuing workflow execution. +- [ ] Add tests for URL/request construction for `http`, `https`, and Unix socket targets. +- [ ] Add tests proving first connection has no `after` query and reconnect includes `after=`. +- [ ] Add tests for reconnect backoff bounds. +- [ ] Add tests for ping/pong timeout behavior using paused Tokio time. +- [ ] Add a local Unix-socket WebSocket test proving the client can complete a handshake against an Axum route. +- [ ] Run `cargo nextest run -p fabro-cli runner`. + +## Task 7: Add Worker-Only Control Stream Route + +**Files:** +- Create: `lib/crates/fabro-server/src/server/handler/worker_control.rs` +- Modify: `lib/crates/fabro-server/src/server/handler/mod.rs` +- Modify: `lib/crates/fabro-server/src/principal_middleware.rs` +- Test: `lib/crates/fabro-server/src/server/tests.rs` + +- [ ] Add a narrow helper or extractor that accepts only authenticated worker principals whose token run id matches the route run id. +- [ ] Add `GET /runs/{id}/worker/control-stream` to real API routes only. +- [ ] Reject missing runs, terminal runs, and archived runs before upgrading. +- [ ] Reject user JWTs and cross-run worker JWTs. +- [ ] Parse absent `after` into `WorkerControlCursor::Start`. +- [ ] Parse present `after` into `WorkerControlCursor::After(id)`. +- [ ] On upgrade, call `worker_control_bus.subscribe(run_id, cursor)`. +- [ ] If `subscribe` returns invalid cursor before upgrade, reject with HTTP 410 Gone. +- [ ] Serialize each `WorkerControlDelivery` to `WorkerControlDeliveryFrame` and send it as a WebSocket text frame. +- [ ] Send server-initiated ping frames every 15s. +- [ ] Respond to received WebSocket ping frames with pong frames. +- [ ] Track pongs for server-initiated pings and close the WebSocket after 45s without a matching pong or other proof of connection liveness. +- [ ] On timeout or disconnect, drop the bus subscription so local resources are released. +- [ ] Do not store the live WebSocket sender in `ManagedRun`; the bus is now the delivery boundary. +- [ ] 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. +- [ ] Run `cargo nextest run -p fabro-server worker_control`. + +## Task 8: Replace Server-Side Stdin Transport with Bus Publishing + +**Files:** +- Modify: `lib/crates/fabro-server/src/server.rs` +- Modify: `lib/crates/fabro-server/src/server/handler/lifecycle.rs` +- Modify: `lib/crates/fabro-server/src/server/handler/pair.rs` +- Test: `lib/crates/fabro-server/src/server/tests.rs` + +- [ ] Replace `RunAnswerTransport::Subprocess { control_tx }` with a bus-backed subprocess/worker variant. +- [ ] Ensure the bus-backed variant has enough context to publish messages for the correct `RunId`. +- [ ] Delete `pump_worker_control_jsonl(...)`. +- [ ] Change `worker_command(...)` so `__run-worker` uses `stdin(Stdio::null())` instead of `stdin(Stdio::piped())`. +- [ ] Remove child-stdin extraction and the control pump task from `execute_run_subprocess(...)`. +- [ ] Keep stderr capture and worker exit handling unchanged. +- [ ] Update `RunAnswerTransport` methods so answer, cancel, steer, interrupt, pair start/message/end all publish the existing envelope to `WorkerControlBus`. +- [ ] Add `pause_run()` and `unpause_run()` methods on `RunAnswerTransport`. +- [ ] Update pause/unpause lifecycle handlers to send `run.pause` and `run.unpause` over the bus for running workers. +- [ ] Keep process signals only for hard cleanup paths such as cancel fallback, shutdown, terminal delete, and force removal. +- [ ] Update existing tests that assert subprocess transport enqueue behavior to assert bus publish behavior instead. +- [ ] Add a `worker_command` test proving stdin is null/not piped and `FABRO_WORKER_TOKEN` still travels only through env. +- [ ] Run `cargo nextest run -p fabro-server worker_command`. + +## Task 9: End-to-End Local Control Flow Regression + +**Files:** +- Test: `lib/crates/fabro-cli/tests/it/cmd/runner.rs` +- Test: `lib/crates/fabro-server/tests/it/scenario/lifecycle.rs` + +- [ ] Add a test run where the worker connects to the control WebSocket and receives a cancel request through `LocalWorkerControlBus`. +- [ ] Add a test where the server publishes a control message before the worker connects and the worker receives it on first connection. +- [ ] Add a reconnect test where the worker applies message A, reconnects with `after=`, and then receives only message B. +- [ ] Add an invalid-cursor test proving the worker reports control-channel loss as infrastructure failure rather than user cancellation. +- [ ] Add a human-interview test proving submitted answers reach the worker through the bus and WebSocket. +- [ ] Add a steer or interrupt test proving live agent controls still reach the worker transport. +- [ ] Add a local Unix-socket server test proving the default local server target works without stdin. +- [ ] Run `cargo nextest run -p fabro-cli --test it runner`. +- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`. + +## Task 10: Final Verification + +**Files:** +- Modify only if failures expose necessary fixes. + +- [ ] Run `cargo nextest run -p fabro-interview control_protocol`. +- [ ] Run `cargo nextest run -p fabro-server worker_control`. +- [ ] Run `cargo nextest run -p fabro-cli runner`. +- [ ] Run `cargo nextest run -p fabro-server worker_command`. +- [ ] Run `cargo nextest run -p fabro-server --test it lifecycle`. +- [ ] Run `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings`. +- [ ] Confirm no code path still writes `WorkerControlEnvelope` to child stdin. +- [ ] Confirm no Redis dependency, Redis config key, or Redis runtime path was added. +- [ ] Confirm there is no `Latest` cursor or wait-for-subscriber behavior in the control bus. +- [ ] Confirm WebSocket ping/pong handling is explicit on both worker and server. +- [ ] Confirm `run.steer` and other controls are protected from duplicate delivery-id application. +- [ ] Confirm `__run-worker` still scrubs `FABRO_WORKER_TOKEN` from process env before launching descendants. + +## Acceptance Criteria + +- All worker control traffic uses the worker control bus plus WebSocket last-mile transport. +- Local Unix-socket server targets and remote HTTP(S) server targets both support worker control without Redis. +- Local and single-node deployments require no external control-channel service. +- The bus API can later be implemented by Redis Streams without changing API handlers or worker message handling. +- First worker connection replays retained messages from the beginning of the run control stream; reconnect resumes after the last fully applied id. +- Invalid cursor is the only fatal control-stream replay failure and is surfaced as infrastructure/control-channel failure, not user cancellation. +- WebSocket liveness is explicit and backend-agnostic. +- Existing run event/blob/artifact HTTP paths are unchanged. +- Existing worker JWT scope rules remain authoritative. +- Temporary WebSocket disconnects reconnect and replay through the bus; only unrecoverable replay loss reports worker-control-unavailable behavior. +- Worker stdout/stderr behavior remains unchanged except that stdin is no longer a control channel. +- Redis is clearly documented as future work and is not required by this plan. + + +## Completed stages +- **toolchain**: succeeded + - 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` + - Output: + ``` + cargo 1.95.0 (f2d3ce0bd 2026-03-21) + ``` +- **preflight_compile**: succeeded + - Script: `cargo check -q --workspace 2>&1` + - Output: (empty) +- **preflight_lint**: succeeded + - Script: `cargo +nightly-2026-04-14 clippy -q --workspace --all-targets -- -D warnings 2>&1` + - Output: (empty) + + +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. \ No newline at end of file diff --git a/stages/005-implement@1/provider_used.json b/stages/005-implement@1/provider_used.json new file mode 100644 index 000000000..c57772db6 --- /dev/null +++ b/stages/005-implement@1/provider_used.json @@ -0,0 +1,6 @@ +{ + "mode": "agent", + "provider": "openai", + "model": "gpt-5.5", + "reasoning_effort": "xhigh" +} \ No newline at end of file diff --git a/stages/005-implement@1/response.md b/stages/005-implement@1/response.md new file mode 100644 index 000000000..3bb9aa078 --- /dev/null +++ b/stages/005-implement@1/response.md @@ -0,0 +1,61 @@ +Implemented the worker control bus plan end-to-end. + +### What changed + +- Added shared protocol support: + - `run.pause` / `run.unpause` + - `WorkerControlDeliveryFrame { id, envelope }` + +- Added server-side worker control bus: + - `WorkerControlBus` abstraction + - opaque `WorkerControlMessageId` + - `WorkerControlCursor::{Start, After(...)}` + - `LocalWorkerControlBus` with ordered replay, `after` resume, bounded retention, invalid-cursor detection, and cleanup + +- Added server WebSocket route: + - `GET /api/v1/runs/{id}/worker/control-stream` + - worker-JWT-only and run-scoped + - rejects missing/terminal/archived runs before upgrade + - invalid cursor returns HTTP 410 before upgrade + - sends JSON delivery frames + - explicit ping/pong liveness + +- Removed worker stdin JSONL runtime transport: + - server now publishes all controls through the bus + - `__run-worker` subprocess stdin is configured as null + - no `control_tx` / stdin control pump remains + +- Added CLI worker WebSocket control manager: + - connects before start/resume for non-terminal runs + - supports HTTP → WS, HTTPS → WSS, and Unix socket targets + - reconnects with exponential backoff capped at 5s + - first connect omits `after`; reconnect uses last fully applied id + - applies frames only once by delivery id + - updates `last_applied_id` only after successful application + - handles pause/unpause/cancel/answer/steer/interrupt/pair controls + - explicit ping/pong liveness and timeout + - HTTP 410 / `invalid_cursor` close is fatal control-channel loss + +- Added/updated tests across protocol, bus, route, CLI, server transport, and lifecycle/cancel flows. + +- Fixed a lifecycle regression where cancellation at a human gate could race with the detached completion guard and emit a duplicate terminal failure. + +### Verification run + +Passed: + +- `cargo nextest run -p fabro-interview control_protocol` +- `cargo nextest run -p fabro-server worker_control` +- `cargo nextest run -p fabro-cli runner` +- `cargo nextest run -p fabro-server worker_command` +- `cargo nextest run -p fabro-cli --test it runner` +- `cargo nextest run -p fabro-server --features test-support --test it lifecycle` +- `cargo +nightly-2026-04-14 fmt --check --all` +- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` +- `git diff --check` + +Also confirmed by grep: + +- no worker-control stdin pump / `control_tx` remnants in server or CLI runtime code +- no `Latest` worker-control cursor +- no Redis dependency/config/runtime path added, only future-backend comments in the bus contract. \ No newline at end of file