From 39b7e5a648d4839a16f8138b24304e854bf9282a Mon Sep 17 00:00:00 2001 From: Fabro Date: Wed, 27 May 2026 19:33:27 -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 | 1119 +++++++++++++------- stages/007-simplify_gpt@1/diff.patch | 321 ++++++ stages/007-simplify_gpt@1/response.md | 28 + stages/007-simplify_gpt@1/status.json | 6 + stages/008-verify@1/script_invocation.json | 5 + 5 files changed, 1071 insertions(+), 408 deletions(-) create mode 100644 stages/007-simplify_gpt@1/diff.patch create mode 100644 stages/007-simplify_gpt@1/response.md create mode 100644 stages/007-simplify_gpt@1/status.json create mode 100644 stages/008-verify@1/script_invocation.json diff --git a/run.json b/run.json index a217ec2fb..ca6c8f9cd 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-27T23:29:21.608742Z", + "last_event_at": "2026-05-27T23:29:25.446517Z", "pending_control": null, "checkpoints": [ { @@ -1140,9 +1140,9 @@ } }, { - "seq": 0, + "seq": 2145, "checkpoint": { - "timestamp": "2026-05-27T23:29:21.639533Z", + "timestamp": "2026-05-27T23:29:25.442840Z", "current_node": "simplify_gpt", "completed_nodes": [ "start", @@ -1155,15 +1155,260 @@ ], "node_retries": {}, "context_values": { + "internal.work_dir": "/home/daytona/workspace/fabro", + "last_response": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has al", + "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.", + "current_node": "simplify_gpt", + "internal.thread_id": "simplify_opus", + "internal.retry_count.simplify_opus": 0, + "last_stage": "simplify_gpt", + "response.simplify_gpt": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has already trimmed messages, rather than silently replaying a truncated stream.\n - Active `Start` subscribers that joined before first publish also fail if messages were trimmed before they could establish a cursor.\n - `After(id)` subscriptions for unknown runs no longer create empty run stream entries.\n - Added tests for these cursor/cleanup cases.\n\n- **Worker-side delivery-id dedupe**\n - Replaced the single `last_applied_id` duplicate check with a bounded recent-id dedupe set.\n - Uses FIFO + `HashSet`, capped at `2048` delivery ids.\n - Still uses the last applied id for reconnect `?after=...`.\n - Prevents repeated non-adjacent duplicate delivery ids from being applied twice without unbounded memory growth.\n\n- **Formatting**\n - `worker_control.rs` had rustfmt-only import formatting changes.\n\nAlready clean / previously addressed in current tree:\n- Worker control route already checks invalid cursor before WebSocket upgrade and returns HTTP `410 Gone`.\n- Invalid-cursor close reason and ping/pong constants were already shared through `fabro-interview`.\n- Terminal-run check already uses `is_terminal()`.\n\nValidation run:\n- `cargo +nightly-2026-04-14 fmt --check --all` ✅\n- `cargo nextest run -p fabro-server worker_control` ✅\n- `cargo nextest run -p fabro-cli runner` ✅\n- `cargo +nightly-2026-04-14 clippy -q -p fabro-server -p fabro-cli --all-targets -- -D warnings` ✅", + "thread.preflight_lint.current_node": "implement", + "graph.rankdir": "LR", + "graph.model_stylesheet": "\n * { model: claude-opus-4-7; }\n ", + "internal.retry_count.toolchain": 0, + "response.simplify_opus": "Net 47 lines deleted. Here's a summary of the cleanup applied based on the three parallel reviews:\n\n## Cleanup Summary\n\n**Correctness fix:**\n- **Lost-wakeup race in `LocalWorkerControlBus`**: `local_subscription_task` could miss a publish that landed between message-read and `notify.notified()` registration. Now registers `notified()` before inspecting the stream and uses the proper enable-before-check pattern.\n\n**Memory / efficiency:**\n- **Removed unbounded `applied_ids: HashSet`** in the CLI worker control task. The server cursor + monotonic IDs already guarantee no replay; kept a simple `last_applied_id == frame.id` check as defense-in-depth.\n- **Dropped `Arc>>` for `last_applied_id`**: it's owned by a single task. Now a plain `&mut Option`, removing two lock-held-across-await sites per loop iteration.\n- **Pinned the `Sleep` once outside the WS loop** in both server and worker handlers, calling `Pin::reset` instead of allocating a new `Sleep` per iteration.\n- **Replaced O(N) `iter().filter().cloned()` over the retained deque** with `partition_point` + tail iteration in `collect_messages_from`.\n\n**Reuse / consolidation:**\n- **Moved `WORKER_CONTROL_WS_PING_INTERVAL`, `WORKER_CONTROL_WS_LIVENESS_TIMEOUT`, `WORKER_CONTROL_INVALID_CURSOR_REASON`, `WORKER_CONTROL_PONG_TIMEOUT_REASON`** into `fabro-interview::control_protocol` so server and worker share one source of truth.\n- **Used `axum::extract::ws::close_code::AWAY` / `POLICY` constants** instead of raw `1001` / `1008` literals on the server.\n\n**Dead code removal:**\n- **Removed `WorkerControlBus::backend_name`** trait method (only used by a tautological test) and the associated `#[allow(dead_code)]` and `Hash` derive on `WorkerControlMessageId`.\n- **Removed the single-variant `WorkerCommandStdin` enum** + constant + tautological test; inlined `Stdio::null()`.\n- **Deleted `#[cfg(test)] parse_worker_control_line` and `apply_worker_control_line`** wrappers; rewrote the four tests to call `apply_worker_control_message` directly with typed envelopes.\n\n**Verification:**\n- `cargo nextest run -p fabro-server worker_control` — 15/15 pass\n- `cargo nextest run -p fabro-interview` — 62/62 pass\n- `cargo nextest run -p fabro-cli runner` — 33/33 pass\n- `cargo nextest run -p fabro-server --test it --features test-support lifecycle` — 3/3 pass\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` — clean\n- Two pre-existing `pause_run_*` / `unpause_run_sets_pending_control` failures verified to exist before this cleanup (confirmed via `git stash`).\n\n**Findings I deliberately skipped** (not high-impact enough or would require API surface changes outside this cleanup):\n- Splitting `WorkerControlBusError` into publish/subscribe error types\n- Replacing `BoxFuture` with `async_trait` (style consistency)\n- Moving WS-URL construction onto `ServerTarget` (touches `fabro-client`, `fabro-sandbox`)\n- Adopting `fabro_util::backoff::BackoffPolicy` (small benefit; existing hand-rolled is correct)\n- Collapsing `WorkerControlSocket` enum into `Box` (works but tungstenite's `from_raw_socket` generic constraints make this fiddly)\n- Per-run inner mutex sharding in the bus (premature for local single-node)\n- `cleanup_worker_control_bus_for_run` spawn-per-call (sync-fast-path optimization; correctness fine)", + "failure_signature": "", + "internal.run_id": "01KSNP2DXVS2TD1HEASQAFBGFK", + "thread.implement.current_node": "simplify_opus", + "internal.retry_count.preflight_compile": 0, + "internal.retry_count.implement": 0, + "thread.toolchain.current_node": "preflight_compile", + "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.fidelity": "compact", + "internal.node_visit_count": 1, + "internal.retry_count.preflight_lint": 0, + "outcome": "succeeded", + "thread.preflight_compile.current_node": "preflight_lint", + "internal.retry_count.simplify_gpt": 0, + "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "thread.simplify_opus.current_node": "simplify_gpt", + "thread.start.current_node": "toolchain", + "failure_class": "", + "internal.retry_count.start": 0 + }, + "node_outcomes": { + "simplify_gpt": { + "status": "succeeded", + "context_updates": { + "last_stage": "simplify_gpt", + "response.simplify_gpt": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has already trimmed messages, rather than silently replaying a truncated stream.\n - Active `Start` subscribers that joined before first publish also fail if messages were trimmed before they could establish a cursor.\n - `After(id)` subscriptions for unknown runs no longer create empty run stream entries.\n - Added tests for these cursor/cleanup cases.\n\n- **Worker-side delivery-id dedupe**\n - Replaced the single `last_applied_id` duplicate check with a bounded recent-id dedupe set.\n - Uses FIFO + `HashSet`, capped at `2048` delivery ids.\n - Still uses the last applied id for reconnect `?after=...`.\n - Prevents repeated non-adjacent duplicate delivery ids from being applied twice without unbounded memory growth.\n\n- **Formatting**\n - `worker_control.rs` had rustfmt-only import formatting changes.\n\nAlready clean / previously addressed in current tree:\n- Worker control route already checks invalid cursor before WebSocket upgrade and returns HTTP `410 Gone`.\n- Invalid-cursor close reason and ping/pong constants were already shared through `fabro-interview`.\n- Terminal-run check already uses `is_terminal()`.\n\nValidation run:\n- `cargo +nightly-2026-04-14 fmt --check --all` ✅\n- `cargo nextest run -p fabro-server worker_control` ✅\n- `cargo nextest run -p fabro-cli runner` ✅\n- `cargo +nightly-2026-04-14 clippy -q -p fabro-server -p fabro-cli --all-targets -- -D warnings` ✅", + "last_response": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has al" + }, + "notes": "Stage completed: simplify_gpt", + "usage": { + "input": { + "usage": { + "model": { + "provider": "openai", + "model_id": "gpt-5.5" + }, + "tokens": { + "input_tokens": 379150, + "output_tokens": 8564, + "reasoning_tokens": 1094, + "cache_read_tokens": 1107968, + "cache_write_tokens": 0 + } + }, + "facts": { + "algorithm": "openai" + } + }, + "total_usd_micros": 2739474 + }, + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 283105, + "tool_time_ms": 155590, + "active_time_ms": 438695 + } + }, + "start": { + "status": "succeeded", + "usage": null + }, + "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 + } + }, + "simplify_opus": { + "status": "succeeded", + "context_updates": { + "last_stage": "simplify_opus", + "response.simplify_opus": "Net 47 lines deleted. Here's a summary of the cleanup applied based on the three parallel reviews:\n\n## Cleanup Summary\n\n**Correctness fix:**\n- **Lost-wakeup race in `LocalWorkerControlBus`**: `local_subscription_task` could miss a publish that landed between message-read and `notify.notified()` registration. Now registers `notified()` before inspecting the stream and uses the proper enable-before-check pattern.\n\n**Memory / efficiency:**\n- **Removed unbounded `applied_ids: HashSet`** in the CLI worker control task. The server cursor + monotonic IDs already guarantee no replay; kept a simple `last_applied_id == frame.id` check as defense-in-depth.\n- **Dropped `Arc>>` for `last_applied_id`**: it's owned by a single task. Now a plain `&mut Option`, removing two lock-held-across-await sites per loop iteration.\n- **Pinned the `Sleep` once outside the WS loop** in both server and worker handlers, calling `Pin::reset` instead of allocating a new `Sleep` per iteration.\n- **Replaced O(N) `iter().filter().cloned()` over the retained deque** with `partition_point` + tail iteration in `collect_messages_from`.\n\n**Reuse / consolidation:**\n- **Moved `WORKER_CONTROL_WS_PING_INTERVAL`, `WORKER_CONTROL_WS_LIVENESS_TIMEOUT`, `WORKER_CONTROL_INVALID_CURSOR_REASON`, `WORKER_CONTROL_PONG_TIMEOUT_REASON`** into `fabro-interview::control_protocol` so server and worker share one source of truth.\n- **Used `axum::extract::ws::close_code::AWAY` / `POLICY` constants** instead of raw `1001` / `1008` literals on the server.\n\n**Dead code removal:**\n- **Removed `WorkerControlBus::backend_name`** trait method (only used by a tautological test) and the associated `#[allow(dead_code)]` and `Hash` derive on `WorkerControlMessageId`.\n- **Removed the single-variant `WorkerCommandStdin` enum** + constant + tautological test; inlined `Stdio::null()`.\n- **Deleted `#[cfg(test)] parse_worker_control_line` and `apply_worker_control_line`** wrappers; rewrote the four tests to call `apply_worker_control_message` directly with typed envelopes.\n\n**Verification:**\n- `cargo nextest run -p fabro-server worker_control` — 15/15 pass\n- `cargo nextest run -p fabro-interview` — 62/62 pass\n- `cargo nextest run -p fabro-cli runner` — 33/33 pass\n- `cargo nextest run -p fabro-server --test it --features test-support lifecycle` — 3/3 pass\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` — clean\n- Two pre-existing `pause_run_*` / `unpause_run_sets_pending_control` failures verified to exist before this cleanup (confirmed via `git stash`).\n\n**Findings I deliberately skipped** (not high-impact enough or would require API surface changes outside this cleanup):\n- Splitting `WorkerControlBusError` into publish/subscribe error types\n- Replacing `BoxFuture` with `async_trait` (style consistency)\n- Moving WS-URL construction onto `ServerTarget` (touches `fabro-client`, `fabro-sandbox`)\n- Adopting `fabro_util::backoff::BackoffPolicy` (small benefit; existing hand-rolled is correct)\n- Collapsing `WorkerControlSocket` enum into `Box` (works but tungstenite's `from_raw_socket` generic constraints make this fiddly)\n- Per-run inner mutex sharding in the bus (premature for local single-node)\n- `cleanup_worker_control_bus_for_run` spawn-per-call (sync-fast-path optimization; correctness fine)", + "last_response": "Net 47 lines deleted. Here's a summary of the cleanup applied based on the three parallel reviews:\n\n## Cleanup Summary\n\n**Correctness fix:**\n- **Lost-wakeup race in `LocalWorkerControlBus`**: `local_s" + }, + "notes": "Stage completed: simplify_opus", + "usage": { + "input": { + "usage": { + "model": { + "provider": "anthropic", + "model_id": "claude-opus-4-7" + }, + "tokens": { + "input_tokens": 148310, + "output_tokens": 40147, + "reasoning_tokens": 0, + "cache_read_tokens": 11102756, + "cache_write_tokens": 713208 + } + }, + "facts": { + "algorithm": "anthropic", + "cache_write_5m_tokens": 713208, + "cache_write_1h_tokens": 0 + } + }, + "total_usd_micros": 11754153 + }, + "files_touched": [ + "/home/daytona/workspace/fabro/lib/crates/fabro-cli/src/commands/run/runner.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-interview/src/control_protocol.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-interview/src/lib.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/server.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/server/handler/worker_control.rs", + "/home/daytona/workspace/fabro/lib/crates/fabro-server/src/server/tests.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" + ], + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 588149, + "tool_time_ms": 786419, + "active_time_ms": 1374568 + } + }, + "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 + } + }, + "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 + } + }, + "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 + } + } + }, + "next_node_id": "verify", + "git_commit_sha": "6722751ca4ea9cfc5e75f77fa219dcd797f8d8a7", + "node_visits": { + "toolchain": 1, + "simplify_gpt": 1, + "preflight_compile": 1, + "simplify_opus": 1, + "implement": 1, + "start": 1, + "preflight_lint": 1 + } + }, + "diff": { + "patch": "diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs\nindex 17c96751d..ac1a2bbd0 100644\n--- a/lib/crates/fabro-cli/src/commands/run/runner.rs\n+++ b/lib/crates/fabro-cli/src/commands/run/runner.rs\n@@ -1,3 +1,4 @@\n+use std::collections::{HashSet, VecDeque};\n use std::path::{Path, PathBuf};\n use std::sync::Arc;\n use std::time::Duration;\n@@ -289,6 +290,38 @@ fn load_worker_vault(storage_dir: Option<&Path>) -> Result,\n+ order: VecDeque,\n+ recent: HashSet,\n+}\n+\n+impl AppliedWorkerControlDeliveryIds {\n+ fn last_applied_id(&self) -> Option<&str> {\n+ self.last.as_deref()\n+ }\n+\n+ fn contains(&self, id: &str) -> bool {\n+ self.recent.contains(id)\n+ }\n+\n+ fn record(&mut self, id: String) {\n+ if !self.recent.insert(id.clone()) {\n+ self.last = Some(id);\n+ return;\n+ }\n+ self.order.push_back(id.clone());\n+ self.last = Some(id);\n+ while self.order.len() > WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY {\n+ if let Some(evicted) = self.order.pop_front() {\n+ self.recent.remove(&evicted);\n+ }\n+ }\n+ }\n+}\n \n struct WorkerControlManagerHandle {\n first_connection: Option>>,\n@@ -439,14 +472,14 @@ async fn run_worker_control_manager(\n let mut first_tx = Some(first_tx);\n let mut fatal_tx = Some(fatal_tx);\n let mut backoff = WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF;\n- let mut last_applied_id: Option = None;\n+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();\n \n while !done.is_cancelled() {\n let request = match build_worker_control_stream_request(\n &target,\n &run_id,\n &worker_token,\n- last_applied_id.as_deref(),\n+ applied_ids.last_applied_id(),\n ) {\n Ok(request) => request,\n Err(err) => {\n@@ -474,7 +507,7 @@ async fn run_worker_control_manager(\n &cancel_token,\n &steering_hub,\n &run_control,\n- &mut last_applied_id,\n+ &mut applied_ids,\n &done,\n )\n .await\n@@ -645,7 +678,7 @@ async fn handle_worker_control_socket(\n cancel_token: &CancellationToken,\n steering_hub: &fabro_workflow::SteeringHub,\n run_control: &RunControlState,\n- last_applied_id: &mut Option,\n+ applied_ids: &mut AppliedWorkerControlDeliveryIds,\n done: &CancellationToken,\n ) -> Result<(), WorkerControlConnectError> {\n let mut ping_interval = time::interval(WORKER_CONTROL_WS_PING_INTERVAL);\n@@ -692,7 +725,7 @@ async fn handle_worker_control_socket(\n cancel_token,\n steering_hub,\n run_control,\n- last_applied_id,\n+ applied_ids,\n frame,\n )\n .await;\n@@ -730,16 +763,16 @@ async fn apply_worker_control_delivery_frame(\n cancel_token: &CancellationToken,\n steering_hub: &fabro_workflow::SteeringHub,\n run_control: &RunControlState,\n- last_applied_id: &mut Option,\n+ applied_ids: &mut AppliedWorkerControlDeliveryIds,\n frame: WorkerControlDeliveryFrame,\n ) -> bool {\n // Duplicate ids cannot reach us under normal operation: the server replays\n- // strictly after `last_applied_id`. Guard against a server-side bug by\n- // ignoring any frame whose id is not strictly newer than what we last\n- // applied.\n- if last_applied_id.as_deref() == Some(frame.id.as_str()) {\n+ // strictly after the last applied id. Guard against a server-side bug or\n+ // reconnect race by ignoring recently-applied delivery ids.\n+ if applied_ids.contains(&frame.id) {\n return false;\n }\n+ let frame_id = frame.id;\n apply_worker_control_message(\n interviewer,\n cancel_token,\n@@ -748,7 +781,7 @@ async fn apply_worker_control_delivery_frame(\n frame.envelope,\n )\n .await;\n- *last_applied_id = Some(frame.id);\n+ applied_ids.record(frame_id);\n true\n }\n \n@@ -1202,8 +1235,8 @@ mod tests {\n use tokio_util::sync::CancellationToken;\n \n use super::{\n- WorkerControlConnectError, WorkerControlSocket, WorkerTitlePhase,\n- apply_worker_control_delivery_frame, apply_worker_control_message,\n+ AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket,\n+ WorkerTitlePhase, apply_worker_control_delivery_frame, apply_worker_control_message,\n build_worker_control_stream_request, connect_worker_control_stream,\n handle_worker_control_socket, initial_worker_title_phase, load_worker_vault,\n next_worker_control_reconnect_backoff, stamp_system_worker, worker_title,\n@@ -1538,7 +1571,7 @@ mod tests {\n let cancel_token = CancellationToken::new();\n let run_control = RunControlState::new();\n let hub = test_steering_hub();\n- let mut last_applied_id: Option = None;\n+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();\n let frame = fabro_interview::WorkerControlDeliveryFrame {\n id: \"local:1\".to_string(),\n envelope: WorkerControlEnvelope::pause_run(),\n@@ -1550,7 +1583,7 @@ mod tests {\n &cancel_token,\n &hub,\n &run_control,\n- &mut last_applied_id,\n+ &mut applied_ids,\n frame.clone(),\n )\n .await\n@@ -1561,13 +1594,13 @@ mod tests {\n &cancel_token,\n &hub,\n &run_control,\n- &mut last_applied_id,\n+ &mut applied_ids,\n frame,\n )\n .await\n );\n \n- assert_eq!(last_applied_id, Some(\"local:1\".to_string()));\n+ assert_eq!(applied_ids.last_applied_id(), Some(\"local:1\"));\n }\n \n #[test]\n@@ -1648,7 +1681,7 @@ mod tests {\n let cancel_token = CancellationToken::new();\n let hub = test_steering_hub();\n let run_control = RunControlState::new();\n- let mut last_applied_id: Option = None;\n+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();\n let done = CancellationToken::new();\n \n let task = tokio::spawn(async move {\n@@ -1658,7 +1691,7 @@ mod tests {\n &cancel_token,\n &hub,\n &run_control,\n- &mut last_applied_id,\n+ &mut applied_ids,\n &done,\n )\n .await\ndiff --git a/lib/crates/fabro-server/src/server/handler/worker_control.rs b/lib/crates/fabro-server/src/server/handler/worker_control.rs\nindex 682107748..d643a904e 100644\n--- a/lib/crates/fabro-server/src/server/handler/worker_control.rs\n+++ b/lib/crates/fabro-server/src/server/handler/worker_control.rs\n@@ -1,6 +1,8 @@\n use std::sync::Arc;\n \n-use axum::extract::ws::{CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code};\n+use axum::extract::ws::{\n+ CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code,\n+};\n use fabro_interview::{\n WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON,\n WORKER_CONTROL_WS_LIVENESS_TIMEOUT, WORKER_CONTROL_WS_PING_INTERVAL,\ndiff --git a/lib/crates/fabro-server/src/worker_control/local.rs b/lib/crates/fabro-server/src/worker_control/local.rs\nindex 51ea1421f..219ff18ec 100644\n--- a/lib/crates/fabro-server/src/worker_control/local.rs\n+++ b/lib/crates/fabro-server/src/worker_control/local.rs\n@@ -23,8 +23,9 @@ pub(crate) struct LocalWorkerControlBus {\n }\n \n struct LocalRunControlStream {\n- messages: VecDeque,\n- notify: Arc,\n+ messages: VecDeque,\n+ notify: Arc,\n+ has_trimmed: bool,\n }\n \n #[derive(Clone)]\n@@ -60,9 +61,17 @@ impl LocalWorkerControlBus {\n .streams\n .lock()\n .expect(\"worker control streams poisoned\");\n- let stream = streams\n- .entry(run_id)\n- .or_insert_with(LocalRunControlStream::new);\n+ let stream = match cursor {\n+ WorkerControlCursor::Start => streams\n+ .entry(run_id)\n+ .or_insert_with(LocalRunControlStream::new),\n+ WorkerControlCursor::After(id) => streams.get_mut(&run_id).ok_or_else(|| {\n+ WorkerControlBusError::invalid_cursor(\n+ id.as_str(),\n+ \"message id is not retained for this run\",\n+ )\n+ })?,\n+ };\n let next_sequence = stream.next_sequence_for_cursor(cursor)?;\n (next_sequence, Arc::clone(&stream.notify))\n };\n@@ -90,8 +99,9 @@ impl LocalWorkerControlBus {\n impl LocalRunControlStream {\n fn new() -> Self {\n Self {\n- messages: VecDeque::new(),\n- notify: Arc::new(Notify::new()),\n+ messages: VecDeque::new(),\n+ notify: Arc::new(Notify::new()),\n+ has_trimmed: false,\n }\n }\n \n@@ -100,6 +110,9 @@ impl LocalRunControlStream {\n cursor: &WorkerControlCursor,\n ) -> Result, WorkerControlBusError> {\n match cursor {\n+ WorkerControlCursor::Start if self.has_trimmed => Err(\n+ WorkerControlBusError::invalid_cursor(\"start\", \"retained local stream is trimmed\"),\n+ ),\n WorkerControlCursor::Start => Ok(self.messages.front().map(|message| message.sequence)),\n WorkerControlCursor::After(id) => {\n let requested_sequence = parse_local_sequence(id)?;\n@@ -125,6 +138,7 @@ impl LocalRunControlStream {\n fn trim_retained(&mut self) {\n while self.messages.len() > LOCAL_WORKER_CONTROL_RETAINED_MESSAGES_PER_RUN {\n self.messages.pop_front();\n+ self.has_trimmed = true;\n }\n }\n }\n@@ -163,12 +177,19 @@ async fn local_subscription_task(\n Some(stream) => {\n // A `Start` subscriber that joined before any publish lazily\n // adopts the first retained message as its cursor.\n- if next_sequence.is_none() {\n- next_sequence = stream.messages.front().map(|message| message.sequence);\n- }\n- match next_sequence {\n- None => Some(Ok(Vec::new())),\n- Some(next) => Some(collect_messages_from(&stream.messages, next)),\n+ if next_sequence.is_none() && stream.has_trimmed {\n+ Some(Err(WorkerControlBusError::invalid_cursor(\n+ \"start\",\n+ \"retained local stream is trimmed\",\n+ )))\n+ } else {\n+ if next_sequence.is_none() {\n+ next_sequence = stream.messages.front().map(|message| message.sequence);\n+ }\n+ match next_sequence {\n+ None => Some(Ok(Vec::new())),\n+ Some(next) => Some(collect_messages_from(&stream.messages, next)),\n+ }\n }\n }\n }\n@@ -399,6 +420,27 @@ mod tests {\n .await\n .unwrap_err();\n assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));\n+\n+ let err = bus\n+ .subscribe(fixtures::RUN_1, WorkerControlCursor::Start)\n+ .await\n+ .unwrap_err();\n+ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));\n+ }\n+\n+ #[tokio::test]\n+ async fn after_cursor_for_unknown_run_does_not_create_stream() {\n+ let bus = LocalWorkerControlBus::new();\n+ let err = bus\n+ .subscribe(\n+ fixtures::RUN_1,\n+ WorkerControlCursor::After(WorkerControlMessageId::new(\"local:1\")),\n+ )\n+ .await\n+ .unwrap_err();\n+\n+ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));\n+ assert_eq!(bus.retained_len(&fixtures::RUN_1), 0);\n }\n \n #[tokio::test]\n", + "summary": { + "files_changed": 20, + "additions": 2313, + "deletions": 331 + } + } + }, + { + "seq": 0, + "checkpoint": { + "timestamp": "2026-05-27T23:33:26.982609Z", + "current_node": "verify", + "completed_nodes": [ + "start", + "toolchain", + "preflight_compile", + "preflight_lint", + "implement", + "simplify_opus", + "simplify_gpt", + "verify" + ], + "node_retries": {}, + "context_values": { + "thread.simplify_gpt.current_node": "verify", "internal.retry_count.start": 0, "internal.retry_count.preflight_lint": 0, "thread.preflight_lint.current_node": "implement", "response.simplify_gpt": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has already trimmed messages, rather than silently replaying a truncated stream.\n - Active `Start` subscribers that joined before first publish also fail if messages were trimmed before they could establish a cursor.\n - `After(id)` subscriptions for unknown runs no longer create empty run stream entries.\n - Added tests for these cursor/cleanup cases.\n\n- **Worker-side delivery-id dedupe**\n - Replaced the single `last_applied_id` duplicate check with a bounded recent-id dedupe set.\n - Uses FIFO + `HashSet`, capped at `2048` delivery ids.\n - Still uses the last applied id for reconnect `?after=...`.\n - Prevents repeated non-adjacent duplicate delivery ids from being applied twice without unbounded memory growth.\n\n- **Formatting**\n - `worker_control.rs` had rustfmt-only import formatting changes.\n\nAlready clean / previously addressed in current tree:\n- Worker control route already checks invalid cursor before WebSocket upgrade and returns HTTP `410 Gone`.\n- Invalid-cursor close reason and ping/pong constants were already shared through `fabro-interview`.\n- Terminal-run check already uses `is_terminal()`.\n\nValidation run:\n- `cargo +nightly-2026-04-14 fmt --check --all` ✅\n- `cargo nextest run -p fabro-server worker_control` ✅\n- `cargo nextest run -p fabro-cli runner` ✅\n- `cargo +nightly-2026-04-14 clippy -q -p fabro-server -p fabro-cli --all-targets -- -D warnings` ✅", - "current_node": "simplify_gpt", + "current_node": "verify", "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": "", + "outcome": "failed", + "failure_class": "transient_infra", "response.simplify_opus": "Net 47 lines deleted. Here's a summary of the cleanup applied based on the three parallel reviews:\n\n## Cleanup Summary\n\n**Correctness fix:**\n- **Lost-wakeup race in `LocalWorkerControlBus`**: `local_subscription_task` could miss a publish that landed between message-read and `notify.notified()` registration. Now registers `notified()` before inspecting the stream and uses the proper enable-before-check pattern.\n\n**Memory / efficiency:**\n- **Removed unbounded `applied_ids: HashSet`** in the CLI worker control task. The server cursor + monotonic IDs already guarantee no replay; kept a simple `last_applied_id == frame.id` check as defense-in-depth.\n- **Dropped `Arc>>` for `last_applied_id`**: it's owned by a single task. Now a plain `&mut Option`, removing two lock-held-across-await sites per loop iteration.\n- **Pinned the `Sleep` once outside the WS loop** in both server and worker handlers, calling `Pin::reset` instead of allocating a new `Sleep` per iteration.\n- **Replaced O(N) `iter().filter().cloned()` over the retained deque** with `partition_point` + tail iteration in `collect_messages_from`.\n\n**Reuse / consolidation:**\n- **Moved `WORKER_CONTROL_WS_PING_INTERVAL`, `WORKER_CONTROL_WS_LIVENESS_TIMEOUT`, `WORKER_CONTROL_INVALID_CURSOR_REASON`, `WORKER_CONTROL_PONG_TIMEOUT_REASON`** into `fabro-interview::control_protocol` so server and worker share one source of truth.\n- **Used `axum::extract::ws::close_code::AWAY` / `POLICY` constants** instead of raw `1001` / `1008` literals on the server.\n\n**Dead code removal:**\n- **Removed `WorkerControlBus::backend_name`** trait method (only used by a tautological test) and the associated `#[allow(dead_code)]` and `Hash` derive on `WorkerControlMessageId`.\n- **Removed the single-variant `WorkerCommandStdin` enum** + constant + tautological test; inlined `Stdio::null()`.\n- **Deleted `#[cfg(test)] parse_worker_control_line` and `apply_worker_control_line`** wrappers; rewrote the four tests to call `apply_worker_control_message` directly with typed envelopes.\n\n**Verification:**\n- `cargo nextest run -p fabro-server worker_control` — 15/15 pass\n- `cargo nextest run -p fabro-interview` — 62/62 pass\n- `cargo nextest run -p fabro-cli runner` — 33/33 pass\n- `cargo nextest run -p fabro-server --test it --features test-support lifecycle` — 3/3 pass\n- `cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings` — clean\n- Two pre-existing `pause_run_*` / `unpause_run_sets_pending_control` failures verified to exist before this cleanup (confirmed via `git stash`).\n\n**Findings I deliberately skipped** (not high-impact enough or would require API surface changes outside this cleanup):\n- Splitting `WorkerControlBusError` into publish/subscribe error types\n- Replacing `BoxFuture` with `async_trait` (style consistency)\n- Moving WS-URL construction onto `ServerTarget` (touches `fabro-client`, `fabro-sandbox`)\n- Adopting `fabro_util::backoff::BackoffPolicy` (small benefit; existing hand-rolled is correct)\n- Collapsing `WorkerControlSocket` enum into `Box` (works but tungstenite's `from_raw_socket` generic constraints make this fiddly)\n- Per-run inner mutex sharding in the bus (premature for local single-node)\n- `cleanup_worker_control_bus_for_run` spawn-per-call (sync-fast-path optimization; correctness fine)", "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", @@ -1171,18 +1416,19 @@ "internal.fidelity": "compact", "internal.retry_count.toolchain": 0, "thread.start.current_node": "toolchain", - "command.output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "command.output": "blob://sha256/10a757ece73e698a21c8f3a03407d736790cedac3bb241c9a7438a47e2097ea3", "internal.retry_count.implement": 0, "last_stage": "simplify_gpt", "internal.retry_count.simplify_gpt": 0, "thread.preflight_compile.current_node": "preflight_lint", - "internal.thread_id": "simplify_opus", + "internal.thread_id": "simplify_gpt", "internal.retry_count.preflight_compile": 0, - "failure_signature": "", + "failure_signature": "verify|transient_infra|script failed with exit code: ## output abro/lib/crates/fabro-variable) compiling fabro-dev v0..-nightly. (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-dev) compiling fabro-oauth v0..-nightly. (/home/daytona/rep", "thread.toolchain.current_node": "preflight_compile", "internal.node_visit_count": 1, "last_response": "Implemented review cleanup after running the three parallel review agents.\n\nFixed:\n- **Local bus replay correctness**\n - `Start` subscriptions now fail with `InvalidCursor` if the local stream has al", "internal.retry_count.simplify_opus": 0, + "internal.retry_count.verify": 0, "thread.implement.current_node": "simplify_opus", "internal.work_dir": "/home/daytona/workspace/fabro", "graph.model_stylesheet": "\n * { model: claude-opus-4-7; }\n " @@ -1359,17 +1605,35 @@ "tool_time_ms": 155590, "active_time_ms": 438695 } + }, + "verify": { + "status": "failed", + "context_updates": { + "command.output": "blob://sha256/10a757ece73e698a21c8f3a03407d736790cedac3bb241c9a7438a47e2097ea3" + }, + "failure": { + "message": "Script failed with exit code: 100\n\n## output\nabro/lib/crates/fabro-variable)\n Compiling fabro-dev v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-dev)\n Compiling fabro-oauth v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-oauth)\n Compiling fabro-agent v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-agent)\n Compiling fabro-acp v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-acp)\n Compiling fabro-validate v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-validate)\n Compiling typify v0.6.1\n Compiling progenitor-impl v0.13.0\n Compiling fabro-hooks v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-hooks)\n Compiling fabro-tracker v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-tracker)\n Compiling progenitor-macro v0.13.0\n Compiling twin-github v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/test/twin/github)\n Compiling fabro-http v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-http)\n Compiling fabro-proc v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-proc)\n Compiling fabro-spa v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-spa)\n Compiling fabro-build-support v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/build-support)\n Compiling fabro-static v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-static)\n Compiling progenitor v0.13.0\n Compiling fabro-api v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-api)\n Compiling fabro-client v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-client)\n Compiling fabro-tool v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-tool)\n Compiling fabro-workflow v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-workflow)\n Compiling fabro-manifest v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-manifest)\n Compiling fabro-slack v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-slack)\n Compiling fabro-mcp-server v0.246.0-nightly.0 (/home/daytona/repos/fabro-sh/fabro/lib/crates/fabro-mcp-server)\n Finished `test` profile [unoptimized + debuginfo] target(s) in 1m 32s\n────────────\n Nextest run ID 62503790-4d0a-499d-abfb-3d7bc927f862 with nextest profile: ci\n Starting 6438 tests across 132 binaries (179 tests skipped)\n FAIL [ 0.118s] (3781/6438) fabro-server server::tests::pause_run_sets_pending_control_on_board_response\n stdout ───\n\n running 1 test\n test server::tests::pause_run_sets_pending_control_on_board_response ... FAILED\n\n failures:\n\n failures:\n server::tests::pause_run_sets_pending_control_on_board_response\n\n test result: FAILED. 0 passed; 1 failed; 0 ignored; 0 measured; 702 filtered out; finished in 0.11s\n\n stderr ───\n\n thread 'server::tests::pause_run_sets_pending_control_on_board_response' (114416) panicked at lib/crates/fabro-test/src/http_assert.rs:25:5:\n lib/crates/fabro-server/src/server/tests.rs:13376\n expected status: 200 OK\n actual status: 409 Conflict\n content-type: application/json\n body:\n {\n \"errors\": [\n {\n \"status\": \"409\",\n \"title\": \"Conflict\",\n \"detail\": \"Run worker is not available.\",\n \"request_id\": \"2d721cd4-36a3-4d8b-bb38-4c1fc1b54ade\"\n }\n ],\n \"request_id\": \"2d721cd4-36a3-4d8b-bb38-4c1fc1b54ade\"\n }\n note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace\n\n Cancelling due to test failure: 7 tests still running\n────────────\n Summary [ 15.811s] 3788/6438 tests run: 3787 passed, 1 failed, 179 skipped\n FAIL [ 0.118s] (3781/6438) fabro-server server::tests::pause_run_sets_pending_control_on_board_response\nwarning: 2650/6438 tests were not run due to test failure (run with --no-fail-fast to run all tests, or run with --max-fail)\nerror: test run failed\n", + "category": "transient_infra" + }, + "usage": null, + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 0, + "tool_time_ms": 241516, + "active_time_ms": 241516 + } } }, - "next_node_id": "verify", + "next_node_id": "fixup", "node_visits": { - "start": 1, - "toolchain": 1, "preflight_lint": 1, + "toolchain": 1, "simplify_gpt": 1, "implement": 1, + "verify": 1, + "preflight_compile": 1, "simplify_opus": 1, - "preflight_compile": 1 + "start": 1 } }, "diff": {} @@ -1401,312 +1665,6 @@ "superseded_by": null, "pending_interviews": {}, "stages": { - "implement@1": { - "first_event_seq": 51, - "prompt": null, - "response": null, - "completion": { - "outcome": "succeeded", - "notes": "Stage completed: implement", - "failure_reason": null, - "timestamp": "2026-05-27T22:58:48.112584Z" - }, - "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", - "timing": { - "wall_time_ms": 4428217, - "inference_time_ms": 2814420, - "tool_time_ms": 1504977, - "active_time_ms": 4319397 - }, - "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": "succeeded" - }, "toolchain@1": { "first_event_seq": 21, "prompt": null, @@ -1803,40 +1761,6 @@ }, "state": "succeeded" }, - "start@1": { - "first_event_seq": 17, - "prompt": null, - "response": null, - "completion": { - "outcome": "succeeded", - "notes": null, - "failure_reason": null, - "timestamp": "2026-05-27T21:40:10.410686Z" - }, - "provider_used": null, - "diff": null, - "script_invocation": null, - "script_timing": null, - "parallel_results": null, - "output": null, - "started_at": "2026-05-27T21:40:10.410597Z", - "handler": "start", - "timing": { - "wall_time_ms": 0, - "inference_time_ms": 0, - "tool_time_ms": 0, - "active_time_ms": 0 - }, - "usage": { - "input_tokens": 0, - "output_tokens": 0, - "total_tokens": 0, - "reasoning_tokens": 0, - "cache_read_tokens": 0, - "cache_write_tokens": 0 - }, - "state": "succeeded" - }, "simplify_opus@1": { "first_event_seq": 1294, "prompt": null, @@ -2169,11 +2093,431 @@ }, "state": "succeeded" }, + "verify@1": { + "first_event_seq": 2148, + "prompt": null, + "response": null, + "completion": null, + "provider_used": null, + "diff": null, + "script_invocation": { + "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", + "command": "exec 2>&1\ngit 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", + "language": "shell" + }, + "script_timing": null, + "parallel_results": null, + "output": null, + "started_at": "2026-05-27T23:29:25.446079Z", + "handler": "command", + "usage": { + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + "reasoning_tokens": 0, + "cache_read_tokens": 0, + "cache_write_tokens": 0 + }, + "state": "running" + }, + "preflight_compile@1": { + "first_event_seq": 31, + "prompt": null, + "response": null, + "completion": { + "outcome": "succeeded", + "notes": "Script completed: cargo check -q --workspace 2>&1", + "failure_reason": null, + "timestamp": "2026-05-27T21:42:27.506687Z" + }, + "provider_used": null, + "diff": null, + "script_invocation": { + "script": "cargo check -q --workspace 2>&1", + "command": "exec 2>&1\ncargo check -q --workspace 2>&1", + "language": "shell" + }, + "script_timing": { + "output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", + "exit_code": 0, + "duration_ms": 130559, + "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:40:16.940647Z", + "handler": "command", + "timing": { + "wall_time_ms": 130565, + "inference_time_ms": 0, + "tool_time_ms": 130559, + "active_time_ms": 130559 + }, + "usage": { + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + "reasoning_tokens": 0, + "cache_read_tokens": 0, + "cache_write_tokens": 0 + }, + "state": "succeeded" + }, + "implement@1": { + "first_event_seq": 51, + "prompt": null, + "response": null, + "completion": { + "outcome": "succeeded", + "notes": "Stage completed: implement", + "failure_reason": null, + "timestamp": "2026-05-27T22:58:48.112584Z" + }, + "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", + "timing": { + "wall_time_ms": 4428217, + "inference_time_ms": 2814420, + "tool_time_ms": 1504977, + "active_time_ms": 4319397 + }, + "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": "succeeded" + }, + "start@1": { + "first_event_seq": 17, + "prompt": null, + "response": null, + "completion": { + "outcome": "succeeded", + "notes": null, + "failure_reason": null, + "timestamp": "2026-05-27T21:40:10.410686Z" + }, + "provider_used": null, + "diff": null, + "script_invocation": null, + "script_timing": null, + "parallel_results": null, + "output": null, + "started_at": "2026-05-27T21:40:10.410597Z", + "handler": "start", + "timing": { + "wall_time_ms": 0, + "inference_time_ms": 0, + "tool_time_ms": 0, + "active_time_ms": 0 + }, + "usage": { + "input_tokens": 0, + "output_tokens": 0, + "total_tokens": 0, + "reasoning_tokens": 0, + "cache_read_tokens": 0, + "cache_write_tokens": 0 + }, + "state": "succeeded" + }, "simplify_gpt@1": { "first_event_seq": 1914, "prompt": null, "response": null, - "completion": null, + "completion": { + "outcome": "succeeded", + "notes": "Stage completed: simplify_gpt", + "failure_reason": null, + "timestamp": "2026-05-27T23:29:21.638800Z" + }, "provider_used": { "mode": "agent", "provider": "openai", @@ -2186,13 +2530,20 @@ "output": null, "started_at": "2026-05-27T23:22:02.256851Z", "handler": "agent", + "timing": { + "wall_time_ms": 439379, + "inference_time_ms": 283105, + "tool_time_ms": 155590, + "active_time_ms": 438695 + }, "usage": { - "input_tokens": 426095, - "output_tokens": 9487, - "total_tokens": 1557024, - "reasoning_tokens": 1186, - "cache_read_tokens": 1120256, - "cache_write_tokens": 0 + "input_tokens": 379150, + "output_tokens": 8564, + "total_tokens": 1496776, + "reasoning_tokens": 1094, + "cache_read_tokens": 1107968, + "cache_write_tokens": 0, + "total_usd_micros": 2739474 }, "model": { "provider": "openai", @@ -2428,54 +2779,6 @@ ], "warnings": [] }, - "state": "running" - }, - "preflight_compile@1": { - "first_event_seq": 31, - "prompt": null, - "response": null, - "completion": { - "outcome": "succeeded", - "notes": "Script completed: cargo check -q --workspace 2>&1", - "failure_reason": null, - "timestamp": "2026-05-27T21:42:27.506687Z" - }, - "provider_used": null, - "diff": null, - "script_invocation": { - "script": "cargo check -q --workspace 2>&1", - "command": "exec 2>&1\ncargo check -q --workspace 2>&1", - "language": "shell" - }, - "script_timing": { - "output": "blob://sha256/12ae32cb1ec02d01eda3581b127c1fee3b0dc53572ed6baf239721a03d82e126", - "exit_code": 0, - "duration_ms": 130559, - "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:40:16.940647Z", - "handler": "command", - "timing": { - "wall_time_ms": 130565, - "inference_time_ms": 0, - "tool_time_ms": 130559, - "active_time_ms": 130559 - }, - "usage": { - "input_tokens": 0, - "output_tokens": 0, - "total_tokens": 0, - "reasoning_tokens": 0, - "cache_read_tokens": 0, - "cache_write_tokens": 0 - }, "state": "succeeded" } } diff --git a/stages/007-simplify_gpt@1/diff.patch b/stages/007-simplify_gpt@1/diff.patch new file mode 100644 index 000000000..5920b9f24 --- /dev/null +++ b/stages/007-simplify_gpt@1/diff.patch @@ -0,0 +1,321 @@ +diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs +index 17c96751d..ac1a2bbd0 100644 +--- a/lib/crates/fabro-cli/src/commands/run/runner.rs ++++ b/lib/crates/fabro-cli/src/commands/run/runner.rs +@@ -1,3 +1,4 @@ ++use std::collections::{HashSet, VecDeque}; + use std::path::{Path, PathBuf}; + use std::sync::Arc; + use std::time::Duration; +@@ -289,6 +290,38 @@ fn load_worker_vault(storage_dir: Option<&Path>) -> Result, ++ order: VecDeque, ++ recent: HashSet, ++} ++ ++impl AppliedWorkerControlDeliveryIds { ++ fn last_applied_id(&self) -> Option<&str> { ++ self.last.as_deref() ++ } ++ ++ fn contains(&self, id: &str) -> bool { ++ self.recent.contains(id) ++ } ++ ++ fn record(&mut self, id: String) { ++ if !self.recent.insert(id.clone()) { ++ self.last = Some(id); ++ return; ++ } ++ self.order.push_back(id.clone()); ++ self.last = Some(id); ++ while self.order.len() > WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY { ++ if let Some(evicted) = self.order.pop_front() { ++ self.recent.remove(&evicted); ++ } ++ } ++ } ++} + + struct WorkerControlManagerHandle { + first_connection: Option>>, +@@ -439,14 +472,14 @@ async fn run_worker_control_manager( + let mut first_tx = Some(first_tx); + let mut fatal_tx = Some(fatal_tx); + let mut backoff = WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF; +- let mut last_applied_id: Option = None; ++ let mut applied_ids = AppliedWorkerControlDeliveryIds::default(); + + while !done.is_cancelled() { + let request = match build_worker_control_stream_request( + &target, + &run_id, + &worker_token, +- last_applied_id.as_deref(), ++ applied_ids.last_applied_id(), + ) { + Ok(request) => request, + Err(err) => { +@@ -474,7 +507,7 @@ async fn run_worker_control_manager( + &cancel_token, + &steering_hub, + &run_control, +- &mut last_applied_id, ++ &mut applied_ids, + &done, + ) + .await +@@ -645,7 +678,7 @@ async fn handle_worker_control_socket( + cancel_token: &CancellationToken, + steering_hub: &fabro_workflow::SteeringHub, + run_control: &RunControlState, +- last_applied_id: &mut Option, ++ applied_ids: &mut AppliedWorkerControlDeliveryIds, + done: &CancellationToken, + ) -> Result<(), WorkerControlConnectError> { + let mut ping_interval = time::interval(WORKER_CONTROL_WS_PING_INTERVAL); +@@ -692,7 +725,7 @@ async fn handle_worker_control_socket( + cancel_token, + steering_hub, + run_control, +- last_applied_id, ++ applied_ids, + frame, + ) + .await; +@@ -730,16 +763,16 @@ async fn apply_worker_control_delivery_frame( + cancel_token: &CancellationToken, + steering_hub: &fabro_workflow::SteeringHub, + run_control: &RunControlState, +- last_applied_id: &mut Option, ++ applied_ids: &mut AppliedWorkerControlDeliveryIds, + frame: WorkerControlDeliveryFrame, + ) -> bool { + // Duplicate ids cannot reach us under normal operation: the server replays +- // strictly after `last_applied_id`. Guard against a server-side bug by +- // ignoring any frame whose id is not strictly newer than what we last +- // applied. +- if last_applied_id.as_deref() == Some(frame.id.as_str()) { ++ // strictly after the last applied id. Guard against a server-side bug or ++ // reconnect race by ignoring recently-applied delivery ids. ++ if applied_ids.contains(&frame.id) { + return false; + } ++ let frame_id = frame.id; + apply_worker_control_message( + interviewer, + cancel_token, +@@ -748,7 +781,7 @@ async fn apply_worker_control_delivery_frame( + frame.envelope, + ) + .await; +- *last_applied_id = Some(frame.id); ++ applied_ids.record(frame_id); + true + } + +@@ -1202,8 +1235,8 @@ mod tests { + use tokio_util::sync::CancellationToken; + + use super::{ +- WorkerControlConnectError, WorkerControlSocket, WorkerTitlePhase, +- apply_worker_control_delivery_frame, apply_worker_control_message, ++ AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket, ++ WorkerTitlePhase, apply_worker_control_delivery_frame, apply_worker_control_message, + build_worker_control_stream_request, connect_worker_control_stream, + handle_worker_control_socket, initial_worker_title_phase, load_worker_vault, + next_worker_control_reconnect_backoff, stamp_system_worker, worker_title, +@@ -1538,7 +1571,7 @@ mod tests { + let cancel_token = CancellationToken::new(); + let run_control = RunControlState::new(); + let hub = test_steering_hub(); +- let mut last_applied_id: Option = None; ++ let mut applied_ids = AppliedWorkerControlDeliveryIds::default(); + let frame = fabro_interview::WorkerControlDeliveryFrame { + id: "local:1".to_string(), + envelope: WorkerControlEnvelope::pause_run(), +@@ -1550,7 +1583,7 @@ mod tests { + &cancel_token, + &hub, + &run_control, +- &mut last_applied_id, ++ &mut applied_ids, + frame.clone(), + ) + .await +@@ -1561,13 +1594,13 @@ mod tests { + &cancel_token, + &hub, + &run_control, +- &mut last_applied_id, ++ &mut applied_ids, + frame, + ) + .await + ); + +- assert_eq!(last_applied_id, Some("local:1".to_string())); ++ assert_eq!(applied_ids.last_applied_id(), Some("local:1")); + } + + #[test] +@@ -1648,7 +1681,7 @@ mod tests { + let cancel_token = CancellationToken::new(); + let hub = test_steering_hub(); + let run_control = RunControlState::new(); +- let mut last_applied_id: Option = None; ++ let mut applied_ids = AppliedWorkerControlDeliveryIds::default(); + let done = CancellationToken::new(); + + let task = tokio::spawn(async move { +@@ -1658,7 +1691,7 @@ mod tests { + &cancel_token, + &hub, + &run_control, +- &mut last_applied_id, ++ &mut applied_ids, + &done, + ) + .await +diff --git a/lib/crates/fabro-server/src/server/handler/worker_control.rs b/lib/crates/fabro-server/src/server/handler/worker_control.rs +index 682107748..d643a904e 100644 +--- a/lib/crates/fabro-server/src/server/handler/worker_control.rs ++++ b/lib/crates/fabro-server/src/server/handler/worker_control.rs +@@ -1,6 +1,8 @@ + use std::sync::Arc; + +-use axum::extract::ws::{CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code}; ++use axum::extract::ws::{ ++ CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code, ++}; + use fabro_interview::{ + WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON, + WORKER_CONTROL_WS_LIVENESS_TIMEOUT, WORKER_CONTROL_WS_PING_INTERVAL, +diff --git a/lib/crates/fabro-server/src/worker_control/local.rs b/lib/crates/fabro-server/src/worker_control/local.rs +index 51ea1421f..219ff18ec 100644 +--- a/lib/crates/fabro-server/src/worker_control/local.rs ++++ b/lib/crates/fabro-server/src/worker_control/local.rs +@@ -23,8 +23,9 @@ pub(crate) struct LocalWorkerControlBus { + } + + struct LocalRunControlStream { +- messages: VecDeque, +- notify: Arc, ++ messages: VecDeque, ++ notify: Arc, ++ has_trimmed: bool, + } + + #[derive(Clone)] +@@ -60,9 +61,17 @@ impl LocalWorkerControlBus { + .streams + .lock() + .expect("worker control streams poisoned"); +- let stream = streams +- .entry(run_id) +- .or_insert_with(LocalRunControlStream::new); ++ let stream = match cursor { ++ WorkerControlCursor::Start => streams ++ .entry(run_id) ++ .or_insert_with(LocalRunControlStream::new), ++ WorkerControlCursor::After(id) => streams.get_mut(&run_id).ok_or_else(|| { ++ WorkerControlBusError::invalid_cursor( ++ id.as_str(), ++ "message id is not retained for this run", ++ ) ++ })?, ++ }; + let next_sequence = stream.next_sequence_for_cursor(cursor)?; + (next_sequence, Arc::clone(&stream.notify)) + }; +@@ -90,8 +99,9 @@ impl LocalWorkerControlBus { + impl LocalRunControlStream { + fn new() -> Self { + Self { +- messages: VecDeque::new(), +- notify: Arc::new(Notify::new()), ++ messages: VecDeque::new(), ++ notify: Arc::new(Notify::new()), ++ has_trimmed: false, + } + } + +@@ -100,6 +110,9 @@ impl LocalRunControlStream { + cursor: &WorkerControlCursor, + ) -> Result, WorkerControlBusError> { + match cursor { ++ WorkerControlCursor::Start if self.has_trimmed => Err( ++ WorkerControlBusError::invalid_cursor("start", "retained local stream is trimmed"), ++ ), + WorkerControlCursor::Start => Ok(self.messages.front().map(|message| message.sequence)), + WorkerControlCursor::After(id) => { + let requested_sequence = parse_local_sequence(id)?; +@@ -125,6 +138,7 @@ impl LocalRunControlStream { + fn trim_retained(&mut self) { + while self.messages.len() > LOCAL_WORKER_CONTROL_RETAINED_MESSAGES_PER_RUN { + self.messages.pop_front(); ++ self.has_trimmed = true; + } + } + } +@@ -163,12 +177,19 @@ async fn local_subscription_task( + Some(stream) => { + // A `Start` subscriber that joined before any publish lazily + // adopts the first retained message as its cursor. +- if next_sequence.is_none() { +- next_sequence = stream.messages.front().map(|message| message.sequence); +- } +- match next_sequence { +- None => Some(Ok(Vec::new())), +- Some(next) => Some(collect_messages_from(&stream.messages, next)), ++ if next_sequence.is_none() && stream.has_trimmed { ++ Some(Err(WorkerControlBusError::invalid_cursor( ++ "start", ++ "retained local stream is trimmed", ++ ))) ++ } else { ++ if next_sequence.is_none() { ++ next_sequence = stream.messages.front().map(|message| message.sequence); ++ } ++ match next_sequence { ++ None => Some(Ok(Vec::new())), ++ Some(next) => Some(collect_messages_from(&stream.messages, next)), ++ } + } + } + } +@@ -399,6 +420,27 @@ mod tests { + .await + .unwrap_err(); + assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. })); ++ ++ let err = bus ++ .subscribe(fixtures::RUN_1, WorkerControlCursor::Start) ++ .await ++ .unwrap_err(); ++ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. })); ++ } ++ ++ #[tokio::test] ++ async fn after_cursor_for_unknown_run_does_not_create_stream() { ++ let bus = LocalWorkerControlBus::new(); ++ let err = bus ++ .subscribe( ++ fixtures::RUN_1, ++ WorkerControlCursor::After(WorkerControlMessageId::new("local:1")), ++ ) ++ .await ++ .unwrap_err(); ++ ++ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. })); ++ assert_eq!(bus.retained_len(&fixtures::RUN_1), 0); + } + + #[tokio::test] diff --git a/stages/007-simplify_gpt@1/response.md b/stages/007-simplify_gpt@1/response.md new file mode 100644 index 000000000..04e7c8402 --- /dev/null +++ b/stages/007-simplify_gpt@1/response.md @@ -0,0 +1,28 @@ +Implemented review cleanup after running the three parallel review agents. + +Fixed: +- **Local bus replay correctness** + - `Start` subscriptions now fail with `InvalidCursor` if the local stream has already trimmed messages, rather than silently replaying a truncated stream. + - Active `Start` subscribers that joined before first publish also fail if messages were trimmed before they could establish a cursor. + - `After(id)` subscriptions for unknown runs no longer create empty run stream entries. + - Added tests for these cursor/cleanup cases. + +- **Worker-side delivery-id dedupe** + - Replaced the single `last_applied_id` duplicate check with a bounded recent-id dedupe set. + - Uses FIFO + `HashSet`, capped at `2048` delivery ids. + - Still uses the last applied id for reconnect `?after=...`. + - Prevents repeated non-adjacent duplicate delivery ids from being applied twice without unbounded memory growth. + +- **Formatting** + - `worker_control.rs` had rustfmt-only import formatting changes. + +Already clean / previously addressed in current tree: +- Worker control route already checks invalid cursor before WebSocket upgrade and returns HTTP `410 Gone`. +- Invalid-cursor close reason and ping/pong constants were already shared through `fabro-interview`. +- Terminal-run check already uses `is_terminal()`. + +Validation run: +- `cargo +nightly-2026-04-14 fmt --check --all` ✅ +- `cargo nextest run -p fabro-server worker_control` ✅ +- `cargo nextest run -p fabro-cli runner` ✅ +- `cargo +nightly-2026-04-14 clippy -q -p fabro-server -p fabro-cli --all-targets -- -D warnings` ✅ \ No newline at end of file diff --git a/stages/007-simplify_gpt@1/status.json b/stages/007-simplify_gpt@1/status.json new file mode 100644 index 000000000..331c17454 --- /dev/null +++ b/stages/007-simplify_gpt@1/status.json @@ -0,0 +1,6 @@ +{ + "outcome": "succeeded", + "notes": "Stage completed: simplify_gpt", + "failure_reason": null, + "timestamp": "2026-05-27T23:29:21.638800Z" +} \ No newline at end of file diff --git a/stages/008-verify@1/script_invocation.json b/stages/008-verify@1/script_invocation.json new file mode 100644 index 000000000..7ad2687d7 --- /dev/null +++ b/stages/008-verify@1/script_invocation.json @@ -0,0 +1,5 @@ +{ + "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", + "command": "exec 2>&1\ngit 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", + "language": "shell" +} \ No newline at end of file