Skip to main content

gate4agent_c2/
runtime.rs

1use crate::protocol::{
2    C2ErrorCategory, C2NodeEvent, C2NodeFailure, C2NodeResponse, C2RelayFailure, C2RelayFailureCode, GapKind, HealthResponse, NodeCursor, NodeFreshness, NodeGap, NodeId,
3    NodeIncarnationId, NodeRequest, ResolvedSpawnReceipt, RoutedNodeEvent, RoutedNodeResponse,
4    ManagedWorktreeLeaseState, ManagedWorktreeSpawnRequest, ManagedWorktreeSpawnRequestV2,
5    SpawnOverride, SpawnSpec,
6    NodeTransportState, ObservedNode, ProviderAdapterContractSupport, ProviderContractSupport,
7    ReadyResponse, SanitizedError, SlimNodeInventory, StatusResponse,
8    C2_API_VERSION, C2_PROVIDER_CONTRACT_MANIFEST_CAPABILITY,
9    MAX_C2_GAPS_PER_NODE, MAX_C2_NODES,
10};
11#[cfg(windows)]
12use crate::protocol::MAX_C2_ENDPOINT_BYTES;
13use gate4agent_node_protocol::{
14    ClientRole, FrameError, NegotiatedNodeCompatibility, NodeEvent, NodeEventEnvelope, NodeFailureCode,
15    NodeResponse, NodeSnapshot, ServerFrame,
16};
17use gate4agent_node_wire::{read_call_home_announce, LocalNodeClient, NodeClientError};
18use std::collections::{BTreeMap, BTreeSet};
19use std::io;
20use std::net::SocketAddr;
21use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
22use std::sync::{Arc, Mutex};
23use std::time::{Duration, SystemTime, UNIX_EPOCH};
24#[cfg(unix)]
25use std::path::{Path, PathBuf};
26#[cfg(unix)]
27use std::os::unix::fs::PermissionsExt;
28use thiserror::Error;
29use tokio::io::{AsyncReadExt, AsyncWriteExt};
30use tokio::net::{TcpListener, TcpStream};
31use tokio::sync::{mpsc, oneshot, watch, Semaphore};
32use tokio::task::{JoinHandle, JoinSet};
33use tokio::time::{sleep, timeout, Instant};
34
35const MANAGED_RESUME_SETTLE_DEADLINE: Duration = Duration::from_secs(30);
36const NODE_REQUEST_IO_HEADROOM: Duration = Duration::from_secs(5);
37const NATIVE_SESSION_REQUEST_DEADLINE: Duration = Duration::from_secs(35);
38const WORKSPACE_ENTRY_CREATE_NODE_SEMANTIC_DEADLINE: Duration = Duration::from_secs(5);
39const WORKSPACE_ENTRY_CREATE_RELAY_DEADLINE: Duration =
40    WORKSPACE_ENTRY_CREATE_NODE_SEMANTIC_DEADLINE.saturating_add(NODE_REQUEST_IO_HEADROOM);
41
42const HEADER_LIMIT_BYTES: usize = 16 * 1024;
43const MAX_HTTP_CONNECTIONS: usize = 16;
44const RESPONSE_BODY_LIMIT_BYTES: usize = 8 * 1024 * 1024;
45
46mod control;
47
48#[cfg(windows)]
49pub const DEFAULT_C2_CONTROL_ENDPOINT: &str = r"\\.\pipe\gate4agent-c2";
50#[cfg(unix)]
51pub const DEFAULT_C2_CONTROL_ENDPOINT: &str = "gate4agent-c2.sock";
52
53#[cfg(windows)]
54pub fn default_c2_control_endpoint() -> Result<String, C2ConfigError> {
55    Ok(DEFAULT_C2_CONTROL_ENDPOINT.to_owned())
56}
57
58#[cfg(unix)]
59pub fn default_c2_control_endpoint() -> Result<String, C2ConfigError> {
60    let root = unix_runtime_root()?;
61    let directory = root.join("gate4agent");
62    let endpoint = directory.join(DEFAULT_C2_CONTROL_ENDPOINT);
63    validate_unix_endpoint(&endpoint).map_err(|_| C2ConfigError::InvalidControlEndpoint)?;
64    Ok(endpoint.to_string_lossy().into_owned())
65}
66
67#[cfg(unix)]
68fn unix_runtime_root() -> Result<PathBuf, C2ConfigError> {
69    if let Some(root) = std::env::var_os("XDG_RUNTIME_DIR").filter(|value| !value.is_empty()) {
70        let root = PathBuf::from(root);
71        if root.is_absolute() { return Ok(root); }
72    }
73    let home = std::env::var_os("HOME").filter(|value| !value.is_empty())
74        .ok_or_else(|| C2ConfigError::RuntimeEndpoint(
75            "neither absolute XDG_RUNTIME_DIR nor HOME is available".to_owned(),
76        ))?;
77    let home = PathBuf::from(home);
78    if !home.is_absolute() {
79        return Err(C2ConfigError::RuntimeEndpoint("HOME is not absolute".to_owned()));
80    }
81    Ok(home.join(".gate4agent").join("run"))
82}
83
84#[derive(Clone)]
85pub struct C2NodeConfig {
86    pub node_id: NodeId,
87    pub endpoint: String,
88    route: C2NodeRoute,
89    token: String,
90}
91
92#[derive(Clone, Copy, Debug, Eq, PartialEq)]
93enum C2NodeRoute {
94    Local,
95    SshForwardedLoopback(SocketAddr),
96    /// This relay does not dial the node; the node dials the relay.
97    ///
98    /// For a node the relay has no way to reach -- behind NAT, in a
99    /// container with no published port, on a machine with no listener the
100    /// relay is allowed to open. The protocol roles are unchanged: the
101    /// node is still the server, this relay is still the client asking it
102    /// for things. Only the direction of the TCP connection moves, which
103    /// is why the whole of this variant's implementation is "wait for a
104    /// socket instead of opening one" and then the identical handshake.
105    CallHome,
106}
107
108impl C2NodeConfig {
109    pub fn new(node_id: NodeId, endpoint: impl Into<String>, token: impl Into<String>) -> Result<Self, C2ConfigError> {
110        let endpoint = endpoint.into();
111        let token = token.into();
112        let (endpoint, route) = parse_node_endpoint(&endpoint)
113            .ok_or_else(|| C2ConfigError::InvalidEndpoint(node_id.clone()))?;
114        validate_token(&token)?;
115        Ok(Self { node_id, endpoint, route, token })
116    }
117
118    fn transport_label(&self) -> &'static str {
119        match self.route {
120            C2NodeRoute::Local => local_transport_label(),
121            C2NodeRoute::SshForwardedLoopback(_) => "ssh-forwarded-loopback",
122            C2NodeRoute::CallHome => "call-home",
123        }
124    }
125}
126
127#[derive(Clone, Copy, Debug)]
128pub struct C2Timings {
129    pub poll_interval: Duration,
130    pub fresh_for: Duration,
131    pub attempt_deadline: Duration,
132    pub transient_backoffs: [Duration; 5],
133    pub parked_backoff: Duration,
134    pub http_io_deadline: Duration,
135}
136
137impl Default for C2Timings {
138    fn default() -> Self {
139        Self {
140            poll_interval: Duration::from_millis(250),
141            fresh_for: Duration::from_secs(10),
142            attempt_deadline: Duration::from_secs(5),
143            transient_backoffs: [
144                Duration::from_millis(500), Duration::from_secs(1), Duration::from_secs(2),
145                Duration::from_secs(4), Duration::from_secs(8),
146            ],
147            parked_backoff: Duration::from_secs(30),
148            http_io_deadline: Duration::from_secs(3),
149        }
150    }
151}
152
153#[derive(Clone)]
154pub struct C2Config {
155    pub api_listen: SocketAddr,
156    pub control_endpoint: String,
157    api_token: String,
158    pub nodes: Vec<C2NodeConfig>,
159    /// Where nodes that cannot be dialled come to announce themselves.
160    /// `None` unless at least one node is configured `accept`.
161    pub node_listen: Option<SocketAddr>,
162    pub timings: C2Timings,
163}
164
165impl C2Config {
166    pub fn new(api_listen: SocketAddr, api_token: impl Into<String>, nodes: Vec<C2NodeConfig>) -> Result<Self, C2ConfigError> {
167        let api_token = api_token.into();
168        if !api_listen.ip().is_loopback() { return Err(C2ConfigError::NonLoopback(api_listen)); }
169        validate_token(&api_token)?;
170        if nodes.is_empty() || nodes.len() > MAX_C2_NODES { return Err(C2ConfigError::NodeCount(nodes.len())); }
171        let mut ids = BTreeSet::new();
172        let mut endpoints = BTreeSet::new();
173        if nodes.iter().any(|node| !ids.insert(node.node_id.clone())) { return Err(C2ConfigError::DuplicateNode); }
174        // Call-home nodes are exempt: they name no endpoint, so every one
175        // of them carries the same `accept` and the uniqueness this check
176        // enforces is `node_id`'s job for them (already done above). The
177        // check exists to stop two dialled nodes pointing at one socket,
178        // which cannot happen to a node that is never dialled.
179        if nodes
180            .iter()
181            .filter(|node| node.route != C2NodeRoute::CallHome)
182            .any(|node| !endpoints.insert(endpoint_key(&node.endpoint)))
183        {
184            return Err(C2ConfigError::DuplicateEndpoint);
185        }
186        let control_endpoint = default_c2_control_endpoint()?;
187        if nodes.iter().any(|node| endpoints_equal(&node.endpoint, &control_endpoint)) {
188            return Err(C2ConfigError::ControlEndpointConflict);
189        }
190        Ok(Self {
191            api_listen,
192            control_endpoint,
193            api_token,
194            nodes,
195            node_listen: None,
196            timings: C2Timings::default(),
197        })
198    }
199
200    /// Binds the address nodes call in on.
201    ///
202    /// Loopback only, deliberately, and for the same reason every other
203    /// listener in this stack is: the node wire authenticates both ends
204    /// with a mutual challenge-response but encrypts nothing, so its
205    /// frames -- terminal contents, keystrokes, file bytes -- are
206    /// plaintext JSON. Accepting node connections from off-box would put
207    /// all of that on the network. Direction and distance are separate
208    /// problems and this change only solves direction.
209    pub fn with_node_listen(mut self, node_listen: SocketAddr) -> Result<Self, C2ConfigError> {
210        if !node_listen.ip().is_loopback() || node_listen.port() == 0 {
211            return Err(C2ConfigError::NonLoopbackNodeListen(node_listen));
212        }
213        if node_listen == self.api_listen {
214            return Err(C2ConfigError::NodeListenConflict);
215        }
216        self.node_listen = Some(node_listen);
217        Ok(self)
218    }
219
220    /// Every node that waits to be called needs somewhere to call, so a
221    /// config that asks for one without the other is refused rather than
222    /// started into a relay that can never reach that node. Checked at
223    /// startup, not at connect time, because the failure is total and
224    /// permanent -- there is nothing to retry.
225    pub fn validate_call_home(&self) -> Result<(), C2ConfigError> {
226        let waiting = self.nodes.iter().find(|node| node.route == C2NodeRoute::CallHome);
227        match (waiting, self.node_listen) {
228            (Some(node), None) => Err(C2ConfigError::CallHomeWithoutListener(node.node_id.clone())),
229            _ => Ok(()),
230        }
231    }
232
233    pub fn with_timings(mut self, timings: C2Timings) -> Self {
234        self.timings = timings;
235        self
236    }
237
238    pub fn with_control_endpoint(mut self, endpoint: impl Into<String>) -> Result<Self, C2ConfigError> {
239        let endpoint = endpoint.into();
240        validate_control_endpoint(&endpoint)?;
241        if self.nodes.iter().any(|node| endpoints_equal(&node.endpoint, &endpoint)) {
242            return Err(C2ConfigError::ControlEndpointConflict);
243        }
244        self.control_endpoint = endpoint;
245        Ok(self)
246    }
247}
248
249fn validate_control_endpoint(endpoint: &str) -> Result<(), C2ConfigError> {
250    if !valid_local_endpoint(endpoint) {
251        return Err(C2ConfigError::InvalidControlEndpoint);
252    }
253    Ok(())
254}
255
256fn parse_node_endpoint(endpoint: &str) -> Option<(String, C2NodeRoute)> {
257    // The one assignment that names no address, because there is nothing
258    // to address: this node will arrive on the call-home listener under
259    // its own name. Spelled as a word rather than an empty value so a
260    // config reader can tell "waits to be called" from "somebody forgot to
261    // fill this in".
262    if endpoint == "accept" {
263        return Some((endpoint.to_owned(), C2NodeRoute::CallHome));
264    }
265    if let Some(authority) = endpoint.strip_prefix("tcp://") {
266        let address = authority.parse::<SocketAddr>().ok()?;
267        let is_exact_loopback = match address.ip() {
268            std::net::IpAddr::V4(ip) => ip == std::net::Ipv4Addr::LOCALHOST,
269            std::net::IpAddr::V6(ip) => ip == std::net::Ipv6Addr::LOCALHOST,
270        };
271        if !is_exact_loopback || address.port() == 0 {
272            return None;
273        }
274        return Some((format!("tcp://{address}"), C2NodeRoute::SshForwardedLoopback(address)));
275    }
276    if endpoint.contains("://") || !valid_local_endpoint(endpoint) {
277        return None;
278    }
279    Some((endpoint.to_owned(), C2NodeRoute::Local))
280}
281
282#[cfg(windows)]
283fn valid_local_endpoint(endpoint: &str) -> bool {
284    endpoint.starts_with(r"\\.\pipe\") && endpoint.len() > r"\\.\pipe\".len()
285        && endpoint.len() <= MAX_C2_ENDPOINT_BYTES
286}
287
288#[cfg(unix)]
289fn valid_local_endpoint(endpoint: &str) -> bool {
290    validate_unix_endpoint(Path::new(endpoint)).is_ok()
291}
292
293#[cfg(unix)]
294fn validate_unix_endpoint(endpoint: &Path) -> Result<(), ()> {
295    const MAX_UNIX_ENDPOINT_BYTES: usize = 103;
296    use std::os::unix::ffi::OsStrExt;
297
298    if !endpoint.is_absolute() || endpoint.file_name().is_none()
299        || endpoint.as_os_str().as_bytes().len() > MAX_UNIX_ENDPOINT_BYTES
300    {
301        return Err(());
302    }
303    Ok(())
304}
305
306#[cfg(windows)]
307fn endpoint_key(endpoint: &str) -> String {
308    if endpoint.starts_with("tcp://") { endpoint.to_owned() } else { endpoint.to_ascii_lowercase() }
309}
310
311#[cfg(unix)]
312fn endpoint_key(endpoint: &str) -> String { endpoint.to_owned() }
313
314fn endpoints_equal(left: &str, right: &str) -> bool { endpoint_key(left) == endpoint_key(right) }
315
316fn validate_token(token: &str) -> Result<(), C2ConfigError> {
317    if token.is_empty() || token.len() > 4096 || !token.bytes().all(|byte| matches!(byte, 0x21..=0x7e)) {
318        return Err(C2ConfigError::InvalidToken);
319    }
320    Ok(())
321}
322
323#[derive(Debug, Error)]
324pub enum C2ConfigError {
325    #[error("C2 tokens must contain 1..=4096 visible ASCII bytes without whitespace")]
326    InvalidToken,
327    #[cfg_attr(windows, error("node '{0}' requires a bounded Windows named pipe or exact loopback TCP endpoint"))]
328    #[cfg_attr(unix, error("node '{0}' requires a bounded local socket or exact loopback TCP endpoint"))]
329    InvalidEndpoint(NodeId),
330    #[error("C2 API listen address must be loopback: {0}")]
331    NonLoopback(SocketAddr),
332    #[error("C2 requires 1..=64 configured nodes; received {0}")]
333    NodeCount(usize),
334    #[error("C2 node IDs must be unique")]
335    DuplicateNode,
336    #[error("C2 node endpoints must be unique")]
337    DuplicateEndpoint,
338    #[cfg_attr(windows, error("C2 control endpoint must be a bounded local Windows named pipe"))]
339    #[cfg_attr(unix, error("C2 control endpoint must be a bounded local endpoint"))]
340    InvalidControlEndpoint,
341    #[error("C2 control endpoint must not equal a configured node endpoint")]
342    ControlEndpointConflict,
343    #[error("C2 node call-home listen address must be loopback with a nonzero port: {0}")]
344    NonLoopbackNodeListen(SocketAddr),
345    #[error("C2 node call-home listen address must not equal the API listen address")]
346    NodeListenConflict,
347    #[error("node '{0}' waits to be called but no --node-listen address was configured")]
348    CallHomeWithoutListener(NodeId),
349    #[cfg(unix)]
350    #[error("C2 default runtime endpoint is unavailable: {0}")]
351    RuntimeEndpoint(String),
352}
353
354type RelayResult = Result<RoutedNodeResponse, C2RelayFailure>;
355
356enum RelayCommand {
357    Request {
358        operator_connection_id: u64,
359        expected_incarnation_id: NodeIncarnationId,
360        request: NodeRequest,
361        reply: oneshot::Sender<RelayResult>,
362    },
363}
364
365#[derive(Clone)]
366struct RelayEndpoint {
367    commands: mpsc::Sender<RelayCommand>,
368    releases: mpsc::Sender<oneshot::Sender<()>>,
369    force_disconnect: watch::Sender<u64>,
370}
371
372#[derive(Clone)]
373struct OperatorHub {
374    sink: Arc<Mutex<Option<OperatorEventSink>>>,
375}
376
377#[derive(Clone)]
378struct OperatorEventSink {
379    connection_id: u64,
380    outbound: mpsc::Sender<control::QueuedFrame>,
381    budget: Arc<AtomicUsize>,
382    disconnect: watch::Sender<bool>,
383}
384
385impl OperatorHub {
386    fn new() -> Self { Self { sink: Arc::new(Mutex::new(None)) } }
387
388    fn attach(&self, sink: OperatorEventSink) {
389        *self.sink.lock().unwrap_or_else(|poisoned| poisoned.into_inner()) = Some(sink);
390    }
391
392    fn detach(&self, connection_id: u64) {
393        let mut sink = self.sink.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
394        if sink.as_ref().is_some_and(|current| current.connection_id == connection_id) { *sink = None; }
395    }
396
397    fn is_active(&self, connection_id: u64) -> bool {
398        self.sink.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
399            .as_ref().is_some_and(|current| current.connection_id == connection_id)
400    }
401
402    fn has_active_operator(&self) -> bool {
403        self.sink.lock().unwrap_or_else(|poisoned| poisoned.into_inner()).is_some()
404    }
405
406    fn publish(&self, event: RoutedNodeEvent) {
407        let sink = self.sink.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
408        if let Some(sink) = sink.as_ref() {
409            if control::queue_operator_event(&sink.outbound, &sink.budget, event).is_err() {
410                let _ = sink.disconnect.send(true);
411            }
412        }
413    }
414}
415
416fn relay_failure(
417    code: C2RelayFailureCode,
418    message: &'static str,
419    current_incarnation_id: Option<NodeIncarnationId>,
420) -> C2RelayFailure {
421    C2RelayFailure { code, message: message.to_owned(), current_incarnation_id }
422}
423
424#[derive(Debug, Error)]
425pub enum C2Error {
426    #[error("C2 API failed: {0}")]
427    Api(#[from] io::Error),
428    #[error("C2 task failed: {0}")]
429    Task(#[from] tokio::task::JoinError),
430}
431
432#[derive(Clone)]
433pub struct C2ShutdownHandle {
434    shutdown: watch::Sender<bool>,
435}
436
437impl C2ShutdownHandle {
438    pub fn shutdown(&self) { let _ = self.shutdown.send(true); }
439}
440
441pub struct C2Running {
442    api_addr: SocketAddr,
443    shutdown: C2ShutdownHandle,
444    task: Option<JoinHandle<Result<(), C2Error>>>,
445}
446
447impl C2Running {
448    pub async fn start(config: C2Config) -> Result<Self, C2Error> {
449        prepare_default_control_parent(&config.control_endpoint)?;
450        let listener = TcpListener::bind(config.api_listen).await?;
451        let api_addr = listener.local_addr()?;
452        let (shutdown_tx, shutdown_rx) = watch::channel(false);
453        let shutdown = C2ShutdownHandle { shutdown: shutdown_tx };
454        let task = tokio::spawn(run_bound(config, listener, shutdown_rx));
455        Ok(Self { api_addr, shutdown, task: Some(task) })
456    }
457
458    pub fn api_addr(&self) -> SocketAddr { self.api_addr }
459    pub fn shutdown_handle(&self) -> C2ShutdownHandle { self.shutdown.clone() }
460    pub async fn wait(mut self) -> Result<(), C2Error> {
461        self.task.take().expect("C2 task is present").await?
462    }
463}
464
465#[cfg(windows)]
466fn prepare_default_control_parent(_endpoint: &str) -> io::Result<()> { Ok(()) }
467
468#[cfg(unix)]
469fn prepare_default_control_parent(endpoint: &str) -> io::Result<()> {
470    let default = default_c2_control_endpoint()
471        .map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error.to_string()))?;
472    if endpoint != default { return Ok(()); }
473    let parent = Path::new(endpoint).parent().ok_or_else(|| {
474        io::Error::new(io::ErrorKind::InvalidInput, "default C2 endpoint has no parent")
475    })?;
476    std::fs::create_dir_all(parent)?;
477    std::fs::set_permissions(parent, std::fs::Permissions::from_mode(0o700))
478}
479
480impl Drop for C2Running {
481    fn drop(&mut self) {
482        self.shutdown.shutdown();
483        if let Some(task) = self.task.take() { task.abort(); }
484    }
485}
486
487async fn run_bound(config: C2Config, listener: TcpListener, mut shutdown: watch::Receiver<bool>) -> Result<(), C2Error> {
488    let now = unix_ms();
489    let nodes = config.nodes.iter().map(|node| {
490        (node.node_id.clone(), initial_observed_node(node))
491    }).collect();
492    let initial = Arc::new(StatusResponse { api_version: C2_API_VERSION, ready: false, observed_at_unix_ms: now, nodes });
493    let (status_tx, status_rx) = watch::channel(initial);
494    let (ingress_tx, ingress_rx) = mpsc::channel(config.nodes.len().saturating_mul(2).max(2));
495    let mut tasks = JoinSet::new();
496    let hub = OperatorHub::new();
497    let mut relay_senders = BTreeMap::new();
498    let mut relay_receivers = Vec::new();
499    // Where the call-home listener hands a freshly announced socket. One
500    // slot per waiting node and no more: a node has one live connection to
501    // this relay, so a second socket arriving while the first is still
502    // being served means the node reconnected and the relay has not
503    // noticed yet -- queueing more of those would only serve them stale.
504    let mut call_home_routes: BTreeMap<NodeId, mpsc::Sender<TcpStream>> = BTreeMap::new();
505    for node in config.nodes.clone() {
506        let (commands_tx, commands_rx) = mpsc::channel(8);
507        let (releases_tx, releases_rx) = mpsc::channel(1);
508        let (force_tx, force_rx) = watch::channel(0_u64);
509        relay_senders.insert(node.node_id.clone(), RelayEndpoint {
510            commands: commands_tx,
511            releases: releases_tx,
512            force_disconnect: force_tx,
513        });
514        let call_home_rx = if node.route == C2NodeRoute::CallHome {
515            let (call_home_tx, call_home_rx) = mpsc::channel(1);
516            call_home_routes.insert(node.node_id.clone(), call_home_tx);
517            Some(call_home_rx)
518        } else {
519            None
520        };
521        relay_receivers.push((node, commands_rx, releases_rx, force_rx, call_home_rx));
522    }
523    let relay_senders = Arc::new(relay_senders);
524    if let Some(node_listen) = config.node_listen {
525        tasks.spawn(accept_call_home(node_listen, call_home_routes, shutdown.clone()));
526    }
527    tasks.spawn(inventory_owner(config.nodes.len(), config.timings.fresh_for, ingress_rx, status_tx, shutdown.clone()));
528    for (node, commands, releases, force_disconnect, call_home) in relay_receivers {
529        tasks.spawn(node_relay_worker(node, config.timings, commands, releases, force_disconnect, call_home, ingress_tx.clone(), status_rx.clone(), hub.clone(), shutdown.clone()));
530    }
531    drop(ingress_tx);
532    tasks.spawn(http_server(listener, config.api_token.clone(), config.timings.http_io_deadline, status_rx.clone(), shutdown.clone()));
533    tasks.spawn(control::run(
534        config.control_endpoint,
535        config.api_token,
536        relay_senders,
537        status_rx,
538        hub,
539        shutdown.clone(),
540    ));
541    loop {
542        tokio::select! {
543            changed = shutdown.changed() => {
544                if changed.is_err() || *shutdown.borrow() { break; }
545            }
546            result = tasks.join_next() => {
547                match result {
548                    Some(Ok(Ok(()))) if *shutdown.borrow() => break,
549                    Some(Ok(Ok(()))) => return Err(C2Error::Api(io::Error::new(io::ErrorKind::Other, "C2 task exited unexpectedly"))),
550                    Some(Ok(Err(error))) => return Err(C2Error::Api(error)),
551                    Some(Err(error)) => return Err(C2Error::Task(error)),
552                    None => break,
553                }
554            }
555        }
556    }
557    tasks.shutdown().await;
558    Ok(())
559}
560
561fn initial_observed_node(node: &C2NodeConfig) -> ObservedNode {
562    ObservedNode {
563        endpoint: node.endpoint.clone(), transport_label: node.transport_label().to_owned(),
564        transport: NodeTransportState::Offline, freshness: NodeFreshness::Unavailable,
565        cursor: None, inventory: None, last_attempt_unix_ms: None, last_success_unix_ms: None,
566        consecutive_failures: 0, last_error: None, gaps: Vec::new(), gaps_truncated: 0,
567    }
568}
569
570#[cfg(windows)]
571fn local_transport_label() -> &'static str { "windows-named-pipe" }
572
573#[cfg(unix)]
574fn local_transport_label() -> &'static str { "unix-domain-socket" }
575
576#[derive(Clone, Default)]
577struct ProviderContractManifest {
578    provider_contracts: Vec<ProviderContractSupport>,
579    provider_adapter_contracts: Vec<ProviderAdapterContractSupport>,
580}
581
582impl ProviderContractManifest {
583    fn from_compatibility(compatibility: Option<&NegotiatedNodeCompatibility>) -> Self {
584        let Some(compatibility) = compatibility.filter(|compatibility| {
585            compatibility.capabilities.iter().any(|capability| {
586                capability.as_str() == C2_PROVIDER_CONTRACT_MANIFEST_CAPABILITY
587            })
588        }) else {
589            return Self::default();
590        };
591        Self {
592            provider_contracts: compatibility.provider_contracts.clone(),
593            provider_adapter_contracts: compatibility.provider_adapter_contracts.clone(),
594        }
595    }
596}
597
598enum AttemptResult {
599    Connected {
600        cursor: NodeCursor,
601        snapshot: NodeSnapshot,
602        gaps: Vec<GapKind>,
603        provider_contract_manifest: ProviderContractManifest,
604    },
605    Success { cursor: NodeCursor, snapshot: NodeSnapshot, gaps: Vec<GapKind> },
606    Cursor {
607        cursor: NodeCursor,
608        gaps: Vec<GapKind>,
609        managed_worktree_events: Vec<NodeEvent>,
610    },
611    Failure { error: SanitizedError, hard: bool },
612}
613
614struct Attempt { node_id: NodeId, at_unix_ms: u64, result: AttemptResult }
615
616/// Accepts node connections and hands each one to the relay worker for
617/// whichever node it says it is.
618///
619/// This is the only place in the relay that learns a node's identity from
620/// the wire rather than from configuration, and it is careful about what
621/// that means: the announced name SELECTS a worker and nothing more. The
622/// worker then runs the same mutual challenge-response it runs on a
623/// dialled connection, against the token configured for that node, so a
624/// caller that announced a name it cannot prove gets no further than the
625/// next frame.
626///
627/// A socket for an unknown node, an unreadable preface, or a node whose
628/// slot is already full is dropped with a line saying which -- never
629/// silently, because "the node is calling but nothing happens" is exactly
630/// the failure that is impossible to diagnose from the node's end.
631async fn accept_call_home(
632    listen: SocketAddr,
633    routes: BTreeMap<NodeId, mpsc::Sender<TcpStream>>,
634    mut shutdown: watch::Receiver<bool>,
635) -> io::Result<()> {
636    let listener = TcpListener::bind(listen).await?;
637    tracing::info!(%listen, waiting_nodes = routes.len(), "call-home listener bound");
638    loop {
639        if *shutdown.borrow() {
640            return Ok(());
641        }
642        let (stream, peer) = tokio::select! {
643            accepted = listener.accept() => accepted?,
644            changed = shutdown.changed() => {
645                if changed.is_err() || *shutdown.borrow() { return Ok(()); }
646                continue;
647            }
648        };
649        let routes = routes.clone();
650        // Off the accept loop: the preface has its own deadline, and one
651        // peer that connects and then says nothing must not stop every
652        // other node from being accepted meanwhile.
653        tokio::spawn(async move {
654            let mut stream = stream;
655            let _ = stream.set_nodelay(true);
656            let node_id = match read_call_home_announce(&mut stream).await {
657                Ok(node_id) => node_id,
658                Err(error) => {
659                    tracing::warn!(%peer, cause = %error, "call-home connection rejected");
660                    return;
661                }
662            };
663            let Some(sender) = routes.get(&node_id) else {
664                tracing::warn!(
665                    %peer,
666                    %node_id,
667                    "call-home connection named a node this relay does not wait for",
668                );
669                return;
670            };
671            match sender.try_send(stream) {
672                Ok(()) => tracing::info!(%peer, %node_id, "node called home"),
673                Err(mpsc::error::TrySendError::Full(_)) => tracing::warn!(
674                    %peer,
675                    %node_id,
676                    "node called home while its previous connection is still being taken up",
677                ),
678                Err(mpsc::error::TrySendError::Closed(_)) => tracing::warn!(
679                    %peer,
680                    %node_id,
681                    "node called home but its relay worker is gone",
682                ),
683            }
684        });
685    }
686}
687
688/// Blocks until this node calls in, or shutdown -- with no deadline of its
689/// own on purpose.
690///
691/// A dialled connection either answers within `attempt_deadline` or has
692/// failed, and treating a slow answer as a failure is right there. Waiting
693/// for a node to call is the opposite: a node that has not called yet has
694/// not failed at anything, and it may be minutes from starting. Putting
695/// that wait under the attempt deadline would turn an idle relay into a
696/// stream of "connection deadline exceeded" warnings about nodes that are
697/// merely not running. The handshake that FOLLOWS is deadlined normally --
698/// a peer that connects and then stalls is a real failure.
699async fn await_call_home(
700    receiver: Option<&mut mpsc::Receiver<TcpStream>>,
701    shutdown: &mut watch::Receiver<bool>,
702) -> Option<TcpStream> {
703    let receiver = receiver?;
704    loop {
705        if *shutdown.borrow() {
706            return None;
707        }
708        tokio::select! {
709            stream = receiver.recv() => return stream,
710            changed = shutdown.changed() => {
711                if changed.is_err() || *shutdown.borrow() { return None; }
712            }
713        }
714    }
715}
716
717async fn node_relay_worker(
718    node: C2NodeConfig,
719    timings: C2Timings,
720    mut commands: mpsc::Receiver<RelayCommand>,
721    mut releases: mpsc::Receiver<oneshot::Sender<()>>,
722    mut force_disconnect: watch::Receiver<u64>,
723    mut call_home: Option<mpsc::Receiver<TcpStream>>,
724    ingress: mpsc::Sender<Attempt>,
725    status: watch::Receiver<Arc<StatusResponse>>,
726    hub: OperatorHub,
727    mut shutdown: watch::Receiver<bool>,
728) -> io::Result<()> {
729    let mut failures = 0_usize;
730    loop {
731        if *shutdown.borrow() { return Ok(()); }
732        let previous = status.borrow().nodes.get(&node.node_id).and_then(|item| item.cursor);
733        // A waiting node's socket is collected BEFORE the deadline starts,
734        // because waiting for a node to call is not an attempt that can
735        // time out -- see `await_call_home`. Once it is here, the
736        // handshake over it is deadlined exactly like a dialled one.
737        let adopted = if node.route == C2NodeRoute::CallHome {
738            match await_call_home(call_home.as_mut(), &mut shutdown).await {
739                Some(stream) => Some(stream),
740                None => return Ok(()),
741            }
742        } else {
743            None
744        };
745        let connected = timeout(timings.attempt_deadline, connect_operator(&node, adopted)).await;
746        let mut client = match connected {
747            Ok(Ok(client)) => client,
748            Ok(Err(error)) => {
749                failures = failures.saturating_add(1);
750                let (error, hard) = sanitize_node_error(&error);
751                tracing::warn!(
752                    node_id = %node.node_id,
753                    cause = %error.message,
754                    category = ?error.category,
755                    hard,
756                    "node connection attempt failed",
757                );
758                ingress_attempt(&ingress, &node.node_id, AttemptResult::Failure { error, hard }).await?;
759                reject_disconnected_commands(&mut commands, previous);
760                acknowledge_disconnected_releases(&mut releases);
761                relay_backoff(&mut shutdown, timings, failures, hard).await?;
762                continue;
763            }
764            Err(_) => {
765                failures = failures.saturating_add(1);
766                tracing::warn!(
767                    node_id = %node.node_id,
768                    cause = "node connection deadline exceeded",
769                    "node connection attempt failed",
770                );
771                ingress_attempt(&ingress, &node.node_id, AttemptResult::Failure {
772                    error: SanitizedError { category: C2ErrorCategory::Timeout, message: "node connection deadline exceeded".to_owned() },
773                    hard: false,
774                }).await?;
775                reject_disconnected_commands(&mut commands, previous);
776                acknowledge_disconnected_releases(&mut releases);
777                relay_backoff(&mut shutdown, timings, failures, false).await?;
778                continue;
779            }
780        };
781        failures = 0;
782        let hello = client.hello().clone();
783        let provider_contract_manifest =
784            ProviderContractManifest::from_compatibility(hello.compatibility.as_ref());
785        let incarnation_id = hello.incarnation_id;
786        let connection_id = hello.connection_id;
787        tracing::info!(
788            node_id = %node.node_id,
789            connection_id,
790            incarnation_id = ?incarnation_id,
791            "node attached to relay",
792        );
793        let mut controller_owned = hello.controller.as_ref().is_some_and(|controller| controller.connection_id == connection_id);
794        let mut cursor = NodeCursor { incarnation_id, sequence: hello.event_sequence };
795        let mut snapshot = hello.snapshot;
796        let mut gaps = Vec::new();
797        let mut did_resync = false;
798        if let Some(previous) = previous {
799            if previous.incarnation_id != incarnation_id {
800                gaps.push(GapKind::IncarnationChanged);
801            } else if cursor.sequence < previous.sequence {
802                gaps.push(GapKind::CursorRegression);
803            } else if cursor.sequence > previous.sequence {
804                match bounded_node_request(&mut client, NodeRequest::Resync { after_sequence: previous.sequence }).await {
805                    Ok(NodeResponse::Resync {
806                        event_sequence,
807                        oldest_available_sequence,
808                        snapshot: resync_snapshot,
809                        events,
810                    }) => {
811                        let resync_gaps = validate_resync(
812                            previous.sequence,
813                            cursor.sequence,
814                            event_sequence,
815                            oldest_available_sequence,
816                            &events,
817                        );
818                        if resync_gaps.is_empty() {
819                            publish_recovered_events(&node.node_id, incarnation_id, &events, &hub);
820                        } else {
821                            hub.publish(RoutedNodeEvent {
822                                node_id: node.node_id.clone(),
823                                cursor: NodeCursor { incarnation_id, sequence: event_sequence },
824                                event: C2NodeEvent::ResyncRequired {
825                                    oldest_available_sequence,
826                                },
827                            });
828                        }
829                        gaps.extend(resync_gaps);
830                        cursor.sequence = event_sequence;
831                        snapshot = resync_snapshot;
832                        did_resync = true;
833                    }
834                    Ok(_) => gaps.push(GapKind::NonContiguousEvents),
835                    Err(error) => {
836                        ingress_attempt(&ingress, &node.node_id, relay_failure_attempt(&error)).await?;
837                        reject_disconnected_commands(&mut commands, Some(cursor));
838                        continue;
839                    }
840                }
841            }
842        }
843        ingress_attempt(&ingress, &node.node_id, AttemptResult::Connected {
844            cursor,
845            snapshot,
846            gaps,
847            provider_contract_manifest,
848        }).await?;
849        if let Err(error) = drain_pending_events(&mut client, &node.node_id, &mut cursor, &hub, &ingress, did_resync, None).await {
850            ingress_attempt(&ingress, &node.node_id, relay_failure_attempt(&error)).await?;
851            reject_disconnected_commands(&mut commands, Some(cursor));
852            continue;
853        }
854
855        let cadence = timings.poll_interval.max(Duration::from_millis(1)).min(Duration::from_millis(250));
856        let mut snapshot_tick = tokio::time::interval(cadence);
857        snapshot_tick.reset();
858        let mut lease_tick = tokio::time::interval(Duration::from_secs(30));
859        lease_tick.reset();
860        let disconnect_error = loop {
861            tokio::select! {
862                command = commands.recv() => {
863                    let Some(command) = command else { return Ok(()); };
864                    tokio::select! {
865                        result = handle_relay_command(
866                            &mut client, command, &node.node_id, incarnation_id, connection_id,
867                            &mut controller_owned, &hub, &mut cursor, &ingress,
868                        ) => match result {
869                            Ok(()) => {}
870                            Err(error) => break error,
871                        },
872                        changed = force_disconnect.changed() => {
873                            let _ = changed;
874                            break NodeClientError::Io(io::Error::new(io::ErrorKind::ConnectionAborted, "node relay cleanup forced reconnect"));
875                        }
876                    }
877                }
878                release = releases.recv() => {
879                    let Some(reply) = release else { return Ok(()); };
880                    let result = release_controller(&mut client, &mut controller_owned).await;
881                    let _ = reply.send(());
882                    if let Err(error) = result { break error; }
883                }
884                changed = force_disconnect.changed() => {
885                    let _ = changed;
886                    break NodeClientError::Io(io::Error::new(io::ErrorKind::ConnectionAborted, "node relay cleanup forced reconnect"));
887                }
888                frame = client.recv() => {
889                    match frame {
890                        Ok(ServerFrame::Event(envelope)) => {
891                            if let Err(error) = handle_live_node_event(
892                                &mut client,
893                                &node.node_id,
894                                envelope,
895                                &mut cursor,
896                                &hub,
897                                &ingress,
898                            ).await {
899                                break error;
900                            }
901                        }
902                        Ok(ServerFrame::Reply(_)
903                            | ServerFrame::Challenge(_)
904                            | ServerFrame::Hello(_)) => {
905                            break NodeClientError::Protocol(
906                                "node sent an unexpected idle frame".to_owned(),
907                            );
908                        }
909                        Err(error) => break error,
910                    }
911                }
912                _ = snapshot_tick.tick() => {
913                    match bounded_node_request(&mut client, NodeRequest::Snapshot).await {
914                        Ok(NodeResponse::Snapshot { event_sequence, snapshot, .. }) => {
915                            if let Err(error) = drain_pending_events(&mut client, &node.node_id, &mut cursor, &hub, &ingress, false, None).await { break error; }
916                            if event_sequence >= cursor.sequence {
917                                cursor.sequence = event_sequence;
918                                ingress_attempt(&ingress, &node.node_id, AttemptResult::Success { cursor, snapshot, gaps: Vec::new() }).await?;
919                            }
920                        }
921                        Ok(_) => break NodeClientError::Protocol("snapshot returned a different response".to_owned()),
922                        Err(error) => break error,
923                    }
924                }
925                _ = lease_tick.tick() => {
926                    if controller_owned {
927                        if hub.has_active_operator() {
928                            match acquire_controller(&mut client, connection_id).await {
929                                Ok(owned) => controller_owned = owned,
930                                Err(error) => break error,
931                            }
932                        } else if let Err(error) = release_controller(&mut client, &mut controller_owned).await {
933                            break error;
934                        }
935                        if let Err(error) = drain_pending_events(&mut client, &node.node_id, &mut cursor, &hub, &ingress, false, None).await { break error; }
936                    }
937                }
938                changed = shutdown.changed() => {
939                    if changed.is_err() || *shutdown.borrow() {
940                        let _ = release_controller(&mut client, &mut controller_owned).await;
941                        return Ok(());
942                    }
943                }
944            }
945        };
946        let (error, hard) = sanitize_node_error(&disconnect_error);
947        tracing::warn!(
948            node_id = %node.node_id,
949            connection_id,
950            cause = %error.message,
951            category = ?error.category,
952            hard,
953            "node dropped from relay",
954        );
955        ingress_attempt(&ingress, &node.node_id, AttemptResult::Failure { error, hard }).await?;
956        reject_disconnected_commands(&mut commands, Some(cursor));
957        failures = failures.saturating_add(1);
958        relay_backoff(&mut shutdown, timings, failures, hard).await?;
959    }
960}
961
962/// Produces an authenticated operator client for `node`, whichever end
963/// opened the socket.
964///
965/// All three arms end in the same handshake. The first two open a
966/// connection and hand it straight to it; the third is handed one that a
967/// node opened and the listener already matched to this node's name. That
968/// symmetry is the point of the change: `LocalNodeClient::adopt` is what
969/// `connect`/`connect_loopback` both do after their own `connect` call, so
970/// a call-home node is authenticated by exactly the same code, against
971/// exactly the same per-node token, as one this relay dialled.
972async fn connect_operator(
973    node: &C2NodeConfig,
974    adopted: Option<TcpStream>,
975) -> Result<LocalNodeClient, NodeClientError> {
976    match node.route {
977        C2NodeRoute::Local => {
978            LocalNodeClient::connect(&node.endpoint, &node.node_id, ClientRole::Operator, &node.token).await
979        }
980        C2NodeRoute::SshForwardedLoopback(endpoint) => {
981            LocalNodeClient::connect_loopback(endpoint, &node.node_id, ClientRole::Operator, &node.token).await
982        }
983        C2NodeRoute::CallHome => {
984            // `node_relay_worker` collects the socket before calling this
985            // and returns rather than calling without one, so `None` here
986            // is a caller bug, not a runtime condition -- reported instead
987            // of panicking because a relay worker is not worth aborting a
988            // whole C2 over.
989            let stream = adopted.ok_or_else(|| {
990                NodeClientError::Protocol(
991                    "call-home node reached the connect step with no adopted stream".to_owned(),
992                )
993            })?;
994            LocalNodeClient::adopt(stream, &node.node_id, ClientRole::Operator, &node.token).await
995        }
996    }
997}
998
999async fn bounded_node_request(
1000    client: &mut LocalNodeClient,
1001    request: NodeRequest,
1002) -> Result<NodeResponse, NodeClientError> {
1003    bounded_node_request_with_deadline(client, request, None).await
1004}
1005
1006async fn bounded_node_request_with_deadline(
1007    client: &mut LocalNodeClient,
1008    request: NodeRequest,
1009    relay_deadline: Option<Instant>,
1010) -> Result<NodeResponse, NodeClientError> {
1011    let deadline = request_budget(&request, relay_deadline, Instant::now());
1012    if deadline.is_zero() {
1013        return Err(NodeClientError::Frame(FrameError::PrefixTimedOut));
1014    }
1015    match timeout(deadline, client.request(request)).await {
1016        Ok(result) => result,
1017        Err(_) => Err(NodeClientError::Frame(FrameError::PrefixTimedOut)),
1018    }
1019}
1020
1021fn relay_request_deadline(request: &NodeRequest, now: Instant) -> Option<Instant> {
1022    matches!(request, NodeRequest::ResumeSessionRecord { .. })
1023        .then(|| now + node_request_deadline(request))
1024}
1025
1026fn request_budget(request: &NodeRequest, relay_deadline: Option<Instant>, now: Instant) -> Duration {
1027    relay_deadline
1028        .map(|deadline| deadline.saturating_duration_since(now))
1029        .unwrap_or_else(|| node_request_deadline(request))
1030}
1031
1032/// Relay bound for a workspace inspection. The node's own inspection
1033/// budget is 8s by default and 11s at most; this must clear that maximum
1034/// plus the round trip, or the relay kills a request the node was still
1035/// legitimately working on -- and killing it drops the node, not just the
1036/// request. See `node_request_deadline`'s own `InspectWorkspace` arm.
1037const WORKSPACE_INSPECTION_RELAY_DEADLINE: Duration = Duration::from_secs(15);
1038
1039fn node_request_deadline(request: &NodeRequest) -> Duration {
1040    match request {
1041        NodeRequest::Snapshot
1042        | NodeRequest::Resync { .. }
1043        | NodeRequest::ArmHarnessMcpReservation { .. }
1044        | NodeRequest::ActivateHarnessMcpReservation { .. }
1045        | NodeRequest::AbortHarnessMcpReservation { .. }
1046        | NodeRequest::PutHarnessMcpReplyChunk { .. }
1047        | NodeRequest::RejectHarnessMcpCall { .. }
1048        | NodeRequest::BrowseHostDirectories { .. }
1049        | NodeRequest::ReadWorkspaceFile { .. }
1050        | NodeRequest::WriteWorkspaceFile { .. }
1051        | NodeRequest::ReadGitHistory { .. }
1052        | NodeRequest::ReadGitDiff { .. }
1053        | NodeRequest::AcquireController { .. }
1054        | NodeRequest::ReleaseController
1055        | NodeRequest::RenameSessionRecord { .. }
1056        | NodeRequest::SetSessionTask { .. }
1057        | NodeRequest::ForgetSessionRecord { .. } => Duration::from_secs(5),
1058        // The node gives its own workspace inspection 8s by default and up
1059        // to 11s, and says so when it uses them
1060        // (`git_time_budget_exceeded`, elapsed over eight seconds on a
1061        // real repository here). Bounding it from out here at five made
1062        // that inner budget unreachable: any workspace big enough to spend
1063        // its own allowance timed out at the relay every single time.
1064        //
1065        // And a relay timeout is not a slow answer, it is a dead node --
1066        // the node is dropped from the relay, its controller lease is
1067        // released, and every read behind it starts answering
1068        // "unavailable" until it reattaches, which it then does, and the
1069        // cycle repeats. One number two seconds too small took the whole
1070        // stack down in a loop: no inventory, so a spawned session could
1071        // never be shown, while its process ran perfectly well.
1072        //
1073        // The outer bound must therefore clear the node's own MAXIMUM, not
1074        // its default, with room for the round trip on top.
1075        NodeRequest::InspectWorkspace { .. } => WORKSPACE_INSPECTION_RELAY_DEADLINE,
1076        NodeRequest::CreateWorkspaceFile { .. }
1077        | NodeRequest::CreateWorkspaceDirectory { .. } => {
1078            WORKSPACE_ENTRY_CREATE_RELAY_DEADLINE
1079        }
1080        NodeRequest::CatalogNativeSessions { .. }
1081        | NodeRequest::PageNativeSessions { .. }
1082        | NodeRequest::PreviewNativeSession { .. }
1083        | NodeRequest::IndexNativeSession { .. }
1084        | NodeRequest::PreviewSessionRecord { .. } => NATIVE_SESSION_REQUEST_DEADLINE,
1085        NodeRequest::CreateStandaloneWorkspace { .. }
1086        | NodeRequest::CreateWorktree { .. }
1087        | NodeRequest::RemoveWorktree { .. }
1088        | NodeRequest::CleanupManagedWorktree { .. } => Duration::from_secs(240),
1089        NodeRequest::Spawn { .. }
1090        | NodeRequest::Resume { .. }
1091        | NodeRequest::Stop { .. } => Duration::from_secs(15),
1092        NodeRequest::SpawnSpec { spec } =>
1093            Duration::from_millis(spec.deadline_ms.get()) + NODE_REQUEST_IO_HEADROOM,
1094        NodeRequest::SpawnSpecWithHarnessMcp { deadline_unix_ms, .. } =>
1095            Duration::from_millis(deadline_unix_ms.saturating_sub(unix_ms()))
1096                + NODE_REQUEST_IO_HEADROOM,
1097        NodeRequest::SpawnManagedWorktree { request } =>
1098            Duration::from_millis(request.spawn_spec.deadline_ms.get())
1099                + NODE_REQUEST_IO_HEADROOM,
1100        NodeRequest::ResumeSessionRecord { .. } =>
1101            MANAGED_RESUME_SETTLE_DEADLINE + NODE_REQUEST_IO_HEADROOM,
1102        _ => Duration::from_secs(10),
1103    }
1104}
1105
1106fn validate_resync(
1107    previous: u64,
1108    hello: u64,
1109    current: u64,
1110    oldest_available_sequence: u64,
1111    events: &[NodeEventEnvelope],
1112) -> Vec<GapKind> {
1113    let mut gaps = validate_events(
1114        previous,
1115        current,
1116        oldest_available_sequence,
1117        events,
1118    );
1119    if current < hello && !gaps.contains(&GapKind::CursorRegression) {
1120        gaps.push(GapKind::CursorRegression);
1121    }
1122    gaps
1123}
1124
1125async fn ingress_attempt(
1126    ingress: &mpsc::Sender<Attempt>,
1127    node_id: &NodeId,
1128    result: AttemptResult,
1129) -> io::Result<()> {
1130    ingress.send(Attempt { node_id: node_id.clone(), at_unix_ms: unix_ms(), result })
1131        .await.map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "inventory owner closed"))
1132}
1133
1134async fn relay_backoff(
1135    shutdown: &mut watch::Receiver<bool>,
1136    timings: C2Timings,
1137    failures: usize,
1138    hard: bool,
1139) -> io::Result<()> {
1140    let delay = if hard || failures >= timings.transient_backoffs.len() {
1141        timings.parked_backoff
1142    } else {
1143        timings.transient_backoffs[failures.saturating_sub(1)]
1144    };
1145    tokio::select! {
1146        _ = sleep(delay) => Ok(()),
1147        _ = shutdown.changed() => {
1148            Ok(())
1149        }
1150    }
1151}
1152
1153fn relay_failure_attempt(error: &NodeClientError) -> AttemptResult {
1154    let (error, hard) = sanitize_node_error(error);
1155    AttemptResult::Failure { error, hard }
1156}
1157
1158async fn handle_relay_command(
1159    client: &mut LocalNodeClient,
1160    command: RelayCommand,
1161    node_id: &NodeId,
1162    incarnation_id: NodeIncarnationId,
1163    connection_id: u64,
1164    controller_owned: &mut bool,
1165    hub: &OperatorHub,
1166    cursor: &mut NodeCursor,
1167    ingress: &mpsc::Sender<Attempt>,
1168) -> Result<(), NodeClientError> {
1169    match command {
1170        RelayCommand::Request { operator_connection_id, expected_incarnation_id, request, reply } => {
1171            let relay_deadline = relay_request_deadline(&request, Instant::now());
1172            let expected_spawn = expected_spawn_request(&request);
1173            let expected_provider_session_index = matches!(
1174                &request,
1175                NodeRequest::IndexProviderSession { .. }
1176            ).then(|| request.clone());
1177            let expected_native_session = matches!(
1178                &request,
1179                NodeRequest::CatalogNativeSessions { .. }
1180                    | NodeRequest::PageNativeSessions { .. }
1181                    | NodeRequest::PreviewNativeSession { .. }
1182                    | NodeRequest::IndexNativeSession { .. }
1183            ).then(|| request.clone());
1184            let expected_workspace_content = matches!(
1185                &request,
1186                NodeRequest::ReadWorkspaceFile { .. }
1187                    | NodeRequest::WriteWorkspaceFile { .. }
1188                    | NodeRequest::CreateWorkspaceFile { .. }
1189                    | NodeRequest::CreateWorkspaceDirectory { .. }
1190                    | NodeRequest::ReadGitHistory { .. }
1191                    | NodeRequest::ReadGitDiff { .. }
1192            ).then(|| request.clone());
1193            let expected_session_task = matches!(
1194                &request,
1195                NodeRequest::SetSessionTask { .. }
1196            ).then(|| request.clone());
1197            let expected_harness_mcp = (request.required_capability()
1198                == Some(gate4agent_node_protocol::NODE_HARNESS_MCP_READ_PROXY_CAPABILITY))
1199                .then(|| request.clone());
1200            if !hub.is_active(operator_connection_id) {
1201                let _ = reply.send(Err(relay_failure(C2RelayFailureCode::ClientLagged, "C2 operator connection is no longer active", Some(incarnation_id))));
1202                return Ok(());
1203            }
1204            if expected_incarnation_id != incarnation_id {
1205                let _ = reply.send(Err(relay_failure(C2RelayFailureCode::StaleNodeIncarnation, "node incarnation changed", Some(incarnation_id))));
1206                return Ok(());
1207            }
1208            if matches!(&request, NodeRequest::IndexNativeSession { selection, .. }
1209                if selection.route.scope
1210                    != gate4agent_node_protocol::NativeSessionCatalogScope::Workspace)
1211            {
1212                let _ = reply.send(Err(relay_failure(
1213                    C2RelayFailureCode::RequestForbidden,
1214                    "external native sessions must be registered as workspaces before indexing",
1215                    Some(incarnation_id),
1216                )));
1217                return Ok(());
1218            }
1219            if !is_read_only_request(&request) && !*controller_owned {
1220                match acquire_controller_with_deadline(client, connection_id, relay_deadline).await {
1221                    Ok(owned) if owned => *controller_owned = true,
1222                    Ok(_) => {
1223                        let _ = reply.send(Err(relay_failure(C2RelayFailureCode::RelayBusy, "node controller lease is unavailable", Some(incarnation_id))));
1224                        return Ok(());
1225                    }
1226                    Err(error) if relay_node_failure(&error).is_some() => {
1227                        let failure = relay_node_failure(&error)
1228                            .expect("guarded node request failure");
1229                        let _ = reply.send(Ok(RoutedNodeResponse {
1230                            node_id: node_id.clone(),
1231                            incarnation_id,
1232                            response: Err(failure),
1233                        }));
1234                        return Ok(());
1235                    }
1236                    Err(error) => {
1237                        let _ = reply.send(Err(relay_failure(C2RelayFailureCode::NodeOffline, "node relay disconnected", Some(incarnation_id))));
1238                        return Err(error);
1239                    }
1240                }
1241            }
1242            let response = match bounded_node_request_with_deadline(client, request, relay_deadline).await {
1243                Ok(response) => Ok(response),
1244                Err(NodeClientError::Node(failure)) => Err(failure),
1245                Err(error @ NodeClientError::UnsupportedCapability(_)) => {
1246                    let failure = relay_node_failure(&error)
1247                        .expect("unsupported capability is a routed node failure");
1248                    let _ = reply.send(Ok(RoutedNodeResponse {
1249                        node_id: node_id.clone(),
1250                        incarnation_id,
1251                        response: Err(failure),
1252                    }));
1253                    return Ok(());
1254                }
1255                Err(error) => {
1256                    let _ = reply.send(Err(relay_failure(C2RelayFailureCode::NodeOffline, "node relay disconnected", Some(incarnation_id))));
1257                    return Err(error);
1258                }
1259            };
1260            if let Err(message) = validate_spawn_spec_response(
1261                expected_spawn.as_ref(),
1262                &response,
1263                incarnation_id,
1264            ) {
1265                let _ = reply.send(Err(relay_failure(
1266                    C2RelayFailureCode::NodeOffline,
1267                    "node relay returned an invalid spawn receipt",
1268                    Some(incarnation_id),
1269                )));
1270                return Err(NodeClientError::Protocol(message.to_owned()));
1271            }
1272            if let Err(message) = validate_provider_session_index_response(
1273                expected_provider_session_index.as_ref(),
1274                &response,
1275            ) {
1276                let _ = reply.send(Err(relay_failure(
1277                    C2RelayFailureCode::NodeOffline,
1278                    "node relay returned an invalid provider session index response",
1279                    Some(incarnation_id),
1280                )));
1281                return Err(NodeClientError::Protocol(message.to_owned()));
1282            }
1283            if let Err(message) = validate_native_session_response(
1284                expected_native_session.as_ref(),
1285                &response,
1286            ) {
1287                let _ = reply.send(Err(relay_failure(
1288                    C2RelayFailureCode::NodeOffline,
1289                    "node relay returned an invalid native session response",
1290                    Some(incarnation_id),
1291                )));
1292                return Err(NodeClientError::Protocol(message.to_owned()));
1293            }
1294            if let Err(message) = validate_workspace_content_response(
1295                expected_workspace_content.as_ref(),
1296                &response,
1297            ) {
1298                let _ = reply.send(Err(relay_failure(
1299                    C2RelayFailureCode::NodeOffline,
1300                    "node relay returned an invalid workspace content response",
1301                    Some(incarnation_id),
1302                )));
1303                return Err(NodeClientError::Protocol(message.to_owned()));
1304            }
1305            if let Err(message) = validate_session_task_response(
1306                expected_session_task.as_ref(),
1307                &response,
1308            ) {
1309                let _ = reply.send(Err(relay_failure(
1310                    C2RelayFailureCode::NodeOffline,
1311                    "node relay returned an invalid session task response",
1312                    Some(incarnation_id),
1313                )));
1314                return Err(NodeClientError::Protocol(message.to_owned()));
1315            }
1316            if let Err(message) = validate_harness_mcp_response(
1317                expected_harness_mcp.as_ref(),
1318                &response,
1319            ) {
1320                let _ = reply.send(Err(relay_failure(
1321                    C2RelayFailureCode::NodeOffline,
1322                    "node relay returned an invalid harness MCP response",
1323                    Some(incarnation_id),
1324                )));
1325                return Err(NodeClientError::Protocol(message.to_owned()));
1326            }
1327            drain_pending_events(client, node_id, cursor, hub, ingress, false, relay_deadline).await?;
1328            update_inventory_from_response(node_id, cursor, &response, ingress).await
1329                .map_err(NodeClientError::Io)?;
1330            let response = response
1331                .map(|response| C2NodeResponse::from_node_response_with_control_detail(&response))
1332                .map_err(|failure| C2NodeFailure::from(&failure));
1333            let _ = reply.send(Ok(RoutedNodeResponse { node_id: node_id.clone(), incarnation_id, response }));
1334        }
1335    }
1336    Ok(())
1337}
1338
1339fn validate_session_task_response(
1340    expected: Option<&NodeRequest>,
1341    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1342) -> Result<(), &'static str> {
1343    match (expected, response) {
1344        (Some(_), Err(_)) | (None, Err(_)) => Ok(()),
1345        (
1346            Some(NodeRequest::SetSessionTask {
1347                record_id,
1348                expected_revision,
1349                target,
1350            }),
1351            Ok(NodeResponse::SessionRecordUpdated { record }),
1352        ) if session_task_record_matches(record, record_id, *expected_revision, target) => Ok(()),
1353        (Some(NodeRequest::SetSessionTask { .. }), Ok(_)) => {
1354            Err("session task response does not match the routed request")
1355        }
1356        (None, Ok(_)) => Ok(()),
1357        (Some(_), Ok(_)) => Ok(()),
1358    }
1359}
1360
1361fn validate_harness_mcp_response(
1362    expected: Option<&NodeRequest>,
1363    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1364) -> Result<(), &'static str> {
1365    use NodeRequest as Request;
1366    use NodeResponse as Response;
1367    let valid = match (expected, response) {
1368        (Some(Request::ArmHarnessMcpReservation { reservation_id, activation_digest, expires_at_unix_ms, .. }),
1369            Ok(Response::Armed { reservation_id: echoed_id, activation_digest: echoed_digest, expires_at_unix_ms: echoed_expiry })) =>
1370            reservation_id == echoed_id && activation_digest == echoed_digest
1371                && expires_at_unix_ms == echoed_expiry,
1372        (Some(Request::SpawnSpecWithHarnessMcp { reservation_id, activation_digest, .. }),
1373            Ok(Response::Spawned { reservation_id: echoed_id, activation_digest: echoed_digest, receipt })) =>
1374            reservation_id == echoed_id && activation_digest == echoed_digest
1375                && receipt.harness_mcp_proxy.as_ref().is_some_and(|proxy| {
1376                    &proxy.reservation_id == reservation_id
1377                        && &proxy.activation_digest == activation_digest
1378                }),
1379        (Some(Request::ActivateHarnessMcpReservation { reservation_id, activation_digest, record_id, session }),
1380            Ok(Response::Activated { reservation_id: echoed_id, activation_digest: echoed_digest, record_id: echoed_record, session: echoed_session })) =>
1381            reservation_id == echoed_id && activation_digest == echoed_digest
1382                && record_id == echoed_record && session == echoed_session,
1383        (Some(Request::AbortHarnessMcpReservation { reservation_id, activation_digest }),
1384            Ok(Response::Aborted { reservation_id: echoed_id, activation_digest: echoed_digest })) =>
1385            reservation_id == echoed_id && activation_digest == echoed_digest,
1386        (Some(Request::PutHarnessMcpReplyChunk { reservation_id, activation_digest, record_id, session, call_id, offset, final_chunk, chunk_hex }),
1387            Ok(Response::ReplyChunkAccepted { reservation_id: echoed_id, activation_digest: echoed_digest, record_id: echoed_record, session: echoed_session, call_id: echoed_call, next_offset, completed })) =>
1388            reservation_id == echoed_id && activation_digest == echoed_digest
1389                && record_id == echoed_record && session == echoed_session && call_id == echoed_call
1390                && offset.checked_add(u32::try_from(chunk_hex.raw_len()).unwrap_or(u32::MAX))
1391                    == Some(*next_offset) && completed == final_chunk,
1392        (Some(Request::RejectHarnessMcpCall { reservation_id, activation_digest, record_id, session, call_id, .. }),
1393            Ok(Response::CallRejected { reservation_id: echoed_id, activation_digest: echoed_digest, record_id: echoed_record, session: echoed_session, call_id: echoed_call })) =>
1394            reservation_id == echoed_id && activation_digest == echoed_digest
1395                && record_id == echoed_record && session == echoed_session && call_id == echoed_call,
1396        (Some(_), Err(_)) | (None, Err(_)) => true,
1397        (Some(_), Ok(_)) => false,
1398        (None, Ok(response)) => !response.requires_harness_mcp_proxy_capability(),
1399    };
1400    if valid { Ok(()) } else { Err("harness MCP response does not match the routed request") }
1401}
1402
1403fn session_task_record_matches(
1404    record: &gate4agent_node_protocol::ManagedSessionRecord,
1405    record_id: &gate4agent_node_protocol::SessionRecordId,
1406    expected_revision: u64,
1407    target: &gate4agent_node_protocol::SessionTaskTargetV1,
1408) -> bool {
1409    if &record.record_id != record_id { return false; }
1410    let next_revision = expected_revision.checked_add(1);
1411    match target {
1412        gate4agent_node_protocol::SessionTaskTargetV1::New => record.task_binding.as_ref()
1413            .is_some_and(|binding| Some(binding.revision) == next_revision && binding.task_id.is_some()),
1414        gate4agent_node_protocol::SessionTaskTargetV1::Existing { task_id } => record.task_binding.as_ref()
1415            .is_some_and(|binding| (binding.revision == expected_revision || Some(binding.revision) == next_revision)
1416                && binding.task_id.as_ref() == Some(task_id)),
1417        gate4agent_node_protocol::SessionTaskTargetV1::Clear => match &record.task_binding {
1418            None => expected_revision == 0,
1419            Some(binding) => binding.task_id.is_none()
1420                && (binding.revision == expected_revision || Some(binding.revision) == next_revision),
1421        },
1422    }
1423}
1424
1425fn validate_workspace_content_response(
1426    expected: Option<&NodeRequest>,
1427    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1428) -> Result<(), &'static str> {
1429    match (expected, response) {
1430        (Some(_), Err(_)) | (None, Err(_)) => Ok(()),
1431        (
1432            Some(NodeRequest::ReadWorkspaceFile { workspace_id, path }),
1433            Ok(NodeResponse::WorkspaceFileRead { file }),
1434        ) if &file.workspace_id == workspace_id && &file.path == path => Ok(()),
1435        (
1436            Some(NodeRequest::WriteWorkspaceFile { workspace_id, path, text, .. }),
1437            Ok(NodeResponse::WorkspaceFileWritten { file }),
1438        ) if &file.workspace_id == workspace_id
1439            && &file.path == path
1440            && matches!(
1441                &file.content,
1442                gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
1443                    text: written,
1444                    byte_len,
1445                } if written == text
1446                    && u32::try_from(text.len()).ok() == Some(*byte_len)
1447            ) => Ok(()),
1448        (
1449            Some(NodeRequest::CreateWorkspaceFile { workspace_id, path }),
1450            Ok(NodeResponse::WorkspaceFileCreated { file }),
1451        ) if &file.workspace_id == workspace_id
1452            && &file.path == path
1453            && file.revision.is_some()
1454            && matches!(
1455                &file.content,
1456                gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
1457                    text,
1458                    byte_len: 0,
1459                } if text.is_empty()
1460            ) => Ok(()),
1461        (
1462            Some(NodeRequest::CreateWorkspaceDirectory { workspace_id, path }),
1463            Ok(NodeResponse::WorkspaceDirectoryCreated {
1464                workspace_id: actual_workspace_id,
1465                entry,
1466            }),
1467        ) if actual_workspace_id == workspace_id
1468            && &entry.relative_path == path
1469            && entry.kind == gate4agent_node_protocol::WorkspaceEntryKind::Directory => Ok(()),
1470        (
1471            Some(NodeRequest::ReadGitHistory { workspace_id, .. }),
1472            Ok(NodeResponse::GitHistoryRead { workspace_id: actual, .. }),
1473        ) if actual == workspace_id => Ok(()),
1474        (
1475            Some(NodeRequest::ReadGitDiff { workspace_id, request }),
1476            Ok(NodeResponse::GitDiffRead { workspace_id: actual, diff }),
1477        ) if actual == workspace_id && diff.mode == request.mode && diff.path == request.path => Ok(()),
1478        (Some(_), Ok(_)) => Err("workspace content response does not match routed request"),
1479        (
1480            None,
1481            Ok(NodeResponse::WorkspaceFileRead { .. }
1482                | NodeResponse::WorkspaceFileWritten { .. }
1483                | NodeResponse::WorkspaceFileCreated { .. }
1484                | NodeResponse::WorkspaceDirectoryCreated { .. }
1485                | NodeResponse::GitHistoryRead { .. }
1486                | NodeResponse::GitDiffRead { .. }),
1487        ) => Err("unexpected workspace content response"),
1488        (None, Ok(_)) => Ok(()),
1489    }
1490}
1491
1492enum ExpectedSpawnRequest {
1493    Spec(SpawnSpec),
1494    ManagedV1(ManagedWorktreeSpawnRequest),
1495    ManagedV2(ManagedWorktreeSpawnRequestV2),
1496}
1497
1498fn expected_spawn_request(request: &NodeRequest) -> Option<ExpectedSpawnRequest> {
1499    match request {
1500        NodeRequest::SpawnSpec { spec } => Some(ExpectedSpawnRequest::Spec(spec.clone())),
1501        NodeRequest::SpawnManagedWorktree { request } => {
1502            Some(ExpectedSpawnRequest::ManagedV1(request.clone()))
1503        }
1504        NodeRequest::SpawnManagedWorktreeV2 { request } => {
1505            Some(ExpectedSpawnRequest::ManagedV2(request.clone()))
1506        }
1507        _ => None,
1508    }
1509}
1510
1511fn validate_provider_session_index_response(
1512    expected: Option<&NodeRequest>,
1513    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1514) -> Result<(), &'static str> {
1515    match (expected, response) {
1516        (
1517            Some(NodeRequest::IndexProviderSession {
1518                workspace_id,
1519                provider,
1520                identity,
1521                ..
1522            }),
1523            Ok(NodeResponse::ProviderSessionIndexed { record }),
1524        ) if &record.workspace_id == workspace_id
1525            && &record.provider == provider
1526            && record.provider_session.as_ref() == Some(identity) => Ok(()),
1527        (Some(_), Ok(NodeResponse::ProviderSessionIndexed { .. })) => {
1528            Err("provider session index response does not match routed request")
1529        }
1530        (Some(_), Ok(_)) => Err("provider session index request returned a different response"),
1531        (None, Ok(NodeResponse::ProviderSessionIndexed { .. })) => {
1532            Err("unexpected provider session index response for a different node request")
1533        }
1534        (Some(_), Err(_)) | (None, _) => Ok(()),
1535    }
1536}
1537
1538fn validate_native_session_response(
1539    expected: Option<&NodeRequest>,
1540    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1541) -> Result<(), &'static str> {
1542    match (expected, response) {
1543        (
1544            Some(NodeRequest::CatalogNativeSessions { route, .. }),
1545            Ok(NodeResponse::NativeSessionsCataloged {
1546                route: echoed_route,
1547                ..
1548            }),
1549        ) if echoed_route == route => Ok(()),
1550        (
1551            Some(NodeRequest::PageNativeSessions {
1552                route,
1553                window,
1554                catalog_revision,
1555                ..
1556            }),
1557            Ok(NodeResponse::NativeSessionsPaged {
1558                route: echoed_route,
1559                page,
1560            }),
1561        ) if echoed_route == route
1562            && page.window == *window
1563            && page.revision == *catalog_revision => Ok(()),
1564        (
1565            Some(NodeRequest::PreviewNativeSession { selection, .. }),
1566            Ok(NodeResponse::NativeSessionPreviewed {
1567                selection: echoed_selection,
1568                ..
1569            }),
1570        ) if echoed_selection == selection => Ok(()),
1571        (
1572            Some(NodeRequest::IndexNativeSession { selection, .. }),
1573            Ok(NodeResponse::NativeSessionIndexed {
1574                selection: echoed_selection,
1575                record,
1576            }),
1577        ) if echoed_selection == selection
1578            && selection.route.scope
1579                == gate4agent_node_protocol::NativeSessionCatalogScope::Workspace
1580            && selection.route.workspace_id.as_ref() == Some(&record.workspace_id)
1581            && selection.route.provider == record.provider => Ok(()),
1582        (Some(_), Err(_)) => Ok(()),
1583        (Some(_), Ok(_)) => Err("native session response does not match routed request"),
1584        (
1585            None,
1586            Ok(
1587                NodeResponse::NativeSessionsCataloged { .. }
1588                | NodeResponse::NativeSessionsPaged { .. }
1589                | NodeResponse::NativeSessionPreviewed { .. }
1590                | NodeResponse::NativeSessionIndexed { .. },
1591            ),
1592        ) => Err("unexpected native session response for a different node request"),
1593        (None, _) => Ok(()),
1594    }
1595}
1596
1597fn validate_spawn_spec_response(
1598    expected: Option<&ExpectedSpawnRequest>,
1599    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1600    relay_incarnation_id: NodeIncarnationId,
1601) -> Result<(), &'static str> {
1602    match (expected, response) {
1603        (Some(ExpectedSpawnRequest::Spec(spec)), Ok(NodeResponse::SpawnSpecAccepted { receipt })) => {
1604            validate_spawn_receipt(spec, receipt, relay_incarnation_id)
1605        }
1606        (
1607            Some(ExpectedSpawnRequest::ManagedV1(request)),
1608            Ok(NodeResponse::ManagedWorktreeSpawnAccepted { receipt }),
1609        ) => validate_managed_spawn_receipt(
1610            &request.spawn_spec,
1611            &request.worktree_profile_id,
1612            receipt,
1613            relay_incarnation_id,
1614        ),
1615        (
1616            Some(ExpectedSpawnRequest::ManagedV2(request)),
1617            Ok(NodeResponse::ManagedWorktreeSpawnAccepted { receipt }),
1618        ) if receipt.lease.profile_revision == request.expected_profile_revision => {
1619            validate_managed_spawn_receipt(
1620                &request.spawn_spec,
1621                &request.worktree_profile_id,
1622                receipt,
1623                relay_incarnation_id,
1624            )
1625        }
1626        (
1627            Some(ExpectedSpawnRequest::ManagedV2(_)),
1628            Ok(NodeResponse::ManagedWorktreeSpawnAccepted { .. }),
1629        ) => Err("managed spawn receipt profile revision does not match routed request"),
1630        (Some(_), Ok(_)) => return Err("spawn spec request returned a different response"),
1631        (Some(_), Err(_)) => return Ok(()),
1632        (None, Ok(NodeResponse::SpawnSpecAccepted { .. }
1633            | NodeResponse::ManagedWorktreeSpawnAccepted { .. })) => {
1634            return Err("unexpected spawn receipt for a different node request");
1635        }
1636        (None, _) => return Ok(()),
1637    }
1638}
1639
1640fn validate_managed_spawn_receipt(
1641    spawn_spec: &SpawnSpec,
1642    worktree_profile_id: &gate4agent_node_protocol::WorktreeProfileId,
1643    receipt: &gate4agent_node_protocol::ManagedWorktreeSpawnReceipt,
1644    relay_incarnation_id: NodeIncarnationId,
1645) -> Result<(), &'static str> {
1646    if receipt.lease.source_workspace_id != spawn_spec.target.workspace_id
1647        || &receipt.lease.profile_id != worktree_profile_id
1648        || receipt.lease.state != ManagedWorktreeLeaseState::InUse
1649        || receipt.lease.cleanup_failure.is_some()
1650        || receipt.lease.active_session_count != 1
1651        || receipt.spawn.target.node_id != spawn_spec.target.node_id
1652        || receipt.spawn.target.workspace_id != spawn_spec.target.workspace_id
1653        || receipt.spawn.target.worktree_id.as_ref() != Some(&receipt.lease.workspace_id)
1654        || receipt.spawn.session.workspace_id != receipt.lease.workspace_id
1655    {
1656        return Err("managed spawn receipt does not match routed request");
1657    }
1658    let mut resolved_spec = spawn_spec.clone();
1659    resolved_spec.target.worktree_id = Some(receipt.lease.workspace_id.clone());
1660    validate_spawn_receipt(&resolved_spec, &receipt.spawn, relay_incarnation_id)
1661}
1662
1663fn validate_spawn_receipt(
1664    spec: &SpawnSpec,
1665    receipt: &ResolvedSpawnReceipt,
1666    relay_incarnation_id: NodeIncarnationId,
1667) -> Result<(), &'static str> {
1668    if receipt.incarnation_id != relay_incarnation_id
1669        || !receipt.context_binding_is_valid()
1670        || receipt.target != spec.target
1671        || receipt.profile_id != spec.profile_id
1672        || receipt.idempotency_key != spec.idempotency_key
1673        || receipt.deadline_ms != spec.deadline_ms
1674        || receipt.required_capabilities != spec.required_capabilities
1675        || receipt.bundle.as_ref().is_some_and(|bundle| {
1676            receipt.bundle_id.as_ref() != Some(&bundle.id)
1677        })
1678        || &receipt.session.workspace_id
1679            != spec
1680                .target
1681                .worktree_id
1682                .as_ref()
1683                .unwrap_or(&spec.target.workspace_id)
1684    {
1685        return Err("spawn receipt does not match routed request");
1686    }
1687    if !required_override_matches(&spec.overrides.provider, &receipt.provider)
1688        || !required_override_matches(&spec.overrides.mode, &receipt.mode)
1689        || !required_override_matches(&spec.overrides.terminal_size, &receipt.terminal_size)
1690        || !optional_override_matches(&spec.overrides.bundle_id, &receipt.bundle_id)
1691        || !optional_override_matches(&spec.overrides.context_id, &receipt.context_id)
1692        || !environment_profile_override_matches(
1693            &spec.overrides.environment_profile_id,
1694            receipt.environment_profile.as_ref(),
1695        )
1696    {
1697        return Err("spawn receipt contradicts explicit overrides");
1698    }
1699    match &spec.overrides.prompt {
1700        SpawnOverride::Inherit => {}
1701        SpawnOverride::Set { value }
1702            if receipt.prompt.present
1703                && receipt.prompt.byte_len == u32::try_from(value.byte_len()).unwrap_or(0) => {}
1704        SpawnOverride::Clear if !receipt.prompt.present && receipt.prompt.byte_len == 0 => {}
1705        SpawnOverride::Set { .. } | SpawnOverride::Clear => {
1706            return Err("spawn receipt contradicts explicit prompt override");
1707        }
1708    }
1709    Ok(())
1710}
1711
1712fn environment_profile_override_matches(
1713    expected: &SpawnOverride<gate4agent_node_protocol::SpawnEnvironmentProfileId>,
1714    actual: Option<&gate4agent_node_protocol::ResolvedEnvironmentProfileReceipt>,
1715) -> bool {
1716    match expected {
1717        SpawnOverride::Inherit => true,
1718        SpawnOverride::Set { value } => {
1719            actual.is_some_and(|receipt| &receipt.profile_id == value)
1720        }
1721        SpawnOverride::Clear => actual.is_none(),
1722    }
1723}
1724
1725fn required_override_matches<T: Eq>(override_value: &SpawnOverride<T>, actual: &T) -> bool {
1726    match override_value {
1727        SpawnOverride::Inherit => true,
1728        SpawnOverride::Set { value } => value == actual,
1729        SpawnOverride::Clear => false,
1730    }
1731}
1732
1733fn optional_override_matches<T: Eq>(
1734    override_value: &SpawnOverride<T>,
1735    actual: &Option<T>,
1736) -> bool {
1737    match override_value {
1738        SpawnOverride::Inherit => true,
1739        SpawnOverride::Set { value } => actual.as_ref() == Some(value),
1740        SpawnOverride::Clear => actual.is_none(),
1741    }
1742}
1743
1744fn relay_node_failure(error: &NodeClientError) -> Option<C2NodeFailure> {
1745    match error {
1746        NodeClientError::Node(failure) => Some(C2NodeFailure::from(failure)),
1747        NodeClientError::UnsupportedCapability(_) => Some(C2NodeFailure {
1748            code: NodeFailureCode::UnsupportedCapability,
1749            message: "required capability unavailable".to_owned(),
1750        }),
1751        NodeClientError::Io(_)
1752        | NodeClientError::Frame(_)
1753        | NodeClientError::Protocol(_)
1754        | NodeClientError::BuildStampMismatch { .. }
1755        | NodeClientError::AuthenticationTimedOut
1756        | NodeClientError::Authentication(_)
1757        | NodeClientError::RequestIdExhausted => None,
1758    }
1759}
1760
1761fn is_read_only_request(request: &NodeRequest) -> bool {
1762    matches!(request,
1763        NodeRequest::Snapshot
1764        | NodeRequest::Resync { .. }
1765        | NodeRequest::BrowseHostDirectories { .. }
1766        | NodeRequest::InspectWorkspace { .. }
1767        | NodeRequest::ReadWorkspaceFile { .. }
1768        | NodeRequest::ReadGitHistory { .. }
1769        | NodeRequest::ReadGitDiff { .. }
1770        | NodeRequest::CatalogNativeSessions { .. }
1771        | NodeRequest::PageNativeSessions { .. }
1772        | NodeRequest::PreviewNativeSession { .. }
1773        | NodeRequest::PreviewSessionRecord { .. }
1774    )
1775}
1776
1777async fn acquire_controller(
1778    client: &mut LocalNodeClient,
1779    connection_id: u64,
1780) -> Result<bool, NodeClientError> {
1781    acquire_controller_with_deadline(client, connection_id, None).await
1782}
1783
1784async fn acquire_controller_with_deadline(
1785    client: &mut LocalNodeClient,
1786    connection_id: u64,
1787    relay_deadline: Option<Instant>,
1788) -> Result<bool, NodeClientError> {
1789    match bounded_node_request_with_deadline(
1790        client,
1791        NodeRequest::AcquireController { lease_ms: gate4agent_node_protocol::MAX_CONTROLLER_LEASE_MS },
1792        relay_deadline,
1793    ).await? {
1794        NodeResponse::Controller { controller } => Ok(controller.as_ref().is_some_and(|state| state.connection_id == connection_id)),
1795        _ => Err(NodeClientError::Protocol("controller acquisition returned a different response".to_owned())),
1796    }
1797}
1798
1799async fn release_controller(
1800    client: &mut LocalNodeClient,
1801    controller_owned: &mut bool,
1802) -> Result<(), NodeClientError> {
1803    if !*controller_owned { return Ok(()); }
1804    match bounded_node_request(client, NodeRequest::ReleaseController).await? {
1805        NodeResponse::Controller { .. } => { *controller_owned = false; Ok(()) }
1806        _ => Err(NodeClientError::Protocol("controller release returned a different response".to_owned())),
1807    }
1808}
1809
1810async fn update_inventory_from_response(
1811    node_id: &NodeId,
1812    cursor: &mut NodeCursor,
1813    response: &Result<NodeResponse, gate4agent_node_protocol::NodeFailure>,
1814    ingress: &mpsc::Sender<Attempt>,
1815) -> io::Result<()> {
1816    match response {
1817        Ok(NodeResponse::Snapshot { event_sequence, snapshot, .. })
1818        | Ok(NodeResponse::Resync { event_sequence, snapshot, .. }) => {
1819            if *event_sequence < cursor.sequence { return Ok(()); }
1820            cursor.sequence = *event_sequence;
1821            ingress_attempt(ingress, node_id, AttemptResult::Success { cursor: *cursor, snapshot: snapshot.clone(), gaps: Vec::new() }).await
1822        }
1823        _ => Ok(()),
1824    }
1825}
1826
1827fn publish_recovered_events(
1828    node_id: &NodeId,
1829    incarnation_id: NodeIncarnationId,
1830    events: &[NodeEventEnvelope],
1831    hub: &OperatorHub,
1832) {
1833    for envelope in events {
1834        if let Some(event) = routed_recovered_node_event(node_id, incarnation_id, envelope) {
1835            hub.publish(event);
1836        }
1837    }
1838}
1839
1840fn routed_recovered_node_event(
1841    node_id: &NodeId,
1842    incarnation_id: NodeIncarnationId,
1843    envelope: &NodeEventEnvelope,
1844) -> Option<RoutedNodeEvent> {
1845    (!matches!(&envelope.event, NodeEvent::HarnessMcpReadCall { .. })).then(|| {
1846        RoutedNodeEvent {
1847            node_id: node_id.clone(),
1848            cursor: NodeCursor { incarnation_id, sequence: envelope.sequence },
1849            event: C2NodeEvent::from_node_event_with_control_detail(&envelope.event),
1850        }
1851    })
1852}
1853
1854fn routed_transient_node_event(
1855    node_id: &NodeId,
1856    cursor: NodeCursor,
1857    envelope: &NodeEventEnvelope,
1858) -> Option<RoutedNodeEvent> {
1859    matches!(&envelope.event, NodeEvent::HarnessMcpReadCall { .. }).then(|| RoutedNodeEvent {
1860        node_id: node_id.clone(),
1861        cursor,
1862        event: C2NodeEvent::from(&envelope.event),
1863    })
1864}
1865
1866/// `NodeEvent::AgentStream` chunks are published unconditionally -- never
1867/// gated on `cursor` contiguity, and never the cause of a
1868/// `CursorRegression`/`NonContiguousEvents` gap -- but, unlike
1869/// `HarnessMcpReadCall` (always `sequence: 0`, see
1870/// `gate4agent-node/src/server.rs::publish_transient`), an
1871/// `AgentStreamChunkV1` envelope draws its `sequence` from the
1872/// SAME counter every durable `NodeEvent` shares
1873/// (`gate4agent-node/src/server.rs::publish_agent_stream_chunk`'s own doc,
1874/// there purely so the node's per-connection discard watermark can compare
1875/// it against everything else). So a chunk really does consume a slot in
1876/// the durable sequence space, and `routed_transient_node_event`'s
1877/// "leave `cursor` untouched" rule -- correct for a call that never had a
1878/// slot to begin with -- would otherwise make the very next durable event
1879/// look one short of contiguous, forcing a resync round trip on every
1880/// single provider turn that streams so much as one chunk.
1881///
1882/// This folds the chunk's sequence into `cursor` with `max` instead: on the
1883/// ordinary, non-bursty path (`Control(N)` then its own `AgentStream(N+1)`,
1884/// each one delivered as they are produced) that keeps `cursor` moving
1885/// exactly as it did before this function existed, so the very next
1886/// `Control(N+2)` still finds `cursor == N+1` and stays contiguous. On the
1887/// bursty path -- the one actually measured live: 19 chunks published
1888/// `source_sequence` 22-40 in a 17ms window, `envelope_sequence` interleaved
1889/// 1:1 with `Control`, only 2 of the 19 ever reached the harness -- the
1890/// node's connection loop can drain its durable channel far enough ahead of
1891/// this dedicated one that a lower-sequence chunk arrives at this relay
1892/// AFTER `cursor` has already been carried past it by a resync recovering
1893/// the higher-sequence `Control` envelopes around it. `max` refuses to let
1894/// that late chunk rewind `cursor` -- it is still published (chunks promise
1895/// no resync, so there is nothing to recover if it were dropped instead:
1896/// `gate4agent-node-protocol::AgentStreamChunkV1`'s own doc, "no
1897/// `ObservationV1` resync promise") -- it simply stops mattering to the
1898/// durable cursor's own contiguity bookkeeping once something newer has
1899/// already passed it by. Ordering and true loss within the agent-stream
1900/// channel itself remain the node's own broadcast `Lagged` warn to name,
1901/// not this cursor's to adjudicate.
1902fn route_agent_stream_event(
1903    node_id: &NodeId,
1904    cursor: &mut NodeCursor,
1905    envelope: &NodeEventEnvelope,
1906) -> Option<RoutedNodeEvent> {
1907    if !matches!(&envelope.event, NodeEvent::AgentStream { .. }) {
1908        return None;
1909    }
1910    let routed = RoutedNodeEvent {
1911        node_id: node_id.clone(),
1912        cursor: *cursor,
1913        event: C2NodeEvent::from(&envelope.event),
1914    };
1915    cursor.sequence = cursor.sequence.max(envelope.sequence);
1916    Some(routed)
1917}
1918
1919async fn drain_pending_events(
1920    client: &mut LocalNodeClient,
1921    node_id: &NodeId,
1922    cursor: &mut NodeCursor,
1923    hub: &OperatorHub,
1924    ingress: &mpsc::Sender<Attempt>,
1925    skip_replayed: bool,
1926    relay_deadline: Option<Instant>,
1927) -> Result<(), NodeClientError> {
1928    let mut skip_replayed = skip_replayed;
1929    for repair_pass in 0..=1 {
1930        let mut gaps = Vec::new();
1931        let mut managed_worktree_events = Vec::new();
1932        let mut changed = false;
1933        let mut repair = false;
1934        while let Some(envelope) = client.take_event() {
1935            if let Some(event) = routed_transient_node_event(node_id, *cursor, &envelope) {
1936                hub.publish(event);
1937                continue;
1938            }
1939            if let Some(event) = route_agent_stream_event(node_id, cursor, &envelope) {
1940                hub.publish(event);
1941                continue;
1942            }
1943            if skip_replayed && envelope.sequence <= cursor.sequence { continue; }
1944            let resync_required = matches!(&envelope.event, gate4agent_node_protocol::NodeEvent::ResyncRequired { .. });
1945            if resync_required || envelope.sequence != cursor.sequence.saturating_add(1) {
1946                gaps.push(if envelope.sequence <= cursor.sequence && !resync_required {
1947                    GapKind::CursorRegression
1948                } else if resync_required {
1949                    GapKind::HistoryEvicted
1950                } else {
1951                    GapKind::NonContiguousEvents
1952                });
1953                repair = true;
1954                continue;
1955            }
1956            let event_cursor = NodeCursor { incarnation_id: cursor.incarnation_id, sequence: envelope.sequence };
1957            if matches!(
1958                &envelope.event,
1959                NodeEvent::ManagedWorktreeUpserted { .. }
1960                    | NodeEvent::ManagedWorktreeRemoved { .. }
1961            ) {
1962                managed_worktree_events.push(envelope.event.clone());
1963            }
1964            hub.publish(RoutedNodeEvent {
1965                node_id: node_id.clone(),
1966                cursor: event_cursor,
1967                event: C2NodeEvent::from_node_event_with_control_detail(&envelope.event),
1968            });
1969            cursor.sequence = envelope.sequence;
1970            changed = true;
1971        }
1972        if !repair {
1973            if changed || !gaps.is_empty() {
1974                ingress_attempt(ingress, node_id, AttemptResult::Cursor {
1975                    cursor: *cursor,
1976                    gaps,
1977                    managed_worktree_events,
1978                }).await
1979                    .map_err(NodeClientError::Io)?;
1980            }
1981            return Ok(());
1982        }
1983        if repair_pass == 1 {
1984            ingress_attempt(ingress, node_id, AttemptResult::Cursor {
1985                cursor: *cursor,
1986                gaps,
1987                managed_worktree_events,
1988            }).await
1989                .map_err(NodeClientError::Io)?;
1990            return Err(NodeClientError::Protocol("node event stream remained noncontiguous after resync".to_owned()));
1991        }
1992        let after_sequence = cursor.sequence;
1993        let response = bounded_node_request_with_deadline(
1994            client,
1995            NodeRequest::Resync { after_sequence },
1996            relay_deadline,
1997        ).await?;
1998        let NodeResponse::Resync {
1999            event_sequence,
2000            oldest_available_sequence,
2001            snapshot,
2002            events,
2003        } = response else {
2004            return Err(NodeClientError::Protocol("event repair resync returned a different response".to_owned()));
2005        };
2006        let repair_gaps = validate_events(
2007            after_sequence,
2008            event_sequence,
2009            oldest_available_sequence,
2010            &events,
2011        );
2012        let contiguous = repair_gaps.is_empty();
2013        gaps = repair_gaps;
2014        if contiguous {
2015            for envelope in events.iter().filter(|event| event.sequence > after_sequence) {
2016                if let Some(event) = routed_recovered_node_event(
2017                    node_id,
2018                    cursor.incarnation_id,
2019                    envelope,
2020                ) {
2021                    hub.publish(event);
2022                }
2023            }
2024        } else {
2025            hub.publish(RoutedNodeEvent {
2026                node_id: node_id.clone(),
2027                cursor: NodeCursor { incarnation_id: cursor.incarnation_id, sequence: event_sequence },
2028                event: C2NodeEvent::ResyncRequired {
2029                    oldest_available_sequence,
2030                },
2031            });
2032        }
2033        cursor.sequence = event_sequence;
2034        ingress_attempt(ingress, node_id, AttemptResult::Success { cursor: *cursor, snapshot, gaps }).await
2035            .map_err(NodeClientError::Io)?;
2036        skip_replayed = true;
2037    }
2038    Ok(())
2039}
2040
2041async fn handle_live_node_event(
2042    client: &mut LocalNodeClient,
2043    node_id: &NodeId,
2044    envelope: NodeEventEnvelope,
2045    cursor: &mut NodeCursor,
2046    hub: &OperatorHub,
2047    ingress: &mpsc::Sender<Attempt>,
2048) -> Result<(), NodeClientError> {
2049    if let Some(event) = routed_transient_node_event(node_id, *cursor, &envelope) {
2050        hub.publish(event);
2051        return Ok(());
2052    }
2053    if let Some(event) = route_agent_stream_event(node_id, cursor, &envelope) {
2054        hub.publish(event);
2055        return Ok(());
2056    }
2057    if live_event_gap(cursor.sequence, &envelope).is_some() {
2058        let after_sequence = cursor.sequence;
2059        let response = bounded_node_request(
2060            client,
2061            NodeRequest::Resync { after_sequence },
2062        ).await?;
2063        let NodeResponse::Resync {
2064            event_sequence,
2065            oldest_available_sequence,
2066            snapshot,
2067            events,
2068        } = response else {
2069            return Err(NodeClientError::Protocol(
2070                "live event repair resync returned a different response".to_owned(),
2071            ));
2072        };
2073        let repair_gaps = validate_events(
2074            after_sequence,
2075            event_sequence,
2076            oldest_available_sequence,
2077            &events,
2078        );
2079        if repair_gaps.is_empty() {
2080            publish_recovered_events(node_id, cursor.incarnation_id, &events, hub);
2081        } else {
2082            hub.publish(RoutedNodeEvent {
2083                node_id: node_id.clone(),
2084                cursor: NodeCursor {
2085                    incarnation_id: cursor.incarnation_id,
2086                    sequence: event_sequence,
2087                },
2088                event: C2NodeEvent::ResyncRequired {
2089                    oldest_available_sequence,
2090                },
2091            });
2092        }
2093        cursor.sequence = event_sequence;
2094        ingress_attempt(
2095            ingress,
2096            node_id,
2097            AttemptResult::Success {
2098                cursor: *cursor,
2099                snapshot,
2100                gaps: repair_gaps,
2101            },
2102        ).await.map_err(NodeClientError::Io)?;
2103        return drain_pending_events(
2104            client,
2105            node_id,
2106            cursor,
2107            hub,
2108            ingress,
2109            true,
2110            None,
2111        ).await;
2112    }
2113
2114    cursor.sequence = envelope.sequence;
2115    let managed_worktree_events = matches!(
2116        &envelope.event,
2117        NodeEvent::ManagedWorktreeUpserted { .. }
2118            | NodeEvent::ManagedWorktreeRemoved { .. }
2119    )
2120    .then(|| vec![envelope.event.clone()])
2121    .unwrap_or_default();
2122    hub.publish(RoutedNodeEvent {
2123        node_id: node_id.clone(),
2124        cursor: *cursor,
2125        event: C2NodeEvent::from_node_event_with_control_detail(&envelope.event),
2126    });
2127    ingress_attempt(
2128        ingress,
2129        node_id,
2130        AttemptResult::Cursor {
2131            cursor: *cursor,
2132            gaps: Vec::new(),
2133            managed_worktree_events,
2134        },
2135    ).await.map_err(NodeClientError::Io)
2136}
2137
2138fn live_event_gap(previous: u64, envelope: &NodeEventEnvelope) -> Option<GapKind> {
2139    if matches!(
2140        &envelope.event,
2141        gate4agent_node_protocol::NodeEvent::ResyncRequired { .. }
2142    ) {
2143        return Some(GapKind::HistoryEvicted);
2144    }
2145    if envelope.sequence <= previous {
2146        return Some(GapKind::CursorRegression);
2147    }
2148    (envelope.sequence != previous.saturating_add(1))
2149        .then_some(GapKind::NonContiguousEvents)
2150}
2151
2152fn reject_disconnected_commands(commands: &mut mpsc::Receiver<RelayCommand>, cursor: Option<NodeCursor>) {
2153    while let Ok(command) = commands.try_recv() {
2154        match command {
2155            RelayCommand::Request { reply, .. } => {
2156                let _ = reply.send(Err(relay_failure(
2157                    C2RelayFailureCode::NodeOffline,
2158                    "node relay disconnected before request dispatch",
2159                    cursor.map(|value| value.incarnation_id),
2160                )));
2161            }
2162        }
2163    }
2164}
2165
2166fn acknowledge_disconnected_releases(releases: &mut mpsc::Receiver<oneshot::Sender<()>>) {
2167    while let Ok(reply) = releases.try_recv() { let _ = reply.send(()); }
2168}
2169
2170fn validate_events(
2171    previous: u64,
2172    current: u64,
2173    oldest_available_sequence: u64,
2174    events: &[NodeEventEnvelope],
2175) -> Vec<GapKind> {
2176    if current < previous { return vec![GapKind::CursorRegression]; }
2177    let max_floor = current.checked_add(1).unwrap_or(u64::MAX);
2178    let minimum_event_sequence = previous
2179        .checked_add(1)
2180        .unwrap_or(u64::MAX)
2181        .max(oldest_available_sequence);
2182    if oldest_available_sequence == 0
2183        || oldest_available_sequence > max_floor
2184        || (current == previous && !events.is_empty())
2185        || events.iter().any(|event| {
2186            event.sequence < minimum_event_sequence || event.sequence > current
2187        })
2188        || events.windows(2).any(|pair| pair[0].sequence >= pair[1].sequence)
2189    {
2190        return vec![GapKind::NonContiguousEvents];
2191    }
2192    if previous.saturating_add(1) < oldest_available_sequence {
2193        return vec![GapKind::HistoryEvicted];
2194    }
2195    Vec::new()
2196}
2197
2198fn sanitize_node_error(error: &NodeClientError) -> (SanitizedError, bool) {
2199    let (category, message, hard) = match error {
2200        NodeClientError::Protocol(message) if message.contains("identity mismatch") =>
2201            (C2ErrorCategory::Identity, "node identity mismatch", true),
2202        NodeClientError::Protocol(message) if message.contains("access-token proof") || message.contains("access denied") =>
2203            (C2ErrorCategory::Authentication, "node authentication failed", true),
2204        NodeClientError::Protocol(_) | NodeClientError::Frame(FrameError::Json(_) | FrameError::InvalidLength { .. }) =>
2205            (C2ErrorCategory::Protocol, "node protocol failed", true),
2206        NodeClientError::BuildStampMismatch { .. } =>
2207            (C2ErrorCategory::Protocol, "node build stamp mismatch", true),
2208        NodeClientError::UnsupportedCapability(_) =>
2209            (C2ErrorCategory::Protocol, "node capability unavailable", true),
2210        NodeClientError::Node(failure) if failure.code == NodeFailureCode::Unauthorized =>
2211            (C2ErrorCategory::Authentication, "node request authentication failed", true),
2212        NodeClientError::Node(_) =>
2213            (C2ErrorCategory::Protocol, "node rejected observer request", true),
2214        NodeClientError::Frame(FrameError::BodyTimedOut { .. } | FrameError::PrefixTimedOut)
2215            | NodeClientError::AuthenticationTimedOut =>
2216            (C2ErrorCategory::Timeout, "node observation deadline exceeded", false),
2217        NodeClientError::Authentication(_) | NodeClientError::RequestIdExhausted =>
2218            (C2ErrorCategory::Internal, "node client failed internally", true),
2219        NodeClientError::Io(_) | NodeClientError::Frame(FrameError::Io(_)) =>
2220            (C2ErrorCategory::Transport, "node transport unavailable", false),
2221    };
2222    (SanitizedError { category, message: message.to_owned() }, hard)
2223}
2224
2225async fn inventory_owner(
2226    configured: usize, fresh_for: Duration, mut ingress: mpsc::Receiver<Attempt>,
2227    status: watch::Sender<Arc<StatusResponse>>, mut shutdown: watch::Receiver<bool>,
2228) -> io::Result<()> {
2229    let mut current = (**status.borrow()).clone();
2230    let mut attempted = BTreeSet::new();
2231    loop {
2232        tokio::select! {
2233            attempt = ingress.recv() => {
2234                let Some(attempt) = attempt else { return Ok(()); };
2235                attempted.insert(attempt.node_id.clone());
2236                let node = current.nodes.get_mut(&attempt.node_id).expect("configured poller node exists");
2237                node.last_attempt_unix_ms = Some(attempt.at_unix_ms);
2238                match attempt.result {
2239                    AttemptResult::Connected {
2240                        cursor,
2241                        snapshot,
2242                        gaps,
2243                        provider_contract_manifest,
2244                    } => {
2245                        let previous = node.cursor;
2246                        let mut inventory = SlimNodeInventory::from_snapshot(&snapshot);
2247                        inventory.provider_contracts =
2248                            provider_contract_manifest.provider_contracts;
2249                        inventory.provider_adapter_contracts =
2250                            provider_contract_manifest.provider_adapter_contracts;
2251                        node.transport = NodeTransportState::Online;
2252                        node.freshness = NodeFreshness::Fresh;
2253                        node.cursor = Some(cursor);
2254                        node.inventory = Some(inventory);
2255                        node.last_success_unix_ms = Some(attempt.at_unix_ms);
2256                        node.consecutive_failures = 0;
2257                        node.last_error = None;
2258                        for kind in gaps {
2259                            if node.gaps.len() == MAX_C2_GAPS_PER_NODE { node.gaps.remove(0); node.gaps_truncated += 1; }
2260                            node.gaps.push(NodeGap { kind, detected_at_unix_ms: attempt.at_unix_ms, previous, observed: cursor });
2261                        }
2262                    }
2263                    AttemptResult::Success { cursor, snapshot, gaps } => {
2264                        let previous = node.cursor;
2265                        let provider_contract_manifest = node.inventory.as_ref().map(|inventory| {
2266                            ProviderContractManifest {
2267                                provider_contracts: inventory.provider_contracts.clone(),
2268                                provider_adapter_contracts: inventory.provider_adapter_contracts.clone(),
2269                            }
2270                        }).unwrap_or_default();
2271                        let mut inventory = SlimNodeInventory::from_snapshot(&snapshot);
2272                        inventory.provider_contracts =
2273                            provider_contract_manifest.provider_contracts;
2274                        inventory.provider_adapter_contracts =
2275                            provider_contract_manifest.provider_adapter_contracts;
2276                        node.transport = NodeTransportState::Online;
2277                        node.freshness = NodeFreshness::Fresh;
2278                        node.cursor = Some(cursor);
2279                        node.inventory = Some(inventory);
2280                        node.last_success_unix_ms = Some(attempt.at_unix_ms);
2281                        node.consecutive_failures = 0;
2282                        node.last_error = None;
2283                        for kind in gaps {
2284                            if node.gaps.len() == MAX_C2_GAPS_PER_NODE { node.gaps.remove(0); node.gaps_truncated += 1; }
2285                            node.gaps.push(NodeGap { kind, detected_at_unix_ms: attempt.at_unix_ms, previous, observed: cursor });
2286                        }
2287                    }
2288                    AttemptResult::Cursor {
2289                        cursor,
2290                        gaps,
2291                        managed_worktree_events,
2292                    } => {
2293                        let previous = node.cursor;
2294                        let incarnation_changed = previous.is_some_and(|previous| {
2295                            previous.incarnation_id != cursor.incarnation_id
2296                        });
2297                        apply_managed_worktree_cursor(
2298                            node.inventory.as_mut(),
2299                            incarnation_changed,
2300                            &managed_worktree_events,
2301                        );
2302                        node.transport = NodeTransportState::Online;
2303                        node.freshness = NodeFreshness::Fresh;
2304                        node.cursor = Some(cursor);
2305                        node.last_success_unix_ms = Some(attempt.at_unix_ms);
2306                        node.consecutive_failures = 0;
2307                        node.last_error = None;
2308                        for kind in gaps {
2309                            if node.gaps.len() == MAX_C2_GAPS_PER_NODE { node.gaps.remove(0); node.gaps_truncated += 1; }
2310                            node.gaps.push(NodeGap { kind, detected_at_unix_ms: attempt.at_unix_ms, previous, observed: cursor });
2311                        }
2312                    }
2313                    AttemptResult::Failure { error, hard } => {
2314                        node.consecutive_failures = node.consecutive_failures.saturating_add(1);
2315                        node.transport = if hard || node.consecutive_failures >= 5 { NodeTransportState::Parked } else { NodeTransportState::Offline };
2316                        node.last_error = Some(error);
2317                    }
2318                }
2319                current.ready = attempted.len() == configured;
2320                refresh_freshness(&mut current, fresh_for);
2321                current.observed_at_unix_ms = unix_ms();
2322                status.send_replace(Arc::new(current.clone()));
2323            }
2324            _ = sleep(Duration::from_millis(250)) => {
2325                refresh_freshness(&mut current, fresh_for);
2326                current.observed_at_unix_ms = unix_ms();
2327                status.send_replace(Arc::new(current.clone()));
2328            }
2329            changed = shutdown.changed() => if changed.is_err() || *shutdown.borrow() { return Ok(()); },
2330        }
2331    }
2332}
2333
2334fn apply_managed_worktree_cursor(
2335    inventory: Option<&mut SlimNodeInventory>,
2336    incarnation_changed: bool,
2337    events: &[NodeEvent],
2338) {
2339    let Some(inventory) = inventory else { return; };
2340    if incarnation_changed {
2341        inventory.provider_runtime_statuses.clear();
2342        inventory.managed_worktrees.clear();
2343        inventory.managed_worktree_count = 0;
2344        inventory.managed_worktrees_truncated = false;
2345        return;
2346    }
2347    for event in events {
2348        match event {
2349            NodeEvent::ManagedWorktreeUpserted { lease } => {
2350                inventory.managed_worktrees.retain(|existing| {
2351                    existing.lease_id != lease.lease_id
2352                        && existing.workspace_id != lease.workspace_id
2353                });
2354                inventory.managed_worktrees.push(lease.clone());
2355                inventory.managed_worktrees.sort_by(|left, right| {
2356                    left.lease_id.cmp(&right.lease_id)
2357                });
2358                inventory.managed_worktree_count = inventory.managed_worktrees.len();
2359                inventory.managed_worktrees.truncate(
2360                    crate::protocol::MAX_C2_MANAGED_WORKTREES_PER_NODE,
2361                );
2362                inventory.managed_worktrees_truncated =
2363                    inventory.managed_worktrees.len() < inventory.managed_worktree_count;
2364            }
2365            NodeEvent::ManagedWorktreeRemoved { lease_id } => {
2366                let before = inventory.managed_worktrees.len();
2367                inventory
2368                    .managed_worktrees
2369                    .retain(|lease| &lease.lease_id != lease_id);
2370                if inventory.managed_worktrees.len() < before
2371                    || inventory.managed_worktrees_truncated
2372                {
2373                    inventory.managed_worktree_count =
2374                        inventory.managed_worktree_count.saturating_sub(1);
2375                }
2376                inventory.managed_worktrees_truncated =
2377                    inventory.managed_worktrees.len() < inventory.managed_worktree_count;
2378            }
2379            _ => {}
2380        }
2381    }
2382}
2383
2384fn refresh_freshness(status: &mut StatusResponse, fresh_for: Duration) {
2385    let now = unix_ms();
2386    let fresh_ms = fresh_for.as_millis().min(u64::MAX as u128) as u64;
2387    for node in status.nodes.values_mut() {
2388        node.freshness = match node.last_success_unix_ms {
2389            None => NodeFreshness::Unavailable,
2390            Some(last) if now.saturating_sub(last) <= fresh_ms => NodeFreshness::Fresh,
2391            Some(_) => NodeFreshness::Stale,
2392        };
2393    }
2394}
2395
2396async fn http_server(
2397    listener: TcpListener, token: String, io_deadline: Duration,
2398    status: watch::Receiver<Arc<StatusResponse>>, mut shutdown: watch::Receiver<bool>,
2399) -> io::Result<()> {
2400    let permits = Arc::new(Semaphore::new(MAX_HTTP_CONNECTIONS));
2401    let mut connections = JoinSet::new();
2402    loop {
2403        tokio::select! {
2404            accepted = listener.accept() => {
2405                let (stream, _) = accepted?;
2406                let Ok(permit) = Arc::clone(&permits).try_acquire_owned() else { drop(stream); continue; };
2407                let token = token.clone();
2408                let status = status.clone();
2409                connections.spawn(async move { let _permit = permit; let _ = serve_http(stream, &token, io_deadline, &status).await; });
2410            }
2411            changed = shutdown.changed() => if changed.is_err() || *shutdown.borrow() { break; },
2412        }
2413        while let Some(result) = connections.try_join_next() { result.map_err(io::Error::other)?; }
2414    }
2415    connections.shutdown().await;
2416    Ok(())
2417}
2418
2419async fn serve_http(mut stream: TcpStream, token: &str, deadline: Duration, status: &watch::Receiver<Arc<StatusResponse>>) -> io::Result<()> {
2420    let request = match timeout(deadline, read_request(&mut stream)).await {
2421        Ok(Ok(request)) => request,
2422        Ok(Err(ReadError::TooLarge)) => return write_response(&mut stream, deadline, Response::plain(413, "Payload Too Large")).await,
2423        Ok(Err(ReadError::Io(error))) => return Err(error),
2424        _ => return Ok(()),
2425    };
2426    let response = route(request, token, status.borrow().as_ref());
2427    write_response(&mut stream, deadline, response).await
2428}
2429
2430struct Request { method: String, path: String, authorization: Option<String> }
2431enum ReadError { Closed, Invalid, TooLarge, Io(io::Error) }
2432
2433async fn read_request(stream: &mut TcpStream) -> Result<Request, ReadError> {
2434    let mut bytes = Vec::with_capacity(1024);
2435    let mut chunk = [0_u8; 1024];
2436    loop {
2437        let count = stream.read(&mut chunk).await.map_err(ReadError::Io)?;
2438        if count == 0 { return Err(ReadError::Closed); }
2439        if bytes.len().saturating_add(count) > HEADER_LIMIT_BYTES { return Err(ReadError::TooLarge); }
2440        bytes.extend_from_slice(&chunk[..count]);
2441        if bytes.windows(4).any(|window| window == b"\r\n\r\n") { break; }
2442    }
2443    let text = std::str::from_utf8(&bytes).map_err(|_| ReadError::Invalid)?;
2444    let mut lines = text.split("\r\n");
2445    let mut first = lines.next().ok_or(ReadError::Invalid)?.split_whitespace();
2446    let method = first.next().ok_or(ReadError::Invalid)?;
2447    let path = first.next().ok_or(ReadError::Invalid)?;
2448    let version = first.next().ok_or(ReadError::Invalid)?;
2449    if first.next().is_some() || !version.starts_with("HTTP/1.") || !path.starts_with('/') { return Err(ReadError::Invalid); }
2450    let mut authorization = None;
2451    for line in lines {
2452        if line.is_empty() { break; }
2453        let (name, value) = line.split_once(':').ok_or(ReadError::Invalid)?;
2454        if name.eq_ignore_ascii_case("authorization") {
2455            if authorization.is_some() { return Err(ReadError::Invalid); }
2456            authorization = Some(value.trim().to_owned());
2457        }
2458    }
2459    Ok(Request { method: method.to_owned(), path: path.to_owned(), authorization })
2460}
2461
2462fn route(request: Request, token: &str, status: &StatusResponse) -> Response {
2463    if request.method != "GET" { return Response::plain(405, "Method Not Allowed").allow_get(); }
2464    let path = request.path.split_once('?').map_or(request.path.as_str(), |pair| pair.0);
2465    match path {
2466        "/health" => Response::json(200, &HealthResponse { ok: true, service: "gate4agent-c2".to_owned(), api_version: C2_API_VERSION, pid: std::process::id(), version: env!("CARGO_PKG_VERSION").to_owned() }),
2467        "/ready" => {
2468            let online_nodes = status.nodes.values().filter(|node| node.transport == NodeTransportState::Online).count();
2469            let offline_nodes = status.nodes.values().filter(|node| node.transport == NodeTransportState::Offline).count();
2470            let parked_nodes = status.nodes.values().filter(|node| node.transport == NodeTransportState::Parked).count();
2471            let body = ReadyResponse { ready: status.ready, api_version: C2_API_VERSION, configured_nodes: status.nodes.len(), attempted_nodes: status.nodes.values().filter(|node| node.last_attempt_unix_ms.is_some()).count(), online_nodes, offline_nodes, parked_nodes };
2472            Response::json(if status.ready { 200 } else { 503 }, &body)
2473        }
2474        "/status" => {
2475            if !authorized(request.authorization.as_deref(), token) { Response::plain(401, "Unauthorized").authenticate() }
2476            else { Response::json(200, status) }
2477        }
2478        _ => Response::plain(404, "Not Found"),
2479    }
2480}
2481
2482fn authorized(header: Option<&str>, token: &str) -> bool {
2483    let Some((scheme, candidate)) = header.and_then(|value| value.split_once(' ')) else { return false; };
2484    scheme.eq_ignore_ascii_case("bearer") && constant_time_eq(candidate.as_bytes(), token.as_bytes())
2485}
2486
2487fn constant_time_eq(left: &[u8], right: &[u8]) -> bool {
2488    if left.len() != right.len() { return false; }
2489    left.iter().zip(right).fold(0_u8, |difference, (left, right)| difference | (left ^ right)) == 0
2490}
2491
2492struct Response { status: u16, reason: &'static str, content_type: &'static str, body: Vec<u8>, headers: Vec<(&'static str, &'static str)> }
2493impl Response {
2494    fn plain(status: u16, reason: &'static str) -> Self { Self { status, reason, content_type: "text/plain; charset=utf-8", body: reason.as_bytes().to_vec(), headers: Vec::new() } }
2495    fn json<T: serde::Serialize>(status: u16, value: &T) -> Self {
2496        let body = serde_json::to_vec(value).expect("C2 DTO must serialize");
2497        if body.len() > RESPONSE_BODY_LIMIT_BYTES { return Self::plain(503, "Service Unavailable"); }
2498        let reason = if status == 200 { "OK" } else { "Service Unavailable" };
2499        Self { status, reason, content_type: "application/json", body, headers: Vec::new() }
2500    }
2501    fn allow_get(mut self) -> Self { self.headers.push(("Allow", "GET")); self }
2502    fn authenticate(mut self) -> Self { self.headers.push(("WWW-Authenticate", "Bearer")); self }
2503}
2504
2505async fn write_response(stream: &mut TcpStream, deadline: Duration, response: Response) -> io::Result<()> {
2506    let mut headers = format!("HTTP/1.1 {} {}\r\nContent-Type: {}\r\nContent-Length: {}\r\nConnection: close\r\n", response.status, response.reason, response.content_type, response.body.len());
2507    for (name, value) in response.headers { headers.push_str(name); headers.push_str(": "); headers.push_str(value); headers.push_str("\r\n"); }
2508    headers.push_str("\r\n");
2509    timeout(deadline, async { stream.write_all(headers.as_bytes()).await?; stream.write_all(&response.body).await?; stream.shutdown().await }).await
2510        .map_err(|_| io::Error::new(io::ErrorKind::TimedOut, "C2 HTTP write timed out"))?
2511}
2512
2513fn unix_ms() -> u64 {
2514    SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_millis().min(u64::MAX as u128) as u64
2515}
2516
2517#[cfg(test)]
2518mod endpoint_tests {
2519    use super::*;
2520    use crate::protocol::{C2RelayRoute, C2Topology};
2521
2522    fn node(endpoint: &str) -> Result<C2NodeConfig, C2ConfigError> {
2523        C2NodeConfig::new(NodeId::new("remote-node").unwrap(), endpoint, "safe-token")
2524    }
2525
2526    #[cfg(windows)]
2527    const LOCAL_ENDPOINT: &str = r"\\.\pipe\relay-route-fact";
2528    #[cfg(unix)]
2529    const LOCAL_ENDPOINT: &str = "/tmp/gate4agent-relay-route-fact.sock";
2530
2531    fn projected_route(node: &C2NodeConfig, transport: NodeTransportState) -> C2RelayRoute {
2532        let mut observed = initial_observed_node(node);
2533        observed.transport = transport;
2534        let status = StatusResponse {
2535            api_version: C2_API_VERSION,
2536            ready: false,
2537            observed_at_unix_ms: 1,
2538            nodes: BTreeMap::from([(node.node_id.clone(), observed)]),
2539        };
2540        C2Topology::from_status(&status).nodes[0].relay_route
2541    }
2542
2543    #[test]
2544    fn relay_route_fact_projects_exactly_and_survives_offline_and_parked() {
2545        let local = node(LOCAL_ENDPOINT).unwrap();
2546        let ssh = node("tcp://127.0.0.1:48100").unwrap();
2547
2548        for transport in [
2549            NodeTransportState::Online,
2550            NodeTransportState::Offline,
2551            NodeTransportState::Parked,
2552        ] {
2553            assert_eq!(projected_route(&local, transport), C2RelayRoute::LocalIpc);
2554            assert_eq!(
2555                projected_route(&ssh, transport),
2556                C2RelayRoute::SshForwardedLoopback,
2557            );
2558        }
2559    }
2560
2561    /// `accept` is the assignment that names no address, and the config
2562    /// has to treat it as a different KIND of route rather than as a
2563    /// malformed endpoint: it is how a node that cannot be dialled is
2564    /// declared.
2565    #[test]
2566    fn a_node_that_waits_to_be_called_is_declared_by_name_not_by_address() {
2567        let waiting = node("accept").unwrap();
2568        assert_eq!(waiting.route, C2NodeRoute::CallHome);
2569        assert_eq!(waiting.transport_label(), "call-home");
2570
2571        // Every waiting node carries the same `accept`, so the
2572        // duplicate-endpoint check must not see them as two nodes fighting
2573        // over one socket. Their uniqueness is `node_id`'s job, and that
2574        // check still applies.
2575        let second = C2NodeConfig::new(
2576            NodeId::new("second-waiting-node").unwrap(),
2577            "accept",
2578            "safe-token",
2579        )
2580        .unwrap();
2581        let config = C2Config::new(
2582            "127.0.0.1:0".parse().unwrap(),
2583            "safe-token",
2584            vec![waiting.clone(), second],
2585        )
2586        .expect("two waiting nodes are not a duplicate endpoint");
2587        assert_eq!(config.nodes.len(), 2);
2588    }
2589
2590    /// A node that waits for a call the relay never listens for is a
2591    /// deployment that can never work, so it is refused at startup instead
2592    /// of becoming a node that is silently offline forever.
2593    #[test]
2594    fn waiting_for_a_call_requires_somewhere_to_be_called() {
2595        let waiting = node("accept").unwrap();
2596        let config =
2597            C2Config::new("127.0.0.1:0".parse().unwrap(), "safe-token", vec![waiting]).unwrap();
2598        assert!(matches!(
2599            config.validate_call_home(),
2600            Err(C2ConfigError::CallHomeWithoutListener(_)),
2601        ));
2602        let config = config.with_node_listen("127.0.0.1:48200".parse().unwrap()).unwrap();
2603        assert!(config.validate_call_home().is_ok());
2604    }
2605
2606    /// The call-home listener holds the same line every other listener in
2607    /// this stack holds. The wire authenticates both ends and encrypts
2608    /// nothing, so accepting node connections from off-box would put
2609    /// terminal contents and keystrokes on the network in the clear.
2610    #[test]
2611    fn the_call_home_listener_refuses_to_leave_loopback() {
2612        let waiting = node("accept").unwrap();
2613        let config =
2614            C2Config::new("127.0.0.1:9000".parse().unwrap(), "safe-token", vec![waiting]).unwrap();
2615        for refused in ["0.0.0.0:48200", "192.168.1.10:48200", "127.0.0.1:0"] {
2616            assert!(
2617                matches!(
2618                    config.clone().with_node_listen(refused.parse().unwrap()),
2619                    Err(C2ConfigError::NonLoopbackNodeListen(_)),
2620                ),
2621                "accepted {refused}",
2622            );
2623        }
2624        // Sharing the API's own address would mean HTTP and node frames
2625        // arriving on one socket; neither parser would survive the other.
2626        assert!(matches!(
2627            config.with_node_listen("127.0.0.1:9000".parse().unwrap()),
2628            Err(C2ConfigError::NodeListenConflict),
2629        ));
2630    }
2631
2632    #[test]
2633    fn ssh_forwarded_loopback_route_is_strict_canonical_and_control_stays_local() {
2634        let ipv4 = node("tcp://127.0.0.1:48100").unwrap();
2635        assert_eq!(ipv4.endpoint, "tcp://127.0.0.1:48100");
2636        assert_eq!(ipv4.route, C2NodeRoute::SshForwardedLoopback("127.0.0.1:48100".parse().unwrap()));
2637        assert_eq!(ipv4.transport_label(), "ssh-forwarded-loopback");
2638
2639        let ipv6 = node("tcp://[0:0:0:0:0:0:0:1]:48100").unwrap();
2640        assert_eq!(ipv6.endpoint, "tcp://[::1]:48100");
2641        assert_eq!(ipv6.transport_label(), "ssh-forwarded-loopback");
2642
2643        for invalid in [
2644            "tcp://localhost:48100",
2645            "tcp://127.0.0.2:48100",
2646            "tcp://0.0.0.0:48100",
2647            "tcp://127.0.0.1:0",
2648            "tcp://user@127.0.0.1:48100",
2649            "tcp://127.0.0.1:48100/path",
2650            "TCP://127.0.0.1:48100",
2651        ] {
2652            assert!(matches!(node(invalid), Err(C2ConfigError::InvalidEndpoint(_))), "accepted {invalid}");
2653        }
2654
2655        assert!(matches!(
2656            C2Config::new(
2657                "127.0.0.1:0".parse().unwrap(),
2658                "safe-token",
2659                vec![ipv4.clone(), C2NodeConfig::new(
2660                    NodeId::new("duplicate-route").unwrap(),
2661                    "tcp://127.0.0.1:48100",
2662                    "safe-token",
2663                ).unwrap()],
2664            ),
2665            Err(C2ConfigError::DuplicateEndpoint)
2666        ));
2667        assert!(matches!(
2668            C2Config::new("127.0.0.1:0".parse().unwrap(), "safe-token", vec![ipv4])
2669                .unwrap()
2670                .with_control_endpoint("tcp://127.0.0.1:48101"),
2671            Err(C2ConfigError::InvalidControlEndpoint)
2672        ));
2673    }
2674}
2675
2676#[cfg(all(test, windows))]
2677mod tests {
2678    use super::*;
2679    use gate4agent_node_protocol::{
2680        AgentId, AgentStreamChunkKindV1, AgentStreamChunkV1, CapabilityId, SessionAddress,
2681        SessionKey, SessionMode, SessionRecordId,
2682        HarnessMcpActivationDigest, HarnessMcpCallId, HarnessMcpContentTypeV1,
2683        HarnessMcpOpaquePayloadV1, HarnessMcpReservationId,
2684        SpawnDeadlineMs, SpawnFieldProvenance, SpawnIdempotencyKey, SpawnOverrides,
2685        SpawnProfileId, SpawnProfileRevision, SpawnPrompt, SpawnPromptMetadata,
2686        SpawnRequiredCapabilities, SpawnResolutionProvenance, SpawnTarget, WorkspaceId,
2687        ManagedWorktreeCleanupFailure, ManagedWorktreeLeaseId,
2688        ManagedWorktreeLeaseSnapshot, ManagedWorktreeRetention, ManagedWorktreeSpawnReceipt,
2689        WorktreeProfileId, WorktreeProfileRevision,
2690    };
2691    use gate4agent_types::{AgentInstanceId, SessionGeneration, TerminalSize};
2692    use std::collections::BTreeMap;
2693
2694    fn agent(value: &str) -> gate4agent_node_protocol::AgentId {
2695        gate4agent_node_protocol::AgentId::new(value).unwrap()
2696    }
2697
2698    fn managed_lease(
2699        lease_id: &str,
2700        workspace_id: &str,
2701        state: ManagedWorktreeLeaseState,
2702    ) -> ManagedWorktreeLeaseSnapshot {
2703        let in_use = state == ManagedWorktreeLeaseState::InUse;
2704        ManagedWorktreeLeaseSnapshot {
2705            lease_id: ManagedWorktreeLeaseId::new(lease_id).unwrap(),
2706            source_workspace_id: WorkspaceId::new("repo").unwrap(),
2707            workspace_id: WorkspaceId::new(workspace_id).unwrap(),
2708            profile_id: WorktreeProfileId::new("review").unwrap(),
2709            profile_revision: WorktreeProfileRevision::new("review.r1").unwrap(),
2710            retention: ManagedWorktreeRetention::RemoveWhenReleased,
2711            state,
2712            active_session_count: u16::from(in_use),
2713            managed_record_count: u16::from(in_use),
2714            cleanup_failure: None::<ManagedWorktreeCleanupFailure>,
2715            created_at_unix_ms: 1,
2716            updated_at_unix_ms: 2,
2717        }
2718    }
2719
2720    #[test]
2721    fn spawn_spec_receipt_correlation_rejects_mismatches_before_forwarding() {
2722        let incarnation_id = NodeIncarnationId::from_bytes([7; 16]);
2723        let terminal_size = TerminalSize { rows: 24, columns: 80 };
2724        let prompt = SpawnPrompt::new("hi").unwrap();
2725        let required_capabilities = SpawnRequiredCapabilities::new([
2726            CapabilityId::new("raw-pty-lifecycle").unwrap(),
2727        ]).unwrap();
2728        let spec = SpawnSpec {
2729            target: SpawnTarget {
2730                node_id: NodeId::new("node-a").unwrap(),
2731                workspace_id: WorkspaceId::new("repo").unwrap(),
2732                worktree_id: None,
2733            },
2734            profile_id: SpawnProfileId::new("default").unwrap(),
2735            expected_profile_revision: SpawnProfileRevision::new("r1").unwrap(),
2736            overrides: SpawnOverrides {
2737                provider: SpawnOverride::Set { value: AgentId::new("codex").unwrap() },
2738                mode: SpawnOverride::Set { value: SessionMode::Pty },
2739                terminal_size: SpawnOverride::Set { value: terminal_size },
2740                prompt: SpawnOverride::Set { value: prompt.clone() },
2741                bundle_id: SpawnOverride::Clear,
2742                context_id: SpawnOverride::Clear,
2743                environment_profile_id: SpawnOverride::Clear,
2744                approval_level: None,
2745                network_allowlist: None,
2746                browser_profile_id: None,
2747            },
2748            deadline_ms: SpawnDeadlineMs::new(5_000).unwrap(),
2749            idempotency_key: SpawnIdempotencyKey::new("spawn-1").unwrap(),
2750            required_capabilities: required_capabilities.clone(),
2751        };
2752        let receipt = ResolvedSpawnReceipt {
2753            incarnation_id,
2754            session: SessionAddress {
2755                workspace_id: spec.target.workspace_id.clone(),
2756                session: SessionKey {
2757                    instance_id: AgentInstanceId(7),
2758                    generation: SessionGeneration(1),
2759                },
2760            },
2761            target: spec.target.clone(),
2762            profile_id: spec.profile_id.clone(),
2763            profile_revision: SpawnProfileRevision::new("r1").unwrap(),
2764            provider: AgentId::new("codex").unwrap(),
2765            mode: SessionMode::Pty,
2766            terminal_size,
2767            prompt: SpawnPromptMetadata::from_prompt(Some(&prompt)),
2768            bundle_id: None,
2769            bundle: None,
2770            context_id: None,
2771            context: None,
2772            environment_profile: None,
2773            deadline_ms: spec.deadline_ms,
2774            idempotency_key: spec.idempotency_key.clone(),
2775            required_capabilities,
2776            provenance: SpawnResolutionProvenance {
2777                provider: SpawnFieldProvenance::Override,
2778                mode: SpawnFieldProvenance::Override,
2779                terminal_size: SpawnFieldProvenance::Override,
2780                prompt: SpawnFieldProvenance::Override,
2781                bundle_id: SpawnFieldProvenance::Cleared,
2782                context_id: SpawnFieldProvenance::Cleared,
2783                environment_profile_id: SpawnFieldProvenance::Cleared,
2784            },
2785            harness_mcp_proxy: None,
2786        };
2787        let response = |receipt| Ok(NodeResponse::SpawnSpecAccepted { receipt });
2788        let expected = ExpectedSpawnRequest::Spec(spec.clone());
2789        assert!(validate_spawn_spec_response(
2790            Some(&expected),
2791            &response(receipt.clone()),
2792            incarnation_id,
2793        ).is_ok());
2794
2795        let mut mismatches = Vec::new();
2796        let mut changed = receipt.clone();
2797        changed.incarnation_id = NodeIncarnationId::from_bytes([8; 16]);
2798        mismatches.push(changed);
2799        let mut changed = receipt.clone();
2800        changed.target.node_id = NodeId::new("node-b").unwrap();
2801        mismatches.push(changed);
2802        let mut changed = receipt.clone();
2803        changed.profile_id = SpawnProfileId::new("other").unwrap();
2804        mismatches.push(changed);
2805        let mut changed = receipt.clone();
2806        changed.idempotency_key = SpawnIdempotencyKey::new("spawn-2").unwrap();
2807        mismatches.push(changed);
2808        let mut changed = receipt.clone();
2809        changed.deadline_ms = SpawnDeadlineMs::new(4_999).unwrap();
2810        mismatches.push(changed);
2811        let mut changed = receipt.clone();
2812        changed.required_capabilities = SpawnRequiredCapabilities::default();
2813        mismatches.push(changed);
2814        let mut changed = receipt.clone();
2815        changed.session.workspace_id = WorkspaceId::new("other").unwrap();
2816        mismatches.push(changed);
2817        let mut changed = receipt.clone();
2818        changed.provider = AgentId::new("claude").unwrap();
2819        mismatches.push(changed);
2820        let mut changed = receipt.clone();
2821        changed.prompt = SpawnPromptMetadata { present: false, byte_len: 0 };
2822        mismatches.push(changed);
2823        let mut changed = receipt.clone();
2824        changed.environment_profile = Some(
2825            gate4agent_node_protocol::ResolvedEnvironmentProfileReceipt {
2826                profile_id:
2827                    gate4agent_node_protocol::SpawnEnvironmentProfileId::new(
2828                        "local-default",
2829                    )
2830                    .unwrap(),
2831                profile_revision:
2832                    gate4agent_node_protocol::SpawnEnvironmentProfileRevision::new(
2833                        "local-default.r1",
2834                    )
2835                    .unwrap(),
2836                    network_allowlist: None,
2837                    browser_profile_id: None,
2838            },
2839        );
2840        mismatches.push(changed);
2841        let mut changed = receipt.clone();
2842        changed.bundle = Some(gate4agent_node_protocol::ResolvedBundleReceipt {
2843            id: gate4agent_node_protocol::SpawnBundleId::new("unexpected-bundle")
2844                .unwrap(),
2845            revision: gate4agent_node_protocol::SpawnBundleRevision::new(
2846                "unexpected-bundle.r1",
2847            )
2848            .unwrap(),
2849            digest: gate4agent_node_protocol::SpawnBundleDigest::new(format!(
2850                "sha256:{}",
2851                "a".repeat(64),
2852            ))
2853            .unwrap(),
2854        });
2855        mismatches.push(changed);
2856        let mut changed = receipt.clone();
2857        changed.context = Some(gate4agent_node_protocol::ResolvedContextPackReceipt {
2858            id: gate4agent_node_protocol::SpawnContextId::new("unexpected-context")
2859                .unwrap(),
2860            digest: gate4agent_node_protocol::SpawnContextDigest::new(format!(
2861                "sha256:{}",
2862                "b".repeat(64),
2863            ))
2864            .unwrap(),
2865            lineage: gate4agent_node_protocol::ContextPackLineageReceipt {
2866                source_node_id: NodeId::new("node-a").unwrap(),
2867                source_session: receipt.session.clone(),
2868                source_provider: AgentId::new("codex").unwrap(),
2869            },
2870            source_message_count: 1,
2871            retained_message_count: 1,
2872            byte_len: 16,
2873            truncated: false,
2874        });
2875        mismatches.push(changed);
2876
2877        for mismatch in mismatches {
2878            assert!(validate_spawn_spec_response(
2879                Some(&expected),
2880                &response(mismatch),
2881                incarnation_id,
2882            ).is_err());
2883        }
2884
2885        let mut environment_spec = spec;
2886        environment_spec.overrides.environment_profile_id = SpawnOverride::Set {
2887            value: gate4agent_node_protocol::SpawnEnvironmentProfileId::new(
2888                "local-default",
2889            )
2890            .unwrap(),
2891        };
2892        let environment_expected = ExpectedSpawnRequest::Spec(environment_spec);
2893        let mut environment_receipt = receipt;
2894        environment_receipt.environment_profile = Some(
2895            gate4agent_node_protocol::ResolvedEnvironmentProfileReceipt {
2896                profile_id:
2897                    gate4agent_node_protocol::SpawnEnvironmentProfileId::new(
2898                        "local-default",
2899                    )
2900                    .unwrap(),
2901                profile_revision:
2902                    gate4agent_node_protocol::SpawnEnvironmentProfileRevision::new(
2903                        "local-default.r1",
2904                    )
2905                    .unwrap(),
2906                    network_allowlist: None,
2907                    browser_profile_id: None,
2908            },
2909        );
2910        assert!(validate_spawn_spec_response(
2911            Some(&environment_expected),
2912            &response(environment_receipt.clone()),
2913            incarnation_id,
2914        )
2915        .is_ok());
2916        environment_receipt.environment_profile.as_mut().unwrap().profile_id =
2917            gate4agent_node_protocol::SpawnEnvironmentProfileId::new("other")
2918                .unwrap();
2919        assert!(validate_spawn_spec_response(
2920            Some(&environment_expected),
2921            &response(environment_receipt),
2922            incarnation_id,
2923        )
2924        .is_err());
2925    }
2926
2927    #[test]
2928    fn provider_session_index_correlation_accepts_exact_identity_and_rejects_mismatches() {
2929        let identity = gate4agent_types::ProviderSessionIdentity {
2930            key: gate4agent_types::ProviderSessionKey::SessionId,
2931            id: "native-session-1".to_owned(),
2932            transcript_path: Some(r"C:\provider\sessions\native-session-1.jsonl".to_owned()),
2933        };
2934        let expected = NodeRequest::IndexProviderSession {
2935            workspace_id: WorkspaceId::new("primary").unwrap(),
2936            provider: AgentId::new("codex").unwrap(),
2937            identity: identity.clone(),
2938            display_name: "release shepherd".to_owned(),
2939        };
2940        let NodeRequest::IndexProviderSession { workspace_id, provider, .. } = &expected else {
2941            unreachable!();
2942        };
2943        let record = gate4agent_node_protocol::ManagedSessionRecord {
2944            record_id: SessionRecordId::new("session-001").unwrap(),
2945            display_name: "release shepherd".to_owned(),
2946            provider: provider.clone(),
2947            mode: SessionMode::Pty,
2948            state: gate4agent_node_protocol::ManagedSessionState::Dormant,
2949            workspace_id: workspace_id.clone(),
2950            canonical_root: gate4agent_node_protocol::OpaqueHostPath::utf8(
2951                r"C:\repo".to_owned(),
2952            ).unwrap(),
2953            provider_session: Some(identity),
2954            active_session: None,
2955            environment_profile: None,
2956            bundle: None,
2957            context_id: None,
2958            context: None,
2959            exported_context: None,
2960            task_binding: None,
2961            created_at_unix_ms: 1,
2962            updated_at_unix_ms: 2,
2963            last_error: None,
2964        };
2965        let response = |record| Ok(NodeResponse::ProviderSessionIndexed { record });
2966
2967        assert!(validate_provider_session_index_response(
2968            Some(&expected),
2969            &response(record.clone()),
2970        ).is_ok());
2971        assert!(validate_provider_session_index_response(
2972            None,
2973            &response(record.clone()),
2974        ).is_err());
2975
2976        let mut identity_mismatch = record.clone();
2977        identity_mismatch.provider_session.as_mut().unwrap().id =
2978            "native-session-2".to_owned();
2979        assert!(validate_provider_session_index_response(
2980            Some(&expected),
2981            &response(identity_mismatch),
2982        ).is_err());
2983
2984        let mut workspace_mismatch = record.clone();
2985        workspace_mismatch.workspace_id = WorkspaceId::new("other").unwrap();
2986        assert!(validate_provider_session_index_response(
2987            Some(&expected),
2988            &response(workspace_mismatch),
2989        ).is_err());
2990
2991        let mut provider_mismatch = record;
2992        provider_mismatch.provider = AgentId::new("claude").unwrap();
2993        assert!(validate_provider_session_index_response(
2994            Some(&expected),
2995            &response(provider_mismatch),
2996        ).is_err());
2997    }
2998
2999    #[test]
3000    fn native_session_route_correlation_rejects_mismatches_before_projection() {
3001        let route = gate4agent_node_protocol::NativeSessionCatalogRoute::workspace(
3002            WorkspaceId::new("primary").unwrap(),
3003            AgentId::new("codex").unwrap(),
3004        );
3005        let selection = gate4agent_node_protocol::NativeSessionSelection {
3006            route: route.clone(),
3007            catalog_revision: 7,
3008            recent_cutoff_unix_ms: 70,
3009            selection_id: "selection-7".to_owned(),
3010        };
3011        let catalog = NodeRequest::CatalogNativeSessions {
3012            route: route.clone(),
3013            limit: 10,
3014        };
3015        let catalog_response = Ok(NodeResponse::NativeSessionsCataloged {
3016            route: route.clone(),
3017            entries: Vec::new(),
3018            summary: None,
3019        });
3020        assert!(validate_native_session_response(Some(&catalog), &catalog_response).is_ok());
3021        let mismatched_route = gate4agent_node_protocol::NativeSessionCatalogRoute::workspace(
3022            WorkspaceId::new("other").unwrap(),
3023            AgentId::new("codex").unwrap(),
3024        );
3025        assert!(validate_native_session_response(
3026            Some(&catalog),
3027            &Ok(NodeResponse::NativeSessionsCataloged {
3028                route: mismatched_route,
3029                entries: Vec::new(),
3030                summary: None,
3031            }),
3032        )
3033        .is_err());
3034
3035        let page = NodeRequest::PageNativeSessions {
3036            route: route.clone(),
3037            window: gate4agent_node_protocol::NativeSessionCatalogWindow::Recent,
3038            catalog_revision: 7,
3039            recent_cutoff_unix_ms: 70,
3040            after_selection_id: None,
3041            limit: 10,
3042        };
3043        let page_response = |revision| {
3044            Ok(NodeResponse::NativeSessionsPaged {
3045                route: route.clone(),
3046                page: gate4agent_node_protocol::NativeSessionCatalogPage {
3047                    window: gate4agent_node_protocol::NativeSessionCatalogWindow::Recent,
3048                    revision,
3049                    entries: Vec::new(),
3050                    next_after_selection_id: None,
3051                    remaining_count: 0,
3052                    has_more: false,
3053                },
3054            })
3055        };
3056        assert!(validate_native_session_response(Some(&page), &page_response(7)).is_ok());
3057        assert!(validate_native_session_response(Some(&page), &page_response(8)).is_err());
3058
3059        let preview = NodeRequest::PreviewNativeSession {
3060            selection: selection.clone(),
3061            message_limit: 10,
3062        };
3063        let preview_response = |selection| {
3064            Ok(NodeResponse::NativeSessionPreviewed {
3065                selection,
3066                preview: gate4agent_node_protocol::SessionRecordPreview {
3067                    title: None,
3068                    modified_at_unix_ms: None,
3069                    model: None,
3070                    message_count: 0,
3071                    message_count_exact: true,
3072                    completed_turn_count: None,
3073                    total_tokens: None,
3074                    truncated: false,
3075                    messages: Vec::new(),
3076                },
3077            })
3078        };
3079        assert!(validate_native_session_response(
3080            Some(&preview),
3081            &preview_response(selection.clone()),
3082        )
3083        .is_ok());
3084        let mut mismatched_selection = selection.clone();
3085        mismatched_selection.catalog_revision = 8;
3086        assert!(validate_native_session_response(
3087            Some(&preview),
3088            &preview_response(mismatched_selection),
3089        )
3090        .is_err());
3091
3092        let index = NodeRequest::IndexNativeSession {
3093            selection: selection.clone(),
3094            display_name: "Indexed".to_owned(),
3095        };
3096        let record = gate4agent_node_protocol::ManagedSessionRecord {
3097            record_id: SessionRecordId::new("session-007").unwrap(),
3098            display_name: "Indexed".to_owned(),
3099            provider: route.provider.clone(),
3100            mode: SessionMode::Pty,
3101            state: gate4agent_node_protocol::ManagedSessionState::Dormant,
3102            workspace_id: route.workspace_id.clone().unwrap(),
3103            canonical_root: gate4agent_node_protocol::OpaqueHostPath::utf8(
3104                r"C:\repo".to_owned(),
3105            )
3106            .unwrap(),
3107            provider_session: None,
3108            active_session: None,
3109            environment_profile: None,
3110            bundle: None,
3111            context_id: None,
3112            context: None,
3113            exported_context: None,
3114            task_binding: None,
3115            created_at_unix_ms: 1,
3116            updated_at_unix_ms: 2,
3117            last_error: None,
3118        };
3119        assert!(validate_native_session_response(
3120            Some(&index),
3121            &Ok(NodeResponse::NativeSessionIndexed {
3122                selection: selection.clone(),
3123                record: record.clone(),
3124            }),
3125        )
3126        .is_ok());
3127        assert!(validate_native_session_response(
3128            Some(&index),
3129            &Ok(NodeResponse::ProviderSessionIndexed {
3130                record: record.clone(),
3131            }),
3132        )
3133        .is_err());
3134        let mut wrong_echo = selection.clone();
3135        wrong_echo.catalog_revision = 8;
3136        assert!(validate_native_session_response(
3137            Some(&index),
3138            &Ok(NodeResponse::NativeSessionIndexed {
3139                selection: wrong_echo,
3140                record: record.clone(),
3141            }),
3142        )
3143        .is_err());
3144        let mut wrong_provider = record.clone();
3145        wrong_provider.provider = AgentId::new("claude").unwrap();
3146        assert!(validate_native_session_response(
3147            Some(&index),
3148            &Ok(NodeResponse::NativeSessionIndexed {
3149                selection: selection.clone(),
3150                record: wrong_provider.clone(),
3151            }),
3152        )
3153        .is_err());
3154        let mut wrong_workspace = record.clone();
3155        wrong_workspace.workspace_id = WorkspaceId::new("other").unwrap();
3156        assert!(validate_native_session_response(
3157            Some(&index),
3158            &Ok(NodeResponse::NativeSessionIndexed {
3159                selection: selection.clone(),
3160                record: wrong_workspace,
3161            }),
3162        )
3163        .is_err());
3164
3165        let external = NodeRequest::IndexNativeSession {
3166            selection: gate4agent_node_protocol::NativeSessionSelection {
3167                route: gate4agent_node_protocol::NativeSessionCatalogRoute::unregistered(
3168                    AgentId::new("codex").unwrap(),
3169                ),
3170                ..selection
3171            },
3172            display_name: "External".to_owned(),
3173        };
3174        let NodeRequest::IndexNativeSession {
3175            selection: external_selection,
3176            ..
3177        } = &external else {
3178            unreachable!();
3179        };
3180        assert!(validate_native_session_response(
3181            Some(&external),
3182            &Ok(NodeResponse::NativeSessionIndexed {
3183                selection: external_selection.clone(),
3184                record: gate4agent_node_protocol::ManagedSessionRecord {
3185                    provider: AgentId::new("codex").unwrap(),
3186                    ..wrong_provider
3187                },
3188            }),
3189        )
3190        .is_err());
3191    }
3192
3193    #[test]
3194    fn managed_spawn_receipt_correlation_covers_legacy_and_v2_requests() {
3195        let incarnation_id = NodeIncarnationId::from_bytes([7; 16]);
3196        let spec = SpawnSpec {
3197            target: SpawnTarget {
3198                node_id: NodeId::new("node-a").unwrap(),
3199                workspace_id: WorkspaceId::new("repo").unwrap(),
3200                worktree_id: None,
3201            },
3202            profile_id: SpawnProfileId::new("default").unwrap(),
3203            expected_profile_revision:
3204                SpawnProfileRevision::new("default.r1").unwrap(),
3205            overrides: SpawnOverrides::default(),
3206            deadline_ms: SpawnDeadlineMs::new(5_000).unwrap(),
3207            idempotency_key: SpawnIdempotencyKey::new("managed-1").unwrap(),
3208            required_capabilities: SpawnRequiredCapabilities::default(),
3209        };
3210        let managed = ManagedWorktreeSpawnRequest {
3211            spawn_spec: spec.clone(),
3212            worktree_profile_id: WorktreeProfileId::new("review").unwrap(),
3213        };
3214        let workspace_id = WorkspaceId::new("managed-a").unwrap();
3215        let spawn = ResolvedSpawnReceipt {
3216            incarnation_id,
3217            session: SessionAddress {
3218                workspace_id: workspace_id.clone(),
3219                session: SessionKey {
3220                    instance_id: AgentInstanceId(8),
3221                    generation: SessionGeneration(1),
3222                },
3223            },
3224            target: SpawnTarget {
3225                node_id: spec.target.node_id.clone(),
3226                workspace_id: spec.target.workspace_id.clone(),
3227                worktree_id: Some(workspace_id.clone()),
3228            },
3229            profile_id: spec.profile_id.clone(),
3230            profile_revision: SpawnProfileRevision::new("default.r1").unwrap(),
3231            provider: AgentId::new("claude").unwrap(),
3232            mode: SessionMode::Pty,
3233            terminal_size: TerminalSize { rows: 24, columns: 80 },
3234            prompt: SpawnPromptMetadata { present: false, byte_len: 0 },
3235            bundle_id: None,
3236            bundle: None,
3237            context_id: None,
3238            context: None,
3239            environment_profile: None,
3240            deadline_ms: spec.deadline_ms,
3241            idempotency_key: spec.idempotency_key.clone(),
3242            required_capabilities: SpawnRequiredCapabilities::default(),
3243            provenance: SpawnResolutionProvenance {
3244                provider: SpawnFieldProvenance::Profile,
3245                mode: SpawnFieldProvenance::Profile,
3246                terminal_size: SpawnFieldProvenance::Profile,
3247                prompt: SpawnFieldProvenance::Profile,
3248                bundle_id: SpawnFieldProvenance::Profile,
3249                context_id: SpawnFieldProvenance::Profile,
3250                environment_profile_id: SpawnFieldProvenance::Profile,
3251            },
3252            harness_mcp_proxy: None,
3253        };
3254        let receipt = ManagedWorktreeSpawnReceipt {
3255            spawn,
3256            lease: managed_lease(
3257                "lease-a",
3258                "managed-a",
3259                ManagedWorktreeLeaseState::InUse,
3260            ),
3261        };
3262        let expected = ExpectedSpawnRequest::ManagedV1(managed.clone());
3263        let response = |receipt| Ok(NodeResponse::ManagedWorktreeSpawnAccepted { receipt });
3264        assert!(validate_spawn_spec_response(
3265            Some(&expected),
3266            &response(receipt.clone()),
3267            incarnation_id,
3268        )
3269        .is_ok());
3270
3271        let mut mismatches = Vec::new();
3272        let mut changed = receipt.clone();
3273        changed.lease.state = ManagedWorktreeLeaseState::Ready;
3274        mismatches.push(changed);
3275        let mut changed = receipt.clone();
3276        changed.lease.cleanup_failure = Some(ManagedWorktreeCleanupFailure::Busy);
3277        mismatches.push(changed);
3278        let mut changed = receipt.clone();
3279        changed.lease.active_session_count = 0;
3280        mismatches.push(changed);
3281        let mut changed = receipt.clone();
3282        changed.lease.profile_id = WorktreeProfileId::new("other").unwrap();
3283        mismatches.push(changed);
3284        let mut changed = receipt.clone();
3285        changed.lease.source_workspace_id = WorkspaceId::new("other").unwrap();
3286        mismatches.push(changed);
3287
3288        for mismatch in mismatches {
3289            assert!(validate_spawn_spec_response(
3290                Some(&expected),
3291                &response(mismatch),
3292                incarnation_id,
3293            )
3294            .is_err());
3295        }
3296
3297        let mut legacy_revision = receipt.clone();
3298        legacy_revision.lease.profile_revision =
3299            WorktreeProfileRevision::new("review.r2").unwrap();
3300        assert!(validate_spawn_spec_response(
3301            Some(&expected),
3302            &response(legacy_revision),
3303            incarnation_id,
3304        )
3305        .is_ok());
3306
3307        let v2_request = ManagedWorktreeSpawnRequestV2 {
3308            spawn_spec: managed.spawn_spec,
3309            worktree_profile_id: managed.worktree_profile_id,
3310            expected_profile_revision: WorktreeProfileRevision::new("review.r1").unwrap(),
3311        };
3312        let routed = NodeRequest::SpawnManagedWorktreeV2 {
3313            request: v2_request.clone(),
3314        };
3315        let captured_expected = expected_spawn_request(&routed).unwrap();
3316        let ExpectedSpawnRequest::ManagedV2(captured_request) = &captured_expected else {
3317            panic!("V2 managed spawn request was not captured separately");
3318        };
3319        assert_eq!(captured_request, &v2_request);
3320        assert!(validate_spawn_spec_response(
3321            Some(&captured_expected),
3322            &response(receipt.clone()),
3323            incarnation_id,
3324        )
3325        .is_ok());
3326
3327        let mut wrong_revision = receipt.clone();
3328        wrong_revision.lease.profile_revision =
3329            WorktreeProfileRevision::new("review.r2").unwrap();
3330        assert!(validate_spawn_spec_response(
3331            Some(&ExpectedSpawnRequest::ManagedV2(v2_request.clone())),
3332            &response(wrong_revision),
3333            incarnation_id,
3334        )
3335        .is_err());
3336        assert!(validate_spawn_spec_response(
3337            Some(&ExpectedSpawnRequest::ManagedV2(v2_request)),
3338            &Ok(NodeResponse::SpawnSpecAccepted {
3339                receipt: receipt.spawn,
3340            }),
3341            incarnation_id,
3342        )
3343        .is_err());
3344    }
3345
3346    #[test]
3347    fn windows_runtime_default_control_endpoint_is_exact_valid_and_distinct_from_a_node() {
3348        assert_eq!(DEFAULT_C2_CONTROL_ENDPOINT, r"\\.\pipe\gate4agent-c2");
3349        validate_control_endpoint(DEFAULT_C2_CONTROL_ENDPOINT).unwrap();
3350        let node = C2NodeConfig::new(
3351            NodeId::new("node-a").unwrap(),
3352            r"\\.\pipe\gate4agent-node",
3353            "safe-token",
3354        )
3355        .unwrap();
3356        let config = C2Config::new(
3357            "127.0.0.1:0".parse().unwrap(),
3358            "safe-token",
3359            vec![node],
3360        )
3361        .unwrap();
3362        assert_eq!(config.control_endpoint, DEFAULT_C2_CONTROL_ENDPOINT);
3363        assert!(!config.nodes[0]
3364            .endpoint
3365            .eq_ignore_ascii_case(&config.control_endpoint));
3366
3367        let conflicting_node = C2NodeConfig::new(
3368            NodeId::new("node-b").unwrap(),
3369            DEFAULT_C2_CONTROL_ENDPOINT,
3370            "safe-token",
3371        )
3372        .unwrap();
3373        assert!(matches!(
3374            C2Config::new(
3375                "127.0.0.1:0".parse().unwrap(),
3376                "safe-token",
3377                vec![conflicting_node],
3378            ),
3379            Err(C2ConfigError::ControlEndpointConflict)
3380        ));
3381    }
3382
3383    #[test]
3384    fn terminal_only_sequence_holes_do_not_mark_c2_partial() {
3385        use gate4agent_node_protocol::{NodeEvent, WorkspaceId};
3386        let event = |sequence| NodeEventEnvelope { sequence, event: NodeEvent::WorkspaceRemoved { workspace_id: WorkspaceId::new("work").unwrap() } };
3387        assert_eq!(validate_events(4, 7, 1, &[event(6)]), Vec::<GapKind>::new());
3388        assert_eq!(validate_events(4, 7, 1, &[]), Vec::<GapKind>::new());
3389        assert_eq!(validate_events(4, 7, 1, &[event(7), event(6)]), vec![GapKind::NonContiguousEvents]);
3390        assert_eq!(validate_events(7, 6, 1, &[]), vec![GapKind::CursorRegression]);
3391        assert_eq!(validate_resync(4, 6, 4, 1, &[]), vec![GapKind::CursorRegression]);
3392    }
3393
3394    #[test]
3395    fn real_durable_eviction_marks_c2_history_evicted() {
3396        use gate4agent_node_protocol::{NodeEvent, WorkspaceId};
3397        let event = |sequence| NodeEventEnvelope {
3398            sequence,
3399            event: NodeEvent::WorkspaceRemoved {
3400                workspace_id: WorkspaceId::new("work").unwrap(),
3401            },
3402        };
3403        assert_eq!(
3404            validate_events(4, 7, 6, &[event(6)]),
3405            vec![GapKind::HistoryEvicted],
3406        );
3407        assert_eq!(validate_events(5, 7, 6, &[event(6)]), Vec::<GapKind>::new());
3408    }
3409
3410    #[test]
3411    fn live_event_gap_preserves_cursor_and_resync_rules() {
3412        use gate4agent_node_protocol::{NodeEvent, WorkspaceId};
3413        let event = |sequence, event| NodeEventEnvelope { sequence, event };
3414        let removed = || NodeEvent::WorkspaceRemoved {
3415            workspace_id: WorkspaceId::new("work").unwrap(),
3416        };
3417
3418        assert_eq!(live_event_gap(4, &event(5, removed())), None);
3419        assert_eq!(
3420            live_event_gap(4, &event(4, removed())),
3421            Some(GapKind::CursorRegression),
3422        );
3423        assert_eq!(
3424            live_event_gap(4, &event(7, removed())),
3425            Some(GapKind::NonContiguousEvents),
3426        );
3427        assert_eq!(
3428            live_event_gap(4, &event(5, NodeEvent::ResyncRequired {
3429                oldest_available_sequence: 3,
3430            })),
3431            Some(GapKind::HistoryEvicted),
3432        );
3433    }
3434
3435    #[test]
3436    fn harness_mcp_transient_live_and_pending_projection_preserves_cursor_without_replay() {
3437        let node_id = NodeId::new("node-a").unwrap();
3438        let cursor = NodeCursor {
3439            incarnation_id: NodeIncarnationId::from_bytes([7; 16]),
3440            sequence: 41,
3441        };
3442        let envelope = NodeEventEnvelope {
3443            sequence: 0,
3444            event: NodeEvent::HarnessMcpReadCall {
3445                reservation_id: HarnessMcpReservationId::new(
3446                    format!("hmcpres_{}", "a".repeat(24)),
3447                ).unwrap(),
3448                activation_digest: HarnessMcpActivationDigest::new(
3449                    format!("sha256:{}", "b".repeat(64)),
3450                ).unwrap(),
3451                record_id: SessionRecordId::new("session-001").unwrap(),
3452                session: SessionAddress {
3453                    workspace_id: WorkspaceId::new("repo").unwrap(),
3454                    session: SessionKey {
3455                        instance_id: AgentInstanceId(8),
3456                        generation: SessionGeneration(1),
3457                    },
3458                },
3459                call_id: HarnessMcpCallId::new(
3460                    format!("hmcpcall_{}", "c".repeat(24)),
3461                ).unwrap(),
3462                request: HarnessMcpOpaquePayloadV1 {
3463                    content_type: HarnessMcpContentTypeV1::HarnessReadRequestJsonV1,
3464                    body: br#"{"kind":"context-get"}"#.to_vec(),
3465                },
3466                deadline_unix_ms: u64::MAX,
3467            },
3468        };
3469
3470        let live = routed_transient_node_event(&node_id, cursor, &envelope).unwrap();
3471        let pending = routed_transient_node_event(&node_id, cursor, &envelope).unwrap();
3472        assert_eq!(live, pending);
3473        assert_eq!(live.cursor, cursor);
3474        assert!(matches!(live.event, C2NodeEvent::HarnessMcpReadCall { .. }));
3475        assert!(routed_recovered_node_event(
3476            &node_id,
3477            cursor.incarnation_id,
3478            &envelope,
3479        ).is_none());
3480        assert_eq!(cursor.sequence, 41);
3481    }
3482
3483    fn agent_stream_test_address() -> SessionAddress {
3484        SessionAddress {
3485            workspace_id: WorkspaceId::new("repo").unwrap(),
3486            session: SessionKey {
3487                instance_id: AgentInstanceId(8),
3488                generation: SessionGeneration(1),
3489            },
3490        }
3491    }
3492
3493    fn agent_stream_test_envelope(sequence: u64, source_sequence: u64) -> NodeEventEnvelope {
3494        NodeEventEnvelope {
3495            sequence,
3496            event: NodeEvent::AgentStream {
3497                address: agent_stream_test_address(),
3498                chunk: AgentStreamChunkV1 {
3499                    source_sequence,
3500                    kind: AgentStreamChunkKindV1::Text {
3501                        text: format!("chunk-{source_sequence}"),
3502                        is_delta: true,
3503                    },
3504                },
3505            },
3506        }
3507    }
3508
3509    #[test]
3510    fn agent_stream_chunk_publishes_unconditionally_and_only_advances_cursor_forward() {
3511        let node_id = NodeId::new("node-a").unwrap();
3512        let mut cursor = NodeCursor {
3513            incarnation_id: NodeIncarnationId::from_bytes([7; 16]),
3514            sequence: 674,
3515        };
3516
3517        // A non-`AgentStream` envelope is not this function's to route.
3518        let control = NodeEventEnvelope {
3519            sequence: 675,
3520            event: NodeEvent::WorkspaceRemoved {
3521                workspace_id: WorkspaceId::new("work").unwrap(),
3522            },
3523        };
3524        assert!(route_agent_stream_event(&node_id, &mut cursor, &control).is_none());
3525        assert_eq!(cursor.sequence, 674);
3526
3527        // A chunk one past the cursor is published, and its OWN cursor
3528        // field on the wire reads the cursor as it stood before this
3529        // chunk (mirroring `HarnessMcpReadCall`'s convention exactly) --
3530        // then the live cursor moves forward to make room for it.
3531        let forward = agent_stream_test_envelope(675, 22);
3532        let routed = route_agent_stream_event(&node_id, &mut cursor, &forward).unwrap();
3533        assert_eq!(routed.cursor.sequence, 674);
3534        assert!(matches!(
3535            &routed.event,
3536            C2NodeEvent::AgentStream { chunk, .. } if chunk.source_sequence == 22
3537        ));
3538        assert_eq!(cursor.sequence, 675);
3539
3540        // A chunk that arrives (or, after a durable resync already carried
3541        // the cursor past it, is only now examined) BEHIND the cursor is
3542        // still published -- chunks promise no resync, so there is nothing
3543        // to recover if this one were dropped instead -- but it must never
3544        // rewind the cursor a genuinely later `Control` envelope already
3545        // advanced past.
3546        cursor.sequence = 711;
3547        let late = agent_stream_test_envelope(677, 23);
3548        let routed = route_agent_stream_event(&node_id, &mut cursor, &late).unwrap();
3549        assert_eq!(routed.cursor.sequence, 711);
3550        assert!(matches!(
3551            &routed.event,
3552            C2NodeEvent::AgentStream { chunk, .. } if chunk.source_sequence == 23
3553        ));
3554        assert_eq!(cursor.sequence, 711, "a late chunk must never rewind the cursor");
3555    }
3556
3557    #[test]
3558    fn agent_stream_burst_survives_a_durable_cursor_that_already_ran_past_it() {
3559        // Reproduces the measured live defect's mechanism at its worst: the
3560        // node's connection loop can drain its durable channel to
3561        // exhaustion across several ticks before this dedicated channel
3562        // gets any budget at all, so a whole burst's `Control` envelopes
3563        // can reach this relay, and carry `cursor` forward, before a
3564        // single one of the SAME burst's interleaved `AgentStream` chunks
3565        // is even looked at. `route_agent_stream_event` must publish every
3566        // one of them regardless.
3567        let node_id = NodeId::new("node-a").unwrap();
3568        let mut cursor = NodeCursor {
3569            incarnation_id: NodeIncarnationId::from_bytes([9; 16]),
3570            sequence: 0,
3571        };
3572        let chunk_count = 200_u64;
3573
3574        for turn in 1..=chunk_count {
3575            // Durable (`Control`) processing advancing the cursor is the
3576            // pre-existing, unchanged path (`live_event_gap` and its
3577            // resync repair, exercised by this module's other tests) --
3578            // simulated here only by its end effect on `cursor`, since
3579            // this test is about the `AgentStream` side's own contract.
3580            cursor.sequence = turn * 2;
3581        }
3582
3583        let mut recovered = Vec::new();
3584        for turn in 1..=chunk_count {
3585            let envelope = agent_stream_test_envelope(turn * 2 - 1, turn);
3586            let routed = route_agent_stream_event(&node_id, &mut cursor, &envelope)
3587                .expect("every chunk in the burst must be published, however far the durable cursor already ran past its sequence");
3588            recovered.push(routed);
3589        }
3590
3591        assert_eq!(recovered.len(), chunk_count as usize);
3592        for (turn, routed) in (1..=chunk_count).zip(recovered.iter()) {
3593            assert!(matches!(
3594                &routed.event,
3595                C2NodeEvent::AgentStream { chunk, .. } if chunk.source_sequence == turn
3596            ));
3597        }
3598        assert_eq!(
3599            cursor.sequence,
3600            chunk_count * 2,
3601            "a whole burst of stale chunks must never rewind the cursor the durable side already advanced",
3602        );
3603    }
3604
3605    #[test]
3606    fn config_rejects_header_injection_duplicate_nodes_and_non_loopback() {
3607        let id = NodeId::new("node-a").unwrap();
3608        assert!(matches!(C2NodeConfig::new(id.clone(), r"\\.\pipe\a", "bad\r\ntoken"), Err(C2ConfigError::InvalidToken)));
3609        let node = C2NodeConfig::new(id, r"\\.\pipe\a", "safe-token").unwrap();
3610        assert!(matches!(C2Config::new("0.0.0.0:0".parse().unwrap(), "safe", vec![node.clone()]), Err(C2ConfigError::NonLoopback(_))));
3611        assert!(matches!(C2Config::new("127.0.0.1:0".parse().unwrap(), "safe", vec![node.clone(), node]), Err(C2ConfigError::DuplicateNode)));
3612        let first = C2NodeConfig::new(NodeId::new("node-a").unwrap(), r"\\.\pipe\same", "safe").unwrap();
3613        let second = C2NodeConfig::new(NodeId::new("node-b").unwrap(), r"\\.\pipe\same", "safe").unwrap();
3614        assert!(matches!(C2Config::new("127.0.0.1:0".parse().unwrap(), "safe", vec![first, second]), Err(C2ConfigError::DuplicateEndpoint)));
3615        let oversized = format!(r"\\.\pipe\{}", "x".repeat(MAX_C2_ENDPOINT_BYTES));
3616        assert!(matches!(C2NodeConfig::new(NodeId::new("node-c").unwrap(), oversized, "safe"), Err(C2ConfigError::InvalidEndpoint(_))));
3617    }
3618
3619    #[test]
3620    fn durable_session_mutations_require_controller_and_use_bounded_deadlines() {
3621        let record_id = SessionRecordId::new("session-001").unwrap();
3622        let rename = NodeRequest::RenameSessionRecord {
3623            record_id: record_id.clone(),
3624            display_name: "release shepherd".to_owned(),
3625        };
3626        let resume = NodeRequest::ResumeSessionRecord {
3627            record_id: record_id.clone(),
3628            terminal_size: TerminalSize { rows: 40, columns: 120 },
3629            initial_prompt: None,
3630        };
3631        let forget = NodeRequest::ForgetSessionRecord { record_id };
3632
3633        assert!(!is_read_only_request(&rename));
3634        assert!(!is_read_only_request(&resume));
3635        assert!(!is_read_only_request(&forget));
3636        assert_eq!(node_request_deadline(&rename), Duration::from_secs(5));
3637        assert_eq!(node_request_deadline(&resume), Duration::from_secs(35));
3638        assert!(node_request_deadline(&resume) > MANAGED_RESUME_SETTLE_DEADLINE);
3639        assert_eq!(
3640            node_request_deadline(&resume) - MANAGED_RESUME_SETTLE_DEADLINE,
3641            NODE_REQUEST_IO_HEADROOM,
3642        );
3643        let started = Instant::now();
3644        let relay_deadline = relay_request_deadline(&resume, started).unwrap();
3645        assert_eq!(
3646            request_budget(&resume, Some(relay_deadline), started),
3647            Duration::from_secs(35),
3648        );
3649        assert_eq!(
3650            request_budget(
3651                &resume,
3652                Some(relay_deadline),
3653                started + Duration::from_secs(5),
3654            ),
3655            MANAGED_RESUME_SETTLE_DEADLINE,
3656        );
3657        assert_eq!(
3658            request_budget(
3659                &resume,
3660                Some(relay_deadline),
3661                started + Duration::from_secs(35),
3662            ),
3663            Duration::ZERO,
3664        );
3665        assert_eq!(node_request_deadline(&forget), Duration::from_secs(5));
3666    }
3667
3668    #[test]
3669    fn workspace_file_reads_are_read_only_with_five_second_deadline() {
3670        let request = NodeRequest::ReadWorkspaceFile {
3671            workspace_id: gate4agent_node_protocol::WorkspaceId::new("primary").unwrap(),
3672            path: gate4agent_node_protocol::RepositoryPath::utf8(
3673                "src/lib.rs".to_owned(),
3674            ).unwrap(),
3675        };
3676
3677        assert!(is_read_only_request(&request));
3678        assert_eq!(node_request_deadline(&request), Duration::from_secs(5));
3679        assert!(relay_request_deadline(&request, Instant::now()).is_none());
3680
3681    }
3682
3683    #[test]
3684    fn workspace_entry_create_relay_has_headroom_and_preserves_semantic_timeout() {
3685        let workspace_id = gate4agent_node_protocol::WorkspaceId::new("primary").unwrap();
3686        let file_path = gate4agent_node_protocol::RepositoryPath::utf8(
3687            "src/new.rs".to_owned(),
3688        ).unwrap();
3689        let create_file = NodeRequest::CreateWorkspaceFile {
3690            workspace_id: workspace_id.clone(),
3691            path: file_path.clone(),
3692        };
3693        assert!(!is_read_only_request(&create_file));
3694        assert_eq!(
3695            node_request_deadline(&create_file),
3696            WORKSPACE_ENTRY_CREATE_RELAY_DEADLINE,
3697        );
3698        assert!(
3699            node_request_deadline(&create_file)
3700                > WORKSPACE_ENTRY_CREATE_NODE_SEMANTIC_DEADLINE,
3701        );
3702        assert!(relay_request_deadline(&create_file, Instant::now()).is_none());
3703
3704        let propagated = relay_node_failure(&NodeClientError::Node(
3705            gate4agent_node_protocol::NodeFailure {
3706                code: NodeFailureCode::RepositoryEntryCreateTimedOut,
3707                message: "private node timeout detail".to_owned(),
3708            },
3709        )).unwrap();
3710        assert_eq!(
3711            propagated.code,
3712            NodeFailureCode::RepositoryEntryCreateTimedOut,
3713        );
3714        assert_eq!(propagated.message, "repository entry creation timed out");
3715
3716        let created_file = gate4agent_node_protocol::WorkspaceFileRead {
3717            workspace_id: workspace_id.clone(),
3718            path: file_path,
3719            content: gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
3720                text: String::new(),
3721                byte_len: 0,
3722            },
3723            revision: Some(
3724                gate4agent_node_protocol::WorkspaceFileRevision::new(
3725                    "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
3726                        .to_owned(),
3727                )
3728                .unwrap(),
3729            ),
3730        };
3731        assert!(validate_workspace_content_response(
3732            Some(&create_file),
3733            &Ok(NodeResponse::WorkspaceFileCreated {
3734                file: created_file.clone(),
3735            }),
3736        ).is_ok());
3737        let mut wrong_content = created_file;
3738        wrong_content.content = gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
3739            text: "unexpected".to_owned(),
3740            byte_len: 10,
3741        };
3742        assert!(validate_workspace_content_response(
3743            Some(&create_file),
3744            &Ok(NodeResponse::WorkspaceFileCreated { file: wrong_content }),
3745        ).is_err());
3746
3747        let directory_path = gate4agent_node_protocol::RepositoryPath::utf8(
3748            "src/new".to_owned(),
3749        ).unwrap();
3750        let create_directory = NodeRequest::CreateWorkspaceDirectory {
3751            workspace_id: workspace_id.clone(),
3752            path: directory_path.clone(),
3753        };
3754        assert!(!is_read_only_request(&create_directory));
3755        assert_eq!(
3756            node_request_deadline(&create_directory),
3757            WORKSPACE_ENTRY_CREATE_RELAY_DEADLINE,
3758        );
3759        assert!(validate_workspace_content_response(
3760            Some(&create_directory),
3761            &Ok(NodeResponse::WorkspaceDirectoryCreated {
3762                workspace_id,
3763                entry: gate4agent_node_protocol::WorkspaceEntry {
3764                    relative_path: directory_path,
3765                    kind: gate4agent_node_protocol::WorkspaceEntryKind::Directory,
3766                },
3767            }),
3768        ).is_ok());
3769        assert!(validate_workspace_content_response(
3770            Some(&create_directory),
3771            &Ok(NodeResponse::Accepted),
3772        ).is_err());
3773    }
3774
3775    #[test]
3776    fn native_session_catalog_is_lease_free_read_only() {
3777        let route = gate4agent_node_protocol::NativeSessionCatalogRoute::workspace(
3778            gate4agent_node_protocol::WorkspaceId::new("primary").unwrap(),
3779            gate4agent_types::AgentId::new("codex").unwrap(),
3780        );
3781        let request = NodeRequest::CatalogNativeSessions {
3782            route: route.clone(),
3783            limit: 8,
3784        };
3785        assert!(is_read_only_request(&request));
3786        assert_eq!(
3787            node_request_deadline(&request),
3788            NATIVE_SESSION_REQUEST_DEADLINE,
3789        );
3790        assert!(relay_request_deadline(&request, Instant::now()).is_none());
3791
3792        let page = NodeRequest::PageNativeSessions {
3793            route: route.clone(),
3794            window: gate4agent_types::NativeSessionCatalogWindow::Older,
3795            catalog_revision: 7,
3796            recent_cutoff_unix_ms: 8,
3797            after_selection_id: Some("hist_selection_1".to_owned()),
3798            limit: 8,
3799        };
3800        assert!(is_read_only_request(&page));
3801        assert_eq!(node_request_deadline(&page), NATIVE_SESSION_REQUEST_DEADLINE);
3802        assert!(relay_request_deadline(&page, Instant::now()).is_none());
3803
3804        let preview = NodeRequest::PreviewNativeSession {
3805            selection: gate4agent_node_protocol::NativeSessionSelection {
3806                route,
3807                catalog_revision: 7,
3808                recent_cutoff_unix_ms: 8,
3809                selection_id: "hist_selection_1".to_owned(),
3810            },
3811            message_limit: 12,
3812        };
3813        assert!(is_read_only_request(&preview));
3814        assert_eq!(
3815            node_request_deadline(&preview),
3816            NATIVE_SESSION_REQUEST_DEADLINE,
3817        );
3818        assert!(relay_request_deadline(&preview, Instant::now()).is_none());
3819
3820        let record_preview = NodeRequest::PreviewSessionRecord {
3821            record_id: gate4agent_node_protocol::SessionRecordId::new("record-1").unwrap(),
3822            message_limit: 12,
3823        };
3824        assert!(is_read_only_request(&record_preview));
3825        assert_eq!(
3826            node_request_deadline(&record_preview),
3827            NATIVE_SESSION_REQUEST_DEADLINE,
3828        );
3829        assert!(relay_request_deadline(&record_preview, Instant::now()).is_none());
3830    }
3831
3832    #[test]
3833    fn standalone_workspace_creation_is_controller_mutation_with_worktree_deadline() {
3834        let request = NodeRequest::CreateStandaloneWorkspace {
3835            workspace_id: gate4agent_node_protocol::WorkspaceId::new("standalone").unwrap(),
3836            root: gate4agent_node_protocol::OpaqueHostPath::utf8(
3837                r"C:\standalone".to_owned(),
3838            ).unwrap(),
3839            initial_branch: Some("main".to_owned()),
3840        };
3841
3842        assert!(!is_read_only_request(&request));
3843        assert_eq!(node_request_deadline(&request), Duration::from_secs(240));
3844        assert!(relay_request_deadline(&request, Instant::now()).is_none());
3845    }
3846
3847    #[test]
3848    fn unsupported_node_capability_is_correlated_without_offline_classification() {
3849        let error = NodeClientError::UnsupportedCapability(
3850            "workspace-file-read-v1-private-detail".to_owned(),
3851        );
3852        let failure = relay_node_failure(&error)
3853            .expect("unsupported node capability must remain an in-band node failure");
3854
3855        assert_eq!(failure.code, NodeFailureCode::UnsupportedCapability);
3856        assert_eq!(failure.message, "required capability unavailable");
3857        assert!(!failure.message.contains("private-detail"));
3858        assert!(relay_node_failure(&NodeClientError::Io(io::Error::new(
3859            io::ErrorKind::BrokenPipe,
3860            "transport closed",
3861        ))).is_none());
3862    }
3863
3864    async fn raw_request(request: Vec<u8>, status: StatusResponse) -> Vec<u8> {
3865        let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
3866        let address = listener.local_addr().unwrap();
3867        let (_status_tx, status_rx) = watch::channel(Arc::new(status));
3868        let server = tokio::spawn(async move {
3869            let (stream, _) = listener.accept().await.unwrap();
3870            serve_http(stream, "api-token", Duration::from_secs(1), &status_rx).await.unwrap();
3871        });
3872        let mut stream = TcpStream::connect(address).await.unwrap();
3873        stream.write_all(&request).await.unwrap();
3874        let mut response = Vec::new();
3875        stream.read_to_end(&mut response).await.unwrap();
3876        server.await.unwrap();
3877        response
3878    }
3879
3880    fn empty_status(ready: bool) -> StatusResponse {
3881        StatusResponse { api_version: C2_API_VERSION, ready, observed_at_unix_ms: 0, nodes: BTreeMap::new() }
3882    }
3883
3884    #[tokio::test]
3885    async fn http_api_enforces_initializing_auth_method_path_and_header_bounds() {
3886        let initializing = raw_request(b"GET /ready HTTP/1.1\r\nHost: localhost\r\n\r\n".to_vec(), empty_status(false)).await;
3887        assert!(initializing.starts_with(b"HTTP/1.1 503 Service Unavailable\r\n"));
3888        assert!(String::from_utf8_lossy(&initializing).contains("\"ready\":false"));
3889
3890        let unauthorized = raw_request(b"GET /status HTTP/1.1\r\nHost: localhost\r\n\r\n".to_vec(), empty_status(true)).await;
3891        assert!(unauthorized.starts_with(b"HTTP/1.1 401 Unauthorized\r\n"));
3892        let method = raw_request(b"POST /health HTTP/1.1\r\nHost: localhost\r\n\r\n".to_vec(), empty_status(true)).await;
3893        assert!(method.starts_with(b"HTTP/1.1 405 Method Not Allowed\r\n"));
3894        let missing = raw_request(b"GET /missing HTTP/1.1\r\nHost: localhost\r\n\r\n".to_vec(), empty_status(true)).await;
3895        assert!(missing.starts_with(b"HTTP/1.1 404 Not Found\r\n"));
3896
3897        let mut oversized = b"GET /health HTTP/1.1\r\nX-Fill: ".to_vec();
3898        oversized.extend(std::iter::repeat(b'x').take(HEADER_LIMIT_BYTES));
3899        oversized.extend_from_slice(b"\r\n\r\n");
3900        let rejected = raw_request(oversized, empty_status(true)).await;
3901        assert!(rejected.starts_with(b"HTTP/1.1 413 Payload Too Large\r\n"));
3902    }
3903
3904    #[tokio::test]
3905    async fn inventory_state_transitions_offline_to_stale_parked_and_recovers() {
3906        let node_id = NodeId::new("node-a").unwrap();
3907        let mut nodes = BTreeMap::new();
3908        nodes.insert(node_id.clone(), ObservedNode {
3909            endpoint: r"\\.\pipe\a".to_owned(), transport_label: "windows-named-pipe".to_owned(),
3910            transport: NodeTransportState::Offline, freshness: NodeFreshness::Unavailable,
3911            cursor: None, inventory: None, last_attempt_unix_ms: None, last_success_unix_ms: None,
3912            consecutive_failures: 0, last_error: None, gaps: Vec::new(), gaps_truncated: 0,
3913        });
3914        let initial = Arc::new(StatusResponse { api_version: C2_API_VERSION, ready: false, observed_at_unix_ms: unix_ms(), nodes });
3915        let (status_tx, mut status_rx) = watch::channel(initial);
3916        let (ingress_tx, ingress_rx) = mpsc::channel(4);
3917        let (shutdown_tx, shutdown_rx) = watch::channel(false);
3918        let owner = tokio::spawn(inventory_owner(1, Duration::from_millis(10), ingress_rx, status_tx, shutdown_rx));
3919        let snapshot = NodeSnapshot {
3920            node_id: node_id.clone(),
3921            enabled_providers: Vec::new(),
3922            provider_runtime_statuses: crate::protocol::ProviderRuntimeStatuses::default(),
3923            workspaces: Vec::new(),
3924            session_records: Vec::new(),
3925            managed_worktrees: Vec::new(),
3926            launch_inventory: None,
3927            agent_progress: Vec::new(),
3928        };
3929        let cursor = NodeCursor { incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([1; 16]), sequence: 0 };
3930        let old_manifest = ProviderContractManifest {
3931            provider_contracts: vec![crate::protocol::ProviderContractSupport {
3932                provider: agent("codex"),
3933                revision: crate::protocol::ProviderContractRevision::new("old-contract").unwrap(),
3934            }],
3935            provider_adapter_contracts: vec![crate::protocol::ProviderAdapterContractSupport {
3936                provider: agent("codex"),
3937                family: crate::protocol::AdapterFamily::PtySemantic,
3938                adapter_id: crate::protocol::AdapterId::new("codex").unwrap(),
3939                revision: crate::protocol::AdapterContractRevision::new("old-adapter").unwrap(),
3940            }],
3941        };
3942        ingress_tx.send(Attempt { node_id: node_id.clone(), at_unix_ms: unix_ms(), result: AttemptResult::Connected {
3943            cursor,
3944            snapshot: snapshot.clone(),
3945            gaps: Vec::new(),
3946            provider_contract_manifest: old_manifest,
3947        } }).await.unwrap();
3948        status_rx.changed().await.unwrap();
3949        assert_eq!(status_rx.borrow().nodes[&node_id].freshness, NodeFreshness::Fresh);
3950        assert_eq!(
3951            status_rx.borrow().nodes[&node_id].inventory.as_ref().unwrap()
3952                .provider_contracts[0].revision.as_str(),
3953            "old-contract",
3954        );
3955
3956        let failure = || AttemptResult::Failure { error: SanitizedError { category: C2ErrorCategory::Transport, message: "node transport unavailable".to_owned() }, hard: false };
3957        ingress_tx.send(Attempt { node_id: node_id.clone(), at_unix_ms: unix_ms(), result: failure() }).await.unwrap();
3958        status_rx.changed().await.unwrap();
3959        assert_eq!(status_rx.borrow().nodes[&node_id].transport, NodeTransportState::Offline);
3960        timeout(Duration::from_secs(1), async {
3961            loop {
3962                status_rx.changed().await.unwrap();
3963                if status_rx.borrow().nodes[&node_id].freshness == NodeFreshness::Stale { break; }
3964            }
3965        }).await.unwrap();
3966        for _ in 0..4 {
3967            ingress_tx.send(Attempt { node_id: node_id.clone(), at_unix_ms: unix_ms(), result: failure() }).await.unwrap();
3968            status_rx.changed().await.unwrap();
3969        }
3970        assert_eq!(status_rx.borrow().nodes[&node_id].transport, NodeTransportState::Parked);
3971        let replacement_manifest = ProviderContractManifest {
3972            provider_contracts: vec![crate::protocol::ProviderContractSupport {
3973                provider: agent("claude"),
3974                revision: crate::protocol::ProviderContractRevision::new("new-contract").unwrap(),
3975            }],
3976            provider_adapter_contracts: Vec::new(),
3977        };
3978        let replacement_cursor = NodeCursor {
3979            incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([2; 16]),
3980            sequence: 0,
3981        };
3982        ingress_tx.send(Attempt { node_id: node_id.clone(), at_unix_ms: unix_ms(), result: AttemptResult::Connected {
3983            cursor: replacement_cursor,
3984            snapshot,
3985            gaps: Vec::new(),
3986            provider_contract_manifest: replacement_manifest,
3987        } }).await.unwrap();
3988        status_rx.changed().await.unwrap();
3989        assert_eq!(status_rx.borrow().nodes[&node_id].transport, NodeTransportState::Online);
3990        assert_eq!(status_rx.borrow().nodes[&node_id].freshness, NodeFreshness::Fresh);
3991        let recovered_inventory = status_rx.borrow().nodes[&node_id].inventory.as_ref().unwrap().clone();
3992        assert_eq!(recovered_inventory.provider_contracts.len(), 1);
3993        assert_eq!(recovered_inventory.provider_contracts[0].provider, agent("claude"));
3994        assert_eq!(recovered_inventory.provider_contracts[0].revision.as_str(), "new-contract");
3995        assert!(recovered_inventory.provider_adapter_contracts.is_empty());
3996        ingress_tx.send(Attempt {
3997            node_id: node_id.clone(),
3998            at_unix_ms: unix_ms(),
3999            result: AttemptResult::Connected {
4000                cursor: NodeCursor {
4001                    incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([3; 16]),
4002                    sequence: 0,
4003                },
4004                snapshot: NodeSnapshot {
4005                    node_id: node_id.clone(),
4006                    enabled_providers: Vec::new(),
4007                    provider_runtime_statuses: crate::protocol::ProviderRuntimeStatuses::default(),
4008                    workspaces: Vec::new(),
4009                    session_records: Vec::new(),
4010                    managed_worktrees: Vec::new(),
4011                    launch_inventory: None,
4012                    agent_progress: Vec::new(),
4013                },
4014                gaps: Vec::new(),
4015                provider_contract_manifest: ProviderContractManifest::default(),
4016            },
4017        }).await.unwrap();
4018        status_rx.changed().await.unwrap();
4019        let unpublished = status_rx.borrow().nodes[&node_id].inventory.as_ref().unwrap().clone();
4020        assert!(unpublished.provider_contracts.is_empty());
4021        assert!(unpublished.provider_adapter_contracts.is_empty());
4022        shutdown_tx.send(true).unwrap();
4023        owner.await.unwrap().unwrap();
4024    }
4025
4026    fn runtime_statuses(
4027        provider: gate4agent_node_protocol::AgentId,
4028        version: &str,
4029    ) -> crate::protocol::ProviderRuntimeStatuses {
4030        crate::protocol::ProviderRuntimeStatuses::new([
4031            crate::protocol::ProviderRuntimeStatus::raw_passthrough(
4032                provider,
4033                Some(crate::protocol::ProviderRuntimeVersion::new(version).unwrap()),
4034            ),
4035        ])
4036        .unwrap()
4037    }
4038
4039    async fn runtime_inventory_owner() -> (
4040        NodeId,
4041        mpsc::Sender<Attempt>,
4042        watch::Receiver<Arc<StatusResponse>>,
4043        watch::Sender<bool>,
4044        tokio::task::JoinHandle<io::Result<()>>,
4045    ) {
4046        let node_id = NodeId::new("runtime-node").unwrap();
4047        let nodes = BTreeMap::from([(
4048            node_id.clone(),
4049            ObservedNode {
4050                endpoint: r"\\.\pipe\runtime-node".to_owned(),
4051                transport_label: "windows-named-pipe".to_owned(),
4052                transport: NodeTransportState::Offline,
4053                freshness: NodeFreshness::Unavailable,
4054                cursor: None,
4055                inventory: None,
4056                last_attempt_unix_ms: None,
4057                last_success_unix_ms: None,
4058                consecutive_failures: 0,
4059                last_error: None,
4060                gaps: Vec::new(),
4061                gaps_truncated: 0,
4062            },
4063        )]);
4064        let initial = Arc::new(StatusResponse {
4065            api_version: C2_API_VERSION,
4066            ready: false,
4067            observed_at_unix_ms: unix_ms(),
4068            nodes,
4069        });
4070        let (status_tx, status_rx) = watch::channel(initial);
4071        let (ingress_tx, ingress_rx) = mpsc::channel(4);
4072        let (shutdown_tx, shutdown_rx) = watch::channel(false);
4073        let owner = tokio::spawn(inventory_owner(
4074            1,
4075            Duration::from_secs(1),
4076            ingress_rx,
4077            status_tx,
4078            shutdown_rx,
4079        ));
4080        (node_id, ingress_tx, status_rx, shutdown_tx, owner)
4081    }
4082
4083    #[tokio::test]
4084    async fn incarnation_change_replaces_runtime_status() {
4085        let (node_id, ingress, mut status, shutdown, owner) = runtime_inventory_owner().await;
4086        for (incarnation, provider, version) in [
4087            (1, agent("claude"), "1.0.0"),
4088            (2, agent("codex"), "2.0.0"),
4089        ] {
4090            ingress
4091                .send(Attempt {
4092                    node_id: node_id.clone(),
4093                    at_unix_ms: unix_ms(),
4094                    result: AttemptResult::Connected {
4095                        cursor: NodeCursor {
4096                            incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([
4097                                incarnation;
4098                                16
4099                            ]),
4100                            sequence: 0,
4101                        },
4102                        snapshot: NodeSnapshot {
4103                            node_id: node_id.clone(),
4104                            enabled_providers: vec![provider.clone()],
4105                            provider_runtime_statuses: runtime_statuses(provider, version),
4106                            workspaces: Vec::new(),
4107                            session_records: Vec::new(),
4108                            managed_worktrees: Vec::new(),
4109                            launch_inventory: None,
4110                            agent_progress: Vec::new(),
4111                        },
4112                        gaps: Vec::new(),
4113                        provider_contract_manifest: ProviderContractManifest::default(),
4114                    },
4115                })
4116                .await
4117                .unwrap();
4118            status.changed().await.unwrap();
4119        }
4120        let current_status = status.borrow();
4121        let statuses = &current_status.nodes[&node_id]
4122            .inventory
4123            .as_ref()
4124            .unwrap()
4125            .provider_runtime_statuses;
4126        assert_eq!(statuses.as_slice().len(), 1);
4127        assert_eq!(
4128            statuses.as_slice()[0].provider(),
4129            &agent("codex"),
4130        );
4131        assert_eq!(statuses.as_slice()[0].version().unwrap().as_str(), "2.0.0");
4132        shutdown.send(true).unwrap();
4133        owner.await.unwrap().unwrap();
4134    }
4135
4136    #[tokio::test]
4137    async fn incarnation_change_without_snapshot_clears_dynamic_inventory() {
4138        let (node_id, ingress, mut status, shutdown, owner) = runtime_inventory_owner().await;
4139        ingress
4140            .send(Attempt {
4141                node_id: node_id.clone(),
4142                at_unix_ms: unix_ms(),
4143                result: AttemptResult::Connected {
4144                    cursor: NodeCursor {
4145                        incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([
4146                            3; 16
4147                        ]),
4148                        sequence: 0,
4149                    },
4150                    snapshot: NodeSnapshot {
4151                        node_id: node_id.clone(),
4152                        enabled_providers: vec![agent("claude")],
4153                        provider_runtime_statuses: runtime_statuses(
4154                            agent("claude"),
4155                            "3.0.0",
4156                        ),
4157                        workspaces: Vec::new(),
4158                        session_records: Vec::new(),
4159                        managed_worktrees: Vec::new(),
4160                        launch_inventory: None,
4161                        agent_progress: Vec::new(),
4162                    },
4163                    gaps: Vec::new(),
4164                    provider_contract_manifest: ProviderContractManifest::default(),
4165                },
4166            })
4167            .await
4168            .unwrap();
4169        status.changed().await.unwrap();
4170        ingress
4171            .send(Attempt {
4172                node_id: node_id.clone(),
4173                at_unix_ms: unix_ms(),
4174                result: AttemptResult::Cursor {
4175                    cursor: NodeCursor {
4176                        incarnation_id: gate4agent_node_protocol::NodeIncarnationId::from_bytes([
4177                            4; 16
4178                        ]),
4179                        sequence: 0,
4180                    },
4181                    gaps: vec![GapKind::IncarnationChanged],
4182                    managed_worktree_events: Vec::new(),
4183                },
4184            })
4185            .await
4186            .unwrap();
4187        status.changed().await.unwrap();
4188        assert!(status.borrow().nodes[&node_id]
4189            .inventory
4190            .as_ref()
4191            .unwrap()
4192            .provider_runtime_statuses
4193            .is_empty());
4194        shutdown.send(true).unwrap();
4195        owner.await.unwrap().unwrap();
4196    }
4197
4198    #[test]
4199    fn managed_worktree_inventory_events_are_exact_bounded_and_incarnation_fenced() {
4200        let mut inventory = SlimNodeInventory::from_snapshot(&NodeSnapshot {
4201            node_id: NodeId::new("node-a").unwrap(),
4202            enabled_providers: Vec::new(),
4203            provider_runtime_statuses: crate::protocol::ProviderRuntimeStatuses::default(),
4204            workspaces: Vec::new(),
4205            session_records: Vec::new(),
4206            managed_worktrees: vec![managed_lease(
4207                "lease-a",
4208                "managed-a",
4209                ManagedWorktreeLeaseState::Ready,
4210            )],
4211            launch_inventory: None,
4212            agent_progress: Vec::new(),
4213        });
4214        apply_managed_worktree_cursor(
4215            Some(&mut inventory),
4216            false,
4217            &[
4218                NodeEvent::ManagedWorktreeUpserted {
4219                    lease: managed_lease(
4220                        "lease-b",
4221                        "managed-b",
4222                        ManagedWorktreeLeaseState::Ready,
4223                    ),
4224                },
4225                NodeEvent::ManagedWorktreeUpserted {
4226                    lease: managed_lease(
4227                        "lease-a",
4228                        "managed-a",
4229                        ManagedWorktreeLeaseState::InUse,
4230                    ),
4231                },
4232            ],
4233        );
4234        assert_eq!(inventory.managed_worktree_count, 2);
4235        assert_eq!(inventory.managed_worktrees[0].lease_id.as_str(), "lease-a");
4236        assert_eq!(
4237            inventory.managed_worktrees[0].state,
4238            ManagedWorktreeLeaseState::InUse,
4239        );
4240
4241        apply_managed_worktree_cursor(
4242            Some(&mut inventory),
4243            false,
4244            &[NodeEvent::ManagedWorktreeRemoved {
4245                lease_id: ManagedWorktreeLeaseId::new("lease-a").unwrap(),
4246            }],
4247        );
4248        assert_eq!(inventory.managed_worktree_count, 1);
4249        assert_eq!(inventory.managed_worktrees[0].lease_id.as_str(), "lease-b");
4250
4251        apply_managed_worktree_cursor(
4252            Some(&mut inventory),
4253            true,
4254            &[NodeEvent::ManagedWorktreeUpserted {
4255                lease: managed_lease(
4256                    "lease-c",
4257                    "managed-c",
4258                    ManagedWorktreeLeaseState::Ready,
4259                ),
4260            }],
4261        );
4262        assert!(inventory.managed_worktrees.is_empty());
4263        assert_eq!(inventory.managed_worktree_count, 0);
4264        assert!(!inventory.managed_worktrees_truncated);
4265    }
4266}
4267
4268#[cfg(test)]
4269mod relay_deadline_tests {
4270    use super::*;
4271
4272    /// The defect this pins, stated as an invariant rather than a number:
4273    /// the relay's bound on a request must never sit below the budget the
4274    /// node itself is allowed to spend answering it. It sat at five
4275    /// seconds against the node's eight-to-eleven, so every inspection of
4276    /// a workspace large enough to use its own allowance timed out at the
4277    /// relay -- and a relay timeout drops the NODE, not the request, which
4278    /// took every read behind it down in a loop.
4279    #[test]
4280    fn workspace_inspection_relay_bound_clears_the_nodes_own_maximum() {
4281        // The node's own ceiling (`WORKSPACE_INSPECTION_TIME_BUDGET_MS_MAX`
4282        // in the node crate) restated here rather than imported: these are
4283        // separate crates by design, and the point of the test is that the
4284        // two numbers must be compared by a human when either moves.
4285        const NODE_INSPECTION_MAX: Duration = Duration::from_millis(11_000);
4286        assert!(
4287            WORKSPACE_INSPECTION_RELAY_DEADLINE > NODE_INSPECTION_MAX,
4288            "relay bound {:?} must exceed the node's own inspection maximum {:?}",
4289            WORKSPACE_INSPECTION_RELAY_DEADLINE,
4290            NODE_INSPECTION_MAX,
4291        );
4292    }
4293}