1use gate4agent_node_protocol::{
2 provider_id_is_legacy,
3 production_node_client_compatibility_offer,
4 read_json_frame_limited_body_timeout, validate_provider_contract_manifest,
5 write_json_frame_limited, CapabilityId,
6 ClientAuthentication, ClientCompatibilityOffer, ClientFrame, ClientHello, ClientRole,
7 FrameError, HarnessMcpLocalReplyV1, HarnessMcpLocalRequestV1, HarnessMcpLocalToken,
8 HarnessMcpOpaquePayloadV1, HarnessMcpRejectReasonV1,
9 NegotiatedNodeCompatibility, NodeEvent, NodeEventEnvelope, NodeFailure,
10 NodeFailureCode, NodeHello, NodeId, NodeRequest, NodeResponse, NodeSnapshot, RequestEnvelope,
11 WorkspaceSnapshot,
12 ServerChallenge, ServerFrame,
13 MAX_HARNESS_MCP_AGGREGATE_REPLY_BYTES, MAX_HARNESS_MCP_LOCAL_REQUEST_BYTES,
14 MAX_NODE_CLIENT_FRAME_BYTES, MAX_NODE_FRAME_BYTES, MAX_NODE_HELLO_FRAME_BYTES,
15 NODE_AUTH_NONCE_BYTES,
16 CAPABILITY_HOST_DIRECTORY_BROWSE_V1,
17 NODE_AGENT_PROGRESS_SNAPSHOT_CAPABILITY,
18 NODE_DELIVERY_BUNDLE_V2_STAGE_COMMIT_CAPABILITY,
19 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY,
20 NODE_HISTORY_CONTEXT_PACK_CAPABILITY,
21 NODE_SESSION_RECORD_CONTEXT_EXPORT_CAPABILITY,
22 NODE_HARNESS_MCP_READ_PROXY_CAPABILITY,
23 NODE_NATIVE_SESSION_CATALOG_CAPABILITY, NODE_NATIVE_SESSION_CATALOG_PAGING_CAPABILITY,
24 NODE_NATIVE_SESSION_INDEX_CAPABILITY, NODE_NATIVE_SESSION_PREVIEW_CAPABILITY,
25 NODE_STANDALONE_WORKSPACE_LIFECYCLE_CAPABILITY,
26 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY,
27 NODE_COMPATIBILITY_METADATA_CAPABILITY, NODE_OPAQUE_UNIX_PATH_CAPABILITY,
28 NODE_PROVIDER_ID_OPEN_CAPABILITY,
29 NODE_PROVIDER_SESSION_REFERENCE_INDEX_CAPABILITY,
30 NODE_SESSION_TASK_CORRELATION_CAPABILITY,
31 NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY, NODE_REPOSITORY_PATH_CAPABILITY,
32 NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY,
33 NODE_SPAWN_PROFILE_REVISION_CAPABILITY,
34 NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY,
35 NODE_TERMINAL_FRAME_EVENTS_CAPABILITY,
36 NODE_GIT_READ_CAPABILITY, NODE_WORKSPACE_FILE_READ_CAPABILITY,
37 NODE_WORKSPACE_FILE_WRITE_CAPABILITY,
38 NODE_WORKSPACE_ENTRY_CREATE_CAPABILITY,
39 NODE_WORKTREE_SELECTION_CAPABILITY,
40 BUILD_STAMP,
41};
42use crate::{
43 connect_local_stream, negotiated_auth_proof, proofs_match, random_nonce, AuthDirection,
44};
45#[cfg(test)]
46use gate4agent_node_protocol::{
47 NodeIncarnationId, ProtocolRange, StateSchemaSupport, NODE_INCARNATION_ID_BYTES,
48 NODE_STATE_SCHEMA_V1, NODE_STATE_SCHEMA_V10,
49};
50#[cfg(test)]
51use crate::auth_proof;
52use std::collections::VecDeque;
53use std::io;
54use std::net::SocketAddr;
55use std::path::{Path, PathBuf};
56use std::pin::Pin;
57use std::task::{Context, Poll};
58use std::time::{Duration, SystemTime, UNIX_EPOCH};
59use thiserror::Error;
60#[cfg(feature = "fixture")]
61use tokio::io::AsyncWriteExt;
62use tokio::io::{AsyncRead, AsyncWrite, ReadBuf, WriteHalf};
63use tokio::net::TcpStream;
64use tokio::sync::mpsc;
65use tokio::task::AbortHandle;
66use tokio::time::timeout;
67
68const AUTH_FRAME_TIMEOUT_MS: u64 = 5_000;
69const FRAME_BODY_TIMEOUT_MS: u64 = 5_000;
70const SERVER_FRAME_QUEUE_CAPACITY: usize = 1;
71const PENDING_EVENTS_MAX: usize = 1_024;
72const PENDING_EVENT_WIRE_BYTES_MAX: usize = 16 * 1024 * 1024;
73const LOCAL_SESSION_HARNESS_MCP_DEADLINE: Duration = Duration::from_secs(3);
74
75#[derive(Clone)]
76pub struct LocalSessionHarnessMcpClient {
77 endpoint: PathBuf,
78 token: HarnessMcpLocalToken,
79}
80
81impl LocalSessionHarnessMcpClient {
82 pub fn new(
83 endpoint: impl AsRef<Path>,
84 token: HarnessMcpLocalToken,
85 ) -> Result<Self, LocalSessionHarnessMcpError> {
86 let endpoint = endpoint.as_ref();
87 if endpoint.as_os_str().is_empty() {
88 return Err(LocalSessionHarnessMcpError::Unavailable);
89 }
90 Ok(Self { endpoint: endpoint.to_path_buf(), token })
91 }
92
93 pub fn send(
94 &self,
95 request: HarnessMcpOpaquePayloadV1,
96 ) -> Result<HarnessMcpOpaquePayloadV1, LocalSessionHarnessMcpError> {
97 let envelope = HarnessMcpLocalRequestV1 {
98 version: 1,
99 token: self.token.clone(),
100 request,
101 };
102 envelope.validate().map_err(|_| LocalSessionHarnessMcpError::InvalidRequest)?;
103 let runtime = tokio::runtime::Builder::new_current_thread()
104 .enable_io()
105 .enable_time()
106 .build()
107 .map_err(|_| LocalSessionHarnessMcpError::Unavailable)?;
108 runtime.block_on(async {
109 timeout(LOCAL_SESSION_HARNESS_MCP_DEADLINE, async {
110 let mut stream = connect_local_stream(&self.endpoint)
111 .await
112 .map_err(map_harness_mcp_io_error)?;
113 write_json_frame_limited(
114 &mut stream,
115 &envelope,
116 MAX_HARNESS_MCP_LOCAL_REQUEST_BYTES,
117 ).await.map_err(map_harness_mcp_frame_error)?;
118 use tokio::io::AsyncWriteExt as _;
119 stream.shutdown().await.map_err(map_harness_mcp_io_error)?;
120 let reply: HarnessMcpLocalReplyV1 = read_json_frame_limited_body_timeout(
121 &mut stream,
122 MAX_HARNESS_MCP_AGGREGATE_REPLY_BYTES,
123 LOCAL_SESSION_HARNESS_MCP_DEADLINE,
124 ).await.map_err(map_harness_mcp_frame_error)?;
125 reply.validate().map_err(|_| LocalSessionHarnessMcpError::InvalidResponse)?;
126 match reply {
127 HarnessMcpLocalReplyV1::Ok { response } => Ok(response),
128 HarnessMcpLocalReplyV1::Rejected {
129 reason: HarnessMcpRejectReasonV1::Unauthorized,
130 } => Err(LocalSessionHarnessMcpError::Unauthorized),
131 HarnessMcpLocalReplyV1::Rejected { reason } => {
132 Err(LocalSessionHarnessMcpError::Rejected(reason))
133 }
134 }
135 }).await.map_err(|_| LocalSessionHarnessMcpError::Deadline)?
136 })
137 }
138}
139
140fn map_harness_mcp_io_error(error: io::Error) -> LocalSessionHarnessMcpError {
141 match error.kind() {
142 io::ErrorKind::PermissionDenied => LocalSessionHarnessMcpError::Unauthorized,
143 io::ErrorKind::TimedOut => LocalSessionHarnessMcpError::Deadline,
144 _ => LocalSessionHarnessMcpError::Unavailable,
145 }
146}
147
148fn map_harness_mcp_frame_error(error: FrameError) -> LocalSessionHarnessMcpError {
149 match error {
150 FrameError::PrefixTimedOut | FrameError::BodyTimedOut { .. } => {
151 LocalSessionHarnessMcpError::Deadline
152 }
153 _ => LocalSessionHarnessMcpError::InvalidResponse,
154 }
155}
156
157#[derive(Debug, Error)]
158pub enum LocalSessionHarnessMcpError {
159 #[error("local session harness MCP request was unauthorized")]
160 Unauthorized,
161 #[error("local session harness MCP endpoint is unavailable")]
162 Unavailable,
163 #[error("local session harness MCP request is invalid")]
164 InvalidRequest,
165 #[error("local session harness MCP response is invalid")]
166 InvalidResponse,
167 #[error("local session harness MCP deadline exceeded")]
168 Deadline,
169 #[error("local session harness MCP host rejected the request: {0:?}")]
170 Rejected(HarnessMcpRejectReasonV1),
171}
172
173pub trait NodeClientStream: AsyncRead + AsyncWrite + Unpin + Send {}
183
184impl<T> NodeClientStream for T where T: AsyncRead + AsyncWrite + Unpin + Send {}
185
186struct CountingReader<R> {
187 inner: R,
188 bytes_read: usize,
189}
190
191impl<R> CountingReader<R> {
192 fn new(inner: R) -> Self {
193 Self { inner, bytes_read: 0 }
194 }
195
196 fn take_bytes_read(&mut self) -> usize {
197 let bytes_read = self.bytes_read;
198 self.bytes_read = 0;
199 bytes_read
200 }
201}
202
203impl<R: AsyncRead + Unpin> AsyncRead for CountingReader<R> {
204 fn poll_read(
205 mut self: Pin<&mut Self>,
206 context: &mut Context<'_>,
207 buffer: &mut ReadBuf<'_>,
208 ) -> Poll<io::Result<()>> {
209 let before = buffer.filled().len();
210 match Pin::new(&mut self.inner).poll_read(context, buffer) {
211 Poll::Ready(Ok(())) => {
212 self.bytes_read = self
213 .bytes_read
214 .saturating_add(buffer.filled().len().saturating_sub(before));
215 Poll::Ready(Ok(()))
216 }
217 result => result,
218 }
219 }
220}
221
222struct ReceivedServerFrame {
223 frame: ServerFrame,
224 wire_bytes: usize,
225}
226
227struct PendingNodeEvent {
228 envelope: NodeEventEnvelope,
229 wire_bytes: usize,
230}
231
232struct AbortReaderOnDrop(AbortHandle);
233
234impl Drop for AbortReaderOnDrop {
235 fn drop(&mut self) {
236 self.0.abort();
237 }
238}
239
240fn start_server_frame_reader(
241 pipe: Box<dyn NodeClientStream>,
242) -> (
243 WriteHalf<Box<dyn NodeClientStream>>,
244 mpsc::Receiver<Result<ReceivedServerFrame, FrameError>>,
245 AbortReaderOnDrop,
246) {
247 let (reader, writer) = tokio::io::split(pipe);
248 let mut reader = CountingReader::new(reader);
249 let (frame_tx, frame_rx) = mpsc::channel(SERVER_FRAME_QUEUE_CAPACITY);
250 let reader_task = tokio::spawn(async move {
251 loop {
252 let frame = read_json_frame_limited_body_timeout(
253 &mut reader,
254 MAX_NODE_FRAME_BYTES,
255 Duration::from_millis(FRAME_BODY_TIMEOUT_MS),
256 )
257 .await
258 .map(|frame| ReceivedServerFrame {
259 frame,
260 wire_bytes: reader.take_bytes_read(),
261 });
262 if frame.is_err() {
263 let _ = reader.take_bytes_read();
264 }
265 let terminal = frame.is_err();
266 if frame_tx.send(frame).await.is_err() || terminal {
267 break;
268 }
269 }
270 });
271 (
272 writer,
273 frame_rx,
274 AbortReaderOnDrop(reader_task.abort_handle()),
275 )
276}
277
278pub struct LocalNodeClient {
279 writer: WriteHalf<Box<dyn NodeClientStream>>,
280 frame_rx: mpsc::Receiver<Result<ReceivedServerFrame, FrameError>>,
281 _reader_abort: AbortReaderOnDrop,
282 hello: NodeHello,
283 opaque_unix_paths_enabled: bool,
284 repository_paths_enabled: bool,
285 open_provider_ids_enabled: bool,
286 terminal_frame_events_enabled: bool,
287 negotiated_capabilities: Vec<CapabilityId>,
288 next_request_id: u64,
289 pending_events: VecDeque<PendingNodeEvent>,
290 pending_event_wire_bytes: usize,
291}
292
293impl LocalNodeClient {
294 pub async fn connect(
295 endpoint: impl AsRef<Path>,
296 expected_node_id: &NodeId,
297 role: ClientRole,
298 access_token: &str,
299 ) -> Result<Self, NodeClientError> {
300 let pipe = connect_local_stream(endpoint).await?;
301 Self::connect_stream(Box::new(pipe), expected_node_id, role, access_token).await
302 }
303
304 pub async fn connect_loopback(
305 endpoint: SocketAddr,
306 expected_node_id: &NodeId,
307 role: ClientRole,
308 access_token: &str,
309 ) -> Result<Self, NodeClientError> {
310 if !endpoint.ip().is_loopback() || endpoint.port() == 0 {
311 return Err(NodeClientError::Io(io::Error::new(
312 io::ErrorKind::InvalidInput,
313 "node TCP endpoint must be loopback with a nonzero port",
314 )));
315 }
316 let stream = TcpStream::connect(endpoint).await?;
317 Self::connect_stream(Box::new(stream), expected_node_id, role, access_token).await
318 }
319
320 pub async fn adopt(
336 stream: impl NodeClientStream + 'static,
337 expected_node_id: &NodeId,
338 role: ClientRole,
339 access_token: &str,
340 ) -> Result<Self, NodeClientError> {
341 Self::connect_stream(Box::new(stream), expected_node_id, role, access_token).await
342 }
343
344 async fn connect_stream(
345 mut pipe: Box<dyn NodeClientStream>,
346 expected_node_id: &NodeId,
347 role: ClientRole,
348 access_token: &str,
349 ) -> Result<Self, NodeClientError> {
350 let client_nonce = random_nonce().map_err(NodeClientError::Authentication)?;
351 let compatibility_offer = client_compatibility_offer()?;
352 write_json_frame_limited(
353 &mut pipe,
354 &ClientFrame::Hello(ClientHello::negotiating(
355 role,
356 client_nonce,
357 compatibility_offer.clone(),
358 )),
359 MAX_NODE_HELLO_FRAME_BYTES,
360 )
361 .await?;
362 let challenge = timeout(
363 Duration::from_millis(AUTH_FRAME_TIMEOUT_MS),
364 read_json_frame_limited_body_timeout(
365 &mut pipe,
366 MAX_NODE_HELLO_FRAME_BYTES,
367 Duration::from_millis(FRAME_BODY_TIMEOUT_MS),
368 ),
369 )
370 .await
371 .map_err(|_| NodeClientError::AuthenticationTimedOut)??;
372 let ServerFrame::Challenge(challenge) = challenge else {
373 return Err(NodeClientError::Protocol(
374 "server did not return an authentication challenge".to_owned(),
375 ));
376 };
377 if challenge.build_stamp != BUILD_STAMP {
378 return Err(NodeClientError::BuildStampMismatch {
379 local: BUILD_STAMP.to_owned(),
380 remote: challenge.build_stamp.clone(),
381 });
382 }
383 let authentication = prepare_negotiated_authentication(
384 &challenge,
385 &compatibility_offer,
386 role,
387 &client_nonce,
388 access_token,
389 )?;
390 write_json_frame_limited(
391 &mut pipe,
392 &ClientFrame::Authenticate(authentication),
393 MAX_NODE_HELLO_FRAME_BYTES,
394 )
395 .await?;
396 let server_hello = timeout(
397 Duration::from_millis(AUTH_FRAME_TIMEOUT_MS),
398 read_json_frame_limited_body_timeout(
399 &mut pipe,
400 MAX_NODE_FRAME_BYTES,
401 Duration::from_millis(FRAME_BODY_TIMEOUT_MS),
402 ),
403 )
404 .await
405 .map_err(|_| NodeClientError::AuthenticationTimedOut)??;
406 let ServerFrame::Hello(hello) = server_hello else {
407 return Err(NodeClientError::Protocol(
408 "server did not return hello".to_owned(),
409 ));
410 };
411 if hello.build_stamp != BUILD_STAMP {
412 return Err(NodeClientError::BuildStampMismatch {
413 local: BUILD_STAMP.to_owned(),
414 remote: hello.build_stamp.clone(),
415 });
416 }
417 validate_authenticated_hello_compatibility(
418 &compatibility_offer,
419 challenge.compatibility.as_ref().ok_or_else(|| {
420 NodeClientError::Protocol(
421 "node omitted the required authenticated compatibility selection".to_owned(),
422 )
423 })?,
424 hello.compatibility.as_ref(),
425 )?;
426 let opaque_unix_paths_enabled = selected_supports_opaque_unix_paths(
427 hello.compatibility.as_ref(),
428 );
429 let repository_paths_enabled = selected_supports_repository_paths(
430 hello.compatibility.as_ref(),
431 );
432 let negotiated_capabilities = hello.compatibility.as_ref()
433 .map(|compatibility| compatibility.capabilities.clone())
434 .unwrap_or_default();
435 let open_provider_ids_enabled = selected_supports_open_provider_ids(
436 hello.compatibility.as_ref(),
437 );
438 let terminal_frame_events_enabled = selected_supports_terminal_frame_events(
439 hello.compatibility.as_ref(),
440 );
441 ensure_node_hello_path_capability(&hello, opaque_unix_paths_enabled)?;
442 ensure_node_hello_provider_capability(&hello, open_provider_ids_enabled)?;
443 ensure_node_hello_environment_profile_capability(
444 &hello,
445 &negotiated_capabilities,
446 )?;
447 ensure_node_hello_bundle_materialization_capability(
448 &hello,
449 &negotiated_capabilities,
450 )?;
451 ensure_node_hello_history_context_pack_capability(
452 &hello,
453 &negotiated_capabilities,
454 )?;
455 ensure_node_snapshot_agent_progress_capability(
456 &hello.snapshot,
457 &negotiated_capabilities,
458 )?;
459 if &hello.snapshot.node_id != expected_node_id {
460 return Err(NodeClientError::Protocol(format!(
461 "node identity mismatch: expected '{}', received '{}'",
462 expected_node_id,
463 hello.snapshot.node_id,
464 )));
465 }
466 let (writer, frame_rx, reader_abort) = start_server_frame_reader(pipe);
467 Ok(Self {
468 writer,
469 frame_rx,
470 _reader_abort: reader_abort,
471 hello,
472 opaque_unix_paths_enabled,
473 repository_paths_enabled,
474 open_provider_ids_enabled,
475 terminal_frame_events_enabled,
476 negotiated_capabilities,
477 next_request_id: 1,
478 pending_events: VecDeque::new(),
479 pending_event_wire_bytes: 0,
480 })
481 }
482
483 pub fn hello(&self) -> &NodeHello {
484 &self.hello
485 }
486
487 pub async fn send(&mut self, request: NodeRequest) -> Result<u64, NodeClientError> {
488 let request_id = reserve_request_id(
489 &mut self.next_request_id,
490 &request,
491 self.opaque_unix_paths_enabled,
492 self.repository_paths_enabled,
493 self.open_provider_ids_enabled,
494 &self.negotiated_capabilities,
495 )?;
496 write_json_frame_limited(
497 &mut self.writer,
498 &ClientFrame::Request(RequestEnvelope {
499 request_id,
500 request,
501 }),
502 MAX_NODE_CLIENT_FRAME_BYTES,
503 )
504 .await?;
505 Ok(request_id)
506 }
507
508 pub async fn recv(&mut self) -> Result<ServerFrame, NodeClientError> {
509 Ok(self.recv_received().await?.frame)
510 }
511
512 async fn recv_received(&mut self) -> Result<ReceivedServerFrame, NodeClientError> {
513 self.recv_received_for_request(None).await
514 }
515
516 async fn recv_received_for_request(
517 &mut self,
518 expected_request: Option<&NodeRequest>,
519 ) -> Result<ReceivedServerFrame, NodeClientError> {
520 let received = self.frame_rx.recv().await.ok_or_else(|| {
521 NodeClientError::Protocol("node frame reader closed".to_owned())
522 })??;
523 let frame = &received.frame;
524 ensure_server_frame_required_capability_for_request(
525 frame,
526 &self.negotiated_capabilities,
527 expected_request,
528 )?;
529 ensure_server_frame_terminal_capability(
530 frame,
531 self.terminal_frame_events_enabled,
532 )?;
533 ensure_server_frame_path_capability(
534 frame,
535 self.opaque_unix_paths_enabled,
536 self.repository_paths_enabled,
537 )?;
538 ensure_server_frame_provider_capability(frame, self.open_provider_ids_enabled)?;
539 Ok(received)
540 }
541
542 pub async fn request(&mut self, request: NodeRequest) -> Result<NodeResponse, NodeClientError> {
543 let expected_request = request.clone();
544 let request_id = self.send(request).await?;
545 loop {
546 let received = self
547 .recv_received_for_request(Some(&expected_request))
548 .await?;
549 match received.frame {
550 ServerFrame::Reply(reply) if reply.request_id == request_id => {
551 validate_provider_session_index_response(&expected_request, &reply.result)?;
552 validate_native_session_response(&expected_request, &reply.result)?;
553 validate_workspace_content_response(&expected_request, &reply.result)?;
554 validate_session_task_response(&expected_request, &reply.result)?;
555 validate_delivery_response(&expected_request, &reply.result)?;
556 validate_harness_mcp_response(&expected_request, &reply.result)?;
557 return reply.result.map_err(NodeClientError::Node);
558 }
559 ServerFrame::Reply(reply) => {
560 return Err(NodeClientError::Protocol(format!(
561 "unexpected response id {} while waiting for {request_id}",
562 reply.request_id,
563 )));
564 }
565 ServerFrame::Event(event) => queue_pending_event_bounded(
566 &mut self.pending_events,
567 &mut self.pending_event_wire_bytes,
568 event,
569 received.wire_bytes,
570 PENDING_EVENTS_MAX,
571 PENDING_EVENT_WIRE_BYTES_MAX,
572 )?,
573 ServerFrame::Challenge(_) => {
574 return Err(NodeClientError::Protocol(
575 "duplicate server challenge".to_owned(),
576 ));
577 }
578 ServerFrame::Hello(_) => {
579 return Err(NodeClientError::Protocol(
580 "duplicate server hello".to_owned(),
581 ));
582 }
583 }
584 }
585 }
586
587 pub fn take_event(&mut self) -> Option<NodeEventEnvelope> {
588 let pending = self.pending_events.pop_front()?;
589 self.pending_event_wire_bytes = self
590 .pending_event_wire_bytes
591 .checked_sub(pending.wire_bytes)
592 .expect("pending node event wire byte accounting diverged");
593 Some(pending.envelope)
594 }
595
596 #[cfg(feature = "fixture")]
597 pub async fn send_malformed_json_frame_for_fixture(&mut self) -> Result<(), NodeClientError> {
598 self.writer.write_u32_le(1).await?;
599 self.writer.write_all(b"{").await?;
600 self.writer.flush().await?;
601 Ok(())
602 }
603}
604
605fn queue_pending_event_bounded(
606 pending: &mut VecDeque<PendingNodeEvent>,
607 pending_wire_bytes: &mut usize,
608 event: NodeEventEnvelope,
609 wire_bytes: usize,
610 max_events: usize,
611 max_wire_bytes: usize,
612) -> Result<(), NodeClientError> {
613 let replacement = if let NodeEvent::TerminalFrame { address, .. } = &event.event {
614 pending.iter().position(|current| {
615 matches!(
616 ¤t.envelope.event,
617 NodeEvent::TerminalFrame {
618 address: current_address,
619 ..
620 } if current_address == address
621 )
622 })
623 } else {
624 None
625 };
626 let replaced_wire_bytes = replacement
627 .and_then(|index| pending.get(index))
628 .map(|pending| pending.wire_bytes)
629 .unwrap_or(0);
630 let retained_wire_bytes = pending_wire_bytes
631 .checked_sub(replaced_wire_bytes)
632 .expect("pending node event wire byte accounting diverged");
633 let next_wire_bytes = retained_wire_bytes
634 .checked_add(wire_bytes)
635 .ok_or_else(|| {
636 NodeClientError::Protocol(
637 "pending node event wire byte accounting overflowed".to_owned(),
638 )
639 })?;
640 if next_wire_bytes > max_wire_bytes {
641 return Err(NodeClientError::Protocol(
642 "pending node events exceeded the bounded wire byte capacity".to_owned(),
643 ));
644 }
645 if replacement.is_none() && pending.len() >= max_events {
646 return Err(NodeClientError::Protocol(
647 "pending node events exceeded the bounded event capacity".to_owned(),
648 ));
649 }
650 if let Some(index) = replacement {
651 pending.remove(index);
652 }
653 pending.push_back(PendingNodeEvent {
654 envelope: event,
655 wire_bytes,
656 });
657 *pending_wire_bytes = next_wire_bytes;
658 Ok(())
659}
660
661fn selected_supports_opaque_unix_paths(
662 selected: Option<&NegotiatedNodeCompatibility>,
663) -> bool {
664 selected.is_some_and(|compatibility| {
665 compatibility.capabilities.iter().any(|capability| {
666 capability.as_str() == NODE_OPAQUE_UNIX_PATH_CAPABILITY
667 })
668 })
669}
670
671fn selected_supports_repository_paths(
672 selected: Option<&NegotiatedNodeCompatibility>,
673) -> bool {
674 selected.is_some_and(|compatibility| {
675 compatibility.capabilities.iter().any(|capability| {
676 capability.as_str() == NODE_REPOSITORY_PATH_CAPABILITY
677 })
678 })
679}
680
681fn selected_supports_open_provider_ids(
682 selected: Option<&NegotiatedNodeCompatibility>,
683) -> bool {
684 selected.is_some_and(|compatibility| {
685 compatibility.capabilities.iter().any(|capability| {
686 capability.as_str() == NODE_PROVIDER_ID_OPEN_CAPABILITY
687 })
688 })
689}
690
691fn selected_supports_terminal_frame_events(
692 selected: Option<&NegotiatedNodeCompatibility>,
693) -> bool {
694 selected.is_some_and(|compatibility| {
695 compatibility.capabilities.iter().any(|capability| {
696 capability.as_str() == NODE_TERMINAL_FRAME_EVENTS_CAPABILITY
697 })
698 })
699}
700
701fn reserve_request_id(
702 next_request_id: &mut u64,
703 request: &NodeRequest,
704 opaque_unix_paths_enabled: bool,
705 repository_paths_enabled: bool,
706 open_provider_ids_enabled: bool,
707 negotiated_capabilities: &[CapabilityId],
708) -> Result<u64, NodeClientError> {
709 let now_unix_ms = current_unix_ms()?;
710 if !request.harness_mcp_contract_is_valid_at(now_unix_ms) {
711 return Err(NodeClientError::Protocol(
712 "invalid harness MCP proxy request".to_owned(),
713 ));
714 }
715 if !request.history_context_pack_contract_is_valid() {
716 return Err(NodeClientError::Protocol(
717 "invalid history context pack request".to_owned(),
718 ));
719 }
720 if !request.native_session_catalog_contract_is_valid() {
721 return Err(NodeClientError::Protocol(
722 "invalid native session catalog request".to_owned(),
723 ));
724 }
725 if !request.native_session_preview_contract_is_valid() {
726 return Err(NodeClientError::Protocol(
727 "invalid native session preview request".to_owned(),
728 ));
729 }
730 ensure_node_request_required_capability(request, negotiated_capabilities)?;
731 ensure_node_request_path_capability(
732 request,
733 opaque_unix_paths_enabled,
734 repository_paths_enabled,
735 )?;
736 ensure_node_request_provider_capability(request, open_provider_ids_enabled)?;
737 let request_id = *next_request_id;
738 *next_request_id = next_request_id
739 .checked_add(1)
740 .ok_or(NodeClientError::RequestIdExhausted)?;
741 Ok(request_id)
742}
743
744fn current_unix_ms() -> Result<u64, NodeClientError> {
745 SystemTime::now()
746 .duration_since(UNIX_EPOCH)
747 .map_err(|_| NodeClientError::Protocol("system clock precedes Unix epoch".to_owned()))?
748 .as_millis()
749 .try_into()
750 .map_err(|_| NodeClientError::Protocol("system clock exceeds protocol range".to_owned()))
751}
752
753fn validate_provider_session_index_response(
754 expected: &NodeRequest,
755 response: &Result<NodeResponse, NodeFailure>,
756) -> Result<(), NodeClientError> {
757 let matches = match (expected, response) {
758 (
759 NodeRequest::IndexProviderSession {
760 workspace_id,
761 provider,
762 identity,
763 ..
764 },
765 Ok(NodeResponse::ProviderSessionIndexed { record }),
766 ) => {
767 &record.workspace_id == workspace_id
768 && &record.provider == provider
769 && record.provider_session.as_ref() == Some(identity)
770 }
771 (NodeRequest::IndexProviderSession { .. }, Err(_)) => true,
772 (NodeRequest::IndexProviderSession { .. }, Ok(_)) => false,
773 (_, Ok(NodeResponse::ProviderSessionIndexed { .. })) => false,
774 _ => true,
775 };
776 if matches {
777 Ok(())
778 } else {
779 Err(NodeClientError::Protocol(
780 "node provider session index response does not match the request".to_owned(),
781 ))
782 }
783}
784
785fn validate_native_session_response(
786 expected: &NodeRequest,
787 response: &Result<NodeResponse, NodeFailure>,
788) -> Result<(), NodeClientError> {
789 let matches = match (expected, response) {
790 (
791 NodeRequest::CatalogNativeSessions { route, .. },
792 Ok(response @ NodeResponse::NativeSessionsCataloged {
793 route: echoed_route,
794 ..
795 }),
796 ) => echoed_route == route && response.native_session_catalog_contract_is_valid(),
797 (
798 NodeRequest::PageNativeSessions {
799 route,
800 window,
801 catalog_revision,
802 ..
803 },
804 Ok(response @ NodeResponse::NativeSessionsPaged {
805 route: echoed_route,
806 page,
807 }),
808 ) => {
809 echoed_route == route
810 && page.window == *window
811 && page.revision == *catalog_revision
812 && response.native_session_catalog_contract_is_valid()
813 }
814 (
815 NodeRequest::PreviewNativeSession { selection, .. },
816 Ok(response @ NodeResponse::NativeSessionPreviewed {
817 selection: echoed_selection,
818 ..
819 }),
820 ) => {
821 echoed_selection == selection
822 && response.native_session_preview_contract_is_valid()
823 }
824 (
825 NodeRequest::IndexNativeSession { selection, .. },
826 Ok(response @ NodeResponse::NativeSessionIndexed {
827 selection: echoed_selection,
828 record,
829 }),
830 ) => {
831 echoed_selection == selection
832 && selection.route.scope
833 == gate4agent_node_protocol::NativeSessionCatalogScope::Workspace
834 && selection.route.workspace_id.as_ref() == Some(&record.workspace_id)
835 && selection.route.provider == record.provider
836 && response.native_session_index_contract_is_valid()
837 }
838 (
839 NodeRequest::CatalogNativeSessions { .. }
840 | NodeRequest::PageNativeSessions { .. }
841 | NodeRequest::PreviewNativeSession { .. }
842 | NodeRequest::IndexNativeSession { .. },
843 Err(_),
844 ) => true,
845 (
846 NodeRequest::CatalogNativeSessions { .. }
847 | NodeRequest::PageNativeSessions { .. }
848 | NodeRequest::PreviewNativeSession { .. }
849 | NodeRequest::IndexNativeSession { .. },
850 Ok(_),
851 ) => false,
852 (
853 _,
854 Ok(
855 NodeResponse::NativeSessionsCataloged { .. }
856 | NodeResponse::NativeSessionsPaged { .. }
857 | NodeResponse::NativeSessionPreviewed { .. }
858 | NodeResponse::NativeSessionIndexed { .. },
859 ),
860 ) => false,
861 _ => true,
862 };
863 if matches {
864 Ok(())
865 } else {
866 Err(NodeClientError::Protocol(
867 "node native session response does not match the request".to_owned(),
868 ))
869 }
870}
871
872fn validate_workspace_content_response(
873 expected: &NodeRequest,
874 response: &Result<NodeResponse, NodeFailure>,
875) -> Result<(), NodeClientError> {
876 let matches = match (expected, response) {
877 (
878 NodeRequest::ReadWorkspaceFile { workspace_id, path },
879 Ok(NodeResponse::WorkspaceFileRead { file }),
880 ) => &file.workspace_id == workspace_id && &file.path == path,
881 (
882 NodeRequest::WriteWorkspaceFile { workspace_id, path, text, .. },
883 Ok(NodeResponse::WorkspaceFileWritten { file }),
884 ) => {
885 &file.workspace_id == workspace_id
886 && &file.path == path
887 && matches!(
888 &file.content,
889 gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
890 text: written,
891 byte_len,
892 } if written == text
893 && u32::try_from(text.len()).ok() == Some(*byte_len)
894 )
895 }
896 (
897 NodeRequest::CreateWorkspaceFile { workspace_id, path },
898 Ok(NodeResponse::WorkspaceFileCreated { file }),
899 ) => {
900 &file.workspace_id == workspace_id
901 && &file.path == path
902 && file.revision.is_some()
903 && matches!(
904 &file.content,
905 gate4agent_node_protocol::WorkspaceFileContent::Utf8 {
906 text,
907 byte_len: 0,
908 } if text.is_empty()
909 )
910 }
911 (
912 NodeRequest::CreateWorkspaceDirectory { workspace_id, path },
913 Ok(NodeResponse::WorkspaceDirectoryCreated {
914 workspace_id: actual_workspace_id,
915 entry,
916 }),
917 ) => {
918 actual_workspace_id == workspace_id
919 && &entry.relative_path == path
920 && entry.kind == gate4agent_node_protocol::WorkspaceEntryKind::Directory
921 }
922 (
923 NodeRequest::ReadGitHistory { workspace_id, .. },
924 Ok(NodeResponse::GitHistoryRead { workspace_id: actual, .. }),
925 ) => actual == workspace_id,
926 (
927 NodeRequest::ReadGitDiff { workspace_id, request },
928 Ok(NodeResponse::GitDiffRead { workspace_id: actual, diff }),
929 ) => actual == workspace_id && diff.mode == request.mode && diff.path == request.path,
930 (
931 NodeRequest::ReadWorkspaceFile { .. }
932 | NodeRequest::WriteWorkspaceFile { .. }
933 | NodeRequest::CreateWorkspaceFile { .. }
934 | NodeRequest::CreateWorkspaceDirectory { .. }
935 | NodeRequest::ReadGitHistory { .. }
936 | NodeRequest::ReadGitDiff { .. },
937 Err(_),
938 ) => true,
939 (
940 NodeRequest::ReadWorkspaceFile { .. }
941 | NodeRequest::WriteWorkspaceFile { .. }
942 | NodeRequest::CreateWorkspaceFile { .. }
943 | NodeRequest::CreateWorkspaceDirectory { .. }
944 | NodeRequest::ReadGitHistory { .. }
945 | NodeRequest::ReadGitDiff { .. },
946 Ok(_),
947 ) => false,
948 (
949 _,
950 Ok(
951 NodeResponse::WorkspaceFileRead { .. }
952 | NodeResponse::WorkspaceFileWritten { .. }
953 | NodeResponse::WorkspaceFileCreated { .. }
954 | NodeResponse::WorkspaceDirectoryCreated { .. }
955 | NodeResponse::GitHistoryRead { .. }
956 | NodeResponse::GitDiffRead { .. },
957 ),
958 ) => false,
959 _ => true,
960 };
961 if matches {
962 Ok(())
963 } else {
964 Err(NodeClientError::Protocol(
965 "node workspace content response does not match the request".to_owned(),
966 ))
967 }
968}
969
970fn validate_session_task_response(
971 expected: &NodeRequest,
972 response: &Result<NodeResponse, NodeFailure>,
973) -> Result<(), NodeClientError> {
974 let matches = match (expected, response) {
975 (
976 NodeRequest::SetSessionTask {
977 record_id,
978 expected_revision,
979 target,
980 },
981 Ok(NodeResponse::SessionRecordUpdated { record }),
982 ) => session_task_record_matches(record, record_id, *expected_revision, target),
983 (NodeRequest::SetSessionTask { .. }, Err(_)) => true,
984 (NodeRequest::SetSessionTask { .. }, Ok(_)) => false,
985 _ => true,
986 };
987 if matches {
988 Ok(())
989 } else {
990 Err(NodeClientError::Protocol(
991 "node session task response does not match the request".to_owned(),
992 ))
993 }
994}
995
996fn validate_delivery_response(
997 expected: &NodeRequest,
998 response: &Result<NodeResponse, NodeFailure>,
999) -> Result<(), NodeClientError> {
1000 let matches = match (expected, response) {
1001 (
1002 NodeRequest::BeginDeliveryStage { manifest },
1003 Ok(NodeResponse::DeliveryStageBegun {
1004 manifest_digest, ..
1005 }),
1006 ) => manifest_digest == &manifest.manifest_digest,
1007 (
1008 NodeRequest::PutDeliveryBlobChunk {
1009 stage_id,
1010 blob_digest,
1011 offset,
1012 chunk_hex,
1013 },
1014 Ok(NodeResponse::DeliveryBlobChunkAccepted {
1015 stage_id: actual_stage_id,
1016 blob_digest: actual_blob_digest,
1017 next_offset,
1018 }),
1019 ) => {
1020 actual_stage_id == stage_id
1021 && actual_blob_digest == blob_digest
1022 && offset
1023 .checked_add(chunk_hex.raw_len() as u64)
1024 .is_some_and(|expected_offset| expected_offset == *next_offset)
1025 }
1026 (
1027 NodeRequest::CommitDeliveryStage { .. },
1028 Ok(NodeResponse::DeliveryCommitted { .. }),
1029 ) => true,
1030 (
1031 NodeRequest::AbortDeliveryStage { stage_id },
1032 Ok(NodeResponse::DeliveryStageAborted {
1033 stage_id: actual_stage_id,
1034 }),
1035 ) => actual_stage_id == stage_id,
1036 (
1037 NodeRequest::BeginDeliveryStage { .. }
1038 | NodeRequest::PutDeliveryBlobChunk { .. }
1039 | NodeRequest::CommitDeliveryStage { .. }
1040 | NodeRequest::AbortDeliveryStage { .. },
1041 Err(_),
1042 ) => true,
1043 (
1044 NodeRequest::BeginDeliveryStage { .. }
1045 | NodeRequest::PutDeliveryBlobChunk { .. }
1046 | NodeRequest::CommitDeliveryStage { .. }
1047 | NodeRequest::AbortDeliveryStage { .. },
1048 Ok(_),
1049 ) => false,
1050 (
1051 _,
1052 Ok(
1053 NodeResponse::DeliveryStageBegun { .. }
1054 | NodeResponse::DeliveryBlobChunkAccepted { .. }
1055 | NodeResponse::DeliveryCommitted { .. }
1056 | NodeResponse::DeliveryStageAborted { .. },
1057 ),
1058 ) => false,
1059 _ => true,
1060 };
1061 if matches {
1062 Ok(())
1063 } else {
1064 Err(NodeClientError::Protocol(
1065 "node delivery response does not match the request".to_owned(),
1066 ))
1067 }
1068}
1069
1070fn validate_harness_mcp_response(
1071 expected: &NodeRequest,
1072 response: &Result<NodeResponse, NodeFailure>,
1073) -> Result<(), NodeClientError> {
1074 use NodeRequest as Request;
1075 use NodeResponse as Response;
1076 let matches = match (expected, response) {
1077 (Request::ArmHarnessMcpReservation { reservation_id, activation_digest, expires_at_unix_ms, .. },
1078 Ok(Response::Armed { reservation_id: echoed_id, activation_digest: echoed_digest, expires_at_unix_ms: echoed_expiry })) =>
1079 reservation_id == echoed_id && activation_digest == echoed_digest
1080 && expires_at_unix_ms == echoed_expiry,
1081 (Request::SpawnSpecWithHarnessMcp { reservation_id, activation_digest, .. },
1082 Ok(Response::Spawned { reservation_id: echoed_id, activation_digest: echoed_digest, receipt })) =>
1083 reservation_id == echoed_id && activation_digest == echoed_digest
1084 && receipt.harness_mcp_proxy.as_ref().is_some_and(|proxy| {
1085 &proxy.reservation_id == reservation_id
1086 && &proxy.activation_digest == activation_digest
1087 }),
1088 (Request::ActivateHarnessMcpReservation { reservation_id, activation_digest, record_id, session },
1089 Ok(Response::Activated { reservation_id: echoed_id, activation_digest: echoed_digest, record_id: echoed_record, session: echoed_session })) =>
1090 reservation_id == echoed_id && activation_digest == echoed_digest
1091 && record_id == echoed_record && session == echoed_session,
1092 (Request::AbortHarnessMcpReservation { reservation_id, activation_digest },
1093 Ok(Response::Aborted { reservation_id: echoed_id, activation_digest: echoed_digest })) =>
1094 reservation_id == echoed_id && activation_digest == echoed_digest,
1095 (Request::PutHarnessMcpReplyChunk { reservation_id, activation_digest, record_id, session, call_id, offset, final_chunk, chunk_hex },
1096 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 })) =>
1097 reservation_id == echoed_id && activation_digest == echoed_digest
1098 && record_id == echoed_record && session == echoed_session && call_id == echoed_call
1099 && offset.checked_add(u32::try_from(chunk_hex.raw_len()).unwrap_or(u32::MAX))
1100 == Some(*next_offset)
1101 && completed == final_chunk,
1102 (Request::RejectHarnessMcpCall { reservation_id, activation_digest, record_id, session, call_id, .. },
1103 Ok(Response::CallRejected { reservation_id: echoed_id, activation_digest: echoed_digest, record_id: echoed_record, session: echoed_session, call_id: echoed_call })) =>
1104 reservation_id == echoed_id && activation_digest == echoed_digest
1105 && record_id == echoed_record && session == echoed_session && call_id == echoed_call,
1106 (request, Err(_)) if request.required_capability()
1107 == Some(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY) => true,
1108 (request, Ok(_)) if request.required_capability()
1109 == Some(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY) => false,
1110 (_, Ok(response)) if response.requires_harness_mcp_proxy_capability() => false,
1111 _ => true,
1112 };
1113 if matches { Ok(()) } else {
1114 Err(NodeClientError::Protocol(
1115 "node harness MCP proxy response does not match the request".to_owned(),
1116 ))
1117 }
1118}
1119
1120fn session_task_record_matches(
1121 record: &gate4agent_node_protocol::ManagedSessionRecord,
1122 record_id: &gate4agent_node_protocol::SessionRecordId,
1123 expected_revision: u64,
1124 target: &gate4agent_node_protocol::SessionTaskTargetV1,
1125) -> bool {
1126 if &record.record_id != record_id { return false; }
1127 let next_revision = expected_revision.checked_add(1);
1128 match target {
1129 gate4agent_node_protocol::SessionTaskTargetV1::New => record
1130 .task_binding
1131 .as_ref()
1132 .is_some_and(|binding| {
1133 Some(binding.revision) == next_revision && binding.task_id.is_some()
1134 }),
1135 gate4agent_node_protocol::SessionTaskTargetV1::Existing { task_id } => record
1136 .task_binding
1137 .as_ref()
1138 .is_some_and(|binding| {
1139 (binding.revision == expected_revision
1140 || Some(binding.revision) == next_revision)
1141 && binding.task_id.as_ref() == Some(task_id)
1142 }),
1143 gate4agent_node_protocol::SessionTaskTargetV1::Clear => match &record.task_binding {
1144 None => expected_revision == 0,
1145 Some(binding) => binding.task_id.is_none()
1146 && (binding.revision == expected_revision
1147 || Some(binding.revision) == next_revision),
1148 },
1149 }
1150}
1151
1152fn ensure_node_request_required_capability(
1153 request: &NodeRequest,
1154 negotiated_capabilities: &[CapabilityId],
1155) -> Result<(), NodeClientError> {
1156 if let Some(required) = request.required_capability() {
1157 if !negotiated_capabilities.iter().any(|capability| capability.as_str() == required) {
1158 return Err(NodeClientError::UnsupportedCapability(required.to_owned()));
1159 }
1160 }
1161 if request.requires_spawn_spec_defaults_overrides_capability()
1162 && !has_capability(
1163 negotiated_capabilities,
1164 NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY,
1165 )
1166 {
1167 return Err(NodeClientError::UnsupportedCapability(
1168 NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY.to_owned(),
1169 ));
1170 }
1171 if request.requires_spawn_profile_revision_capability()
1172 && !has_capability(
1173 negotiated_capabilities,
1174 NODE_SPAWN_PROFILE_REVISION_CAPABILITY,
1175 )
1176 {
1177 return Err(NodeClientError::UnsupportedCapability(
1178 NODE_SPAWN_PROFILE_REVISION_CAPABILITY.to_owned(),
1179 ));
1180 }
1181 if request.requires_worktree_selection_capability()
1182 && !negotiated_capabilities.iter().any(|capability| {
1183 capability.as_str() == NODE_WORKTREE_SELECTION_CAPABILITY
1184 })
1185 {
1186 return Err(NodeClientError::UnsupportedCapability(
1187 NODE_WORKTREE_SELECTION_CAPABILITY.to_owned(),
1188 ));
1189 }
1190 if request.requires_child_environment_profile_capability()
1191 && !has_capability(
1192 negotiated_capabilities,
1193 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY,
1194 )
1195 {
1196 return Err(NodeClientError::UnsupportedCapability(
1197 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY.to_owned(),
1198 ));
1199 }
1200 if request.requires_session_bundle_materialization_capability()
1201 && !has_capability(
1202 negotiated_capabilities,
1203 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY,
1204 )
1205 {
1206 return Err(NodeClientError::UnsupportedCapability(
1207 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY.to_owned(),
1208 ));
1209 }
1210 if request.requires_history_context_pack_capability()
1211 && !has_capability(
1212 negotiated_capabilities,
1213 NODE_HISTORY_CONTEXT_PACK_CAPABILITY,
1214 )
1215 {
1216 return Err(NodeClientError::UnsupportedCapability(
1217 NODE_HISTORY_CONTEXT_PACK_CAPABILITY.to_owned(),
1218 ));
1219 }
1220 Ok(())
1221}
1222
1223#[cfg(test)]
1224fn ensure_server_frame_required_capability(
1225 frame: &ServerFrame,
1226 negotiated_capabilities: &[CapabilityId],
1227) -> Result<(), NodeClientError> {
1228 ensure_server_frame_required_capability_for_request(
1229 frame,
1230 negotiated_capabilities,
1231 None,
1232 )
1233}
1234
1235fn ensure_server_frame_required_capability_for_request(
1236 frame: &ServerFrame,
1237 negotiated_capabilities: &[CapabilityId],
1238 expected_request: Option<&NodeRequest>,
1239) -> Result<(), NodeClientError> {
1240 if matches!(frame, ServerFrame::Event(event)
1241 if !event.event.harness_mcp_contract_is_valid_at(current_unix_ms()?))
1242 {
1243 return Err(NodeClientError::Protocol(
1244 "node sent an invalid harness MCP read call".to_owned(),
1245 ));
1246 }
1247 if matches!(frame, ServerFrame::Reply(reply)
1248 if reply.result.as_ref().is_ok_and(|response| {
1249 !response.native_session_catalog_contract_is_valid()
1250 }))
1251 {
1252 return Err(NodeClientError::Protocol(
1253 "node sent an invalid native session catalog response".to_owned(),
1254 ));
1255 }
1256 if matches!(frame, ServerFrame::Reply(reply)
1257 if reply.result.as_ref().is_ok_and(|response| {
1258 response.requires_native_session_preview_capability()
1259 && !response.native_session_preview_contract_is_valid()
1260 }))
1261 {
1262 return Err(NodeClientError::Protocol(
1263 "node sent an invalid native session preview response".to_owned(),
1264 ));
1265 }
1266 if matches!(frame, ServerFrame::Reply(reply)
1267 if reply.result.as_ref().is_ok_and(|response| {
1268 response.requires_native_session_index_capability()
1269 && !response.native_session_index_contract_is_valid()
1270 }))
1271 {
1272 return Err(NodeClientError::Protocol(
1273 "node sent an invalid native session index response".to_owned(),
1274 ));
1275 }
1276 let required = match frame {
1277 ServerFrame::Reply(reply) => match reply.result.as_ref() {
1278 Ok(NodeResponse::DeliveryStageBegun { .. })
1279 | Ok(NodeResponse::DeliveryBlobChunkAccepted { .. })
1280 | Ok(NodeResponse::DeliveryCommitted { .. })
1281 | Ok(NodeResponse::DeliveryStageAborted { .. }) => {
1282 Some(NODE_DELIVERY_BUNDLE_V2_STAGE_COMMIT_CAPABILITY)
1283 }
1284 Err(NodeFailure {
1285 code:
1286 NodeFailureCode::DeliveryManifestInvalid
1287 | NodeFailureCode::UnknownDeliveryStage
1288 | NodeFailureCode::DeliveryStageConflict
1289 | NodeFailureCode::DeliveryBlobUnexpected
1290 | NodeFailureCode::DeliveryChunkOutOfOrder
1291 | NodeFailureCode::DeliveryBlobDigestMismatch
1292 | NodeFailureCode::DeliveryBundleDigestMismatch
1293 | NodeFailureCode::DeliveryStageIncomplete
1294 | NodeFailureCode::DeliveryStageStorageFailed,
1295 ..
1296 }) => Some(NODE_DELIVERY_BUNDLE_V2_STAGE_COMMIT_CAPABILITY),
1297 Ok(response) if response.requires_harness_mcp_proxy_capability() => {
1298 Some(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY)
1299 }
1300 Err(NodeFailure { code, .. }) if matches!(code,
1301 NodeFailureCode::HarnessMcpUnavailable
1302 | NodeFailureCode::ReservationNotFound
1303 | NodeFailureCode::ReservationConflict
1304 | NodeFailureCode::ReservationExpired
1305 | NodeFailureCode::BindingMismatch
1306 | NodeFailureCode::NotActivated
1307 | NodeFailureCode::CallNotFound
1308 | NodeFailureCode::ChunkOutOfOrder
1309 | NodeFailureCode::ResponseTooLarge) => {
1310 Some(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY)
1311 }
1312 Ok(NodeResponse::WorkspaceFileRead { .. }) => {
1313 Some(NODE_WORKSPACE_FILE_READ_CAPABILITY)
1314 }
1315 Ok(NodeResponse::WorkspaceFileWritten { .. }) => {
1316 Some(NODE_WORKSPACE_FILE_WRITE_CAPABILITY)
1317 }
1318 Ok(NodeResponse::WorkspaceFileCreated { .. })
1319 | Ok(NodeResponse::WorkspaceDirectoryCreated { .. }) => {
1320 Some(NODE_WORKSPACE_ENTRY_CREATE_CAPABILITY)
1321 }
1322 Ok(NodeResponse::GitHistoryRead { .. }) | Ok(NodeResponse::GitDiffRead { .. }) => {
1323 Some(NODE_GIT_READ_CAPABILITY)
1324 }
1325 Ok(NodeResponse::HostDirectoriesBrowsed { .. }) => {
1326 Some(CAPABILITY_HOST_DIRECTORY_BROWSE_V1)
1327 }
1328 Ok(NodeResponse::StandaloneWorkspaceCreated { .. }) => {
1329 Some(NODE_STANDALONE_WORKSPACE_LIFECYCLE_CAPABILITY)
1330 }
1331 Ok(NodeResponse::ProviderSessionIndexed { .. }) => {
1332 Some(NODE_PROVIDER_SESSION_REFERENCE_INDEX_CAPABILITY)
1333 }
1334 Ok(NodeResponse::NativeSessionIndexed { .. }) => {
1335 Some(NODE_NATIVE_SESSION_INDEX_CAPABILITY)
1336 }
1337 Ok(NodeResponse::NativeSessionsCataloged { .. }) => {
1338 Some(NODE_NATIVE_SESSION_CATALOG_CAPABILITY)
1339 }
1340 Ok(NodeResponse::NativeSessionsPaged { .. }) => {
1341 Some(NODE_NATIVE_SESSION_CATALOG_PAGING_CAPABILITY)
1342 }
1343 Err(NodeFailure {
1344 code: NodeFailureCode::StaleNativeSessionCatalog,
1345 ..
1346 }) => Some(NODE_NATIVE_SESSION_CATALOG_PAGING_CAPABILITY),
1347 Ok(NodeResponse::NativeSessionPreviewed { .. })
1348 | Ok(NodeResponse::SessionRecordPreviewed { .. }) => {
1349 Some(NODE_NATIVE_SESSION_PREVIEW_CAPABILITY)
1350 }
1351 Ok(NodeResponse::ContextPackForSessionRecordExported { .. }) => {
1352 Some(NODE_SESSION_RECORD_CONTEXT_EXPORT_CAPABILITY)
1353 }
1354 Ok(NodeResponse::SessionRecordUpdated { .. })
1355 if matches!(expected_request, Some(NodeRequest::SetSessionTask { .. })) => {
1356 Some(NODE_SESSION_TASK_CORRELATION_CAPABILITY)
1357 }
1358 Ok(response) if response.requires_spawn_spec_defaults_overrides_capability() => {
1359 Some(NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY)
1360 }
1361 _ => None,
1362 },
1363 ServerFrame::Event(event) if event.event.requires_harness_mcp_proxy_capability() => {
1364 Some(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY)
1365 }
1366 _ => None,
1367 };
1368 if let Some(required) = required {
1369 if !negotiated_capabilities
1370 .iter()
1371 .any(|capability| capability.as_str() == required)
1372 {
1373 return Err(NodeClientError::UnsupportedCapability(required.to_owned()));
1374 }
1375 }
1376 if matches!(frame, ServerFrame::Reply(reply)
1377 if reply.result.as_ref().is_ok_and(|response|
1378 response.requires_spawn_profile_revision_capability()))
1379 && !has_capability(
1380 negotiated_capabilities,
1381 NODE_SPAWN_PROFILE_REVISION_CAPABILITY,
1382 )
1383 {
1384 return Err(NodeClientError::UnsupportedCapability(
1385 NODE_SPAWN_PROFILE_REVISION_CAPABILITY.to_owned(),
1386 ));
1387 }
1388 let contains_managed_worktree = server_frame_contains_managed_worktree(frame);
1389 if contains_managed_worktree
1390 && !has_capability(
1391 negotiated_capabilities,
1392 NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY,
1393 )
1394 {
1395 return Err(NodeClientError::UnsupportedCapability(
1396 NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY.to_owned(),
1397 ));
1398 }
1399 let requires_worktree_selection = contains_managed_worktree || match frame {
1400 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
1401 NodeResponse::requires_worktree_selection_capability,
1402 ),
1403 _ => false,
1404 };
1405 if requires_worktree_selection
1406 && !negotiated_capabilities.iter().any(|capability| {
1407 capability.as_str() == NODE_WORKTREE_SELECTION_CAPABILITY
1408 })
1409 {
1410 return Err(NodeClientError::UnsupportedCapability(
1411 NODE_WORKTREE_SELECTION_CAPABILITY.to_owned(),
1412 ));
1413 }
1414 let contains_environment_profile = match frame {
1415 ServerFrame::Hello(hello) => hello
1416 .snapshot
1417 .requires_child_environment_profile_capability(),
1418 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
1419 NodeResponse::requires_child_environment_profile_capability,
1420 ),
1421 ServerFrame::Event(event) => event
1422 .event
1423 .requires_child_environment_profile_capability(),
1424 ServerFrame::Challenge(_) => false,
1425 };
1426 if contains_environment_profile
1427 && !has_capability(
1428 negotiated_capabilities,
1429 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY,
1430 )
1431 {
1432 return Err(NodeClientError::Protocol(
1433 "node sent child environment profile metadata without negotiating the capability"
1434 .to_owned(),
1435 ));
1436 }
1437 let contains_bundle = match frame {
1438 ServerFrame::Hello(hello) => hello
1439 .snapshot
1440 .requires_session_bundle_materialization_capability(),
1441 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
1442 NodeResponse::requires_session_bundle_materialization_capability,
1443 ),
1444 ServerFrame::Event(event) => event
1445 .event
1446 .requires_session_bundle_materialization_capability(),
1447 ServerFrame::Challenge(_) => false,
1448 };
1449 if contains_bundle
1450 && !has_capability(
1451 negotiated_capabilities,
1452 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY,
1453 )
1454 {
1455 return Err(NodeClientError::Protocol(
1456 "node sent session bundle materialization metadata without negotiating the capability"
1457 .to_owned(),
1458 ));
1459 }
1460 let contains_history_context_pack = match frame {
1461 ServerFrame::Hello(hello) => hello
1462 .snapshot
1463 .requires_history_context_pack_capability(),
1464 ServerFrame::Reply(reply) => {
1465 reply.result.as_ref().ok().is_some_and(
1466 NodeResponse::requires_history_context_pack_capability,
1467 ) || reply.result.as_ref().err().is_some_and(|failure| {
1468 matches!(
1469 failure.code,
1470 NodeFailureCode::UnknownContextPack
1471 | NodeFailureCode::ContextPackBusy
1472 | NodeFailureCode::ContextPackMaterializationFailed
1473 )
1474 })
1475 }
1476 ServerFrame::Event(event) => event
1477 .event
1478 .requires_history_context_pack_capability(),
1479 ServerFrame::Challenge(_) => false,
1480 };
1481 if contains_history_context_pack
1482 && !has_capability(
1483 negotiated_capabilities,
1484 NODE_HISTORY_CONTEXT_PACK_CAPABILITY,
1485 )
1486 {
1487 return Err(NodeClientError::Protocol(
1488 "node sent history context pack metadata without negotiating the capability"
1489 .to_owned(),
1490 ));
1491 }
1492 let contains_session_task_binding = match frame {
1493 ServerFrame::Hello(hello) => hello
1494 .snapshot
1495 .session_records
1496 .iter()
1497 .any(|record| record.task_binding.is_some()),
1498 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(|response| {
1499 match response {
1500 NodeResponse::Snapshot { snapshot, .. } => snapshot
1501 .session_records
1502 .iter()
1503 .any(|record| record.task_binding.is_some()),
1504 NodeResponse::Resync { snapshot, events, .. } => snapshot
1505 .session_records
1506 .iter()
1507 .any(|record| record.task_binding.is_some())
1508 || events.iter().any(|event| matches!(&event.event,
1509 NodeEvent::SessionRecordUpserted { record }
1510 if record.task_binding.is_some())),
1511 NodeResponse::SessionRecordUpdated { record }
1512 | NodeResponse::ProviderSessionIndexed { record }
1513 | NodeResponse::NativeSessionIndexed { record, .. }
1514 | NodeResponse::SessionRecordResumed { record, .. } => {
1515 record.task_binding.is_some()
1516 }
1517 _ => false,
1518 }
1519 }),
1520 ServerFrame::Event(event) => matches!(&event.event,
1521 NodeEvent::SessionRecordUpserted { record } if record.task_binding.is_some()),
1522 ServerFrame::Challenge(_) => false,
1523 };
1524 if contains_session_task_binding
1525 && !has_capability(
1526 negotiated_capabilities,
1527 NODE_SESSION_TASK_CORRELATION_CAPABILITY,
1528 )
1529 {
1530 return Err(NodeClientError::Protocol(
1531 "node sent session task correlation metadata without negotiating the capability"
1532 .to_owned(),
1533 ));
1534 }
1535 if server_frame_contains_agent_progress(frame)
1536 && !has_capability(
1537 negotiated_capabilities,
1538 NODE_AGENT_PROGRESS_SNAPSHOT_CAPABILITY,
1539 )
1540 {
1541 return Err(NodeClientError::Protocol(
1542 "node sent agent progress metadata without negotiating the capability"
1543 .to_owned(),
1544 ));
1545 }
1546 Ok(())
1547}
1548
1549fn has_capability(capabilities: &[CapabilityId], required: &str) -> bool {
1550 capabilities
1551 .iter()
1552 .any(|capability| capability.as_str() == required)
1553}
1554
1555fn ensure_node_snapshot_agent_progress_capability(
1556 snapshot: &NodeSnapshot,
1557 negotiated_capabilities: &[CapabilityId],
1558) -> Result<(), NodeClientError> {
1559 if !snapshot.agent_progress.is_empty()
1560 && !has_capability(
1561 negotiated_capabilities,
1562 NODE_AGENT_PROGRESS_SNAPSHOT_CAPABILITY,
1563 )
1564 {
1565 return Err(NodeClientError::Protocol(
1566 "node sent agent progress metadata without negotiating the capability"
1567 .to_owned(),
1568 ));
1569 }
1570 Ok(())
1571}
1572
1573fn server_frame_contains_agent_progress(frame: &ServerFrame) -> bool {
1574 match frame {
1575 ServerFrame::Hello(hello) => !hello.snapshot.agent_progress.is_empty(),
1576 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(|response| {
1577 match response {
1578 NodeResponse::Snapshot { snapshot, .. }
1579 | NodeResponse::Resync { snapshot, .. } => !snapshot.agent_progress.is_empty(),
1580 _ => false,
1581 }
1582 }),
1583 ServerFrame::Challenge(_) | ServerFrame::Event(_) => false,
1584 }
1585}
1586
1587fn server_frame_contains_managed_worktree(frame: &ServerFrame) -> bool {
1588 match frame {
1589 ServerFrame::Hello(hello) => node_snapshot_contains_managed_worktree(&hello.snapshot),
1590 ServerFrame::Reply(reply) => reply
1591 .result
1592 .as_ref()
1593 .ok()
1594 .is_some_and(node_response_contains_managed_worktree),
1595 ServerFrame::Event(event) => node_event_is_managed_worktree(&event.event),
1596 ServerFrame::Challenge(_) => false,
1597 }
1598}
1599
1600fn node_response_contains_managed_worktree(response: &NodeResponse) -> bool {
1601 match response {
1602 NodeResponse::Snapshot { snapshot, .. } => node_snapshot_contains_managed_worktree(snapshot),
1603 NodeResponse::Resync { snapshot, events, .. } => {
1604 node_snapshot_contains_managed_worktree(snapshot)
1605 || events
1606 .iter()
1607 .any(|event| node_event_is_managed_worktree(&event.event))
1608 }
1609 NodeResponse::WorkspaceRegistered { workspace }
1610 | NodeResponse::StandaloneWorkspaceCreated { workspace }
1611 | NodeResponse::WorktreeCreated { workspace, .. } => {
1612 workspace_contains_managed_worktree_metadata(workspace)
1613 }
1614 NodeResponse::ManagedWorktreeSpawnAccepted { .. }
1615 | NodeResponse::ManagedWorktreeCleanup { .. } => true,
1616 _ => false,
1617 }
1618}
1619
1620fn node_event_is_managed_worktree(event: &NodeEvent) -> bool {
1621 match event {
1622 NodeEvent::WorkspaceAdded { workspace } => {
1623 workspace_contains_managed_worktree_metadata(workspace)
1624 }
1625 NodeEvent::ManagedWorktreeUpserted { .. }
1626 | NodeEvent::ManagedWorktreeRemoved { .. } => true,
1627 _ => false,
1628 }
1629}
1630
1631fn workspace_contains_managed_worktree_metadata(workspace: &WorkspaceSnapshot) -> bool {
1632 workspace.worktree_service_mode.is_some()
1633 || workspace.managed_worktree_profiles.is_some()
1634}
1635
1636fn node_snapshot_contains_managed_worktree(snapshot: &NodeSnapshot) -> bool {
1637 !snapshot.managed_worktrees.is_empty()
1638 || snapshot
1639 .workspaces
1640 .iter()
1641 .any(workspace_contains_managed_worktree_metadata)
1642}
1643
1644fn ensure_server_frame_terminal_capability(
1645 frame: &ServerFrame,
1646 terminal_frame_events_enabled: bool,
1647) -> Result<(), NodeClientError> {
1648 let contains_terminal_frame_event = match frame {
1649 ServerFrame::Event(NodeEventEnvelope {
1650 event: NodeEvent::TerminalFrame { .. },
1651 ..
1652 }) => true,
1653 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(|response| {
1654 matches!(response, NodeResponse::Resync { events, .. }
1655 if events.iter().any(|event| {
1656 matches!(&event.event, NodeEvent::TerminalFrame { .. })
1657 }))
1658 }),
1659 ServerFrame::Challenge(_) | ServerFrame::Hello(_) | ServerFrame::Event(_) => false,
1660 };
1661 if contains_terminal_frame_event && !terminal_frame_events_enabled {
1662 return Err(NodeClientError::UnsupportedCapability(
1663 NODE_TERMINAL_FRAME_EVENTS_CAPABILITY.to_owned(),
1664 ));
1665 }
1666 Ok(())
1667}
1668
1669fn ensure_node_hello_path_capability(
1670 hello: &NodeHello,
1671 opaque_unix_paths_enabled: bool,
1672) -> Result<(), NodeClientError> {
1673 ensure_opaque_unix_path_capability(
1674 node_snapshot_contains_opaque_unix_path(&hello.snapshot),
1675 opaque_unix_paths_enabled,
1676 )
1677}
1678
1679fn ensure_node_request_path_capability(
1680 request: &NodeRequest,
1681 opaque_unix_paths_enabled: bool,
1682 repository_paths_enabled: bool,
1683) -> Result<(), NodeClientError> {
1684 ensure_opaque_unix_path_capability(
1685 node_request_contains_opaque_unix_path(request),
1686 opaque_unix_paths_enabled,
1687 )?;
1688 ensure_repository_path_capability(
1689 node_request_contains_tagged_repository_path(request),
1690 repository_paths_enabled,
1691 )
1692}
1693
1694fn ensure_server_frame_path_capability(
1695 frame: &ServerFrame,
1696 opaque_unix_paths_enabled: bool,
1697 repository_paths_enabled: bool,
1698) -> Result<(), NodeClientError> {
1699 ensure_opaque_unix_path_capability(
1700 server_frame_contains_opaque_unix_path(frame),
1701 opaque_unix_paths_enabled,
1702 )?;
1703 ensure_repository_path_capability(
1704 server_frame_contains_tagged_repository_path(frame),
1705 repository_paths_enabled,
1706 )
1707}
1708
1709fn ensure_node_hello_provider_capability(
1710 hello: &NodeHello,
1711 open_provider_ids_enabled: bool,
1712) -> Result<(), NodeClientError> {
1713 ensure_inbound_open_provider_capability(
1714 node_snapshot_contains_open_provider_id(&hello.snapshot),
1715 open_provider_ids_enabled,
1716 )
1717}
1718
1719fn ensure_node_hello_environment_profile_capability(
1720 hello: &NodeHello,
1721 negotiated_capabilities: &[CapabilityId],
1722) -> Result<(), NodeClientError> {
1723 if hello
1724 .snapshot
1725 .requires_child_environment_profile_capability()
1726 && !has_capability(
1727 negotiated_capabilities,
1728 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY,
1729 )
1730 {
1731 return Err(NodeClientError::Protocol(
1732 "node sent child environment profile metadata without negotiating the capability"
1733 .to_owned(),
1734 ));
1735 }
1736 Ok(())
1737}
1738
1739fn ensure_node_hello_bundle_materialization_capability(
1740 hello: &NodeHello,
1741 negotiated_capabilities: &[CapabilityId],
1742) -> Result<(), NodeClientError> {
1743 if hello
1744 .snapshot
1745 .requires_session_bundle_materialization_capability()
1746 && !has_capability(
1747 negotiated_capabilities,
1748 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY,
1749 )
1750 {
1751 return Err(NodeClientError::Protocol(
1752 "node sent session bundle materialization metadata without negotiating the capability"
1753 .to_owned(),
1754 ));
1755 }
1756 Ok(())
1757}
1758
1759fn ensure_node_hello_history_context_pack_capability(
1760 hello: &NodeHello,
1761 negotiated_capabilities: &[CapabilityId],
1762) -> Result<(), NodeClientError> {
1763 if hello
1764 .snapshot
1765 .requires_history_context_pack_capability()
1766 && !has_capability(
1767 negotiated_capabilities,
1768 NODE_HISTORY_CONTEXT_PACK_CAPABILITY,
1769 )
1770 {
1771 return Err(NodeClientError::Protocol(
1772 "node sent history context pack metadata without negotiating the capability"
1773 .to_owned(),
1774 ));
1775 }
1776 Ok(())
1777}
1778
1779fn ensure_node_request_provider_capability(
1780 request: &NodeRequest,
1781 open_provider_ids_enabled: bool,
1782) -> Result<(), NodeClientError> {
1783 match request {
1784 NodeRequest::Spawn { provider, .. } => {
1785 ensure_outbound_provider_id_capability(provider, open_provider_ids_enabled)?;
1786 }
1787 NodeRequest::IndexProviderSession { provider, .. } => {
1788 ensure_outbound_provider_id_capability(provider, open_provider_ids_enabled)?;
1789 }
1790 NodeRequest::CatalogNativeSessions { route, .. }
1791 | NodeRequest::PageNativeSessions { route, .. } => {
1792 ensure_outbound_provider_id_capability(&route.provider, open_provider_ids_enabled)?;
1793 }
1794 NodeRequest::PreviewNativeSession { selection, .. }
1795 | NodeRequest::IndexNativeSession { selection, .. } => {
1796 ensure_outbound_provider_id_capability(
1797 &selection.route.provider,
1798 open_provider_ids_enabled,
1799 )?;
1800 }
1801 NodeRequest::SpawnSpec { spec }
1802 | NodeRequest::ArmHarnessMcpReservation { spawn_spec: spec, .. }
1803 | NodeRequest::SpawnSpecWithHarnessMcp { spec, .. }
1804 | NodeRequest::SpawnManagedWorktree {
1805 request: gate4agent_node_protocol::ManagedWorktreeSpawnRequest {
1806 spawn_spec: spec,
1807 ..
1808 },
1809 }
1810 | NodeRequest::SpawnManagedWorktreeV2 {
1811 request: gate4agent_node_protocol::ManagedWorktreeSpawnRequestV2 {
1812 spawn_spec: spec,
1813 ..
1814 },
1815 } => {
1816 if let gate4agent_node_protocol::SpawnOverride::Set { value: provider } =
1817 &spec.overrides.provider
1818 {
1819 ensure_outbound_provider_id_capability(provider, open_provider_ids_enabled)?;
1820 }
1821 }
1822 NodeRequest::ForgetContextPack { .. }
1823 | NodeRequest::ResolveDurableContextPack { .. }
1824 | NodeRequest::ReadContextPack { .. } => {
1825 if !open_provider_ids_enabled {
1826 return Err(NodeClientError::UnsupportedCapability(
1827 NODE_PROVIDER_ID_OPEN_CAPABILITY.to_owned(),
1828 ));
1829 }
1830 }
1831 NodeRequest::Snapshot
1832 | NodeRequest::Resync { .. }
1833 | NodeRequest::ActivateHarnessMcpReservation { .. }
1834 | NodeRequest::AbortHarnessMcpReservation { .. }
1835 | NodeRequest::PutHarnessMcpReplyChunk { .. }
1836 | NodeRequest::RejectHarnessMcpCall { .. }
1837 | NodeRequest::BeginDeliveryStage { .. }
1838 | NodeRequest::PutDeliveryBlobChunk { .. }
1839 | NodeRequest::CommitDeliveryStage { .. }
1840 | NodeRequest::AbortDeliveryStage { .. }
1841 | NodeRequest::BrowseHostDirectories { .. }
1842 | NodeRequest::InspectWorkspace { .. }
1843 | NodeRequest::ReadWorkspaceFile { .. }
1844 | NodeRequest::WriteWorkspaceFile { .. }
1845 | NodeRequest::CreateWorkspaceFile { .. }
1846 | NodeRequest::CreateWorkspaceDirectory { .. }
1847 | NodeRequest::ReadGitHistory { .. }
1848 | NodeRequest::ReadGitDiff { .. }
1849 | NodeRequest::AcquireController { .. }
1850 | NodeRequest::ReleaseController
1851 | NodeRequest::RegisterWorkspace { .. }
1852 | NodeRequest::CreateStandaloneWorkspace { .. }
1853 | NodeRequest::UnregisterWorkspace { .. }
1854 | NodeRequest::CreateWorktree { .. }
1855 | NodeRequest::RemoveWorktree { .. }
1856 | NodeRequest::CleanupManagedWorktree { .. }
1857 | NodeRequest::Resume { .. }
1858 | NodeRequest::RenameSessionRecord { .. }
1859 | NodeRequest::SetSessionTask { .. }
1860 | NodeRequest::ResumeSessionRecord { .. }
1861 | NodeRequest::ForgetSessionRecord { .. }
1862 | NodeRequest::PreviewSessionRecord { .. }
1863 | NodeRequest::DiscoverHistory { .. }
1864 | NodeRequest::LoadHistory { .. }
1865 | NodeRequest::ExportContextPackForSessionRecord { .. }
1866 | NodeRequest::ExportContextPack { .. }
1867 | NodeRequest::Prompt { .. }
1868 | NodeRequest::Paste { .. }
1869 | NodeRequest::Input { .. }
1870 | NodeRequest::TerminalBytes { .. }
1871 | NodeRequest::TerminalControl { .. }
1872 | NodeRequest::Resize { .. }
1873 | NodeRequest::Interrupt { .. }
1874 | NodeRequest::Stop { .. }
1875 | NodeRequest::Remove { .. }
1876 | NodeRequest::ResolveInteraction { .. }
1877 | NodeRequest::SetSessionMode { .. }
1878 | NodeRequest::SetSessionConfigOption { .. }
1879 | NodeRequest::SetSessionModel { .. }
1880 | NodeRequest::Shutdown => {}
1881 }
1882 Ok(())
1883}
1884
1885fn ensure_outbound_provider_id_capability(
1886 provider: &gate4agent_node_protocol::AgentId,
1887 open_provider_ids_enabled: bool,
1888) -> Result<(), NodeClientError> {
1889 if !provider_id_is_legacy(provider) && !open_provider_ids_enabled {
1890 return Err(NodeClientError::UnsupportedCapability(
1891 NODE_PROVIDER_ID_OPEN_CAPABILITY.to_owned(),
1892 ));
1893 }
1894 Ok(())
1895}
1896
1897fn ensure_server_frame_provider_capability(
1898 frame: &ServerFrame,
1899 open_provider_ids_enabled: bool,
1900) -> Result<(), NodeClientError> {
1901 ensure_inbound_open_provider_capability(
1902 server_frame_contains_open_provider_id(frame),
1903 open_provider_ids_enabled,
1904 )
1905}
1906
1907fn ensure_inbound_open_provider_capability(
1908 contains_open_provider_id: bool,
1909 open_provider_ids_enabled: bool,
1910) -> Result<(), NodeClientError> {
1911 if contains_open_provider_id && !open_provider_ids_enabled {
1912 return Err(NodeClientError::Protocol(
1913 "node sent an open provider ID without negotiating the capability".to_owned(),
1914 ));
1915 }
1916 Ok(())
1917}
1918
1919fn ensure_opaque_unix_path_capability(
1920 contains_opaque_unix_path: bool,
1921 opaque_unix_paths_enabled: bool,
1922
1923) -> Result<(), NodeClientError> {
1924 if contains_opaque_unix_path && !opaque_unix_paths_enabled {
1925 return Err(NodeClientError::Protocol(
1926 "node sent or received opaque Unix path bytes without negotiating the capability"
1927 .to_owned(),
1928 ));
1929 }
1930 Ok(())
1931}
1932
1933fn ensure_repository_path_capability(
1934 contains_tagged_repository_path: bool,
1935 repository_paths_enabled: bool,
1936) -> Result<(), NodeClientError> {
1937 if contains_tagged_repository_path && !repository_paths_enabled {
1938 return Err(NodeClientError::Protocol(
1939 "node sent tagged repository path bytes without negotiating the capability"
1940 .to_owned(),
1941 ));
1942 }
1943 Ok(())
1944}
1945
1946fn node_request_contains_opaque_unix_path(request: &NodeRequest) -> bool {
1947 match request {
1948 NodeRequest::RegisterWorkspace { root, .. }
1949 | NodeRequest::CreateStandaloneWorkspace { root, .. } => {
1950 root.as_unix_bytes().is_some()
1951 }
1952 NodeRequest::BrowseHostDirectories { directory, after } => {
1953 directory.as_ref().is_some_and(|path| path.as_unix_bytes().is_some())
1954 || after.as_ref().is_some_and(|path| path.as_unix_bytes().is_some())
1955 }
1956 NodeRequest::CreateWorktree { target_root, .. }
1957 | NodeRequest::RemoveWorktree { target_root, .. } => {
1958 target_root.as_unix_bytes().is_some()
1959 }
1960 NodeRequest::Snapshot
1961 | NodeRequest::Resync { .. }
1962 | NodeRequest::ArmHarnessMcpReservation { .. }
1963 | NodeRequest::SpawnSpecWithHarnessMcp { .. }
1964 | NodeRequest::ActivateHarnessMcpReservation { .. }
1965 | NodeRequest::AbortHarnessMcpReservation { .. }
1966 | NodeRequest::PutHarnessMcpReplyChunk { .. }
1967 | NodeRequest::RejectHarnessMcpCall { .. }
1968 | NodeRequest::BeginDeliveryStage { .. }
1969 | NodeRequest::PutDeliveryBlobChunk { .. }
1970 | NodeRequest::CommitDeliveryStage { .. }
1971 | NodeRequest::AbortDeliveryStage { .. }
1972 | NodeRequest::InspectWorkspace { .. }
1973 | NodeRequest::ReadWorkspaceFile { .. }
1974 | NodeRequest::WriteWorkspaceFile { .. }
1975 | NodeRequest::CreateWorkspaceFile { .. }
1976 | NodeRequest::CreateWorkspaceDirectory { .. }
1977 | NodeRequest::ReadGitHistory { .. }
1978 | NodeRequest::ReadGitDiff { .. }
1979 | NodeRequest::AcquireController { .. }
1980 | NodeRequest::ReleaseController
1981 | NodeRequest::UnregisterWorkspace { .. }
1982 | NodeRequest::Spawn { .. }
1983 | NodeRequest::SpawnSpec { .. }
1984 | NodeRequest::SpawnManagedWorktree { .. }
1985 | NodeRequest::SpawnManagedWorktreeV2 { .. }
1986 | NodeRequest::CleanupManagedWorktree { .. }
1987 | NodeRequest::Resume { .. }
1988 | NodeRequest::RenameSessionRecord { .. }
1989 | NodeRequest::SetSessionTask { .. }
1990 | NodeRequest::IndexProviderSession { .. }
1991 | NodeRequest::IndexNativeSession { .. }
1992 | NodeRequest::ResumeSessionRecord { .. }
1993 | NodeRequest::ForgetSessionRecord { .. }
1994 | NodeRequest::CatalogNativeSessions { .. }
1995 | NodeRequest::PageNativeSessions { .. }
1996 | NodeRequest::PreviewNativeSession { .. }
1997 | NodeRequest::PreviewSessionRecord { .. }
1998 | NodeRequest::DiscoverHistory { .. }
1999 | NodeRequest::LoadHistory { .. }
2000 | NodeRequest::ExportContextPackForSessionRecord { .. }
2001 | NodeRequest::ExportContextPack { .. }
2002 | NodeRequest::ForgetContextPack { .. }
2003 | NodeRequest::ResolveDurableContextPack { .. }
2004 | NodeRequest::ReadContextPack { .. }
2005 | NodeRequest::Prompt { .. }
2006 | NodeRequest::Paste { .. }
2007 | NodeRequest::Input { .. }
2008 | NodeRequest::TerminalBytes { .. }
2009 | NodeRequest::TerminalControl { .. }
2010 | NodeRequest::Resize { .. }
2011 | NodeRequest::Interrupt { .. }
2012 | NodeRequest::Stop { .. }
2013 | NodeRequest::Remove { .. }
2014 | NodeRequest::ResolveInteraction { .. }
2015 | NodeRequest::SetSessionMode { .. }
2016 | NodeRequest::SetSessionConfigOption { .. }
2017 | NodeRequest::SetSessionModel { .. }
2018 | NodeRequest::Shutdown => false,
2019 }
2020}
2021
2022fn server_frame_contains_open_provider_id(frame: &ServerFrame) -> bool {
2023 match frame {
2024 ServerFrame::Hello(hello) => node_snapshot_contains_open_provider_id(&hello.snapshot),
2025 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
2026 node_response_contains_open_provider_id,
2027 ),
2028 ServerFrame::Event(event) => node_event_contains_open_provider_id(&event.event),
2029 ServerFrame::Challenge(_) => false,
2030 }
2031}
2032
2033fn node_response_contains_open_provider_id(response: &NodeResponse) -> bool {
2034 match response {
2035 NodeResponse::Snapshot { snapshot, .. } => {
2036 node_snapshot_contains_open_provider_id(snapshot)
2037 }
2038 NodeResponse::Resync { snapshot, events, .. } => {
2039 node_snapshot_contains_open_provider_id(snapshot)
2040 || events.iter().any(|event| {
2041 node_event_contains_open_provider_id(&event.event)
2042 })
2043 }
2044 NodeResponse::SessionRecordUpdated { record }
2045 | NodeResponse::ProviderSessionIndexed { record }
2046 | NodeResponse::NativeSessionIndexed { record, .. }
2047 | NodeResponse::SessionRecordResumed { record, .. } => {
2048 managed_session_record_contains_open_provider_id(record)
2049 }
2050 NodeResponse::WorkspaceRegistered { workspace }
2051 | NodeResponse::StandaloneWorkspaceCreated { workspace }
2052 | NodeResponse::WorktreeCreated { workspace, .. } => {
2053 workspace_contains_open_provider_id(workspace)
2054 }
2055 NodeResponse::SpawnSpecAccepted { receipt }
2056 | NodeResponse::Spawned { receipt, .. } => {
2057 resolved_spawn_receipt_contains_open_provider_id(receipt)
2058 }
2059 NodeResponse::ManagedWorktreeSpawnAccepted { receipt } => {
2060 resolved_spawn_receipt_contains_open_provider_id(&receipt.spawn)
2061 }
2062 NodeResponse::ContextPackExported { context }
2063 | NodeResponse::ContextPackForSessionRecordExported { context, .. }
2064 | NodeResponse::DurableContextPackResolved { context } => {
2065 context_pack_contains_open_provider_id(context)
2066 }
2067 NodeResponse::NativeSessionsCataloged { route, .. }
2068 | NodeResponse::NativeSessionsPaged { route, .. } => {
2069 !provider_id_is_legacy(&route.provider)
2070 }
2071 NodeResponse::NativeSessionPreviewed { selection, .. } => {
2072 !provider_id_is_legacy(&selection.route.provider)
2073 }
2074 NodeResponse::Armed { .. }
2075 | NodeResponse::Activated { .. }
2076 | NodeResponse::Aborted { .. }
2077 | NodeResponse::ReplyChunkAccepted { .. }
2078 | NodeResponse::CallRejected { .. }
2079 | NodeResponse::WorkspaceInspected { .. }
2080 | NodeResponse::DeliveryStageBegun { .. }
2081 | NodeResponse::DeliveryBlobChunkAccepted { .. }
2082 | NodeResponse::DeliveryCommitted { .. }
2083 | NodeResponse::DeliveryStageAborted { .. }
2084 | NodeResponse::HostDirectoriesBrowsed { .. }
2085 | NodeResponse::WorkspaceFileRead { .. }
2086 | NodeResponse::WorkspaceFileWritten { .. }
2087 | NodeResponse::WorkspaceFileCreated { .. }
2088 | NodeResponse::WorkspaceDirectoryCreated { .. }
2089 | NodeResponse::GitHistoryRead { .. }
2090 | NodeResponse::GitDiffRead { .. }
2091 | NodeResponse::Controller { .. }
2092 | NodeResponse::SpawnAccepted { .. }
2093 | NodeResponse::ManagedWorktreeCleanup { .. }
2094 | NodeResponse::SessionRecordForgotten { .. }
2095 | NodeResponse::SessionRecordPreviewed { .. }
2096 | NodeResponse::HistoryDiscovered { .. }
2097 | NodeResponse::HistoryLoaded { .. }
2098 | NodeResponse::ContextPackForgotten { .. }
2099 | NodeResponse::ContextPackBytesRead { .. }
2100 | NodeResponse::WorkspaceUnregistered { .. }
2101 | NodeResponse::WorktreeRemoved { .. }
2102 | NodeResponse::Accepted
2103 | NodeResponse::ShuttingDown => false,
2104 }
2105}
2106
2107fn node_snapshot_contains_open_provider_id(snapshot: &NodeSnapshot) -> bool {
2108 snapshot.enabled_providers.iter().any(|provider| {
2109 !provider_id_is_legacy(provider)
2110 }) || snapshot.provider_runtime_statuses.iter().any(|status| {
2111 !provider_id_is_legacy(status.provider())
2112 }) || snapshot.session_records.iter().any(
2113 managed_session_record_contains_open_provider_id,
2114 ) || snapshot.workspaces.iter().any(workspace_contains_open_provider_id)
2115}
2116
2117fn managed_session_record_contains_open_provider_id(
2118 record: &gate4agent_node_protocol::ManagedSessionRecord,
2119) -> bool {
2120 !provider_id_is_legacy(&record.provider)
2121 || record.context.as_ref().is_some_and(
2122 context_pack_contains_open_provider_id,
2123 )
2124}
2125
2126fn resolved_spawn_receipt_contains_open_provider_id(
2127 receipt: &gate4agent_node_protocol::ResolvedSpawnReceipt,
2128) -> bool {
2129 !provider_id_is_legacy(&receipt.provider)
2130 || receipt.context.as_ref().is_some_and(
2131 context_pack_contains_open_provider_id,
2132 )
2133}
2134
2135fn context_pack_contains_open_provider_id(
2136 context: &gate4agent_node_protocol::ResolvedContextPackReceipt,
2137) -> bool {
2138 !provider_id_is_legacy(&context.lineage.source_provider)
2139}
2140
2141fn workspace_contains_open_provider_id(workspace: &WorkspaceSnapshot) -> bool {
2142 workspace.sessions.iter().any(|session| {
2143 !provider_id_is_legacy(&session.agent_id)
2144 })
2145}
2146
2147fn node_event_contains_open_provider_id(event: &NodeEvent) -> bool {
2148 match event {
2149 NodeEvent::WorkspaceAdded { workspace } => {
2150 workspace_contains_open_provider_id(workspace)
2151 }
2152 NodeEvent::SessionRecordUpserted { record } => {
2153 managed_session_record_contains_open_provider_id(record)
2154 }
2155 NodeEvent::HarnessMcpReadCall { .. }
2156 | NodeEvent::Control { .. }
2157 | NodeEvent::SessionRecordHistorySummarized { .. }
2158 | NodeEvent::TerminalFrame { .. }
2159 | NodeEvent::ControllerChanged { .. }
2160 | NodeEvent::WorkspaceRemoved { .. }
2161 | NodeEvent::SessionRecordRemoved { .. }
2162 | NodeEvent::ManagedWorktreeUpserted { .. }
2163 | NodeEvent::ManagedWorktreeRemoved { .. }
2164 | NodeEvent::AgentStream { .. }
2165 | NodeEvent::ResyncRequired { .. } => false,
2166 }
2167}
2168
2169fn node_request_contains_tagged_repository_path(request: &NodeRequest) -> bool {
2170 match request {
2171 NodeRequest::ReadWorkspaceFile { path, .. }
2172 | NodeRequest::WriteWorkspaceFile { path, .. }
2173 | NodeRequest::CreateWorkspaceFile { path, .. }
2174 | NodeRequest::CreateWorkspaceDirectory { path, .. } => path.as_unix_bytes().is_some(),
2175 NodeRequest::ReadGitDiff { request, .. } => request
2176 .path
2177 .as_ref()
2178 .is_some_and(|path| path.as_unix_bytes().is_some()),
2179 NodeRequest::Snapshot
2180 | NodeRequest::Resync { .. }
2181 | NodeRequest::ArmHarnessMcpReservation { .. }
2182 | NodeRequest::SpawnSpecWithHarnessMcp { .. }
2183 | NodeRequest::ActivateHarnessMcpReservation { .. }
2184 | NodeRequest::AbortHarnessMcpReservation { .. }
2185 | NodeRequest::PutHarnessMcpReplyChunk { .. }
2186 | NodeRequest::RejectHarnessMcpCall { .. }
2187 | NodeRequest::BeginDeliveryStage { .. }
2188 | NodeRequest::PutDeliveryBlobChunk { .. }
2189 | NodeRequest::CommitDeliveryStage { .. }
2190 | NodeRequest::AbortDeliveryStage { .. }
2191 | NodeRequest::BrowseHostDirectories { .. }
2192 | NodeRequest::InspectWorkspace { .. }
2193 | NodeRequest::ReadGitHistory { .. }
2194 | NodeRequest::AcquireController { .. }
2195 | NodeRequest::ReleaseController
2196 | NodeRequest::RegisterWorkspace { .. }
2197 | NodeRequest::CreateStandaloneWorkspace { .. }
2198 | NodeRequest::UnregisterWorkspace { .. }
2199 | NodeRequest::CreateWorktree { .. }
2200 | NodeRequest::RemoveWorktree { .. }
2201 | NodeRequest::Spawn { .. }
2202 | NodeRequest::SpawnSpec { .. }
2203 | NodeRequest::SpawnManagedWorktree { .. }
2204 | NodeRequest::SpawnManagedWorktreeV2 { .. }
2205 | NodeRequest::CleanupManagedWorktree { .. }
2206 | NodeRequest::Resume { .. }
2207 | NodeRequest::RenameSessionRecord { .. }
2208 | NodeRequest::SetSessionTask { .. }
2209 | NodeRequest::IndexProviderSession { .. }
2210 | NodeRequest::IndexNativeSession { .. }
2211 | NodeRequest::ResumeSessionRecord { .. }
2212 | NodeRequest::ForgetSessionRecord { .. }
2213 | NodeRequest::CatalogNativeSessions { .. }
2214 | NodeRequest::PageNativeSessions { .. }
2215 | NodeRequest::PreviewNativeSession { .. }
2216 | NodeRequest::PreviewSessionRecord { .. }
2217 | NodeRequest::DiscoverHistory { .. }
2218 | NodeRequest::LoadHistory { .. }
2219 | NodeRequest::ExportContextPackForSessionRecord { .. }
2220 | NodeRequest::ExportContextPack { .. }
2221 | NodeRequest::ForgetContextPack { .. }
2222 | NodeRequest::ResolveDurableContextPack { .. }
2223 | NodeRequest::ReadContextPack { .. }
2224 | NodeRequest::Prompt { .. }
2225 | NodeRequest::Paste { .. }
2226 | NodeRequest::Input { .. }
2227 | NodeRequest::TerminalBytes { .. }
2228 | NodeRequest::TerminalControl { .. }
2229 | NodeRequest::Resize { .. }
2230 | NodeRequest::Interrupt { .. }
2231 | NodeRequest::Stop { .. }
2232 | NodeRequest::Remove { .. }
2233 | NodeRequest::ResolveInteraction { .. }
2234 | NodeRequest::SetSessionMode { .. }
2235 | NodeRequest::SetSessionConfigOption { .. }
2236 | NodeRequest::SetSessionModel { .. }
2237 | NodeRequest::Shutdown => false,
2238 }
2239}
2240
2241fn server_frame_contains_opaque_unix_path(frame: &ServerFrame) -> bool {
2242 match frame {
2243 ServerFrame::Hello(hello) => node_snapshot_contains_opaque_unix_path(&hello.snapshot),
2244 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
2245 node_response_contains_opaque_unix_path,
2246 ),
2247 ServerFrame::Event(event) => node_event_contains_opaque_unix_path(&event.event),
2248 ServerFrame::Challenge(_) => false,
2249 }
2250}
2251
2252fn server_frame_contains_tagged_repository_path(frame: &ServerFrame) -> bool {
2253 match frame {
2254 ServerFrame::Reply(reply) => reply.result.as_ref().ok().is_some_and(
2255 node_response_contains_tagged_repository_path,
2256 ),
2257 ServerFrame::Challenge(_) | ServerFrame::Hello(_) | ServerFrame::Event(_) => false,
2258 }
2259}
2260
2261fn node_response_contains_tagged_repository_path(response: &NodeResponse) -> bool {
2262 match response {
2263 NodeResponse::WorkspaceInspected { inspection } => {
2264 inspection.entries.iter().any(|entry| {
2265 entry.relative_path.as_unix_bytes().is_some()
2266 }) || inspection.git.status.iter().any(|status| {
2267 status.path.as_unix_bytes().is_some()
2268 || status.previous_path.as_ref().is_some_and(|path| {
2269 path.as_unix_bytes().is_some()
2270 })
2271 })
2272 }
2273 NodeResponse::WorkspaceFileRead { file }
2274 | NodeResponse::WorkspaceFileWritten { file }
2275 | NodeResponse::WorkspaceFileCreated { file } => file.path.as_unix_bytes().is_some(),
2276 NodeResponse::WorkspaceDirectoryCreated { entry, .. } => {
2277 entry.relative_path.as_unix_bytes().is_some()
2278 }
2279 NodeResponse::GitDiffRead { diff, .. } => diff
2280 .path
2281 .as_ref()
2282 .is_some_and(|path| path.as_unix_bytes().is_some()),
2283 NodeResponse::Snapshot { .. }
2284 | NodeResponse::Resync { .. }
2285 | NodeResponse::Armed { .. }
2286 | NodeResponse::Spawned { .. }
2287 | NodeResponse::Activated { .. }
2288 | NodeResponse::Aborted { .. }
2289 | NodeResponse::ReplyChunkAccepted { .. }
2290 | NodeResponse::CallRejected { .. }
2291 | NodeResponse::DeliveryStageBegun { .. }
2292 | NodeResponse::DeliveryBlobChunkAccepted { .. }
2293 | NodeResponse::DeliveryCommitted { .. }
2294 | NodeResponse::DeliveryStageAborted { .. }
2295 | NodeResponse::HostDirectoriesBrowsed { .. }
2296 | NodeResponse::GitHistoryRead { .. }
2297 | NodeResponse::Controller { .. }
2298 | NodeResponse::SpawnAccepted { .. }
2299 | NodeResponse::SpawnSpecAccepted { .. }
2300 | NodeResponse::ManagedWorktreeSpawnAccepted { .. }
2301 | NodeResponse::ManagedWorktreeCleanup { .. }
2302 | NodeResponse::SessionRecordUpdated { .. }
2303 | NodeResponse::ProviderSessionIndexed { .. }
2304 | NodeResponse::NativeSessionIndexed { .. }
2305 | NodeResponse::SessionRecordResumed { .. }
2306 | NodeResponse::SessionRecordForgotten { .. }
2307 | NodeResponse::NativeSessionsCataloged { .. }
2308 | NodeResponse::NativeSessionsPaged { .. }
2309 | NodeResponse::NativeSessionPreviewed { .. }
2310 | NodeResponse::SessionRecordPreviewed { .. }
2311 | NodeResponse::HistoryDiscovered { .. }
2312 | NodeResponse::HistoryLoaded { .. }
2313 | NodeResponse::ContextPackForSessionRecordExported { .. }
2314 | NodeResponse::ContextPackExported { .. }
2315 | NodeResponse::ContextPackForgotten { .. }
2316 | NodeResponse::DurableContextPackResolved { .. }
2317 | NodeResponse::ContextPackBytesRead { .. }
2318 | NodeResponse::WorkspaceRegistered { .. }
2319 | NodeResponse::StandaloneWorkspaceCreated { .. }
2320 | NodeResponse::WorkspaceUnregistered { .. }
2321 | NodeResponse::WorktreeCreated { .. }
2322 | NodeResponse::WorktreeRemoved { .. }
2323 | NodeResponse::Accepted
2324 | NodeResponse::ShuttingDown => false,
2325 }
2326}
2327
2328fn node_response_contains_opaque_unix_path(response: &NodeResponse) -> bool {
2329 match response {
2330 NodeResponse::Snapshot { snapshot, .. } => {
2331 node_snapshot_contains_opaque_unix_path(snapshot)
2332 }
2333 NodeResponse::Resync { snapshot, events, .. } => {
2334 node_snapshot_contains_opaque_unix_path(snapshot)
2335 || events.iter().any(|event| {
2336 node_event_contains_opaque_unix_path(&event.event)
2337 })
2338 }
2339 NodeResponse::WorkspaceInspected { inspection } => inspection.git.worktrees
2340 .iter()
2341 .any(|worktree| worktree.path.as_unix_bytes().is_some()),
2342 NodeResponse::HostDirectoriesBrowsed { listing } => {
2343 listing.directory.as_ref().is_some_and(|path| path.as_unix_bytes().is_some())
2344 || listing.parent.as_ref().is_some_and(|path| path.as_unix_bytes().is_some())
2345 || listing.entries.iter().any(|entry| entry.path.as_unix_bytes().is_some())
2346 || listing.next_after.as_ref().is_some_and(|path| path.as_unix_bytes().is_some())
2347 }
2348 NodeResponse::WorkspaceFileRead { .. }
2349 | NodeResponse::WorkspaceFileWritten { .. }
2350 | NodeResponse::WorkspaceFileCreated { .. }
2351 | NodeResponse::WorkspaceDirectoryCreated { .. }
2352 | NodeResponse::GitHistoryRead { .. }
2353 | NodeResponse::GitDiffRead { .. } => false,
2354 NodeResponse::SessionRecordUpdated { record }
2355 | NodeResponse::ProviderSessionIndexed { record }
2356 | NodeResponse::NativeSessionIndexed { record, .. }
2357 | NodeResponse::SessionRecordResumed { record, .. } => {
2358 record.canonical_root.as_unix_bytes().is_some()
2359 }
2360 NodeResponse::WorkspaceRegistered { workspace }
2361 | NodeResponse::StandaloneWorkspaceCreated { workspace } => {
2362 workspace.canonical_root.as_unix_bytes().is_some()
2363 }
2364 NodeResponse::WorktreeCreated { worktree, workspace } => {
2365 worktree.path.as_unix_bytes().is_some()
2366 || workspace.canonical_root.as_unix_bytes().is_some()
2367 }
2368 NodeResponse::WorktreeRemoved { target_root, .. } => {
2369 target_root.as_unix_bytes().is_some()
2370 }
2371 NodeResponse::Armed { .. }
2372 | NodeResponse::Spawned { .. }
2373 | NodeResponse::Activated { .. }
2374 | NodeResponse::Aborted { .. }
2375 | NodeResponse::ReplyChunkAccepted { .. }
2376 | NodeResponse::CallRejected { .. }
2377 | NodeResponse::Controller { .. }
2378 | NodeResponse::DeliveryStageBegun { .. }
2379 | NodeResponse::DeliveryBlobChunkAccepted { .. }
2380 | NodeResponse::DeliveryCommitted { .. }
2381 | NodeResponse::DeliveryStageAborted { .. }
2382 | NodeResponse::SpawnAccepted { .. }
2383 | NodeResponse::SpawnSpecAccepted { .. }
2384 | NodeResponse::ManagedWorktreeSpawnAccepted { .. }
2385 | NodeResponse::ManagedWorktreeCleanup { .. }
2386 | NodeResponse::SessionRecordForgotten { .. }
2387 | NodeResponse::NativeSessionsCataloged { .. }
2388 | NodeResponse::NativeSessionsPaged { .. }
2389 | NodeResponse::NativeSessionPreviewed { .. }
2390 | NodeResponse::SessionRecordPreviewed { .. }
2391 | NodeResponse::HistoryDiscovered { .. }
2392 | NodeResponse::HistoryLoaded { .. }
2393 | NodeResponse::ContextPackForSessionRecordExported { .. }
2394 | NodeResponse::ContextPackExported { .. }
2395 | NodeResponse::ContextPackForgotten { .. }
2396 | NodeResponse::DurableContextPackResolved { .. }
2397 | NodeResponse::ContextPackBytesRead { .. }
2398 | NodeResponse::WorkspaceUnregistered { .. }
2399 | NodeResponse::Accepted
2400 | NodeResponse::ShuttingDown => false,
2401 }
2402}
2403
2404fn node_snapshot_contains_opaque_unix_path(snapshot: &NodeSnapshot) -> bool {
2405 snapshot.workspaces.iter().any(|workspace| {
2406 workspace.canonical_root.as_unix_bytes().is_some()
2407 }) || snapshot.session_records.iter().any(|record| {
2408 record.canonical_root.as_unix_bytes().is_some()
2409 })
2410}
2411
2412fn node_event_contains_opaque_unix_path(event: &NodeEvent) -> bool {
2413 match event {
2414 NodeEvent::WorkspaceAdded { workspace } => {
2415 workspace.canonical_root.as_unix_bytes().is_some()
2416 }
2417 NodeEvent::SessionRecordUpserted { record } => {
2418 record.canonical_root.as_unix_bytes().is_some()
2419 }
2420 NodeEvent::HarnessMcpReadCall { .. }
2421 | NodeEvent::Control { .. }
2422 | NodeEvent::SessionRecordHistorySummarized { .. }
2423 | NodeEvent::TerminalFrame { .. }
2424 | NodeEvent::ControllerChanged { .. }
2425 | NodeEvent::WorkspaceRemoved { .. }
2426 | NodeEvent::SessionRecordRemoved { .. }
2427 | NodeEvent::ManagedWorktreeUpserted { .. }
2428 | NodeEvent::ManagedWorktreeRemoved { .. }
2429 | NodeEvent::AgentStream { .. }
2430 | NodeEvent::ResyncRequired { .. } => false,
2431 }
2432}
2433
2434fn client_compatibility_offer() -> Result<ClientCompatibilityOffer, NodeClientError> {
2435 let mut offer = production_node_client_compatibility_offer();
2436 let history_context_pack =
2437 CapabilityId::new(NODE_HISTORY_CONTEXT_PACK_CAPABILITY).map_err(|error| {
2438 NodeClientError::Protocol(error.to_string())
2439 })?;
2440 if !offer.capabilities.contains(&history_context_pack) {
2441 offer.capabilities.push(history_context_pack);
2442 }
2443 let session_record_context_export =
2444 CapabilityId::new(NODE_SESSION_RECORD_CONTEXT_EXPORT_CAPABILITY).map_err(|error| {
2445 NodeClientError::Protocol(error.to_string())
2446 })?;
2447 if !offer.capabilities.contains(&session_record_context_export) {
2448 offer.capabilities.push(session_record_context_export);
2449 }
2450 Ok(offer)
2451}
2452
2453fn prepare_negotiated_authentication(
2454 challenge: &ServerChallenge,
2455 offer: &ClientCompatibilityOffer,
2456 role: ClientRole,
2457 client_nonce: &[u8; NODE_AUTH_NONCE_BYTES],
2458 access_token: &str,
2459) -> Result<ClientAuthentication, NodeClientError> {
2460 let selected = challenge.compatibility.as_ref().ok_or_else(|| {
2461 NodeClientError::Protocol(
2462 "node omitted the required authenticated compatibility selection".to_owned(),
2463 )
2464 })?;
2465 validate_selected_compatibility(offer, selected)?;
2466 let expected_server_proof = negotiated_auth_proof(
2467 access_token.as_bytes(),
2468 AuthDirection::Server,
2469 role,
2470 client_nonce,
2471 &challenge.server_nonce,
2472 offer,
2473 selected,
2474 )
2475 .map_err(NodeClientError::Authentication)?;
2476 if !proofs_match(&challenge.server_proof, &expected_server_proof) {
2477 return Err(NodeClientError::Protocol(
2478 "server failed access-token proof".to_owned(),
2479 ));
2480 }
2481 let client_proof = negotiated_auth_proof(
2482 access_token.as_bytes(),
2483 AuthDirection::Client,
2484 role,
2485 client_nonce,
2486 &challenge.server_nonce,
2487 offer,
2488 selected,
2489 )
2490 .map_err(NodeClientError::Authentication)?;
2491 Ok(ClientAuthentication { client_proof })
2492}
2493
2494fn validate_authenticated_hello_compatibility(
2495 offer: &ClientCompatibilityOffer,
2496 challenge: &NegotiatedNodeCompatibility,
2497 hello: Option<&NegotiatedNodeCompatibility>,
2498) -> Result<(), NodeClientError> {
2499 let hello = hello.ok_or_else(|| {
2500 NodeClientError::Protocol(
2501 "node hello omitted the required authenticated compatibility selection".to_owned(),
2502 )
2503 })?;
2504 if hello != challenge {
2505 return Err(NodeClientError::Protocol(
2506 "node compatibility selection changed during authentication".to_owned(),
2507 ));
2508 }
2509 validate_selected_compatibility(offer, hello)
2510}
2511
2512fn validate_selected_compatibility(
2513 offer: &ClientCompatibilityOffer,
2514 selected: &NegotiatedNodeCompatibility,
2515) -> Result<(), NodeClientError> {
2516 if selected.build_stamp != BUILD_STAMP {
2517 return Err(NodeClientError::Protocol(format!(
2518 "node selected build stamp {} for active wire build {}",
2519 selected.build_stamp,
2520 BUILD_STAMP,
2521 )));
2522 }
2523 if selected
2524 .capabilities
2525 .iter()
2526 .any(|capability| !offer.capabilities.contains(capability))
2527 {
2528 return Err(NodeClientError::Protocol(
2529 "node selected a capability outside the client offer".to_owned(),
2530 ));
2531 }
2532 if !selected.capabilities.iter().any(|capability| {
2533 capability.as_str() == NODE_COMPATIBILITY_METADATA_CAPABILITY
2534 }) {
2535 return Err(NodeClientError::Protocol(
2536 "node omitted the required compatibility metadata capability".to_owned(),
2537 ));
2538 }
2539 let open_provider_ids_selected = selected.capabilities.iter().any(|capability| {
2540 capability.as_str() == NODE_PROVIDER_ID_OPEN_CAPABILITY
2541 });
2542 if !open_provider_ids_selected
2543 && (selected.provider_contracts.iter().any(|contract| {
2544 !provider_id_is_legacy(&contract.provider)
2545 }) || selected.provider_adapter_contracts.iter().any(|contract| {
2546 !provider_id_is_legacy(&contract.provider)
2547 }))
2548 {
2549 return Err(NodeClientError::Protocol(
2550 "node published an open provider ID without negotiating the capability".to_owned(),
2551 ));
2552 }
2553 let provider_manifest_selected = selected.capabilities.iter().any(|capability| {
2554 capability.as_str() == NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY
2555 });
2556 if !provider_manifest_selected
2557 && (!selected.provider_contracts.is_empty()
2558 || !selected.provider_adapter_contracts.is_empty())
2559 {
2560 return Err(NodeClientError::Protocol(
2561 "node published a provider contract manifest without negotiating the capability"
2562 .to_owned(),
2563 ));
2564 }
2565 if provider_manifest_selected {
2566 validate_provider_contract_manifest(
2567 &selected.provider_contracts,
2568 &selected.provider_adapter_contracts,
2569 )
2570 .map_err(|error| NodeClientError::Protocol(error.to_string()))?;
2571 }
2572 if let Some(state_schema_version) = selected.state_schema_version {
2573 let Some(state_schema) = offer.state_schema else {
2574 return Err(NodeClientError::Protocol(
2575 "node selected a state schema that the client did not offer".to_owned(),
2576 ));
2577 };
2578 if !state_schema.versions.contains(state_schema_version) {
2579 return Err(NodeClientError::Protocol(format!(
2580 "node selected state schema version {state_schema_version} outside the client offer",
2581 )));
2582 }
2583 }
2584 Ok(())
2585}
2586
2587#[derive(Debug, Error)]
2588pub enum NodeClientError {
2589 #[error("local IPC I/O failed: {0}")]
2590 Io(#[from] io::Error),
2591 #[error(transparent)]
2592 Frame(#[from] FrameError),
2593 #[error("node rejected request: {0:?}")]
2594 Node(NodeFailure),
2595 #[error("node protocol failed: {0}")]
2596 Protocol(String),
2597 #[error("build stamp mismatch: local={local} remote={remote}")]
2598 BuildStampMismatch { local: String, remote: String },
2599 #[error("node capability was not negotiated: {0}")]
2600 UnsupportedCapability(String),
2601 #[error("node authentication frame was not received before the bounded deadline")]
2602 AuthenticationTimedOut,
2603 #[error("node authentication primitive failed: {0}")]
2604 Authentication(String),
2605 #[error("request id counter is exhausted")]
2606 RequestIdExhausted,
2607}
2608
2609#[cfg(test)]
2610mod tests {
2611 use super::*;
2612 use gate4agent_node_protocol::{
2613 AdapterContractRevision, AdapterFamily, AdapterId, AgentId, AgentProgressCurrentV1,
2614 AgentProgressV1, ArchitectureId,
2615 ContextPackLineageReceipt,
2616 GitSnapshot, GitStatusEntry, GitWorktreeSnapshot, HostDescriptor, LocalTransportKind,
2617 ManagedSessionRecord, ManagedSessionState, NodeCompatibilitySupport, OpaqueHostPath,
2618 ManagedWorktreeCleanupFailure, ManagedWorktreeLeaseId,
2619 ManagedWorktreeLeaseSnapshot, ManagedWorktreeLeaseState, ManagedWorktreeRetention,
2620 ManagedWorktreeSpawnReceipt, ManagedWorktreeSpawnRequest,
2621 HostDirectoryEntry, HostDirectoryListing, OperatingSystemId, PathEncoding,
2622 PathSemantics, PathStyle,
2623 ProviderAdapterContractSupport, ProviderContractRevision, ProviderContractSupport,
2624 ProviderRuntimeStatus, ProviderRuntimeStatuses, RepositoryPath, ResponseEnvelope,
2625 ResolvedBundleReceipt, ResolvedContextPackReceipt, ResolvedEnvironmentProfileReceipt,
2626 ResolvedSpawnReceipt,
2627 SessionAddress, SessionAgentProgress, SessionKey,
2628 SessionMode, SessionRecordId,
2629 SpawnBundleDigest, SpawnBundleId, SpawnBundleRevision, SpawnDeadlineMs,
2630 SpawnContextDigest, SpawnContextId,
2631 SpawnFieldProvenance, SpawnIdempotencyKey, SpawnOverrides,
2632 SpawnEnvironmentProfileId, SpawnEnvironmentProfileRevision, SpawnProfileId,
2633 SpawnProfileRevision, SpawnPromptMetadata, SpawnRequiredCapabilities,
2634 SpawnResolutionProvenance, SpawnSpec, SpawnTarget, WorkspaceEntry,
2635 WorkspaceEntryKind, WorkspaceFileContent, WorkspaceFileRead, WorkspaceFileRevision,
2636 WorkspaceId,
2637 WorkspaceInspection, WorkspaceSnapshot,
2638 WorktreeProfileId, WorktreeProfileRevision, WorktreeServiceMode,
2639 };
2640
2641 #[test]
2642 fn harness_mcp_wire_capability_and_exact_response_correlation() {
2643 let reservation_id = gate4agent_node_protocol::HarnessMcpReservationId::new(
2644 format!("hmcpres_{}", "a".repeat(24)),
2645 ).unwrap();
2646 let activation_digest = gate4agent_node_protocol::HarnessMcpActivationDigest::new(
2647 format!("sha256:{}", "b".repeat(64)),
2648 ).unwrap();
2649 let request = NodeRequest::AbortHarnessMcpReservation {
2650 reservation_id: reservation_id.clone(),
2651 activation_digest: activation_digest.clone(),
2652 };
2653 assert!(matches!(
2654 ensure_node_request_required_capability(&request, &[]),
2655 Err(NodeClientError::UnsupportedCapability(capability))
2656 if capability == NODE_HARNESS_MCP_READ_PROXY_CAPABILITY
2657 ));
2658 let capability = CapabilityId::new(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY).unwrap();
2659 assert!(ensure_node_request_required_capability(
2660 &request,
2661 std::slice::from_ref(&capability),
2662 ).is_ok());
2663 let exact = Ok(NodeResponse::Aborted {
2664 reservation_id: reservation_id.clone(),
2665 activation_digest: activation_digest.clone(),
2666 });
2667 assert!(validate_harness_mcp_response(&request, &exact).is_ok());
2668 let mismatch = Ok(NodeResponse::Aborted {
2669 reservation_id: gate4agent_node_protocol::HarnessMcpReservationId::new(
2670 format!("hmcpres_{}", "c".repeat(24)),
2671 ).unwrap(),
2672 activation_digest,
2673 });
2674 assert!(validate_harness_mcp_response(&request, &mismatch).is_err());
2675 let frame = ServerFrame::Reply(ResponseEnvelope {
2676 request_id: 1,
2677 result: exact,
2678 });
2679 assert!(ensure_server_frame_required_capability(&frame, &[]).is_err());
2680 assert!(ensure_server_frame_required_capability(&frame, &[capability]).is_ok());
2681 }
2682 use gate4agent_types::{
2683 AgentInstanceId, CapabilitySnapshot, ForegroundSnapshot, HistorySnapshot,
2684 ProviderSnapshot, PtyScreenState, ResumeSnapshot, SessionGeneration, SessionSnapshot,
2685 SessionStatus, TerminalFrame, TerminalSize, TransportKind,
2686 };
2687
2688 fn agent(value: &str) -> AgentId {
2689 AgentId::new(value).unwrap()
2690 }
2691
2692 fn negotiated_fixture() -> (ClientCompatibilityOffer, NegotiatedNodeCompatibility) {
2693 let offer = ClientCompatibilityOffer {
2694 build_stamp: BUILD_STAMP.to_owned(),
2695 capabilities: vec![CapabilityId::new(NODE_COMPATIBILITY_METADATA_CAPABILITY).unwrap()],
2696 state_schema: Some(StateSchemaSupport {
2697 versions: ProtocolRange::exact(1).unwrap(),
2698 }),
2699 };
2700 let support = NodeCompatibilitySupport {
2701 build_stamp: BUILD_STAMP.to_owned(),
2702 capabilities: vec![CapabilityId::new(NODE_COMPATIBILITY_METADATA_CAPABILITY).unwrap()],
2703 host: HostDescriptor {
2704 operating_system: OperatingSystemId::new("windows").unwrap(),
2705 architecture: ArchitectureId::new("x86_64").unwrap(),
2706 },
2707 path_semantics: PathSemantics {
2708 style: PathStyle::Windows,
2709 encoding: PathEncoding::Utf8,
2710 },
2711 local_transport: LocalTransportKind::WindowsNamedPipe,
2712 state_schema: StateSchemaSupport {
2713 versions: ProtocolRange::exact(1).unwrap(),
2714 },
2715 provider_contracts: Vec::new(),
2716 provider_adapter_contracts: Vec::new(),
2717 };
2718 let selected = support.negotiate(&offer).unwrap();
2719 (offer, selected)
2720 }
2721
2722 fn unix_path() -> OpaqueHostPath {
2723 OpaqueHostPath::unix_bytes(vec![b'/', b's', b'r', b'v', b'/', 0xff]).unwrap()
2724 }
2725
2726 fn utf8_path() -> OpaqueHostPath {
2727 OpaqueHostPath::utf8(r"C:\repo".to_owned()).unwrap()
2728 }
2729
2730 fn tagged_repository_path(value: &[u8]) -> RepositoryPath {
2731 RepositoryPath::unix_bytes(value.to_vec()).unwrap()
2732 }
2733
2734 fn utf8_repository_path(value: &str) -> RepositoryPath {
2735 RepositoryPath::utf8(value.to_owned()).unwrap()
2736 }
2737
2738 fn workspace_with_path(canonical_root: OpaqueHostPath) -> WorkspaceSnapshot {
2739 WorkspaceSnapshot {
2740 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
2741 canonical_root,
2742 sessions: Vec::new(),
2743 worktree_service_mode: None,
2744 managed_worktree_profiles: None,
2745 }
2746 }
2747
2748 fn session_snapshot(provider: &str) -> SessionSnapshot {
2749 SessionSnapshot {
2750 instance_id: AgentInstanceId(7),
2751 agent_id: agent(provider),
2752 transport: TransportKind::Pty,
2753 generation: SessionGeneration(2),
2754 status: SessionStatus::Running,
2755 pending_operation: None,
2756 pending_input: None,
2757 process_id: None,
2758 terminal_size: None,
2759 terminal_frame: None,
2760 terminal_stale: None,
2761 session_options: None,
2762 capabilities: CapabilitySnapshot::default(),
2763 history: HistorySnapshot::default(),
2764 resume: ResumeSnapshot::default(),
2765 foreground: ForegroundSnapshot::default(),
2766 provider: ProviderSnapshot::default(),
2767 screen_state: PtyScreenState::default(),
2768 }
2769 }
2770
2771 fn session_record_with_path(canonical_root: OpaqueHostPath) -> ManagedSessionRecord {
2772 ManagedSessionRecord {
2773 record_id: SessionRecordId::new("session-a").unwrap(),
2774 display_name: "session a".to_owned(),
2775 provider: agent("claude"),
2776 mode: SessionMode::Pty,
2777 state: ManagedSessionState::Dormant,
2778 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
2779 canonical_root,
2780 provider_session: None,
2781 active_session: None,
2782 environment_profile: None,
2783 bundle: None,
2784 context_id: None,
2785 context: None,
2786 exported_context: None,
2787 task_binding: None,
2788 created_at_unix_ms: 1,
2789 updated_at_unix_ms: 2,
2790 last_error: None,
2791 }
2792 }
2793
2794 #[test]
2795 fn session_task_capability_and_idempotent_response_correlation_are_fail_closed() {
2796 let task_id = gate4agent_node_protocol::TaskId::from_nonce([7; 12]);
2797 let request = NodeRequest::SetSessionTask {
2798 record_id: SessionRecordId::new("session-a").unwrap(),
2799 expected_revision: 7,
2800 target: gate4agent_node_protocol::SessionTaskTargetV1::Existing {
2801 task_id: task_id.clone(),
2802 },
2803 };
2804 assert!(ensure_node_request_required_capability(&request, &[]).is_err());
2805 let capability = CapabilityId::new(NODE_SESSION_TASK_CORRELATION_CAPABILITY).unwrap();
2806 assert!(ensure_node_request_required_capability(
2807 &request,
2808 std::slice::from_ref(&capability),
2809 ).is_ok());
2810
2811 let mut record = session_record_with_path(utf8_path());
2812 record.task_binding = Some(gate4agent_node_protocol::SessionTaskBindingV1 {
2813 revision: 7,
2814 task_id: Some(task_id.clone()),
2815 changed_at_unix_ms: 2,
2816 });
2817 assert!(validate_session_task_response(
2818 &request,
2819 &Ok(NodeResponse::SessionRecordUpdated { record: record.clone() }),
2820 ).is_ok());
2821 record.task_binding.as_mut().unwrap().revision = 8;
2822 assert!(validate_session_task_response(
2823 &request,
2824 &Ok(NodeResponse::SessionRecordUpdated { record: record.clone() }),
2825 ).is_ok());
2826 record.task_binding.as_mut().unwrap().revision = 9;
2827 assert!(validate_session_task_response(
2828 &request,
2829 &Ok(NodeResponse::SessionRecordUpdated { record: record.clone() }),
2830 ).is_err());
2831
2832 let clear = NodeRequest::SetSessionTask {
2833 record_id: SessionRecordId::new("session-a").unwrap(),
2834 expected_revision: 0,
2835 target: gate4agent_node_protocol::SessionTaskTargetV1::Clear,
2836 };
2837 record.task_binding = None;
2838 assert!(validate_session_task_response(
2839 &clear,
2840 &Ok(NodeResponse::SessionRecordUpdated { record: record.clone() }),
2841 ).is_ok());
2842
2843 let new = NodeRequest::SetSessionTask {
2844 record_id: SessionRecordId::new("session-a").unwrap(),
2845 expected_revision: 7,
2846 target: gate4agent_node_protocol::SessionTaskTargetV1::New,
2847 };
2848 record.task_binding = Some(gate4agent_node_protocol::SessionTaskBindingV1 {
2849 revision: 7,
2850 task_id: Some(task_id),
2851 changed_at_unix_ms: 2,
2852 });
2853 assert!(validate_session_task_response(
2854 &new,
2855 &Ok(NodeResponse::SessionRecordUpdated { record: record.clone() }),
2856 ).is_err());
2857 record.task_binding.as_mut().unwrap().revision = 8;
2858 let frame = response_frame(NodeResponse::SessionRecordUpdated { record: record.clone() });
2859 assert!(ensure_server_frame_required_capability(&frame, &[]).is_err());
2860 assert!(ensure_server_frame_required_capability(&frame, &[capability]).is_ok());
2861 assert!(validate_session_task_response(
2862 &new,
2863 &Ok(NodeResponse::SessionRecordUpdated { record }),
2864 ).is_ok());
2865 }
2866
2867 fn worktree_with_path(path: OpaqueHostPath) -> GitWorktreeSnapshot {
2868 GitWorktreeSnapshot {
2869 path,
2870 head: "abcdef".to_owned(),
2871 branch: Some("main".to_owned()),
2872 is_bare: false,
2873 is_main: true,
2874 locked: false,
2875 lock_reason: None,
2876 prunable: false,
2877 prunable_reason: None,
2878 workspace_id: Some(WorkspaceId::new("workspace-a").unwrap()),
2879 }
2880 }
2881
2882 fn empty_snapshot() -> NodeSnapshot {
2883 NodeSnapshot {
2884 node_id: NodeId::new("node-a").unwrap(),
2885 enabled_providers: Vec::new(),
2886 provider_runtime_statuses: ProviderRuntimeStatuses::default(),
2887 workspaces: Vec::new(),
2888 session_records: Vec::new(),
2889 managed_worktrees: Vec::new(),
2890 launch_inventory: None,
2891 agent_progress: Vec::new(),
2892 }
2893 }
2894
2895 fn hello_with_snapshot(snapshot: NodeSnapshot) -> NodeHello {
2896 NodeHello {
2897 build_stamp: BUILD_STAMP.to_owned(),
2898 incarnation_id: NodeIncarnationId::from_bytes([3; NODE_INCARNATION_ID_BYTES]),
2899 connection_id: 7,
2900 role: ClientRole::Operator,
2901 event_sequence: 0,
2902 controller: None,
2903 snapshot,
2904 compatibility: None,
2905 }
2906 }
2907
2908 #[test]
2909 fn agent_progress_requires_negotiated_capability_on_hello_snapshot_and_reply() {
2910 let mut snapshot = empty_snapshot();
2911 snapshot.agent_progress.push(SessionAgentProgress {
2912 address: SessionAddress {
2913 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
2914 session: SessionKey {
2915 instance_id: AgentInstanceId(7),
2916 generation: SessionGeneration(2),
2917 },
2918 },
2919 progress: AgentProgressV1 {
2920 provider_sequence: 3,
2921 activity: gate4agent_types::ProviderActivity::Idle,
2922 completed_turns: 0,
2923 usage: None,
2924 current: AgentProgressCurrentV1::Idle,
2925 active_tool_labels: Vec::new(),
2926 active_tool_count: 0,
2927 attention: None,
2928 subagent_count: 0,
2929 last_event_kind: None,
2930 gap_count: 0,
2931 stale: false,
2932 truncated: false,
2933 },
2934 });
2935 let capabilities = Vec::new();
2936 assert!(matches!(
2937 ensure_node_snapshot_agent_progress_capability(&snapshot, &capabilities),
2938 Err(NodeClientError::Protocol(message))
2939 if message.contains("without negotiating the capability")
2940 ));
2941 let frame = response_frame(NodeResponse::Snapshot {
2942 event_sequence: 0,
2943 controller: None,
2944 snapshot: snapshot.clone(),
2945 });
2946 assert!(matches!(
2947 ensure_server_frame_required_capability(&frame, &capabilities),
2948 Err(NodeClientError::Protocol(message))
2949 if message.contains("without negotiating the capability")
2950 ));
2951 let negotiated = vec![
2952 CapabilityId::new(NODE_AGENT_PROGRESS_SNAPSHOT_CAPABILITY).unwrap(),
2953 ];
2954 assert!(ensure_node_snapshot_agent_progress_capability(
2955 &snapshot,
2956 &negotiated,
2957 )
2958 .is_ok());
2959 assert!(ensure_server_frame_required_capability(&frame, &negotiated).is_ok());
2960 }
2961
2962 fn response_frame(response: NodeResponse) -> ServerFrame {
2963 ServerFrame::Reply(ResponseEnvelope {
2964 request_id: 1,
2965 result: Ok(response),
2966 })
2967 }
2968
2969 const SPAWN_GATE_CAPABILITIES: [&str; 6] = [
2986 NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY,
2987 NODE_SPAWN_PROFILE_REVISION_CAPABILITY,
2988 NODE_WORKTREE_SELECTION_CAPABILITY,
2989 NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY,
2990 NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY,
2991 NODE_HISTORY_CONTEXT_PACK_CAPABILITY,
2992 ];
2993
2994 fn spawn_gate_capabilities_without(withheld: &str) -> Vec<CapabilityId> {
2997 assert!(
2998 SPAWN_GATE_CAPABILITIES.contains(&withheld),
2999 "{withheld} is not one of the capabilities this gate demands",
3000 );
3001 SPAWN_GATE_CAPABILITIES
3002 .into_iter()
3003 .filter(|capability| *capability != withheld)
3004 .map(|capability| CapabilityId::new(capability).unwrap())
3005 .collect()
3006 }
3007
3008 fn spawn_spec_request() -> NodeRequest {
3009 let mut overrides = SpawnOverrides::default();
3010 overrides.context_id = gate4agent_node_protocol::SpawnOverride::Clear;
3011 NodeRequest::SpawnSpec {
3012 spec: SpawnSpec {
3013 target: SpawnTarget {
3014 node_id: NodeId::new("node-a").unwrap(),
3015 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3016 worktree_id: None,
3017 },
3018 profile_id: SpawnProfileId::new("default").unwrap(),
3019 expected_profile_revision:
3020 SpawnProfileRevision::new("default.r1").unwrap(),
3021 overrides,
3022 deadline_ms: SpawnDeadlineMs::new(5_000).unwrap(),
3023 idempotency_key: SpawnIdempotencyKey::new("request-a").unwrap(),
3024 required_capabilities: SpawnRequiredCapabilities::default(),
3025 },
3026 }
3027 }
3028
3029 fn spawn_spec_receipt() -> ResolvedSpawnReceipt {
3030 ResolvedSpawnReceipt {
3031 incarnation_id: NodeIncarnationId::from_bytes([3; NODE_INCARNATION_ID_BYTES]),
3032 session: SessionAddress {
3033 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3034 session: SessionKey {
3035 instance_id: AgentInstanceId(9),
3036 generation: SessionGeneration(1),
3037 },
3038 },
3039 target: SpawnTarget {
3040 node_id: NodeId::new("node-a").unwrap(),
3041 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3042 worktree_id: None,
3043 },
3044 profile_id: SpawnProfileId::new("default").unwrap(),
3045 profile_revision: SpawnProfileRevision::new("default.r1").unwrap(),
3046 provider: agent("claude"),
3047 mode: SessionMode::Pty,
3048 terminal_size: TerminalSize {
3049 rows: 24,
3050 columns: 80,
3051 },
3052 prompt: SpawnPromptMetadata {
3053 present: false,
3054 byte_len: 0,
3055 },
3056 bundle_id: None,
3057 bundle: None,
3058 context_id: None,
3059 context: None,
3060 environment_profile: None,
3061 deadline_ms: SpawnDeadlineMs::new(5_000).unwrap(),
3062 idempotency_key: SpawnIdempotencyKey::new("request-a").unwrap(),
3063 required_capabilities: SpawnRequiredCapabilities::default(),
3064 provenance: SpawnResolutionProvenance {
3065 provider: SpawnFieldProvenance::Profile,
3066 mode: SpawnFieldProvenance::Profile,
3067 terminal_size: SpawnFieldProvenance::Profile,
3068 prompt: SpawnFieldProvenance::Profile,
3069 bundle_id: SpawnFieldProvenance::Profile,
3070 context_id: SpawnFieldProvenance::Profile,
3071 environment_profile_id: SpawnFieldProvenance::Profile,
3072 },
3073 harness_mcp_proxy: None,
3074 }
3075 }
3076
3077 fn context_pack_receipt() -> ResolvedContextPackReceipt {
3078 ResolvedContextPackReceipt {
3079 id: SpawnContextId::new("context-a").unwrap(),
3080 digest: SpawnContextDigest::new(format!("sha256:{}", "c".repeat(64)))
3081 .unwrap(),
3082 lineage: ContextPackLineageReceipt {
3083 source_node_id: NodeId::new("node-a").unwrap(),
3084 source_session: SessionAddress {
3085 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3086 session: SessionKey {
3087 instance_id: AgentInstanceId(7),
3088 generation: SessionGeneration(2),
3089 },
3090 },
3091 source_provider: agent("claude"),
3092 },
3093 source_message_count: 4,
3094 retained_message_count: 3,
3095 byte_len: 128,
3096 truncated: true,
3097 }
3098 }
3099
3100 #[test]
3101 fn native_session_response_correlation_is_fail_closed() {
3102 let route = gate4agent_node_protocol::NativeSessionCatalogRoute::workspace(
3103 WorkspaceId::new("workspace-a").unwrap(),
3104 agent("codex"),
3105 );
3106 let selection = gate4agent_node_protocol::NativeSessionSelection {
3107 route: route.clone(),
3108 catalog_revision: 7,
3109 recent_cutoff_unix_ms: 70,
3110 selection_id: "selection-7".to_owned(),
3111 };
3112
3113 let catalog = NodeRequest::CatalogNativeSessions {
3114 route: route.clone(),
3115 limit: 10,
3116 };
3117 assert!(validate_native_session_response(
3118 &catalog,
3119 &Ok(NodeResponse::NativeSessionsCataloged {
3120 route: route.clone(),
3121 entries: Vec::new(),
3122 summary: None,
3123 }),
3124 )
3125 .is_ok());
3126 let wrong_route = gate4agent_node_protocol::NativeSessionCatalogRoute::workspace(
3127 WorkspaceId::new("workspace-b").unwrap(),
3128 agent("codex"),
3129 );
3130 assert!(validate_native_session_response(
3131 &catalog,
3132 &Ok(NodeResponse::NativeSessionsCataloged {
3133 route: wrong_route,
3134 entries: Vec::new(),
3135 summary: None,
3136 }),
3137 )
3138 .is_err());
3139
3140 let page = NodeRequest::PageNativeSessions {
3141 route: route.clone(),
3142 window: gate4agent_node_protocol::NativeSessionCatalogWindow::Recent,
3143 catalog_revision: 7,
3144 recent_cutoff_unix_ms: 70,
3145 after_selection_id: None,
3146 limit: 10,
3147 };
3148 assert!(validate_native_session_response(
3149 &page,
3150 &Ok(NodeResponse::NativeSessionsPaged {
3151 route: route.clone(),
3152 page: gate4agent_node_protocol::NativeSessionCatalogPage {
3153 window: gate4agent_node_protocol::NativeSessionCatalogWindow::Recent,
3154 revision: 8,
3155 entries: Vec::new(),
3156 next_after_selection_id: None,
3157 remaining_count: 0,
3158 has_more: false,
3159 },
3160 }),
3161 )
3162 .is_err());
3163
3164 let preview = NodeRequest::PreviewNativeSession {
3165 selection: selection.clone(),
3166 message_limit: 10,
3167 };
3168 let mut wrong_selection = selection.clone();
3169 wrong_selection.catalog_revision = 8;
3170 assert!(validate_native_session_response(
3171 &preview,
3172 &Ok(NodeResponse::NativeSessionPreviewed {
3173 selection: wrong_selection,
3174 preview: gate4agent_node_protocol::SessionRecordPreview {
3175 title: None,
3176 modified_at_unix_ms: None,
3177 model: None,
3178 message_count: 0,
3179 message_count_exact: true,
3180 completed_turn_count: None,
3181 total_tokens: None,
3182 truncated: false,
3183 messages: Vec::new(),
3184 },
3185 }),
3186 )
3187 .is_err());
3188
3189 let index = NodeRequest::IndexNativeSession {
3190 selection: selection.clone(),
3191 display_name: "Indexed".to_owned(),
3192 };
3193 let mut record = session_record_with_path(utf8_path());
3194 record.provider = agent("codex");
3195 assert!(validate_native_session_response(
3196 &index,
3197 &Ok(NodeResponse::NativeSessionIndexed {
3198 selection: selection.clone(),
3199 record: record.clone(),
3200 }),
3201 )
3202 .is_ok());
3203 assert!(validate_native_session_response(
3204 &index,
3205 &Ok(NodeResponse::ProviderSessionIndexed {
3206 record: record.clone(),
3207 }),
3208 )
3209 .is_err());
3210 let mut wrong_selection = selection.clone();
3211 wrong_selection.catalog_revision = 8;
3212 assert!(validate_native_session_response(
3213 &index,
3214 &Ok(NodeResponse::NativeSessionIndexed {
3215 selection: wrong_selection,
3216 record: record.clone(),
3217 }),
3218 )
3219 .is_err());
3220 let mut wrong_record = record.clone();
3221 wrong_record.provider = agent("claude");
3222 assert!(validate_native_session_response(
3223 &index,
3224 &Ok(NodeResponse::NativeSessionIndexed {
3225 selection: selection.clone(),
3226 record: wrong_record,
3227 }),
3228 )
3229 .is_err());
3230
3231 let external_index = NodeRequest::IndexNativeSession {
3232 selection: gate4agent_node_protocol::NativeSessionSelection {
3233 route: gate4agent_node_protocol::NativeSessionCatalogRoute::unregistered(
3234 agent("codex"),
3235 ),
3236 ..selection
3237 },
3238 display_name: "External".to_owned(),
3239 };
3240 assert!(validate_native_session_response(
3241 &external_index,
3242 &Ok(NodeResponse::NativeSessionIndexed {
3243 selection: match &external_index {
3244 NodeRequest::IndexNativeSession { selection, .. } => selection.clone(),
3245 _ => unreachable!(),
3246 },
3247 record,
3248 }),
3249 )
3250 .is_err());
3251 }
3252
3253 #[test]
3254 fn the_client_offer_carries_the_open_provider_capabilities_and_the_whole_schema_range() {
3267 let offer = client_compatibility_offer().unwrap();
3268 assert_eq!(offer.build_stamp, BUILD_STAMP);
3269 assert!(offer.capabilities.contains(
3270 &CapabilityId::new(NODE_OPAQUE_UNIX_PATH_CAPABILITY).unwrap(),
3271 ));
3272 assert!(offer.capabilities.contains(
3273 &CapabilityId::new(NODE_REPOSITORY_PATH_CAPABILITY).unwrap(),
3274 ));
3275 assert!(offer.capabilities.contains(
3276 &CapabilityId::new(NODE_WORKSPACE_FILE_READ_CAPABILITY).unwrap(),
3277 ));
3278 assert!(offer.capabilities.contains(
3279 &CapabilityId::new(NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY).unwrap(),
3280 ));
3281 assert!(offer.capabilities.contains(
3282 &CapabilityId::new(NODE_PROVIDER_ID_OPEN_CAPABILITY).unwrap(),
3283 ));
3284 assert!(offer.capabilities.contains(
3285 &CapabilityId::new(NODE_TERMINAL_FRAME_EVENTS_CAPABILITY).unwrap(),
3286 ));
3287 assert!(offer.capabilities.contains(
3288 &CapabilityId::new(NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY).unwrap(),
3289 ));
3290 assert!(offer.capabilities.contains(
3291 &CapabilityId::new(NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY).unwrap(),
3292 ));
3293 assert!(offer.capabilities.contains(
3294 &CapabilityId::new(NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY).unwrap(),
3295 ));
3296 assert!(offer.capabilities.contains(
3297 &CapabilityId::new(NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY).unwrap(),
3298 ));
3299 assert!(offer.capabilities.contains(
3300 &CapabilityId::new(NODE_HISTORY_CONTEXT_PACK_CAPABILITY).unwrap(),
3301 ));
3302 let versions = offer.state_schema.unwrap().versions;
3303 assert_eq!(
3304 versions.minimum(),
3305 NODE_STATE_SCHEMA_V1,
3306 "dropping the oldest state schema is a compatibility break, not a bump",
3307 );
3308 assert_eq!(
3309 versions.maximum(),
3310 NODE_STATE_SCHEMA_V10,
3311 "the client must offer every schema this crate family can write",
3312 );
3313 }
3314
3315 #[test]
3316 fn hmac_valid_legacy_challenge_is_rejected_before_authenticate() {
3317 let offer = client_compatibility_offer().unwrap();
3318 let client_nonce = [3; NODE_AUTH_NONCE_BYTES];
3319 let server_nonce = [7; NODE_AUTH_NONCE_BYTES];
3320 let access_token = "strict-negotiation-token";
3321 let server_proof = auth_proof(
3322 access_token.as_bytes(),
3323 AuthDirection::Server,
3324 ClientRole::Observer,
3325 &client_nonce,
3326 &server_nonce,
3327 )
3328 .unwrap();
3329 let challenge = ServerChallenge {
3330 build_stamp: BUILD_STAMP.to_owned(),
3331 server_nonce,
3332 server_proof,
3333 compatibility: None,
3334 };
3335
3336 let result = prepare_negotiated_authentication(
3337 &challenge,
3338 &offer,
3339 ClientRole::Observer,
3340 &client_nonce,
3341 access_token,
3342 );
3343 assert!(matches!(
3344 result,
3345 Err(NodeClientError::Protocol(message))
3346 if message.contains("required authenticated compatibility selection")
3347 ));
3348 }
3349
3350 #[test]
3351 fn hmac_valid_selection_without_compatibility_metadata_is_rejected() {
3352 let (offer, mut selected) = negotiated_fixture();
3353 selected.capabilities.clear();
3354 let client_nonce = [3; NODE_AUTH_NONCE_BYTES];
3355 let server_nonce = [7; NODE_AUTH_NONCE_BYTES];
3356 let access_token = "strict-negotiation-token";
3357 let server_proof = negotiated_auth_proof(
3358 access_token.as_bytes(),
3359 AuthDirection::Server,
3360 ClientRole::Observer,
3361 &client_nonce,
3362 &server_nonce,
3363 &offer,
3364 &selected,
3365 )
3366 .unwrap();
3367 let challenge = ServerChallenge {
3368 build_stamp: BUILD_STAMP.to_owned(),
3369 server_nonce,
3370 server_proof,
3371 compatibility: Some(selected),
3372 };
3373
3374 let result = prepare_negotiated_authentication(
3375 &challenge,
3376 &offer,
3377 ClientRole::Observer,
3378 &client_nonce,
3379 access_token,
3380 );
3381 assert!(matches!(
3382 result,
3383 Err(NodeClientError::Protocol(message))
3384 if message.contains("required compatibility metadata capability")
3385 ));
3386 }
3387
3388 #[test]
3389 fn selected_manifest_requires_capability_and_valid_provider_linkage() {
3390 let (offer, mut selected) = negotiated_fixture();
3391 selected.provider_contracts.push(ProviderContractSupport {
3392 provider: agent("codex"),
3393 revision: ProviderContractRevision::new("codex.2026-08").unwrap(),
3394 });
3395 assert!(matches!(
3396 validate_selected_compatibility(&offer, &selected),
3397 Err(NodeClientError::Protocol(message))
3398 if message.contains("without negotiating the capability")
3399 ));
3400
3401 let mut offer = offer;
3402 let manifest_capability =
3403 CapabilityId::new(NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY).unwrap();
3404 offer.capabilities.push(manifest_capability.clone());
3405 selected.capabilities.push(manifest_capability);
3406 selected.provider_adapter_contracts.push(ProviderAdapterContractSupport {
3407 provider: agent("claude"),
3408 family: AdapterFamily::PtySemantic,
3409 adapter_id: AdapterId::new("claude-code").unwrap(),
3410 revision: AdapterContractRevision::new("pty-semantic-v1").unwrap(),
3411 });
3412 assert!(matches!(
3413 validate_selected_compatibility(&offer, &selected),
3414 Err(NodeClientError::Protocol(message))
3415 if message.contains("has no provider contract")
3416 ));
3417 }
3418
3419 #[test]
3420 fn authenticated_manifest_revision_tampering_is_rejected() {
3421 let (mut offer, mut selected) = negotiated_fixture();
3422 let manifest_capability =
3423 CapabilityId::new(NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY).unwrap();
3424 offer.capabilities.push(manifest_capability.clone());
3425 selected.capabilities.push(manifest_capability);
3426 selected.provider_contracts.push(ProviderContractSupport {
3427 provider: agent("codex"),
3428 revision: ProviderContractRevision::new("codex.2026-08").unwrap(),
3429 });
3430 selected.provider_adapter_contracts.push(ProviderAdapterContractSupport {
3431 provider: agent("codex"),
3432 family: AdapterFamily::PtySemantic,
3433 adapter_id: AdapterId::new("codex-cli").unwrap(),
3434 revision: AdapterContractRevision::new("pty-semantic-v1").unwrap(),
3435 });
3436 let client_nonce = [3; NODE_AUTH_NONCE_BYTES];
3437 let server_nonce = [7; NODE_AUTH_NONCE_BYTES];
3438 let access_token = "strict-negotiation-token";
3439 let server_proof = negotiated_auth_proof(
3440 access_token.as_bytes(),
3441 AuthDirection::Server,
3442 ClientRole::Observer,
3443 &client_nonce,
3444 &server_nonce,
3445 &offer,
3446 &selected,
3447 )
3448 .unwrap();
3449 selected.provider_adapter_contracts[0].revision =
3450 AdapterContractRevision::new("pty-semantic-v2").unwrap();
3451 let challenge = ServerChallenge {
3452 build_stamp: BUILD_STAMP.to_owned(),
3453 server_nonce,
3454 server_proof,
3455 compatibility: Some(selected),
3456 };
3457 assert!(matches!(
3458 prepare_negotiated_authentication(
3459 &challenge,
3460 &offer,
3461 ClientRole::Observer,
3462 &client_nonce,
3463 access_token,
3464 ),
3465 Err(NodeClientError::Protocol(message))
3466 if message.contains("server failed access-token proof")
3467 ));
3468 }
3469
3470 #[test]
3471 fn spawn_spec_capability_is_bound_to_authentication_proof() {
3472 let (mut offer, mut selected) = negotiated_fixture();
3473 let capability =
3474 CapabilityId::new(NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY).unwrap();
3475 offer.capabilities.push(capability.clone());
3476 selected.capabilities.push(capability);
3477 let client_nonce = [5; NODE_AUTH_NONCE_BYTES];
3478 let server_nonce = [8; NODE_AUTH_NONCE_BYTES];
3479 let access_token = "spawn-spec-auth-token";
3480 let server_proof = negotiated_auth_proof(
3481 access_token.as_bytes(),
3482 AuthDirection::Server,
3483 ClientRole::Operator,
3484 &client_nonce,
3485 &server_nonce,
3486 &offer,
3487 &selected,
3488 )
3489 .unwrap();
3490 selected.capabilities.retain(|candidate| {
3491 candidate.as_str() != NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY
3492 });
3493 let challenge = ServerChallenge {
3494 build_stamp: BUILD_STAMP.to_owned(),
3495 server_nonce,
3496 server_proof,
3497 compatibility: Some(selected),
3498 };
3499 assert!(matches!(
3500 prepare_negotiated_authentication(
3501 &challenge,
3502 &offer,
3503 ClientRole::Operator,
3504 &client_nonce,
3505 access_token,
3506 ),
3507 Err(NodeClientError::Protocol(message))
3508 if message.contains("server failed access-token proof")
3509 ));
3510 }
3511
3512 #[test]
3513 fn history_context_pack_capability_is_bound_to_authentication_proof() {
3514 let (mut offer, mut selected) = negotiated_fixture();
3515 let capability = CapabilityId::new(NODE_HISTORY_CONTEXT_PACK_CAPABILITY).unwrap();
3516 offer.capabilities.push(capability.clone());
3517 selected.capabilities.push(capability);
3518 let client_nonce = [6; NODE_AUTH_NONCE_BYTES];
3519 let server_nonce = [9; NODE_AUTH_NONCE_BYTES];
3520 let access_token = "history-context-auth-token";
3521 let server_proof = negotiated_auth_proof(
3522 access_token.as_bytes(),
3523 AuthDirection::Server,
3524 ClientRole::Operator,
3525 &client_nonce,
3526 &server_nonce,
3527 &offer,
3528 &selected,
3529 )
3530 .unwrap();
3531 selected.capabilities.retain(|candidate| {
3532 candidate.as_str() != NODE_HISTORY_CONTEXT_PACK_CAPABILITY
3533 });
3534 let challenge = ServerChallenge {
3535 build_stamp: BUILD_STAMP.to_owned(),
3536 server_nonce,
3537 server_proof,
3538 compatibility: Some(selected),
3539 };
3540 assert!(matches!(
3541 prepare_negotiated_authentication(
3542 &challenge,
3543 &offer,
3544 ClientRole::Operator,
3545 &client_nonce,
3546 access_token,
3547 ),
3548 Err(NodeClientError::Protocol(message))
3549 if message.contains("server failed access-token proof")
3550 ));
3551 }
3552
3553 #[test]
3554 fn node_hello_requires_the_same_nonempty_compatibility_selection() {
3555 let (offer, selected) = negotiated_fixture();
3556 assert!(validate_authenticated_hello_compatibility(
3557 &offer,
3558 &selected,
3559 None,
3560 )
3561 .is_err());
3562
3563 let mut changed = selected.clone();
3564 changed.path_semantics.style = PathStyle::Posix;
3565 assert!(validate_authenticated_hello_compatibility(
3566 &offer,
3567 &selected,
3568 Some(&changed),
3569 )
3570 .is_err());
3571 assert!(validate_authenticated_hello_compatibility(
3572 &offer,
3573 &selected,
3574 Some(&selected),
3575 )
3576 .is_ok());
3577 }
3578
3579 #[test]
3580 fn opaque_unix_path_gate_requires_explicit_authenticated_selection() {
3581 let (_, mut selected) = negotiated_fixture();
3582 assert!(!selected_supports_opaque_unix_paths(None));
3583 assert!(!selected_supports_opaque_unix_paths(Some(&selected)));
3584
3585 selected.capabilities.push(
3586 CapabilityId::new(NODE_OPAQUE_UNIX_PATH_CAPABILITY).unwrap(),
3587 );
3588 assert!(selected_supports_opaque_unix_paths(Some(&selected)));
3589 }
3590
3591 #[test]
3592 fn repository_path_gate_requires_explicit_authenticated_selection() {
3593 let (_, mut selected) = negotiated_fixture();
3594 assert!(!selected_supports_repository_paths(None));
3595 assert!(!selected_supports_repository_paths(Some(&selected)));
3596
3597 selected.capabilities.push(
3598 CapabilityId::new(NODE_REPOSITORY_PATH_CAPABILITY).unwrap(),
3599 );
3600 assert!(selected_supports_repository_paths(Some(&selected)));
3601 }
3602
3603 #[test]
3604 fn open_provider_gate_requires_explicit_authenticated_selection() {
3605 let (_, mut selected) = negotiated_fixture();
3606 assert!(!selected_supports_open_provider_ids(None));
3607 assert!(!selected_supports_open_provider_ids(Some(&selected)));
3608
3609 selected.capabilities.push(
3610 CapabilityId::new(NODE_PROVIDER_ID_OPEN_CAPABILITY).unwrap(),
3611 );
3612 assert!(selected_supports_open_provider_ids(Some(&selected)));
3613 }
3614
3615 #[test]
3616 fn inbound_terminal_frame_events_require_authenticated_selection() {
3617 let address = SessionAddress {
3618 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3619 session: SessionKey {
3620 instance_id: AgentInstanceId(7),
3621 generation: SessionGeneration(2),
3622 },
3623 };
3624 let event = ServerFrame::Event(NodeEventEnvelope {
3625 sequence: 11,
3626 event: NodeEvent::TerminalFrame {
3627 address,
3628 frame: TerminalFrame {
3629 sequence: 3,
3630 size: TerminalSize {
3631 rows: 24,
3632 columns: 80,
3633 },
3634 cursor_row: 4,
3635 cursor_column: 5,
3636 contents: "ready".to_owned(),
3637 formatted: vec![1, 2, 3],
3638 scrollback_formatted: Vec::new(),
3639 alternate_screen: false,
3640 mouse_protocol_enabled: false,
3641 mouse_protocol_encoding: Default::default(),
3642 produced_at_unix_ms: 0,
3643 screen_state: PtyScreenState::default(),
3644 bracketed_paste: None,
3645 },
3646 },
3647 });
3648
3649 assert!(matches!(
3650 ensure_server_frame_terminal_capability(&event, false),
3651 Err(NodeClientError::UnsupportedCapability(capability))
3652 if capability == NODE_TERMINAL_FRAME_EVENTS_CAPABILITY
3653 ));
3654 assert!(ensure_server_frame_terminal_capability(&event, true).is_ok());
3655
3656 let (_, mut selected) = negotiated_fixture();
3657 assert!(!selected_supports_terminal_frame_events(Some(&selected)));
3658 selected.capabilities.push(
3659 CapabilityId::new(NODE_TERMINAL_FRAME_EVENTS_CAPABILITY).unwrap(),
3660 );
3661 assert!(selected_supports_terminal_frame_events(Some(&selected)));
3662 }
3663
3664 #[tokio::test]
3665 async fn cancelled_recv_preserves_the_next_complete_server_frame() {
3666 let (client_stream, mut server_stream) = tokio::io::duplex(4_096);
3667 let (writer, frame_rx, reader_abort) =
3668 start_server_frame_reader(Box::new(client_stream));
3669 let mut client = LocalNodeClient {
3670 writer,
3671 frame_rx,
3672 _reader_abort: reader_abort,
3673 hello: hello_with_snapshot(empty_snapshot()),
3674 opaque_unix_paths_enabled: false,
3675 repository_paths_enabled: false,
3676 open_provider_ids_enabled: false,
3677 terminal_frame_events_enabled: false,
3678 negotiated_capabilities: Vec::new(),
3679 next_request_id: 1,
3680 pending_events: VecDeque::new(),
3681 pending_event_wire_bytes: 0,
3682 };
3683
3684 tokio::select! {
3685 _ = async { tokio::task::yield_now().await } => {}
3686 result = client.recv() => panic!("empty receive completed unexpectedly: {result:?}"),
3687 }
3688
3689 let expected = ServerFrame::Event(NodeEventEnvelope {
3690 sequence: 1,
3691 event: NodeEvent::ControllerChanged { controller: None },
3692 });
3693 write_json_frame_limited(
3694 &mut server_stream,
3695 &expected,
3696 MAX_NODE_FRAME_BYTES,
3697 )
3698 .await
3699 .unwrap();
3700 assert_eq!(client.recv().await.unwrap(), expected);
3701 }
3702
3703 #[test]
3704 fn pending_events_coalesce_exact_terminal_address_and_enforce_wire_byte_bound() {
3705 let address = SessionAddress {
3706 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3707 session: SessionKey {
3708 instance_id: AgentInstanceId(7),
3709 generation: SessionGeneration(2),
3710 },
3711 };
3712 let terminal_event = |event_sequence, address: SessionAddress, frame_sequence| {
3713 NodeEventEnvelope {
3714 sequence: event_sequence,
3715 event: NodeEvent::TerminalFrame {
3716 address,
3717 frame: TerminalFrame {
3718 sequence: frame_sequence,
3719 size: TerminalSize {
3720 rows: 24,
3721 columns: 80,
3722 },
3723 cursor_row: 0,
3724 cursor_column: 0,
3725 contents: format!("frame-{frame_sequence}"),
3726 formatted: vec![frame_sequence as u8],
3727 scrollback_formatted: Vec::new(),
3728 alternate_screen: false,
3729 mouse_protocol_enabled: false,
3730 mouse_protocol_encoding: Default::default(),
3731 produced_at_unix_ms: 0,
3732 screen_state: PtyScreenState::default(),
3733 bracketed_paste: None,
3734 },
3735 },
3736 }
3737 };
3738 let mut pending = VecDeque::new();
3739 let mut wire_bytes = 0;
3740 queue_pending_event_bounded(
3741 &mut pending,
3742 &mut wire_bytes,
3743 terminal_event(1, address.clone(), 1),
3744 4,
3745 4,
3746 9,
3747 )
3748 .unwrap();
3749 queue_pending_event_bounded(
3750 &mut pending,
3751 &mut wire_bytes,
3752 NodeEventEnvelope {
3753 sequence: 2,
3754 event: NodeEvent::ControllerChanged { controller: None },
3755 },
3756 1,
3757 4,
3758 9,
3759 )
3760 .unwrap();
3761 queue_pending_event_bounded(
3762 &mut pending,
3763 &mut wire_bytes,
3764 terminal_event(3, address.clone(), 2),
3765 5,
3766 4,
3767 9,
3768 )
3769 .unwrap();
3770 assert_eq!(pending.len(), 2);
3771 assert_eq!(wire_bytes, 6);
3772 assert_eq!(pending[0].envelope.sequence, 2);
3773 assert_eq!(pending[1].envelope.sequence, 3);
3774
3775 let mut next_generation = address.clone();
3776 next_generation.session.generation = SessionGeneration(3);
3777 queue_pending_event_bounded(
3778 &mut pending,
3779 &mut wire_bytes,
3780 terminal_event(4, next_generation, 1),
3781 3,
3782 4,
3783 9,
3784 )
3785 .unwrap();
3786 assert_eq!(pending.len(), 3);
3787 assert_eq!(wire_bytes, 9);
3788 let mut third_address = address.clone();
3789 third_address.session.instance_id = AgentInstanceId(8);
3790 assert!(matches!(
3791 queue_pending_event_bounded(
3792 &mut pending,
3793 &mut wire_bytes,
3794 terminal_event(5, third_address, 1),
3795 1,
3796 4,
3797 9,
3798 ),
3799 Err(NodeClientError::Protocol(message))
3800 if message.contains("wire byte capacity")
3801 ));
3802 assert_eq!(pending.len(), 3);
3803 assert_eq!(wire_bytes, 9);
3804
3805 let mut regular = VecDeque::new();
3806 let mut regular_bytes = 0;
3807 for sequence in 1..=2 {
3808 queue_pending_event_bounded(
3809 &mut regular,
3810 &mut regular_bytes,
3811 NodeEventEnvelope {
3812 sequence,
3813 event: NodeEvent::ControllerChanged { controller: None },
3814 },
3815 1,
3816 2,
3817 8,
3818 )
3819 .unwrap();
3820 }
3821 assert!(matches!(
3822 queue_pending_event_bounded(
3823 &mut regular,
3824 &mut regular_bytes,
3825 NodeEventEnvelope {
3826 sequence: 3,
3827 event: NodeEvent::ControllerChanged { controller: None },
3828 },
3829 1,
3830 2,
3831 8,
3832 ),
3833 Err(NodeClientError::Protocol(message))
3834 if message.contains("event capacity")
3835 ));
3836 assert_eq!(
3837 regular
3838 .iter()
3839 .map(|pending| pending.envelope.sequence)
3840 .collect::<Vec<_>>(),
3841 vec![1, 2],
3842 );
3843 }
3844
3845 #[test]
3846 fn outbound_open_provider_gate_preserves_all_legacy_ids() {
3847 for provider in [agent("claude"), agent("codex"), agent("kimi")] {
3848 assert!(ensure_outbound_provider_id_capability(&provider, false).is_ok());
3849 }
3850 let error = ensure_outbound_provider_id_capability(&agent("third-party-agent"), false)
3851 .unwrap_err();
3852 assert!(matches!(
3853 error,
3854 NodeClientError::UnsupportedCapability(capability)
3855 if capability == NODE_PROVIDER_ID_OPEN_CAPABILITY
3856 ));
3857 assert!(ensure_outbound_provider_id_capability(&agent("third-party-agent"), true).is_ok());
3858 }
3859
3860 #[test]
3861 fn inbound_open_provider_payloads_require_the_authenticated_capability() {
3862 let mut snapshot = empty_snapshot();
3863 snapshot.enabled_providers.push(agent("third-party-agent"));
3864 assert!(ensure_node_hello_provider_capability(
3865 &hello_with_snapshot(snapshot.clone()),
3866 false,
3867 )
3868 .is_err());
3869 assert!(ensure_node_hello_provider_capability(
3870 &hello_with_snapshot(snapshot),
3871 true,
3872 )
3873 .is_ok());
3874
3875 let mut runtime_snapshot = empty_snapshot();
3876 runtime_snapshot.provider_runtime_statuses = ProviderRuntimeStatuses::new([
3877 ProviderRuntimeStatus::unavailable(agent("grok")),
3878 ])
3879 .unwrap();
3880 assert!(ensure_node_hello_provider_capability(
3881 &hello_with_snapshot(runtime_snapshot),
3882 false,
3883 )
3884 .is_err());
3885
3886 let mut record = session_record_with_path(utf8_path());
3887 record.provider = agent("third-party-agent");
3888 let reply = response_frame(NodeResponse::SessionRecordUpdated {
3889 record: record.clone(),
3890 });
3891 assert!(ensure_server_frame_provider_capability(&reply, false).is_err());
3892 let event = ServerFrame::Event(NodeEventEnvelope {
3893 sequence: 1,
3894 event: NodeEvent::SessionRecordUpserted { record },
3895 });
3896 assert!(ensure_server_frame_provider_capability(&event, false).is_err());
3897 assert!(ensure_server_frame_provider_capability(&event, true).is_ok());
3898
3899 let mut open_workspace = workspace_with_path(utf8_path());
3900 open_workspace.sessions.push(session_snapshot("third-party-agent"));
3901 let mut nested_snapshot = empty_snapshot();
3902 nested_snapshot.workspaces.push(open_workspace.clone());
3903 let nested_frames = [
3904 ServerFrame::Hello(hello_with_snapshot(nested_snapshot)),
3905 ServerFrame::Event(NodeEventEnvelope {
3906 sequence: 2,
3907 event: NodeEvent::WorkspaceAdded {
3908 workspace: open_workspace.clone(),
3909 },
3910 }),
3911 response_frame(NodeResponse::WorkspaceRegistered {
3912 workspace: open_workspace.clone(),
3913 }),
3914 response_frame(NodeResponse::WorktreeCreated {
3915 worktree: worktree_with_path(utf8_path()),
3916 workspace: open_workspace,
3917 }),
3918 ];
3919 for frame in nested_frames {
3920 assert!(ensure_server_frame_provider_capability(&frame, false).is_err());
3921 assert!(ensure_server_frame_provider_capability(&frame, true).is_ok());
3922 }
3923
3924 let mut legacy_workspace = workspace_with_path(utf8_path());
3925 legacy_workspace.sessions.push(session_snapshot("claude"));
3926 let legacy_frame = response_frame(NodeResponse::WorkspaceRegistered {
3927 workspace: legacy_workspace,
3928 });
3929 assert!(ensure_server_frame_provider_capability(&legacy_frame, false).is_ok());
3930
3931 let mut context = context_pack_receipt();
3932 context.lineage.source_provider = agent("third-party-agent");
3933 let context_frame = response_frame(NodeResponse::ContextPackExported { context });
3934 assert!(ensure_server_frame_provider_capability(&context_frame, false).is_err());
3935 assert!(ensure_server_frame_provider_capability(&context_frame, true).is_ok());
3936 }
3937
3938 #[test]
3939 fn selected_open_provider_manifest_requires_the_open_id_capability() {
3940 let (mut offer, mut selected) = negotiated_fixture();
3941 let manifest_capability =
3942 CapabilityId::new(NODE_PROVIDER_CONTRACT_MANIFEST_CAPABILITY).unwrap();
3943 offer.capabilities.push(manifest_capability.clone());
3944 selected.capabilities.push(manifest_capability);
3945 selected.provider_contracts.push(ProviderContractSupport {
3946 provider: agent("third-party-agent"),
3947 revision: ProviderContractRevision::new("third-party.2026-08").unwrap(),
3948 });
3949 assert!(matches!(
3950 validate_selected_compatibility(&offer, &selected),
3951 Err(NodeClientError::Protocol(message))
3952 if message.contains("open provider ID")
3953 ));
3954
3955 let open_capability = CapabilityId::new(NODE_PROVIDER_ID_OPEN_CAPABILITY).unwrap();
3956 offer.capabilities.push(open_capability.clone());
3957 selected.capabilities.push(open_capability);
3958 assert!(validate_selected_compatibility(&offer, &selected).is_ok());
3959 }
3960
3961 #[test]
3962 fn unnegotiated_workspace_file_read_is_rejected_before_consuming_a_request_id() {
3963 let request = NodeRequest::ReadWorkspaceFile {
3964 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
3965 path: utf8_repository_path("src/lib.rs"),
3966 };
3967 let mut next_request_id = 41;
3968 let no_capabilities = Vec::new();
3969 let error = reserve_request_id(
3970 &mut next_request_id,
3971 &request,
3972 false,
3973 false,
3974 false,
3975 &no_capabilities,
3976 )
3977 .unwrap_err();
3978 assert!(matches!(
3979 error,
3980 NodeClientError::UnsupportedCapability(capability)
3981 if capability == NODE_WORKSPACE_FILE_READ_CAPABILITY
3982 ));
3983 assert_eq!(next_request_id, 41);
3984
3985 let file_read_capabilities = vec![
3986 CapabilityId::new(NODE_WORKSPACE_FILE_READ_CAPABILITY).unwrap(),
3987 ];
3988 let request_id = reserve_request_id(
3989 &mut next_request_id,
3990 &request,
3991 false,
3992 false,
3993 false,
3994 &file_read_capabilities,
3995 )
3996 .unwrap();
3997 assert_eq!(request_id, 41);
3998 assert_eq!(next_request_id, 42);
3999 }
4000
4001 #[test]
4002 fn unnegotiated_workspace_file_response_is_rejected_before_exposure() {
4003 let frame = response_frame(NodeResponse::WorkspaceFileRead {
4004 file: WorkspaceFileRead {
4005 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4006 path: utf8_repository_path("src/lib.rs"),
4007 content: WorkspaceFileContent::Utf8 {
4008 text: "hello".to_owned(),
4009 byte_len: 5,
4010 },
4011 revision: None,
4012 },
4013 });
4014 assert!(ensure_server_frame_required_capability(&frame, &[]).is_err());
4015 let capabilities = vec![
4016 CapabilityId::new(NODE_WORKSPACE_FILE_READ_CAPABILITY).unwrap(),
4017 ];
4018 assert!(ensure_server_frame_required_capability(&frame, &capabilities).is_ok());
4019 }
4020
4021 #[test]
4022 fn workspace_entry_create_is_capability_gated_and_response_correlated_fail_closed() {
4023 let workspace_id = WorkspaceId::new("workspace-a").unwrap();
4024 let file_path = utf8_repository_path("src/new.rs");
4025 let file_request = NodeRequest::CreateWorkspaceFile {
4026 workspace_id: workspace_id.clone(),
4027 path: file_path.clone(),
4028 };
4029 let mut next_request_id = 73;
4030 let error = reserve_request_id(
4031 &mut next_request_id,
4032 &file_request,
4033 false,
4034 false,
4035 false,
4036 &[],
4037 )
4038 .unwrap_err();
4039 assert!(matches!(
4040 error,
4041 NodeClientError::UnsupportedCapability(capability)
4042 if capability == NODE_WORKSPACE_ENTRY_CREATE_CAPABILITY
4043 ));
4044 assert_eq!(next_request_id, 73);
4045
4046 let capabilities = vec![
4047 CapabilityId::new(NODE_WORKSPACE_ENTRY_CREATE_CAPABILITY).unwrap(),
4048 ];
4049 assert_eq!(
4050 reserve_request_id(
4051 &mut next_request_id,
4052 &file_request,
4053 false,
4054 false,
4055 false,
4056 &capabilities,
4057 )
4058 .unwrap(),
4059 73,
4060 );
4061
4062 let file = WorkspaceFileRead {
4063 workspace_id: workspace_id.clone(),
4064 path: file_path.clone(),
4065 content: WorkspaceFileContent::Utf8 {
4066 text: String::new(),
4067 byte_len: 0,
4068 },
4069 revision: Some(WorkspaceFileRevision::new(
4070 "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"
4071 .to_owned(),
4072 ).unwrap()),
4073 };
4074 let file_frame = response_frame(NodeResponse::WorkspaceFileCreated {
4075 file: file.clone(),
4076 });
4077 assert!(ensure_server_frame_required_capability(&file_frame, &[]).is_err());
4078 assert!(ensure_server_frame_required_capability(&file_frame, &capabilities).is_ok());
4079 assert!(validate_workspace_content_response(
4080 &file_request,
4081 &Ok(NodeResponse::WorkspaceFileCreated { file: file.clone() }),
4082 ).is_ok());
4083
4084 let mut wrong_content = file;
4085 wrong_content.content = WorkspaceFileContent::Utf8 {
4086 text: "not-empty".to_owned(),
4087 byte_len: 9,
4088 };
4089 assert!(validate_workspace_content_response(
4090 &file_request,
4091 &Ok(NodeResponse::WorkspaceFileCreated { file: wrong_content }),
4092 ).is_err());
4093
4094 let directory_path = utf8_repository_path("src/new");
4095 let directory_request = NodeRequest::CreateWorkspaceDirectory {
4096 workspace_id: workspace_id.clone(),
4097 path: directory_path.clone(),
4098 };
4099 let directory_response = NodeResponse::WorkspaceDirectoryCreated {
4100 workspace_id,
4101 entry: WorkspaceEntry {
4102 relative_path: directory_path,
4103 kind: WorkspaceEntryKind::Directory,
4104 },
4105 };
4106 assert!(validate_workspace_content_response(
4107 &directory_request,
4108 &Ok(directory_response),
4109 ).is_ok());
4110 let tagged_directory_request = NodeRequest::CreateWorkspaceDirectory {
4111 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4112 path: tagged_repository_path(b"src/\xff"),
4113 };
4114 assert!(ensure_node_request_path_capability(
4115 &tagged_directory_request,
4116 false,
4117 false,
4118 ).is_err());
4119 assert!(ensure_node_request_path_capability(
4120 &tagged_directory_request,
4121 false,
4122 true,
4123 ).is_ok());
4124 let tagged_directory_frame = response_frame(
4125 NodeResponse::WorkspaceDirectoryCreated {
4126 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4127 entry: WorkspaceEntry {
4128 relative_path: tagged_repository_path(b"src/\xff"),
4129 kind: WorkspaceEntryKind::Directory,
4130 },
4131 },
4132 );
4133 assert!(ensure_server_frame_path_capability(
4134 &tagged_directory_frame,
4135 false,
4136 false,
4137 ).is_err());
4138 assert!(ensure_server_frame_path_capability(
4139 &tagged_directory_frame,
4140 false,
4141 true,
4142 ).is_ok());
4143 assert!(validate_workspace_content_response(
4144 &directory_request,
4145 &Ok(NodeResponse::WorkspaceFileCreated {
4146 file: WorkspaceFileRead {
4147 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4148 path: utf8_repository_path("src/new"),
4149 content: WorkspaceFileContent::Utf8 {
4150 text: String::new(),
4151 byte_len: 0,
4152 },
4153 revision: None,
4154 },
4155 }),
4156 ).is_err());
4157 }
4158
4159 #[test]
4160 fn ordinary_startup_snapshot_is_not_rejected_as_native_session_preview() {
4161 let frame = response_frame(NodeResponse::Snapshot {
4162 event_sequence: 0,
4163 controller: None,
4164 snapshot: empty_snapshot(),
4165 });
4166
4167 assert!(ensure_server_frame_required_capability(&frame, &[]).is_ok());
4168 }
4169
4170 #[test]
4171 fn host_directory_browse_requires_capability_and_preserves_opaque_paths() {
4172 let unix = OpaqueHostPath::unix_bytes(b"/srv/\xff".to_vec()).unwrap();
4173 let request = NodeRequest::BrowseHostDirectories {
4174 directory: None,
4175 after: Some(unix.clone()),
4176 };
4177 assert!(ensure_node_request_required_capability(&request, &[]).is_err());
4178 assert!(node_request_contains_opaque_unix_path(&request));
4179 let capability = CapabilityId::new(CAPABILITY_HOST_DIRECTORY_BROWSE_V1).unwrap();
4180 assert!(ensure_node_request_required_capability(
4181 &request,
4182 &[capability.clone()],
4183 ).is_ok());
4184
4185 let frame = response_frame(NodeResponse::HostDirectoriesBrowsed {
4186 listing: HostDirectoryListing {
4187 directory: Some(unix.clone()),
4188 parent: None,
4189 entries: vec![HostDirectoryEntry {
4190 path: unix,
4191 display_name: "opaque".to_owned(),
4192 is_link: false,
4193 }],
4194 next_after: None,
4195 incomplete: false,
4196 },
4197 });
4198 assert!(ensure_server_frame_required_capability(&frame, &[]).is_err());
4199 assert!(ensure_server_frame_required_capability(&frame, &[capability]).is_ok());
4200 assert!(server_frame_contains_opaque_unix_path(&frame));
4201 }
4202
4203 #[test]
4204 fn standalone_workspace_lifecycle_is_exactly_capability_and_path_gated() {
4205 let root = OpaqueHostPath::unix_bytes(b"/srv/standalone".to_vec()).unwrap();
4206 let request = NodeRequest::CreateStandaloneWorkspace {
4207 workspace_id: WorkspaceId::new("standalone").unwrap(),
4208 root: root.clone(),
4209 initial_branch: Some("main".to_owned()),
4210 };
4211 let capability = CapabilityId::new(
4212 NODE_STANDALONE_WORKSPACE_LIFECYCLE_CAPABILITY,
4213 ).unwrap();
4214
4215 assert!(matches!(
4216 ensure_node_request_required_capability(&request, &[]),
4217 Err(NodeClientError::UnsupportedCapability(required))
4218 if required == NODE_STANDALONE_WORKSPACE_LIFECYCLE_CAPABILITY
4219 ));
4220 assert!(ensure_node_request_required_capability(
4221 &request,
4222 std::slice::from_ref(&capability),
4223 ).is_ok());
4224 assert!(node_request_contains_opaque_unix_path(&request));
4225
4226 let frame = response_frame(NodeResponse::StandaloneWorkspaceCreated {
4227 workspace: workspace_with_path(root),
4228 });
4229 assert!(matches!(
4230 ensure_server_frame_required_capability(&frame, &[]),
4231 Err(NodeClientError::UnsupportedCapability(required))
4232 if required == NODE_STANDALONE_WORKSPACE_LIFECYCLE_CAPABILITY
4233 ));
4234 assert!(ensure_server_frame_required_capability(&frame, &[capability]).is_ok());
4235 assert!(server_frame_contains_opaque_unix_path(&frame));
4236 }
4237
4238 #[test]
4239 fn spawn_profile_revision_capability_is_advertised_and_gates_all_spawn_spec_paths() {
4240 let offer = client_compatibility_offer().unwrap();
4241 let profile_revision =
4242 CapabilityId::new(NODE_SPAWN_PROFILE_REVISION_CAPABILITY).unwrap();
4243 assert!(offer.capabilities.contains(&profile_revision));
4244
4245 let NodeRequest::SpawnSpec { spec } = spawn_spec_request() else {
4246 unreachable!("spawn spec helper changed variant");
4247 };
4248 let reservation_id = gate4agent_node_protocol::HarnessMcpReservationId::new(
4249 format!("hmcpres_{}", "a".repeat(24)),
4250 ).unwrap();
4251 let activation_digest = gate4agent_node_protocol::HarnessMcpActivationDigest::new(
4252 format!("sha256:{}", "b".repeat(64)),
4253 ).unwrap();
4254 let requests = [
4255 NodeRequest::SpawnSpec { spec: spec.clone() },
4256 NodeRequest::SpawnManagedWorktree {
4257 request: ManagedWorktreeSpawnRequest {
4258 spawn_spec: spec.clone(),
4259 worktree_profile_id: WorktreeProfileId::new("review").unwrap(),
4260 },
4261 },
4262 NodeRequest::ArmHarnessMcpReservation {
4263 reservation_id: reservation_id.clone(),
4264 activation_digest: activation_digest.clone(),
4265 spawn_spec: spec.clone(),
4266 launch: gate4agent_node_protocol::HarnessMcpLaunchV1 {
4267 server_name: "fixture-harness".to_owned(),
4268 args: vec!["--session-proxy".to_owned()],
4269 endpoint_env: "FIXTURE_SESSION_ENDPOINT".to_owned(),
4270 token_env: "FIXTURE_SESSION_TOKEN".to_owned(),
4271 program_env: "FIXTURE_MCP_PROGRAM".to_owned(),
4272 trace: None,
4273 scrub_env: Vec::new(),
4274 },
4275 expires_at_unix_ms: 10_000,
4276 },
4277 NodeRequest::SpawnSpecWithHarnessMcp {
4278 reservation_id,
4279 activation_digest,
4280 spec,
4281 deadline_unix_ms: 10_000,
4282 },
4283 ];
4284 let mut capabilities = vec![
4285 CapabilityId::new(NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY).unwrap(),
4286 CapabilityId::new(NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY).unwrap(),
4287 CapabilityId::new(NODE_HARNESS_MCP_READ_PROXY_CAPABILITY).unwrap(),
4288 ];
4289 for request in &requests {
4290 assert!(matches!(
4291 ensure_node_request_required_capability(request, &capabilities),
4292 Err(NodeClientError::UnsupportedCapability(capability))
4293 if capability == NODE_SPAWN_PROFILE_REVISION_CAPABILITY
4294 ));
4295 }
4296 capabilities.push(profile_revision.clone());
4297 for request in &requests {
4298 let result = ensure_node_request_required_capability(request, &capabilities);
4299 assert!(!matches!(
4300 result,
4301 Err(NodeClientError::UnsupportedCapability(capability))
4302 if capability == NODE_SPAWN_PROFILE_REVISION_CAPABILITY
4303 ));
4304 }
4305
4306 let frame = response_frame(NodeResponse::SpawnSpecAccepted {
4307 receipt: spawn_spec_receipt(),
4308 });
4309 capabilities.retain(|capability| capability != &profile_revision);
4310 assert!(matches!(
4311 ensure_server_frame_required_capability(&frame, &capabilities),
4312 Err(NodeClientError::UnsupportedCapability(capability))
4313 if capability == NODE_SPAWN_PROFILE_REVISION_CAPABILITY
4314 ));
4315 capabilities.push(profile_revision);
4316 assert!(ensure_server_frame_required_capability(&frame, &capabilities).is_ok());
4317 }
4318
4319 #[test]
4320 fn unnegotiated_spawn_spec_request_and_receipt_fail_closed() {
4321 let request = spawn_spec_request();
4322 let mut next_request_id = 73;
4323 let error = reserve_request_id(
4324 &mut next_request_id,
4325 &request,
4326 false,
4327 false,
4328 false,
4329 &[],
4330 )
4331 .unwrap_err();
4332 assert!(matches!(
4333 error,
4334 NodeClientError::UnsupportedCapability(capability)
4335 if capability == NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY
4336 ));
4337 assert_eq!(next_request_id, 73);
4338
4339 let capabilities = spawn_gate_capabilities_without(NODE_WORKTREE_SELECTION_CAPABILITY);
4342 assert_eq!(
4343 reserve_request_id(
4344 &mut next_request_id,
4345 &request,
4346 false,
4347 false,
4348 false,
4349 &capabilities,
4350 )
4351 .unwrap(),
4352 73,
4353 );
4354 assert_eq!(next_request_id, 74);
4355
4356 let frame = response_frame(NodeResponse::SpawnSpecAccepted {
4357 receipt: spawn_spec_receipt(),
4358 });
4359 assert!(matches!(
4360 ensure_server_frame_required_capability(&frame, &[]),
4361 Err(NodeClientError::UnsupportedCapability(capability))
4362 if capability == NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY
4363 ));
4364 assert!(ensure_server_frame_required_capability(&frame, &capabilities).is_ok());
4365
4366 let mut worktree_request = spawn_spec_request();
4367 let NodeRequest::SpawnSpec { spec } = &mut worktree_request else {
4368 unreachable!("spawn spec helper changed variant");
4369 };
4370 spec.target.worktree_id = Some(WorkspaceId::new("review-tree").unwrap());
4371 let mut worktree_request_id = 81;
4372 assert!(matches!(
4373 reserve_request_id(
4374 &mut worktree_request_id,
4375 &worktree_request,
4376 false,
4377 false,
4378 false,
4379 &capabilities,
4380 ),
4381 Err(NodeClientError::UnsupportedCapability(capability))
4382 if capability == NODE_WORKTREE_SELECTION_CAPABILITY
4383 ));
4384 assert_eq!(worktree_request_id, 81);
4385 let mut worktree_capabilities = capabilities.clone();
4386 worktree_capabilities.push(
4387 CapabilityId::new(NODE_WORKTREE_SELECTION_CAPABILITY).unwrap(),
4388 );
4389 assert_eq!(
4390 reserve_request_id(
4391 &mut worktree_request_id,
4392 &worktree_request,
4393 false,
4394 false,
4395 false,
4396 &worktree_capabilities,
4397 )
4398 .unwrap(),
4399 81,
4400 );
4401
4402 let mut worktree_receipt = spawn_spec_receipt();
4403 worktree_receipt.target.worktree_id =
4404 Some(WorkspaceId::new("review-tree").unwrap());
4405 let worktree_frame = response_frame(NodeResponse::SpawnSpecAccepted {
4406 receipt: worktree_receipt,
4407 });
4408 assert!(matches!(
4409 ensure_server_frame_required_capability(
4410 &worktree_frame,
4411 &capabilities,
4412 ),
4413 Err(NodeClientError::UnsupportedCapability(capability))
4414 if capability == NODE_WORKTREE_SELECTION_CAPABILITY
4415 ));
4416 assert!(ensure_server_frame_required_capability(
4417 &worktree_frame,
4418 &worktree_capabilities,
4419 )
4420 .is_ok());
4421 }
4422
4423 #[test]
4424 fn child_environment_profile_is_gated_before_write_and_on_recursive_read() {
4425 let mut request = spawn_spec_request();
4426 let NodeRequest::SpawnSpec { spec } = &mut request else {
4427 unreachable!("spawn spec helper changed variant");
4428 };
4429 spec.overrides.environment_profile_id =
4430 gate4agent_node_protocol::SpawnOverride::Set {
4431 value: SpawnEnvironmentProfileId::new("local-default").unwrap(),
4432 };
4433 let granted =
4434 spawn_gate_capabilities_without(NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY);
4435 let mut next_request_id = 91;
4436 assert!(matches!(
4437 reserve_request_id(&mut next_request_id, &request, false, false, false, &granted),
4438 Err(NodeClientError::UnsupportedCapability(capability))
4439 if capability == NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY
4440 ));
4441 assert_eq!(next_request_id, 91);
4442
4443 let mut cleared = spawn_spec_request();
4444 let NodeRequest::SpawnSpec { spec } = &mut cleared else {
4445 unreachable!("spawn spec helper changed variant");
4446 };
4447 spec.overrides.environment_profile_id =
4448 gate4agent_node_protocol::SpawnOverride::Clear;
4449 let mut clear_request_id = 101;
4450 assert_eq!(
4451 reserve_request_id(
4452 &mut clear_request_id,
4453 &cleared,
4454 false,
4455 false,
4456 false,
4457 &granted,
4458 )
4459 .unwrap(),
4460 101,
4461 );
4462
4463 let environment_capability =
4464 CapabilityId::new(NODE_CHILD_ENVIRONMENT_PROFILE_CAPABILITY).unwrap();
4465 let mut fully_granted = granted.clone();
4466 fully_granted.push(environment_capability.clone());
4467 assert_eq!(
4468 reserve_request_id(
4469 &mut next_request_id,
4470 &request,
4471 false,
4472 false,
4473 false,
4474 &fully_granted,
4475 )
4476 .unwrap(),
4477 91,
4478 );
4479
4480 let mut receipt = spawn_spec_receipt();
4481 receipt.environment_profile = Some(ResolvedEnvironmentProfileReceipt {
4482 profile_id: SpawnEnvironmentProfileId::new("local-default").unwrap(),
4483 profile_revision: SpawnEnvironmentProfileRevision::new("local-default.r1")
4484 .unwrap(),
4485 network_allowlist: None,
4486 browser_profile_id: None,
4487 });
4488 let frame = response_frame(NodeResponse::SpawnSpecAccepted { receipt });
4489 assert!(matches!(
4490 ensure_server_frame_required_capability(&frame, &granted),
4491 Err(NodeClientError::Protocol(ref message))
4492 if message.contains("child environment profile metadata")
4493 ));
4494 assert!(ensure_server_frame_required_capability(&frame, &fully_granted).is_ok());
4495 }
4496
4497 #[test]
4498 fn session_bundle_materialization_is_gated_before_write_and_on_recursive_read() {
4499 let bundle_capability =
4500 CapabilityId::new(NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY).unwrap();
4501 let granted =
4502 spawn_gate_capabilities_without(NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY);
4503 let mut fully_granted = granted.clone();
4504 fully_granted.push(bundle_capability.clone());
4505 let request = spawn_spec_request();
4506 let mut next_request_id = 111;
4507 assert!(matches!(
4508 reserve_request_id(&mut next_request_id, &request, false, false, false, &granted),
4509 Err(NodeClientError::UnsupportedCapability(capability))
4510 if capability == NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY
4511 ));
4512 assert_eq!(next_request_id, 111);
4513 assert_eq!(
4514 reserve_request_id(
4515 &mut next_request_id,
4516 &request,
4517 false,
4518 false,
4519 false,
4520 &fully_granted,
4521 )
4522 .unwrap(),
4523 111,
4524 );
4525
4526 let mut cleared = spawn_spec_request();
4527 let NodeRequest::SpawnSpec { spec } = &mut cleared else {
4528 unreachable!("spawn spec helper changed variant");
4529 };
4530 spec.overrides.bundle_id = gate4agent_node_protocol::SpawnOverride::Clear;
4531 let mut clear_request_id = 121;
4532 assert_eq!(
4533 reserve_request_id(
4534 &mut clear_request_id,
4535 &cleared,
4536 false,
4537 false,
4538 false,
4539 &granted,
4540 )
4541 .unwrap(),
4542 121,
4543 );
4544
4545 let mut receipt = spawn_spec_receipt();
4546 let bundle_id = SpawnBundleId::new("review-bundle").unwrap();
4547 receipt.bundle_id = Some(bundle_id.clone());
4548 receipt.bundle = Some(ResolvedBundleReceipt {
4549 id: bundle_id,
4550 revision: SpawnBundleRevision::new("review-bundle.r1").unwrap(),
4551 digest: SpawnBundleDigest::new(format!("sha256:{}", "a".repeat(64)))
4552 .unwrap(),
4553 });
4554 let frame = response_frame(NodeResponse::SpawnSpecAccepted { receipt });
4555 assert!(matches!(
4556 ensure_server_frame_required_capability(&frame, &granted),
4557 Err(NodeClientError::Protocol(ref message))
4558 if message.contains("session bundle materialization metadata")
4559 ));
4560 assert!(ensure_server_frame_required_capability(&frame, &fully_granted).is_ok());
4561 }
4562
4563 #[test]
4564 fn history_context_pack_requests_and_responses_fail_closed() {
4565 let session = SessionAddress {
4566 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4567 session: SessionKey {
4568 instance_id: AgentInstanceId(7),
4569 generation: SessionGeneration(2),
4570 },
4571 };
4572 let context_id = SpawnContextId::new("context-a").unwrap();
4573 let requests = [
4574 NodeRequest::DiscoverHistory {
4575 session: session.clone(),
4576 limit: 8,
4577 },
4578 NodeRequest::LoadHistory {
4579 session: session.clone(),
4580 candidate_id: "candidate-a".to_owned(),
4581 },
4582 NodeRequest::ExportContextPack {
4583 session: session.clone(),
4584 },
4585 NodeRequest::ForgetContextPack {
4586 context_id: context_id.clone(),
4587 },
4588 ];
4589 let capability = CapabilityId::new(NODE_HISTORY_CONTEXT_PACK_CAPABILITY).unwrap();
4590 for request in requests {
4591 assert!(matches!(
4592 ensure_node_request_required_capability(&request, &[]),
4593 Err(NodeClientError::UnsupportedCapability(ref required))
4594 if required == NODE_HISTORY_CONTEXT_PACK_CAPABILITY
4595 ));
4596 assert!(ensure_node_request_required_capability(
4597 &request,
4598 std::slice::from_ref(&capability),
4599 )
4600 .is_ok());
4601 }
4602
4603 let mut next_request_id = 41;
4604 let invalid = NodeRequest::DiscoverHistory {
4605 session: session.clone(),
4606 limit: 0,
4607 };
4608 assert!(matches!(
4609 reserve_request_id(
4610 &mut next_request_id,
4611 &invalid,
4612 false,
4613 false,
4614 true,
4615 std::slice::from_ref(&capability),
4616 ),
4617 Err(NodeClientError::Protocol(ref message))
4618 if message == "invalid history context pack request"
4619 ));
4620 assert_eq!(next_request_id, 41);
4621
4622 let forget = NodeRequest::ForgetContextPack {
4623 context_id: context_id.clone(),
4624 };
4625 assert!(matches!(
4626 reserve_request_id(
4627 &mut next_request_id,
4628 &forget,
4629 false,
4630 false,
4631 false,
4632 std::slice::from_ref(&capability),
4633 ),
4634 Err(NodeClientError::UnsupportedCapability(ref required))
4635 if required == NODE_PROVIDER_ID_OPEN_CAPABILITY
4636 ));
4637 assert_eq!(next_request_id, 41);
4638 assert_eq!(
4639 reserve_request_id(
4640 &mut next_request_id,
4641 &forget,
4642 false,
4643 false,
4644 true,
4645 std::slice::from_ref(&capability),
4646 )
4647 .unwrap(),
4648 41,
4649 );
4650
4651 let responses = [
4652 NodeResponse::HistoryDiscovered {
4653 session: session.clone(),
4654 candidates: Vec::new(),
4655 },
4656 NodeResponse::HistoryLoaded {
4657 session,
4658 session_id: "provider-session-a".to_owned(),
4659 message_count: 4,
4660 completed_turn_count: None,
4661 },
4662 NodeResponse::ContextPackExported {
4663 context: context_pack_receipt(),
4664 },
4665 NodeResponse::ContextPackForgotten { context_id },
4666 ];
4667 for response in responses {
4668 let frame = response_frame(response);
4669 assert!(matches!(
4670 ensure_server_frame_required_capability(&frame, &[]),
4671 Err(NodeClientError::Protocol(ref message))
4672 if message.contains("history context pack metadata")
4673 ));
4674 assert!(ensure_server_frame_required_capability(
4675 &frame,
4676 std::slice::from_ref(&capability),
4677 )
4678 .is_ok());
4679 }
4680
4681 for code in [
4682 NodeFailureCode::UnknownContextPack,
4683 NodeFailureCode::ContextPackBusy,
4684 NodeFailureCode::ContextPackMaterializationFailed,
4685 ] {
4686 let frame = ServerFrame::Reply(ResponseEnvelope {
4687 request_id: 1,
4688 result: Err(NodeFailure {
4689 code,
4690 message: "typed F7 failure".to_owned(),
4691 }),
4692 });
4693 assert!(matches!(
4694 ensure_server_frame_required_capability(&frame, &[]),
4695 Err(NodeClientError::Protocol(ref message))
4696 if message.contains("history context pack metadata")
4697 ));
4698 assert!(ensure_server_frame_required_capability(
4699 &frame,
4700 std::slice::from_ref(&capability),
4701 )
4702 .is_ok());
4703 }
4704 }
4705
4706 #[test]
4707 fn history_context_pack_nested_metadata_fails_closed() {
4708 let history_capability =
4709 CapabilityId::new(NODE_HISTORY_CONTEXT_PACK_CAPABILITY).unwrap();
4710 let granted = spawn_gate_capabilities_without(NODE_HISTORY_CONTEXT_PACK_CAPABILITY);
4711 let mut request = spawn_spec_request();
4712 let NodeRequest::SpawnSpec { spec } = &mut request else {
4713 unreachable!("spawn spec helper changed variant");
4714 };
4715 spec.overrides.context_id = gate4agent_node_protocol::SpawnOverride::Set {
4716 value: SpawnContextId::new("context-a").unwrap(),
4717 };
4718 assert!(matches!(
4719 ensure_node_request_required_capability(&request, &granted),
4720 Err(NodeClientError::UnsupportedCapability(ref required))
4721 if required == NODE_HISTORY_CONTEXT_PACK_CAPABILITY
4722 ));
4723 let mut fully_granted = granted.clone();
4724 fully_granted.push(history_capability.clone());
4725 assert!(ensure_node_request_required_capability(&request, &fully_granted).is_ok());
4726
4727 let context = context_pack_receipt();
4728 let mut receipt = spawn_spec_receipt();
4729 receipt.context_id = Some(context.id.clone());
4730 receipt.context = Some(context.clone());
4731 let receipt_frame = response_frame(NodeResponse::SpawnSpecAccepted { receipt });
4732 assert!(matches!(
4733 ensure_server_frame_required_capability(&receipt_frame, &granted),
4734 Err(NodeClientError::Protocol(ref message))
4735 if message.contains("history context pack metadata")
4736 ));
4737 assert!(
4738 ensure_server_frame_required_capability(&receipt_frame, &fully_granted).is_ok(),
4739 );
4740
4741 let mut record = session_record_with_path(utf8_path());
4742 record.context_id = Some(context.id.clone());
4743 record.context = Some(context);
4744 let event_frame = ServerFrame::Event(NodeEventEnvelope {
4745 sequence: 1,
4746 event: NodeEvent::SessionRecordUpserted {
4747 record: record.clone(),
4748 },
4749 });
4750 assert!(ensure_server_frame_required_capability(&event_frame, &[]).is_err());
4751 assert!(ensure_server_frame_required_capability(
4752 &event_frame,
4753 std::slice::from_ref(&history_capability),
4754 )
4755 .is_ok());
4756
4757 let mut snapshot = empty_snapshot();
4758 snapshot.session_records.push(record);
4759 let hello = hello_with_snapshot(snapshot);
4760 assert!(ensure_node_hello_history_context_pack_capability(&hello, &[]).is_err());
4761 assert!(ensure_node_hello_history_context_pack_capability(
4762 &hello,
4763 std::slice::from_ref(&history_capability),
4764 )
4765 .is_ok());
4766 }
4767
4768 #[test]
4769 fn managed_worktree_partial_capability_intersections_fail_closed() {
4770 let NodeRequest::SpawnSpec { spec } = spawn_spec_request() else {
4771 unreachable!("spawn spec helper changed variant");
4772 };
4773 let request = NodeRequest::SpawnManagedWorktree {
4774 request: ManagedWorktreeSpawnRequest {
4775 spawn_spec: spec,
4776 worktree_profile_id: WorktreeProfileId::new("review").unwrap(),
4777 },
4778 };
4779 let managed = CapabilityId::new(NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY).unwrap();
4780 let spawn = CapabilityId::new(NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY).unwrap();
4781 let revision = CapabilityId::new(NODE_SPAWN_PROFILE_REVISION_CAPABILITY).unwrap();
4782 let worktree = CapabilityId::new(NODE_WORKTREE_SELECTION_CAPABILITY).unwrap();
4783 let bundle =
4784 CapabilityId::new(NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY).unwrap();
4785 let mut next_request_id = 91;
4786 assert!(matches!(
4787 reserve_request_id(&mut next_request_id, &request, false, false, false, &[]),
4788 Err(NodeClientError::UnsupportedCapability(capability))
4789 if capability == NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY
4790 ));
4791 assert!(matches!(
4792 reserve_request_id(
4793 &mut next_request_id,
4794 &request,
4795 false,
4796 false,
4797 false,
4798 &[managed.clone()],
4799 ),
4800 Err(NodeClientError::UnsupportedCapability(capability))
4801 if capability == NODE_SPAWN_SPEC_DEFAULTS_OVERRIDES_CAPABILITY
4802 ));
4803 assert!(matches!(
4809 reserve_request_id(
4810 &mut next_request_id,
4811 &request,
4812 false,
4813 false,
4814 false,
4815 &[managed.clone(), spawn.clone()],
4816 ),
4817 Err(NodeClientError::UnsupportedCapability(capability))
4818 if capability == NODE_SPAWN_PROFILE_REVISION_CAPABILITY
4819 ));
4820 assert!(matches!(
4821 reserve_request_id(
4822 &mut next_request_id,
4823 &request,
4824 false,
4825 false,
4826 false,
4827 &[managed.clone(), spawn.clone(), revision.clone()],
4828 ),
4829 Err(NodeClientError::UnsupportedCapability(capability))
4830 if capability == NODE_WORKTREE_SELECTION_CAPABILITY
4831 ));
4832 assert!(matches!(reserve_request_id(
4833 &mut next_request_id,
4834 &request,
4835 false,
4836 false,
4837 false,
4838 &[managed.clone(), spawn.clone(), revision.clone(), worktree.clone()],
4839 ), Err(NodeClientError::UnsupportedCapability(capability))
4840 if capability == NODE_SESSION_BUNDLE_MATERIALIZATION_CAPABILITY));
4841 assert!(reserve_request_id(
4842 &mut next_request_id,
4843 &request,
4844 false,
4845 false,
4846 false,
4847 &[managed.clone(), spawn.clone(), revision.clone(), worktree.clone(), bundle],
4848 )
4849 .is_ok());
4850
4851 let cleanup = NodeRequest::CleanupManagedWorktree {
4852 lease_id: ManagedWorktreeLeaseId::new("lease-a").unwrap(),
4853 };
4854 assert!(ensure_node_request_required_capability(&cleanup, &[managed.clone()]).is_err());
4855 assert!(ensure_node_request_required_capability(
4856 &cleanup,
4857 &[managed.clone(), worktree.clone()],
4858 )
4859 .is_ok());
4860
4861 let event = ServerFrame::Event(NodeEventEnvelope {
4862 sequence: 4,
4863 event: NodeEvent::ManagedWorktreeRemoved {
4864 lease_id: ManagedWorktreeLeaseId::new("lease-a").unwrap(),
4865 },
4866 });
4867 assert!(ensure_server_frame_required_capability(&event, &[managed.clone()]).is_err());
4868 assert!(ensure_server_frame_required_capability(&event, &[worktree.clone()]).is_err());
4869 assert!(ensure_server_frame_required_capability(
4870 &event,
4871 &[managed.clone(), worktree.clone()],
4872 )
4873 .is_ok());
4874
4875 let mut spawn_receipt = spawn_spec_receipt();
4876 let workspace_id = WorkspaceId::new("managed-a").unwrap();
4877 spawn_receipt.target.worktree_id = Some(workspace_id.clone());
4878 spawn_receipt.session.workspace_id = workspace_id.clone();
4879 let receipt = ManagedWorktreeSpawnReceipt {
4880 spawn: spawn_receipt,
4881 lease: ManagedWorktreeLeaseSnapshot {
4882 lease_id: ManagedWorktreeLeaseId::new("lease-a").unwrap(),
4883 source_workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4884 workspace_id,
4885 profile_id: WorktreeProfileId::new("review").unwrap(),
4886 profile_revision: WorktreeProfileRevision::new("review.r1").unwrap(),
4887 retention: ManagedWorktreeRetention::RemoveWhenReleased,
4888 state: ManagedWorktreeLeaseState::InUse,
4889 active_session_count: 1,
4890 managed_record_count: 1,
4891 cleanup_failure: None::<ManagedWorktreeCleanupFailure>,
4892 created_at_unix_ms: 1,
4893 updated_at_unix_ms: 2,
4894 },
4895 };
4896 let reply = response_frame(NodeResponse::ManagedWorktreeSpawnAccepted { receipt });
4897 assert!(ensure_server_frame_required_capability(
4898 &reply,
4899 &[managed.clone(), worktree.clone()],
4900 )
4901 .is_err());
4902 assert!(ensure_server_frame_required_capability(
4906 &reply,
4907 &[managed, worktree, spawn, revision],
4908 )
4909 .is_ok());
4910 }
4911
4912 #[test]
4913 fn worktree_service_mode_metadata_requires_managed_worktree_capabilities() {
4914 let managed = CapabilityId::new(NODE_MANAGED_WORKTREE_LIFECYCLE_CAPABILITY).unwrap();
4915 let worktree = CapabilityId::new(NODE_WORKTREE_SELECTION_CAPABILITY).unwrap();
4916 let mut workspace = workspace_with_path(utf8_path());
4917 workspace.worktree_service_mode = Some(WorktreeServiceMode::Manual);
4918 let reply = response_frame(NodeResponse::WorkspaceRegistered {
4919 workspace: workspace.clone(),
4920 });
4921
4922 assert!(ensure_server_frame_required_capability(&reply, &[]).is_err());
4923 assert!(ensure_server_frame_required_capability(
4924 &reply,
4925 std::slice::from_ref(&managed),
4926 )
4927 .is_err());
4928 assert!(ensure_server_frame_required_capability(
4929 &reply,
4930 std::slice::from_ref(&worktree),
4931 )
4932 .is_err());
4933 assert!(ensure_server_frame_required_capability(
4934 &reply,
4935 &[managed.clone(), worktree.clone()],
4936 )
4937 .is_ok());
4938
4939 let event = ServerFrame::Event(NodeEventEnvelope {
4940 sequence: 5,
4941 event: NodeEvent::WorkspaceAdded { workspace },
4942 });
4943 assert!(ensure_server_frame_required_capability(&event, &[]).is_err());
4944 assert!(ensure_server_frame_required_capability(
4945 &event,
4946 &[managed.clone(), worktree.clone()],
4947 )
4948 .is_ok());
4949
4950 let mut snapshot = empty_snapshot();
4951 let mut profile_only = workspace_with_path(utf8_path());
4952 profile_only.managed_worktree_profiles = Some(
4953 gate4agent_node_protocol::WorktreeProfileInventory {
4954 profiles: Vec::new(),
4955 },
4956 );
4957 snapshot.workspaces.push(profile_only);
4958 let hello = ServerFrame::Hello(hello_with_snapshot(snapshot));
4959 assert!(ensure_server_frame_required_capability(&hello, &[]).is_err());
4960 assert!(ensure_server_frame_required_capability(
4961 &hello,
4962 &[managed.clone(), worktree.clone()],
4963 )
4964 .is_ok());
4965
4966 let legacy = response_frame(NodeResponse::WorkspaceRegistered {
4967 workspace: workspace_with_path(utf8_path()),
4968 });
4969 assert!(ensure_server_frame_required_capability(&legacy, &[]).is_ok());
4970 }
4971
4972 #[test]
4973 fn malicious_legacy_hello_with_unix_path_is_rejected_before_exposure() {
4974 let mut snapshot = empty_snapshot();
4975 snapshot.workspaces.push(workspace_with_path(unix_path()));
4976 let hello = hello_with_snapshot(snapshot);
4977
4978 assert!(ensure_node_hello_path_capability(&hello, false).is_err());
4979 assert!(ensure_server_frame_path_capability(
4980 &ServerFrame::Hello(hello),
4981 false,
4982 false,
4983 )
4984 .is_err());
4985 }
4986
4987 #[test]
4988 fn legacy_utf8_hello_and_payloads_remain_accepted() {
4989 let mut snapshot = empty_snapshot();
4990 snapshot.workspaces.push(workspace_with_path(utf8_path()));
4991 let hello = hello_with_snapshot(snapshot);
4992
4993 assert!(ensure_node_hello_path_capability(&hello, false).is_ok());
4994 assert!(ensure_node_request_path_capability(
4995 &NodeRequest::RegisterWorkspace {
4996 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
4997 root: utf8_path(),
4998 },
4999 false,
5000 false,
5001 )
5002 .is_ok());
5003 assert!(ensure_server_frame_path_capability(
5004 &response_frame(NodeResponse::WorkspaceRegistered {
5005 workspace: workspace_with_path(utf8_path()),
5006 }),
5007 false,
5008 false,
5009 )
5010 .is_ok());
5011 }
5012
5013 #[test]
5014 fn outbound_guard_covers_every_path_bearing_request_variant() {
5015 let requests = [
5016 NodeRequest::RegisterWorkspace {
5017 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5018 root: unix_path(),
5019 },
5020 NodeRequest::CreateWorktree {
5021 source_workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5022 workspace_id: WorkspaceId::new("workspace-b").unwrap(),
5023 target_root: unix_path(),
5024 branch: "feature/a".to_owned(),
5025 base: None,
5026 },
5027 NodeRequest::RemoveWorktree {
5028 source_workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5029 target_root: unix_path(),
5030 },
5031 ];
5032
5033 for request in requests {
5034 assert!(ensure_node_request_path_capability(&request, false, false).is_err());
5035 assert!(ensure_node_request_path_capability(&request, true, false).is_ok());
5036 }
5037 assert!(ensure_node_request_path_capability(&NodeRequest::Snapshot, false, false).is_ok());
5038
5039 let tagged_file_read = NodeRequest::ReadWorkspaceFile {
5040 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5041 path: tagged_repository_path(b"src/\xff"),
5042 };
5043 assert!(ensure_node_request_path_capability(&tagged_file_read, false, false).is_err());
5044 assert!(ensure_node_request_path_capability(&tagged_file_read, false, true).is_ok());
5045 }
5046
5047 #[test]
5048 fn inbound_guard_covers_path_bearing_response_variants() {
5049 let mut snapshot = empty_snapshot();
5050 snapshot.workspaces.push(workspace_with_path(unix_path()));
5051 let inspection = WorkspaceInspection {
5052 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5053 entries: Vec::new(),
5054 tree_truncated: false,
5055 git: GitSnapshot {
5056 is_repository: true,
5057 branch: Some("main".to_owned()),
5058 status: Vec::new(),
5059 recent_commits: Vec::new(),
5060 worktrees: vec![worktree_with_path(unix_path())],
5061 managed_worktree: None,
5062 truncated: false,
5063 diagnostic: None,
5064 },
5065 truncation: None,
5066 };
5067 let responses = vec![
5068 NodeResponse::Snapshot {
5069 event_sequence: 1,
5070 controller: None,
5071 snapshot: snapshot.clone(),
5072 },
5073 NodeResponse::Resync {
5074 event_sequence: 1,
5075 oldest_available_sequence: 1,
5076 snapshot: empty_snapshot(),
5077 events: vec![NodeEventEnvelope {
5078 sequence: 1,
5079 event: NodeEvent::WorkspaceAdded {
5080 workspace: workspace_with_path(unix_path()),
5081 },
5082 }],
5083 },
5084 NodeResponse::WorkspaceInspected { inspection },
5085 NodeResponse::SessionRecordUpdated {
5086 record: session_record_with_path(unix_path()),
5087 },
5088 NodeResponse::WorkspaceRegistered {
5089 workspace: workspace_with_path(unix_path()),
5090 },
5091 NodeResponse::WorktreeCreated {
5092 worktree: worktree_with_path(unix_path()),
5093 workspace: workspace_with_path(utf8_path()),
5094 },
5095 NodeResponse::WorktreeRemoved {
5096 target_root: unix_path(),
5097 workspace_id: None,
5098 },
5099 ];
5100
5101 for response in responses {
5102 let frame = response_frame(response);
5103 assert!(ensure_server_frame_path_capability(&frame, false, false).is_err());
5104 assert!(ensure_server_frame_path_capability(&frame, true, false).is_ok());
5105 }
5106 }
5107
5108 #[test]
5109 fn inbound_guard_covers_path_bearing_event_variants() {
5110 let events = [
5111 NodeEvent::WorkspaceAdded {
5112 workspace: workspace_with_path(unix_path()),
5113 },
5114 NodeEvent::SessionRecordUpserted {
5115 record: session_record_with_path(unix_path()),
5116 },
5117 ];
5118
5119 for event in events {
5120 let frame = ServerFrame::Event(NodeEventEnvelope {
5121 sequence: 1,
5122 event,
5123 });
5124 assert!(ensure_server_frame_path_capability(&frame, false, false).is_err());
5125 assert!(ensure_server_frame_path_capability(&frame, true, false).is_ok());
5126 }
5127 }
5128
5129 #[test]
5130 fn inbound_guard_covers_every_tagged_repository_path_location() {
5131 let inspections = [
5132 WorkspaceInspection {
5133 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5134 entries: vec![WorkspaceEntry {
5135 relative_path: tagged_repository_path(b"src/\xff"),
5136 kind: WorkspaceEntryKind::File,
5137 }],
5138 tree_truncated: false,
5139 git: GitSnapshot {
5140 is_repository: false,
5141 branch: None,
5142 status: Vec::new(),
5143 recent_commits: Vec::new(),
5144 worktrees: Vec::new(),
5145 managed_worktree: None,
5146 truncated: false,
5147 diagnostic: None,
5148 },
5149 truncation: None,
5150 },
5151 WorkspaceInspection {
5152 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5153 entries: Vec::new(),
5154 tree_truncated: false,
5155 git: GitSnapshot {
5156 is_repository: true,
5157 branch: None,
5158 status: vec![GitStatusEntry {
5159 index_status: "M".to_owned(),
5160 worktree_status: " ".to_owned(),
5161 path: tagged_repository_path(b"src/\xff"),
5162 previous_path: None,
5163 }],
5164 recent_commits: Vec::new(),
5165 worktrees: Vec::new(),
5166 managed_worktree: None,
5167 truncated: false,
5168 diagnostic: None,
5169 },
5170 truncation: None,
5171 },
5172 WorkspaceInspection {
5173 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5174 entries: Vec::new(),
5175 tree_truncated: false,
5176 git: GitSnapshot {
5177 is_repository: true,
5178 branch: None,
5179 status: vec![GitStatusEntry {
5180 index_status: "R".to_owned(),
5181 worktree_status: " ".to_owned(),
5182 path: utf8_repository_path("src/new.rs"),
5183 previous_path: Some(tagged_repository_path(b"src/\xff")),
5184 }],
5185 recent_commits: Vec::new(),
5186 worktrees: Vec::new(),
5187 managed_worktree: None,
5188 truncated: false,
5189 diagnostic: None,
5190 },
5191 truncation: None,
5192 },
5193 ];
5194
5195 for inspection in inspections {
5196 let frame = response_frame(NodeResponse::WorkspaceInspected { inspection });
5197 assert!(ensure_server_frame_path_capability(&frame, false, false).is_err());
5198 assert!(ensure_server_frame_path_capability(&frame, false, true).is_ok());
5199 }
5200
5201 let file_frame = response_frame(NodeResponse::WorkspaceFileRead {
5202 file: WorkspaceFileRead {
5203 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5204 path: tagged_repository_path(b"src/\xff"),
5205 content: WorkspaceFileContent::NonUtf8 { byte_len: 3 },
5206 revision: None,
5207 },
5208 });
5209 assert!(ensure_server_frame_path_capability(&file_frame, false, false).is_err());
5210 assert!(ensure_server_frame_path_capability(&file_frame, false, true).is_ok());
5211 }
5212
5213 #[test]
5214 fn legacy_utf8_repository_paths_remain_accepted_without_capability() {
5215 let inspection = WorkspaceInspection {
5216 workspace_id: WorkspaceId::new("workspace-a").unwrap(),
5217 entries: vec![WorkspaceEntry {
5218 relative_path: utf8_repository_path("src/lib.rs"),
5219 kind: WorkspaceEntryKind::File,
5220 }],
5221 tree_truncated: false,
5222 git: GitSnapshot {
5223 is_repository: true,
5224 branch: None,
5225 status: vec![GitStatusEntry {
5226 index_status: "R".to_owned(),
5227 worktree_status: " ".to_owned(),
5228 path: utf8_repository_path("src/new.rs"),
5229 previous_path: Some(utf8_repository_path("src/old.rs")),
5230 }],
5231 recent_commits: Vec::new(),
5232 worktrees: Vec::new(),
5233 managed_worktree: None,
5234 truncated: false,
5235 diagnostic: None,
5236 },
5237 truncation: None,
5238 };
5239 let frame = response_frame(NodeResponse::WorkspaceInspected { inspection });
5240
5241 assert!(ensure_server_frame_path_capability(&frame, false, false).is_ok());
5242 }
5243
5244 #[test]
5245 fn delivery_wire_capability_and_exact_reply_validation() {
5246 use gate4agent_node_protocol::{
5247 DeliveryBlobChunkHexV1, DeliveryBlobDigestV1, DeliveryManifestDigestV2,
5248 DeliveryStageId,
5249 };
5250
5251 let stage_id = DeliveryStageId::new(format!(
5252 "delivery-stage-{}",
5253 "1".repeat(32),
5254 ))
5255 .unwrap();
5256 let digest = DeliveryBlobDigestV1::new(format!("sha256:{}", "a".repeat(64)))
5257 .unwrap();
5258 let request = NodeRequest::PutDeliveryBlobChunk {
5259 stage_id: stage_id.clone(),
5260 blob_digest: digest.clone(),
5261 offset: 7,
5262 chunk_hex: DeliveryBlobChunkHexV1::new("00ff").unwrap(),
5263 };
5264 let mut next_request_id = 1;
5265 assert!(matches!(
5266 reserve_request_id(&mut next_request_id, &request, false, false, false, &[]),
5267 Err(NodeClientError::UnsupportedCapability(capability))
5268 if capability == NODE_DELIVERY_BUNDLE_V2_STAGE_COMMIT_CAPABILITY
5269 ));
5270 let capability = CapabilityId::new(
5271 NODE_DELIVERY_BUNDLE_V2_STAGE_COMMIT_CAPABILITY,
5272 )
5273 .unwrap();
5274 assert!(reserve_request_id(
5275 &mut next_request_id,
5276 &request,
5277 false,
5278 false,
5279 false,
5280 std::slice::from_ref(&capability),
5281 )
5282 .is_ok());
5283
5284 let exact = Ok(NodeResponse::DeliveryBlobChunkAccepted {
5285 stage_id: stage_id.clone(),
5286 blob_digest: digest.clone(),
5287 next_offset: 9,
5288 });
5289 assert!(validate_delivery_response(&request, &exact).is_ok());
5290 let wrong_offset = Ok(NodeResponse::DeliveryBlobChunkAccepted {
5291 stage_id: stage_id.clone(),
5292 blob_digest: digest.clone(),
5293 next_offset: 8,
5294 });
5295 assert!(validate_delivery_response(&request, &wrong_offset).is_err());
5296 let unexpected = Ok(NodeResponse::DeliveryStageBegun {
5297 stage_id,
5298 manifest_digest: DeliveryManifestDigestV2::new(format!(
5299 "sha256:{}",
5300 "b".repeat(64),
5301 ))
5302 .unwrap(),
5303 missing_blobs: vec![digest],
5304 });
5305 assert!(validate_delivery_response(&request, &unexpected).is_err());
5306 assert!(ensure_server_frame_required_capability(
5307 &response_frame(unexpected.unwrap()),
5308 &[],
5309 )
5310 .is_err());
5311 }
5312
5313}