diff --git a/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.json b/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.json new file mode 100644 index 0000000000..26e59f2a11 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.json @@ -0,0 +1,50 @@ +{ + "id": "compact_evj1c14lvl5i", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-09T18:43:45.715Z", + "sourceTrajectories": [ + "traj_yzfk3s3ayol8" + ], + "dateRange": { + "start": "2026-09-09T18:16:43.036Z", + "end": "2026-09-09T18:43:04.851Z" + }, + "summary": { + "totalDecisions": 1, + "totalEvents": 1, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "api", + "decisions": [ + { + "question": "Keep local-only capabilities fixed through reconciliation", + "chosen": "Keep local-only capabilities fixed through reconciliation", + "reasoning": "Recovered connectivity drains durable audit records without republishing DMs, which could execute work twice or route local names to another machine. Fleet capabilities require a deliberate normal restart.", + "fromTrajectory": "traj_yzfk3s3ayol8" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [ + "Explicit --local-only starts without Relaycast, confines the API to loopback, suppresses fleet connections and worker credentials/MCP injection, and reports degradation in startup, health, session, connection metadata and CLI status.", + "Local delivery is saved before acceptance. Pending work waits for a restarted recipient; a bounded, destination-pinned audit outbox survives restart and reconciles with stable event IDs without replaying DMs.", + "Validation: 1065 Rust tests passed (4 ignored), 107 CLI tests passed, TypeScript typecheck and Clippy passed. Compiled base/head proof verified outage startup failure on base and local spawn, send, terminal IO, restart retention and reconnect audit replay on head.", + "Inherited GIT_CONFIG_COUNT and RELAY_ATTEST_SESSION_ID alter isolated Git hook fixtures. Unsetting only those variables for the test subprocess yields a passing full suite; repository commit hooks remain enabled.", + "GitHub API authentication returns HTTP 401; SSH access works. PR creation and subsequent CI/review follow-through require restored API authentication. No merge authorized." + ], + "filesAffected": [ + "crates/broker/src/runtime/degraded.rs", + "crates/broker/src/runtime/init.rs", + "crates/broker/src/runtime/delivery.rs", + "crates/broker/src/listen_api.rs", + "packages/cli/src/cli/lib/broker-lifecycle.ts", + "tests/relayflows/cases/broker-local-only/run.mjs" + ], + "commits": [] +} diff --git a/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.md b/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.md new file mode 100644 index 0000000000..39af39f331 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.md @@ -0,0 +1,13 @@ +# Broker local-only operation + +Explicit --local-only starts without Relaycast, confines the API to loopback, suppresses fleet connections and worker credentials/MCP injection, and reports degradation in startup, health, session, connection metadata and CLI status. + +Local delivery is saved before acceptance. Pending work waits for a restarted recipient; a bounded, destination-pinned audit outbox survives restart and reconciles with stable event IDs without replaying DMs. + +Validation: 1065 Rust tests passed (4 ignored), 107 CLI tests passed, TypeScript typecheck and Clippy passed. Compiled base/head proof verified outage startup failure on base and local spawn, send, terminal IO, restart retention and reconnect audit replay on head. + +Inherited GIT_CONFIG_COUNT and RELAY_ATTEST_SESSION_ID alter isolated Git hook fixtures. Unsetting only those variables for the test subprocess yields a passing full suite; repository commit hooks remain enabled. + +GitHub API authentication returns HTTP 401; SSH access works. PR creation and subsequent CI/review follow-through require restored API authentication. No merge authorized. + +Decision: Recovered connectivity drains durable audit records without republishing DMs, which could execute work twice or route local names to another machine. Fleet capabilities require a deliberate normal restart. diff --git a/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.json b/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.json new file mode 100644 index 0000000000..5ffdb3ceed --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.json @@ -0,0 +1,75 @@ +{ + "id": "compact_fzz26f8hgo9y", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-10T11:35:10.845Z", + "sourceTrajectories": [ + "traj_bl7fngii53oq" + ], + "dateRange": { + "start": "2026-09-10T05:19:58.099Z", + "end": "2026-09-10T11:35:09.983Z" + }, + "summary": { + "totalDecisions": 2, + "totalEvents": 4, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "testing", + "decisions": [ + { + "question": "Pin unscoped nonempty outboxes and retain exhausted local deliveries during recipient absence", + "chosen": "Pin unscoped nonempty outboxes and retain exhausted local deliveries during recipient absence", + "reasoning": "CodeRabbit identified valid privacy and data-integrity holes. Regression tests failed on the prior implementation. Reset the handoff failure budget while waiting so a respawn can actually receive retained work, and preserve original outbox bytes on destination mismatch.", + "fromTrajectory": "traj_bl7fngii53oq" + }, + { + "question": "Preserve roster and identity safety across the main rebase", + "chosen": "Preserve roster and identity safety across the main rebase", + "reasoning": "Read merged software-garden#510 and relaycast#401 diffs. Stale roster snapshots are dispatch-only evidence; offline owners and all durable node associations remain protected. Combined upstream status timeout warnings and startup/cleanup tests with local degraded behavior. Local workers create no Relaycast identity, so exclude them from remote owned-generation cleanup and prove exited names can respawn during an outage.", + "fromTrajectory": "traj_bl7fngii53oq" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [], + "filesAffected": [ + ".agentworkforce/trajectories/active/traj_bl7fngii53oq/trajectory.json", + ".agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.json", + ".agentworkforce/trajectories/compacted/compact_ah4a5rcyq7wq_2026-09-09.md", + "CHANGELOG.md", + "crates/broker/src/cli/mod.rs", + "crates/broker/src/listen_api.rs", + "crates/broker/src/relaycast/ws.rs", + "crates/broker/src/runtime/api.rs", + "crates/broker/src/runtime/degraded.rs", + "crates/broker/src/runtime/delivery.rs", + "crates/broker/src/runtime/event_loop.rs", + "crates/broker/src/runtime/fleet.rs", + "crates/broker/src/runtime/init.rs", + "crates/broker/src/runtime/mod.rs", + "crates/broker/src/runtime/session.rs", + "crates/broker/src/runtime/tests.rs", + "crates/broker/src/worker.rs", + "packages/cli/README.md", + "packages/cli/src/cli/commands/core.test.ts", + "packages/cli/src/cli/commands/core.ts", + "packages/cli/src/cli/commands/node.ts", + "packages/cli/src/cli/commands/status.test.ts", + "packages/cli/src/cli/commands/status.ts", + "packages/cli/src/cli/lib/broker-lifecycle.ts", + "packages/harness-driver/src/protocol.ts", + "tests/relayflows/cases/broker-local-only/case.json", + "tests/relayflows/cases/broker-local-only/run.mjs" + ], + "commits": [ + "8f3bb080b0593815f7512d6f47ee77c876c158cb", + "bc4c177dcfab13777dd859f11e4723deb7f3c10c", + "5717d1b9bcbf1df81f2978883a0325b99af1d280" + ] +} diff --git a/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.md b/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.md new file mode 100644 index 0000000000..d8fffbbfcb --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_fzz26f8hgo9y_2026-09-10.md @@ -0,0 +1,9 @@ +# PR 1726 review fixes and main rebase + +The four review findings are fixed: nonempty unscoped outboxes cannot acquire an upload destination; exhausted local deliveries wait for absent recipients and retry after reconnect; taskless spawns remain idle; IPv6 listener and discovery addresses are bracketed. The outage proof records unexpected traffic during the outage as well as after recovery. + +Read the merged software-garden#510 and relaycast#401 diffs before resolving the rebase. Stale roster snapshots remain dispatch-only evidence; offline or inactive ownership never proves an identity can be deleted. Preserved main's bounded status probes, startup ordering, cleanup safeguards, and release notes. Local workers now stay out of remote identity ownership bookkeeping; the process proof rejects a remote identity deletion request without stopping the local worker, and verifies local exit/respawn. + +Rebased onto main b90248a39 (v12.0.0). Validation passed: 1,099 Rust tests (four ignored), 163 CLI tests, typecheck, strict Clippy, and local red/green process proofs using the exact base and rebased harness. Both the outbox/retry regressions and the local identity-ownership regression were observed failing before their fixes. + +GitHub reports the PR mergeable. CI run 34470354002 passed proof metadata validation and exact artifact verification, then Cloud failed before either proof arm: registering base-prover returned workspace_busy with a 60-second retry interval. The uploaded cloud.log confirms this is a pre-test infrastructure failure. A normal Actions rerun was rejected by GitHub authentication; this durable record also supplies the branch update for a fresh CI run after cooldown. CI is still pending at the time of this record. No merge was performed. diff --git a/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.json b/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.json new file mode 100644 index 0000000000..1d93a5fb27 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.json @@ -0,0 +1,37 @@ +{ + "id": "compact_kr0gviod7ik8", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-10T12:42:58.526Z", + "sourceTrajectories": [ + "traj_4rcc69em81p7" + ], + "dateRange": { + "start": "2026-09-10T12:36:24.740Z", + "end": "2026-09-10T12:42:57.920Z" + }, + "summary": { + "totalDecisions": 1, + "totalEvents": 1, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "testing", + "decisions": [ + { + "question": "Preserve privacy opt-outs and existing identity recovery; constrain audit transport and normalize effective persistence", + "chosen": "Preserve privacy opt-outs and existing identity recovery; constrain audit transport and normalize effective persistence", + "reasoning": "The delayed review correctly identified state-dir lease expiry and bind trimming. Privacy opt-outs must stay enabled rather than be removed. Audit payloads must never follow redirects; the existing registration SDK strips cross-origin Authorization, verified by a regression test. A separate trajectory data directory avoids modifying an unrelated active task.", + "fromTrajectory": "traj_4rcc69em81p7" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [], + "filesAffected": [], + "commits": [] +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.md b/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.md new file mode 100644 index 0000000000..9b81ee7a1c --- /dev/null +++ b/.agentworkforce/trajectories/compacted/compact_kr0gviod7ik8_2026-09-10.md @@ -0,0 +1,18 @@ +# Trajectory Compaction: Sep 10, 2026 - Sep 10, 2026 + +## Summary +- Sessions: 1 +- Decisions: 1 +- Events: 1 +- Agents: default +- Files: 0 +- Commits: 0 + +## Testing +- Preserve privacy opt-outs and existing identity recovery; constrain audit transport and normalize effective persistence -> Preserve privacy opt-outs and existing identity recovery; constrain audit transport and normalize effective persistence (traj_4rcc69em81p7) + +## Key Learnings +- None + +## Key Findings +- None \ No newline at end of file diff --git a/.agentworkforce/trajectories/compacted/pr1726-followup-review.md b/.agentworkforce/trajectories/compacted/pr1726-followup-review.md new file mode 100644 index 0000000000..17dc4f24d2 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/pr1726-followup-review.md @@ -0,0 +1,52 @@ +# PR 1726 follow-up review — 2026-09-10 + +New CodeRabbit and Cursor findings arrived after the rebase validation. + +- Reset restored local deliveries' failed transport budget at load time. The new regression reproduced a dropped delivery when the worker was already registered before the first retry, then passed after the fix. Remote retry budgets remain unchanged. +- Normalize bracketed IPv6 before local-only loopback validation; the process proof now starts with `[::1]`. +- Remove Relaycast identity and telemetry variables after all child environment injection. A real shell child proves the variables and credentials are absent. +- Clarify that absent local recipients retain queued work, restarted recipients receive a fresh transport budget, and only a configured digest-matching audit destination can drain a backlog. +- Keep `[Unreleased - Minor]`: AGENTS.md explicitly requires a release level for pending user-visible changes. The review suggestion to remove the level conflicts with that instruction. + +Validation: 1,100 Rust tests passed (4 ignored), strict Clippy passed, formatting and diff checks passed, and the updated local process proof passed. The preceding head's Cloud red-green proof passed in run 34472224310. New-head CI remains pending at this record's creation. + +The delayed review was recorded through the trajectory tool as `traj_4rcc69em81p7`, using an isolated data directory so the unrelated subscription-demo trajectory remains untouched. Its completed, compacted record is tracked alongside this note. + + +## Hosted validation follow-through + +Head 64b34f7cb passed every code/build/lint/security/smoke check. The standalone +macOS smoke also passed locally, including workspace reuse and confirmed cleanup. +All nine code/documentation review threads are resolved. The remaining changelog +thread contradicts AGENTS.md's explicit pending-release-level rule; its prepared +reply could not be posted because GitHub write authentication is unavailable. + +Cloud run 34473632153 failed before executing either proof arm: the executor +could not register `base-prover` after transient retries because Relaycast returned +`workspace_busy` with Retry-After 60 seconds. This is the same pre-case service +failure seen before the successful Cloud run 34472224310. The implementation and +proof assertions remain unchanged; this evidence update triggers a fresh CI run. + + +## Delayed review fixes + +Use effective `paths.persist` for owner leases and the renew-lease response; the +process proof now starts using only `--state-dir` and asserts persistent state +with no expiry. Normalize padded IPv6 consistently before binding and discovery. + +Audit endpoints require HTTPS except for literal loopback HTTP development +endpoints. Audit records use a pooled HTTP client that refuses redirects, with +a regression proving a redirected destination receives neither credentials nor +private work. The existing registration SDK's cross-host/port Authorization +stripping has a separate regression; identity ownership and recovery stay intact. + +Keep `AGENT_RELAY_TELEMETRY_DISABLED`, `DO_NOT_TRACK`, and +`AGENT_RELAY_NO_DEBUG_FILES` set to 1 in local children. They are privacy opt-outs, +not Relaycast identity; the review's assertion that they must be absent is not +part of the contract. The real-child proof now asserts those values explicitly. + +Validation: 1,102 full-suite Rust tests passed (4 ignored), plus the new SDK +redirect regression passed; strict Clippy and the updated process proof passed. +Cloud attempts 34475219721 and 34476269829 failed before case execution, at +sandbox launch and workspace registration respectively. Hosted CI must be +rechecked after this change; no green result is claimed here. diff --git a/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.json b/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.json new file mode 100644 index 0000000000..6da4057df9 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.json @@ -0,0 +1,37 @@ +{ + "id": "compact_vx8338g80q6c", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-10T19:36:29.991Z", + "sourceTrajectories": [ + "traj_sn0ch8h9fihu" + ], + "dateRange": { + "start": "2026-09-10T19:32:41.173Z", + "end": "2026-09-10T19:36:29.408Z" + }, + "summary": { + "totalDecisions": 1, + "totalEvents": 1, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "testing", + "decisions": [ + { + "question": "Treat prior local reconciliation as optional during normal startup", + "chosen": "Treat prior local reconciliation as optional during normal startup", + "reasoning": "The real-process regression exits before normal readiness when an unscoped backlog meets a configured destination. Preserve strict local-only startup checks, but retain an unmatched backlog with a warning and continue normal mode. Add process assertions for unscoped and differently scoped files, unchanged bytes, and no private audit upload. Reviewed the related Garden roster-cache and Relaycast retention changes; this does not change either ownership or roster policy.", + "fromTrajectory": "traj_sn0ch8h9fihu" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [], + "filesAffected": [], + "commits": [] +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.md b/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.md new file mode 100644 index 0000000000..ccd94678c4 --- /dev/null +++ b/.agentworkforce/trajectories/compacted/pr1726-normal-recovery.md @@ -0,0 +1,18 @@ +# Trajectory Compaction: Sep 10, 2026 - Sep 10, 2026 + +## Summary +- Sessions: 1 +- Decisions: 1 +- Events: 1 +- Agents: default +- Files: 0 +- Commits: 0 + +## Testing +- Treat prior local reconciliation as optional during normal startup -> Treat prior local reconciliation as optional during normal startup (traj_sn0ch8h9fihu) + +## Key Learnings +- None + +## Key Findings +- None \ No newline at end of file diff --git a/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.json b/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.json new file mode 100644 index 0000000000..160c5f454d --- /dev/null +++ b/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.json @@ -0,0 +1,69 @@ +{ + "id": "compact_drrcx6ih8b5e", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-10T19:24:28.621Z", + "sourceTrajectories": [ + "traj_glcziiq16vc3" + ], + "dateRange": { + "start": "2026-09-10T19:16:15.407Z", + "end": "2026-09-10T19:24:27.985Z" + }, + "summary": { + "totalDecisions": 2, + "totalEvents": 3, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "testing", + "decisions": [ + { + "question": "Preserve both changelog entries and strengthen delivery proof", + "chosen": "Preserve both changelog entries and strengthen delivery proof", + "reasoning": "Only CHANGELOG conflicted on latest main. All four pending-review replies were already posted and resolved. The queue regression must test live exhausted state before load resets the budget; the process proof must verify retained worker delivery separately from audit reconciliation.", + "fromTrajectory": "traj_glcziiq16vc3" + } + ] + }, + { + "category": "other", + "decisions": [ + { + "question": "Validate both live retry ordering and crash recovery against pre-fix delivery code", + "chosen": "Validate both live retry ordering and crash recovery against pre-fix delivery code", + "reasoning": "The strengthened unit regression fails with exhaustion before absence handling (exit 101). The expanded real-process proof passes fixed code and fails with delivery.rs from initial feature commit 0aaac651e (exit 1, absent work dead-lettered). The marker was shortened to fit the PTY line width after confirming payload rendering.", + "fromTrajectory": "traj_glcziiq16vc3" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [], + "filesAffected": [ + "CHANGELOG.md", + "packages/cli/src/cli/commands/fleet.test.ts", + "packages/cli/src/cli/commands/fleet.ts", + "packages/cloud/src/fleet-sandbox.test.ts", + "packages/cloud/src/fleet-sandbox.ts", + "tests/relayflows/cases/1630-scoped-relayfile-sandbox-mount/run.mjs", + "tests/relayflows/cases/1732-daytona-provider-sandbox-identity/case.json", + "tests/relayflows/cases/1732-daytona-provider-sandbox-identity/run.mjs" + ], + "commits": [ + "4ea088d93", + "10a7bac1c", + "2bac41acb", + "94138ffbc", + "0930799d6", + "6c816f0bb", + "7c9aa708b", + "0531e38f5", + "0aaac651e", + "6e44912d9", + "39ca68f5d" + ] +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.md b/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.md new file mode 100644 index 0000000000..74bdf883cc --- /dev/null +++ b/.agentworkforce/trajectories/compacted/pr1726-refresh-validation.md @@ -0,0 +1,21 @@ +# Trajectory Compaction: Sep 10, 2026 - Sep 10, 2026 + +## Summary +- Sessions: 1 +- Decisions: 2 +- Events: 3 +- Agents: default +- Files: 8 +- Commits: 11 + +## Testing +- Preserve both changelog entries and strengthen delivery proof -> Preserve both changelog entries and strengthen delivery proof (traj_glcziiq16vc3) + +## Other +- Validate both live retry ordering and crash recovery against pre-fix delivery code -> Validate both live retry ordering and crash recovery against pre-fix delivery code (traj_glcziiq16vc3) + +## Key Learnings +- None + +## Key Findings +- None \ No newline at end of file diff --git a/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.json b/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.json new file mode 100644 index 0000000000..2289c20feb --- /dev/null +++ b/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.json @@ -0,0 +1,37 @@ +{ + "id": "compact_x2t2gg1so54f", + "version": 1, + "type": "compacted", + "compactedAt": "2026-09-11T11:59:19.711Z", + "sourceTrajectories": [ + "traj_4fb0m0ojackj" + ], + "dateRange": { + "start": "2026-09-11T11:57:23.186Z", + "end": "2026-09-11T11:59:19.144Z" + }, + "summary": { + "totalDecisions": 1, + "totalEvents": 1, + "uniqueAgents": [ + "default" + ] + }, + "decisionGroups": [ + { + "category": "architecture", + "decisions": [ + { + "question": "Combine platform ownership check with explicit local-only startup; preserve all merged release notes", + "chosen": "Combine platform ownership check with explicit local-only startup; preserve all merged release notes", + "reasoning": "Main added fail-closed macOS/Linux startup ownership verification. Rebase retains both behaviors. Delivery implementation, degraded module, recovery tests and proof case are byte-identical to reviewed head d7a627553.", + "fromTrajectory": "traj_4fb0m0ojackj" + } + ] + } + ], + "keyLearnings": [], + "keyFindings": [], + "filesAffected": [], + "commits": [] +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.md b/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.md new file mode 100644 index 0000000000..ac3480c3be --- /dev/null +++ b/.agentworkforce/trajectories/rebase-pr1726/compacted/compact_x2t2gg1so54f_2026-09-11.md @@ -0,0 +1,18 @@ +# Trajectory Compaction: Sep 11, 2026 - Sep 11, 2026 + +## Summary +- Sessions: 1 +- Decisions: 1 +- Events: 1 +- Agents: default +- Files: 0 +- Commits: 0 + +## Architecture +- Combine platform ownership check with explicit local-only startup; preserve all merged release notes -> Combine platform ownership check with explicit local-only startup; preserve all merged release notes (traj_4fb0m0ojackj) + +## Key Learnings +- None + +## Key Findings +- None \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index d559eeb477..43203fac32 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased - Major] +### Added + +- `agent-relay node up --local-only` runs local agents during Relaycast outages, visibly reports degraded capabilities, and retains local delivery records for reconciliation after reconnect. + ### Fixed - `agent-relay node up` retries the narrowly transient Relaycast `workspace_busy` admission response while keeping unrelated rate limits terminal and preserving bounded startup diagnostics. @@ -20,6 +24,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `fleet spawn --sandbox` uses the provider-neutral durable profile for explicit Daytona and E2B sandboxes, so their measured resource envelopes are routable while Agent37 retains its heavy profile. +- Normal broker restarts preserve unmatched local audit backlogs without blocking fleet recovery. - Cloud Daytona Fleet provisioning now requires and returns the exact provider sandbox UUID alongside the stable Cloud sandbox ID, enabling ID-bound inspection and cleanup after interrupted launches. - Node startup recovery and shutdown require a persisted process and runtime-lock identity for the selected state directory, preserving unrelated agents. diff --git a/crates/broker/src/cli/mod.rs b/crates/broker/src/cli/mod.rs index 328f30d5dc..8b6f3250da 100644 --- a/crates/broker/src/cli/mod.rs +++ b/crates/broker/src/cli/mod.rs @@ -270,6 +270,11 @@ pub(crate) struct McpArgsCommand { #[derive(Debug, clap::Args)] pub(crate) struct InitCommand { + /// Run local agents only; retain delivery records for background reconciliation. + /// Fleet routing and remote attachment remain disabled until a normal restart. + #[arg(long)] + pub(crate) local_only: bool, + /// Legacy broker instance name flag. Prefer --instance-name. #[arg(long, default_value = "", alias = "broker-name")] pub(crate) name: String, @@ -371,6 +376,7 @@ mod tests { fn init_command(name: &str, instance_name: Option<&str>) -> InitCommand { InitCommand { + local_only: false, name: name.to_string(), instance_name: instance_name.map(ToOwned::to_owned), workspace_key: None, diff --git a/crates/broker/src/listen_api.rs b/crates/broker/src/listen_api.rs index addbbe0fc9..74f9c2aa09 100644 --- a/crates/broker/src/listen_api.rs +++ b/crates/broker/src/listen_api.rs @@ -259,6 +259,7 @@ pub enum ListenApiRequest { /// so they don't need a variant here. #[derive(Debug, Clone, PartialEq, Eq)] pub enum DeliveryRouteError { + CapabilityDisabled, /// No worker with that name is currently registered with the broker. WorkerNotFound(WorkerName), } @@ -266,6 +267,7 @@ pub enum DeliveryRouteError { impl std::fmt::Display for DeliveryRouteError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { + DeliveryRouteError::CapabilityDisabled => write!(f, "DEGRADED: manual flush is unavailable in local-only mode; local deliveries use the durable automatic queue"), DeliveryRouteError::WorkerNotFound(name) => { write!(f, "agent_not_found: no worker named '{name}'") } @@ -363,6 +365,7 @@ impl SnapshotFormat { #[derive(Clone)] struct ListenApiState { + local_only: bool, tx: mpsc::Sender, events_tx: broadcast::Sender, broker_api_key: Option, @@ -412,6 +415,7 @@ impl ListenReplayQuery { // --------------------------------------------------------------------------- pub struct ListenApiConfig { + pub local_only: bool, pub tx: mpsc::Sender, pub events_tx: broadcast::Sender, pub replay_buffer: ReplayBuffer, @@ -443,6 +447,7 @@ fn listen_api_router_with_auth( use axum::{middleware, routing, Router}; let state = ListenApiState { + local_only: config.local_only, tx: config.tx, events_tx: config.events_tx, broker_api_key: broker_api_key @@ -611,7 +616,13 @@ pub(crate) fn listen_api_health_payload( async fn listen_api_health( axum::extract::State(state): axum::extract::State, ) -> axum::Json { + let local_only = state.local_only; let mut payload = listen_api_health_payload(state.default_workspace_id, state.memberships); + if local_only { + payload["status"] = json!("degraded"); + payload["mode"] = json!("local_only"); + payload["relaycastConnected"] = json!(false); + } if let Some(status) = fetch_status_for_health(&state.tx).await { merge_status_into_health_payload(&mut payload, &status); } @@ -633,6 +644,12 @@ fn merge_status_into_health_payload(payload: &mut Value, status: &Value) { let Some(object) = payload.as_object_mut() else { return; }; + if status.get("mode").and_then(Value::as_str) == Some("local_only") { + object.insert("status".into(), json!("degraded")); + object.insert("mode".into(), json!("local_only")); + object.insert("relaycastConnected".into(), json!(false)); + object.insert("degraded".into(), status["degraded"].clone()); + } if let Some(agent_count) = status.get("agent_count").and_then(Value::as_u64) { object.insert("agentCount".to_string(), json!(agent_count)); } @@ -681,6 +698,8 @@ async fn listen_api_session( "broker_version": state.broker_version, "spawn_capabilities": {"explicit_empty_channels": true, "create_only_identity": true}, "protocol_version": 2, + "operation_mode": if state.local_only { "local_only" } else { "normal" }, + "degraded": state.local_only, "workspace_key": state.workspace_key, "relay_base_url": state.relay_base_url, "default_workspace_id": state.default_workspace_id, @@ -2610,6 +2629,11 @@ fn delivery_route_error_to_response( err: &DeliveryRouteError, ) -> (axum::http::StatusCode, axum::Json) { match err { + DeliveryRouteError::CapabilityDisabled => api_error( + axum::http::StatusCode::CONFLICT, + "capability_disabled", + err.to_string(), + ), DeliveryRouteError::WorkerNotFound(_) => api_error( axum::http::StatusCode::NOT_FOUND, "agent_not_found", @@ -3875,6 +3899,13 @@ mod auth_tests { fn test_router( broker_api_key: Option<&str>, + ) -> (axum::Router, mpsc::Receiver) { + test_router_with_mode(broker_api_key, false) + } + + fn test_router_with_mode( + broker_api_key: Option<&str>, + local_only: bool, ) -> (axum::Router, mpsc::Receiver) { let (tx, rx) = mpsc::channel(8); let (events_tx, _events_rx) = broadcast::channel(8); @@ -3882,6 +3913,7 @@ mod auth_tests { ( listen_api_router_with_auth( ListenApiConfig { + local_only, tx, events_tx, replay_buffer, @@ -3907,6 +3939,26 @@ mod auth_tests { serde_json::from_slice(&body).expect("response body should be json") } + #[tokio::test] + async fn local_only_health_stays_degraded_without_a_runtime_status_reply() { + let (router, rx) = test_router_with_mode(Some("test"), true); + drop(rx); + let response = router + .oneshot( + Request::builder() + .uri("/health") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = response_json(response).await; + assert_eq!(body["status"], "degraded"); + assert_eq!(body["mode"], "local_only"); + assert_eq!(body["relaycastConnected"], false); + } + #[tokio::test] async fn health_route_is_public_even_when_auth_enabled() { let (router, _rx) = test_router(Some("secret")); diff --git a/crates/broker/src/relaycast/ws.rs b/crates/broker/src/relaycast/ws.rs index 9d838b6d23..8b7f0ea2f6 100644 --- a/crates/broker/src/relaycast/ws.rs +++ b/crates/broker/src/relaycast/ws.rs @@ -174,6 +174,20 @@ impl RecipientReachability { } impl RelaycastHttpClient { + /// A client with no transport or credentials. Local-only runtime paths cannot + /// accidentally publish presence, register workers, or route remote messages. + pub fn local_only(agent_name: impl Into) -> Self { + Self { + base_url: None, + api_key: String::new(), + relay: Arc::new(None), + registration: Arc::new(None), + takeover_locks: Arc::new(StdMutex::new(HashMap::new())), + agent_name: agent_name.into(), + default_cli: String::new(), + } + } + pub fn new( base_url: Option, api_key: impl Into, @@ -1963,7 +1977,7 @@ mod tests { let result = retry_agent_registration_with_budget( "worker-a", Duration::from_millis(10), - || std::future::pending::>(), + std::future::pending::>, |_| std::future::ready(()), ) .await; diff --git a/crates/broker/src/runtime/api.rs b/crates/broker/src/runtime/api.rs index 3af6873dc1..8908d1c9bd 100644 --- a/crates/broker/src/runtime/api.rs +++ b/crates/broker/src/runtime/api.rs @@ -270,6 +270,15 @@ fn observer_token_filters_are_empty(filters: &ObserverTokenFilters) -> bool { impl BrokerRuntime { pub(super) async fn handle_api_request(&mut self, req: ListenApiRequest) { + let req = if self.degraded.is_some() { + match self.handle_local_request(req).await { + Some(req) => req, + None => return, + } + } else { + req + }; + let local_only = self.degraded.is_some(); let paths = &self.paths; let state = &mut self.state; let workspaces = &self.workspaces; @@ -354,7 +363,7 @@ impl BrokerRuntime { let _ = reply.send(Err(format!("agent '{name}' already exists"))); return; } - let owns_identity = agent_token.is_none(); + let owns_identity = !local_only && agent_token.is_none(); let effective_channels = channels.unwrap_or_else(default_spawn_channels); let effective_channels = match super::relaycast_events::relaycast_spawn_channels( &json!({"channels": effective_channels}), @@ -386,6 +395,10 @@ impl BrokerRuntime { return; } }; + if local_only && agent_token.is_some() { + let _ = reply.send(Err("DEGRADED: supplied Relaycast agent tokens are unsupported in local-only mode".into())); + return; + } let mut preregistration_warning: Option = None; // Caller-supplied agent_token is authoritative. In fleet mode it // was minted by the node control connection, and the worker must @@ -396,7 +409,10 @@ impl BrokerRuntime { // so the worker MCP never re-registers over HTTP. let mut fleet_registration = None; let session_ref = super::fleet::fleet_initial_session_ref(&spec); - let worker_relay_key = if let Some(token) = agent_token { + let worker_relay_key = if local_only { + preregistration_warning = Some(super::degraded::WARNING.into()); + None + } else if let Some(token) = agent_token { seed_supplied_agent_token(relaycast_http, &name, &token); match super::fleet::resolve_fleet_agent_token_identity( relaycast_http, @@ -531,6 +547,13 @@ impl BrokerRuntime { } } + let skip_relay_prompt = skip_relay_prompt || local_only; + let task = if local_only { + normalize_initial_task(task) + .map(|task| format!("{}\n\n{}", super::degraded::WARNING, task)) + } else { + task + }; let mut effective_task = if exit_after_task { Some(apply_exit_after_task_instruction(task)) } else { @@ -2041,13 +2064,14 @@ impl BrokerRuntime { .collect(); let auth_workspaces: Vec = workspaces .iter() + .filter(|_| !local_only) .map(|workspace| { json!({ "workspace_id": workspace.workspace_id, "workspace_alias": workspace.workspace_alias, "self_name": workspace.self_name, "self_agent_id": workspace.self_agent_id, - "authenticated": true, + "authenticated": !local_only, "default": default_workspace_id .as_deref() .is_some_and(|id| id == workspace.workspace_id), @@ -2055,6 +2079,9 @@ impl BrokerRuntime { }) .collect(); let _ = reply.send(Ok(json!({ + "mode": if local_only { "local_only" } else { "normal" }, + "status": if local_only { "degraded" } else { "running" }, + "degraded": self.degraded.as_ref().map(|state| state.status()), "agent_count": workers.workers.len(), "agents": workers.list(&super::delivery::pending_message_counts( delivery_states, @@ -2069,7 +2096,7 @@ impl BrokerRuntime { }, "dead_letter_count": dead_letters.len(), "auth": { - "authenticated": !auth_workspaces.is_empty(), + "authenticated": !local_only && !auth_workspaces.is_empty(), "workspace_count": auth_workspaces.len(), "default_workspace_id": default_workspace_id, "workspaces": auth_workspaces, diff --git a/crates/broker/src/runtime/degraded.rs b/crates/broker/src/runtime/degraded.rs new file mode 100644 index 0000000000..59e7c30a7a --- /dev/null +++ b/crates/broker/src/runtime/degraded.rs @@ -0,0 +1,653 @@ +//! Deliberately reduced operation. Local delivery never uses Relaycast routing. +//! Reconciliation publishes audit records, not DMs: replaying a DM would execute +//! work twice and could route an old local name to a different machine. +use super::*; +use sha2::{Digest, Sha256}; +use std::sync::Mutex; + +pub(super) const WARNING: &str = "[agent-relay] DEGRADED — LOCAL ONLY: local spawn, attach and durable delivery are available. Cross-machine routing, worker presence, Relaycast messaging tools and remote attachment are DISABLED. Delivery records reconcile in the background; restart without --local-only to enable fleet capabilities."; +const MAX_OUTBOX_RECORDS: usize = 10_000; +const MAX_OUTBOX_BYTES: usize = 32 * 1024 * 1024; + +pub(super) fn local_session(name: &str) -> RelaySession { + let (ws_control_tx, rx) = mpsc::channel(1); + drop(rx); + let (tx, ws_inbound_rx) = mpsc::channel(1); + drop(tx); + let workspace = RelayWorkspace { + workspace_id: WorkspaceId::new("local"), + workspace_alias: None, + relay_workspace_key: String::new(), + self_name: name.into(), + self_agent_id: AgentId::new("local"), + self_names: HashSet::from([name.to_string()]), + self_agent_ids: HashSet::new(), + http_client: RelaycastHttpClient::local_only(name), + ws_control_tx, + }; + RelaySession { + configured_base: None, + default_workspace_id: Some(workspace.workspace_id.clone()), + workspaces: vec![workspace], + ws_inbound_rx, + } +} + +#[derive(Default, Serialize, Deserialize)] +struct Outbox { + // Digest pins destination without persisting credentials. Refuse to send a + // backlog to a newly selected workspace or service after a restart. + scope: Option, + records: VecDeque, +} + +struct Journal { + path: PathBuf, + outbox: Outbox, + connected: bool, + configured: bool, +} + +impl Journal { + fn open(path: PathBuf, scope: Option) -> Result { + let mut outbox: Outbox = match std::fs::read(&path) { + Ok(bytes) => serde_json::from_slice(&bytes) + .context("local delivery outbox is corrupt; preserve it for recovery")?, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Outbox::default(), + Err(error) => return Err(error).context("cannot read local delivery outbox"), + }; + if !outbox.records.is_empty() { + anyhow::ensure!(outbox.scope == scope, "local delivery outbox belongs to another reconciliation destination; restore the original configuration or use a different state directory"); + } + outbox.scope = scope.clone(); + let journal = Self { + path, + outbox, + connected: false, + configured: scope.is_some(), + }; + journal.save()?; + Ok(journal) + } + + fn save(&self) -> Result<()> { + crate::util::fs::write_json_atomic(&self.path, &self.outbox) + } + + fn enqueue(&mut self, record: Value) -> Result<()> { + anyhow::ensure!( + self.outbox.records.len() < MAX_OUTBOX_RECORDS, + "local delivery reconciliation outbox is full" + ); + self.outbox.records.push_back(record); + let result = (|| { + anyhow::ensure!( + serde_json::to_vec(&self.outbox)?.len() <= MAX_OUTBOX_BYTES, + "local delivery reconciliation outbox is full" + ); + self.save() + })(); + if result.is_err() { + self.outbox.records.pop_back(); + } + result + } + + fn acknowledge(&mut self, event_id: &str) -> Result<()> { + if self + .outbox + .records + .front() + .and_then(|r| r["event_id"].as_str()) + != Some(event_id) + { + return Ok(()); + } + let record = self.outbox.records.pop_front().expect("front was checked"); + if let Err(error) = self.save() { + self.outbox.records.push_front(record); + return Err(error); + } + Ok(()) + } +} + +pub(super) struct DegradedState { + journal: Arc>, + task: tokio::task::JoinHandle<()>, +} + +impl Drop for DegradedState { + fn drop(&mut self) { + self.task.abort(); + } +} + +fn validate_audit_endpoint(base: Option<&str>) -> Result<()> { + let url = reqwest::Url::parse(base.unwrap_or("https://cast.agentrelay.com")) + .context("invalid local reconciliation endpoint")?; + let loopback = url.host_str().is_some_and(|host| { + host.trim_matches(['[', ']']) + .parse::() + .is_ok_and(|ip| ip.is_loopback()) + }); + anyhow::ensure!( + url.scheme() == "https" || (url.scheme() == "http" && loopback), + "local reconciliation requires HTTPS, except for literal loopback HTTP endpoints" + ); + anyhow::ensure!( + url.username().is_empty() && url.password().is_none(), + "local reconciliation endpoint must not contain URL credentials" + ); + Ok(()) +} + +impl DegradedState { + pub(super) fn start(paths: &RuntimePaths, name: &str) -> Result { + let key = std::env::var("AGENT_RELAY_WORKSPACE_KEY") + .ok() + .or_else(|| std::env::var("RELAY_WORKSPACE_KEY").ok()) + .filter(|key| !key.trim().is_empty()); + let base = std::env::var("RELAYCAST_BASE_URL") + .ok() + .or_else(|| std::env::var("RELAY_BASE_URL").ok()) + .filter(|base| !base.trim().is_empty()); + if key.is_some() { + validate_audit_endpoint(base.as_deref())?; + } + let scope = key.as_ref().map(|key| { + format!( + "{:x}", + Sha256::digest(format!( + "{}\0{}", + base.as_deref() + .unwrap_or("https://cast.agentrelay.com") + .trim_end_matches('/'), + key + )) + ) + }); + let journal = Arc::new(Mutex::new(Journal::open( + paths.state.with_extension("local-outbox.json"), + scope, + )?)); + let audit_http = audit_http_client()?; + let task_journal = journal.clone(); + let identity = stable_node_identity_key(&paths.state); + // A separate identity cannot rotate the normal broker's token or claim + // that the local workers are reachable through a fleet node. + let name = format!( + "{name}-local-{}", + &identity[identity.len().saturating_sub(12)..] + ); + let task = tokio::spawn(async move { + let Some(key) = key else { + return; + }; + let auth = AuthClient::new(base.clone()); + let mut client = None; + let mut retry_delay = true; + loop { + if retry_delay { + tokio::time::sleep(Duration::from_secs(5)).await; + } + retry_delay = true; + let pending = task_journal.lock().unwrap().outbox.records.front().cloned(); + let Some(record) = pending else { + continue; + }; + if client.is_none() { + let result = timeout( + Duration::from_secs(12), + auth.startup_session_set_with_identity( + Some(&name), + true, + Some("agent"), + Some(&identity), + ), + ) + .await; + if let Ok(Ok(sessions)) = result { + if let Some(session) = sessions.default_session() { + // No websocket, node registration, or worker presence + // is enabled by an audit connection. + let http = + RelaycastHttpClient::new(base.clone(), key.clone(), &name, "local"); + if let Some(token) = session.credentials.agent_token.as_deref() { + http.seed_agent_token(&name, token); + } + client = Some(http); + } + } + } + let Some(http) = client.as_ref() else { + task_journal.lock().unwrap().connected = false; + continue; + }; + if reconcile_record(&audit_http, http, &name, &record) + .await + .is_ok() + { + let mut journal = task_journal.lock().unwrap(); + if !journal.connected { + eprintln!("[agent-relay] local delivery audit reconciliation connected (broker operating mode unchanged)"); + } + journal.connected = true; + retry_delay = false; + if journal + .acknowledge(record["event_id"].as_str().unwrap_or_default()) + .is_err() + { + retry_delay = true; + eprintln!("[agent-relay] could not persist local reconciliation acknowledgement; record retained for retry"); + } + } else { + task_journal.lock().unwrap().connected = false; + client = None; + } + } + }); + Ok(Self { journal, task }) + } + + pub(super) fn status(&self) -> Value { + let journal = self.journal.lock().unwrap(); + json!({ + "mode": "local_only", "status": "degraded", + "capabilities": {"local_spawn": true, "local_attach": true, "local_queue": true, + "cross_machine_routing": false, "worker_presence": false, "remote_delivery": false, + "remote_attach": false, "relaycast_tools": false}, + "reconciliation": {"configured": journal.configured, "connected": journal.connected, + "pending_records": journal.outbox.records.len(), "delivery_semantics": "audit_only_at_least_once"} + }) + } + + pub(super) fn enqueue(&self, record: Value) -> Result<()> { + self.journal.lock().unwrap().enqueue(record) + } +} + +fn audit_http_client() -> Result { + Ok(reqwest::Client::builder() + .redirect(reqwest::redirect::Policy::none()) + .timeout(Duration::from_secs(5)) + .build()?) +} + +async fn reconcile_record( + client: &reqwest::Client, + http: &RelaycastHttpClient, + name: &str, + record: &Value, +) -> Result<()> { + validate_audit_endpoint(http.base_url.as_deref())?; + // A redirect must not send the retained body or workspace credential to + // an endpoint that was never pinned by the journal's destination digest. + let response = client + .post(format!( + "{}/v1/agents/{}/events", + http.base_url + .as_deref() + .unwrap_or("https://cast.agentrelay.com") + .trim_end_matches('/'), + urlencoding::encode(name) + )) + .bearer_auth(&http.api_key) + .header( + "X-Relaycast-Origin-Actor", + crate::telemetry::BROKER_ORIGIN_ACTOR, + ) + .json(&json!({"type": "local.delivery.queued", "payload": record})) + .send() + .await?; + anyhow::ensure!( + response.status().is_success(), + "local audit event was not accepted" + ); + let envelope: Value = response.json().await?; + anyhow::ensure!(envelope["ok"] == true, "local audit event was rejected"); + let _: relaycast::SessionEvent = serde_json::from_value(envelope["data"].clone()) + .context("invalid local audit acknowledgement")?; + Ok(()) +} + +impl BrokerRuntime { + pub(super) async fn handle_local_request( + &mut self, + req: ListenApiRequest, + ) -> Option { + let req = match req { + ListenApiRequest::SetInboundDeliveryMode { + mode: InboundDeliveryMode::ManualFlush, + reply, + .. + } => { + let _ = reply.send(Err(DeliveryRouteError::CapabilityDisabled)); + return None; + } + req => req, + }; + let ListenApiRequest::Send { + to, + text, + from, + thread_id, + workspace_id, + workspace_alias, + mode, + reply, + } = req + else { + return Some(req); + }; + let result = self + .queue_local_delivery( + to, + text, + from, + thread_id, + workspace_id, + workspace_alias, + mode, + ) + .await; + let _ = reply.send(result.map_err(|error| error.to_string())); + None + } + + #[allow(clippy::too_many_arguments)] + async fn queue_local_delivery( + &mut self, + to: MessageTarget, + text: String, + from: Option, + thread_id: Option, + workspace_id: Option, + workspace_alias: Option, + mode: MessageInjectionMode, + ) -> Result { + anyhow::ensure!( + workspace_alias.is_none() && workspace_id.as_deref().is_none_or(|id| id == "local"), + "DEGRADED: workspace routing is disabled in local-only mode" + ); + let name = to.trim().trim_start_matches('@'); + anyhow::ensure!( + self.workers.has_worker(name), + "DEGRADED: only agents running on this broker can receive local deliveries" + ); + let delivery_id = DeliveryId::new(format!("del_{}", Uuid::new_v4().simple())); + let event_id = EventId::new(format!("local_{}", Uuid::new_v4().simple())); + let from = normalize_sender(from); + let queued_at_ms = unix_timestamp_millis(); + let delivery = RelayDelivery { + delivery_id: delivery_id.clone(), + event_id: event_id.clone(), + from: from.clone(), + target: MessageTarget::new(name), + body: text.clone(), + thread_id: thread_id.clone(), + workspace_id: Some(WorkspaceId::new("local")), + workspace_alias: None, + priority: Some(1), + injection_mode: mode, + }; + // Commit locally before any handoff or success response. The normal + // pending-delivery lifecycle retains/acks/dead-letters this work, just + // as it does engine deliveries; it survives broker restarts. + self.pending_deliveries.insert( + delivery_id.clone(), + PendingDelivery { + worker_name: WorkerName::new(name), + delivery, + attempts: 0, + failed_attempts: 0, + next_retry_at: Instant::now(), + queued_at_ms, + last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, + }, + ); + if let Err(error) = save_pending_deliveries(&self.paths.pending, &self.pending_deliveries) { + self.pending_deliveries.remove(&delivery_id); + return Err(error).context("local delivery was not accepted: persistence failed"); + } + let record = json!({"event_id": event_id, "delivery_id": delivery_id, + "from": from, "to": name, "body": text, "thread_id": thread_id, + "queued_at_ms": queued_at_ms, "delivery_status": "queued_local", "mode": "local_only"}); + if let Err(error) = self.degraded.as_ref().expect("local mode").enqueue(record) { + self.pending_deliveries.remove(&delivery_id); + save_pending_deliveries(&self.paths.pending, &self.pending_deliveries) + .context("local delivery acceptance uncertain: rollback persistence failed")?; + return Err(error).context("local delivery was not accepted"); + } + // Maintenance performs handoff, including after a reconnect/restart. + // Never label this delivered: the PTY still owes an acknowledgement. + Ok( + json!({"success": true, "event_id": event_id, "delivery_id": delivery_id, + "local": true, "mode": "local_only", "relaycast_published": false, + "delivery_status": "queued_local", "reconciliation_pending": true, + "workspace_id": "local"}), + ) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn outbox_survives_restart_and_refuses_a_different_destination() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("outbox.json"); + let record = json!({"event_id": "local_one", "body": "retained work"}); + let mut journal = Journal::open(path.clone(), Some("destination-a".into())).unwrap(); + journal.enqueue(record.clone()).unwrap(); + drop(journal); + assert!(Journal::open(path.clone(), Some("destination-b".into())).is_err()); + let mut recovered = Journal::open(path.clone(), Some("destination-a".into())).unwrap(); + assert_eq!(recovered.outbox.records.front(), Some(&record)); + recovered.acknowledge("unrelated").unwrap(); + assert_eq!(recovered.outbox.records.len(), 1); + recovered.acknowledge("local_one").unwrap(); + assert!(Journal::open(path, Some("destination-a".into())) + .unwrap() + .outbox + .records + .is_empty()); + } + + #[test] + fn unscoped_backlog_cannot_acquire_an_upload_destination() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("outbox.json"); + let record = json!({"event_id": "local_private", "body": "local-only work"}); + let mut journal = Journal::open(path.clone(), None).unwrap(); + journal.enqueue(record.clone()).unwrap(); + drop(journal); + let original = std::fs::read(&path).unwrap(); + + assert!(Journal::open(path.clone(), Some("new-destination".into())).is_err()); + assert_eq!(std::fs::read(&path).unwrap(), original); + let recovered = Journal::open(path.clone(), None).unwrap(); + assert_eq!(recovered.outbox.scope, None); + assert_eq!(recovered.outbox.records.front(), Some(&record)); + assert!(!recovered.configured); + + // A fresh, empty journal may still be configured normally. + let empty = dir.path().join("empty.json"); + Journal::open(empty.clone(), None).unwrap(); + assert!(Journal::open(empty, Some("new-destination".into())).is_ok()); + } + + #[test] + fn corrupt_or_full_outbox_is_never_silently_discarded() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("outbox.json"); + std::fs::write(&path, "{").unwrap(); + assert!(Journal::open(path.clone(), None).is_err()); + std::fs::remove_file(&path).unwrap(); + let mut journal = Journal::open(path, None).unwrap(); + journal.outbox.records = (0..MAX_OUTBOX_RECORDS) + .map(|id| json!({"event_id": id})) + .collect(); + assert!(journal.enqueue(json!({"event_id": "overflow"})).is_err()); + assert_eq!(journal.outbox.records.len(), MAX_OUTBOX_RECORDS); + } + + #[tokio::test] + async fn reconnect_replays_the_retained_audit_record_without_remote_delivery() { + use httpmock::Method::POST; + let server = httpmock::MockServer::start_async().await; + let unavailable = server + .mock_async(|when, then| { + when.method(POST).path("/v1/agents/local-broker/events"); + then.status(503).json_body( + json!({"ok":false,"error":{"code":"unavailable","message":"test outage"}}), + ); + }) + .await; + let remote_delivery = server + .mock_async(|when, then| { + when.method(POST).path("/v1/dm"); + then.status(500); + }) + .await; + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("outbox.json"); + let record = json!({"event_id":"local_replay", "delivery_id":"del_replay", "body":"local work", "delivery_status":"queued_local"}); + let mut journal = Journal::open(path.clone(), None).unwrap(); + journal.enqueue(record.clone()).unwrap(); + let http = RelaycastHttpClient::new( + Some(server.base_url()), + "rk_live_test", + "local-broker", + "local", + ); + assert!(reconcile_record( + &audit_http_client().unwrap(), + &http, + "local-broker", + &record + ) + .await + .is_err()); + assert_eq!( + Journal::open(path.clone(), None) + .unwrap() + .outbox + .records + .front(), + Some(&record) + ); + unavailable.delete_async().await; + let accepted = server.mock_async(|when, then| { + when.method(POST).path("/v1/agents/local-broker/events") + .json_body(json!({"type":"local.delivery.queued", "payload": record})); + then.status(200).json_body(json!({"ok":true,"data":{"id":"audit_one","agent_id":"local-broker","type":"local.delivery.queued","payload":record,"created_at":"2026-09-09T00:00:00Z"}})); + }).await; + reconcile_record( + &audit_http_client().unwrap(), + &http, + "local-broker", + &record, + ) + .await + .unwrap(); + journal.acknowledge("local_replay").unwrap(); + accepted.assert_hits_async(1).await; + remote_delivery.assert_hits_async(0).await; + assert!(Journal::open(path, None).unwrap().outbox.records.is_empty()); + } + #[test] + fn audit_endpoint_requires_tls_outside_literal_loopback() { + for url in [ + "https://cast.agentrelay.com", + "http://127.0.0.1:8787", + "http://[::1]:8787", + ] { + assert!(validate_audit_endpoint(Some(url)).is_ok()); + } + for url in [ + "http://remote.example", + "http://localhost:8787", + "ftp://127.0.0.1", + "https://user:password@example.com", + ] { + assert!(validate_audit_endpoint(Some(url)).is_err()); + } + } + + #[tokio::test] + async fn audit_redirect_cannot_forward_private_record_to_another_destination() { + use httpmock::Method::POST; + let source = httpmock::MockServer::start_async().await; + let destination = httpmock::MockServer::start_async().await; + let leak = destination + .mock_async(|when, then| { + when.method(POST); + then.status(200).json_body(json!({"ok": true})); + }) + .await; + let redirect = source + .mock_async(|when, then| { + when.method(POST).path("/v1/agents/local-broker/events"); + then.status(307) + .header("Location", destination.url("/private-data")); + }) + .await; + let http = RelaycastHttpClient::new( + Some(source.base_url()), + "rk_live_test", + "local-broker", + "local", + ); + assert!(reconcile_record( + &audit_http_client().unwrap(), + &http, + "local-broker", + &json!({"event_id": "local_private", "body": "private work"}) + ) + .await + .is_err()); + redirect.assert_hits_async(1).await; + leak.assert_hits_async(0).await; + } + #[tokio::test] + async fn registration_sdk_strips_workspace_credential_on_cross_origin_redirect() { + use httpmock::Method::POST; + let source = httpmock::MockServer::start_async().await; + let destination = httpmock::MockServer::start_async().await; + let leaked_key = destination + .mock_async(|when, then| { + when.method(POST) + .header("authorization", "Bearer rk_live_test"); + then.status(401); + }) + .await; + let redirected = destination.mock_async(|when, then| { + when.method(POST); + then.status(401).json_body(json!({"ok": false, "error": {"code": "unauthorized", "message": "no credential"}})); + }).await; + source + .mock_async(|when, then| { + when.method(POST).path("/v1/agents"); + then.status(307) + .header("Location", destination.url("/v1/agents")); + }) + .await; + let relay = relaycast::RelayCast::new( + relaycast::RelayCastOptions::new("rk_live_test").with_base_url(source.base_url()), + ) + .unwrap(); + let request = relaycast::CreateAgentRequest { + name: "local-broker".into(), + agent_type: Some("agent".into()), + persona: None, + metadata: None, + }; + assert!(relay.register_agent(request).await.is_err()); + leaked_key.assert_hits_async(0).await; + redirected.assert_hits_async(1).await; + } +} diff --git a/crates/broker/src/runtime/delivery.rs b/crates/broker/src/runtime/delivery.rs index 159c9ef9a0..1299bf8742 100644 --- a/crates/broker/src/runtime/delivery.rs +++ b/crates/broker/src/runtime/delivery.rs @@ -258,13 +258,20 @@ pub(crate) fn load_pending_deliveries(path: &Path) -> HashMap Option<& if event_id.is_empty() { return Some("blank_event_id"); } + if event_id.starts_with("local_") { + return Some("local_only_synthetic_event_id"); + } if event_id.starts_with("http_") { return Some("http_api_synthetic_event_id"); } @@ -935,6 +945,20 @@ pub(crate) async fn retry_pending_delivery( None => return Ok(DeliveryAttemptOutcome::Noop), }; + // A local queue can outlive its broker and worker. Check absence before + // retry exhaustion, and give a respawned recipient a fresh handoff budget. + // Explicit release still moves its pending deliveries to dead letters. + if pending.delivery.event_id.as_str().starts_with("local_") + && !workers.has_worker(&pending.worker_name) + { + if let Some(current) = pending_deliveries.get_mut(delivery_id) { + current.failed_attempts = 0; + current.next_retry_at = Instant::now() + retry_interval; + current.last_error = Some("waiting for local recipient to reconnect".into()); + } + return Ok(DeliveryAttemptOutcome::Noop); + } + if pending.failed_attempts >= MAX_DELIVERY_RETRIES { let removed = pending_deliveries.remove(delivery_id).unwrap_or(pending); let last_error = removed diff --git a/crates/broker/src/runtime/event_loop.rs b/crates/broker/src/runtime/event_loop.rs index f55943dbad..585fd7a15c 100644 --- a/crates/broker/src/runtime/event_loop.rs +++ b/crates/broker/src/runtime/event_loop.rs @@ -190,6 +190,7 @@ pub(crate) fn commit_resize_ownership( } pub(crate) struct BrokerRuntime { + pub(super) degraded: Option, pub(super) persist: bool, pub(super) broker_start: Instant, pub(super) agent_spawn_count: u32, diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index a8b3323fd0..18c6770b42 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -622,6 +622,14 @@ impl BrokerRuntime { blocked_reason: ok.blocked_reason, }); } + Ok(Err(error @ DeliveryRouteError::CapabilityDisabled)) => { + self.send_terminal(TerminalToCloud::Error { + session_id, + code: "capability_disabled".into(), + message: error.to_string(), + request_id, + }); + } Ok(Err(DeliveryRouteError::WorkerNotFound(name))) => { self.send_terminal(TerminalToCloud::Error { session_id, @@ -726,6 +734,14 @@ impl BrokerRuntime { revision: ok.revision.to_string(), }); } + Ok(Err(error @ DeliveryRouteError::CapabilityDisabled)) => { + self.send_terminal(TerminalToCloud::Error { + session_id, + code: "capability_disabled".into(), + message: error.to_string(), + request_id, + }); + } Ok(Err(DeliveryRouteError::WorkerNotFound(name))) => { self.send_terminal(TerminalToCloud::Error { session_id, diff --git a/crates/broker/src/runtime/init.rs b/crates/broker/src/runtime/init.rs index 12058b3864..7d907f67c0 100644 --- a/crates/broker/src/runtime/init.rs +++ b/crates/broker/src/runtime/init.rs @@ -2,6 +2,26 @@ use super::*; use std::net::{IpAddr, SocketAddr}; pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Result<()> { + let local_only = + cmd.local_only || std::env::var("AGENT_RELAY_LOCAL_ONLY").as_deref() == Ok("1"); + if local_only { + anyhow::ensure!( + cmd.persist || cmd.state_dir.is_some(), + "local-only mode requires --persist or --state-dir for durable delivery state" + ); + let bind: IpAddr = unbracket_ipv6(cmd.api_bind.trim()) + .parse() + .context("local-only mode requires a loopback IP bind address")?; + anyhow::ensure!( + bind.is_loopback(), + "local-only mode requires a loopback API bind address" + ); + anyhow::ensure!( + std::env::var("RELAY_WORKSPACES_JSON").map_or(true, |v| v.trim().is_empty()), + "local-only mode supports one workspace key; unset RELAY_WORKSPACES_JSON" + ); + eprintln!("{}", super::degraded::WARNING); + } let broker_start = Instant::now(); let startup_debug = startup_debug_enabled(); let agent_spawn_count: u32 = 0; @@ -127,7 +147,8 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let relay_ready = Arc::new(Notify::new()); let relay_ready_state: Arc>> = Arc::new(RwLock::new(None)); let (api_tx, api_rx) = mpsc::channel::(32); - let bind_addr = format!("{}:{}", cmd.api_bind, cmd.api_port); + let api_host = bracket_ipv6_host(unbracket_ipv6(cmd.api_bind.trim())); + let bind_addr = format!("{}:{}", api_host, cmd.api_port); log_startup_phase( startup_debug, broker_start, @@ -141,23 +162,24 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re log_startup_phase( startup_debug, broker_start, - format!("API listener bound on {}:{}", cmd.api_bind, actual_port), + format!("API listener bound on {}:{}", api_host, actual_port), ); // Machine-readable on stdout (SDK parses this to discover the port). // Diagnostic logs stay on stderr via tracing/eprintln. println!( "[agent-relay] API listening on http://{}:{}", - cmd.api_bind, actual_port + api_host, actual_port ); // Write connection file so CLI commands can find this broker. let connection_dir = paths.state.parent().unwrap(); let connection_path = connection_dir.join("connection.json"); let connection = json!({ - "url": format!("http://{}:{}", cmd.api_bind, actual_port), + "url": format!("http://{}:{}", api_host, actual_port), "port": actual_port, "api_key": &api_key, "pid": std::process::id(), + "operation_mode": if local_only { "local_only" } else { "normal" }, }); if let Ok(json_str) = serde_json::to_string_pretty(&connection) { if let Ok(mut tmp) = tempfile::NamedTempFile::new_in(connection_dir) { @@ -173,24 +195,29 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re tokio::sync::oneshot::channel::(); let relay_ready_for_startup = relay_ready.clone(); tokio::spawn(async move { - let listener = serve_startup_api_until_ready(listener, relay_ready_for_startup).await; + let listener = + serve_startup_api_until_ready(listener, relay_ready_for_startup, local_only).await; let _ = startup_listener_tx.send(listener); }); log_startup_phase(startup_debug, broker_start, "calling connect_relay"); - let relay = connect_relay(RelaySessionOptions { - paths: &paths, - requested_name: &resolved_name, - channels: channels_from_csv(&cmd.channels), - // Ephemeral brokers are short-lived and frequently restarted by tests/SDK - // callers. Use non-strict registration so stale Relaycast identities from - // prior runs don't hard-fail startup. - strict_name: cmd.persist, - agent_type: Some(agent_type_ref), - read_mcp_identity: true, - runtime_cwd: &runtime_cwd, - }) - .await?; + let relay = if local_only { + super::degraded::local_session(&resolved_name) + } else { + connect_relay(RelaySessionOptions { + paths: &paths, + requested_name: &resolved_name, + channels: channels_from_csv(&cmd.channels), + // Ephemeral brokers are short-lived and frequently restarted by tests/SDK + // callers. Use non-strict registration so stale Relaycast identities from + // prior runs don't hard-fail startup. + strict_name: cmd.persist, + agent_type: Some(agent_type_ref), + read_mcp_identity: true, + runtime_cwd: &runtime_cwd, + }) + .await? + }; log_startup_phase(startup_debug, broker_start, "connect_relay completed"); let RelaySession { @@ -235,7 +262,11 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re hosted_agent_event_rx, )); let node_workspace_id = default_workspace.workspace_id.as_str().to_string(); - let node_id = resolve_broker_node_id(&node_workspace_id); + let node_id = if local_only { + format!("local-{}", Uuid::new_v4().simple()) + } else { + resolve_broker_node_id(&node_workspace_id) + }; // The node registers under its resolved instance name (--instance-name, the // legacy --name/--broker-name alias, or AGENT_RELAY_BROKER_NAME), falling back // to the machine hostname only when none is set. Deriving this from the raw @@ -271,8 +302,11 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re // node-control client mints one in the background (it holds the same minter) // and publishes it to `session_node_token`, so realtime delivery still comes // online without gating startup on it. - let node_token = - resolve_cached_node_token(&node_id, &node_workspace_id, node_base_url.as_deref()); + let node_token = if local_only { + None + } else { + resolve_cached_node_token(&node_id, &node_workspace_id, node_base_url.as_deref()) + }; let node_manifest = bootstrap_node_manifest(&node_name, &node_id, &broker_version); // Retain the node name for the runtime: the HTTP `bind_agent_to_node` // fallback (used when node-control `agent.register` is unavailable) binds @@ -320,40 +354,47 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let (terminal_event_tx, terminal_event_rx) = mpsc::channel::(1024); let node_delivery_token_present = node_token.is_some(); - tokio::spawn(crate::node_control::run_node_control_client( - crate::node_control::FleetControlConfig { - ws_url: fleet_ws_url, - node_token, - node_id, - node_name, - broker_version, - token_minter, - session_token: Some(session_node_token.clone()), - read_idle_timeout: None, - }, - fleet_control_rx, - fleet_event_tx, - )); - tokio::spawn(crate::terminal_control::run_terminal_control_client( - crate::terminal_control::TerminalControlConfig { - ws_url: terminal_ws_url, - session_token: session_node_token.clone(), - read_idle_timeout: None, - }, - terminal_control_rx, - terminal_event_tx, - )); - // Register this node unconditionally on connect (no sidecar required). This - // is the only command that flips the control client out of its idle state - // and into the connect loop, so the broker enrolls every startup. - if let Err(error) = fleet_control_tx - .send(FleetControlCommand::RegisterNode { - manifest: node_manifest, - resume_cursor: None, - }) - .await - { - tracing::warn!(error = %error, "failed to queue node.register at startup"); + if !local_only { + tokio::spawn(crate::node_control::run_node_control_client( + crate::node_control::FleetControlConfig { + ws_url: fleet_ws_url, + node_token, + node_id, + node_name, + broker_version, + token_minter, + session_token: Some(session_node_token.clone()), + read_idle_timeout: None, + }, + fleet_control_rx, + fleet_event_tx, + )); + tokio::spawn(crate::terminal_control::run_terminal_control_client( + crate::terminal_control::TerminalControlConfig { + ws_url: terminal_ws_url, + session_token: session_node_token.clone(), + read_idle_timeout: None, + }, + terminal_control_rx, + terminal_event_tx, + )); + // Register this node unconditionally on connect (no sidecar required). This + // is the only command that flips the control client out of its idle state + // and into the connect loop, so the broker enrolls every startup. + if let Err(error) = fleet_control_tx + .send(FleetControlCommand::RegisterNode { + manifest: node_manifest, + resume_cursor: None, + }) + .await + { + tracing::warn!(error = %error, "failed to queue node.register at startup"); + } + } else { + drop(fleet_control_rx); + drop(fleet_event_tx); + drop(terminal_control_rx); + drop(terminal_event_tx); } let workspace_memberships: Vec = workspaces .iter() @@ -384,18 +425,52 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let (events_tx, _events_rx) = broadcast::channel::(512); let replay_buffer = ReplayBuffer::new(DEFAULT_REPLAY_CAPACITY); + let degraded = if local_only { + Some(super::degraded::DegradedState::start( + &paths, + &resolved_name, + )?) + } else { + None + }; + // A normal restart also drains retained audit records without replaying + // them as messages or changing the normal broker capability state. A + // backlog that cannot be reconciled must not block normal startup. + let _previous_local_reconciliation = + if !local_only && paths.state.with_extension("local-outbox.json").exists() { + match super::degraded::DegradedState::start(&paths, &resolved_name) { + Ok(reconciliation) => Some(reconciliation), + Err(error) => { + eprintln!( + "[agent-relay] local audit backlog retained without reconciliation: {error}" + ); + None + } + } + } else { + None + }; let ready_router = listen_api_router(ListenApiConfig { + local_only, tx: api_tx.clone(), events_tx: events_tx.clone(), replay_buffer: replay_buffer.clone(), - workspace_key: Some(relay_workspace_key.clone()), + workspace_key: (!local_only).then(|| relay_workspace_key.clone()), relay_base_url: configured_base.clone(), - memberships: workspace_memberships.clone(), - default_workspace_id: default_workspace_id.clone(), + memberships: if local_only { + vec![] + } else { + workspace_memberships.clone() + }, + default_workspace_id: if local_only { + None + } else { + default_workspace_id.clone() + }, node_id: session_node_id, node_name: session_node_name, node_token: session_node_token, - persist: cmd.persist, + persist: paths.persist, }); { let mut ready = relay_ready_state.write().await; @@ -524,6 +599,11 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re )); } + if local_only { + worker_env.retain(|(key, _)| key == "AGENT_RELAY_RESULT_URL"); + worker_env.push(("AGENT_RELAY_LOCAL_ONLY".into(), "1".into())); + } + let (sdk_out_tx, mut sdk_out_rx) = mpsc::channel::>(1024); let events_tx_for_stdout = events_tx.clone(); let replay_buffer_for_stdout = replay_buffer.clone(); @@ -643,7 +723,7 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re // Owner lease: in ephemeral mode, the broker shuts down if the SDK // doesn't renew the lease within this duration. Replaces stdin EOF // detection. Disabled in persist mode. - let lease_duration = if cmd.persist { + let lease_duration = if paths.persist { None } else { Some(Duration::from_secs(120)) @@ -661,7 +741,8 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let mut sigterm = tokio::signal::windows::ctrl_shutdown()?; let runtime = BrokerRuntime { - persist: cmd.persist, + degraded, + persist: paths.persist, broker_start, agent_spawn_count, paths, diff --git a/crates/broker/src/runtime/mod.rs b/crates/broker/src/runtime/mod.rs index f817e2b1e3..af0e59e4a4 100644 --- a/crates/broker/src/runtime/mod.rs +++ b/crates/broker/src/runtime/mod.rs @@ -72,6 +72,7 @@ mod api; mod app_server; mod connection; mod dead_letter; +mod degraded; mod delivery; mod event_loop; mod fleet; diff --git a/crates/broker/src/runtime/session.rs b/crates/broker/src/runtime/session.rs index 366c4c9e1e..13cd74df0c 100644 --- a/crates/broker/src/runtime/session.rs +++ b/crates/broker/src/runtime/session.rs @@ -31,6 +31,7 @@ pub(crate) struct RelayReadyState { pub(crate) async fn serve_startup_api_until_ready( listener: tokio::net::TcpListener, relay_ready: Arc, + local_only: bool, ) -> tokio::net::TcpListener { loop { tokio::select! { @@ -40,7 +41,7 @@ pub(crate) async fn serve_startup_api_until_ready( accepted = listener.accept() => { match accepted { Ok((stream, _addr)) => { - tokio::spawn(handle_startup_api_connection(stream)); + tokio::spawn(handle_startup_api_connection(stream, local_only)); } Err(error) => { tracing::warn!(error = %error, "startup API accept failed"); @@ -52,7 +53,10 @@ pub(crate) async fn serve_startup_api_until_ready( } } -pub(crate) async fn handle_startup_api_connection(mut stream: tokio::net::TcpStream) { +pub(crate) async fn handle_startup_api_connection( + mut stream: tokio::net::TcpStream, + local_only: bool, +) { let mut buffer = [0_u8; 1024]; let read = match timeout(Duration::from_secs(5), stream.read(&mut buffer)).await { Ok(Ok(read)) => read, @@ -70,11 +74,15 @@ pub(crate) async fn handle_startup_api_connection(mut stream: tokio::net::TcpStr .and_then(|line| line.split_whitespace().nth(1)) .unwrap_or("/"); let (status, content_type, body) = if path == "/health" { - ( - "200 OK", - "application/json", - listen_api::listen_api_health_payload(None, vec![]).to_string(), - ) + ("200 OK", "application/json", { + let mut payload = listen_api::listen_api_health_payload(None, vec![]); + if local_only { + payload["status"] = json!("degraded"); + payload["mode"] = json!("local_only"); + payload["relaycastConnected"] = json!(false); + } + payload.to_string() + }) } else { ( "503 Service Unavailable", diff --git a/crates/broker/src/runtime/tests.rs b/crates/broker/src/runtime/tests.rs index 15eaffe6f0..d96722a563 100644 --- a/crates/broker/src/runtime/tests.rs +++ b/crates/broker/src/runtime/tests.rs @@ -608,6 +608,7 @@ fn worker_event_runtime_fixture( tokio::signal::windows::ctrl_shutdown().expect("install test Ctrl+Shutdown listener"); let runtime = BrokerRuntime { + degraded: None, persist: false, broker_start: Instant::now(), agent_spawn_count: 0, @@ -6347,3 +6348,103 @@ async fn owned_cleanup_retries_delete_without_repeating_acknowledged_deregistrat .contains_key(&name)); fixture.runtime.workers.release("unrelated").await.unwrap(); } + +#[tokio::test] +async fn local_only_queued_work_survives_restart_and_replays_when_recipient_reconnects() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("pending.json"); + let id = DeliveryId::new("del_local_reconnect"); + let mut pending = HashMap::from([( + id.clone(), + pending_delivery("local-worker", id.as_str(), "local_reconnect"), + )]); + pending.get_mut(&id).unwrap().attempts = 0; + pending.get_mut(&id).unwrap().delivery.workspace_id = Some(WorkspaceId::new("local")); + pending.get_mut(&id).unwrap().delivery.workspace_alias = None; + super::save_pending_deliveries(&path, &pending).unwrap(); + pending = load_pending_deliveries(&path); + let (tx, _rx) = mpsc::channel(8); + let mut absent = WorkerRegistry::new(tx, vec![], dir.path().join("logs"), Instant::now()); + assert!(matches!( + retry_pending_delivery(&id, &mut absent, &mut pending, Duration::from_secs(1)) + .await + .unwrap(), + DeliveryAttemptOutcome::Noop + )); + assert_eq!(pending[&id].attempts, 0); + assert!(pending[&id] + .last_error + .as_deref() + .unwrap() + .contains("reconnect")); + let mut reconnected = make_worker_registry_with_worker("local-worker").await; + assert!(matches!( + retry_pending_delivery(&id, &mut reconnected, &mut pending, Duration::from_secs(1)) + .await + .unwrap(), + DeliveryAttemptOutcome::Attempted { .. } + )); + assert_eq!(pending[&id].attempts, 1); + assert_eq!(pending[&id].delivery.event_id.as_str(), "local_reconnect"); + cleanup_worker_registry(reconnected).await; +} + +#[tokio::test] +async fn local_only_exhausted_delivery_survives_absence_and_replays_after_restart() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("pending.json"); + let id = DeliveryId::new("del_local_exhausted"); + let mut entry = pending_delivery("local-worker", id.as_str(), "local_exhausted"); + entry.attempts = MAX_DELIVERY_RETRIES; + entry.failed_attempts = MAX_DELIVERY_RETRIES; + let expected_delivery = entry.delivery.clone(); + let mut pending = HashMap::from([(id.clone(), entry)]); + // Exercise the live exhausted queue first: loading a snapshot resets the + // failure budget and would hide an exhaustion check before absence handling. + let (tx, _rx) = mpsc::channel(8); + let mut absent = WorkerRegistry::new(tx, vec![], dir.path().join("logs"), Instant::now()); + for _ in 0..2 { + assert!(matches!( + retry_pending_delivery(&id, &mut absent, &mut pending, Duration::from_secs(1)) + .await + .unwrap(), + DeliveryAttemptOutcome::Noop + )); + assert_eq!(pending[&id].delivery, expected_delivery); + assert_eq!(pending[&id].attempts, MAX_DELIVERY_RETRIES); + super::save_pending_deliveries(&path, &pending).unwrap(); + pending = load_pending_deliveries(&path); + } + let mut reconnected = make_worker_registry_with_worker("local-worker").await; + let outcome = + retry_pending_delivery(&id, &mut reconnected, &mut pending, Duration::from_secs(1)) + .await + .unwrap(); + cleanup_worker_registry(reconnected).await; + assert!(matches!(outcome, DeliveryAttemptOutcome::Attempted { .. })); + assert_eq!(pending[&id].delivery, expected_delivery); + assert_eq!(pending[&id].attempts, MAX_DELIVERY_RETRIES + 1); + assert_eq!(pending[&id].failed_attempts, 0); +} + +#[tokio::test] +async fn local_only_restored_exhausted_delivery_replays_to_already_registered_worker() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("pending.json"); + let id = DeliveryId::new("del_local_present_on_restart"); + let mut entry = pending_delivery("local-worker", id.as_str(), "local_present_on_restart"); + entry.attempts = MAX_DELIVERY_RETRIES; + entry.failed_attempts = MAX_DELIVERY_RETRIES; + let expected_delivery = entry.delivery.clone(); + super::save_pending_deliveries(&path, &HashMap::from([(id.clone(), entry)])).unwrap(); + let mut pending = load_pending_deliveries(&path); + let mut workers = make_worker_registry_with_worker("local-worker").await; + let outcome = retry_pending_delivery(&id, &mut workers, &mut pending, Duration::from_secs(1)) + .await + .unwrap(); + cleanup_worker_registry(workers).await; + assert!(matches!(outcome, DeliveryAttemptOutcome::Attempted { .. })); + assert_eq!(pending[&id].delivery, expected_delivery); + assert_eq!(pending[&id].attempts, MAX_DELIVERY_RETRIES + 1); + assert_eq!(pending[&id].failed_attempts, 0); +} diff --git a/crates/broker/src/worker.rs b/crates/broker/src/worker.rs index 885c271209..b582f1c269 100644 --- a/crates/broker/src/worker.rs +++ b/crates/broker/src/worker.rs @@ -456,7 +456,7 @@ impl WorkerRegistry { // save tokens. We honor that even when `agent_result` is configured — // `AGENT_RELAY_RESULT_*` env vars are still set on the worker process // below, so a separately-configured Agent Relay MCP can pick them up. - if skip_relay_prompt { + if skip_relay_prompt || self.env_value("AGENT_RELAY_LOCAL_ONLY") == Some("1") { return Ok(Vec::new()); } configure_agent_relay_mcp_with_result( @@ -1253,8 +1253,27 @@ impl WorkerRegistry { command.env("RELAY_AGENT_TYPE", "agent"); command.env("RELAY_STRICT_AGENT_NAME", "1"); } - // Remove CLAUDECODE from child env to prevent nested Claude Code instances - // from interfering with the parent's session management + // Local-only workers must not bootstrap a separate Relaycast session. + if self.env_value("AGENT_RELAY_LOCAL_ONLY") == Some("1") { + for key in [ + "AGENT_RELAY_ORIGIN_ACTOR", + "RELAY_AGENT_NAME", + "RELAY_AGENT_TYPE", + "RELAY_STRICT_AGENT_NAME", + "AGENT_RELAY_WORKSPACE_KEY", + "RELAY_WORKSPACE_KEY", + "RELAY_API_KEY", + "RELAY_AGENT_TOKEN", + "RELAY_NODE_TOKEN", + "RELAY_WORKSPACES_JSON", + "RELAY_DEFAULT_WORKSPACE", + "RELAY_WORKSPACE_ID", + ] { + command.env_remove(key); + } + command.env("AGENT_RELAY_LOCAL_ONLY", "1"); + } + // Prevent nested Claude Code instances from sharing the parent session. command.env_remove("CLAUDECODE"); if let Some(cwd) = spec.cwd.as_ref() { command.current_dir(cwd); diff --git a/packages/cli/README.md b/packages/cli/README.md index 01990324f6..757f2dd910 100644 --- a/packages/cli/README.md +++ b/packages/cli/README.md @@ -48,6 +48,70 @@ agent-relay node agent release For AI SDK native harnesses, attach renders structured activity, text, tools, approvals, files, usage, and lifecycle events. Add `--json` for NDJSON, `--reasoning` for reasoning events, or `--diagnostics` for sidecar diagnostics. Native harness `drive` is line-oriented and acknowledged; native harness `passthrough` is unsupported because no terminal stream exists. PTY attach behavior is unchanged. +### Local operation during a Relaycast outage + +```bash +agent-relay node up --local-only +agent-relay node status +agent-relay node agent spawn claude --runtime pty +agent-relay node agent attach --mode view +``` + +`--local-only` deliberately starts a **DEGRADED** broker without waiting for +Relaycast. Startup, `/health`, `/api/status`, and `node status` report that mode. +The authenticated `/api/session` exposes `operation_mode: "local_only"` and +`degraded: true`; its existing `mode` still describes persistence. +The standalone broker accepts `init --local-only --persist`; SDK callers can +set `AGENT_RELAY_LOCAL_ONLY=1` and enable persistence. The API must bind to a +loopback IP, and the mode supports a single workspace key. + +Local spawn, terminal view/input, and the durable automatic delivery queue +remain available. `POST /api/send` accepts only a worker currently running on +this broker; channel, cross-workspace, and remote destinations are rejected. +It reports `delivery_status: "queued_local"`, `local: true`, and +`relaycast_published: false`. Acceptance means the work was saved, not that an +agent has read it. Pending work survives restart and waits for an absent local +recipient to respawn without exhausting retries. A restarted recipient gets a +fresh transport retry budget. Explicit release and exhausted transport failures +while the recipient is present still use the dead-letter lifecycle. Manual-flush mode is +unavailable and returns `capability_disabled` explicitly. + +Fleet routing, worker presence, remote terminal attachment, node capability +providers, and injected Relaycast messaging tools are disabled. Local agents +receive a degraded-mode notice. Their model provider and any tools they +configure themselves still have their own connectivity requirements. + +When a workspace key is configured through the normal workspace selection, +the broker retries an independent audit connection in the background. Audit +endpoints require HTTPS, with HTTP allowed only for literal loopback IPs used +in local development. Audit records never follow HTTP redirects. Queued +local delivery records are persisted in `state-.local-outbox.json` beside +broker state, then reconciled as `local.delivery.queued` events under a separate +broker audit identity when Relaycast responds. These events contain the +original delivery/event IDs, sender, recipient, body, and queue timestamp. +They record local acceptance, not model completion. They are **audit replay**, +not re-sent messages: replay must never execute the work twice or address a +local worker name on another machine. Reconciliation is at least once; consumers +can deduplicate by `event_id` after an ambiguous response or crash. + +`node status` reports the reconciliation backlog and the last connection +result. Without a workspace key, local work still runs and records remain on +disk with no upload destination; no workspace is created. A nonempty unscoped +backlog cannot acquire a destination on restart. Keep its state directory and +continue without a key, or select a different state directory for new work +with a configured destination. The outbox +is bounded to 10,000 records / 32 MiB and rejects new sends when full. A digest +pins a configured backlog to its original workspace key and Relaycast base URL; +restore that configuration to drain it before rotating keys or changing the +destination. Corrupt outboxes cause an explicit startup failure rather than +being discarded. Preserve the state directory until reconciliation completes. + +Recovery never silently enables fleet capabilities. Stop the broker and start +normally (without the flag or environment opt-in) to enable them; a normal +restart drains retained audit backlog only when its configured destination and +digest match; an unscoped backlog remains on disk. Restart local workers as needed +to give them Relaycast messaging tools and registered identities. + ### Workspace binding and recovery `agent-relay up` and `agent-relay node up` resolve the workspace through one diff --git a/packages/cli/src/cli/commands/core.test.ts b/packages/cli/src/cli/commands/core.test.ts index 23744b887a..984153d812 100644 --- a/packages/cli/src/cli/commands/core.test.ts +++ b/packages/cli/src/cli/commands/core.test.ts @@ -2233,6 +2233,47 @@ describe('registerCoreCommands', () => { expect(sdkStatusClient.disconnect).toHaveBeenCalled(); }); + it('status visibly reports local-only degradation and reconciliation backlog', async () => { + sdkStatusClient.getStatus.mockResolvedValue({ + agent_count: 1, + pending_delivery_count: 2, + mode: 'local_only', + degraded: { reconciliation: { configured: true, connected: false, pending_records: 3 } }, + } as Awaited>); + const fs = createFsMock({ '/tmp/project/.agentworkforce/relay/connection.json': connectionFile(4242) }); + const { program, deps } = createHarness({ fs }); + await runCommand(program, ['status']); + expect(deps.log).toHaveBeenCalledWith('Status: DEGRADED (LOCAL ONLY)'); + expect(deps.log).not.toHaveBeenCalledWith('Status: RUNNING'); + expect(deps.warn).toHaveBeenCalledWith( + expect.stringContaining('remote delivery and remote attachment: DISABLED') + ); + expect(deps.log).toHaveBeenCalledWith('Reconciliation: disconnected; pending records: 3'); + }); + + it('status retains visible degradation when the runtime status request fails', async () => { + sdkStatusClient.getStatus.mockRejectedValue(new Error('runtime busy')); + const connection = { ...JSON.parse(connectionFile(4242)), operation_mode: 'local_only' }; + const fs = createFsMock({ + '/tmp/project/.agentworkforce/relay/connection.json': JSON.stringify(connection), + }); + const { program, deps } = createHarness({ fs }); + await runCommand(program, ['status']); + expect(deps.log).toHaveBeenCalledWith('Status: DEGRADED (LOCAL ONLY)'); + expect(deps.log).not.toHaveBeenCalledWith('Status: RUNNING'); + }); + + it('up --local-only skips fleet readiness and capability providers while spawning local agents', async () => { + const { program, deps, relay } = createHarness({ + teamsConfig: { agents: [{ name: 'local-worker', cli: 'cat' }] }, + }); + await runCommand(program, ['up', '--local-only', '--spawn']); + expect(deps.env.AGENT_RELAY_LOCAL_ONLY).toBe('1'); + expect(deps.warn).toHaveBeenCalledWith(expect.stringContaining('DEGRADED')); + expect(relay.spawn).toHaveBeenCalledWith(expect.objectContaining({ name: 'local-worker' })); + expect(deps.error).not.toHaveBeenCalledWith(expect.stringContaining('Refusing to auto-spawn')); + }); + it('status cleans stale connection metadata when broker is not running', async () => { const connectionPath = '/tmp/project/.agentworkforce/relay/connection.json'; const fs = createFsMock({ [connectionPath]: connectionFile(9999) }); diff --git a/packages/cli/src/cli/commands/core.ts b/packages/cli/src/cli/commands/core.ts index 515fcdef11..5b0fa7de46 100644 --- a/packages/cli/src/cli/commands/core.ts +++ b/packages/cli/src/cli/commands/core.ts @@ -288,6 +288,7 @@ export function withDefaults(overrides: Partial = {}): CoreDep /** Options accepted by the `up` command action (shared by `local`/`node`). */ export interface UpCommandOptions { + localOnly?: boolean; spawn?: boolean; background?: boolean; /** Internal marker set only on the detached child re-exec. */ @@ -309,6 +310,7 @@ export interface UpCommandOptions { export function addUpCommandOptions(command: Command): Command { return ( command + .option('--local-only', 'DEGRADED: local agents and durable delivery only; disable fleet routing') .option('--spawn', 'Force spawn all agents from teams.json') .option('--no-spawn', 'Do not auto-spawn agents (just start broker)') .option('--background', 'Run broker in the background (detached)') diff --git a/packages/cli/src/cli/commands/node.ts b/packages/cli/src/cli/commands/node.ts index 0fa0008321..423ec5dfbb 100644 --- a/packages/cli/src/cli/commands/node.ts +++ b/packages/cli/src/cli/commands/node.ts @@ -242,6 +242,11 @@ function applyResolvedNodeSession( * delegating to the shared broker `up` flow. */ async function runNodeUp(options: UpCommandOptions, deps: NodeCommandDependencies): Promise { + if (options.localOnly || deps.core.env.AGENT_RELAY_LOCAL_ONLY === '1') { + await runUpCommand({ ...options, localOnly: true, discoverConfig: false }, deps.core); + return; + } + const env = deps.core.env; // Fleet nodes may be started concurrently on one machine. Let the broker // bind an ephemeral API port atomically unless the operator explicitly diff --git a/packages/cli/src/cli/commands/status.test.ts b/packages/cli/src/cli/commands/status.test.ts index 31c658acbc..d2b25fb383 100644 --- a/packages/cli/src/cli/commands/status.test.ts +++ b/packages/cli/src/cli/commands/status.test.ts @@ -31,6 +31,12 @@ describe('relay status (composite)', () => { expect(lines).toContainEqual(expect.stringContaining('logged in (https://cloud.example)')); }); + it('reports degraded health without calling the broker healthy', async () => { + const { program, log } = harness({ probe: vi.fn(async () => 'degraded' as const) }); + await program.parseAsync(['status'], { from: 'user' }); + expect(log).toHaveBeenCalledWith('Local broker: DEGRADED (LOCAL ONLY) (http://localhost:4123)'); + }); + it('reports stopped broker and not-logged-in', async () => { const { program, log } = harness({ getBrokerConnection: () => null, diff --git a/packages/cli/src/cli/commands/status.ts b/packages/cli/src/cli/commands/status.ts index b7cfbdd421..b54bbc7358 100644 --- a/packages/cli/src/cli/commands/status.ts +++ b/packages/cli/src/cli/commands/status.ts @@ -11,7 +11,7 @@ type ExitFn = (code: number) => never; export interface StatusDependencies { getProjectRoot: () => string; getBrokerConnection: () => { url: string; apiKey?: string } | null; - probe: (url: string) => Promise; + probe: (url: string) => Promise; getCloudAuth: () => Promise<{ apiUrl: string } | null>; log: (...args: unknown[]) => void; error: (...args: unknown[]) => void; @@ -28,7 +28,9 @@ function withDefaults(overrides: Partial = {}): StatusDepend probe: async (url: string) => { try { const res = await fetch(new URL('/health', url)); - return res.ok; + if (!res.ok) return false; + const health = (await res.json()) as { status?: string }; + return health.status === 'degraded' ? 'degraded' : true; } catch { return false; } @@ -59,8 +61,9 @@ export function registerStatusCommand(program: Command, overrides: Partial capacity for this set. A pre-set @@ -1976,7 +1985,11 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies): const joinedWorkspaceId = relay.workspaceId ?? 'unknown'; // The multi-workspace session always joins a configured membership; it // never mints a new workspace the way an unresolved single key does. - if (workspaceSelection || joinsMultiWorkspaceSession) { + if (localOnly) { + deps.log( + 'Workspace: local only; configured credentials are used only for delivery-record reconciliation' + ); + } else if (workspaceSelection || joinsMultiWorkspaceSession) { deps.log(`Workspace: joined ${joinedWorkspaceId}`); } else { deps.log(`Workspace: created new workspace ${joinedWorkspaceId}`); @@ -2010,7 +2023,9 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies): } vlog(deps, options.verbose, 'Starting node capability providers (if any)...'); - nodeProviders = await startNodeCapabilityProviders(paths, relay, options, deps, nodePlan); + nodeProviders = localOnly + ? undefined + : await startNodeCapabilityProviders(paths, relay, options, deps, nodePlan); // When Reflex is enabled, periodically sync + push local session history to // relayhistory-cloud in-process via the ai-hist-native addon (no subprocess). // No-op when disabled or the addon isn't available. @@ -2023,11 +2038,9 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies): // Node delivery can't connect until the broker mints its node token, which // now happens in the background after `Broker started.`. Budget for that // mint window plus the connect so a slow mint doesn't abort auto-spawn. - const delivery = await waitForNodeDelivery( - relay, - deps, - NODE_TOKEN_WAIT_MS + NODE_DELIVERY_READY_TIMEOUT_MS - ); + const delivery = localOnly + ? { ready: true, status: null } + : await waitForNodeDelivery(relay, deps, NODE_TOKEN_WAIT_MS + NODE_DELIVERY_READY_TIMEOUT_MS); if (!delivery.ready) { deps.error('Refusing to auto-spawn agents because broker node delivery is not connected.'); deps.error(`Node delivery: ${formatNodeDeliveryStatus(delivery.status)}`); @@ -2253,8 +2266,21 @@ export async function runStatusCommand( return; } - deps.log('Status: RUNNING'); - deps.log('Mode: broker (stdio)'); + const statusDetails = + readiness.statusDetails ?? (waitMs > 0 ? null : await readBrokerStatusDetails(readiness.conn)); + const localOnly = + statusDetails?.status.mode === 'local_only' || readiness.conn.operation_mode === 'local_only'; + deps.log(localOnly ? 'Status: DEGRADED (LOCAL ONLY)' : 'Status: RUNNING'); + deps.log(localOnly ? 'Mode: local only' : 'Mode: broker (stdio)'); + if (localOnly) { + deps.warn('Cross-machine routing, worker presence, remote delivery and remote attachment: DISABLED'); + const reconciliation = statusDetails?.status.degraded?.reconciliation; + if (reconciliation) { + deps.log( + `Reconciliation: ${reconciliation.connected ? 'connected (audit only)' : reconciliation.configured ? 'disconnected' : 'not configured'}; pending records: ${reconciliation.pending_records}` + ); + } + } deps.log(`PID: ${readiness.conn.pid}`); deps.log(`Project: ${paths.projectRoot}`); const source = workspaceBindingSource(readiness.conn.workspace_source); @@ -2264,9 +2290,7 @@ export async function runStatusCommand( : 'Workspace source: unknown (startup provenance was not recorded)' ); - // Query the running broker for additional status info - const statusDetails = - readiness.statusDetails ?? (waitMs > 0 ? null : await readBrokerStatusDetails(readiness.conn)); + // Additional runtime details use the same bounded snapshot as the mode above. if (!statusDetails || statusDetails.session === null) { deps.warn('Broker API details unavailable (request failed or exceeded the 2s limit).'); } diff --git a/packages/harness-driver/src/protocol.ts b/packages/harness-driver/src/protocol.ts index 0e82a7ea85..51b4d7fea4 100644 --- a/packages/harness-driver/src/protocol.ts +++ b/packages/harness-driver/src/protocol.ts @@ -272,6 +272,12 @@ export interface FleetInventoryAgent { } export interface BrokerStatus { + mode?: 'normal' | 'local_only'; + status?: 'running' | 'degraded'; + degraded?: { + capabilities: Record; + reconciliation: { configured: boolean; connected: boolean; pending_records: number }; + } | null; agent_count: number; agents: ListAgent[]; pending_delivery_count: number; diff --git a/tests/relayflows/cases/broker-local-only/case.json b/tests/relayflows/cases/broker-local-only/case.json new file mode 100644 index 0000000000..0db3ae6b1c --- /dev/null +++ b/tests/relayflows/cases/broker-local-only/case.json @@ -0,0 +1,13 @@ +{ + "version": 1, + "id": "broker-local-only", + "kind": "feature", + "title": "Local agents remain available during Relaycast outages with visible degradation and durable reconciliation", + "requirements": ["broker-linux-x64"], + "runner": { "command": ["node", "tests/relayflows/cases/broker-local-only/run.mjs"] }, + "timeoutSeconds": 900, + "expected": { + "base": { "outcome": "absent", "signature": "relaycast_outage_blocks_local_runtime" }, + "head": { "outcome": "fixed", "signature": "visible_local_runtime_and_reconciled_delivery" } + } +} diff --git a/tests/relayflows/cases/broker-local-only/run.mjs b/tests/relayflows/cases/broker-local-only/run.mjs new file mode 100644 index 0000000000..06aec68199 --- /dev/null +++ b/tests/relayflows/cases/broker-local-only/run.mjs @@ -0,0 +1,437 @@ +import assert from 'node:assert/strict'; +import { spawn, execFileSync } from 'node:child_process'; +import { once } from 'node:events'; +import http from 'node:http'; +import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; + +const caseId = 'broker-local-only'; +const binary = process.env.RELAY_PR_PROOF_BROKER_BINARY; +assert(binary, 'Set RELAY_PR_PROOF_BROKER_BINARY to the compiled broker'); +const arm = process.env.RELAY_PR_PROOF_ARM ?? 'head'; +const target = process.env.RELAY_PR_PROOF_TARGET_DIR; +if (target) { + const sha = execFileSync('git', ['-C', target, 'rev-parse', 'HEAD'], { encoding: 'utf8' }).trim(); + assert.equal(sha, process.env[arm === 'base' ? 'RELAY_PR_PROOF_BASE_SHA' : 'RELAY_PR_PROOF_HEAD_SHA']); + const harness = path.resolve(process.env.RELAY_PR_PROOF_HARNESS_DIR); + const relative = path.relative(harness, fileURLToPath(import.meta.url)); + assert(!relative.startsWith('..') && !path.isAbsolute(relative), 'Run the exact-head harness'); +} +const dir = await mkdtemp(path.join(tmpdir(), 'broker-local-only-')); +const stateDir = path.join(dir, 'state'); +const homeDir = path.join(dir, 'home'); +await mkdir(homeDir); +const observations = { unavailable: 0, registrations: [], events: [], unexpected: [] }; +let online = false; +const server = http.createServer(async (req, res) => { + const chunks = []; + for await (const chunk of req) chunks.push(chunk); + const body = chunks.length ? JSON.parse(Buffer.concat(chunks).toString()) : {}; + const reply = (status, data) => { + res.writeHead(status, { 'content-type': 'application/json' }); + res.end(JSON.stringify(data)); + }; + if (!online) { + observations.unavailable++; + if ( + req.method !== 'POST' || + (req.url !== '/v1/agents' && !/^\/v1\/agents\/[^/]+\/events$/.test(req.url)) + ) { + observations.unexpected.push(`${req.method} ${req.url}`); + } + reply(503, { ok: false, error: { code: 'database_overloaded', message: 'deterministic test outage' } }); + } else if (req.method === 'POST' && req.url === '/v1/agents') { + observations.registrations.push(body.name); + reply(200, { + ok: true, + data: { + id: 'audit-agent', + workspace_id: 'test-workspace', + name: body.name, + token: 'at_test_local_only', + status: 'active', + created_at: '2026-09-09T00:00:00Z', + }, + }); + } else if (req.method === 'POST' && /^\/v1\/agents\/[^/]+\/events$/.test(req.url)) { + observations.events.push(body); + reply(200, { + ok: true, + data: { + id: `audit-${observations.events.length}`, + agent_id: 'audit-agent', + type: body.type, + payload: body.payload, + created_at: '2026-09-09T00:00:00Z', + }, + }); + } else { + observations.unexpected.push(`${req.method} ${req.url}`); + reply(404, { ok: false, error: { code: 'not_found', message: 'unexpected remote operation' } }); + } +}); +server.listen(0, '127.0.0.1'); +await once(server, 'listening'); +const base = `http://127.0.0.1:${server.address().port}`; +// An allowlist avoids inheriting live workspace, fleet, observer, or harness credentials. +const env = { + PATH: process.env.PATH, + HOME: homeDir, + TMPDIR: dir, + AGENT_RELAY_WORKSPACE_KEY: 'rk_live_test_local_only', + RELAYCAST_BASE_URL: base, + RELAY_BASE_URL: base, + RELAYCAST_WS_URL: base, + RELAY_BROKER_API_KEY: 'br_test_local_only', + AGENT_RELAY_TELEMETRY_DISABLED: '1', + DO_NOT_TRACK: '1', + AGENT_RELAY_NO_DEBUG_FILES: '1', + AGENT_RELAY_ORIGIN_ACTOR: 'local-proof-parent', + RELAY_AGENT_NAME: 'local-proof-parent', + RELAY_AGENT_TYPE: 'agent', + RELAY_STRICT_AGENT_NAME: '1', +}; +let child; +let logs = ''; +let url; +let exited; +async function start(localOnly, apiBind = '127.0.0.1') { + logs = ''; + url = undefined; + child = spawn( + path.resolve(binary), + [ + 'init', + '--state-dir', + stateDir, + '--instance-name', + 'outage-test', + '--api-bind', + apiBind, + ...(localOnly ? ['--local-only'] : []), + ], + { cwd: dir, env, stdio: ['pipe', 'pipe', 'pipe'] } + ); + exited = once(child, 'exit'); + child.stdout.on('data', (chunk) => { + logs += chunk.toString(); + url ??= logs.match(/API listening on (http:\/\/[^\s]+)/)?.[1]; + }); + child.stderr.on('data', (chunk) => { + logs += chunk.toString(); + }); + child.on('error', () => {}); +} +async function poll(probe, description) { + const deadline = Date.now() + 60_000; + while (Date.now() < deadline) { + const result = await probe(); + if (result) return result; + if (child.exitCode !== null) throw new Error(`Broker exited before ${description}`); + await new Promise((resolve) => setTimeout(resolve, 50)); + } + throw new Error(`Missing published signal: ${description}`); +} +async function request(route, body, method = body ? 'POST' : 'GET') { + const response = await fetch(`${url}${route}`, { + method, + headers: { 'x-api-key': env.RELAY_BROKER_API_KEY, 'content-type': 'application/json' }, + ...(body ? { body: JSON.stringify(body) } : {}), + signal: AbortSignal.timeout(10_000), + }); + const data = await response.json(); + return { response, data }; +} +async function ready(mode = 'local_only') { + await poll(async () => { + if (!url) return false; + try { + const { response, data } = await request('/api/status'); + return response.ok && data.mode === mode; + } catch { + return false; + } + }, `${mode} runtime readiness`); +} +async function stop() { + if (child && child.exitCode === null && child.signalCode === null) { + const stopping = child; + stopping.stdin.end(); + stopping.kill('SIGTERM'); + const timer = setTimeout(() => stopping.kill('SIGKILL'), 5_000); + try { + await exited; + } finally { + clearTimeout(timer); + } + } +} + +let result; +try { + if (arm === 'base') { + await start(false); + const code = await Promise.race([ + exited, + new Promise((_, reject) => + setTimeout(() => reject(new Error('Base did not surface startup failure')), 60_000).unref() + ), + ]); + assert.notEqual(code[0], 0); + assert(observations.unavailable > 0); + assert.match(logs, /failed to initialize relaycast session|failed registering agent/); + result = { outcome: 'absent', signature: 'relaycast_outage_blocks_local_runtime' }; + } else { + await start(true, ' [::1] '); + await ready(); + assert.equal(new URL(url).hostname, '[::1]', 'IPv6 API discovery URL is bracketed'); + const connection = JSON.parse(await readFile(path.join(stateDir, 'connection.json'), 'utf8')); + assert.equal(connection.url, url, 'Persisted discovery URL matches the listening IPv6 API'); + assert.match(logs, /DEGRADED.*LOCAL ONLY/); + const lease = await request('/api/session/renew', {}); + assert(lease.response.ok); + assert.equal(lease.data.persist, true, '--state-dir enables effective persistence'); + assert.equal(lease.data.expires_in_secs, 0, 'Durable local work has no owner lease expiry'); + let health = (await request('/health')).data; + assert.equal(health.status, 'degraded'); + assert.equal(health.relaycastConnected, false); + const session = (await request('/api/session')).data; + assert.equal(session.operation_mode, 'local_only'); + assert.equal(session.workspace_key, null); + assert.equal(session.node_token, null); + const isolated = await request('/api/spawn', { + name: 'local-env-worker', + cli: 'sh', + cwd: dir, + args: [ + '-c', + 'for key in AGENT_RELAY_ORIGIN_ACTOR RELAY_AGENT_NAME RELAY_AGENT_TYPE RELAY_STRICT_AGENT_NAME AGENT_RELAY_WORKSPACE_KEY RELAY_WORKSPACE_KEY RELAY_API_KEY RELAY_AGENT_TOKEN RELAY_NODE_TOKEN; do if printenv "$key" >/dev/null; then exit 9; fi; done; test "$AGENT_RELAY_TELEMETRY_DISABLED" = 1 && test "$DO_NOT_TRACK" = 1 && test "$AGENT_RELAY_NO_DEBUG_FILES" = 1 || exit 10; printf clean > local-env-proof; read -r line', + ], + channels: [], + }); + assert(isolated.response.ok); + await poll(async () => { + try { + return (await readFile(path.join(dir, 'local-env-proof'), 'utf8')) === 'clean'; + } catch (error) { + if (error.code === 'ENOENT') return false; + throw error; + } + }, 'local child has no Relaycast credentials or identity variables'); + const exiting = await request('/api/spawn', { + name: 'local-exit-worker', + cli: 'sh', + cwd: dir, + args: ['-c', 'read -r line'], + channels: [], + }); + assert(exiting.response.ok); + assert((await request('/api/input/local-exit-worker', { data: 'EXIT_LOCAL_WORKER\r' })).response.ok); + await poll(async () => { + const { data } = await request('/api/spawned'); + return !data.agents.some((agent) => agent.name === 'local-exit-worker'); + }, 'exited local worker reaped'); + const respawned = await request('/api/spawn', { + name: 'local-exit-worker', + cli: 'cat', + cwd: dir, + args: [], + channels: [], + }); + assert(respawned.response.ok, 'Local exit must not reserve a name for remote identity cleanup'); + const deletion = await request( + '/api/spawned/local-exit-worker', + { expected_generation: respawned.data.generation, delete_identity: true }, + 'DELETE' + ); + assert(!deletion.response.ok, 'Local generations cannot claim ownership of a remote identity'); + assert.match(JSON.stringify(deletion.data), /identity was not created by this worker generation/); + assert( + (await request('/api/spawned')).data.agents.some((agent) => agent.name === 'local-exit-worker'), + 'Rejected remote identity deletion preserves the local worker' + ); + const spawned = await request('/api/spawn', { + name: 'local-worker', + cli: 'cat', + cwd: dir, + args: [], + channels: [], + }); + assert(spawned.response.ok, 'Local spawn succeeds during outage'); + assert.match(spawned.data.warning, /DEGRADED.*LOCAL ONLY/); + let state = JSON.parse(await readFile(path.join(stateDir, 'state-outage-test.json'), 'utf8')); + assert.equal(state.agents['local-worker'].initial_task ?? null, null, 'Taskless spawn stays idle'); + for (const [name, task, exitAfterTask] of [ + ['empty-task-worker', ' ', false], + ['task-worker', 'EXPLICIT_INITIAL_TASK', true], + ]) { + const response = await request('/api/spawn', { + name, + cli: 'cat', + cwd: dir, + args: [], + channels: [], + task, + exit_after_task: exitAfterTask, + }); + assert(response.response.ok); + state = JSON.parse(await readFile(path.join(stateDir, 'state-outage-test.json'), 'utf8')); + if (exitAfterTask) { + assert.match(state.agents[name].initial_task, /DEGRADED.*LOCAL ONLY/); + assert.match(state.agents[name].initial_task, /EXPLICIT_INITIAL_TASK/); + assert.match(state.agents[name].initial_task, /Post-task exit/); + } else { + assert.equal(state.agents[name].initial_task ?? null, null, 'Blank task stays idle'); + } + } + const unknown = await request('/api/send', { to: 'remote-worker', text: 'must not route' }); + assert(!unknown.response.ok, 'Remote destination is explicitly rejected'); + const sent = await request('/api/send', { to: 'local-worker', text: 'LOCAL_WORK_PROOF', mode: 'steer' }); + assert(sent.response.ok, 'Local work accepted'); + assert.equal(sent.data.delivery_status, 'queued_local'); + assert.equal(sent.data.relaycast_published, false); + await poll(async () => { + const snapshot = await request('/api/spawned/local-worker/snapshot'); + return snapshot.response.ok && snapshot.data.screen?.includes('LOCAL_WORK_PROOF'); + }, 'local work rendered in attached PTY'); + const input = await request('/api/input/local-worker', { data: 'LOCAL_ATTACH_PROOF\r' }); + assert(input.response.ok, 'Local terminal input works during outage'); + await poll( + async () => + (await request('/api/spawned/local-worker/snapshot')).data.screen?.includes('LOCAL_ATTACH_PROOF'), + 'local attachment output' + ); + await poll(async () => observations.unavailable > 0, 'failed background registration'); + const before = (await request('/api/status')).data; + assert.equal(before.degraded.reconciliation.pending_records, 1); + assert.equal(before.auth.authenticated, false); + assert.equal(before.degraded.capabilities.remote_delivery, false); + // Reopen the same durable state while still disconnected. This checks the + // record survived a process boundary, not just an in-memory retry. + await stop(); + // Model a crash snapshot with exhausted local handoffs. Audit reconciliation + // alone cannot prove that the separate worker-delivery queue survives. + const retainedDelivery = { + worker_name: 'recovered-worker', + delivery: { + delivery_id: 'del_local_retention_proof', + event_id: 'local_retention_proof', + workspace_id: 'local', + from: 'sender', + target: 'recovered-worker', + body: 'LOCAL_REPLAY_PROOF', + injection_mode: 'steer', + }, + attempts: 100, + failed_attempts: 100, + queued_at_ms: Date.now(), + last_error: 'prior transport failures', + }; + await writeFile(path.join(stateDir, 'pending-outage-test.json'), JSON.stringify([retainedDelivery])); + await start(true); + await ready(); + await poll(async () => { + const status = (await request('/api/status')).data; + assert.equal(status.dead_letter_count, 0, 'Absent local work must not be dead-lettered'); + return status.pending_deliveries.some( + (delivery) => + delivery.delivery_id === retainedDelivery.delivery.delivery_id && + delivery.last_error === 'waiting for local recipient to reconnect' + ); + }, 'exhausted local delivery retained while recipient is absent'); + // Persist the retained queue across another real broker restart before + // introducing the recipient, then observe its original payload in the PTY. + await stop(); + await start(true); + await ready(); + assert.equal((await request('/api/status')).data.dead_letter_count, 0); + assert( + ( + await request('/api/spawn', { + name: 'recovered-worker', + cli: 'cat', + cwd: dir, + args: [], + channels: [], + }) + ).response.ok + ); + await poll( + async () => + (await request('/api/spawned/recovered-worker/snapshot')).data.screen?.includes( + retainedDelivery.delivery.body + ), + 'retained local work rendered after recipient respawn' + ); + assert.equal((await request('/api/status')).data.degraded.reconciliation.pending_records, 1); + online = true; + await poll( + async () => (await request('/api/status')).data.degraded.reconciliation.pending_records === 0, + 'reconciliation acknowledgement' + ); + assert.equal(observations.events.length, 1); + assert.equal(observations.events[0].type, 'local.delivery.queued'); + assert.equal(observations.events[0].payload.event_id, sent.data.event_id); + assert.equal(observations.events[0].payload.body, 'LOCAL_WORK_PROOF'); + assert(observations.registrations.every((name) => name.startsWith('outage-test-local-'))); + assert.deepEqual(observations.unexpected, [], 'No presence, fleet, message, or remote attach traffic'); + health = (await request('/health')).data; + assert.equal(health.status, 'degraded'); + assert.equal(health.nodeConnected, false); + assert.equal(health.relaycastConnected, false); + const persisted = JSON.parse( + await readFile(path.join(stateDir, 'state-outage-test.local-outbox.json'), 'utf8') + ); + assert.equal(persisted.records.length, 0); + assert(!JSON.stringify(persisted).includes(env.AGENT_RELAY_WORKSPACE_KEY)); + // Recovery may use a different workspace, or follow a local-only session + // with no key at all. Neither backlog may be uploaded or block normal mode. + for (const scope of [null, 'different-destination']) { + await stop(); + const backlog = JSON.stringify({ + scope, + records: [{ event_id: 'local_foreign_scope', body: 'PRIVATE_RETAINED_AUDIT' }], + }); + const outboxPath = path.join(stateDir, 'state-outage-test.local-outbox.json'); + await writeFile(outboxPath, backlog); + await start(false); + await ready('normal'); + assert.equal(await readFile(outboxPath, 'utf8'), backlog, 'Unmatched backlog stays byte-identical'); + assert.match(logs, /local audit backlog retained without reconciliation/); + assert( + !observations.events.some((event) => event.payload?.event_id === 'local_foreign_scope'), + 'Unmatched local audit work must not be published' + ); + } + result = { outcome: 'fixed', signature: 'visible_local_runtime_and_reconciled_delivery' }; + } + const report = { + version: 1, + caseId, + arm, + ...result, + details: + 'Published startup/status, local spawn/send, destination rejection, exhausted local queue retention and respawn replay across restarts, audit replay, and normal recovery with unmatched backlogs checked; no elapsed-time assertions.', + }; + if (process.env.RELAY_PR_PROOF_RESULT_PATH) { + await mkdir(path.dirname(process.env.RELAY_PR_PROOF_RESULT_PATH), { recursive: true }); + await writeFile(process.env.RELAY_PR_PROOF_RESULT_PATH, JSON.stringify(report) + '\n'); + } + console.log(JSON.stringify(report)); +} catch (error) { + // Logs use only isolated test credentials; still redact them on failure. + console.error( + logs + .replaceAll(env.AGENT_RELAY_WORKSPACE_KEY, '[REDACTED]') + .replaceAll(env.RELAY_BROKER_API_KEY, '[REDACTED]') + .slice(-5000) + ); + throw error; +} finally { + await stop(); + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); + await rm(dir, { recursive: true, force: true }); +}