Skip to main content

roder_api/
subagents.rs

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>,
16    pub tools: Option<Vec<String>>,
17    #[serde(default, skip_serializing_if = "Option::is_none")]
18    pub lane: Option<SubagentLane>,
19    #[serde(default, skip_serializing_if = "Option::is_none")]
20    pub max_concurrent: Option<usize>,
21    #[serde(default, skip_serializing_if = "Option::is_none")]
22    pub allowed_tools: Option<Vec<String>>,
23    #[serde(default, skip_serializing_if = "Option::is_none")]
24    pub parent_deadline_seconds: Option<u64>,
25    #[serde(default, skip_serializing_if = "Option::is_none")]
26    pub inputs: Option<serde_json::Value>,
27    pub timeout_seconds: Option<u64>,
28}
29
30#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)]
31#[serde(rename_all = "snake_case")]
32pub enum SubagentLane {
33    Scout,
34    Editor,
35    Reviewer,
36    Runner,
37}
38
39impl SubagentLane {
40    pub fn as_str(self) -> &'static str {
41        match self {
42            Self::Scout => "scout",
43            Self::Editor => "editor",
44            Self::Reviewer => "reviewer",
45            Self::Runner => "runner",
46        }
47    }
48
49    pub fn preset(self) -> SubagentLanePreset {
50        match self {
51            Self::Scout => SubagentLanePreset {
52                lane: self,
53                description: "Read and search without changing state.",
54                max_concurrent: 4,
55                timeout_seconds: 120,
56                allowed_tools: &[
57                    "Read",
58                    "Grep",
59                    "Glob",
60                    "read_file",
61                    "grep",
62                    "glob",
63                    "list_files",
64                ],
65            },
66            Self::Editor => SubagentLanePreset {
67                lane: self,
68                description: "Make a bounded file-change slice.",
69                max_concurrent: 2,
70                timeout_seconds: 180,
71                allowed_tools: &[
72                    "Read",
73                    "Grep",
74                    "Glob",
75                    "read_file",
76                    "grep",
77                    "glob",
78                    "list_files",
79                    "write_file",
80                    "edit",
81                    "multi_edit",
82                    "apply_patch",
83                ],
84            },
85            Self::Reviewer => SubagentLanePreset {
86                lane: self,
87                description: "Review and verify with evidence.",
88                max_concurrent: 2,
89                timeout_seconds: 120,
90                allowed_tools: &[
91                    "Read",
92                    "Grep",
93                    "Glob",
94                    "read_file",
95                    "grep",
96                    "glob",
97                    "list_files",
98                ],
99            },
100            Self::Runner => SubagentLanePreset {
101                lane: self,
102                description: "Run commands or tests when process policy allows it.",
103                max_concurrent: 1,
104                timeout_seconds: 120,
105                allowed_tools: &["Shell", "shell", "exec_command", "run_command"],
106            },
107        }
108    }
109}
110
111#[derive(Debug, Clone, Copy, Serialize, PartialEq, Eq)]
112#[serde(rename_all = "camelCase")]
113pub struct SubagentLanePreset {
114    pub lane: SubagentLane,
115    pub description: &'static str,
116    pub max_concurrent: usize,
117    pub timeout_seconds: u64,
118    pub allowed_tools: &'static [&'static str],
119}
120
121pub fn built_in_subagent_lane_presets() -> [SubagentLanePreset; 4] {
122    [
123        SubagentLane::Scout.preset(),
124        SubagentLane::Editor.preset(),
125        SubagentLane::Reviewer.preset(),
126        SubagentLane::Runner.preset(),
127    ]
128}
129
130pub const SUBAGENT_SUMMARY_CONTRACT: &str = "Child summary must include these labels: Conclusion, Evidence, Files inspected, Files changed, Remaining uncertainty.";
131
132#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
133pub struct SubagentDefinition {
134    pub agent_type: String,
135    pub description: String,
136    pub tools: Vec<String>,
137    pub model: Option<String>,
138    pub system_prompt: Option<String>,
139    pub permission_mode: SubagentPermissionMode,
140    pub max_turns: Option<u32>,
141    pub max_result_chars: Option<usize>,
142}
143
144#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
145#[serde(rename_all = "snake_case")]
146pub enum SubagentPermissionMode {
147    ReadOnly,
148    #[default]
149    Default,
150    AutoEdit,
151}
152
153#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
154pub struct SubagentResult {
155    pub thread_id: ThreadId,
156    pub turn_id: TurnId,
157    pub agent_type: String,
158    pub model: Option<String>,
159    pub final_message: String,
160    pub usage: Option<TokenUsage>,
161    pub exit_reason: SubagentExitReason,
162    #[serde(default, skip_serializing_if = "Option::is_none")]
163    pub transcript: Option<serde_json::Value>,
164    #[serde(default)]
165    pub metadata: serde_json::Value,
166}
167
168#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
169#[serde(rename_all = "snake_case")]
170pub enum SubagentExitReason {
171    Completed,
172    MaxTurns,
173    Timeout,
174    Cancelled,
175    Failed,
176}
177
178#[async_trait::async_trait]
179pub trait SubagentDispatcher: Send + Sync + 'static {
180    fn id(&self) -> SubagentDispatcherId;
181
182    fn definitions(&self) -> Vec<SubagentDefinition>;
183
184    async fn dispatch(
185        &self,
186        parent_thread_id: ThreadId,
187        parent_turn_id: TurnId,
188        request: SubagentRequest,
189    ) -> anyhow::Result<SubagentResult>;
190
191    async fn dispatch_traced(
192        &self,
193        parent_thread_id: ThreadId,
194        parent_turn_id: TurnId,
195        request: SubagentRequest,
196        trace_sink: Option<std::sync::Arc<dyn SubagentTraceSink>>,
197    ) -> anyhow::Result<SubagentResult> {
198        let _ = trace_sink;
199        self.dispatch(parent_thread_id, parent_turn_id, request)
200            .await
201    }
202
203    /// Dispatch a child carrying the parent tool call's execution handles
204    /// (workspace, remote workspace, process runner, context artifacts). This
205    /// is how a subagent or swarm child operates on the same workspace as the
206    /// lead thread instead of failing with "workspace handle is not available".
207    ///
208    /// The default implementation ignores the handles and falls back to
209    /// [`Self::dispatch_traced`], so non-workspace dispatchers (e.g. remote or
210    /// fake ones) keep working unchanged.
211    async fn dispatch_with_context(
212        &self,
213        parent_thread_id: ThreadId,
214        parent_turn_id: TurnId,
215        request: SubagentRequest,
216        trace_sink: Option<std::sync::Arc<dyn SubagentTraceSink>>,
217        handles: crate::tools::ToolExecutionHandles,
218    ) -> anyhow::Result<SubagentResult> {
219        let _ = handles;
220        self.dispatch_traced(parent_thread_id, parent_turn_id, request, trace_sink)
221            .await
222    }
223}
224
225// ---------------------------------------------------------------------------
226// Agent-swarm mode (roadmap phase 104)
227//
228// Swarm is a Roder-native composition over the canonical subagent dispatch
229// surface above. It lets a lead model (or `/agent-swarm` command) launch many
230// homogeneous child tasks from one prompt template, resume unfinished children,
231// and collect an ordered, machine-readable result. The types here are
232// provider-neutral and do not vendor any external implementation.
233// ---------------------------------------------------------------------------
234
235/// Canonical model-facing swarm tool name. Single source of truth shared by
236/// the tool registration (`roder-ext-subagents`) and the core turn loop's
237/// exclusivity enforcement (`roder-core`).
238pub const AGENT_SWARM_TOOL_NAME: &str = "agent_swarm";
239
240/// Literal placeholder replaced with each `agent_swarm` item value.
241pub const AGENT_SWARM_PROMPT_PLACEHOLDER: &str = "{{item}}";
242
243/// Default upper bound on swarm child count. Config may lower but not exceed it.
244pub const AGENT_SWARM_MAX_SUBAGENTS: usize = 128;
245/// Default number of children that may start immediately before pacing applies.
246pub const AGENT_SWARM_INITIAL_LAUNCH_LIMIT: usize = 5;
247/// Default pacing interval between additional child launches, in milliseconds.
248pub const AGENT_SWARM_LAUNCH_INTERVAL_MS: u64 = 700;
249/// Default number of times a child is retried after a provider rate limit
250/// before it is reported as failed.
251pub const AGENT_SWARM_RATE_LIMIT_MAX_RETRIES: usize = 4;
252/// Default base backoff between rate-limit retries, in milliseconds (doubled
253/// each subsequent attempt: 3s, 6s, 12s, ...).
254pub const AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS: u64 = 3_000;
255/// Hard cap on rate-limit retries so config/env can never make a swarm wait
256/// unboundedly.
257pub const AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP: usize = 8;
258/// Default minimum interval between successive global rate-limit capacity
259/// shrinks, in milliseconds. While the provider keeps rate-limiting, the swarm
260/// shrinks its global capacity by one at most this often, so it throttles
261/// smoothly instead of collapsing on the first burst of 429s.
262pub const AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS: u64 = 2_000;
263/// Default quiet window after which the swarm recovers one unit of global
264/// rate-limit capacity, in milliseconds (3 minutes). Any fresh rate limit
265/// restarts the window, and recovery happens at most once per quiet window.
266pub const AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS: u64 = 180_000;
267
268fn default_rate_limit_max_retries() -> usize {
269    AGENT_SWARM_RATE_LIMIT_MAX_RETRIES
270}
271
272fn default_rate_limit_base_backoff_ms() -> u64 {
273    AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS
274}
275
276fn default_rate_limit_shrink_interval_ms() -> u64 {
277    AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS
278}
279
280fn default_rate_limit_recovery_interval_ms() -> u64 {
281    AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS
282}
283
284/// Canonical swarm-mode reminder injected into a turn's developer instructions
285/// (server-side) while agent-swarm mode is active, so the model reaches for the
286/// `agent_swarm` fanout tool. Shared by the runtime injection and the TUI label
287/// so the wording stays in one place.
288pub const AGENT_SWARM_MODE_REMINDER: &str = "Agent-swarm mode is active. When the task splits into \
289several similarly-shaped subtasks over different inputs, call the agent_swarm tool exactly once \
290with a prompt_template containing {{item}} and an items array (or resume_agent_ids), and make it \
291the only tool call in that response.";
292
293/// Emitted when agent-swarm mode is toggled on a runtime/thread.
294#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
295pub struct AgentSwarmModeChanged {
296    pub thread_id: ThreadId,
297    #[serde(default, skip_serializing_if = "Option::is_none")]
298    pub turn_id: Option<TurnId>,
299    pub enabled: bool,
300    pub trigger: AgentSwarmModeTrigger,
301    #[serde(with = "time::serde::rfc3339")]
302    pub timestamp: time::OffsetDateTime,
303}
304
305/// Emitted when the `agent_swarm` tool begins fanning out children, so any
306/// app-server/SDK/TUI client can observe the swarm as a whole. Per-child
307/// progress flows through the existing `Subagent*` trace events.
308#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
309pub struct AgentSwarmStarted {
310    pub thread_id: ThreadId,
311    pub turn_id: TurnId,
312    pub tool_id: String,
313    /// Children the swarm will launch (item-based spawns plus resumes).
314    pub child_count: usize,
315    #[serde(with = "time::serde::rfc3339")]
316    pub timestamp: time::OffsetDateTime,
317}
318
319/// Emitted when the `agent_swarm` tool finishes, carrying the aggregate child
320/// outcome counts.
321#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
322pub struct AgentSwarmCompleted {
323    pub thread_id: ThreadId,
324    pub turn_id: TurnId,
325    pub tool_id: String,
326    pub completed: usize,
327    pub failed: usize,
328    pub aborted: usize,
329    #[serde(with = "time::serde::rfc3339")]
330    pub timestamp: time::OffsetDateTime,
331}
332
333/// A running tally of swarm children as they resolve, used for live progress.
334/// `resolved` is `completed + failed + aborted`; clients can render
335/// `resolved/total` as a progress bar without tracking individual children.
336#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
337pub struct AgentSwarmProgressSnapshot {
338    pub total: usize,
339    pub completed: usize,
340    pub failed: usize,
341    pub aborted: usize,
342}
343
344impl AgentSwarmProgressSnapshot {
345    /// Children that have finished (in any outcome).
346    pub fn resolved(&self) -> usize {
347        self.completed + self.failed + self.aborted
348    }
349}
350
351/// Emitted each time a swarm child resolves, so a client can render a live
352/// "N/total done" progress tick between `AgentSwarmStarted` and
353/// `AgentSwarmCompleted` rather than only the final result.
354#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
355pub struct AgentSwarmProgress {
356    pub thread_id: ThreadId,
357    pub turn_id: TurnId,
358    pub tool_id: String,
359    pub snapshot: AgentSwarmProgressSnapshot,
360    #[serde(with = "time::serde::rfc3339")]
361    pub timestamp: time::OffsetDateTime,
362}
363
364/// Sink the runtime supplies on a tool-execution context so the `agent_swarm`
365/// tool can publish live progress snapshots (carrying the runtime's thread/turn
366/// ids and the bus emitter) without the tool depending on `roder-core`.
367#[async_trait::async_trait]
368pub trait AgentSwarmProgressSink: Send + Sync {
369    async fn emit_progress(
370        &self,
371        thread_id: &str,
372        turn_id: &str,
373        tool_id: &str,
374        snapshot: AgentSwarmProgressSnapshot,
375    );
376}
377
378/// Why a batch of tool calls violates the `agent_swarm` exclusivity rule:
379/// `agent_swarm` must be the only tool call in a single model response.
380#[derive(Debug, Clone, Copy, PartialEq, Eq)]
381pub enum AgentSwarmBatchViolation {
382    /// More than one `agent_swarm` call appeared in the same response.
383    MultipleSwarms {
384        /// Whether non-swarm tool calls were also present in the batch.
385        has_other_tools: bool,
386    },
387    /// Exactly one `agent_swarm` call was mixed with other tool calls.
388    MixedWithOtherTools,
389}
390
391impl AgentSwarmBatchViolation {
392    /// Actionable retry text returned to the model for every call in the
393    /// denied batch, so it re-issues `agent_swarm` by itself.
394    pub fn deny_message(self) -> String {
395        match self {
396            Self::MultipleSwarms { has_other_tools } => {
397                let mut message = String::from(
398                    "agent_swarm must be called one swarm at a time. Multiple agent_swarm calls \
399                     are not forbidden, but issue them sequentially: call one agent_swarm, wait \
400                     for its result, then call the next; or merge the work into a single \
401                     agent_swarm when one swarm can cover it.",
402                );
403                if has_other_tools {
404                    message.push_str(
405                        " agent_swarm also must not be combined with other tools in the same \
406                         response.",
407                    );
408                }
409                message
410            }
411            Self::MixedWithOtherTools => String::from(
412                "agent_swarm must be the only tool call in a model response. Retry with a single \
413                 agent_swarm call by itself, then call any other tools after it returns.",
414            ),
415        }
416    }
417}
418
419/// Detect whether a batch of tool-call names violates the `agent_swarm`
420/// exclusivity rule. Returns `None` when the batch is valid: no `agent_swarm`
421/// call, or exactly one `agent_swarm` call by itself.
422pub fn agent_swarm_batch_violation<'a>(
423    tool_names: impl Iterator<Item = &'a str>,
424) -> Option<AgentSwarmBatchViolation> {
425    let mut total = 0usize;
426    let mut swarm = 0usize;
427    for name in tool_names {
428        total += 1;
429        if name == AGENT_SWARM_TOOL_NAME {
430            swarm += 1;
431        }
432    }
433    if swarm == 0 || (swarm == 1 && total == 1) {
434        return None;
435    }
436    if swarm > 1 {
437        Some(AgentSwarmBatchViolation::MultipleSwarms {
438            has_other_tools: total > swarm,
439        })
440    } else {
441        Some(AgentSwarmBatchViolation::MixedWithOtherTools)
442    }
443}
444
445/// How a swarm-mode session was entered.
446///
447/// `manual` is a persistent toggle (`/agent-swarm on`); `task` is a one-shot
448/// `/agent-swarm <prompt>`; `tool` is implicit entry from the `agent_swarm`
449/// tool call. Only `manual` stays active across turns.
450#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
451#[serde(rename_all = "snake_case")]
452pub enum AgentSwarmModeTrigger {
453    Manual,
454    Task,
455    Tool,
456}
457
458impl AgentSwarmModeTrigger {
459    pub fn as_str(self) -> &'static str {
460        match self {
461            Self::Manual => "manual",
462            Self::Task => "task",
463            Self::Tool => "tool",
464        }
465    }
466
467    /// One-shot triggers auto-exit at the end of the relevant turn.
468    pub fn should_auto_exit(self) -> bool {
469        matches!(self, Self::Task | Self::Tool)
470    }
471}
472
473/// Whether a swarm child is a fresh spawn or a resume of an existing agent.
474#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
475#[serde(rename_all = "snake_case")]
476pub enum AgentSwarmChildKind {
477    Spawn,
478    Resume,
479}
480
481impl AgentSwarmChildKind {
482    pub fn as_str(self) -> &'static str {
483        match self {
484            Self::Spawn => "spawn",
485            Self::Resume => "resume",
486        }
487    }
488}
489
490/// Final outcome of a single swarm child.
491#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
492#[serde(rename_all = "snake_case")]
493pub enum AgentSwarmChildOutcome {
494    Completed,
495    Failed,
496    Aborted,
497}
498
499impl AgentSwarmChildOutcome {
500    pub fn as_str(self) -> &'static str {
501        match self {
502            Self::Completed => "completed",
503            Self::Failed => "failed",
504            Self::Aborted => "aborted",
505        }
506    }
507}
508
509/// Whether an aborted child had started running before cancellation.
510#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
511#[serde(rename_all = "snake_case")]
512pub enum AgentSwarmChildState {
513    Started,
514    NotStarted,
515}
516
517impl AgentSwarmChildState {
518    pub fn as_str(self) -> &'static str {
519        match self {
520            Self::Started => "started",
521            Self::NotStarted => "not_started",
522        }
523    }
524}
525
526/// Parsed `agent_swarm` tool input, before any child is dispatched.
527#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
528pub struct AgentSwarmRequest {
529    pub description: String,
530    #[serde(default, skip_serializing_if = "Option::is_none")]
531    pub subagent_type: Option<String>,
532    #[serde(default, skip_serializing_if = "Option::is_none")]
533    pub prompt_template: Option<String>,
534    #[serde(default, skip_serializing_if = "Vec::is_empty")]
535    pub items: Vec<String>,
536    /// Map of existing subagent agent_id to the prompt used to resume it.
537    /// Resumed children are dispatched before new item-based spawns.
538    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
539    pub resume_agent_ids: BTreeMap<String, String>,
540}
541
542/// One ordered child to dispatch as part of a swarm.
543#[derive(Debug, Clone, PartialEq, Eq)]
544pub struct AgentSwarmChildSpec {
545    /// 1-based ordering index, stable across the whole swarm.
546    pub index: usize,
547    pub kind: AgentSwarmChildKind,
548    /// The item value (for spawns) used to render the prompt, when present.
549    pub item: Option<String>,
550    /// The fully-rendered prompt for this child.
551    pub prompt: String,
552    /// For resumes, the existing agent id to continue.
553    pub resume_agent_id: Option<String>,
554}
555
556/// Tunable scheduler/bounds for swarm fanout.
557#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
558pub struct AgentSwarmConfig {
559    /// Hard upper bound on children per swarm. Never exceeds
560    /// [`AGENT_SWARM_MAX_SUBAGENTS`].
561    pub max_subagents: usize,
562    /// Children allowed to start immediately before pacing applies.
563    pub initial_launch_limit: usize,
564    /// Pacing interval between additional launches, in milliseconds.
565    pub launch_interval_ms: u64,
566    /// Optional cap on simultaneously-active children (normal phase).
567    #[serde(default, skip_serializing_if = "Option::is_none")]
568    pub max_concurrency: Option<usize>,
569    /// Optional per-child timeout override, in seconds.
570    #[serde(default, skip_serializing_if = "Option::is_none")]
571    pub child_timeout_seconds: Option<u64>,
572    /// Times a child is retried after a provider rate limit before failing.
573    #[serde(default = "default_rate_limit_max_retries")]
574    pub rate_limit_max_retries: usize,
575    /// Base backoff between rate-limit retries, in milliseconds (doubled each
576    /// subsequent attempt).
577    #[serde(default = "default_rate_limit_base_backoff_ms")]
578    pub rate_limit_base_backoff_ms: u64,
579    /// Minimum interval between successive global rate-limit capacity shrinks,
580    /// in milliseconds. Sustained rate limits shrink global capacity by one at
581    /// most this often.
582    #[serde(default = "default_rate_limit_shrink_interval_ms")]
583    pub rate_limit_shrink_interval_ms: u64,
584    /// Quiet window with no rate limit after which the swarm recovers one unit
585    /// of global rate-limit capacity, in milliseconds. A fresh rate limit
586    /// restarts the window.
587    #[serde(default = "default_rate_limit_recovery_interval_ms")]
588    pub rate_limit_recovery_interval_ms: u64,
589}
590
591impl Default for AgentSwarmConfig {
592    fn default() -> Self {
593        Self {
594            max_subagents: AGENT_SWARM_MAX_SUBAGENTS,
595            initial_launch_limit: AGENT_SWARM_INITIAL_LAUNCH_LIMIT,
596            launch_interval_ms: AGENT_SWARM_LAUNCH_INTERVAL_MS,
597            max_concurrency: None,
598            child_timeout_seconds: None,
599            rate_limit_max_retries: AGENT_SWARM_RATE_LIMIT_MAX_RETRIES,
600            rate_limit_base_backoff_ms: AGENT_SWARM_RATE_LIMIT_BASE_BACKOFF_MS,
601            rate_limit_shrink_interval_ms: AGENT_SWARM_RATE_LIMIT_SHRINK_INTERVAL_MS,
602            rate_limit_recovery_interval_ms: AGENT_SWARM_RATE_LIMIT_RECOVERY_INTERVAL_MS,
603        }
604    }
605}
606
607impl AgentSwarmConfig {
608    /// Clamp config into a bounded, deterministic range so config (or env) can
609    /// never request unbounded fanout, a zero-sized ramp, or an unbounded
610    /// rate-limit wait.
611    pub fn clamped(mut self) -> Self {
612        self.max_subagents = self.max_subagents.clamp(1, AGENT_SWARM_MAX_SUBAGENTS);
613        self.initial_launch_limit = self.initial_launch_limit.max(1);
614        if let Some(cap) = self.max_concurrency {
615            self.max_concurrency = Some(cap.max(1));
616        }
617        if let Some(timeout) = self.child_timeout_seconds {
618            self.child_timeout_seconds = Some(timeout.max(1));
619        }
620        self.rate_limit_max_retries = self
621            .rate_limit_max_retries
622            .min(AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP);
623        self
624    }
625}
626
627/// Outcome of a single resolved swarm child entry.
628#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
629pub struct AgentSwarmChildResult {
630    pub index: usize,
631    pub kind: AgentSwarmChildKind,
632    #[serde(default, skip_serializing_if = "Option::is_none")]
633    pub item: Option<String>,
634    #[serde(default, skip_serializing_if = "Option::is_none")]
635    pub agent_id: Option<String>,
636    pub outcome: AgentSwarmChildOutcome,
637    #[serde(default, skip_serializing_if = "Option::is_none")]
638    pub state: Option<AgentSwarmChildState>,
639    pub body: String,
640    #[serde(default, skip_serializing_if = "Option::is_none")]
641    pub usage: Option<TokenUsage>,
642}
643
644/// Aggregated, ordered swarm result returned to the lead model.
645#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
646pub struct AgentSwarmResult {
647    pub completed: usize,
648    pub failed: usize,
649    pub aborted: usize,
650    pub children: Vec<AgentSwarmChildResult>,
651}
652
653impl AgentSwarmResult {
654    pub fn from_children(children: Vec<AgentSwarmChildResult>) -> Self {
655        let mut completed = 0;
656        let mut failed = 0;
657        let mut aborted = 0;
658        for child in &children {
659            match child.outcome {
660                AgentSwarmChildOutcome::Completed => completed += 1,
661                AgentSwarmChildOutcome::Failed => failed += 1,
662                AgentSwarmChildOutcome::Aborted => aborted += 1,
663            }
664        }
665        Self {
666            completed,
667            failed,
668            aborted,
669            children,
670        }
671    }
672
673    /// `completed: 2, failed: 1` style summary; omits zero buckets.
674    pub fn summary_line(&self) -> String {
675        let mut parts = Vec::new();
676        if self.completed > 0 {
677            parts.push(format!("completed: {}", self.completed));
678        }
679        if self.failed > 0 {
680            parts.push(format!("failed: {}", self.failed));
681        }
682        if self.aborted > 0 {
683            parts.push(format!("aborted: {}", self.aborted));
684        }
685        if parts.is_empty() {
686            "completed: 0".to_string()
687        } else {
688            parts.join(", ")
689        }
690    }
691
692    /// A resume hint is useful when unfinished children carry agent ids.
693    pub fn needs_resume_hint(&self) -> bool {
694        self.children.iter().any(|child| {
695            child.outcome != AgentSwarmChildOutcome::Completed && child.agent_id.is_some()
696        })
697    }
698
699    /// Render the durable, transcript-safe `<agent_swarm_result>` text block.
700    pub fn render_text(&self) -> String {
701        let mut lines = vec![
702            "<agent_swarm_result>".to_string(),
703            format!("<summary>{}</summary>", self.summary_line()),
704        ];
705        if self.needs_resume_hint() {
706            lines.push(
707                "<resume_hint>Call agent_swarm with resume_agent_ids using the agent_id values in this result to continue unfinished work.</resume_hint>"
708                    .to_string(),
709            );
710        }
711        for child in &self.children {
712            let mode = if child.kind == AgentSwarmChildKind::Resume {
713                " mode=\"resume\"".to_string()
714            } else {
715                String::new()
716            };
717            let agent_id = child
718                .agent_id
719                .as_deref()
720                .map(|id| format!(" agent_id=\"{}\"", escape_xml_attr(id)))
721                .unwrap_or_default();
722            let item = child
723                .item
724                .as_deref()
725                .map(|item| format!(" item=\"{}\"", escape_xml_attr(item)))
726                .unwrap_or_default();
727            let state = child
728                .state
729                .map(|state| format!(" state=\"{}\"", state.as_str()))
730                .unwrap_or_default();
731            lines.push(format!(
732                "<subagent{mode}{agent_id}{item}{state} outcome=\"{}\">{}</subagent>",
733                child.outcome.as_str(),
734                escape_xml_text(&child.body)
735            ));
736        }
737        lines.push("</agent_swarm_result>".to_string());
738        lines.join("\n")
739    }
740}
741
742fn escape_xml_attr(value: &str) -> String {
743    value
744        .replace('&', "&amp;")
745        .replace('"', "&quot;")
746        .replace('<', "&lt;")
747        .replace('>', "&gt;")
748}
749
750fn escape_xml_text(value: &str) -> String {
751    value
752        .replace('&', "&amp;")
753        .replace('<', "&lt;")
754        .replace('>', "&gt;")
755}
756
757/// Validate an [`AgentSwarmRequest`] and build the ordered child specs.
758///
759/// Mirrors the documented contract: at least two `items` unless
760/// `resume_agent_ids` is present; `prompt_template` is required whenever
761/// `items` are present and must contain the `{{item}}` placeholder; filled
762/// prompts must be distinct; resumed children are ordered before spawns; and
763/// the total child count must not exceed `config.max_subagents`. No child is
764/// dispatched if validation fails.
765pub fn build_agent_swarm_specs(
766    request: &AgentSwarmRequest,
767    config: &AgentSwarmConfig,
768) -> Result<Vec<AgentSwarmChildSpec>, AgentSwarmValidationError> {
769    if request.description.trim().is_empty() {
770        return Err(AgentSwarmValidationError::EmptyDescription);
771    }
772
773    let items: Vec<String> = request.items.iter().map(|item| item.trim().to_string()).collect();
774    if items.iter().any(|item| item.is_empty()) {
775        return Err(AgentSwarmValidationError::EmptyItem);
776    }
777
778    let resume_entries: Vec<(String, String)> = request
779        .resume_agent_ids
780        .iter()
781        .map(|(id, prompt)| (id.trim().to_string(), prompt.trim().to_string()))
782        .collect();
783    if resume_entries
784        .iter()
785        .any(|(id, prompt)| id.is_empty() || prompt.is_empty())
786    {
787        return Err(AgentSwarmValidationError::EmptyResumeEntry);
788    }
789
790    let item_count = items.len();
791    let resume_count = resume_entries.len();
792    let total = item_count + resume_count;
793
794    if resume_count == 0 && item_count < 2 {
795        return Err(AgentSwarmValidationError::TooFewItems);
796    }
797    let max = config.max_subagents.clamp(1, AGENT_SWARM_MAX_SUBAGENTS);
798    if total > max {
799        return Err(AgentSwarmValidationError::TooManySubagents { total, max });
800    }
801
802    let prompt_template = request
803        .prompt_template
804        .as_ref()
805        .map(|template| template.trim().to_string())
806        .filter(|template| !template.is_empty());
807
808    if item_count > 0 {
809        let Some(template) = prompt_template.as_ref() else {
810            return Err(AgentSwarmValidationError::MissingPromptTemplate);
811        };
812        if !template.contains(AGENT_SWARM_PROMPT_PLACEHOLDER) {
813            return Err(AgentSwarmValidationError::MissingPlaceholder);
814        }
815    }
816
817    let mut specs = Vec::with_capacity(total);
818    for (agent_id, prompt) in &resume_entries {
819        specs.push(AgentSwarmChildSpec {
820            index: specs.len() + 1,
821            kind: AgentSwarmChildKind::Resume,
822            item: None,
823            prompt: prompt.clone(),
824            resume_agent_id: Some(agent_id.clone()),
825        });
826    }
827
828    if item_count > 0 {
829        let template = prompt_template.expect("prompt template validated above");
830        let mut seen: BTreeMap<String, usize> = BTreeMap::new();
831        for (offset, item) in items.iter().enumerate() {
832            let prompt = template.replace(AGENT_SWARM_PROMPT_PLACEHOLDER, item);
833            if let Some(previous) = seen.get(&prompt) {
834                return Err(AgentSwarmValidationError::DuplicatePrompt {
835                    first: *previous,
836                    second: offset + 1,
837                });
838            }
839            seen.insert(prompt.clone(), offset + 1);
840            specs.push(AgentSwarmChildSpec {
841                index: specs.len() + 1,
842                kind: AgentSwarmChildKind::Spawn,
843                item: Some(item.clone()),
844                prompt,
845                resume_agent_id: None,
846            });
847        }
848    }
849
850    Ok(specs)
851}
852
853/// Reasons an `agent_swarm` request is rejected before any child starts.
854#[derive(Debug, Clone, PartialEq, Eq)]
855pub enum AgentSwarmValidationError {
856    EmptyDescription,
857    EmptyItem,
858    EmptyResumeEntry,
859    TooFewItems,
860    TooManySubagents { total: usize, max: usize },
861    MissingPromptTemplate,
862    MissingPlaceholder,
863    DuplicatePrompt { first: usize, second: usize },
864}
865
866impl std::fmt::Display for AgentSwarmValidationError {
867    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
868        match self {
869            Self::EmptyDescription => write!(f, "description must not be empty"),
870            Self::EmptyItem => write!(f, "items must not contain empty values"),
871            Self::EmptyResumeEntry => {
872                write!(f, "resume_agent_ids entries must have non-empty ids and prompts")
873            }
874            Self::TooFewItems => write!(
875                f,
876                "agent_swarm requires at least 2 items unless resume_agent_ids is provided"
877            ),
878            Self::TooManySubagents { total, max } => {
879                write!(f, "agent_swarm supports at most {max} subagents (got {total})")
880            }
881            Self::MissingPromptTemplate => {
882                write!(f, "prompt_template is required when items are provided")
883            }
884            Self::MissingPlaceholder => write!(
885                f,
886                "prompt_template must include the {AGENT_SWARM_PROMPT_PLACEHOLDER} placeholder"
887            ),
888            Self::DuplicatePrompt { first, second } => write!(
889                f,
890                "duplicate subagent prompts from items {first} and {second}; agent_swarm requires distinct subagents"
891            ),
892        }
893    }
894}
895
896impl std::error::Error for AgentSwarmValidationError {}
897
898#[cfg(test)]
899mod tests {
900    use std::sync::Arc;
901
902    use super::*;
903
904    struct NoopDispatcher;
905
906    #[async_trait::async_trait]
907    impl SubagentDispatcher for NoopDispatcher {
908        fn id(&self) -> SubagentDispatcherId {
909            "noop".to_string()
910        }
911
912        fn definitions(&self) -> Vec<SubagentDefinition> {
913            vec![SubagentDefinition {
914                agent_type: "explore".to_string(),
915                description: "Explore the workspace".to_string(),
916                tools: vec!["Read".to_string()],
917                model: Some("test-model".to_string()),
918                system_prompt: Some("Report findings only".to_string()),
919                permission_mode: SubagentPermissionMode::ReadOnly,
920                max_turns: Some(4),
921                max_result_chars: Some(4000),
922            }]
923        }
924
925        async fn dispatch(
926            &self,
927            _parent_thread_id: ThreadId,
928            _parent_turn_id: TurnId,
929            request: SubagentRequest,
930        ) -> anyhow::Result<SubagentResult> {
931            Ok(SubagentResult {
932                thread_id: "child-thread".to_string(),
933                turn_id: "child-turn".to_string(),
934                agent_type: request
935                    .subagent_type
936                    .unwrap_or_else(|| "explore".to_string()),
937                model: request.model,
938                final_message: "done".to_string(),
939                usage: None,
940                exit_reason: SubagentExitReason::Completed,
941                transcript: None,
942                metadata: serde_json::json!({}),
943            })
944        }
945    }
946
947    #[tokio::test]
948    async fn subagent_dispatcher_trait_is_object_safe() {
949        let dispatcher: Arc<dyn SubagentDispatcher> = Arc::new(NoopDispatcher);
950
951        assert_eq!(dispatcher.id(), "noop");
952        assert_eq!(dispatcher.definitions()[0].agent_type, "explore");
953
954        let result = dispatcher
955            .dispatch(
956                "parent-thread".to_string(),
957                "parent-turn".to_string(),
958                SubagentRequest {
959                    description: "Check files".to_string(),
960                    prompt: "Find the API entrypoint".to_string(),
961                    subagent_type: Some("explore".to_string()),
962                    model: Some("test-model".to_string()),
963                    tools: Some(vec!["Read".to_string()]),
964                    lane: None,
965                    max_concurrent: None,
966                    allowed_tools: None,
967                    parent_deadline_seconds: None,
968                    inputs: None,
969                    timeout_seconds: Some(10),
970                },
971            )
972            .await
973            .unwrap();
974
975        assert_eq!(result.thread_id, "child-thread");
976        assert_eq!(result.exit_reason, SubagentExitReason::Completed);
977    }
978
979    fn swarm_request(items: &[&str]) -> AgentSwarmRequest {
980        AgentSwarmRequest {
981            description: "inspect files".to_string(),
982            subagent_type: Some("explore".to_string()),
983            prompt_template: Some("Read {{item}} and report.".to_string()),
984            items: items.iter().map(|item| item.to_string()).collect(),
985            resume_agent_ids: BTreeMap::new(),
986        }
987    }
988
989    #[test]
990    fn agent_swarm_config_clamps_into_bounds() {
991        let clamped = AgentSwarmConfig {
992            max_subagents: 9001,
993            initial_launch_limit: 0,
994            launch_interval_ms: 700,
995            max_concurrency: Some(0),
996            child_timeout_seconds: Some(0),
997            rate_limit_max_retries: 9001,
998            ..AgentSwarmConfig::default()
999        }
1000        .clamped();
1001        assert_eq!(clamped.max_subagents, AGENT_SWARM_MAX_SUBAGENTS);
1002        assert_eq!(clamped.initial_launch_limit, 1);
1003        assert_eq!(clamped.max_concurrency, Some(1));
1004        assert_eq!(clamped.child_timeout_seconds, Some(1));
1005        assert_eq!(
1006            clamped.rate_limit_max_retries,
1007            AGENT_SWARM_RATE_LIMIT_MAX_RETRIES_CAP
1008        );
1009    }
1010
1011    #[test]
1012    fn build_specs_expands_items_in_order() {
1013        let specs =
1014            build_agent_swarm_specs(&swarm_request(&["a.rs", "b.rs"]), &AgentSwarmConfig::default())
1015                .unwrap();
1016        assert_eq!(specs.len(), 2);
1017        assert_eq!(specs[0].index, 1);
1018        assert_eq!(specs[0].kind, AgentSwarmChildKind::Spawn);
1019        assert_eq!(specs[0].prompt, "Read a.rs and report.");
1020        assert_eq!(specs[1].prompt, "Read b.rs and report.");
1021    }
1022
1023    #[test]
1024    fn build_specs_orders_resumes_before_spawns() {
1025        let mut request = swarm_request(&["a.rs"]);
1026        request
1027            .resume_agent_ids
1028            .insert("agent-9".to_string(), "continue".to_string());
1029        let specs = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap();
1030        assert_eq!(specs.len(), 2);
1031        assert_eq!(specs[0].kind, AgentSwarmChildKind::Resume);
1032        assert_eq!(specs[0].resume_agent_id.as_deref(), Some("agent-9"));
1033        assert_eq!(specs[1].kind, AgentSwarmChildKind::Spawn);
1034        assert_eq!(specs[1].index, 2);
1035    }
1036
1037    #[test]
1038    fn build_specs_rejects_single_item_without_resume() {
1039        let err = build_agent_swarm_specs(&swarm_request(&["only.rs"]), &AgentSwarmConfig::default())
1040            .unwrap_err();
1041        assert_eq!(err, AgentSwarmValidationError::TooFewItems);
1042    }
1043
1044    #[test]
1045    fn build_specs_rejects_missing_placeholder() {
1046        let mut request = swarm_request(&["a.rs", "b.rs"]);
1047        request.prompt_template = Some("no placeholder here".to_string());
1048        let err = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap_err();
1049        assert_eq!(err, AgentSwarmValidationError::MissingPlaceholder);
1050    }
1051
1052    #[test]
1053    fn build_specs_rejects_duplicate_prompts() {
1054        let request = swarm_request(&["dup", "dup"]);
1055        let err = build_agent_swarm_specs(&request, &AgentSwarmConfig::default()).unwrap_err();
1056        assert!(matches!(
1057            err,
1058            AgentSwarmValidationError::DuplicatePrompt { .. }
1059        ));
1060    }
1061
1062    #[test]
1063    fn build_specs_enforces_max_subagents() {
1064        let config = AgentSwarmConfig {
1065            max_subagents: 2,
1066            ..AgentSwarmConfig::default()
1067        };
1068        let err = build_agent_swarm_specs(&swarm_request(&["a", "b", "c"]), &config).unwrap_err();
1069        assert_eq!(
1070            err,
1071            AgentSwarmValidationError::TooManySubagents { total: 3, max: 2 }
1072        );
1073    }
1074
1075    #[test]
1076    fn agent_swarm_result_renders_summary_and_resume_hint() {
1077        let result = AgentSwarmResult::from_children(vec![
1078            AgentSwarmChildResult {
1079                index: 1,
1080                kind: AgentSwarmChildKind::Spawn,
1081                item: Some("a.rs".to_string()),
1082                agent_id: Some("agent-1".to_string()),
1083                outcome: AgentSwarmChildOutcome::Completed,
1084                state: None,
1085                body: "ok".to_string(),
1086                usage: None,
1087            },
1088            AgentSwarmChildResult {
1089                index: 2,
1090                kind: AgentSwarmChildKind::Spawn,
1091                item: Some("b & c.rs".to_string()),
1092                agent_id: Some("agent-2".to_string()),
1093                outcome: AgentSwarmChildOutcome::Failed,
1094                state: Some(AgentSwarmChildState::Started),
1095                body: "boom <fatal>".to_string(),
1096                usage: None,
1097            },
1098        ]);
1099        assert_eq!(result.completed, 1);
1100        assert_eq!(result.failed, 1);
1101        assert!(result.needs_resume_hint());
1102        let text = result.render_text();
1103        assert!(text.contains("<summary>completed: 1, failed: 1</summary>"));
1104        assert!(text.contains("<resume_hint>"));
1105        assert!(text.contains("item=\"b &amp; c.rs\""));
1106        assert!(text.contains("boom &lt;fatal&gt;"));
1107        assert!(text.contains("outcome=\"failed\""));
1108    }
1109
1110    #[test]
1111    fn agent_swarm_dtos_round_trip_json() {
1112        let result = AgentSwarmResult::from_children(vec![AgentSwarmChildResult {
1113            index: 1,
1114            kind: AgentSwarmChildKind::Resume,
1115            item: None,
1116            agent_id: Some("agent-1".to_string()),
1117            outcome: AgentSwarmChildOutcome::Aborted,
1118            state: Some(AgentSwarmChildState::NotStarted),
1119            body: "cancelled".to_string(),
1120            usage: None,
1121        }]);
1122        let json = serde_json::to_value(&result).unwrap();
1123        assert_eq!(json["aborted"], 1);
1124        assert_eq!(json["children"][0]["kind"], "resume");
1125        assert_eq!(json["children"][0]["outcome"], "aborted");
1126        assert_eq!(json["children"][0]["state"], "not_started");
1127        let round: AgentSwarmResult = serde_json::from_value(json).unwrap();
1128        assert_eq!(round, result);
1129    }
1130
1131    #[test]
1132    fn agent_swarm_trigger_auto_exit_rules() {
1133        assert!(!AgentSwarmModeTrigger::Manual.should_auto_exit());
1134        assert!(AgentSwarmModeTrigger::Task.should_auto_exit());
1135        assert!(AgentSwarmModeTrigger::Tool.should_auto_exit());
1136    }
1137
1138    #[test]
1139    fn batch_violation_allows_single_swarm_alone() {
1140        assert_eq!(
1141            agent_swarm_batch_violation(["agent_swarm"].into_iter()),
1142            None
1143        );
1144    }
1145
1146    #[test]
1147    fn batch_violation_allows_batches_without_swarm() {
1148        assert_eq!(
1149            agent_swarm_batch_violation(["read_file", "write_file"].into_iter()),
1150            None
1151        );
1152    }
1153
1154    #[test]
1155    fn batch_violation_flags_swarm_mixed_with_other_tools() {
1156        assert_eq!(
1157            agent_swarm_batch_violation(["agent_swarm", "read_file"].into_iter()),
1158            Some(AgentSwarmBatchViolation::MixedWithOtherTools)
1159        );
1160        let message = AgentSwarmBatchViolation::MixedWithOtherTools.deny_message();
1161        assert!(message.contains("only tool call"));
1162    }
1163
1164    #[test]
1165    fn batch_violation_flags_multiple_swarms() {
1166        assert_eq!(
1167            agent_swarm_batch_violation(["agent_swarm", "agent_swarm"].into_iter()),
1168            Some(AgentSwarmBatchViolation::MultipleSwarms {
1169                has_other_tools: false
1170            })
1171        );
1172        assert_eq!(
1173            agent_swarm_batch_violation(["agent_swarm", "agent_swarm", "read_file"].into_iter()),
1174            Some(AgentSwarmBatchViolation::MultipleSwarms {
1175                has_other_tools: true
1176            })
1177        );
1178        let message = AgentSwarmBatchViolation::MultipleSwarms {
1179            has_other_tools: true,
1180        }
1181        .deny_message();
1182        assert!(message.contains("one swarm at a time"));
1183        assert!(message.contains("combined with other tools"));
1184    }
1185}