1use 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 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
366struct 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 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 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 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 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 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 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 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}