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
418// A trusted envelope carries the rendered text twice — once as the outer MCP
419// `content`, once inside `structuredContent` — and for reads the raw `content`
420// data field rides along a third time, so a read body crosses the connection
421// roughly 3x. This is deliberate, not an oversight: the bridge re-lifts
422// `structuredContent` to reconstruct the flat response, and dropping the
423// duplicate is a forward-incompatible wire change (`reliftReply` rejects a
424// reply without `structuredContent.text`, so a module-first rollout breaks
425// every installed plugin). Collapsing it safely means emitting both shapes,
426// waiting for plugins to update, then dropping one behind a version floor.
427// The bridge now accepts replies that omit `structuredContent.text`, and the module may omit that field after the minimum supported plugin version includes this compatibility behavior.
428//
429// Measured on a live daemon: the largest real frames were ~200 KB, with zero
430// egress-write time, writer queue depth 1, never full, and no reserve
431// timeouts — the amplification costs nothing observable. Revisit if
432// `egress_write` on tool-call phase traces becomes nonzero, if the writer
433// queue starts backing up, or if typical frames grow well past a few hundred
434// KB; at that point the two-step migration earns its risk.
435
436impl Serialize for ToolResponseEnvelope<'_> {
437    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
438    where
439        S: Serializer,
440    {
441        let fields = if self.include_structured { 3 } else { 2 };
442        let mut envelope = serializer.serialize_struct("ToolResponseEnvelope", fields)?;
443        envelope.serialize_field(
444            "content",
445            &[TextContent {
446                kind: "text",
447                text: &self.result.text,
448            }],
449        )?;
450        envelope.serialize_field("isError", &!self.result.response.success)?;
451        if self.include_structured {
452            envelope.serialize_field(
453                "structuredContent",
454                &FlatToolResponse {
455                    response: &self.result.response,
456                    text: &self.result.text,
457                },
458            )?;
459        }
460        envelope.end()
461    }
462}
463
464#[derive(Serialize)]
465struct TextContent<'a> {
466    #[serde(rename = "type")]
467    kind: &'static str,
468    text: &'a str,
469}
470
471pub(super) fn build_tool_response_frame(
472    ver: u8,
473    route: RouteChannel,
474    corr: u64,
475    flags: Flags,
476    result: &ToolCallResult,
477    trust: BindTrust,
478) -> Result<Frame, SubcError> {
479    // `content`/`isError` is the MCP-native surface a GENERIC host reads. The
480    // FIRST-PARTY AFT plugin instead reads `structuredContent`, which carries
481    // the full flat standalone shape ({id, success, ...data, text}) so every
482    // structured sidecar the plugin drives UI from — status_bar, bg_completions
483    // (in-band drain), preview_diff, code, message, attachments — survives the
484    // route. subc relays the body byte-for-byte, so this reaches the plugin
485    // unchanged. SubcTransport.toolCall re-lifts `structuredContent` straight to
486    // the flat ToolCallResult, so nothing downstream of the transport differs
487    // from the NDJSON path.
488    //
489    // UNTRUSTED binds (MCP hosts via subc-mcp) get text-only replies: they
490    // have no re-lift layer, we declare no outputSchema (so omitting is
491    // MCP-spec-clean), and hosts like Claude Code prefer `structuredContent`
492    // for model input when present — feeding the model a raw JSON dump with
493    // the rendered text buried inside it, at a multiple of the token cost.
494    let include_structured = !matches!(trust, BindTrust::Untrusted);
495    let body = serde_json::to_vec(&ToolResponseEnvelope {
496        result,
497        include_structured,
498    })
499    .map_err(SubcError::Json)?;
500
501    Frame::build_with_version(
502        ver,
503        FrameType::Response,
504        flags,
505        route.channel,
506        route.epoch,
507        corr,
508        body,
509    )
510    .map_err(SubcError::FrameBuild)
511}
512
513pub(super) fn build_error_frame(
514    ver: u8,
515    channel: u16,
516    epoch: u32,
517    corr: u64,
518    flags: Flags,
519    code: &str,
520    message: &str,
521) -> Result<Frame, SubcError> {
522    let body = serde_json::to_vec(&ErrorBody {
523        code: code.to_string(),
524        message: message.to_string(),
525    })
526    .map_err(SubcError::Json)?;
527    Frame::build_with_version(ver, FrameType::Error, flags, channel, epoch, corr, body)
528        .map_err(SubcError::FrameBuild)
529}
530
531pub(super) fn build_goodbye_frame(
532    ver: u8,
533    channel: u16,
534    epoch: u32,
535    corr: u64,
536) -> Result<Frame, SubcError> {
537    Frame::build_with_version(
538        ver,
539        FrameType::Goodbye,
540        control_flags(),
541        channel,
542        epoch,
543        corr,
544        Vec::new(),
545    )
546    .map_err(SubcError::FrameBuild)
547}
548
549pub(super) fn response_message(response: &Response, fallback: &str) -> String {
550    response
551        .data
552        .get("message")
553        .and_then(Value::as_str)
554        .map(ToOwned::to_owned)
555        .unwrap_or_else(|| fallback.to_string())
556}
557
558pub(super) fn response_is_fatal_panic(response: &Response) -> bool {
559    !response.success && response.data.get("code").and_then(Value::as_str) == Some("actor_fatal")
560}
561
562#[derive(Debug)]
563pub enum SubcError {
564    Runtime(std::io::Error),
565    ConnectionFile {
566        path: PathBuf,
567        source: subc_transport::ConnectionFileError,
568    },
569    NoEndpoint {
570        path: PathBuf,
571    },
572    InvalidEndpoint {
573        path: PathBuf,
574        endpoint: String,
575    },
576    Connect {
577        endpoint: String,
578        source: std::io::Error,
579    },
580    Auth {
581        endpoint: String,
582        source: subc_transport::AuthError,
583    },
584    FrameIo(subc_transport::FrameIoError),
585    FrameBuild(subc_protocol::FrameBuildError),
586    WriterClosed,
587    WriterBackpressureTimeout,
588    WriterJoin(tokio::task::JoinError),
589    Json(serde_json::Error),
590    ClosedBeforeHelloAck,
591    HelloRejected {
592        body: Option<ErrorBody>,
593    },
594    UnexpectedFrame {
595        ty: FrameType,
596    },
597}
598
599impl fmt::Display for SubcError {
600    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
601        match self {
602            Self::Runtime(e) => write!(f, "failed to build subc tokio runtime: {e}"),
603            Self::ConnectionFile { path, source } => {
604                write!(f, "failed to read subc connection file {path:?}: {source}")
605            }
606            Self::NoEndpoint { path } => {
607                write!(f, "subc connection file {path:?} has no endpoints")
608            }
609            Self::InvalidEndpoint { path, endpoint } => {
610                write!(
611                    f,
612                    "subc connection file {path:?} has invalid endpoint {endpoint}"
613                )
614            }
615            Self::Connect { endpoint, source } => {
616                write!(f, "failed to connect to subc endpoint {endpoint}: {source}")
617            }
618            Self::Auth { endpoint, source } => {
619                write!(
620                    f,
621                    "failed to authenticate to subc endpoint {endpoint}: {source}"
622                )
623            }
624            Self::FrameIo(e) => write!(f, "subc frame I/O error: {e}"),
625            Self::FrameBuild(e) => write!(f, "subc frame build error: {e}"),
626            Self::WriterClosed => write!(f, "subc writer task closed"),
627            Self::WriterBackpressureTimeout => write!(
628                f,
629                "subc writer task stayed backpressured while sending a control frame"
630            ),
631            Self::WriterJoin(e) => write!(f, "subc writer task join error: {e}"),
632            Self::Json(e) => write!(f, "subc JSON error: {e}"),
633            Self::ClosedBeforeHelloAck => {
634                write!(f, "subc daemon closed the connection before HelloAck")
635            }
636            Self::HelloRejected { body } => match body {
637                Some(b) => write!(f, "subc rejected ModuleHello: {} ({})", b.code, b.message),
638                None => write!(f, "subc rejected ModuleHello (unparseable error body)"),
639            },
640            Self::UnexpectedFrame { ty } => {
641                write!(f, "subc sent unexpected frame in place of HelloAck: {ty:?}")
642            }
643        }
644    }
645}
646
647impl std::error::Error for SubcError {}
648
649#[cfg(test)]
650mod tests {
651    use super::*;
652    use crate::subc::route_key;
653    use serde_json::json;
654    use std::sync::Arc;
655    use std::time::{Duration, Instant};
656    use subc_protocol::PROTOCOL_VERSION;
657
658    #[test]
659    fn writer_depth_counter_tracks_enqueued_frames_until_drain() {
660        let metrics = DispatchPathMetrics::new();
661        let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(8);
662
663        for corr in 1..=3 {
664            let frame = Frame::build(FrameType::Ping, control_flags(), 0, 0, corr, Vec::new())
665                .expect("test frame");
666            assert!(try_enqueue_writer_frame(&writer_tx, &metrics, frame).is_enqueued());
667        }
668        assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 3);
669
670        for _ in 0..3 {
671            writer_rx.try_recv().expect("queued writer frame");
672            decrement_counted_channel(&metrics.writer_queued);
673        }
674        assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 0);
675    }
676
677    #[tokio::test]
678    async fn reliable_writer_send_retries_after_timeout_and_preserves_frame() {
679        let metrics = Arc::new(DispatchPathMetrics::new());
680        let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(1);
681        writer_tx
682            .try_send(WriterFrame::plain(
683                Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
684            ))
685            .expect("prefill writer queue");
686
687        let metrics_for_task = Arc::clone(&metrics);
688        let tx_for_task = writer_tx.clone();
689        let send_task = tokio::spawn(async move {
690            send_reliable_writer_frame(
691                &tx_for_task,
692                &metrics_for_task,
693                Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
694                "test reliable frame",
695            )
696            .await
697        });
698
699        tokio::time::timeout(Duration::from_secs(2), async {
700            while metrics.writer_saturation_count.load(Ordering::Relaxed) < 2 {
701                tokio::time::sleep(Duration::from_millis(10)).await;
702            }
703        })
704        .await
705        .expect("reliable send should observe a timed-out full writer queue");
706
707        let prefilled = writer_rx.recv().await.expect("prefilled frame");
708        assert_eq!(prefilled.header.corr, 1);
709        let result = tokio::time::timeout(Duration::from_secs(2), send_task)
710            .await
711            .expect("reliable send should finish after writer drains")
712            .expect("reliable send task should not panic");
713        assert!(result.is_ok());
714        let delivered = writer_rx.recv().await.expect("retried reliable frame");
715        assert_eq!(delivered.header.corr, 2);
716    }
717
718    #[test]
719    fn response_is_fatal_panic_only_matches_panic_exclusive_code() {
720        let tool_error = Response::error("request-1", "internal_error", "ordinary tool error");
721        let panic_error = Response::error("request-2", "actor_fatal", "mutating panic");
722
723        assert!(!response_is_fatal_panic(&tool_error));
724        assert!(response_is_fatal_panic(&panic_error));
725    }
726
727    #[tokio::test]
728    async fn control_send_times_out_when_writer_queue_remains_full() {
729        let (writer_tx, _writer_rx) = mpsc::channel::<WriterFrame>(1);
730        let metrics = DispatchPathMetrics::new();
731        writer_tx
732            .try_send(WriterFrame::plain(
733                Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
734            ))
735            .expect("prefill writer queue");
736        let started = Instant::now();
737
738        let result = send_frame(
739            &writer_tx,
740            &metrics,
741            Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
742        )
743        .await;
744
745        assert!(matches!(result, Err(SubcError::WriterBackpressureTimeout)));
746        assert!(
747            started.elapsed() < Duration::from_secs(2),
748            "control send guard should be bounded"
749        );
750    }
751
752    #[test]
753    fn tool_response_frame_carries_flat_standalone_shape_in_structured_content() {
754        use crate::protocol::Response;
755
756        // A response with sidecars the FIRST-PARTY plugin drives UI from
757        // (status_bar, bg_completions, code) plus a normal result field.
758        let response = Response::success(
759            "req-7",
760            json!({
761                "complete": true,
762                "matches": 3,
763                "status_bar": { "errors": 0, "warnings": 1 },
764                "bg_completions": [{ "task_id": "bash-abc" }],
765            }),
766        );
767        let result = ToolCallResult {
768            text: "rendered text".to_string(),
769            response,
770        };
771
772        // The flat shape must equal the standalone NDJSON `tool_call` body:
773        // {id, success, ...data, text}. Build the standalone expectation the
774        // same way commands::tool_call::response_with_text does.
775        let expected_flat = json!({
776            "id": "req-7",
777            "success": true,
778            "complete": true,
779            "matches": 3,
780            "status_bar": { "errors": 0, "warnings": 1 },
781            "bg_completions": [{ "task_id": "bash-abc" }],
782            "text": "rendered text",
783        });
784        assert_eq!(
785            serde_json::to_value(FlatToolResponse {
786                response: &result.response,
787                text: &result.text,
788            })
789            .unwrap(),
790            expected_flat,
791            "structuredContent must be byte-identical to the standalone flat response"
792        );
793
794        // The frame body carries the MCP surface for generic hosts AND the flat
795        // sidecar shape under structuredContent for the first-party plugin.
796        let frame = build_tool_response_frame(
797            PROTOCOL_VERSION,
798            route_key(1, 1),
799            42,
800            control_flags(),
801            &result,
802            BindTrust::FirstParty,
803        )
804        .unwrap();
805        let expected_body = serde_json::to_vec(&json!({
806            "content": [{ "type": "text", "text": "rendered text" }],
807            "isError": false,
808            "structuredContent": expected_flat.clone(),
809        }))
810        .unwrap();
811        assert_eq!(
812            frame.body, expected_body,
813            "tool response wire bytes drifted"
814        );
815        let body: Value = serde_json::from_slice(&frame.body).unwrap();
816        assert_eq!(body["isError"], json!(false));
817        assert_eq!(body["content"][0]["type"], json!("text"));
818        assert_eq!(body["content"][0]["text"], json!("rendered text"));
819        assert_eq!(body["structuredContent"], expected_flat);
820
821        // A failed response flips isError and still carries the flat shape
822        // (with success:false + code) for the plugin's error path.
823        let err = Response::error_with_data(
824            "req-8",
825            "ambiguous_match",
826            "batch: edits[0] match 'same' is ambiguous (2 occurrences, expected 1). Use 'occurrence' (1-based) to select one, or 'replaceAll': true to replace every occurrence.",
827            json!({
828                "occurrences": [
829                    { "occurrence": 1, "line": 1, "context": "same same" },
830                    { "occurrence": 2, "line": 1, "context": "same same" }
831                ]
832            }),
833        );
834        let err_result = ToolCallResult {
835            text: "batch: edits[0] match 'same' is ambiguous (2 occurrences, expected 1). Use 'occurrence' (1-based) to select one, or 'replaceAll': true to replace every occurrence.".to_string(),
836            response: err,
837        };
838        let err_frame = build_tool_response_frame(
839            PROTOCOL_VERSION,
840            route_key(1, 1),
841            43,
842            control_flags(),
843            &err_result,
844            BindTrust::FirstParty,
845        )
846        .unwrap();
847        let err_body: Value = serde_json::from_slice(&err_frame.body).unwrap();
848        assert_eq!(err_body["isError"], json!(true));
849        assert_eq!(err_body["structuredContent"]["success"], json!(false));
850        assert_eq!(
851            err_body["structuredContent"]["code"],
852            json!("ambiguous_match")
853        );
854        assert_eq!(
855            err_body["structuredContent"]["occurrences"],
856            json!([
857                { "occurrence": 1, "line": 1, "context": "same same" },
858                { "occurrence": 2, "line": 1, "context": "same same" }
859            ])
860        );
861        let err_message = err_body["structuredContent"]["message"]
862            .as_str()
863            .expect("structured contract message");
864        assert!(err_message.contains("occurrence"));
865        assert!(err_message.contains("1-based"));
866        assert!(!err_message.contains("0-based"));
867        assert!(!err_message.contains("0-indexed"));
868        assert_eq!(err_body["structuredContent"]["text"], json!(err_message));
869
870        // UNTRUSTED (MCP) binds get text-only replies: no structuredContent
871        // key at all. Generic MCP hosts have no re-lift layer, and hosts like
872        // Claude Code feed structuredContent to the model verbatim when
873        // present, a raw JSON dump at a multiple of the token cost.
874        let untrusted_frame = build_tool_response_frame(
875            PROTOCOL_VERSION,
876            route_key(1, 1),
877            44,
878            control_flags(),
879            &err_result,
880            BindTrust::Untrusted,
881        )
882        .unwrap();
883        let untrusted_body: Value = serde_json::from_slice(&untrusted_frame.body).unwrap();
884        let untrusted_message = untrusted_body["content"][0]["text"]
885            .as_str()
886            .expect("untrusted contract message");
887        assert!(untrusted_message.contains("occurrence"));
888        assert!(untrusted_message.contains("1-based"));
889        assert!(!untrusted_message.contains("0-based"));
890        assert!(!untrusted_message.contains("0-indexed"));
891        assert_eq!(untrusted_body["isError"], json!(true));
892        assert!(
893            untrusted_body.get("structuredContent").is_none(),
894            "untrusted binds must not receive structuredContent: {untrusted_body}"
895        );
896    }
897}