1use std::collections::BTreeMap;
2
3use serde::{Deserialize, Serialize};
4
5use crate::events::{ThreadId, TurnId};
6use crate::extension::SubagentDispatcherId;
7use crate::inference::TokenUsage;
8use crate::trace::SubagentTraceSink;
9
10#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
11pub struct SubagentRequest {
12 pub description: String,
13 pub prompt: String,
14 pub subagent_type: Option<String>,
15 pub model: Option<String>,
18 #[serde(default, skip_serializing_if = "Option::is_none")]
22 pub provider: Option<String>,
23 pub tools: Option<Vec<String>>,
24 #[serde(default, skip_serializing_if = "Option::is_none")]
25 pub lane: Option<SubagentLane>,
26 #[serde(default, skip_serializing_if = "Option::is_none")]
27 pub max_concurrent: Option<usize>,
28 #[serde(default, skip_serializing_if = "Option::is_none")]
29 pub allowed_tools: Option<Vec<String>>,
30 #[serde(default, skip_serializing_if = "Option::is_none")]
31 pub parent_deadline_seconds: Option<u64>,
32 #[serde(default, skip_serializing_if = "Option::is_none")]
33 pub inputs: Option<serde_json::Value>,
34 pub timeout_seconds: Option<u64>,
35}
36
37#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)]
38#[serde(rename_all = "snake_case")]
39pub enum SubagentLane {
40 Scout,
41 Editor,
42 Reviewer,
43 Runner,
44}
45
46impl SubagentLane {
47 pub fn as_str(self) -> &'static str {
48 match self {
49 Self::Scout => "scout",
50 Self::Editor => "editor",
51 Self::Reviewer => "reviewer",
52 Self::Runner => "runner",
53 }
54 }
55
56 pub fn preset(self) -> SubagentLanePreset {
57 match self {
58 Self::Scout => SubagentLanePreset {
59 lane: self,
60 description: "Read and search without changing state.",
61 max_concurrent: 12,
64 timeout_seconds: 120,
65 allowed_tools: &[
66 "Read",
67 "Grep",
68 "Glob",
69 "read_file",
70 "grep",
71 "glob",
72 "list_files",
73 ],
74 },
75 Self::Editor => SubagentLanePreset {
76 lane: self,
77 description: "Make a bounded file-change slice.",
78 max_concurrent: 4,
79 timeout_seconds: 180,
80 allowed_tools: &[
81 "Read",
82 "Grep",
83 "Glob",
84 "read_file",
85 "grep",
86 "glob",
87 "list_files",
88 "write_file",
89 "edit",
90 "multi_edit",
91 "apply_patch",
92 ],
93 },
94 Self::Reviewer => SubagentLanePreset {
95 lane: self,
96 description: "Review and verify with evidence.",
97 max_concurrent: 8,
98 timeout_seconds: 120,
99 allowed_tools: &[
100 "Read",
101 "Grep",
102 "Glob",
103 "read_file",
104 "grep",
105 "glob",
106 "list_files",
107 ],
108 },
109 Self::Runner => SubagentLanePreset {
110 lane: self,
111 description: "Run commands or tests when process policy allows it.",
112 max_concurrent: 1,
113 timeout_seconds: 120,
114 allowed_tools: &["Shell", "shell", "exec_command", "run_command"],
115 },
116 }
117 }
118}
119
120#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
121#[serde(rename_all = "camelCase")]
122pub struct SubagentLanePreset {
123 pub lane: SubagentLane,
124 pub description: &'static str,
125 pub max_concurrent: usize,
126 pub timeout_seconds: u64,
127 pub allowed_tools: &'static [&'static str],
128}
129
130pub fn built_in_subagent_lane_presets() -> [SubagentLanePreset; 4] {
131 [
132 SubagentLane::Scout.preset(),
133 SubagentLane::Editor.preset(),
134 SubagentLane::Reviewer.preset(),
135 SubagentLane::Runner.preset(),
136 ]
137}
138
139pub const SUBAGENT_SUMMARY_CONTRACT: &str = "Child summary must include these labels: Conclusion, Evidence, Files inspected, Files changed, Remaining uncertainty.";
140
141#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
142pub struct SubagentDefinition {
143 pub agent_type: String,
144 pub description: String,
145 pub tools: Vec<String>,
146 pub model: Option<String>,
147 pub system_prompt: Option<String>,
148 pub permission_mode: SubagentPermissionMode,
149 pub max_turns: Option<u32>,
150 pub max_result_chars: Option<usize>,
151}
152
153#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
154#[serde(rename_all = "snake_case")]
155pub enum SubagentPermissionMode {
156 ReadOnly,
157 #[default]
158 Default,
159 AutoEdit,
160}
161
162#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
163pub struct SubagentResult {
164 pub thread_id: ThreadId,
165 pub turn_id: TurnId,
166 pub agent_type: String,
167 pub model: Option<String>,
168 pub final_message: String,
169 pub usage: Option<TokenUsage>,
170 pub exit_reason: SubagentExitReason,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
172 pub transcript: Option<serde_json::Value>,
173 #[serde(default)]
174 pub metadata: serde_json::Value,
175}
176
177#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
178#[serde(rename_all = "snake_case")]
179pub enum SubagentExitReason {
180 Completed,
181 MaxTurns,
182 Timeout,
183 Cancelled,
184 Failed,
185}
186
187#[async_trait::async_trait]
188pub trait SubagentDispatcher: Send + Sync + 'static {
189 fn id(&self) -> SubagentDispatcherId;
190
191 fn definitions(&self) -> Vec<SubagentDefinition>;
192
193 async fn dispatch(
194 &self,
195 parent_thread_id: ThreadId,
196 parent_turn_id: TurnId,
197 request: SubagentRequest,
198 ) -> anyhow::Result<SubagentResult>;
199
200 async fn dispatch_traced(
201 &self,
202 parent_thread_id: ThreadId,
203 parent_turn_id: TurnId,
204 request: SubagentRequest,
205 trace_sink: Option<std::sync::Arc<dyn SubagentTraceSink>>,
206 ) -> anyhow::Result<SubagentResult> {
207 let _ = trace_sink;
208 self.dispatch(parent_thread_id, parent_turn_id, request)
209 .await
210 }
211
212 async fn dispatch_with_context(
221 &self,
222 parent_thread_id: ThreadId,
223 parent_turn_id: TurnId,
224 request: SubagentRequest,
225 trace_sink: Option<std::sync::Arc<dyn SubagentTraceSink>>,
226 handles: crate::tools::ToolExecutionHandles,
227 ) -> anyhow::Result<SubagentResult> {
228 let _ = handles;
229 self.dispatch_traced(parent_thread_id, parent_turn_id, request, trace_sink)
230 .await
231 }
232}
233
234pub const AGENT_SWARM_TOOL_NAME: &str = "agent_swarm";
248
249pub const AGENT_SWARM_PROMPT_PLACEHOLDER: &str = "{{item}}";
251
252pub const AGENT_SWARM_MAX_SUBAGENTS: usize = 128;
254pub const AGENT_SWARM_INITIAL_LAUNCH_LIMIT: usize = 5;
256pub const AGENT_SWARM_LAUNCH_INTERVAL_MS: u64 = 700;
258pub const AGENT_SWARM_RATE_LIMIT_MAX_RETRIES: usize = 4;
261pub const AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS: u64 = 3_000;
264pub const AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP: usize = 8;
267pub const AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS: u64 = 2_000;
272pub const AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS: u64 = 180_000;
276
277fn default_rate_limit_max_retries() -> usize {
278 AGENT_SWARM_RATE_LIMIT_MAX_RETRIES
279}
280
281fn default_rate_limit_base_backoff_ms() -> u64 {
282 AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS
283}
284
285fn default_rate_limit_shrink_interval_ms() -> u64 {
286 AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS
287}
288
289fn default_rate_limit_recovery_interval_ms() -> u64 {
290 AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS
291}
292
293pub const AGENT_SWARM_MODE_REMINDER: &str = "Agent-swarm mode is active. When the task splits into \
298several similarly-shaped subtasks over different inputs, call the agent_swarm tool exactly once \
299with a prompt_template containing {{item}} and an items array (or resume_agent_ids), and make it \
300the only tool call in that response. agent_swarm dispatches configured roles only: its \
301subagent_type must exactly match a role advertised in the tool schema, and lane names such as \
302scout are not role IDs. Do not rely on a lane to add missing tools; for generic repository work \
303use spawn_agent when available.";
304
305pub const ULTRA_MODE_REMINDER: &str = "Ultra mode is active for this task. Prefer proactive, \
309bounded multi-agent delegation via spawn_agent (and agent_swarm when many similar items apply) \
310whenever parallel work would materially improve speed or quality. Keep the lead focused on \
311orchestration, synthesis, and decisions the workers cannot make alone.";
312
313#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
315pub struct AgentSwarmModeChanged {
316 pub thread_id: ThreadId,
317 #[serde(default, skip_serializing_if = "Option::is_none")]
318 pub turn_id: Option<TurnId>,
319 pub enabled: bool,
320 pub trigger: AgentSwarmModeTrigger,
321 #[serde(with = "time::serde::rfc3339")]
322 pub timestamp: time::OffsetDateTime,
323}
324
325#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
329pub struct UltraModeChanged {
330 pub thread_id: ThreadId,
331 #[serde(default, skip_serializing_if = "Option::is_none")]
332 pub turn_id: Option<TurnId>,
333 pub enabled: bool,
334 pub trigger: UltraModeTrigger,
335 #[serde(with = "time::serde::rfc3339")]
336 pub timestamp: time::OffsetDateTime,
337}
338
339#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
343pub struct AgentSwarmStarted {
344 pub thread_id: ThreadId,
345 pub turn_id: TurnId,
346 pub tool_id: String,
347 pub child_count: usize,
349 #[serde(with = "time::serde::rfc3339")]
350 pub timestamp: time::OffsetDateTime,
351}
352
353#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
356pub struct AgentSwarmCompleted {
357 pub thread_id: ThreadId,
358 pub turn_id: TurnId,
359 pub tool_id: String,
360 pub completed: usize,
361 pub failed: usize,
362 pub aborted: usize,
363 #[serde(with = "time::serde::rfc3339")]
364 pub timestamp: time::OffsetDateTime,
365}
366
367#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
371pub struct AgentSwarmProgressSnapshot {
372 pub total: usize,
373 pub completed: usize,
374 pub failed: usize,
375 pub aborted: usize,
376}
377
378impl AgentSwarmProgressSnapshot {
379 pub fn resolved(&self) -> usize {
381 self.completed + self.failed + self.aborted
382 }
383}
384
385#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
389pub struct AgentSwarmProgress {
390 pub thread_id: ThreadId,
391 pub turn_id: TurnId,
392 pub tool_id: String,
393 pub snapshot: AgentSwarmProgressSnapshot,
394 #[serde(with = "time::serde::rfc3339")]
395 pub timestamp: time::OffsetDateTime,
396}
397
398#[async_trait::async_trait]
402pub trait AgentSwarmProgressSink: Send + Sync {
403 async fn emit_progress(
404 &self,
405 thread_id: &str,
406 turn_id: &str,
407 tool_id: &str,
408 snapshot: AgentSwarmProgressSnapshot,
409 );
410}
411
412#[derive(Debug, Clone, Copy, PartialEq, Eq)]
415pub enum AgentSwarmBatchViolation {
416 MultipleSwarms {
418 has_other_tools: bool,
420 },
421 MixedWithOtherTools,
423}
424
425impl AgentSwarmBatchViolation {
426 pub fn deny_message(self) -> String {
429 match self {
430 Self::MultipleSwarms { has_other_tools } => {
431 let mut message = String::from(
432 "agent_swarm must be called one swarm at a time. Multiple agent_swarm calls \
433 are not forbidden, but issue them sequentially: call one agent_swarm, wait \
434 for its result, then call the next; or merge the work into a single \
435 agent_swarm when one swarm can cover it.",
436 );
437 if has_other_tools {
438 message.push_str(
439 " agent_swarm also must not be combined with other tools in the same \
440 response.",
441 );
442 }
443 message
444 }
445 Self::MixedWithOtherTools => String::from(
446 "agent_swarm must be the only tool call in a model response. Retry with a single \
447 agent_swarm call by itself, then call any other tools after it returns.",
448 ),
449 }
450 }
451}
452
453pub fn agent_swarm_batch_violation<'a>(
457 tool_names: impl Iterator<Item = &'a str>,
458) -> Option<AgentSwarmBatchViolation> {
459 let mut total = 0usize;
460 let mut swarm = 0usize;
461 for name in tool_names {
462 total += 1;
463 if name == AGENT_SWARM_TOOL_NAME {
464 swarm += 1;
465 }
466 }
467 if swarm == 0 || (swarm == 1 && total == 1) {
468 return None;
469 }
470 if swarm > 1 {
471 Some(AgentSwarmBatchViolation::MultipleSwarms {
472 has_other_tools: total > swarm,
473 })
474 } else {
475 Some(AgentSwarmBatchViolation::MixedWithOtherTools)
476 }
477}
478
479#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
485#[serde(rename_all = "snake_case")]
486pub enum AgentSwarmModeTrigger {
487 Manual,
488 Task,
489 Tool,
490}
491
492impl AgentSwarmModeTrigger {
493 pub fn as_str(self) -> &'static str {
494 match self {
495 Self::Manual => "manual",
496 Self::Task => "task",
497 Self::Tool => "tool",
498 }
499 }
500
501 pub fn should_auto_exit(self) -> bool {
503 matches!(self, Self::Task | Self::Tool)
504 }
505}
506
507#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
512#[serde(rename_all = "snake_case")]
513pub enum UltraModeTrigger {
514 Manual,
515 Task,
516}
517
518impl UltraModeTrigger {
519 pub fn as_str(self) -> &'static str {
520 match self {
521 Self::Manual => "manual",
522 Self::Task => "task",
523 }
524 }
525
526 pub fn should_auto_exit(self) -> bool {
528 matches!(self, Self::Task)
529 }
530}
531
532#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
534#[serde(rename_all = "snake_case")]
535pub enum AgentSwarmChildKind {
536 Spawn,
537 Resume,
538}
539
540impl AgentSwarmChildKind {
541 pub fn as_str(self) -> &'static str {
542 match self {
543 Self::Spawn => "spawn",
544 Self::Resume => "resume",
545 }
546 }
547}
548
549#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
551#[serde(rename_all = "snake_case")]
552pub enum AgentSwarmChildOutcome {
553 Completed,
554 Failed,
555 Aborted,
556}
557
558impl AgentSwarmChildOutcome {
559 pub fn as_str(self) -> &'static str {
560 match self {
561 Self::Completed => "completed",
562 Self::Failed => "failed",
563 Self::Aborted => "aborted",
564 }
565 }
566}
567
568#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
570#[serde(rename_all = "snake_case")]
571pub enum AgentSwarmChildState {
572 Started,
573 NotStarted,
574}
575
576impl AgentSwarmChildState {
577 pub fn as_str(self) -> &'static str {
578 match self {
579 Self::Started => "started",
580 Self::NotStarted => "not_started",
581 }
582 }
583}
584
585#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
587pub struct AgentSwarmRequest {
588 pub description: String,
589 #[serde(default, skip_serializing_if = "Option::is_none")]
590 pub subagent_type: Option<String>,
591 #[serde(default, skip_serializing_if = "Option::is_none")]
592 pub prompt_template: Option<String>,
593 #[serde(default, skip_serializing_if = "Vec::is_empty")]
594 pub items: Vec<String>,
595 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
598 pub resume_agent_ids: BTreeMap<String, String>,
599}
600
601#[derive(Debug, Clone, PartialEq, Eq)]
603pub struct AgentSwarmChildSpec {
604 pub index: usize,
606 pub kind: AgentSwarmChildKind,
607 pub item: Option<String>,
609 pub prompt: String,
611 pub resume_agent_id: Option<String>,
613}
614
615#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
617pub struct AgentSwarmConfig {
618 pub max_subagents: usize,
621 pub initial_launch_limit: usize,
623 pub launch_interval_ms: u64,
625 #[serde(default, skip_serializing_if = "Option::is_none")]
627 pub max_concurrency: Option<usize>,
628 #[serde(default, skip_serializing_if = "Option::is_none")]
630 pub child_timeout_seconds: Option<u64>,
631 #[serde(default = "default_rate_limit_max_retries")]
633 pub rate_limit_max_retries: usize,
634 #[serde(default = "default_rate_limit_base_backoff_ms")]
637 pub rate_limit_base_backoff_ms: u64,
638 #[serde(default = "default_rate_limit_shrink_interval_ms")]
642 pub rate_limit_shrink_interval_ms: u64,
643 #[serde(default = "default_rate_limit_recovery_interval_ms")]
647 pub rate_limit_recovery_interval_ms: u64,
648}
649
650impl Default for AgentSwarmConfig {
651 fn default() -> Self {
652 Self {
653 max_subagents: AGENT_SWARM_MAX_SUBAGENTS,
654 initial_launch_limit: AGENT_SWARM_INITIAL_LAUNCH_LIMIT,
655 launch_interval_ms: AGENT_SWARM_LAUNCH_INTERVAL_MS,
656 max_concurrency: None,
657 child_timeout_seconds: None,
658 rate_limit_max_retries: AGENT_SWARM_RATE_LIMIT_MAX_RETRIES,
659 rate_limit_base_backoff_ms: AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS,
660 rate_limit_shrink_interval_ms: AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS,
661 rate_limit_recovery_interval_ms: AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS,
662 }
663 }
664}
665
666impl AgentSwarmConfig {
667 pub fn clamped(mut self) -> Self {
671 self.max_subagents = self.max_subagents.clamp(1, AGENT_SWARM_MAX_SUBAGENTS);
672 self.initial_launch_limit = self.initial_launch_limit.max(1);
673 if let Some(cap) = self.max_concurrency {
674 self.max_concurrency = Some(cap.max(1));
675 }
676 if let Some(timeout) = self.child_timeout_seconds {
677 self.child_timeout_seconds = Some(timeout.max(1));
678 }
679 self.rate_limit_max_retries = self
680 .rate_limit_max_retries
681 .min(AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP);
682 self
683 }
684}
685
686#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
688pub struct AgentSwarmChildResult {
689 pub index: usize,
690 pub kind: AgentSwarmChildKind,
691 #[serde(default, skip_serializing_if = "Option::is_none")]
692 pub item: Option<String>,
693 #[serde(default, skip_serializing_if = "Option::is_none")]
694 pub agent_id: Option<String>,
695 pub outcome: AgentSwarmChildOutcome,
696 #[serde(default, skip_serializing_if = "Option::is_none")]
697 pub state: Option<AgentSwarmChildState>,
698 pub body: String,
699 #[serde(default, skip_serializing_if = "Option::is_none")]
700 pub usage: Option<TokenUsage>,
701}
702
703#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
705pub struct AgentSwarmResult {
706 pub completed: usize,
707 pub failed: usize,
708 pub aborted: usize,
709 pub children: Vec<AgentSwarmChildResult>,
710}
711
712impl AgentSwarmResult {
713 pub fn from_children(children: Vec<AgentSwarmChildResult>) -> Self {
714 let mut completed = 0;
715 let mut failed = 0;
716 let mut aborted = 0;
717 for child in &children {
718 match child.outcome {
719 AgentSwarmChildOutcome::Completed => completed += 1,
720 AgentSwarmChildOutcome::Failed => failed += 1,
721 AgentSwarmChildOutcome::Aborted => aborted += 1,
722 }
723 }
724 Self {
725 completed,
726 failed,
727 aborted,
728 children,
729 }
730 }
731
732 pub fn summary_line(&self) -> String {
734 let mut parts = Vec::new();
735 if self.completed > 0 {
736 parts.push(format!("completed: {}", self.completed));
737 }
738 if self.failed > 0 {
739 parts.push(format!("failed: {}", self.failed));
740 }
741 if self.aborted > 0 {
742 parts.push(format!("aborted: {}", self.aborted));
743 }
744 if parts.is_empty() {
745 "completed: 0".to_string()
746 } else {
747 parts.join(", ")
748 }
749 }
750
751 pub fn needs_resume_hint(&self) -> bool {
753 self.children.iter().any(|child| {
754 child.outcome != AgentSwarmChildOutcome::Completed && child.agent_id.is_some()
755 })
756 }
757
758 pub fn render_text(&self) -> String {
760 let mut lines = vec![
761 "<agent_swarm_result>".to_string(),
762 format!("<summary>{}</summary>", self.summary_line()),
763 ];
764 if self.needs_resume_hint() {
765 lines.push(
766 "<resume_hint>Call agent_swarm with resume_agent_ids using the agent_id values in this result to continue unfinished work.</resume_hint>"
767 .to_string(),
768 );
769 }
770 for child in &self.children {
771 let mode = if child.kind == AgentSwarmChildKind::Resume {
772 " mode=\"resume\"".to_string()
773 } else {
774 String::new()
775 };
776 let agent_id = child
777 .agent_id
778 .as_deref()
779 .map(|id| format!(" agent_id=\"{}\"", escape_xml_attr(id)))
780 .unwrap_or_default();
781 let item = child
782 .item
783 .as_deref()
784 .map(|item| format!(" item=\"{}\"", escape_xml_attr(item)))
785 .unwrap_or_default();
786 let state = child
787 .state
788 .map(|state| format!(" state=\"{}\"", state.as_str()))
789 .unwrap_or_default();
790 lines.push(format!(
791 "<subagent{mode}{agent_id}{item}{state} outcome=\"{}\">{}</subagent>",
792 child.outcome.as_str(),
793 escape_xml_text(&child.body)
794 ));
795 }
796 lines.push("</agent_swarm_result>".to_string());
797 lines.join("\n")
798 }
799}
800
801fn escape_xml_attr(value: &str) -> String {
802 value
803 .replace('&', "&")
804 .replace('"', """)
805 .replace('<', "<")
806 .replace('>', ">")
807}
808
809fn escape_xml_text(value: &str) -> String {
810 value
811 .replace('&', "&")
812 .replace('<', "<")
813 .replace('>', ">")
814}
815
816pub fn build_agent_swarm_specs(
825 request: &AgentSwarmRequest,
826 config: &AgentSwarmConfig,
827) -> Result<Vec<AgentSwarmChildSpec>, AgentSwarmValidationError> {
828 if request.description.trim().is_empty() {
829 return Err(AgentSwarmValidationError::EmptyDescription);
830 }
831
832 let items: Vec<String> = request
833 .items
834 .iter()
835 .map(|item| item.trim().to_string())
836 .collect();
837 if items.iter().any(|item| item.is_empty()) {
838 return Err(AgentSwarmValidationError::EmptyItem);
839 }
840
841 let resume_entries: Vec<(String, String)> = request
842 .resume_agent_ids
843 .iter()
844 .map(|(id, prompt)| (id.trim().to_string(), prompt.trim().to_string()))
845 .collect();
846 if resume_entries
847 .iter()
848 .any(|(id, prompt)| id.is_empty() || prompt.is_empty())
849 {
850 return Err(AgentSwarmValidationError::EmptyResumeEntry);
851 }
852
853 let item_count = items.len();
854 let resume_count = resume_entries.len();
855 let total = item_count + resume_count;
856
857 if resume_count == 0 && item_count < 2 {
858 return Err(AgentSwarmValidationError::TooFewItems);
859 }
860 let max = config.max_subagents.clamp(1, AGENT_SWARM_MAX_SUBAGENTS);
861 if total > max {
862 return Err(AgentSwarmValidationError::TooManySubagents { total, max });
863 }
864
865 let prompt_template = request
866 .prompt_template
867 .as_ref()
868 .map(|template| template.trim().to_string())
869 .filter(|template| !template.is_empty());
870
871 if item_count > 0 {
872 let Some(template) = prompt_template.as_ref() else {
873 return Err(AgentSwarmValidationError::MissingPromptTemplate);
874 };
875 if !template.contains(AGENT_SWARM_PROMPT_PLACEHOLDER) {
876 return Err(AgentSwarmValidationError::MissingPlaceholder);
877 }
878 }
879
880 let mut specs = Vec::with_capacity(total);
881 for (agent_id, prompt) in &resume_entries {
882 specs.push(AgentSwarmChildSpec {
883 index: specs.len() + 1,
884 kind: AgentSwarmChildKind::Resume,
885 item: None,
886 prompt: prompt.clone(),
887 resume_agent_id: Some(agent_id.clone()),
888 });
889 }
890
891 if item_count > 0 {
892 let template = prompt_template.expect("prompt template validated above");
893 let mut seen: BTreeMap<String, usize> = BTreeMap::new();
894 for (offset, item) in items.iter().enumerate() {
895 let prompt = template.replace(AGENT_SWARM_PROMPT_PLACEHOLDER, item);
896 if let Some(previous) = seen.get(&prompt) {
897 return Err(AgentSwarmValidationError::DuplicatePrompt {
898 first: *previous,
899 second: offset + 1,
900 });
901 }
902 seen.insert(prompt.clone(), offset + 1);
903 specs.push(AgentSwarmChildSpec {
904 index: specs.len() + 1,
905 kind: AgentSwarmChildKind::Spawn,
906 item: Some(item.clone()),
907 prompt,
908 resume_agent_id: None,
909 });
910 }
911 }
912
913 Ok(specs)
914}
915
916#[derive(Debug, Clone, PartialEq, Eq)]
918pub enum AgentSwarmValidationError {
919 EmptyDescription,
920 EmptyItem,
921 EmptyResumeEntry,
922 TooFewItems,
923 TooManySubagents { total: usize, max: usize },
924 MissingPromptTemplate,
925 MissingPlaceholder,
926 DuplicatePrompt { first: usize, second: usize },
927}
928
929impl std::fmt::Display for AgentSwarmValidationError {
930 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
931 match self {
932 Self::EmptyDescription => write!(f, "description must not be empty"),
933 Self::EmptyItem => write!(f, "items must not contain empty values"),
934 Self::EmptyResumeEntry => {
935 write!(
936 f,
937 "resume_agent_ids entries must have non-empty ids and prompts"
938 )
939 }
940 Self::TooFewItems => write!(
941 f,
942 "agent_swarm requires at least 2 items unless resume_agent_ids is provided"
943 ),
944 Self::TooManySubagents { total, max } => {
945 write!(
946 f,
947 "agent_swarm supports at most {max} subagents (got {total})"
948 )
949 }
950 Self::MissingPromptTemplate => {
951 write!(f, "prompt_template is required when items are provided")
952 }
953 Self::MissingPlaceholder => write!(
954 f,
955 "prompt_template must include the {AGENT_SWARM_PROMPT_PLACEHOLDER} placeholder"
956 ),
957 Self::DuplicatePrompt { first, second } => write!(
958 f,
959 "duplicate subagent prompts from items {first} and {second}; agent_swarm requires distinct subagents"
960 ),
961 }
962 }
963}
964
965impl std::error::Error for AgentSwarmValidationError {}
966
967#[cfg(test)]
968mod tests {
969 use std::sync::Arc;
970
971 use super::*;
972
973 struct NoopDispatcher;
974
975 #[async_trait::async_trait]
976 impl SubagentDispatcher for NoopDispatcher {
977 fn id(&self) -> SubagentDispatcherId {
978 "noop".to_string()
979 }
980
981 fn definitions(&self) -> Vec<SubagentDefinition> {
982 vec![SubagentDefinition {
983 agent_type: "explore".to_string(),
984 description: "Explore the workspace".to_string(),
985 tools: vec!["Read".to_string()],
986 model: Some("test-model".to_string()),
987 system_prompt: Some("Report findings only".to_string()),
988 permission_mode: SubagentPermissionMode::ReadOnly,
989 max_turns: Some(4),
990 max_result_chars: Some(4000),
991 }]
992 }
993
994 async fn dispatch(
995 &self,
996 _parent_thread_id: ThreadId,
997 _parent_turn_id: TurnId,
998 request: SubagentRequest,
999 ) -> anyhow::Result<SubagentResult> {
1000 Ok(SubagentResult {
1001 thread_id: "child-thread".to_string(),
1002 turn_id: "child-turn".to_string(),
1003 agent_type: request
1004 .subagent_type
1005 .unwrap_or_else(|| "explore".to_string()),
1006 model: request.model,
1007 final_message: "done".to_string(),
1008 usage: None,
1009 exit_reason: SubagentExitReason::Completed,
1010 transcript: None,
1011 metadata: serde_json::json!({}),
1012 })
1013 }
1014 }
1015
1016 #[tokio::test]
1017 async fn subagent_dispatcher_trait_is_object_safe() {
1018 let dispatcher: Arc<dyn SubagentDispatcher> = Arc::new(NoopDispatcher);
1019
1020 assert_eq!(dispatcher.id(), "noop");
1021 assert_eq!(dispatcher.definitions()[0].agent_type, "explore");
1022
1023 let result = dispatcher
1024 .dispatch(
1025 "parent-thread".to_string(),
1026 "parent-turn".to_string(),
1027 SubagentRequest {
1028 description: "Check files".to_string(),
1029 prompt: "Find the API entrypoint".to_string(),
1030 subagent_type: Some("explore".to_string()),
1031 model: Some("test-model".to_string()),
1032 provider: None,
1033 tools: Some(vec!["Read".to_string()]),
1034 lane: None,
1035 max_concurrent: None,
1036 allowed_tools: None,
1037 parent_deadline_seconds: None,
1038 inputs: None,
1039 timeout_seconds: Some(10),
1040 },
1041 )
1042 .await
1043 .unwrap();
1044
1045 assert_eq!(result.thread_id, "child-thread");
1046 assert_eq!(result.exit_reason, SubagentExitReason::Completed);
1047 }
1048
1049 fn swarm_request(items: &[&str]) -> AgentSwarmRequest {
1050 AgentSwarmRequest {
1051 description: "inspect files".to_string(),
1052 subagent_type: Some("explore".to_string()),
1053 prompt_template: Some("Read {{item}} and report.".to_string()),
1054 items: items.iter().map(|item| item.to_string()).collect(),
1055 resume_agent_ids: BTreeMap::new(),
1056 }
1057 }
1058
1059 #[test]
1060 fn agent_swarm_config_clamps_into_bounds() {
1061 let clamped = AgentSwarmConfig {
1062 max_subagents: 9001,
1063 initial_launch_limit: 0,
1064 launch_interval_ms: 700,
1065 max_concurrency: Some(0),
1066 child_timeout_seconds: Some(0),
1067 rate_limit_max_retries: 9001,
1068 ..AgentSwarmConfig::default()
1069 }
1070 .clamped();
1071 assert_eq!(clamped.max_subagents, AGENT_SWARM_MAX_SUBAGENTS);
1072 assert_eq!(clamped.initial_launch_limit, 1);
1073 assert_eq!(clamped.max_concurrency, Some(1));
1074 assert_eq!(clamped.child_timeout_seconds, Some(1));
1075 assert_eq!(
1076 clamped.rate_limit_max_retries,
1077 AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP
1078 );
1079 }
1080
1081 #[test]
1082 fn build_specs_expands_items_in_order() {
1083 let specs = build_agent_swarm_specs(
1084 &swarm_request(&["a.rs", "b.rs"]),
1085 &AgentSwarmConfig::default(),
1086 )
1087 .unwrap();
1088 assert_eq!(specs.len(), 2);
1089 assert_eq!(specs[0].index, 1);
1090 assert_eq!(specs[0].kind, AgentSwarmChildKind::Spawn);
1091 assert_eq!(specs[0].prompt, "Read a.rs and report.");
1092 assert_eq!(specs[1].prompt, "Read b.rs and report.");
1093 }
1094
1095 #[test]
1096 fn build_specs_orders_resumes_before_spawns() {
1097 let mut request = swarm_request(&["a.rs"]);
1098 request
1099 .resume_agent_ids
1100 .insert("agent-9".to_string(), "continue".to_string());
1101 let specs = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap();
1102 assert_eq!(specs.len(), 2);
1103 assert_eq!(specs[0].kind, AgentSwarmChildKind::Resume);
1104 assert_eq!(specs[0].resume_agent_id.as_deref(), Some("agent-9"));
1105 assert_eq!(specs[1].kind, AgentSwarmChildKind::Spawn);
1106 assert_eq!(specs[1].index, 2);
1107 }
1108
1109 #[test]
1110 fn build_specs_rejects_single_item_without_resume() {
1111 let err =
1112 build_agent_swarm_specs(&swarm_request(&["only.rs"]), &AgentSwarmConfig::default())
1113 .unwrap_err();
1114 assert_eq!(err, AgentSwarmValidationError::TooFewItems);
1115 }
1116
1117 #[test]
1118 fn build_specs_rejects_missing_placeholder() {
1119 let mut request = swarm_request(&["a.rs", "b.rs"]);
1120 request.prompt_template = Some("no placeholder here".to_string());
1121 let err = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap_err();
1122 assert_eq!(err, AgentSwarmValidationError::MissingPlaceholder);
1123 }
1124
1125 #[test]
1126 fn build_specs_rejects_duplicate_prompts() {
1127 let request = swarm_request(&["dup", "dup"]);
1128 let err = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap_err();
1129 assert!(matches!(
1130 err,
1131 AgentSwarmValidationError::DuplicatePrompt { .. }
1132 ));
1133 }
1134
1135 #[test]
1136 fn build_specs_enforces_max_subagents() {
1137 let config = AgentSwarmConfig {
1138 max_subagents: 2,
1139 ..AgentSwarmConfig::default()
1140 };
1141 let err = build_agent_swarm_specs(&swarm_request(&["a", "b", "c"]), &config).unwrap_err();
1142 assert_eq!(
1143 err,
1144 AgentSwarmValidationError::TooManySubagents { total: 3, max: 2 }
1145 );
1146 }
1147
1148 #[test]
1149 fn agent_swarm_result_renders_summary_and_resume_hint() {
1150 let result = AgentSwarmResult::from_children(vec![
1151 AgentSwarmChildResult {
1152 index: 1,
1153 kind: AgentSwarmChildKind::Spawn,
1154 item: Some("a.rs".to_string()),
1155 agent_id: Some("agent-1".to_string()),
1156 outcome: AgentSwarmChildOutcome::Completed,
1157 state: None,
1158 body: "ok".to_string(),
1159 usage: None,
1160 },
1161 AgentSwarmChildResult {
1162 index: 2,
1163 kind: AgentSwarmChildKind::Spawn,
1164 item: Some("b & c.rs".to_string()),
1165 agent_id: Some("agent-2".to_string()),
1166 outcome: AgentSwarmChildOutcome::Failed,
1167 state: Some(AgentSwarmChildState::Started),
1168 body: "boom <fatal>".to_string(),
1169 usage: None,
1170 },
1171 ]);
1172 assert_eq!(result.completed, 1);
1173 assert_eq!(result.failed, 1);
1174 assert!(result.needs_resume_hint());
1175 let text = result.render_text();
1176 assert!(text.contains("<summary>completed: 1, failed: 1</summary>"));
1177 assert!(text.contains("<resume_hint>"));
1178 assert!(text.contains("item=\"b & c.rs\""));
1179 assert!(text.contains("boom <fatal>"));
1180 assert!(text.contains("outcome=\"failed\""));
1181 }
1182
1183 #[test]
1184 fn agent_swarm_dtos_round_trip_json() {
1185 let result = AgentSwarmResult::from_children(vec![AgentSwarmChildResult {
1186 index: 1,
1187 kind: AgentSwarmChildKind::Resume,
1188 item: None,
1189 agent_id: Some("agent-1".to_string()),
1190 outcome: AgentSwarmChildOutcome::Aborted,
1191 state: Some(AgentSwarmChildState::NotStarted),
1192 body: "cancelled".to_string(),
1193 usage: None,
1194 }]);
1195 let json = serde_json::to_value(&result).unwrap();
1196 assert_eq!(json["aborted"], 1);
1197 assert_eq!(json["children"][0]["kind"], "resume");
1198 assert_eq!(json["children"][0]["outcome"], "aborted");
1199 assert_eq!(json["children"][0]["state"], "not_started");
1200 let round: AgentSwarmResult = serde_json::from_value(json).unwrap();
1201 assert_eq!(round, result);
1202 }
1203
1204 #[test]
1205 fn agent_swarm_trigger_auto_exit_rules() {
1206 assert!(!AgentSwarmModeTrigger::Manual.should_auto_exit());
1207 assert!(AgentSwarmModeTrigger::Task.should_auto_exit());
1208 assert!(AgentSwarmModeTrigger::Tool.should_auto_exit());
1209 }
1210
1211 #[test]
1212 fn batch_violation_allows_single_swarm_alone() {
1213 assert_eq!(
1214 agent_swarm_batch_violation(["agent_swarm"].into_iter()),
1215 None
1216 );
1217 }
1218
1219 #[test]
1220 fn batch_violation_allows_batches_without_swarm() {
1221 assert_eq!(
1222 agent_swarm_batch_violation(["read_file", "write_file"].into_iter()),
1223 None
1224 );
1225 }
1226
1227 #[test]
1228 fn batch_violation_flags_swarm_mixed_with_other_tools() {
1229 assert_eq!(
1230 agent_swarm_batch_violation(["agent_swarm", "read_file"].into_iter()),
1231 Some(AgentSwarmBatchViolation::MixedWithOtherTools)
1232 );
1233 let message = AgentSwarmBatchViolation::MixedWithOtherTools.deny_message();
1234 assert!(message.contains("only tool call"));
1235 }
1236
1237 #[test]
1238 fn batch_violation_flags_multiple_swarms() {
1239 assert_eq!(
1240 agent_swarm_batch_violation(["agent_swarm", "agent_swarm"].into_iter()),
1241 Some(AgentSwarmBatchViolation::MultipleSwarms {
1242 has_other_tools: false
1243 })
1244 );
1245 assert_eq!(
1246 agent_swarm_batch_violation(["agent_swarm", "agent_swarm", "read_file"].into_iter()),
1247 Some(AgentSwarmBatchViolation::MultipleSwarms {
1248 has_other_tools: true
1249 })
1250 );
1251 let message = AgentSwarmBatchViolation::MultipleSwarms {
1252 has_other_tools: true,
1253 }
1254 .deny_message();
1255 assert!(message.contains("one swarm at a time"));
1256 assert!(message.contains("combined with other tools"));
1257 }
1258}