Skip to main content

gate4agent_node_wire/
client.rs

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
173/// Any byte stream this wire's CLIENT side can run on.
174///
175/// Public so a caller can hand in a stream it established itself. The
176/// handshake never cared what carried it -- it reads and writes framed
177/// JSON -- but until [`LocalNodeClient::adopt`] existed the only way in
178/// was through a constructor that also opened the connection, which welded
179/// "who dials" to "who is the protocol client". Those are separate
180/// questions: a node behind NAT must dial out, and must still be the
181/// protocol's server when it gets there.
182pub 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    /// Runs the client handshake over a stream the caller already opened
321    /// -- or, on a call-home relay, one it ACCEPTED.
322    ///
323    /// Identical to what [`Self::connect`] and [`Self::connect_loopback`]
324    /// do after their own `connect` call: same hello, same mutual
325    /// challenge-response over both nonces, same compatibility
326    /// negotiation. The direction the socket was opened in is not an input
327    /// to any of it, and this constructor is the proof of that -- it is
328    /// the same code path, entered one step later.
329    ///
330    /// `expected_node_id` still has to be known up front, because the
331    /// server proof is computed from that node's own access token. On a
332    /// call-home connection the relay learns it from the node's
333    /// `NodeCallHomeAnnounce` preface, which selects a token and nothing
334    /// else -- see that type's own doc comment.
335    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                &current.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    /// Every capability `ensure_node_request_required_capability` can
2970    /// demand of a spawn request, in the order that function checks them.
2971    ///
2972    /// The tests below each assert that ONE withheld capability is the one
2973    /// reported, and that only isolates the capability under test if every
2974    /// other one is granted. Written as hand-listed pairs, they did not:
2975    /// when `requires_spawn_profile_revision_capability` became true for
2976    /// every spawn request, the gate started reporting the profile-revision
2977    /// capability first and four tests began asserting against a name they
2978    /// were not about. They stayed red.
2979    ///
2980    /// Nothing was ever unguarded -- the gate refused those requests the
2981    /// whole time, just for a different missing capability -- but a test
2982    /// that cannot say which guard fired is not testing that guard. Listed
2983    /// once here so the next capability added to the gate is added in one
2984    /// place instead of silently retargeting every assertion.
2985    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    /// The whole spawn-gate set with exactly one capability withheld --
2995    /// "a fully negotiated client, except for this".
2996    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    /// The version was in this test's NAME and pinned in its body, and
3255    /// both went stale: the offer advertises through `NODE_STATE_SCHEMA_
3256    /// V10` and the assertion still demanded V8, so it has been red since
3257    /// the schema moved. A range that grows by design cannot be pinned by
3258    /// a literal without going red on every growth.
3259    ///
3260    /// The invariant that does not move is the MINIMUM. A node that stops
3261    /// offering `NODE_STATE_SCHEMA_V1` has dropped support for every state
3262    /// file written before it, and that is a compatibility break rather
3263    /// than a version bump -- so that end is asserted against the constant
3264    /// it must never leave, and the other against the newest schema this
3265    /// crate family defines.
3266    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        // Everything but worktree selection, which the second half of this
4340        // test withholds on purpose to check the worktree guard.
4341        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        // Every spawn request needs the profile-revision capability, so it
4804        // is demanded here -- between the overrides capability and worktree
4805        // selection. This rung was missing, which is not a gap in the gate
4806        // but a gap in the ladder: the two steps below it were asserting
4807        // against the name this one reports.
4808        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        // The receipt carries a resolved profile revision, so reading it
4903        // needs that capability too -- the response gate mirrors the
4904        // request gate, and this assertion had the same missing rung.
4905        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}