返回 CodeWhale
coord.rs
根目录 / crates / tui / src / tools / subagent / coord.rs
1 //! Narrow model-facing agent coordination tools.
2 //!
3 //! Keeps `agent` as the creation surface. These five tools wrap existing
4 //! SubAgentManager / mailbox / checkpoint machinery without restoring the
5 //! retired lifecycle theater (`agent_open` / `agent_eval` / …).
6
7 use std::sync::Arc;
8 use std::time::{Duration, Instant};
9
10 use async_trait::async_trait;
11 use serde_json::{Value, json};
12
13 use super::{
14 COMPLETED_AGENT_RETENTION, NEEDS_PERSON_WAIT_NOTE, ParentMailReceipt, SharedSubAgentManager,
15 SubAgentRuntime, SubAgentStatus, parse_agent_ref, subagent_session_projection,
16 subagent_status_name, take_new_needs_person, wait_for_subagents_from_input,
17 };
18 use crate::tools::registry::ToolRegistryBuilder;
19 use crate::tools::spec::{
20 ApprovalRequirement, ToolCapability, ToolContext, ToolError, ToolResult, ToolSpec,
21 };
22
23 /// Bounds for `agents/wait`. Short on purpose: a blocked wait makes the
24 /// session deaf to typed input, and settled children already report back as
25 /// `<codewhale:subagent.done>` sentinels that start a fresh turn (#4097).
26 /// `pub(crate)` so the description-pinning test can tie the advertised
27 /// numbers to the runtime constants.
28 pub(crate) const COORD_WAIT_DEFAULT_TIMEOUT_SECS: u64 = 30;
29 const COORD_WAIT_MIN_TIMEOUT_SECS: u64 = 1;
30 pub(crate) const COORD_WAIT_MAX_TIMEOUT_SECS: u64 = 120;
31 const COORD_WAIT_CHECK_INTERVAL: Duration = Duration::from_millis(250);
32 const RECENT_PROGRESS_LIMIT: usize = 8;
33 pub(super) const COORDINATION_RECORD_LIMIT: usize = 128;
34 const COORDINATION_INSPECT_LIMIT: usize = 24;
35 pub(super) const COORDINATION_PROJECTION_DECISION_LIMIT: usize = 8;
36 pub(super) const COORDINATION_PROJECTION_BYTE_LIMIT: usize = 4096;
37
38 mod ledger;
39
40 // The ledger types moved to `ledger` unchanged and are re-published here, so
41 // `crate::tools::subagent::coord::{DecisionRecord, …}` still resolves for every
42 // consumer that never had reason to know where the definitions sit — the point
43 // of the split was to shorten two files, not to make four other files import
44 // differently.
45 //
46 // A glob, not a list: several of these types are named only from `cfg(test)`
47 // code in other modules, and an explicit `pub use` of those reads as an unused
48 // import in a release build. The glob also keeps each item's own visibility, so
49 // `MAX_RECONCILIATION_RETRIES` stays reachable here without becoming part of
50 // the module's public surface.
51 pub use ledger::*;
52
53 // ── agents/list ──────────────────────────────────────────────────────────
54
55 pub struct AgentsListTool {
56 manager: SharedSubAgentManager,
57 }
58
59 impl AgentsListTool {
60 #[must_use]
61 pub fn new(manager: SharedSubAgentManager) -> Self {
62 Self { manager }
63 }
64 }
65
66 #[async_trait]
67 impl ToolSpec for AgentsListTool {
68 fn model_visible(&self) -> bool {
69 // #5462: `agent` is the sole model-facing sub-agent surface. These
70 // narrow tools stay registered and executable by name so a persisted
71 // transcript replays byte-for-byte, but they are never advertised in
72 // the catalog and can never be returned by `tool_search` — the same
73 // shape `rlm` and `exec_shell` already use.
74 false
75 }
76
77 fn name(&self) -> &'static str {
78 "agents/list"
79 }
80
81 fn description(&self) -> &'static str {
82 "List child agents: ids, parent hierarchy, state, bounded recent progress, and token budget. Read-only coordination view — does not spawn or wake workers."
83 }
84
85 fn input_schema(&self) -> Value {
86 json!({
87 "type": "object",
88 "properties": {
89 "include_archived": {
90 "type": "boolean",
91 "description": "Include prior-session / archived agents. Default false."
92 },
93 "agent_id": {
94 "type": "string",
95 "description": "Optional single agent id or session name to inspect."
96 }
97 },
98 "required": []
99 })
100 }
101
102 fn capabilities(&self) -> Vec<ToolCapability> {
103 vec![ToolCapability::ReadOnly]
104 }
105
106 fn approval_requirement(&self) -> ApprovalRequirement {
107 ApprovalRequirement::Auto
108 }
109
110 fn is_read_only_for(&self, _input: &Value) -> bool {
111 true
112 }
113
114 fn supports_parallel_for(&self, _input: &Value) -> bool {
115 true
116 }
117
118 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
119 let include_archived = input
120 .get("include_archived")
121 .and_then(Value::as_bool)
122 .unwrap_or(false);
123 let agent_ref = parse_agent_ref(&input)?;
124
125 let mut manager = self.manager.write().await;
126 manager.cleanup_for_session(&context.state_namespace, COMPLETED_AGENT_RETENTION);
127 let summaries = if let Some(agent_ref) = agent_ref {
128 let summary = manager
129 .coordination_summary_for_session(
130 &context.state_namespace,
131 &agent_ref,
132 RECENT_PROGRESS_LIMIT,
133 )
134 .map_err(|err| ToolError::invalid_input(err.to_string()))?;
135 vec![summary]
136 } else {
137 manager.list_coordination_summaries_for_session(
138 &context.state_namespace,
139 include_archived,
140 RECENT_PROGRESS_LIMIT,
141 )
142 };
143 drop(manager);
144
145 let payload = json!({
146 "action": "list",
147 "count": summaries.len(),
148 "agents": summaries,
149 });
150 let mut tool_result = ToolResult::json(&payload)
151 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
152 tool_result.metadata = Some(json!({
153 "action": "list",
154 "count": summaries.len(),
155 }));
156 Ok(tool_result)
157 }
158 }
159
160 // ── agents/message ───────────────────────────────────────────────────────
161
162 pub struct AgentsMessageTool {
163 manager: SharedSubAgentManager,
164 caller_agent_id: Option<String>,
165 }
166
167 impl AgentsMessageTool {
168 #[must_use]
169 pub fn new(manager: SharedSubAgentManager) -> Self {
170 Self {
171 manager,
172 caller_agent_id: None,
173 }
174 }
175
176 #[must_use]
177 pub(crate) fn with_optional_caller(mut self, caller_agent_id: Option<String>) -> Self {
178 self.caller_agent_id = caller_agent_id;
179 self
180 }
181 }
182
183 #[async_trait]
184 impl ToolSpec for AgentsMessageTool {
185 fn model_visible(&self) -> bool {
186 // #5462: `agent` is the sole model-facing sub-agent surface. These
187 // narrow tools stay registered and executable by name so a persisted
188 // transcript replays byte-for-byte, but they are never advertised in
189 // the catalog and can never be returned by `tool_search` — the same
190 // shape `rlm` and `exec_shell` already use.
191 false
192 }
193
194 fn name(&self) -> &'static str {
195 "agents/message"
196 }
197
198 fn description(&self) -> &'static str {
199 "Queue a parent message onto a running child without waking it. The message stays queued until a later agents/followup delivers it through the child's live input channel. Use agents/followup directly when you want immediate delivery."
200 }
201
202 fn input_schema(&self) -> Value {
203 json!({
204 "type": "object",
205 "properties": {
206 "agent_id": {
207 "type": "string",
208 "description": "Target child agent id or session name."
209 },
210 "message": {
211 "type": "string",
212 "description": "Message text to queue."
213 }
214 },
215 "required": ["agent_id", "message"]
216 })
217 }
218
219 fn capabilities(&self) -> Vec<ToolCapability> {
220 vec![ToolCapability::RequiresApproval]
221 }
222
223 fn approval_requirement(&self) -> ApprovalRequirement {
224 ApprovalRequirement::Required
225 }
226
227 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
228 let agent_ref =
229 parse_agent_ref(&input)?.ok_or_else(|| ToolError::missing_field("agent_id"))?;
230 let message = input
231 .get("message")
232 .or_else(|| input.get("text"))
233 .and_then(Value::as_str)
234 .map(str::trim)
235 .filter(|s| !s.is_empty())
236 .ok_or_else(|| ToolError::missing_field("message"))?
237 .to_string();
238
239 let receipt = {
240 let mut manager = self.manager.write().await;
241 manager
242 .ensure_caller_controls_descendant_for_session(
243 &context.state_namespace,
244 &agent_ref,
245 self.caller_agent_id.as_deref(),
246 "agents/message",
247 )
248 .map_err(|err| ToolError::invalid_input(err.to_string()))?;
249 manager
250 .queue_running_parent_message_for_session(
251 &context.state_namespace,
252 &agent_ref,
253 message,
254 )
255 .map_err(|err| ToolError::invalid_input(err.to_string()))?
256 };
257
258 let payload = json!({
259 "action": "message",
260 "agent_id": receipt.agent_id,
261 "queued": true,
262 "woke": false,
263 "queue_depth": receipt.queue_depth,
264 "status": receipt.status,
265 "note": "Message queued without waking the child.",
266 });
267 let mut tool_result = ToolResult::json(&payload)
268 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
269 tool_result.metadata = Some(json!({
270 "action": "message",
271 "agent_id": receipt.agent_id,
272 "woke": false,
273 "queue_depth": receipt.queue_depth,
274 }));
275 Ok(tool_result)
276 }
277 }
278
279 // ── agents/followup ──────────────────────────────────────────────────────
280
281 pub struct AgentsFollowupTool {
282 manager: SharedSubAgentManager,
283 caller_agent_id: Option<String>,
284 /// Runtime for checkpoint resume. `None` (legacy/test construction)
285 /// keeps the queue-only followup behavior.
286 runtime: Option<SubAgentRuntime>,
287 }
288
289 impl AgentsFollowupTool {
290 #[must_use]
291 pub fn new(manager: SharedSubAgentManager) -> Self {
292 Self {
293 manager,
294 caller_agent_id: None,
295 runtime: None,
296 }
297 }
298
299 #[must_use]
300 pub fn with_runtime(mut self, runtime: SubAgentRuntime) -> Self {
301 self.runtime = Some(runtime);
302 self
303 }
304
305 #[must_use]
306 pub(crate) fn with_optional_caller(mut self, caller_agent_id: Option<String>) -> Self {
307 self.caller_agent_id = caller_agent_id;
308 self
309 }
310 }
311
312 impl AgentsFollowupTool {
313 async fn followup_one(
314 &self,
315 agent_ref: &str,
316 message: &str,
317 context: &ToolContext,
318 ) -> Result<Value, ToolError> {
319 let mut manager = self.manager.write().await;
320 let (source, target) = manager
321 .continuation_target_for_caller(
322 &context.state_namespace,
323 agent_ref,
324 self.caller_agent_id.as_deref(),
325 "agents/followup",
326 )
327 .map_err(|error| ToolError::invalid_input(error.to_string()))?;
328 let snapshot = manager
329 .get_result(&target)
330 .map_err(|error| ToolError::invalid_input(error.to_string()))?;
331 let resumed_already = source != target;
332 let receipt = if super::subagent_checkpoint_is_continuable(&snapshot)
333 && self.runtime.is_some()
334 {
335 let snapshot = manager
336 .resume_from_checkpoint_for_session(
337 &context.state_namespace,
338 Arc::clone(&self.manager),
339 self.runtime.clone().expect("runtime checked"),
340 &target,
341 message,
342 )
343 .map_err(|error| ToolError::execution_failed(error.to_string()))?;
344 ParentMailReceipt {
345 agent_id: snapshot.agent_id.clone(),
346 status: subagent_status_name(&snapshot.status).to_string(),
347 queue_depth: 0,
348 woke: true,
349 continued_from_checkpoint: true,
350 continuation_handle: None,
351 note: format!(
352 "resumed from checkpoint {source} as {}; original receipt retained",
353 snapshot.agent_id
354 ),
355 }
356 } else if resumed_already && snapshot.status != SubAgentStatus::Running {
357 ParentMailReceipt {
358 agent_id: target, status: subagent_status_name(&snapshot.status).to_string(),
359 queue_depth: 0, woke: false, continued_from_checkpoint: true,
360 continuation_handle: None,
361 note: "Existing continuation has settled; no duplicate worker was started and no message was delivered.".to_string(),
362 }
363 } else {
364 manager
365 .followup_child_for_session(&context.state_namespace, &target, message.to_string())
366 .map_err(|error| ToolError::invalid_input(error.to_string()))?
367 };
368 let child_route = manager
369 .get_worker_record_for_session(&context.state_namespace, &receipt.agent_id)
370 .and_then(|record| record.spec.child_route);
371 Ok(json!({
372 "action": "followup", "from": source, "to": receipt.agent_id,
373 "agent_id": receipt.agent_id, "queued": receipt.woke || receipt.queue_depth > 0,
374 "woke": receipt.woke, "queue_depth": receipt.queue_depth, "status": receipt.status,
375 "continued_from_checkpoint": receipt.continued_from_checkpoint || resumed_already,
376 "continuation_handle": receipt.continuation_handle, "note": receipt.note,
377 "child_route": child_route,
378 }))
379 }
380 }
381
382 #[async_trait]
383 impl ToolSpec for AgentsFollowupTool {
384 fn model_visible(&self) -> bool {
385 // #5462: `agent` is the sole model-facing sub-agent surface. These
386 // narrow tools stay registered and executable by name so a persisted
387 // transcript replays byte-for-byte, but they are never advertised in
388 // the catalog and can never be returned by `tool_search` — the same
389 // shape `rlm` and `exec_shell` already use.
390 false
391 }
392
393 fn name(&self) -> &'static str {
394 "agents/followup"
395 }
396
397 fn description(&self) -> &'static str {
398 "Queue a message and attempt to resume an idle or interrupted child. Running children receive the message on their next step; interrupted_continuable children are resumed from their checkpoint into a fresh agent loop (new agent id, original prompt plus prior conversation tail) when a runtime is attached, and otherwise keep queue-only semantics with the continuation_handle returned."
399 }
400
401 fn input_schema(&self) -> Value {
402 json!({
403 "type": "object",
404 "properties": {
405 "agent_id": {"type": "string"},
406 "agent_ids": {"type": "array", "items": {"type": "string", "minLength": 1}, "minItems": 1, "maxItems": 32},
407 "all_parked": {"type": "boolean", "description": "Continue every owned parked child, up to 32."},
408 "message": {"type": "string", "minLength": 1}
409 },
410 "required": ["message"],
411 "oneOf": [{"required": ["agent_id"]}, {"required": ["agent_ids"]}, {"required": ["all_parked"]}]
412 })
413 }
414
415 fn capabilities(&self) -> Vec<ToolCapability> {
416 vec![ToolCapability::RequiresApproval]
417 }
418
419 fn approval_requirement(&self) -> ApprovalRequirement {
420 ApprovalRequirement::Required
421 }
422
423 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
424 let message = input
425 .get("message")
426 .or_else(|| input.get("text"))
427 .and_then(Value::as_str)
428 .map(str::trim)
429 .filter(|value| !value.is_empty())
430 .ok_or_else(|| ToolError::missing_field("message"))?;
431 let single = parse_agent_ref(&input)?;
432 let all_parked = super::parse_optional_bool(&input, &["all_parked"])?.unwrap_or(false);
433 let batch = input.get("agent_ids");
434 if usize::from(single.is_some()) + usize::from(batch.is_some()) + usize::from(all_parked)
435 != 1
436 {
437 return Err(ToolError::invalid_input(
438 "followup requires exactly one of agent_id, agent_ids, or all_parked=true",
439 ));
440 }
441 let mut targets = if let Some(ids) = batch {
442 let ids = ids
443 .as_array()
444 .filter(|ids| !ids.is_empty() && ids.len() <= 32)
445 .ok_or_else(|| {
446 ToolError::invalid_input("agent_ids must contain 1..32 nonempty strings")
447 })?;
448 ids.iter()
449 .map(|id| {
450 id.as_str()
451 .map(str::trim)
452 .filter(|id| !id.is_empty())
453 .map(str::to_string)
454 .ok_or_else(|| {
455 ToolError::invalid_input("agent_ids must contain nonempty strings")
456 })
457 })
458 .collect::<Result<Vec<_>, _>>()?
459 } else if let Some(id) = single.as_ref() {
460 vec![id.clone()]
461 } else {
462 let manager = self.manager.read().await;
463 let mut ids = manager
464 .agents
465 .values()
466 .filter(|agent| {
467 manager.agent_is_owned_by_session(agent, &context.state_namespace)
468 && agent
469 .checkpoint
470 .as_ref()
471 .is_some_and(|checkpoint| checkpoint.parked_at_turn_end)
472 && matches!(agent.status, SubAgentStatus::Interrupted(_))
473 && manager
474 .ensure_caller_controls_descendant(
475 &agent.id,
476 self.caller_agent_id.as_deref(),
477 "agents/followup",
478 )
479 .is_ok()
480 && manager
481 .continuation_target(&agent.id)
482 .is_ok_and(|target| target == agent.id)
483 })
484 .map(|agent| agent.id.clone())
485 .collect::<Vec<_>>();
486 ids.sort();
487 if ids.len() > 32 {
488 return Err(ToolError::invalid_input(
489 "More than 32 parked children; use explicit agent_ids batches",
490 ));
491 }
492 ids
493 };
494 let mut seen = std::collections::HashSet::new();
495 targets.retain(|target| seen.insert(target.clone()));
496 let mut results = Vec::new();
497 let mut errors = Vec::new();
498 // Each mutation and both hierarchy checks share the manager write lock.
499 // A target failure cannot erase successful results from another target.
500 for target in targets {
501 match self.followup_one(&target, message, context).await {
502 Ok(payload) => results.push(payload),
503 Err(error) if single.is_some() => return Err(error),
504 Err(error) => errors.push(json!({"from": target, "error": error.to_string()})),
505 }
506 }
507 let payload = if single.is_some() {
508 results.pop().expect("single target returned a result")
509 } else {
510 json!({"action": "followup", "results": results, "errors": errors})
511 };
512 let mut result = ToolResult::json(&payload)
513 .map_err(|error| ToolError::execution_failed(error.to_string()))?;
514 result.metadata = Some(payload.clone());
515 Ok(result)
516 }
517 }
518
519 // ── agents/interrupt ─────────────────────────────────────────────────────
520
521 pub struct AgentsInterruptTool {
522 manager: SharedSubAgentManager,
523 /// Optional caller identity for fail-closed self-interrupt checks.
524 caller_agent_id: Option<String>,
525 }
526
527 impl AgentsInterruptTool {
528 #[must_use]
529 pub fn new(manager: SharedSubAgentManager) -> Self {
530 Self {
531 manager,
532 caller_agent_id: None,
533 }
534 }
535
536 #[must_use]
537 #[allow(dead_code)] // arms self-interrupt fail-closed when child registries thread caller (P1.2)
538 pub fn with_caller(mut self, caller_agent_id: impl Into<String>) -> Self {
539 self.caller_agent_id = Some(caller_agent_id.into());
540 self
541 }
542
543 #[must_use]
544 pub(crate) fn with_optional_caller(mut self, caller_agent_id: Option<String>) -> Self {
545 self.caller_agent_id = caller_agent_id;
546 self
547 }
548 }
549
550 #[async_trait]
551 impl ToolSpec for AgentsInterruptTool {
552 fn model_visible(&self) -> bool {
553 // #5462: `agent` is the sole model-facing sub-agent surface. These
554 // narrow tools stay registered and executable by name so a persisted
555 // transcript replays byte-for-byte, but they are never advertised in
556 // the catalog and can never be returned by `tool_search` — the same
557 // shape `rlm` and `exec_shell` already use.
558 false
559 }
560
561 fn name(&self) -> &'static str {
562 "agents/interrupt"
563 }
564
565 fn description(&self) -> &'static str {
566 "Interrupt a running child agent, preserve its checkpoint, and return the prior state. Fails closed on root or self targets. Prefer this over cancel when you may resume later."
567 }
568
569 fn input_schema(&self) -> Value {
570 json!({
571 "type": "object",
572 "properties": {
573 "agent_id": {
574 "type": "string",
575 "description": "Child agent id or session name to interrupt."
576 },
577 "reason": {
578 "type": "string",
579 "description": "Optional interrupt reason recorded on the checkpoint."
580 }
581 },
582 "required": ["agent_id"]
583 })
584 }
585
586 fn capabilities(&self) -> Vec<ToolCapability> {
587 vec![ToolCapability::RequiresApproval]
588 }
589
590 fn approval_requirement(&self) -> ApprovalRequirement {
591 ApprovalRequirement::Required
592 }
593
594 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
595 let agent_ref =
596 parse_agent_ref(&input)?.ok_or_else(|| ToolError::missing_field("agent_id"))?;
597 let reason = input
598 .get("reason")
599 .and_then(Value::as_str)
600 .map(str::trim)
601 .filter(|s| !s.is_empty())
602 .unwrap_or("interrupted by parent via agents/interrupt")
603 .to_string();
604
605 let (prior, snapshot) = {
606 let mut manager = self.manager.write().await;
607 manager
608 .interrupt_child_for_session(
609 &context.state_namespace,
610 &agent_ref,
611 self.caller_agent_id.as_deref(),
612 reason,
613 )
614 .map_err(|err| ToolError::invalid_input(err.to_string()))?
615 };
616
617 let snapshot = super::settle_requested_child(&self.manager, snapshot).await;
618 let worker_record = {
619 let manager = self.manager.read().await;
620 manager.get_worker_record_for_session(&context.state_namespace, &snapshot.agent_id)
621 };
622 let projection =
623 subagent_session_projection(&self.manager, snapshot, false, context, worker_record)
624 .await;
625 let payload = json!({
626 "action": "interrupt",
627 "agent_id": projection.agent_id,
628 "prior_status": subagent_status_name(&prior.status),
629 "prior_steps_taken": prior.steps_taken,
630 "status": projection.status,
631 "checkpoint_preserved": projection.checkpoint.is_some(),
632 "continuable": projection.continuable,
633 "projection": projection,
634 "child_route": projection.child_route,
635 });
636 let mut tool_result = ToolResult::json(&payload)
637 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
638 tool_result.metadata = Some(json!({
639 "action": "interrupt",
640 "agent_id": payload["agent_id"],
641 "checkpoint_preserved": payload["checkpoint_preserved"],
642 "child_route": payload["child_route"],
643 }));
644 Ok(tool_result)
645 }
646 }
647
648 // ── agents/wait ──────────────────────────────────────────────────────────
649
650 pub struct AgentsWaitTool {
651 manager: SharedSubAgentManager,
652 }
653
654 impl AgentsWaitTool {
655 #[must_use]
656 pub fn new(manager: SharedSubAgentManager) -> Self {
657 Self { manager }
658 }
659 }
660
661 #[async_trait]
662 impl ToolSpec for AgentsWaitTool {
663 fn model_visible(&self) -> bool {
664 // #5462: `agent` is the sole model-facing sub-agent surface. These
665 // narrow tools stay registered and executable by name so a persisted
666 // transcript replays byte-for-byte, but they are never advertised in
667 // the catalog and can never be returned by `tool_search` — the same
668 // shape `rlm` and `exec_shell` already use.
669 false
670 }
671
672 fn name(&self) -> &'static str {
673 "agents/wait"
674 }
675
676 fn description(&self) -> &'static str {
677 "Block briefly until one child settles or timeout_secs (default 30, max 120) elapses; on timeout the receipt reports timed_out=true and any settled children. Keep waits short: on timeout, end your turn — settled children wake you automatically as completion sentinels; polling agents/list in a loop is not the right shape either. until=all is the fan-out join: it returns only when every child running at call time has left running, with each child's outcome. until=completion (default) returns as soon as any one child settles. until=activity also returns on progress."
678 }
679
680 fn input_schema(&self) -> Value {
681 json!({
682 "type": "object",
683 "properties": {
684 "agent_id": {
685 "type": "string",
686 "description": "Optional specific child. When omitted, watches every child running at call time."
687 },
688 "timeout_secs": {
689 "type": "integer",
690 "minimum": 1,
691 "maximum": 120,
692 "description": "Maximum seconds to block. Default 30. Keep it short — on timeout, end your turn; settled children report back as completion sentinels."
693 },
694 "until": {
695 "type": "string",
696 "enum": ["completion", "all", "activity"],
697 "description": "completion (default): return when any one child leaves running. all: return only when every watched child has left running — use this after a fan-out so one wait covers the whole batch. activity: also return when recent progress changes. Children spawned after the call are not watched; no children means an immediate return."
698 }
699 },
700 "required": []
701 })
702 }
703
704 fn capabilities(&self) -> Vec<ToolCapability> {
705 vec![ToolCapability::ReadOnly]
706 }
707
708 fn approval_requirement(&self) -> ApprovalRequirement {
709 ApprovalRequirement::Auto
710 }
711
712 fn is_read_only_for(&self, _input: &Value) -> bool {
713 true
714 }
715
716 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
717 dispatch_wait(&input, Arc::clone(&self.manager), context).await
718 }
719 }
720
721 /// Single entry point for every blocking wait, shared by `agents/wait` and
722 /// `agent(action="wait")` so the two surfaces cannot drift.
723 ///
724 /// `until` selects the join shape:
725 /// - `completion` (default) — return as soon as any one watched child settles.
726 /// - `all` — return only when every watched child has settled (the fan-out
727 /// join the parent should use after dispatching a batch).
728 /// - `activity` — also return when a running child makes visible progress.
729 pub(super) async fn dispatch_wait(
730 input: &Value,
731 manager: SharedSubAgentManager,
732 context: &ToolContext,
733 ) -> Result<ToolResult, ToolError> {
734 let until = input
735 .get("until")
736 .and_then(Value::as_str)
737 .unwrap_or("completion")
738 .trim()
739 .to_ascii_lowercase();
740
741 match until.as_str() {
742 "" | "completion" => {
743 let mut wait_input = input.clone();
744 if wait_input.get("action").is_none() {
745 wait_input["action"] = json!("wait");
746 }
747 wait_for_subagents_from_input(&wait_input, manager, context).await
748 }
749 "all" => wait_for_all_children(input, manager, context).await,
750 "activity" => wait_for_activity(input, manager, context).await,
751 other => Err(ToolError::invalid_input(format!(
752 "Invalid until '{other}'. Use completion, all, or activity."
753 ))),
754 }
755 }
756
757 /// `until=all`: block until every child that was running when the call was
758 /// made has left `Running`.
759 ///
760 /// The watch set is fixed at call time. A child spawned while this wait is
761 /// blocked is deliberately **not** joined — the parent asked to join the batch
762 /// it had just dispatched, and silently extending the set would make the call
763 /// unbounded in a way the caller never asked for. Callers that fan out again
764 /// simply issue another wait.
765 ///
766 /// Cancel-safe (no lock is held across an await), honours `timeout_secs`, and
767 /// returns immediately with `all_settled: true` when nothing is running.
768 async fn wait_for_all_children(
769 input: &Value,
770 manager: SharedSubAgentManager,
771 context: &ToolContext,
772 ) -> Result<ToolResult, ToolError> {
773 let timeout_secs = input
774 .get("timeout_secs")
775 .or_else(|| input.get("timeout"))
776 .and_then(Value::as_u64)
777 .unwrap_or(COORD_WAIT_DEFAULT_TIMEOUT_SECS)
778 .clamp(COORD_WAIT_MIN_TIMEOUT_SECS, COORD_WAIT_MAX_TIMEOUT_SECS);
779 let timeout = Duration::from_secs(timeout_secs);
780 let agent_ref = parse_agent_ref(input)?;
781
782 // Resolve the watch set up front so a bad reference fails immediately
783 // rather than blocking for the whole timeout.
784 let watched: Vec<String> = {
785 let manager = manager.read().await;
786 if let Some(agent_ref) = &agent_ref {
787 let snapshot = manager
788 .get_result_by_ref_for_session(&context.state_namespace, agent_ref)
789 .map_err(|err| ToolError::invalid_input(err.to_string()))?;
790 if snapshot.status != SubAgentStatus::Running {
791 // Already settled: hand back its outcome rather than an empty
792 // "nothing to join" that hides what the caller asked about.
793 let settled = json!({
794 "agent_id": snapshot.agent_id,
795 "name": snapshot.name,
796 "status": subagent_status_name(&snapshot.status),
797 "steps_taken": snapshot.steps_taken,
798 });
799 drop(manager);
800 return wait_all_payload(&[settled], &[], &[], 0, false);
801 }
802 vec![snapshot.agent_id]
803 } else {
804 manager
805 .list_filtered_for_session(&context.state_namespace, false)
806 .into_iter()
807 .filter(|snapshot| snapshot.status == SubAgentStatus::Running)
808 .map(|snapshot| snapshot.agent_id)
809 .collect()
810 }
811 };
812
813 // Zero children is an immediate return, never a hang.
814 if watched.is_empty() {
815 return wait_all_payload(&[], &[], &[], 0, false);
816 }
817
818 let started = Instant::now();
819 let cancelled = async {
820 match &context.cancel_token {
821 Some(token) => token.cancelled().await,
822 None => std::future::pending().await,
823 }
824 };
825 tokio::pin!(cancelled);
826
827 loop {
828 let (settled, still_running) = {
829 let manager = manager.read().await;
830 let mut settled = Vec::new();
831 let mut still_running = Vec::new();
832 for agent_id in &watched {
833 match manager.get_result_by_ref_for_session(&context.state_namespace, agent_id) {
834 Ok(snapshot) if snapshot.status == SubAgentStatus::Running => {
835 still_running.push(json!({
836 "agent_id": snapshot.agent_id,
837 "name": snapshot.name,
838 "status": "running",
839 }));
840 }
841 Ok(snapshot) => settled.push(json!({
842 "agent_id": snapshot.agent_id,
843 "name": snapshot.name,
844 "status": subagent_status_name(&snapshot.status),
845 "steps_taken": snapshot.steps_taken,
846 })),
847 // A watched child that vanished from the ledger (retention
848 // cleanup) is no longer running; report it rather than
849 // blocking on a record that will never settle.
850 Err(_) => settled.push(json!({
851 "agent_id": agent_id,
852 "status": "gone",
853 })),
854 }
855 }
856 (settled, still_running)
857 };
858
859 if still_running.is_empty() {
860 return wait_all_payload(&settled, &[], &[], started.elapsed().as_millis(), false);
861 }
862 // A child blocked on a person ends the join early (approvals C2).
863 let needs_person = take_new_needs_person(&manager, &watched).await;
864 if !needs_person.is_empty() {
865 return wait_all_payload(
866 &settled,
867 &still_running,
868 &needs_person,
869 started.elapsed().as_millis(),
870 false,
871 );
872 }
873 if started.elapsed() >= timeout {
874 return wait_all_payload(
875 &settled,
876 &still_running,
877 &[],
878 started.elapsed().as_millis(),
879 true,
880 );
881 }
882
883 tokio::select! {
884 biased;
885 () = &mut cancelled => {
886 return Err(ToolError::cancelled(
887 "Wait interrupted by user cancellation before every child settled.".to_string(),
888 ));
889 }
890 () = tokio::time::sleep(COORD_WAIT_CHECK_INTERVAL) => {}
891 }
892 }
893 }
894
895 /// `until=all` result: every watched child with its own outcome, so the parent
896 /// can synthesize from one return instead of re-inspecting each child.
897 fn wait_all_payload(
898 settled: &[Value],
899 still_running: &[Value],
900 needs_person: &[Value],
901 waited_ms: u128,
902 timed_out: bool,
903 ) -> Result<ToolResult, ToolError> {
904 let note = if !needs_person.is_empty() {
905 NEEDS_PERSON_WAIT_NOTE
906 } else if timed_out {
907 "The wait interval ended; the children are still running. You may answer the user or continue other work. Ordinary turn completion keeps them running; results arrive as <codewhale:subagent.done> sentinels. Use followup only when a child actually needs continuation."
908 } else if settled.is_empty() {
909 "No sub-agents were running; nothing to join."
910 } else {
911 "Every watched child has settled. Full results arrive as <codewhale:subagent.done> sentinels — synthesize from those."
912 };
913 let mut payload = json!({
914 "action": "wait",
915 "until": "all",
916 "all_settled": still_running.is_empty(),
917 "settled": settled,
918 "still_running": still_running,
919 "waited_ms": u64::try_from(waited_ms).unwrap_or(u64::MAX),
920 "timed_out": timed_out,
921 "note": note,
922 });
923 if !needs_person.is_empty() {
924 payload["needs_person"] = json!(needs_person);
925 }
926 let mut tool_result =
927 ToolResult::json(&payload).map_err(|err| ToolError::execution_failed(err.to_string()))?;
928 tool_result.metadata = Some(json!({
929 "action": "wait",
930 "until": "all",
931 "all_settled": still_running.is_empty(),
932 "settled": settled.len(),
933 "running": still_running.len(),
934 "timed_out": timed_out,
935 }));
936 Ok(tool_result)
937 }
938
939 async fn wait_for_activity(
940 input: &Value,
941 manager: SharedSubAgentManager,
942 context: &ToolContext,
943 ) -> Result<ToolResult, ToolError> {
944 let timeout_secs = input
945 .get("timeout_secs")
946 .or_else(|| input.get("timeout"))
947 .and_then(Value::as_u64)
948 .unwrap_or(COORD_WAIT_DEFAULT_TIMEOUT_SECS)
949 .clamp(COORD_WAIT_MIN_TIMEOUT_SECS, COORD_WAIT_MAX_TIMEOUT_SECS);
950 let timeout = Duration::from_secs(timeout_secs);
951 let agent_ref = parse_agent_ref(input)?;
952
953 let (watched, baseline): (Vec<String>, Vec<(String, u64)>) = {
954 let manager = manager.read().await;
955 if let Some(agent_ref) = &agent_ref {
956 let snap = manager
957 .get_result_by_ref_for_session(&context.state_namespace, agent_ref)
958 .map_err(|err| ToolError::invalid_input(err.to_string()))?;
959 let fp = manager.activity_fingerprint(&snap.agent_id).unwrap_or(0);
960 if snap.status != SubAgentStatus::Running {
961 let payload = json!({
962 "action": "wait",
963 "until": "activity",
964 "reason": "already_settled",
965 "timed_out": false,
966 "agent_id": snap.agent_id,
967 "status": subagent_status_name(&snap.status),
968 });
969 let mut tool_result = ToolResult::json(&payload)
970 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
971 tool_result.metadata = Some(json!({ "action": "wait", "timed_out": false }));
972 return Ok(tool_result);
973 }
974 (vec![snap.agent_id.clone()], vec![(snap.agent_id, fp)])
975 } else {
976 let running = manager
977 .list_filtered_for_session(&context.state_namespace, false)
978 .into_iter()
979 .filter(|s| s.status == SubAgentStatus::Running)
980 .map(|s| s.agent_id)
981 .collect::<Vec<_>>();
982 let baseline = running
983 .iter()
984 .map(|id| {
985 let fp = manager.activity_fingerprint(id).unwrap_or(0);
986 (id.clone(), fp)
987 })
988 .collect();
989 (running, baseline)
990 }
991 };
992
993 if watched.is_empty() {
994 let payload = json!({
995 "action": "wait",
996 "until": "activity",
997 "note": "No running sub-agents; nothing to wait for.",
998 "timed_out": false,
999 });
1000 let mut tool_result = ToolResult::json(&payload)
1001 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
1002 tool_result.metadata = Some(json!({ "action": "wait", "timed_out": false }));
1003 return Ok(tool_result);
1004 }
1005
1006 let started = Instant::now();
1007 let cancelled = async {
1008 match &context.cancel_token {
1009 Some(token) => token.cancelled().await,
1010 None => std::future::pending().await,
1011 }
1012 };
1013 tokio::pin!(cancelled);
1014
1015 loop {
1016 let outcome = {
1017 let manager = manager.read().await;
1018 let mut settled = Vec::new();
1019 let mut activity = Vec::new();
1020 for (id, base_fp) in &baseline {
1021 if let Ok(snap) =
1022 manager.get_result_by_ref_for_session(&context.state_namespace, id)
1023 {
1024 if snap.status != SubAgentStatus::Running {
1025 settled.push(snap);
1026 continue;
1027 }
1028 let fp = manager.activity_fingerprint(id).unwrap_or(0);
1029 if fp != *base_fp {
1030 activity.push(json!({
1031 "agent_id": id,
1032 "status": "running",
1033 "activity_fingerprint": fp,
1034 }));
1035 }
1036 }
1037 }
1038 (
1039 settled,
1040 activity,
1041 manager.running_count_for_session(&context.state_namespace),
1042 )
1043 };
1044
1045 if !outcome.0.is_empty() || !outcome.1.is_empty() {
1046 let payload = json!({
1047 "action": "wait",
1048 "until": "activity",
1049 "settled": outcome.0.iter().map(|s| json!({
1050 "agent_id": s.agent_id,
1051 "status": subagent_status_name(&s.status),
1052 })).collect::<Vec<_>>(),
1053 "activity": outcome.1,
1054 "running": outcome.2,
1055 "elapsed_ms": started.elapsed().as_millis(),
1056 "timed_out": false,
1057 });
1058 let mut tool_result = ToolResult::json(&payload)
1059 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
1060 tool_result.metadata = Some(json!({
1061 "action": "wait",
1062 "timed_out": false,
1063 "settled": outcome.0.len(),
1064 "activity": outcome.1.len(),
1065 }));
1066 return Ok(tool_result);
1067 }
1068
1069 if started.elapsed() >= timeout {
1070 let payload = json!({
1071 "action": "wait",
1072 "until": "activity",
1073 "settled": [],
1074 "activity": [],
1075 "running": outcome.2,
1076 "elapsed_ms": started.elapsed().as_millis(),
1077 "timed_out": true,
1078 "note": "The wait interval ended without new child activity. Children keep running after ordinary turn completion and report through <codewhale:subagent.done> sentinels.",
1079 });
1080 let mut tool_result = ToolResult::json(&payload)
1081 .map_err(|err| ToolError::execution_failed(err.to_string()))?;
1082 tool_result.metadata = Some(json!({ "action": "wait", "timed_out": true }));
1083 return Ok(tool_result);
1084 }
1085
1086 tokio::select! {
1087 biased;
1088 () = &mut cancelled => {
1089 return Err(ToolError::cancelled(
1090 "Wait interrupted by user cancellation before child activity.".to_string(),
1091 ));
1092 }
1093 () = tokio::time::sleep(COORD_WAIT_CHECK_INTERVAL) => {}
1094 }
1095 }
1096 }
1097
1098 /// Register the narrow coordination tools alongside `agent`.
1099 pub fn register_coordination_tools(
1100 builder: ToolRegistryBuilder,
1101 manager: SharedSubAgentManager,
1102 runtime: SubAgentRuntime,
1103 ) -> ToolRegistryBuilder {
1104 // `runtime.parent_agent_id` is the identity of the agent this registry is
1105 // being built FOR: `runtime_for_nested_agent_tools` stamps the child's own
1106 // id there before `new_with_owner` registers tools, so anything that agent
1107 // spawns records it as parent. Thread that identity through every mutating
1108 // hierarchy tool: a child may control only its own descendants, while the
1109 // root registry (`None`) may control any child (TUI-DOG-017).
1110 let caller = runtime.parent_agent_id.clone();
1111 let message = AgentsMessageTool::new(Arc::clone(&manager)).with_optional_caller(caller.clone());
1112 let followup = AgentsFollowupTool::new(Arc::clone(&manager))
1113 .with_optional_caller(caller.clone())
1114 .with_runtime(runtime.clone());
1115 let interrupt =
1116 AgentsInterruptTool::new(Arc::clone(&manager)).with_optional_caller(caller.clone());
1117 let coordinate = AgentsCoordinateTool::new(Arc::clone(&manager), caller);
1118 builder
1119 .with_tool(Arc::new(AgentsListTool::new(Arc::clone(&manager))))
1120 .with_tool(Arc::new(message))
1121 .with_tool(Arc::new(followup))
1122 .with_tool(Arc::new(interrupt))
1123 .with_tool(Arc::new(coordinate))
1124 .with_tool(Arc::new(AgentsWaitTool::new(manager)))
1125 }
1126
1127 pub struct AgentsCoordinateTool {
1128 manager: SharedSubAgentManager,
1129 caller: Option<String>,
1130 }
1131
1132 impl AgentsCoordinateTool {
1133 #[must_use]
1134 pub fn new(manager: SharedSubAgentManager, caller: Option<String>) -> Self {
1135 Self { manager, caller }
1136 }
1137 fn mutate_coordination(
1138 caller: Option<&str>,
1139 session: &str,
1140 workspace: &std::path::Path,
1141 manager: &mut super::SubAgentManager,
1142 input: Value,
1143 ) -> Result<ToolResult, ToolError> {
1144 let action = input
1145 .get("action")
1146 .and_then(Value::as_str)
1147 .unwrap_or("inspect");
1148 let owner = caller.map(str::to_string).unwrap_or_else(|| "root".into());
1149 let bounded_text = |key: &str| {
1150 input
1151 .get(key)
1152 .and_then(Value::as_str)
1153 .map(|value| value.chars().take(512).collect::<String>())
1154 };
1155 let strings = |key: &str| {
1156 input
1157 .get(key)
1158 .and_then(Value::as_array)
1159 .map(|items| {
1160 items
1161 .iter()
1162 .take(24)
1163 .filter_map(Value::as_str)
1164 .map(|value| value.chars().take(512).collect::<String>())
1165 .collect::<Vec<_>>()
1166 })
1167 .unwrap_or_default()
1168 };
1169 if let Some(caller) = caller {
1170 manager
1171 .get_result_by_ref_for_session(session, caller)
1172 .map_err(|_| {
1173 ToolError::invalid_input("Agent not found in the active session".to_string())
1174 })?;
1175 }
1176 if matches!(action, "accept" | "supersede") {
1177 let decision_id = bounded_text("decision_id").unwrap_or_default();
1178 if !manager.coordination_decision_is_owned_by_session(session, &decision_id) {
1179 return Err(ToolError::invalid_input(
1180 "Coordination decision not found in the active session".to_string(),
1181 ));
1182 }
1183 }
1184 if action == "reconcile"
1185 && strings("input_decisions").iter().any(|decision_id| {
1186 !manager.coordination_decision_is_owned_by_session(session, decision_id)
1187 })
1188 {
1189 return Err(ToolError::invalid_input(
1190 "One or more coordination decisions were not found in the active session"
1191 .to_string(),
1192 ));
1193 }
1194 let coordination_before = manager.coordination.clone();
1195 let mutation = match action {
1196 "propose" => manager
1197 .record_coordination_decision_in_workspace(
1198 DecisionRecord {
1199 decision_id: bounded_text("decision_id").unwrap_or_default(),
1200 subject: bounded_text("subject").unwrap_or_default(),
1201 status: DecisionStatus::Proposed,
1202 owner,
1203 scope: strings("scope"),
1204 constraints: strings("constraints"),
1205 evidence_handles: strings("evidence_handles"),
1206 version: 1,
1207 sequence: 0,
1208 },
1209 workspace,
1210 )
1211 .map_err(ToolError::invalid_input)
1212 .and_then(|record| {
1213 serde_json::to_value(record)
1214 .map_err(|e| ToolError::execution_failed(e.to_string()))
1215 }),
1216 "accept" | "supersede" => input
1217 .get("expected_version")
1218 .and_then(Value::as_u64)
1219 .and_then(|value| u32::try_from(value).ok())
1220 .ok_or_else(|| {
1221 ToolError::invalid_input(
1222 "accept/supersede requires expected_version".to_string(),
1223 )
1224 })
1225 .and_then(|expected_version| {
1226 manager
1227 .update_coordination_decision(
1228 &bounded_text("decision_id").unwrap_or_default(),
1229 if action == "accept" {
1230 DecisionStatus::Accepted
1231 } else {
1232 DecisionStatus::Superseded
1233 },
1234 &owner,
1235 expected_version,
1236 )
1237 .map_err(ToolError::invalid_input)
1238 })
1239 .and_then(|record| {
1240 serde_json::to_value(record)
1241 .map_err(|e| ToolError::execution_failed(e.to_string()))
1242 }),
1243 "claim" => manager
1244 .expand_write_claim(
1245 &owner,
1246 strings("roots"),
1247 strings("exact_files"),
1248 strings("contracts"),
1249 )
1250 .map_err(ToolError::invalid_input)
1251 .and_then(|claim| {
1252 serde_json::to_value(claim)
1253 .map_err(|e| ToolError::execution_failed(e.to_string()))
1254 }),
1255 "reconcile" => manager
1256 .reconcile_coordination(
1257 bounded_text("subject").unwrap_or_default(),
1258 owner,
1259 strings("input_decisions"),
1260 bounded_text("outcome").unwrap_or_default(),
1261 strings("evidence_handles"),
1262 strings("candidate_handles"),
1263 input
1264 .get("retry_count")
1265 .and_then(Value::as_u64)
1266 .and_then(|value| u32::try_from(value).ok())
1267 .unwrap_or_default(),
1268 input
1269 .get("retry_limit")
1270 .and_then(Value::as_u64)
1271 .and_then(|value| u32::try_from(value).ok())
1272 .unwrap_or(MAX_RECONCILIATION_RETRIES),
1273 strings("reviewer_evidence_handles"),
1274 strings("verifier_evidence_handles"),
1275 bounded_text("verification_outcome").unwrap_or_default(),
1276 )
1277 .map_err(ToolError::invalid_input)
1278 .and_then(|receipt| {
1279 serde_json::to_value(receipt)
1280 .map_err(|e| ToolError::execution_failed(e.to_string()))
1281 }),
1282 "release" => {
1283 let owner = input
1284 .get("owner")
1285 .and_then(Value::as_str)
1286 .map(|value| value.to_string());
1287 let released = manager
1288 .release_stale_write_claims(owner)
1289 .map_err(ToolError::invalid_input)?;
1290 let sequence = manager.coordination.sequence;
1291 serde_json::to_value(json!({
1292 "released": released.len(),
1293 "owners": released,
1294 "sequence": sequence
1295 }))
1296 .map_err(|e| ToolError::execution_failed(e.to_string()))
1297 }
1298 _ => unreachable!("coordination action validated above"),
1299 };
1300 let value = match mutation {
1301 Ok(value) => value,
1302 Err(error) => {
1303 // Contention failures deliberately append a durable receipt.
1304 // Stamp and persist every sequence allocated by the failed
1305 // action before returning its error; validation failures that
1306 // did not mutate the ledger allocate nothing.
1307 let first_new_sequence = coordination_before.sequence.saturating_add(1);
1308 let last_new_sequence = manager.coordination.sequence;
1309 for sequence in first_new_sequence..=last_new_sequence {
1310 if let Err(stamp_error) =
1311 manager.stamp_coordination_sequence_for_session(sequence, session)
1312 {
1313 manager.coordination = coordination_before;
1314 return Err(ToolError::execution_failed(format!(
1315 "{error}; additionally failed to stamp coordination receipt: {stamp_error}"
1316 )));
1317 }
1318 }
1319 if last_new_sequence >= first_new_sequence
1320 && let Err(persist_error) = manager.persist_state_synchronously()
1321 {
1322 manager.coordination = coordination_before;
1323 return Err(ToolError::execution_failed(format!(
1324 "{error}; additionally failed to persist coordination receipt: {persist_error}"
1325 )));
1326 }
1327 return Err(error);
1328 }
1329 };
1330 let Some(sequence) = value.get("sequence").and_then(Value::as_u64) else {
1331 manager.coordination = coordination_before;
1332 return Err(ToolError::execution_failed(format!(
1333 "coordination action '{action}' produced no durable sequence"
1334 )));
1335 };
1336 if let Err(error) = manager.stamp_coordination_sequence_for_session(sequence, session) {
1337 manager.coordination = coordination_before;
1338 return Err(ToolError::execution_failed(error));
1339 }
1340 if let Err(error) = manager.persist_state_synchronously() {
1341 manager.coordination = coordination_before;
1342 return Err(ToolError::execution_failed(format!(
1343 "failed to persist coordination action '{action}': {error}"
1344 )));
1345 }
1346 ToolResult::json(&value).map_err(|e| ToolError::execution_failed(e.to_string()))
1347 }
1348 }
1349
1350 #[async_trait]
1351 impl ToolSpec for AgentsCoordinateTool {
1352 fn model_visible(&self) -> bool {
1353 // #5462: `agent` is the sole model-facing sub-agent surface. These
1354 // narrow tools stay registered and executable by name so a persisted
1355 // transcript replays byte-for-byte, but they are never advertised in
1356 // the catalog and can never be returned by `tool_search` — the same
1357 // shape `rlm` and `exec_shell` already use.
1358 false
1359 }
1360
1361 fn name(&self) -> &'static str {
1362 "agents/coordinate"
1363 }
1364
1365 fn description(&self) -> &'static str {
1366 "Record or inspect bounded coordination state: propose/accept/supersede decisions, expand the caller's write claim before mutation, reconcile multiple decision records into one neutral fan-in receipt, or release stale write-claims whose owner is no longer running."
1367 }
1368
1369 fn input_schema(&self) -> Value {
1370 json!({
1371 "type": "object",
1372 "properties": {
1373 "action": { "type": "string", "enum": ["inspect", "propose", "accept", "supersede", "claim", "reconcile", "release"] },
1374 "decision_id": { "type": "string" },
1375 "subject": { "type": "string" },
1376 "expected_version": { "type": "integer", "minimum": 1 },
1377 "scope": { "type": "array", "items": { "type": "string" } },
1378 "constraints": { "type": "array", "items": { "type": "string" } },
1379 "evidence_handles": { "type": "array", "items": { "type": "string" } },
1380 "roots": { "type": "array", "items": { "type": "string" } },
1381 "exact_files": { "type": "array", "items": { "type": "string" } },
1382 "contracts": { "type": "array", "items": { "type": "string" } },
1383 "owner": { "type": "string" },
1384 "input_decisions": { "type": "array", "items": { "type": "string" } },
1385 "outcome": { "type": "string" },
1386 "candidate_handles": { "type": "array", "items": { "type": "string" } },
1387 "retry_count": { "type": "integer", "minimum": 0, "maximum": 3 },
1388 "retry_limit": { "type": "integer", "minimum": 1, "maximum": 3 },
1389 "reviewer_evidence_handles": { "type": "array", "items": { "type": "string" } },
1390 "verifier_evidence_handles": { "type": "array", "items": { "type": "string" } },
1391 "verification_outcome": { "type": "string" },
1392 "limit": { "type": "integer", "minimum": 1, "maximum": 24 }
1393 },
1394 "required": ["action"]
1395 })
1396 }
1397
1398 fn capabilities(&self) -> Vec<ToolCapability> {
1399 // #5123-class: this tool mutates the coordination ledger and expands
1400 // the caller's write claim (actions propose/accept/supersede/claim/
1401 // reconcile) — declaring ReadOnly was a lie that let policy layers
1402 // treat a mutating call as a safe read. Only `inspect` is read-only,
1403 // which is what is_read_only_for reports.
1404 vec![ToolCapability::WritesFiles]
1405 }
1406 fn approval_requirement(&self) -> ApprovalRequirement {
1407 // Stays Auto: coordination records are session-scoped in-memory
1408 // state, and gating them would deadlock autonomous sub-agent fan-in.
1409 ApprovalRequirement::Auto
1410 }
1411 fn is_read_only_for(&self, input: &Value) -> bool {
1412 input.get("action").and_then(Value::as_str) == Some("inspect")
1413 }
1414
1415 async fn execute(&self, input: Value, context: &ToolContext) -> Result<ToolResult, ToolError> {
1416 let action = input
1417 .get("action")
1418 .and_then(Value::as_str)
1419 .unwrap_or("inspect");
1420 let bounded_text = |key: &str| {
1421 input
1422 .get(key)
1423 .and_then(Value::as_str)
1424 .map(|value| value.chars().take(512).collect::<String>())
1425 };
1426 if action == "inspect" {
1427 let manager = self.manager.read().await;
1428 let value = manager.inspect_coordination_for_session(
1429 &context.state_namespace,
1430 bounded_text("subject").as_deref(),
1431 input
1432 .get("limit")
1433 .and_then(Value::as_u64)
1434 .unwrap_or(COORDINATION_INSPECT_LIMIT as u64) as usize,
1435 );
1436 return ToolResult::json(&value)
1437 .map_err(|e| ToolError::execution_failed(e.to_string()));
1438 }
1439 if !matches!(
1440 action,
1441 "propose" | "accept" | "supersede" | "claim" | "reconcile" | "release"
1442 ) {
1443 return Err(ToolError::invalid_input(format!(
1444 "unknown coordination action '{action}'"
1445 )));
1446 }
1447
1448 let mut manager = self.manager.clone().write_owned().await;
1449 if matches!(action, "claim" | "propose") {
1450 let caller = self.caller.clone();
1451 let session = context.state_namespace.clone();
1452 let workspace = context.workspace.clone();
1453 return codewhale_app_server::daemon_socket::owner_work(move || {
1454 Ok(Self::mutate_coordination(
1455 caller.as_deref(),
1456 &session,
1457 &workspace,
1458 &mut manager,
1459 input,
1460 ))
1461 })
1462 .await
1463 .map_err(|error| ToolError::execution_failed(error.to_string()))?;
1464 }
1465 Self::mutate_coordination(
1466 self.caller.as_deref(),
1467 &context.state_namespace,
1468 &context.workspace,
1469 &mut manager,
1470 input,
1471 )
1472 }
1473 }
1474
1475 #[cfg(test)]
1476 mod tests {
1477 use super::*;
1478 use crate::tools::spec::ToolContext;
1479 use codewhale_models::Role;
1480 use std::collections::BTreeSet;
1481 use tempfile::tempdir;
1482
1483 #[test]
1484 fn coordinate_tool_does_not_declare_read_only() {
1485 // #5123-class: the tool mutates the coordination ledger and expands
1486 // write claims; its declared capabilities must not say ReadOnly.
1487 let manager = Arc::new(tokio::sync::RwLock::new(
1488 super::super::SubAgentManager::new(std::path::PathBuf::from("."), 1),
1489 ));
1490 let tool = AgentsCoordinateTool::new(manager, None);
1491 let capabilities = ToolSpec::capabilities(&tool);
1492 assert!(
1493 !capabilities.contains(&ToolCapability::ReadOnly),
1494 "agents/coordinate mutates the ledger — ReadOnly is a lie: {capabilities:?}"
1495 );
1496 // …but the dynamic check still marks inspect as read-only.
1497 assert!(tool.is_read_only_for(&json!({"action": "inspect"})));
1498 assert!(!tool.is_read_only_for(&json!({"action": "propose"})));
1499 }
1500
1501 #[test]
1502 fn coordination_descriptions_match_implemented_resume_behavior() {
1503 // Checkpoint resume is implemented (#5242): the descriptions must
1504 // describe the real behavior, including the honest queue-only
1505 // fallback when no runtime is attached.
1506 let manager = Arc::new(tokio::sync::RwLock::new(
1507 super::super::SubAgentManager::new(std::env::temp_dir(), 1),
1508 ));
1509 let message = AgentsMessageTool::new(Arc::clone(&manager));
1510 let followup = AgentsFollowupTool::new(manager);
1511
1512 assert!(!message.description().contains("natural resume"));
1513 assert!(message.description().contains("stays queued"));
1514 assert!(followup.description().contains("attempt to resume"));
1515 assert!(
1516 followup
1517 .description()
1518 .contains("resumed from their checkpoint")
1519 );
1520 assert!(followup.description().contains("queue-only semantics"));
1521 }
1522
1523 async fn manager_with_running_child(
1524 workspace: &std::path::Path,
1525 ) -> (SharedSubAgentManager, String) {
1526 let manager = Arc::new(tokio::sync::RwLock::new(
1527 super::super::SubAgentManager::new(workspace.to_path_buf(), 4),
1528 ));
1529 let agent_id = {
1530 let mut guard = manager.write().await;
1531 guard.insert_test_running_agent("coord_child", workspace)
1532 };
1533 (manager, agent_id)
1534 }
1535
1536 async fn manager_with_agent_hierarchy(
1537 workspace: &std::path::Path,
1538 ) -> (SharedSubAgentManager, String, String, String) {
1539 let manager = Arc::new(tokio::sync::RwLock::new(
1540 super::super::SubAgentManager::new(workspace.to_path_buf(), 8),
1541 ));
1542 let (parent, child, sibling) = {
1543 let mut guard = manager.write().await;
1544 let parent = guard.insert_test_running_agent("hierarchy_parent", workspace);
1545 let child = guard.insert_test_running_agent("hierarchy_child", workspace);
1546 let sibling = guard.insert_test_running_agent("hierarchy_sibling", workspace);
1547 for (agent_id, parent_id) in [
1548 (&parent, "root"),
1549 (&child, parent.as_str()),
1550 (&sibling, "root"),
1551 ] {
1552 let record = guard
1553 .worker_records
1554 .get_mut(agent_id)
1555 .expect("hierarchy worker record");
1556 record.parent_run_id = Some(parent_id.to_string());
1557 record.spec.parent_run_id = Some(parent_id.to_string());
1558 }
1559 (parent, child, sibling)
1560 };
1561 (manager, parent, child, sibling)
1562 }
1563
1564 #[tokio::test]
1565 async fn message_queues_without_waking() {
1566 let tmp = tempdir().unwrap();
1567 let (manager, agent_id) = manager_with_running_child(tmp.path()).await;
1568 let tool = AgentsMessageTool::new(Arc::clone(&manager));
1569 let result = tool
1570 .execute(
1571 json!({ "agent_id": agent_id, "message": "hold this" }),
1572 &ToolContext::new(tmp.path()),
1573 )
1574 .await
1575 .expect("message ok");
1576 let body: Value = serde_json::from_str(&result.content).unwrap();
1577 assert_eq!(body["woke"], json!(false));
1578 assert_eq!(body["queued"], json!(true));
1579 assert_eq!(body["queue_depth"], json!(1));
1580
1581 let guard = manager.read().await;
1582 let depth = guard.queued_mail_depth(&agent_id).unwrap();
1583 assert_eq!(depth, 1);
1584 assert!(!guard.child_was_woken(&agent_id));
1585 }
1586
1587 #[tokio::test]
1588 async fn followup_does_not_claim_wake_when_live_channel_is_closed() {
1589 let tmp = tempdir().unwrap();
1590 let (manager, agent_id) = manager_with_running_child(tmp.path()).await;
1591 let result = AgentsFollowupTool::new(Arc::clone(&manager))
1592 .execute(
1593 json!({ "agent_id": agent_id, "message": "try to wake" }),
1594 &ToolContext::new(tmp.path()),
1595 )
1596 .await
1597 .expect("truthful closed-channel receipt");
1598 let body: Value = serde_json::from_str(&result.content).unwrap();
1599 assert_eq!(body["woke"], json!(false));
1600 assert_eq!(body["queue_depth"], json!(1));
1601 assert!(
1602 body["note"].as_str().unwrap_or_default().contains("closed"),
1603 "{body}"
1604 );
1605
1606 let guard = manager.read().await;
1607 assert_eq!(guard.queued_mail_depth(&agent_id), Some(1));
1608 assert!(!guard.child_was_woken(&agent_id));
1609 }
1610
1611 #[tokio::test]
1612 async fn hierarchy_mutations_allow_own_descendants_and_deny_siblings_or_ancestors() {
1613 let tmp = tempdir().unwrap();
1614 let (manager, parent, child, sibling) = manager_with_agent_hierarchy(tmp.path()).await;
1615 let context = ToolContext::new(tmp.path());
1616
1617 AgentsMessageTool::new(Arc::clone(&manager))
1618 .with_optional_caller(Some(parent.clone()))
1619 .execute(
1620 json!({ "agent_id": child, "message": "bounded parent note" }),
1621 &context,
1622 )
1623 .await
1624 .expect("parent may message its own child");
1625 AgentsFollowupTool::new(Arc::clone(&manager))
1626 .with_optional_caller(Some(parent.clone()))
1627 .execute(
1628 json!({ "agent_id": child, "message": "resume own child" }),
1629 &context,
1630 )
1631 .await
1632 .expect("parent may follow up its own child");
1633
1634 let sibling_message = AgentsMessageTool::new(Arc::clone(&manager))
1635 .with_optional_caller(Some(parent.clone()))
1636 .execute(
1637 json!({ "agent_id": sibling, "message": "cross branch" }),
1638 &context,
1639 )
1640 .await
1641 .expect_err("sibling message must fail closed")
1642 .to_string();
1643 assert!(
1644 sibling_message.contains("own descendants"),
1645 "{sibling_message}"
1646 );
1647
1648 let ancestor_followup = AgentsFollowupTool::new(Arc::clone(&manager))
1649 .with_optional_caller(Some(child.clone()))
1650 .execute(
1651 json!({ "agent_id": parent, "message": "wake ancestor" }),
1652 &context,
1653 )
1654 .await
1655 .expect_err("ancestor followup must fail closed")
1656 .to_string();
1657 assert!(
1658 ancestor_followup.contains("own descendants"),
1659 "{ancestor_followup}"
1660 );
1661
1662 let sibling_interrupt = AgentsInterruptTool::new(Arc::clone(&manager))
1663 .with_optional_caller(Some(parent.clone()))
1664 .execute(json!({ "agent_id": sibling }), &context)
1665 .await
1666 .expect_err("sibling interrupt must fail closed")
1667 .to_string();
1668 assert!(
1669 sibling_interrupt.contains("own descendants"),
1670 "{sibling_interrupt}"
1671 );
1672
1673 let interrupted = AgentsInterruptTool::new(Arc::clone(&manager))
1674 .with_optional_caller(Some(parent))
1675 .execute(json!({ "agent_id": child }), &context)
1676 .await
1677 .expect("parent may interrupt its own child");
1678 let body: Value = serde_json::from_str(&interrupted.content).unwrap();
1679 assert_eq!(body["status"], json!("interrupted"));
1680 }
1681
1682 #[tokio::test]
1683 async fn coordinate_inspect_is_side_effect_free_and_mutations_are_synchronously_durable() {
1684 let tmp = tempdir().unwrap();
1685 let blocked_state_path = tmp.path().join("blocked-state");
1686 std::fs::create_dir(&blocked_state_path).unwrap();
1687 let blocked_manager = Arc::new(tokio::sync::RwLock::new(
1688 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1689 .with_state_path(blocked_state_path),
1690 ));
1691 let blocked_tool = AgentsCoordinateTool::new(Arc::clone(&blocked_manager), None);
1692
1693 blocked_tool
1694 .execute(
1695 json!({ "action": "inspect" }),
1696 &ToolContext::new(tmp.path()),
1697 )
1698 .await
1699 .expect("read-only inspect must not attempt persistence");
1700 let error = blocked_tool
1701 .execute(
1702 json!({
1703 "action": "propose",
1704 "decision_id": "durable-decision",
1705 "subject": "durability",
1706 "constraints": ["persist before acknowledgement"]
1707 }),
1708 &ToolContext::new(tmp.path()),
1709 )
1710 .await
1711 .expect_err("mutation must fail when its receipt cannot persist")
1712 .to_string();
1713 assert!(error.contains("failed to persist"), "{error}");
1714 assert!(
1715 blocked_manager
1716 .read()
1717 .await
1718 .coordination
1719 .decisions
1720 .is_empty(),
1721 "failed persistence must roll the in-memory decision back"
1722 );
1723
1724 let durable_workspace = tempdir().unwrap();
1725 let state_path = durable_workspace.path().join("subagents.v1.json");
1726 let manager = Arc::new(tokio::sync::RwLock::new(
1727 super::super::SubAgentManager::new(durable_workspace.path().to_path_buf(), 4)
1728 .with_state_path(state_path.clone()),
1729 ));
1730 AgentsCoordinateTool::new(Arc::clone(&manager), None)
1731 .execute(
1732 json!({
1733 "action": "propose",
1734 "decision_id": "durable-decision",
1735 "subject": "durability",
1736 "constraints": ["persist before acknowledgement"]
1737 }),
1738 &ToolContext::new(durable_workspace.path()),
1739 )
1740 .await
1741 .expect("durable mutation");
1742 let mut replayed =
1743 super::super::SubAgentManager::new(durable_workspace.path().to_path_buf(), 4)
1744 .with_state_path(state_path);
1745 replayed.load_state().expect("reload durable action");
1746 assert_eq!(replayed.coordination.decisions.len(), 1);
1747 assert_eq!(
1748 replayed.coordination.decisions[0].decision_id,
1749 "durable-decision"
1750 );
1751 }
1752
1753 #[tokio::test]
1754 async fn release_action_clears_only_stale_write_claims_and_persists() {
1755 let tmp = tempdir().unwrap();
1756 let state_path = tmp.path().join("subagents.v1.json");
1757 let manager = Arc::new(tokio::sync::RwLock::new(
1758 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1759 .with_state_path(state_path.clone()),
1760 ));
1761 let live = {
1762 let mut guard = manager.write().await;
1763 let live = guard.insert_test_running_agent("live-builder", tmp.path());
1764 let active = [live.clone()].into_iter().collect::<BTreeSet<_>>();
1765 for claim in [
1766 WriteScopeClaim {
1767 owner: live.clone(),
1768 roots: vec!["src/live".into()],
1769 exact_files: Vec::new(),
1770 contracts: Vec::new(),
1771 },
1772 WriteScopeClaim {
1773 owner: "zombie-builder".into(),
1774 roots: vec!["src/zombie".into()],
1775 exact_files: Vec::new(),
1776 contracts: Vec::new(),
1777 },
1778 ] {
1779 guard
1780 .coordination
1781 .register_claim(claim, false, |candidate| active.contains(candidate))
1782 .expect("initial claim");
1783 }
1784 let _ = guard.persist_state_synchronously();
1785 live
1786 };
1787
1788 let result = AgentsCoordinateTool::new(Arc::clone(&manager), None)
1789 .execute(
1790 json!({ "action": "release" }),
1791 &ToolContext::new(tmp.path()),
1792 )
1793 .await
1794 .expect("release sweeps stale claims");
1795 let body: Value = serde_json::from_str(&result.content).unwrap();
1796 assert_eq!(body["released"], json!(1));
1797 assert_eq!(body["owners"], json!(["zombie-builder"]));
1798
1799 // The sweep is durable: reloading the ledger keeps only the live claim.
1800 let mut replayed = super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1801 .with_state_path(state_path);
1802 replayed.load_state().expect("reload released ledger");
1803 assert_eq!(replayed.coordination.write_claims.len(), 1);
1804 assert_eq!(replayed.coordination.write_claims[0].claim.owner, live);
1805
1806 // Named-owner release of a live claimant is a no-op and still succeeds.
1807 let noop = AgentsCoordinateTool::new(Arc::clone(&manager), None)
1808 .execute(
1809 json!({ "action": "release", "owner": live }),
1810 &ToolContext::new(tmp.path()),
1811 )
1812 .await
1813 .expect("live release is a no-op");
1814 let body: Value = serde_json::from_str(&noop.content).unwrap();
1815 assert_eq!(body["released"], json!(0));
1816 }
1817
1818 #[tokio::test]
1819 async fn rejected_claim_contention_is_persisted_before_returning_the_error() {
1820 let tmp = tempdir().unwrap();
1821 let state_path = tmp.path().join("subagents.v1.json");
1822 let manager = Arc::new(tokio::sync::RwLock::new(
1823 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1824 .with_state_path(state_path.clone()),
1825 ));
1826 let (claimant, owner) = {
1827 let mut guard = manager.write().await;
1828 let claimant = guard.insert_test_running_agent("claimant", tmp.path());
1829 let owner = guard.insert_test_running_agent("owner", tmp.path());
1830 let active = [claimant.clone(), owner.clone()]
1831 .into_iter()
1832 .collect::<BTreeSet<_>>();
1833 for claim in [
1834 WriteScopeClaim {
1835 owner: claimant.clone(),
1836 roots: vec!["src/claimant".into()],
1837 exact_files: Vec::new(),
1838 contracts: Vec::new(),
1839 },
1840 WriteScopeClaim {
1841 owner: owner.clone(),
1842 roots: vec!["src/shared".into()],
1843 exact_files: Vec::new(),
1844 contracts: Vec::new(),
1845 },
1846 ] {
1847 guard
1848 .coordination
1849 .register_claim(claim, false, |candidate| active.contains(candidate))
1850 .expect("initial non-overlapping claim");
1851 }
1852 (claimant, owner)
1853 };
1854
1855 let error = AgentsCoordinateTool::new(Arc::clone(&manager), Some(claimant.clone()))
1856 .execute(
1857 json!({ "action": "claim", "roots": ["src/shared/nested"] }),
1858 &ToolContext::new(tmp.path()),
1859 )
1860 .await
1861 .expect_err("overlap must block")
1862 .to_string();
1863 assert!(
1864 error.contains(&owner) && error.contains("contention"),
1865 "{error}"
1866 );
1867
1868 let mut replayed = super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1869 .with_state_path(state_path);
1870 replayed.load_state().expect("reload contention receipt");
1871 assert_eq!(replayed.coordination.contentions.len(), 1);
1872 assert_eq!(replayed.coordination.contentions[0].claimant, claimant);
1873 assert_eq!(
1874 replayed.coordination.contentions[0].conflicting_owner,
1875 owner
1876 );
1877 }
1878
1879 #[tokio::test]
1880 async fn coordination_resolution_survives_reload_and_resolving_claim_eviction() {
1881 let tmp = tempdir().unwrap();
1882 let state_path = tmp.path().join("subagents.v1.json");
1883 let manager = Arc::new(tokio::sync::RwLock::new(
1884 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1885 .with_state_path(state_path.clone()),
1886 ));
1887 let claimant = {
1888 let mut guard = manager.write().await;
1889 let claimant = guard.insert_test_running_agent("claimant", tmp.path());
1890 let owner = guard.insert_test_running_agent("owner", tmp.path());
1891 let active = [claimant.clone(), owner.clone()]
1892 .into_iter()
1893 .collect::<BTreeSet<_>>();
1894 for claim in [
1895 WriteScopeClaim {
1896 owner: claimant.clone(),
1897 roots: vec!["src/claimant".into()],
1898 exact_files: Vec::new(),
1899 contracts: Vec::new(),
1900 },
1901 WriteScopeClaim {
1902 owner: owner.clone(),
1903 roots: vec!["src/shared".into()],
1904 exact_files: Vec::new(),
1905 contracts: Vec::new(),
1906 },
1907 ] {
1908 guard
1909 .coordination
1910 .register_claim(claim, false, |candidate| active.contains(candidate))
1911 .expect("initial non-overlapping claim");
1912 }
1913 claimant
1914 };
1915
1916 AgentsCoordinateTool::new(Arc::clone(&manager), Some(claimant.clone()))
1917 .execute(
1918 json!({ "action": "claim", "roots": ["src/shared/nested"] }),
1919 &ToolContext::new(tmp.path()),
1920 )
1921 .await
1922 .expect_err("overlap must block and persist its receipt");
1923
1924 let resolution_sequence = {
1925 let mut guard = manager.write().await;
1926 let record = guard
1927 .coordination
1928 .register_claim(
1929 WriteScopeClaim {
1930 owner: claimant.clone(),
1931 roots: vec!["src/isolated".into()],
1932 exact_files: Vec::new(),
1933 contracts: Vec::new(),
1934 },
1935 true,
1936 |_| true,
1937 )
1938 .expect("later isolated claim resolves contention");
1939 guard
1940 .persist_state_synchronously()
1941 .expect("persist resolved contention");
1942 record.sequence
1943 };
1944
1945 let mut replayed = super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
1946 .with_state_path(state_path.clone());
1947 replayed.load_state().expect("reload resolved contention");
1948 assert_eq!(replayed.coordination.contentions.len(), 1);
1949 assert_eq!(
1950 replayed.coordination.contentions[0].disposition,
1951 WriteContentionDisposition::ResolvedBySuccessfulClaim
1952 );
1953 assert_eq!(
1954 replayed.coordination.contentions[0].resolution_sequence,
1955 Some(resolution_sequence)
1956 );
1957
1958 let slots = COORDINATION_RECORD_LIMIT - replayed.coordination.write_claims.len();
1959 for index in 0..slots {
1960 replayed
1961 .coordination
1962 .register_claim(
1963 WriteScopeClaim {
1964 owner: format!("inactive-fill-{index:03}"),
1965 roots: vec![format!("pkg/fill-{index:03}")],
1966 exact_files: Vec::new(),
1967 contracts: Vec::new(),
1968 },
1969 true,
1970 |_| false,
1971 )
1972 .expect("fill inactive claim capacity");
1973 }
1974 for index in 0..2 {
1975 replayed
1976 .coordination
1977 .register_claim(
1978 WriteScopeClaim {
1979 owner: format!("inactive-overflow-{index}"),
1980 roots: vec![format!("pkg/overflow-{index}")],
1981 exact_files: Vec::new(),
1982 contracts: Vec::new(),
1983 },
1984 true,
1985 |_| false,
1986 )
1987 .expect("evict oldest inactive claim at capacity");
1988 }
1989 assert!(
1990 !replayed
1991 .coordination
1992 .write_claims
1993 .iter()
1994 .any(|claim| claim.claim.owner == claimant),
1995 "the resolving claimant claim must be evicted for the durability regression"
1996 );
1997 replayed
1998 .persist_state_synchronously()
1999 .expect("persist after inactive claim eviction");
2000
2001 let mut final_replay = super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4)
2002 .with_state_path(state_path);
2003 final_replay
2004 .load_state()
2005 .expect("reload after resolving claim eviction");
2006 let projection = final_replay.coordination_detail_projection(None, 24);
2007 assert!(
2008 !projection
2009 .write_claims
2010 .iter()
2011 .any(|claim| claim.claim.owner == claimant)
2012 );
2013 assert_eq!(projection.contentions.len(), 1);
2014 assert_eq!(
2015 projection.contentions[0].disposition,
2016 WriteContentionDisposition::ResolvedBySuccessfulClaim
2017 );
2018 assert_eq!(
2019 projection.contentions[0].resolution_sequence,
2020 Some(resolution_sequence)
2021 );
2022 assert!(!crate::tui::coordination_detail::needs_attention(
2023 &projection
2024 ));
2025 let pager = crate::tui::coordination_detail::format(
2026 codewhale_localization::Locale::En,
2027 &projection,
2028 );
2029 assert!(
2030 pager.contains("disposition resolved_by_successful_claim"),
2031 "{pager}"
2032 );
2033 assert!(!pager.contains("disposition blocked_pending"), "{pager}");
2034 }
2035
2036 #[tokio::test]
2037 async fn interrupt_fails_closed_on_self() {
2038 let tmp = tempdir().unwrap();
2039 let (manager, agent_id) = manager_with_running_child(tmp.path()).await;
2040 let tool = AgentsInterruptTool::new(Arc::clone(&manager)).with_caller(agent_id.clone());
2041 let err = tool
2042 .execute(
2043 json!({ "agent_id": agent_id }),
2044 &ToolContext::new(tmp.path()),
2045 )
2046 .await
2047 .expect_err("self interrupt must fail");
2048 let msg = err.to_string().to_ascii_lowercase();
2049 assert!(
2050 msg.contains("self") || msg.contains("own"),
2051 "unexpected error: {err}"
2052 );
2053 }
2054
2055 #[tokio::test]
2056 async fn interrupt_fails_closed_on_missing_target() {
2057 let tmp = tempdir().unwrap();
2058 let manager = Arc::new(tokio::sync::RwLock::new(
2059 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 2),
2060 ));
2061 let tool = AgentsInterruptTool::new(manager);
2062 let err = tool
2063 .execute(
2064 json!({ "agent_id": "agent_missing" }),
2065 &ToolContext::new(tmp.path()),
2066 )
2067 .await
2068 .expect_err("missing target");
2069 assert!(err.to_string().contains("not found") || err.to_string().contains("Agent"));
2070 }
2071
2072 #[tokio::test]
2073 async fn wait_times_out_when_child_stays_running() {
2074 let tmp = tempdir().unwrap();
2075 let (manager, agent_id) = manager_with_running_child(tmp.path()).await;
2076 let tool = AgentsWaitTool::new(manager);
2077 let result = tool
2078 .execute(
2079 json!({
2080 "agent_id": agent_id,
2081 "timeout_secs": 1,
2082 "until": "activity"
2083 }),
2084 &ToolContext::new(tmp.path()),
2085 )
2086 .await
2087 .expect("wait returns");
2088 let body: Value = serde_json::from_str(&result.content).unwrap();
2089 assert_eq!(body["timed_out"], json!(true));
2090 }
2091
2092 #[tokio::test]
2093 async fn list_resolves_target_and_reports_queue() {
2094 let tmp = tempdir().unwrap();
2095 let (manager, agent_id) = manager_with_running_child(tmp.path()).await;
2096 {
2097 let mut guard = manager.write().await;
2098 guard
2099 .queue_parent_message(&agent_id, "note".into(), false)
2100 .unwrap();
2101 }
2102 let tool = AgentsListTool::new(manager);
2103 let result = tool
2104 .execute(
2105 json!({ "agent_id": agent_id }),
2106 &ToolContext::new(tmp.path()),
2107 )
2108 .await
2109 .expect("list ok");
2110 let body: Value = serde_json::from_str(&result.content).unwrap();
2111 assert_eq!(body["count"], json!(1));
2112 assert_eq!(body["agents"][0]["agent_id"], json!(agent_id));
2113 assert!(body["agents"][0]["queued_mail"].as_u64().unwrap_or(0) >= 1);
2114 }
2115
2116 #[tokio::test]
2117 async fn followup_interrupted_continuable_without_runtime_queues_honestly() {
2118 let tmp = tempdir().unwrap();
2119 let manager = Arc::new(tokio::sync::RwLock::new(
2120 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4),
2121 ));
2122 let (agent_id, handle) = {
2123 let mut guard = manager.write().await;
2124 guard.insert_test_interrupted_continuable_agent(
2125 "paused_child",
2126 tmp.path(),
2127 vec![codewhale_models::Message {
2128 role: Role::User,
2129 content: vec![codewhale_models::ContentBlock::Text {
2130 text: "prior work".to_string(),
2131 cache_control: None,
2132 }],
2133 }],
2134 )
2135 };
2136 // No runtime attached: checkpoint resume is unavailable, so followup
2137 // keeps the honest queue-only semantics with the continuation handle.
2138 let tool = AgentsFollowupTool::new(Arc::clone(&manager));
2139 let result = tool
2140 .execute(
2141 json!({ "agent_id": agent_id, "message": "please continue" }),
2142 &ToolContext::new(tmp.path()),
2143 )
2144 .await
2145 .expect("followup ok");
2146 let body: Value = serde_json::from_str(&result.content).unwrap();
2147 assert_eq!(body["queued"], json!(true));
2148 assert_eq!(body["woke"], json!(false));
2149 assert_eq!(body["continued_from_checkpoint"], json!(false));
2150 assert_eq!(body["continuation_handle"], json!(handle));
2151 let note = body["note"].as_str().unwrap_or_default();
2152 assert!(
2153 note.contains("attach a runtime") && note.contains(&handle),
2154 "note must point at the resume path with the continuation handle: {note}"
2155 );
2156
2157 let guard = manager.read().await;
2158 assert_eq!(guard.queued_mail_depth(&agent_id).unwrap(), 1);
2159 assert!(!guard.child_was_woken(&agent_id));
2160 }
2161
2162 // === until="all": the fan-out join ===================================
2163 //
2164 // Before this existed a parent with five children had to issue five
2165 // waits — while the prompt told it not to poll. These lock the join in.
2166
2167 fn empty_manager(workspace: &std::path::Path) -> SharedSubAgentManager {
2168 Arc::new(tokio::sync::RwLock::new(
2169 super::super::SubAgentManager::new(workspace.to_path_buf(), 8),
2170 ))
2171 }
2172
2173 async fn settle(manager: &SharedSubAgentManager, agent_id: &str, status: SubAgentStatus) {
2174 let mut guard = manager.write().await;
2175 if let Some(agent) = guard.agents.get_mut(agent_id) {
2176 agent.status = status;
2177 }
2178 }
2179
2180 #[test]
2181 fn wait_schema_offers_all_as_a_first_class_until() {
2182 let tmp = tempdir().unwrap();
2183 let tool = AgentsWaitTool::new(empty_manager(tmp.path()));
2184 let schema = tool.input_schema();
2185 let until = &schema["properties"]["until"];
2186 assert_eq!(
2187 until["enum"],
2188 json!(["completion", "all", "activity"]),
2189 "until must expose all alongside completion/activity: {schema}"
2190 );
2191 let described = until["description"].as_str().unwrap_or_default();
2192 assert!(
2193 described.contains("every watched child") && described.contains("any one child"),
2194 "the schema must make completion vs all unmistakable: {described}"
2195 );
2196 }
2197
2198 #[tokio::test]
2199 async fn wait_until_all_on_an_already_settled_child_reports_its_outcome() {
2200 let tmp = tempdir().unwrap();
2201 let manager = empty_manager(tmp.path());
2202 let agent_id = {
2203 let mut guard = manager.write().await;
2204 guard.insert_test_running_agent("all_already_done", tmp.path())
2205 };
2206 settle(&manager, &agent_id, SubAgentStatus::Completed).await;
2207
2208 let result = dispatch_wait(
2209 &json!({ "until": "all", "agent_id": agent_id, "timeout_secs": 60 }),
2210 Arc::clone(&manager),
2211 &ToolContext::new(tmp.path()),
2212 )
2213 .await
2214 .expect("a settled child is an immediate return");
2215 let body: Value = serde_json::from_str(&result.content).unwrap();
2216 assert_eq!(body["all_settled"], json!(true), "{body}");
2217 let settled = body["settled"].as_array().unwrap();
2218 assert_eq!(settled.len(), 1, "{body}");
2219 assert_eq!(settled[0]["status"], json!("completed"), "{body}");
2220 }
2221
2222 #[tokio::test]
2223 async fn foreign_session_is_excluded_from_default_and_explicit_list_waits() {
2224 let tmp = tempdir().unwrap();
2225 let manager = empty_manager(tmp.path());
2226 let agent_a = {
2227 let mut guard = manager.write().await;
2228 let agent_id = guard.insert_test_running_agent("foreign_wait_a", tmp.path());
2229 guard.assign_test_session_owner(&agent_id, "session-a");
2230 agent_id
2231 };
2232 let context_b = ToolContext::new(tmp.path()).with_state_namespace("session-b");
2233
2234 let listed = AgentsListTool::new(Arc::clone(&manager))
2235 .execute(json!({}), &context_b)
2236 .await
2237 .expect("B list");
2238 let listed: Value = serde_json::from_str(&listed.content).unwrap();
2239 assert_eq!(listed["count"], json!(0));
2240
2241 let default_wait = dispatch_wait(
2242 &json!({ "until": "all", "timeout_secs": 60 }),
2243 Arc::clone(&manager),
2244 &context_b,
2245 )
2246 .await
2247 .expect("B default wait has no visible children");
2248 let default_wait: Value = serde_json::from_str(&default_wait.content).unwrap();
2249 assert!(default_wait["settled"].as_array().unwrap().is_empty());
2250
2251 let error = dispatch_wait(
2252 &json!({ "until": "all", "agent_id": agent_a, "timeout_secs": 60 }),
2253 manager,
2254 &context_b,
2255 )
2256 .await
2257 .expect_err("B explicit wait must reject A")
2258 .to_string();
2259 assert!(
2260 error.contains("Agent not found in the active session"),
2261 "{error}"
2262 );
2263 }
2264
2265 #[tokio::test]
2266 async fn wait_until_all_returns_immediately_with_no_children() {
2267 let tmp = tempdir().unwrap();
2268 let started = Instant::now();
2269 let result = dispatch_wait(
2270 &json!({ "until": "all", "timeout_secs": 60 }),
2271 empty_manager(tmp.path()),
2272 &ToolContext::new(tmp.path()),
2273 )
2274 .await
2275 .expect("wait-for-all with zero children must return, not hang");
2276 assert!(
2277 started.elapsed() < Duration::from_secs(5),
2278 "zero children must not burn the timeout"
2279 );
2280 let body: Value = serde_json::from_str(&result.content).unwrap();
2281 assert_eq!(body["all_settled"], json!(true));
2282 assert_eq!(body["timed_out"], json!(false));
2283 assert!(body["settled"].as_array().unwrap().is_empty(), "{body}");
2284 }
2285
2286 #[tokio::test]
2287 async fn wait_until_all_blocks_until_every_child_settles() {
2288 let tmp = tempdir().unwrap();
2289 let manager = empty_manager(tmp.path());
2290 let (first, second, third) = {
2291 let mut guard = manager.write().await;
2292 (
2293 guard.insert_test_running_agent("all_first", tmp.path()),
2294 guard.insert_test_running_agent("all_second", tmp.path()),
2295 guard.insert_test_running_agent("all_third", tmp.path()),
2296 )
2297 };
2298
2299 // Staggered settles: an `until=completion` wait would return after the
2300 // first one. `until=all` must stay blocked through the last.
2301 let flip = Arc::clone(&manager);
2302 let (a, b, c) = (first.clone(), second.clone(), third.clone());
2303 tokio::spawn(async move {
2304 tokio::time::sleep(Duration::from_millis(50)).await;
2305 settle(&flip, &a, SubAgentStatus::Completed).await;
2306 tokio::time::sleep(Duration::from_millis(150)).await;
2307 settle(&flip, &b, SubAgentStatus::Failed("boom".to_string())).await;
2308 tokio::time::sleep(Duration::from_millis(150)).await;
2309 settle(&flip, &c, SubAgentStatus::Cancelled).await;
2310 });
2311
2312 let result = dispatch_wait(
2313 &json!({ "until": "all", "timeout_secs": 30 }),
2314 Arc::clone(&manager),
2315 &ToolContext::new(tmp.path()),
2316 )
2317 .await
2318 .expect("wait-for-all should succeed");
2319 let body: Value = serde_json::from_str(&result.content).unwrap();
2320 assert_eq!(body["all_settled"], json!(true), "{body}");
2321 assert_eq!(body["timed_out"], json!(false), "{body}");
2322 assert!(
2323 body["still_running"].as_array().unwrap().is_empty(),
2324 "{body}"
2325 );
2326
2327 // Per-child outcomes come back on the single return.
2328 let settled = body["settled"].as_array().unwrap();
2329 assert_eq!(settled.len(), 3, "{body}");
2330 let outcomes: std::collections::BTreeMap<&str, &str> = settled
2331 .iter()
2332 .map(|entry| {
2333 (
2334 entry["agent_id"].as_str().unwrap(),
2335 entry["status"].as_str().unwrap(),
2336 )
2337 })
2338 .collect();
2339 assert_eq!(outcomes.get(first.as_str()), Some(&"completed"), "{body}");
2340 assert_eq!(outcomes.get(second.as_str()), Some(&"failed"), "{body}");
2341 assert_eq!(outcomes.get(third.as_str()), Some(&"cancelled"), "{body}");
2342 }
2343
2344 #[tokio::test]
2345 async fn wait_until_all_times_out_reporting_settled_and_still_running() {
2346 let tmp = tempdir().unwrap();
2347 let manager = empty_manager(tmp.path());
2348 let (done, stuck) = {
2349 let mut guard = manager.write().await;
2350 (
2351 guard.insert_test_running_agent("all_done", tmp.path()),
2352 guard.insert_test_running_agent("all_stuck", tmp.path()),
2353 )
2354 };
2355
2356 let request = json!({ "until": "all", "timeout_secs": 1 });
2357 let context = ToolContext::new(tmp.path());
2358 let wait = dispatch_wait(&request, Arc::clone(&manager), &context);
2359 tokio::pin!(wait);
2360 // Capture both running children before completing one. The timeout
2361 // receipt must not depend on a background task winning a 50ms race.
2362 assert!(futures_util::poll!(wait.as_mut()).is_pending());
2363 settle(&manager, &done, SubAgentStatus::Completed).await;
2364
2365 let result = wait
2366 .await
2367 .expect("a timeout is a partial receipt, not an error");
2368 let body: Value = serde_json::from_str(&result.content).unwrap();
2369 assert_eq!(body["timed_out"], json!(true), "{body}");
2370 assert_eq!(body["all_settled"], json!(false), "{body}");
2371
2372 let settled = body["settled"].as_array().unwrap();
2373 assert_eq!(settled.len(), 1, "{body}");
2374 assert_eq!(settled[0]["agent_id"], json!(done), "{body}");
2375 assert_eq!(settled[0]["status"], json!("completed"), "{body}");
2376
2377 let running = body["still_running"].as_array().unwrap();
2378 assert_eq!(running.len(), 1, "{body}");
2379 assert_eq!(running[0]["agent_id"], json!(stuck), "{body}");
2380 }
2381
2382 #[tokio::test]
2383 async fn wait_until_all_ignores_children_spawned_mid_wait() {
2384 let tmp = tempdir().unwrap();
2385 let manager = empty_manager(tmp.path());
2386 let original = {
2387 let mut guard = manager.write().await;
2388 guard.insert_test_running_agent("all_original", tmp.path())
2389 };
2390
2391 let flip = Arc::clone(&manager);
2392 let tmp_path = tmp.path().to_path_buf();
2393 let original_id = original.clone();
2394 tokio::spawn(async move {
2395 tokio::time::sleep(Duration::from_millis(50)).await;
2396 {
2397 let mut guard = flip.write().await;
2398 guard.insert_test_running_agent("all_latecomer", &tmp_path);
2399 }
2400 settle(&flip, &original_id, SubAgentStatus::Completed).await;
2401 });
2402
2403 let result = dispatch_wait(
2404 &json!({ "until": "all", "timeout_secs": 30 }),
2405 Arc::clone(&manager),
2406 &ToolContext::new(tmp.path()),
2407 )
2408 .await
2409 .expect("wait-for-all should succeed");
2410 let body: Value = serde_json::from_str(&result.content).unwrap();
2411 // The watch set is the batch as of call time: the latecomer must not
2412 // extend a wait the caller never asked to include it in.
2413 assert_eq!(body["all_settled"], json!(true), "{body}");
2414 assert_eq!(body["timed_out"], json!(false), "{body}");
2415 let settled = body["settled"].as_array().unwrap();
2416 assert_eq!(settled.len(), 1, "{body}");
2417 assert_eq!(settled[0]["agent_id"], json!(original), "{body}");
2418 }
2419
2420 #[tokio::test]
2421 async fn wait_rejects_unknown_until_naming_every_supported_mode() {
2422 let tmp = tempdir().unwrap();
2423 let error = dispatch_wait(
2424 &json!({ "until": "forever" }),
2425 empty_manager(tmp.path()),
2426 &ToolContext::new(tmp.path()),
2427 )
2428 .await
2429 .expect_err("an unknown until must fail loudly");
2430 let message = error.to_string();
2431 for mode in ["completion", "all", "activity"] {
2432 assert!(message.contains(mode), "{message}");
2433 }
2434 }
2435
2436 #[tokio::test]
2437 async fn followup_interrupted_continuable_resumes_with_runtime() {
2438 let tmp = tempdir().unwrap();
2439 let manager = Arc::new(tokio::sync::RwLock::new(
2440 super::super::SubAgentManager::new(tmp.path().to_path_buf(), 4),
2441 ));
2442 let (agent_id, _handle) = {
2443 let mut guard = manager.write().await;
2444 let interrupted = guard.insert_test_interrupted_continuable_agent(
2445 "paused_child",
2446 tmp.path(),
2447 vec![codewhale_models::Message {
2448 role: Role::User,
2449 content: vec![codewhale_models::ContentBlock::Text {
2450 text: "prior work".to_string(),
2451 cache_control: None,
2452 }],
2453 }],
2454 );
2455 // This fixture is resumed by root, rather than a fabricated parent_session agent.
2456 let record = guard.worker_records.get_mut(&interrupted.0).unwrap();
2457 record.parent_run_id = None;
2458 record.spec.parent_run_id = None;
2459 interrupted
2460 };
2461 let mut runtime = super::super::tests::stub_runtime();
2462 runtime.manager = Arc::clone(&manager);
2463 runtime.context = ToolContext::new(tmp.path());
2464 let tool = AgentsFollowupTool::new(Arc::clone(&manager)).with_runtime(runtime);
2465 let result = tool
2466 .execute(
2467 json!({ "agent_id": agent_id, "message": "please continue" }),
2468 &ToolContext::new(tmp.path()),
2469 )
2470 .await
2471 .expect("followup ok");
2472 let body: Value = serde_json::from_str(&result.content).unwrap();
2473 assert_eq!(body["queued"], json!(true));
2474 assert_eq!(body["woke"], json!(true));
2475 assert_eq!(body["continued_from_checkpoint"], json!(true));
2476 let note = body["note"].as_str().unwrap_or_default();
2477 assert!(note.contains("resumed from checkpoint"), "{note}");
2478 let resumed_id = body["agent_id"].as_str().unwrap_or_default();
2479 assert_ne!(
2480 resumed_id, agent_id,
2481 "resume re-dispatches under a new agent id"
2482 );
2483
2484 // A fresh record exists for the resumed session; the prior terminal
2485 // record stays immutable (receipts are never rewritten).
2486 let guard = manager.read().await;
2487 guard.get_result(resumed_id).expect("resumed agent exists");
2488 let prior = guard.get_result(&agent_id).expect("prior record");
2489 assert!(matches!(prior.status, SubAgentStatus::Interrupted(_)));
2490 }
2491 }
2492
2492 lines RUST