Skip to main content

aft/subc/
wire.rs

1//! Frame encoding and writer-queue helpers used by the subc transport edge.
2
3use serde::ser::{SerializeMap, SerializeStruct};
4use serde::{Serialize, Serializer};
5
6use super::{
7    control_flags, fmt, mpsc, Arc, AtomicUsize, BindTrust, DispatchPathMetrics, ErrorBody, Flags,
8    Frame, FrameType, Ordering, PathBuf, Response, RouteChannel, ToolCallResult, Value,
9    CONTROL_SEND_TIMEOUT, RELIABLE_WRITER_RETRY_INITIAL_BACKOFF, RELIABLE_WRITER_RETRY_MAX_BACKOFF,
10};
11use crate::run_tool_call::{PhaseTrace, ToolCallEgressTiming, ToolCallPhaseDurations};
12
13pub(super) type WriterSender = mpsc::Sender<WriterFrame>;
14
15pub(super) struct ToolResponseWriteTrace {
16    phase_trace: PhaseTrace,
17    name: String,
18    root: PathBuf,
19    session: String,
20    channel: u16,
21    corr: u64,
22    enqueued_at: Option<std::time::Instant>,
23    queue_depth: usize,
24    writer_active_at_enqueue: bool,
25    writer_queue_was_full: bool,
26    reserve_timeouts: u32,
27}
28
29impl ToolResponseWriteTrace {
30    pub(super) fn new(
31        phase_trace: PhaseTrace,
32        name: String,
33        root: PathBuf,
34        session: String,
35        channel: u16,
36        corr: u64,
37    ) -> Self {
38        Self {
39            phase_trace,
40            name,
41            root,
42            session,
43            channel,
44            corr,
45            enqueued_at: None,
46            queue_depth: 0,
47            writer_active_at_enqueue: false,
48            writer_queue_was_full: false,
49            reserve_timeouts: 0,
50        }
51    }
52
53    fn mark_writer_queue_full(&mut self) {
54        self.writer_queue_was_full = true;
55    }
56
57    fn mark_reserve_timeout(&mut self) {
58        self.reserve_timeouts = self.reserve_timeouts.saturating_add(1);
59    }
60
61    fn mark_enqueued(&mut self, queue_depth: usize, writer_active: bool) {
62        self.enqueued_at = Some(std::time::Instant::now());
63        self.queue_depth = queue_depth;
64        self.writer_active_at_enqueue = writer_active;
65    }
66
67    pub(super) fn finish(
68        self,
69        dequeued: std::time::Instant,
70        write_started: std::time::Instant,
71        write_finished: std::time::Instant,
72        frame_bytes: usize,
73    ) -> Option<CompletedToolResponseTrace> {
74        let phases = self.phase_trace.finish(ToolCallEgressTiming {
75            enqueued: self.enqueued_at?,
76            dequeued,
77            write_started,
78            write_finished,
79            frame_bytes,
80            queue_depth: self.queue_depth,
81            writer_active_at_enqueue: self.writer_active_at_enqueue,
82            writer_queue_was_full: self.writer_queue_was_full,
83            reserve_timeouts: self.reserve_timeouts,
84        })?;
85        Some(CompletedToolResponseTrace {
86            name: self.name,
87            root: self.root,
88            session: self.session,
89            channel: self.channel,
90            corr: self.corr,
91            phases,
92        })
93    }
94}
95
96pub(super) struct CompletedToolResponseTrace {
97    pub(super) name: String,
98    pub(super) root: PathBuf,
99    pub(super) session: String,
100    pub(super) channel: u16,
101    pub(super) corr: u64,
102    pub(super) phases: ToolCallPhaseDurations,
103}
104
105enum WriterFrameBody {
106    Owned,
107    SharedPush(Arc<Vec<u8>>),
108}
109
110pub(super) struct WriterFrame {
111    /// The route-specific header. Shared Push bodies remain outside this frame.
112    pub(super) frame: Frame,
113    body: WriterFrameBody,
114    pub(super) tool_response_trace: Option<ToolResponseWriteTrace>,
115}
116
117impl std::ops::Deref for WriterFrame {
118    type Target = Frame;
119
120    fn deref(&self) -> &Self::Target {
121        &self.frame
122    }
123}
124
125impl WriterFrame {
126    pub(super) fn plain(frame: Frame) -> Self {
127        Self {
128            frame,
129            body: WriterFrameBody::Owned,
130            tool_response_trace: None,
131        }
132    }
133
134    pub(super) fn shared_push(frame: Frame, body: Arc<Vec<u8>>) -> Self {
135        debug_assert_eq!(frame.header.len as usize, body.len());
136        Self {
137            frame,
138            body: WriterFrameBody::SharedPush(body),
139            tool_response_trace: None,
140        }
141    }
142
143    fn traced_tool_response(frame: Frame, trace: ToolResponseWriteTrace) -> Self {
144        Self {
145            frame,
146            body: WriterFrameBody::Owned,
147            tool_response_trace: Some(trace),
148        }
149    }
150
151    pub(super) fn frame(&self) -> &Frame {
152        &self.frame
153    }
154
155    pub(super) fn body(&self) -> &[u8] {
156        match &self.body {
157            WriterFrameBody::Owned => &self.frame.body,
158            WriterFrameBody::SharedPush(body) => body,
159        }
160    }
161
162    #[cfg(test)]
163    pub(super) fn shared_push_body_strong_count(&self) -> Option<usize> {
164        match &self.body {
165            WriterFrameBody::Owned => None,
166            WriterFrameBody::SharedPush(body) => Some(Arc::strong_count(body)),
167        }
168    }
169
170    fn mark_writer_queue_full(&mut self) {
171        if let Some(trace) = self.tool_response_trace.as_mut() {
172            trace.mark_writer_queue_full();
173        }
174    }
175
176    fn mark_reserve_timeout(&mut self) {
177        if let Some(trace) = self.tool_response_trace.as_mut() {
178            trace.mark_reserve_timeout();
179        }
180    }
181
182    fn mark_enqueued(&mut self, queue_depth: usize, writer_active: bool) {
183        if let Some(trace) = self.tool_response_trace.as_mut() {
184            trace.mark_enqueued(queue_depth, writer_active);
185        }
186    }
187}
188
189pub(super) enum WriterEnqueueOutcome {
190    Enqueued,
191    Full(WriterFrame),
192    Closed,
193}
194
195impl WriterEnqueueOutcome {
196    #[cfg(test)]
197    pub(super) fn is_enqueued(&self) -> bool {
198        matches!(self, Self::Enqueued)
199    }
200}
201
202pub(super) fn decrement_counted_channel(counter: &AtomicUsize) {
203    let previous = counter.fetch_sub(1, Ordering::Relaxed);
204    debug_assert!(previous > 0, "counted channel depth underflow");
205}
206
207pub(super) async fn send_counted_channel<T>(
208    tx: &mpsc::Sender<T>,
209    counter: &AtomicUsize,
210    item: T,
211) -> Result<(), mpsc::error::SendError<T>> {
212    counter.fetch_add(1, Ordering::Relaxed);
213    match tx.send(item).await {
214        Ok(()) => Ok(()),
215        Err(error) => {
216            decrement_counted_channel(counter);
217            Err(error)
218        }
219    }
220}
221
222fn enqueue_writer_item(
223    permit: mpsc::Permit<'_, WriterFrame>,
224    metrics: &DispatchPathMetrics,
225    mut item: WriterFrame,
226) {
227    let queue_depth = metrics.writer_queued.fetch_add(1, Ordering::Relaxed) + 1;
228    item.mark_enqueued(queue_depth, metrics.writer_active.load(Ordering::Relaxed));
229    permit.send(item);
230}
231
232fn try_enqueue_writer_item(
233    tx: &WriterSender,
234    metrics: &DispatchPathMetrics,
235    mut item: WriterFrame,
236) -> WriterEnqueueOutcome {
237    match tx.try_reserve() {
238        Ok(permit) => {
239            enqueue_writer_item(permit, metrics, item);
240            WriterEnqueueOutcome::Enqueued
241        }
242        Err(mpsc::error::TrySendError::Full(())) => {
243            metrics
244                .writer_saturation_count
245                .fetch_add(1, Ordering::Relaxed);
246            item.mark_writer_queue_full();
247            WriterEnqueueOutcome::Full(item)
248        }
249        Err(mpsc::error::TrySendError::Closed(())) => {
250            drop(item);
251            WriterEnqueueOutcome::Closed
252        }
253    }
254}
255
256pub(super) fn try_enqueue_writer_frame(
257    tx: &WriterSender,
258    metrics: &DispatchPathMetrics,
259    frame: Frame,
260) -> WriterEnqueueOutcome {
261    try_enqueue_writer_item(tx, metrics, WriterFrame::plain(frame))
262}
263
264pub(super) fn try_enqueue_shared_push_frame(
265    tx: &WriterSender,
266    metrics: &DispatchPathMetrics,
267    frame: Frame,
268    body: Arc<Vec<u8>>,
269) -> WriterEnqueueOutcome {
270    try_enqueue_writer_item(tx, metrics, WriterFrame::shared_push(frame, body))
271}
272
273async fn send_reliable_writer_item(
274    tx: &WriterSender,
275    metrics: &DispatchPathMetrics,
276    mut item: WriterFrame,
277    context: &'static str,
278) -> Result<(), SubcError> {
279    let mut warned = false;
280    let mut backoff = RELIABLE_WRITER_RETRY_INITIAL_BACKOFF;
281
282    loop {
283        match try_enqueue_writer_item(tx, metrics, item) {
284            WriterEnqueueOutcome::Enqueued => return Ok(()),
285            WriterEnqueueOutcome::Closed => return Err(SubcError::WriterClosed),
286            WriterEnqueueOutcome::Full(returned_item) => {
287                item = returned_item;
288            }
289        }
290
291        match tokio::time::timeout(CONTROL_SEND_TIMEOUT, tx.reserve()).await {
292            Ok(Ok(permit)) => {
293                enqueue_writer_item(permit, metrics, item);
294                return Ok(());
295            }
296            Ok(Err(_)) => return Err(SubcError::WriterClosed),
297            Err(_) => {
298                metrics
299                    .writer_saturation_count
300                    .fetch_add(1, Ordering::Relaxed);
301                item.mark_reserve_timeout();
302                if !warned {
303                    log::warn!(
304                        "subc attach: writer queue stayed full while sending {context}; retrying reliable frame"
305                    );
306                    warned = true;
307                }
308                tokio::time::sleep(backoff).await;
309                backoff =
310                    std::cmp::min(backoff.saturating_mul(2), RELIABLE_WRITER_RETRY_MAX_BACKOFF);
311            }
312        }
313    }
314}
315
316pub(super) async fn send_reliable_writer_frame(
317    tx: &WriterSender,
318    metrics: &DispatchPathMetrics,
319    frame: Frame,
320    context: &'static str,
321) -> Result<(), SubcError> {
322    send_reliable_writer_item(tx, metrics, WriterFrame::plain(frame), context).await
323}
324
325pub(super) async fn send_traced_tool_response_frame(
326    tx: &WriterSender,
327    metrics: &DispatchPathMetrics,
328    frame: Frame,
329    trace: ToolResponseWriteTrace,
330) -> Result<(), SubcError> {
331    send_reliable_writer_item(
332        tx,
333        metrics,
334        WriterFrame::traced_tool_response(frame, trace),
335        "tool response",
336    )
337    .await
338}
339
340pub(super) async fn send_frame(
341    tx: &WriterSender,
342    metrics: &DispatchPathMetrics,
343    frame: Frame,
344) -> Result<(), SubcError> {
345    match try_enqueue_writer_item(tx, metrics, WriterFrame::plain(frame)) {
346        WriterEnqueueOutcome::Enqueued => Ok(()),
347        WriterEnqueueOutcome::Closed => Err(SubcError::WriterClosed),
348        WriterEnqueueOutcome::Full(item) => {
349            match tokio::time::timeout(CONTROL_SEND_TIMEOUT, tx.reserve()).await {
350                Ok(Ok(permit)) => {
351                    enqueue_writer_item(permit, metrics, item);
352                    Ok(())
353                }
354                Ok(Err(_)) => Err(SubcError::WriterClosed),
355                Err(_) => {
356                    metrics
357                        .writer_saturation_count
358                        .fetch_add(1, Ordering::Relaxed);
359                    Err(SubcError::WriterBackpressureTimeout)
360                }
361            }
362        }
363    }
364}
365
366/// Borrowed flat response matching the standalone NDJSON shape without cloning
367/// the response id or any structured data values.
368struct FlatToolResponse<'a> {
369    response: &'a crate::protocol::Response,
370    text: &'a str,
371}
372
373impl Serialize for FlatToolResponse<'_> {
374    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
375    where
376        S: Serializer,
377    {
378        let data = self.response.data.as_object();
379        let has_text = data.is_some_and(|data| data.contains_key("text"));
380        let field_count =
381            2 + data.map_or(0, |data| {
382                data.len()
383                    - usize::from(data.contains_key("id"))
384                    - usize::from(data.contains_key("success"))
385            }) + usize::from(!has_text);
386        let mut map = serializer.serialize_map(Some(field_count))?;
387        match data.and_then(|data| data.get("id")) {
388            Some(value) => map.serialize_entry("id", value)?,
389            None => map.serialize_entry("id", &self.response.id)?,
390        }
391        match data.and_then(|data| data.get("success")) {
392            Some(value) => map.serialize_entry("success", value)?,
393            None => map.serialize_entry("success", &self.response.success)?,
394        }
395        if let Some(data) = data {
396            for (key, value) in data {
397                match key.as_str() {
398                    "id" | "success" => {}
399                    "text" => map.serialize_entry(key, self.text)?,
400                    _ => map.serialize_entry(key, value)?,
401                }
402            }
403        }
404        if !has_text {
405            map.serialize_entry("text", self.text)?;
406        }
407        map.end()
408    }
409}
410
411struct ToolResponseEnvelope<'a> {
412    result: &'a ToolCallResult,
413    /// First-party binds get the full flat response in `structuredContent`
414    /// for the plugin re-lift; untrusted (MCP) binds get text-only.
415    include_structured: bool,
416}
417
418impl Serialize for ToolResponseEnvelope<'_> {
419    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
420    where
421        S: Serializer,
422    {
423        let fields = if self.include_structured { 3 } else { 2 };
424        let mut envelope = serializer.serialize_struct("ToolResponseEnvelope", fields)?;
425        envelope.serialize_field(
426            "content",
427            &[TextContent {
428                kind: "text",
429                text: &self.result.text,
430            }],
431        )?;
432        envelope.serialize_field("isError", &!self.result.response.success)?;
433        if self.include_structured {
434            envelope.serialize_field(
435                "structuredContent",
436                &FlatToolResponse {
437                    response: &self.result.response,
438                    text: &self.result.text,
439                },
440            )?;
441        }
442        envelope.end()
443    }
444}
445
446#[derive(Serialize)]
447struct TextContent<'a> {
448    #[serde(rename = "type")]
449    kind: &'static str,
450    text: &'a str,
451}
452
453pub(super) fn build_tool_response_frame(
454    ver: u8,
455    route: RouteChannel,
456    corr: u64,
457    flags: Flags,
458    result: &ToolCallResult,
459    trust: BindTrust,
460) -> Result<Frame, SubcError> {
461    // `content`/`isError` is the MCP-native surface a GENERIC host reads. The
462    // FIRST-PARTY AFT plugin instead reads `structuredContent`, which carries
463    // the full flat standalone shape ({id, success, ...data, text}) so every
464    // structured sidecar the plugin drives UI from — status_bar, bg_completions
465    // (in-band drain), preview_diff, code, message, attachments — survives the
466    // route. subc relays the body byte-for-byte, so this reaches the plugin
467    // unchanged. SubcTransport.toolCall re-lifts `structuredContent` straight to
468    // the flat ToolCallResult, so nothing downstream of the transport differs
469    // from the NDJSON path.
470    //
471    // UNTRUSTED binds (MCP hosts via subc-mcp) get text-only replies: they
472    // have no re-lift layer, we declare no outputSchema (so omitting is
473    // MCP-spec-clean), and hosts like Claude Code prefer `structuredContent`
474    // for model input when present — feeding the model a raw JSON dump with
475    // the rendered text buried inside it, at a multiple of the token cost.
476    let include_structured = !matches!(trust, BindTrust::Untrusted);
477    let body = serde_json::to_vec(&ToolResponseEnvelope {
478        result,
479        include_structured,
480    })
481    .map_err(SubcError::Json)?;
482
483    Frame::build_with_version(
484        ver,
485        FrameType::Response,
486        flags,
487        route.channel,
488        route.epoch,
489        corr,
490        body,
491    )
492    .map_err(SubcError::FrameBuild)
493}
494
495pub(super) fn build_error_frame(
496    ver: u8,
497    channel: u16,
498    epoch: u32,
499    corr: u64,
500    flags: Flags,
501    code: &str,
502    message: &str,
503) -> Result<Frame, SubcError> {
504    let body = serde_json::to_vec(&ErrorBody {
505        code: code.to_string(),
506        message: message.to_string(),
507    })
508    .map_err(SubcError::Json)?;
509    Frame::build_with_version(ver, FrameType::Error, flags, channel, epoch, corr, body)
510        .map_err(SubcError::FrameBuild)
511}
512
513pub(super) fn build_goodbye_frame(
514    ver: u8,
515    channel: u16,
516    epoch: u32,
517    corr: u64,
518) -> Result<Frame, SubcError> {
519    Frame::build_with_version(
520        ver,
521        FrameType::Goodbye,
522        control_flags(),
523        channel,
524        epoch,
525        corr,
526        Vec::new(),
527    )
528    .map_err(SubcError::FrameBuild)
529}
530
531pub(super) fn response_message(response: &Response, fallback: &str) -> String {
532    response
533        .data
534        .get("message")
535        .and_then(Value::as_str)
536        .map(ToOwned::to_owned)
537        .unwrap_or_else(|| fallback.to_string())
538}
539
540pub(super) fn response_is_fatal_panic(response: &Response) -> bool {
541    !response.success && response.data.get("code").and_then(Value::as_str) == Some("actor_fatal")
542}
543
544#[derive(Debug)]
545pub enum SubcError {
546    Runtime(std::io::Error),
547    ConnectionFile {
548        path: PathBuf,
549        source: subc_transport::ConnectionFileError,
550    },
551    NoEndpoint {
552        path: PathBuf,
553    },
554    InvalidEndpoint {
555        path: PathBuf,
556        endpoint: String,
557    },
558    Connect {
559        endpoint: String,
560        source: std::io::Error,
561    },
562    Auth {
563        endpoint: String,
564        source: subc_transport::AuthError,
565    },
566    FrameIo(subc_transport::FrameIoError),
567    FrameBuild(subc_protocol::FrameBuildError),
568    WriterClosed,
569    WriterBackpressureTimeout,
570    WriterJoin(tokio::task::JoinError),
571    Json(serde_json::Error),
572    ClosedBeforeHelloAck,
573    HelloRejected {
574        body: Option<ErrorBody>,
575    },
576    UnexpectedFrame {
577        ty: FrameType,
578    },
579}
580
581impl fmt::Display for SubcError {
582    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
583        match self {
584            Self::Runtime(e) => write!(f, "failed to build subc tokio runtime: {e}"),
585            Self::ConnectionFile { path, source } => {
586                write!(f, "failed to read subc connection file {path:?}: {source}")
587            }
588            Self::NoEndpoint { path } => {
589                write!(f, "subc connection file {path:?} has no endpoints")
590            }
591            Self::InvalidEndpoint { path, endpoint } => {
592                write!(
593                    f,
594                    "subc connection file {path:?} has invalid endpoint {endpoint}"
595                )
596            }
597            Self::Connect { endpoint, source } => {
598                write!(f, "failed to connect to subc endpoint {endpoint}: {source}")
599            }
600            Self::Auth { endpoint, source } => {
601                write!(
602                    f,
603                    "failed to authenticate to subc endpoint {endpoint}: {source}"
604                )
605            }
606            Self::FrameIo(e) => write!(f, "subc frame I/O error: {e}"),
607            Self::FrameBuild(e) => write!(f, "subc frame build error: {e}"),
608            Self::WriterClosed => write!(f, "subc writer task closed"),
609            Self::WriterBackpressureTimeout => write!(
610                f,
611                "subc writer task stayed backpressured while sending a control frame"
612            ),
613            Self::WriterJoin(e) => write!(f, "subc writer task join error: {e}"),
614            Self::Json(e) => write!(f, "subc JSON error: {e}"),
615            Self::ClosedBeforeHelloAck => {
616                write!(f, "subc daemon closed the connection before HelloAck")
617            }
618            Self::HelloRejected { body } => match body {
619                Some(b) => write!(f, "subc rejected ModuleHello: {} ({})", b.code, b.message),
620                None => write!(f, "subc rejected ModuleHello (unparseable error body)"),
621            },
622            Self::UnexpectedFrame { ty } => {
623                write!(f, "subc sent unexpected frame in place of HelloAck: {ty:?}")
624            }
625        }
626    }
627}
628
629impl std::error::Error for SubcError {}
630
631#[cfg(test)]
632mod tests {
633    use super::*;
634    use crate::subc::route_key;
635    use serde_json::json;
636    use std::sync::Arc;
637    use std::time::{Duration, Instant};
638    use subc_protocol::PROTOCOL_VERSION;
639
640    #[test]
641    fn writer_depth_counter_tracks_enqueued_frames_until_drain() {
642        let metrics = DispatchPathMetrics::new();
643        let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(8);
644
645        for corr in 1..=3 {
646            let frame = Frame::build(FrameType::Ping, control_flags(), 0, 0, corr, Vec::new())
647                .expect("test frame");
648            assert!(try_enqueue_writer_frame(&writer_tx, &metrics, frame).is_enqueued());
649        }
650        assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 3);
651
652        for _ in 0..3 {
653            writer_rx.try_recv().expect("queued writer frame");
654            decrement_counted_channel(&metrics.writer_queued);
655        }
656        assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 0);
657    }
658
659    #[tokio::test]
660    async fn reliable_writer_send_retries_after_timeout_and_preserves_frame() {
661        let metrics = Arc::new(DispatchPathMetrics::new());
662        let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(1);
663        writer_tx
664            .try_send(WriterFrame::plain(
665                Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
666            ))
667            .expect("prefill writer queue");
668
669        let metrics_for_task = Arc::clone(&metrics);
670        let tx_for_task = writer_tx.clone();
671        let send_task = tokio::spawn(async move {
672            send_reliable_writer_frame(
673                &tx_for_task,
674                &metrics_for_task,
675                Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
676                "test reliable frame",
677            )
678            .await
679        });
680
681        tokio::time::timeout(Duration::from_secs(2), async {
682            while metrics.writer_saturation_count.load(Ordering::Relaxed) < 2 {
683                tokio::time::sleep(Duration::from_millis(10)).await;
684            }
685        })
686        .await
687        .expect("reliable send should observe a timed-out full writer queue");
688
689        let prefilled = writer_rx.recv().await.expect("prefilled frame");
690        assert_eq!(prefilled.header.corr, 1);
691        let result = tokio::time::timeout(Duration::from_secs(2), send_task)
692            .await
693            .expect("reliable send should finish after writer drains")
694            .expect("reliable send task should not panic");
695        assert!(result.is_ok());
696        let delivered = writer_rx.recv().await.expect("retried reliable frame");
697        assert_eq!(delivered.header.corr, 2);
698    }
699
700    #[test]
701    fn response_is_fatal_panic_only_matches_panic_exclusive_code() {
702        let tool_error = Response::error("request-1", "internal_error", "ordinary tool error");
703        let panic_error = Response::error("request-2", "actor_fatal", "mutating panic");
704
705        assert!(!response_is_fatal_panic(&tool_error));
706        assert!(response_is_fatal_panic(&panic_error));
707    }
708
709    #[tokio::test]
710    async fn control_send_times_out_when_writer_queue_remains_full() {
711        let (writer_tx, _writer_rx) = mpsc::channel::<WriterFrame>(1);
712        let metrics = DispatchPathMetrics::new();
713        writer_tx
714            .try_send(WriterFrame::plain(
715                Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
716            ))
717            .expect("prefill writer queue");
718        let started = Instant::now();
719
720        let result = send_frame(
721            &writer_tx,
722            &metrics,
723            Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
724        )
725        .await;
726
727        assert!(matches!(result, Err(SubcError::WriterBackpressureTimeout)));
728        assert!(
729            started.elapsed() < Duration::from_secs(2),
730            "control send guard should be bounded"
731        );
732    }
733
734    #[test]
735    fn tool_response_frame_carries_flat_standalone_shape_in_structured_content() {
736        use crate::protocol::Response;
737
738        // A response with sidecars the FIRST-PARTY plugin drives UI from
739        // (status_bar, bg_completions, code) plus a normal result field.
740        let response = Response::success(
741            "req-7",
742            json!({
743                "complete": true,
744                "matches": 3,
745                "status_bar": { "errors": 0, "warnings": 1 },
746                "bg_completions": [{ "task_id": "bash-abc" }],
747            }),
748        );
749        let result = ToolCallResult {
750            text: "rendered text".to_string(),
751            response,
752        };
753
754        // The flat shape must equal the standalone NDJSON `tool_call` body:
755        // {id, success, ...data, text}. Build the standalone expectation the
756        // same way commands::tool_call::response_with_text does.
757        let expected_flat = json!({
758            "id": "req-7",
759            "success": true,
760            "complete": true,
761            "matches": 3,
762            "status_bar": { "errors": 0, "warnings": 1 },
763            "bg_completions": [{ "task_id": "bash-abc" }],
764            "text": "rendered text",
765        });
766        assert_eq!(
767            serde_json::to_value(FlatToolResponse {
768                response: &result.response,
769                text: &result.text,
770            })
771            .unwrap(),
772            expected_flat,
773            "structuredContent must be byte-identical to the standalone flat response"
774        );
775
776        // The frame body carries the MCP surface for generic hosts AND the flat
777        // sidecar shape under structuredContent for the first-party plugin.
778        let frame = build_tool_response_frame(
779            PROTOCOL_VERSION,
780            route_key(1, 1),
781            42,
782            control_flags(),
783            &result,
784            BindTrust::FirstParty,
785        )
786        .unwrap();
787        let expected_body = serde_json::to_vec(&json!({
788            "content": [{ "type": "text", "text": "rendered text" }],
789            "isError": false,
790            "structuredContent": expected_flat.clone(),
791        }))
792        .unwrap();
793        assert_eq!(
794            frame.body, expected_body,
795            "tool response wire bytes drifted"
796        );
797        let body: Value = serde_json::from_slice(&frame.body).unwrap();
798        assert_eq!(body["isError"], json!(false));
799        assert_eq!(body["content"][0]["type"], json!("text"));
800        assert_eq!(body["content"][0]["text"], json!("rendered text"));
801        assert_eq!(body["structuredContent"], expected_flat);
802
803        // A failed response flips isError and still carries the flat shape
804        // (with success:false + code) for the plugin's error path.
805        let err = Response::error_with_data(
806            "req-8",
807            "ambiguous_match",
808            "too many matches",
809            json!({ "candidates": ["a", "b"] }),
810        );
811        let err_result = ToolCallResult {
812            text: "error text".to_string(),
813            response: err,
814        };
815        let err_frame = build_tool_response_frame(
816            PROTOCOL_VERSION,
817            route_key(1, 1),
818            43,
819            control_flags(),
820            &err_result,
821            BindTrust::FirstParty,
822        )
823        .unwrap();
824        let err_body: Value = serde_json::from_slice(&err_frame.body).unwrap();
825        assert_eq!(err_body["isError"], json!(true));
826        assert_eq!(err_body["structuredContent"]["success"], json!(false));
827        assert_eq!(
828            err_body["structuredContent"]["code"],
829            json!("ambiguous_match")
830        );
831        assert_eq!(
832            err_body["structuredContent"]["candidates"],
833            json!(["a", "b"])
834        );
835        assert_eq!(err_body["structuredContent"]["text"], json!("error text"));
836
837        // UNTRUSTED (MCP) binds get text-only replies: no structuredContent
838        // key at all. Generic MCP hosts have no re-lift layer, and hosts like
839        // Claude Code feed structuredContent to the model verbatim when
840        // present, a raw JSON dump at a multiple of the token cost.
841        let untrusted_frame = build_tool_response_frame(
842            PROTOCOL_VERSION,
843            route_key(1, 1),
844            44,
845            control_flags(),
846            &result,
847            BindTrust::Untrusted,
848        )
849        .unwrap();
850        let untrusted_body: Value = serde_json::from_slice(&untrusted_frame.body).unwrap();
851        assert_eq!(untrusted_body["content"][0]["text"], json!("rendered text"));
852        assert_eq!(untrusted_body["isError"], json!(false));
853        assert!(
854            untrusted_body.get("structuredContent").is_none(),
855            "untrusted binds must not receive structuredContent: {untrusted_body}"
856        );
857    }
858}