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