From 3d9871a5882c5339732b5f0b150b875b3dfc9e56 Mon Sep 17 00:00:00 2001 From: Proactive Runtime Bot Date: Tue, 30 Jun 2026 09:18:37 -0700 Subject: [PATCH 1/4] fix(broker): workspace-scope auto node id so up injection works MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The broker injects messages into agents solely over the node-control /v1/node/ws channel (node-only delivery, #1201), which requires a node token minted via create_node. The auto-derived node id was sha256(machine_seed, cwd) — workspace-independent — but the engine scopes nodes globally. So a project directory re-pointed at a new workspace (e.g. `agent-relay up` minting a fresh workspace) reused a node id already owned by the old workspace; create_node's INSERT hit a unique-constraint violation (surfaced as a misclassified HTTP 500), the retry wrapper collapsed it to "Max retries exceeded", no token was minted, and every spawned agent was silently undeliverable while /health showed relaycastConnected:true. Root fix: - derive_node_id now folds workspace_id into the hash, so each (machine, cwd, workspace) gets a distinct node. Pinned RELAY_NODE_TOKEN (operator-enrolled / fleet) path is unchanged. Hardening so this can never be silent again: - create_node mint failures log the real HTTP status + body instead of "Max retries exceeded"; non-retryable 4xx are not retried. - `agent-relay up` refuses to auto-spawn when node delivery is down (exit 1 + guidance) instead of producing a dead session. - broker /health and /api/status expose nodeConnected/nodeDelivery; `status` and `doctor` report it. relaycastConnected alone was misleading. Tests: cargo test -p agent-relay-broker --lib (794 passed); CLI vitest + tsc + build:cli green. Regression test asserts same (seed,cwd) yields distinct node ids per workspace, stable within one. Co-Authored-By: Claude Opus 4.8 --- .../active/traj_lbmykojlyhtm/trajectory.json | 46 ++ .../2026-06/traj_eowv9c937zq9/summary.md | 14 + .../2026-06/traj_eowv9c937zq9/trajectory.json | 25 ++ .../2026-06/traj_wwmbn0x18dso/summary.md | 31 ++ .../2026-06/traj_wwmbn0x18dso/trajectory.json | 53 +++ CHANGELOG.md | 5 +- crates/broker/src/listen_api.rs | 66 ++- crates/broker/src/node_control.rs | 423 +++++++++++++++--- crates/broker/src/runtime/api.rs | 7 + crates/broker/src/runtime/event_loop.rs | 2 + crates/broker/src/runtime/fleet.rs | 2 + crates/broker/src/runtime/init.rs | 73 ++- packages/cli/src/cli/commands/core.test.ts | 39 +- .../cli/src/cli/lib/broker-lifecycle.test.ts | 30 +- packages/cli/src/cli/lib/broker-lifecycle.ts | 73 +++ packages/cli/src/cli/lib/doctor.ts | 18 + packages/harness-driver/src/protocol.ts | 5 + 17 files changed, 816 insertions(+), 96 deletions(-) create mode 100644 .agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json create mode 100644 .agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/summary.md create mode 100644 .agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json create mode 100644 .agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md create mode 100644 .agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json diff --git a/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json b/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json new file mode 100644 index 000000000..05451eb40 --- /dev/null +++ b/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json @@ -0,0 +1,46 @@ +{ + "id": "traj_lbmykojlyhtm", + "version": 1, + "task": { + "title": "Implement workspace-scoped node delivery fix" + }, + "status": "active", + "startedAt": "2026-06-30T15:37:17.052Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-06-30T15:39:47.866Z" + } + ], + "chapters": [ + { + "id": "chap_w24sy92u02nm", + "title": "Work", + "agentName": "default", + "startedAt": "2026-06-30T15:39:47.866Z", + "events": [ + { + "ts": 1782833987867, + "type": "decision", + "content": "Workspace-scope auto-derived node IDs only: Workspace-scope auto-derived node IDs only", + "raw": { + "question": "Workspace-scope auto-derived node IDs only", + "chosen": "Workspace-scope auto-derived node IDs only", + "alternatives": [], + "reasoning": "The pinned RELAY_NODE_TOKEN path must keep using the enrolled machine seed verbatim, while create_node auto-mint needs node IDs unique across workspaces for the same cwd." + }, + "significance": "high" + } + ] + } + ], + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", + "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" + } +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/summary.md b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/summary.md new file mode 100644 index 000000000..eeee2d568 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/summary.md @@ -0,0 +1,14 @@ +# Trajectory: Probe Relay create_node HTTP status for local-up node conflict + +> **Status:** ✅ Completed +> **Confidence:** 90% +> **Started:** June 30, 2026 at 08:35 AM +> **Completed:** June 30, 2026 at 08:37 AM + +--- + +## Summary + +Captured create_node probe: current restarted local-up workspace returned HTTP 500 internal_error with failed nodes insert for reused node_id; reported status/body to lead. + +**Approach:** Standard approach diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json new file mode 100644 index 000000000..ca32f5cc5 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json @@ -0,0 +1,25 @@ +{ + "id": "traj_eowv9c937zq9", + "version": 1, + "task": { + "title": "Probe Relay create_node HTTP status for local-up node conflict" + }, + "status": "completed", + "startedAt": "2026-06-30T15:35:59.618Z", + "completedAt": "2026-06-30T15:37:06.842Z", + "agents": [], + "chapters": [], + "retrospective": { + "summary": "Captured create_node probe: current restarted local-up workspace returned HTTP 500 internal_error with failed nodes insert for reused node_id; reported status/body to lead.", + "approach": "Standard approach", + "confidence": 0.9 + }, + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", + "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" + } +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md new file mode 100644 index 000000000..34234f3f6 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md @@ -0,0 +1,31 @@ +# Trajectory: Investigate local up injection path + +> **Status:** ✅ Completed +> **Confidence:** 90% +> **Started:** June 30, 2026 at 08:26 AM +> **Completed:** June 30, 2026 at 08:28 AM + +--- + +## Summary + +Investigated local-up runtime config and PTY broker logs. Found local-up bound to workspace 197732795721383936, home cloud/workspace config lacks explicit persisted coverage for that workspace, and parent broker logs show node_token_missing followed by node_not_found binding failures for claude and codex. + +**Approach:** Standard approach + +--- + +## Key Decisions + +### Treat parent relay-hyperagent log as authoritative for local-up injection failure +- **Chose:** Treat parent relay-hyperagent log as authoritative for local-up injection failure +- **Reasoning:** Child PTY logs were empty or only had readiness warnings, while the parent broker log records node token mint failure and node binding failures for claude and codex. + +--- + +## Chapters + +### 1. Work +*Agent: default* + +- Treat parent relay-hyperagent log as authoritative for local-up injection failure: Treat parent relay-hyperagent log as authoritative for local-up injection failure diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json new file mode 100644 index 000000000..a01a0bebb --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json @@ -0,0 +1,53 @@ +{ + "id": "traj_wwmbn0x18dso", + "version": 1, + "task": { + "title": "Investigate local up injection path" + }, + "status": "completed", + "startedAt": "2026-06-30T15:26:27.723Z", + "completedAt": "2026-06-30T15:28:57.709Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-06-30T15:28:57.629Z" + } + ], + "chapters": [ + { + "id": "chap_iaz4nysvb3oe", + "title": "Work", + "agentName": "default", + "startedAt": "2026-06-30T15:28:57.629Z", + "endedAt": "2026-06-30T15:28:57.709Z", + "events": [ + { + "ts": 1782833337630, + "type": "decision", + "content": "Treat parent relay-hyperagent log as authoritative for local-up injection failure: Treat parent relay-hyperagent log as authoritative for local-up injection failure", + "raw": { + "question": "Treat parent relay-hyperagent log as authoritative for local-up injection failure", + "chosen": "Treat parent relay-hyperagent log as authoritative for local-up injection failure", + "alternatives": [], + "reasoning": "Child PTY logs were empty or only had readiness warnings, while the parent broker log records node token mint failure and node binding failures for claude and codex." + }, + "significance": "high" + } + ] + } + ], + "retrospective": { + "summary": "Investigated local-up runtime config and PTY broker logs. Found local-up bound to workspace 197732795721383936, home cloud/workspace config lacks explicit persisted coverage for that workspace, and parent broker logs show node_token_missing followed by node_not_found binding failures for claude and codex.", + "approach": "Standard approach", + "confidence": 0.9 + }, + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", + "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" + } +} \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index 8f284c45d..e8deaea31 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -28,7 +28,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `agent-relay integration subscribe` now resolves provider-native `--resource` values through relayfile before binding, so Slack channel names, GitHub repos, Linear team keys, and Telegram chats bind to matching relayfile VFS globs while explicit `/`-prefixed globs still work. - `agent-relay integration subscribe` is now idempotent and supports multiple resources/channels per provider. Each inbound webhook is scoped to its `(provider, resource)` binding (not one-per-provider), so subscribing a second Slack channel — or two sources into the same relay channel — no longer collides on the unique `(workspace, webhook name)` index or clobbers the other binding's webhook. Re-subscribing creates the replacement webhook/subscription before retiring the old one, so a transient failure can't leave you with no working binding; a failed cleanup now warns instead of being silently swallowed. The relay channel id is normalized (`#general` → `general`) consistently across the webhook, subscription filter, relayfile bind, and writeback-secret lookup, and `listBindings` now maps relayfile's `pathGlob` field so unsubscribe/replace match correctly. - `agent-relay-broker` bootstrap `node.register` no longer advertises a generic `"spawn"` capability. Because the engine does not treat bare `"spawn"` as a placement capability (only `spawn:*`), it materialized a `spawn` action pinned to whichever node bootstrapped first, which then hijacked capability-based spawn placement for the whole workspace — every `spawn` invoke was dispatched to that node, ignoring `cli`/`target_node`/least-loaded routing. The pre-sidecar descriptor now carries no capabilities; real `spawn:*`/action capabilities arrive on the sidecar's `node.register`. -- `agent-relay-broker` node id, when the broker auto-mints its node, is derived from the machine-id seed plus a hash of the working directory, so multiple brokers / workspaces on one host no longer collide on `create_node` (stable across restarts in the same directory, distinct across directories). When an explicit `RELAY_NODE_TOKEN` is supplied (operator-enrolled / fleet nodes), the pinned node id is used verbatim so `node.register` matches the token's node instead of being rejected with `node_id_mismatch`. +- `agent-relay-broker` node id, when the broker auto-mints its node, is derived from the machine-id seed plus a hash of the working directory **and the workspace id**, so the same project directory re-pointed at a different workspace (e.g. `agent-relay up` minting a fresh workspace) mints a distinct node instead of reusing a node id already owned by the old workspace. Previously that reuse made `create_node` fail the mint and silently disabled all realtime injection. Stable across restarts in the same directory+workspace; distinct across directories or workspaces. When an explicit `RELAY_NODE_TOKEN` is supplied (operator-enrolled / fleet nodes), the pinned node id is used verbatim so `node.register` matches the token's node instead of being rejected with `node_id_mismatch`. +- `agent-relay-broker` `create_node` mint failures now log the real HTTP status and response body instead of collapsing to `Max retries exceeded`, and non-retryable `4xx` responses are no longer retried — so a node-token mint failure is diagnosable from the broker log. +- `agent-relay up` refuses to auto-spawn agents when broker node delivery (`/v1/node/ws`) is not connected, exiting non-zero with guidance, instead of spawning agents that can never receive realtime injection. +- `agent-relay status` and `agent-relay doctor` report node-delivery health (node token present + `/v1/node/ws` connected); the broker `/health` and `/api/status` expose `nodeConnected`/`nodeDelivery`, so a broker that is relaycast-connected but cannot inject is no longer indistinguishable from a healthy one. - `agent-relay-broker` node-only delivery: agents spawned via the HTTP-register fallback (when node-control `agent.register` is unavailable) are now bound to the broker's node so the engine delivers to them; a failed bind surfaces a loud `registration_warning` instead of silently producing an undeliverable agent. A missing node token is now logged as a hard fault (realtime delivery disabled) rather than a quiet warning. - `agent-relay-broker` node-only delivery: the persisted node token is now scoped to the workspace (and engine base URL) it was minted for, so a token cached for one workspace/engine is no longer reused against another and rejected with HTTP 401. A node-control `/v1/node/ws` 401 now discards the stale token and re-mints a fresh one (bounded re-minting) instead of looping forever on the rejected token. - `agent-relay-broker` node token cache is now scoped per `node_id` (`node-tokens/{node_id}.json`) instead of one host-wide file, so two brokers in different working directories on one host no longer overwrite each other's cached token and re-mint on every restart. diff --git a/crates/broker/src/listen_api.rs b/crates/broker/src/listen_api.rs index 65bcc07d9..4fd032fb3 100644 --- a/crates/broker/src/listen_api.rs +++ b/crates/broker/src/listen_api.rs @@ -26,6 +26,7 @@ use uuid::Uuid; use crate::worker_request::{RequestWorkerError, DEFAULT_REQUEST_TIMEOUT}; const LISTEN_API_SEND_TIMEOUT: Duration = Duration::from_secs(30); +const HEALTH_STATUS_TIMEOUT: Duration = Duration::from_millis(100); type PtyInputSerializers = Arc>>>>; @@ -499,16 +500,73 @@ pub(crate) fn listen_api_health_payload( "wsConnections": 0, "memoryMb": 0, "relaycastConnected": startup_error_code.is_none(), + "nodeConnected": false, + "nodeDelivery": { + "tokenPresent": false, + "connected": false, + }, }) } async fn listen_api_health( axum::extract::State(state): axum::extract::State, ) -> axum::Json { - axum::Json(listen_api_health_payload( - state.default_workspace_id, - state.memberships, - )) + let mut payload = listen_api_health_payload(state.default_workspace_id, state.memberships); + if let Some(status) = fetch_status_for_health(&state.tx).await { + merge_status_into_health_payload(&mut payload, &status); + } + axum::Json(payload) +} + +async fn fetch_status_for_health(tx: &mpsc::Sender) -> Option { + let (reply_tx, reply_rx) = tokio::sync::oneshot::channel(); + tx.send(ListenApiRequest::GetStatus { reply: reply_tx }) + .await + .ok()?; + timeout(HEALTH_STATUS_TIMEOUT, reply_rx) + .await + .ok()? + .ok()? + .ok() +} + +fn merge_status_into_health_payload(payload: &mut Value, status: &Value) { + let Some(object) = payload.as_object_mut() else { + return; + }; + if let Some(agent_count) = status.get("agent_count").and_then(Value::as_u64) { + object.insert("agentCount".to_string(), json!(agent_count)); + } + if let Some(pending_count) = status.get("pending_delivery_count").and_then(Value::as_u64) { + object.insert("pendingDeliveryCount".to_string(), json!(pending_count)); + } + let token_present = status + .get("node_delivery") + .and_then(|value| value.get("token_present")) + .and_then(Value::as_bool) + .unwrap_or(false); + let connected = status + .get("node_connected") + .and_then(Value::as_bool) + .or_else(|| { + status + .get("node_delivery") + .and_then(|value| value.get("connected")) + .and_then(Value::as_bool) + }) + .unwrap_or(false); + object.insert("nodeConnected".to_string(), json!(connected)); + object.insert( + "nodeDelivery".to_string(), + json!({ + "tokenPresent": token_present, + "connected": connected, + }), + ); + object.insert( + "wsConnections".to_string(), + json!(if connected { 1 } else { 0 }), + ); } /// Authenticated endpoint that returns broker configuration, including the diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index df98ad9c2..b1ef58e60 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -7,6 +7,7 @@ use std::{ use anyhow::{Context, Result}; use futures_util::{Sink, SinkExt, StreamExt}; +use serde::Deserialize; use serde_json::Value; use sha2::{Digest, Sha256}; use tokio::sync::{mpsc, oneshot}; @@ -27,6 +28,8 @@ const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(12); const INITIAL_RECONNECT_DELAY: Duration = Duration::from_secs(1); const MAX_RECONNECT_DELAY: Duration = Duration::from_secs(30); const REGISTER_AGENT_PENDING_TTL: Duration = Duration::from_secs(300); +const RELAYCAST_DEFAULT_BASE_URL: &str = "https://cast.agentrelay.com"; +const CREATE_NODE_RETRY_BACKOFFS_MS: [u64; 3] = [200, 400, 800]; /// How many consecutive `/v1/node/ws` 401s to tolerate (each triggering a /// re-mint) before giving up and surfacing a hard error instead of looping. const MAX_UNAUTHORIZED_BEFORE_GIVING_UP: u32 = 5; @@ -50,17 +53,17 @@ pub(crate) struct FleetControlConfig { pub(crate) broker_version: String, /// Optional facility to re-mint a node token when the engine rejects the /// current one with HTTP 401 on the `/v1/node/ws` handshake. Absent in tests - /// and when no workspace `RelayCast` client is available. + /// and when no workspace key is available. pub(crate) token_minter: Option, } -/// Re-mints a fresh node token via `RelayCast::create_node` and rewrites the +/// Re-mints a fresh node token via `POST /v1/nodes` and rewrites the /// workspace-scoped cache, used to recover from a stale/rejected cached token /// (HTTP 401 on the node-control handshake) instead of looping forever on the /// same token. Mirrors the initial mint wired in `runtime::init::resolve_node_token`. #[derive(Clone)] pub(crate) struct NodeTokenMinter { - pub(crate) relay_client: relaycast::RelayCast, + pub(crate) workspace_key: String, pub(crate) workspace_id: String, pub(crate) base_url: Option, pub(crate) node_id: String, @@ -90,29 +93,19 @@ impl NodeTokenMinter { } } } - let request = relaycast::CreateNodeRequest { - node_id: Some(self.node_id.clone()), - name: self.node_name.clone(), - kind: Some("ws".to_string()), - role: Some("broker".to_string()), - delivery_adapter: None, - delivery: None, - capabilities: None, - max_agents: None, - tags: None, - version: Some(self.broker_version.clone()), - }; - match self.relay_client.create_node(request).await { - Ok(response) => { - let token = response.token.trim().to_string(); - if token.is_empty() { - tracing::warn!( - target = "relay_broker::fleet", - node_id = %self.node_id, - "re-mint create_node returned an empty token" - ); - return None; - } + let request = create_node_request(&self.node_id, &self.node_name, &self.broker_version); + match mint_node_token( + &self.workspace_key, + self.base_url.as_deref(), + request, + MintNodeTokenLogContext { + node_id: &self.node_id, + workspace_id: &self.workspace_id, + }, + ) + .await + { + Ok(token) => { if let Some(path) = self.token_path.as_deref() { if let Err(error) = persist_node_token( path, @@ -138,11 +131,12 @@ impl NodeTokenMinter { Some(token) } Err(error) => { - tracing::error!( - target = "relay_broker::fleet", - node_id = %self.node_id, - error = %error, - "failed to re-mint node token after node-control 401" + log_create_node_mint_error( + "relay_broker::fleet", + &self.node_id, + &self.workspace_id, + &error, + "failed to re-mint node token after node-control 401", ); None } @@ -150,6 +144,228 @@ impl NodeTokenMinter { } } +pub(crate) fn create_node_request( + node_id: &str, + node_name: &str, + broker_version: &str, +) -> relaycast::CreateNodeRequest { + relaycast::CreateNodeRequest { + node_id: Some(node_id.to_string()), + name: node_name.to_string(), + kind: Some("ws".to_string()), + role: Some("broker".to_string()), + delivery_adapter: None, + delivery: None, + capabilities: None, + max_agents: None, + tags: None, + version: Some(broker_version.to_string()), + } +} + +#[derive(Debug, thiserror::Error)] +pub(crate) enum CreateNodeMintError { + #[error("HTTP request failed: {0}")] + Http(#[from] reqwest::Error), + #[error("create_node returned HTTP {status} ({code}): {message}")] + Api { + status: u16, + code: String, + message: String, + body: String, + }, + #[error("create_node returned invalid HTTP {status} response: {reason}")] + InvalidResponse { + status: u16, + reason: String, + body: String, + }, +} + +impl CreateNodeMintError { + pub(crate) fn status(&self) -> Option { + match self { + Self::Api { status, .. } | Self::InvalidResponse { status, .. } => Some(*status), + Self::Http(error) => error.status().map(|status| status.as_u16()), + } + } + + pub(crate) fn code(&self) -> Option<&str> { + match self { + Self::Api { code, .. } => Some(code.as_str()), + _ => None, + } + } + + pub(crate) fn response_body(&self) -> Option<&str> { + match self { + Self::Api { body, .. } | Self::InvalidResponse { body, .. } => Some(body.as_str()), + Self::Http(_) => None, + } + } +} + +pub(crate) struct MintNodeTokenLogContext<'a> { + pub(crate) node_id: &'a str, + pub(crate) workspace_id: &'a str, +} + +#[derive(Debug, Deserialize)] +struct CreateNodeApiResponse { + ok: bool, + data: Option, + error: Option, +} + +#[derive(Debug, Deserialize)] +struct CreateNodeApiError { + code: String, + message: String, +} + +pub(crate) async fn mint_node_token( + workspace_key: &str, + base_url: Option<&str>, + request: relaycast::CreateNodeRequest, + context: MintNodeTokenLogContext<'_>, +) -> std::result::Result { + let url = format!( + "{}/v1/nodes", + normalized_relaycast_base_url(base_url).trim_end_matches('/') + ); + let client = reqwest::Client::new(); + let mut last_error: Option = None; + + for attempt in 0..=CREATE_NODE_RETRY_BACKOFFS_MS.len() { + let response = client + .post(&url) + .bearer_auth(workspace_key) + .header("X-SDK-Version", crate::util::version::broker_version()) + .header("X-Relaycast-Origin-Client", "agent-relay-broker") + .header( + "X-Relaycast-Origin-Version", + crate::util::version::broker_version(), + ) + .header( + "X-Relaycast-Origin-Actor", + crate::telemetry::BROKER_ORIGIN_ACTOR, + ) + .json(&request) + .send() + .await?; + let status = response.status().as_u16(); + let body = response.text().await?; + let envelope = serde_json::from_str::(&body).map_err(|error| { + CreateNodeMintError::InvalidResponse { + status, + reason: format!("failed to parse JSON: {error}"), + body: body.clone(), + } + })?; + + if envelope.ok { + let token = envelope + .data + .ok_or_else(|| CreateNodeMintError::InvalidResponse { + status, + reason: "response missing data field".to_string(), + body: body.clone(), + })? + .token + .trim() + .to_string(); + if token.is_empty() { + return Err(CreateNodeMintError::InvalidResponse { + status, + reason: "response returned an empty token".to_string(), + body, + }); + } + return Ok(token); + } + + let error = envelope.error.unwrap_or(CreateNodeApiError { + code: "unknown_error".to_string(), + message: "Unknown error".to_string(), + }); + let mint_error = CreateNodeMintError::Api { + status, + code: error.code, + message: error.message, + body, + }; + log_create_node_mint_error( + "relay_broker::fleet", + context.node_id, + context.workspace_id, + &mint_error, + "create_node mint attempt failed", + ); + + if !(500..=599).contains(&status) || attempt >= CREATE_NODE_RETRY_BACKOFFS_MS.len() { + return Err(mint_error); + } + last_error = Some(mint_error); + tokio::time::sleep(Duration::from_millis( + CREATE_NODE_RETRY_BACKOFFS_MS[attempt], + )) + .await; + } + + Err( + last_error.unwrap_or_else(|| CreateNodeMintError::InvalidResponse { + status: 0, + reason: "create_node retry loop exhausted without a response".to_string(), + body: String::new(), + }), + ) +} + +pub(crate) fn log_create_node_mint_error( + target: &'static str, + node_id: &str, + workspace_id: &str, + error: &CreateNodeMintError, + message: &'static str, +) { + let status = error.status(); + let code = error.code(); + let response_body = error.response_body().map(redact_relaycast_secrets); + tracing::warn!( + target = target, + node_id = %node_id, + workspace_id = %workspace_id, + http_status = ?status, + error_code = ?code, + response_body = ?response_body, + error = %error, + "{message}" + ); +} + +fn normalized_relaycast_base_url(base_url: Option<&str>) -> String { + base_url + .map(str::trim) + .filter(|value| !value.is_empty()) + .unwrap_or(RELAYCAST_DEFAULT_BASE_URL) + .to_string() +} + +fn redact_relaycast_secrets(input: &str) -> String { + let mut redacted = input.to_string(); + if let Ok(re) = + regex::Regex::new(r#"(rk_live_|at_live_|nt_live_|br_|Bearer\s+)[A-Za-z0-9._-]+"#) + { + redacted = re.replace_all(&redacted, "${1}[REDACTED]").into_owned(); + } + if let Ok(re) = regex::Regex::new(r#""token"\s*:\s*"[^"]+""#) { + redacted = re + .replace_all(&redacted, r#""token":"[REDACTED]""#) + .into_owned(); + } + redacted +} + #[derive(Debug, Clone, Default, PartialEq, Eq)] pub(crate) struct FleetLoadSnapshot { pub(crate) active_agents: u32, @@ -581,24 +797,31 @@ pub(crate) fn load_or_create_machine_seed(path: &Path) -> Result { Ok(seed) } -/// Derive a node id from a stable machine `seed` and the broker's working -/// directory `cwd`. +/// Derive a node id from a stable machine `seed`, the broker's working +/// directory `cwd`, and the relaycast `workspace_id`. /// /// The engine scopes nodes globally, so a single host that runs brokers for /// two different workspaces (each in its own per-project working directory) /// must present a distinct node id per broker — otherwise the second -/// `create_node` collides with the first. Deriving the id from -/// `(seed, cwd)` keeps it: +/// `create_node` collides with the first. The same conflict can happen when +/// `agent-relay up` creates a fresh workspace for the same project directory: +/// `(seed, cwd)` stays identical while the workspace changes, so the global +/// `node_id` would still collide. Deriving the id from +/// `(seed, cwd, workspace_id)` keeps it: /// -/// - **stable across restarts** in the same working directory, and -/// - **distinct across different working directories** on the same machine. +/// - **stable across restarts** in the same working directory and workspace, +/// - **distinct across different working directories** on the same machine, +/// and +/// - **distinct across different workspaces** in the same working directory. /// /// The result has the same `node_<32 hex>` shape as a freshly minted id. -pub(crate) fn derive_node_id(seed: &str, cwd: &str) -> String { +pub(crate) fn derive_node_id(seed: &str, cwd: &str, workspace_id: &str) -> String { let mut hasher = Sha256::new(); hasher.update(seed.as_bytes()); hasher.update([0u8]); hasher.update(cwd.as_bytes()); + hasher.update([0u8]); + hasher.update(workspace_id.as_bytes()); let digest = hasher.finalize(); let hex = digest .iter() @@ -610,14 +833,14 @@ pub(crate) fn derive_node_id(seed: &str, cwd: &str) -> String { } /// Resolve the node id for this broker: load (or create) the machine seed at -/// `seed_path`, then derive a per-working-directory node id from it. +/// `seed_path`, then derive a per-working-directory, per-workspace node id +/// from it. /// /// The working directory is read from [`std::env::current_dir`] and /// canonicalized when possible so equivalent paths resolve to the same id. If -/// the cwd cannot be read, the id is derived from the seed alone rather than -/// panicking (every broker on the host then shares one id, matching the old -/// global-machine-id behavior). -pub(crate) fn load_or_create_node_id(seed_path: &Path) -> Result { +/// the cwd cannot be read, the id is derived from the seed plus workspace +/// rather than panicking. +pub(crate) fn load_or_create_node_id(seed_path: &Path, workspace_id: &str) -> Result { let seed = load_or_create_machine_seed(seed_path)?; let cwd = std::env::current_dir() .map(|path| { @@ -627,7 +850,7 @@ pub(crate) fn load_or_create_node_id(seed_path: &Path) -> Result { .into_owned() }) .unwrap_or_default(); - Ok(derive_node_id(&seed, &cwd)) + Ok(derive_node_id(&seed, &cwd, workspace_id)) } pub(crate) fn default_node_id_path() -> Option { @@ -1300,6 +1523,7 @@ pub(crate) fn delivery_ack(agent: impl Into, up_to_seq: u64) -> BrokerTo #[cfg(test)] mod tests { + use httpmock::{Method::POST, MockServer}; use serde_json::json; use tokio::net::TcpListener; use tokio_tungstenite::accept_async; @@ -1309,27 +1533,44 @@ mod tests { #[test] fn derive_node_id_is_stable_for_same_seed_and_cwd() { - let a = derive_node_id("node_seed123", "/Users/will/Projects/relay"); - let b = derive_node_id("node_seed123", "/Users/will/Projects/relay"); - assert_eq!(a, b, "same (seed, cwd) must derive an identical node id"); + let a = derive_node_id("node_seed123", "/Users/will/Projects/relay", "workspace-a"); + let b = derive_node_id("node_seed123", "/Users/will/Projects/relay", "workspace-a"); + assert_eq!( + a, b, + "same (seed, cwd, workspace) must derive an identical node id" + ); } #[test] fn derive_node_id_differs_across_working_directories() { let seed = "node_seed123"; - let a = derive_node_id(seed, "/Users/will/Projects/workspace-a"); - let b = derive_node_id(seed, "/Users/will/Projects/workspace-b"); + let workspace_id = "workspace-a"; + let a = derive_node_id(seed, "/Users/will/Projects/workspace-a", workspace_id); + let b = derive_node_id(seed, "/Users/will/Projects/workspace-b", workspace_id); assert_ne!( a, b, "different cwd on the same machine must derive distinct node ids" ); } + #[test] + fn derive_node_id_differs_across_workspaces() { + let seed = "node_seed123"; + let cwd = "/Users/will/Projects/relay"; + let a = derive_node_id(seed, cwd, "workspace-a"); + let b = derive_node_id(seed, cwd, "workspace-b"); + assert_ne!( + a, b, + "same seed and cwd under different workspaces must derive distinct node ids" + ); + } + #[test] fn derive_node_id_differs_across_seeds() { let cwd = "/Users/will/Projects/relay"; - let a = derive_node_id("seed-a", cwd); - let b = derive_node_id("seed-b", cwd); + let workspace_id = "workspace-a"; + let a = derive_node_id("seed-a", cwd, workspace_id); + let b = derive_node_id("seed-b", cwd, workspace_id); assert_ne!( a, b, "different machine seeds must derive distinct node ids" @@ -1338,7 +1579,7 @@ mod tests { #[test] fn derive_node_id_has_node_prefix_and_32_hex_suffix() { - let id = derive_node_id("node_seed123", "/Users/will/Projects/relay"); + let id = derive_node_id("node_seed123", "/Users/will/Projects/relay", "workspace-a"); let suffix = id .strip_prefix("node_") .expect("derived node id must start with node_"); @@ -1353,6 +1594,84 @@ mod tests { ); } + #[tokio::test] + async fn mint_node_token_does_not_retry_non_retryable_4xx() { + let server = MockServer::start(); + let create_node = server.mock(|when, then| { + when.method(POST) + .path("/v1/nodes") + .header("authorization", "Bearer rk_live_test"); + then.status(403).json_body(json!({ + "ok": false, + "error": { + "code": "forbidden", + "message": "node creation is not allowed" + } + })); + }); + + let error = mint_node_token( + "rk_live_test", + Some(&server.base_url()), + create_node_request("node_abc", "local-node", "relay-broker/test"), + MintNodeTokenLogContext { + node_id: "node_abc", + workspace_id: "ws_test", + }, + ) + .await + .expect_err("403 create_node response should fail"); + + create_node.assert_hits(1); + assert_eq!(error.status(), Some(403)); + assert_eq!(error.code(), Some("forbidden")); + assert!( + error + .response_body() + .is_some_and(|body| body.contains("node creation is not allowed")), + "error should preserve response body: {error:?}" + ); + } + + #[tokio::test] + async fn mint_node_token_preserves_5xx_status_and_body_after_retries() { + let server = MockServer::start(); + let create_node = server.mock(|when, then| { + when.method(POST) + .path("/v1/nodes") + .header("authorization", "Bearer rk_live_test"); + then.status(500).json_body(json!({ + "ok": false, + "error": { + "code": "internal_error", + "message": "Failed query: insert into \"nodes\" ..." + } + })); + }); + + let error = mint_node_token( + "rk_live_test", + Some(&server.base_url()), + create_node_request("node_abc", "local-node", "relay-broker/test"), + MintNodeTokenLogContext { + node_id: "node_abc", + workspace_id: "ws_test", + }, + ) + .await + .expect_err("500 create_node response should fail"); + + create_node.assert_hits(CREATE_NODE_RETRY_BACKOFFS_MS.len() + 1); + assert_eq!(error.status(), Some(500)); + assert_eq!(error.code(), Some("internal_error")); + assert!( + error + .response_body() + .is_some_and(|body| body.contains("insert into")), + "error should preserve response body: {error:?}" + ); + } + #[test] fn delivery_book_dedups_and_tracks_cumulative_ack() { let mut book = FleetDeliveryBook::default(); diff --git a/crates/broker/src/runtime/api.rs b/crates/broker/src/runtime/api.rs index 9a9d30090..f9a591fe4 100644 --- a/crates/broker/src/runtime/api.rs +++ b/crates/broker/src/runtime/api.rs @@ -32,6 +32,8 @@ impl BrokerRuntime { let workers = &mut self.workers; let fleet_control_tx = &self.fleet_control_tx; let fleet_node_name = self.fleet_node_name.as_str(); + let node_delivery_token_present = self.node_delivery_token_present; + let node_delivery_connected = self.node_delivery_connected; let fleet_inventory = &mut self.fleet_inventory; let fleet_delivery_book = &mut self.fleet_delivery_book; let fleet_max_agents = self.fleet_max_agents; @@ -1318,6 +1320,11 @@ impl BrokerRuntime { "agents": workers.list(), "pending_delivery_count": pending.len(), "pending_deliveries": pending, + "node_connected": node_delivery_connected, + "node_delivery": { + "token_present": node_delivery_token_present, + "connected": node_delivery_connected, + }, "auth": { "authenticated": !auth_workspaces.is_empty(), "workspace_count": auth_workspaces.len(), diff --git a/crates/broker/src/runtime/event_loop.rs b/crates/broker/src/runtime/event_loop.rs index 8ff034699..58137c497 100644 --- a/crates/broker/src/runtime/event_loop.rs +++ b/crates/broker/src/runtime/event_loop.rs @@ -21,6 +21,8 @@ pub(crate) struct BrokerRuntime { /// This broker's relaycast node name, used to bind agents to the node over /// HTTP when the node-control `agent.register` path is unavailable. pub(super) fleet_node_name: String, + pub(super) node_delivery_token_present: bool, + pub(super) node_delivery_connected: bool, pub(super) fleet_event_rx: mpsc::Receiver, pub(super) fleet_control_open: bool, pub(super) fleet_delivery_book: FleetDeliveryBook, diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index 96b015a15..377b2efb3 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -314,6 +314,7 @@ impl BrokerRuntime { pub(super) async fn handle_fleet_control_event(&mut self, event: FleetControlEvent) { match event { FleetControlEvent::Connected => { + self.node_delivery_connected = true; // Node delivery is live: message delivery flows solely over // /v1/node/ws. The workspace firehose delivery path was removed, // so there is no firehose injection to suppress here. @@ -323,6 +324,7 @@ impl BrokerRuntime { ); } FleetControlEvent::Disconnected => { + self.node_delivery_connected = false; tracing::warn!( target = "relay_broker::fleet", "fleet node control disconnected" diff --git a/crates/broker/src/runtime/init.rs b/crates/broker/src/runtime/init.rs index c934f00e2..5374030c3 100644 --- a/crates/broker/src/runtime/init.rs +++ b/crates/broker/src/runtime/init.rs @@ -218,6 +218,7 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let self_names = default_workspace.self_names.clone(); let ws_control_tx = default_workspace.ws_control_tx.clone(); let relaycast_http = default_workspace.http_client.clone(); + let node_workspace_id = default_workspace.workspace_id.as_str().to_string(); let node_id = match crate::node_control::default_node_id_path() { Some(path) => { // Node-id mode — explicit vs auto: @@ -235,7 +236,7 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re let loaded = if pinned { crate::node_control::load_or_create_machine_seed(&path) } else { - crate::node_control::load_or_create_node_id(&path) + crate::node_control::load_or_create_node_id(&path, &node_workspace_id) }; loaded.unwrap_or_else(|error| { tracing::warn!(error = %error, "failed to load fleet node machine id; using ephemeral id"); @@ -258,15 +259,14 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re // against. Thread the resolved workspace id and base URL through so the // cached token is only reused when both match, and so a re-mint after a // node-control 401 rewrites the correctly-scoped cache. - let node_workspace_id = default_workspace.workspace_id.as_str().to_string(); let node_base_url = configured_base.clone(); let node_token = resolve_node_token( &node_id, &node_name, &broker_version, &node_workspace_id, + &relay_workspace_key, node_base_url.as_deref(), - relaycast_http.relay_client(), ) .await; if node_token.is_none() { @@ -299,20 +299,18 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re // forever on the rejected token. Mirrors the initial mint above. Absent when // no workspace RelayCast client is available (then a 401 surfaces a hard // error rather than recovering). - let token_minter = - relaycast_http - .relay_client() - .map(|client| crate::node_control::NodeTokenMinter { - relay_client: client.clone(), - workspace_id: node_workspace_id.clone(), - base_url: node_base_url.clone(), - node_id: node_id.clone(), - node_name: node_name.clone(), - broker_version: broker_version.clone(), - token_path: crate::node_control::default_node_token_path(&node_id), - }); + let token_minter = Some(crate::node_control::NodeTokenMinter { + workspace_key: relay_workspace_key.clone(), + workspace_id: node_workspace_id.clone(), + base_url: node_base_url.clone(), + node_id: node_id.clone(), + node_name: node_name.clone(), + broker_version: broker_version.clone(), + token_path: crate::node_control::default_node_token_path(&node_id), + }); let (fleet_control_tx, fleet_control_rx) = mpsc::channel::(256); let (fleet_event_tx, fleet_event_rx) = mpsc::channel::(256); + 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, @@ -612,6 +610,8 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re relaycast_open: true, fleet_control_tx, fleet_node_name, + node_delivery_token_present, + node_delivery_connected: false, fleet_event_rx, fleet_control_open: true, fleet_delivery_book: FleetDeliveryBook::default(), @@ -667,8 +667,8 @@ async fn resolve_node_token( node_name: &str, broker_version: &str, workspace_id: &str, + workspace_key: &str, base_url: Option<&str>, - relay_client: Option<&relaycast::RelayCast>, ) -> Option { if let Some(token) = std::env::var("RELAY_NODE_TOKEN") .ok() @@ -686,26 +686,19 @@ async fn resolve_node_token( return Some(token); } - let relay_client = relay_client?; - let request = relaycast::CreateNodeRequest { - node_id: Some(node_id.to_string()), - name: node_name.to_string(), - kind: Some("ws".to_string()), - role: Some("broker".to_string()), - delivery_adapter: None, - delivery: None, - capabilities: None, - max_agents: None, - tags: None, - version: Some(broker_version.to_string()), - }; - match relay_client.create_node(request).await { - Ok(response) => { - let token = response.token.trim().to_string(); - if token.is_empty() { - tracing::warn!(node_id = %node_id, "create_node returned an empty token"); - return None; - } + let request = crate::node_control::create_node_request(node_id, node_name, broker_version); + match crate::node_control::mint_node_token( + workspace_key, + base_url, + request, + crate::node_control::MintNodeTokenLogContext { + node_id, + workspace_id, + }, + ) + .await + { + Ok(token) => { if let Some(path) = token_path.as_deref() { if let Err(error) = crate::node_control::persist_node_token( path, @@ -721,7 +714,13 @@ async fn resolve_node_token( Some(token) } Err(error) => { - tracing::warn!(node_id = %node_id, error = %error, "failed to mint node token via create_node"); + crate::node_control::log_create_node_mint_error( + "relay_broker::fleet", + node_id, + workspace_id, + &error, + "failed to mint node token via create_node", + ); None } } diff --git a/packages/cli/src/cli/commands/core.test.ts b/packages/cli/src/cli/commands/core.test.ts index dd59a0634..4c44ba3d0 100644 --- a/packages/cli/src/cli/commands/core.test.ts +++ b/packages/cli/src/cli/commands/core.test.ts @@ -263,7 +263,14 @@ describe('registerCoreCommands', () => { }); it('up auto-spawns agents from teams config', async () => { - const relay = createRelayMock(); + const relay = createRelayMock({ + getStatus: vi.fn(async () => ({ + agent_count: 0, + pending_delivery_count: 0, + node_connected: true, + node_delivery: { token_present: true, connected: true }, + })), + }); const { program } = createHarness({ relay, teamsConfig: { @@ -285,6 +292,36 @@ describe('registerCoreCommands', () => { }); }); + it('up refuses auto-spawn when node delivery is down', async () => { + const relay = createRelayMock({ + getStatus: vi.fn(async () => ({ + agent_count: 0, + pending_delivery_count: 0, + node_connected: false, + node_delivery: { token_present: false, connected: false }, + })), + }); + const { program, deps } = createHarness({ + relay, + teamsConfig: { + team: 'platform', + autoSpawn: true, + agents: [{ name: 'WorkerA', cli: 'codex', task: 'Ship tests' }], + }, + nowImpl: vi.fn().mockReturnValueOnce(0).mockReturnValueOnce(10_000).mockReturnValue(10_000), + }); + + const exitCode = await runCommand(program, ['up']); + + expect(exitCode).toBe(1); + expect(relay.spawn).not.toHaveBeenCalled(); + expect(relay.shutdown).toHaveBeenCalled(); + expect(deps.error).toHaveBeenCalledWith( + 'Refusing to auto-spawn agents because broker node delivery is not connected.' + ); + expect(deps.error).toHaveBeenCalledWith('Node delivery: DOWN (no node token)'); + }); + it('up probes for a free API port before spawning the broker', async () => { const relay = createRelayMock(); const { program, deps } = createHarness({ relay }); diff --git a/packages/cli/src/cli/lib/broker-lifecycle.test.ts b/packages/cli/src/cli/lib/broker-lifecycle.test.ts index c07f026d7..2b3edec5d 100644 --- a/packages/cli/src/cli/lib/broker-lifecycle.test.ts +++ b/packages/cli/src/cli/lib/broker-lifecycle.test.ts @@ -1,6 +1,11 @@ import { describe, expect, it } from 'vitest'; -import { classifyBrokerStartError, classifyBrokerStartStage, describeError } from './broker-lifecycle.js'; +import { + classifyBrokerStartError, + classifyBrokerStartStage, + describeError, + readNodeDeliveryStatus, +} from './broker-lifecycle.js'; describe('describeError', () => { it('returns plain message for a bare Error', () => { @@ -87,3 +92,26 @@ describe('classifyBrokerStartStage', () => { expect(classifyBrokerStartStage(new Error('???'), '???')).toBe('startup'); }); }); + +describe('readNodeDeliveryStatus', () => { + it('reads the canonical snake_case broker status shape', () => { + expect( + readNodeDeliveryStatus({ + node_connected: true, + node_delivery: { token_present: true, connected: true }, + }) + ).toEqual({ tokenPresent: true, connected: true }); + }); + + it('defaults absent node delivery fields to false', () => { + expect(readNodeDeliveryStatus({ agent_count: 0 })).toEqual({ + tokenPresent: false, + connected: false, + }); + }); + + it('rejects non-object status values', () => { + expect(readNodeDeliveryStatus(null)).toBeNull(); + expect(readNodeDeliveryStatus('nope')).toBeNull(); + }); +}); diff --git a/packages/cli/src/cli/lib/broker-lifecycle.ts b/packages/cli/src/cli/lib/broker-lifecycle.ts index 7697a9079..d8c64fbae 100644 --- a/packages/cli/src/cli/lib/broker-lifecycle.ts +++ b/packages/cli/src/cli/lib/broker-lifecycle.ts @@ -43,6 +43,7 @@ const DEFAULT_BROKER_BASE_PORT = 3888; const CONNECTION_FILENAME = 'connection.json'; const STATUS_POLL_INTERVAL_MS = 500; const DETACHED_START_READY_TIMEOUT_MS = 10_000; +const NODE_DELIVERY_READY_TIMEOUT_MS = 10_000; export interface BrokerConnection { url: string; @@ -56,6 +57,11 @@ type BrokerStatusDetails = { session: Awaited> | null; }; +type NodeDeliveryStatus = { + tokenPresent: boolean; + connected: boolean; +}; + type BrokerReadiness = | { state: 'running'; @@ -119,6 +125,41 @@ function toErrorMessage(err: unknown): string { type ErrorWithCode = { code?: unknown }; +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null; +} + +export function readNodeDeliveryStatus(status: unknown): NodeDeliveryStatus | null { + if (!isRecord(status)) { + return null; + } + const snake = isRecord(status.node_delivery) ? status.node_delivery : null; + const tokenPresent = typeof snake?.token_present === 'boolean' ? snake.token_present : false; + const connected = + typeof status.node_connected === 'boolean' + ? status.node_connected + : typeof snake?.connected === 'boolean' + ? snake.connected + : false; + return { tokenPresent, connected }; +} + +function nodeDeliveryReady(status: unknown): boolean { + const delivery = readNodeDeliveryStatus(status); + return Boolean(delivery?.tokenPresent && delivery.connected); +} + +function formatNodeDeliveryStatus(status: unknown): string { + const delivery = readNodeDeliveryStatus(status); + if (!delivery) { + return 'unknown'; + } + if (!delivery.tokenPresent) { + return 'DOWN (no node token)'; + } + return delivery.connected ? 'CONNECTED' : 'DOWN (node websocket disconnected)'; +} + function errorCode(err: unknown): string | undefined { if (!err || typeof err !== 'object') return undefined; const code = (err as ErrorWithCode).code; @@ -697,6 +738,26 @@ async function waitForBrokerReadiness( return latest; } +async function waitForNodeDelivery( + relay: CoreRelay, + deps: CoreDependencies, + waitMs = NODE_DELIVERY_READY_TIMEOUT_MS +): Promise<{ ready: boolean; status: unknown }> { + const deadline = deps.now() + waitMs; + let latest: unknown = null; + + while (true) { + latest = await relay.getStatus(); + if (nodeDeliveryReady(latest)) { + return { ready: true, status: latest }; + } + if (waitMs <= 0 || deps.now() >= deadline) { + return { ready: false, status: latest }; + } + await deps.sleep(Math.min(STATUS_POLL_INTERVAL_MS, Math.max(0, deadline - deps.now()))); + } +} + async function shutdownUpResources(relay: CoreRelay, dataDir: string, deps: CoreDependencies): Promise { await relay.shutdown().catch(() => undefined); safeUnlink(path.join(dataDir, CONNECTION_FILENAME), deps); @@ -849,6 +910,17 @@ export async function runUpCommand(options: UpOptions, deps: CoreDependencies): options.spawn === true ? true : options.spawn === false ? false : Boolean(teamsConfig?.autoSpawn); if (shouldSpawn && teamsConfig && teamsConfig.agents.length > 0) { + const delivery = await waitForNodeDelivery(relay, deps); + if (!delivery.ready) { + deps.error('Refusing to auto-spawn agents because broker node delivery is not connected.'); + deps.error(`Node delivery: ${formatNodeDeliveryStatus(delivery.status)}`); + deps.error( + 'Realtime injection depends on /v1/node/ws. Check broker logs for create_node/node token errors, then retry `agent-relay up --spawn`.' + ); + await shutdownOnce(); + deps.exit(1); + return; + } for (const agent of teamsConfig.agents) { await relay.spawn({ name: agent.name, @@ -1066,6 +1138,7 @@ export async function runStatusCommand( if (typeof status.pending_delivery_count === 'number' && status.pending_delivery_count > 0) { deps.log(`Pending deliveries: ${status.pending_delivery_count}`); } + deps.log(`Node delivery: ${formatNodeDeliveryStatus(status)}`); if (session?.workspace_key) { deps.log(`Workspace Key: ${session.workspace_key}`); deps.log(`Observer: https://agentrelay.com/observer?key=${session.workspace_key}`); diff --git a/packages/cli/src/cli/lib/doctor.ts b/packages/cli/src/cli/lib/doctor.ts index 3368dee4d..29ccac6e9 100644 --- a/packages/cli/src/cli/lib/doctor.ts +++ b/packages/cli/src/cli/lib/doctor.ts @@ -4,6 +4,7 @@ import { execFileSync } from 'node:child_process'; import { createRequire } from 'node:module'; import { getProjectPaths } from '@agent-relay/config'; import { HarnessDriverClient, type BrokerStatus } from '@agent-relay/harness-driver'; +import { readNodeDeliveryStatus } from './broker-lifecycle.js'; type SqliteDriver = 'better-sqlite3' | 'node'; @@ -236,6 +237,15 @@ async function checkBrokerReliability(): Promise { : 'No authenticated Relaycast workspace reported by broker'; const pending = typedStatus.pending_deliveries ?? []; const stuck = pending.filter((delivery) => (delivery.age_ms ?? 0) >= 10_000 || delivery.last_error); + const nodeDelivery = readNodeDeliveryStatus(typedStatus); + const nodeDeliveryOk = Boolean(nodeDelivery?.tokenPresent && nodeDelivery.connected); + const nodeDeliveryMessage = !nodeDelivery + ? 'Node delivery status unavailable' + : nodeDelivery.connected + ? 'Node delivery connected' + : nodeDelivery.tokenPresent + ? 'Node token present, but /v1/node/ws is disconnected' + : 'No node token; /v1/node/ws cannot connect'; return [ ...(relayKeyTemplateResult ? [relayKeyTemplateResult] : []), { @@ -248,6 +258,14 @@ async function checkBrokerReliability(): Promise { ok: true, message: authMessage, }, + { + name: 'Node delivery', + ok: nodeDeliveryOk, + message: nodeDeliveryMessage, + remediation: !nodeDeliveryOk + ? 'Check broker logs for create_node/node token errors; realtime injection requires an active /v1/node/ws connection.' + : undefined, + }, { name: 'Outbound queues', ok: stuck.length === 0, diff --git a/packages/harness-driver/src/protocol.ts b/packages/harness-driver/src/protocol.ts index 04ced52be..1674d234e 100644 --- a/packages/harness-driver/src/protocol.ts +++ b/packages/harness-driver/src/protocol.ts @@ -289,6 +289,11 @@ export interface BrokerStatus { }>; pending_delivery_count: number; pending_deliveries: PendingDeliveryInfo[]; + node_connected?: boolean; + node_delivery?: { + token_present?: boolean; + connected?: boolean; + }; auth?: BrokerAuthStatus; } From 6a0ff63f2907f5444b406596d432d20ebadf6e45 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Tue, 30 Jun 2026 17:12:01 +0000 Subject: [PATCH 2/4] style: auto-format with Prettier --- .../trajectories/active/traj_lbmykojlyhtm/trajectory.json | 2 +- .../completed/2026-06/traj_eowv9c937zq9/trajectory.json | 2 +- .../completed/2026-06/traj_wwmbn0x18dso/summary.md | 4 +++- .../completed/2026-06/traj_wwmbn0x18dso/trajectory.json | 2 +- 4 files changed, 6 insertions(+), 4 deletions(-) diff --git a/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json b/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json index 05451eb40..0d302ee63 100644 --- a/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json +++ b/.agentworkforce/trajectories/active/traj_lbmykojlyhtm/trajectory.json @@ -43,4 +43,4 @@ "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" } -} \ No newline at end of file +} diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json index ca32f5cc5..6c832289f 100644 --- a/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json +++ b/.agentworkforce/trajectories/completed/2026-06/traj_eowv9c937zq9/trajectory.json @@ -22,4 +22,4 @@ "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" } -} \ No newline at end of file +} diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md index 34234f3f6..67b8348c0 100644 --- a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md +++ b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/summary.md @@ -18,6 +18,7 @@ Investigated local-up runtime config and PTY broker logs. Found local-up bound t ## Key Decisions ### Treat parent relay-hyperagent log as authoritative for local-up injection failure + - **Chose:** Treat parent relay-hyperagent log as authoritative for local-up injection failure - **Reasoning:** Child PTY logs were empty or only had readiness warnings, while the parent broker log records node token mint failure and node binding failures for claude and codex. @@ -26,6 +27,7 @@ Investigated local-up runtime config and PTY broker logs. Found local-up bound t ## Chapters ### 1. Work -*Agent: default* + +_Agent: default_ - Treat parent relay-hyperagent log as authoritative for local-up injection failure: Treat parent relay-hyperagent log as authoritative for local-up injection failure diff --git a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json index a01a0bebb..c15afe6f9 100644 --- a/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json +++ b/.agentworkforce/trajectories/completed/2026-06/traj_wwmbn0x18dso/trajectory.json @@ -50,4 +50,4 @@ "startRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e", "endRef": "05d586b789623067f35a1fd2e9c0a4c1cdc1a96e" } -} \ No newline at end of file +} From cc066f2b089b54ab217a737a7c587d510976f632 Mon Sep 17 00:00:00 2001 From: Proactive Runtime Bot Date: Tue, 30 Jun 2026 11:13:00 -0700 Subject: [PATCH 3/4] fix(broker): address PR #1216 review feedback - mint_node_token: keep transient send()/text() failures inside the retry loop (log + backoff + continue) instead of early-returning via `?`, so a network blip no longer bypasses retries. - listen_api /health: use try_send for the status fetch so a full channel degrades gracefully instead of blocking; reply still bounded by the 100ms timeout. - up waitForNodeDelivery: catch transient getStatus() errors and keep polling to the deadline instead of crashing the command. - fleet Connected: set node_delivery_token_present=true on connect so /health reports tokenPresent correctly after a later re-mint. - redact_relaycast_secrets: compile regexes once via OnceLock. Tests: cargo test -p agent-relay-broker --lib (795 passed, +retry test); CLI vitest (16) + tsc green. Co-Authored-By: Claude Opus 4.8 --- crates/broker/src/listen_api.rs | 3 +- crates/broker/src/node_control.rs | 132 ++++++++++++++++-- crates/broker/src/runtime/fleet.rs | 1 + .../cli/src/cli/lib/broker-lifecycle.test.ts | 35 +++++ packages/cli/src/cli/lib/broker-lifecycle.ts | 8 +- 5 files changed, 163 insertions(+), 16 deletions(-) diff --git a/crates/broker/src/listen_api.rs b/crates/broker/src/listen_api.rs index 4fd032fb3..ade82e5b7 100644 --- a/crates/broker/src/listen_api.rs +++ b/crates/broker/src/listen_api.rs @@ -520,8 +520,7 @@ async fn listen_api_health( async fn fetch_status_for_health(tx: &mpsc::Sender) -> Option { let (reply_tx, reply_rx) = tokio::sync::oneshot::channel(); - tx.send(ListenApiRequest::GetStatus { reply: reply_tx }) - .await + tx.try_send(ListenApiRequest::GetStatus { reply: reply_tx }) .ok()?; timeout(HEALTH_STATUS_TIMEOUT, reply_rx) .await diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index b1ef58e60..b9b36c894 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -2,6 +2,7 @@ use std::{ collections::{BTreeMap, HashMap, HashSet, VecDeque}, fs, path::Path, + sync::OnceLock, time::{Duration, Instant}, }; @@ -237,7 +238,7 @@ pub(crate) async fn mint_node_token( let mut last_error: Option = None; for attempt in 0..=CREATE_NODE_RETRY_BACKOFFS_MS.len() { - let response = client + let response = match client .post(&url) .bearer_auth(workspace_key) .header("X-SDK-Version", crate::util::version::broker_version()) @@ -252,9 +253,52 @@ pub(crate) async fn mint_node_token( ) .json(&request) .send() - .await?; + .await + { + Ok(response) => response, + Err(error) => { + let mint_error = CreateNodeMintError::Http(error); + log_create_node_mint_error( + "relay_broker::fleet", + context.node_id, + context.workspace_id, + &mint_error, + "create_node mint attempt failed", + ); + if attempt >= CREATE_NODE_RETRY_BACKOFFS_MS.len() { + return Err(mint_error); + } + last_error = Some(mint_error); + tokio::time::sleep(Duration::from_millis( + CREATE_NODE_RETRY_BACKOFFS_MS[attempt], + )) + .await; + continue; + } + }; let status = response.status().as_u16(); - let body = response.text().await?; + let body = match response.text().await { + Ok(body) => body, + Err(error) => { + let mint_error = CreateNodeMintError::Http(error); + log_create_node_mint_error( + "relay_broker::fleet", + context.node_id, + context.workspace_id, + &mint_error, + "create_node response body read failed", + ); + if attempt >= CREATE_NODE_RETRY_BACKOFFS_MS.len() { + return Err(mint_error); + } + last_error = Some(mint_error); + tokio::time::sleep(Duration::from_millis( + CREATE_NODE_RETRY_BACKOFFS_MS[attempt], + )) + .await; + continue; + } + }; let envelope = serde_json::from_str::(&body).map_err(|error| { CreateNodeMintError::InvalidResponse { status, @@ -352,17 +396,25 @@ fn normalized_relaycast_base_url(base_url: Option<&str>) -> String { } fn redact_relaycast_secrets(input: &str) -> String { + static BEARER_SECRET_RE: OnceLock = OnceLock::new(); + static TOKEN_FIELD_RE: OnceLock = OnceLock::new(); + let mut redacted = input.to_string(); - if let Ok(re) = + let bearer_re = BEARER_SECRET_RE.get_or_init(|| { regex::Regex::new(r#"(rk_live_|at_live_|nt_live_|br_|Bearer\s+)[A-Za-z0-9._-]+"#) - { - redacted = re.replace_all(&redacted, "${1}[REDACTED]").into_owned(); - } - if let Ok(re) = regex::Regex::new(r#""token"\s*:\s*"[^"]+""#) { - redacted = re - .replace_all(&redacted, r#""token":"[REDACTED]""#) - .into_owned(); - } + .expect("bearer secret redaction regex must compile") + }); + redacted = bearer_re + .replace_all(&redacted, "${1}[REDACTED]") + .into_owned(); + + let token_field_re = TOKEN_FIELD_RE.get_or_init(|| { + regex::Regex::new(r#""token"\s*:\s*"[^"]+""#) + .expect("token field redaction regex must compile") + }); + redacted = token_field_re + .replace_all(&redacted, r#""token":"[REDACTED]""#) + .into_owned(); redacted } @@ -1523,8 +1575,14 @@ pub(crate) fn delivery_ack(agent: impl Into, up_to_seq: u64) -> BrokerTo #[cfg(test)] mod tests { + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }; + use httpmock::{Method::POST, MockServer}; use serde_json::json; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; use tokio_tungstenite::accept_async; @@ -1672,6 +1730,56 @@ mod tests { ); } + #[tokio::test] + async fn mint_node_token_retries_transient_send_error_then_succeeds() { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("test server should bind"); + let base_url = format!("http://{}", listener.local_addr().expect("local addr")); + let attempts = Arc::new(AtomicUsize::new(0)); + let server_attempts = Arc::clone(&attempts); + let server = tokio::spawn(async move { + loop { + let (mut socket, _) = listener.accept().await.expect("accept request"); + let attempt = server_attempts.fetch_add(1, Ordering::SeqCst) + 1; + if attempt == 1 { + drop(socket); + continue; + } + + let mut buffer = [0_u8; 4096]; + let _ = socket.read(&mut buffer).await.expect("read request"); + let body = r#"{"ok":true,"data":{"id":"node_abc","name":"local-node","kind":"ws","role":"broker","version":"relay-broker/test","status":"online","live":true,"handlers_live":true,"load":0.0,"active_agents":0,"max_agents":0,"created_at":"2026-06-30T00:00:00Z","token":"nt_live_retry_success"}}"#; + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}", + body.len(), + body + ); + socket + .write_all(response.as_bytes()) + .await + .expect("write response"); + break; + } + }); + + let token = mint_node_token( + "rk_live_test", + Some(&base_url), + create_node_request("node_abc", "local-node", "relay-broker/test"), + MintNodeTokenLogContext { + node_id: "node_abc", + workspace_id: "ws_test", + }, + ) + .await + .expect("transient send failure should be retried"); + + server.await.expect("test server should finish"); + assert_eq!(token, "nt_live_retry_success"); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } + #[test] fn delivery_book_dedups_and_tracks_cumulative_ack() { let mut book = FleetDeliveryBook::default(); diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index 377b2efb3..cb79bc3f2 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -314,6 +314,7 @@ impl BrokerRuntime { pub(super) async fn handle_fleet_control_event(&mut self, event: FleetControlEvent) { match event { FleetControlEvent::Connected => { + self.node_delivery_token_present = true; self.node_delivery_connected = true; // Node delivery is live: message delivery flows solely over // /v1/node/ws. The workspace firehose delivery path was removed, diff --git a/packages/cli/src/cli/lib/broker-lifecycle.test.ts b/packages/cli/src/cli/lib/broker-lifecycle.test.ts index 2b3edec5d..a40fbed65 100644 --- a/packages/cli/src/cli/lib/broker-lifecycle.test.ts +++ b/packages/cli/src/cli/lib/broker-lifecycle.test.ts @@ -5,6 +5,7 @@ import { classifyBrokerStartStage, describeError, readNodeDeliveryStatus, + waitForNodeDelivery, } from './broker-lifecycle.js'; describe('describeError', () => { @@ -115,3 +116,37 @@ describe('readNodeDeliveryStatus', () => { expect(readNodeDeliveryStatus('nope')).toBeNull(); }); }); + +describe('waitForNodeDelivery', () => { + it('continues polling after a transient status failure', async () => { + let now = 0; + let calls = 0; + const relay = { + async getStatus() { + calls += 1; + if (calls === 1) { + throw new Error('broker not ready yet'); + } + return { + node_connected: calls >= 3, + node_delivery: { token_present: true, connected: calls >= 3 }, + }; + }, + }; + const deps = { + now: () => now, + sleep: async (ms: number) => { + now += ms; + }, + }; + + await expect(waitForNodeDelivery(relay as never, deps as never, 1_000)).resolves.toEqual({ + ready: true, + status: { + node_connected: true, + node_delivery: { token_present: true, connected: true }, + }, + }); + expect(calls).toBe(3); + }); +}); diff --git a/packages/cli/src/cli/lib/broker-lifecycle.ts b/packages/cli/src/cli/lib/broker-lifecycle.ts index d8c64fbae..dceb483ce 100644 --- a/packages/cli/src/cli/lib/broker-lifecycle.ts +++ b/packages/cli/src/cli/lib/broker-lifecycle.ts @@ -738,7 +738,7 @@ async function waitForBrokerReadiness( return latest; } -async function waitForNodeDelivery( +export async function waitForNodeDelivery( relay: CoreRelay, deps: CoreDependencies, waitMs = NODE_DELIVERY_READY_TIMEOUT_MS @@ -747,7 +747,11 @@ async function waitForNodeDelivery( let latest: unknown = null; while (true) { - latest = await relay.getStatus(); + try { + latest = await relay.getStatus(); + } catch { + latest = null; + } if (nodeDeliveryReady(latest)) { return { ready: true, status: latest }; } From cfe8acc24769bc119d4d681639bc45204e58da90 Mon Sep 17 00:00:00 2001 From: Proactive Runtime Bot Date: Tue, 30 Jun 2026 11:20:44 -0700 Subject: [PATCH 4/4] fix(broker): appease clippy for create-node retry loop --- crates/broker/src/node_control.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index b9b36c894..98f5f9093 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -237,6 +237,9 @@ pub(crate) async fn mint_node_token( let client = reqwest::Client::new(); let mut last_error: Option = None; + // This loop intentionally runs one more time than there are backoffs: the + // final iteration returns the last error instead of sleeping again. + #[allow(clippy::needless_range_loop)] for attempt in 0..=CREATE_NODE_RETRY_BACKOFFS_MS.len() { let response = match client .post(&url)