返回 CodeWhale
mailbox.rs
根目录 / crates / tui / src / tools / subagent / mailbox.rs
1 //! Mailbox abstraction for sub-agent runtime coordination.
2 //!
3 //! Monotonic sequence numbers give every consumer a consistent ordering even
4 //! when multiple subscribers (e.g. UI card + parent agent) drain
5 //! independently; close-as-cancel lets a single signal both stop new mail and
6 //! propagate cancellation through nested children.
7
8 use std::collections::VecDeque;
9 use std::sync::Arc;
10 use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
11 #[cfg(test)]
12 use std::time::Duration;
13
14 use serde::{Deserialize, Serialize};
15 use tokio::sync::{mpsc, watch};
16 use tokio_util::sync::CancellationToken;
17
18 #[cfg(test)]
19 use crate::config::ProviderKind;
20 use crate::tools::todo::TodoListSnapshot;
21 use codewhale_models::Usage;
22
23 use super::FleetRole;
24
25 /// Stable, structured progress envelope shared across the sub-agent surface.
26 ///
27 /// Tracks the lifecycle of a single agent (identified by `agent_id`) end to
28 /// end: spawn, per-step progress, tool execution, completion / failure /
29 /// cancellation, and parent → child topology so consumers can render trees.
30 #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
31 #[serde(tag = "kind", rename_all = "snake_case")]
32 pub enum MailboxMessage {
33 /// Agent has been started (background task is running).
34 Started {
35 agent_id: String,
36 agent_type: String,
37 },
38 /// Free-form human-readable progress (mirrors `Event::AgentProgress`).
39 Progress { agent_id: String, status: String },
40 /// A tool call inside the agent has started.
41 ToolCallStarted {
42 agent_id: String,
43 tool_name: String,
44 step: u32,
45 },
46 /// A tool call inside the agent has finished.
47 ToolCallCompleted {
48 agent_id: String,
49 tool_name: String,
50 step: u32,
51 ok: bool,
52 },
53 /// A child agent was spawned by this agent.
54 ChildSpawned { parent_id: String, child_id: String },
55 /// Agent completed successfully (carries the summary line shown in the
56 /// transcript; full result is still available through the transcript handle).
57 Completed { agent_id: String, summary: String },
58 /// Agent failed with the carried error message.
59 Failed { agent_id: String, error: String },
60 /// Agent was interrupted (e.g. API timeout) with a continuable
61 /// checkpoint; the worker is parked waiting for continuation input.
62 Interrupted { agent_id: String, reason: String },
63 /// Cancellation propagated to this agent.
64 Cancelled { agent_id: String },
65 /// This agent's **own** bounded To-do snapshot (#4810).
66 ///
67 /// Published by the agent that owns the ledger, from its private list, so a
68 /// consumer keyed on `agent_id` can never attribute a parent's or a
69 /// sibling's work to this agent. Emitted only when the snapshot actually
70 /// changes; the payload is the canonical [`TodoListSnapshot`], not a second
71 /// ledger.
72 WorkState {
73 agent_id: String,
74 /// Absent in older persisted payloads, which predate per-agent Work
75 /// state; those deserialize to an empty (no work stated) snapshot.
76 #[serde(default)]
77 todo: TodoListSnapshot,
78 },
79 /// Incremental token usage from a sub-agent's API call.
80 /// Published after each turn so the parent's cost counter updates live.
81 TokenUsage {
82 agent_id: String,
83 /// Stable identity of the provider response. Runtime accounting uses
84 /// this across direct durability, mailbox replay, and restart dedupe.
85 source_id: String,
86 /// Immutable provider/model/billing evidence captured before the
87 /// child request was sent. Boxed: the envelope dwarfs every other
88 /// variant, and mailboxes queue many messages.
89 route: Box<crate::cost_status::EffectiveRouteEnvelope>,
90 /// Provider usage payload, including cache-hit/cache-miss fields.
91 usage: Usage,
92 },
93 }
94
95 impl MailboxMessage {
96 /// `agent_id` of the message subject (for `ChildSpawned` this is the
97 /// child, since that's the new lifecycle being announced).
98 #[must_use]
99 pub fn agent_id(&self) -> &str {
100 match self {
101 Self::Started { agent_id, .. }
102 | Self::Progress { agent_id, .. }
103 | Self::ToolCallStarted { agent_id, .. }
104 | Self::ToolCallCompleted { agent_id, .. }
105 | Self::Completed { agent_id, .. }
106 | Self::Failed { agent_id, .. }
107 | Self::Interrupted { agent_id, .. }
108 | Self::Cancelled { agent_id }
109 | Self::WorkState { agent_id, .. }
110 | Self::TokenUsage { agent_id, .. } => agent_id,
111 Self::ChildSpawned { child_id, .. } => child_id,
112 }
113 }
114
115 pub(crate) fn started(agent_id: impl Into<String>, agent_type: FleetRole) -> Self {
116 Self::Started {
117 agent_id: agent_id.into(),
118 agent_type: agent_type.as_str().to_string(),
119 }
120 }
121
122 pub(crate) fn progress(agent_id: impl Into<String>, status: impl Into<String>) -> Self {
123 Self::Progress {
124 agent_id: agent_id.into(),
125 status: status.into(),
126 }
127 }
128
129 #[cfg(test)]
130 pub(crate) fn work_state(agent_id: impl Into<String>, todo: TodoListSnapshot) -> Self {
131 Self::WorkState {
132 agent_id: agent_id.into(),
133 todo,
134 }
135 }
136
137 pub(crate) fn token_usage(
138 agent_id: impl Into<String>,
139 source_id: impl Into<String>,
140 route: crate::cost_status::EffectiveRouteEnvelope,
141 usage: Usage,
142 ) -> Self {
143 Self::TokenUsage {
144 agent_id: agent_id.into(),
145 source_id: source_id.into(),
146 route: Box::new(route),
147 usage,
148 }
149 }
150 }
151
152 /// One delivery: a sequence number plus the message. The sequence is
153 /// monotonic across the entire mailbox (not per-agent) so a single ordering
154 /// is well-defined even when multiple sub-agents share one mailbox.
155 #[derive(Debug, Clone, PartialEq, Eq)]
156 pub struct MailboxEnvelope {
157 pub seq: u64,
158 pub message: MailboxMessage,
159 }
160
161 /// Sender side of the mailbox.
162 ///
163 /// Cheaply cloneable (everything inside is `Arc`/atomic). Cloning a
164 /// `Mailbox` shares the same delivery channel, sequence counter, watch
165 /// notifier, and close/cancel state — so a child runtime that clones its
166 /// parent's `Mailbox` participates in the same stream.
167 #[derive(Clone)]
168 pub struct Mailbox {
169 inner: Arc<MailboxInner>,
170 }
171
172 struct MailboxInner {
173 tx: mpsc::UnboundedSender<MailboxEnvelope>,
174 next_seq: AtomicU64,
175 seq_tx: watch::Sender<u64>,
176 closed: AtomicBool,
177 /// Linearizes publication against turn-end sealing. Without this gate a
178 /// producer could observe `closed = false`, lose the race to the turn
179 /// completion barrier, and enqueue usage after `TurnComplete`.
180 send_gate: std::sync::Mutex<()>,
181 #[cfg(test)]
182 cancel_token: CancellationToken,
183 }
184
185 /// Receiver side of the mailbox. Not `Clone` — only the original creator
186 /// can drain. Use `Mailbox::subscribe()` for fanout (UI cards + parent both
187 /// observing the same stream).
188 pub struct MailboxReceiver {
189 rx: mpsc::UnboundedReceiver<MailboxEnvelope>,
190 pending: VecDeque<MailboxEnvelope>,
191 }
192
193 impl Mailbox {
194 /// Create a new mailbox bound to the given cancellation token. Closing
195 /// the mailbox (or dropping the last sender) cancels this token. Runtimes
196 /// that derive from the same token observe that cancellation; detached
197 /// background `agent` sessions use their own runtime token.
198 #[must_use]
199 pub fn new(cancel_token: CancellationToken) -> (Self, MailboxReceiver) {
200 #[cfg(not(test))]
201 let _ = cancel_token;
202 let (tx, rx) = mpsc::unbounded_channel();
203 let (seq_tx, _) = watch::channel(0);
204 let inner = MailboxInner {
205 tx,
206 next_seq: AtomicU64::new(0),
207 seq_tx,
208 closed: AtomicBool::new(false),
209 send_gate: std::sync::Mutex::new(()),
210 #[cfg(test)]
211 cancel_token,
212 };
213 (
214 Self {
215 inner: Arc::new(inner),
216 },
217 MailboxReceiver {
218 rx,
219 pending: VecDeque::new(),
220 },
221 )
222 }
223
224 /// Subscribe to seq-bump notifications. Each `recv()` returns when the
225 /// sequence counter advances, signaling new mail without copying it —
226 /// the consumer then calls `drain` (or `recv_one` on its own receiver).
227 /// Multiple subscribers may exist; this is the fanout primitive.
228 #[cfg(test)]
229 #[must_use]
230 pub fn subscribe(&self) -> watch::Receiver<u64> {
231 self.inner.seq_tx.subscribe()
232 }
233
234 /// Send a message; returns `Some(seq)` on success, `None` if the
235 /// mailbox is already closed (callers should treat this as "the
236 /// receiver is gone, stop publishing").
237 pub fn send(&self, message: MailboxMessage) -> Option<u64> {
238 let _send_gate = self
239 .inner
240 .send_gate
241 .lock()
242 .unwrap_or_else(|error| error.into_inner());
243 if self.inner.closed.load(Ordering::Acquire) {
244 return None;
245 }
246 let seq = self.inner.next_seq.fetch_add(1, Ordering::Relaxed) + 1;
247 let envelope = MailboxEnvelope { seq, message };
248 if self.inner.tx.send(envelope).is_err() {
249 return None;
250 }
251 let _ = self.inner.seq_tx.send_replace(seq);
252 Some(seq)
253 }
254
255 /// Stop publication for this turn without cancelling detached workers.
256 ///
257 /// The engine seals, drains, and awaits the mailbox before it emits
258 /// `TurnComplete`. The send gate makes this a hard ordering boundary:
259 /// once this returns, every accepted envelope is already in the receiver
260 /// and no later worker message can attach itself to the completed turn.
261 pub(crate) fn seal(&self) {
262 let _send_gate = self
263 .inner
264 .send_gate
265 .lock()
266 .unwrap_or_else(|error| error.into_inner());
267 self.inner.closed.store(true, Ordering::Release);
268 }
269
270 /// Whether the mailbox has been closed.
271 #[cfg(test)]
272 #[must_use]
273 pub fn is_closed(&self) -> bool {
274 self.inner.closed.load(Ordering::Acquire)
275 }
276
277 /// Close the mailbox AND cancel the bound cancellation token.
278 ///
279 /// "Close-as-cancel": there's no useful state where the consumer is gone
280 /// but producers bound to this mailbox token should keep publishing.
281 /// Closing cancels the bound token; directly derived `child_runtime()`
282 /// children observe it, while detached `agent` sessions rely on their
283 /// own explicit cancellation.
284 #[cfg(test)]
285 pub fn close(&self) {
286 let was_closed = self.inner.closed.load(Ordering::Acquire);
287 self.seal();
288 if !was_closed {
289 self.inner.cancel_token.cancel();
290 }
291 }
292 }
293
294 impl MailboxReceiver {
295 #[cfg(test)]
296 fn sync_pending(&mut self) {
297 while let Ok(env) = self.rx.try_recv() {
298 self.pending.push_back(env);
299 }
300 }
301
302 /// Whether any envelopes are buffered (or arrived since last check).
303 #[cfg(test)]
304 pub fn has_pending(&mut self) -> bool {
305 self.sync_pending();
306 !self.pending.is_empty()
307 }
308
309 /// Drain all currently available envelopes, in delivery order.
310 #[cfg(test)]
311 pub fn drain(&mut self) -> Vec<MailboxEnvelope> {
312 self.sync_pending();
313 self.pending.drain(..).collect()
314 }
315
316 /// Await the next envelope, with backpressure-aware blocking. Returns
317 /// `None` when every sender has been dropped and the buffer is drained.
318 pub async fn recv(&mut self) -> Option<MailboxEnvelope> {
319 if let Some(env) = self.pending.pop_front() {
320 return Some(env);
321 }
322 self.rx.recv().await
323 }
324
325 /// Drain all envelopes accepted before a mailbox was sealed.
326 pub(crate) fn drain_available(&mut self) -> Vec<MailboxEnvelope> {
327 while let Ok(envelope) = self.rx.try_recv() {
328 self.pending.push_back(envelope);
329 }
330 self.pending.drain(..).collect()
331 }
332
333 /// Awaits the next envelope with a timeout. Useful in tests.
334 #[cfg(test)]
335 pub async fn recv_timeout(&mut self, timeout: Duration) -> Option<MailboxEnvelope> {
336 tokio::time::timeout(timeout, self.recv())
337 .await
338 .ok()
339 .flatten()
340 }
341 }
342
343 #[cfg(test)]
344 mod tests {
345 use super::*;
346 use tokio::time::Duration;
347
348 fn open() -> (Mailbox, MailboxReceiver, CancellationToken) {
349 let token = CancellationToken::new();
350 let (mb, rx) = Mailbox::new(token.clone());
351 (mb, rx, token)
352 }
353
354 fn test_route(
355 provider: ProviderKind,
356 model: &str,
357 ) -> crate::cost_status::EffectiveRouteEnvelope {
358 crate::cost_status::EffectiveRouteEnvelope::capture(
359 None,
360 provider,
361 provider.as_str(),
362 model,
363 Some(provider.provider().default_base_url()),
364 chrono::Utc::now(),
365 )
366 }
367
368 #[tokio::test]
369 async fn mailbox_assigns_monotonic_sequence_numbers() {
370 let (mb, _rx, _tok) = open();
371 let s1 = mb
372 .send(MailboxMessage::progress("a", "one"))
373 .expect("seq 1");
374 let s2 = mb
375 .send(MailboxMessage::progress("a", "two"))
376 .expect("seq 2");
377 let s3 = mb
378 .send(MailboxMessage::progress("b", "three"))
379 .expect("seq 3");
380 assert_eq!(s1, 1);
381 assert_eq!(s2, 2);
382 assert_eq!(s3, 3);
383 assert!(s2 > s1 && s3 > s2);
384 }
385
386 #[tokio::test]
387 async fn mailbox_drains_in_delivery_order() {
388 let (mb, mut rx, _tok) = open();
389 mb.send(MailboxMessage::progress("a", "first"));
390 mb.send(MailboxMessage::progress("a", "second"));
391 mb.send(MailboxMessage::Completed {
392 agent_id: "a".into(),
393 summary: "done".into(),
394 });
395 let drained = rx.drain();
396 assert_eq!(drained.len(), 3);
397 assert_eq!(drained[0].seq, 1);
398 assert_eq!(drained[1].seq, 2);
399 assert_eq!(drained[2].seq, 3);
400 assert!(matches!(
401 drained[0].message,
402 MailboxMessage::Progress { .. }
403 ));
404 assert!(matches!(
405 drained[2].message,
406 MailboxMessage::Completed { .. }
407 ));
408 assert!(!rx.has_pending());
409 }
410
411 #[tokio::test]
412 async fn subscribers_receive_seq_bumps_for_backpressure() {
413 let (mb, _rx, _tok) = open();
414 let mut sub_a = mb.subscribe();
415 let mut sub_b = mb.subscribe();
416 // Initial state: both at 0.
417 assert_eq!(*sub_a.borrow(), 0);
418 assert_eq!(*sub_b.borrow(), 0);
419
420 mb.send(MailboxMessage::progress("x", "tick"));
421 sub_a.changed().await.expect("subscriber a sees bump");
422 sub_b.changed().await.expect("subscriber b sees bump");
423 assert_eq!(*sub_a.borrow(), 1);
424 assert_eq!(*sub_b.borrow(), 1);
425
426 // A second send updates both subscribers' watch values too — even
427 // though they share a single watch channel, fanout is N-to-many.
428 mb.send(MailboxMessage::progress("x", "tick2"));
429 sub_a.changed().await.expect("a sees second bump");
430 assert_eq!(*sub_a.borrow(), 2);
431 }
432
433 #[tokio::test]
434 async fn close_cancels_bound_token_and_blocks_further_sends() {
435 let (mb, _rx, token) = open();
436 assert!(!token.is_cancelled());
437 mb.send(MailboxMessage::progress("a", "before close"));
438 mb.close();
439 assert!(token.is_cancelled(), "close-as-cancel: token must fire");
440 assert!(mb.is_closed());
441 // Further sends are no-ops, returning None instead of poisoning seq.
442 assert!(
443 mb.send(MailboxMessage::progress("a", "after close"))
444 .is_none()
445 );
446 }
447
448 #[test]
449 fn turn_end_seal_forms_a_flush_barrier_without_cancelling_worker() {
450 let (mb, mut rx, token) = open();
451 assert_eq!(
452 mb.send(MailboxMessage::progress("a", "accepted before barrier")),
453 Some(1)
454 );
455 mb.seal();
456
457 assert!(!token.is_cancelled(), "detached worker is not cancelled");
458 assert!(
459 mb.send(MailboxMessage::progress("a", "too late")).is_none(),
460 "no event may be accepted after the completion barrier"
461 );
462 let drained = rx.drain_available();
463 assert_eq!(drained.len(), 1);
464 assert_eq!(drained[0].seq, 1);
465 }
466
467 #[tokio::test]
468 async fn close_propagates_to_child_tokens_across_max_spawn_depth() {
469 // Mirror the runtime: root → child → grandchild (default depth 3).
470 let root = CancellationToken::new();
471 let child = root.child_token();
472 let grandchild = child.child_token();
473 let (mb, _rx) = Mailbox::new(root.clone());
474
475 assert!(!child.is_cancelled());
476 assert!(!grandchild.is_cancelled());
477 mb.close();
478 assert!(child.is_cancelled(), "child inherits root close");
479 assert!(
480 grandchild.is_cancelled(),
481 "grandchild inherits too — covers default max_spawn_depth = 3"
482 );
483 }
484
485 #[tokio::test]
486 async fn recv_returns_envelope_then_none_after_close_and_drop() {
487 let (mb, mut rx, _tok) = open();
488 mb.send(MailboxMessage::progress("a", "queued"));
489 let env = rx.recv().await.expect("buffered envelope");
490 assert_eq!(env.seq, 1);
491
492 // After closing AND dropping the sender, recv must yield None.
493 mb.close();
494 drop(mb);
495 let next = rx.recv_timeout(Duration::from_millis(100)).await;
496 assert!(next.is_none(), "drained + dropped → recv yields None");
497 }
498
499 #[tokio::test]
500 async fn cloned_mailbox_shares_sequence_and_close_state() {
501 let (mb, mut rx, token) = open();
502 let mb_clone = mb.clone();
503 let s1 = mb
504 .send(MailboxMessage::progress("a", "from original"))
505 .unwrap();
506 let s2 = mb_clone
507 .send(MailboxMessage::progress("a", "from clone"))
508 .unwrap();
509 assert_eq!(s1, 1);
510 assert_eq!(s2, 2, "clones share the seq counter");
511
512 let drained = rx.drain();
513 assert_eq!(drained.len(), 2);
514
515 // Closing through one clone closes them all (the AtomicBool is shared).
516 mb_clone.close();
517 assert!(mb.is_closed());
518 assert!(token.is_cancelled());
519 }
520
521 #[test]
522 fn work_state_payload_round_trips_and_tolerates_a_missing_snapshot() {
523 use crate::tools::todo::{TodoItem, TodoStatus};
524
525 let message = MailboxMessage::work_state(
526 "agent_child",
527 TodoListSnapshot {
528 items: vec![TodoItem {
529 id: 2,
530 content: "write the projection".to_string(),
531 status: TodoStatus::InProgress,
532 }],
533 completion_pct: 50,
534 in_progress_id: Some(2),
535 },
536 );
537 let encoded = serde_json::to_string(&message).expect("encode");
538 let decoded: MailboxMessage = serde_json::from_str(&encoded).expect("decode");
539 assert_eq!(decoded, message);
540
541 // An older payload that predates per-agent Work state decodes to an
542 // empty snapshot rather than failing the whole stream.
543 let legacy: MailboxMessage =
544 serde_json::from_str(r#"{"kind":"work_state","agent_id":"agent_old"}"#)
545 .expect("legacy");
546 assert_eq!(legacy.agent_id(), "agent_old");
547 match legacy {
548 MailboxMessage::WorkState { todo, .. } => assert!(todo.is_empty()),
549 other => panic!("expected work state, got {other:?}"),
550 }
551
552 // Every pre-existing variant still decodes unchanged.
553 let started: MailboxMessage =
554 serde_json::from_str(r#"{"kind":"started","agent_id":"a","agent_type":"worker"}"#)
555 .expect("started");
556 assert_eq!(started.agent_id(), "a");
557 }
558
559 #[tokio::test]
560 async fn agent_id_is_extractable_from_every_variant() {
561 let cases: Vec<(MailboxMessage, &str)> = vec![
562 (MailboxMessage::started("a1", FleetRole::Worker), "a1"),
563 (MailboxMessage::progress("a2", "x"), "a2"),
564 (
565 MailboxMessage::ToolCallStarted {
566 agent_id: "a3".into(),
567 tool_name: "read_file".into(),
568 step: 1,
569 },
570 "a3",
571 ),
572 (
573 MailboxMessage::ToolCallCompleted {
574 agent_id: "a4".into(),
575 tool_name: "read_file".into(),
576 step: 1,
577 ok: true,
578 },
579 "a4",
580 ),
581 (
582 MailboxMessage::ChildSpawned {
583 parent_id: "parent".into(),
584 child_id: "a5".into(),
585 },
586 "a5",
587 ),
588 (
589 MailboxMessage::Completed {
590 agent_id: "a6".into(),
591 summary: "done".into(),
592 },
593 "a6",
594 ),
595 (
596 MailboxMessage::Failed {
597 agent_id: "a7".into(),
598 error: "boom".into(),
599 },
600 "a7",
601 ),
602 (
603 MailboxMessage::Cancelled {
604 agent_id: "a8".into(),
605 },
606 "a8",
607 ),
608 (
609 MailboxMessage::Interrupted {
610 agent_id: "a10".into(),
611 reason: "API call timed out".into(),
612 },
613 "a10",
614 ),
615 (
616 MailboxMessage::work_state("a11", TodoListSnapshot::default()),
617 "a11",
618 ),
619 (
620 MailboxMessage::TokenUsage {
621 agent_id: "a9".into(),
622 source_id: "response-a9".into(),
623 route: Box::new(test_route(ProviderKind::Deepseek, "deepseek-v4-flash")),
624 usage: Usage {
625 input_tokens: 100,
626 output_tokens: 50,
627 ..Default::default()
628 },
629 },
630 "a9",
631 ),
632 ];
633 for (msg, expected) in cases {
634 assert_eq!(msg.agent_id(), expected, "extract failed for {msg:?}");
635 }
636 }
637
638 #[test]
639 fn token_usage_serde_round_trip_preserves_immutable_route_evidence() {
640 let route = crate::cost_status::EffectiveRouteEnvelope {
641 openrouter_vendor: None,
642 provider: ProviderKind::Moonshot,
643 provider_identity: ProviderKind::Moonshot.as_str().to_string(),
644 model: "k3".to_string(),
645 billing_surface: Some(crate::pricing::MOONSHOT_KIMI_CODE_BILLING_SURFACE.to_string()),
646 endpoint_fingerprint: Some("a".repeat(64)),
647 provider_live_pricing: None,
648 billing_mode: crate::cost_status::RouteBillingMode::Subscription,
649 dispatched_at: chrono::DateTime::<chrono::Utc>::from_timestamp(1_234, 0)
650 .expect("timestamp"),
651 };
652 let message =
653 MailboxMessage::token_usage("agent-k3", "response-k3", route, Usage::default());
654 let json = serde_json::to_string(&message).expect("serialize token usage");
655 let restored: MailboxMessage =
656 serde_json::from_str(&json).expect("deserialize token usage");
657 assert_eq!(restored, message);
658 }
659 }
660
660 lines RUST