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