Skip to main content

codeswarm_adapters/
adapters.rs

1//! Protocol adapters for CodeSwarm.
2//!
3//! ACP and native CLI protocols are intentionally peers here. They emit the
4//! same core events and expose capabilities through the same adapter boundary.
5
6use std::collections::{BTreeMap, VecDeque};
7use std::io::Read;
8use std::path::{Path, PathBuf};
9use std::process::Stdio;
10use std::sync::{
11    Arc, Mutex,
12    atomic::{AtomicBool, AtomicUsize, Ordering},
13};
14
15use crate::{
16    AgentCapabilities, AgentCommand, AgentEvent, Effect, EventLog, Mode, PermissionAnswer,
17    PermissionRequest, RosterSlot, RosterUpdate, SessionState, TerminalEvent, ToolStatus,
18    ToolUpdate, UsageUpdate,
19    persistence::{BufferedSessionMetadataStore, SessionMetadata},
20    reduce,
21    relay::{
22        CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, Relay, RelayDecision, STOP_TOKEN,
23        control_token_visible_end, is_usage_limit_response, requested_next_slot,
24        strip_control_tokens, strip_stop_token,
25    },
26    resources,
27};
28use async_trait::async_trait;
29use base64::{Engine, engine::general_purpose::STANDARD as BASE64};
30use serde_json::Value;
31use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt, BufReader};
32use tokio::process::{Child, ChildStdout, Command};
33use tokio::sync::{Mutex as AsyncMutex, Notify, mpsc};
34
35#[path = "codex.rs"]
36mod codex;
37#[path = "native.rs"]
38mod native;
39pub use codex::CodexAdapter;
40#[path = "claude.rs"]
41mod claude;
42pub use claude::ClaudeAdapter;
43
44pub type AdapterResult<T> = Result<T, AdapterError>;
45
46/// Keep a peer from allocating unbounded memory for one newline-delimited
47/// protocol frame. The Python client applied the same boundary (10 MiB) to
48/// its asyncio reader; ACP messages are normally much smaller than this.
49const MAX_ACP_LINE_BYTES: usize = 10 * 1024 * 1024;
50const MAX_FILE_READ_BYTES: usize = 4 * 1024 * 1024;
51const MAX_TERMINAL_OUTPUT_BYTES: usize = 1024 * 1024;
52
53#[derive(Clone, Debug)]
54struct TerminalProcess {
55    child: Arc<AsyncMutex<Option<Child>>>,
56    output: Arc<Mutex<Vec<u8>>>,
57    truncated: Arc<AtomicBool>,
58    output_readers: Arc<AtomicUsize>,
59}
60
61impl TerminalProcess {
62    async fn kill(&self) {
63        if let Some(child) = self.child.lock().await.as_mut() {
64            #[cfg(unix)]
65            if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
66                tokio::time::sleep(std::time::Duration::from_millis(100)).await;
67                signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
68            }
69            let _ = child.start_kill();
70        }
71    }
72
73    async fn stop(&self) {
74        if let Some(mut child) = self.child.lock().await.take() {
75            let _ = terminate_child(&mut child).await;
76        }
77    }
78
79    async fn wait(&self) -> Option<i32> {
80        loop {
81            let code = {
82                let mut child = self.child.lock().await;
83                match child.as_mut() {
84                    None => Some(-1),
85                    Some(child) => match child.try_wait() {
86                        Ok(Some(status)) => Some(status.code().unwrap_or(-1)),
87                        Ok(None) => None,
88                        // A terminal whose child handle can no longer be
89                        // polled must not leave `wait_for_exit` spinning
90                        // forever. Preserve the protocol's integer exit
91                        // contract with an unknown/failure sentinel.
92                        Err(_) => Some(-1),
93                    },
94                }
95            };
96            if code.is_some() {
97                while self.output_readers.load(Ordering::Acquire) != 0 {
98                    tokio::task::yield_now().await;
99                }
100                return code;
101            }
102            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
103        }
104    }
105
106    async fn exit_code(&self) -> Option<i32> {
107        let mut child = self.child.lock().await;
108        child
109            .as_mut()
110            .and_then(|child| child.try_wait().ok().flatten())
111            .map(|status| status.code().unwrap_or(-1))
112    }
113}
114
115#[derive(Clone, Debug)]
116pub struct HostUpdate {
117    pub event: AgentEvent,
118    pub effects: Vec<Effect>,
119}
120
121#[derive(Clone, Debug, Eq, PartialEq)]
122pub enum AdapterError {
123    Unsupported(&'static str),
124    Spawn(String),
125    Transport(String),
126    Protocol(String),
127}
128
129impl std::fmt::Display for AdapterError {
130    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
131        match self {
132            Self::Unsupported(operation) => write!(formatter, "unsupported operation: {operation}"),
133            Self::Spawn(error) => write!(formatter, "unable to launch agent: {error}"),
134            Self::Transport(error) => write!(formatter, "agent transport error: {error}"),
135            Self::Protocol(error) => write!(formatter, "agent protocol error: {error}"),
136        }
137    }
138}
139
140impl std::error::Error for AdapterError {}
141
142fn floor_char_boundary(text: &str, index: usize) -> usize {
143    let mut index = index.min(text.len());
144    while index > 0 && !text.is_char_boundary(index) {
145        index -= 1;
146    }
147    index
148}
149
150/// A shell-free argv parser for commands stored in the agent catalog.
151///
152/// Agent commands are configuration data, not shell snippets: expansion and
153/// pipelines are intentionally unsupported. We still accept the quoting users
154/// expect when entering a command (`'...'`, `"..."`, and backslash escapes),
155/// so a configured executable or argument containing whitespace is passed to
156/// `Command` as one argument and cannot be re-split accidentally.
157#[derive(Clone, Debug, Eq, PartialEq)]
158pub enum CommandParseError {
159    Empty,
160    UnterminatedQuote,
161    TrailingEscape,
162}
163
164impl std::fmt::Display for CommandParseError {
165    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
166        let message = match self {
167            Self::Empty => "command is empty",
168            Self::UnterminatedQuote => "command contains an unterminated quote",
169            Self::TrailingEscape => "command ends with an incomplete escape",
170        };
171        formatter.write_str(message)
172    }
173}
174
175impl std::error::Error for CommandParseError {}
176
177/// Parse a configured command into an executable and argv without invoking a
178/// shell. This is shared by native and ACP adapters.
179pub fn parse_command_line(command: &str) -> Result<(String, Vec<String>), CommandParseError> {
180    let mut argv = Vec::new();
181    let mut argument = String::new();
182    let mut quoted = None;
183    let mut escaped = false;
184    let mut started = false;
185
186    for character in command.chars() {
187        if escaped {
188            argument.push(character);
189            escaped = false;
190            started = true;
191            continue;
192        }
193        match (quoted, character) {
194            (_, '\\') if quoted != Some('\'') => {
195                escaped = true;
196                started = true;
197            }
198            (None, '\'' | '"') => {
199                quoted = Some(character);
200                started = true;
201            }
202            (Some(quote), character) if character == quote => quoted = None,
203            (None, character) if character.is_whitespace() => {
204                if started {
205                    argv.push(std::mem::take(&mut argument));
206                    started = false;
207                }
208            }
209            (_, character) => {
210                argument.push(character);
211                started = true;
212            }
213        }
214    }
215
216    if escaped {
217        return Err(CommandParseError::TrailingEscape);
218    }
219    if quoted.is_some() {
220        return Err(CommandParseError::UnterminatedQuote);
221    }
222    if started {
223        argv.push(argument);
224    }
225    let Some((program, args)) = argv.split_first() else {
226        return Err(CommandParseError::Empty);
227    };
228    Ok((program.clone(), args.to_vec()))
229}
230
231/// Kill and reap a child process. Tokio intentionally does not reap a child
232/// when its handle is dropped, so every adapter shutdown path must await this
233/// helper before releasing the handle.
234async fn terminate_child(child: &mut Child) -> AdapterResult<()> {
235    #[cfg(unix)]
236    if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
237        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
238        signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
239    }
240    let kill_error = child.start_kill().err();
241    let wait_error = child.wait().await.err();
242    if let Some(error) = kill_error.or(wait_error) {
243        return Err(AdapterError::Transport(error.to_string()));
244    }
245    Ok(())
246}
247
248#[cfg(unix)]
249fn signal_isolated_process_group(child: &Child, signal: nix::sys::signal::Signal) -> bool {
250    use nix::{
251        sys::signal::killpg,
252        unistd::{Pid, getpgid, getpgrp},
253    };
254
255    let Some(raw_pid) = child.id().and_then(|pid| i32::try_from(pid).ok()) else {
256        return false;
257    };
258    let pid = Pid::from_raw(raw_pid);
259    // `process_group(0)` makes the child its own group leader. Verify that
260    // invariant at signal time and also reject CodeSwarm's own group. This
261    // keeps descendant cleanup without any possibility of reaching the
262    // containing shell or terminal multiplexer.
263    if getpgid(Some(pid)).ok() == Some(pid) && pid != getpgrp() {
264        let _ = killpg(pid, signal);
265        true
266    } else {
267        false
268    }
269}
270
271fn isolate_process_group(command: &mut Command) {
272    #[cfg(unix)]
273    command.process_group(0);
274}
275
276/// Drain diagnostics without allowing a noisy peer to block on stderr or
277/// retain an unbounded failure report.
278async fn drain_bounded<R>(mut reader: R, limit: usize) -> String
279where
280    R: AsyncRead + Unpin,
281{
282    let mut bytes = Vec::new();
283    let mut chunk = [0_u8; 4096];
284    while let Ok(count) = reader.read(&mut chunk).await {
285        if count == 0 {
286            break;
287        }
288        bytes.extend_from_slice(&chunk[..count]);
289        if bytes.len() > limit {
290            let keep_from = bytes.len() - limit;
291            bytes.drain(..keep_from);
292        }
293    }
294    String::from_utf8_lossy(&bytes).trim().to_owned()
295}
296
297/// Read one JSONL frame without allowing a peer to allocate an unbounded
298/// string before the size check runs. `AsyncBufReadExt::read_line` checks only
299/// after it has appended the complete line, so the bounded fill/consume loop
300/// below is intentional.
301async fn read_bounded_line<R>(reader: &mut R) -> AdapterResult<String>
302where
303    R: AsyncBufRead + Unpin,
304{
305    let mut bytes = Vec::with_capacity(4096);
306    loop {
307        let buffer = reader
308            .fill_buf()
309            .await
310            .map_err(|error| AdapterError::Transport(error.to_string()))?;
311        if buffer.is_empty() {
312            if bytes.is_empty() {
313                return Err(AdapterError::Transport("ACP stream closed".into()));
314            }
315            break;
316        }
317        let newline = buffer.iter().position(|byte| *byte == b'\n');
318        let available = newline.map_or(buffer.len(), |index| index + 1);
319        let remaining = MAX_ACP_LINE_BYTES
320            .saturating_add(1)
321            .saturating_sub(bytes.len());
322        if available > remaining {
323            reader.consume(remaining);
324            return Err(AdapterError::Protocol(format!(
325                "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
326            )));
327        }
328        bytes.extend_from_slice(&buffer[..available]);
329        reader.consume(available);
330        if newline.is_some() {
331            break;
332        }
333    }
334    if bytes.len() > MAX_ACP_LINE_BYTES {
335        return Err(AdapterError::Protocol(format!(
336            "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
337        )));
338    }
339    String::from_utf8(bytes).map_err(|error| AdapterError::Protocol(error.to_string()))
340}
341
342async fn drain_terminal_output<R>(
343    mut reader: R,
344    output: Arc<Mutex<Vec<u8>>>,
345    truncated: Arc<AtomicBool>,
346    output_readers: Arc<AtomicUsize>,
347    limit: usize,
348) where
349    R: AsyncRead + Unpin,
350{
351    let mut chunk = [0_u8; 4096];
352    while let Ok(count) = reader.read(&mut chunk).await {
353        if count == 0 {
354            break;
355        }
356        if let Ok(mut bytes) = output.lock() {
357            let remaining = limit.saturating_sub(bytes.len());
358            if count > remaining {
359                bytes.extend_from_slice(&chunk[..remaining]);
360                truncated.store(true, Ordering::Release);
361            } else {
362                bytes.extend_from_slice(&chunk[..count]);
363            }
364        }
365    }
366    output_readers.fetch_sub(1, Ordering::AcqRel);
367}
368
369/// Uniform control plane for ACP and custom command-line adapters.
370#[async_trait]
371pub trait AgentAdapter: Send {
372    fn slot(&self) -> RosterSlot;
373    /// Human-readable adapter identity used in the first-turn roster context.
374    /// Catalog-backed launchers can override this with their display name;
375    /// direct command adapters still get a deterministic fallback.
376    fn display_name(&self) -> String {
377        format!("Agent {}", self.slot().saturating_add(1))
378    }
379    /// Return the protocol session handle when this adapter has one. Custom
380    /// adapters may omit it; runtime metadata then still preserves roster
381    /// identity without claiming the session is resumable.
382    fn session_id(&self) -> Option<String> {
383        None
384    }
385    /// Stable protocol label recorded in the runtime session snapshot.
386    /// Custom adapters may override this when they expose a resumable
387    /// protocol distinct from the built-in native and ACP bridges.
388    fn protocol(&self) -> &'static str {
389        "custom"
390    }
391    /// Whether the adapter lost its underlying transport and needs a
392    /// coordinator-owned reload before the next prompt. Most adapters keep
393    /// their transport alive between turns; ACP overrides this for EOF or
394    /// broken-pipe recovery.
395    fn needs_restart(&self) -> bool {
396        false
397    }
398    fn capabilities(&self) -> AgentCapabilities;
399    async fn start(&mut self) -> AdapterResult<()>;
400    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()>;
401    async fn cancel(&mut self) -> AdapterResult<bool>;
402    async fn answer_permission(
403        &mut self,
404        request_id: String,
405        answer: PermissionAnswer,
406    ) -> AdapterResult<()>;
407    async fn set_mode(&mut self, mode: String) -> AdapterResult<()>;
408    async fn set_model(&mut self, _model: String) -> AdapterResult<()> {
409        Err(AdapterError::Unsupported("set_model"))
410    }
411    async fn reload(&mut self) -> AdapterResult<()>;
412    async fn stop(&mut self) -> AdapterResult<()>;
413    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>>;
414}
415
416/// Preserve a stable CodeSwarm roster slot when an already-running adapter is
417/// moved to another logical position (for example, a roster reorder). Native
418/// protocol implementations keep their original slot internally, while this
419/// boundary rewrites the normalized event identity seen by the reducer.
420struct SlotMappedAdapter {
421    logical_slot: RosterSlot,
422    inner: Box<dyn AgentAdapter>,
423}
424
425impl std::fmt::Debug for SlotMappedAdapter {
426    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
427        formatter
428            .debug_struct("SlotMappedAdapter")
429            .field("logical_slot", &self.logical_slot)
430            .field("inner_slot", &self.inner.slot())
431            .finish_non_exhaustive()
432    }
433}
434
435fn restored_history_event(event: AgentEvent) -> AgentEvent {
436    use crate::HistoryContent;
437    let (slot, content) = match event {
438        AgentEvent::UserText { slot, text } => (slot, HistoryContent::UserText(text)),
439        AgentEvent::Text { slot, text } => (slot, HistoryContent::Text(text)),
440        AgentEvent::Thought { slot, text } => (slot, HistoryContent::Thought(text)),
441        AgentEvent::Tool { slot, update } => (slot, HistoryContent::Tool(update)),
442        other => return other,
443    };
444    AgentEvent::History { slot, content }
445}
446
447fn map_event_slot(event: AgentEvent, slot: RosterSlot) -> AgentEvent {
448    match event {
449        AgentEvent::SessionMetadataUpdated { metadata } => {
450            AgentEvent::SessionMetadataUpdated { metadata }
451        }
452        AgentEvent::History { content, .. } => AgentEvent::History { slot, content },
453        AgentEvent::GoalUpdated { goal } => AgentEvent::GoalUpdated { goal },
454        AgentEvent::RosterUpdated { update } => AgentEvent::RosterUpdated { update },
455        AgentEvent::Ready { capabilities, .. } => AgentEvent::Ready { slot, capabilities },
456        AgentEvent::TurnStarted { .. } => AgentEvent::TurnStarted { slot },
457        AgentEvent::ModesReplaced {
458            modes,
459            current_mode,
460            ..
461        } => AgentEvent::ModesReplaced {
462            slot,
463            modes,
464            current_mode,
465        },
466        AgentEvent::ModeUpdated { current_mode, .. } => {
467            AgentEvent::ModeUpdated { slot, current_mode }
468        }
469        AgentEvent::ModelsReplaced {
470            config_id,
471            models,
472            current_model,
473            ..
474        } => AgentEvent::ModelsReplaced {
475            slot,
476            config_id,
477            models,
478            current_model,
479        },
480        AgentEvent::ModelUpdated { current_model, .. } => AgentEvent::ModelUpdated {
481            slot,
482            current_model,
483        },
484        AgentEvent::UserText { text, .. } => AgentEvent::UserText { slot, text },
485        AgentEvent::CommandsReplaced { commands, .. } => {
486            AgentEvent::CommandsReplaced { slot, commands }
487        }
488        AgentEvent::UsageUpdated { usage, .. } => AgentEvent::UsageUpdated { slot, usage },
489        AgentEvent::Text { text, .. } => AgentEvent::Text { slot, text },
490        AgentEvent::Thought { text, .. } => AgentEvent::Thought { slot, text },
491        AgentEvent::Tool { update, .. } => AgentEvent::Tool { slot, update },
492        AgentEvent::Permission { request, .. } => AgentEvent::Permission { slot, request },
493        AgentEvent::Terminal { event, .. } => AgentEvent::Terminal { slot, event },
494        AgentEvent::TurnComplete { .. } => AgentEvent::TurnComplete { slot },
495        AgentEvent::BatchComplete { elapsed } => AgentEvent::BatchComplete { elapsed },
496        AgentEvent::UsageLimitReached { detail, .. } => {
497            AgentEvent::UsageLimitReached { slot, detail }
498        }
499        AgentEvent::Failed {
500            started, detail, ..
501        } => AgentEvent::Failed {
502            slot,
503            started,
504            detail,
505        },
506    }
507}
508
509#[async_trait]
510impl AgentAdapter for SlotMappedAdapter {
511    fn slot(&self) -> RosterSlot {
512        self.logical_slot
513    }
514
515    fn display_name(&self) -> String {
516        self.inner.display_name()
517    }
518
519    fn session_id(&self) -> Option<String> {
520        self.inner.session_id()
521    }
522
523    fn protocol(&self) -> &'static str {
524        self.inner.protocol()
525    }
526
527    fn capabilities(&self) -> AgentCapabilities {
528        self.inner.capabilities()
529    }
530
531    fn needs_restart(&self) -> bool {
532        self.inner.needs_restart()
533    }
534
535    async fn start(&mut self) -> AdapterResult<()> {
536        self.inner.start().await
537    }
538
539    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
540        self.inner.send_prompt(prompt).await
541    }
542
543    async fn cancel(&mut self) -> AdapterResult<bool> {
544        self.inner.cancel().await
545    }
546
547    async fn answer_permission(
548        &mut self,
549        request_id: String,
550        answer: PermissionAnswer,
551    ) -> AdapterResult<()> {
552        self.inner.answer_permission(request_id, answer).await
553    }
554
555    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
556        self.inner.set_mode(mode).await
557    }
558
559    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
560        self.inner.set_model(model).await
561    }
562
563    async fn reload(&mut self) -> AdapterResult<()> {
564        self.inner.reload().await
565    }
566
567    async fn stop(&mut self) -> AdapterResult<()> {
568        self.inner.stop().await
569    }
570
571    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
572        let slot = self.logical_slot;
573        self.inner
574            .next_event()
575            .await
576            .map(|result| result.map(|event| map_event_slot(event, slot)))
577    }
578}
579
580/// Owns one adapter and feeds normalized events through the deterministic core
581/// reducer. The UI consumes effects and state snapshots.
582pub struct AdapterHost {
583    adapter: Box<dyn AgentAdapter>,
584    pub state: SessionState,
585    pub last_error: Option<String>,
586    event_log: Option<EventLog>,
587}
588
589impl std::fmt::Debug for AdapterHost {
590    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
591        formatter
592            .debug_struct("AdapterHost")
593            .field("state", &self.state)
594            .field("last_error", &self.last_error)
595            .field("event_log", &self.event_log)
596            .finish_non_exhaustive()
597    }
598}
599
600impl AdapterHost {
601    pub fn new(adapter: Box<dyn AgentAdapter>, event_log: Option<EventLog>) -> Self {
602        let slot = adapter.slot();
603        Self {
604            adapter,
605            state: SessionState::new(slot.saturating_add(1)),
606            last_error: None,
607            event_log,
608        }
609    }
610
611    pub async fn start(&mut self) -> AdapterResult<()> {
612        self.adapter.start().await
613    }
614
615    pub async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
616        self.adapter.send_prompt(prompt).await
617    }
618
619    pub async fn cancel(&mut self) -> AdapterResult<bool> {
620        self.adapter.cancel().await
621    }
622
623    pub async fn answer_permission(
624        &mut self,
625        request_id: String,
626        answer: PermissionAnswer,
627    ) -> AdapterResult<()> {
628        self.adapter.answer_permission(request_id, answer).await
629    }
630
631    pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
632        self.adapter.set_mode(mode).await
633    }
634
635    pub async fn set_model(&mut self, model: String) -> AdapterResult<()> {
636        self.adapter.set_model(model).await
637    }
638
639    pub async fn reload(&mut self) -> AdapterResult<()> {
640        self.adapter.reload().await?;
641        let slot = self.adapter.slot();
642        if let Some(agent) = self.state.slots.get_mut(slot) {
643            agent.active = true;
644            agent.capabilities = self.adapter.capabilities();
645        }
646        self.last_error = None;
647        Ok(())
648    }
649
650    pub async fn stop(&mut self) -> AdapterResult<()> {
651        self.adapter.stop().await
652    }
653
654    pub async fn next_effects(&mut self) -> Option<AdapterResult<Vec<Effect>>> {
655        Some(self.next_update().await?.map(|update| update.effects))
656    }
657
658    pub async fn next_update(&mut self) -> Option<AdapterResult<HostUpdate>> {
659        let event = match self.adapter.next_event().await {
660            None => return None,
661            Some(Err(error)) => {
662                self.last_error = Some(error.to_string());
663                let slot = self.adapter.slot();
664                let failure = AgentEvent::Failed {
665                    slot,
666                    started: true,
667                    detail: error.to_string(),
668                };
669                let effects = reduce(&mut self.state, failure.clone());
670                return Some(Ok(HostUpdate {
671                    event: failure,
672                    effects,
673                }));
674            }
675            Some(Ok(event)) => event,
676        };
677        if let Some(log) = &self.event_log
678            && let Err(error) = log.append(&event)
679        {
680            return Some(Err(AdapterError::Transport(error.to_string())));
681        }
682        let effects = reduce(&mut self.state, event.clone());
683        Some(Ok(HostUpdate { event, effects }))
684    }
685
686    pub fn adapter(&self) -> &dyn AgentAdapter {
687        &*self.adapter
688    }
689
690    pub fn session_id(&self) -> Option<String> {
691        self.adapter.session_id()
692    }
693
694    /// Move this host to a new logical roster slot. The adapter process is not
695    /// restarted; only normalized event identities and reducer state move.
696    fn remap(self, logical_slot: RosterSlot) -> Self {
697        let old_slot = self.adapter.slot();
698        if old_slot == logical_slot {
699            return self;
700        }
701
702        let mut state = self.state;
703        if state.slots.len() <= logical_slot {
704            state
705                .slots
706                .resize(logical_slot.saturating_add(1), Default::default());
707        }
708        if let Some(agent) = state.slots.get(old_slot).cloned() {
709            state.slots[logical_slot] = agent;
710        }
711        if state.active_slot == Some(old_slot) {
712            state.active_slot = Some(logical_slot);
713        }
714        for (slot, _) in &mut state.queued_prompts {
715            if *slot == old_slot {
716                *slot = logical_slot;
717            }
718        }
719        for (slot, _) in &mut state.public_text {
720            if *slot == old_slot {
721                *slot = logical_slot;
722            }
723        }
724
725        Self {
726            adapter: Box::new(SlotMappedAdapter {
727                logical_slot,
728                inner: self.adapter,
729            }),
730            state,
731            last_error: self.last_error,
732            event_log: self.event_log,
733        }
734    }
735}
736
737/// Sequential multi-adapter runner. It intentionally never polls two
738/// adapters concurrently: the next prompt depends on the prior response.
739pub struct RelayHost {
740    goal: Option<crate::goal::Goal>,
741    hosts: Vec<AdapterHost>,
742    relay: Relay,
743    introduced: Vec<bool>,
744    roster_names: Vec<String>,
745    roster_identities: Vec<String>,
746    roster_launch_specs: Vec<(String, String)>,
747    desired_policy: String,
748    metadata_writer: Option<BufferedSessionMetadataStore>,
749    metadata_workspace: Option<String>,
750    dispatches: Vec<(RosterSlot, String)>,
751    last_public_dispatch: Option<RosterSlot>,
752    pair_implementer: Option<RosterSlot>,
753    event_sink: Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
754    cancel_requested: Arc<AtomicBool>,
755    cancel_notify: Arc<Notify>,
756}
757
758/// A clonable signal used by a terminal control loop to interrupt the active
759/// relay turn without borrowing the relay while its adapter is being polled.
760#[derive(Clone, Debug)]
761pub struct RelayCancellation {
762    requested: Arc<AtomicBool>,
763    notify: Arc<Notify>,
764}
765
766/// A permission response delivered while a relay turn is still streaming.
767///
768/// Relay turns own the active adapter mutably, so permission answers must be
769/// consumed inside that turn rather than deferred by the outer coordinator.
770#[derive(Debug)]
771pub struct RelayPermissionAnswer {
772    pub slot: RosterSlot,
773    pub request_id: String,
774    pub answer: PermissionAnswer,
775}
776
777#[cfg(test)]
778const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
779#[cfg(not(test))]
780const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
781
782#[cfg(test)]
783const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(20);
784#[cfg(not(test))]
785const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
786
787/// Bound third-party cancellation hooks so a broken adapter cannot freeze the
788/// terminal control loop.  Native and ACP adapters normally return promptly;
789/// the timeout is for custom adapters whose cancellation future may never
790/// resolve.
791async fn cancel_with_timeout(host: &mut AdapterHost) -> AdapterResult<bool> {
792    tokio::time::timeout(CANCEL_TIMEOUT, host.cancel())
793        .await
794        .map_err(|_| AdapterError::Transport("adapter cancellation timed out".into()))?
795}
796
797fn canonical_policy_id(policy: &str) -> &str {
798    match policy {
799        "plan" => "codeswarm:mode:plan",
800        "default" | "manual" => "codeswarm:mode:manual",
801        "accept-edits" => "codeswarm:mode:accept-edits",
802        "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
803        other => other,
804    }
805}
806
807async fn apply_policy_to_host(host: &mut AdapterHost, policy: &str) -> AdapterResult<()> {
808    if !host.adapter().capabilities().supports_modes {
809        return Ok(());
810    }
811    let policy_id = canonical_policy_id(policy);
812    let slot = host.adapter().slot();
813    let advertised = host
814        .state
815        .slots
816        .get(slot)
817        .map(|agent| agent.modes.as_slice())
818        .unwrap_or_default();
819    let native = if advertised.is_empty() {
820        match policy_id {
821            "codeswarm:mode:plan" => "plan".into(),
822            "codeswarm:mode:manual" => "default".into(),
823            "codeswarm:mode:accept-edits" => "accept-edits".into(),
824            "codeswarm:mode:full-access" => "full-access".into(),
825            other => other.into(),
826        }
827    } else {
828        crate::policy::resolve(policy_id, advertised)
829            .map(|mode| mode.id)
830            .ok_or(AdapterError::Unsupported(
831                "desired policy is unavailable for adapter",
832            ))?
833    };
834    host.set_mode(native).await
835}
836
837async fn refresh_mode_catalog(
838    host: &mut AdapterHost,
839    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
840) -> AdapterResult<bool> {
841    if !host.adapter().capabilities().supports_modes {
842        return Ok(false);
843    }
844    tokio::time::timeout(std::time::Duration::from_secs(2), async {
845        let mut ready_seen = false;
846        loop {
847            let update = host.next_update().await.ok_or_else(|| {
848                AdapterError::Transport("adapter ended before advertising modes".into())
849            })??;
850            let catalog_ready = matches!(update.event, AgentEvent::ModesReplaced { .. });
851            ready_seen |= matches!(update.event, AgentEvent::Ready { .. });
852            if let Some(sink) = event_sink {
853                sink(update.event);
854            }
855            if catalog_ready {
856                return Ok(ready_seen);
857            }
858        }
859    })
860    .await
861    .map_err(|_| AdapterError::Transport("adapter mode catalog timed out".into()))?
862}
863
864async fn refresh_adapter_startup(
865    host: &mut AdapterHost,
866    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
867) -> AdapterResult<()> {
868    let ready_seen = refresh_mode_catalog(host, event_sink).await?;
869    if ready_seen || !matches!(host.adapter().protocol(), "native" | "acp") {
870        return Ok(());
871    }
872    tokio::time::timeout(std::time::Duration::from_secs(2), async {
873        loop {
874            let update = host.next_update().await.ok_or_else(|| {
875                AdapterError::Transport("adapter ended before becoming ready".into())
876            })??;
877            let ready = matches!(update.event, AgentEvent::Ready { .. });
878            if let Some(sink) = event_sink {
879                sink(update.event);
880            }
881            if ready {
882                return Ok(());
883            }
884        }
885    })
886    .await
887    .map_err(|_| AdapterError::Transport("adapter ready handshake timed out".into()))?
888}
889
890fn public_context_speaker(name: &str) -> String {
891    let now = time::OffsetDateTime::now_local().unwrap_or_else(|_| time::OffsetDateTime::now_utc());
892    format!("{name} {:02}:{:02}", now.hour(), now.minute())
893}
894
895impl RelayCancellation {
896    pub fn request(&self) {
897        self.requested.store(true, Ordering::Release);
898        self.notify.notify_one();
899    }
900}
901
902impl std::fmt::Debug for RelayHost {
903    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
904        formatter
905            .debug_struct("RelayHost")
906            .field("hosts", &self.hosts)
907            .field("relay", &self.relay)
908            .field("dispatches", &self.dispatches)
909            .field("event_sink", &self.event_sink.is_some())
910            .field(
911                "cancel_requested",
912                &self.cancel_requested.load(Ordering::Acquire),
913            )
914            .finish()
915    }
916}
917
918impl RelayHost {
919    pub fn new(hosts: Vec<AdapterHost>, max_rounds: usize) -> Result<Self, AdapterError> {
920        if hosts.is_empty() {
921            return Err(AdapterError::Unsupported("relay requires an adapter"));
922        }
923        Ok(Self {
924            relay: Relay::new(hosts.len(), max_rounds),
925            introduced: vec![false; hosts.len()],
926            roster_names: hosts
927                .iter()
928                .map(|host| host.adapter().display_name())
929                .collect(),
930            roster_identities: hosts
931                .iter()
932                .map(|host| host.adapter().display_name())
933                .collect(),
934            roster_launch_specs: Vec::new(),
935            desired_policy: crate::policy::DEFAULT_POLICY_ID.into(),
936            metadata_writer: None,
937            metadata_workspace: None,
938            hosts,
939            dispatches: Vec::new(),
940            last_public_dispatch: None,
941            pair_implementer: None,
942            goal: None,
943            event_sink: None,
944            cancel_requested: Arc::new(AtomicBool::new(false)),
945            cancel_notify: Arc::new(Notify::new()),
946        })
947    }
948
949    /// Send each normalized event to a client while a turn is being drained.
950    /// The callback runs synchronously on the relay task and should only
951    /// enqueue the event; expensive rendering must happen outside the callback.
952    pub fn set_event_sink<F>(&mut self, sink: F)
953    where
954        F: Fn(AgentEvent) + Send + Sync + 'static,
955    {
956        self.event_sink = Some(Arc::new(sink));
957    }
958
959    /// Install catalog-backed names for first-turn introductions. Direct
960    /// scripted/custom adapters retain deterministic `Agent N` fallbacks.
961    pub fn set_roster_names(&mut self, names: Vec<String>) {
962        if names.len() == self.hosts.len() {
963            self.roster_names = names;
964        }
965    }
966
967    /// Install catalog identities used by the launcher and persisted session
968    /// metadata. Names remain a separate display concern, so custom adapters
969    /// can retain friendly labels while still restoring by stable identity.
970    pub fn set_roster_identities(&mut self, identities: Vec<String>) {
971        if identities.len() == self.hosts.len() {
972            self.roster_identities = identities;
973        }
974    }
975
976    pub fn set_roster_launch_specs(&mut self, specs: Vec<(String, String)>) {
977        if specs.len() == self.hosts.len() {
978            self.roster_launch_specs = specs;
979        }
980    }
981
982    /// Attach a background metadata writer. Runtime changes enqueue complete
983    /// snapshots; no metadata filesystem work runs on the terminal thread.
984    pub fn set_session_metadata_writer(&mut self, writer: BufferedSessionMetadataStore) {
985        self.metadata_writer = Some(writer);
986    }
987
988    pub fn set_session_metadata_workspace(&mut self, workspace: impl Into<String>) {
989        self.metadata_workspace = Some(workspace.into());
990    }
991
992    /// Build the coordinator-owned session snapshot for the active roster.
993    pub fn session_metadata(&self) -> SessionMetadata {
994        let active = self.relay.active_slots().collect::<Vec<_>>();
995        let mut data = serde_json::Map::new();
996        data.insert(
997            "goal".into(),
998            serde_json::to_value(&self.goal).expect("goal serializes"),
999        );
1000        if let Some(workspace) = &self.metadata_workspace {
1001            data.insert("cwd".into(), serde_json::Value::String(workspace.clone()));
1002        }
1003        data.insert(
1004            "title".into(),
1005            serde_json::Value::String("CodeSwarm".into()),
1006        );
1007        data.insert(
1008            "agents".into(),
1009            serde_json::Value::Array(
1010                active
1011                    .into_iter()
1012                    .filter_map(|slot| {
1013                        let host = self.hosts.get(slot)?;
1014                        let (protocol, command) = self.roster_launch_specs.get(slot)?;
1015                        let mut agent = serde_json::Map::new();
1016                        agent.insert("slot".into(), serde_json::json!(slot));
1017                        agent.insert(
1018                            "name".into(),
1019                            serde_json::Value::String(
1020                                self.roster_names
1021                                    .get(slot)
1022                                    .cloned()
1023                                    .unwrap_or_else(|| host.adapter().display_name()),
1024                            ),
1025                        );
1026                        agent.insert(
1027                            "identity".into(),
1028                            serde_json::Value::String(
1029                                self.roster_identities
1030                                    .get(slot)
1031                                    .cloned()
1032                                    .unwrap_or_else(|| host.adapter().display_name()),
1033                            ),
1034                        );
1035                        agent.insert(
1036                            "protocol".into(),
1037                            serde_json::Value::String(protocol.clone()),
1038                        );
1039                        agent.insert("command".into(), serde_json::Value::String(command.clone()));
1040                        agent.insert(
1041                            "supports_load_session".into(),
1042                            serde_json::Value::Bool(
1043                                host.adapter().capabilities().supports_session_load,
1044                            ),
1045                        );
1046                        if let Some(session_id) = host.session_id() {
1047                            agent
1048                                .insert("session_id".into(), serde_json::Value::String(session_id));
1049                        }
1050                        Some(serde_json::Value::Object(agent))
1051                    })
1052                    .collect(),
1053            ),
1054        );
1055        SessionMetadata::new(data)
1056    }
1057
1058    fn queue_session_metadata(&self) -> AdapterResult<()> {
1059        let metadata = self.session_metadata();
1060        if let Some(sink) = &self.event_sink {
1061            sink(AgentEvent::SessionMetadataUpdated {
1062                metadata: metadata.to_value(),
1063            });
1064        }
1065        if let Some(writer) = &self.metadata_writer {
1066            writer
1067                .write(metadata)
1068                .map_err(|error| AdapterError::Transport(error.to_string()))?;
1069        }
1070        Ok(())
1071    }
1072
1073    pub fn restore_goal(&mut self, goal: Option<crate::goal::Goal>) {
1074        self.goal = goal;
1075        if let Some(sink) = &self.event_sink {
1076            sink(AgentEvent::GoalUpdated {
1077                goal: self.goal.clone(),
1078            });
1079        }
1080    }
1081
1082    pub fn apply_goal(
1083        &mut self,
1084        command: crate::goal::GoalCommand,
1085    ) -> Result<Option<String>, String> {
1086        let task = crate::goal::apply(&mut self.goal, command)?;
1087        if let Some(sink) = &self.event_sink {
1088            sink(AgentEvent::GoalUpdated {
1089                goal: self.goal.clone(),
1090            });
1091        }
1092        self.queue_session_metadata()
1093            .map_err(|error| error.to_string())?;
1094        Ok(task)
1095    }
1096
1097    pub fn roster_names(&self) -> &[String] {
1098        &self.roster_names
1099    }
1100
1101    pub fn session_ids(&self) -> Vec<Option<String>> {
1102        self.hosts.iter().map(AdapterHost::session_id).collect()
1103    }
1104
1105    pub async fn start(&mut self) -> AdapterResult<()> {
1106        self.start_isolating_failures("agent could not start").await
1107    }
1108
1109    /// Restore each provider independently. An expired saved handle must not
1110    /// tear down healthy peers; failed slots are immediately unroutable.
1111    pub async fn start_resuming(&mut self) -> AdapterResult<()> {
1112        self.start_isolating_failures("saved session could not be restored")
1113            .await
1114    }
1115
1116    async fn start_isolating_failures(&mut self, failure_prefix: &str) -> AdapterResult<()> {
1117        let event_sink = self.event_sink.clone();
1118        let policy = self.desired_policy.clone();
1119        let starts = self.hosts.iter_mut().map(|host| {
1120            let sink = event_sink.clone();
1121            let policy = policy.clone();
1122            async move {
1123                host.start().await?;
1124                refresh_adapter_startup(host, &sink).await?;
1125                apply_policy_to_host(host, &policy).await
1126            }
1127        });
1128        let results = futures::future::join_all(starts).await;
1129        let mut first_error = None;
1130        for (slot, result) in results.into_iter().enumerate() {
1131            if let Err(error) = result {
1132                let _ = self.relay.tombstone(slot);
1133                let _ = self.hosts[slot].stop().await;
1134                if let Some(sink) = &event_sink {
1135                    sink(AgentEvent::Failed {
1136                        slot,
1137                        started: false,
1138                        detail: format!("{failure_prefix}: {error}"),
1139                    });
1140                }
1141                if first_error.is_none() {
1142                    first_error = Some(error);
1143                }
1144            }
1145        }
1146        if self.relay.active_slots().next().is_none() {
1147            return Err(first_error
1148                .unwrap_or_else(|| AdapterError::Transport("no agents could be started".into())));
1149        }
1150        let _ = self.queue_session_metadata();
1151        Ok(())
1152    }
1153
1154    pub async fn stop(&mut self) -> AdapterResult<()> {
1155        // A third-party adapter can fail during shutdown (for example after
1156        // its transport has already disappeared). Always give every still-active
1157        // member a chance to clean up, then return the first error so callers
1158        // still get an actionable failure without leaking later processes.
1159        let mut first_error = None;
1160        let active = self.relay.active_slots().collect::<Vec<_>>();
1161        for slot in active {
1162            let host = &mut self.hosts[slot];
1163            if let Err(error) = host.stop().await
1164                && first_error.is_none()
1165            {
1166                first_error = Some(error);
1167            }
1168        }
1169        if let Some(writer) = &self.metadata_writer
1170            && let Err(error) = writer.flush()
1171            && first_error.is_none()
1172        {
1173            first_error = Some(AdapterError::Transport(error.to_string()));
1174        }
1175        first_error.map_or(Ok(()), Err)
1176    }
1177
1178    /// Forward a normalized permission answer to the adapter owning `slot`.
1179    /// This keeps protocol-specific response framing out of the relay and UI.
1180    pub async fn answer_permission(
1181        &mut self,
1182        slot: RosterSlot,
1183        request_id: String,
1184        answer: PermissionAnswer,
1185    ) -> AdapterResult<()> {
1186        let host = self
1187            .hosts
1188            .get_mut(slot)
1189            .ok_or_else(|| AdapterError::Transport("permission target is missing".into()))?;
1190        host.answer_permission(request_id, answer).await
1191    }
1192
1193    pub fn pause(&mut self) {
1194        self.relay.pause();
1195    }
1196
1197    pub fn resume(&mut self) {
1198        self.relay.resume();
1199    }
1200
1201    /// Select how future non-direct turns are routed. Queued prompts and the
1202    /// shared context remain intact when a user changes this setting.
1203    pub fn set_strategy(&mut self, strategy: CollaborationStrategy) {
1204        self.relay.set_strategy(strategy);
1205        if strategy != CollaborationStrategy::Pair {
1206            self.pair_implementer = None;
1207        }
1208        let _ = self.queue_session_metadata();
1209    }
1210
1211    pub fn strategy(&self) -> CollaborationStrategy {
1212        self.relay.strategy()
1213    }
1214
1215    pub fn roster_identity(&self, slot: RosterSlot) -> Option<&str> {
1216        self.roster_identities.get(slot).map(String::as_str)
1217    }
1218
1219    pub fn active_slot_for_identity(&self, identity: &str) -> Option<RosterSlot> {
1220        self.relay.active_slots().find(|slot| {
1221            self.roster_identity(*slot)
1222                .is_some_and(|candidate| candidate.eq_ignore_ascii_case(identity))
1223        })
1224    }
1225
1226    /// Apply one semantic mode to every active adapter that advertises mode
1227    /// support. Native adapters may translate the policy to their own IDs.
1228    pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
1229        let active = self.relay.active_slots().collect::<Vec<_>>();
1230        for slot in active {
1231            let Some(host) = self.hosts.get_mut(slot) else {
1232                continue;
1233            };
1234            if host.adapter().capabilities().supports_modes {
1235                host.set_mode(mode.clone()).await?;
1236            }
1237        }
1238        Ok(())
1239    }
1240
1241    pub async fn set_model(&mut self, slot: RosterSlot, model: String) -> AdapterResult<()> {
1242        let host = self
1243            .hosts
1244            .get_mut(slot)
1245            .ok_or_else(|| AdapterError::Transport("model target is missing".into()))?;
1246        if !host.adapter().capabilities().supports_models {
1247            return Err(AdapterError::Unsupported("set_model"));
1248        }
1249        host.set_model(model.clone()).await?;
1250        if let Some(sink) = &self.event_sink {
1251            sink(AgentEvent::ModelUpdated {
1252                slot,
1253                current_model: model,
1254            });
1255        }
1256        Ok(())
1257    }
1258
1259    /// Translate a CodeSwarm semantic policy to each adapter's currently
1260    /// advertised native mode, falling back to conventional IDs before the
1261    /// first catalog update arrives.
1262    pub async fn set_policy(&mut self, policy: String) -> AdapterResult<()> {
1263        let desired_policy = canonical_policy_id(&policy).to_owned();
1264        // Preflight every active catalog before mutating any adapter. A
1265        // mixed roster must never end up with only some agents on the new
1266        // policy when another selected agent cannot represent it.
1267        let active = self.relay.active_slots().collect::<Vec<_>>();
1268        for active_slot in &active {
1269            let Some(host) = self.hosts.get(*active_slot) else {
1270                continue;
1271            };
1272            if !host.adapter().capabilities().supports_modes {
1273                continue;
1274            }
1275            let advertised = host
1276                .state
1277                .slots
1278                .get(*active_slot)
1279                .map(|agent| agent.modes.as_slice())
1280                .unwrap_or_default();
1281            if !advertised.is_empty()
1282                && crate::policy::resolve(&desired_policy, advertised).is_none()
1283            {
1284                return Err(AdapterError::Unsupported(
1285                    "desired policy is unavailable for an active adapter",
1286                ));
1287            }
1288        }
1289        self.desired_policy = desired_policy.clone();
1290        let mut first_error = None;
1291        for active_slot in active {
1292            let Some(host) = self.hosts.get_mut(active_slot) else {
1293                continue;
1294            };
1295            if let Err(error) = apply_policy_to_host(host, &desired_policy).await
1296                && first_error.is_none()
1297            {
1298                first_error = Some(error);
1299            }
1300        }
1301        first_error.map_or(Ok(()), Err)
1302    }
1303
1304    pub async fn reload(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1305        let desired_policy = self.desired_policy.clone();
1306        let _ = self.queue_session_metadata();
1307        let event_sink = self.event_sink.clone();
1308        let host = self
1309            .hosts
1310            .get_mut(slot)
1311            .ok_or_else(|| AdapterError::Transport("reload target is missing".into()))?;
1312        host.reload().await?;
1313        refresh_adapter_startup(host, &event_sink).await?;
1314        apply_policy_to_host(host, &desired_policy).await?;
1315        if let Some(introduced) = self.introduced.get_mut(slot) {
1316            *introduced = false;
1317        }
1318        self.relay
1319            .reactivate(slot)
1320            .map_err(|error| AdapterError::Transport(error.into()))?;
1321        let _ = self.queue_session_metadata();
1322        if let Some(sink) = &self.event_sink {
1323            sink(AgentEvent::RosterUpdated {
1324                update: RosterUpdate::Reloaded { slot },
1325            });
1326        }
1327        let _ = self.relay.clear_limited(slot);
1328        Ok(())
1329    }
1330
1331    /// Stop and tombstone an agent while preserving its stable roster slot.
1332    pub async fn drop_agent(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1333        self.relay
1334            .drop_agent(slot)
1335            .map_err(|error| AdapterError::Transport(error.into()))?;
1336        let _stop_result = if let Some(host) = self.hosts.get_mut(slot) {
1337            host.stop().await
1338        } else {
1339            Ok(())
1340        };
1341        // Persist the tombstone even when a third-party process reports a
1342        // shutdown error; otherwise the next launch can resurrect a slot the
1343        // user explicitly removed.
1344        let _ = self.queue_session_metadata();
1345        if let Some(sink) = &self.event_sink {
1346            sink(AgentEvent::RosterUpdated {
1347                update: RosterUpdate::Dropped { slot },
1348            });
1349        }
1350        Ok(())
1351    }
1352
1353    /// Start and append a new adapter in the next stable roster slot. The
1354    /// adapter is started before it becomes visible to the relay, so a failed
1355    /// add cannot leave a half-active slot behind.
1356    pub async fn add_agent(
1357        &mut self,
1358        mut host: AdapterHost,
1359        name: impl Into<String>,
1360        identity: impl Into<String>,
1361        command: impl Into<String>,
1362    ) -> AdapterResult<RosterSlot> {
1363        let slot = self.hosts.len();
1364        if host.adapter().slot() != slot {
1365            return Err(AdapterError::Transport(
1366                "new adapter slot must append after the existing roster".into(),
1367            ));
1368        }
1369        if let Err(error) = host.start().await {
1370            let _ = host.stop().await;
1371            return Err(error);
1372        }
1373        if let Err(error) = refresh_adapter_startup(&mut host, &self.event_sink).await {
1374            let _ = host.stop().await;
1375            return Err(error);
1376        }
1377        if let Err(error) = apply_policy_to_host(&mut host, &self.desired_policy).await {
1378            let _ = host.stop().await;
1379            return Err(error);
1380        }
1381        let capabilities = host.adapter().capabilities();
1382        self.hosts.push(host);
1383        self.relay.add_agent();
1384        self.introduced.push(false);
1385        let name = name.into();
1386        let identity = identity.into();
1387        self.roster_names.push(name.clone());
1388        self.roster_identities.push(identity.clone());
1389        self.roster_launch_specs
1390            .push((self.hosts[slot].adapter().protocol().into(), command.into()));
1391        let _ = self.queue_session_metadata();
1392        if let Some(sink) = &self.event_sink {
1393            sink(AgentEvent::RosterUpdated {
1394                update: RosterUpdate::Added {
1395                    slot,
1396                    name,
1397                    identity,
1398                },
1399            });
1400            sink(AgentEvent::Ready { slot, capabilities });
1401        }
1402        Ok(slot)
1403    }
1404
1405    /// Reorder two live adapters without restarting either process. The
1406    /// logical slot wrapper keeps normalized events, reducer state, and
1407    /// permission routing aligned with the new roster order.
1408    pub fn swap_agents(&mut self, first: RosterSlot, second: RosterSlot) -> AdapterResult<()> {
1409        if first == second {
1410            return Ok(());
1411        }
1412        if first >= self.hosts.len() || second >= self.hosts.len() {
1413            return Err(AdapterError::Transport("roster slot out of range".into()));
1414        }
1415        self.relay
1416            .swap_agents(first, second)
1417            .map_err(|error| AdapterError::Transport(error.into()))?;
1418        let low = first.min(second);
1419        let high = first.max(second);
1420        let high_host = self.hosts.remove(high);
1421        let low_host = self.hosts.remove(low);
1422        self.hosts.insert(low, high_host.remap(low));
1423        self.hosts.insert(high, low_host.remap(high));
1424        self.roster_names.swap(first, second);
1425        self.roster_identities.swap(first, second);
1426        if self.roster_launch_specs.len() == self.hosts.len() {
1427            self.roster_launch_specs.swap(first, second);
1428        }
1429        self.introduced.swap(first, second);
1430        if self.pair_implementer == Some(first) {
1431            self.pair_implementer = Some(second);
1432        } else if self.pair_implementer == Some(second) {
1433            self.pair_implementer = Some(first);
1434        }
1435        if self.last_public_dispatch == Some(first) {
1436            self.last_public_dispatch = Some(second);
1437        } else if self.last_public_dispatch == Some(second) {
1438            self.last_public_dispatch = Some(first);
1439        }
1440        if let Some(sink) = &self.event_sink {
1441            sink(AgentEvent::RosterUpdated {
1442                update: RosterUpdate::Swapped { first, second },
1443            });
1444            sink(AgentEvent::Ready {
1445                slot: first,
1446                capabilities: self.hosts[first].adapter().capabilities(),
1447            });
1448            sink(AgentEvent::Ready {
1449                slot: second,
1450                capabilities: self.hosts[second].adapter().capabilities(),
1451            });
1452        }
1453        let _ = self.queue_session_metadata();
1454        Ok(())
1455    }
1456
1457    pub fn relay(&self) -> &Relay {
1458        &self.relay
1459    }
1460
1461    pub fn next_slot(&self) -> RosterSlot {
1462        self.hosts.len()
1463    }
1464
1465    pub fn relay_mut(&mut self) -> &mut Relay {
1466        &mut self.relay
1467    }
1468
1469    pub fn cancellation(&self) -> RelayCancellation {
1470        RelayCancellation {
1471            requested: Arc::clone(&self.cancel_requested),
1472            notify: Arc::clone(&self.cancel_notify),
1473        }
1474    }
1475
1476    /// Prompts sent to adapters, in causal dispatch order. This is useful to
1477    /// diagnostics and makes the context-routing boundary observable without
1478    /// exposing protocol-specific adapter internals.
1479    pub fn dispatches(&self) -> &[(RosterSlot, String)] {
1480        &self.dispatches
1481    }
1482
1483    pub async fn run_turn(
1484        &mut self,
1485        task: impl Into<String>,
1486        first_slot: RosterSlot,
1487    ) -> AdapterResult<RelayDecision> {
1488        self.run_turn_inner(task.into(), first_slot, None).await
1489    }
1490
1491    pub async fn run_turn_with_permissions(
1492        &mut self,
1493        task: impl Into<String>,
1494        first_slot: RosterSlot,
1495        permissions: &mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>,
1496    ) -> AdapterResult<RelayDecision> {
1497        self.run_turn_inner(task.into(), first_slot, Some(permissions))
1498            .await
1499    }
1500
1501    async fn run_turn_inner(
1502        &mut self,
1503        task: String,
1504        first_slot: RosterSlot,
1505        mut permissions: Option<&mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>>,
1506    ) -> AdapterResult<RelayDecision> {
1507        let decision = self.relay.begin(task, first_slot);
1508        let RelayDecision::Dispatch {
1509            slot,
1510            prompt,
1511            direct,
1512            can_stop,
1513        } = &decision
1514        else {
1515            return Ok(decision);
1516        };
1517        let event_sink = self.event_sink.clone();
1518        // Recover a broken ACP transport before collecting context or building
1519        // the prompt. Reload rewinds the slot's context watermark and clears
1520        // its introduction flag, so those values must be read afterward.
1521        if self
1522            .hosts
1523            .get(*slot)
1524            .is_some_and(|host| host.adapter().needs_restart())
1525            && let Err(error) = self.reload(*slot).await
1526        {
1527            let limited =
1528                report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1529            if limited {
1530                self.relay.finish(*slot, *direct, false);
1531            }
1532            let _ = self.queue_session_metadata();
1533            if limited {
1534                return Ok(decision);
1535            }
1536            return Err(error);
1537        }
1538        let speaker_name = self
1539            .roster_names
1540            .get(*slot)
1541            .cloned()
1542            .unwrap_or_else(|| self.hosts[*slot].adapter().display_name());
1543        let unseen = self.relay.unseen_context(*slot);
1544        // A non-direct prompt with text is a new public human turn. Record it
1545        // after collecting the recipient's existing unseen context so that
1546        // the selected agent receives the prompt once, while every peer gets
1547        // it through the shared journal on its next turn. Recording before
1548        // adapter I/O also preserves the job when the selected turn is
1549        // cancelled.
1550        if !*direct && !prompt.trim().is_empty() {
1551            if self.relay.shared_task().is_none() {
1552                self.relay.set_shared_task(prompt.clone());
1553            }
1554            self.relay
1555                .record_public(public_context_speaker("User"), prompt.clone());
1556        }
1557        let prompt = if unseen.is_empty() {
1558            prompt.clone()
1559        } else {
1560            format!("{prompt}\n\nPublic updates:\n{unseen}")
1561        };
1562        let introduction = if !self.introduced.get(*slot).copied().unwrap_or(false) {
1563            let self_name = speaker_name.clone();
1564            let roster = self
1565                .relay
1566                .active_slots()
1567                .map(|candidate| {
1568                    let name = self
1569                        .roster_names
1570                        .get(candidate)
1571                        .cloned()
1572                        .unwrap_or_else(|| self.hosts[candidate].adapter().display_name());
1573                    let role = if candidate == *slot { " — you" } else { "" };
1574                    format!("{}. {name}{role}", candidate.saturating_add(1))
1575                })
1576                .collect::<Vec<_>>();
1577            let shared_task = self
1578                .relay
1579                .shared_task()
1580                .filter(|task| *task != prompt)
1581                .map(|task| format!("\n\nShared task:\n{task}"))
1582                .unwrap_or_default();
1583            format!(
1584                "You are {self_name}.\nCodeSwarm roster (ordered):\n{}\n\
1585                 Turns relay sequentially through this roster. Treat the user request as the shared task; use timestamped public updates as conversation context.{shared_task}",
1586                roster.join("\n")
1587            )
1588        } else {
1589            String::new()
1590        };
1591        // Pair strategy only: make the implementer/reviewer handoff explicit.
1592        // Solo, roster, manual, and direct turns keep their existing prompt
1593        // text, and stop eligibility stays governed by the footer below.
1594        if !*direct && !*can_stop && self.relay.strategy() == CollaborationStrategy::Pair {
1595            self.pair_implementer = Some(*slot);
1596        }
1597        let role_block = if !*direct
1598            && self.relay.strategy() == CollaborationStrategy::Pair
1599            && self.relay.active_slots().count() >= 2
1600        {
1601            crate::workflow::pair_role(self.pair_implementer, *slot)
1602                .map(|role| {
1603                    let peer = match role {
1604                        crate::workflow::PairRole::Reviewer => self
1605                            .last_public_dispatch
1606                            .and_then(|previous| self.roster_names.get(previous).cloned()),
1607                        crate::workflow::PairRole::Implementer => None,
1608                    };
1609                    crate::workflow::role_fragment(role, peer.as_deref())
1610                })
1611                .unwrap_or_default()
1612        } else {
1613            String::new()
1614        };
1615        let effective_can_stop = *can_stop
1616            && !(self.relay.strategy() == CollaborationStrategy::Pair
1617                && self.relay.routable_slots().count() >= 2
1618                && self.pair_implementer == Some(*slot));
1619        let role_separator =
1620            if role_block.is_empty() || (introduction.is_empty() && prompt.is_empty()) {
1621                ""
1622            } else {
1623                "\n\n"
1624            };
1625        // Refresh target numbers every turn: roster members may have been
1626        // added, replaced, reordered, dropped, or usage-limited since introduction.
1627        let handoff_block = if self.relay.strategy() == CollaborationStrategy::Roster && !*direct {
1628            let targets = self
1629                .relay
1630                .routable_slots()
1631                .filter(|candidate| candidate != slot)
1632                .map(|candidate| {
1633                    format!(
1634                        "[CODESWARM:NEXT:{}] → {}",
1635                        candidate + 1,
1636                        self.roster_names
1637                            .get(candidate)
1638                            .cloned()
1639                            .unwrap_or_else(|| self.hosts[candidate].adapter().display_name())
1640                    )
1641                })
1642                .collect::<Vec<_>>()
1643                .join("\n");
1644            format!(
1645                "\n\nRoster handoff: to choose the next agent instead of normal roster order, end your final message with exactly one of the following markers:\n{targets}\nUse only a listed target, never yourself. The marker must follow all text, reasoning, and tool activity; only trailing whitespace is allowed. CodeSwarm hides it and routes at turn completion. Without a valid marker, normal roster order applies. Unavailable targets are ignored. Queued user input takes priority; turn limits and review-stop rules still apply. Choose either a handoff marker or the stop marker, not both."
1646            )
1647        } else {
1648            "\n\nAgent-directed handoff markers are disabled on this turn; they only route public turns in Roster mode.".to_owned()
1649        };
1650        let prompt = format!(
1651            "{introduction}{separator}{prompt}{role_separator}{role_block}\n\n{}{handoff_block}",
1652            if effective_can_stop {
1653                format!(
1654                    "You are reviewing another agent. If no meaningful correction is needed,\nend your final response with {STOP_TOKEN}, optionally preceded by an emoji.\nOnly a terminal marker after all reasoning and tool activity requests a stop. A marker followed by more output or activity is non-stopping reasoning. Trailing whitespace is allowed.\nCodeSwarm hides the token and evaluates it only when your turn is complete."
1655                )
1656            } else {
1657                format!(
1658                    "Do not use {STOP_TOKEN} on this turn. Your response must be offered to another agent for review."
1659                )
1660            },
1661            separator = if introduction.is_empty() { "" } else { "\n\n" },
1662        );
1663        let host = self
1664            .hosts
1665            .get_mut(*slot)
1666            .ok_or_else(|| AdapterError::Transport("relay selected missing adapter".into()))?;
1667        let prompt = crate::goal::prompt(self.goal.as_ref(), &prompt);
1668        if let Err(error) = host.send_prompt(prompt.clone()).await {
1669            let limited =
1670                report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1671            if limited {
1672                self.relay.finish(*slot, *direct, false);
1673            }
1674            let _ = self.queue_session_metadata();
1675            if limited {
1676                return Ok(decision);
1677            }
1678            return Err(error);
1679        }
1680        if let Some(sink) = &self.event_sink {
1681            sink(AgentEvent::TurnStarted { slot: *slot });
1682        }
1683        if let Some(introduced) = self.introduced.get_mut(*slot) {
1684            *introduced = true;
1685        }
1686        self.dispatches.push((*slot, prompt));
1687        if !*direct {
1688            self.last_public_dispatch = Some(*slot);
1689        }
1690        let mut response = String::new();
1691        // Only the final contiguous message segment can request a stop or handoff.
1692        // Later reasoning/tool activity invalidates an earlier marker.
1693        let mut stop_segment_start = 0;
1694        let mut emitted_text = 0usize;
1695        let completion_event = loop {
1696            if self.cancel_requested.swap(false, Ordering::AcqRel) {
1697                if let Err(error) = cancel_with_timeout(host).await {
1698                    report_relay_failure(
1699                        &mut self.relay,
1700                        &event_sink,
1701                        *slot,
1702                        true,
1703                        error.to_string(),
1704                    );
1705                    let _ = self.queue_session_metadata();
1706                    return Err(error);
1707                }
1708                return Err(AdapterError::Transport("relay turn cancelled".into()));
1709            }
1710            let update = tokio::select! {
1711                update = host.next_update() => match update {
1712                    Some(Ok(update)) => update,
1713                    Some(Err(error)) => {
1714                        let limited = report_relay_failure(
1715                            &mut self.relay,
1716                            &event_sink,
1717                            *slot,
1718                            true,
1719                            error.to_string(),
1720                        );
1721                        if limited {
1722                            self.relay.finish(*slot, *direct, false);
1723                        }
1724                        let _ = self.queue_session_metadata();
1725                        if limited {
1726                            return Ok(decision);
1727                        }
1728                        return Err(error);
1729                    }
1730                    None => {
1731                        let error = AdapterError::Transport("adapter ended during turn".into());
1732                        let limited = report_relay_failure(
1733                            &mut self.relay,
1734                            &event_sink,
1735                            *slot,
1736                            true,
1737                            error.to_string(),
1738                        );
1739                        if limited {
1740                            self.relay.finish(*slot, *direct, false);
1741                        }
1742                        let _ = self.queue_session_metadata();
1743                        if limited {
1744                            return Ok(decision);
1745                        }
1746                        return Err(error);
1747                    }
1748                },
1749                _ = self.cancel_notify.notified() => {
1750                    if !self.cancel_requested.swap(false, Ordering::AcqRel) {
1751                        continue;
1752                    }
1753                    if let Err(error) = cancel_with_timeout(host).await {
1754                        report_relay_failure(
1755                            &mut self.relay,
1756                            &event_sink,
1757                            *slot,
1758                            true,
1759                            error.to_string(),
1760                        );
1761                        let _ = self.queue_session_metadata();
1762                        return Err(error);
1763                    }
1764                    return Err(AdapterError::Transport("relay turn cancelled".into()));
1765                },
1766                permission = async {
1767                    match permissions.as_mut() {
1768                        Some(receiver) => receiver.recv().await,
1769                        None => std::future::pending().await,
1770                    }
1771                } => {
1772                    let Some(permission) = permission else {
1773                        permissions = None;
1774                        continue;
1775                    };
1776                    if permission.slot != *slot {
1777                        return Err(AdapterError::Transport(
1778                            "permission response targets an inactive relay slot".into(),
1779                        ));
1780                    }
1781                    host.answer_permission(permission.request_id, permission.answer).await?;
1782                    continue;
1783                },
1784            };
1785            match &update.event {
1786                AgentEvent::Text { text, .. } => response.push_str(text),
1787                AgentEvent::Thought { text, .. } | AgentEvent::UserText { text, .. }
1788                    if !text.trim().is_empty() =>
1789                {
1790                    stop_segment_start = response.len();
1791                }
1792                AgentEvent::Tool { .. }
1793                | AgentEvent::Permission { .. }
1794                | AgentEvent::Terminal { .. } => {
1795                    stop_segment_start = response.len();
1796                }
1797                AgentEvent::TurnComplete { .. } => {
1798                    let visible_response = strip_control_tokens(&response);
1799                    let visible_start = emitted_text.min(visible_response.len());
1800                    let visible_start = floor_char_boundary(&visible_response, visible_start);
1801                    if visible_start < visible_response.len()
1802                        && let Some(sink) = &self.event_sink
1803                    {
1804                        sink(AgentEvent::Text {
1805                            slot: *slot,
1806                            text: visible_response[visible_start..].to_owned(),
1807                        });
1808                    }
1809                    self.cancel_requested.store(false, Ordering::Release);
1810                    break update.event.clone();
1811                }
1812                AgentEvent::Failed {
1813                    started, detail, ..
1814                } => {
1815                    let limited = report_relay_failure(
1816                        &mut self.relay,
1817                        &event_sink,
1818                        *slot,
1819                        *started,
1820                        detail.clone(),
1821                    );
1822                    if limited {
1823                        self.relay.finish(*slot, *direct, false);
1824                    }
1825                    let _ = self.queue_session_metadata();
1826                    if limited {
1827                        return Ok(decision);
1828                    }
1829                    return Err(AdapterError::Transport(detail.clone()));
1830                }
1831                _ => {}
1832            }
1833            if let AgentEvent::Text { .. } = &update.event {
1834                // Only a possible split marker needs to wait for another chunk.
1835                let visible_response = strip_control_tokens(&response);
1836                let visible_end = control_token_visible_end(&visible_response);
1837                if emitted_text < visible_end {
1838                    if let Some(sink) = &self.event_sink {
1839                        sink(AgentEvent::Text {
1840                            slot: *slot,
1841                            text: visible_response[emitted_text..visible_end].to_owned(),
1842                        });
1843                    }
1844                    emitted_text = visible_end;
1845                }
1846            } else if let Some(sink) = &self.event_sink {
1847                sink(update.event.clone());
1848            }
1849        };
1850        let requested_stop = response[stop_segment_start..]
1851            .trim_end()
1852            .ends_with(STOP_TOKEN);
1853        let next_slot = requested_next_slot(&response[stop_segment_start..]);
1854        let (response, _) = strip_stop_token(&response);
1855        let response = strip_control_tokens(&response);
1856        let accepted_stop = requested_stop && effective_can_stop;
1857        let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1858        let response = if needs_stop_acknowledgment {
1859            DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1860        } else {
1861            response
1862        };
1863        // A token-only reviewer response is intentionally hidden, but the
1864        // UI still needs the documented visible acknowledgement. Streamed
1865        // text is normally emitted above; this synthetic acknowledgement is
1866        // the one case where there was no visible adapter chunk to forward.
1867        if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1868            sink(AgentEvent::Text {
1869                slot: *slot,
1870                text: response.clone(),
1871            });
1872        }
1873        if let Some(sink) = &self.event_sink {
1874            sink(completion_event);
1875        }
1876        // A provider plan that ran out mid-turn routes future turns around
1877        // the agent instead of back into the exhausted quota.
1878        if is_usage_limit_response(&response) {
1879            let detail = response.clone();
1880            let _ = self.relay.mark_limited(*slot);
1881            // A normal finish: the ring cursor advances past the limited
1882            // agent so the batch continues with a healthy peer.
1883            self.relay.finish(*slot, *direct, false);
1884            self.queue_session_metadata()?;
1885            if let Some(sink) = &self.event_sink {
1886                sink(AgentEvent::UsageLimitReached {
1887                    slot: *slot,
1888                    detail,
1889                });
1890            }
1891            return Ok(decision);
1892        }
1893        if !*direct && !response.is_empty() {
1894            self.relay
1895                .record_public(public_context_speaker(&speaker_name), response);
1896        }
1897        self.relay.mark_context_seen(*slot);
1898        self.relay
1899            .finish_with_handoff(*slot, *direct, accepted_stop, next_slot);
1900        self.queue_session_metadata()?;
1901        Ok(decision)
1902    }
1903}
1904
1905fn report_relay_failure(
1906    relay: &mut Relay,
1907    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1908    slot: RosterSlot,
1909    started: bool,
1910    detail: String,
1911) -> bool {
1912    // A post-start failure is scoped to this automated batch. Keep the
1913    // user's roster selection and provider handle so the next human prompt
1914    // retries it; repeated failures remain for the user to resolve.
1915    if is_usage_limit_response(&detail) {
1916        let _ = relay.mark_limited(slot);
1917        if let Some(sink) = event_sink {
1918            sink(AgentEvent::UsageLimitReached { slot, detail });
1919        }
1920        return true;
1921    }
1922    if started {
1923        let _ = relay.mark_limited(slot);
1924    } else {
1925        let _ = relay.tombstone(slot);
1926    }
1927    if let Some(sink) = event_sink {
1928        sink(AgentEvent::Failed {
1929            slot,
1930            started,
1931            detail,
1932        });
1933    }
1934    started
1935}
1936
1937/// Deterministic in-memory adapter used for contract and relay tests.
1938#[derive(Debug)]
1939pub struct ScriptedAdapter {
1940    slot: RosterSlot,
1941    capabilities: AgentCapabilities,
1942    events: VecDeque<AdapterResult<AgentEvent>>,
1943    prompts: Vec<String>,
1944}
1945
1946impl ScriptedAdapter {
1947    pub fn new(
1948        slot: RosterSlot,
1949        capabilities: AgentCapabilities,
1950        events: impl IntoIterator<Item = AgentEvent>,
1951    ) -> Self {
1952        Self {
1953            slot,
1954            capabilities,
1955            events: events.into_iter().map(Ok).collect(),
1956            prompts: Vec::new(),
1957        }
1958    }
1959
1960    pub fn prompts(&self) -> &[String] {
1961        &self.prompts
1962    }
1963}
1964
1965#[async_trait]
1966impl AgentAdapter for ScriptedAdapter {
1967    fn slot(&self) -> RosterSlot {
1968        self.slot
1969    }
1970
1971    fn capabilities(&self) -> AgentCapabilities {
1972        self.capabilities.clone()
1973    }
1974
1975    async fn start(&mut self) -> AdapterResult<()> {
1976        Ok(())
1977    }
1978
1979    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1980        self.prompts.push(prompt);
1981        Ok(())
1982    }
1983
1984    async fn cancel(&mut self) -> AdapterResult<bool> {
1985        Ok(self.capabilities.supports_cancel)
1986    }
1987
1988    async fn answer_permission(
1989        &mut self,
1990        _request_id: String,
1991        _answer: PermissionAnswer,
1992    ) -> AdapterResult<()> {
1993        if self.capabilities.supports_permissions {
1994            Ok(())
1995        } else {
1996            Err(AdapterError::Unsupported("permission answer"))
1997        }
1998    }
1999
2000    async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
2001        if self.capabilities.supports_modes {
2002            Ok(())
2003        } else {
2004            Err(AdapterError::Unsupported("set_mode"))
2005        }
2006    }
2007
2008    async fn reload(&mut self) -> AdapterResult<()> {
2009        Ok(())
2010    }
2011
2012    async fn stop(&mut self) -> AdapterResult<()> {
2013        Ok(())
2014    }
2015
2016    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2017        self.events.pop_front()
2018    }
2019}
2020
2021/// A direct stream-JSON adapter for Antigravity. It deliberately does not
2022/// pretend to be ACP; it translates its documented events into core events.
2023#[derive(Debug)]
2024pub struct AgyAdapter {
2025    slot: RosterSlot,
2026    cwd: PathBuf,
2027    command: String,
2028    mode: String,
2029    mode_policy: String,
2030    session_id: Option<String>,
2031    child: Option<Child>,
2032    sender: mpsc::Sender<AdapterResult<AgentEvent>>,
2033    receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
2034    /// Antigravity announces its conversation id in the `init` event. Keep it
2035    /// outside the stream task so the next prompt can resume the same native
2036    /// conversation without making the adapter reader/UI share a borrow.
2037    announced_session: Arc<Mutex<Option<String>>>,
2038    cancel_requested: Arc<AtomicBool>,
2039}
2040
2041impl AgyAdapter {
2042    pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
2043        let (sender, receiver) = mpsc::channel(256);
2044        Self {
2045            slot,
2046            cwd,
2047            command: command.into(),
2048            mode: "default".into(),
2049            mode_policy: "agy:full-access".into(),
2050            session_id: None,
2051            child: None,
2052            sender,
2053            receiver,
2054            announced_session: Arc::new(Mutex::new(None)),
2055            cancel_requested: Arc::new(AtomicBool::new(false)),
2056        }
2057    }
2058
2059    pub fn with_session_id(
2060        slot: RosterSlot,
2061        cwd: PathBuf,
2062        command: impl Into<String>,
2063        session_id: impl Into<String>,
2064    ) -> Self {
2065        let mut adapter = Self::new(slot, cwd, command);
2066        adapter.session_id = Some(session_id.into());
2067        adapter
2068    }
2069
2070    fn modes() -> Vec<Mode> {
2071        vec![
2072            Mode {
2073                id: "agy:full-access".into(),
2074                label: "Auto pilot".into(),
2075            },
2076            Mode {
2077                id: "agy:manual".into(),
2078                label: "Manual".into(),
2079            },
2080            Mode {
2081                id: "accept-edits".into(),
2082                label: "Accept Edits".into(),
2083            },
2084            Mode {
2085                id: "plan".into(),
2086                label: "Plan".into(),
2087            },
2088        ]
2089    }
2090
2091    async fn emit(&self, event: AdapterResult<AgentEvent>) {
2092        let _ = self.sender.send(event).await;
2093    }
2094}
2095
2096#[async_trait]
2097impl AgentAdapter for AgyAdapter {
2098    fn slot(&self) -> RosterSlot {
2099        self.slot
2100    }
2101
2102    fn session_id(&self) -> Option<String> {
2103        self.session_id.clone()
2104    }
2105
2106    fn protocol(&self) -> &'static str {
2107        "native"
2108    }
2109
2110    fn capabilities(&self) -> AgentCapabilities {
2111        AgentCapabilities {
2112            supports_cancel: true,
2113            supports_modes: true,
2114            supports_permissions: false,
2115            supports_terminals: true,
2116            supports_session_load: true,
2117            supports_models: false,
2118        }
2119    }
2120
2121    async fn start(&mut self) -> AdapterResult<()> {
2122        self.cancel_requested.store(false, Ordering::Release);
2123        self.emit(Ok(AgentEvent::Ready {
2124            slot: self.slot,
2125            capabilities: self.capabilities(),
2126        }))
2127        .await;
2128        self.emit(Ok(AgentEvent::ModesReplaced {
2129            slot: self.slot,
2130            modes: Self::modes(),
2131            current_mode: Some(self.mode_policy.clone()),
2132        }))
2133        .await;
2134        Ok(())
2135    }
2136
2137    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
2138        if self.child.is_some() {
2139            return Err(AdapterError::Transport(
2140                "agent is already handling a turn".into(),
2141            ));
2142        }
2143        if self.session_id.is_none()
2144            && let Ok(session) = self.announced_session.lock()
2145        {
2146            self.session_id = session.clone();
2147        }
2148        self.cancel_requested.store(false, Ordering::Release);
2149        let (program, args) = parse_command_line(&self.command)
2150            .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
2151        let mut command = Command::new(program);
2152        isolate_process_group(&mut command);
2153        command
2154            .args(args)
2155            .arg("--print")
2156            .arg(prompt)
2157            .arg("--print-timeout")
2158            // Allow day-long work; cancellation remains under user control.
2159            .arg("1440m")
2160            .arg("--output-format")
2161            .arg("stream-json")
2162            .current_dir(&self.cwd)
2163            .env("CODESWARM_CWD", &self.cwd)
2164            .stdout(Stdio::piped())
2165            .stderr(Stdio::piped());
2166        if let Some(session_id) = &self.session_id {
2167            command.arg("--conversation").arg(session_id);
2168        }
2169        if self.mode != "default" {
2170            command.arg("--mode").arg(&self.mode);
2171        }
2172        let mut child = command
2173            .spawn()
2174            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2175        let stdout = match child.stdout.take() {
2176            Some(stdout) => stdout,
2177            None => {
2178                let _ = terminate_child(&mut child).await;
2179                return Err(AdapterError::Transport("agent has no stdout".into()));
2180            }
2181        };
2182        let stderr = match child.stderr.take() {
2183            Some(stderr) => stderr,
2184            None => {
2185                let _ = terminate_child(&mut child).await;
2186                return Err(AdapterError::Transport("agent has no stderr".into()));
2187            }
2188        };
2189        let sender = self.sender.clone();
2190        let slot = self.slot;
2191        let announced_session = Arc::clone(&self.announced_session);
2192        let cancel_requested = Arc::clone(&self.cancel_requested);
2193        tokio::spawn(async move {
2194            let stderr_task = tokio::spawn(async move {
2195                const MAX_STDERR: usize = 32 * 1024;
2196                let mut stderr = BufReader::new(stderr);
2197                let mut bytes = Vec::new();
2198                let mut chunk = [0_u8; 4096];
2199                while let Ok(count) = stderr.read(&mut chunk).await {
2200                    if count == 0 {
2201                        break;
2202                    }
2203                    bytes.extend_from_slice(&chunk[..count]);
2204                    if bytes.len() > MAX_STDERR {
2205                        let keep_from = bytes.len() - MAX_STDERR;
2206                        bytes.drain(..keep_from);
2207                    }
2208                }
2209                String::from_utf8_lossy(&bytes).trim().to_owned()
2210            });
2211            let mut lines = BufReader::new(stdout).lines();
2212            let mut result: Option<Value> = None;
2213            let mut streamed_response = false;
2214            while let Ok(Some(line)) = lines.next_line().await {
2215                let value = match serde_json::from_str::<Value>(&line) {
2216                    Ok(value) => value,
2217                    Err(_) => {
2218                        // Native stream-json can contain diagnostic junk on
2219                        // stdout. Match the Python adapter's tolerant stream
2220                        // behavior and wait for the final result instead of
2221                        // turning one malformed line into a dead turn.
2222                        continue;
2223                    }
2224                };
2225                if value.get("event").and_then(Value::as_str) == Some("init")
2226                    && let Some(session_id) = value
2227                        .get("conversation_id")
2228                        .or_else(|| value.get("conversationId"))
2229                        .and_then(Value::as_str)
2230                        .filter(|id| !id.is_empty())
2231                    && let Ok(mut announced) = announced_session.lock()
2232                {
2233                    *announced = Some(session_id.to_owned());
2234                }
2235                if value.get("event").and_then(Value::as_str) == Some("result") {
2236                    result = value.get("result").cloned();
2237                }
2238                match parse_agy_value(slot, &value) {
2239                    Ok(Some(event)) => {
2240                        if matches!(event, AgentEvent::Text { .. }) {
2241                            streamed_response = true;
2242                        }
2243                        if sender.send(Ok(event)).await.is_err() {
2244                            break;
2245                        }
2246                    }
2247                    Ok(None) => {}
2248                    Err(error) => {
2249                        let _ = sender.send(Err(error)).await;
2250                    }
2251                }
2252            }
2253            let stderr = stderr_task.await.ok().unwrap_or_default();
2254            let succeeded = cancel_requested.load(Ordering::Acquire)
2255                || result
2256                    .as_ref()
2257                    .and_then(|result| result.get("status"))
2258                    .and_then(Value::as_str)
2259                    == Some("SUCCESS");
2260            if succeeded {
2261                // Some native stream-json wrappers emit only lifecycle events
2262                // and put the complete answer in the final result object. The
2263                // Python adapter surfaced that answer; do the same, while
2264                // avoiding duplication when token chunks were already sent.
2265                if !streamed_response
2266                    && let Some(response) = result
2267                        .as_ref()
2268                        .and_then(|result| result.get("response"))
2269                        .and_then(Value::as_str)
2270                        .filter(|response| !response.is_empty())
2271                {
2272                    let _ = sender
2273                        .send(Ok(AgentEvent::Text {
2274                            slot,
2275                            text: response.to_owned(),
2276                        }))
2277                        .await;
2278                }
2279                let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2280            } else {
2281                let detail = result
2282                    .as_ref()
2283                    .and_then(|result| result.get("error"))
2284                    .and_then(Value::as_str)
2285                    .filter(|detail| !detail.is_empty())
2286                    .map(str::to_owned)
2287                    .or_else(|| (!stderr.is_empty()).then_some(stderr))
2288                    .unwrap_or_else(|| "native stream ended before a successful result".into());
2289                let _ = sender
2290                    .send(Ok(AgentEvent::Failed {
2291                        slot,
2292                        started: true,
2293                        detail,
2294                    }))
2295                    .await;
2296            }
2297        });
2298        self.child = Some(child);
2299        Ok(())
2300    }
2301
2302    async fn cancel(&mut self) -> AdapterResult<bool> {
2303        self.cancel_requested.store(true, Ordering::Release);
2304        let Some(mut child) = self.child.take() else {
2305            return Ok(false);
2306        };
2307        // `Child` does not reap itself when dropped. Awaiting wait after the
2308        // kill keeps repeated prompts from accumulating zombies, especially
2309        // when cancellation happens before the stream reader observes EOF.
2310        terminate_child(&mut child).await?;
2311        let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2312            while let Some(event) = self.receiver.recv().await {
2313                if matches!(
2314                    event,
2315                    Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2316                ) {
2317                    break;
2318                }
2319            }
2320        })
2321        .await;
2322        Ok(true)
2323    }
2324
2325    async fn answer_permission(
2326        &mut self,
2327        _request_id: String,
2328        _answer: PermissionAnswer,
2329    ) -> AdapterResult<()> {
2330        Err(AdapterError::Unsupported("permission answer"))
2331    }
2332
2333    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2334        let (mode, mode_policy) = match mode.as_str() {
2335            "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2336                ("default".to_owned(), "agy:full-access".to_owned())
2337            }
2338            "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2339            "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2340                ("accept-edits".to_owned(), "accept-edits".to_owned())
2341            }
2342            "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2343                ("default".to_owned(), "agy:manual".to_owned())
2344            }
2345            "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2346            "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2347            _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2348        };
2349        self.mode = mode;
2350        self.mode_policy = mode_policy.clone();
2351        self.emit(Ok(AgentEvent::ModesReplaced {
2352            slot: self.slot,
2353            modes: Self::modes(),
2354            current_mode: Some(mode_policy),
2355        }))
2356        .await;
2357        Ok(())
2358    }
2359
2360    async fn reload(&mut self) -> AdapterResult<()> {
2361        self.stop().await?;
2362        self.start().await
2363    }
2364
2365    async fn stop(&mut self) -> AdapterResult<()> {
2366        let _ = self.cancel().await?;
2367        Ok(())
2368    }
2369
2370    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2371        let event = self.receiver.recv().await;
2372        if matches!(
2373            event.as_ref(),
2374            Some(Ok(
2375                AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. }
2376            ))
2377        ) {
2378            if self.session_id.is_none()
2379                && let Ok(session) = self.announced_session.lock()
2380            {
2381                self.session_id = session.clone();
2382            }
2383            // A failed turn is terminal for this process too. Reap it here
2384            // so the next user prompt starts a fresh native process instead
2385            // of trying to write to a dead child.
2386            if let Some(mut child) = self.child.take() {
2387                let _ = child.wait().await;
2388            }
2389        }
2390        event
2391    }
2392}
2393
2394#[cfg(test)]
2395fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2396    let value: Value =
2397        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2398    parse_agy_value(slot, &value)
2399}
2400
2401fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2402    let event = value.get("event").and_then(Value::as_str);
2403    if let Some(terminal) = parse_terminal_event(value, event) {
2404        return Ok(Some(AgentEvent::Terminal {
2405            slot,
2406            event: terminal,
2407        }));
2408    }
2409    match event {
2410        Some("step_update") => {
2411            let Some(update) = value.get("step_update") else {
2412                return Ok(None);
2413            };
2414            let is_response = update
2415                .get("step_type")
2416                .and_then(Value::as_str)
2417                .is_some_and(|kind| kind == "agent_response");
2418            let text = update.get("text_delta").and_then(Value::as_str);
2419            let response = is_response
2420                .then(|| text.map(str::to_owned))
2421                .flatten()
2422                .filter(|text| !text.is_empty())
2423                .map(|text| AgentEvent::Text { slot, text });
2424            Ok(response.or_else(|| parse_agy_tool(slot, value)))
2425        }
2426        _ => Ok(None),
2427    }
2428}
2429
2430fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2431    let update = value.get("step_update")?;
2432    if update.get("step_type")?.as_str()? != "tool" {
2433        return None;
2434    }
2435    let step_index = update.get("step_index")?.as_i64()?;
2436    let title = update
2437        .get("tool_name")
2438        .and_then(Value::as_str)
2439        .unwrap_or("Tool call")
2440        .replace('_', " ");
2441    let status = match update.get("state").and_then(Value::as_str) {
2442        Some("DONE") => ToolStatus::Completed,
2443        Some("FAILED") => ToolStatus::Failed,
2444        Some("ACTIVE") => ToolStatus::Running,
2445        _ => ToolStatus::Pending,
2446    };
2447    let detail = update
2448        .get("tool_info")
2449        .and_then(|info| info.get("output"))
2450        .and_then(Value::as_str)
2451        .map(str::to_owned);
2452    Some(AgentEvent::Tool {
2453        slot,
2454        update: ToolUpdate {
2455            id: format!("agy-tool-{step_index}"),
2456            title,
2457            status,
2458            detail,
2459        },
2460    })
2461}
2462
2463/// Stdio ACP transport. Protocol-specific response handling belongs here,
2464/// keeping JSON-RPC framing outside the core and terminal renderer.
2465#[derive(Debug)]
2466pub struct AcpAdapter {
2467    slot: RosterSlot,
2468    program: String,
2469    args: Vec<String>,
2470    cwd: PathBuf,
2471    child: Option<Child>,
2472    reader: Option<BufReader<ChildStdout>>,
2473    capabilities: AgentCapabilities,
2474    modes: Vec<Mode>,
2475    models: Vec<Mode>,
2476    model_config_id: Option<String>,
2477    session_id: Option<String>,
2478    next_request_id: u64,
2479    prompt_request_id: Option<u64>,
2480    queued_events: VecDeque<AdapterResult<AgentEvent>>,
2481    tool_updates: BTreeMap<String, ToolUpdate>,
2482    stderr_task: Option<tokio::task::JoinHandle<String>>,
2483    terminals: BTreeMap<String, TerminalProcess>,
2484    next_terminal_id: u64,
2485}
2486
2487impl AcpAdapter {
2488    pub fn new(
2489        slot: RosterSlot,
2490        cwd: PathBuf,
2491        program: impl Into<String>,
2492        args: Vec<String>,
2493    ) -> Self {
2494        Self {
2495            slot,
2496            program: program.into(),
2497            args,
2498            cwd,
2499            child: None,
2500            reader: None,
2501            capabilities: AgentCapabilities::default(),
2502            modes: Vec::new(),
2503            models: Vec::new(),
2504            model_config_id: None,
2505            session_id: None,
2506            next_request_id: 1,
2507            prompt_request_id: None,
2508            queued_events: VecDeque::new(),
2509            tool_updates: BTreeMap::new(),
2510            stderr_task: None,
2511            terminals: BTreeMap::new(),
2512            next_terminal_id: 1,
2513        }
2514    }
2515
2516    pub fn with_session_id(
2517        slot: RosterSlot,
2518        cwd: PathBuf,
2519        program: impl Into<String>,
2520        args: Vec<String>,
2521        session_id: impl Into<String>,
2522    ) -> Self {
2523        let mut adapter = Self::new(slot, cwd, program, args);
2524        adapter.session_id = Some(session_id.into());
2525        adapter
2526    }
2527
2528    async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2529        self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2530            .await
2531    }
2532
2533    async fn request_with_timeout(
2534        &mut self,
2535        method: &str,
2536        params: Value,
2537        deadline: std::time::Duration,
2538    ) -> AdapterResult<Value> {
2539        tokio::time::timeout(deadline, self.request_inner(method, params))
2540            .await
2541            .map_err(|_| {
2542                AdapterError::Transport(format!(
2543                    "ACP {method} timed out; reload the agent to retry"
2544                ))
2545            })?
2546    }
2547
2548    async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2549        let request_id = self.next_request_id;
2550        self.next_request_id += 1;
2551        self.write_json(serde_json::json!({
2552            "jsonrpc": "2.0",
2553            "id": request_id,
2554            "method": method,
2555            "params": params,
2556        }))
2557        .await?;
2558        loop {
2559            let line = self.read_line().await?;
2560            let value: Value = match serde_json::from_str(&line) {
2561                Ok(value) => value,
2562                Err(_) => {
2563                    // Keep the transport alive when a peer writes a stray
2564                    // diagnostic line. This is common with CLI wrappers and
2565                    // matches the baseline client's tolerant stream loop.
2566                    continue;
2567                }
2568            };
2569            if self.reject_empty_permission_request(&value).await? {
2570                continue;
2571            }
2572            if self.handle_client_request(&value).await? {
2573                continue;
2574            }
2575            if value
2576                .get("id")
2577                .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2578            {
2579                if let Some(error) = value.get("error") {
2580                    return Err(AdapterError::Protocol(error.to_string()));
2581                }
2582                return value
2583                    .get("result")
2584                    .cloned()
2585                    .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2586            }
2587            if let Some(event) = parse_acp_value(self.slot, &value, &mut self.tool_updates)? {
2588                let event = if method == "session/load" {
2589                    restored_history_event(event)
2590                } else {
2591                    event
2592                };
2593                self.queued_events.push_back(Ok(event));
2594            }
2595        }
2596    }
2597
2598    async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2599        let child = self
2600            .child
2601            .as_mut()
2602            .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2603        let stdin = child
2604            .stdin
2605            .as_mut()
2606            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2607        stdin
2608            .write_all(value.to_string().as_bytes())
2609            .await
2610            .map_err(|error| AdapterError::Transport(error.to_string()))?;
2611        stdin
2612            .write_all(b"\n")
2613            .await
2614            .map_err(|error| AdapterError::Transport(error.to_string()))
2615    }
2616
2617    /// Tear down a broken ACP stream while retaining the provider session
2618    /// handle. The coordinator reloads the slot before its next prompt.
2619    async fn reset_transport(&mut self) {
2620        let terminals = std::mem::take(&mut self.terminals);
2621        for terminal in terminals.values() {
2622            terminal.stop().await;
2623        }
2624        self.queued_events.clear();
2625        self.tool_updates.clear();
2626        if let Some(mut child) = self.child.take() {
2627            let _ = terminate_child(&mut child).await;
2628        }
2629        self.reader = None;
2630        if let Some(task) = self.stderr_task.take() {
2631            task.abort();
2632        }
2633        self.prompt_request_id = None;
2634    }
2635
2636    /// ACP permission requests are JSON-RPC requests, not fire-and-forget
2637    /// notifications. An empty option list is invalid and must be answered
2638    /// with an error so the peer does not wait forever for a decision. This is
2639    /// the same validation performed by the Python ACP server.
2640    async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2641        if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2642            || value.get("id").is_none()
2643        {
2644            return Ok(false);
2645        }
2646        let valid = value
2647            .get("params")
2648            .and_then(|params| params.get("options"))
2649            .and_then(Value::as_array)
2650            .is_some_and(|options| !options.is_empty());
2651        if valid {
2652            return Ok(false);
2653        }
2654        self.write_json(serde_json::json!({
2655            "jsonrpc": "2.0",
2656            "id": value.get("id").cloned().unwrap_or(Value::Null),
2657            "error": {
2658                "code": -32602,
2659                "message": "Permission request requires at least one option",
2660            },
2661        }))
2662        .await?;
2663        Ok(true)
2664    }
2665
2666    fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2667        let root = self
2668            .cwd
2669            .canonicalize()
2670            .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2671        let requested = Path::new(path);
2672        let candidate = if requested.is_absolute() {
2673            requested.to_path_buf()
2674        } else {
2675            root.join(requested)
2676        };
2677        let resolved = if !candidate.exists() {
2678            let parent = candidate
2679                .parent()
2680                .ok_or_else(|| "file path has no parent".to_owned())?
2681                .canonicalize()
2682                .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2683            parent.join(
2684                candidate
2685                    .file_name()
2686                    .ok_or_else(|| "file path has no filename".to_owned())?,
2687            )
2688        } else {
2689            candidate
2690                .canonicalize()
2691                .map_err(|error| format!("unable to resolve file path: {error}"))?
2692        };
2693        if !resolved.starts_with(&root) {
2694            return Err("file path is outside the project".into());
2695        }
2696        Ok(resolved)
2697    }
2698
2699    fn read_workspace_text(
2700        &self,
2701        path: &str,
2702        line: Option<i64>,
2703        limit: Option<i64>,
2704    ) -> Result<String, String> {
2705        if line.is_some_and(|line| line < 1) {
2706            return Err("line must be positive".into());
2707        }
2708        if limit.is_some_and(|limit| limit < 0) {
2709            return Err("limit must not be negative".into());
2710        }
2711        let path = self.workspace_path(path)?;
2712        let mut bytes = Vec::new();
2713        let mut source = match std::fs::File::open(path) {
2714            Ok(source) => source,
2715            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2716            Err(error) => return Err(error.to_string()),
2717        };
2718        source
2719            .by_ref()
2720            .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2721            .read_to_end(&mut bytes)
2722            .map_err(|error| error.to_string())?;
2723        bytes.truncate(MAX_FILE_READ_BYTES);
2724        let text = String::from_utf8_lossy(&bytes);
2725        if line.is_none() && limit.is_none() {
2726            return Ok(text.into_owned());
2727        }
2728        let start = line.map_or(0, |line| line as usize - 1);
2729        let limit = limit.unwrap_or(i64::MAX) as usize;
2730        let selected = text
2731            .split_inclusive('\n')
2732            .skip(start)
2733            .take(limit)
2734            .collect::<String>();
2735        if line.is_some() {
2736            Ok(selected.trim_end_matches('\n').to_owned())
2737        } else {
2738            Ok(selected)
2739        }
2740    }
2741
2742    fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2743        let path = params
2744            .get("path")
2745            .and_then(Value::as_str)
2746            .filter(|path| !path.is_empty())
2747            .ok_or("path must be a non-empty string")?;
2748        let content = params
2749            .get("content")
2750            .and_then(Value::as_str)
2751            .ok_or("content must be a string")?;
2752        let path = self.workspace_path(path)?;
2753        std::fs::write(path, content).map_err(|error| error.to_string())
2754    }
2755
2756    async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2757        let command = params
2758            .get("command")
2759            .and_then(Value::as_str)
2760            .filter(|command| !command.trim().is_empty())
2761            .ok_or_else(|| "terminal command is required".to_owned())?;
2762        let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2763        let cwd = self.workspace_path(cwd)?;
2764        if !cwd.is_dir() {
2765            return Err("terminal cwd is not a directory".into());
2766        }
2767        let mut process = Command::new(command);
2768        isolate_process_group(&mut process);
2769        if let Some(args) = params.get("args").and_then(Value::as_array) {
2770            process.args(args.iter().filter_map(Value::as_str));
2771        }
2772        process
2773            .current_dir(&cwd)
2774            .stdin(Stdio::null())
2775            .stdout(Stdio::piped())
2776            .stderr(Stdio::piped());
2777        if let Some(env) = params.get("env") {
2778            if let Some(entries) = env.as_array() {
2779                for entry in entries {
2780                    if let (Some(name), Some(value)) = (
2781                        entry.get("name").and_then(Value::as_str),
2782                        entry.get("value").and_then(Value::as_str),
2783                    ) {
2784                        process.env(name, value);
2785                    }
2786                }
2787            } else if let Some(entries) = env.as_object() {
2788                for (name, value) in entries {
2789                    if let Some(value) = value.as_str() {
2790                        process.env(name, value);
2791                    }
2792                }
2793            }
2794        }
2795        let mut child = process.spawn().map_err(|error| error.to_string())?;
2796        let stdout = child.stdout.take();
2797        let stderr = child.stderr.take();
2798        let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2799            let _ = terminate_child(&mut child).await;
2800            return Err("terminal has no output pipes".into());
2801        };
2802        let output = Arc::new(Mutex::new(Vec::new()));
2803        let truncated = Arc::new(AtomicBool::new(false));
2804        let output_readers = Arc::new(AtomicUsize::new(2));
2805        let output_limit = params
2806            .get("outputByteLimit")
2807            .and_then(Value::as_u64)
2808            .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2809                usize::try_from(limit)
2810                    .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2811                    .min(MAX_TERMINAL_OUTPUT_BYTES)
2812            });
2813        tokio::spawn(drain_terminal_output(
2814            stdout,
2815            Arc::clone(&output),
2816            Arc::clone(&truncated),
2817            Arc::clone(&output_readers),
2818            output_limit,
2819        ));
2820        tokio::spawn(drain_terminal_output(
2821            stderr,
2822            Arc::clone(&output),
2823            Arc::clone(&truncated),
2824            Arc::clone(&output_readers),
2825            output_limit,
2826        ));
2827        let id = format!("terminal-{}", self.next_terminal_id);
2828        self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2829        let state = TerminalProcess {
2830            child: Arc::new(AsyncMutex::new(Some(child))),
2831            output,
2832            truncated,
2833            output_readers,
2834        };
2835        self.terminals.insert(id.clone(), state);
2836        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2837            slot: self.slot,
2838            event: TerminalEvent::Created {
2839                id: id.clone(),
2840                command: std::iter::once(command)
2841                    .chain(
2842                        params
2843                            .get("args")
2844                            .and_then(Value::as_array)
2845                            .into_iter()
2846                            .flatten()
2847                            .filter_map(Value::as_str),
2848                    )
2849                    .collect::<Vec<_>>()
2850                    .join(" "),
2851            },
2852        }));
2853        Ok(serde_json::json!({"terminalId": id}))
2854    }
2855
2856    async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2857        let terminal = self
2858            .terminals
2859            .get(id)
2860            .ok_or_else(|| "terminal not found".to_owned())?
2861            .clone();
2862        let output = terminal
2863            .output
2864            .lock()
2865            .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2866            .unwrap_or_default();
2867        let exit_code = terminal.exit_code().await;
2868        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2869            slot: self.slot,
2870            event: TerminalEvent::Output {
2871                id: id.to_owned(),
2872                text: output.clone(),
2873            },
2874        }));
2875        let mut response = serde_json::json!({
2876            "output": output,
2877            "truncated": terminal.truncated.load(Ordering::Acquire),
2878        });
2879        if let Some(code) = exit_code {
2880            response["exitStatus"] = serde_json::json!({"exitCode": code});
2881        }
2882        Ok(response)
2883    }
2884
2885    async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2886        let terminal = self
2887            .terminals
2888            .get(id)
2889            .ok_or_else(|| "terminal not found".to_owned())?
2890            .clone();
2891        let exit_code = terminal.wait().await;
2892        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2893            slot: self.slot,
2894            event: TerminalEvent::Exited {
2895                id: id.to_owned(),
2896                code: exit_code.unwrap_or(-1),
2897            },
2898        }));
2899        Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2900    }
2901
2902    /// Handle requests initiated by an ACP agent against the client. File
2903    /// access is mediated through the configured workspace root; unsupported
2904    /// requests receive a JSON-RPC error instead of hanging the agent.
2905    async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2906        let Some(method) = value.get("method").and_then(Value::as_str) else {
2907            return Ok(false);
2908        };
2909        let Some(id) = value.get("id").cloned() else {
2910            return Ok(false);
2911        };
2912        // Permission requests are consumed by the normalized event parser and
2913        // answered later through the focused UI action, not here.
2914        if method == "session/request_permission" {
2915            return Ok(false);
2916        }
2917        let params = value.get("params").cloned().unwrap_or(Value::Null);
2918        let response = match method {
2919            "fs/read_text_file" => {
2920                let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2921                let line = params.get("line").and_then(Value::as_i64);
2922                let limit = params.get("limit").and_then(Value::as_i64);
2923                match self.read_workspace_text(path, line, limit) {
2924                    Ok(content) => serde_json::json!({
2925                        "jsonrpc": "2.0",
2926                        "id": id,
2927                        "result": {"content": content},
2928                    }),
2929                    Err(message) => serde_json::json!({
2930                        "jsonrpc": "2.0",
2931                        "id": id,
2932                        "error": {"code": -32602, "message": message},
2933                    }),
2934                }
2935            }
2936            "fs/write_text_file" => {
2937                let result = self.write_workspace_text(&params);
2938                match result {
2939                    Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2940                    Err(message) => serde_json::json!({
2941                        "jsonrpc": "2.0",
2942                        "id": id,
2943                        "error": {"code": -32602, "message": message},
2944                    }),
2945                }
2946            }
2947            "terminal/create" => match self.terminal_create(&params).await {
2948                Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2949                Err(message) => serde_json::json!({
2950                    "jsonrpc": "2.0",
2951                    "id": id,
2952                    "error": {"code": -32602, "message": message},
2953                }),
2954            },
2955            "terminal/output" => {
2956                let terminal_id = params
2957                    .get("terminalId")
2958                    .and_then(Value::as_str)
2959                    .unwrap_or("");
2960                match self.terminal_output(terminal_id).await {
2961                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2962                    Err(message) => serde_json::json!({
2963                        "jsonrpc": "2.0",
2964                        "id": id,
2965                        "error": {"code": -32602, "message": message},
2966                    }),
2967                }
2968            }
2969            "terminal/wait_for_exit" => {
2970                let terminal_id = params
2971                    .get("terminalId")
2972                    .and_then(Value::as_str)
2973                    .unwrap_or("");
2974                match self.terminal_wait(terminal_id).await {
2975                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2976                    Err(message) => serde_json::json!({
2977                        "jsonrpc": "2.0",
2978                        "id": id,
2979                        "error": {"code": -32602, "message": message},
2980                    }),
2981                }
2982            }
2983            "terminal/kill" => {
2984                let terminal_id = params
2985                    .get("terminalId")
2986                    .and_then(Value::as_str)
2987                    .unwrap_or("");
2988                if let Some(terminal) = self.terminals.get(terminal_id) {
2989                    terminal.kill().await;
2990                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2991                } else {
2992                    serde_json::json!({
2993                        "jsonrpc": "2.0",
2994                        "id": id,
2995                        "error": {"code": -32602, "message": "terminal not found"},
2996                    })
2997                }
2998            }
2999            "terminal/release" => {
3000                let terminal_id = params
3001                    .get("terminalId")
3002                    .and_then(Value::as_str)
3003                    .unwrap_or("");
3004                if let Some(terminal) = self.terminals.remove(terminal_id) {
3005                    terminal.stop().await;
3006                    self.queued_events.push_back(Ok(AgentEvent::Terminal {
3007                        slot: self.slot,
3008                        event: TerminalEvent::Released {
3009                            id: terminal_id.to_owned(),
3010                        },
3011                    }));
3012                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3013                } else {
3014                    serde_json::json!({
3015                        "jsonrpc": "2.0",
3016                        "id": id,
3017                        "error": {"code": -32602, "message": "terminal not found"},
3018                    })
3019                }
3020            }
3021            _ => serde_json::json!({
3022                "jsonrpc": "2.0",
3023                "id": id,
3024                "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
3025            }),
3026        };
3027        self.write_json(response).await?;
3028        Ok(true)
3029    }
3030
3031    async fn read_line(&mut self) -> AdapterResult<String> {
3032        let reader = self
3033            .reader
3034            .as_mut()
3035            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
3036        read_bounded_line(reader).await
3037    }
3038
3039    async fn start(&mut self) -> AdapterResult<()> {
3040        self.modes.clear();
3041        // Starting an adapter twice must never orphan the first transport.
3042        if self.child.is_some() {
3043            self.stop().await?;
3044        }
3045        let mut command = Command::new(&self.program);
3046        isolate_process_group(&mut command);
3047        command
3048            .args(&self.args)
3049            .current_dir(&self.cwd)
3050            .stdin(Stdio::piped())
3051            .stdout(Stdio::piped())
3052            .stderr(Stdio::piped())
3053            .env("CODESWARM_CWD", &self.cwd);
3054        if self.program.to_ascii_lowercase().contains("gemini")
3055            || self
3056                .args
3057                .iter()
3058                .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
3059        {
3060            command.env("GEMINI_TELEMETRY_ENABLED", "false");
3061        }
3062        let mut child = command
3063            .spawn()
3064            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
3065        let stdout = match child.stdout.take() {
3066            Some(stdout) => stdout,
3067            None => {
3068                let _ = terminate_child(&mut child).await;
3069                return Err(AdapterError::Transport("ACP agent has no stdout".into()));
3070            }
3071        };
3072        let stderr = match child.stderr.take() {
3073            Some(stderr) => stderr,
3074            None => {
3075                let _ = terminate_child(&mut child).await;
3076                return Err(AdapterError::Transport("ACP agent has no stderr".into()));
3077            }
3078        };
3079        self.child = Some(child);
3080        self.reader = Some(BufReader::new(stdout));
3081        self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
3082
3083        let initialize = match self
3084            .request(
3085                "initialize",
3086                serde_json::json!({
3087                    "protocolVersion": 1,
3088                    "clientCapabilities": {
3089                        "fs": {"readTextFile": true, "writeTextFile": true},
3090                        "terminal": true,
3091                    },
3092                    "clientInfo": {
3093                        "name": "CodeSwarm",
3094                        "title": "CodeSwarm",
3095                        "version": env!("CARGO_PKG_VERSION"),
3096                    },
3097                }),
3098            )
3099            .await
3100        {
3101            Ok(value) => value,
3102            Err(error) => {
3103                let _ = self.stop().await;
3104                return Err(error);
3105            }
3106        };
3107        let agent_capabilities = initialize
3108            .get("agentCapabilities")
3109            .cloned()
3110            .unwrap_or(Value::Null);
3111        self.capabilities = AgentCapabilities {
3112            supports_cancel: true,
3113            supports_modes: true,
3114            supports_permissions: true,
3115            supports_terminals: true,
3116            supports_session_load: agent_capabilities
3117                .get("loadSession")
3118                .and_then(Value::as_bool)
3119                .unwrap_or(false),
3120            supports_models: false,
3121        };
3122        let session = if let Some(session_id) = self.session_id.clone() {
3123            if !self.capabilities.supports_session_load {
3124                let _ = self.stop().await;
3125                return Err(AdapterError::Unsupported("session/load"));
3126            }
3127            match self
3128                .request(
3129                    "session/load",
3130                    serde_json::json!({
3131                        "cwd": self.cwd,
3132                        "mcpServers": [],
3133                        "sessionId": session_id,
3134                    }),
3135                )
3136                .await
3137            {
3138                Ok(value) => value,
3139                Err(error) => {
3140                    let _ = self.stop().await;
3141                    return Err(error);
3142                }
3143            }
3144        } else {
3145            let session = match self
3146                .request(
3147                    "session/new",
3148                    serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
3149                )
3150                .await
3151            {
3152                Ok(value) => value,
3153                Err(error) => {
3154                    let _ = self.stop().await;
3155                    return Err(error);
3156                }
3157            };
3158            self.session_id = session
3159                .get("sessionId")
3160                .and_then(Value::as_str)
3161                .map(str::to_owned);
3162            if self.session_id.is_none() {
3163                let _ = self.stop().await;
3164                return Err(AdapterError::Protocol(
3165                    "session/new returned no sessionId".into(),
3166                ));
3167            }
3168            session
3169        };
3170        self.capabilities.supports_modes = false;
3171        if let Some(modes) = session.get("modes") {
3172            let available = modes
3173                .get("availableModes")
3174                .and_then(Value::as_array)
3175                .map(|modes| {
3176                    modes
3177                        .iter()
3178                        .filter_map(|mode| {
3179                            Some(Mode {
3180                                id: mode.get("id")?.as_str()?.to_owned(),
3181                                label: mode.get("name")?.as_str()?.to_owned(),
3182                            })
3183                        })
3184                        .collect::<Vec<_>>()
3185                })
3186                .unwrap_or_default();
3187            self.modes = available.clone();
3188            self.capabilities.supports_modes = !available.is_empty();
3189            self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
3190                slot: self.slot,
3191                modes: available,
3192                current_mode: modes
3193                    .get("currentModeId")
3194                    .and_then(Value::as_str)
3195                    .map(str::to_owned),
3196            }));
3197        }
3198        self.models.clear();
3199        self.model_config_id = None;
3200        let current_model =
3201            parse_model_config(&session).and_then(|(config_id, models, current)| {
3202                self.model_config_id = Some(config_id);
3203                self.models = models;
3204                current
3205            });
3206        self.capabilities.supports_models =
3207            self.model_config_id.is_some() && !self.models.is_empty();
3208        if let Some(config_id) = self.model_config_id.clone()
3209            && !self.models.is_empty()
3210        {
3211            self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
3212                slot: self.slot,
3213                config_id,
3214                models: self.models.clone(),
3215                current_model,
3216            }));
3217        }
3218        self.queued_events.push_back(Ok(AgentEvent::Ready {
3219            slot: self.slot,
3220            capabilities: self.capabilities(),
3221        }));
3222        Ok(())
3223    }
3224}
3225
3226fn prompt_resource_paths(prompt: &str) -> Vec<String> {
3227    let characters = prompt.chars().collect::<Vec<_>>();
3228    let mut paths = Vec::new();
3229    let mut index = 0;
3230    while index < characters.len() {
3231        if characters[index] != '@' {
3232            index += 1;
3233            continue;
3234        }
3235        index += 1;
3236        let quoted = characters.get(index) == Some(&'"');
3237        if quoted {
3238            index += 1;
3239        }
3240        let start = index;
3241        while index < characters.len()
3242            && if quoted {
3243                characters[index] != '"'
3244            } else {
3245                !characters[index].is_whitespace()
3246            }
3247        {
3248            index += 1;
3249        }
3250        if index > start {
3251            paths.push(characters[start..index].iter().collect());
3252        }
3253        if quoted && index < characters.len() {
3254            index += 1;
3255        }
3256    }
3257    paths
3258}
3259
3260fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3261    let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3262    for path in prompt_resource_paths(prompt) {
3263        if path.ends_with('/') {
3264            continue;
3265        }
3266        let Ok(resource) = resources::load(cwd, &path) else {
3267            continue;
3268        };
3269        let uri = format!("file://{}", resource.path.display());
3270        let resource_value = if let Some(text) = resource.text {
3271            serde_json::json!({
3272                "uri": uri,
3273                "text": text,
3274                "mimeType": resource.mime_type,
3275            })
3276        } else if let Some(data) = resource.data {
3277            serde_json::json!({
3278                "uri": uri,
3279                "blob": BASE64.encode(data),
3280                "mimeType": resource.mime_type,
3281            })
3282        } else {
3283            continue;
3284        };
3285        blocks.push(serde_json::json!({
3286            "type": "resource",
3287            "resource": resource_value,
3288        }));
3289    }
3290    blocks
3291}
3292
3293#[async_trait]
3294impl AgentAdapter for AcpAdapter {
3295    fn slot(&self) -> RosterSlot {
3296        self.slot
3297    }
3298
3299    fn session_id(&self) -> Option<String> {
3300        self.session_id.clone()
3301    }
3302
3303    fn protocol(&self) -> &'static str {
3304        "acp"
3305    }
3306
3307    fn needs_restart(&self) -> bool {
3308        self.reader.is_none() || self.child.is_none()
3309    }
3310
3311    fn capabilities(&self) -> AgentCapabilities {
3312        self.capabilities.clone()
3313    }
3314
3315    async fn start(&mut self) -> AdapterResult<()> {
3316        // Use the inherent implementation, which owns the cleanup boundary
3317        // around the multi-step ACP handshake.
3318        AcpAdapter::start(self).await
3319    }
3320
3321    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3322        if self.needs_restart() {
3323            return Err(AdapterError::Transport(
3324                "ACP agent transport is not running; reload the agent before retrying".into(),
3325            ));
3326        }
3327        let session_id = self
3328            .session_id
3329            .as_ref()
3330            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3331        self.tool_updates.clear();
3332        let request_id = self.next_request_id;
3333        self.next_request_id += 1;
3334        let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3335        let write_result = self
3336            .write_json(serde_json::json!({
3337                "jsonrpc": "2.0",
3338                "id": request_id,
3339                "method": "session/prompt",
3340                "params": {
3341                    "sessionId": session_id,
3342                    "prompt": prompt_blocks,
3343                },
3344            }))
3345            .await;
3346        if let Err(error) = write_result {
3347            self.reset_transport().await;
3348            return Err(error);
3349        }
3350        self.prompt_request_id = Some(request_id);
3351        Ok(())
3352    }
3353
3354    async fn cancel(&mut self) -> AdapterResult<bool> {
3355        let Some(session_id) = &self.session_id else {
3356            return Ok(false);
3357        };
3358        self.write_json(serde_json::json!({
3359            "jsonrpc": "2.0",
3360            "method": "session/cancel",
3361            "params": {"sessionId": session_id, "_meta": {}},
3362        }))
3363        .await?;
3364        let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3365            loop {
3366                match <Self as AgentAdapter>::next_event(self).await {
3367                    Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3368                    Some(Ok(_)) => {}
3369                    Some(Err(_)) => break,
3370                }
3371            }
3372        })
3373        .await
3374        .is_ok();
3375        if !settled {
3376            // A peer that never acknowledges cancellation cannot safely share
3377            // its stream with the next prompt. Restart the transport while
3378            // preserving a loadable provider session when supported.
3379            self.reload().await?;
3380        }
3381        Ok(true)
3382    }
3383
3384    async fn answer_permission(
3385        &mut self,
3386        request_id: String,
3387        answer: PermissionAnswer,
3388    ) -> AdapterResult<()> {
3389        let id = request_id
3390            .parse::<u64>()
3391            .map(Value::from)
3392            .unwrap_or_else(|_| Value::String(request_id));
3393        let outcome = match answer {
3394            PermissionAnswer::Selected { option_id } => {
3395                serde_json::json!({"outcome": "selected", "optionId": option_id})
3396            }
3397            PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3398        };
3399        self.write_json(serde_json::json!({
3400            "jsonrpc": "2.0",
3401            "id": id,
3402            // RequestPermissionResponse wraps the selected/cancelled
3403            // discriminator in its `outcome` field. Keep this nested shape
3404            // compatible with the Python ACP server and ACP schema.
3405            "result": {"outcome": outcome},
3406        }))
3407        .await
3408    }
3409
3410    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3411        let session_id = self
3412            .session_id
3413            .as_ref()
3414            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3415        let policy = match mode.as_str() {
3416            "plan" => "codeswarm:mode:plan",
3417            "default" | "manual" => "codeswarm:mode:manual",
3418            "accept-edits" => "codeswarm:mode:accept-edits",
3419            "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3420            other => other,
3421        };
3422        let native_mode = crate::policy::resolve(policy, &self.modes)
3423            .map(|mode| mode.id)
3424            .unwrap_or(mode);
3425        let _ = self
3426            .request(
3427                "session/set_mode",
3428                serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3429            )
3430            .await?;
3431        self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3432            slot: self.slot,
3433            current_mode: native_mode,
3434        }));
3435        Ok(())
3436    }
3437
3438    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3439        let session_id = self
3440            .session_id
3441            .clone()
3442            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3443        let config_id = self
3444            .model_config_id
3445            .clone()
3446            .ok_or(AdapterError::Unsupported("set_model"))?;
3447        if !self.models.iter().any(|candidate| candidate.id == model) {
3448            return Err(AdapterError::Protocol(
3449                "model is not advertised by the agent".into(),
3450            ));
3451        }
3452        let _ = self
3453            .request(
3454                "session/set_config_option",
3455                serde_json::json!({
3456                    "sessionId": session_id,
3457                    "configId": config_id,
3458                    "value": model,
3459                }),
3460            )
3461            .await?;
3462        Ok(())
3463    }
3464
3465    async fn reload(&mut self) -> AdapterResult<()> {
3466        // `stop` tears down the process and clears its transport-owned
3467        // session handle. A reload is different from a final shutdown: ACP
3468        // peers advertising `loadSession` must receive the prior ID so the
3469        // replacement process can resume the same conversation.
3470        let session_id = self
3471            .capabilities
3472            .supports_session_load
3473            .then(|| self.session_id.clone())
3474            .flatten();
3475        self.stop().await?;
3476        self.session_id = session_id.clone();
3477        let result = self.start().await;
3478        if result.is_err() {
3479            // `start` cleans up a partially initialized transport by calling
3480            // `stop`, which also clears the handle. Keep it available for a
3481            // subsequent retry after the coordinator reports the failure.
3482            self.session_id = session_id;
3483        }
3484        result
3485    }
3486
3487    async fn stop(&mut self) -> AdapterResult<()> {
3488        let terminals = std::mem::take(&mut self.terminals);
3489        self.queued_events.clear();
3490        self.tool_updates.clear();
3491        for terminal in terminals.values() {
3492            terminal.stop().await;
3493        }
3494        if let Some(mut child) = self.child.take() {
3495            terminate_child(&mut child).await?;
3496        }
3497        self.reader = None;
3498        self.session_id = None;
3499        self.prompt_request_id = None;
3500        if let Some(task) = self.stderr_task.take() {
3501            task.abort();
3502            let _ = task.await;
3503        }
3504        Ok(())
3505    }
3506
3507    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3508        if let Some(event) = self.queued_events.pop_front() {
3509            return Some(event);
3510        }
3511        loop {
3512            let line = match self.read_line().await {
3513                Ok(line) => line,
3514                Err(error) => {
3515                    // EOF or a broken pipe ends this transport. Reap the
3516                    // process now; `send_prompt` will reconnect it for the
3517                    // next human turn while preserving `session_id`.
3518                    self.reset_transport().await;
3519                    return Some(Err(error));
3520                }
3521            };
3522            let value: Value = match serde_json::from_str(&line) {
3523                Ok(value) => value,
3524                Err(_) => continue,
3525            };
3526            match self.reject_empty_permission_request(&value).await {
3527                Ok(true) => continue,
3528                Ok(false) => {}
3529                Err(error) => return Some(Err(error)),
3530            }
3531            match self.handle_client_request(&value).await {
3532                Ok(true) => continue,
3533                Ok(false) => {}
3534                Err(error) => return Some(Err(error)),
3535            }
3536            match parse_acp_value(self.slot, &value, &mut self.tool_updates) {
3537                Ok(Some(event)) => {
3538                    if let AgentEvent::ModelsReplaced {
3539                        config_id, models, ..
3540                    } = &event
3541                    {
3542                        self.model_config_id = Some(config_id.clone());
3543                        self.models = models.clone();
3544                        self.capabilities.supports_models = !models.is_empty();
3545                    }
3546                    return Some(Ok(event));
3547                }
3548                Ok(None) => {}
3549                Err(error) => return Some(Err(error)),
3550            }
3551            if value.get("id").is_some_and(|id| {
3552                self.prompt_request_id
3553                    .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3554            }) {
3555                if let Some(error) = value.get("error") {
3556                    self.prompt_request_id = None;
3557                    return Some(Err(AdapterError::Protocol(error.to_string())));
3558                }
3559                self.prompt_request_id = None;
3560                return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3561            }
3562        }
3563    }
3564}
3565
3566#[cfg(test)]
3567fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3568    let value: Value =
3569        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3570    parse_acp_value(slot, &value, &mut BTreeMap::new())
3571}
3572
3573fn parse_acp_value(
3574    slot: RosterSlot,
3575    value: &Value,
3576    tools: &mut BTreeMap<String, ToolUpdate>,
3577) -> AdapterResult<Option<AgentEvent>> {
3578    let method = value.get("method").and_then(Value::as_str);
3579    if method == Some("session/request_permission") {
3580        let params = value.get("params").cloned().unwrap_or(Value::Null);
3581        let request_id = value
3582            .get("id")
3583            .map(rpc_id_to_string)
3584            .unwrap_or_else(|| "permission".into());
3585        return Ok(parse_permission_event(
3586            slot,
3587            &params,
3588            &request_id,
3589            params.get("options"),
3590        ));
3591    }
3592    if method != Some("session/update") {
3593        return Ok(None);
3594    }
3595    let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3596        return Ok(None);
3597    };
3598    let kind = update.get("sessionUpdate").and_then(Value::as_str);
3599    if kind == Some("config_option_update")
3600        && let Some((config_id, models, current_model)) = parse_model_config(update)
3601    {
3602        return Ok(Some(AgentEvent::ModelsReplaced {
3603            slot,
3604            config_id,
3605            models,
3606            current_model,
3607        }));
3608    }
3609    if kind == Some("request_permission") {
3610        let request_id = update
3611            .get("toolCall")
3612            .and_then(|tool| tool.get("toolCallId"))
3613            .and_then(Value::as_str)
3614            .unwrap_or("permission");
3615        return Ok(parse_permission_event(
3616            slot,
3617            update,
3618            request_id,
3619            update.get("options"),
3620        ));
3621    }
3622    if kind == Some("available_commands_update") {
3623        let commands = update
3624            .get("availableCommands")
3625            .and_then(Value::as_array)
3626            .map(|commands| {
3627                commands
3628                    .iter()
3629                    .filter_map(|command| {
3630                        let name = command.get("name").and_then(Value::as_str)?.trim();
3631                        (!name.is_empty()).then(|| AgentCommand {
3632                            name: name.to_owned(),
3633                        })
3634                    })
3635                    .collect::<Vec<_>>()
3636            })
3637            .unwrap_or_default();
3638        return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3639    }
3640    if kind == Some("current_mode_update") {
3641        if let Some(mode) = update
3642            .get("currentModeId")
3643            .and_then(Value::as_str)
3644            .filter(|mode| !mode.trim().is_empty())
3645        {
3646            return Ok(Some(AgentEvent::ModeUpdated {
3647                slot,
3648                current_mode: mode.to_owned(),
3649            }));
3650        }
3651        return Ok(None);
3652    }
3653    if kind == Some("usage_update") {
3654        let Some(used) = update.get("used").and_then(Value::as_u64) else {
3655            return Ok(None);
3656        };
3657        let Some(size) = update.get("size").and_then(Value::as_u64) else {
3658            return Ok(None);
3659        };
3660        return Ok(Some(AgentEvent::UsageUpdated {
3661            slot,
3662            usage: UsageUpdate { used, size },
3663        }));
3664    }
3665    if let Some(terminal) = parse_terminal_event(update, kind) {
3666        return Ok(Some(AgentEvent::Terminal {
3667            slot,
3668            event: terminal,
3669        }));
3670    }
3671    let text = update
3672        .get("content")
3673        .and_then(|content| content.get("text"))
3674        .and_then(Value::as_str)
3675        .map(str::to_owned);
3676    if kind == Some("user_message_chunk") {
3677        return Ok(text
3678            .filter(|text| !text.is_empty())
3679            .map(|text| AgentEvent::UserText { slot, text }));
3680    }
3681    if kind == Some("agent_message_chunk")
3682        && let Some(mode) = text
3683            .as_deref()
3684            .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3685            .map(str::trim)
3686            .filter(|mode| !mode.is_empty())
3687    {
3688        // Gemini's native ACP bridge historically encoded a mode change as a
3689        // control marker in the message stream. It is state, not transcript
3690        // content; expose it as a normalized catalog replacement instead of
3691        // leaking the marker into the conversation.
3692        return Ok(Some(AgentEvent::ModesReplaced {
3693            slot,
3694            modes: vec![Mode {
3695                id: mode.to_owned(),
3696                label: mode.to_owned(),
3697            }],
3698            current_mode: Some(mode.to_owned()),
3699        }));
3700    }
3701    match (kind, text) {
3702        (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3703            Ok(Some(AgentEvent::Text { slot, text }))
3704        }
3705        (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3706            Ok(Some(AgentEvent::Thought { slot, text }))
3707        }
3708        (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3709            Ok(normalize_acp_tool(update, tools).map(|update| AgentEvent::Tool { slot, update }))
3710        }
3711        _ => Ok(None),
3712    }
3713}
3714
3715/// ACP updates are patches: omitted/invalid fields retain their last valid
3716/// values, whereas an explicit empty content array clears previous output.
3717fn normalize_acp_tool(
3718    value: &Value,
3719    tools: &mut BTreeMap<String, ToolUpdate>,
3720) -> Option<ToolUpdate> {
3721    let id = value.get("toolCallId")?.as_str()?;
3722    if id.trim().is_empty() {
3723        return None;
3724    }
3725    if value.get("sessionUpdate").and_then(Value::as_str) == Some("tool_call") {
3726        tools.remove(id);
3727    }
3728    let tool = tools.entry(id.to_owned()).or_insert_with(|| ToolUpdate {
3729        id: id.to_owned(),
3730        title: "Tool call".into(),
3731        status: ToolStatus::Pending,
3732        detail: None,
3733    });
3734    if let Some(title) = value.get("title").and_then(Value::as_str) {
3735        tool.title = title.to_owned();
3736    }
3737    if let Some(status) =
3738        value
3739            .get("status")
3740            .and_then(Value::as_str)
3741            .and_then(|status| match status {
3742                "pending" => Some(ToolStatus::Pending),
3743                "in_progress" => Some(ToolStatus::Running),
3744                "completed" => Some(ToolStatus::Completed),
3745                "failed" => Some(ToolStatus::Failed),
3746                _ => None,
3747            })
3748    {
3749        tool.status = status;
3750    }
3751    if let Some(content) = value.get("content").and_then(Value::as_array) {
3752        let text = content
3753            .iter()
3754            .filter_map(|entry| match entry.get("type").and_then(Value::as_str) {
3755                Some("content") => entry.get("content")?.get("text")?.as_str(),
3756                Some("diff") => entry.get("newText")?.as_str(),
3757                _ => None,
3758            })
3759            .collect::<Vec<_>>()
3760            .join("\n");
3761        tool.detail = (!text.is_empty()).then_some(text);
3762    } else if let Some(output) = value.get("rawOutput").filter(|output| !output.is_null()) {
3763        tool.detail = Some(
3764            output
3765                .as_str()
3766                .map(str::to_owned)
3767                .unwrap_or_else(|| output.to_string()),
3768        );
3769    }
3770    Some(tool.clone())
3771}
3772
3773fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3774    let config = value
3775        .get("configOptions")?
3776        .as_array()?
3777        .iter()
3778        .find(|option| {
3779            option.get("category").and_then(Value::as_str) == Some("model")
3780                && matches!(
3781                    option.get("type").and_then(Value::as_str),
3782                    Some("select" | "enum")
3783                )
3784        })?;
3785    let config_id = config.get("id")?.as_str()?.to_owned();
3786    let models = config
3787        .get("options")?
3788        .as_array()?
3789        .iter()
3790        .filter_map(|option| {
3791            let id = option.get("value")?.as_str()?.to_owned();
3792            let label = option
3793                .get("name")
3794                .or_else(|| option.get("label"))
3795                .and_then(Value::as_str)
3796                .unwrap_or(&id)
3797                .to_owned();
3798            Some(Mode { id, label })
3799        })
3800        .collect::<Vec<_>>();
3801    (!models.is_empty()).then(|| {
3802        let current = config
3803            .get("currentValue")
3804            .and_then(Value::as_str)
3805            .map(str::to_owned);
3806        (config_id, models, current)
3807    })
3808}
3809
3810/// Normalize terminal lifecycle updates emitted by ACP-compatible bridges and
3811/// native stream adapters. Protocols have used both snake_case update names
3812/// and a nested `terminal` object, so accept either without leaking that
3813/// shape beyond the adapter boundary.
3814fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3815    let nested = value.get("terminal").unwrap_or(value);
3816    let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3817    let id = nested
3818        .get("terminalId")
3819        .or_else(|| nested.get("terminal_id"))
3820        .or_else(|| nested.get("id"))
3821        .and_then(Value::as_str)
3822        .unwrap_or("terminal")
3823        .to_owned();
3824    match kind {
3825        "terminal_created" | "terminal_create" | "terminal_started" => {
3826            let command = nested
3827                .get("command")
3828                .and_then(Value::as_str)
3829                .unwrap_or("")
3830                .to_owned();
3831            Some(TerminalEvent::Created { id, command })
3832        }
3833        "terminal_output" | "terminal_output_chunk" => {
3834            let text = nested
3835                .get("output")
3836                .or_else(|| nested.get("text"))
3837                .and_then(Value::as_str)
3838                .unwrap_or("")
3839                .to_owned();
3840            Some(TerminalEvent::Output { id, text })
3841        }
3842        "terminal_exited" | "terminal_exit" => {
3843            let code = nested
3844                .get("exitCode")
3845                .or_else(|| nested.get("exit_code"))
3846                .or_else(|| nested.get("code"))
3847                .and_then(Value::as_i64)
3848                .unwrap_or(0) as i32;
3849            Some(TerminalEvent::Exited { id, code })
3850        }
3851        "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3852        _ => None,
3853    }
3854}
3855
3856fn parse_permission_event(
3857    slot: RosterSlot,
3858    value: &Value,
3859    request_id: &str,
3860    options: Option<&Value>,
3861) -> Option<AgentEvent> {
3862    let tool = value.get("toolCall").unwrap_or(value);
3863    let title = tool
3864        .get("title")
3865        .and_then(Value::as_str)
3866        .unwrap_or("Agent requests permission")
3867        .to_owned();
3868    let (options, option_ids): (Vec<String>, Vec<String>) = options
3869        .and_then(Value::as_array)
3870        .map(|options| {
3871            options
3872                .iter()
3873                .filter_map(|option| {
3874                    let label = option
3875                        .get("name")
3876                        .or_else(|| option.get("optionId"))
3877                        .and_then(Value::as_str)?
3878                        .to_owned();
3879                    let option_id = option
3880                        .get("optionId")
3881                        .or_else(|| option.get("id"))
3882                        .and_then(Value::as_str)
3883                        .map(str::to_owned)
3884                        .unwrap_or_else(|| label.clone());
3885                    Some((label, option_id))
3886                })
3887                .unzip()
3888        })
3889        .unwrap_or_default();
3890    if options.is_empty() {
3891        return None;
3892    }
3893    Some(AgentEvent::Permission {
3894        slot,
3895        request: PermissionRequest {
3896            id: request_id.to_owned(),
3897            title,
3898            options,
3899            option_ids,
3900        },
3901    })
3902}
3903
3904fn rpc_id_to_string(value: &Value) -> String {
3905    value
3906        .as_str()
3907        .map(str::to_owned)
3908        .or_else(|| value.as_u64().map(|id| id.to_string()))
3909        .unwrap_or_else(|| value.to_string())
3910}
3911
3912#[cfg(test)]
3913mod tests {
3914    use super::{
3915        AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3916        RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3917        parse_model_config, prompt_content_blocks, read_bounded_line,
3918    };
3919    #[cfg(target_os = "linux")]
3920    use super::{isolate_process_group, terminate_child};
3921    use crate::TerminalEvent;
3922    use crate::{
3923        AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3924        persistence::SessionMetadataStore,
3925        relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3926    };
3927    use async_trait::async_trait;
3928    use serde_json::Value;
3929    use std::collections::VecDeque;
3930    use std::sync::{
3931        Arc, Mutex,
3932        atomic::{AtomicUsize, Ordering},
3933    };
3934
3935    fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3936        let nonce = std::time::SystemTime::now()
3937            .duration_since(std::time::UNIX_EPOCH)
3938            .expect("clock")
3939            .as_nanos();
3940        std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3941    }
3942
3943    #[test]
3944    fn malformed_file_writes_preserve_existing_content() {
3945        let root = unique_test_path("codeswarm-write-validation", "dir");
3946        std::fs::create_dir_all(&root).unwrap();
3947        let file = root.join("keep.txt");
3948        std::fs::write(&file, "valuable content").unwrap();
3949        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3950        for content in [
3951            Value::Null,
3952            serde_json::json!(false),
3953            serde_json::json!(42),
3954            serde_json::json!([]),
3955        ] {
3956            assert!(
3957                adapter
3958                    .write_workspace_text(
3959                        &serde_json::json!({"path":"keep.txt", "content": content})
3960                    )
3961                    .is_err()
3962            );
3963            assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3964        }
3965        assert!(
3966            adapter
3967                .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
3968                .is_err()
3969        );
3970        assert!(
3971            adapter
3972                .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
3973                .is_err()
3974        );
3975        assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3976        adapter
3977            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
3978            .unwrap();
3979        assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
3980        adapter
3981            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
3982            .unwrap();
3983        assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
3984        std::fs::remove_dir_all(root).unwrap();
3985    }
3986
3987    #[tokio::test]
3988    async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
3989        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
3990        let mut adapter = AcpAdapter::new(
3991            0,
3992            std::env::current_dir().unwrap(),
3993            "sh",
3994            vec!["-c".into(), script.into()],
3995        );
3996        adapter.start().await.unwrap();
3997        let error = adapter
3998            .request_with_timeout(
3999                "session/set_mode",
4000                serde_json::json!({}),
4001                std::time::Duration::from_millis(10),
4002            )
4003            .await
4004            .unwrap_err();
4005        assert!(error.to_string().contains("session/set_mode timed out"));
4006        adapter.stop().await.unwrap();
4007        assert!(adapter.child.is_none());
4008    }
4009
4010    #[tokio::test]
4011    async fn goals_reach_every_roster_slot_without_native_goal_support() {
4012        use crate::goal::GoalCommand;
4013        let hosts = (0..3)
4014            .map(|slot| {
4015                AdapterHost::new(
4016                    Box::new(ScriptedAdapter::new(
4017                        slot,
4018                        AgentCapabilities::default(),
4019                        [
4020                            AgentEvent::TurnComplete { slot },
4021                            AgentEvent::TurnComplete { slot },
4022                        ],
4023                    )),
4024                    None,
4025                )
4026            })
4027            .collect();
4028        let mut relay = RelayHost::new(hosts, 10).unwrap();
4029        relay.start().await.unwrap();
4030        let task = relay
4031            .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4032            .unwrap()
4033            .unwrap();
4034        relay.relay_mut().enqueue_human(task, Some(0));
4035        for slot in 0..3 {
4036            relay.run_turn("", 0).await.unwrap();
4037            let (actual, prompt) = relay.dispatches().last().unwrap();
4038            assert_eq!(*actual, slot);
4039            assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4040        }
4041        let snapshot = relay.session_metadata();
4042        let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4043        assert!(restored.is_some());
4044        relay.restore_goal(restored);
4045        relay.reload(0).await.unwrap();
4046        relay.run_turn("", 0).await.unwrap();
4047        assert!(
4048            relay
4049                .dispatches()
4050                .last()
4051                .unwrap()
4052                .1
4053                .contains("Active shared goal: Ship the settings screen")
4054        );
4055        relay.apply_goal(GoalCommand::Done).unwrap();
4056        relay.run_turn("", 0).await.unwrap();
4057        assert!(
4058            relay
4059                .dispatches()
4060                .last()
4061                .unwrap()
4062                .1
4063                .contains("No active shared goal")
4064        );
4065        relay.apply_goal(GoalCommand::Clear).unwrap();
4066        assert!(relay.session_metadata().get("goal").unwrap().is_null());
4067    }
4068
4069    #[tokio::test]
4070    async fn replacement_agent_receives_task_after_public_journal_pruning() {
4071        let hosts = (0..2)
4072            .map(|slot| {
4073                AdapterHost::new(
4074                    Box::new(ScriptedAdapter::new(
4075                        slot,
4076                        AgentCapabilities::default(),
4077                        [
4078                            AgentEvent::Text {
4079                                slot,
4080                                text: "progress".into(),
4081                            },
4082                            AgentEvent::TurnComplete { slot },
4083                            AgentEvent::Text {
4084                                slot,
4085                                text: "more progress".into(),
4086                            },
4087                            AgentEvent::TurnComplete { slot },
4088                        ],
4089                    )),
4090                    None,
4091                )
4092            })
4093            .collect();
4094        let mut relay = RelayHost::new(hosts, 10).unwrap();
4095        relay.start().await.unwrap();
4096        relay
4097            .relay_mut()
4098            .enqueue_human("Fix the login bug", Some(0));
4099        relay.run_turn("", 0).await.unwrap();
4100        relay.run_turn("", 0).await.unwrap();
4101        relay.run_turn("", 0).await.unwrap();
4102        relay.reload(1).await.unwrap();
4103        assert!(
4104            !relay
4105                .relay_mut()
4106                .unseen_context(1)
4107                .contains("Fix the login bug")
4108        );
4109        relay.run_turn("", 0).await.unwrap();
4110        assert!(
4111            relay
4112                .dispatches()
4113                .last()
4114                .unwrap()
4115                .1
4116                .contains("Shared task:\nFix the login bug")
4117        );
4118    }
4119
4120    #[cfg(target_os = "linux")]
4121    #[tokio::test]
4122    async fn termination_kills_only_the_verified_isolated_child_group() {
4123        use nix::unistd::{Pid, getpgid, getpgrp};
4124        use tokio::io::{AsyncBufReadExt, BufReader};
4125
4126        let own_group = getpgrp();
4127        let mut command = tokio::process::Command::new("sh");
4128        isolate_process_group(&mut command);
4129        command
4130            .arg("-c")
4131            .arg("sleep 60 & echo $!; wait")
4132            .stdout(std::process::Stdio::piped());
4133        let mut child = command.spawn().expect("spawn isolated shell");
4134        let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4135        assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4136        assert_ne!(leader, own_group);
4137
4138        let stdout = child.stdout.take().expect("child stdout");
4139        let mut lines = BufReader::new(stdout).lines();
4140        let descendant = lines
4141            .next_line()
4142            .await
4143            .expect("read descendant pid")
4144            .expect("descendant pid")
4145            .parse::<i32>()
4146            .expect("numeric descendant pid");
4147        let descendant = Pid::from_raw(descendant);
4148        assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4149
4150        terminate_child(&mut child).await.expect("terminate group");
4151        for _ in 0..100 {
4152            if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4153                return;
4154            }
4155            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4156        }
4157        panic!("descendant {descendant} survived isolated group termination");
4158    }
4159
4160    #[test]
4161    fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4162        assert_eq!(
4163            parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4164            Ok((
4165                "agent".into(),
4166                vec![
4167                    "--name".into(),
4168                    "local bridge".into(),
4169                    "--flag".into(),
4170                    "two words".into(),
4171                ]
4172            ),)
4173        );
4174        assert_eq!(
4175            parse_command_line(r#"agent "" escaped\ argument"#),
4176            Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4177        );
4178    }
4179
4180    #[test]
4181    fn acp_prompt_expands_safe_at_path_resources() {
4182        let root = unique_test_path("codeswarm-prompt-resource", "dir");
4183        std::fs::create_dir_all(&root).expect("workspace");
4184        std::fs::write(root.join("note.md"), "resource text").expect("resource");
4185        let blocks = prompt_content_blocks(&root, "inspect @note.md");
4186        assert_eq!(blocks[0]["type"], "text");
4187        assert_eq!(blocks[0]["text"], "inspect @note.md");
4188        assert_eq!(blocks[1]["type"], "resource");
4189        assert_eq!(blocks[1]["resource"]["text"], "resource text");
4190        assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4191        std::fs::remove_dir_all(root).expect("cleanup workspace");
4192    }
4193
4194    #[tokio::test]
4195    async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4196        let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4197        bytes.push(b'\n');
4198        let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4199        assert!(matches!(
4200            read_bounded_line(&mut reader).await,
4201            Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4202        ));
4203    }
4204
4205    #[test]
4206    fn rejects_malformed_configured_commands_before_spawn() {
4207        assert_eq!(
4208            parse_command_line("agent 'unfinished"),
4209            Err(super::CommandParseError::UnterminatedQuote)
4210        );
4211        assert_eq!(
4212            parse_command_line("agent\\"),
4213            Err(super::CommandParseError::TrailingEscape)
4214        );
4215        assert_eq!(
4216            parse_command_line("   \t"),
4217            Err(super::CommandParseError::Empty)
4218        );
4219    }
4220
4221    #[derive(Debug)]
4222    struct PendingAdapter {
4223        slot: usize,
4224        hang_on_cancel: bool,
4225    }
4226
4227    #[derive(Debug)]
4228    struct ConcurrentStartAdapter {
4229        slot: usize,
4230        barrier: Arc<tokio::sync::Barrier>,
4231    }
4232
4233    #[derive(Debug)]
4234    struct ReloadProbeAdapter {
4235        slot: usize,
4236        crashed: bool,
4237        reloaded: bool,
4238        events: VecDeque<AgentEvent>,
4239        prompts: Arc<Mutex<Vec<String>>>,
4240    }
4241
4242    #[async_trait]
4243    impl AgentAdapter for ReloadProbeAdapter {
4244        fn slot(&self) -> usize {
4245            self.slot
4246        }
4247
4248        fn display_name(&self) -> String {
4249            "Reload probe".into()
4250        }
4251
4252        fn protocol(&self) -> &'static str {
4253            "native"
4254        }
4255
4256        fn capabilities(&self) -> AgentCapabilities {
4257            AgentCapabilities {
4258                supports_modes: true,
4259                ..AgentCapabilities::default()
4260            }
4261        }
4262
4263        fn needs_restart(&self) -> bool {
4264            self.crashed && !self.reloaded
4265        }
4266
4267        async fn start(&mut self) -> super::AdapterResult<()> {
4268            self.events.push_back(AgentEvent::ModesReplaced {
4269                slot: self.slot,
4270                modes: vec![
4271                    Mode {
4272                        id: "codeswarm:mode:full-access".into(),
4273                        label: "Auto pilot".into(),
4274                    },
4275                    Mode {
4276                        id: "codeswarm:mode:plan".into(),
4277                        label: "Plan".into(),
4278                    },
4279                ],
4280                current_mode: Some("codeswarm:mode:full-access".into()),
4281            });
4282            self.events.push_back(AgentEvent::Ready {
4283                slot: self.slot,
4284                capabilities: self.capabilities(),
4285            });
4286            Ok(())
4287        }
4288
4289        async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4290            self.prompts.lock().expect("prompts").push(prompt);
4291            if !self.crashed {
4292                self.crashed = true;
4293                self.events.push_back(AgentEvent::Failed {
4294                    slot: self.slot,
4295                    started: true,
4296                    detail: "probe crashed".into(),
4297                });
4298            } else {
4299                self.events.push_back(AgentEvent::Text {
4300                    slot: self.slot,
4301                    text: "recovered".into(),
4302                });
4303                self.events
4304                    .push_back(AgentEvent::TurnComplete { slot: self.slot });
4305            }
4306            Ok(())
4307        }
4308
4309        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4310            Ok(false)
4311        }
4312
4313        async fn answer_permission(
4314            &mut self,
4315            _request_id: String,
4316            _answer: PermissionAnswer,
4317        ) -> super::AdapterResult<()> {
4318            Err(super::AdapterError::Unsupported("permission answer"))
4319        }
4320
4321        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4322            Ok(())
4323        }
4324
4325        async fn reload(&mut self) -> super::AdapterResult<()> {
4326            self.reloaded = true;
4327            self.start().await
4328        }
4329
4330        async fn stop(&mut self) -> super::AdapterResult<()> {
4331            Ok(())
4332        }
4333
4334        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4335            self.events.pop_front().map(Ok)
4336        }
4337    }
4338
4339    #[async_trait]
4340    impl AgentAdapter for ConcurrentStartAdapter {
4341        fn slot(&self) -> usize {
4342            self.slot
4343        }
4344
4345        fn capabilities(&self) -> AgentCapabilities {
4346            AgentCapabilities::default()
4347        }
4348
4349        async fn start(&mut self) -> super::AdapterResult<()> {
4350            self.barrier.wait().await;
4351            Ok(())
4352        }
4353
4354        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4355            Ok(())
4356        }
4357
4358        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4359            Ok(true)
4360        }
4361
4362        async fn answer_permission(
4363            &mut self,
4364            _request_id: String,
4365            _answer: PermissionAnswer,
4366        ) -> super::AdapterResult<()> {
4367            Ok(())
4368        }
4369
4370        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4371            Ok(())
4372        }
4373
4374        async fn reload(&mut self) -> super::AdapterResult<()> {
4375            Ok(())
4376        }
4377
4378        async fn stop(&mut self) -> super::AdapterResult<()> {
4379            Ok(())
4380        }
4381
4382        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4383            std::future::pending().await
4384        }
4385    }
4386
4387    #[derive(Debug)]
4388    struct PermissionBlockingAdapter {
4389        slot: usize,
4390        phase: u8,
4391    }
4392
4393    #[async_trait]
4394    impl AgentAdapter for PermissionBlockingAdapter {
4395        fn slot(&self) -> usize {
4396            self.slot
4397        }
4398
4399        fn capabilities(&self) -> AgentCapabilities {
4400            AgentCapabilities {
4401                supports_permissions: true,
4402                ..AgentCapabilities::default()
4403            }
4404        }
4405
4406        async fn start(&mut self) -> super::AdapterResult<()> {
4407            Ok(())
4408        }
4409
4410        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4411            Ok(())
4412        }
4413
4414        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4415            Ok(true)
4416        }
4417
4418        async fn answer_permission(
4419            &mut self,
4420            request_id: String,
4421            answer: PermissionAnswer,
4422        ) -> super::AdapterResult<()> {
4423            if self.phase != 1 || request_id != "permission-1" {
4424                return Err(super::AdapterError::Protocol(
4425                    "unexpected permission response".into(),
4426                ));
4427            }
4428            assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4429            self.phase = 2;
4430            Ok(())
4431        }
4432
4433        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4434            Ok(())
4435        }
4436
4437        async fn reload(&mut self) -> super::AdapterResult<()> {
4438            Ok(())
4439        }
4440
4441        async fn stop(&mut self) -> super::AdapterResult<()> {
4442            Ok(())
4443        }
4444
4445        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4446            match self.phase {
4447                0 => {
4448                    self.phase = 1;
4449                    Some(Ok(AgentEvent::Permission {
4450                        slot: self.slot,
4451                        request: crate::PermissionRequest {
4452                            id: "permission-1".into(),
4453                            title: "Allow?".into(),
4454                            options: vec!["Allow".into()],
4455                            option_ids: vec!["allow".into()],
4456                        },
4457                    }))
4458                }
4459                1 => std::future::pending().await,
4460                _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4461            }
4462        }
4463    }
4464
4465    #[async_trait]
4466    impl AgentAdapter for PendingAdapter {
4467        fn slot(&self) -> usize {
4468            self.slot
4469        }
4470
4471        fn capabilities(&self) -> AgentCapabilities {
4472            AgentCapabilities {
4473                supports_cancel: true,
4474                ..AgentCapabilities::default()
4475            }
4476        }
4477
4478        async fn start(&mut self) -> super::AdapterResult<()> {
4479            Ok(())
4480        }
4481
4482        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4483            Ok(())
4484        }
4485
4486        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4487            if self.hang_on_cancel {
4488                return std::future::pending().await;
4489            }
4490            Ok(true)
4491        }
4492
4493        async fn answer_permission(
4494            &mut self,
4495            _request_id: String,
4496            _answer: PermissionAnswer,
4497        ) -> super::AdapterResult<()> {
4498            Ok(())
4499        }
4500
4501        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4502            Ok(())
4503        }
4504
4505        async fn reload(&mut self) -> super::AdapterResult<()> {
4506            Ok(())
4507        }
4508
4509        async fn stop(&mut self) -> super::AdapterResult<()> {
4510            Ok(())
4511        }
4512
4513        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4514            std::future::pending().await
4515        }
4516    }
4517
4518    #[derive(Debug)]
4519    struct StopTrackingAdapter {
4520        slot: usize,
4521        stopped: Arc<AtomicUsize>,
4522        fail_stop: bool,
4523    }
4524
4525    #[derive(Debug)]
4526    struct ModeOrderAdapter {
4527        slot: usize,
4528        log: Arc<Mutex<Vec<String>>>,
4529        phase: u8,
4530    }
4531
4532    #[derive(Debug)]
4533    struct StartupAcpAdapter {
4534        slot: usize,
4535        events: std::collections::VecDeque<AgentEvent>,
4536    }
4537
4538    impl StartupAcpAdapter {
4539        fn new(slot: usize) -> Self {
4540            Self {
4541                slot,
4542                events: [
4543                    AgentEvent::ModesReplaced {
4544                        slot,
4545                        modes: vec![Mode {
4546                            id: "full-access".into(),
4547                            label: "Auto pilot".into(),
4548                        }],
4549                        current_mode: Some("full-access".into()),
4550                    },
4551                    AgentEvent::Ready {
4552                        slot,
4553                        capabilities: AgentCapabilities {
4554                            supports_modes: true,
4555                            ..AgentCapabilities::default()
4556                        },
4557                    },
4558                ]
4559                .into(),
4560            }
4561        }
4562    }
4563
4564    #[async_trait]
4565    impl AgentAdapter for StartupAcpAdapter {
4566        fn slot(&self) -> usize {
4567            self.slot
4568        }
4569
4570        fn protocol(&self) -> &'static str {
4571            "acp"
4572        }
4573
4574        fn capabilities(&self) -> AgentCapabilities {
4575            AgentCapabilities {
4576                supports_modes: true,
4577                ..AgentCapabilities::default()
4578            }
4579        }
4580
4581        async fn start(&mut self) -> super::AdapterResult<()> {
4582            Ok(())
4583        }
4584
4585        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4586            Ok(())
4587        }
4588
4589        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4590            Ok(true)
4591        }
4592
4593        async fn answer_permission(
4594            &mut self,
4595            _request_id: String,
4596            _answer: PermissionAnswer,
4597        ) -> super::AdapterResult<()> {
4598            Ok(())
4599        }
4600
4601        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4602            Ok(())
4603        }
4604
4605        async fn reload(&mut self) -> super::AdapterResult<()> {
4606            self.events = Self::new(self.slot).events;
4607            Ok(())
4608        }
4609
4610        async fn stop(&mut self) -> super::AdapterResult<()> {
4611            Ok(())
4612        }
4613
4614        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4615            self.events.pop_front().map(Ok)
4616        }
4617    }
4618
4619    #[async_trait]
4620    impl AgentAdapter for ModeOrderAdapter {
4621        fn slot(&self) -> usize {
4622            self.slot
4623        }
4624
4625        fn capabilities(&self) -> AgentCapabilities {
4626            AgentCapabilities {
4627                supports_modes: true,
4628                ..AgentCapabilities::default()
4629            }
4630        }
4631
4632        async fn start(&mut self) -> super::AdapterResult<()> {
4633            self.log.lock().expect("log").push("start".into());
4634            Ok(())
4635        }
4636
4637        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4638            self.log.lock().expect("log").push("prompt".into());
4639            Ok(())
4640        }
4641
4642        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4643            Ok(true)
4644        }
4645
4646        async fn answer_permission(
4647            &mut self,
4648            _request_id: String,
4649            _answer: PermissionAnswer,
4650        ) -> super::AdapterResult<()> {
4651            Ok(())
4652        }
4653
4654        async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4655            self.log.lock().expect("log").push(format!("mode:{mode}"));
4656            Ok(())
4657        }
4658
4659        async fn reload(&mut self) -> super::AdapterResult<()> {
4660            self.log.lock().expect("log").push("reload".into());
4661            self.phase = 0;
4662            Ok(())
4663        }
4664
4665        async fn stop(&mut self) -> super::AdapterResult<()> {
4666            Ok(())
4667        }
4668
4669        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4670            match self.phase {
4671                0 => {
4672                    self.phase = 1;
4673                    Some(Ok(AgentEvent::ModesReplaced {
4674                        slot: self.slot,
4675                        modes: vec![Mode {
4676                            id: "yolo".into(),
4677                            label: "YOLO".into(),
4678                        }],
4679                        current_mode: None,
4680                    }))
4681                }
4682                1 => {
4683                    self.phase = 2;
4684                    Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4685                }
4686                _ => std::future::pending().await,
4687            }
4688        }
4689    }
4690
4691    #[async_trait]
4692    impl AgentAdapter for StopTrackingAdapter {
4693        fn slot(&self) -> usize {
4694            self.slot
4695        }
4696
4697        fn capabilities(&self) -> AgentCapabilities {
4698            AgentCapabilities::default()
4699        }
4700
4701        async fn start(&mut self) -> super::AdapterResult<()> {
4702            Ok(())
4703        }
4704
4705        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4706            Ok(())
4707        }
4708
4709        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4710            Ok(false)
4711        }
4712
4713        async fn answer_permission(
4714            &mut self,
4715            _request_id: String,
4716            _answer: PermissionAnswer,
4717        ) -> super::AdapterResult<()> {
4718            Ok(())
4719        }
4720
4721        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4722            Ok(())
4723        }
4724
4725        async fn reload(&mut self) -> super::AdapterResult<()> {
4726            Ok(())
4727        }
4728
4729        async fn stop(&mut self) -> super::AdapterResult<()> {
4730            self.stopped.fetch_add(1, Ordering::Relaxed);
4731            if self.fail_stop {
4732                Err(super::AdapterError::Transport("stop failed".into()))
4733            } else {
4734                Ok(())
4735            }
4736        }
4737
4738        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4739            None
4740        }
4741    }
4742
4743    /// Startup can fail after an adapter has allocated resources.  Keep a
4744    /// fixture that records whether the failing adapter itself receives the
4745    /// cleanup call, not just the already-started peers.
4746    #[derive(Debug)]
4747    struct FailingStartAdapter {
4748        slot: usize,
4749        stopped: Arc<AtomicUsize>,
4750    }
4751
4752    #[async_trait]
4753    impl AgentAdapter for FailingStartAdapter {
4754        fn slot(&self) -> usize {
4755            self.slot
4756        }
4757
4758        fn capabilities(&self) -> AgentCapabilities {
4759            AgentCapabilities::default()
4760        }
4761
4762        async fn start(&mut self) -> super::AdapterResult<()> {
4763            Err(super::AdapterError::Spawn("startup failed".into()))
4764        }
4765
4766        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4767            Ok(())
4768        }
4769
4770        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4771            Ok(false)
4772        }
4773
4774        async fn answer_permission(
4775            &mut self,
4776            _request_id: String,
4777            _answer: PermissionAnswer,
4778        ) -> super::AdapterResult<()> {
4779            Ok(())
4780        }
4781
4782        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4783            Ok(())
4784        }
4785
4786        async fn reload(&mut self) -> super::AdapterResult<()> {
4787            Ok(())
4788        }
4789
4790        async fn stop(&mut self) -> super::AdapterResult<()> {
4791            self.stopped.fetch_add(1, Ordering::Relaxed);
4792            Ok(())
4793        }
4794
4795        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4796            None
4797        }
4798    }
4799
4800    #[tokio::test]
4801    async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4802        let stopped = Arc::new(AtomicUsize::new(0));
4803        let relay = RelayHost::new(
4804            vec![
4805                AdapterHost::new(
4806                    Box::new(StopTrackingAdapter {
4807                        slot: 0,
4808                        stopped: Arc::clone(&stopped),
4809                        fail_stop: true,
4810                    }),
4811                    None,
4812                ),
4813                AdapterHost::new(
4814                    Box::new(StopTrackingAdapter {
4815                        slot: 1,
4816                        stopped: Arc::clone(&stopped),
4817                        fail_stop: false,
4818                    }),
4819                    None,
4820                ),
4821            ],
4822            4,
4823        )
4824        .expect("relay");
4825        let mut relay = relay;
4826
4827        let error = relay.stop().await.expect_err("first stop failure");
4828        assert!(error.to_string().contains("stop failed"));
4829        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4830    }
4831
4832    #[tokio::test]
4833    async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4834        let stopped = Arc::new(AtomicUsize::new(0));
4835        let mut relay = RelayHost::new(
4836            vec![
4837                AdapterHost::new(
4838                    Box::new(StopTrackingAdapter {
4839                        slot: 0,
4840                        stopped: Arc::clone(&stopped),
4841                        fail_stop: false,
4842                    }),
4843                    None,
4844                ),
4845                AdapterHost::new(
4846                    Box::new(FailingStartAdapter {
4847                        slot: 1,
4848                        stopped: Arc::clone(&stopped),
4849                    }),
4850                    None,
4851                ),
4852            ],
4853            4,
4854        )
4855        .expect("relay");
4856
4857        relay.start().await.expect("healthy peer remains available");
4858        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4859        assert_eq!(stopped.load(Ordering::Relaxed), 1);
4860        relay.stop().await.unwrap();
4861        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4862    }
4863
4864    #[test]
4865    fn parses_acp_text_without_ui_dependency() {
4866        let event = parse_acp_notification(
4867            2,
4868            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4869        )
4870        .expect("valid ACP")
4871        .expect("text event");
4872        assert_eq!(
4873            event,
4874            AgentEvent::Text {
4875                slot: 2,
4876                text: "hello".into(),
4877            }
4878        );
4879    }
4880
4881    #[test]
4882    fn parses_acp_state_notifications_at_the_adapter_boundary() {
4883        let commands = parse_acp_notification(
4884            3,
4885            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4886        )
4887        .expect("valid ACP")
4888        .expect("commands event");
4889        assert_eq!(
4890            commands,
4891            AgentEvent::CommandsReplaced {
4892                slot: 3,
4893                commands: vec![crate::AgentCommand {
4894                    name: "review".into()
4895                }]
4896            }
4897        );
4898
4899        let mode = parse_acp_notification(
4900            3,
4901            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4902        )
4903        .expect("valid ACP")
4904        .expect("mode event");
4905        assert_eq!(
4906            mode,
4907            AgentEvent::ModeUpdated {
4908                slot: 3,
4909                current_mode: "review".into()
4910            }
4911        );
4912
4913        let usage = parse_acp_notification(
4914            3,
4915            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4916        )
4917        .expect("valid ACP")
4918        .expect("usage event");
4919        assert_eq!(
4920            usage,
4921            AgentEvent::UsageUpdated {
4922                slot: 3,
4923                usage: crate::UsageUpdate {
4924                    used: 4200,
4925                    size: 128000
4926                }
4927            }
4928        );
4929
4930        let models = parse_acp_notification(
4931            3,
4932            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"config_option_update","configOptions":[{"id":"model","category":"model","type":"select","currentValue":"smart","options":[{"value":"fast","name":"Fast"},{"value":"smart","name":"Smart"}]}]}}}"#,
4933        )
4934        .expect("valid ACP")
4935        .expect("models event");
4936        assert!(matches!(
4937            models,
4938            AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4939                if models.len() == 2 && current_model.as_deref() == Some("smart")
4940        ));
4941        assert_eq!(
4942            parse_model_config(&serde_json::json!({
4943                "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4944            })),
4945            None
4946        );
4947
4948        let user = parse_acp_notification(
4949            3,
4950            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4951        )
4952        .expect("valid ACP")
4953        .expect("user event");
4954        assert_eq!(
4955            user,
4956            AgentEvent::UserText {
4957                slot: 3,
4958                text: "context".into()
4959            }
4960        );
4961    }
4962
4963    #[test]
4964    fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4965        let event = parse_acp_notification(
4966            0,
4967            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4968        )
4969        .expect("valid ACP")
4970        .expect("mode event");
4971        assert!(matches!(
4972            event,
4973            AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4974                if mode == "yolo" && modes[0].id == "yolo"
4975        ));
4976    }
4977
4978    #[test]
4979    fn parses_native_agy_text_without_acp_bridge() {
4980        let event = parse_agy_line(
4981            1,
4982            r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4983        )
4984        .expect("valid stream-json")
4985        .expect("text event");
4986        assert_eq!(
4987            event,
4988            AgentEvent::Text {
4989                slot: 1,
4990                text: "hello".into(),
4991            }
4992        );
4993    }
4994
4995    #[test]
4996    fn parses_tool_lifecycle_from_each_protocol() {
4997        let agy = parse_agy_line(
4998            1,
4999            r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5000        )
5001        .expect("valid native tool")
5002        .expect("tool event");
5003        assert!(matches!(
5004            agy,
5005            AgentEvent::Tool {
5006                update: crate::ToolUpdate {
5007                    status: ToolStatus::Completed,
5008                    ..
5009                },
5010                ..
5011            }
5012        ));
5013
5014        let acp = parse_acp_notification(
5015            1,
5016            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5017        )
5018        .expect("valid ACP tool")
5019        .expect("tool event");
5020        assert!(matches!(
5021            acp,
5022            AgentEvent::Tool {
5023                update: crate::ToolUpdate {
5024                    status: ToolStatus::Failed,
5025                    ..
5026                },
5027                ..
5028            }
5029        ));
5030    }
5031
5032    #[test]
5033    fn parses_terminal_lifecycle_from_acp_and_native_events() {
5034        let created = parse_acp_notification(
5035            0,
5036            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5037        )
5038        .expect("valid ACP terminal")
5039        .expect("terminal event");
5040        assert_eq!(
5041            created,
5042            AgentEvent::Terminal {
5043                slot: 0,
5044                event: TerminalEvent::Created {
5045                    id: "term-1".into(),
5046                    command: "cargo test".into(),
5047                },
5048            }
5049        );
5050        let output = parse_agy_line(
5051            1,
5052            r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5053        )
5054        .expect("valid native terminal")
5055        .expect("terminal event");
5056        assert_eq!(
5057            output,
5058            AgentEvent::Terminal {
5059                slot: 1,
5060                event: TerminalEvent::Output {
5061                    id: "term-1".into(),
5062                    text: "ok\n".into(),
5063                },
5064            }
5065        );
5066        let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5067            .expect("valid native release")
5068            .expect("terminal event");
5069        assert!(matches!(
5070            released,
5071            AgentEvent::Terminal {
5072                event: TerminalEvent::Released { id },
5073                ..
5074            } if id == "term-1"
5075        ));
5076    }
5077
5078    #[test]
5079    fn parses_acp_permission_requests() {
5080        let event = parse_acp_notification(
5081            0,
5082            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"request_permission","toolCall":{"toolCallId":"t1","title":"Write file"},"options":[{"name":"Allow once","optionId":"allow-once"},{"name":"Reject","optionId":"reject"}]}}}"#,
5083        )
5084        .expect("valid permission")
5085        .expect("permission event");
5086        assert!(matches!(
5087            event,
5088            AgentEvent::Permission { request, .. }
5089                if request.id == "t1"
5090                    && request.title == "Write file"
5091                    && request.options == ["Allow once", "Reject"]
5092                    && request.option_ids == ["allow-once", "reject"]
5093        ));
5094    }
5095
5096    #[test]
5097    fn parses_acp_permission_request_as_json_rpc_request() {
5098        let event = parse_acp_notification(
5099            2,
5100            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5101        )
5102        .expect("valid permission request")
5103        .expect("permission event");
5104        assert!(matches!(
5105            event,
5106            AgentEvent::Permission { request, .. }
5107                if request.id == "17"
5108                    && request.title == "Write file"
5109                    && request.options == ["allow-once", "reject"]
5110                    && request.option_ids == ["allow-once", "reject"]
5111        ));
5112    }
5113
5114    #[tokio::test]
5115    async fn native_adapter_explicitly_rejects_permission_answers() {
5116        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5117        assert_eq!(
5118            adapter
5119                .answer_permission(
5120                    "request".into(),
5121                    PermissionAnswer::Selected {
5122                        option_id: "allow".into()
5123                    },
5124                )
5125                .await,
5126            Err(super::AdapterError::Unsupported("permission answer"))
5127        );
5128    }
5129
5130    #[tokio::test]
5131    async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5132        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5133        adapter
5134            .set_mode("full-access".into())
5135            .await
5136            .expect("auto-pilot alias");
5137        assert!(matches!(
5138            adapter.next_event().await,
5139            Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5140        ));
5141    }
5142
5143    #[tokio::test]
5144    async fn native_turns_receive_a_twenty_four_hour_timeout() {
5145        let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5146        std::fs::write(
5147            &script_path,
5148            r#"#!/bin/sh
5149seen=0
5150while [ "$#" -gt 0 ]; do
5151    case "$1" in
5152        --print-timeout)
5153            shift
5154            [ "$1" = "1440m" ] || exit 2
5155            seen=$((seen + 1))
5156            ;;
5157    esac
5158    shift
5159done
5160[ "$seen" = 1 ] || exit 3
5161printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5162"#,
5163        )
5164        .unwrap();
5165        let mut adapter = AgyAdapter::with_session_id(
5166            0,
5167            std::env::current_dir().unwrap(),
5168            format!("sh {}", script_path.display()),
5169            "saved-session",
5170        );
5171        adapter.start().await.unwrap();
5172        adapter.next_event().await.unwrap().unwrap();
5173        adapter.next_event().await.unwrap().unwrap();
5174        for prompt in ["first task", "follow-up task"] {
5175            adapter.send_prompt(prompt.into()).await.unwrap();
5176            assert!(
5177                matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5178            );
5179            assert!(matches!(
5180                adapter.next_event().await,
5181                Some(Ok(AgentEvent::TurnComplete { .. }))
5182            ));
5183        }
5184        adapter.stop().await.unwrap();
5185        std::fs::remove_file(script_path).unwrap();
5186    }
5187
5188    #[tokio::test]
5189    async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5190        let script_path = unique_test_path("codeswarm-agy-session", "sh");
5191        std::fs::write(
5192            &script_path,
5193            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"init\",\"conversation_id\":\"native-session\"}' '{\"event\":\"step_update\",\"step_update\":{\"step_type\":\"agent_response\",\"text_delta\":\"ok\"}}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"ok\"}}'\n",
5194        )
5195        .expect("write native test script");
5196        let mut adapter = AgyAdapter::new(
5197            0,
5198            std::env::current_dir().expect("cwd"),
5199            format!("sh {}", script_path.display()),
5200        );
5201        adapter.start().await.expect("start native adapter");
5202        // Startup emits its mode catalog and readiness before a turn.
5203        assert!(adapter.next_event().await.is_some());
5204        assert!(adapter.next_event().await.is_some());
5205        adapter
5206            .send_prompt("first".into())
5207            .await
5208            .expect("first prompt");
5209        while !matches!(
5210            adapter.next_event().await,
5211            Some(Ok(AgentEvent::TurnComplete { .. }))
5212        ) {}
5213        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5214        adapter
5215            .send_prompt("follow up".into())
5216            .await
5217            .expect("follow-up prompt");
5218        while !matches!(
5219            adapter.next_event().await,
5220            Some(Ok(AgentEvent::TurnComplete { .. }))
5221        ) {}
5222        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5223        adapter.stop().await.expect("stop native adapter");
5224        std::fs::remove_file(script_path).expect("cleanup native script");
5225    }
5226
5227    #[tokio::test]
5228    async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5229        let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5230        std::fs::write(
5231            &script_path,
5232            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5233        )
5234        .expect("write native test script");
5235        let mut adapter = AgyAdapter::new(
5236            0,
5237            std::env::current_dir().expect("cwd"),
5238            format!("sh {}", script_path.display()),
5239        );
5240        adapter.start().await.expect("start native adapter");
5241        assert!(adapter.next_event().await.is_some());
5242        assert!(adapter.next_event().await.is_some());
5243        adapter.send_prompt("fail".into()).await.expect("prompt");
5244        assert!(matches!(
5245            adapter.next_event().await,
5246            Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5247                if detail == "agent failed"
5248        ));
5249        adapter.stop().await.expect("stop native adapter");
5250        std::fs::remove_file(script_path).expect("cleanup native script");
5251    }
5252
5253    #[tokio::test]
5254    async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5255        let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5256        let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5257        std::fs::write(
5258            &script_path,
5259            format!(
5260                "#!/bin/sh\ncount=0\nif [ -f '{}' ]; then count=$(cat '{}'); fi\ncount=$((count + 1))\nprintf '%s' \"$count\" > '{}'\nif [ \"$count\" = 1 ]; then\n  printf '%s\\n' '{{\"event\":\"result\",\"result\":{{\"status\":\"FAILURE\",\"error\":\"first crash\"}}}}'\nelse\n  printf '%s\\n' '{{\"event\":\"result\",\"result\":{{\"status\":\"SUCCESS\",\"response\":\"recovered\"}}}}'\nfi\n",
5261                marker_path.display(),
5262                marker_path.display(),
5263                marker_path.display(),
5264            ),
5265        )
5266        .expect("write retry script");
5267        let mut adapter = AgyAdapter::new(
5268            0,
5269            std::env::current_dir().expect("cwd"),
5270            format!("sh {}", script_path.display()),
5271        );
5272        adapter.start().await.expect("start native adapter");
5273        assert!(adapter.next_event().await.is_some());
5274        assert!(adapter.next_event().await.is_some());
5275        adapter
5276            .send_prompt("first".into())
5277            .await
5278            .expect("first prompt");
5279        assert!(matches!(
5280            adapter.next_event().await,
5281            Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5282        ));
5283        adapter
5284            .send_prompt("retry".into())
5285            .await
5286            .expect("retry prompt starts a fresh process");
5287        assert!(matches!(
5288            adapter.next_event().await,
5289            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5290        ));
5291        assert!(matches!(
5292            adapter.next_event().await,
5293            Some(Ok(AgentEvent::TurnComplete { .. }))
5294        ));
5295        adapter.stop().await.expect("stop native adapter");
5296        std::fs::remove_file(script_path).expect("cleanup retry script");
5297        std::fs::remove_file(marker_path).expect("cleanup retry marker");
5298    }
5299
5300    #[tokio::test]
5301    async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5302        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"}]}}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}'; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'"#;
5303        let cwd = std::env::current_dir().expect("cwd");
5304        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5305        adapter.start().await.expect("initialize");
5306        assert!(matches!(
5307            adapter.next_event().await,
5308            Some(Ok(AgentEvent::ModesReplaced { .. }))
5309        ));
5310        assert!(matches!(
5311            adapter.next_event().await,
5312            Some(Ok(AgentEvent::Ready { .. }))
5313        ));
5314        adapter.send_prompt("hello".into()).await.expect("prompt");
5315        assert!(matches!(
5316            adapter.next_event().await,
5317            Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5318        ));
5319        assert!(matches!(
5320            adapter.next_event().await,
5321            Some(Ok(AgentEvent::TurnComplete { .. }))
5322        ));
5323    }
5324
5325    #[tokio::test]
5326    async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5327        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"}]}}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"first"}}}}'; echo '{"jsonrpc":"2.0","id":"3","result":{"stopReason":"end_turn"}}'; read _; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"second"}}}}'; echo '{"jsonrpc":"2.0","id":"4","result":{"stopReason":"end_turn"}}'"#;
5328        let cwd = std::env::current_dir().expect("cwd");
5329        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5330        adapter.start().await.expect("initialize");
5331        assert!(adapter.next_event().await.is_some());
5332        assert!(adapter.next_event().await.is_some());
5333
5334        for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5335            adapter.send_prompt(prompt.into()).await.expect("prompt");
5336            assert!(matches!(
5337                adapter.next_event().await,
5338                Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5339            ));
5340            assert!(matches!(
5341                adapter.next_event().await,
5342                Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5343            ));
5344        }
5345        adapter.stop().await.expect("stop");
5346    }
5347
5348    #[tokio::test]
5349    async fn empty_acp_mode_catalog_disables_mode_control() {
5350        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"availableModes":[]}}}'"#;
5351        let cwd = std::env::current_dir().expect("cwd");
5352        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5353        adapter.start().await.expect("initialize");
5354        assert!(!adapter.capabilities().supports_modes);
5355        assert!(matches!(
5356            adapter.next_event().await,
5357            Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5358        ));
5359        assert!(matches!(
5360            adapter.next_event().await,
5361            Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5362        ));
5363        adapter.stop().await.expect("stop");
5364    }
5365
5366    #[tokio::test]
5367    async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5368        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","configOptions":[{"id":"model","category":"model","type":"select","currentValue":"fast","options":[{"value":"fast","name":"Fast"},{"value":"smart","name":"Smart"}]}]}}'; read request; case "$request" in *session/set_config_option*\"value\":\"smart\"*) echo '{"jsonrpc":"2.0","id":3,"result":{}}';; *) echo '{"jsonrpc":"2.0","id":3,"error":{"code":-32602,"message":"wrong model request"}}';; esac"#;
5369        let cwd = std::env::current_dir().expect("cwd");
5370        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5371        adapter.start().await.expect("initialize");
5372        assert!(adapter.capabilities().supports_models);
5373        assert!(matches!(
5374            adapter.next_event().await,
5375            Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5376                if config_id == "model"
5377                    && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5378                    && current_model.as_deref() == Some("fast")
5379        ));
5380        assert!(matches!(
5381            adapter.next_event().await,
5382            Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5383        ));
5384        adapter.set_model("smart".into()).await.expect("set model");
5385        assert!(adapter.set_model("invented".into()).await.is_err());
5386        adapter.stop().await.expect("stop");
5387    }
5388
5389    #[tokio::test]
5390    async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5391        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1","modes":{"currentModeId":"plan","availableModes":[{"id":"plan","name":"Plan"},{"id":"yolo","name":"YOLO"}]}}}'; read _; echo '{"jsonrpc":"2.0","id":3,"result":{}}'"#;
5392        let cwd = std::env::current_dir().expect("cwd");
5393        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5394        adapter.start().await.expect("initialize");
5395        adapter
5396            .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5397            .await
5398            .expect("set mode");
5399        assert!(matches!(
5400            adapter.next_event().await,
5401            Some(Ok(AgentEvent::ModesReplaced { .. }))
5402        ));
5403        assert!(matches!(
5404            adapter.next_event().await,
5405            Some(Ok(AgentEvent::Ready { .. }))
5406        ));
5407        assert!(matches!(
5408            adapter.next_event().await,
5409            Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5410        ));
5411        adapter.stop().await.expect("stop");
5412    }
5413
5414    #[tokio::test]
5415    async fn acp_reload_preserves_a_loadable_session_id() {
5416        let cwd = std::env::current_dir().expect("cwd");
5417        let mut adapter = AcpAdapter::with_session_id(
5418            0,
5419            cwd,
5420            "__codeswarm_missing_acp_for_reload_test__",
5421            Vec::new(),
5422            "saved-session",
5423        );
5424        adapter.capabilities.supports_session_load = true;
5425        // A failed replacement process still must not erase the session ID:
5426        // the coordinator can report the startup error and offer another
5427        // reload, preserving the only handle that can resume the conversation.
5428        assert!(adapter.reload().await.is_err());
5429        assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5430    }
5431
5432    #[tokio::test]
5433    async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5434        let cwd = std::env::current_dir().expect("cwd");
5435        let mut adapter = AcpAdapter::with_session_id(
5436            0,
5437            cwd,
5438            "__codeswarm_missing_nonloadable_acp__",
5439            Vec::new(),
5440            "stale-session",
5441        );
5442        adapter.capabilities.supports_session_load = false;
5443        assert!(adapter.reload().await.is_err());
5444        assert_eq!(adapter.session_id, None);
5445    }
5446
5447    #[tokio::test]
5448    async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5449        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"s1"}}'; read _; echo 'diagnostic from wrapper'; echo '{"jsonrpc":"2.0","method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"partial"}}}}'; echo '{"jsonrpc":"2.0","id":3,"error":{"code":-32000,"message":"capacity"}}'"#;
5450        let cwd = std::env::current_dir().expect("cwd");
5451        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5452        adapter.start().await.expect("initialize");
5453        assert!(matches!(
5454            adapter.next_event().await,
5455            Some(Ok(AgentEvent::Ready { .. }))
5456        ));
5457        adapter.send_prompt("hello".into()).await.expect("prompt");
5458        assert!(matches!(
5459            adapter.next_event().await,
5460            Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5461        ));
5462        assert!(matches!(
5463            adapter.next_event().await,
5464            Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5465        ));
5466    }
5467
5468    #[test]
5469    fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5470        let mut tools = std::collections::BTreeMap::new();
5471        let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5472            "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5473        let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5474        assert_eq!(initial.detail.as_deref(), Some("old output"));
5475        let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5476        assert_eq!(completed.title, "Read config");
5477        assert_eq!(completed.detail.as_deref(), Some("old output"));
5478        assert_eq!(completed.status, ToolStatus::Completed);
5479        let malformed = super::normalize_acp_tool(
5480            &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5481            &mut tools,
5482        )
5483        .unwrap();
5484        assert_eq!(malformed, completed);
5485        let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5486        assert_eq!(replaced.detail.as_deref(), Some("new output"));
5487        let cleared = super::normalize_acp_tool(
5488            &serde_json::json!({"toolCallId":"read","content":[]}),
5489            &mut tools,
5490        )
5491        .unwrap();
5492        assert_eq!(cleared.detail, None);
5493        let raw = super::normalize_acp_tool(
5494            &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5495            &mut tools,
5496        )
5497        .unwrap();
5498        assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5499        let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5500        assert_eq!(fresh.status, ToolStatus::Pending);
5501        assert_eq!(fresh.detail, None);
5502        for invalid in [
5503            serde_json::json!({}),
5504            serde_json::json!({"toolCallId":7}),
5505            serde_json::json!({"toolCallId":" "}),
5506        ] {
5507            assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5508        }
5509        assert_eq!(tools.len(), 1);
5510        // IDs are opaque, not whitespace-normalized aliases of another tool.
5511        super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5512        assert_eq!(tools.len(), 2);
5513    }
5514
5515    #[tokio::test]
5516    async fn acp_tool_status_only_notifications_retain_name_and_output() {
5517        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5518read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5519read _
5520echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"r","title":"Read config","status":"in_progress","content":[{"type":"content","content":{"type":"text","text":"file content"}}]}}}'
5521echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5522echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5523        let mut adapter = AcpAdapter::new(
5524            0,
5525            std::env::current_dir().unwrap(),
5526            "sh",
5527            vec!["-c".into(), script.into()],
5528        );
5529        adapter.start().await.unwrap();
5530        adapter.next_event().await.unwrap().unwrap();
5531        adapter.send_prompt("read".into()).await.unwrap();
5532        for status in [ToolStatus::Running, ToolStatus::Completed] {
5533            let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5534                panic!("tool event");
5535            };
5536            assert_eq!(update.status, status);
5537            assert_eq!(update.title, "Read config");
5538            assert_eq!(update.detail.as_deref(), Some("file content"));
5539        }
5540        assert!(matches!(
5541            adapter.next_event().await,
5542            Some(Ok(AgentEvent::TurnComplete { .. }))
5543        ));
5544        adapter.stop().await.unwrap();
5545    }
5546
5547    #[tokio::test]
5548    async fn acp_reload_discards_old_queued_events_and_catalogs() {
5549        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5550        let mut adapter = AcpAdapter::new(
5551            0,
5552            std::env::current_dir().unwrap(),
5553            "sh",
5554            vec!["-c".into(), script.into()],
5555        );
5556        adapter.start().await.unwrap();
5557        adapter.queued_events.push_back(Ok(AgentEvent::Text {
5558            slot: 0,
5559            text: "stale".into(),
5560        }));
5561        adapter.modes = vec![Mode {
5562            id: "stale".into(),
5563            label: "Stale".into(),
5564        }];
5565        // Real peers echo each request's ID. Reset for this fixed-ID test script.
5566        adapter.next_request_id = 1;
5567        adapter.reload().await.unwrap();
5568        assert!(adapter.modes.is_empty());
5569        assert_eq!(adapter.queued_events.len(), 1);
5570        assert!(matches!(
5571            adapter.next_event().await,
5572            Some(Ok(AgentEvent::Ready { .. }))
5573        ));
5574        adapter.stop().await.unwrap();
5575        assert!(adapter.queued_events.is_empty());
5576    }
5577
5578    #[tokio::test]
5579    async fn acp_load_replays_history_without_starting_a_turn() {
5580        let script = r#"
5581read _
5582echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5583read request
5584case "$request" in *session/load*) ;; *) exit 2;; esac
5585echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5586echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5587echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5588echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5589echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5590echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5591read request
5592case "$request" in *session/prompt*) ;; *) exit 3;; esac
5593echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5594echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5595"#;
5596        let mut adapter = AcpAdapter::with_session_id(
5597            2,
5598            std::env::current_dir().unwrap(),
5599            "sh",
5600            vec!["-c".into(), script.into()],
5601            "saved",
5602        );
5603        adapter.start().await.unwrap();
5604        let mut state = crate::SessionState::new(3);
5605        for _ in 0..5 {
5606            let event = adapter.next_event().await.unwrap().unwrap();
5607            assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5608            crate::reduce(&mut state, event);
5609            assert_eq!(state.active_slot, None);
5610        }
5611        assert!(matches!(
5612            adapter.next_event().await,
5613            Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5614        ));
5615        adapter.send_prompt("new question".into()).await.unwrap();
5616        assert!(
5617            matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5618        );
5619        assert!(matches!(
5620            adapter.next_event().await,
5621            Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5622        ));
5623        adapter.stop().await.unwrap();
5624    }
5625
5626    #[tokio::test]
5627    async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5628        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5629        let cwd = std::env::current_dir().expect("cwd");
5630        let mut adapter = AcpAdapter::with_session_id(
5631            0,
5632            cwd,
5633            "sh",
5634            vec!["-c".into(), script.into()],
5635            "existing-session",
5636        );
5637        adapter.start().await.expect("load existing session");
5638        assert!(matches!(
5639            adapter.next_event().await,
5640            Some(Ok(AgentEvent::Ready { .. }))
5641        ));
5642    }
5643
5644    #[tokio::test]
5645    async fn acp_start_failure_reaps_transport_process() {
5646        // The child emits an invalid initialize response and exits. The
5647        // adapter must not retain a live child after protocol startup fails;
5648        // this is the path used when a configured ACP command is unavailable
5649        // or speaks a different protocol.
5650        let mut adapter = AcpAdapter::new(
5651            0,
5652            std::env::current_dir().expect("cwd"),
5653            "sh",
5654            vec!["-c".into(), "printf 'not-json\\n'".into()],
5655        );
5656        assert!(adapter.start().await.is_err());
5657        assert!(adapter.child.is_none());
5658        assert!(adapter.reader.is_none());
5659    }
5660
5661    #[tokio::test]
5662    async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5663        let marker = unique_test_path("codeswarm-acp-retry", "count");
5664        let script = format!(
5665            r#"count=0
5666if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5667count=$((count + 1))
5668printf '%s' "$count" > '{0}'
5669while IFS= read -r request; do
5670  id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5671  case "$request" in
5672    *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5673    *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5674    *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5675    *session/prompt*)
5676      if [ "$count" = 1 ]; then exit 0; fi
5677      printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5678      printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5679      ;;
5680  esac
5681done
5682"#,
5683            marker.display()
5684        );
5685        let mut adapter = AcpAdapter::new(
5686            0,
5687            std::env::current_dir().expect("cwd"),
5688            "sh",
5689            vec!["-c".into(), script],
5690        );
5691        adapter.start().await.expect("initial ACP startup");
5692        assert!(matches!(
5693            adapter.next_event().await,
5694            Some(Ok(AgentEvent::Ready { .. }))
5695        ));
5696        adapter
5697            .send_prompt("first".into())
5698            .await
5699            .expect("first prompt");
5700        assert!(matches!(
5701            adapter.next_event().await,
5702            Some(Err(AdapterError::Transport(_)))
5703        ));
5704        assert!(adapter.child.is_none());
5705        assert!(adapter.reader.is_none());
5706        assert_eq!(adapter.session_id(), Some("saved-session".into()));
5707
5708        adapter.reload().await.expect("reload ACP transport");
5709        assert!(matches!(
5710            adapter.next_event().await,
5711            Some(Ok(AgentEvent::Ready { .. }))
5712        ));
5713        adapter
5714            .send_prompt("retry".into())
5715            .await
5716            .expect("retry prompt");
5717        assert!(matches!(
5718            adapter.next_event().await,
5719            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5720        ));
5721        assert!(matches!(
5722            adapter.next_event().await,
5723            Some(Ok(AgentEvent::TurnComplete { .. }))
5724        ));
5725        adapter.stop().await.expect("stop ACP");
5726        std::fs::remove_file(marker).expect("cleanup marker");
5727    }
5728
5729    #[tokio::test]
5730    async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5731        let prompts = Arc::new(Mutex::new(Vec::new()));
5732        let healthy = ScriptedAdapter::new(
5733            0,
5734            AgentCapabilities::default(),
5735            [
5736                AgentEvent::Text {
5737                    slot: 0,
5738                    text: "peer context".into(),
5739                },
5740                AgentEvent::TurnComplete { slot: 0 },
5741            ],
5742        );
5743        let probe = ReloadProbeAdapter {
5744            slot: 1,
5745            crashed: false,
5746            reloaded: false,
5747            events: VecDeque::new(),
5748            prompts: Arc::clone(&prompts),
5749        };
5750        let mut relay = RelayHost::new(
5751            vec![
5752                AdapterHost::new(Box::new(healthy), None),
5753                AdapterHost::new(Box::new(probe), None),
5754            ],
5755            8,
5756        )
5757        .expect("relay");
5758        relay.start().await.expect("start");
5759        relay
5760            .run_turn("original task", 0)
5761            .await
5762            .expect("first turn");
5763        relay.run_turn("", 0).await.expect("crashed turn");
5764        relay.relay_mut().enqueue_human("retry", Some(1));
5765        relay.run_turn("", 0).await.expect("reloaded turn");
5766
5767        let prompts = prompts.lock().expect("prompts");
5768        let retry = prompts.last().expect("retry prompt");
5769        assert!(retry.contains("You are Reload probe"), "{retry}");
5770        assert!(retry.contains("original task"), "{retry}");
5771        assert!(retry.contains("peer context"), "{retry}");
5772        assert!(retry.contains("retry"), "{retry}");
5773    }
5774
5775    #[tokio::test]
5776    async fn acp_adapter_answers_permission_json_rpc_requests() {
5777        let path = std::env::temp_dir().join(format!(
5778            "codeswarm-permission-answer-{}",
5779            std::process::id()
5780        ));
5781        let script = format!(
5782            r#"read _; echo '{{"jsonrpc":"2.0","id":1,"result":{{"agentCapabilities":{{}}}}}}'; read _; echo '{{"jsonrpc":"2.0","id":2,"result":{{"sessionId":"s1"}}}}'; read _; echo '{{"jsonrpc":"2.0","id":9,"method":"session/request_permission","params":{{"toolCall":{{"title":"Write file"}},"options":[{{"optionId":"allow-once","name":"Allow once"}}]}}}}'; read answer; printf '%s' "$answer" > '{}'; echo '{{"jsonrpc":"2.0","id":3,"result":{{"stopReason":"end_turn"}}}}'"#,
5783            path.display()
5784        );
5785        let mut adapter = AcpAdapter::new(
5786            0,
5787            std::env::current_dir().expect("cwd"),
5788            "sh",
5789            vec!["-c".into(), script],
5790        );
5791        adapter.start().await.expect("start ACP");
5792        assert!(matches!(
5793            adapter.next_event().await,
5794            Some(Ok(AgentEvent::Ready { .. }))
5795        ));
5796        adapter.send_prompt("do it".into()).await.expect("prompt");
5797        assert!(matches!(
5798            adapter.next_event().await,
5799            Some(Ok(AgentEvent::Permission { request, .. }))
5800                if request.id == "9"
5801                    && request.options == ["Allow once"]
5802                    && request.option_ids == ["allow-once"]
5803        ));
5804        adapter
5805            .answer_permission(
5806                "9".into(),
5807                PermissionAnswer::Selected {
5808                    option_id: "allow-once".into(),
5809                },
5810            )
5811            .await
5812            .expect("permission answer");
5813        assert!(matches!(
5814            adapter.next_event().await,
5815            Some(Ok(AgentEvent::TurnComplete { .. }))
5816        ));
5817        let answer: Value = serde_json::from_str(
5818            &std::fs::read_to_string(&path).expect("captured permission answer"),
5819        )
5820        .expect("valid JSON-RPC answer");
5821        assert_eq!(answer["id"], 9);
5822        assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5823        assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5824        std::fs::remove_file(path).expect("cleanup");
5825    }
5826
5827    #[test]
5828    fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5829        let event = parse_acp_notification(
5830            0,
5831            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5832        )
5833        .expect("valid JSON-RPC request");
5834        assert!(event.is_none());
5835    }
5836
5837    #[tokio::test]
5838    async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5839        let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5840        std::fs::write(
5841            &script_path,
5842            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5843        )
5844        .expect("write native test script");
5845        let mut adapter = AgyAdapter::new(
5846            0,
5847            std::env::current_dir().expect("cwd"),
5848            format!("sh {}", script_path.display()),
5849        );
5850        adapter.start().await.expect("start native adapter");
5851        assert!(adapter.next_event().await.is_some());
5852        assert!(adapter.next_event().await.is_some());
5853        adapter
5854            .send_prompt("continue".into())
5855            .await
5856            .expect("prompt");
5857        assert!(matches!(
5858            adapter.next_event().await,
5859            Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5860        ));
5861        assert!(matches!(
5862            adapter.next_event().await,
5863            Some(Ok(AgentEvent::TurnComplete { .. }))
5864        ));
5865        adapter.stop().await.expect("stop native adapter");
5866        std::fs::remove_file(script_path).expect("cleanup native script");
5867    }
5868
5869    #[test]
5870    fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5871        let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5872        let _ = std::fs::remove_dir_all(&root);
5873        std::fs::create_dir_all(&root).expect("workspace");
5874        std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5875        let outside =
5876            std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5877        std::fs::write(&outside, "secret").expect("outside file");
5878        let link = root.join("outside-link");
5879        #[cfg(unix)]
5880        std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5881        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5882
5883        assert_eq!(
5884            adapter
5885                .read_workspace_text("inside.txt", Some(2), Some(1))
5886                .expect("read inside"),
5887            "two"
5888        );
5889        std::fs::write(
5890            root.join("large.txt"),
5891            vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5892        )
5893        .expect("large file");
5894        let bounded = adapter
5895            .read_workspace_text("large.txt", None, None)
5896            .expect("bounded read");
5897        assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5898        #[cfg(unix)]
5899        {
5900            std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5901                .expect("internal symlink");
5902            assert_eq!(
5903                adapter
5904                    .read_workspace_text("inside-link", None, None)
5905                    .expect("read internal symlink"),
5906                "one\ntwo\nthree\n"
5907            );
5908        }
5909        assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5910        assert!(
5911            adapter
5912                .workspace_path(&outside.display().to_string())
5913                .is_err()
5914        );
5915        #[cfg(unix)]
5916        assert!(adapter.workspace_path("outside-link").is_err());
5917        #[cfg(unix)]
5918        std::fs::remove_file(link).expect("cleanup symlink");
5919        #[cfg(unix)]
5920        std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5921        std::fs::remove_file(outside).expect("cleanup outside");
5922        std::fs::remove_dir_all(root).expect("cleanup workspace");
5923    }
5924
5925    #[tokio::test]
5926    async fn running_terminal_output_omits_exit_status_until_completion() {
5927        let root = unique_test_path("codeswarm-terminal-output", "dir");
5928        std::fs::create_dir_all(&root).expect("workspace");
5929        let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5930        let result = adapter
5931            .terminal_create(&serde_json::json!({
5932                "command": "sh",
5933                "args": ["-c", "sleep 0.2; printf done"],
5934                "cwd": ".",
5935            }))
5936            .await
5937            .expect("terminal create");
5938        let id = result["terminalId"].as_str().expect("terminal id");
5939        let output = adapter.terminal_output(id).await.expect("terminal output");
5940        assert!(output.get("exitStatus").is_none());
5941        if let Some(terminal) = adapter.terminals.remove(id) {
5942            terminal.stop().await;
5943        }
5944        std::fs::remove_dir_all(root).expect("cleanup workspace");
5945    }
5946
5947    #[tokio::test]
5948    async fn acp_adapter_answers_workspace_read_requests() {
5949        let root =
5950            std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
5951        let _ = std::fs::remove_dir_all(&root);
5952        std::fs::create_dir_all(&root).expect("workspace");
5953        let source = root.join("inside.txt");
5954        let answer = root.join("answer.json");
5955        std::fs::write(&source, "workspace content").expect("source");
5956        let script = format!(
5957            r#"read _; echo '{{"jsonrpc":"2.0","id":1,"result":{{"agentCapabilities":{{}}}}}}'; read _; echo '{{"jsonrpc":"2.0","id":2,"result":{{"sessionId":"s1"}}}}'; read _; echo '{{"jsonrpc":"2.0","id":9,"method":"fs/read_text_file","params":{{"sessionId":"s1","path":"{}"}}}}'; read response; printf '%s' "$response" > '{}'; echo '{{"jsonrpc":"2.0","id":3,"result":{{"stopReason":"end_turn"}}}}'"#,
5958            source.display(),
5959            answer.display(),
5960        );
5961        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5962        adapter.start().await.expect("start ACP");
5963        assert!(matches!(
5964            adapter.next_event().await,
5965            Some(Ok(AgentEvent::Ready { .. }))
5966        ));
5967        adapter.send_prompt("read it".into()).await.expect("prompt");
5968        assert!(matches!(
5969            adapter.next_event().await,
5970            Some(Ok(AgentEvent::TurnComplete { .. }))
5971        ));
5972        let response: Value =
5973            serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
5974                .expect("response JSON");
5975        assert_eq!(response["id"], 9);
5976        assert_eq!(response["result"]["content"], "workspace content");
5977        adapter.stop().await.expect("stop ACP");
5978        std::fs::remove_dir_all(root).expect("cleanup workspace");
5979    }
5980
5981    #[tokio::test]
5982    async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
5983        let root =
5984            std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
5985        let _ = std::fs::remove_dir_all(&root);
5986        std::fs::create_dir_all(&root).expect("workspace");
5987        let create_request = serde_json::json!({
5988            "jsonrpc": "2.0",
5989            "id": 9,
5990            "method": "terminal/create",
5991            "params": {
5992                "sessionId": "s1",
5993                "command": "sh",
5994                "args": ["-c", "sleep 0.1; printf terminal-ok"],
5995                "cwd": ".",
5996            },
5997        });
5998        let wait_request = serde_json::json!({
5999            "jsonrpc": "2.0",
6000            "id": 10,
6001            "method": "terminal/wait_for_exit",
6002            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6003        });
6004        let output_request = serde_json::json!({
6005            "jsonrpc": "2.0",
6006            "id": 11,
6007            "method": "terminal/output",
6008            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6009        });
6010        let create_answer = root.join("create-answer.json");
6011        let wait_answer = root.join("wait-answer.json");
6012        let output_answer = root.join("output-answer.json");
6013        let script = format!(
6014            "read _; echo '{{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":{{\"agentCapabilities\":{{}}}}}}'; read _; echo '{{\"jsonrpc\":\"2.0\",\"id\":2,\"result\":{{\"sessionId\":\"s1\"}}}}'; read _; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{}'; read response; printf '%s' \"$response\" > '{}'; echo '{{\"jsonrpc\":\"2.0\",\"id\":3,\"result\":{{\"stopReason\":\"end_turn\"}}}}'",
6015            create_request,
6016            create_answer.display(),
6017            wait_request,
6018            wait_answer.display(),
6019            output_request,
6020            output_answer.display(),
6021        );
6022        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6023        adapter.start().await.expect("start ACP");
6024        assert!(matches!(
6025            adapter.next_event().await,
6026            Some(Ok(AgentEvent::Ready { .. }))
6027        ));
6028        adapter
6029            .send_prompt("run terminal".into())
6030            .await
6031            .expect("prompt");
6032        let mut saw_complete = false;
6033        for _ in 0..6 {
6034            match adapter.next_event().await {
6035                Some(Ok(AgentEvent::TurnComplete { .. })) => {
6036                    saw_complete = true;
6037                    break;
6038                }
6039                Some(_) => {}
6040                None => break,
6041            }
6042        }
6043        assert!(saw_complete, "terminal requests should not stall ACP");
6044        let create: Value = serde_json::from_str(
6045            &std::fs::read_to_string(&create_answer).expect("captured create response"),
6046        )
6047        .expect("create JSON");
6048        assert_eq!(create["result"]["terminalId"], "terminal-1");
6049        let output: Value = serde_json::from_str(
6050            &std::fs::read_to_string(&output_answer).expect("captured output response"),
6051        )
6052        .expect("output JSON");
6053        assert!(
6054            output["result"]["output"]
6055                .as_str()
6056                .unwrap_or_default()
6057                .contains("terminal-ok"),
6058            "output response: {output}"
6059        );
6060        adapter.stop().await.expect("stop ACP");
6061        std::fs::remove_dir_all(root).expect("cleanup workspace");
6062    }
6063
6064    #[tokio::test]
6065    async fn host_reduces_and_persists_adapter_events() {
6066        let path =
6067            std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6068        let adapter = ScriptedAdapter::new(
6069            0,
6070            AgentCapabilities::default(),
6071            [AgentEvent::Text {
6072                slot: 0,
6073                text: "hello".into(),
6074            }],
6075        );
6076        let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6077        host.start().await.expect("start");
6078        host.next_effects()
6079            .await
6080            .expect("event")
6081            .expect("valid event");
6082        assert_eq!(host.state.public_text[0].1, "hello");
6083        assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6084        std::fs::remove_file(path).expect("cleanup");
6085    }
6086
6087    #[tokio::test]
6088    async fn relay_applies_default_policy_before_the_first_prompt() {
6089        let first_log = Arc::new(Mutex::new(Vec::new()));
6090        let second_log = Arc::new(Mutex::new(Vec::new()));
6091        let first = AdapterHost::new(
6092            Box::new(ModeOrderAdapter {
6093                slot: 0,
6094                log: Arc::clone(&first_log),
6095                phase: 0,
6096            }),
6097            None,
6098        );
6099        let second = AdapterHost::new(
6100            Box::new(ModeOrderAdapter {
6101                slot: 1,
6102                log: Arc::clone(&second_log),
6103                phase: 0,
6104            }),
6105            None,
6106        );
6107        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6108        relay.start().await.expect("start and synchronize policy");
6109        relay.run_turn("task", 0).await.expect("first turn");
6110        {
6111            let log = first_log.lock().expect("log");
6112            assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6113        }
6114        assert_eq!(
6115            second_log.lock().expect("log").as_slice(),
6116            ["start", "mode:yolo"]
6117        );
6118        let added_log = Arc::new(Mutex::new(Vec::new()));
6119        relay
6120            .add_agent(
6121                AdapterHost::new(
6122                    Box::new(ModeOrderAdapter {
6123                        slot: 2,
6124                        log: Arc::clone(&added_log),
6125                        phase: 0,
6126                    }),
6127                    None,
6128                ),
6129                "Added",
6130                "added.example",
6131                "added-agent",
6132            )
6133            .await
6134            .expect("add with synchronized policy");
6135        assert_eq!(
6136            added_log.lock().expect("log").as_slice(),
6137            ["start", "mode:yolo"]
6138        );
6139        relay.drop_agent(2).await.expect("drop added agent");
6140        added_log.lock().expect("log").clear();
6141        relay
6142            .reload(2)
6143            .await
6144            .expect("reload with synchronized policy");
6145        assert_eq!(
6146            added_log.lock().expect("log").as_slice(),
6147            ["reload", "mode:yolo"]
6148        );
6149    }
6150
6151    #[tokio::test]
6152    async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6153        let hosts = (0..2)
6154            .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6155            .collect::<Vec<_>>();
6156        let startup_events = Arc::new(Mutex::new(Vec::new()));
6157        let captured = Arc::clone(&startup_events);
6158        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6159        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6160
6161        relay.start().await.expect("complete startup handshake");
6162
6163        assert!(relay.dispatches().is_empty());
6164        let ready_slots = startup_events
6165            .lock()
6166            .expect("events")
6167            .iter()
6168            .filter_map(|event| match event {
6169                AgentEvent::Ready { slot, .. } => Some(*slot),
6170                _ => None,
6171            })
6172            .collect::<Vec<_>>();
6173        assert_eq!(ready_slots, vec![0, 1]);
6174    }
6175
6176    #[tokio::test]
6177    async fn independent_roster_adapters_start_concurrently() {
6178        let barrier = Arc::new(tokio::sync::Barrier::new(2));
6179        let hosts = (0..2)
6180            .map(|slot| {
6181                AdapterHost::new(
6182                    Box::new(ConcurrentStartAdapter {
6183                        slot,
6184                        barrier: Arc::clone(&barrier),
6185                    }),
6186                    None,
6187                )
6188            })
6189            .collect::<Vec<_>>();
6190        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6191        tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6192            .await
6193            .expect("startup should not serialize barrier participants")
6194            .expect("startup succeeds");
6195    }
6196
6197    #[tokio::test]
6198    async fn relay_host_dispatches_turns_sequentially() {
6199        let capabilities = AgentCapabilities {
6200            supports_cancel: true,
6201            ..AgentCapabilities::default()
6202        };
6203        let first = ScriptedAdapter::new(
6204            0,
6205            capabilities.clone(),
6206            [
6207                AgentEvent::Text {
6208                    slot: 0,
6209                    text: "first".into(),
6210                },
6211                AgentEvent::TurnComplete { slot: 0 },
6212            ],
6213        );
6214        let second = ScriptedAdapter::new(
6215            1,
6216            capabilities,
6217            [
6218                AgentEvent::Text {
6219                    slot: 1,
6220                    text: "review".into(),
6221                },
6222                AgentEvent::TurnComplete { slot: 1 },
6223            ],
6224        );
6225        let hosts = vec![
6226            AdapterHost::new(Box::new(first), None),
6227            AdapterHost::new(Box::new(second), None),
6228        ];
6229        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6230        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6231        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6232        let captured = std::sync::Arc::clone(&events);
6233        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6234        relay.start().await.expect("start");
6235        events.lock().expect("events").clear();
6236        assert!(matches!(
6237            relay.run_turn("task", 0).await.expect("first turn"),
6238            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6239        ));
6240        assert!(matches!(
6241            relay.run_turn("first", 0).await.expect("second turn"),
6242            crate::relay::RelayDecision::Dispatch {
6243                slot: 1,
6244                can_stop: true,
6245                ..
6246            }
6247        ));
6248        assert_eq!(
6249            relay
6250                .dispatches()
6251                .iter()
6252                .map(|(slot, _)| *slot)
6253                .collect::<Vec<_>>(),
6254            [0, 1]
6255        );
6256        assert!(relay.dispatches()[0].1.contains("You are Claude"));
6257        assert!(
6258            relay.dispatches()[0]
6259                .1
6260                .contains("CodeSwarm roster (ordered)")
6261        );
6262        assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6263        assert!(relay.dispatches()[0].1.contains("2. Codex"));
6264        assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6265        assert!(relay.dispatches()[0].1.contains("Do not use"));
6266        let lifecycle = events.lock().expect("events");
6267        let positions = lifecycle
6268            .iter()
6269            .filter_map(|event| match event {
6270                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6271                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6272                _ => None,
6273            })
6274            .collect::<Vec<_>>();
6275        assert_eq!(
6276            positions,
6277            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6278        );
6279    }
6280
6281    #[tokio::test]
6282    async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6283        let stops = Arc::new(AtomicUsize::new(0));
6284        let failed = FailingStartAdapter {
6285            slot: 0,
6286            stopped: stops.clone(),
6287        };
6288        let healthy = ScriptedAdapter::new(
6289            1,
6290            AgentCapabilities::default(),
6291            [
6292                AgentEvent::Text {
6293                    slot: 1,
6294                    text: "healthy response".into(),
6295                },
6296                AgentEvent::TurnComplete { slot: 1 },
6297            ],
6298        );
6299        let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6300        let captured = events.clone();
6301        let mut relay = RelayHost::new(
6302            vec![
6303                AdapterHost::new(Box::new(failed), None),
6304                AdapterHost::new(Box::new(healthy), None),
6305            ],
6306            4,
6307        )
6308        .unwrap();
6309        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6310        relay.start_resuming().await.unwrap();
6311        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6312        assert!(relay.dispatches().is_empty());
6313        assert_eq!(stops.load(Ordering::Relaxed), 1);
6314        assert!(
6315            events
6316                .lock()
6317                .unwrap()
6318                .iter()
6319                .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6320        );
6321        assert!(
6322            !events
6323                .lock()
6324                .unwrap()
6325                .iter()
6326                .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6327        );
6328        assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6329        assert!(
6330            relay
6331                .relay_mut()
6332                .enqueue_human("explicit healthy target", Some(1))
6333        );
6334        assert!(matches!(
6335            relay.run_turn("", 1).await.unwrap(),
6336            RelayDecision::Dispatch { slot: 1, .. }
6337        ));
6338        relay.stop().await.unwrap();
6339    }
6340
6341    #[tokio::test]
6342    async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6343        let capabilities = AgentCapabilities::default();
6344        let first = ScriptedAdapter::new(
6345            0,
6346            capabilities.clone(),
6347            [
6348                AgentEvent::Text {
6349                    slot: 0,
6350                    text: "implemented".into(),
6351                },
6352                AgentEvent::TurnComplete { slot: 0 },
6353                AgentEvent::Text {
6354                    slot: 0,
6355                    text: format!("fixed review findings {STOP_TOKEN}"),
6356                },
6357                AgentEvent::TurnComplete { slot: 0 },
6358            ],
6359        );
6360        let second = ScriptedAdapter::new(
6361            1,
6362            capabilities,
6363            [
6364                AgentEvent::Text {
6365                    slot: 1,
6366                    text: "reviewed".into(),
6367                },
6368                AgentEvent::TurnComplete { slot: 1 },
6369                AgentEvent::Text {
6370                    slot: 1,
6371                    text: format!("approved {STOP_TOKEN}"),
6372                },
6373                AgentEvent::TurnComplete { slot: 1 },
6374                AgentEvent::Text {
6375                    slot: 1,
6376                    text: "new task".into(),
6377                },
6378                AgentEvent::TurnComplete { slot: 1 },
6379            ],
6380        );
6381        let hosts = vec![
6382            AdapterHost::new(Box::new(first), None),
6383            AdapterHost::new(Box::new(second), None),
6384        ];
6385        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6386        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6387        relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6388        relay.start().await.expect("start");
6389        assert!(matches!(
6390            relay.run_turn("task", 0).await.expect("implementer turn"),
6391            RelayDecision::Dispatch {
6392                slot: 0,
6393                can_stop: false,
6394                ..
6395            }
6396        ));
6397        let implementer_prompt = &relay.dispatches()[0].1;
6398        assert!(implementer_prompt.contains("you are the implementer"));
6399        assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6400        assert!(!implementer_prompt.contains("you are the reviewer"));
6401        assert!(implementer_prompt.contains("Do not use"));
6402        assert!(matches!(
6403            relay.run_turn("", 0).await.expect("reviewer turn"),
6404            RelayDecision::Dispatch {
6405                slot: 1,
6406                can_stop: true,
6407                ..
6408            }
6409        ));
6410        let reviewer_prompt = &relay.dispatches()[1].1;
6411        assert!(reviewer_prompt.contains("you are the reviewer"));
6412        assert!(reviewer_prompt.contains("Claude handed off"));
6413        assert!(reviewer_prompt.contains("concrete defects"));
6414        assert!(reviewer_prompt.contains("concise approval"));
6415        assert!(reviewer_prompt.contains(STOP_TOKEN));
6416        assert!(!reviewer_prompt.contains("you are the implementer"));
6417        relay.run_turn("", 0).await.unwrap();
6418        assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6419        assert!(relay.dispatches()[2].1.contains("Do not use"));
6420        assert!(matches!(
6421            relay.run_turn("", 0).await.unwrap(),
6422            RelayDecision::Dispatch { slot: 1, .. }
6423        ));
6424        assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6425        assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6426        relay.run_turn("", 1).await.unwrap();
6427        assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6428    }
6429
6430    #[tokio::test]
6431    async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6432        let solo = ScriptedAdapter::new(
6433            0,
6434            AgentCapabilities::default(),
6435            [
6436                AgentEvent::Text {
6437                    slot: 0,
6438                    text: "solo".into(),
6439                },
6440                AgentEvent::TurnComplete { slot: 0 },
6441            ],
6442        );
6443        let mut solo_relay =
6444            RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6445        solo_relay
6446            .relay_mut()
6447            .set_strategy(CollaborationStrategy::Pair);
6448        solo_relay.start().await.expect("start");
6449        solo_relay.run_turn("task", 0).await.expect("solo turn");
6450        assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6451
6452        let roster_first = ScriptedAdapter::new(
6453            0,
6454            AgentCapabilities::default(),
6455            [AgentEvent::TurnComplete { slot: 0 }],
6456        );
6457        let roster_second = ScriptedAdapter::new(
6458            1,
6459            AgentCapabilities::default(),
6460            [AgentEvent::TurnComplete { slot: 1 }],
6461        );
6462        let mut roster = RelayHost::new(
6463            vec![
6464                AdapterHost::new(Box::new(roster_first), None),
6465                AdapterHost::new(Box::new(roster_second), None),
6466            ],
6467            4,
6468        )
6469        .expect("relay");
6470        roster.start().await.expect("start");
6471        roster.run_turn("task", 0).await.expect("first turn");
6472        roster.run_turn("", 0).await.expect("second turn");
6473        assert!(!roster.dispatches()[0].1.contains("Pair role"));
6474        assert!(!roster.dispatches()[1].1.contains("Pair role"));
6475
6476        let pair_first = ScriptedAdapter::new(
6477            0,
6478            AgentCapabilities::default(),
6479            [AgentEvent::TurnComplete { slot: 0 }],
6480        );
6481        let pair_second = ScriptedAdapter::new(
6482            1,
6483            AgentCapabilities::default(),
6484            [AgentEvent::TurnComplete { slot: 1 }],
6485        );
6486        let mut pair = RelayHost::new(
6487            vec![
6488                AdapterHost::new(Box::new(pair_first), None),
6489                AdapterHost::new(Box::new(pair_second), None),
6490            ],
6491            4,
6492        )
6493        .expect("relay");
6494        pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6495        assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6496        pair.start().await.expect("start");
6497        assert!(matches!(
6498            pair.run_turn("ignored", 0).await.expect("direct turn"),
6499            RelayDecision::Dispatch {
6500                slot: 1,
6501                direct: true,
6502                ..
6503            }
6504        ));
6505        let direct_prompt = &pair.dispatches()[0].1;
6506        assert!(direct_prompt.contains("private"));
6507        assert!(!direct_prompt.contains("Pair role"));
6508    }
6509
6510    #[tokio::test]
6511    async fn relay_host_routes_around_a_usage_limited_agent() {
6512        let capabilities = AgentCapabilities::default();
6513        let first = ScriptedAdapter::new(
6514            0,
6515            capabilities.clone(),
6516            [
6517                AgentEvent::Text {
6518                    slot: 0,
6519                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6520                           credits or try again later."
6521                        .into(),
6522                },
6523                AgentEvent::TurnComplete { slot: 0 },
6524            ],
6525        );
6526        let second = ScriptedAdapter::new(
6527            1,
6528            capabilities,
6529            [
6530                AgentEvent::Text {
6531                    slot: 1,
6532                    text: "review done".into(),
6533                },
6534                AgentEvent::TurnComplete { slot: 1 },
6535            ],
6536        );
6537        let hosts = vec![
6538            AdapterHost::new(Box::new(first), None),
6539            AdapterHost::new(Box::new(second), None),
6540        ];
6541        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6542        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6543        let captured = std::sync::Arc::clone(&events);
6544        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6545        relay.start().await.expect("start");
6546        events.lock().expect("events").clear();
6547        assert!(matches!(
6548            relay.run_turn("task", 0).await.expect("limited turn"),
6549            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6550        ));
6551        assert!(
6552            events
6553                .lock()
6554                .expect("events")
6555                .iter()
6556                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6557        );
6558        // The next automatic turn skips the limited agent entirely.
6559        assert!(matches!(
6560            relay.run_turn("", 0).await.expect("next turn"),
6561            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6562        ));
6563        assert!(relay.relay().is_limited(0));
6564        // A reload restores the agent to the ring. (ScriptedAdapter cannot
6565        // feed further turns, so the restored routing itself is covered by
6566        // the relay unit tests.)
6567        relay.reload(0).await.expect("reload");
6568        assert!(!relay.relay().is_limited(0));
6569    }
6570
6571    #[tokio::test]
6572    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6573        let limited = ScriptedAdapter::new(
6574            0,
6575            AgentCapabilities::default(),
6576            [AgentEvent::Failed {
6577                slot: 0,
6578                started: true,
6579                detail: "request failed: insufficient_quota".into(),
6580            }],
6581        );
6582        let healthy = ScriptedAdapter::new(
6583            1,
6584            AgentCapabilities::default(),
6585            [AgentEvent::TurnComplete { slot: 1 }],
6586        );
6587        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6588        let captured = std::sync::Arc::clone(&events);
6589        let mut relay = RelayHost::new(
6590            vec![
6591                AdapterHost::new(Box::new(limited), None),
6592                AdapterHost::new(Box::new(healthy), None),
6593            ],
6594            4,
6595        )
6596        .expect("relay");
6597        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6598        relay.start().await.expect("start");
6599        events.lock().expect("events").clear();
6600
6601        assert!(matches!(
6602            relay.run_turn("task", 0).await.expect("limited failure"),
6603            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6604        ));
6605        assert!(relay.relay().is_limited(0));
6606        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6607        {
6608            let events = events.lock().expect("events");
6609            assert!(
6610                events
6611                    .iter()
6612                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6613            );
6614            assert!(
6615                !events
6616                    .iter()
6617                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
6618            );
6619        }
6620
6621        assert!(matches!(
6622            relay.run_turn("", 0).await.expect("healthy peer"),
6623            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6624        ));
6625    }
6626
6627    #[tokio::test]
6628    async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6629        let failed = ScriptedAdapter::new(
6630            0,
6631            AgentCapabilities::default(),
6632            [AgentEvent::Failed {
6633                slot: 0,
6634                started: true,
6635                detail: "connection lost".into(),
6636            }],
6637        );
6638        let healthy = ScriptedAdapter::new(
6639            1,
6640            AgentCapabilities::default(),
6641            [AgentEvent::TurnComplete { slot: 1 }],
6642        );
6643        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6644        let captured = std::sync::Arc::clone(&events);
6645        let mut relay = RelayHost::new(
6646            vec![
6647                AdapterHost::new(Box::new(failed), None),
6648                AdapterHost::new(Box::new(healthy), None),
6649            ],
6650            4,
6651        )
6652        .expect("relay");
6653        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6654        relay.start().await.expect("start");
6655
6656        assert!(matches!(
6657            relay.run_turn("task", 0).await.expect("handled failure"),
6658            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6659        ));
6660        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6661        assert!(relay.relay().is_limited(0));
6662        assert!(events.lock().expect("lock").iter().any(|event| {
6663            matches!(
6664                event,
6665                AgentEvent::Failed {
6666                    slot: 0,
6667                    started: true,
6668                    ..
6669                }
6670            )
6671        }));
6672        assert!(matches!(
6673            relay.run_turn("", 0).await.expect("healthy peer"),
6674            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6675        ));
6676    }
6677
6678    #[tokio::test]
6679    async fn codex_stop_does_not_skip_later_roster_reviewers() {
6680        let hosts = (0..3)
6681            .map(|slot| {
6682                AdapterHost::new(
6683                    Box::new(ScriptedAdapter::new(
6684                        slot,
6685                        AgentCapabilities::default(),
6686                        [
6687                            AgentEvent::Text {
6688                                slot,
6689                                text: STOP_TOKEN.into(),
6690                            },
6691                            AgentEvent::TurnComplete { slot },
6692                        ],
6693                    )),
6694                    None,
6695                )
6696            })
6697            .collect();
6698        let mut relay = RelayHost::new(hosts, 10).expect("relay");
6699        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6700        relay.start().await.expect("start");
6701        for expected in 0..3 {
6702            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6703                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6704        }
6705        assert_eq!(
6706            relay.run_turn("", 0).await.expect("complete"),
6707            RelayDecision::Complete
6708        );
6709    }
6710
6711    #[tokio::test]
6712    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6713        let first = ScriptedAdapter::new(
6714            0,
6715            AgentCapabilities::default(),
6716            [
6717                AgentEvent::Text {
6718                    slot: 0,
6719                    text: "done".into(),
6720                },
6721                AgentEvent::TurnComplete { slot: 0 },
6722            ],
6723        );
6724        let reviewer = ScriptedAdapter::new(
6725            1,
6726            AgentCapabilities::default(),
6727            [
6728                AgentEvent::Text {
6729                    slot: 1,
6730                    text: STOP_TOKEN.into(),
6731                },
6732                AgentEvent::TurnComplete { slot: 1 },
6733            ],
6734        );
6735        let mut relay = RelayHost::new(
6736            vec![
6737                AdapterHost::new(Box::new(first), None),
6738                AdapterHost::new(Box::new(reviewer), None),
6739            ],
6740            10,
6741        )
6742        .expect("relay");
6743        relay.start().await.expect("start");
6744        let first_decision = relay.run_turn("task", 0).await.expect("first");
6745        assert!(matches!(
6746            first_decision,
6747            RelayDecision::Dispatch { slot: 0, .. }
6748        ));
6749        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6750        assert!(matches!(
6751            reviewer_decision,
6752            RelayDecision::Dispatch {
6753                slot: 1,
6754                can_stop: true,
6755                ..
6756            }
6757        ));
6758        assert_eq!(
6759            relay.run_turn("", 0).await.expect("complete"),
6760            RelayDecision::Complete
6761        );
6762    }
6763
6764    #[tokio::test]
6765    async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6766        let tool = AgentEvent::Tool {
6767            slot: 0,
6768            update: crate::ToolUpdate {
6769                id: "read".into(),
6770                title: "Read file".into(),
6771                status: ToolStatus::Running,
6772                detail: None,
6773            },
6774        };
6775        let updates = vec![
6776            AgentEvent::Thought {
6777                slot: 0,
6778                text: "Check the buffer. ✈".into(),
6779            },
6780            AgentEvent::Text {
6781                slot: 0,
6782                text: "Let me check.".into(),
6783            },
6784            tool.clone(),
6785            AgentEvent::Text {
6786                slot: 0,
6787                text: "[CODE".into(),
6788            },
6789            AgentEvent::Text {
6790                slot: 0,
6791                text: " is ordinary.".into(),
6792            },
6793            AgentEvent::Text {
6794                slot: 0,
6795                text: "[CODESWARM:".into(),
6796            },
6797            AgentEvent::Text {
6798                slot: 0,
6799                text: "STOP] Done.".into(),
6800            },
6801            tool,
6802            AgentEvent::TurnComplete { slot: 0 },
6803        ];
6804        let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6805        let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6806        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6807        let captured = std::sync::Arc::clone(&events);
6808        let mut relay = RelayHost::new(
6809            vec![
6810                AdapterHost::new(Box::new(first), None),
6811                AdapterHost::new(Box::new(reviewer), None),
6812            ],
6813            2,
6814        )
6815        .expect("relay");
6816        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6817        relay.start().await.unwrap();
6818        relay.run_turn("task", 0).await.unwrap();
6819        let captured = events.lock().unwrap();
6820        let visible: Vec<_> = captured
6821            .iter()
6822            .filter(|event| {
6823                matches!(
6824                    event,
6825                    AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
6826                )
6827            })
6828            .cloned()
6829            .collect();
6830        assert_eq!(
6831            visible,
6832            vec![
6833                updates[0].clone(),
6834                updates[1].clone(),
6835                updates[2].clone(),
6836                AgentEvent::Text {
6837                    slot: 0,
6838                    text: "[CODE is ordinary.".into()
6839                },
6840                AgentEvent::Text {
6841                    slot: 0,
6842                    text: " Done.".into()
6843                },
6844                updates[7].clone(),
6845            ]
6846        );
6847    }
6848
6849    #[tokio::test]
6850    async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
6851        let text = |value: &str| AgentEvent::Text {
6852            slot: 0,
6853            text: value.into(),
6854        };
6855        let thought = || AgentEvent::Thought {
6856            slot: 0,
6857            text: "still checking".into(),
6858        };
6859        let tool = || AgentEvent::Tool {
6860            slot: 0,
6861            update: crate::ToolUpdate {
6862                id: "read".into(),
6863                title: "Read file".into(),
6864                status: ToolStatus::Running,
6865                detail: None,
6866            },
6867        };
6868        let cases = vec![
6869            (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
6870            (
6871                vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
6872                2,
6873            ),
6874            (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
6875            (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
6876            (
6877                vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
6878                1,
6879            ),
6880            (
6881                vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
6882                1,
6883            ),
6884            (
6885                vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
6886                2,
6887            ),
6888            (vec![text("result [CODESWARM:NEXT:1]")], 1),
6889            (vec![text("result [CODESWARM:NEXT:0]")], 1),
6890            (vec![text("result [CODESWARM:NEXT:99]")], 1),
6891            (
6892                vec![
6893                    text("result [CODESWARM:NEXT:3]"),
6894                    AgentEvent::UsageUpdated {
6895                        slot: 0,
6896                        usage: crate::UsageUpdate { used: 1, size: 100 },
6897                    },
6898                ],
6899                2,
6900            ),
6901        ];
6902        for (mut updates, expected) in cases {
6903            updates.push(AgentEvent::TurnComplete { slot: 0 });
6904            let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
6905            let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
6906                .chain((1..3).map(|slot| {
6907                    AdapterHost::new(
6908                        Box::new(ScriptedAdapter::new(
6909                            slot,
6910                            AgentCapabilities::default(),
6911                            [AgentEvent::TurnComplete { slot }],
6912                        )),
6913                        None,
6914                    )
6915                }))
6916                .collect();
6917            let mut relay = RelayHost::new(hosts, 10).unwrap();
6918            relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
6919            let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6920            let captured = events.clone();
6921            relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6922            relay.start().await.unwrap();
6923            relay.run_turn("task", 0).await.unwrap();
6924            let prompt = &relay.dispatches()[0].1;
6925            assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
6926            assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
6927            assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
6928            // The new names must appear even when introduction has already run.
6929            relay.introduced.fill(true);
6930            relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
6931            let next = relay.run_turn("", 0).await.unwrap();
6932            assert!(
6933                matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
6934                "{next:?}"
6935            );
6936            let prompt = &relay.dispatches()[1].1;
6937            assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
6938            assert!(prompt.contains("result"));
6939            let public = prompt
6940                .split("Public updates:\n")
6941                .nth(1)
6942                .unwrap()
6943                .split("\n\nDo not use")
6944                .next()
6945                .unwrap();
6946            assert!(!public.contains("[CODESWARM:NEXT:"));
6947            let visible = events
6948                .lock()
6949                .unwrap()
6950                .iter()
6951                .filter_map(|event| match event {
6952                    AgentEvent::Text { text, .. } => Some(text.clone()),
6953                    _ => None,
6954                })
6955                .collect::<String>();
6956            assert!(visible.contains("result"));
6957            assert!(!visible.contains("[CODESWARM:"), "{visible}");
6958        }
6959    }
6960
6961    #[tokio::test]
6962    async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
6963        let text = |value: &str| AgentEvent::Text {
6964            slot: 1,
6965            text: value.into(),
6966        };
6967        let thought = || AgentEvent::Thought {
6968            slot: 1,
6969            text: "still checking".into(),
6970        };
6971        let tool = || AgentEvent::Tool {
6972            slot: 1,
6973            update: crate::ToolUpdate {
6974                id: "read".into(),
6975                title: "Read file".into(),
6976                status: ToolStatus::Running,
6977                detail: None,
6978            },
6979        };
6980        let cases = vec![
6981            (vec![text(&format!("done {STOP_TOKEN}"))], true),
6982            (vec![text(STOP_TOKEN), text("\n  ")], true),
6983            (vec![text(STOP_TOKEN), text(" actually keep going")], false),
6984            (vec![text(STOP_TOKEN), thought()], false),
6985            (vec![text(STOP_TOKEN), tool()], false),
6986            (vec![text(STOP_TOKEN), tool(), text(" ")], false),
6987            (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
6988            (vec![text("[CODESWARM:"), text("STOP]")], true),
6989            (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
6990            (
6991                vec![AgentEvent::Thought {
6992                    slot: 1,
6993                    text: STOP_TOKEN.into(),
6994                }],
6995                false,
6996            ),
6997            (
6998                vec![
6999                    text(STOP_TOKEN),
7000                    AgentEvent::UsageUpdated {
7001                        slot: 1,
7002                        usage: crate::UsageUpdate { used: 1, size: 100 },
7003                    },
7004                ],
7005                true,
7006            ),
7007        ];
7008        for (mut events, stop) in cases {
7009            let first = ScriptedAdapter::new(
7010                0,
7011                AgentCapabilities::default(),
7012                [
7013                    AgentEvent::Text {
7014                        slot: 0,
7015                        text: "initial response".into(),
7016                    },
7017                    AgentEvent::TurnComplete { slot: 0 },
7018                    AgentEvent::TurnComplete { slot: 0 },
7019                ],
7020            );
7021            events.push(AgentEvent::TurnComplete { slot: 1 });
7022            let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7023            let mut relay = RelayHost::new(
7024                vec![
7025                    AdapterHost::new(Box::new(first), None),
7026                    AdapterHost::new(Box::new(reviewer), None),
7027                ],
7028                4,
7029            )
7030            .unwrap();
7031            relay.start().await.unwrap();
7032            relay.run_turn("task", 0).await.unwrap();
7033            relay.run_turn("", 0).await.unwrap();
7034            let next = relay.run_turn("", 0).await.unwrap();
7035            assert_eq!(
7036                matches!(next, RelayDecision::Complete),
7037                stop,
7038                "events={events:?}"
7039            );
7040            relay.stop().await.unwrap();
7041        }
7042    }
7043
7044    #[tokio::test]
7045    async fn stop_token_is_filtered_from_streamed_ui_events() {
7046        let first = ScriptedAdapter::new(
7047            0,
7048            AgentCapabilities::default(),
7049            [
7050                AgentEvent::Text {
7051                    slot: 0,
7052                    text: format!("visible {STOP_TOKEN} trailing"),
7053                },
7054                AgentEvent::TurnComplete { slot: 0 },
7055            ],
7056        );
7057        let reviewer = ScriptedAdapter::new(
7058            1,
7059            AgentCapabilities::default(),
7060            [AgentEvent::TurnComplete { slot: 1 }],
7061        );
7062        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7063        let captured = std::sync::Arc::clone(&events);
7064        let mut relay = RelayHost::new(
7065            vec![
7066                AdapterHost::new(Box::new(first), None),
7067                AdapterHost::new(Box::new(reviewer), None),
7068            ],
7069            2,
7070        )
7071        .expect("relay");
7072        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7073        relay.start().await.expect("start");
7074        relay.run_turn("task", 0).await.expect("turn");
7075        let captured = events.lock().expect("lock");
7076        assert!(captured.iter().all(|event| match event {
7077            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7078            _ => true,
7079        }));
7080        let visible = captured
7081            .iter()
7082            .filter_map(|event| match event {
7083                AgentEvent::Text { text, .. } => Some(text.as_str()),
7084                _ => None,
7085            })
7086            .collect::<String>();
7087        assert_eq!(visible, "visible  trailing");
7088    }
7089
7090    #[tokio::test]
7091    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7092        let first = ScriptedAdapter::new(
7093            0,
7094            AgentCapabilities::default(),
7095            [
7096                AgentEvent::Text {
7097                    slot: 0,
7098                    text: "done".into(),
7099                },
7100                AgentEvent::TurnComplete { slot: 0 },
7101            ],
7102        );
7103        let reviewer = ScriptedAdapter::new(
7104            1,
7105            AgentCapabilities::default(),
7106            [
7107                AgentEvent::Text {
7108                    slot: 1,
7109                    text: STOP_TOKEN.into(),
7110                },
7111                AgentEvent::TurnComplete { slot: 1 },
7112            ],
7113        );
7114        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7115        let captured = std::sync::Arc::clone(&events);
7116        let mut relay = RelayHost::new(
7117            vec![
7118                AdapterHost::new(Box::new(first), None),
7119                AdapterHost::new(Box::new(reviewer), None),
7120            ],
7121            4,
7122        )
7123        .expect("relay");
7124        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7125        relay.start().await.expect("start");
7126        relay.run_turn("task", 0).await.expect("first turn");
7127        relay.run_turn("", 0).await.expect("review turn");
7128        let captured = events.lock().expect("lock");
7129        assert!(captured.iter().any(|event| {
7130            matches!(
7131                event,
7132                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7133            )
7134        }));
7135        assert!(captured.iter().all(|event| match event {
7136            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7137            _ => true,
7138        }));
7139        let acknowledgment = captured
7140            .iter()
7141            .position(|event| {
7142                matches!(
7143                    event,
7144                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7145                )
7146            })
7147            .expect("visible acknowledgment");
7148        let completion = captured
7149            .iter()
7150            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7151            .expect("reviewer completion");
7152        assert!(acknowledgment < completion);
7153    }
7154
7155    #[tokio::test]
7156    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7157        let first = ScriptedAdapter::new(
7158            0,
7159            AgentCapabilities::default(),
7160            [
7161                AgentEvent::Text {
7162                    slot: 0,
7163                    text: "done".into(),
7164                },
7165                AgentEvent::TurnComplete { slot: 0 },
7166            ],
7167        );
7168        let reviewer = ScriptedAdapter::new(
7169            1,
7170            AgentCapabilities::default(),
7171            [
7172                AgentEvent::Text {
7173                    slot: 1,
7174                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7175                },
7176                AgentEvent::TurnComplete { slot: 1 },
7177            ],
7178        );
7179        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7180        let captured = std::sync::Arc::clone(&events);
7181        let mut relay = RelayHost::new(
7182            vec![
7183                AdapterHost::new(Box::new(first), None),
7184                AdapterHost::new(Box::new(reviewer), None),
7185            ],
7186            4,
7187        )
7188        .expect("relay");
7189        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7190        relay.start().await.expect("start");
7191        relay.run_turn("task", 0).await.expect("first turn");
7192        relay.run_turn("", 0).await.expect("review turn");
7193
7194        let visible = events
7195            .lock()
7196            .expect("lock")
7197            .iter()
7198            .filter_map(|event| match event {
7199                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7200                _ => None,
7201            })
7202            .collect::<String>();
7203        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7204        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7205    }
7206
7207    #[tokio::test]
7208    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7209        let first = AdapterHost::new(
7210            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7211            None,
7212        );
7213        let second = AdapterHost::new(
7214            Box::new(ScriptedAdapter::new(
7215                1,
7216                AgentCapabilities::default(),
7217                [AgentEvent::TurnComplete { slot: 1 }],
7218            )),
7219            None,
7220        );
7221        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7222        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7223        relay.set_event_sink(move |event| {
7224            if matches!(event, AgentEvent::Permission { .. }) {
7225                let _ = seen_sender.send(());
7226            }
7227        });
7228        relay.start().await.expect("start");
7229        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7230        let answer = async move {
7231            seen_receiver.recv().await.expect("permission request");
7232            sender
7233                .send(super::RelayPermissionAnswer {
7234                    slot: 0,
7235                    request_id: "permission-1".into(),
7236                    answer: PermissionAnswer::Selected {
7237                        option_id: "allow".into(),
7238                    },
7239                })
7240                .expect("queue permission answer");
7241        };
7242        tokio::time::timeout(std::time::Duration::from_millis(100), async {
7243            let ((), result) = tokio::join!(
7244                answer,
7245                relay.run_turn_with_permissions("task", 0, &mut receiver)
7246            );
7247            result
7248        })
7249        .await
7250        .expect("permission-gated turn should not deadlock")
7251        .expect("turn completes");
7252    }
7253
7254    #[tokio::test]
7255    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7256        let first = AdapterHost::new(
7257            Box::new(PendingAdapter {
7258                slot: 0,
7259                hang_on_cancel: false,
7260            }),
7261            None,
7262        );
7263        let second = AdapterHost::new(
7264            Box::new(ScriptedAdapter::new(
7265                1,
7266                AgentCapabilities::default(),
7267                [AgentEvent::TurnComplete { slot: 1 }],
7268            )),
7269            None,
7270        );
7271        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7272        relay.start().await.expect("start");
7273        let cancellation = relay.cancellation();
7274        let error = {
7275            let turn = relay.run_turn("task", 0);
7276            tokio::pin!(turn);
7277            cancellation.request();
7278            turn.await.expect_err("cancellation should stop turn")
7279        };
7280        assert!(error.to_string().contains("relay turn cancelled"));
7281
7282        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7283        relay
7284            .run_turn("", 1)
7285            .await
7286            .expect("replacement job reaches the selected peer");
7287        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7288        assert!(replacement.contains("replacement job"));
7289        assert!(replacement.contains("User "));
7290        assert!(replacement.contains(":\ntask"));
7291        let owner_updates = relay.relay_mut().unseen_context(0);
7292        assert!(owner_updates.contains("User "));
7293        assert!(owner_updates.contains(":\ntask"));
7294        assert!(owner_updates.contains(":\nreplacement job"));
7295    }
7296
7297    #[tokio::test]
7298    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7299        let first = AdapterHost::new(
7300            Box::new(PendingAdapter {
7301                slot: 0,
7302                hang_on_cancel: true,
7303            }),
7304            None,
7305        );
7306        let second = AdapterHost::new(
7307            Box::new(PendingAdapter {
7308                slot: 1,
7309                hang_on_cancel: false,
7310            }),
7311            None,
7312        );
7313        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7314        relay.start().await.expect("start");
7315        let cancellation = relay.cancellation();
7316        let turn = relay.run_turn("task", 0);
7317        tokio::pin!(turn);
7318        cancellation.request();
7319        let error = turn.await.expect_err("cancellation should stop turn");
7320        assert!(error.to_string().contains("timed out"));
7321    }
7322
7323    #[tokio::test]
7324    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7325        let event = [AgentEvent::TurnComplete { slot: 0 }];
7326        let first = AdapterHost::new(
7327            Box::new(ScriptedAdapter::new(
7328                0,
7329                AgentCapabilities::default(),
7330                event.clone(),
7331            )),
7332            None,
7333        );
7334        let second = AdapterHost::new(
7335            Box::new(ScriptedAdapter::new(
7336                1,
7337                AgentCapabilities::default(),
7338                [AgentEvent::TurnComplete { slot: 1 }],
7339            )),
7340            None,
7341        );
7342        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7343        relay.start().await.expect("start");
7344
7345        relay.pause();
7346        assert_eq!(
7347            relay.run_turn("paused", 0).await.expect("paused turn"),
7348            crate::relay::RelayDecision::Paused
7349        );
7350        assert!(relay.dispatches().is_empty());
7351
7352        relay.resume();
7353        relay.relay_mut().drop_agent(1).expect("drop reviewer");
7354        assert!(matches!(
7355            relay
7356                .run_turn("solo follow-up", 0)
7357                .await
7358                .expect("solo turn"),
7359            crate::relay::RelayDecision::Dispatch {
7360                slot: 0,
7361                can_stop: false,
7362                ..
7363            }
7364        ));
7365        assert_eq!(relay.dispatches().len(), 1);
7366    }
7367
7368    #[tokio::test]
7369    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7370        let first = AdapterHost::new(
7371            Box::new(ScriptedAdapter::new(
7372                0,
7373                AgentCapabilities::default(),
7374                [AgentEvent::TurnComplete { slot: 0 }],
7375            )),
7376            None,
7377        );
7378        let second = AdapterHost::new(
7379            Box::new(ScriptedAdapter::new(
7380                1,
7381                AgentCapabilities::default(),
7382                [AgentEvent::TurnComplete { slot: 1 }],
7383            )),
7384            None,
7385        );
7386        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7387        relay.set_roster_names(vec!["First".into(), "Second".into()]);
7388        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7389        relay.set_roster_launch_specs(vec![
7390            ("custom".into(), "owner".into()),
7391            ("custom".into(), "peer".into()),
7392        ]);
7393        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7394        let captured = std::sync::Arc::clone(&events);
7395        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7396        relay.start().await.expect("start");
7397        let slot = relay
7398            .add_agent(
7399                AdapterHost::new(
7400                    Box::new(ScriptedAdapter::new(
7401                        2,
7402                        AgentCapabilities::default(),
7403                        [AgentEvent::TurnComplete { slot: 2 }],
7404                    )),
7405                    None,
7406                ),
7407                "Reviewer",
7408                "reviewer.example",
7409                "reviewer --acp",
7410            )
7411            .await
7412            .expect("append agent");
7413        assert_eq!(slot, 2);
7414        assert_eq!(
7415            relay.relay().active_slots().collect::<Vec<_>>(),
7416            vec![0, 1, 2]
7417        );
7418        assert_eq!(
7419            relay
7420                .session_metadata()
7421                .get("agents")
7422                .and_then(|value| value.as_array())
7423                .map(Vec::len),
7424            Some(3)
7425        );
7426        relay.drop_agent(1).await.expect("drop middle peer");
7427        let metadata = relay.session_metadata();
7428        assert_eq!(
7429            metadata.get("agents"),
7430            Some(&serde_json::json!([
7431                {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7432                {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7433            ]))
7434        );
7435        assert!(
7436            events
7437                .lock()
7438                .expect("lock")
7439                .iter()
7440                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7441        );
7442    }
7443
7444    #[tokio::test]
7445    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7446        let path = unique_test_path("codeswarm-session-metadata", "json");
7447        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7448        let writer = metadata_store.buffered().expect("metadata writer");
7449        let first = AdapterHost::new(
7450            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7451            None,
7452        );
7453        let second = AdapterHost::new(
7454            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7455            None,
7456        );
7457        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7458        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7459        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7460        relay.set_roster_launch_specs(vec![
7461            ("custom".into(), "claude".into()),
7462            ("custom".into(), "codex".into()),
7463        ]);
7464        relay.set_session_metadata_writer(writer);
7465        relay.start().await.expect("start");
7466        relay.drop_agent(0).await.expect("drop first agent");
7467        relay.stop().await.expect("stop");
7468
7469        let loaded = metadata_store
7470            .read()
7471            .expect("read metadata")
7472            .expect("metadata snapshot");
7473        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7474        assert_eq!(
7475            loaded.get("agents"),
7476            Some(&serde_json::json!([{
7477                "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7478                "command": "codex", "supports_load_session": false
7479            }]))
7480        );
7481        assert!(loaded.get("owner").is_none());
7482        let _ = std::fs::remove_file(path);
7483    }
7484
7485    #[tokio::test]
7486    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7487        let first = AdapterHost::new(
7488            Box::new(ScriptedAdapter::new(
7489                0,
7490                AgentCapabilities::default(),
7491                [
7492                    AgentEvent::Text {
7493                        slot: 0,
7494                        text: "owner stream".into(),
7495                    },
7496                    AgentEvent::TurnComplete { slot: 0 },
7497                ],
7498            )),
7499            None,
7500        );
7501        let second = AdapterHost::new(
7502            Box::new(ScriptedAdapter::new(
7503                1,
7504                AgentCapabilities::default(),
7505                [
7506                    AgentEvent::Text {
7507                        slot: 1,
7508                        text: "peer stream".into(),
7509                    },
7510                    AgentEvent::TurnComplete { slot: 1 },
7511                ],
7512            )),
7513            None,
7514        );
7515        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7516        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7517        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7518        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7519        let captured = std::sync::Arc::clone(&events);
7520        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7521        relay.start().await.expect("start");
7522
7523        relay.swap_agents(0, 1).expect("swap peers");
7524        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7525        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7526        relay.run_turn("task", 0).await.expect("swapped turn");
7527        let events = events.lock().expect("events");
7528        assert!(events.iter().any(|event| {
7529            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7530        }));
7531        assert!(relay.dispatches()[0].1.contains("You are Peer"));
7532    }
7533
7534    #[tokio::test]
7535    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7536        let path = unique_test_path("codeswarm-session-metadata", "json");
7537        let first = AdapterHost::new(
7538            Box::new(ScriptedAdapter::new(
7539                0,
7540                AgentCapabilities::default(),
7541                [AgentEvent::TurnComplete { slot: 0 }],
7542            )),
7543            None,
7544        );
7545        let second = AdapterHost::new(
7546            Box::new(ScriptedAdapter::new(
7547                1,
7548                AgentCapabilities::default(),
7549                [AgentEvent::TurnComplete { slot: 1 }],
7550            )),
7551            None,
7552        );
7553        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7554        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7555        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7556        relay.set_roster_launch_specs(vec![
7557            ("custom".into(), "claude".into()),
7558            ("custom".into(), "codex".into()),
7559        ]);
7560        let writer = SessionMetadataStore::open(&path)
7561            .buffered()
7562            .expect("metadata writer");
7563        relay.set_session_metadata_writer(writer);
7564        relay.start().await.expect("start");
7565        relay.stop().await.expect("stop");
7566        let loaded = SessionMetadataStore::open(&path)
7567            .read()
7568            .expect("read metadata")
7569            .expect("metadata snapshot");
7570        let agents = loaded
7571            .get("agents")
7572            .and_then(|value| value.as_array())
7573            .expect("agents");
7574        assert_eq!(agents.len(), 2);
7575        assert_eq!(agents[0]["identity"], "claude.com");
7576        assert_eq!(agents[1]["identity"], "openai.com");
7577        let _ = std::fs::remove_file(path);
7578    }
7579
7580    #[tokio::test]
7581    async fn relay_host_routes_unseen_public_context_to_next_agent() {
7582        let first = AdapterHost::new(
7583            Box::new(ScriptedAdapter::new(
7584                0,
7585                AgentCapabilities::default(),
7586                [
7587                    AgentEvent::Text {
7588                        slot: 0,
7589                        text: "implemented the fix".into(),
7590                    },
7591                    AgentEvent::TurnComplete { slot: 0 },
7592                ],
7593            )),
7594            None,
7595        );
7596        let second = AdapterHost::new(
7597            Box::new(ScriptedAdapter::new(
7598                1,
7599                AgentCapabilities::default(),
7600                [AgentEvent::TurnComplete { slot: 1 }],
7601            )),
7602            None,
7603        );
7604        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7605        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7606        relay.start().await.expect("start");
7607        relay.run_turn("task", 0).await.expect("first turn");
7608        relay.run_turn("review this", 0).await.expect("review turn");
7609
7610        assert_eq!(relay.dispatches().len(), 2);
7611        assert_eq!(relay.dispatches()[0].0, 0);
7612        assert!(relay.dispatches()[0].1.contains("task"));
7613        assert!(relay.dispatches()[0].1.contains("You are Codex"));
7614        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7615        assert_eq!(relay.dispatches()[1].0, 1);
7616        assert!(relay.dispatches()[1].1.contains("review this"));
7617        let public = relay.dispatches()[1]
7618            .1
7619            .split_once("Public updates:\n")
7620            .map(|(_, updates)| updates)
7621            .expect("review receives public context");
7622        let header = public
7623            .lines()
7624            .find(|line| line.starts_with("Codex "))
7625            .expect("named previous agent");
7626        let timestamp = header
7627            .strip_prefix("Codex ")
7628            .and_then(|value| value.strip_suffix(':'))
7629            .expect("timestamped header");
7630        assert_eq!(timestamp.len(), 5);
7631        assert_eq!(timestamp.as_bytes()[2], b':');
7632        assert!(
7633            timestamp
7634                .bytes()
7635                .enumerate()
7636                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7637        );
7638        assert!(public.contains("implemented the fix"));
7639        assert!(!public.contains("Agent 0"));
7640    }
7641}