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                if let Some(reason) = value
3588                    .get("result")
3589                    .and_then(|result| result.get("stopReason"))
3590                    .and_then(Value::as_str)
3591                    .filter(|reason| !matches!(*reason, "end_turn" | "cancelled"))
3592                {
3593                    let detail = match reason {
3594                        "max_tokens" => {
3595                            "ACP turn stopped because the output token limit was reached".into()
3596                        }
3597                        "max_turn_requests" => {
3598                            "ACP turn stopped because the tool-turn limit was reached".into()
3599                        }
3600                        other => format!("ACP turn stopped before responding: {other}"),
3601                    };
3602                    return Some(Ok(AgentEvent::Failed {
3603                        slot: self.slot,
3604                        started: true,
3605                        detail,
3606                    }));
3607                }
3608                return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3609            }
3610        }
3611    }
3612}
3613
3614#[cfg(test)]
3615fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3616    let value: Value =
3617        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3618    parse_acp_value(slot, &value, &mut BTreeMap::new())
3619}
3620
3621fn parse_acp_value(
3622    slot: RosterSlot,
3623    value: &Value,
3624    tools: &mut BTreeMap<String, ToolUpdate>,
3625) -> AdapterResult<Option<AgentEvent>> {
3626    let method = value.get("method").and_then(Value::as_str);
3627    if method == Some("session/request_permission") {
3628        let params = value.get("params").cloned().unwrap_or(Value::Null);
3629        let request_id = value
3630            .get("id")
3631            .map(rpc_id_to_string)
3632            .unwrap_or_else(|| "permission".into());
3633        return Ok(parse_permission_event(
3634            slot,
3635            &params,
3636            &request_id,
3637            params.get("options"),
3638        ));
3639    }
3640    if method != Some("session/update") {
3641        return Ok(None);
3642    }
3643    let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3644        return Ok(None);
3645    };
3646    let kind = update.get("sessionUpdate").and_then(Value::as_str);
3647    if kind == Some("config_option_update")
3648        && let Some((config_id, models, current_model)) = parse_model_config(update)
3649    {
3650        return Ok(Some(AgentEvent::ModelsReplaced {
3651            slot,
3652            config_id,
3653            models,
3654            current_model,
3655        }));
3656    }
3657    if kind == Some("request_permission") {
3658        let request_id = update
3659            .get("toolCall")
3660            .and_then(|tool| tool.get("toolCallId"))
3661            .and_then(Value::as_str)
3662            .unwrap_or("permission");
3663        return Ok(parse_permission_event(
3664            slot,
3665            update,
3666            request_id,
3667            update.get("options"),
3668        ));
3669    }
3670    if kind == Some("available_commands_update") {
3671        let commands = update
3672            .get("availableCommands")
3673            .and_then(Value::as_array)
3674            .map(|commands| {
3675                commands
3676                    .iter()
3677                    .filter_map(|command| {
3678                        let name = command.get("name").and_then(Value::as_str)?.trim();
3679                        (!name.is_empty()).then(|| AgentCommand {
3680                            name: name.to_owned(),
3681                        })
3682                    })
3683                    .collect::<Vec<_>>()
3684            })
3685            .unwrap_or_default();
3686        return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3687    }
3688    if kind == Some("current_mode_update") {
3689        if let Some(mode) = update
3690            .get("currentModeId")
3691            .and_then(Value::as_str)
3692            .filter(|mode| !mode.trim().is_empty())
3693        {
3694            return Ok(Some(AgentEvent::ModeUpdated {
3695                slot,
3696                current_mode: mode.to_owned(),
3697            }));
3698        }
3699        return Ok(None);
3700    }
3701    if kind == Some("usage_update") {
3702        let Some(used) = update.get("used").and_then(Value::as_u64) else {
3703            return Ok(None);
3704        };
3705        let Some(size) = update.get("size").and_then(Value::as_u64) else {
3706            return Ok(None);
3707        };
3708        return Ok(Some(AgentEvent::UsageUpdated {
3709            slot,
3710            usage: UsageUpdate { used, size },
3711        }));
3712    }
3713    if let Some(terminal) = parse_terminal_event(update, kind) {
3714        return Ok(Some(AgentEvent::Terminal {
3715            slot,
3716            event: terminal,
3717        }));
3718    }
3719    let text = update
3720        .get("content")
3721        .and_then(|content| content.get("text"))
3722        .and_then(Value::as_str)
3723        .map(str::to_owned);
3724    if kind == Some("user_message_chunk") {
3725        return Ok(text
3726            .filter(|text| !text.is_empty())
3727            .map(|text| AgentEvent::UserText { slot, text }));
3728    }
3729    if kind == Some("agent_message_chunk")
3730        && let Some(mode) = text
3731            .as_deref()
3732            .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3733            .map(str::trim)
3734            .filter(|mode| !mode.is_empty())
3735    {
3736        // Gemini's native ACP bridge historically encoded a mode change as a
3737        // control marker in the message stream. It is state, not transcript
3738        // content; expose it as a normalized catalog replacement instead of
3739        // leaking the marker into the conversation.
3740        return Ok(Some(AgentEvent::ModesReplaced {
3741            slot,
3742            modes: vec![Mode {
3743                id: mode.to_owned(),
3744                label: mode.to_owned(),
3745            }],
3746            current_mode: Some(mode.to_owned()),
3747        }));
3748    }
3749    match (kind, text) {
3750        (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3751            Ok(Some(AgentEvent::Text { slot, text }))
3752        }
3753        (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3754            Ok(Some(AgentEvent::Thought { slot, text }))
3755        }
3756        (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3757            Ok(normalize_acp_tool(update, tools).map(|update| AgentEvent::Tool { slot, update }))
3758        }
3759        _ => Ok(None),
3760    }
3761}
3762
3763/// ACP updates are patches: omitted/invalid fields retain their last valid
3764/// values, whereas an explicit empty content array clears previous output.
3765fn normalize_acp_tool(
3766    value: &Value,
3767    tools: &mut BTreeMap<String, ToolUpdate>,
3768) -> Option<ToolUpdate> {
3769    let id = value.get("toolCallId")?.as_str()?;
3770    if id.trim().is_empty() {
3771        return None;
3772    }
3773    if value.get("sessionUpdate").and_then(Value::as_str) == Some("tool_call") {
3774        tools.remove(id);
3775    }
3776    let tool = tools.entry(id.to_owned()).or_insert_with(|| ToolUpdate {
3777        id: id.to_owned(),
3778        title: "Tool call".into(),
3779        status: ToolStatus::Pending,
3780        detail: None,
3781    });
3782    if let Some(title) = value.get("title").and_then(Value::as_str) {
3783        tool.title = title.to_owned();
3784    }
3785    if let Some(status) =
3786        value
3787            .get("status")
3788            .and_then(Value::as_str)
3789            .and_then(|status| match status {
3790                "pending" => Some(ToolStatus::Pending),
3791                "in_progress" => Some(ToolStatus::Running),
3792                "completed" => Some(ToolStatus::Completed),
3793                "failed" => Some(ToolStatus::Failed),
3794                _ => None,
3795            })
3796    {
3797        tool.status = status;
3798    }
3799    if let Some(content) = value.get("content").and_then(Value::as_array) {
3800        let text = content
3801            .iter()
3802            .filter_map(|entry| match entry.get("type").and_then(Value::as_str) {
3803                Some("content") => entry.get("content")?.get("text")?.as_str(),
3804                Some("diff") => entry.get("newText")?.as_str(),
3805                _ => None,
3806            })
3807            .collect::<Vec<_>>()
3808            .join("\n");
3809        tool.detail = (!text.is_empty()).then_some(text);
3810    } else if let Some(output) = value.get("rawOutput").filter(|output| !output.is_null()) {
3811        tool.detail = Some(
3812            output
3813                .as_str()
3814                .map(str::to_owned)
3815                .unwrap_or_else(|| output.to_string()),
3816        );
3817    }
3818    Some(tool.clone())
3819}
3820
3821fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3822    let config = value
3823        .get("configOptions")?
3824        .as_array()?
3825        .iter()
3826        .find(|option| {
3827            option.get("category").and_then(Value::as_str) == Some("model")
3828                && matches!(
3829                    option.get("type").and_then(Value::as_str),
3830                    Some("select" | "enum")
3831                )
3832        })?;
3833    let config_id = config.get("id")?.as_str()?.to_owned();
3834    let models = config
3835        .get("options")?
3836        .as_array()?
3837        .iter()
3838        .filter_map(|option| {
3839            let id = option.get("value")?.as_str()?.to_owned();
3840            let label = option
3841                .get("name")
3842                .or_else(|| option.get("label"))
3843                .and_then(Value::as_str)
3844                .unwrap_or(&id)
3845                .to_owned();
3846            Some(Mode { id, label })
3847        })
3848        .collect::<Vec<_>>();
3849    (!models.is_empty()).then(|| {
3850        let current = config
3851            .get("currentValue")
3852            .and_then(Value::as_str)
3853            .map(str::to_owned);
3854        (config_id, models, current)
3855    })
3856}
3857
3858/// Normalize terminal lifecycle updates emitted by ACP-compatible bridges and
3859/// native stream adapters. Protocols have used both snake_case update names
3860/// and a nested `terminal` object, so accept either without leaking that
3861/// shape beyond the adapter boundary.
3862fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3863    let nested = value.get("terminal").unwrap_or(value);
3864    let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3865    let id = nested
3866        .get("terminalId")
3867        .or_else(|| nested.get("terminal_id"))
3868        .or_else(|| nested.get("id"))
3869        .and_then(Value::as_str)
3870        .unwrap_or("terminal")
3871        .to_owned();
3872    match kind {
3873        "terminal_created" | "terminal_create" | "terminal_started" => {
3874            let command = nested
3875                .get("command")
3876                .and_then(Value::as_str)
3877                .unwrap_or("")
3878                .to_owned();
3879            Some(TerminalEvent::Created { id, command })
3880        }
3881        "terminal_output" | "terminal_output_chunk" => {
3882            let text = nested
3883                .get("output")
3884                .or_else(|| nested.get("text"))
3885                .and_then(Value::as_str)
3886                .unwrap_or("")
3887                .to_owned();
3888            Some(TerminalEvent::Output { id, text })
3889        }
3890        "terminal_exited" | "terminal_exit" => {
3891            let code = nested
3892                .get("exitCode")
3893                .or_else(|| nested.get("exit_code"))
3894                .or_else(|| nested.get("code"))
3895                .and_then(Value::as_i64)
3896                .unwrap_or(0) as i32;
3897            Some(TerminalEvent::Exited { id, code })
3898        }
3899        "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3900        _ => None,
3901    }
3902}
3903
3904fn parse_permission_event(
3905    slot: RosterSlot,
3906    value: &Value,
3907    request_id: &str,
3908    options: Option<&Value>,
3909) -> Option<AgentEvent> {
3910    let tool = value.get("toolCall").unwrap_or(value);
3911    let title = tool
3912        .get("title")
3913        .and_then(Value::as_str)
3914        .unwrap_or("Agent requests permission")
3915        .to_owned();
3916    let (options, option_ids): (Vec<String>, Vec<String>) = options
3917        .and_then(Value::as_array)
3918        .map(|options| {
3919            options
3920                .iter()
3921                .filter_map(|option| {
3922                    let label = option
3923                        .get("name")
3924                        .or_else(|| option.get("optionId"))
3925                        .and_then(Value::as_str)?
3926                        .to_owned();
3927                    let option_id = option
3928                        .get("optionId")
3929                        .or_else(|| option.get("id"))
3930                        .and_then(Value::as_str)
3931                        .map(str::to_owned)
3932                        .unwrap_or_else(|| label.clone());
3933                    Some((label, option_id))
3934                })
3935                .unzip()
3936        })
3937        .unwrap_or_default();
3938    if options.is_empty() {
3939        return None;
3940    }
3941    Some(AgentEvent::Permission {
3942        slot,
3943        request: PermissionRequest {
3944            id: request_id.to_owned(),
3945            title,
3946            options,
3947            option_ids,
3948        },
3949    })
3950}
3951
3952fn rpc_id_to_string(value: &Value) -> String {
3953    value
3954        .as_str()
3955        .map(str::to_owned)
3956        .or_else(|| value.as_u64().map(|id| id.to_string()))
3957        .unwrap_or_else(|| value.to_string())
3958}
3959
3960#[cfg(test)]
3961mod tests {
3962    use super::{
3963        AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3964        RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3965        parse_model_config, prompt_content_blocks, read_bounded_line,
3966    };
3967    #[cfg(target_os = "linux")]
3968    use super::{isolate_process_group, terminate_child};
3969    use crate::TerminalEvent;
3970    use crate::{
3971        AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3972        persistence::SessionMetadataStore,
3973        relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3974    };
3975    use async_trait::async_trait;
3976    use serde_json::Value;
3977    use std::collections::VecDeque;
3978    use std::sync::{
3979        Arc, Mutex,
3980        atomic::{AtomicUsize, Ordering},
3981    };
3982
3983    fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3984        let nonce = std::time::SystemTime::now()
3985            .duration_since(std::time::UNIX_EPOCH)
3986            .expect("clock")
3987            .as_nanos();
3988        std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3989    }
3990
3991    #[test]
3992    fn malformed_file_writes_preserve_existing_content() {
3993        let root = unique_test_path("codeswarm-write-validation", "dir");
3994        std::fs::create_dir_all(&root).unwrap();
3995        let file = root.join("keep.txt");
3996        std::fs::write(&file, "valuable content").unwrap();
3997        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3998        for content in [
3999            Value::Null,
4000            serde_json::json!(false),
4001            serde_json::json!(42),
4002            serde_json::json!([]),
4003        ] {
4004            assert!(
4005                adapter
4006                    .write_workspace_text(
4007                        &serde_json::json!({"path":"keep.txt", "content": content})
4008                    )
4009                    .is_err()
4010            );
4011            assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4012        }
4013        assert!(
4014            adapter
4015                .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
4016                .is_err()
4017        );
4018        assert!(
4019            adapter
4020                .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
4021                .is_err()
4022        );
4023        assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4024        adapter
4025            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
4026            .unwrap();
4027        assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
4028        adapter
4029            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
4030            .unwrap();
4031        assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
4032        std::fs::remove_dir_all(root).unwrap();
4033    }
4034
4035    #[tokio::test]
4036    async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
4037        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
4038        let mut adapter = AcpAdapter::new(
4039            0,
4040            std::env::current_dir().unwrap(),
4041            "sh",
4042            vec!["-c".into(), script.into()],
4043        );
4044        adapter.start().await.unwrap();
4045        let error = adapter
4046            .request_with_timeout(
4047                "session/set_mode",
4048                serde_json::json!({}),
4049                std::time::Duration::from_millis(10),
4050            )
4051            .await
4052            .unwrap_err();
4053        assert!(error.to_string().contains("session/set_mode timed out"));
4054        adapter.stop().await.unwrap();
4055        assert!(adapter.child.is_none());
4056    }
4057
4058    #[tokio::test]
4059    async fn goals_reach_every_roster_slot_without_native_goal_support() {
4060        use crate::goal::GoalCommand;
4061        let hosts = (0..3)
4062            .map(|slot| {
4063                AdapterHost::new(
4064                    Box::new(ScriptedAdapter::new(
4065                        slot,
4066                        AgentCapabilities::default(),
4067                        [
4068                            AgentEvent::TurnComplete { slot },
4069                            AgentEvent::TurnComplete { slot },
4070                        ],
4071                    )),
4072                    None,
4073                )
4074            })
4075            .collect();
4076        let mut relay = RelayHost::new(hosts, 10).unwrap();
4077        relay.start().await.unwrap();
4078        let task = relay
4079            .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4080            .unwrap()
4081            .unwrap();
4082        relay.relay_mut().enqueue_human(task, Some(0));
4083        for slot in 0..3 {
4084            relay.run_turn("", 0).await.unwrap();
4085            let (actual, prompt) = relay.dispatches().last().unwrap();
4086            assert_eq!(*actual, slot);
4087            assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4088        }
4089        let snapshot = relay.session_metadata();
4090        let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4091        assert!(restored.is_some());
4092        relay.restore_goal(restored);
4093        relay.reload(0).await.unwrap();
4094        relay.run_turn("", 0).await.unwrap();
4095        assert!(
4096            relay
4097                .dispatches()
4098                .last()
4099                .unwrap()
4100                .1
4101                .contains("Active shared goal: Ship the settings screen")
4102        );
4103        relay.apply_goal(GoalCommand::Done).unwrap();
4104        relay.run_turn("", 0).await.unwrap();
4105        assert!(
4106            relay
4107                .dispatches()
4108                .last()
4109                .unwrap()
4110                .1
4111                .contains("No active shared goal")
4112        );
4113        relay.apply_goal(GoalCommand::Clear).unwrap();
4114        assert!(relay.session_metadata().get("goal").unwrap().is_null());
4115    }
4116
4117    #[tokio::test]
4118    async fn replacement_agent_receives_task_after_public_journal_pruning() {
4119        let hosts = (0..2)
4120            .map(|slot| {
4121                AdapterHost::new(
4122                    Box::new(ScriptedAdapter::new(
4123                        slot,
4124                        AgentCapabilities::default(),
4125                        [
4126                            AgentEvent::Text {
4127                                slot,
4128                                text: "progress".into(),
4129                            },
4130                            AgentEvent::TurnComplete { slot },
4131                            AgentEvent::Text {
4132                                slot,
4133                                text: "more progress".into(),
4134                            },
4135                            AgentEvent::TurnComplete { slot },
4136                        ],
4137                    )),
4138                    None,
4139                )
4140            })
4141            .collect();
4142        let mut relay = RelayHost::new(hosts, 10).unwrap();
4143        relay.start().await.unwrap();
4144        relay
4145            .relay_mut()
4146            .enqueue_human("Fix the login bug", Some(0));
4147        relay.run_turn("", 0).await.unwrap();
4148        relay.run_turn("", 0).await.unwrap();
4149        relay.run_turn("", 0).await.unwrap();
4150        relay.reload(1).await.unwrap();
4151        assert!(
4152            !relay
4153                .relay_mut()
4154                .unseen_context(1)
4155                .contains("Fix the login bug")
4156        );
4157        relay.run_turn("", 0).await.unwrap();
4158        assert!(
4159            relay
4160                .dispatches()
4161                .last()
4162                .unwrap()
4163                .1
4164                .contains("Shared task:\nFix the login bug")
4165        );
4166    }
4167
4168    #[cfg(target_os = "linux")]
4169    #[tokio::test]
4170    async fn termination_kills_only_the_verified_isolated_child_group() {
4171        use nix::unistd::{Pid, getpgid, getpgrp};
4172        use tokio::io::{AsyncBufReadExt, BufReader};
4173
4174        let own_group = getpgrp();
4175        let mut command = tokio::process::Command::new("sh");
4176        isolate_process_group(&mut command);
4177        command
4178            .arg("-c")
4179            .arg("sleep 60 & echo $!; wait")
4180            .stdout(std::process::Stdio::piped());
4181        let mut child = command.spawn().expect("spawn isolated shell");
4182        let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4183        assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4184        assert_ne!(leader, own_group);
4185
4186        let stdout = child.stdout.take().expect("child stdout");
4187        let mut lines = BufReader::new(stdout).lines();
4188        let descendant = lines
4189            .next_line()
4190            .await
4191            .expect("read descendant pid")
4192            .expect("descendant pid")
4193            .parse::<i32>()
4194            .expect("numeric descendant pid");
4195        let descendant = Pid::from_raw(descendant);
4196        assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4197
4198        terminate_child(&mut child).await.expect("terminate group");
4199        for _ in 0..100 {
4200            if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4201                return;
4202            }
4203            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4204        }
4205        panic!("descendant {descendant} survived isolated group termination");
4206    }
4207
4208    #[test]
4209    fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4210        assert_eq!(
4211            parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4212            Ok((
4213                "agent".into(),
4214                vec![
4215                    "--name".into(),
4216                    "local bridge".into(),
4217                    "--flag".into(),
4218                    "two words".into(),
4219                ]
4220            ),)
4221        );
4222        assert_eq!(
4223            parse_command_line(r#"agent "" escaped\ argument"#),
4224            Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4225        );
4226    }
4227
4228    #[test]
4229    fn acp_prompt_expands_safe_at_path_resources() {
4230        let root = unique_test_path("codeswarm-prompt-resource", "dir");
4231        std::fs::create_dir_all(&root).expect("workspace");
4232        std::fs::write(root.join("note.md"), "resource text").expect("resource");
4233        let blocks = prompt_content_blocks(&root, "inspect @note.md");
4234        assert_eq!(blocks[0]["type"], "text");
4235        assert_eq!(blocks[0]["text"], "inspect @note.md");
4236        assert_eq!(blocks[1]["type"], "resource");
4237        assert_eq!(blocks[1]["resource"]["text"], "resource text");
4238        assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4239        std::fs::remove_dir_all(root).expect("cleanup workspace");
4240    }
4241
4242    #[tokio::test]
4243    async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4244        let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4245        bytes.push(b'\n');
4246        let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4247        assert!(matches!(
4248            read_bounded_line(&mut reader).await,
4249            Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4250        ));
4251    }
4252
4253    #[test]
4254    fn rejects_malformed_configured_commands_before_spawn() {
4255        assert_eq!(
4256            parse_command_line("agent 'unfinished"),
4257            Err(super::CommandParseError::UnterminatedQuote)
4258        );
4259        assert_eq!(
4260            parse_command_line("agent\\"),
4261            Err(super::CommandParseError::TrailingEscape)
4262        );
4263        assert_eq!(
4264            parse_command_line("   \t"),
4265            Err(super::CommandParseError::Empty)
4266        );
4267    }
4268
4269    #[derive(Debug)]
4270    struct PendingAdapter {
4271        slot: usize,
4272        hang_on_cancel: bool,
4273    }
4274
4275    #[derive(Debug)]
4276    struct ConcurrentStartAdapter {
4277        slot: usize,
4278        barrier: Arc<tokio::sync::Barrier>,
4279    }
4280
4281    #[derive(Debug)]
4282    struct ReloadProbeAdapter {
4283        slot: usize,
4284        crashed: bool,
4285        reloaded: bool,
4286        events: VecDeque<AgentEvent>,
4287        prompts: Arc<Mutex<Vec<String>>>,
4288    }
4289
4290    #[async_trait]
4291    impl AgentAdapter for ReloadProbeAdapter {
4292        fn slot(&self) -> usize {
4293            self.slot
4294        }
4295
4296        fn display_name(&self) -> String {
4297            "Reload probe".into()
4298        }
4299
4300        fn protocol(&self) -> &'static str {
4301            "native"
4302        }
4303
4304        fn capabilities(&self) -> AgentCapabilities {
4305            AgentCapabilities {
4306                supports_modes: true,
4307                ..AgentCapabilities::default()
4308            }
4309        }
4310
4311        fn needs_restart(&self) -> bool {
4312            self.crashed && !self.reloaded
4313        }
4314
4315        async fn start(&mut self) -> super::AdapterResult<()> {
4316            self.events.push_back(AgentEvent::ModesReplaced {
4317                slot: self.slot,
4318                modes: vec![
4319                    Mode {
4320                        id: "codeswarm:mode:full-access".into(),
4321                        label: "Auto pilot".into(),
4322                    },
4323                    Mode {
4324                        id: "codeswarm:mode:plan".into(),
4325                        label: "Plan".into(),
4326                    },
4327                ],
4328                current_mode: Some("codeswarm:mode:full-access".into()),
4329            });
4330            self.events.push_back(AgentEvent::Ready {
4331                slot: self.slot,
4332                capabilities: self.capabilities(),
4333            });
4334            Ok(())
4335        }
4336
4337        async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4338            self.prompts.lock().expect("prompts").push(prompt);
4339            if !self.crashed {
4340                self.crashed = true;
4341                self.events.push_back(AgentEvent::Failed {
4342                    slot: self.slot,
4343                    started: true,
4344                    detail: "probe crashed".into(),
4345                });
4346            } else {
4347                self.events.push_back(AgentEvent::Text {
4348                    slot: self.slot,
4349                    text: "recovered".into(),
4350                });
4351                self.events
4352                    .push_back(AgentEvent::TurnComplete { slot: self.slot });
4353            }
4354            Ok(())
4355        }
4356
4357        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4358            Ok(false)
4359        }
4360
4361        async fn answer_permission(
4362            &mut self,
4363            _request_id: String,
4364            _answer: PermissionAnswer,
4365        ) -> super::AdapterResult<()> {
4366            Err(super::AdapterError::Unsupported("permission answer"))
4367        }
4368
4369        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4370            Ok(())
4371        }
4372
4373        async fn reload(&mut self) -> super::AdapterResult<()> {
4374            self.reloaded = true;
4375            self.start().await
4376        }
4377
4378        async fn stop(&mut self) -> super::AdapterResult<()> {
4379            Ok(())
4380        }
4381
4382        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4383            self.events.pop_front().map(Ok)
4384        }
4385    }
4386
4387    #[async_trait]
4388    impl AgentAdapter for ConcurrentStartAdapter {
4389        fn slot(&self) -> usize {
4390            self.slot
4391        }
4392
4393        fn capabilities(&self) -> AgentCapabilities {
4394            AgentCapabilities::default()
4395        }
4396
4397        async fn start(&mut self) -> super::AdapterResult<()> {
4398            self.barrier.wait().await;
4399            Ok(())
4400        }
4401
4402        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4403            Ok(())
4404        }
4405
4406        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4407            Ok(true)
4408        }
4409
4410        async fn answer_permission(
4411            &mut self,
4412            _request_id: String,
4413            _answer: PermissionAnswer,
4414        ) -> super::AdapterResult<()> {
4415            Ok(())
4416        }
4417
4418        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4419            Ok(())
4420        }
4421
4422        async fn reload(&mut self) -> super::AdapterResult<()> {
4423            Ok(())
4424        }
4425
4426        async fn stop(&mut self) -> super::AdapterResult<()> {
4427            Ok(())
4428        }
4429
4430        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4431            std::future::pending().await
4432        }
4433    }
4434
4435    #[derive(Debug)]
4436    struct PermissionBlockingAdapter {
4437        slot: usize,
4438        phase: u8,
4439    }
4440
4441    #[async_trait]
4442    impl AgentAdapter for PermissionBlockingAdapter {
4443        fn slot(&self) -> usize {
4444            self.slot
4445        }
4446
4447        fn capabilities(&self) -> AgentCapabilities {
4448            AgentCapabilities {
4449                supports_permissions: true,
4450                ..AgentCapabilities::default()
4451            }
4452        }
4453
4454        async fn start(&mut self) -> super::AdapterResult<()> {
4455            Ok(())
4456        }
4457
4458        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4459            Ok(())
4460        }
4461
4462        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4463            Ok(true)
4464        }
4465
4466        async fn answer_permission(
4467            &mut self,
4468            request_id: String,
4469            answer: PermissionAnswer,
4470        ) -> super::AdapterResult<()> {
4471            if self.phase != 1 || request_id != "permission-1" {
4472                return Err(super::AdapterError::Protocol(
4473                    "unexpected permission response".into(),
4474                ));
4475            }
4476            assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4477            self.phase = 2;
4478            Ok(())
4479        }
4480
4481        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4482            Ok(())
4483        }
4484
4485        async fn reload(&mut self) -> super::AdapterResult<()> {
4486            Ok(())
4487        }
4488
4489        async fn stop(&mut self) -> super::AdapterResult<()> {
4490            Ok(())
4491        }
4492
4493        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4494            match self.phase {
4495                0 => {
4496                    self.phase = 1;
4497                    Some(Ok(AgentEvent::Permission {
4498                        slot: self.slot,
4499                        request: crate::PermissionRequest {
4500                            id: "permission-1".into(),
4501                            title: "Allow?".into(),
4502                            options: vec!["Allow".into()],
4503                            option_ids: vec!["allow".into()],
4504                        },
4505                    }))
4506                }
4507                1 => std::future::pending().await,
4508                _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4509            }
4510        }
4511    }
4512
4513    #[async_trait]
4514    impl AgentAdapter for PendingAdapter {
4515        fn slot(&self) -> usize {
4516            self.slot
4517        }
4518
4519        fn capabilities(&self) -> AgentCapabilities {
4520            AgentCapabilities {
4521                supports_cancel: true,
4522                ..AgentCapabilities::default()
4523            }
4524        }
4525
4526        async fn start(&mut self) -> super::AdapterResult<()> {
4527            Ok(())
4528        }
4529
4530        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4531            Ok(())
4532        }
4533
4534        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4535            if self.hang_on_cancel {
4536                return std::future::pending().await;
4537            }
4538            Ok(true)
4539        }
4540
4541        async fn answer_permission(
4542            &mut self,
4543            _request_id: String,
4544            _answer: PermissionAnswer,
4545        ) -> super::AdapterResult<()> {
4546            Ok(())
4547        }
4548
4549        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4550            Ok(())
4551        }
4552
4553        async fn reload(&mut self) -> super::AdapterResult<()> {
4554            Ok(())
4555        }
4556
4557        async fn stop(&mut self) -> super::AdapterResult<()> {
4558            Ok(())
4559        }
4560
4561        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4562            std::future::pending().await
4563        }
4564    }
4565
4566    #[derive(Debug)]
4567    struct StopTrackingAdapter {
4568        slot: usize,
4569        stopped: Arc<AtomicUsize>,
4570        fail_stop: bool,
4571    }
4572
4573    #[derive(Debug)]
4574    struct ModeOrderAdapter {
4575        slot: usize,
4576        log: Arc<Mutex<Vec<String>>>,
4577        phase: u8,
4578    }
4579
4580    #[derive(Debug)]
4581    struct StartupAcpAdapter {
4582        slot: usize,
4583        events: std::collections::VecDeque<AgentEvent>,
4584    }
4585
4586    impl StartupAcpAdapter {
4587        fn new(slot: usize) -> Self {
4588            Self {
4589                slot,
4590                events: [
4591                    AgentEvent::ModesReplaced {
4592                        slot,
4593                        modes: vec![Mode {
4594                            id: "full-access".into(),
4595                            label: "Auto pilot".into(),
4596                        }],
4597                        current_mode: Some("full-access".into()),
4598                    },
4599                    AgentEvent::Ready {
4600                        slot,
4601                        capabilities: AgentCapabilities {
4602                            supports_modes: true,
4603                            ..AgentCapabilities::default()
4604                        },
4605                    },
4606                ]
4607                .into(),
4608            }
4609        }
4610    }
4611
4612    #[async_trait]
4613    impl AgentAdapter for StartupAcpAdapter {
4614        fn slot(&self) -> usize {
4615            self.slot
4616        }
4617
4618        fn protocol(&self) -> &'static str {
4619            "acp"
4620        }
4621
4622        fn capabilities(&self) -> AgentCapabilities {
4623            AgentCapabilities {
4624                supports_modes: true,
4625                ..AgentCapabilities::default()
4626            }
4627        }
4628
4629        async fn start(&mut self) -> super::AdapterResult<()> {
4630            Ok(())
4631        }
4632
4633        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4634            Ok(())
4635        }
4636
4637        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4638            Ok(true)
4639        }
4640
4641        async fn answer_permission(
4642            &mut self,
4643            _request_id: String,
4644            _answer: PermissionAnswer,
4645        ) -> super::AdapterResult<()> {
4646            Ok(())
4647        }
4648
4649        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4650            Ok(())
4651        }
4652
4653        async fn reload(&mut self) -> super::AdapterResult<()> {
4654            self.events = Self::new(self.slot).events;
4655            Ok(())
4656        }
4657
4658        async fn stop(&mut self) -> super::AdapterResult<()> {
4659            Ok(())
4660        }
4661
4662        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4663            self.events.pop_front().map(Ok)
4664        }
4665    }
4666
4667    #[async_trait]
4668    impl AgentAdapter for ModeOrderAdapter {
4669        fn slot(&self) -> usize {
4670            self.slot
4671        }
4672
4673        fn capabilities(&self) -> AgentCapabilities {
4674            AgentCapabilities {
4675                supports_modes: true,
4676                ..AgentCapabilities::default()
4677            }
4678        }
4679
4680        async fn start(&mut self) -> super::AdapterResult<()> {
4681            self.log.lock().expect("log").push("start".into());
4682            Ok(())
4683        }
4684
4685        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4686            self.log.lock().expect("log").push("prompt".into());
4687            Ok(())
4688        }
4689
4690        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4691            Ok(true)
4692        }
4693
4694        async fn answer_permission(
4695            &mut self,
4696            _request_id: String,
4697            _answer: PermissionAnswer,
4698        ) -> super::AdapterResult<()> {
4699            Ok(())
4700        }
4701
4702        async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4703            self.log.lock().expect("log").push(format!("mode:{mode}"));
4704            Ok(())
4705        }
4706
4707        async fn reload(&mut self) -> super::AdapterResult<()> {
4708            self.log.lock().expect("log").push("reload".into());
4709            self.phase = 0;
4710            Ok(())
4711        }
4712
4713        async fn stop(&mut self) -> super::AdapterResult<()> {
4714            Ok(())
4715        }
4716
4717        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4718            match self.phase {
4719                0 => {
4720                    self.phase = 1;
4721                    Some(Ok(AgentEvent::ModesReplaced {
4722                        slot: self.slot,
4723                        modes: vec![Mode {
4724                            id: "yolo".into(),
4725                            label: "YOLO".into(),
4726                        }],
4727                        current_mode: None,
4728                    }))
4729                }
4730                1 => {
4731                    self.phase = 2;
4732                    Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4733                }
4734                _ => std::future::pending().await,
4735            }
4736        }
4737    }
4738
4739    #[async_trait]
4740    impl AgentAdapter for StopTrackingAdapter {
4741        fn slot(&self) -> usize {
4742            self.slot
4743        }
4744
4745        fn capabilities(&self) -> AgentCapabilities {
4746            AgentCapabilities::default()
4747        }
4748
4749        async fn start(&mut self) -> super::AdapterResult<()> {
4750            Ok(())
4751        }
4752
4753        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4754            Ok(())
4755        }
4756
4757        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4758            Ok(false)
4759        }
4760
4761        async fn answer_permission(
4762            &mut self,
4763            _request_id: String,
4764            _answer: PermissionAnswer,
4765        ) -> super::AdapterResult<()> {
4766            Ok(())
4767        }
4768
4769        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4770            Ok(())
4771        }
4772
4773        async fn reload(&mut self) -> super::AdapterResult<()> {
4774            Ok(())
4775        }
4776
4777        async fn stop(&mut self) -> super::AdapterResult<()> {
4778            self.stopped.fetch_add(1, Ordering::Relaxed);
4779            if self.fail_stop {
4780                Err(super::AdapterError::Transport("stop failed".into()))
4781            } else {
4782                Ok(())
4783            }
4784        }
4785
4786        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4787            None
4788        }
4789    }
4790
4791    /// Startup can fail after an adapter has allocated resources.  Keep a
4792    /// fixture that records whether the failing adapter itself receives the
4793    /// cleanup call, not just the already-started peers.
4794    #[derive(Debug)]
4795    struct FailingStartAdapter {
4796        slot: usize,
4797        stopped: Arc<AtomicUsize>,
4798    }
4799
4800    #[async_trait]
4801    impl AgentAdapter for FailingStartAdapter {
4802        fn slot(&self) -> usize {
4803            self.slot
4804        }
4805
4806        fn capabilities(&self) -> AgentCapabilities {
4807            AgentCapabilities::default()
4808        }
4809
4810        async fn start(&mut self) -> super::AdapterResult<()> {
4811            Err(super::AdapterError::Spawn("startup failed".into()))
4812        }
4813
4814        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4815            Ok(())
4816        }
4817
4818        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4819            Ok(false)
4820        }
4821
4822        async fn answer_permission(
4823            &mut self,
4824            _request_id: String,
4825            _answer: PermissionAnswer,
4826        ) -> super::AdapterResult<()> {
4827            Ok(())
4828        }
4829
4830        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4831            Ok(())
4832        }
4833
4834        async fn reload(&mut self) -> super::AdapterResult<()> {
4835            Ok(())
4836        }
4837
4838        async fn stop(&mut self) -> super::AdapterResult<()> {
4839            self.stopped.fetch_add(1, Ordering::Relaxed);
4840            Ok(())
4841        }
4842
4843        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4844            None
4845        }
4846    }
4847
4848    #[tokio::test]
4849    async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4850        let stopped = Arc::new(AtomicUsize::new(0));
4851        let relay = RelayHost::new(
4852            vec![
4853                AdapterHost::new(
4854                    Box::new(StopTrackingAdapter {
4855                        slot: 0,
4856                        stopped: Arc::clone(&stopped),
4857                        fail_stop: true,
4858                    }),
4859                    None,
4860                ),
4861                AdapterHost::new(
4862                    Box::new(StopTrackingAdapter {
4863                        slot: 1,
4864                        stopped: Arc::clone(&stopped),
4865                        fail_stop: false,
4866                    }),
4867                    None,
4868                ),
4869            ],
4870            4,
4871        )
4872        .expect("relay");
4873        let mut relay = relay;
4874
4875        let error = relay.stop().await.expect_err("first stop failure");
4876        assert!(error.to_string().contains("stop failed"));
4877        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4878    }
4879
4880    #[tokio::test]
4881    async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4882        let stopped = Arc::new(AtomicUsize::new(0));
4883        let mut relay = RelayHost::new(
4884            vec![
4885                AdapterHost::new(
4886                    Box::new(StopTrackingAdapter {
4887                        slot: 0,
4888                        stopped: Arc::clone(&stopped),
4889                        fail_stop: false,
4890                    }),
4891                    None,
4892                ),
4893                AdapterHost::new(
4894                    Box::new(FailingStartAdapter {
4895                        slot: 1,
4896                        stopped: Arc::clone(&stopped),
4897                    }),
4898                    None,
4899                ),
4900            ],
4901            4,
4902        )
4903        .expect("relay");
4904
4905        relay.start().await.expect("healthy peer remains available");
4906        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4907        assert_eq!(stopped.load(Ordering::Relaxed), 1);
4908        relay.stop().await.unwrap();
4909        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4910    }
4911
4912    #[test]
4913    fn parses_acp_text_without_ui_dependency() {
4914        let event = parse_acp_notification(
4915            2,
4916            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4917        )
4918        .expect("valid ACP")
4919        .expect("text event");
4920        assert_eq!(
4921            event,
4922            AgentEvent::Text {
4923                slot: 2,
4924                text: "hello".into(),
4925            }
4926        );
4927    }
4928
4929    #[test]
4930    fn parses_acp_state_notifications_at_the_adapter_boundary() {
4931        let commands = parse_acp_notification(
4932            3,
4933            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4934        )
4935        .expect("valid ACP")
4936        .expect("commands event");
4937        assert_eq!(
4938            commands,
4939            AgentEvent::CommandsReplaced {
4940                slot: 3,
4941                commands: vec![crate::AgentCommand {
4942                    name: "review".into()
4943                }]
4944            }
4945        );
4946
4947        let mode = parse_acp_notification(
4948            3,
4949            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4950        )
4951        .expect("valid ACP")
4952        .expect("mode event");
4953        assert_eq!(
4954            mode,
4955            AgentEvent::ModeUpdated {
4956                slot: 3,
4957                current_mode: "review".into()
4958            }
4959        );
4960
4961        let usage = parse_acp_notification(
4962            3,
4963            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4964        )
4965        .expect("valid ACP")
4966        .expect("usage event");
4967        assert_eq!(
4968            usage,
4969            AgentEvent::UsageUpdated {
4970                slot: 3,
4971                usage: crate::UsageUpdate {
4972                    used: 4200,
4973                    size: 128000
4974                }
4975            }
4976        );
4977
4978        let models = parse_acp_notification(
4979            3,
4980            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"}]}]}}}"#,
4981        )
4982        .expect("valid ACP")
4983        .expect("models event");
4984        assert!(matches!(
4985            models,
4986            AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4987                if models.len() == 2 && current_model.as_deref() == Some("smart")
4988        ));
4989        assert_eq!(
4990            parse_model_config(&serde_json::json!({
4991                "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4992            })),
4993            None
4994        );
4995
4996        let user = parse_acp_notification(
4997            3,
4998            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4999        )
5000        .expect("valid ACP")
5001        .expect("user event");
5002        assert_eq!(
5003            user,
5004            AgentEvent::UserText {
5005                slot: 3,
5006                text: "context".into()
5007            }
5008        );
5009    }
5010
5011    #[test]
5012    fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
5013        let event = parse_acp_notification(
5014            0,
5015            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
5016        )
5017        .expect("valid ACP")
5018        .expect("mode event");
5019        assert!(matches!(
5020            event,
5021            AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
5022                if mode == "yolo" && modes[0].id == "yolo"
5023        ));
5024    }
5025
5026    #[test]
5027    fn parses_native_agy_text_without_acp_bridge() {
5028        let event = parse_agy_line(
5029            1,
5030            r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
5031        )
5032        .expect("valid stream-json")
5033        .expect("text event");
5034        assert_eq!(
5035            event,
5036            AgentEvent::Text {
5037                slot: 1,
5038                text: "hello".into(),
5039            }
5040        );
5041    }
5042
5043    #[test]
5044    fn parses_tool_lifecycle_from_each_protocol() {
5045        let agy = parse_agy_line(
5046            1,
5047            r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5048        )
5049        .expect("valid native tool")
5050        .expect("tool event");
5051        assert!(matches!(
5052            agy,
5053            AgentEvent::Tool {
5054                update: crate::ToolUpdate {
5055                    status: ToolStatus::Completed,
5056                    ..
5057                },
5058                ..
5059            }
5060        ));
5061
5062        let acp = parse_acp_notification(
5063            1,
5064            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5065        )
5066        .expect("valid ACP tool")
5067        .expect("tool event");
5068        assert!(matches!(
5069            acp,
5070            AgentEvent::Tool {
5071                update: crate::ToolUpdate {
5072                    status: ToolStatus::Failed,
5073                    ..
5074                },
5075                ..
5076            }
5077        ));
5078    }
5079
5080    #[test]
5081    fn parses_terminal_lifecycle_from_acp_and_native_events() {
5082        let created = parse_acp_notification(
5083            0,
5084            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5085        )
5086        .expect("valid ACP terminal")
5087        .expect("terminal event");
5088        assert_eq!(
5089            created,
5090            AgentEvent::Terminal {
5091                slot: 0,
5092                event: TerminalEvent::Created {
5093                    id: "term-1".into(),
5094                    command: "cargo test".into(),
5095                },
5096            }
5097        );
5098        let output = parse_agy_line(
5099            1,
5100            r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5101        )
5102        .expect("valid native terminal")
5103        .expect("terminal event");
5104        assert_eq!(
5105            output,
5106            AgentEvent::Terminal {
5107                slot: 1,
5108                event: TerminalEvent::Output {
5109                    id: "term-1".into(),
5110                    text: "ok\n".into(),
5111                },
5112            }
5113        );
5114        let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5115            .expect("valid native release")
5116            .expect("terminal event");
5117        assert!(matches!(
5118            released,
5119            AgentEvent::Terminal {
5120                event: TerminalEvent::Released { id },
5121                ..
5122            } if id == "term-1"
5123        ));
5124    }
5125
5126    #[test]
5127    fn parses_acp_permission_requests() {
5128        let event = parse_acp_notification(
5129            0,
5130            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"}]}}}"#,
5131        )
5132        .expect("valid permission")
5133        .expect("permission event");
5134        assert!(matches!(
5135            event,
5136            AgentEvent::Permission { request, .. }
5137                if request.id == "t1"
5138                    && request.title == "Write file"
5139                    && request.options == ["Allow once", "Reject"]
5140                    && request.option_ids == ["allow-once", "reject"]
5141        ));
5142    }
5143
5144    #[test]
5145    fn parses_acp_permission_request_as_json_rpc_request() {
5146        let event = parse_acp_notification(
5147            2,
5148            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5149        )
5150        .expect("valid permission request")
5151        .expect("permission event");
5152        assert!(matches!(
5153            event,
5154            AgentEvent::Permission { request, .. }
5155                if request.id == "17"
5156                    && request.title == "Write file"
5157                    && request.options == ["allow-once", "reject"]
5158                    && request.option_ids == ["allow-once", "reject"]
5159        ));
5160    }
5161
5162    #[tokio::test]
5163    async fn native_adapter_explicitly_rejects_permission_answers() {
5164        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5165        assert_eq!(
5166            adapter
5167                .answer_permission(
5168                    "request".into(),
5169                    PermissionAnswer::Selected {
5170                        option_id: "allow".into()
5171                    },
5172                )
5173                .await,
5174            Err(super::AdapterError::Unsupported("permission answer"))
5175        );
5176    }
5177
5178    #[tokio::test]
5179    async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5180        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5181        adapter
5182            .set_mode("full-access".into())
5183            .await
5184            .expect("auto-pilot alias");
5185        assert!(matches!(
5186            adapter.next_event().await,
5187            Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5188        ));
5189    }
5190
5191    #[tokio::test]
5192    async fn native_turns_receive_a_twenty_four_hour_timeout() {
5193        let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5194        std::fs::write(
5195            &script_path,
5196            r#"#!/bin/sh
5197seen=0
5198while [ "$#" -gt 0 ]; do
5199    case "$1" in
5200        --print-timeout)
5201            shift
5202            [ "$1" = "1440m" ] || exit 2
5203            seen=$((seen + 1))
5204            ;;
5205    esac
5206    shift
5207done
5208[ "$seen" = 1 ] || exit 3
5209printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5210"#,
5211        )
5212        .unwrap();
5213        let mut adapter = AgyAdapter::with_session_id(
5214            0,
5215            std::env::current_dir().unwrap(),
5216            format!("sh {}", script_path.display()),
5217            "saved-session",
5218        );
5219        adapter.start().await.unwrap();
5220        adapter.next_event().await.unwrap().unwrap();
5221        adapter.next_event().await.unwrap().unwrap();
5222        for prompt in ["first task", "follow-up task"] {
5223            adapter.send_prompt(prompt.into()).await.unwrap();
5224            assert!(
5225                matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5226            );
5227            assert!(matches!(
5228                adapter.next_event().await,
5229                Some(Ok(AgentEvent::TurnComplete { .. }))
5230            ));
5231        }
5232        adapter.stop().await.unwrap();
5233        std::fs::remove_file(script_path).unwrap();
5234    }
5235
5236    #[tokio::test]
5237    async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5238        let script_path = unique_test_path("codeswarm-agy-session", "sh");
5239        std::fs::write(
5240            &script_path,
5241            "#!/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",
5242        )
5243        .expect("write native test script");
5244        let mut adapter = AgyAdapter::new(
5245            0,
5246            std::env::current_dir().expect("cwd"),
5247            format!("sh {}", script_path.display()),
5248        );
5249        adapter.start().await.expect("start native adapter");
5250        // Startup emits its mode catalog and readiness before a turn.
5251        assert!(adapter.next_event().await.is_some());
5252        assert!(adapter.next_event().await.is_some());
5253        adapter
5254            .send_prompt("first".into())
5255            .await
5256            .expect("first prompt");
5257        while !matches!(
5258            adapter.next_event().await,
5259            Some(Ok(AgentEvent::TurnComplete { .. }))
5260        ) {}
5261        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5262        adapter
5263            .send_prompt("follow up".into())
5264            .await
5265            .expect("follow-up prompt");
5266        while !matches!(
5267            adapter.next_event().await,
5268            Some(Ok(AgentEvent::TurnComplete { .. }))
5269        ) {}
5270        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5271        adapter.stop().await.expect("stop native adapter");
5272        std::fs::remove_file(script_path).expect("cleanup native script");
5273    }
5274
5275    #[tokio::test]
5276    async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5277        let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5278        std::fs::write(
5279            &script_path,
5280            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5281        )
5282        .expect("write native test script");
5283        let mut adapter = AgyAdapter::new(
5284            0,
5285            std::env::current_dir().expect("cwd"),
5286            format!("sh {}", script_path.display()),
5287        );
5288        adapter.start().await.expect("start native adapter");
5289        assert!(adapter.next_event().await.is_some());
5290        assert!(adapter.next_event().await.is_some());
5291        adapter.send_prompt("fail".into()).await.expect("prompt");
5292        assert!(matches!(
5293            adapter.next_event().await,
5294            Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5295                if detail == "agent failed"
5296        ));
5297        adapter.stop().await.expect("stop native adapter");
5298        std::fs::remove_file(script_path).expect("cleanup native script");
5299    }
5300
5301    #[tokio::test]
5302    async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5303        let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5304        let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5305        std::fs::write(
5306            &script_path,
5307            format!(
5308                "#!/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",
5309                marker_path.display(),
5310                marker_path.display(),
5311                marker_path.display(),
5312            ),
5313        )
5314        .expect("write retry script");
5315        let mut adapter = AgyAdapter::new(
5316            0,
5317            std::env::current_dir().expect("cwd"),
5318            format!("sh {}", script_path.display()),
5319        );
5320        adapter.start().await.expect("start native adapter");
5321        assert!(adapter.next_event().await.is_some());
5322        assert!(adapter.next_event().await.is_some());
5323        adapter
5324            .send_prompt("first".into())
5325            .await
5326            .expect("first prompt");
5327        assert!(matches!(
5328            adapter.next_event().await,
5329            Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5330        ));
5331        adapter
5332            .send_prompt("retry".into())
5333            .await
5334            .expect("retry prompt starts a fresh process");
5335        assert!(matches!(
5336            adapter.next_event().await,
5337            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5338        ));
5339        assert!(matches!(
5340            adapter.next_event().await,
5341            Some(Ok(AgentEvent::TurnComplete { .. }))
5342        ));
5343        adapter.stop().await.expect("stop native adapter");
5344        std::fs::remove_file(script_path).expect("cleanup retry script");
5345        std::fs::remove_file(marker_path).expect("cleanup retry marker");
5346    }
5347
5348    #[tokio::test]
5349    async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5350        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"}}'"#;
5351        let cwd = std::env::current_dir().expect("cwd");
5352        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5353        adapter.start().await.expect("initialize");
5354        assert!(matches!(
5355            adapter.next_event().await,
5356            Some(Ok(AgentEvent::ModesReplaced { .. }))
5357        ));
5358        assert!(matches!(
5359            adapter.next_event().await,
5360            Some(Ok(AgentEvent::Ready { .. }))
5361        ));
5362        adapter.send_prompt("hello".into()).await.expect("prompt");
5363        assert!(matches!(
5364            adapter.next_event().await,
5365            Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5366        ));
5367        assert!(matches!(
5368            adapter.next_event().await,
5369            Some(Ok(AgentEvent::TurnComplete { .. }))
5370        ));
5371    }
5372
5373    #[tokio::test]
5374    async fn acp_output_token_limit_is_reported_as_a_failed_turn() {
5375        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{"sessionId":"session-1"}}'; read _; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"max_tokens"}}'"#;
5376        let cwd = std::env::current_dir().expect("cwd");
5377        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5378        adapter.start().await.expect("initialize");
5379        assert!(matches!(
5380            adapter.next_event().await,
5381            Some(Ok(AgentEvent::Ready { .. }))
5382        ));
5383        adapter.send_prompt("hello".into()).await.expect("prompt");
5384        assert!(matches!(
5385            adapter.next_event().await,
5386            Some(Ok(AgentEvent::Failed {
5387                slot: 1,
5388                started: true,
5389                detail,
5390            })) if detail.contains("output token limit")
5391        ));
5392    }
5393
5394    #[tokio::test]
5395    async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5396        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"}}'"#;
5397        let cwd = std::env::current_dir().expect("cwd");
5398        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5399        adapter.start().await.expect("initialize");
5400        assert!(adapter.next_event().await.is_some());
5401        assert!(adapter.next_event().await.is_some());
5402
5403        for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5404            adapter.send_prompt(prompt.into()).await.expect("prompt");
5405            assert!(matches!(
5406                adapter.next_event().await,
5407                Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5408            ));
5409            assert!(matches!(
5410                adapter.next_event().await,
5411                Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5412            ));
5413        }
5414        adapter.stop().await.expect("stop");
5415    }
5416
5417    #[tokio::test]
5418    async fn empty_acp_mode_catalog_disables_mode_control() {
5419        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":[]}}}'"#;
5420        let cwd = std::env::current_dir().expect("cwd");
5421        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5422        adapter.start().await.expect("initialize");
5423        assert!(!adapter.capabilities().supports_modes);
5424        assert!(matches!(
5425            adapter.next_event().await,
5426            Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5427        ));
5428        assert!(matches!(
5429            adapter.next_event().await,
5430            Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5431        ));
5432        adapter.stop().await.expect("stop");
5433    }
5434
5435    #[tokio::test]
5436    async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5437        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"#;
5438        let cwd = std::env::current_dir().expect("cwd");
5439        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5440        adapter.start().await.expect("initialize");
5441        assert!(adapter.capabilities().supports_models);
5442        assert!(matches!(
5443            adapter.next_event().await,
5444            Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5445                if config_id == "model"
5446                    && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5447                    && current_model.as_deref() == Some("fast")
5448        ));
5449        assert!(matches!(
5450            adapter.next_event().await,
5451            Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5452        ));
5453        adapter.set_model("smart".into()).await.expect("set model");
5454        assert!(adapter.set_model("invented".into()).await.is_err());
5455        adapter.stop().await.expect("stop");
5456    }
5457
5458    #[tokio::test]
5459    async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5460        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":{}}'"#;
5461        let cwd = std::env::current_dir().expect("cwd");
5462        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5463        adapter.start().await.expect("initialize");
5464        adapter
5465            .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5466            .await
5467            .expect("set mode");
5468        assert!(matches!(
5469            adapter.next_event().await,
5470            Some(Ok(AgentEvent::ModesReplaced { .. }))
5471        ));
5472        assert!(matches!(
5473            adapter.next_event().await,
5474            Some(Ok(AgentEvent::Ready { .. }))
5475        ));
5476        assert!(matches!(
5477            adapter.next_event().await,
5478            Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5479        ));
5480        adapter.stop().await.expect("stop");
5481    }
5482
5483    #[tokio::test]
5484    async fn acp_reload_preserves_a_loadable_session_id() {
5485        let cwd = std::env::current_dir().expect("cwd");
5486        let mut adapter = AcpAdapter::with_session_id(
5487            0,
5488            cwd,
5489            "__codeswarm_missing_acp_for_reload_test__",
5490            Vec::new(),
5491            "saved-session",
5492        );
5493        adapter.capabilities.supports_session_load = true;
5494        // A failed replacement process still must not erase the session ID:
5495        // the coordinator can report the startup error and offer another
5496        // reload, preserving the only handle that can resume the conversation.
5497        assert!(adapter.reload().await.is_err());
5498        assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5499    }
5500
5501    #[tokio::test]
5502    async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5503        let cwd = std::env::current_dir().expect("cwd");
5504        let mut adapter = AcpAdapter::with_session_id(
5505            0,
5506            cwd,
5507            "__codeswarm_missing_nonloadable_acp__",
5508            Vec::new(),
5509            "stale-session",
5510        );
5511        adapter.capabilities.supports_session_load = false;
5512        assert!(adapter.reload().await.is_err());
5513        assert_eq!(adapter.session_id, None);
5514    }
5515
5516    #[tokio::test]
5517    async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5518        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"}}'"#;
5519        let cwd = std::env::current_dir().expect("cwd");
5520        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5521        adapter.start().await.expect("initialize");
5522        assert!(matches!(
5523            adapter.next_event().await,
5524            Some(Ok(AgentEvent::Ready { .. }))
5525        ));
5526        adapter.send_prompt("hello".into()).await.expect("prompt");
5527        assert!(matches!(
5528            adapter.next_event().await,
5529            Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5530        ));
5531        assert!(matches!(
5532            adapter.next_event().await,
5533            Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5534        ));
5535    }
5536
5537    #[test]
5538    fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5539        let mut tools = std::collections::BTreeMap::new();
5540        let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5541            "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5542        let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5543        assert_eq!(initial.detail.as_deref(), Some("old output"));
5544        let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5545        assert_eq!(completed.title, "Read config");
5546        assert_eq!(completed.detail.as_deref(), Some("old output"));
5547        assert_eq!(completed.status, ToolStatus::Completed);
5548        let malformed = super::normalize_acp_tool(
5549            &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5550            &mut tools,
5551        )
5552        .unwrap();
5553        assert_eq!(malformed, completed);
5554        let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5555        assert_eq!(replaced.detail.as_deref(), Some("new output"));
5556        let cleared = super::normalize_acp_tool(
5557            &serde_json::json!({"toolCallId":"read","content":[]}),
5558            &mut tools,
5559        )
5560        .unwrap();
5561        assert_eq!(cleared.detail, None);
5562        let raw = super::normalize_acp_tool(
5563            &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5564            &mut tools,
5565        )
5566        .unwrap();
5567        assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5568        let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5569        assert_eq!(fresh.status, ToolStatus::Pending);
5570        assert_eq!(fresh.detail, None);
5571        for invalid in [
5572            serde_json::json!({}),
5573            serde_json::json!({"toolCallId":7}),
5574            serde_json::json!({"toolCallId":" "}),
5575        ] {
5576            assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5577        }
5578        assert_eq!(tools.len(), 1);
5579        // IDs are opaque, not whitespace-normalized aliases of another tool.
5580        super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5581        assert_eq!(tools.len(), 2);
5582    }
5583
5584    #[tokio::test]
5585    async fn acp_tool_status_only_notifications_retain_name_and_output() {
5586        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5587read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5588read _
5589echo '{"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"}}]}}}'
5590echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5591echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5592        let mut adapter = AcpAdapter::new(
5593            0,
5594            std::env::current_dir().unwrap(),
5595            "sh",
5596            vec!["-c".into(), script.into()],
5597        );
5598        adapter.start().await.unwrap();
5599        adapter.next_event().await.unwrap().unwrap();
5600        adapter.send_prompt("read".into()).await.unwrap();
5601        for status in [ToolStatus::Running, ToolStatus::Completed] {
5602            let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5603                panic!("tool event");
5604            };
5605            assert_eq!(update.status, status);
5606            assert_eq!(update.title, "Read config");
5607            assert_eq!(update.detail.as_deref(), Some("file content"));
5608        }
5609        assert!(matches!(
5610            adapter.next_event().await,
5611            Some(Ok(AgentEvent::TurnComplete { .. }))
5612        ));
5613        adapter.stop().await.unwrap();
5614    }
5615
5616    #[tokio::test]
5617    async fn acp_reload_discards_old_queued_events_and_catalogs() {
5618        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5619        let mut adapter = AcpAdapter::new(
5620            0,
5621            std::env::current_dir().unwrap(),
5622            "sh",
5623            vec!["-c".into(), script.into()],
5624        );
5625        adapter.start().await.unwrap();
5626        adapter.queued_events.push_back(Ok(AgentEvent::Text {
5627            slot: 0,
5628            text: "stale".into(),
5629        }));
5630        adapter.modes = vec![Mode {
5631            id: "stale".into(),
5632            label: "Stale".into(),
5633        }];
5634        // Real peers echo each request's ID. Reset for this fixed-ID test script.
5635        adapter.next_request_id = 1;
5636        adapter.reload().await.unwrap();
5637        assert!(adapter.modes.is_empty());
5638        assert_eq!(adapter.queued_events.len(), 1);
5639        assert!(matches!(
5640            adapter.next_event().await,
5641            Some(Ok(AgentEvent::Ready { .. }))
5642        ));
5643        adapter.stop().await.unwrap();
5644        assert!(adapter.queued_events.is_empty());
5645    }
5646
5647    #[tokio::test]
5648    async fn acp_load_replays_history_without_starting_a_turn() {
5649        let script = r#"
5650read _
5651echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5652read request
5653case "$request" in *session/load*) ;; *) exit 2;; esac
5654echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5655echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5656echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5657echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5658echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5659echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5660read request
5661case "$request" in *session/prompt*) ;; *) exit 3;; esac
5662echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5663echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5664"#;
5665        let mut adapter = AcpAdapter::with_session_id(
5666            2,
5667            std::env::current_dir().unwrap(),
5668            "sh",
5669            vec!["-c".into(), script.into()],
5670            "saved",
5671        );
5672        adapter.start().await.unwrap();
5673        let mut state = crate::SessionState::new(3);
5674        for _ in 0..5 {
5675            let event = adapter.next_event().await.unwrap().unwrap();
5676            assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5677            crate::reduce(&mut state, event);
5678            assert_eq!(state.active_slot, None);
5679        }
5680        assert!(matches!(
5681            adapter.next_event().await,
5682            Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5683        ));
5684        adapter.send_prompt("new question".into()).await.unwrap();
5685        assert!(
5686            matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5687        );
5688        assert!(matches!(
5689            adapter.next_event().await,
5690            Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5691        ));
5692        adapter.stop().await.unwrap();
5693    }
5694
5695    #[tokio::test]
5696    async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5697        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5698        let cwd = std::env::current_dir().expect("cwd");
5699        let mut adapter = AcpAdapter::with_session_id(
5700            0,
5701            cwd,
5702            "sh",
5703            vec!["-c".into(), script.into()],
5704            "existing-session",
5705        );
5706        adapter.start().await.expect("load existing session");
5707        assert!(matches!(
5708            adapter.next_event().await,
5709            Some(Ok(AgentEvent::Ready { .. }))
5710        ));
5711    }
5712
5713    #[tokio::test]
5714    async fn acp_start_failure_reaps_transport_process() {
5715        // The child emits an invalid initialize response and exits. The
5716        // adapter must not retain a live child after protocol startup fails;
5717        // this is the path used when a configured ACP command is unavailable
5718        // or speaks a different protocol.
5719        let mut adapter = AcpAdapter::new(
5720            0,
5721            std::env::current_dir().expect("cwd"),
5722            "sh",
5723            vec!["-c".into(), "printf 'not-json\\n'".into()],
5724        );
5725        assert!(adapter.start().await.is_err());
5726        assert!(adapter.child.is_none());
5727        assert!(adapter.reader.is_none());
5728    }
5729
5730    #[tokio::test]
5731    async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5732        let marker = unique_test_path("codeswarm-acp-retry", "count");
5733        let script = format!(
5734            r#"count=0
5735if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5736count=$((count + 1))
5737printf '%s' "$count" > '{0}'
5738while IFS= read -r request; do
5739  id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5740  case "$request" in
5741    *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5742    *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5743    *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5744    *session/prompt*)
5745      if [ "$count" = 1 ]; then exit 0; fi
5746      printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5747      printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5748      ;;
5749  esac
5750done
5751"#,
5752            marker.display()
5753        );
5754        let mut adapter = AcpAdapter::new(
5755            0,
5756            std::env::current_dir().expect("cwd"),
5757            "sh",
5758            vec!["-c".into(), script],
5759        );
5760        adapter.start().await.expect("initial ACP startup");
5761        assert!(matches!(
5762            adapter.next_event().await,
5763            Some(Ok(AgentEvent::Ready { .. }))
5764        ));
5765        adapter
5766            .send_prompt("first".into())
5767            .await
5768            .expect("first prompt");
5769        assert!(matches!(
5770            adapter.next_event().await,
5771            Some(Err(AdapterError::Transport(_)))
5772        ));
5773        assert!(adapter.child.is_none());
5774        assert!(adapter.reader.is_none());
5775        assert_eq!(adapter.session_id(), Some("saved-session".into()));
5776
5777        adapter.reload().await.expect("reload ACP transport");
5778        assert!(matches!(
5779            adapter.next_event().await,
5780            Some(Ok(AgentEvent::Ready { .. }))
5781        ));
5782        adapter
5783            .send_prompt("retry".into())
5784            .await
5785            .expect("retry prompt");
5786        assert!(matches!(
5787            adapter.next_event().await,
5788            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5789        ));
5790        assert!(matches!(
5791            adapter.next_event().await,
5792            Some(Ok(AgentEvent::TurnComplete { .. }))
5793        ));
5794        adapter.stop().await.expect("stop ACP");
5795        std::fs::remove_file(marker).expect("cleanup marker");
5796    }
5797
5798    #[tokio::test]
5799    async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5800        let prompts = Arc::new(Mutex::new(Vec::new()));
5801        let healthy = ScriptedAdapter::new(
5802            0,
5803            AgentCapabilities::default(),
5804            [
5805                AgentEvent::Text {
5806                    slot: 0,
5807                    text: "peer context".into(),
5808                },
5809                AgentEvent::TurnComplete { slot: 0 },
5810            ],
5811        );
5812        let probe = ReloadProbeAdapter {
5813            slot: 1,
5814            crashed: false,
5815            reloaded: false,
5816            events: VecDeque::new(),
5817            prompts: Arc::clone(&prompts),
5818        };
5819        let mut relay = RelayHost::new(
5820            vec![
5821                AdapterHost::new(Box::new(healthy), None),
5822                AdapterHost::new(Box::new(probe), None),
5823            ],
5824            8,
5825        )
5826        .expect("relay");
5827        relay.start().await.expect("start");
5828        relay
5829            .run_turn("original task", 0)
5830            .await
5831            .expect("first turn");
5832        relay.run_turn("", 0).await.expect("crashed turn");
5833        relay.relay_mut().enqueue_human("retry", Some(1));
5834        relay.run_turn("", 0).await.expect("reloaded turn");
5835
5836        let prompts = prompts.lock().expect("prompts");
5837        let retry = prompts.last().expect("retry prompt");
5838        assert!(retry.contains("You are Reload probe"), "{retry}");
5839        assert!(retry.contains("original task"), "{retry}");
5840        assert!(retry.contains("peer context"), "{retry}");
5841        assert!(retry.contains("retry"), "{retry}");
5842    }
5843
5844    #[tokio::test]
5845    async fn acp_adapter_answers_permission_json_rpc_requests() {
5846        let path = std::env::temp_dir().join(format!(
5847            "codeswarm-permission-answer-{}",
5848            std::process::id()
5849        ));
5850        let script = format!(
5851            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"}}}}'"#,
5852            path.display()
5853        );
5854        let mut adapter = AcpAdapter::new(
5855            0,
5856            std::env::current_dir().expect("cwd"),
5857            "sh",
5858            vec!["-c".into(), script],
5859        );
5860        adapter.start().await.expect("start ACP");
5861        assert!(matches!(
5862            adapter.next_event().await,
5863            Some(Ok(AgentEvent::Ready { .. }))
5864        ));
5865        adapter.send_prompt("do it".into()).await.expect("prompt");
5866        assert!(matches!(
5867            adapter.next_event().await,
5868            Some(Ok(AgentEvent::Permission { request, .. }))
5869                if request.id == "9"
5870                    && request.options == ["Allow once"]
5871                    && request.option_ids == ["allow-once"]
5872        ));
5873        adapter
5874            .answer_permission(
5875                "9".into(),
5876                PermissionAnswer::Selected {
5877                    option_id: "allow-once".into(),
5878                },
5879            )
5880            .await
5881            .expect("permission answer");
5882        assert!(matches!(
5883            adapter.next_event().await,
5884            Some(Ok(AgentEvent::TurnComplete { .. }))
5885        ));
5886        let answer: Value = serde_json::from_str(
5887            &std::fs::read_to_string(&path).expect("captured permission answer"),
5888        )
5889        .expect("valid JSON-RPC answer");
5890        assert_eq!(answer["id"], 9);
5891        assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5892        assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5893        std::fs::remove_file(path).expect("cleanup");
5894    }
5895
5896    #[test]
5897    fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5898        let event = parse_acp_notification(
5899            0,
5900            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5901        )
5902        .expect("valid JSON-RPC request");
5903        assert!(event.is_none());
5904    }
5905
5906    #[tokio::test]
5907    async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5908        let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5909        std::fs::write(
5910            &script_path,
5911            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5912        )
5913        .expect("write native test script");
5914        let mut adapter = AgyAdapter::new(
5915            0,
5916            std::env::current_dir().expect("cwd"),
5917            format!("sh {}", script_path.display()),
5918        );
5919        adapter.start().await.expect("start native adapter");
5920        assert!(adapter.next_event().await.is_some());
5921        assert!(adapter.next_event().await.is_some());
5922        adapter
5923            .send_prompt("continue".into())
5924            .await
5925            .expect("prompt");
5926        assert!(matches!(
5927            adapter.next_event().await,
5928            Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5929        ));
5930        assert!(matches!(
5931            adapter.next_event().await,
5932            Some(Ok(AgentEvent::TurnComplete { .. }))
5933        ));
5934        adapter.stop().await.expect("stop native adapter");
5935        std::fs::remove_file(script_path).expect("cleanup native script");
5936    }
5937
5938    #[test]
5939    fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5940        let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5941        let _ = std::fs::remove_dir_all(&root);
5942        std::fs::create_dir_all(&root).expect("workspace");
5943        std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5944        let outside =
5945            std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5946        std::fs::write(&outside, "secret").expect("outside file");
5947        let link = root.join("outside-link");
5948        #[cfg(unix)]
5949        std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5950        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5951
5952        assert_eq!(
5953            adapter
5954                .read_workspace_text("inside.txt", Some(2), Some(1))
5955                .expect("read inside"),
5956            "two"
5957        );
5958        std::fs::write(
5959            root.join("large.txt"),
5960            vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5961        )
5962        .expect("large file");
5963        let bounded = adapter
5964            .read_workspace_text("large.txt", None, None)
5965            .expect("bounded read");
5966        assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5967        #[cfg(unix)]
5968        {
5969            std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5970                .expect("internal symlink");
5971            assert_eq!(
5972                adapter
5973                    .read_workspace_text("inside-link", None, None)
5974                    .expect("read internal symlink"),
5975                "one\ntwo\nthree\n"
5976            );
5977        }
5978        assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5979        assert!(
5980            adapter
5981                .workspace_path(&outside.display().to_string())
5982                .is_err()
5983        );
5984        #[cfg(unix)]
5985        assert!(adapter.workspace_path("outside-link").is_err());
5986        #[cfg(unix)]
5987        std::fs::remove_file(link).expect("cleanup symlink");
5988        #[cfg(unix)]
5989        std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5990        std::fs::remove_file(outside).expect("cleanup outside");
5991        std::fs::remove_dir_all(root).expect("cleanup workspace");
5992    }
5993
5994    #[tokio::test]
5995    async fn running_terminal_output_omits_exit_status_until_completion() {
5996        let root = unique_test_path("codeswarm-terminal-output", "dir");
5997        std::fs::create_dir_all(&root).expect("workspace");
5998        let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5999        let result = adapter
6000            .terminal_create(&serde_json::json!({
6001                "command": "sh",
6002                "args": ["-c", "sleep 0.2; printf done"],
6003                "cwd": ".",
6004            }))
6005            .await
6006            .expect("terminal create");
6007        let id = result["terminalId"].as_str().expect("terminal id");
6008        let output = adapter.terminal_output(id).await.expect("terminal output");
6009        assert!(output.get("exitStatus").is_none());
6010        if let Some(terminal) = adapter.terminals.remove(id) {
6011            terminal.stop().await;
6012        }
6013        std::fs::remove_dir_all(root).expect("cleanup workspace");
6014    }
6015
6016    #[tokio::test]
6017    async fn acp_adapter_answers_workspace_read_requests() {
6018        let root =
6019            std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
6020        let _ = std::fs::remove_dir_all(&root);
6021        std::fs::create_dir_all(&root).expect("workspace");
6022        let source = root.join("inside.txt");
6023        let answer = root.join("answer.json");
6024        std::fs::write(&source, "workspace content").expect("source");
6025        let script = format!(
6026            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"}}}}'"#,
6027            source.display(),
6028            answer.display(),
6029        );
6030        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6031        adapter.start().await.expect("start ACP");
6032        assert!(matches!(
6033            adapter.next_event().await,
6034            Some(Ok(AgentEvent::Ready { .. }))
6035        ));
6036        adapter.send_prompt("read it".into()).await.expect("prompt");
6037        assert!(matches!(
6038            adapter.next_event().await,
6039            Some(Ok(AgentEvent::TurnComplete { .. }))
6040        ));
6041        let response: Value =
6042            serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
6043                .expect("response JSON");
6044        assert_eq!(response["id"], 9);
6045        assert_eq!(response["result"]["content"], "workspace content");
6046        adapter.stop().await.expect("stop ACP");
6047        std::fs::remove_dir_all(root).expect("cleanup workspace");
6048    }
6049
6050    #[tokio::test]
6051    async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
6052        let root =
6053            std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
6054        let _ = std::fs::remove_dir_all(&root);
6055        std::fs::create_dir_all(&root).expect("workspace");
6056        let create_request = serde_json::json!({
6057            "jsonrpc": "2.0",
6058            "id": 9,
6059            "method": "terminal/create",
6060            "params": {
6061                "sessionId": "s1",
6062                "command": "sh",
6063                "args": ["-c", "sleep 0.1; printf terminal-ok"],
6064                "cwd": ".",
6065            },
6066        });
6067        let wait_request = serde_json::json!({
6068            "jsonrpc": "2.0",
6069            "id": 10,
6070            "method": "terminal/wait_for_exit",
6071            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6072        });
6073        let output_request = serde_json::json!({
6074            "jsonrpc": "2.0",
6075            "id": 11,
6076            "method": "terminal/output",
6077            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6078        });
6079        let create_answer = root.join("create-answer.json");
6080        let wait_answer = root.join("wait-answer.json");
6081        let output_answer = root.join("output-answer.json");
6082        let script = format!(
6083            "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\"}}}}'",
6084            create_request,
6085            create_answer.display(),
6086            wait_request,
6087            wait_answer.display(),
6088            output_request,
6089            output_answer.display(),
6090        );
6091        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6092        adapter.start().await.expect("start ACP");
6093        assert!(matches!(
6094            adapter.next_event().await,
6095            Some(Ok(AgentEvent::Ready { .. }))
6096        ));
6097        adapter
6098            .send_prompt("run terminal".into())
6099            .await
6100            .expect("prompt");
6101        let mut saw_complete = false;
6102        for _ in 0..6 {
6103            match adapter.next_event().await {
6104                Some(Ok(AgentEvent::TurnComplete { .. })) => {
6105                    saw_complete = true;
6106                    break;
6107                }
6108                Some(_) => {}
6109                None => break,
6110            }
6111        }
6112        assert!(saw_complete, "terminal requests should not stall ACP");
6113        let create: Value = serde_json::from_str(
6114            &std::fs::read_to_string(&create_answer).expect("captured create response"),
6115        )
6116        .expect("create JSON");
6117        assert_eq!(create["result"]["terminalId"], "terminal-1");
6118        let output: Value = serde_json::from_str(
6119            &std::fs::read_to_string(&output_answer).expect("captured output response"),
6120        )
6121        .expect("output JSON");
6122        assert!(
6123            output["result"]["output"]
6124                .as_str()
6125                .unwrap_or_default()
6126                .contains("terminal-ok"),
6127            "output response: {output}"
6128        );
6129        adapter.stop().await.expect("stop ACP");
6130        std::fs::remove_dir_all(root).expect("cleanup workspace");
6131    }
6132
6133    #[tokio::test]
6134    async fn host_reduces_and_persists_adapter_events() {
6135        let path =
6136            std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6137        let adapter = ScriptedAdapter::new(
6138            0,
6139            AgentCapabilities::default(),
6140            [AgentEvent::Text {
6141                slot: 0,
6142                text: "hello".into(),
6143            }],
6144        );
6145        let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6146        host.start().await.expect("start");
6147        host.next_effects()
6148            .await
6149            .expect("event")
6150            .expect("valid event");
6151        assert_eq!(host.state.public_text[0].1, "hello");
6152        assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6153        std::fs::remove_file(path).expect("cleanup");
6154    }
6155
6156    #[tokio::test]
6157    async fn relay_applies_default_policy_before_the_first_prompt() {
6158        let first_log = Arc::new(Mutex::new(Vec::new()));
6159        let second_log = Arc::new(Mutex::new(Vec::new()));
6160        let first = AdapterHost::new(
6161            Box::new(ModeOrderAdapter {
6162                slot: 0,
6163                log: Arc::clone(&first_log),
6164                phase: 0,
6165            }),
6166            None,
6167        );
6168        let second = AdapterHost::new(
6169            Box::new(ModeOrderAdapter {
6170                slot: 1,
6171                log: Arc::clone(&second_log),
6172                phase: 0,
6173            }),
6174            None,
6175        );
6176        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6177        relay.start().await.expect("start and synchronize policy");
6178        relay.run_turn("task", 0).await.expect("first turn");
6179        {
6180            let log = first_log.lock().expect("log");
6181            assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6182        }
6183        assert_eq!(
6184            second_log.lock().expect("log").as_slice(),
6185            ["start", "mode:yolo"]
6186        );
6187        let added_log = Arc::new(Mutex::new(Vec::new()));
6188        relay
6189            .add_agent(
6190                AdapterHost::new(
6191                    Box::new(ModeOrderAdapter {
6192                        slot: 2,
6193                        log: Arc::clone(&added_log),
6194                        phase: 0,
6195                    }),
6196                    None,
6197                ),
6198                "Added",
6199                "added.example",
6200                "added-agent",
6201            )
6202            .await
6203            .expect("add with synchronized policy");
6204        assert_eq!(
6205            added_log.lock().expect("log").as_slice(),
6206            ["start", "mode:yolo"]
6207        );
6208        relay.drop_agent(2).await.expect("drop added agent");
6209        added_log.lock().expect("log").clear();
6210        relay
6211            .reload(2)
6212            .await
6213            .expect("reload with synchronized policy");
6214        assert_eq!(
6215            added_log.lock().expect("log").as_slice(),
6216            ["reload", "mode:yolo"]
6217        );
6218    }
6219
6220    #[tokio::test]
6221    async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6222        let hosts = (0..2)
6223            .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6224            .collect::<Vec<_>>();
6225        let startup_events = Arc::new(Mutex::new(Vec::new()));
6226        let captured = Arc::clone(&startup_events);
6227        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6228        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6229
6230        relay.start().await.expect("complete startup handshake");
6231
6232        assert!(relay.dispatches().is_empty());
6233        let ready_slots = startup_events
6234            .lock()
6235            .expect("events")
6236            .iter()
6237            .filter_map(|event| match event {
6238                AgentEvent::Ready { slot, .. } => Some(*slot),
6239                _ => None,
6240            })
6241            .collect::<Vec<_>>();
6242        assert_eq!(ready_slots, vec![0, 1]);
6243    }
6244
6245    #[tokio::test]
6246    async fn independent_roster_adapters_start_concurrently() {
6247        let barrier = Arc::new(tokio::sync::Barrier::new(2));
6248        let hosts = (0..2)
6249            .map(|slot| {
6250                AdapterHost::new(
6251                    Box::new(ConcurrentStartAdapter {
6252                        slot,
6253                        barrier: Arc::clone(&barrier),
6254                    }),
6255                    None,
6256                )
6257            })
6258            .collect::<Vec<_>>();
6259        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6260        tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6261            .await
6262            .expect("startup should not serialize barrier participants")
6263            .expect("startup succeeds");
6264    }
6265
6266    #[tokio::test]
6267    async fn relay_host_dispatches_turns_sequentially() {
6268        let capabilities = AgentCapabilities {
6269            supports_cancel: true,
6270            ..AgentCapabilities::default()
6271        };
6272        let first = ScriptedAdapter::new(
6273            0,
6274            capabilities.clone(),
6275            [
6276                AgentEvent::Text {
6277                    slot: 0,
6278                    text: "first".into(),
6279                },
6280                AgentEvent::TurnComplete { slot: 0 },
6281            ],
6282        );
6283        let second = ScriptedAdapter::new(
6284            1,
6285            capabilities,
6286            [
6287                AgentEvent::Text {
6288                    slot: 1,
6289                    text: "review".into(),
6290                },
6291                AgentEvent::TurnComplete { slot: 1 },
6292            ],
6293        );
6294        let hosts = vec![
6295            AdapterHost::new(Box::new(first), None),
6296            AdapterHost::new(Box::new(second), None),
6297        ];
6298        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6299        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6300        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6301        let captured = std::sync::Arc::clone(&events);
6302        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6303        relay.start().await.expect("start");
6304        events.lock().expect("events").clear();
6305        assert!(matches!(
6306            relay.run_turn("task", 0).await.expect("first turn"),
6307            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6308        ));
6309        assert!(matches!(
6310            relay.run_turn("first", 0).await.expect("second turn"),
6311            crate::relay::RelayDecision::Dispatch {
6312                slot: 1,
6313                can_stop: true,
6314                ..
6315            }
6316        ));
6317        assert_eq!(
6318            relay
6319                .dispatches()
6320                .iter()
6321                .map(|(slot, _)| *slot)
6322                .collect::<Vec<_>>(),
6323            [0, 1]
6324        );
6325        assert!(relay.dispatches()[0].1.contains("You are Claude"));
6326        assert!(
6327            relay.dispatches()[0]
6328                .1
6329                .contains("CodeSwarm roster (ordered)")
6330        );
6331        assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6332        assert!(relay.dispatches()[0].1.contains("2. Codex"));
6333        assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6334        assert!(relay.dispatches()[0].1.contains("Do not use"));
6335        let lifecycle = events.lock().expect("events");
6336        let positions = lifecycle
6337            .iter()
6338            .filter_map(|event| match event {
6339                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6340                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6341                _ => None,
6342            })
6343            .collect::<Vec<_>>();
6344        assert_eq!(
6345            positions,
6346            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6347        );
6348    }
6349
6350    #[tokio::test]
6351    async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6352        let stops = Arc::new(AtomicUsize::new(0));
6353        let failed = FailingStartAdapter {
6354            slot: 0,
6355            stopped: stops.clone(),
6356        };
6357        let healthy = ScriptedAdapter::new(
6358            1,
6359            AgentCapabilities::default(),
6360            [
6361                AgentEvent::Text {
6362                    slot: 1,
6363                    text: "healthy response".into(),
6364                },
6365                AgentEvent::TurnComplete { slot: 1 },
6366            ],
6367        );
6368        let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6369        let captured = events.clone();
6370        let mut relay = RelayHost::new(
6371            vec![
6372                AdapterHost::new(Box::new(failed), None),
6373                AdapterHost::new(Box::new(healthy), None),
6374            ],
6375            4,
6376        )
6377        .unwrap();
6378        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6379        relay.start_resuming().await.unwrap();
6380        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6381        assert!(relay.dispatches().is_empty());
6382        assert_eq!(stops.load(Ordering::Relaxed), 1);
6383        assert!(
6384            events
6385                .lock()
6386                .unwrap()
6387                .iter()
6388                .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6389        );
6390        assert!(
6391            !events
6392                .lock()
6393                .unwrap()
6394                .iter()
6395                .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6396        );
6397        assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6398        assert!(
6399            relay
6400                .relay_mut()
6401                .enqueue_human("explicit healthy target", Some(1))
6402        );
6403        assert!(matches!(
6404            relay.run_turn("", 1).await.unwrap(),
6405            RelayDecision::Dispatch { slot: 1, .. }
6406        ));
6407        relay.stop().await.unwrap();
6408    }
6409
6410    #[tokio::test]
6411    async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6412        let capabilities = AgentCapabilities::default();
6413        let first = ScriptedAdapter::new(
6414            0,
6415            capabilities.clone(),
6416            [
6417                AgentEvent::Text {
6418                    slot: 0,
6419                    text: "implemented".into(),
6420                },
6421                AgentEvent::TurnComplete { slot: 0 },
6422                AgentEvent::Text {
6423                    slot: 0,
6424                    text: format!("fixed review findings {STOP_TOKEN}"),
6425                },
6426                AgentEvent::TurnComplete { slot: 0 },
6427            ],
6428        );
6429        let second = ScriptedAdapter::new(
6430            1,
6431            capabilities,
6432            [
6433                AgentEvent::Text {
6434                    slot: 1,
6435                    text: "reviewed".into(),
6436                },
6437                AgentEvent::TurnComplete { slot: 1 },
6438                AgentEvent::Text {
6439                    slot: 1,
6440                    text: format!("approved {STOP_TOKEN}"),
6441                },
6442                AgentEvent::TurnComplete { slot: 1 },
6443                AgentEvent::Text {
6444                    slot: 1,
6445                    text: "new task".into(),
6446                },
6447                AgentEvent::TurnComplete { slot: 1 },
6448            ],
6449        );
6450        let hosts = vec![
6451            AdapterHost::new(Box::new(first), None),
6452            AdapterHost::new(Box::new(second), None),
6453        ];
6454        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6455        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6456        relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6457        relay.start().await.expect("start");
6458        assert!(matches!(
6459            relay.run_turn("task", 0).await.expect("implementer turn"),
6460            RelayDecision::Dispatch {
6461                slot: 0,
6462                can_stop: false,
6463                ..
6464            }
6465        ));
6466        let implementer_prompt = &relay.dispatches()[0].1;
6467        assert!(implementer_prompt.contains("you are the implementer"));
6468        assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6469        assert!(!implementer_prompt.contains("you are the reviewer"));
6470        assert!(implementer_prompt.contains("Do not use"));
6471        assert!(matches!(
6472            relay.run_turn("", 0).await.expect("reviewer turn"),
6473            RelayDecision::Dispatch {
6474                slot: 1,
6475                can_stop: true,
6476                ..
6477            }
6478        ));
6479        let reviewer_prompt = &relay.dispatches()[1].1;
6480        assert!(reviewer_prompt.contains("you are the reviewer"));
6481        assert!(reviewer_prompt.contains("Claude handed off"));
6482        assert!(reviewer_prompt.contains("concrete defects"));
6483        assert!(reviewer_prompt.contains("concise approval"));
6484        assert!(reviewer_prompt.contains(STOP_TOKEN));
6485        assert!(!reviewer_prompt.contains("you are the implementer"));
6486        relay.run_turn("", 0).await.unwrap();
6487        assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6488        assert!(relay.dispatches()[2].1.contains("Do not use"));
6489        assert!(matches!(
6490            relay.run_turn("", 0).await.unwrap(),
6491            RelayDecision::Dispatch { slot: 1, .. }
6492        ));
6493        assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6494        assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6495        relay.run_turn("", 1).await.unwrap();
6496        assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6497    }
6498
6499    #[tokio::test]
6500    async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6501        let solo = ScriptedAdapter::new(
6502            0,
6503            AgentCapabilities::default(),
6504            [
6505                AgentEvent::Text {
6506                    slot: 0,
6507                    text: "solo".into(),
6508                },
6509                AgentEvent::TurnComplete { slot: 0 },
6510            ],
6511        );
6512        let mut solo_relay =
6513            RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6514        solo_relay
6515            .relay_mut()
6516            .set_strategy(CollaborationStrategy::Pair);
6517        solo_relay.start().await.expect("start");
6518        solo_relay.run_turn("task", 0).await.expect("solo turn");
6519        assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6520
6521        let roster_first = ScriptedAdapter::new(
6522            0,
6523            AgentCapabilities::default(),
6524            [AgentEvent::TurnComplete { slot: 0 }],
6525        );
6526        let roster_second = ScriptedAdapter::new(
6527            1,
6528            AgentCapabilities::default(),
6529            [AgentEvent::TurnComplete { slot: 1 }],
6530        );
6531        let mut roster = RelayHost::new(
6532            vec![
6533                AdapterHost::new(Box::new(roster_first), None),
6534                AdapterHost::new(Box::new(roster_second), None),
6535            ],
6536            4,
6537        )
6538        .expect("relay");
6539        roster.start().await.expect("start");
6540        roster.run_turn("task", 0).await.expect("first turn");
6541        roster.run_turn("", 0).await.expect("second turn");
6542        assert!(!roster.dispatches()[0].1.contains("Pair role"));
6543        assert!(!roster.dispatches()[1].1.contains("Pair role"));
6544
6545        let pair_first = ScriptedAdapter::new(
6546            0,
6547            AgentCapabilities::default(),
6548            [AgentEvent::TurnComplete { slot: 0 }],
6549        );
6550        let pair_second = ScriptedAdapter::new(
6551            1,
6552            AgentCapabilities::default(),
6553            [AgentEvent::TurnComplete { slot: 1 }],
6554        );
6555        let mut pair = RelayHost::new(
6556            vec![
6557                AdapterHost::new(Box::new(pair_first), None),
6558                AdapterHost::new(Box::new(pair_second), None),
6559            ],
6560            4,
6561        )
6562        .expect("relay");
6563        pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6564        assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6565        pair.start().await.expect("start");
6566        assert!(matches!(
6567            pair.run_turn("ignored", 0).await.expect("direct turn"),
6568            RelayDecision::Dispatch {
6569                slot: 1,
6570                direct: true,
6571                ..
6572            }
6573        ));
6574        let direct_prompt = &pair.dispatches()[0].1;
6575        assert!(direct_prompt.contains("private"));
6576        assert!(!direct_prompt.contains("Pair role"));
6577    }
6578
6579    #[tokio::test]
6580    async fn relay_host_routes_around_a_usage_limited_agent() {
6581        let capabilities = AgentCapabilities::default();
6582        let first = ScriptedAdapter::new(
6583            0,
6584            capabilities.clone(),
6585            [
6586                AgentEvent::Text {
6587                    slot: 0,
6588                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6589                           credits or try again later."
6590                        .into(),
6591                },
6592                AgentEvent::TurnComplete { slot: 0 },
6593            ],
6594        );
6595        let second = ScriptedAdapter::new(
6596            1,
6597            capabilities,
6598            [
6599                AgentEvent::Text {
6600                    slot: 1,
6601                    text: "review done".into(),
6602                },
6603                AgentEvent::TurnComplete { slot: 1 },
6604            ],
6605        );
6606        let hosts = vec![
6607            AdapterHost::new(Box::new(first), None),
6608            AdapterHost::new(Box::new(second), None),
6609        ];
6610        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6611        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6612        let captured = std::sync::Arc::clone(&events);
6613        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6614        relay.start().await.expect("start");
6615        events.lock().expect("events").clear();
6616        assert!(matches!(
6617            relay.run_turn("task", 0).await.expect("limited turn"),
6618            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6619        ));
6620        assert!(
6621            events
6622                .lock()
6623                .expect("events")
6624                .iter()
6625                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6626        );
6627        // The next automatic turn skips the limited agent entirely.
6628        assert!(matches!(
6629            relay.run_turn("", 0).await.expect("next turn"),
6630            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6631        ));
6632        assert!(relay.relay().is_limited(0));
6633        // A reload restores the agent to the ring. (ScriptedAdapter cannot
6634        // feed further turns, so the restored routing itself is covered by
6635        // the relay unit tests.)
6636        relay.reload(0).await.expect("reload");
6637        assert!(!relay.relay().is_limited(0));
6638    }
6639
6640    #[tokio::test]
6641    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6642        let limited = ScriptedAdapter::new(
6643            0,
6644            AgentCapabilities::default(),
6645            [AgentEvent::Failed {
6646                slot: 0,
6647                started: true,
6648                detail: "request failed: insufficient_quota".into(),
6649            }],
6650        );
6651        let healthy = ScriptedAdapter::new(
6652            1,
6653            AgentCapabilities::default(),
6654            [AgentEvent::TurnComplete { slot: 1 }],
6655        );
6656        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6657        let captured = std::sync::Arc::clone(&events);
6658        let mut relay = RelayHost::new(
6659            vec![
6660                AdapterHost::new(Box::new(limited), None),
6661                AdapterHost::new(Box::new(healthy), None),
6662            ],
6663            4,
6664        )
6665        .expect("relay");
6666        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6667        relay.start().await.expect("start");
6668        events.lock().expect("events").clear();
6669
6670        assert!(matches!(
6671            relay.run_turn("task", 0).await.expect("limited failure"),
6672            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6673        ));
6674        assert!(relay.relay().is_limited(0));
6675        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6676        {
6677            let events = events.lock().expect("events");
6678            assert!(
6679                events
6680                    .iter()
6681                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6682            );
6683            assert!(
6684                !events
6685                    .iter()
6686                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
6687            );
6688        }
6689
6690        assert!(matches!(
6691            relay.run_turn("", 0).await.expect("healthy peer"),
6692            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6693        ));
6694    }
6695
6696    #[tokio::test]
6697    async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6698        let failed = ScriptedAdapter::new(
6699            0,
6700            AgentCapabilities::default(),
6701            [AgentEvent::Failed {
6702                slot: 0,
6703                started: true,
6704                detail: "connection lost".into(),
6705            }],
6706        );
6707        let healthy = ScriptedAdapter::new(
6708            1,
6709            AgentCapabilities::default(),
6710            [AgentEvent::TurnComplete { slot: 1 }],
6711        );
6712        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6713        let captured = std::sync::Arc::clone(&events);
6714        let mut relay = RelayHost::new(
6715            vec![
6716                AdapterHost::new(Box::new(failed), None),
6717                AdapterHost::new(Box::new(healthy), None),
6718            ],
6719            4,
6720        )
6721        .expect("relay");
6722        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6723        relay.start().await.expect("start");
6724
6725        assert!(matches!(
6726            relay.run_turn("task", 0).await.expect("handled failure"),
6727            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6728        ));
6729        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6730        assert!(relay.relay().is_limited(0));
6731        assert!(events.lock().expect("lock").iter().any(|event| {
6732            matches!(
6733                event,
6734                AgentEvent::Failed {
6735                    slot: 0,
6736                    started: true,
6737                    ..
6738                }
6739            )
6740        }));
6741        assert!(matches!(
6742            relay.run_turn("", 0).await.expect("healthy peer"),
6743            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6744        ));
6745    }
6746
6747    #[tokio::test]
6748    async fn codex_stop_does_not_skip_later_roster_reviewers() {
6749        let hosts = (0..3)
6750            .map(|slot| {
6751                AdapterHost::new(
6752                    Box::new(ScriptedAdapter::new(
6753                        slot,
6754                        AgentCapabilities::default(),
6755                        [
6756                            AgentEvent::Text {
6757                                slot,
6758                                text: STOP_TOKEN.into(),
6759                            },
6760                            AgentEvent::TurnComplete { slot },
6761                        ],
6762                    )),
6763                    None,
6764                )
6765            })
6766            .collect();
6767        let mut relay = RelayHost::new(hosts, 10).expect("relay");
6768        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6769        relay.start().await.expect("start");
6770        for expected in 0..3 {
6771            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6772                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6773        }
6774        assert_eq!(
6775            relay.run_turn("", 0).await.expect("complete"),
6776            RelayDecision::Complete
6777        );
6778    }
6779
6780    #[tokio::test]
6781    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6782        let first = ScriptedAdapter::new(
6783            0,
6784            AgentCapabilities::default(),
6785            [
6786                AgentEvent::Text {
6787                    slot: 0,
6788                    text: "done".into(),
6789                },
6790                AgentEvent::TurnComplete { slot: 0 },
6791            ],
6792        );
6793        let reviewer = ScriptedAdapter::new(
6794            1,
6795            AgentCapabilities::default(),
6796            [
6797                AgentEvent::Text {
6798                    slot: 1,
6799                    text: STOP_TOKEN.into(),
6800                },
6801                AgentEvent::TurnComplete { slot: 1 },
6802            ],
6803        );
6804        let mut relay = RelayHost::new(
6805            vec![
6806                AdapterHost::new(Box::new(first), None),
6807                AdapterHost::new(Box::new(reviewer), None),
6808            ],
6809            10,
6810        )
6811        .expect("relay");
6812        relay.start().await.expect("start");
6813        let first_decision = relay.run_turn("task", 0).await.expect("first");
6814        assert!(matches!(
6815            first_decision,
6816            RelayDecision::Dispatch { slot: 0, .. }
6817        ));
6818        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6819        assert!(matches!(
6820            reviewer_decision,
6821            RelayDecision::Dispatch {
6822                slot: 1,
6823                can_stop: true,
6824                ..
6825            }
6826        ));
6827        assert_eq!(
6828            relay.run_turn("", 0).await.expect("complete"),
6829            RelayDecision::Complete
6830        );
6831    }
6832
6833    #[tokio::test]
6834    async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6835        let tool = AgentEvent::Tool {
6836            slot: 0,
6837            update: crate::ToolUpdate {
6838                id: "read".into(),
6839                title: "Read file".into(),
6840                status: ToolStatus::Running,
6841                detail: None,
6842            },
6843        };
6844        let updates = vec![
6845            AgentEvent::Thought {
6846                slot: 0,
6847                text: "Check the buffer. ✈".into(),
6848            },
6849            AgentEvent::Text {
6850                slot: 0,
6851                text: "Let me check.".into(),
6852            },
6853            tool.clone(),
6854            AgentEvent::Text {
6855                slot: 0,
6856                text: "[CODE".into(),
6857            },
6858            AgentEvent::Text {
6859                slot: 0,
6860                text: " is ordinary.".into(),
6861            },
6862            AgentEvent::Text {
6863                slot: 0,
6864                text: "[CODESWARM:".into(),
6865            },
6866            AgentEvent::Text {
6867                slot: 0,
6868                text: "STOP] Done.".into(),
6869            },
6870            tool,
6871            AgentEvent::TurnComplete { slot: 0 },
6872        ];
6873        let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6874        let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6875        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6876        let captured = std::sync::Arc::clone(&events);
6877        let mut relay = RelayHost::new(
6878            vec![
6879                AdapterHost::new(Box::new(first), None),
6880                AdapterHost::new(Box::new(reviewer), None),
6881            ],
6882            2,
6883        )
6884        .expect("relay");
6885        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6886        relay.start().await.unwrap();
6887        relay.run_turn("task", 0).await.unwrap();
6888        let captured = events.lock().unwrap();
6889        let visible: Vec<_> = captured
6890            .iter()
6891            .filter(|event| {
6892                matches!(
6893                    event,
6894                    AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
6895                )
6896            })
6897            .cloned()
6898            .collect();
6899        assert_eq!(
6900            visible,
6901            vec![
6902                updates[0].clone(),
6903                updates[1].clone(),
6904                updates[2].clone(),
6905                AgentEvent::Text {
6906                    slot: 0,
6907                    text: "[CODE is ordinary.".into()
6908                },
6909                AgentEvent::Text {
6910                    slot: 0,
6911                    text: " Done.".into()
6912                },
6913                updates[7].clone(),
6914            ]
6915        );
6916    }
6917
6918    #[tokio::test]
6919    async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
6920        let text = |value: &str| AgentEvent::Text {
6921            slot: 0,
6922            text: value.into(),
6923        };
6924        let thought = || AgentEvent::Thought {
6925            slot: 0,
6926            text: "still checking".into(),
6927        };
6928        let tool = || AgentEvent::Tool {
6929            slot: 0,
6930            update: crate::ToolUpdate {
6931                id: "read".into(),
6932                title: "Read file".into(),
6933                status: ToolStatus::Running,
6934                detail: None,
6935            },
6936        };
6937        let cases = vec![
6938            (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
6939            (
6940                vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
6941                2,
6942            ),
6943            (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
6944            (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
6945            (
6946                vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
6947                1,
6948            ),
6949            (
6950                vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
6951                1,
6952            ),
6953            (
6954                vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
6955                2,
6956            ),
6957            (vec![text("result [CODESWARM:NEXT:1]")], 1),
6958            (vec![text("result [CODESWARM:NEXT:0]")], 1),
6959            (vec![text("result [CODESWARM:NEXT:99]")], 1),
6960            (
6961                vec![
6962                    text("result [CODESWARM:NEXT:3]"),
6963                    AgentEvent::UsageUpdated {
6964                        slot: 0,
6965                        usage: crate::UsageUpdate { used: 1, size: 100 },
6966                    },
6967                ],
6968                2,
6969            ),
6970        ];
6971        for (mut updates, expected) in cases {
6972            updates.push(AgentEvent::TurnComplete { slot: 0 });
6973            let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
6974            let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
6975                .chain((1..3).map(|slot| {
6976                    AdapterHost::new(
6977                        Box::new(ScriptedAdapter::new(
6978                            slot,
6979                            AgentCapabilities::default(),
6980                            [AgentEvent::TurnComplete { slot }],
6981                        )),
6982                        None,
6983                    )
6984                }))
6985                .collect();
6986            let mut relay = RelayHost::new(hosts, 10).unwrap();
6987            relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
6988            let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6989            let captured = events.clone();
6990            relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6991            relay.start().await.unwrap();
6992            relay.run_turn("task", 0).await.unwrap();
6993            let prompt = &relay.dispatches()[0].1;
6994            assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
6995            assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
6996            assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
6997            // The new names must appear even when introduction has already run.
6998            relay.introduced.fill(true);
6999            relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7000            let next = relay.run_turn("", 0).await.unwrap();
7001            assert!(
7002                matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7003                "{next:?}"
7004            );
7005            let prompt = &relay.dispatches()[1].1;
7006            assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7007            assert!(prompt.contains("result"));
7008            let public = prompt
7009                .split("Public updates:\n")
7010                .nth(1)
7011                .unwrap()
7012                .split("\n\nDo not use")
7013                .next()
7014                .unwrap();
7015            assert!(!public.contains("[CODESWARM:NEXT:"));
7016            let visible = events
7017                .lock()
7018                .unwrap()
7019                .iter()
7020                .filter_map(|event| match event {
7021                    AgentEvent::Text { text, .. } => Some(text.clone()),
7022                    _ => None,
7023                })
7024                .collect::<String>();
7025            assert!(visible.contains("result"));
7026            assert!(!visible.contains("[CODESWARM:"), "{visible}");
7027        }
7028    }
7029
7030    #[tokio::test]
7031    async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7032        let text = |value: &str| AgentEvent::Text {
7033            slot: 1,
7034            text: value.into(),
7035        };
7036        let thought = || AgentEvent::Thought {
7037            slot: 1,
7038            text: "still checking".into(),
7039        };
7040        let tool = || AgentEvent::Tool {
7041            slot: 1,
7042            update: crate::ToolUpdate {
7043                id: "read".into(),
7044                title: "Read file".into(),
7045                status: ToolStatus::Running,
7046                detail: None,
7047            },
7048        };
7049        let cases = vec![
7050            (vec![text(&format!("done {STOP_TOKEN}"))], true),
7051            (vec![text(STOP_TOKEN), text("\n  ")], true),
7052            (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7053            (vec![text(STOP_TOKEN), thought()], false),
7054            (vec![text(STOP_TOKEN), tool()], false),
7055            (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7056            (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7057            (vec![text("[CODESWARM:"), text("STOP]")], true),
7058            (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7059            (
7060                vec![AgentEvent::Thought {
7061                    slot: 1,
7062                    text: STOP_TOKEN.into(),
7063                }],
7064                false,
7065            ),
7066            (
7067                vec![
7068                    text(STOP_TOKEN),
7069                    AgentEvent::UsageUpdated {
7070                        slot: 1,
7071                        usage: crate::UsageUpdate { used: 1, size: 100 },
7072                    },
7073                ],
7074                true,
7075            ),
7076        ];
7077        for (mut events, stop) in cases {
7078            let first = ScriptedAdapter::new(
7079                0,
7080                AgentCapabilities::default(),
7081                [
7082                    AgentEvent::Text {
7083                        slot: 0,
7084                        text: "initial response".into(),
7085                    },
7086                    AgentEvent::TurnComplete { slot: 0 },
7087                    AgentEvent::TurnComplete { slot: 0 },
7088                ],
7089            );
7090            events.push(AgentEvent::TurnComplete { slot: 1 });
7091            let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7092            let mut relay = RelayHost::new(
7093                vec![
7094                    AdapterHost::new(Box::new(first), None),
7095                    AdapterHost::new(Box::new(reviewer), None),
7096                ],
7097                4,
7098            )
7099            .unwrap();
7100            relay.start().await.unwrap();
7101            relay.run_turn("task", 0).await.unwrap();
7102            relay.run_turn("", 0).await.unwrap();
7103            let next = relay.run_turn("", 0).await.unwrap();
7104            assert_eq!(
7105                matches!(next, RelayDecision::Complete),
7106                stop,
7107                "events={events:?}"
7108            );
7109            relay.stop().await.unwrap();
7110        }
7111    }
7112
7113    #[tokio::test]
7114    async fn stop_token_is_filtered_from_streamed_ui_events() {
7115        let first = ScriptedAdapter::new(
7116            0,
7117            AgentCapabilities::default(),
7118            [
7119                AgentEvent::Text {
7120                    slot: 0,
7121                    text: format!("visible {STOP_TOKEN} trailing"),
7122                },
7123                AgentEvent::TurnComplete { slot: 0 },
7124            ],
7125        );
7126        let reviewer = ScriptedAdapter::new(
7127            1,
7128            AgentCapabilities::default(),
7129            [AgentEvent::TurnComplete { slot: 1 }],
7130        );
7131        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7132        let captured = std::sync::Arc::clone(&events);
7133        let mut relay = RelayHost::new(
7134            vec![
7135                AdapterHost::new(Box::new(first), None),
7136                AdapterHost::new(Box::new(reviewer), None),
7137            ],
7138            2,
7139        )
7140        .expect("relay");
7141        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7142        relay.start().await.expect("start");
7143        relay.run_turn("task", 0).await.expect("turn");
7144        let captured = events.lock().expect("lock");
7145        assert!(captured.iter().all(|event| match event {
7146            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7147            _ => true,
7148        }));
7149        let visible = captured
7150            .iter()
7151            .filter_map(|event| match event {
7152                AgentEvent::Text { text, .. } => Some(text.as_str()),
7153                _ => None,
7154            })
7155            .collect::<String>();
7156        assert_eq!(visible, "visible  trailing");
7157    }
7158
7159    #[tokio::test]
7160    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7161        let first = ScriptedAdapter::new(
7162            0,
7163            AgentCapabilities::default(),
7164            [
7165                AgentEvent::Text {
7166                    slot: 0,
7167                    text: "done".into(),
7168                },
7169                AgentEvent::TurnComplete { slot: 0 },
7170            ],
7171        );
7172        let reviewer = ScriptedAdapter::new(
7173            1,
7174            AgentCapabilities::default(),
7175            [
7176                AgentEvent::Text {
7177                    slot: 1,
7178                    text: STOP_TOKEN.into(),
7179                },
7180                AgentEvent::TurnComplete { slot: 1 },
7181            ],
7182        );
7183        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7184        let captured = std::sync::Arc::clone(&events);
7185        let mut relay = RelayHost::new(
7186            vec![
7187                AdapterHost::new(Box::new(first), None),
7188                AdapterHost::new(Box::new(reviewer), None),
7189            ],
7190            4,
7191        )
7192        .expect("relay");
7193        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7194        relay.start().await.expect("start");
7195        relay.run_turn("task", 0).await.expect("first turn");
7196        relay.run_turn("", 0).await.expect("review turn");
7197        let captured = events.lock().expect("lock");
7198        assert!(captured.iter().any(|event| {
7199            matches!(
7200                event,
7201                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7202            )
7203        }));
7204        assert!(captured.iter().all(|event| match event {
7205            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7206            _ => true,
7207        }));
7208        let acknowledgment = captured
7209            .iter()
7210            .position(|event| {
7211                matches!(
7212                    event,
7213                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7214                )
7215            })
7216            .expect("visible acknowledgment");
7217        let completion = captured
7218            .iter()
7219            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7220            .expect("reviewer completion");
7221        assert!(acknowledgment < completion);
7222    }
7223
7224    #[tokio::test]
7225    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7226        let first = ScriptedAdapter::new(
7227            0,
7228            AgentCapabilities::default(),
7229            [
7230                AgentEvent::Text {
7231                    slot: 0,
7232                    text: "done".into(),
7233                },
7234                AgentEvent::TurnComplete { slot: 0 },
7235            ],
7236        );
7237        let reviewer = ScriptedAdapter::new(
7238            1,
7239            AgentCapabilities::default(),
7240            [
7241                AgentEvent::Text {
7242                    slot: 1,
7243                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7244                },
7245                AgentEvent::TurnComplete { slot: 1 },
7246            ],
7247        );
7248        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7249        let captured = std::sync::Arc::clone(&events);
7250        let mut relay = RelayHost::new(
7251            vec![
7252                AdapterHost::new(Box::new(first), None),
7253                AdapterHost::new(Box::new(reviewer), None),
7254            ],
7255            4,
7256        )
7257        .expect("relay");
7258        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7259        relay.start().await.expect("start");
7260        relay.run_turn("task", 0).await.expect("first turn");
7261        relay.run_turn("", 0).await.expect("review turn");
7262
7263        let visible = events
7264            .lock()
7265            .expect("lock")
7266            .iter()
7267            .filter_map(|event| match event {
7268                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7269                _ => None,
7270            })
7271            .collect::<String>();
7272        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7273        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7274    }
7275
7276    #[tokio::test]
7277    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7278        let first = AdapterHost::new(
7279            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7280            None,
7281        );
7282        let second = AdapterHost::new(
7283            Box::new(ScriptedAdapter::new(
7284                1,
7285                AgentCapabilities::default(),
7286                [AgentEvent::TurnComplete { slot: 1 }],
7287            )),
7288            None,
7289        );
7290        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7291        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7292        relay.set_event_sink(move |event| {
7293            if matches!(event, AgentEvent::Permission { .. }) {
7294                let _ = seen_sender.send(());
7295            }
7296        });
7297        relay.start().await.expect("start");
7298        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7299        let answer = async move {
7300            seen_receiver.recv().await.expect("permission request");
7301            sender
7302                .send(super::RelayPermissionAnswer {
7303                    slot: 0,
7304                    request_id: "permission-1".into(),
7305                    answer: PermissionAnswer::Selected {
7306                        option_id: "allow".into(),
7307                    },
7308                })
7309                .expect("queue permission answer");
7310        };
7311        tokio::time::timeout(std::time::Duration::from_millis(100), async {
7312            let ((), result) = tokio::join!(
7313                answer,
7314                relay.run_turn_with_permissions("task", 0, &mut receiver)
7315            );
7316            result
7317        })
7318        .await
7319        .expect("permission-gated turn should not deadlock")
7320        .expect("turn completes");
7321    }
7322
7323    #[tokio::test]
7324    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7325        let first = AdapterHost::new(
7326            Box::new(PendingAdapter {
7327                slot: 0,
7328                hang_on_cancel: false,
7329            }),
7330            None,
7331        );
7332        let second = AdapterHost::new(
7333            Box::new(ScriptedAdapter::new(
7334                1,
7335                AgentCapabilities::default(),
7336                [AgentEvent::TurnComplete { slot: 1 }],
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 error = {
7344            let turn = relay.run_turn("task", 0);
7345            tokio::pin!(turn);
7346            cancellation.request();
7347            turn.await.expect_err("cancellation should stop turn")
7348        };
7349        assert!(error.to_string().contains("relay turn cancelled"));
7350
7351        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7352        relay
7353            .run_turn("", 1)
7354            .await
7355            .expect("replacement job reaches the selected peer");
7356        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7357        assert!(replacement.contains("replacement job"));
7358        assert!(replacement.contains("User "));
7359        assert!(replacement.contains(":\ntask"));
7360        let owner_updates = relay.relay_mut().unseen_context(0);
7361        assert!(owner_updates.contains("User "));
7362        assert!(owner_updates.contains(":\ntask"));
7363        assert!(owner_updates.contains(":\nreplacement job"));
7364    }
7365
7366    #[tokio::test]
7367    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7368        let first = AdapterHost::new(
7369            Box::new(PendingAdapter {
7370                slot: 0,
7371                hang_on_cancel: true,
7372            }),
7373            None,
7374        );
7375        let second = AdapterHost::new(
7376            Box::new(PendingAdapter {
7377                slot: 1,
7378                hang_on_cancel: false,
7379            }),
7380            None,
7381        );
7382        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7383        relay.start().await.expect("start");
7384        let cancellation = relay.cancellation();
7385        let turn = relay.run_turn("task", 0);
7386        tokio::pin!(turn);
7387        cancellation.request();
7388        let error = turn.await.expect_err("cancellation should stop turn");
7389        assert!(error.to_string().contains("timed out"));
7390    }
7391
7392    #[tokio::test]
7393    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7394        let event = [AgentEvent::TurnComplete { slot: 0 }];
7395        let first = AdapterHost::new(
7396            Box::new(ScriptedAdapter::new(
7397                0,
7398                AgentCapabilities::default(),
7399                event.clone(),
7400            )),
7401            None,
7402        );
7403        let second = AdapterHost::new(
7404            Box::new(ScriptedAdapter::new(
7405                1,
7406                AgentCapabilities::default(),
7407                [AgentEvent::TurnComplete { slot: 1 }],
7408            )),
7409            None,
7410        );
7411        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7412        relay.start().await.expect("start");
7413
7414        relay.pause();
7415        assert_eq!(
7416            relay.run_turn("paused", 0).await.expect("paused turn"),
7417            crate::relay::RelayDecision::Paused
7418        );
7419        assert!(relay.dispatches().is_empty());
7420
7421        relay.resume();
7422        relay.relay_mut().drop_agent(1).expect("drop reviewer");
7423        assert!(matches!(
7424            relay
7425                .run_turn("solo follow-up", 0)
7426                .await
7427                .expect("solo turn"),
7428            crate::relay::RelayDecision::Dispatch {
7429                slot: 0,
7430                can_stop: false,
7431                ..
7432            }
7433        ));
7434        assert_eq!(relay.dispatches().len(), 1);
7435    }
7436
7437    #[tokio::test]
7438    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7439        let first = AdapterHost::new(
7440            Box::new(ScriptedAdapter::new(
7441                0,
7442                AgentCapabilities::default(),
7443                [AgentEvent::TurnComplete { slot: 0 }],
7444            )),
7445            None,
7446        );
7447        let second = AdapterHost::new(
7448            Box::new(ScriptedAdapter::new(
7449                1,
7450                AgentCapabilities::default(),
7451                [AgentEvent::TurnComplete { slot: 1 }],
7452            )),
7453            None,
7454        );
7455        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7456        relay.set_roster_names(vec!["First".into(), "Second".into()]);
7457        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7458        relay.set_roster_launch_specs(vec![
7459            ("custom".into(), "owner".into()),
7460            ("custom".into(), "peer".into()),
7461        ]);
7462        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7463        let captured = std::sync::Arc::clone(&events);
7464        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7465        relay.start().await.expect("start");
7466        let slot = relay
7467            .add_agent(
7468                AdapterHost::new(
7469                    Box::new(ScriptedAdapter::new(
7470                        2,
7471                        AgentCapabilities::default(),
7472                        [AgentEvent::TurnComplete { slot: 2 }],
7473                    )),
7474                    None,
7475                ),
7476                "Reviewer",
7477                "reviewer.example",
7478                "reviewer --acp",
7479            )
7480            .await
7481            .expect("append agent");
7482        assert_eq!(slot, 2);
7483        assert_eq!(
7484            relay.relay().active_slots().collect::<Vec<_>>(),
7485            vec![0, 1, 2]
7486        );
7487        assert_eq!(
7488            relay
7489                .session_metadata()
7490                .get("agents")
7491                .and_then(|value| value.as_array())
7492                .map(Vec::len),
7493            Some(3)
7494        );
7495        relay.drop_agent(1).await.expect("drop middle peer");
7496        let metadata = relay.session_metadata();
7497        assert_eq!(
7498            metadata.get("agents"),
7499            Some(&serde_json::json!([
7500                {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7501                {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7502            ]))
7503        );
7504        assert!(
7505            events
7506                .lock()
7507                .expect("lock")
7508                .iter()
7509                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7510        );
7511    }
7512
7513    #[tokio::test]
7514    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7515        let path = unique_test_path("codeswarm-session-metadata", "json");
7516        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7517        let writer = metadata_store.buffered().expect("metadata writer");
7518        let first = AdapterHost::new(
7519            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7520            None,
7521        );
7522        let second = AdapterHost::new(
7523            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7524            None,
7525        );
7526        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7527        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7528        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7529        relay.set_roster_launch_specs(vec![
7530            ("custom".into(), "claude".into()),
7531            ("custom".into(), "codex".into()),
7532        ]);
7533        relay.set_session_metadata_writer(writer);
7534        relay.start().await.expect("start");
7535        relay.drop_agent(0).await.expect("drop first agent");
7536        relay.stop().await.expect("stop");
7537
7538        let loaded = metadata_store
7539            .read()
7540            .expect("read metadata")
7541            .expect("metadata snapshot");
7542        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7543        assert_eq!(
7544            loaded.get("agents"),
7545            Some(&serde_json::json!([{
7546                "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7547                "command": "codex", "supports_load_session": false
7548            }]))
7549        );
7550        assert!(loaded.get("owner").is_none());
7551        let _ = std::fs::remove_file(path);
7552    }
7553
7554    #[tokio::test]
7555    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7556        let first = AdapterHost::new(
7557            Box::new(ScriptedAdapter::new(
7558                0,
7559                AgentCapabilities::default(),
7560                [
7561                    AgentEvent::Text {
7562                        slot: 0,
7563                        text: "owner stream".into(),
7564                    },
7565                    AgentEvent::TurnComplete { slot: 0 },
7566                ],
7567            )),
7568            None,
7569        );
7570        let second = AdapterHost::new(
7571            Box::new(ScriptedAdapter::new(
7572                1,
7573                AgentCapabilities::default(),
7574                [
7575                    AgentEvent::Text {
7576                        slot: 1,
7577                        text: "peer stream".into(),
7578                    },
7579                    AgentEvent::TurnComplete { slot: 1 },
7580                ],
7581            )),
7582            None,
7583        );
7584        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7585        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7586        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7587        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7588        let captured = std::sync::Arc::clone(&events);
7589        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7590        relay.start().await.expect("start");
7591
7592        relay.swap_agents(0, 1).expect("swap peers");
7593        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7594        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7595        relay.run_turn("task", 0).await.expect("swapped turn");
7596        let events = events.lock().expect("events");
7597        assert!(events.iter().any(|event| {
7598            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7599        }));
7600        assert!(relay.dispatches()[0].1.contains("You are Peer"));
7601    }
7602
7603    #[tokio::test]
7604    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7605        let path = unique_test_path("codeswarm-session-metadata", "json");
7606        let first = AdapterHost::new(
7607            Box::new(ScriptedAdapter::new(
7608                0,
7609                AgentCapabilities::default(),
7610                [AgentEvent::TurnComplete { slot: 0 }],
7611            )),
7612            None,
7613        );
7614        let second = AdapterHost::new(
7615            Box::new(ScriptedAdapter::new(
7616                1,
7617                AgentCapabilities::default(),
7618                [AgentEvent::TurnComplete { slot: 1 }],
7619            )),
7620            None,
7621        );
7622        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7623        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7624        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7625        relay.set_roster_launch_specs(vec![
7626            ("custom".into(), "claude".into()),
7627            ("custom".into(), "codex".into()),
7628        ]);
7629        let writer = SessionMetadataStore::open(&path)
7630            .buffered()
7631            .expect("metadata writer");
7632        relay.set_session_metadata_writer(writer);
7633        relay.start().await.expect("start");
7634        relay.stop().await.expect("stop");
7635        let loaded = SessionMetadataStore::open(&path)
7636            .read()
7637            .expect("read metadata")
7638            .expect("metadata snapshot");
7639        let agents = loaded
7640            .get("agents")
7641            .and_then(|value| value.as_array())
7642            .expect("agents");
7643        assert_eq!(agents.len(), 2);
7644        assert_eq!(agents[0]["identity"], "claude.com");
7645        assert_eq!(agents[1]["identity"], "openai.com");
7646        let _ = std::fs::remove_file(path);
7647    }
7648
7649    #[tokio::test]
7650    async fn relay_host_routes_unseen_public_context_to_next_agent() {
7651        let first = AdapterHost::new(
7652            Box::new(ScriptedAdapter::new(
7653                0,
7654                AgentCapabilities::default(),
7655                [
7656                    AgentEvent::Text {
7657                        slot: 0,
7658                        text: "implemented the fix".into(),
7659                    },
7660                    AgentEvent::TurnComplete { slot: 0 },
7661                ],
7662            )),
7663            None,
7664        );
7665        let second = AdapterHost::new(
7666            Box::new(ScriptedAdapter::new(
7667                1,
7668                AgentCapabilities::default(),
7669                [AgentEvent::TurnComplete { slot: 1 }],
7670            )),
7671            None,
7672        );
7673        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7674        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7675        relay.start().await.expect("start");
7676        relay.run_turn("task", 0).await.expect("first turn");
7677        relay.run_turn("review this", 0).await.expect("review turn");
7678
7679        assert_eq!(relay.dispatches().len(), 2);
7680        assert_eq!(relay.dispatches()[0].0, 0);
7681        assert!(relay.dispatches()[0].1.contains("task"));
7682        assert!(relay.dispatches()[0].1.contains("You are Codex"));
7683        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7684        assert_eq!(relay.dispatches()[1].0, 1);
7685        assert!(relay.dispatches()[1].1.contains("review this"));
7686        let public = relay.dispatches()[1]
7687            .1
7688            .split_once("Public updates:\n")
7689            .map(|(_, updates)| updates)
7690            .expect("review receives public context");
7691        let header = public
7692            .lines()
7693            .find(|line| line.starts_with("Codex "))
7694            .expect("named previous agent");
7695        let timestamp = header
7696            .strip_prefix("Codex ")
7697            .and_then(|value| value.strip_suffix(':'))
7698            .expect("timestamped header");
7699        assert_eq!(timestamp.len(), 5);
7700        assert_eq!(timestamp.as_bytes()[2], b':');
7701        assert!(
7702            timestamp
7703                .bytes()
7704                .enumerate()
7705                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7706        );
7707        assert!(public.contains("implemented the fix"));
7708        assert!(!public.contains("Agent 0"));
7709    }
7710}