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 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 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 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 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 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 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 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
616async 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 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
688async 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 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
962async 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 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
1032const 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 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
1866fn 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 #[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 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 #[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 #[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 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 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 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 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 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 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 = ¤t_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 #[test]
4280 fn workspace_inspection_relay_bound_clears_the_nodes_own_maximum() {
4281 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}