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. {STOP_TOKEN} is a global batch stop: it stops all other agents and ends the entire automated relay, not just your turn. Use it with extreme care.\nUse it only when the shared task is fully complete, no meaningful correction is needed, and no other agent should continue working. If there is any uncertainty, do not use it; state what remains and let the relay continue.\nWhen—and only when—those conditions are met, end 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!(
6335            relay.dispatches()[1]
6336                .1
6337                .contains("stops all other agents and ends the entire automated relay")
6338        );
6339        assert!(relay.dispatches()[1].1.contains("Use it with extreme care"));
6340        assert!(
6341            relay.dispatches()[1]
6342                .1
6343                .contains("If there is any uncertainty, do not use it")
6344        );
6345        assert!(relay.dispatches()[0].1.contains("Do not use"));
6346        let lifecycle = events.lock().expect("events");
6347        let positions = lifecycle
6348            .iter()
6349            .filter_map(|event| match event {
6350                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6351                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6352                _ => None,
6353            })
6354            .collect::<Vec<_>>();
6355        assert_eq!(
6356            positions,
6357            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6358        );
6359    }
6360
6361    #[tokio::test]
6362    async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6363        let stops = Arc::new(AtomicUsize::new(0));
6364        let failed = FailingStartAdapter {
6365            slot: 0,
6366            stopped: stops.clone(),
6367        };
6368        let healthy = ScriptedAdapter::new(
6369            1,
6370            AgentCapabilities::default(),
6371            [
6372                AgentEvent::Text {
6373                    slot: 1,
6374                    text: "healthy response".into(),
6375                },
6376                AgentEvent::TurnComplete { slot: 1 },
6377            ],
6378        );
6379        let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6380        let captured = events.clone();
6381        let mut relay = RelayHost::new(
6382            vec![
6383                AdapterHost::new(Box::new(failed), None),
6384                AdapterHost::new(Box::new(healthy), None),
6385            ],
6386            4,
6387        )
6388        .unwrap();
6389        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6390        relay.start_resuming().await.unwrap();
6391        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6392        assert!(relay.dispatches().is_empty());
6393        assert_eq!(stops.load(Ordering::Relaxed), 1);
6394        assert!(
6395            events
6396                .lock()
6397                .unwrap()
6398                .iter()
6399                .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6400        );
6401        assert!(
6402            !events
6403                .lock()
6404                .unwrap()
6405                .iter()
6406                .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6407        );
6408        assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6409        assert!(
6410            relay
6411                .relay_mut()
6412                .enqueue_human("explicit healthy target", Some(1))
6413        );
6414        assert!(matches!(
6415            relay.run_turn("", 1).await.unwrap(),
6416            RelayDecision::Dispatch { slot: 1, .. }
6417        ));
6418        relay.stop().await.unwrap();
6419    }
6420
6421    #[tokio::test]
6422    async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6423        let capabilities = AgentCapabilities::default();
6424        let first = ScriptedAdapter::new(
6425            0,
6426            capabilities.clone(),
6427            [
6428                AgentEvent::Text {
6429                    slot: 0,
6430                    text: "implemented".into(),
6431                },
6432                AgentEvent::TurnComplete { slot: 0 },
6433                AgentEvent::Text {
6434                    slot: 0,
6435                    text: format!("fixed review findings {STOP_TOKEN}"),
6436                },
6437                AgentEvent::TurnComplete { slot: 0 },
6438            ],
6439        );
6440        let second = ScriptedAdapter::new(
6441            1,
6442            capabilities,
6443            [
6444                AgentEvent::Text {
6445                    slot: 1,
6446                    text: "reviewed".into(),
6447                },
6448                AgentEvent::TurnComplete { slot: 1 },
6449                AgentEvent::Text {
6450                    slot: 1,
6451                    text: format!("approved {STOP_TOKEN}"),
6452                },
6453                AgentEvent::TurnComplete { slot: 1 },
6454                AgentEvent::Text {
6455                    slot: 1,
6456                    text: "new task".into(),
6457                },
6458                AgentEvent::TurnComplete { slot: 1 },
6459            ],
6460        );
6461        let hosts = vec![
6462            AdapterHost::new(Box::new(first), None),
6463            AdapterHost::new(Box::new(second), None),
6464        ];
6465        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6466        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6467        relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6468        relay.start().await.expect("start");
6469        assert!(matches!(
6470            relay.run_turn("task", 0).await.expect("implementer turn"),
6471            RelayDecision::Dispatch {
6472                slot: 0,
6473                can_stop: false,
6474                ..
6475            }
6476        ));
6477        let implementer_prompt = &relay.dispatches()[0].1;
6478        assert!(implementer_prompt.contains("you are the implementer"));
6479        assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6480        assert!(!implementer_prompt.contains("you are the reviewer"));
6481        assert!(implementer_prompt.contains("Do not use"));
6482        assert!(matches!(
6483            relay.run_turn("", 0).await.expect("reviewer turn"),
6484            RelayDecision::Dispatch {
6485                slot: 1,
6486                can_stop: true,
6487                ..
6488            }
6489        ));
6490        let reviewer_prompt = &relay.dispatches()[1].1;
6491        assert!(reviewer_prompt.contains("you are the reviewer"));
6492        assert!(reviewer_prompt.contains("Claude handed off"));
6493        assert!(reviewer_prompt.contains("concrete defects"));
6494        assert!(reviewer_prompt.contains("concise approval"));
6495        assert!(reviewer_prompt.contains(STOP_TOKEN));
6496        assert!(!reviewer_prompt.contains("you are the implementer"));
6497        relay.run_turn("", 0).await.unwrap();
6498        assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6499        assert!(relay.dispatches()[2].1.contains("Do not use"));
6500        assert!(matches!(
6501            relay.run_turn("", 0).await.unwrap(),
6502            RelayDecision::Dispatch { slot: 1, .. }
6503        ));
6504        assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6505        assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6506        relay.run_turn("", 1).await.unwrap();
6507        assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6508    }
6509
6510    #[tokio::test]
6511    async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6512        let solo = ScriptedAdapter::new(
6513            0,
6514            AgentCapabilities::default(),
6515            [
6516                AgentEvent::Text {
6517                    slot: 0,
6518                    text: "solo".into(),
6519                },
6520                AgentEvent::TurnComplete { slot: 0 },
6521            ],
6522        );
6523        let mut solo_relay =
6524            RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6525        solo_relay
6526            .relay_mut()
6527            .set_strategy(CollaborationStrategy::Pair);
6528        solo_relay.start().await.expect("start");
6529        solo_relay.run_turn("task", 0).await.expect("solo turn");
6530        assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6531
6532        let roster_first = ScriptedAdapter::new(
6533            0,
6534            AgentCapabilities::default(),
6535            [AgentEvent::TurnComplete { slot: 0 }],
6536        );
6537        let roster_second = ScriptedAdapter::new(
6538            1,
6539            AgentCapabilities::default(),
6540            [AgentEvent::TurnComplete { slot: 1 }],
6541        );
6542        let mut roster = RelayHost::new(
6543            vec![
6544                AdapterHost::new(Box::new(roster_first), None),
6545                AdapterHost::new(Box::new(roster_second), None),
6546            ],
6547            4,
6548        )
6549        .expect("relay");
6550        roster.start().await.expect("start");
6551        roster.run_turn("task", 0).await.expect("first turn");
6552        roster.run_turn("", 0).await.expect("second turn");
6553        assert!(!roster.dispatches()[0].1.contains("Pair role"));
6554        assert!(!roster.dispatches()[1].1.contains("Pair role"));
6555
6556        let pair_first = ScriptedAdapter::new(
6557            0,
6558            AgentCapabilities::default(),
6559            [AgentEvent::TurnComplete { slot: 0 }],
6560        );
6561        let pair_second = ScriptedAdapter::new(
6562            1,
6563            AgentCapabilities::default(),
6564            [AgentEvent::TurnComplete { slot: 1 }],
6565        );
6566        let mut pair = RelayHost::new(
6567            vec![
6568                AdapterHost::new(Box::new(pair_first), None),
6569                AdapterHost::new(Box::new(pair_second), None),
6570            ],
6571            4,
6572        )
6573        .expect("relay");
6574        pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6575        assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6576        pair.start().await.expect("start");
6577        assert!(matches!(
6578            pair.run_turn("ignored", 0).await.expect("direct turn"),
6579            RelayDecision::Dispatch {
6580                slot: 1,
6581                direct: true,
6582                ..
6583            }
6584        ));
6585        let direct_prompt = &pair.dispatches()[0].1;
6586        assert!(direct_prompt.contains("private"));
6587        assert!(!direct_prompt.contains("Pair role"));
6588    }
6589
6590    #[tokio::test]
6591    async fn relay_host_routes_around_a_usage_limited_agent() {
6592        let capabilities = AgentCapabilities::default();
6593        let first = ScriptedAdapter::new(
6594            0,
6595            capabilities.clone(),
6596            [
6597                AgentEvent::Text {
6598                    slot: 0,
6599                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6600                           credits or try again later."
6601                        .into(),
6602                },
6603                AgentEvent::TurnComplete { slot: 0 },
6604            ],
6605        );
6606        let second = ScriptedAdapter::new(
6607            1,
6608            capabilities,
6609            [
6610                AgentEvent::Text {
6611                    slot: 1,
6612                    text: "review done".into(),
6613                },
6614                AgentEvent::TurnComplete { slot: 1 },
6615            ],
6616        );
6617        let hosts = vec![
6618            AdapterHost::new(Box::new(first), None),
6619            AdapterHost::new(Box::new(second), None),
6620        ];
6621        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6622        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6623        let captured = std::sync::Arc::clone(&events);
6624        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6625        relay.start().await.expect("start");
6626        events.lock().expect("events").clear();
6627        assert!(matches!(
6628            relay.run_turn("task", 0).await.expect("limited turn"),
6629            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6630        ));
6631        assert!(
6632            events
6633                .lock()
6634                .expect("events")
6635                .iter()
6636                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6637        );
6638        // The next automatic turn skips the limited agent entirely.
6639        assert!(matches!(
6640            relay.run_turn("", 0).await.expect("next turn"),
6641            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6642        ));
6643        assert!(relay.relay().is_limited(0));
6644        // A reload restores the agent to the ring. (ScriptedAdapter cannot
6645        // feed further turns, so the restored routing itself is covered by
6646        // the relay unit tests.)
6647        relay.reload(0).await.expect("reload");
6648        assert!(!relay.relay().is_limited(0));
6649    }
6650
6651    #[tokio::test]
6652    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6653        let limited = ScriptedAdapter::new(
6654            0,
6655            AgentCapabilities::default(),
6656            [AgentEvent::Failed {
6657                slot: 0,
6658                started: true,
6659                detail: "request failed: insufficient_quota".into(),
6660            }],
6661        );
6662        let healthy = ScriptedAdapter::new(
6663            1,
6664            AgentCapabilities::default(),
6665            [AgentEvent::TurnComplete { slot: 1 }],
6666        );
6667        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6668        let captured = std::sync::Arc::clone(&events);
6669        let mut relay = RelayHost::new(
6670            vec![
6671                AdapterHost::new(Box::new(limited), None),
6672                AdapterHost::new(Box::new(healthy), None),
6673            ],
6674            4,
6675        )
6676        .expect("relay");
6677        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6678        relay.start().await.expect("start");
6679        events.lock().expect("events").clear();
6680
6681        assert!(matches!(
6682            relay.run_turn("task", 0).await.expect("limited failure"),
6683            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6684        ));
6685        assert!(relay.relay().is_limited(0));
6686        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6687        {
6688            let events = events.lock().expect("events");
6689            assert!(
6690                events
6691                    .iter()
6692                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6693            );
6694            assert!(
6695                !events
6696                    .iter()
6697                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
6698            );
6699        }
6700
6701        assert!(matches!(
6702            relay.run_turn("", 0).await.expect("healthy peer"),
6703            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6704        ));
6705    }
6706
6707    #[tokio::test]
6708    async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6709        let failed = ScriptedAdapter::new(
6710            0,
6711            AgentCapabilities::default(),
6712            [AgentEvent::Failed {
6713                slot: 0,
6714                started: true,
6715                detail: "connection lost".into(),
6716            }],
6717        );
6718        let healthy = ScriptedAdapter::new(
6719            1,
6720            AgentCapabilities::default(),
6721            [AgentEvent::TurnComplete { slot: 1 }],
6722        );
6723        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6724        let captured = std::sync::Arc::clone(&events);
6725        let mut relay = RelayHost::new(
6726            vec![
6727                AdapterHost::new(Box::new(failed), None),
6728                AdapterHost::new(Box::new(healthy), None),
6729            ],
6730            4,
6731        )
6732        .expect("relay");
6733        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6734        relay.start().await.expect("start");
6735
6736        assert!(matches!(
6737            relay.run_turn("task", 0).await.expect("handled failure"),
6738            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6739        ));
6740        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6741        assert!(relay.relay().is_limited(0));
6742        assert!(events.lock().expect("lock").iter().any(|event| {
6743            matches!(
6744                event,
6745                AgentEvent::Failed {
6746                    slot: 0,
6747                    started: true,
6748                    ..
6749                }
6750            )
6751        }));
6752        assert!(matches!(
6753            relay.run_turn("", 0).await.expect("healthy peer"),
6754            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6755        ));
6756    }
6757
6758    #[tokio::test]
6759    async fn codex_stop_does_not_skip_later_roster_reviewers() {
6760        let hosts = (0..3)
6761            .map(|slot| {
6762                AdapterHost::new(
6763                    Box::new(ScriptedAdapter::new(
6764                        slot,
6765                        AgentCapabilities::default(),
6766                        [
6767                            AgentEvent::Text {
6768                                slot,
6769                                text: STOP_TOKEN.into(),
6770                            },
6771                            AgentEvent::TurnComplete { slot },
6772                        ],
6773                    )),
6774                    None,
6775                )
6776            })
6777            .collect();
6778        let mut relay = RelayHost::new(hosts, 10).expect("relay");
6779        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6780        relay.start().await.expect("start");
6781        for expected in 0..3 {
6782            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6783                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6784        }
6785        assert_eq!(
6786            relay.run_turn("", 0).await.expect("complete"),
6787            RelayDecision::Complete
6788        );
6789    }
6790
6791    #[tokio::test]
6792    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6793        let first = ScriptedAdapter::new(
6794            0,
6795            AgentCapabilities::default(),
6796            [
6797                AgentEvent::Text {
6798                    slot: 0,
6799                    text: "done".into(),
6800                },
6801                AgentEvent::TurnComplete { slot: 0 },
6802            ],
6803        );
6804        let reviewer = ScriptedAdapter::new(
6805            1,
6806            AgentCapabilities::default(),
6807            [
6808                AgentEvent::Text {
6809                    slot: 1,
6810                    text: STOP_TOKEN.into(),
6811                },
6812                AgentEvent::TurnComplete { slot: 1 },
6813            ],
6814        );
6815        let mut relay = RelayHost::new(
6816            vec![
6817                AdapterHost::new(Box::new(first), None),
6818                AdapterHost::new(Box::new(reviewer), None),
6819            ],
6820            10,
6821        )
6822        .expect("relay");
6823        relay.start().await.expect("start");
6824        let first_decision = relay.run_turn("task", 0).await.expect("first");
6825        assert!(matches!(
6826            first_decision,
6827            RelayDecision::Dispatch { slot: 0, .. }
6828        ));
6829        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6830        assert!(matches!(
6831            reviewer_decision,
6832            RelayDecision::Dispatch {
6833                slot: 1,
6834                can_stop: true,
6835                ..
6836            }
6837        ));
6838        assert_eq!(
6839            relay.run_turn("", 0).await.expect("complete"),
6840            RelayDecision::Complete
6841        );
6842    }
6843
6844    #[tokio::test]
6845    async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6846        let tool = AgentEvent::Tool {
6847            slot: 0,
6848            update: crate::ToolUpdate {
6849                id: "read".into(),
6850                title: "Read file".into(),
6851                status: ToolStatus::Running,
6852                detail: None,
6853            },
6854        };
6855        let updates = vec![
6856            AgentEvent::Thought {
6857                slot: 0,
6858                text: "Check the buffer. ✈".into(),
6859            },
6860            AgentEvent::Text {
6861                slot: 0,
6862                text: "Let me check.".into(),
6863            },
6864            tool.clone(),
6865            AgentEvent::Text {
6866                slot: 0,
6867                text: "[CODE".into(),
6868            },
6869            AgentEvent::Text {
6870                slot: 0,
6871                text: " is ordinary.".into(),
6872            },
6873            AgentEvent::Text {
6874                slot: 0,
6875                text: "[CODESWARM:".into(),
6876            },
6877            AgentEvent::Text {
6878                slot: 0,
6879                text: "STOP] Done.".into(),
6880            },
6881            tool,
6882            AgentEvent::TurnComplete { slot: 0 },
6883        ];
6884        let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
6885        let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
6886        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6887        let captured = std::sync::Arc::clone(&events);
6888        let mut relay = RelayHost::new(
6889            vec![
6890                AdapterHost::new(Box::new(first), None),
6891                AdapterHost::new(Box::new(reviewer), None),
6892            ],
6893            2,
6894        )
6895        .expect("relay");
6896        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6897        relay.start().await.unwrap();
6898        relay.run_turn("task", 0).await.unwrap();
6899        let captured = events.lock().unwrap();
6900        let visible: Vec<_> = captured
6901            .iter()
6902            .filter(|event| {
6903                matches!(
6904                    event,
6905                    AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
6906                )
6907            })
6908            .cloned()
6909            .collect();
6910        assert_eq!(
6911            visible,
6912            vec![
6913                updates[0].clone(),
6914                updates[1].clone(),
6915                updates[2].clone(),
6916                AgentEvent::Text {
6917                    slot: 0,
6918                    text: "[CODE is ordinary.".into()
6919                },
6920                AgentEvent::Text {
6921                    slot: 0,
6922                    text: " Done.".into()
6923                },
6924                updates[7].clone(),
6925            ]
6926        );
6927    }
6928
6929    #[tokio::test]
6930    async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
6931        let text = |value: &str| AgentEvent::Text {
6932            slot: 0,
6933            text: value.into(),
6934        };
6935        let thought = || AgentEvent::Thought {
6936            slot: 0,
6937            text: "still checking".into(),
6938        };
6939        let tool = || AgentEvent::Tool {
6940            slot: 0,
6941            update: crate::ToolUpdate {
6942                id: "read".into(),
6943                title: "Read file".into(),
6944                status: ToolStatus::Running,
6945                detail: None,
6946            },
6947        };
6948        let cases = vec![
6949            (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
6950            (
6951                vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
6952                2,
6953            ),
6954            (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
6955            (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
6956            (
6957                vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
6958                1,
6959            ),
6960            (
6961                vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
6962                1,
6963            ),
6964            (
6965                vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
6966                2,
6967            ),
6968            (vec![text("result [CODESWARM:NEXT:1]")], 1),
6969            (vec![text("result [CODESWARM:NEXT:0]")], 1),
6970            (vec![text("result [CODESWARM:NEXT:99]")], 1),
6971            (
6972                vec![
6973                    text("result [CODESWARM:NEXT:3]"),
6974                    AgentEvent::UsageUpdated {
6975                        slot: 0,
6976                        usage: crate::UsageUpdate { used: 1, size: 100 },
6977                    },
6978                ],
6979                2,
6980            ),
6981        ];
6982        for (mut updates, expected) in cases {
6983            updates.push(AgentEvent::TurnComplete { slot: 0 });
6984            let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
6985            let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
6986                .chain((1..3).map(|slot| {
6987                    AdapterHost::new(
6988                        Box::new(ScriptedAdapter::new(
6989                            slot,
6990                            AgentCapabilities::default(),
6991                            [AgentEvent::TurnComplete { slot }],
6992                        )),
6993                        None,
6994                    )
6995                }))
6996                .collect();
6997            let mut relay = RelayHost::new(hosts, 10).unwrap();
6998            relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
6999            let events = Arc::new(std::sync::Mutex::new(Vec::new()));
7000            let captured = events.clone();
7001            relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7002            relay.start().await.unwrap();
7003            relay.run_turn("task", 0).await.unwrap();
7004            let prompt = &relay.dispatches()[0].1;
7005            assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
7006            assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
7007            assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
7008            // The new names must appear even when introduction has already run.
7009            relay.introduced.fill(true);
7010            relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7011            let next = relay.run_turn("", 0).await.unwrap();
7012            assert!(
7013                matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7014                "{next:?}"
7015            );
7016            let prompt = &relay.dispatches()[1].1;
7017            assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7018            assert!(prompt.contains("result"));
7019            let public = prompt
7020                .split("Public updates:\n")
7021                .nth(1)
7022                .unwrap()
7023                .split("\n\nDo not use")
7024                .next()
7025                .unwrap();
7026            assert!(!public.contains("[CODESWARM:NEXT:"));
7027            let visible = events
7028                .lock()
7029                .unwrap()
7030                .iter()
7031                .filter_map(|event| match event {
7032                    AgentEvent::Text { text, .. } => Some(text.clone()),
7033                    _ => None,
7034                })
7035                .collect::<String>();
7036            assert!(visible.contains("result"));
7037            assert!(!visible.contains("[CODESWARM:"), "{visible}");
7038        }
7039    }
7040
7041    #[tokio::test]
7042    async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7043        let text = |value: &str| AgentEvent::Text {
7044            slot: 1,
7045            text: value.into(),
7046        };
7047        let thought = || AgentEvent::Thought {
7048            slot: 1,
7049            text: "still checking".into(),
7050        };
7051        let tool = || AgentEvent::Tool {
7052            slot: 1,
7053            update: crate::ToolUpdate {
7054                id: "read".into(),
7055                title: "Read file".into(),
7056                status: ToolStatus::Running,
7057                detail: None,
7058            },
7059        };
7060        let cases = vec![
7061            (vec![text(&format!("done {STOP_TOKEN}"))], true),
7062            (vec![text(STOP_TOKEN), text("\n  ")], true),
7063            (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7064            (vec![text(STOP_TOKEN), thought()], false),
7065            (vec![text(STOP_TOKEN), tool()], false),
7066            (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7067            (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7068            (vec![text("[CODESWARM:"), text("STOP]")], true),
7069            (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7070            (
7071                vec![AgentEvent::Thought {
7072                    slot: 1,
7073                    text: STOP_TOKEN.into(),
7074                }],
7075                false,
7076            ),
7077            (
7078                vec![
7079                    text(STOP_TOKEN),
7080                    AgentEvent::UsageUpdated {
7081                        slot: 1,
7082                        usage: crate::UsageUpdate { used: 1, size: 100 },
7083                    },
7084                ],
7085                true,
7086            ),
7087        ];
7088        for (mut events, stop) in cases {
7089            let first = ScriptedAdapter::new(
7090                0,
7091                AgentCapabilities::default(),
7092                [
7093                    AgentEvent::Text {
7094                        slot: 0,
7095                        text: "initial response".into(),
7096                    },
7097                    AgentEvent::TurnComplete { slot: 0 },
7098                    AgentEvent::TurnComplete { slot: 0 },
7099                ],
7100            );
7101            events.push(AgentEvent::TurnComplete { slot: 1 });
7102            let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7103            let mut relay = RelayHost::new(
7104                vec![
7105                    AdapterHost::new(Box::new(first), None),
7106                    AdapterHost::new(Box::new(reviewer), None),
7107                ],
7108                4,
7109            )
7110            .unwrap();
7111            relay.start().await.unwrap();
7112            relay.run_turn("task", 0).await.unwrap();
7113            relay.run_turn("", 0).await.unwrap();
7114            let next = relay.run_turn("", 0).await.unwrap();
7115            assert_eq!(
7116                matches!(next, RelayDecision::Complete),
7117                stop,
7118                "events={events:?}"
7119            );
7120            relay.stop().await.unwrap();
7121        }
7122    }
7123
7124    #[tokio::test]
7125    async fn stop_token_is_filtered_from_streamed_ui_events() {
7126        let first = ScriptedAdapter::new(
7127            0,
7128            AgentCapabilities::default(),
7129            [
7130                AgentEvent::Text {
7131                    slot: 0,
7132                    text: format!("visible {STOP_TOKEN} trailing"),
7133                },
7134                AgentEvent::TurnComplete { slot: 0 },
7135            ],
7136        );
7137        let reviewer = ScriptedAdapter::new(
7138            1,
7139            AgentCapabilities::default(),
7140            [AgentEvent::TurnComplete { slot: 1 }],
7141        );
7142        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7143        let captured = std::sync::Arc::clone(&events);
7144        let mut relay = RelayHost::new(
7145            vec![
7146                AdapterHost::new(Box::new(first), None),
7147                AdapterHost::new(Box::new(reviewer), None),
7148            ],
7149            2,
7150        )
7151        .expect("relay");
7152        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7153        relay.start().await.expect("start");
7154        relay.run_turn("task", 0).await.expect("turn");
7155        let captured = events.lock().expect("lock");
7156        assert!(captured.iter().all(|event| match event {
7157            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7158            _ => true,
7159        }));
7160        let visible = captured
7161            .iter()
7162            .filter_map(|event| match event {
7163                AgentEvent::Text { text, .. } => Some(text.as_str()),
7164                _ => None,
7165            })
7166            .collect::<String>();
7167        assert_eq!(visible, "visible  trailing");
7168    }
7169
7170    #[tokio::test]
7171    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7172        let first = ScriptedAdapter::new(
7173            0,
7174            AgentCapabilities::default(),
7175            [
7176                AgentEvent::Text {
7177                    slot: 0,
7178                    text: "done".into(),
7179                },
7180                AgentEvent::TurnComplete { slot: 0 },
7181            ],
7182        );
7183        let reviewer = ScriptedAdapter::new(
7184            1,
7185            AgentCapabilities::default(),
7186            [
7187                AgentEvent::Text {
7188                    slot: 1,
7189                    text: STOP_TOKEN.into(),
7190                },
7191                AgentEvent::TurnComplete { slot: 1 },
7192            ],
7193        );
7194        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7195        let captured = std::sync::Arc::clone(&events);
7196        let mut relay = RelayHost::new(
7197            vec![
7198                AdapterHost::new(Box::new(first), None),
7199                AdapterHost::new(Box::new(reviewer), None),
7200            ],
7201            4,
7202        )
7203        .expect("relay");
7204        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7205        relay.start().await.expect("start");
7206        relay.run_turn("task", 0).await.expect("first turn");
7207        relay.run_turn("", 0).await.expect("review turn");
7208        let captured = events.lock().expect("lock");
7209        assert!(captured.iter().any(|event| {
7210            matches!(
7211                event,
7212                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7213            )
7214        }));
7215        assert!(captured.iter().all(|event| match event {
7216            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7217            _ => true,
7218        }));
7219        let acknowledgment = captured
7220            .iter()
7221            .position(|event| {
7222                matches!(
7223                    event,
7224                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7225                )
7226            })
7227            .expect("visible acknowledgment");
7228        let completion = captured
7229            .iter()
7230            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7231            .expect("reviewer completion");
7232        assert!(acknowledgment < completion);
7233    }
7234
7235    #[tokio::test]
7236    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7237        let first = ScriptedAdapter::new(
7238            0,
7239            AgentCapabilities::default(),
7240            [
7241                AgentEvent::Text {
7242                    slot: 0,
7243                    text: "done".into(),
7244                },
7245                AgentEvent::TurnComplete { slot: 0 },
7246            ],
7247        );
7248        let reviewer = ScriptedAdapter::new(
7249            1,
7250            AgentCapabilities::default(),
7251            [
7252                AgentEvent::Text {
7253                    slot: 1,
7254                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7255                },
7256                AgentEvent::TurnComplete { slot: 1 },
7257            ],
7258        );
7259        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7260        let captured = std::sync::Arc::clone(&events);
7261        let mut relay = RelayHost::new(
7262            vec![
7263                AdapterHost::new(Box::new(first), None),
7264                AdapterHost::new(Box::new(reviewer), None),
7265            ],
7266            4,
7267        )
7268        .expect("relay");
7269        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7270        relay.start().await.expect("start");
7271        relay.run_turn("task", 0).await.expect("first turn");
7272        relay.run_turn("", 0).await.expect("review turn");
7273
7274        let visible = events
7275            .lock()
7276            .expect("lock")
7277            .iter()
7278            .filter_map(|event| match event {
7279                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7280                _ => None,
7281            })
7282            .collect::<String>();
7283        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7284        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7285    }
7286
7287    #[tokio::test]
7288    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7289        let first = AdapterHost::new(
7290            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7291            None,
7292        );
7293        let second = AdapterHost::new(
7294            Box::new(ScriptedAdapter::new(
7295                1,
7296                AgentCapabilities::default(),
7297                [AgentEvent::TurnComplete { slot: 1 }],
7298            )),
7299            None,
7300        );
7301        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7302        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7303        relay.set_event_sink(move |event| {
7304            if matches!(event, AgentEvent::Permission { .. }) {
7305                let _ = seen_sender.send(());
7306            }
7307        });
7308        relay.start().await.expect("start");
7309        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7310        let answer = async move {
7311            seen_receiver.recv().await.expect("permission request");
7312            sender
7313                .send(super::RelayPermissionAnswer {
7314                    slot: 0,
7315                    request_id: "permission-1".into(),
7316                    answer: PermissionAnswer::Selected {
7317                        option_id: "allow".into(),
7318                    },
7319                })
7320                .expect("queue permission answer");
7321        };
7322        tokio::time::timeout(std::time::Duration::from_millis(100), async {
7323            let ((), result) = tokio::join!(
7324                answer,
7325                relay.run_turn_with_permissions("task", 0, &mut receiver)
7326            );
7327            result
7328        })
7329        .await
7330        .expect("permission-gated turn should not deadlock")
7331        .expect("turn completes");
7332    }
7333
7334    #[tokio::test]
7335    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7336        let first = AdapterHost::new(
7337            Box::new(PendingAdapter {
7338                slot: 0,
7339                hang_on_cancel: false,
7340            }),
7341            None,
7342        );
7343        let second = AdapterHost::new(
7344            Box::new(ScriptedAdapter::new(
7345                1,
7346                AgentCapabilities::default(),
7347                [AgentEvent::TurnComplete { slot: 1 }],
7348            )),
7349            None,
7350        );
7351        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7352        relay.start().await.expect("start");
7353        let cancellation = relay.cancellation();
7354        let error = {
7355            let turn = relay.run_turn("task", 0);
7356            tokio::pin!(turn);
7357            cancellation.request();
7358            turn.await.expect_err("cancellation should stop turn")
7359        };
7360        assert!(error.to_string().contains("relay turn cancelled"));
7361
7362        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7363        relay
7364            .run_turn("", 1)
7365            .await
7366            .expect("replacement job reaches the selected peer");
7367        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7368        assert!(replacement.contains("replacement job"));
7369        assert!(replacement.contains("User "));
7370        assert!(replacement.contains(":\ntask"));
7371        let owner_updates = relay.relay_mut().unseen_context(0);
7372        assert!(owner_updates.contains("User "));
7373        assert!(owner_updates.contains(":\ntask"));
7374        assert!(owner_updates.contains(":\nreplacement job"));
7375    }
7376
7377    #[tokio::test]
7378    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7379        let first = AdapterHost::new(
7380            Box::new(PendingAdapter {
7381                slot: 0,
7382                hang_on_cancel: true,
7383            }),
7384            None,
7385        );
7386        let second = AdapterHost::new(
7387            Box::new(PendingAdapter {
7388                slot: 1,
7389                hang_on_cancel: false,
7390            }),
7391            None,
7392        );
7393        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7394        relay.start().await.expect("start");
7395        let cancellation = relay.cancellation();
7396        let turn = relay.run_turn("task", 0);
7397        tokio::pin!(turn);
7398        cancellation.request();
7399        let error = turn.await.expect_err("cancellation should stop turn");
7400        assert!(error.to_string().contains("timed out"));
7401    }
7402
7403    #[tokio::test]
7404    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7405        let event = [AgentEvent::TurnComplete { slot: 0 }];
7406        let first = AdapterHost::new(
7407            Box::new(ScriptedAdapter::new(
7408                0,
7409                AgentCapabilities::default(),
7410                event.clone(),
7411            )),
7412            None,
7413        );
7414        let second = AdapterHost::new(
7415            Box::new(ScriptedAdapter::new(
7416                1,
7417                AgentCapabilities::default(),
7418                [AgentEvent::TurnComplete { slot: 1 }],
7419            )),
7420            None,
7421        );
7422        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7423        relay.start().await.expect("start");
7424
7425        relay.pause();
7426        assert_eq!(
7427            relay.run_turn("paused", 0).await.expect("paused turn"),
7428            crate::relay::RelayDecision::Paused
7429        );
7430        assert!(relay.dispatches().is_empty());
7431
7432        relay.resume();
7433        relay.relay_mut().drop_agent(1).expect("drop reviewer");
7434        assert!(matches!(
7435            relay
7436                .run_turn("solo follow-up", 0)
7437                .await
7438                .expect("solo turn"),
7439            crate::relay::RelayDecision::Dispatch {
7440                slot: 0,
7441                can_stop: false,
7442                ..
7443            }
7444        ));
7445        assert_eq!(relay.dispatches().len(), 1);
7446    }
7447
7448    #[tokio::test]
7449    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7450        let first = AdapterHost::new(
7451            Box::new(ScriptedAdapter::new(
7452                0,
7453                AgentCapabilities::default(),
7454                [AgentEvent::TurnComplete { slot: 0 }],
7455            )),
7456            None,
7457        );
7458        let second = AdapterHost::new(
7459            Box::new(ScriptedAdapter::new(
7460                1,
7461                AgentCapabilities::default(),
7462                [AgentEvent::TurnComplete { slot: 1 }],
7463            )),
7464            None,
7465        );
7466        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7467        relay.set_roster_names(vec!["First".into(), "Second".into()]);
7468        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7469        relay.set_roster_launch_specs(vec![
7470            ("custom".into(), "owner".into()),
7471            ("custom".into(), "peer".into()),
7472        ]);
7473        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7474        let captured = std::sync::Arc::clone(&events);
7475        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7476        relay.start().await.expect("start");
7477        let slot = relay
7478            .add_agent(
7479                AdapterHost::new(
7480                    Box::new(ScriptedAdapter::new(
7481                        2,
7482                        AgentCapabilities::default(),
7483                        [AgentEvent::TurnComplete { slot: 2 }],
7484                    )),
7485                    None,
7486                ),
7487                "Reviewer",
7488                "reviewer.example",
7489                "reviewer --acp",
7490            )
7491            .await
7492            .expect("append agent");
7493        assert_eq!(slot, 2);
7494        assert_eq!(
7495            relay.relay().active_slots().collect::<Vec<_>>(),
7496            vec![0, 1, 2]
7497        );
7498        assert_eq!(
7499            relay
7500                .session_metadata()
7501                .get("agents")
7502                .and_then(|value| value.as_array())
7503                .map(Vec::len),
7504            Some(3)
7505        );
7506        relay.drop_agent(1).await.expect("drop middle peer");
7507        let metadata = relay.session_metadata();
7508        assert_eq!(
7509            metadata.get("agents"),
7510            Some(&serde_json::json!([
7511                {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7512                {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7513            ]))
7514        );
7515        assert!(
7516            events
7517                .lock()
7518                .expect("lock")
7519                .iter()
7520                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7521        );
7522    }
7523
7524    #[tokio::test]
7525    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7526        let path = unique_test_path("codeswarm-session-metadata", "json");
7527        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7528        let writer = metadata_store.buffered().expect("metadata writer");
7529        let first = AdapterHost::new(
7530            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7531            None,
7532        );
7533        let second = AdapterHost::new(
7534            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7535            None,
7536        );
7537        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7538        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7539        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7540        relay.set_roster_launch_specs(vec![
7541            ("custom".into(), "claude".into()),
7542            ("custom".into(), "codex".into()),
7543        ]);
7544        relay.set_session_metadata_writer(writer);
7545        relay.start().await.expect("start");
7546        relay.drop_agent(0).await.expect("drop first agent");
7547        relay.stop().await.expect("stop");
7548
7549        let loaded = metadata_store
7550            .read()
7551            .expect("read metadata")
7552            .expect("metadata snapshot");
7553        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7554        assert_eq!(
7555            loaded.get("agents"),
7556            Some(&serde_json::json!([{
7557                "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7558                "command": "codex", "supports_load_session": false
7559            }]))
7560        );
7561        assert!(loaded.get("owner").is_none());
7562        let _ = std::fs::remove_file(path);
7563    }
7564
7565    #[tokio::test]
7566    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7567        let first = AdapterHost::new(
7568            Box::new(ScriptedAdapter::new(
7569                0,
7570                AgentCapabilities::default(),
7571                [
7572                    AgentEvent::Text {
7573                        slot: 0,
7574                        text: "owner stream".into(),
7575                    },
7576                    AgentEvent::TurnComplete { slot: 0 },
7577                ],
7578            )),
7579            None,
7580        );
7581        let second = AdapterHost::new(
7582            Box::new(ScriptedAdapter::new(
7583                1,
7584                AgentCapabilities::default(),
7585                [
7586                    AgentEvent::Text {
7587                        slot: 1,
7588                        text: "peer stream".into(),
7589                    },
7590                    AgentEvent::TurnComplete { slot: 1 },
7591                ],
7592            )),
7593            None,
7594        );
7595        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7596        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7597        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7598        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7599        let captured = std::sync::Arc::clone(&events);
7600        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7601        relay.start().await.expect("start");
7602
7603        relay.swap_agents(0, 1).expect("swap peers");
7604        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7605        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7606        relay.run_turn("task", 0).await.expect("swapped turn");
7607        let events = events.lock().expect("events");
7608        assert!(events.iter().any(|event| {
7609            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7610        }));
7611        assert!(relay.dispatches()[0].1.contains("You are Peer"));
7612    }
7613
7614    #[tokio::test]
7615    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7616        let path = unique_test_path("codeswarm-session-metadata", "json");
7617        let first = AdapterHost::new(
7618            Box::new(ScriptedAdapter::new(
7619                0,
7620                AgentCapabilities::default(),
7621                [AgentEvent::TurnComplete { slot: 0 }],
7622            )),
7623            None,
7624        );
7625        let second = AdapterHost::new(
7626            Box::new(ScriptedAdapter::new(
7627                1,
7628                AgentCapabilities::default(),
7629                [AgentEvent::TurnComplete { slot: 1 }],
7630            )),
7631            None,
7632        );
7633        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7634        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7635        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7636        relay.set_roster_launch_specs(vec![
7637            ("custom".into(), "claude".into()),
7638            ("custom".into(), "codex".into()),
7639        ]);
7640        let writer = SessionMetadataStore::open(&path)
7641            .buffered()
7642            .expect("metadata writer");
7643        relay.set_session_metadata_writer(writer);
7644        relay.start().await.expect("start");
7645        relay.stop().await.expect("stop");
7646        let loaded = SessionMetadataStore::open(&path)
7647            .read()
7648            .expect("read metadata")
7649            .expect("metadata snapshot");
7650        let agents = loaded
7651            .get("agents")
7652            .and_then(|value| value.as_array())
7653            .expect("agents");
7654        assert_eq!(agents.len(), 2);
7655        assert_eq!(agents[0]["identity"], "claude.com");
7656        assert_eq!(agents[1]["identity"], "openai.com");
7657        let _ = std::fs::remove_file(path);
7658    }
7659
7660    #[tokio::test]
7661    async fn relay_host_routes_unseen_public_context_to_next_agent() {
7662        let first = AdapterHost::new(
7663            Box::new(ScriptedAdapter::new(
7664                0,
7665                AgentCapabilities::default(),
7666                [
7667                    AgentEvent::Text {
7668                        slot: 0,
7669                        text: "implemented the fix".into(),
7670                    },
7671                    AgentEvent::TurnComplete { slot: 0 },
7672                ],
7673            )),
7674            None,
7675        );
7676        let second = AdapterHost::new(
7677            Box::new(ScriptedAdapter::new(
7678                1,
7679                AgentCapabilities::default(),
7680                [AgentEvent::TurnComplete { slot: 1 }],
7681            )),
7682            None,
7683        );
7684        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7685        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7686        relay.start().await.expect("start");
7687        relay.run_turn("task", 0).await.expect("first turn");
7688        relay.run_turn("review this", 0).await.expect("review turn");
7689
7690        assert_eq!(relay.dispatches().len(), 2);
7691        assert_eq!(relay.dispatches()[0].0, 0);
7692        assert!(relay.dispatches()[0].1.contains("task"));
7693        assert!(relay.dispatches()[0].1.contains("You are Codex"));
7694        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7695        assert_eq!(relay.dispatches()[1].0, 1);
7696        assert!(relay.dispatches()[1].1.contains("review this"));
7697        let public = relay.dispatches()[1]
7698            .1
7699            .split_once("Public updates:\n")
7700            .map(|(_, updates)| updates)
7701            .expect("review receives public context");
7702        let header = public
7703            .lines()
7704            .find(|line| line.starts_with("Codex "))
7705            .expect("named previous agent");
7706        let timestamp = header
7707            .strip_prefix("Codex ")
7708            .and_then(|value| value.strip_suffix(':'))
7709            .expect("timestamped header");
7710        assert_eq!(timestamp.len(), 5);
7711        assert_eq!(timestamp.as_bytes()[2], b':');
7712        assert!(
7713            timestamp
7714                .bytes()
7715                .enumerate()
7716                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7717        );
7718        assert!(public.contains("implemented the fix"));
7719        assert!(!public.contains("Agent 0"));
7720    }
7721}