1#[cfg(test)]
4use serde::ser::{SerializeMap, SerializeStruct};
5#[cfg(test)]
6use serde::{Serialize, Serializer};
7
8use super::{
9 control_flags, fmt, mpsc, Arc, AtomicUsize, BindTrust, DispatchPathMetrics, ErrorBody, Flags,
10 Frame, FrameType, Ordering, PathBuf, Response, RouteChannel, ToolCallResult, Value,
11 CONTROL_SEND_TIMEOUT, RELIABLE_WRITER_RETRY_INITIAL_BACKOFF, RELIABLE_WRITER_RETRY_MAX_BACKOFF,
12};
13use crate::run_tool_call::{PhaseTrace, ToolCallEgressTiming, ToolCallPhaseDurations};
14use std::borrow::Cow;
15use subc_protocol::{FrameBuildError, MAX_FRAME_BODY_LEN};
16
17pub(super) type WriterSender = mpsc::Sender<WriterFrame>;
18
19pub(super) struct ToolResponseWriteTrace {
20 phase_trace: PhaseTrace,
21 name: String,
22 root: PathBuf,
23 session: String,
24 channel: u16,
25 corr: u64,
26 enqueued_at: Option<std::time::Instant>,
27 queue_depth: usize,
28 writer_active_at_enqueue: bool,
29 writer_queue_was_full: bool,
30 reserve_timeouts: u32,
31}
32
33impl ToolResponseWriteTrace {
34 pub(super) fn new(
35 phase_trace: PhaseTrace,
36 name: String,
37 root: PathBuf,
38 session: String,
39 channel: u16,
40 corr: u64,
41 ) -> Self {
42 Self {
43 phase_trace,
44 name,
45 root,
46 session,
47 channel,
48 corr,
49 enqueued_at: None,
50 queue_depth: 0,
51 writer_active_at_enqueue: false,
52 writer_queue_was_full: false,
53 reserve_timeouts: 0,
54 }
55 }
56
57 fn mark_writer_queue_full(&mut self) {
58 self.writer_queue_was_full = true;
59 }
60
61 fn mark_reserve_timeout(&mut self) {
62 self.reserve_timeouts = self.reserve_timeouts.saturating_add(1);
63 }
64
65 fn mark_enqueued(&mut self, queue_depth: usize, writer_active: bool) {
66 self.enqueued_at = Some(std::time::Instant::now());
67 self.queue_depth = queue_depth;
68 self.writer_active_at_enqueue = writer_active;
69 }
70
71 pub(super) fn finish(
72 self,
73 dequeued: std::time::Instant,
74 write_started: std::time::Instant,
75 write_finished: std::time::Instant,
76 frame_bytes: usize,
77 ) -> Option<CompletedToolResponseTrace> {
78 let phases = self.phase_trace.finish(ToolCallEgressTiming {
79 enqueued: self.enqueued_at?,
80 dequeued,
81 write_started,
82 write_finished,
83 frame_bytes,
84 queue_depth: self.queue_depth,
85 writer_active_at_enqueue: self.writer_active_at_enqueue,
86 writer_queue_was_full: self.writer_queue_was_full,
87 reserve_timeouts: self.reserve_timeouts,
88 })?;
89 Some(CompletedToolResponseTrace {
90 name: self.name,
91 root: self.root,
92 session: self.session,
93 channel: self.channel,
94 corr: self.corr,
95 phases,
96 })
97 }
98}
99
100pub(super) struct CompletedToolResponseTrace {
101 pub(super) name: String,
102 pub(super) root: PathBuf,
103 pub(super) session: String,
104 pub(super) channel: u16,
105 pub(super) corr: u64,
106 pub(super) phases: ToolCallPhaseDurations,
107}
108
109enum WriterFrameBody {
110 Owned,
111 SharedPush(Arc<Vec<u8>>),
112}
113
114pub(super) struct WriterFrame {
115 pub(super) frame: Frame,
117 body: WriterFrameBody,
118 pub(super) tool_response_trace: Option<ToolResponseWriteTrace>,
119}
120
121impl std::ops::Deref for WriterFrame {
122 type Target = Frame;
123
124 fn deref(&self) -> &Self::Target {
125 &self.frame
126 }
127}
128
129impl WriterFrame {
130 pub(super) fn plain(frame: Frame) -> Self {
131 Self {
132 frame,
133 body: WriterFrameBody::Owned,
134 tool_response_trace: None,
135 }
136 }
137
138 pub(super) fn shared_push(frame: Frame, body: Arc<Vec<u8>>) -> Self {
139 debug_assert_eq!(frame.header.len as usize, body.len());
140 Self {
141 frame,
142 body: WriterFrameBody::SharedPush(body),
143 tool_response_trace: None,
144 }
145 }
146
147 fn traced_tool_response(frame: Frame, trace: ToolResponseWriteTrace) -> Self {
148 Self {
149 frame,
150 body: WriterFrameBody::Owned,
151 tool_response_trace: Some(trace),
152 }
153 }
154
155 pub(super) fn frame(&self) -> &Frame {
156 &self.frame
157 }
158
159 pub(super) fn body(&self) -> &[u8] {
160 match &self.body {
161 WriterFrameBody::Owned => &self.frame.body,
162 WriterFrameBody::SharedPush(body) => body,
163 }
164 }
165
166 #[cfg(test)]
167 pub(super) fn shared_push_body_strong_count(&self) -> Option<usize> {
168 match &self.body {
169 WriterFrameBody::Owned => None,
170 WriterFrameBody::SharedPush(body) => Some(Arc::strong_count(body)),
171 }
172 }
173
174 fn mark_writer_queue_full(&mut self) {
175 if let Some(trace) = self.tool_response_trace.as_mut() {
176 trace.mark_writer_queue_full();
177 }
178 }
179
180 fn mark_reserve_timeout(&mut self) {
181 if let Some(trace) = self.tool_response_trace.as_mut() {
182 trace.mark_reserve_timeout();
183 }
184 }
185
186 fn mark_enqueued(&mut self, queue_depth: usize, writer_active: bool) {
187 if let Some(trace) = self.tool_response_trace.as_mut() {
188 trace.mark_enqueued(queue_depth, writer_active);
189 }
190 }
191}
192
193pub(super) enum WriterEnqueueOutcome {
194 Enqueued,
195 Full(WriterFrame),
196 Closed,
197}
198
199impl WriterEnqueueOutcome {
200 #[cfg(test)]
201 pub(super) fn is_enqueued(&self) -> bool {
202 matches!(self, Self::Enqueued)
203 }
204}
205
206pub(super) fn decrement_counted_channel(counter: &AtomicUsize) {
207 let previous = counter.fetch_sub(1, Ordering::Relaxed);
208 debug_assert!(previous > 0, "counted channel depth underflow");
209}
210
211pub(super) async fn send_counted_channel<T>(
212 tx: &mpsc::Sender<T>,
213 counter: &AtomicUsize,
214 item: T,
215) -> Result<(), mpsc::error::SendError<T>> {
216 counter.fetch_add(1, Ordering::Relaxed);
217 match tx.send(item).await {
218 Ok(()) => Ok(()),
219 Err(error) => {
220 decrement_counted_channel(counter);
221 Err(error)
222 }
223 }
224}
225
226fn enqueue_writer_item(
227 permit: mpsc::Permit<'_, WriterFrame>,
228 metrics: &DispatchPathMetrics,
229 mut item: WriterFrame,
230) {
231 let queue_depth = metrics.writer_queued.fetch_add(1, Ordering::Relaxed) + 1;
232 item.mark_enqueued(queue_depth, metrics.writer_active.load(Ordering::Relaxed));
233 permit.send(item);
234}
235
236fn try_enqueue_writer_item(
237 tx: &WriterSender,
238 metrics: &DispatchPathMetrics,
239 mut item: WriterFrame,
240) -> WriterEnqueueOutcome {
241 match tx.try_reserve() {
242 Ok(permit) => {
243 enqueue_writer_item(permit, metrics, item);
244 WriterEnqueueOutcome::Enqueued
245 }
246 Err(mpsc::error::TrySendError::Full(())) => {
247 metrics
248 .writer_saturation_count
249 .fetch_add(1, Ordering::Relaxed);
250 item.mark_writer_queue_full();
251 WriterEnqueueOutcome::Full(item)
252 }
253 Err(mpsc::error::TrySendError::Closed(())) => {
254 drop(item);
255 WriterEnqueueOutcome::Closed
256 }
257 }
258}
259
260pub(super) fn try_enqueue_writer_frame(
261 tx: &WriterSender,
262 metrics: &DispatchPathMetrics,
263 frame: Frame,
264) -> WriterEnqueueOutcome {
265 try_enqueue_writer_item(tx, metrics, WriterFrame::plain(frame))
266}
267
268pub(super) fn try_enqueue_shared_push_frame(
269 tx: &WriterSender,
270 metrics: &DispatchPathMetrics,
271 frame: Frame,
272 body: Arc<Vec<u8>>,
273) -> WriterEnqueueOutcome {
274 try_enqueue_writer_item(tx, metrics, WriterFrame::shared_push(frame, body))
275}
276
277async fn send_reliable_writer_item(
278 tx: &WriterSender,
279 metrics: &DispatchPathMetrics,
280 mut item: WriterFrame,
281 context: &'static str,
282) -> Result<(), SubcError> {
283 let mut warned = false;
284 let mut backoff = RELIABLE_WRITER_RETRY_INITIAL_BACKOFF;
285
286 loop {
287 match try_enqueue_writer_item(tx, metrics, item) {
288 WriterEnqueueOutcome::Enqueued => return Ok(()),
289 WriterEnqueueOutcome::Closed => return Err(SubcError::WriterClosed),
290 WriterEnqueueOutcome::Full(returned_item) => {
291 item = returned_item;
292 }
293 }
294
295 match tokio::time::timeout(CONTROL_SEND_TIMEOUT, tx.reserve()).await {
296 Ok(Ok(permit)) => {
297 enqueue_writer_item(permit, metrics, item);
298 return Ok(());
299 }
300 Ok(Err(_)) => return Err(SubcError::WriterClosed),
301 Err(_) => {
302 metrics
303 .writer_saturation_count
304 .fetch_add(1, Ordering::Relaxed);
305 item.mark_reserve_timeout();
306 if !warned {
307 log::warn!(
308 "subc attach: writer queue stayed full while sending {context}; retrying reliable frame"
309 );
310 warned = true;
311 }
312 tokio::time::sleep(backoff).await;
313 backoff =
314 std::cmp::min(backoff.saturating_mul(2), RELIABLE_WRITER_RETRY_MAX_BACKOFF);
315 }
316 }
317 }
318}
319
320pub(super) async fn send_reliable_writer_frame(
321 tx: &WriterSender,
322 metrics: &DispatchPathMetrics,
323 frame: Frame,
324 context: &'static str,
325) -> Result<(), SubcError> {
326 send_reliable_writer_item(tx, metrics, WriterFrame::plain(frame), context).await
327}
328
329pub(super) async fn send_traced_tool_response_frame(
330 tx: &WriterSender,
331 metrics: &DispatchPathMetrics,
332 frame: Frame,
333 trace: ToolResponseWriteTrace,
334) -> Result<(), SubcError> {
335 send_reliable_writer_item(
336 tx,
337 metrics,
338 WriterFrame::traced_tool_response(frame, trace),
339 "tool response",
340 )
341 .await
342}
343
344pub(super) async fn send_frame(
345 tx: &WriterSender,
346 metrics: &DispatchPathMetrics,
347 frame: Frame,
348) -> Result<(), SubcError> {
349 match try_enqueue_writer_item(tx, metrics, WriterFrame::plain(frame)) {
350 WriterEnqueueOutcome::Enqueued => Ok(()),
351 WriterEnqueueOutcome::Closed => Err(SubcError::WriterClosed),
352 WriterEnqueueOutcome::Full(item) => {
353 match tokio::time::timeout(CONTROL_SEND_TIMEOUT, tx.reserve()).await {
354 Ok(Ok(permit)) => {
355 enqueue_writer_item(permit, metrics, item);
356 Ok(())
357 }
358 Ok(Err(_)) => Err(SubcError::WriterClosed),
359 Err(_) => {
360 metrics
361 .writer_saturation_count
362 .fetch_add(1, Ordering::Relaxed);
363 Err(SubcError::WriterBackpressureTimeout)
364 }
365 }
366 }
367 }
368}
369
370#[cfg(test)]
373struct FlatToolResponse<'a> {
374 response: &'a crate::protocol::Response,
375 text: &'a str,
376}
377
378#[cfg(test)]
379impl Serialize for FlatToolResponse<'_> {
380 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
381 where
382 S: Serializer,
383 {
384 let data = self.response.data.as_object();
385 let has_text = data.is_some_and(|data| data.contains_key("text"));
386 let field_count =
387 2 + data.map_or(0, |data| {
388 data.len()
389 - usize::from(data.contains_key("id"))
390 - usize::from(data.contains_key("success"))
391 }) + usize::from(!has_text);
392 let mut map = serializer.serialize_map(Some(field_count))?;
393 match data.and_then(|data| data.get("id")) {
394 Some(value) => map.serialize_entry("id", value)?,
395 None => map.serialize_entry("id", &self.response.id)?,
396 }
397 match data.and_then(|data| data.get("success")) {
398 Some(value) => map.serialize_entry("success", value)?,
399 None => map.serialize_entry("success", &self.response.success)?,
400 }
401 if let Some(data) = data {
402 for (key, value) in data {
403 match key.as_str() {
404 "id" | "success" => {}
405 "text" => map.serialize_entry(key, self.text)?,
406 _ => map.serialize_entry(key, value)?,
407 }
408 }
409 }
410 if !has_text {
411 map.serialize_entry("text", self.text)?;
412 }
413 map.end()
414 }
415}
416
417#[cfg(test)]
418struct ToolResponseEnvelope<'a> {
419 result: &'a ToolCallResult,
420 include_structured: bool,
423}
424
425#[cfg(test)]
443impl Serialize for ToolResponseEnvelope<'_> {
444 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
445 where
446 S: Serializer,
447 {
448 let fields = if self.include_structured { 3 } else { 2 };
449 let mut envelope = serializer.serialize_struct("ToolResponseEnvelope", fields)?;
450 envelope.serialize_field(
451 "content",
452 &[TextContent {
453 kind: "text",
454 text: &self.result.text,
455 }],
456 )?;
457 envelope.serialize_field("isError", &!self.result.response.success)?;
458 if self.include_structured {
459 envelope.serialize_field(
460 "structuredContent",
461 &FlatToolResponse {
462 response: &self.result.response,
463 text: &self.result.text,
464 },
465 )?;
466 }
467 envelope.end()
468 }
469}
470
471#[cfg(test)]
472#[derive(Serialize)]
473struct TextContent<'a> {
474 #[serde(rename = "type")]
475 kind: &'static str,
476 text: &'a str,
477}
478
479const TRANSPORT_MIB: usize = 1024 * 1024;
480const TOOL_RESPONSE_ENVELOPE_MARGIN_DIVISOR: usize = 8;
481const RESPONSE_TOO_LARGE_CODE: &str = "response_too_large";
482const TRANSPORT_TRUNCATION_REASON: &str = "transport_frame_limit";
483
484fn transport_limit_mib(bytes: usize) -> usize {
485 bytes.div_ceil(TRANSPORT_MIB).max(1)
486}
487
488fn tool_response_text_limit(max_body_len: usize, include_structured: bool) -> usize {
489 let margin = (max_body_len / TOOL_RESPONSE_ENVELOPE_MARGIN_DIVISOR)
493 .max(1_024)
494 .min(max_body_len / 2);
495 let rendered_copies = if include_structured { 2 } else { 1 };
496 max_body_len.saturating_sub(margin) / rendered_copies
497}
498
499fn truncate_rendered_text(text: &str, limit: usize) -> Cow<'_, str> {
500 if text.len() <= limit {
501 return Cow::Borrowed(text);
502 }
503
504 let notice = format!(
505 "[response truncated at {} MiB: full output exceeds the transport frame limit; use offset/limit or write to a file]",
506 transport_limit_mib(limit)
507 );
508 let suffix = format!("\n{notice}");
509 let mut prefix_end = text.len().min(limit.saturating_sub(suffix.len()));
510 while !text.is_char_boundary(prefix_end) {
511 prefix_end -= 1;
512 }
513
514 let mut truncated = String::with_capacity(prefix_end.saturating_add(suffix.len()));
515 truncated.push_str(&text[..prefix_end]);
516 truncated.push_str(&suffix);
517 Cow::Owned(truncated)
518}
519
520fn serialize_tool_response_body(
521 result: &ToolCallResult,
522 include_structured: bool,
523) -> Result<Vec<u8>, serde_json::Error> {
524 serialize_tool_response_body_with_text(result, include_structured, &result.text, false)
525}
526
527fn serialize_tool_response_body_with_text(
528 result: &ToolCallResult,
529 include_structured: bool,
530 rendered_text: &str,
531 transport_truncated: bool,
532) -> Result<Vec<u8>, serde_json::Error> {
533 let data_capacity = result
534 .response
535 .data
536 .as_object()
537 .and_then(|data| {
538 ["content", "output", "preview_diff"]
539 .into_iter()
540 .find_map(|key| data.get(key).and_then(Value::as_str))
541 .map(|data_text| {
542 if transport_truncated && data_text == result.text.as_str() {
543 rendered_text.len()
544 } else {
545 data_text.len()
546 }
547 })
548 })
549 .unwrap_or(0);
550 let capacity = rendered_text
551 .len()
552 .saturating_add(include_structured.then_some(data_capacity).unwrap_or(0))
553 .saturating_add(
554 include_structured
555 .then_some(rendered_text.len())
556 .unwrap_or(0),
557 )
558 .saturating_add(512);
559 let mut body = Vec::with_capacity(capacity);
560
561 body.extend_from_slice(b"{\"content\":[{\"type\":\"text\",\"text\":");
562 let encoded_text_start = body.len();
563 serde_json::to_writer(&mut body, rendered_text)?;
564 let encoded_text_end = body.len();
565 body.extend_from_slice(b"}],\"isError\":");
566 body.extend_from_slice(if result.response.success {
567 b"false"
568 } else {
569 b"true"
570 });
571
572 if include_structured {
573 body.extend_from_slice(b",\"structuredContent\":{\"id\":");
574 let data = result.response.data.as_object();
575 match data.and_then(|data| data.get("id")) {
576 Some(value) => serde_json::to_writer(&mut body, value)?,
577 None => serde_json::to_writer(&mut body, &result.response.id)?,
578 }
579 body.extend_from_slice(b",\"success\":");
580 match data.and_then(|data| data.get("success")) {
581 Some(value) => serde_json::to_writer(&mut body, value)?,
582 None => body.extend_from_slice(if result.response.success {
583 b"true"
584 } else {
585 b"false"
586 }),
587 }
588
589 let mut has_text = false;
590 let mut has_complete = false;
591 let mut has_truncated = false;
592 let mut has_truncation_reason = false;
593 if let Some(data) = data {
594 for (key, value) in data {
595 match key.as_str() {
596 "id" | "success" => continue,
597 "text" => has_text = true,
598 "complete" => has_complete = true,
599 "truncated" => has_truncated = true,
600 "truncation_reason" => has_truncation_reason = true,
601 _ => {}
602 }
603 body.push(b',');
604 serde_json::to_writer(&mut body, key)?;
605 body.push(b':');
606 match key.as_str() {
607 "text" => body.extend_from_within(encoded_text_start..encoded_text_end),
608 "complete" if transport_truncated => body.extend_from_slice(b"false"),
609 "truncated" if transport_truncated => body.extend_from_slice(b"true"),
610 "truncation_reason" if transport_truncated => {
611 serde_json::to_writer(&mut body, TRANSPORT_TRUNCATION_REASON)?;
612 }
613 _ if value.as_str() == Some(result.text.as_str()) => {
614 body.extend_from_within(encoded_text_start..encoded_text_end);
615 }
616 _ => serde_json::to_writer(&mut body, value)?,
617 }
618 }
619 }
620 if !has_text {
621 body.extend_from_slice(b",\"text\":");
622 body.extend_from_within(encoded_text_start..encoded_text_end);
623 }
624 if transport_truncated {
625 if !has_complete {
626 body.extend_from_slice(b",\"complete\":false");
627 }
628 if !has_truncated {
629 body.extend_from_slice(b",\"truncated\":true");
630 }
631 if !has_truncation_reason {
632 body.extend_from_slice(b",\"truncation_reason\":\"transport_frame_limit\"");
633 }
634 }
635 body.push(b'}');
636 }
637 body.push(b'}');
638 Ok(body)
639}
640
641fn response_too_large_frame(
642 ver: u8,
643 route: RouteChannel,
644 corr: u64,
645 flags: Flags,
646 body_len: usize,
647 max_body_len: usize,
648 include_structured: bool,
649) -> Frame {
650 let message = format!(
651 "tool response serialized to {body_len} bytes, exceeding the daemon transport limit of {max_body_len} bytes; re-run with a narrower range, a smaller limit, or offset+limit paging; output over {} MiB cannot cross the daemon transport",
652 transport_limit_mib(max_body_len)
653 );
654 let response = Response::error_with_data(
657 format!("subc-{}-{corr}", route.channel),
658 RESPONSE_TOO_LARGE_CODE,
659 message.clone(),
660 serde_json::json!({
661 "complete": false,
662 "truncated": true,
663 "truncation_reason": TRANSPORT_TRUNCATION_REASON,
664 }),
665 );
666 let result = ToolCallResult {
667 text: message,
668 response,
669 };
670 let body = serialize_tool_response_body(&result, include_structured)
671 .expect("fixed response_too_large envelope must serialize");
672 debug_assert!(
673 body.len() <= max_body_len,
674 "fixed response_too_large envelope must fit the effective body limit"
675 );
676 Frame::build_with_version(
677 ver,
678 FrameType::Response,
679 flags,
680 route.channel,
681 route.epoch,
682 corr,
683 body,
684 )
685 .expect("fixed response_too_large envelope must fit the protocol frame limit")
686}
687
688pub(super) fn build_tool_response_frame(
689 ver: u8,
690 route: RouteChannel,
691 corr: u64,
692 flags: Flags,
693 result: &ToolCallResult,
694 trust: BindTrust,
695) -> Result<Frame, SubcError> {
696 build_tool_response_frame_with_limit(
697 ver,
698 route,
699 corr,
700 flags,
701 result,
702 trust,
703 MAX_FRAME_BODY_LEN as usize,
704 )
705}
706
707pub(super) fn build_tool_response_frame_with_limit(
708 ver: u8,
709 route: RouteChannel,
710 corr: u64,
711 flags: Flags,
712 result: &ToolCallResult,
713 trust: BindTrust,
714 max_body_len: usize,
715) -> Result<Frame, SubcError> {
716 let include_structured = !matches!(trust, BindTrust::Untrusted);
732 let effective_max_body_len = max_body_len.min(MAX_FRAME_BODY_LEN as usize);
733 let text_limit = tool_response_text_limit(effective_max_body_len, include_structured);
734 let rendered_text = truncate_rendered_text(&result.text, text_limit);
735 let transport_truncated = matches!(rendered_text, Cow::Owned(_));
736 let body = serialize_tool_response_body_with_text(
737 result,
738 include_structured,
739 &rendered_text,
740 transport_truncated,
741 )
742 .map_err(SubcError::Json)?;
743
744 if effective_max_body_len < MAX_FRAME_BODY_LEN as usize && body.len() > effective_max_body_len {
745 return Ok(response_too_large_frame(
746 ver,
747 route,
748 corr,
749 flags,
750 body.len(),
751 effective_max_body_len,
752 include_structured,
753 ));
754 }
755
756 match Frame::build_with_version(
757 ver,
758 FrameType::Response,
759 flags,
760 route.channel,
761 route.epoch,
762 corr,
763 body,
764 ) {
765 Ok(frame) => Ok(frame),
766 Err(FrameBuildError::BodyExceedsMax { body_len, max }) => Ok(response_too_large_frame(
767 ver,
768 route,
769 corr,
770 flags,
771 body_len,
772 max as usize,
773 include_structured,
774 )),
775 Err(FrameBuildError::BodyTooLarge { body_len }) => Ok(response_too_large_frame(
776 ver,
777 route,
778 corr,
779 flags,
780 body_len,
781 MAX_FRAME_BODY_LEN as usize,
782 include_structured,
783 )),
784 }
785}
786
787pub(super) fn build_error_frame(
788 ver: u8,
789 channel: u16,
790 epoch: u32,
791 corr: u64,
792 flags: Flags,
793 code: &str,
794 message: &str,
795) -> Result<Frame, SubcError> {
796 let body = serde_json::to_vec(&ErrorBody {
797 code: code.to_string(),
798 message: message.to_string(),
799 })
800 .map_err(SubcError::Json)?;
801 Frame::build_with_version(ver, FrameType::Error, flags, channel, epoch, corr, body)
802 .map_err(SubcError::FrameBuild)
803}
804
805pub(super) fn build_goodbye_frame(
806 ver: u8,
807 channel: u16,
808 epoch: u32,
809 corr: u64,
810) -> Result<Frame, SubcError> {
811 Frame::build_with_version(
812 ver,
813 FrameType::Goodbye,
814 control_flags(),
815 channel,
816 epoch,
817 corr,
818 Vec::new(),
819 )
820 .map_err(SubcError::FrameBuild)
821}
822
823pub(super) fn response_message(response: &Response, fallback: &str) -> String {
824 response
825 .data
826 .get("message")
827 .and_then(Value::as_str)
828 .map(ToOwned::to_owned)
829 .unwrap_or_else(|| fallback.to_string())
830}
831
832pub(super) fn response_is_fatal_panic(response: &Response) -> bool {
833 !response.success && response.data.get("code").and_then(Value::as_str) == Some("actor_fatal")
834}
835
836#[derive(Debug)]
837pub enum SubcError {
838 Runtime(std::io::Error),
839 ConnectionFile {
840 path: PathBuf,
841 source: subc_transport::ConnectionFileError,
842 },
843 NoEndpoint {
844 path: PathBuf,
845 },
846 InvalidEndpoint {
847 path: PathBuf,
848 endpoint: String,
849 },
850 Connect {
851 endpoint: String,
852 source: std::io::Error,
853 },
854 Auth {
855 endpoint: String,
856 source: subc_transport::AuthError,
857 },
858 FrameIo(subc_transport::FrameIoError),
859 FrameBuild(subc_protocol::FrameBuildError),
860 WriterClosed,
861 WriterBackpressureTimeout,
862 WriterJoin(tokio::task::JoinError),
863 Json(serde_json::Error),
864 ClosedBeforeHelloAck,
865 HelloRejected {
866 body: Option<ErrorBody>,
867 },
868 UnexpectedFrame {
869 ty: FrameType,
870 },
871}
872
873impl fmt::Display for SubcError {
874 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
875 match self {
876 Self::Runtime(e) => write!(f, "failed to build subc tokio runtime: {e}"),
877 Self::ConnectionFile { path, source } => {
878 write!(f, "failed to read subc connection file {path:?}: {source}")
879 }
880 Self::NoEndpoint { path } => {
881 write!(f, "subc connection file {path:?} has no endpoints")
882 }
883 Self::InvalidEndpoint { path, endpoint } => {
884 write!(
885 f,
886 "subc connection file {path:?} has invalid endpoint {endpoint}"
887 )
888 }
889 Self::Connect { endpoint, source } => {
890 write!(f, "failed to connect to subc endpoint {endpoint}: {source}")
891 }
892 Self::Auth { endpoint, source } => {
893 write!(
894 f,
895 "failed to authenticate to subc endpoint {endpoint}: {source}"
896 )
897 }
898 Self::FrameIo(e) => write!(f, "subc frame I/O error: {e}"),
899 Self::FrameBuild(e) => write!(f, "subc frame build error: {e}"),
900 Self::WriterClosed => write!(f, "subc writer task closed"),
901 Self::WriterBackpressureTimeout => write!(
902 f,
903 "subc writer task stayed backpressured while sending a control frame"
904 ),
905 Self::WriterJoin(e) => write!(f, "subc writer task join error: {e}"),
906 Self::Json(e) => write!(f, "subc JSON error: {e}"),
907 Self::ClosedBeforeHelloAck => {
908 write!(f, "subc daemon closed the connection before HelloAck")
909 }
910 Self::HelloRejected { body } => match body {
911 Some(b) => write!(f, "subc rejected ModuleHello: {} ({})", b.code, b.message),
912 None => write!(f, "subc rejected ModuleHello (unparseable error body)"),
913 },
914 Self::UnexpectedFrame { ty } => {
915 write!(f, "subc sent unexpected frame in place of HelloAck: {ty:?}")
916 }
917 }
918 }
919}
920
921impl std::error::Error for SubcError {}
922
923#[cfg(test)]
924mod tests {
925 use super::*;
926 use crate::subc::route_key;
927 use serde_json::json;
928 use std::sync::Arc;
929 use std::time::{Duration, Instant};
930 use subc_protocol::PROTOCOL_VERSION;
931
932 #[test]
933 fn writer_depth_counter_tracks_enqueued_frames_until_drain() {
934 let metrics = DispatchPathMetrics::new();
935 let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(8);
936
937 for corr in 1..=3 {
938 let frame = Frame::build(FrameType::Ping, control_flags(), 0, 0, corr, Vec::new())
939 .expect("test frame");
940 assert!(try_enqueue_writer_frame(&writer_tx, &metrics, frame).is_enqueued());
941 }
942 assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 3);
943
944 for _ in 0..3 {
945 writer_rx.try_recv().expect("queued writer frame");
946 decrement_counted_channel(&metrics.writer_queued);
947 }
948 assert_eq!(metrics.writer_queued.load(Ordering::Relaxed), 0);
949 }
950
951 #[tokio::test]
952 async fn reliable_writer_send_retries_after_timeout_and_preserves_frame() {
953 let metrics = Arc::new(DispatchPathMetrics::new());
954 let (writer_tx, mut writer_rx) = mpsc::channel::<WriterFrame>(1);
955 writer_tx
956 .try_send(WriterFrame::plain(
957 Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
958 ))
959 .expect("prefill writer queue");
960
961 let metrics_for_task = Arc::clone(&metrics);
962 let tx_for_task = writer_tx.clone();
963 let send_task = tokio::spawn(async move {
964 send_reliable_writer_frame(
965 &tx_for_task,
966 &metrics_for_task,
967 Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
968 "test reliable frame",
969 )
970 .await
971 });
972
973 tokio::time::timeout(Duration::from_secs(2), async {
974 while metrics.writer_saturation_count.load(Ordering::Relaxed) < 2 {
975 tokio::time::sleep(Duration::from_millis(10)).await;
976 }
977 })
978 .await
979 .expect("reliable send should observe a timed-out full writer queue");
980
981 let prefilled = writer_rx.recv().await.expect("prefilled frame");
982 assert_eq!(prefilled.header.corr, 1);
983 let result = tokio::time::timeout(Duration::from_secs(2), send_task)
984 .await
985 .expect("reliable send should finish after writer drains")
986 .expect("reliable send task should not panic");
987 assert!(result.is_ok());
988 let delivered = writer_rx.recv().await.expect("retried reliable frame");
989 assert_eq!(delivered.header.corr, 2);
990 }
991
992 #[test]
993 fn response_is_fatal_panic_only_matches_panic_exclusive_code() {
994 let tool_error = Response::error("request-1", "internal_error", "ordinary tool error");
995 let panic_error = Response::error("request-2", "actor_fatal", "mutating panic");
996
997 assert!(!response_is_fatal_panic(&tool_error));
998 assert!(response_is_fatal_panic(&panic_error));
999 }
1000
1001 #[tokio::test]
1002 async fn control_send_times_out_when_writer_queue_remains_full() {
1003 let (writer_tx, _writer_rx) = mpsc::channel::<WriterFrame>(1);
1004 let metrics = DispatchPathMetrics::new();
1005 writer_tx
1006 .try_send(WriterFrame::plain(
1007 Frame::build(FrameType::Ping, control_flags(), 0, 0, 1, Vec::new()).unwrap(),
1008 ))
1009 .expect("prefill writer queue");
1010 let started = Instant::now();
1011
1012 let result = send_frame(
1013 &writer_tx,
1014 &metrics,
1015 Frame::build(FrameType::Pong, control_flags(), 0, 0, 2, Vec::new()).unwrap(),
1016 )
1017 .await;
1018
1019 assert!(matches!(result, Err(SubcError::WriterBackpressureTimeout)));
1020 assert!(
1021 started.elapsed() < Duration::from_secs(2),
1022 "control send guard should be bounded"
1023 );
1024 }
1025
1026 fn legacy_tool_response_body(result: &ToolCallResult, include_structured: bool) -> Vec<u8> {
1027 serde_json::to_vec(&ToolResponseEnvelope {
1028 result,
1029 include_structured,
1030 })
1031 .expect("serialize legacy tool response envelope")
1032 }
1033
1034 fn assert_tool_response_frame_matches_legacy(result: &ToolCallResult, trust: BindTrust) {
1035 let include_structured = !matches!(trust, BindTrust::Untrusted);
1036 let legacy = Frame::build_with_version(
1037 PROTOCOL_VERSION,
1038 FrameType::Response,
1039 control_flags(),
1040 7,
1041 3,
1042 42,
1043 legacy_tool_response_body(result, include_structured),
1044 )
1045 .expect("build legacy tool response frame");
1046 let optimized = build_tool_response_frame(
1047 PROTOCOL_VERSION,
1048 route_key(7, 3),
1049 42,
1050 control_flags(),
1051 result,
1052 trust,
1053 )
1054 .expect("build optimized tool response frame");
1055 assert_eq!(optimized.header.encode(), legacy.header.encode());
1056 assert_eq!(optimized.body, legacy.body);
1057 }
1058
1059 fn production_shape_result(tool: &str, payload_bytes: usize) -> ToolCallResult {
1060 let payload = "x".repeat(payload_bytes);
1061 let response = match tool {
1062 "read" => Response::success(
1063 "subc-7-42",
1064 json!({
1065 "content": payload,
1066 "path": "/workspace/src/fixture.rs",
1067 "start_line": 1,
1068 "end_line": payload_bytes / 40 + 1,
1069 "total_lines": payload_bytes / 40 + 1,
1070 "truncated": false,
1071 }),
1072 ),
1073 "edit" => Response::success(
1074 "subc-7-42",
1075 json!({
1076 "path": "/workspace/src/fixture.rs",
1077 "edits_applied": 1,
1078 "diff": { "additions": 12, "deletions": 8 },
1079 "preview_diff": payload,
1080 }),
1081 ),
1082 "bash" => Response::success(
1083 "subc-7-42",
1084 json!({
1085 "output": payload,
1086 "exit_code": 0,
1087 "timed_out": false,
1088 "status": "completed",
1089 }),
1090 ),
1091 _ => unreachable!("production-shape probe tool"),
1092 };
1093 let text = crate::subc_format::format_response_with_context(
1094 tool,
1095 &response,
1096 &crate::subc_format::FormatContext::default(),
1097 );
1098 ToolCallResult { text, response }
1099 }
1100
1101 #[test]
1102 fn optimized_tool_response_body_matches_legacy_wire_at_production_shapes() {
1103 for tool in ["read", "edit", "bash"] {
1104 for payload_bytes in [1_024, 10 * 1_024, 50 * 1_024] {
1105 let result = production_shape_result(tool, payload_bytes);
1106 for trust in [BindTrust::Untrusted, BindTrust::FirstParty] {
1107 assert_tool_response_frame_matches_legacy(&result, trust);
1108 }
1109 }
1110 }
1111
1112 let result = ToolCallResult {
1113 text: r#"replacement text with \"escapes\"
1114"#
1115 .to_string(),
1116 response: Response::success(
1117 "outer-id",
1118 json!({
1119 "id": "data-id",
1120 "success": false,
1121 "text": "ignored source text",
1122 "nested": [null, true, -7, "replacement text with \"escapes\"\n"],
1123 }),
1124 ),
1125 };
1126 assert_tool_response_frame_matches_legacy(&result, BindTrust::FirstParty);
1127 }
1128
1129 fn median_duration(samples: &mut [Duration]) -> Duration {
1130 samples.sort_unstable();
1131 samples[samples.len() / 2]
1132 }
1133
1134 fn time_envelope_builds(
1135 results: &[ToolCallResult],
1136 weights: &[usize],
1137 legacy: bool,
1138 iterations: usize,
1139 ) -> Duration {
1140 let started = Instant::now();
1141 for _ in 0..iterations {
1142 for (result, &weight) in results.iter().zip(weights) {
1143 for _ in 0..weight {
1144 let body = if legacy {
1145 legacy_tool_response_body(std::hint::black_box(result), true)
1146 } else {
1147 serialize_tool_response_body(std::hint::black_box(result), true)
1148 .expect("serialize optimized tool response envelope")
1149 };
1150 std::hint::black_box(body);
1151 }
1152 }
1153 }
1154 started.elapsed()
1155 }
1156
1157 #[test]
1158 #[ignore = "manual release-mode serving-plane performance probe"]
1159 fn tool_response_envelope_perf_probe() {
1160 let shapes = [
1161 ("read", 1_024),
1162 ("read", 10 * 1_024),
1163 ("read", 50 * 1_024),
1164 ("edit", 1_024),
1165 ("edit", 10 * 1_024),
1166 ("edit", 50 * 1_024),
1167 ("bash", 1_024),
1168 ("bash", 10 * 1_024),
1169 ("bash", 50 * 1_024),
1170 ];
1171 let weights = [20, 15, 5, 15, 10, 5, 15, 10, 5];
1174 let results: Vec<_> = shapes
1175 .iter()
1176 .map(|&(tool, payload_bytes)| production_shape_result(tool, payload_bytes))
1177 .collect();
1178 for result in &results {
1179 assert_eq!(
1180 serialize_tool_response_body(result, true).unwrap(),
1181 legacy_tool_response_body(result, true)
1182 );
1183 }
1184
1185 let iterations = 40;
1186 let calls_per_sample = iterations * weights.iter().sum::<usize>();
1187 let mut legacy_samples = Vec::new();
1188 let mut optimized_samples = Vec::new();
1189 for _ in 0..7 {
1190 legacy_samples.push(time_envelope_builds(&results, &weights, true, iterations));
1191 optimized_samples.push(time_envelope_builds(&results, &weights, false, iterations));
1192 }
1193 let legacy = median_duration(&mut legacy_samples);
1194 let optimized = median_duration(&mut optimized_samples);
1195 let legacy_us = legacy.as_secs_f64() * 1_000_000.0 / calls_per_sample as f64;
1196 let optimized_us = optimized.as_secs_f64() * 1_000_000.0 / calls_per_sample as f64;
1197 let saved_percent = (legacy_us - optimized_us) / legacy_us * 100.0;
1198 eprintln!(
1199 "trusted mixed tool-response envelope: before={legacy_us:.3} us/call after={optimized_us:.3} us/call saved={saved_percent:.1}%; 7 run medians, {calls_per_sample} calls/run"
1200 );
1201 assert!(optimized < legacy, "optimized envelope regressed");
1202 }
1203
1204 #[test]
1205 fn tool_response_frame_carries_flat_standalone_shape_in_structured_content() {
1206 use crate::protocol::Response;
1207
1208 let response = Response::success(
1211 "req-7",
1212 json!({
1213 "complete": true,
1214 "matches": 3,
1215 "status_bar": { "errors": 0, "warnings": 1 },
1216 "bg_completions": [{ "task_id": "bash-abc" }],
1217 }),
1218 );
1219 let result = ToolCallResult {
1220 text: "rendered text".to_string(),
1221 response,
1222 };
1223
1224 let expected_flat = json!({
1228 "id": "req-7",
1229 "success": true,
1230 "complete": true,
1231 "matches": 3,
1232 "status_bar": { "errors": 0, "warnings": 1 },
1233 "bg_completions": [{ "task_id": "bash-abc" }],
1234 "text": "rendered text",
1235 });
1236 assert_eq!(
1237 serde_json::to_value(FlatToolResponse {
1238 response: &result.response,
1239 text: &result.text,
1240 })
1241 .unwrap(),
1242 expected_flat,
1243 "structuredContent must be byte-identical to the standalone flat response"
1244 );
1245
1246 let frame = build_tool_response_frame(
1249 PROTOCOL_VERSION,
1250 route_key(1, 1),
1251 42,
1252 control_flags(),
1253 &result,
1254 BindTrust::FirstParty,
1255 )
1256 .unwrap();
1257 let expected_body = serde_json::to_vec(&json!({
1258 "content": [{ "type": "text", "text": "rendered text" }],
1259 "isError": false,
1260 "structuredContent": expected_flat.clone(),
1261 }))
1262 .unwrap();
1263 assert_eq!(
1264 frame.body, expected_body,
1265 "tool response wire bytes drifted"
1266 );
1267 let body: Value = serde_json::from_slice(&frame.body).unwrap();
1268 assert_eq!(body["isError"], json!(false));
1269 assert_eq!(body["content"][0]["type"], json!("text"));
1270 assert_eq!(body["content"][0]["text"], json!("rendered text"));
1271 assert_eq!(body["structuredContent"], expected_flat);
1272
1273 let err = Response::error_with_data(
1276 "req-8",
1277 "ambiguous_match",
1278 "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.",
1279 json!({
1280 "occurrences": [
1281 { "occurrence": 1, "line": 1, "context": "same same" },
1282 { "occurrence": 2, "line": 1, "context": "same same" }
1283 ]
1284 }),
1285 );
1286 let err_result = ToolCallResult {
1287 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(),
1288 response: err,
1289 };
1290 let err_frame = build_tool_response_frame(
1291 PROTOCOL_VERSION,
1292 route_key(1, 1),
1293 43,
1294 control_flags(),
1295 &err_result,
1296 BindTrust::FirstParty,
1297 )
1298 .unwrap();
1299 let err_body: Value = serde_json::from_slice(&err_frame.body).unwrap();
1300 assert_eq!(err_body["isError"], json!(true));
1301 assert_eq!(err_body["structuredContent"]["success"], json!(false));
1302 assert_eq!(
1303 err_body["structuredContent"]["code"],
1304 json!("ambiguous_match")
1305 );
1306 assert_eq!(
1307 err_body["structuredContent"]["occurrences"],
1308 json!([
1309 { "occurrence": 1, "line": 1, "context": "same same" },
1310 { "occurrence": 2, "line": 1, "context": "same same" }
1311 ])
1312 );
1313 let err_message = err_body["structuredContent"]["message"]
1314 .as_str()
1315 .expect("structured contract message");
1316 assert!(err_message.contains("occurrence"));
1317 assert!(err_message.contains("1-based"));
1318 assert!(!err_message.contains("0-based"));
1319 assert!(!err_message.contains("0-indexed"));
1320 assert_eq!(err_body["structuredContent"]["text"], json!(err_message));
1321
1322 let untrusted_frame = build_tool_response_frame(
1327 PROTOCOL_VERSION,
1328 route_key(1, 1),
1329 44,
1330 control_flags(),
1331 &err_result,
1332 BindTrust::Untrusted,
1333 )
1334 .unwrap();
1335 let untrusted_body: Value = serde_json::from_slice(&untrusted_frame.body).unwrap();
1336 let untrusted_message = untrusted_body["content"][0]["text"]
1337 .as_str()
1338 .expect("untrusted contract message");
1339 assert!(untrusted_message.contains("occurrence"));
1340 assert!(untrusted_message.contains("1-based"));
1341 assert!(!untrusted_message.contains("0-based"));
1342 assert!(!untrusted_message.contains("0-indexed"));
1343 assert_eq!(untrusted_body["isError"], json!(true));
1344 assert!(
1345 untrusted_body.get("structuredContent").is_none(),
1346 "untrusted binds must not receive structuredContent: {untrusted_body}"
1347 );
1348 }
1349
1350 #[test]
1351 fn normal_response_is_byte_identical_with_the_size_guard_enabled() {
1352 let result = production_shape_result("read", 1_024);
1353 let default = build_tool_response_frame(
1354 PROTOCOL_VERSION,
1355 route_key(5, 2),
1356 76,
1357 control_flags(),
1358 &result,
1359 BindTrust::FirstParty,
1360 )
1361 .expect("default response frame");
1362 let guarded = build_tool_response_frame_with_limit(
1363 PROTOCOL_VERSION,
1364 route_key(5, 2),
1365 76,
1366 control_flags(),
1367 &result,
1368 BindTrust::FirstParty,
1369 8 * 1_024,
1370 )
1371 .expect("guarded response frame");
1372
1373 assert_eq!(guarded.header.encode(), default.header.encode());
1374 assert_eq!(guarded.body, default.body);
1375 }
1376
1377 #[test]
1378 fn oversized_rendered_text_is_utf8_safely_truncated_with_an_explicit_gap() {
1379 const TEST_BODY_LIMIT: usize = 8 * 1_024;
1380 let text = "é".repeat(3_000);
1381 let result = ToolCallResult {
1382 text: text.clone(),
1383 response: Response::success("large-text", json!({ "text": text })),
1384 };
1385
1386 let frame = build_tool_response_frame_with_limit(
1387 PROTOCOL_VERSION,
1388 route_key(5, 2),
1389 77,
1390 control_flags(),
1391 &result,
1392 BindTrust::FirstParty,
1393 TEST_BODY_LIMIT,
1394 )
1395 .expect("truncated response frame");
1396
1397 assert!(frame.body.len() <= TEST_BODY_LIMIT);
1398 let body: Value =
1399 serde_json::from_slice(&frame.body).expect("valid truncated response JSON");
1400 let rendered = body["content"][0]["text"]
1401 .as_str()
1402 .expect("rendered response text");
1403 assert!(rendered.ends_with(
1404 "[response truncated at 1 MiB: full output exceeds the transport frame limit; use offset/limit or write to a file]"
1405 ));
1406 assert!(rendered.len() < result.text.len());
1407 assert_eq!(body["isError"], json!(false));
1408 assert_eq!(body["structuredContent"]["success"], json!(true));
1409 assert_eq!(body["structuredContent"]["complete"], json!(false));
1410 assert_eq!(body["structuredContent"]["truncated"], json!(true));
1411 assert_eq!(
1412 body["structuredContent"]["truncation_reason"],
1413 json!(TRANSPORT_TRUNCATION_REASON)
1414 );
1415 assert_eq!(body["structuredContent"]["text"], json!(rendered));
1416 }
1417
1418 #[test]
1419 fn oversized_structured_data_gets_a_correlated_response_too_large_fallback() {
1420 const TEST_BODY_LIMIT: usize = 8 * 1_024;
1421 let text_limit = tool_response_text_limit(TEST_BODY_LIMIT, true);
1422 let between_threshold_and_limit = text_limit + 512;
1423 assert!(between_threshold_and_limit < TEST_BODY_LIMIT);
1424 let result = ToolCallResult {
1425 text: "r".repeat(between_threshold_and_limit),
1426 response: Response::success(
1427 "large-structured",
1428 json!({ "payload": "p".repeat(between_threshold_and_limit) }),
1429 ),
1430 };
1431
1432 let frame = build_tool_response_frame_with_limit(
1433 PROTOCOL_VERSION,
1434 route_key(9, 4),
1435 88,
1436 control_flags(),
1437 &result,
1438 BindTrust::FirstParty,
1439 TEST_BODY_LIMIT,
1440 )
1441 .expect("response_too_large fallback frame");
1442
1443 assert_eq!(frame.header.ty, FrameType::Response);
1444 assert_eq!(frame.header.channel, 9);
1445 assert_eq!(frame.header.epoch, 4);
1446 assert_eq!(frame.header.corr, 88);
1447 assert!(frame.body.len() <= TEST_BODY_LIMIT);
1448 let body: Value =
1449 serde_json::from_slice(&frame.body).expect("valid fallback response JSON");
1450 assert_eq!(body["isError"], json!(true));
1451 assert_eq!(body["structuredContent"]["success"], json!(false));
1452 assert_eq!(
1453 body["structuredContent"]["code"],
1454 json!(RESPONSE_TOO_LARGE_CODE)
1455 );
1456 assert_eq!(body["structuredContent"]["complete"], json!(false));
1457 assert_eq!(body["structuredContent"]["truncated"], json!(true));
1458 let message = body["structuredContent"]["message"]
1459 .as_str()
1460 .expect("fallback message");
1461 assert!(message.contains("serialized to "));
1462 assert!(message.contains("8192 bytes"));
1463 assert!(message.contains("narrower range"));
1464 assert!(message.contains("offset+limit paging"));
1465 assert!(message.contains("output over 1 MiB cannot cross the daemon transport"));
1466 }
1467}