diff --git a/CHANGELOG.md b/CHANGELOG.md index e0d0ab2bc..1ada22779 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Broker node connections now require an accepted registration before reporting readiness or publishing inventory and heartbeats; rejected and unanswered registrations reconnect with bounded backoff. - Direct-message CLI and MCP receipts separate directory name matches from unconfirmed recipient delivery, preserving the queued message ID without reporting reachability from a roster match. ## [12.1.0] - 2026-09-12 diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index 2ae45d3d3..c5663bbfb 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -1909,6 +1909,72 @@ impl Drop for ProbeSessionGuard<'_> { } } +/// An authenticated socket is not yet the authoritative broker provider. The +/// engine can reject node.register while leaving the socket open; sending +/// inventory or heartbeats then updates a fallback provider and masks the loss. +/// Only this request's successful reply opens the application delivery path. +async fn register_node_session( + sink: &mut S, + stream: &mut R, + registration: &mut NodeRegister, + config: &FleetControlConfig, +) -> bool +where + S: Sink + Unpin, + S::Error: std::error::Error + Send + Sync + 'static, + R: futures_util::Stream> + Unpin, +{ + let id = format!("node_register_{}", Uuid::new_v4().simple()); + registration.id = Some(id.clone()); + // Bound both the write and the response, including peers which keep sending + // pongs or unrelated replies. Tests may use their shorter transport budget. + let deadline = config + .read_idle_timeout + .unwrap_or(Duration::from_secs(10)) + .min(Duration::from_secs(10)); + let accepted = tokio::time::timeout(deadline, async { + if send_wire(sink, &BrokerToRelaycast::NodeRegister(registration.clone())).await.is_err() { + return false; + } + while let Some(Ok(message)) = stream.next().await { + match message { + Message::Text(text) => { + if let Some(probe) = config.probe.as_ref() { probe.record_text_frame(); } + match serde_json::from_str::(&text) { + Ok(frame) => { + if let Some(probe) = config.probe.as_ref() { probe.record_frame(&frame); } + match frame { + RelaycastToBroker::Reply(reply) if reply.id == id => return reply.ok, + RelaycastToBroker::Error(error) => { + tracing::error!(code = %error.code, "node registration rejected; reconnecting without advertising delivery readiness"); + return false; + } + // The engine replies before replaying deliveries. Never + // acknowledge or inject a frame on an unaccepted provider. + RelaycastToBroker::Deliver(_) | RelaycastToBroker::ActionInvoke(_) => return false, + _ => {}, + } + } + Err(error) => { + if let Some(probe) = config.probe.as_ref() { probe.record_parse_failure(&error.to_string(), &text); } + } + } + } + Message::Ping(payload) => { + if sink.send(Message::Pong(payload)).await.is_err() { return false; } + } + Message::Close(_) => return false, + _ => {}, + } + } + false + }).await.unwrap_or(false); + if !accepted { + tracing::warn!("node registration was not accepted within its deadline; delivery unavailable, reconnecting"); + } + accepted +} + async fn run_connected_once( config: &FleetControlConfig, command_rx: &mut mpsc::Receiver, @@ -1981,21 +2047,16 @@ async fn run_connected_once( }; // Socket-owned connectivity: armed here, released by Drop on every exit. let _probe_session = ProbeSessionGuard::enter(config.probe.as_ref()); - let _ = event_tx.send(FleetControlEvent::Connected).await; let (mut sink, mut stream) = ws.split(); let mut pending_agent_registrations: HashMap = HashMap::new(); let mut pending_deregistrations: HashMap>> = HashMap::new(); - if send_wire( - &mut sink, - &BrokerToRelaycast::NodeRegister(node_register.clone()), - ) - .await - .is_err() - { + if !register_node_session(&mut sink, &mut stream, &mut node_register, config).await { return ControlRunResult::Disconnected; } + *registration = Some(node_register.clone()); + let _ = event_tx.send(FleetControlEvent::Connected).await; if !send_inventory_sync(&mut sink, inventory, &mut pending_agent_registrations).await { return ControlRunResult::Disconnected; } @@ -2032,11 +2093,15 @@ async fn run_connected_once( load.handlers_live = true; let mut next = build_node_register(&manifest, &config.node_id, &config.node_name, &config.broker_version, resume_cursor); next.provider = Some(provider.clone()); - node_register = next.clone(); - *registration = Some(next.clone()); - if send_wire(&mut sink, &BrokerToRelaycast::NodeRegister(next)).await.is_err() { - return ControlRunResult::Disconnected; + // Reopen the provider session for a new manifest. Running + // the registration gate inside this active socket would + // consume replies belonging to in-flight agent requests. + *registration = Some(next); + drain_agent_registrations(&mut pending_agent_registrations, "node_control_reconfiguring"); + for (_, pending) in pending_deregistrations.drain() { + let _ = pending.send(Err("node_control_reconfiguring".to_string())); } + return ControlRunResult::Disconnected; } Some(FleetControlCommand::UpdateInventory(next)) => { *inventory = next; @@ -4241,6 +4306,7 @@ mod tests { let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.unwrap(); let mut ws = accept_async(stream).await.unwrap(); + let _ = next_node_to_server(&mut ws).await; while ws.next().await.is_some() {} }); @@ -4609,7 +4675,7 @@ mod tests { ); let driver = async { // Both frames counted -> the client has drained the read side. - wait_for_probe(&probe, |snapshot| snapshot["socket"]["text_frames"] == 2).await; + wait_for_probe(&probe, |snapshot| snapshot["socket"]["text_frames"] == 3).await; command_tx .send(FleetControlCommand::Shutdown) .await @@ -4624,10 +4690,10 @@ mod tests { server.abort(); let snapshot = probe.snapshot_with_token(true); - // BOTH frames are counted as arrived, including the one that could not - // be parsed. If the count moved after the parse, this would read 1. + // Registration reply plus BOTH test frames count as arrived, including + // the unparseable one. Counting after parse would incorrectly read 2. assert_eq!( - snapshot["socket"]["text_frames"], 2, + snapshot["socket"]["text_frames"], 3, "an unparseable frame must still count as having arrived: {snapshot}" ); assert_eq!(snapshot["socket"]["parse_failures"], 1); @@ -5079,11 +5145,21 @@ mod tests { async fn next_node_to_server(ws: &mut S) -> BrokerToRelaycast where S: futures_util::Stream> + + Sink + Unpin, { loop { if let Message::Text(text) = ws.next().await.unwrap().unwrap() { - return serde_json::from_str(&text).unwrap(); + let frame = serde_json::from_str(&text).unwrap(); + if let BrokerToRelaycast::NodeRegister(register) = &frame { + ws.send(Message::Text( + json!({"v":1,"type":"reply","id":register.id,"ok":true,"data":{}}) + .to_string(), + )) + .await + .unwrap(); + } + return frame; } } } @@ -5091,6 +5167,7 @@ mod tests { async fn next_non_heartbeat_node_to_server(ws: &mut S) -> BrokerToRelaycast where S: futures_util::Stream> + + Sink + Unpin, { loop { @@ -5336,3 +5413,6 @@ mod tests { )); } } + +#[cfg(test)] +mod registration_tests; diff --git a/crates/broker/src/node_control/registration_tests.rs b/crates/broker/src/node_control/registration_tests.rs new file mode 100644 index 000000000..efb19788b --- /dev/null +++ b/crates/broker/src/node_control/registration_tests.rs @@ -0,0 +1,284 @@ +//! Real loopback WS regression for an authenticated but unregistered provider. +use super::*; +use serde_json::{json, Value}; +use tokio::net::TcpListener; +use tokio_tungstenite::accept_async; + +async fn registration_gate_case(response: &str) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let ws_url = format!("ws://{}/v1/node/ws", listener.local_addr().unwrap()); + let (command_tx, mut command_rx) = mpsc::channel(8); + let (event_tx, mut event_rx) = mpsc::channel(8); + let mut registration = Some(NodeRegister { + v: FLEET_WIRE_VERSION, + id: None, + name: "test-node".into(), + node_id: "node-test".into(), + provider: None, + capabilities: vec![], + max_agents: 8, + tags: vec![], + repo_keys: None, + version: "test".into(), + machine_id: None, + resume_cursor: None, + }); + let mut inventory = vec![InventoryAgent { + name: "old-worker".into(), + agent_id: "old-worker-id".into(), + invocation_id: None, + session_ref: Some("old-session".into()), + }]; + let mut load = FleetLoadSnapshot { + active_agents: 1, + max_agents: 8, + handlers_live: true, + active_agent_names: vec!["old-worker".into()], + }; + let accepted = response == "accept" || response == "reconfigure"; + let reconfigure = response == "reconfigure"; + let server_command_tx = command_tx.clone(); + let (forwarded_tx, mut forwarded_rx) = oneshot::channel(); + let response = response.to_owned(); + let server = tokio::spawn(async move { + let (tcp, _) = listener.accept().await.unwrap(); + let mut ws = accept_async(tcp).await.unwrap(); + let Message::Text(raw) = ws.next().await.unwrap().unwrap() else { + panic!("node.register expected") + }; + let frame: Value = serde_json::from_str(&raw).unwrap(); + assert_eq!(frame["type"], "node.register"); + let id = frame["id"].as_str().unwrap_or("uncorrelated-base-request"); + if accepted { + // Keep the socket live without accepting the provider. Neither an + // unrelated success nor transport traffic may open the gate. + ws.send(Message::Text( + json!({"v":1,"type":"reply","ok":true,"id":"unrelated","data":{}}).to_string(), + )) + .await + .unwrap(); + // A pong proves the client processed the preceding unrelated reply. + // No dependent text frame may precede this transport-only barrier. + ws.send(Message::Ping(b"registration-barrier".to_vec())) + .await + .unwrap(); + assert!( + matches!(ws.next().await, Some(Ok(Message::Pong(payload))) if payload == b"registration-barrier") + ); + ws.send(Message::Text( + json!({"v":1,"type":"reply","ok":true,"id":id,"data":{}}).to_string(), + )) + .await + .unwrap(); + for expected in ["inventory.sync", "node.heartbeat"] { + let Message::Text(raw) = ws.next().await.unwrap().unwrap() else { + panic!("text frame expected"); + }; + assert_eq!( + serde_json::from_str::(&raw).unwrap()["type"], + expected + ); + } + if reconfigure { + let (register_reply, register_result) = oneshot::channel(); + let (deregister_reply, deregister_result) = oneshot::channel(); + server_command_tx + .send(FleetControlCommand::RegisterAgent { + request: serde_json::from_value(json!({"v":1,"name":"pending-worker"})) + .unwrap(), + reply: register_reply, + }) + .await + .unwrap(); + server_command_tx + .send(FleetControlCommand::DeregisterAgent { + request: serde_json::from_value( + json!({"v":1,"agent_id":"retiring-worker-id","name":"retiring-worker"}), + ) + .unwrap(), + reply: deregister_reply, + }) + .await + .unwrap(); + for expected in ["agent.register", "agent.deregister"] { + loop { + let Message::Text(raw) = ws.next().await.unwrap().unwrap() else { + continue; + }; + let frame: Value = serde_json::from_str(&raw).unwrap(); + if frame["type"] == "node.heartbeat" { + continue; + } + assert_eq!(frame["type"], expected); + break; + } + } + server_command_tx + .send(FleetControlCommand::RegisterNode { + manifest: NodeManifest { + name: "updated-node".into(), + node_id: None, + capabilities: vec![], + max_agents: None, + tags: None, + repo_keys: None, + version: None, + }, + resume_cursor: None, + }) + .await + .unwrap(); + assert_eq!( + register_result.await.unwrap().unwrap_err(), + "node_control_reconfiguring" + ); + assert_eq!( + deregister_result.await.unwrap().unwrap_err(), + "node_control_reconfiguring" + ); + return; + } + for name in ["old-worker", "fresh-worker"] { + ws.send(Message::Text(json!({"v":1,"type":"deliver","agent":name,"agent_id":format!("{name}-id"),"delivery_id":format!("delivery-{name}"),"msg_id":format!("message-{name}"),"seq":1,"mode":"wait","payload":{"type":"dm.received","text":"local probe"}}).to_string())).await.unwrap(); + } + // Keep the peer polling (and answering pings) until the runtime + // has forwarded both events. Closing immediately after send races + // the client's heartbeat write against draining buffered deliveries. + loop { + tokio::select! { + result = &mut forwarded_rx => { + result.expect("paired runtime observations completed"); + ws.close(None).await.unwrap(); + return; + } + message = ws.next() => { + assert!(matches!(message, Some(Ok(_))), "accepted client disconnected before forwarding paired events: {message:?}"); + } + } + } + } + if response != "timeout" { + let reply = match response.as_str() { + "error" => { + json!({"v":1,"type":"error","ok":false,"id":id,"code":"provider_instance_conflict","message":"incumbent provider still live"}) + } + "false" => json!({"v":1,"type":"reply","ok":false,"id":id,"data":{}}), + "uncorrelated" => { + json!({"v":1,"type":"reply","ok":true,"id":"some-other-request","data":{}}) + } + _ => panic!("unknown arm"), + }; + ws.send(Message::Text(reply.to_string())).await.unwrap(); + } + // A failed registration must never be followed by inventory, heartbeat, + // agent registration or an ACK on the unauthoritative socket. + while let Ok(Some(Ok(frame))) = + tokio::time::timeout(Duration::from_secs(1), ws.next()).await + { + match frame { + Message::Text(raw) => panic!( + "registration-dependent frame escaped gate: {}", + serde_json::from_str::(&raw).unwrap()["type"] + ), + Message::Close(_) => break, + _ => {} + } + } + }); + let config = FleetControlConfig { + ws_url, + node_token: Some("nt_test".into()), + node_id: "node-test".into(), + node_name: "test-node".into(), + broker_version: "test".into(), + token_minter: None, + session_token: None, + // Only the rejection/timeout arms test the short registration deadline. + // Positive delivery completion is event-coordinated and bounded below. + read_idle_timeout: if accepted { + None + } else { + Some(Duration::from_millis(150)) + }, + probe: None, + }; + let observe = async { + if accepted { + let connected = event_rx.recv().await; + assert!( + matches!(connected, Some(FleetControlEvent::Connected)), + "accepted registration must connect: {connected:?}" + ); + if !reconfigure { + for expected in ["old-worker", "fresh-worker"] { + let event = event_rx.recv().await; + let Some(FleetControlEvent::Message(RelaycastToBroker::Deliver(deliver))) = + event + else { + panic!("accepted provider must forward {expected}; observed {event:?}"); + }; + assert_eq!(deliver.agent, expected); + assert_eq!(deliver.msg_id, format!("message-{expected}")); + } + forwarded_tx + .send(()) + .expect("accepted peer remains open until observations complete"); + } + } + }; + let (result, ()) = tokio::time::timeout(Duration::from_secs(10), async { + tokio::join!( + run_connected_once( + &config, + &mut command_rx, + &event_tx, + &mut registration, + &mut inventory, + &mut load, + if accepted { + Duration::from_secs(60) + } else { + Duration::from_millis(50) + }, + ), + observe, + ) + }) + .await + .expect("registration case must terminate after its bounded protocol exchange"); + assert_eq!(result, ControlRunResult::Disconnected); + assert!( + event_rx.try_recv().is_err(), + "no unexpected or duplicate runtime events" + ); + server + .await + .expect("the rejected socket emitted no dependent frames"); +} + +#[tokio::test] +async fn rejected_registration_never_advertises_or_syncs() { + registration_gate_case("error").await; +} +#[tokio::test] +async fn unsuccessful_registration_reply_never_advertises_or_syncs() { + registration_gate_case("false").await; +} +#[tokio::test] +async fn silent_registration_never_advertises_or_syncs() { + registration_gate_case("timeout").await; +} +#[tokio::test] +async fn unrelated_reply_cannot_open_registration_gate() { + registration_gate_case("uncorrelated").await; +} + +#[tokio::test] +async fn accepted_registration_forwards_paired_deliveries_once() { + registration_gate_case("accept").await; +} + +#[tokio::test] +async fn manifest_change_fails_pending_requests_before_reconnect() { + registration_gate_case("reconfigure").await; +} diff --git a/packages/cloud/src/credential-directory-windows.native.test.ts b/packages/cloud/src/credential-directory-windows.native.test.ts index 209bcfe86..b33e1676e 100644 --- a/packages/cloud/src/credential-directory-windows.native.test.ts +++ b/packages/cloud/src/credential-directory-windows.native.test.ts @@ -15,33 +15,51 @@ afterEach(() => { directory = undefined; }); +// assertWindowsCredentialDirectory() shells out to a cold powershell.exe process +// and allows up to WINDOWS_ACL_TIMEOUT_MS (15_000ms) in credential-directory-windows.ts. +// Give these tests a per-test timeout comfortably above that so a cold PowerShell +// start doesn't trip vitest's default 5000ms timeout. +const WINDOWS_ACL_TEST_TIMEOUT_MS = 20_000; + describeWindows('native Windows credential directory ACL validation', () => { - it('accepts a private directory under the user profile without changing its ACL', () => { - directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-')); - expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); - }); - - it('rejects an untrusted read grant on the credential directory itself', () => { - directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-unsafe-')); - expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); - execFileSync('icacls.exe', [directory, '/grant', '*S-1-1-0:(R)'], { - stdio: 'ignore', - windowsHide: true, - }); - expect(() => assertWindowsCredentialDirectory(directory!)).toThrow( - 'Windows Relaycast credential storage requires a private directory' - ); - }); - - it('rejects an untrusted generic-all grant on the credential directory itself', () => { - directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-generic-unsafe-')); - expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); - execFileSync('icacls.exe', [directory, '/grant', '*S-1-1-0:(GA)'], { - stdio: 'ignore', - windowsHide: true, - }); - expect(() => assertWindowsCredentialDirectory(directory!)).toThrow( - 'Windows Relaycast credential storage requires a private directory' - ); - }); + it( + 'accepts a private directory under the user profile without changing its ACL', + () => { + directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-')); + expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); + }, + WINDOWS_ACL_TEST_TIMEOUT_MS + ); + + it( + 'rejects an untrusted read grant on the credential directory itself', + () => { + directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-unsafe-')); + expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); + execFileSync('icacls.exe', [directory, '/grant', '*S-1-1-0:(R)'], { + stdio: 'ignore', + windowsHide: true, + }); + expect(() => assertWindowsCredentialDirectory(directory!)).toThrow( + 'Windows Relaycast credential storage requires a private directory' + ); + }, + WINDOWS_ACL_TEST_TIMEOUT_MS + ); + + it( + 'rejects an untrusted generic-all grant on the credential directory itself', + () => { + directory = fs.mkdtempSync(path.join(os.homedir(), '.relay-acl-native-generic-unsafe-')); + expect(() => assertWindowsCredentialDirectory(directory!)).not.toThrow(); + execFileSync('icacls.exe', [directory, '/grant', '*S-1-1-0:(GA)'], { + stdio: 'ignore', + windowsHide: true, + }); + expect(() => assertWindowsCredentialDirectory(directory!)).toThrow( + 'Windows Relaycast credential storage requires a private directory' + ); + }, + WINDOWS_ACL_TEST_TIMEOUT_MS + ); }); diff --git a/tests/relayflows/cases/1593-node-registration-gate/case.json b/tests/relayflows/cases/1593-node-registration-gate/case.json new file mode 100644 index 000000000..38a62e80a --- /dev/null +++ b/tests/relayflows/cases/1593-node-registration-gate/case.json @@ -0,0 +1,13 @@ +{ + "version": 1, + "id": "1593-node-registration-gate", + "kind": "bugfix", + "title": "Reject an unauthoritative node session before inventory or readiness", + "runner": { "command": ["node", "tests/relayflows/cases/1593-node-registration-gate/run.mjs"] }, + "requirements": ["broker-linux-x64"], + "timeoutSeconds": 900, + "expected": { + "base": { "outcome": "bug", "signature": "rejected_registration_still_publishes_inventory" }, + "head": { "outcome": "fixed", "signature": "rejected_registration_never_opens_delivery_path" } + } +} diff --git a/tests/relayflows/cases/1593-node-registration-gate/run.mjs b/tests/relayflows/cases/1593-node-registration-gate/run.mjs new file mode 100644 index 000000000..949ef9331 --- /dev/null +++ b/tests/relayflows/cases/1593-node-registration-gate/run.mjs @@ -0,0 +1,279 @@ +// Reuses the dependency-free HTTP/RFC6455 stand-in from relay#1636. +// Exercise the exact provided broker artifact; no source compilation or edits. +import assert from 'node:assert/strict'; +import crypto from 'node:crypto'; +import http from 'node:http'; +import { execFileSync, spawn } from 'node:child_process'; +import { mkdtemp, mkdir, writeFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +const required = (name) => { + const v = process.env[name]; + if (!v) throw Error(`Missing ${name}`); + return v; +}; +const target = required('RELAY_PR_PROOF_TARGET_DIR'); +const harness = required('RELAY_PR_PROOF_HARNESS_DIR'); +const binary = required('RELAY_PR_PROOF_BROKER_BINARY'); +const resultPath = required('RELAY_PR_PROOF_RESULT_PATH'); +const arm = required('RELAY_PR_PROOF_ARM'); +assert.ok(['base', 'head'].includes(arm)); +const gitSha = (dir) => execFileSync('git', ['-C', dir, 'rev-parse', 'HEAD'], { encoding: 'utf8' }).trim(); +assert.equal( + gitSha(target), + required(arm === 'base' ? 'RELAY_PR_PROOF_BASE_SHA' : 'RELAY_PR_PROOF_HEAD_SHA') +); +assert.equal(gitSha(harness), required('RELAY_PR_PROOF_HEAD_SHA')); +const relative = path.relative(harness, fileURLToPath(import.meta.url)); +assert.ok(relative && !relative.startsWith('..') && !path.isAbsolute(relative)); +const root = await mkdtemp(path.join(tmpdir(), 'relayflow-registration-')); +const state = path.join(root, 'state'); +await mkdir(state); +const sockets = new Set(); +const sessions = []; +let broker; +const WS_GUID = '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'; +function encodeTextFrame(text) { + const payload = Buffer.from(text, 'utf8'); + const length = payload.length; + let header; + if (length < 126) { + header = Buffer.from([0x81, length]); + } else if (length < 65536) { + header = Buffer.alloc(4); + header[0] = 0x81; + header[1] = 126; + header.writeUInt16BE(length, 2); + } else { + header = Buffer.alloc(10); + header[0] = 0x81; + header[1] = 127; + header.writeBigUInt64BE(BigInt(length), 2); + } + return Buffer.concat([header, payload]); +} + +function createFrameReader(onText) { + let buffer = Buffer.alloc(0); + return (chunk) => { + buffer = Buffer.concat([buffer, chunk]); + for (;;) { + if (buffer.length < 2) return; + const opcode = buffer[0] & 0x0f; + const masked = (buffer[1] & 0x80) !== 0; + let length = buffer[1] & 0x7f; + let offset = 2; + if (length === 126) { + if (buffer.length < offset + 2) return; + length = buffer.readUInt16BE(offset); + offset += 2; + } else if (length === 127) { + if (buffer.length < offset + 8) return; + length = Number(buffer.readBigUInt64BE(offset)); + offset += 8; + } + let mask = null; + if (masked) { + if (buffer.length < offset + 4) return; + mask = buffer.subarray(offset, offset + 4); + offset += 4; + } + if (buffer.length < offset + length) return; + const payload = Buffer.from(buffer.subarray(offset, offset + length)); + buffer = buffer.subarray(offset + length); + if (mask) for (let i = 0; i < payload.length; i += 1) payload[i] ^= mask[i % 4]; + if (opcode === 0x1) onText(payload.toString('utf8')); + } + }; +} + +const server = http.createServer((request, response) => { + let body = ''; + request.on('data', (chunk) => { + body += chunk; + }); + request.on('end', () => { + const url = request.url.split('?')[0]; + const send = (data) => { + response.writeHead(200, { 'content-type': 'application/json' }); + response.end(JSON.stringify({ ok: true, data })); + }; + if (request.method === 'POST' && url === '/v1/agents') { + let parsed = {}; + try { + parsed = JSON.parse(body || '{}'); + } catch {} + send({ + id: 'agt_relayflow_broker', + workspace_id: 'ws_relayflow', + name: parsed.name ?? 'broker', + token: 'at_relayflow_broker', + status: 'online', + created_at: '2026-09-01T00:00:00.000Z', + }); + return; + } + if (url === '/v1/agents' || url === '/v1/channels') { + send([]); + return; + } + if (url.startsWith('/v1/agents/')) { + send({ id: 'agt_relayflow_other', name: 'other', status: 'offline', metadata: {} }); + return; + } + send({}); + }); +}); + +server.on('upgrade', (request, socket) => { + const key = request.headers['sec-websocket-key']; + if (!key) { + socket.destroy(); + return; + } + const accept = crypto + .createHash('sha1') + .update(key + WS_GUID) + .digest('base64'); + socket.write( + 'HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: ' + + accept + + '\r\n\r\n' + ); + sockets.add(socket); + socket.on('error', () => {}); + socket.on('close', () => sockets.delete(socket)); + if (request.url.split('?')[0] !== '/v1/node/ws') return; + const session = { + registration: false, + rejected: sessions.length === 0, + closed: false, + inventory: 0, + heartbeat: 0, + other: 0, + }; + sessions.push(session); + socket.on('close', () => { + session.closed = true; + }); + // Upgraded HTTP sockets can remain writable after peer FIN; read EOF is + // the actual client-disconnect signal, independent of our writable half. + socket.on('end', () => { + session.closed = true; + socket.end(); + }); + socket.on( + 'data', + createFrameReader((text) => { + const frame = JSON.parse(text); + if (frame.type === 'node.register') { + session.registration = true; + const reply = session.rejected + ? { + v: 1, + type: 'error', + id: frame.id ?? 'base-registration', + ok: false, + code: 'provider_instance_conflict', + message: 'incumbent provider still live', + } + : { v: 1, type: 'reply', id: frame.id, ok: true, data: {} }; + socket.write(encodeTextFrame(JSON.stringify(reply))); + } else if (frame.type === 'inventory.sync') session.inventory++; + else if (frame.type === 'node.heartbeat') session.heartbeat++; + else session.other++; + }) + ); +}); +try { + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const port = server.address().port; + broker = spawn( + binary, + [ + 'init', + '--instance-name', + 'registration-proof', + '--api-port', + '0', + '--api-bind', + '127.0.0.1', + '--state-dir', + state, + ], + { + cwd: root, + stdio: ['ignore', 'pipe', 'pipe'], + env: { + PATH: process.env.PATH, + HOME: root, + TMPDIR: root, + RELAY_API_KEY: 'rk_local_registration_proof', + RELAYCAST_BASE_URL: `http://127.0.0.1:${port}`, + RELAY_BASE_URL: `http://127.0.0.1:${port}`, + RELAY_NODE_TOKEN: 'nt_local_registration_proof', + RELAY_TELEMETRY_DISABLED: '1', + RELAY_SKIP_TELEMETRY: '1', + }, + } + ); + let spawnFailed = false; + broker.once('error', () => { + spawnFailed = true; + }); + // Drain output without retaining peer bodies, credentials or other runtime data. + broker.stdout.resume(); + broker.stderr.resume(); + let outcome, signature; + const deadline = Date.now() + 120000; + while (Date.now() < deadline) { + if (spawnFailed) throw Error('Broker artifact could not be started.'); + if (broker.exitCode !== null) throw Error(`Broker exited before evidence: ${broker.exitCode}`); + const first = sessions[0]; + if (first?.registration && (first.inventory || first.heartbeat || first.other)) { + outcome = 'bug'; + signature = 'rejected_registration_still_publishes_inventory'; + break; + } + if ( + first?.registration && + first.closed && + sessions.slice(1).some((s) => s.registration && s.inventory > 0 && s.heartbeat > 0) + ) { + assert.equal(first.inventory + first.heartbeat + first.other, 0); + outcome = 'fixed'; + signature = 'rejected_registration_never_opens_delivery_path'; + break; + } + await new Promise((resolve) => setTimeout(resolve, 25)); + } + if (!outcome) throw Error('No terminal registration discriminator observed within120s.'); + await mkdir(path.dirname(resultPath), { recursive: true }); + await writeFile( + resultPath, + JSON.stringify({ + version: 1, + caseId: '1593-node-registration-gate', + arm, + outcome, + signature, + details: + 'The peer rejects the first node.register and leaves its socket open. Base sends dependent frames anyway. Head closes that socket without inventory/heartbeat/ACK and reconnects; an accepted second registration then publishes inventory and heartbeat.', + sessions, + }) + '\n' + ); + console.log(signature); +} finally { + if (broker && broker.exitCode === null) { + broker.kill('SIGTERM'); + await Promise.race([ + new Promise((resolve) => broker.once('exit', resolve)), + new Promise((resolve) => setTimeout(resolve, 3000)), + ]); + if (broker.exitCode === null) broker.kill('SIGKILL'); + } + for (const socket of sockets) socket.destroy(); + await new Promise((resolve) => server.close(resolve)); + await rm(root, { recursive: true, force: true }); +}