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    prompt_had_output: bool,
2508    queued_events: VecDeque<AdapterResult<AgentEvent>>,
2509    tool_updates: BTreeMap<String, ToolUpdate>,
2510    stderr_task: Option<tokio::task::JoinHandle<String>>,
2511    terminals: BTreeMap<String, TerminalProcess>,
2512    next_terminal_id: u64,
2513}
2514
2515impl AcpAdapter {
2516    pub fn new(
2517        slot: RosterSlot,
2518        cwd: PathBuf,
2519        program: impl Into<String>,
2520        args: Vec<String>,
2521    ) -> Self {
2522        Self {
2523            slot,
2524            program: program.into(),
2525            args,
2526            cwd,
2527            child: None,
2528            reader: None,
2529            capabilities: AgentCapabilities::default(),
2530            modes: Vec::new(),
2531            models: Vec::new(),
2532            model_config_id: None,
2533            session_id: None,
2534            next_request_id: 1,
2535            prompt_request_id: None,
2536            prompt_had_output: false,
2537            queued_events: VecDeque::new(),
2538            tool_updates: BTreeMap::new(),
2539            stderr_task: None,
2540            terminals: BTreeMap::new(),
2541            next_terminal_id: 1,
2542        }
2543    }
2544
2545    pub fn with_session_id(
2546        slot: RosterSlot,
2547        cwd: PathBuf,
2548        program: impl Into<String>,
2549        args: Vec<String>,
2550        session_id: impl Into<String>,
2551    ) -> Self {
2552        let mut adapter = Self::new(slot, cwd, program, args);
2553        adapter.session_id = Some(session_id.into());
2554        adapter
2555    }
2556
2557    async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2558        self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2559            .await
2560    }
2561
2562    async fn request_with_timeout(
2563        &mut self,
2564        method: &str,
2565        params: Value,
2566        deadline: std::time::Duration,
2567    ) -> AdapterResult<Value> {
2568        tokio::time::timeout(deadline, self.request_inner(method, params))
2569            .await
2570            .map_err(|_| {
2571                AdapterError::Transport(format!(
2572                    "ACP {method} timed out; reload the agent to retry"
2573                ))
2574            })?
2575    }
2576
2577    async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2578        let request_id = self.next_request_id;
2579        self.next_request_id += 1;
2580        self.write_json(serde_json::json!({
2581            "jsonrpc": "2.0",
2582            "id": request_id,
2583            "method": method,
2584            "params": params,
2585        }))
2586        .await?;
2587        loop {
2588            let line = self.read_line().await?;
2589            let value: Value = match serde_json::from_str(&line) {
2590                Ok(value) => value,
2591                Err(_) => {
2592                    // Keep the transport alive when a peer writes a stray
2593                    // diagnostic line. This is common with CLI wrappers and
2594                    // matches the baseline client's tolerant stream loop.
2595                    continue;
2596                }
2597            };
2598            if self.reject_empty_permission_request(&value).await? {
2599                continue;
2600            }
2601            if self.handle_client_request(&value).await? {
2602                continue;
2603            }
2604            if value
2605                .get("id")
2606                .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2607            {
2608                if let Some(error) = value.get("error") {
2609                    return Err(AdapterError::Protocol(error.to_string()));
2610                }
2611                return value
2612                    .get("result")
2613                    .cloned()
2614                    .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2615            }
2616            if let Some(event) = parse_acp_value(self.slot, &value, &mut self.tool_updates)? {
2617                let event = if method == "session/load" {
2618                    restored_history_event(event)
2619                } else {
2620                    event
2621                };
2622                self.queued_events.push_back(Ok(event));
2623            }
2624        }
2625    }
2626
2627    async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2628        let child = self
2629            .child
2630            .as_mut()
2631            .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2632        let stdin = child
2633            .stdin
2634            .as_mut()
2635            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2636        stdin
2637            .write_all(value.to_string().as_bytes())
2638            .await
2639            .map_err(|error| AdapterError::Transport(error.to_string()))?;
2640        stdin
2641            .write_all(b"\n")
2642            .await
2643            .map_err(|error| AdapterError::Transport(error.to_string()))
2644    }
2645
2646    /// Tear down a broken ACP stream while retaining the provider session
2647    /// handle. The coordinator reloads the slot before its next prompt.
2648    async fn reset_transport(&mut self) {
2649        let terminals = std::mem::take(&mut self.terminals);
2650        for terminal in terminals.values() {
2651            terminal.stop().await;
2652        }
2653        self.queued_events.clear();
2654        self.tool_updates.clear();
2655        if let Some(mut child) = self.child.take() {
2656            let _ = terminate_child(&mut child).await;
2657        }
2658        self.reader = None;
2659        if let Some(task) = self.stderr_task.take() {
2660            task.abort();
2661        }
2662        self.prompt_request_id = None;
2663    }
2664
2665    /// ACP permission requests are JSON-RPC requests, not fire-and-forget
2666    /// notifications. An empty option list is invalid and must be answered
2667    /// with an error so the peer does not wait forever for a decision. This is
2668    /// the same validation performed by the Python ACP server.
2669    async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2670        if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2671            || value.get("id").is_none()
2672        {
2673            return Ok(false);
2674        }
2675        let valid = value
2676            .get("params")
2677            .and_then(|params| params.get("options"))
2678            .and_then(Value::as_array)
2679            .is_some_and(|options| !options.is_empty());
2680        if valid {
2681            return Ok(false);
2682        }
2683        self.write_json(serde_json::json!({
2684            "jsonrpc": "2.0",
2685            "id": value.get("id").cloned().unwrap_or(Value::Null),
2686            "error": {
2687                "code": -32602,
2688                "message": "Permission request requires at least one option",
2689            },
2690        }))
2691        .await?;
2692        Ok(true)
2693    }
2694
2695    fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2696        let root = self
2697            .cwd
2698            .canonicalize()
2699            .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2700        let requested = Path::new(path);
2701        let candidate = if requested.is_absolute() {
2702            requested.to_path_buf()
2703        } else {
2704            root.join(requested)
2705        };
2706        let resolved = if !candidate.exists() {
2707            let parent = candidate
2708                .parent()
2709                .ok_or_else(|| "file path has no parent".to_owned())?
2710                .canonicalize()
2711                .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2712            parent.join(
2713                candidate
2714                    .file_name()
2715                    .ok_or_else(|| "file path has no filename".to_owned())?,
2716            )
2717        } else {
2718            candidate
2719                .canonicalize()
2720                .map_err(|error| format!("unable to resolve file path: {error}"))?
2721        };
2722        if !resolved.starts_with(&root) {
2723            return Err("file path is outside the project".into());
2724        }
2725        Ok(resolved)
2726    }
2727
2728    fn read_workspace_text(
2729        &self,
2730        path: &str,
2731        line: Option<i64>,
2732        limit: Option<i64>,
2733    ) -> Result<String, String> {
2734        if line.is_some_and(|line| line < 1) {
2735            return Err("line must be positive".into());
2736        }
2737        if limit.is_some_and(|limit| limit < 0) {
2738            return Err("limit must not be negative".into());
2739        }
2740        let path = self.workspace_path(path)?;
2741        let mut bytes = Vec::new();
2742        let mut source = match std::fs::File::open(path) {
2743            Ok(source) => source,
2744            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2745            Err(error) => return Err(error.to_string()),
2746        };
2747        source
2748            .by_ref()
2749            .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2750            .read_to_end(&mut bytes)
2751            .map_err(|error| error.to_string())?;
2752        bytes.truncate(MAX_FILE_READ_BYTES);
2753        let text = String::from_utf8_lossy(&bytes);
2754        if line.is_none() && limit.is_none() {
2755            return Ok(text.into_owned());
2756        }
2757        let start = line.map_or(0, |line| line as usize - 1);
2758        let limit = limit.unwrap_or(i64::MAX) as usize;
2759        let selected = text
2760            .split_inclusive('\n')
2761            .skip(start)
2762            .take(limit)
2763            .collect::<String>();
2764        if line.is_some() {
2765            Ok(selected.trim_end_matches('\n').to_owned())
2766        } else {
2767            Ok(selected)
2768        }
2769    }
2770
2771    fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2772        let path = params
2773            .get("path")
2774            .and_then(Value::as_str)
2775            .filter(|path| !path.is_empty())
2776            .ok_or("path must be a non-empty string")?;
2777        let content = params
2778            .get("content")
2779            .and_then(Value::as_str)
2780            .ok_or("content must be a string")?;
2781        let path = self.workspace_path(path)?;
2782        std::fs::write(path, content).map_err(|error| error.to_string())
2783    }
2784
2785    async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2786        let command = params
2787            .get("command")
2788            .and_then(Value::as_str)
2789            .filter(|command| !command.trim().is_empty())
2790            .ok_or_else(|| "terminal command is required".to_owned())?;
2791        let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2792        let cwd = self.workspace_path(cwd)?;
2793        if !cwd.is_dir() {
2794            return Err("terminal cwd is not a directory".into());
2795        }
2796        let mut process = Command::new(command);
2797        isolate_process_group(&mut process);
2798        if let Some(args) = params.get("args").and_then(Value::as_array) {
2799            process.args(args.iter().filter_map(Value::as_str));
2800        }
2801        process
2802            .current_dir(&cwd)
2803            .stdin(Stdio::null())
2804            .stdout(Stdio::piped())
2805            .stderr(Stdio::piped());
2806        if let Some(env) = params.get("env") {
2807            if let Some(entries) = env.as_array() {
2808                for entry in entries {
2809                    if let (Some(name), Some(value)) = (
2810                        entry.get("name").and_then(Value::as_str),
2811                        entry.get("value").and_then(Value::as_str),
2812                    ) {
2813                        process.env(name, value);
2814                    }
2815                }
2816            } else if let Some(entries) = env.as_object() {
2817                for (name, value) in entries {
2818                    if let Some(value) = value.as_str() {
2819                        process.env(name, value);
2820                    }
2821                }
2822            }
2823        }
2824        let mut child = process.spawn().map_err(|error| error.to_string())?;
2825        let stdout = child.stdout.take();
2826        let stderr = child.stderr.take();
2827        let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2828            let _ = terminate_child(&mut child).await;
2829            return Err("terminal has no output pipes".into());
2830        };
2831        let output = Arc::new(Mutex::new(Vec::new()));
2832        let truncated = Arc::new(AtomicBool::new(false));
2833        let output_readers = Arc::new(AtomicUsize::new(2));
2834        let output_limit = params
2835            .get("outputByteLimit")
2836            .and_then(Value::as_u64)
2837            .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2838                usize::try_from(limit)
2839                    .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2840                    .min(MAX_TERMINAL_OUTPUT_BYTES)
2841            });
2842        tokio::spawn(drain_terminal_output(
2843            stdout,
2844            Arc::clone(&output),
2845            Arc::clone(&truncated),
2846            Arc::clone(&output_readers),
2847            output_limit,
2848        ));
2849        tokio::spawn(drain_terminal_output(
2850            stderr,
2851            Arc::clone(&output),
2852            Arc::clone(&truncated),
2853            Arc::clone(&output_readers),
2854            output_limit,
2855        ));
2856        let id = format!("terminal-{}", self.next_terminal_id);
2857        self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2858        let state = TerminalProcess {
2859            child: Arc::new(AsyncMutex::new(Some(child))),
2860            output,
2861            truncated,
2862            output_readers,
2863        };
2864        self.terminals.insert(id.clone(), state);
2865        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2866            slot: self.slot,
2867            event: TerminalEvent::Created {
2868                id: id.clone(),
2869                command: std::iter::once(command)
2870                    .chain(
2871                        params
2872                            .get("args")
2873                            .and_then(Value::as_array)
2874                            .into_iter()
2875                            .flatten()
2876                            .filter_map(Value::as_str),
2877                    )
2878                    .collect::<Vec<_>>()
2879                    .join(" "),
2880            },
2881        }));
2882        Ok(serde_json::json!({"terminalId": id}))
2883    }
2884
2885    async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2886        let terminal = self
2887            .terminals
2888            .get(id)
2889            .ok_or_else(|| "terminal not found".to_owned())?
2890            .clone();
2891        let output = terminal
2892            .output
2893            .lock()
2894            .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2895            .unwrap_or_default();
2896        let exit_code = terminal.exit_code().await;
2897        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2898            slot: self.slot,
2899            event: TerminalEvent::Output {
2900                id: id.to_owned(),
2901                text: output.clone(),
2902            },
2903        }));
2904        let mut response = serde_json::json!({
2905            "output": output,
2906            "truncated": terminal.truncated.load(Ordering::Acquire),
2907        });
2908        if let Some(code) = exit_code {
2909            response["exitStatus"] = serde_json::json!({"exitCode": code});
2910        }
2911        Ok(response)
2912    }
2913
2914    async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2915        let terminal = self
2916            .terminals
2917            .get(id)
2918            .ok_or_else(|| "terminal not found".to_owned())?
2919            .clone();
2920        let exit_code = terminal.wait().await;
2921        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2922            slot: self.slot,
2923            event: TerminalEvent::Exited {
2924                id: id.to_owned(),
2925                code: exit_code.unwrap_or(-1),
2926            },
2927        }));
2928        Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2929    }
2930
2931    /// Handle requests initiated by an ACP agent against the client. File
2932    /// access is mediated through the configured workspace root; unsupported
2933    /// requests receive a JSON-RPC error instead of hanging the agent.
2934    async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2935        let Some(method) = value.get("method").and_then(Value::as_str) else {
2936            return Ok(false);
2937        };
2938        let Some(id) = value.get("id").cloned() else {
2939            return Ok(false);
2940        };
2941        // Permission requests are consumed by the normalized event parser and
2942        // answered later through the focused UI action, not here.
2943        if method == "session/request_permission" {
2944            return Ok(false);
2945        }
2946        let params = value.get("params").cloned().unwrap_or(Value::Null);
2947        let response = match method {
2948            "fs/read_text_file" => {
2949                let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2950                let line = params.get("line").and_then(Value::as_i64);
2951                let limit = params.get("limit").and_then(Value::as_i64);
2952                match self.read_workspace_text(path, line, limit) {
2953                    Ok(content) => serde_json::json!({
2954                        "jsonrpc": "2.0",
2955                        "id": id,
2956                        "result": {"content": content},
2957                    }),
2958                    Err(message) => serde_json::json!({
2959                        "jsonrpc": "2.0",
2960                        "id": id,
2961                        "error": {"code": -32602, "message": message},
2962                    }),
2963                }
2964            }
2965            "fs/write_text_file" => {
2966                let result = self.write_workspace_text(&params);
2967                match result {
2968                    Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2969                    Err(message) => serde_json::json!({
2970                        "jsonrpc": "2.0",
2971                        "id": id,
2972                        "error": {"code": -32602, "message": message},
2973                    }),
2974                }
2975            }
2976            "terminal/create" => match self.terminal_create(&params).await {
2977                Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2978                Err(message) => serde_json::json!({
2979                    "jsonrpc": "2.0",
2980                    "id": id,
2981                    "error": {"code": -32602, "message": message},
2982                }),
2983            },
2984            "terminal/output" => {
2985                let terminal_id = params
2986                    .get("terminalId")
2987                    .and_then(Value::as_str)
2988                    .unwrap_or("");
2989                match self.terminal_output(terminal_id).await {
2990                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2991                    Err(message) => serde_json::json!({
2992                        "jsonrpc": "2.0",
2993                        "id": id,
2994                        "error": {"code": -32602, "message": message},
2995                    }),
2996                }
2997            }
2998            "terminal/wait_for_exit" => {
2999                let terminal_id = params
3000                    .get("terminalId")
3001                    .and_then(Value::as_str)
3002                    .unwrap_or("");
3003                match self.terminal_wait(terminal_id).await {
3004                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
3005                    Err(message) => serde_json::json!({
3006                        "jsonrpc": "2.0",
3007                        "id": id,
3008                        "error": {"code": -32602, "message": message},
3009                    }),
3010                }
3011            }
3012            "terminal/kill" => {
3013                let terminal_id = params
3014                    .get("terminalId")
3015                    .and_then(Value::as_str)
3016                    .unwrap_or("");
3017                if let Some(terminal) = self.terminals.get(terminal_id) {
3018                    terminal.kill().await;
3019                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3020                } else {
3021                    serde_json::json!({
3022                        "jsonrpc": "2.0",
3023                        "id": id,
3024                        "error": {"code": -32602, "message": "terminal not found"},
3025                    })
3026                }
3027            }
3028            "terminal/release" => {
3029                let terminal_id = params
3030                    .get("terminalId")
3031                    .and_then(Value::as_str)
3032                    .unwrap_or("");
3033                if let Some(terminal) = self.terminals.remove(terminal_id) {
3034                    terminal.stop().await;
3035                    self.queued_events.push_back(Ok(AgentEvent::Terminal {
3036                        slot: self.slot,
3037                        event: TerminalEvent::Released {
3038                            id: terminal_id.to_owned(),
3039                        },
3040                    }));
3041                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
3042                } else {
3043                    serde_json::json!({
3044                        "jsonrpc": "2.0",
3045                        "id": id,
3046                        "error": {"code": -32602, "message": "terminal not found"},
3047                    })
3048                }
3049            }
3050            _ => serde_json::json!({
3051                "jsonrpc": "2.0",
3052                "id": id,
3053                "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
3054            }),
3055        };
3056        self.write_json(response).await?;
3057        Ok(true)
3058    }
3059
3060    async fn read_line(&mut self) -> AdapterResult<String> {
3061        let reader = self
3062            .reader
3063            .as_mut()
3064            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
3065        read_bounded_line(reader).await
3066    }
3067
3068    async fn start(&mut self) -> AdapterResult<()> {
3069        self.modes.clear();
3070        // Starting an adapter twice must never orphan the first transport.
3071        if self.child.is_some() {
3072            self.stop().await?;
3073        }
3074        let mut command = Command::new(&self.program);
3075        isolate_process_group(&mut command);
3076        command
3077            .args(&self.args)
3078            .current_dir(&self.cwd)
3079            .stdin(Stdio::piped())
3080            .stdout(Stdio::piped())
3081            .stderr(Stdio::piped())
3082            .env("CODESWARM_CWD", &self.cwd);
3083        if self.program.to_ascii_lowercase().contains("gemini")
3084            || self
3085                .args
3086                .iter()
3087                .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
3088        {
3089            command.env("GEMINI_TELEMETRY_ENABLED", "false");
3090        }
3091        let mut child = command
3092            .spawn()
3093            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
3094        let stdout = match child.stdout.take() {
3095            Some(stdout) => stdout,
3096            None => {
3097                let _ = terminate_child(&mut child).await;
3098                return Err(AdapterError::Transport("ACP agent has no stdout".into()));
3099            }
3100        };
3101        let stderr = match child.stderr.take() {
3102            Some(stderr) => stderr,
3103            None => {
3104                let _ = terminate_child(&mut child).await;
3105                return Err(AdapterError::Transport("ACP agent has no stderr".into()));
3106            }
3107        };
3108        self.child = Some(child);
3109        self.reader = Some(BufReader::new(stdout));
3110        self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
3111
3112        let initialize = match self
3113            .request(
3114                "initialize",
3115                serde_json::json!({
3116                    "protocolVersion": 1,
3117                    "clientCapabilities": {
3118                        "fs": {"readTextFile": true, "writeTextFile": true},
3119                        "terminal": true,
3120                    },
3121                    "clientInfo": {
3122                        "name": "CodeSwarm",
3123                        "title": "CodeSwarm",
3124                        "version": env!("CARGO_PKG_VERSION"),
3125                    },
3126                }),
3127            )
3128            .await
3129        {
3130            Ok(value) => value,
3131            Err(error) => {
3132                let _ = self.stop().await;
3133                return Err(error);
3134            }
3135        };
3136        let agent_capabilities = initialize
3137            .get("agentCapabilities")
3138            .cloned()
3139            .unwrap_or(Value::Null);
3140        self.capabilities = AgentCapabilities {
3141            supports_cancel: true,
3142            supports_modes: true,
3143            supports_permissions: true,
3144            supports_terminals: true,
3145            supports_session_load: agent_capabilities
3146                .get("loadSession")
3147                .and_then(Value::as_bool)
3148                .unwrap_or(false),
3149            supports_models: false,
3150        };
3151        let session = if let Some(session_id) = self.session_id.clone() {
3152            if !self.capabilities.supports_session_load {
3153                let _ = self.stop().await;
3154                return Err(AdapterError::Unsupported("session/load"));
3155            }
3156            match self
3157                .request(
3158                    "session/load",
3159                    serde_json::json!({
3160                        "cwd": self.cwd,
3161                        "mcpServers": [],
3162                        "sessionId": session_id,
3163                    }),
3164                )
3165                .await
3166            {
3167                Ok(value) => value,
3168                Err(error) => {
3169                    let _ = self.stop().await;
3170                    return Err(error);
3171                }
3172            }
3173        } else {
3174            let session = match self
3175                .request(
3176                    "session/new",
3177                    serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
3178                )
3179                .await
3180            {
3181                Ok(value) => value,
3182                Err(error) => {
3183                    let _ = self.stop().await;
3184                    return Err(error);
3185                }
3186            };
3187            self.session_id = session
3188                .get("sessionId")
3189                .and_then(Value::as_str)
3190                .map(str::to_owned);
3191            if self.session_id.is_none() {
3192                let _ = self.stop().await;
3193                return Err(AdapterError::Protocol(
3194                    "session/new returned no sessionId".into(),
3195                ));
3196            }
3197            session
3198        };
3199        self.capabilities.supports_modes = false;
3200        if let Some(modes) = session.get("modes") {
3201            let available = modes
3202                .get("availableModes")
3203                .and_then(Value::as_array)
3204                .map(|modes| {
3205                    modes
3206                        .iter()
3207                        .filter_map(|mode| {
3208                            Some(Mode {
3209                                id: mode.get("id")?.as_str()?.to_owned(),
3210                                label: mode.get("name")?.as_str()?.to_owned(),
3211                            })
3212                        })
3213                        .collect::<Vec<_>>()
3214                })
3215                .unwrap_or_default();
3216            self.modes = available.clone();
3217            self.capabilities.supports_modes = !available.is_empty();
3218            self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
3219                slot: self.slot,
3220                modes: available,
3221                current_mode: modes
3222                    .get("currentModeId")
3223                    .and_then(Value::as_str)
3224                    .map(str::to_owned),
3225            }));
3226        }
3227        self.models.clear();
3228        self.model_config_id = None;
3229        let current_model =
3230            parse_model_config(&session).and_then(|(config_id, models, current)| {
3231                self.model_config_id = Some(config_id);
3232                self.models = models;
3233                current
3234            });
3235        self.capabilities.supports_models =
3236            self.model_config_id.is_some() && !self.models.is_empty();
3237        if let Some(config_id) = self.model_config_id.clone()
3238            && !self.models.is_empty()
3239        {
3240            self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
3241                slot: self.slot,
3242                config_id,
3243                models: self.models.clone(),
3244                current_model,
3245            }));
3246        }
3247        self.queued_events.push_back(Ok(AgentEvent::Ready {
3248            slot: self.slot,
3249            capabilities: self.capabilities(),
3250        }));
3251        Ok(())
3252    }
3253}
3254
3255fn prompt_resource_paths(prompt: &str) -> Vec<String> {
3256    let characters = prompt.chars().collect::<Vec<_>>();
3257    let mut paths = Vec::new();
3258    let mut index = 0;
3259    while index < characters.len() {
3260        if characters[index] != '@' {
3261            index += 1;
3262            continue;
3263        }
3264        index += 1;
3265        let quoted = characters.get(index) == Some(&'"');
3266        if quoted {
3267            index += 1;
3268        }
3269        let start = index;
3270        while index < characters.len()
3271            && if quoted {
3272                characters[index] != '"'
3273            } else {
3274                !characters[index].is_whitespace()
3275            }
3276        {
3277            index += 1;
3278        }
3279        if index > start {
3280            paths.push(characters[start..index].iter().collect());
3281        }
3282        if quoted && index < characters.len() {
3283            index += 1;
3284        }
3285    }
3286    paths
3287}
3288
3289fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3290    let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3291    for path in prompt_resource_paths(prompt) {
3292        if path.ends_with('/') {
3293            continue;
3294        }
3295        let Ok(resource) = resources::load(cwd, &path) else {
3296            continue;
3297        };
3298        let uri = format!("file://{}", resource.path.display());
3299        let resource_value = if let Some(text) = resource.text {
3300            serde_json::json!({
3301                "uri": uri,
3302                "text": text,
3303                "mimeType": resource.mime_type,
3304            })
3305        } else if let Some(data) = resource.data {
3306            serde_json::json!({
3307                "uri": uri,
3308                "blob": BASE64.encode(data),
3309                "mimeType": resource.mime_type,
3310            })
3311        } else {
3312            continue;
3313        };
3314        blocks.push(serde_json::json!({
3315            "type": "resource",
3316            "resource": resource_value,
3317        }));
3318    }
3319    blocks
3320}
3321
3322#[async_trait]
3323impl AgentAdapter for AcpAdapter {
3324    fn slot(&self) -> RosterSlot {
3325        self.slot
3326    }
3327
3328    fn session_id(&self) -> Option<String> {
3329        self.session_id.clone()
3330    }
3331
3332    fn protocol(&self) -> &'static str {
3333        "acp"
3334    }
3335
3336    fn needs_restart(&self) -> bool {
3337        self.reader.is_none() || self.child.is_none()
3338    }
3339
3340    fn capabilities(&self) -> AgentCapabilities {
3341        self.capabilities.clone()
3342    }
3343
3344    async fn start(&mut self) -> AdapterResult<()> {
3345        // Use the inherent implementation, which owns the cleanup boundary
3346        // around the multi-step ACP handshake.
3347        AcpAdapter::start(self).await
3348    }
3349
3350    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3351        if self.needs_restart() {
3352            return Err(AdapterError::Transport(
3353                "ACP agent transport is not running; reload the agent before retrying".into(),
3354            ));
3355        }
3356        let session_id = self
3357            .session_id
3358            .as_ref()
3359            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3360        self.tool_updates.clear();
3361        let request_id = self.next_request_id;
3362        self.next_request_id += 1;
3363        let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3364        let write_result = self
3365            .write_json(serde_json::json!({
3366                "jsonrpc": "2.0",
3367                "id": request_id,
3368                "method": "session/prompt",
3369                "params": {
3370                    "sessionId": session_id,
3371                    "prompt": prompt_blocks,
3372                },
3373            }))
3374            .await;
3375        if let Err(error) = write_result {
3376            self.reset_transport().await;
3377            return Err(error);
3378        }
3379        self.prompt_request_id = Some(request_id);
3380        self.prompt_had_output = false;
3381        Ok(())
3382    }
3383
3384    async fn cancel(&mut self) -> AdapterResult<bool> {
3385        let Some(session_id) = &self.session_id else {
3386            return Ok(false);
3387        };
3388        self.write_json(serde_json::json!({
3389            "jsonrpc": "2.0",
3390            "method": "session/cancel",
3391            "params": {"sessionId": session_id, "_meta": {}},
3392        }))
3393        .await?;
3394        let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3395            loop {
3396                match <Self as AgentAdapter>::next_event(self).await {
3397                    Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3398                    Some(Ok(_)) => {}
3399                    Some(Err(_)) => break,
3400                }
3401            }
3402        })
3403        .await
3404        .is_ok();
3405        if !settled {
3406            // A peer that never acknowledges cancellation cannot safely share
3407            // its stream with the next prompt. Restart the transport while
3408            // preserving a loadable provider session when supported.
3409            self.reload().await?;
3410        }
3411        Ok(true)
3412    }
3413
3414    async fn answer_permission(
3415        &mut self,
3416        request_id: String,
3417        answer: PermissionAnswer,
3418    ) -> AdapterResult<()> {
3419        let id = request_id
3420            .parse::<u64>()
3421            .map(Value::from)
3422            .unwrap_or_else(|_| Value::String(request_id));
3423        let outcome = match answer {
3424            PermissionAnswer::Selected { option_id } => {
3425                serde_json::json!({"outcome": "selected", "optionId": option_id})
3426            }
3427            PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3428        };
3429        self.write_json(serde_json::json!({
3430            "jsonrpc": "2.0",
3431            "id": id,
3432            // RequestPermissionResponse wraps the selected/cancelled
3433            // discriminator in its `outcome` field. Keep this nested shape
3434            // compatible with the Python ACP server and ACP schema.
3435            "result": {"outcome": outcome},
3436        }))
3437        .await
3438    }
3439
3440    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3441        let session_id = self
3442            .session_id
3443            .as_ref()
3444            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3445        let policy = match mode.as_str() {
3446            "plan" => "codeswarm:mode:plan",
3447            "default" | "manual" => "codeswarm:mode:manual",
3448            "accept-edits" => "codeswarm:mode:accept-edits",
3449            "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3450            other => other,
3451        };
3452        let native_mode = crate::policy::resolve(policy, &self.modes)
3453            .map(|mode| mode.id)
3454            .unwrap_or(mode);
3455        let _ = self
3456            .request(
3457                "session/set_mode",
3458                serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3459            )
3460            .await?;
3461        self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3462            slot: self.slot,
3463            current_mode: native_mode,
3464        }));
3465        Ok(())
3466    }
3467
3468    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3469        let session_id = self
3470            .session_id
3471            .clone()
3472            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3473        let config_id = self
3474            .model_config_id
3475            .clone()
3476            .ok_or(AdapterError::Unsupported("set_model"))?;
3477        if !self.models.iter().any(|candidate| candidate.id == model) {
3478            return Err(AdapterError::Protocol(
3479                "model is not advertised by the agent".into(),
3480            ));
3481        }
3482        let _ = self
3483            .request(
3484                "session/set_config_option",
3485                serde_json::json!({
3486                    "sessionId": session_id,
3487                    "configId": config_id,
3488                    "value": model,
3489                }),
3490            )
3491            .await?;
3492        Ok(())
3493    }
3494
3495    async fn reload(&mut self) -> AdapterResult<()> {
3496        // `stop` tears down the process and clears its transport-owned
3497        // session handle. A reload is different from a final shutdown: ACP
3498        // peers advertising `loadSession` must receive the prior ID so the
3499        // replacement process can resume the same conversation.
3500        let session_id = self
3501            .capabilities
3502            .supports_session_load
3503            .then(|| self.session_id.clone())
3504            .flatten();
3505        self.stop().await?;
3506        self.session_id = session_id.clone();
3507        let result = self.start().await;
3508        if result.is_err() {
3509            // `start` cleans up a partially initialized transport by calling
3510            // `stop`, which also clears the handle. Keep it available for a
3511            // subsequent retry after the coordinator reports the failure.
3512            self.session_id = session_id;
3513        }
3514        result
3515    }
3516
3517    async fn stop(&mut self) -> AdapterResult<()> {
3518        let terminals = std::mem::take(&mut self.terminals);
3519        self.queued_events.clear();
3520        self.tool_updates.clear();
3521        for terminal in terminals.values() {
3522            terminal.stop().await;
3523        }
3524        if let Some(mut child) = self.child.take() {
3525            terminate_child(&mut child).await?;
3526        }
3527        self.reader = None;
3528        self.session_id = None;
3529        self.prompt_request_id = None;
3530        if let Some(task) = self.stderr_task.take() {
3531            task.abort();
3532            let _ = task.await;
3533        }
3534        Ok(())
3535    }
3536
3537    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3538        if let Some(event) = self.queued_events.pop_front() {
3539            return Some(event);
3540        }
3541        loop {
3542            let line = match self.read_line().await {
3543                Ok(line) => line,
3544                Err(error) => {
3545                    // EOF or a broken pipe ends this transport. Reap the
3546                    // process now; `send_prompt` will reconnect it for the
3547                    // next human turn while preserving `session_id`.
3548                    self.reset_transport().await;
3549                    return Some(Err(error));
3550                }
3551            };
3552            let value: Value = match serde_json::from_str(&line) {
3553                Ok(value) => value,
3554                Err(_) => continue,
3555            };
3556            match self.reject_empty_permission_request(&value).await {
3557                Ok(true) => continue,
3558                Ok(false) => {}
3559                Err(error) => return Some(Err(error)),
3560            }
3561            match self.handle_client_request(&value).await {
3562                Ok(true) => {
3563                    // Any agent-initiated client request (fs, terminal) proves
3564                    // the model is alive, even before it streams text.
3565                    if self.prompt_request_id.is_some() {
3566                        self.prompt_had_output = true;
3567                    }
3568                    continue;
3569                }
3570                Ok(false) => {}
3571                Err(error) => return Some(Err(error)),
3572            }
3573            match parse_acp_value(self.slot, &value, &mut self.tool_updates) {
3574                Ok(Some(event)) => {
3575                    if let AgentEvent::ModelsReplaced {
3576                        config_id, models, ..
3577                    } = &event
3578                    {
3579                        self.model_config_id = Some(config_id.clone());
3580                        self.models = models.clone();
3581                        self.capabilities.supports_models = !models.is_empty();
3582                    }
3583                    if self.prompt_request_id.is_some()
3584                        && matches!(
3585                            event,
3586                            AgentEvent::Text { .. }
3587                                | AgentEvent::Thought { .. }
3588                                | AgentEvent::Tool { .. }
3589                                | AgentEvent::Permission { .. }
3590                                | AgentEvent::Terminal { .. }
3591                        )
3592                    {
3593                        self.prompt_had_output = true;
3594                    }
3595                    return Some(Ok(event));
3596                }
3597                Ok(None) => {}
3598                Err(error) => return Some(Err(error)),
3599            }
3600            if value.get("id").is_some_and(|id| {
3601                self.prompt_request_id
3602                    .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3603            }) {
3604                if let Some(error) = value.get("error") {
3605                    self.prompt_request_id = None;
3606                    return Some(Err(AdapterError::Protocol(error.to_string())));
3607                }
3608                self.prompt_request_id = None;
3609                let stop_reason = value
3610                    .get("result")
3611                    .and_then(|result| result.get("stopReason"))
3612                    .and_then(Value::as_str);
3613                if let Some(reason) =
3614                    stop_reason.filter(|reason| !matches!(*reason, "end_turn" | "cancelled"))
3615                {
3616                    let detail = match reason {
3617                        "max_tokens" => {
3618                            "ACP turn stopped because the output token limit was reached".into()
3619                        }
3620                        "max_turn_requests" => {
3621                            "ACP turn stopped because the tool-turn limit was reached".into()
3622                        }
3623                        other => format!("ACP turn stopped before responding: {other}"),
3624                    };
3625                    return Some(Ok(AgentEvent::Failed {
3626                        slot: self.slot,
3627                        started: true,
3628                        detail,
3629                    }));
3630                }
3631                // OpenCode ends the turn with only a `user_message_chunk` echo
3632                // and zero usage when the configured model cannot run (bad
3633                // API key, quota, unknown model). Surface that as a failure
3634                // instead of a silent empty turn.
3635                if !self.prompt_had_output && !matches!(stop_reason, Some("cancelled")) {
3636                    return Some(Ok(AgentEvent::Failed {
3637                        slot: self.slot,
3638                        started: true,
3639                        detail: "ACP turn ended with no agent output; the configured model may be unauthorized or out of quota — pick another model with /model or verify with `opencode run`".into(),
3640                    }));
3641                }
3642                return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3643            }
3644        }
3645    }
3646}
3647
3648#[cfg(test)]
3649fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3650    let value: Value =
3651        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3652    parse_acp_value(slot, &value, &mut BTreeMap::new())
3653}
3654
3655fn parse_acp_value(
3656    slot: RosterSlot,
3657    value: &Value,
3658    tools: &mut BTreeMap<String, ToolUpdate>,
3659) -> AdapterResult<Option<AgentEvent>> {
3660    let method = value.get("method").and_then(Value::as_str);
3661    if method == Some("session/request_permission") {
3662        let params = value.get("params").cloned().unwrap_or(Value::Null);
3663        let request_id = value
3664            .get("id")
3665            .map(rpc_id_to_string)
3666            .unwrap_or_else(|| "permission".into());
3667        return Ok(parse_permission_event(
3668            slot,
3669            &params,
3670            &request_id,
3671            params.get("options"),
3672        ));
3673    }
3674    if method != Some("session/update") {
3675        return Ok(None);
3676    }
3677    let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3678        return Ok(None);
3679    };
3680    let kind = update.get("sessionUpdate").and_then(Value::as_str);
3681    if kind == Some("config_option_update")
3682        && let Some((config_id, models, current_model)) = parse_model_config(update)
3683    {
3684        return Ok(Some(AgentEvent::ModelsReplaced {
3685            slot,
3686            config_id,
3687            models,
3688            current_model,
3689        }));
3690    }
3691    if kind == Some("request_permission") {
3692        let request_id = update
3693            .get("toolCall")
3694            .and_then(|tool| tool.get("toolCallId"))
3695            .and_then(Value::as_str)
3696            .unwrap_or("permission");
3697        return Ok(parse_permission_event(
3698            slot,
3699            update,
3700            request_id,
3701            update.get("options"),
3702        ));
3703    }
3704    if kind == Some("available_commands_update") {
3705        let commands = update
3706            .get("availableCommands")
3707            .and_then(Value::as_array)
3708            .map(|commands| {
3709                commands
3710                    .iter()
3711                    .filter_map(|command| {
3712                        let name = command.get("name").and_then(Value::as_str)?.trim();
3713                        (!name.is_empty()).then(|| AgentCommand {
3714                            name: name.to_owned(),
3715                        })
3716                    })
3717                    .collect::<Vec<_>>()
3718            })
3719            .unwrap_or_default();
3720        return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3721    }
3722    if kind == Some("current_mode_update") {
3723        if let Some(mode) = update
3724            .get("currentModeId")
3725            .and_then(Value::as_str)
3726            .filter(|mode| !mode.trim().is_empty())
3727        {
3728            return Ok(Some(AgentEvent::ModeUpdated {
3729                slot,
3730                current_mode: mode.to_owned(),
3731            }));
3732        }
3733        return Ok(None);
3734    }
3735    if kind == Some("usage_update") {
3736        let Some(used) = update.get("used").and_then(Value::as_u64) else {
3737            return Ok(None);
3738        };
3739        let Some(size) = update.get("size").and_then(Value::as_u64) else {
3740            return Ok(None);
3741        };
3742        return Ok(Some(AgentEvent::UsageUpdated {
3743            slot,
3744            usage: UsageUpdate { used, size },
3745        }));
3746    }
3747    if let Some(terminal) = parse_terminal_event(update, kind) {
3748        return Ok(Some(AgentEvent::Terminal {
3749            slot,
3750            event: terminal,
3751        }));
3752    }
3753    let text = update
3754        .get("content")
3755        .and_then(|content| content.get("text"))
3756        .and_then(Value::as_str)
3757        .map(str::to_owned);
3758    if kind == Some("user_message_chunk") {
3759        return Ok(text
3760            .filter(|text| !text.is_empty())
3761            .map(|text| AgentEvent::UserText { slot, text }));
3762    }
3763    if kind == Some("agent_message_chunk")
3764        && let Some(mode) = text
3765            .as_deref()
3766            .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3767            .map(str::trim)
3768            .filter(|mode| !mode.is_empty())
3769    {
3770        // Gemini's native ACP bridge historically encoded a mode change as a
3771        // control marker in the message stream. It is state, not transcript
3772        // content; expose it as a normalized catalog replacement instead of
3773        // leaking the marker into the conversation.
3774        return Ok(Some(AgentEvent::ModesReplaced {
3775            slot,
3776            modes: vec![Mode {
3777                id: mode.to_owned(),
3778                label: mode.to_owned(),
3779            }],
3780            current_mode: Some(mode.to_owned()),
3781        }));
3782    }
3783    match (kind, text) {
3784        (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3785            Ok(Some(AgentEvent::Text { slot, text }))
3786        }
3787        (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3788            Ok(Some(AgentEvent::Thought { slot, text }))
3789        }
3790        (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3791            Ok(normalize_acp_tool(update, tools).map(|update| AgentEvent::Tool { slot, update }))
3792        }
3793        _ => Ok(None),
3794    }
3795}
3796
3797/// ACP updates are patches: omitted/invalid fields retain their last valid
3798/// values, whereas an explicit empty content array clears previous output.
3799fn normalize_acp_tool(
3800    value: &Value,
3801    tools: &mut BTreeMap<String, ToolUpdate>,
3802) -> Option<ToolUpdate> {
3803    let id = value.get("toolCallId")?.as_str()?;
3804    if id.trim().is_empty() {
3805        return None;
3806    }
3807    if value.get("sessionUpdate").and_then(Value::as_str) == Some("tool_call") {
3808        tools.remove(id);
3809    }
3810    let tool = tools.entry(id.to_owned()).or_insert_with(|| ToolUpdate {
3811        id: id.to_owned(),
3812        title: "Tool call".into(),
3813        status: ToolStatus::Pending,
3814        detail: None,
3815    });
3816    if let Some(title) = value.get("title").and_then(Value::as_str) {
3817        tool.title = title.to_owned();
3818    }
3819    if let Some(status) =
3820        value
3821            .get("status")
3822            .and_then(Value::as_str)
3823            .and_then(|status| match status {
3824                "pending" => Some(ToolStatus::Pending),
3825                "in_progress" => Some(ToolStatus::Running),
3826                "completed" => Some(ToolStatus::Completed),
3827                "failed" => Some(ToolStatus::Failed),
3828                _ => None,
3829            })
3830    {
3831        tool.status = status;
3832    }
3833    if let Some(content) = value.get("content").and_then(Value::as_array) {
3834        let text = content
3835            .iter()
3836            .filter_map(|entry| match entry.get("type").and_then(Value::as_str) {
3837                Some("content") => entry.get("content")?.get("text")?.as_str(),
3838                Some("diff") => entry.get("newText")?.as_str(),
3839                _ => None,
3840            })
3841            .collect::<Vec<_>>()
3842            .join("\n");
3843        tool.detail = (!text.is_empty()).then_some(text);
3844    } else if let Some(output) = value.get("rawOutput").filter(|output| !output.is_null()) {
3845        tool.detail = Some(
3846            output
3847                .as_str()
3848                .map(str::to_owned)
3849                .unwrap_or_else(|| output.to_string()),
3850        );
3851    }
3852    Some(tool.clone())
3853}
3854
3855fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3856    // Newer agents (OpenCode 1.4+) also advertise `models` with
3857    // `currentModelId`/`availableModels`; prefer the legacy `configOptions`
3858    // select when present so `session/set_config_option` keeps working.
3859    if let Some(object) = value.get("models").and_then(Value::as_object)
3860        && value.get("configOptions").is_none()
3861    {
3862        let available = object.get("availableModels")?.as_array()?;
3863        let modes = available
3864            .iter()
3865            .filter_map(|option| {
3866                let id = option.get("modelId")?.as_str()?.to_owned();
3867                let label = option
3868                    .get("name")
3869                    .or_else(|| option.get("label"))
3870                    .and_then(Value::as_str)
3871                    .unwrap_or(&id)
3872                    .to_owned();
3873                Some(Mode { id, label })
3874            })
3875            .collect::<Vec<_>>();
3876        return (!modes.is_empty()).then(|| {
3877            let current = object
3878                .get("currentModelId")
3879                .and_then(Value::as_str)
3880                .map(str::to_owned);
3881            ("model".to_owned(), modes, current)
3882        });
3883    }
3884    let config = value
3885        .get("configOptions")?
3886        .as_array()?
3887        .iter()
3888        .find(|option| {
3889            option.get("category").and_then(Value::as_str) == Some("model")
3890                && matches!(
3891                    option.get("type").and_then(Value::as_str),
3892                    Some("select" | "enum")
3893                )
3894        })?;
3895    let config_id = config.get("id")?.as_str()?.to_owned();
3896    let models = config
3897        .get("options")?
3898        .as_array()?
3899        .iter()
3900        .filter_map(|option| {
3901            let id = option.get("value")?.as_str()?.to_owned();
3902            let label = option
3903                .get("name")
3904                .or_else(|| option.get("label"))
3905                .and_then(Value::as_str)
3906                .unwrap_or(&id)
3907                .to_owned();
3908            Some(Mode { id, label })
3909        })
3910        .collect::<Vec<_>>();
3911    (!models.is_empty()).then(|| {
3912        let current = config
3913            .get("currentValue")
3914            .and_then(Value::as_str)
3915            .map(str::to_owned);
3916        (config_id, models, current)
3917    })
3918}
3919
3920/// Normalize terminal lifecycle updates emitted by ACP-compatible bridges and
3921/// native stream adapters. Protocols have used both snake_case update names
3922/// and a nested `terminal` object, so accept either without leaking that
3923/// shape beyond the adapter boundary.
3924fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3925    let nested = value.get("terminal").unwrap_or(value);
3926    let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3927    let id = nested
3928        .get("terminalId")
3929        .or_else(|| nested.get("terminal_id"))
3930        .or_else(|| nested.get("id"))
3931        .and_then(Value::as_str)
3932        .unwrap_or("terminal")
3933        .to_owned();
3934    match kind {
3935        "terminal_created" | "terminal_create" | "terminal_started" => {
3936            let command = nested
3937                .get("command")
3938                .and_then(Value::as_str)
3939                .unwrap_or("")
3940                .to_owned();
3941            Some(TerminalEvent::Created { id, command })
3942        }
3943        "terminal_output" | "terminal_output_chunk" => {
3944            let text = nested
3945                .get("output")
3946                .or_else(|| nested.get("text"))
3947                .and_then(Value::as_str)
3948                .unwrap_or("")
3949                .to_owned();
3950            Some(TerminalEvent::Output { id, text })
3951        }
3952        "terminal_exited" | "terminal_exit" => {
3953            let code = nested
3954                .get("exitCode")
3955                .or_else(|| nested.get("exit_code"))
3956                .or_else(|| nested.get("code"))
3957                .and_then(Value::as_i64)
3958                .unwrap_or(0) as i32;
3959            Some(TerminalEvent::Exited { id, code })
3960        }
3961        "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3962        _ => None,
3963    }
3964}
3965
3966fn parse_permission_event(
3967    slot: RosterSlot,
3968    value: &Value,
3969    request_id: &str,
3970    options: Option<&Value>,
3971) -> Option<AgentEvent> {
3972    let tool = value.get("toolCall").unwrap_or(value);
3973    let title = tool
3974        .get("title")
3975        .and_then(Value::as_str)
3976        .unwrap_or("Agent requests permission")
3977        .to_owned();
3978    let (options, option_ids, option_kinds): (Vec<String>, Vec<String>, Vec<String>) = options
3979        .and_then(Value::as_array)
3980        .map(|options| {
3981            let parsed = options
3982                .iter()
3983                .filter_map(|option| {
3984                    let label = option
3985                        .get("name")
3986                        .or_else(|| option.get("optionId"))
3987                        .and_then(Value::as_str)?
3988                        .to_owned();
3989                    let option_id = option
3990                        .get("optionId")
3991                        .or_else(|| option.get("id"))
3992                        .and_then(Value::as_str)
3993                        .map(str::to_owned)
3994                        .unwrap_or_else(|| label.clone());
3995                    let kind = option
3996                        .get("kind")
3997                        .and_then(Value::as_str)
3998                        .unwrap_or_default()
3999                        .to_owned();
4000                    Some((label, option_id, kind))
4001                })
4002                .collect::<Vec<_>>();
4003            let mut labels = Vec::with_capacity(parsed.len());
4004            let mut ids = Vec::with_capacity(parsed.len());
4005            let mut kinds = Vec::with_capacity(parsed.len());
4006            for (label, id, kind) in parsed {
4007                labels.push(label);
4008                ids.push(id);
4009                kinds.push(kind);
4010            }
4011            (labels, ids, kinds)
4012        })
4013        .unwrap_or_default();
4014    if options.is_empty() {
4015        return None;
4016    }
4017    Some(AgentEvent::Permission {
4018        slot,
4019        request: PermissionRequest {
4020            id: request_id.to_owned(),
4021            title,
4022            options,
4023            option_ids,
4024            option_kinds,
4025        },
4026    })
4027}
4028
4029fn rpc_id_to_string(value: &Value) -> String {
4030    value
4031        .as_str()
4032        .map(str::to_owned)
4033        .or_else(|| value.as_u64().map(|id| id.to_string()))
4034        .unwrap_or_else(|| value.to_string())
4035}
4036
4037#[cfg(test)]
4038mod tests {
4039    use super::{
4040        AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
4041        RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
4042        parse_model_config, prompt_content_blocks, read_bounded_line,
4043    };
4044    #[cfg(target_os = "linux")]
4045    use super::{isolate_process_group, terminate_child};
4046    use crate::TerminalEvent;
4047    use crate::{
4048        AdapterError, AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
4049        persistence::SessionMetadataStore,
4050        relay::{CollaborationStrategy, DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
4051    };
4052    use async_trait::async_trait;
4053    use serde_json::Value;
4054    use std::collections::VecDeque;
4055    use std::sync::{
4056        Arc, Mutex,
4057        atomic::{AtomicUsize, Ordering},
4058    };
4059
4060    fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
4061        let nonce = std::time::SystemTime::now()
4062            .duration_since(std::time::UNIX_EPOCH)
4063            .expect("clock")
4064            .as_nanos();
4065        std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
4066    }
4067
4068    #[test]
4069    fn malformed_file_writes_preserve_existing_content() {
4070        let root = unique_test_path("codeswarm-write-validation", "dir");
4071        std::fs::create_dir_all(&root).unwrap();
4072        let file = root.join("keep.txt");
4073        std::fs::write(&file, "valuable content").unwrap();
4074        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
4075        for content in [
4076            Value::Null,
4077            serde_json::json!(false),
4078            serde_json::json!(42),
4079            serde_json::json!([]),
4080        ] {
4081            assert!(
4082                adapter
4083                    .write_workspace_text(
4084                        &serde_json::json!({"path":"keep.txt", "content": content})
4085                    )
4086                    .is_err()
4087            );
4088            assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4089        }
4090        assert!(
4091            adapter
4092                .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
4093                .is_err()
4094        );
4095        assert!(
4096            adapter
4097                .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
4098                .is_err()
4099        );
4100        assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
4101        adapter
4102            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
4103            .unwrap();
4104        assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
4105        adapter
4106            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
4107            .unwrap();
4108        assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
4109        std::fs::remove_dir_all(root).unwrap();
4110    }
4111
4112    #[tokio::test]
4113    async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
4114        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
4115        let mut adapter = AcpAdapter::new(
4116            0,
4117            std::env::current_dir().unwrap(),
4118            "sh",
4119            vec!["-c".into(), script.into()],
4120        );
4121        adapter.start().await.unwrap();
4122        let error = adapter
4123            .request_with_timeout(
4124                "session/set_mode",
4125                serde_json::json!({}),
4126                std::time::Duration::from_millis(10),
4127            )
4128            .await
4129            .unwrap_err();
4130        assert!(error.to_string().contains("session/set_mode timed out"));
4131        adapter.stop().await.unwrap();
4132        assert!(adapter.child.is_none());
4133    }
4134
4135    #[tokio::test]
4136    async fn goals_reach_every_roster_slot_without_native_goal_support() {
4137        use crate::goal::GoalCommand;
4138        let hosts = (0..3)
4139            .map(|slot| {
4140                AdapterHost::new(
4141                    Box::new(ScriptedAdapter::new(
4142                        slot,
4143                        AgentCapabilities::default(),
4144                        [
4145                            AgentEvent::TurnComplete { slot },
4146                            AgentEvent::TurnComplete { slot },
4147                        ],
4148                    )),
4149                    None,
4150                )
4151            })
4152            .collect();
4153        let mut relay = RelayHost::new(hosts, 10).unwrap();
4154        relay.start().await.unwrap();
4155        let task = relay
4156            .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
4157            .unwrap()
4158            .unwrap();
4159        relay.relay_mut().enqueue_human(task, Some(0));
4160        for slot in 0..3 {
4161            relay.run_turn("", 0).await.unwrap();
4162            let (actual, prompt) = relay.dispatches().last().unwrap();
4163            assert_eq!(*actual, slot);
4164            assert!(prompt.contains("Active shared goal: Ship the settings screen"));
4165        }
4166        let snapshot = relay.session_metadata();
4167        let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
4168        assert!(restored.is_some());
4169        relay.restore_goal(restored);
4170        relay.reload(0).await.unwrap();
4171        relay.run_turn("", 0).await.unwrap();
4172        assert!(
4173            relay
4174                .dispatches()
4175                .last()
4176                .unwrap()
4177                .1
4178                .contains("Active shared goal: Ship the settings screen")
4179        );
4180        relay.apply_goal(GoalCommand::Done).unwrap();
4181        relay.run_turn("", 0).await.unwrap();
4182        assert!(
4183            relay
4184                .dispatches()
4185                .last()
4186                .unwrap()
4187                .1
4188                .contains("No active shared goal")
4189        );
4190        relay.apply_goal(GoalCommand::Clear).unwrap();
4191        assert!(relay.session_metadata().get("goal").unwrap().is_null());
4192    }
4193
4194    #[tokio::test]
4195    async fn replacement_agent_receives_task_after_public_journal_pruning() {
4196        let hosts = (0..2)
4197            .map(|slot| {
4198                AdapterHost::new(
4199                    Box::new(ScriptedAdapter::new(
4200                        slot,
4201                        AgentCapabilities::default(),
4202                        [
4203                            AgentEvent::Text {
4204                                slot,
4205                                text: "progress".into(),
4206                            },
4207                            AgentEvent::TurnComplete { slot },
4208                            AgentEvent::Text {
4209                                slot,
4210                                text: "more progress".into(),
4211                            },
4212                            AgentEvent::TurnComplete { slot },
4213                        ],
4214                    )),
4215                    None,
4216                )
4217            })
4218            .collect();
4219        let mut relay = RelayHost::new(hosts, 10).unwrap();
4220        relay.start().await.unwrap();
4221        relay
4222            .relay_mut()
4223            .enqueue_human("Fix the login bug", Some(0));
4224        relay.run_turn("", 0).await.unwrap();
4225        relay.run_turn("", 0).await.unwrap();
4226        relay.run_turn("", 0).await.unwrap();
4227        relay.reload(1).await.unwrap();
4228        assert!(
4229            !relay
4230                .relay_mut()
4231                .unseen_context(1)
4232                .contains("Fix the login bug")
4233        );
4234        relay.run_turn("", 0).await.unwrap();
4235        assert!(
4236            relay
4237                .dispatches()
4238                .last()
4239                .unwrap()
4240                .1
4241                .contains("Shared task:\nFix the login bug")
4242        );
4243    }
4244
4245    #[cfg(target_os = "linux")]
4246    #[tokio::test]
4247    async fn termination_kills_only_the_verified_isolated_child_group() {
4248        use nix::unistd::{Pid, getpgid, getpgrp};
4249        use tokio::io::{AsyncBufReadExt, BufReader};
4250
4251        let own_group = getpgrp();
4252        let mut command = tokio::process::Command::new("sh");
4253        isolate_process_group(&mut command);
4254        command
4255            .arg("-c")
4256            .arg("sleep 60 & echo $!; wait")
4257            .stdout(std::process::Stdio::piped());
4258        let mut child = command.spawn().expect("spawn isolated shell");
4259        let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
4260        assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
4261        assert_ne!(leader, own_group);
4262
4263        let stdout = child.stdout.take().expect("child stdout");
4264        let mut lines = BufReader::new(stdout).lines();
4265        let descendant = lines
4266            .next_line()
4267            .await
4268            .expect("read descendant pid")
4269            .expect("descendant pid")
4270            .parse::<i32>()
4271            .expect("numeric descendant pid");
4272        let descendant = Pid::from_raw(descendant);
4273        assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
4274
4275        terminate_child(&mut child).await.expect("terminate group");
4276        for _ in 0..100 {
4277            if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
4278                return;
4279            }
4280            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
4281        }
4282        panic!("descendant {descendant} survived isolated group termination");
4283    }
4284
4285    #[test]
4286    fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
4287        assert_eq!(
4288            parse_command_line(r#"agent --name "local bridge" --flag 'two words'"#),
4289            Ok((
4290                "agent".into(),
4291                vec![
4292                    "--name".into(),
4293                    "local bridge".into(),
4294                    "--flag".into(),
4295                    "two words".into(),
4296                ]
4297            ),)
4298        );
4299        assert_eq!(
4300            parse_command_line(r#"agent "" escaped\ argument"#),
4301            Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
4302        );
4303    }
4304
4305    #[test]
4306    fn acp_prompt_expands_safe_at_path_resources() {
4307        let root = unique_test_path("codeswarm-prompt-resource", "dir");
4308        std::fs::create_dir_all(&root).expect("workspace");
4309        std::fs::write(root.join("note.md"), "resource text").expect("resource");
4310        let blocks = prompt_content_blocks(&root, "inspect @note.md");
4311        assert_eq!(blocks[0]["type"], "text");
4312        assert_eq!(blocks[0]["text"], "inspect @note.md");
4313        assert_eq!(blocks[1]["type"], "resource");
4314        assert_eq!(blocks[1]["resource"]["text"], "resource text");
4315        assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
4316        std::fs::remove_dir_all(root).expect("cleanup workspace");
4317    }
4318
4319    #[tokio::test]
4320    async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
4321        let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
4322        bytes.push(b'\n');
4323        let mut reader = tokio::io::BufReader::new(bytes.as_slice());
4324        assert!(matches!(
4325            read_bounded_line(&mut reader).await,
4326            Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
4327        ));
4328    }
4329
4330    #[test]
4331    fn rejects_malformed_configured_commands_before_spawn() {
4332        assert_eq!(
4333            parse_command_line("agent 'unfinished"),
4334            Err(super::CommandParseError::UnterminatedQuote)
4335        );
4336        assert_eq!(
4337            parse_command_line("agent\\"),
4338            Err(super::CommandParseError::TrailingEscape)
4339        );
4340        assert_eq!(
4341            parse_command_line("   \t"),
4342            Err(super::CommandParseError::Empty)
4343        );
4344    }
4345
4346    #[derive(Debug)]
4347    struct PendingAdapter {
4348        slot: usize,
4349        hang_on_cancel: bool,
4350    }
4351
4352    #[derive(Debug)]
4353    struct ConcurrentStartAdapter {
4354        slot: usize,
4355        barrier: Arc<tokio::sync::Barrier>,
4356    }
4357
4358    #[derive(Debug)]
4359    struct ReloadProbeAdapter {
4360        slot: usize,
4361        crashed: bool,
4362        reloaded: bool,
4363        events: VecDeque<AgentEvent>,
4364        prompts: Arc<Mutex<Vec<String>>>,
4365    }
4366
4367    #[async_trait]
4368    impl AgentAdapter for ReloadProbeAdapter {
4369        fn slot(&self) -> usize {
4370            self.slot
4371        }
4372
4373        fn display_name(&self) -> String {
4374            "Reload probe".into()
4375        }
4376
4377        fn protocol(&self) -> &'static str {
4378            "native"
4379        }
4380
4381        fn capabilities(&self) -> AgentCapabilities {
4382            AgentCapabilities {
4383                supports_modes: true,
4384                ..AgentCapabilities::default()
4385            }
4386        }
4387
4388        fn needs_restart(&self) -> bool {
4389            self.crashed && !self.reloaded
4390        }
4391
4392        async fn start(&mut self) -> super::AdapterResult<()> {
4393            self.events.push_back(AgentEvent::ModesReplaced {
4394                slot: self.slot,
4395                modes: vec![
4396                    Mode {
4397                        id: "codeswarm:mode:full-access".into(),
4398                        label: "Auto pilot".into(),
4399                    },
4400                    Mode {
4401                        id: "codeswarm:mode:plan".into(),
4402                        label: "Plan".into(),
4403                    },
4404                ],
4405                current_mode: Some("codeswarm:mode:full-access".into()),
4406            });
4407            self.events.push_back(AgentEvent::Ready {
4408                slot: self.slot,
4409                capabilities: self.capabilities(),
4410            });
4411            Ok(())
4412        }
4413
4414        async fn send_prompt(&mut self, prompt: String) -> super::AdapterResult<()> {
4415            self.prompts.lock().expect("prompts").push(prompt);
4416            if !self.crashed {
4417                self.crashed = true;
4418                self.events.push_back(AgentEvent::Failed {
4419                    slot: self.slot,
4420                    started: true,
4421                    detail: "probe crashed".into(),
4422                });
4423            } else {
4424                self.events.push_back(AgentEvent::Text {
4425                    slot: self.slot,
4426                    text: "recovered".into(),
4427                });
4428                self.events
4429                    .push_back(AgentEvent::TurnComplete { slot: self.slot });
4430            }
4431            Ok(())
4432        }
4433
4434        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4435            Ok(false)
4436        }
4437
4438        async fn answer_permission(
4439            &mut self,
4440            _request_id: String,
4441            _answer: PermissionAnswer,
4442        ) -> super::AdapterResult<()> {
4443            Err(super::AdapterError::Unsupported("permission answer"))
4444        }
4445
4446        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4447            Ok(())
4448        }
4449
4450        async fn reload(&mut self) -> super::AdapterResult<()> {
4451            self.reloaded = true;
4452            self.start().await
4453        }
4454
4455        async fn stop(&mut self) -> super::AdapterResult<()> {
4456            Ok(())
4457        }
4458
4459        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4460            self.events.pop_front().map(Ok)
4461        }
4462    }
4463
4464    #[async_trait]
4465    impl AgentAdapter for ConcurrentStartAdapter {
4466        fn slot(&self) -> usize {
4467            self.slot
4468        }
4469
4470        fn capabilities(&self) -> AgentCapabilities {
4471            AgentCapabilities::default()
4472        }
4473
4474        async fn start(&mut self) -> super::AdapterResult<()> {
4475            self.barrier.wait().await;
4476            Ok(())
4477        }
4478
4479        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4480            Ok(())
4481        }
4482
4483        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4484            Ok(true)
4485        }
4486
4487        async fn answer_permission(
4488            &mut self,
4489            _request_id: String,
4490            _answer: PermissionAnswer,
4491        ) -> super::AdapterResult<()> {
4492            Ok(())
4493        }
4494
4495        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4496            Ok(())
4497        }
4498
4499        async fn reload(&mut self) -> super::AdapterResult<()> {
4500            Ok(())
4501        }
4502
4503        async fn stop(&mut self) -> super::AdapterResult<()> {
4504            Ok(())
4505        }
4506
4507        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4508            std::future::pending().await
4509        }
4510    }
4511
4512    #[derive(Debug)]
4513    struct PermissionBlockingAdapter {
4514        slot: usize,
4515        phase: u8,
4516    }
4517
4518    #[async_trait]
4519    impl AgentAdapter for PermissionBlockingAdapter {
4520        fn slot(&self) -> usize {
4521            self.slot
4522        }
4523
4524        fn capabilities(&self) -> AgentCapabilities {
4525            AgentCapabilities {
4526                supports_permissions: true,
4527                ..AgentCapabilities::default()
4528            }
4529        }
4530
4531        async fn start(&mut self) -> super::AdapterResult<()> {
4532            Ok(())
4533        }
4534
4535        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4536            Ok(())
4537        }
4538
4539        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4540            Ok(true)
4541        }
4542
4543        async fn answer_permission(
4544            &mut self,
4545            request_id: String,
4546            answer: PermissionAnswer,
4547        ) -> super::AdapterResult<()> {
4548            if self.phase != 1 || request_id != "permission-1" {
4549                return Err(super::AdapterError::Protocol(
4550                    "unexpected permission response".into(),
4551                ));
4552            }
4553            assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4554            self.phase = 2;
4555            Ok(())
4556        }
4557
4558        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4559            Ok(())
4560        }
4561
4562        async fn reload(&mut self) -> super::AdapterResult<()> {
4563            Ok(())
4564        }
4565
4566        async fn stop(&mut self) -> super::AdapterResult<()> {
4567            Ok(())
4568        }
4569
4570        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4571            match self.phase {
4572                0 => {
4573                    self.phase = 1;
4574                    Some(Ok(AgentEvent::Permission {
4575                        slot: self.slot,
4576                        request: crate::PermissionRequest {
4577                            id: "permission-1".into(),
4578                            title: "Allow?".into(),
4579                            options: vec!["Allow".into()],
4580                            option_ids: vec!["allow".into()],
4581                            option_kinds: Vec::new(),
4582                        },
4583                    }))
4584                }
4585                1 => std::future::pending().await,
4586                _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4587            }
4588        }
4589    }
4590
4591    #[async_trait]
4592    impl AgentAdapter for PendingAdapter {
4593        fn slot(&self) -> usize {
4594            self.slot
4595        }
4596
4597        fn capabilities(&self) -> AgentCapabilities {
4598            AgentCapabilities {
4599                supports_cancel: true,
4600                ..AgentCapabilities::default()
4601            }
4602        }
4603
4604        async fn start(&mut self) -> super::AdapterResult<()> {
4605            Ok(())
4606        }
4607
4608        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4609            Ok(())
4610        }
4611
4612        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4613            if self.hang_on_cancel {
4614                return std::future::pending().await;
4615            }
4616            Ok(true)
4617        }
4618
4619        async fn answer_permission(
4620            &mut self,
4621            _request_id: String,
4622            _answer: PermissionAnswer,
4623        ) -> super::AdapterResult<()> {
4624            Ok(())
4625        }
4626
4627        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4628            Ok(())
4629        }
4630
4631        async fn reload(&mut self) -> super::AdapterResult<()> {
4632            Ok(())
4633        }
4634
4635        async fn stop(&mut self) -> super::AdapterResult<()> {
4636            Ok(())
4637        }
4638
4639        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4640            std::future::pending().await
4641        }
4642    }
4643
4644    #[derive(Debug)]
4645    struct StopTrackingAdapter {
4646        slot: usize,
4647        stopped: Arc<AtomicUsize>,
4648        fail_stop: bool,
4649    }
4650
4651    #[derive(Debug)]
4652    struct ModeOrderAdapter {
4653        slot: usize,
4654        log: Arc<Mutex<Vec<String>>>,
4655        phase: u8,
4656    }
4657
4658    #[derive(Debug)]
4659    struct StartupAcpAdapter {
4660        slot: usize,
4661        events: std::collections::VecDeque<AgentEvent>,
4662    }
4663
4664    impl StartupAcpAdapter {
4665        fn new(slot: usize) -> Self {
4666            Self {
4667                slot,
4668                events: [
4669                    AgentEvent::ModesReplaced {
4670                        slot,
4671                        modes: vec![Mode {
4672                            id: "full-access".into(),
4673                            label: "Auto pilot".into(),
4674                        }],
4675                        current_mode: Some("full-access".into()),
4676                    },
4677                    AgentEvent::Ready {
4678                        slot,
4679                        capabilities: AgentCapabilities {
4680                            supports_modes: true,
4681                            ..AgentCapabilities::default()
4682                        },
4683                    },
4684                ]
4685                .into(),
4686            }
4687        }
4688    }
4689
4690    #[async_trait]
4691    impl AgentAdapter for StartupAcpAdapter {
4692        fn slot(&self) -> usize {
4693            self.slot
4694        }
4695
4696        fn protocol(&self) -> &'static str {
4697            "acp"
4698        }
4699
4700        fn capabilities(&self) -> AgentCapabilities {
4701            AgentCapabilities {
4702                supports_modes: true,
4703                ..AgentCapabilities::default()
4704            }
4705        }
4706
4707        async fn start(&mut self) -> super::AdapterResult<()> {
4708            Ok(())
4709        }
4710
4711        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4712            Ok(())
4713        }
4714
4715        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4716            Ok(true)
4717        }
4718
4719        async fn answer_permission(
4720            &mut self,
4721            _request_id: String,
4722            _answer: PermissionAnswer,
4723        ) -> super::AdapterResult<()> {
4724            Ok(())
4725        }
4726
4727        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4728            Ok(())
4729        }
4730
4731        async fn reload(&mut self) -> super::AdapterResult<()> {
4732            self.events = Self::new(self.slot).events;
4733            Ok(())
4734        }
4735
4736        async fn stop(&mut self) -> super::AdapterResult<()> {
4737            Ok(())
4738        }
4739
4740        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4741            self.events.pop_front().map(Ok)
4742        }
4743    }
4744
4745    #[async_trait]
4746    impl AgentAdapter for ModeOrderAdapter {
4747        fn slot(&self) -> usize {
4748            self.slot
4749        }
4750
4751        fn capabilities(&self) -> AgentCapabilities {
4752            AgentCapabilities {
4753                supports_modes: true,
4754                ..AgentCapabilities::default()
4755            }
4756        }
4757
4758        async fn start(&mut self) -> super::AdapterResult<()> {
4759            self.log.lock().expect("log").push("start".into());
4760            Ok(())
4761        }
4762
4763        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4764            self.log.lock().expect("log").push("prompt".into());
4765            Ok(())
4766        }
4767
4768        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4769            Ok(true)
4770        }
4771
4772        async fn answer_permission(
4773            &mut self,
4774            _request_id: String,
4775            _answer: PermissionAnswer,
4776        ) -> super::AdapterResult<()> {
4777            Ok(())
4778        }
4779
4780        async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4781            self.log.lock().expect("log").push(format!("mode:{mode}"));
4782            Ok(())
4783        }
4784
4785        async fn reload(&mut self) -> super::AdapterResult<()> {
4786            self.log.lock().expect("log").push("reload".into());
4787            self.phase = 0;
4788            Ok(())
4789        }
4790
4791        async fn stop(&mut self) -> super::AdapterResult<()> {
4792            Ok(())
4793        }
4794
4795        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4796            match self.phase {
4797                0 => {
4798                    self.phase = 1;
4799                    Some(Ok(AgentEvent::ModesReplaced {
4800                        slot: self.slot,
4801                        modes: vec![Mode {
4802                            id: "yolo".into(),
4803                            label: "YOLO".into(),
4804                        }],
4805                        current_mode: None,
4806                    }))
4807                }
4808                1 => {
4809                    self.phase = 2;
4810                    Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4811                }
4812                _ => std::future::pending().await,
4813            }
4814        }
4815    }
4816
4817    #[async_trait]
4818    impl AgentAdapter for StopTrackingAdapter {
4819        fn slot(&self) -> usize {
4820            self.slot
4821        }
4822
4823        fn capabilities(&self) -> AgentCapabilities {
4824            AgentCapabilities::default()
4825        }
4826
4827        async fn start(&mut self) -> super::AdapterResult<()> {
4828            Ok(())
4829        }
4830
4831        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4832            Ok(())
4833        }
4834
4835        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4836            Ok(false)
4837        }
4838
4839        async fn answer_permission(
4840            &mut self,
4841            _request_id: String,
4842            _answer: PermissionAnswer,
4843        ) -> super::AdapterResult<()> {
4844            Ok(())
4845        }
4846
4847        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4848            Ok(())
4849        }
4850
4851        async fn reload(&mut self) -> super::AdapterResult<()> {
4852            Ok(())
4853        }
4854
4855        async fn stop(&mut self) -> super::AdapterResult<()> {
4856            self.stopped.fetch_add(1, Ordering::Relaxed);
4857            if self.fail_stop {
4858                Err(super::AdapterError::Transport("stop failed".into()))
4859            } else {
4860                Ok(())
4861            }
4862        }
4863
4864        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4865            None
4866        }
4867    }
4868
4869    /// Startup can fail after an adapter has allocated resources.  Keep a
4870    /// fixture that records whether the failing adapter itself receives the
4871    /// cleanup call, not just the already-started peers.
4872    #[derive(Debug)]
4873    struct FailingStartAdapter {
4874        slot: usize,
4875        stopped: Arc<AtomicUsize>,
4876    }
4877
4878    #[async_trait]
4879    impl AgentAdapter for FailingStartAdapter {
4880        fn slot(&self) -> usize {
4881            self.slot
4882        }
4883
4884        fn capabilities(&self) -> AgentCapabilities {
4885            AgentCapabilities::default()
4886        }
4887
4888        async fn start(&mut self) -> super::AdapterResult<()> {
4889            Err(super::AdapterError::Spawn("startup failed".into()))
4890        }
4891
4892        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4893            Ok(())
4894        }
4895
4896        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4897            Ok(false)
4898        }
4899
4900        async fn answer_permission(
4901            &mut self,
4902            _request_id: String,
4903            _answer: PermissionAnswer,
4904        ) -> super::AdapterResult<()> {
4905            Ok(())
4906        }
4907
4908        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4909            Ok(())
4910        }
4911
4912        async fn reload(&mut self) -> super::AdapterResult<()> {
4913            Ok(())
4914        }
4915
4916        async fn stop(&mut self) -> super::AdapterResult<()> {
4917            self.stopped.fetch_add(1, Ordering::Relaxed);
4918            Ok(())
4919        }
4920
4921        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4922            None
4923        }
4924    }
4925
4926    #[tokio::test]
4927    async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4928        let stopped = Arc::new(AtomicUsize::new(0));
4929        let relay = RelayHost::new(
4930            vec![
4931                AdapterHost::new(
4932                    Box::new(StopTrackingAdapter {
4933                        slot: 0,
4934                        stopped: Arc::clone(&stopped),
4935                        fail_stop: true,
4936                    }),
4937                    None,
4938                ),
4939                AdapterHost::new(
4940                    Box::new(StopTrackingAdapter {
4941                        slot: 1,
4942                        stopped: Arc::clone(&stopped),
4943                        fail_stop: false,
4944                    }),
4945                    None,
4946                ),
4947            ],
4948            4,
4949        )
4950        .expect("relay");
4951        let mut relay = relay;
4952
4953        let error = relay.stop().await.expect_err("first stop failure");
4954        assert!(error.to_string().contains("stop failed"));
4955        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4956    }
4957
4958    #[tokio::test]
4959    async fn relay_start_isolates_a_failed_adapter_and_keeps_healthy_peers() {
4960        let stopped = Arc::new(AtomicUsize::new(0));
4961        let mut relay = RelayHost::new(
4962            vec![
4963                AdapterHost::new(
4964                    Box::new(StopTrackingAdapter {
4965                        slot: 0,
4966                        stopped: Arc::clone(&stopped),
4967                        fail_stop: false,
4968                    }),
4969                    None,
4970                ),
4971                AdapterHost::new(
4972                    Box::new(FailingStartAdapter {
4973                        slot: 1,
4974                        stopped: Arc::clone(&stopped),
4975                    }),
4976                    None,
4977                ),
4978            ],
4979            4,
4980        )
4981        .expect("relay");
4982
4983        relay.start().await.expect("healthy peer remains available");
4984        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0]);
4985        assert_eq!(stopped.load(Ordering::Relaxed), 1);
4986        relay.stop().await.unwrap();
4987        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4988    }
4989
4990    #[test]
4991    fn parses_acp_text_without_ui_dependency() {
4992        let event = parse_acp_notification(
4993            2,
4994            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4995        )
4996        .expect("valid ACP")
4997        .expect("text event");
4998        assert_eq!(
4999            event,
5000            AgentEvent::Text {
5001                slot: 2,
5002                text: "hello".into(),
5003            }
5004        );
5005    }
5006
5007    #[test]
5008    fn parses_acp_state_notifications_at_the_adapter_boundary() {
5009        let commands = parse_acp_notification(
5010            3,
5011            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
5012        )
5013        .expect("valid ACP")
5014        .expect("commands event");
5015        assert_eq!(
5016            commands,
5017            AgentEvent::CommandsReplaced {
5018                slot: 3,
5019                commands: vec![crate::AgentCommand {
5020                    name: "review".into()
5021                }]
5022            }
5023        );
5024
5025        let mode = parse_acp_notification(
5026            3,
5027            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
5028        )
5029        .expect("valid ACP")
5030        .expect("mode event");
5031        assert_eq!(
5032            mode,
5033            AgentEvent::ModeUpdated {
5034                slot: 3,
5035                current_mode: "review".into()
5036            }
5037        );
5038
5039        let usage = parse_acp_notification(
5040            3,
5041            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
5042        )
5043        .expect("valid ACP")
5044        .expect("usage event");
5045        assert_eq!(
5046            usage,
5047            AgentEvent::UsageUpdated {
5048                slot: 3,
5049                usage: crate::UsageUpdate {
5050                    used: 4200,
5051                    size: 128000
5052                }
5053            }
5054        );
5055
5056        let models = parse_acp_notification(
5057            3,
5058            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"}]}]}}}"#,
5059        )
5060        .expect("valid ACP")
5061        .expect("models event");
5062        assert!(matches!(
5063            models,
5064            AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
5065                if models.len() == 2 && current_model.as_deref() == Some("smart")
5066        ));
5067        assert_eq!(
5068            parse_model_config(&serde_json::json!({
5069                "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
5070            })),
5071            None
5072        );
5073
5074        let user = parse_acp_notification(
5075            3,
5076            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
5077        )
5078        .expect("valid ACP")
5079        .expect("user event");
5080        assert_eq!(
5081            user,
5082            AgentEvent::UserText {
5083                slot: 3,
5084                text: "context".into()
5085            }
5086        );
5087    }
5088
5089    #[test]
5090    fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
5091        let event = parse_acp_notification(
5092            0,
5093            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
5094        )
5095        .expect("valid ACP")
5096        .expect("mode event");
5097        assert!(matches!(
5098            event,
5099            AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
5100                if mode == "yolo" && modes[0].id == "yolo"
5101        ));
5102    }
5103
5104    #[test]
5105    fn parses_native_agy_text_without_acp_bridge() {
5106        let event = parse_agy_line(
5107            1,
5108            r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
5109        )
5110        .expect("valid stream-json")
5111        .expect("text event");
5112        assert_eq!(
5113            event,
5114            AgentEvent::Text {
5115                slot: 1,
5116                text: "hello".into(),
5117            }
5118        );
5119    }
5120
5121    #[test]
5122    fn parses_tool_lifecycle_from_each_protocol() {
5123        let agy = parse_agy_line(
5124            1,
5125            r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
5126        )
5127        .expect("valid native tool")
5128        .expect("tool event");
5129        assert!(matches!(
5130            agy,
5131            AgentEvent::Tool {
5132                update: crate::ToolUpdate {
5133                    status: ToolStatus::Completed,
5134                    ..
5135                },
5136                ..
5137            }
5138        ));
5139
5140        let acp = parse_acp_notification(
5141            1,
5142            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
5143        )
5144        .expect("valid ACP tool")
5145        .expect("tool event");
5146        assert!(matches!(
5147            acp,
5148            AgentEvent::Tool {
5149                update: crate::ToolUpdate {
5150                    status: ToolStatus::Failed,
5151                    ..
5152                },
5153                ..
5154            }
5155        ));
5156    }
5157
5158    #[test]
5159    fn parses_terminal_lifecycle_from_acp_and_native_events() {
5160        let created = parse_acp_notification(
5161            0,
5162            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
5163        )
5164        .expect("valid ACP terminal")
5165        .expect("terminal event");
5166        assert_eq!(
5167            created,
5168            AgentEvent::Terminal {
5169                slot: 0,
5170                event: TerminalEvent::Created {
5171                    id: "term-1".into(),
5172                    command: "cargo test".into(),
5173                },
5174            }
5175        );
5176        let output = parse_agy_line(
5177            1,
5178            r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
5179        )
5180        .expect("valid native terminal")
5181        .expect("terminal event");
5182        assert_eq!(
5183            output,
5184            AgentEvent::Terminal {
5185                slot: 1,
5186                event: TerminalEvent::Output {
5187                    id: "term-1".into(),
5188                    text: "ok\n".into(),
5189                },
5190            }
5191        );
5192        let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
5193            .expect("valid native release")
5194            .expect("terminal event");
5195        assert!(matches!(
5196            released,
5197            AgentEvent::Terminal {
5198                event: TerminalEvent::Released { id },
5199                ..
5200            } if id == "term-1"
5201        ));
5202    }
5203
5204    #[test]
5205    fn parses_acp_permission_requests() {
5206        let event = parse_acp_notification(
5207            0,
5208            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"request_permission","toolCall":{"toolCallId":"t1","title":"Write file"},"options":[{"name":"Yes","optionId":"opaque-approval","kind":"allow_once"},{"name":"No","optionId":"reject","kind":"reject_once"}]}}}"#,
5209        )
5210        .expect("valid permission")
5211        .expect("permission event");
5212        assert!(matches!(
5213            event,
5214            AgentEvent::Permission { request, .. }
5215                if request.id == "t1"
5216                    && request.title == "Write file"
5217                    && request.options == ["Yes", "No"]
5218                    && request.option_ids == ["opaque-approval", "reject"]
5219                    && request.option_kinds == ["allow_once", "reject_once"]
5220        ));
5221    }
5222
5223    #[test]
5224    fn parses_acp_permission_request_as_json_rpc_request() {
5225        let event = parse_acp_notification(
5226            2,
5227            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
5228        )
5229        .expect("valid permission request")
5230        .expect("permission event");
5231        assert!(matches!(
5232            event,
5233            AgentEvent::Permission { request, .. }
5234                if request.id == "17"
5235                    && request.title == "Write file"
5236                    && request.options == ["allow-once", "reject"]
5237                    && request.option_ids == ["allow-once", "reject"]
5238                    && request.option_kinds == ["", ""]
5239        ));
5240    }
5241
5242    #[tokio::test]
5243    async fn native_adapter_explicitly_rejects_permission_answers() {
5244        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5245        assert_eq!(
5246            adapter
5247                .answer_permission(
5248                    "request".into(),
5249                    PermissionAnswer::Selected {
5250                        option_id: "allow".into()
5251                    },
5252                )
5253                .await,
5254            Err(super::AdapterError::Unsupported("permission answer"))
5255        );
5256    }
5257
5258    #[tokio::test]
5259    async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
5260        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
5261        adapter
5262            .set_mode("full-access".into())
5263            .await
5264            .expect("auto-pilot alias");
5265        assert!(matches!(
5266            adapter.next_event().await,
5267            Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
5268        ));
5269    }
5270
5271    #[tokio::test]
5272    async fn native_turns_receive_a_twenty_four_hour_timeout() {
5273        let script_path = unique_test_path("codeswarm-native-timeout", "sh");
5274        std::fs::write(
5275            &script_path,
5276            r#"#!/bin/sh
5277seen=0
5278while [ "$#" -gt 0 ]; do
5279    case "$1" in
5280        --print-timeout)
5281            shift
5282            [ "$1" = "1440m" ] || exit 2
5283            seen=$((seen + 1))
5284            ;;
5285    esac
5286    shift
5287done
5288[ "$seen" = 1 ] || exit 3
5289printf '%s\n' '{"event":"result","result":{"status":"SUCCESS","response":"timeout accepted"}}'
5290"#,
5291        )
5292        .unwrap();
5293        let mut adapter = AgyAdapter::with_session_id(
5294            0,
5295            std::env::current_dir().unwrap(),
5296            format!("sh {}", script_path.display()),
5297            "saved-session",
5298        );
5299        adapter.start().await.unwrap();
5300        adapter.next_event().await.unwrap().unwrap();
5301        adapter.next_event().await.unwrap().unwrap();
5302        for prompt in ["first task", "follow-up task"] {
5303            adapter.send_prompt(prompt.into()).await.unwrap();
5304            assert!(
5305                matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "timeout accepted")
5306            );
5307            assert!(matches!(
5308                adapter.next_event().await,
5309                Some(Ok(AgentEvent::TurnComplete { .. }))
5310            ));
5311        }
5312        adapter.stop().await.unwrap();
5313        std::fs::remove_file(script_path).unwrap();
5314    }
5315
5316    #[tokio::test]
5317    async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
5318        let script_path = unique_test_path("codeswarm-agy-session", "sh");
5319        std::fs::write(
5320            &script_path,
5321            "#!/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",
5322        )
5323        .expect("write native test script");
5324        let mut adapter = AgyAdapter::new(
5325            0,
5326            std::env::current_dir().expect("cwd"),
5327            format!("sh {}", script_path.display()),
5328        );
5329        adapter.start().await.expect("start native adapter");
5330        // Startup emits its mode catalog and readiness before a turn.
5331        assert!(adapter.next_event().await.is_some());
5332        assert!(adapter.next_event().await.is_some());
5333        adapter
5334            .send_prompt("first".into())
5335            .await
5336            .expect("first prompt");
5337        while !matches!(
5338            adapter.next_event().await,
5339            Some(Ok(AgentEvent::TurnComplete { .. }))
5340        ) {}
5341        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5342        adapter
5343            .send_prompt("follow up".into())
5344            .await
5345            .expect("follow-up prompt");
5346        while !matches!(
5347            adapter.next_event().await,
5348            Some(Ok(AgentEvent::TurnComplete { .. }))
5349        ) {}
5350        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
5351        adapter.stop().await.expect("stop native adapter");
5352        std::fs::remove_file(script_path).expect("cleanup native script");
5353    }
5354
5355    #[tokio::test]
5356    async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
5357        let script_path = unique_test_path("codeswarm-agy-failure", "sh");
5358        std::fs::write(
5359            &script_path,
5360            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
5361        )
5362        .expect("write native test script");
5363        let mut adapter = AgyAdapter::new(
5364            0,
5365            std::env::current_dir().expect("cwd"),
5366            format!("sh {}", script_path.display()),
5367        );
5368        adapter.start().await.expect("start native adapter");
5369        assert!(adapter.next_event().await.is_some());
5370        assert!(adapter.next_event().await.is_some());
5371        adapter.send_prompt("fail".into()).await.expect("prompt");
5372        assert!(matches!(
5373            adapter.next_event().await,
5374            Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
5375                if detail == "agent failed"
5376        ));
5377        adapter.stop().await.expect("stop native adapter");
5378        std::fs::remove_file(script_path).expect("cleanup native script");
5379    }
5380
5381    #[tokio::test]
5382    async fn native_crash_reaps_process_and_retries_on_next_prompt() {
5383        let script_path = unique_test_path("codeswarm-agy-retry", "sh");
5384        let marker_path = unique_test_path("codeswarm-agy-retry-marker", "txt");
5385        std::fs::write(
5386            &script_path,
5387            format!(
5388                "#!/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",
5389                marker_path.display(),
5390                marker_path.display(),
5391                marker_path.display(),
5392            ),
5393        )
5394        .expect("write retry script");
5395        let mut adapter = AgyAdapter::new(
5396            0,
5397            std::env::current_dir().expect("cwd"),
5398            format!("sh {}", script_path.display()),
5399        );
5400        adapter.start().await.expect("start native adapter");
5401        assert!(adapter.next_event().await.is_some());
5402        assert!(adapter.next_event().await.is_some());
5403        adapter
5404            .send_prompt("first".into())
5405            .await
5406            .expect("first prompt");
5407        assert!(matches!(
5408            adapter.next_event().await,
5409            Some(Ok(AgentEvent::Failed { detail, .. })) if detail == "first crash"
5410        ));
5411        adapter
5412            .send_prompt("retry".into())
5413            .await
5414            .expect("retry prompt starts a fresh process");
5415        assert!(matches!(
5416            adapter.next_event().await,
5417            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5418        ));
5419        assert!(matches!(
5420            adapter.next_event().await,
5421            Some(Ok(AgentEvent::TurnComplete { .. }))
5422        ));
5423        adapter.stop().await.expect("stop native adapter");
5424        std::fs::remove_file(script_path).expect("cleanup retry script");
5425        std::fs::remove_file(marker_path).expect("cleanup retry marker");
5426    }
5427
5428    #[tokio::test]
5429    async fn acp_adapter_initializes_session_and_completes_a_prompt() {
5430        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"}}'"#;
5431        let cwd = std::env::current_dir().expect("cwd");
5432        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5433        adapter.start().await.expect("initialize");
5434        assert!(matches!(
5435            adapter.next_event().await,
5436            Some(Ok(AgentEvent::ModesReplaced { .. }))
5437        ));
5438        assert!(matches!(
5439            adapter.next_event().await,
5440            Some(Ok(AgentEvent::Ready { .. }))
5441        ));
5442        adapter.send_prompt("hello".into()).await.expect("prompt");
5443        assert!(matches!(
5444            adapter.next_event().await,
5445            Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
5446        ));
5447        assert!(matches!(
5448            adapter.next_event().await,
5449            Some(Ok(AgentEvent::TurnComplete { .. }))
5450        ));
5451    }
5452
5453    #[tokio::test]
5454    async fn acp_empty_end_turn_is_a_failed_turn_not_an_echo() {
5455        // OpenCode echoes the prompt as `user_message_chunk` and ends with
5456        // zero usage when the model cannot run. That must fail loudly.
5457        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","method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"say hello"}}}}'; echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'"#;
5458        let cwd = std::env::current_dir().expect("cwd");
5459        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5460        adapter.start().await.expect("initialize");
5461        assert!(matches!(
5462            adapter.next_event().await,
5463            Some(Ok(AgentEvent::Ready { .. }))
5464        ));
5465        adapter
5466            .send_prompt("say hello".into())
5467            .await
5468            .expect("prompt");
5469        assert!(matches!(
5470            adapter.next_event().await,
5471            Some(Ok(AgentEvent::UserText { .. }))
5472        ));
5473        assert!(matches!(
5474            adapter.next_event().await,
5475            Some(Ok(AgentEvent::Failed {
5476                slot: 1,
5477                started: true,
5478                detail,
5479            })) if detail.contains("no agent output")
5480        ));
5481    }
5482
5483    #[test]
5484    fn parses_opencode_models_shape_without_config_options() {
5485        let session = serde_json::json!({
5486            "sessionId": "s",
5487            "models": {
5488                "currentModelId": "opencode-go/m",
5489                "availableModels": [
5490                    {"modelId": "opencode-go/m", "name": "M"},
5491                    {"modelId": "other/m2", "name": "M2"},
5492                ],
5493            },
5494            "modes": {"currentModeId": "build", "availableModes": []},
5495        });
5496        let (config_id, models, current) =
5497            super::parse_model_config(&session).expect("models shape");
5498        assert_eq!(config_id, "model");
5499        assert_eq!(models.len(), 2);
5500        assert_eq!(current.as_deref(), Some("opencode-go/m"));
5501    }
5502
5503    #[tokio::test]
5504    async fn acp_output_token_limit_is_reported_as_a_failed_turn() {
5505        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"}}'"#;
5506        let cwd = std::env::current_dir().expect("cwd");
5507        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5508        adapter.start().await.expect("initialize");
5509        assert!(matches!(
5510            adapter.next_event().await,
5511            Some(Ok(AgentEvent::Ready { .. }))
5512        ));
5513        adapter.send_prompt("hello".into()).await.expect("prompt");
5514        assert!(matches!(
5515            adapter.next_event().await,
5516            Some(Ok(AgentEvent::Failed {
5517                slot: 1,
5518                started: true,
5519                detail,
5520            })) if detail.contains("output token limit")
5521        ));
5522    }
5523
5524    #[tokio::test]
5525    async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
5526        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"}}'"#;
5527        let cwd = std::env::current_dir().expect("cwd");
5528        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
5529        adapter.start().await.expect("initialize");
5530        assert!(adapter.next_event().await.is_some());
5531        assert!(adapter.next_event().await.is_some());
5532
5533        for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
5534            adapter.send_prompt(prompt.into()).await.expect("prompt");
5535            assert!(matches!(
5536                adapter.next_event().await,
5537                Some(Ok(AgentEvent::Text { text, .. })) if text == expected
5538            ));
5539            assert!(matches!(
5540                adapter.next_event().await,
5541                Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
5542            ));
5543        }
5544        adapter.stop().await.expect("stop");
5545    }
5546
5547    #[tokio::test]
5548    async fn empty_acp_mode_catalog_disables_mode_control() {
5549        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":[]}}}'"#;
5550        let cwd = std::env::current_dir().expect("cwd");
5551        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5552        adapter.start().await.expect("initialize");
5553        assert!(!adapter.capabilities().supports_modes);
5554        assert!(matches!(
5555            adapter.next_event().await,
5556            Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
5557        ));
5558        assert!(matches!(
5559            adapter.next_event().await,
5560            Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
5561        ));
5562        adapter.stop().await.expect("stop");
5563    }
5564
5565    #[tokio::test]
5566    async fn acp_models_are_discovered_live_and_changed_through_session_config() {
5567        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"#;
5568        let cwd = std::env::current_dir().expect("cwd");
5569        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5570        adapter.start().await.expect("initialize");
5571        assert!(adapter.capabilities().supports_models);
5572        assert!(matches!(
5573            adapter.next_event().await,
5574            Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
5575                if config_id == "model"
5576                    && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
5577                    && current_model.as_deref() == Some("fast")
5578        ));
5579        assert!(matches!(
5580            adapter.next_event().await,
5581            Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
5582        ));
5583        adapter.set_model("smart".into()).await.expect("set model");
5584        assert!(adapter.set_model("invented".into()).await.is_err());
5585        adapter.stop().await.expect("stop");
5586    }
5587
5588    #[tokio::test]
5589    async fn acp_mode_change_is_acknowledged_without_provider_notification() {
5590        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":{}}'"#;
5591        let cwd = std::env::current_dir().expect("cwd");
5592        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5593        adapter.start().await.expect("initialize");
5594        adapter
5595            .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
5596            .await
5597            .expect("set mode");
5598        assert!(matches!(
5599            adapter.next_event().await,
5600            Some(Ok(AgentEvent::ModesReplaced { .. }))
5601        ));
5602        assert!(matches!(
5603            adapter.next_event().await,
5604            Some(Ok(AgentEvent::Ready { .. }))
5605        ));
5606        assert!(matches!(
5607            adapter.next_event().await,
5608            Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
5609        ));
5610        adapter.stop().await.expect("stop");
5611    }
5612
5613    #[tokio::test]
5614    async fn acp_reload_preserves_a_loadable_session_id() {
5615        let cwd = std::env::current_dir().expect("cwd");
5616        let mut adapter = AcpAdapter::with_session_id(
5617            0,
5618            cwd,
5619            "__codeswarm_missing_acp_for_reload_test__",
5620            Vec::new(),
5621            "saved-session",
5622        );
5623        adapter.capabilities.supports_session_load = true;
5624        // A failed replacement process still must not erase the session ID:
5625        // the coordinator can report the startup error and offer another
5626        // reload, preserving the only handle that can resume the conversation.
5627        assert!(adapter.reload().await.is_err());
5628        assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
5629    }
5630
5631    #[tokio::test]
5632    async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
5633        let cwd = std::env::current_dir().expect("cwd");
5634        let mut adapter = AcpAdapter::with_session_id(
5635            0,
5636            cwd,
5637            "__codeswarm_missing_nonloadable_acp__",
5638            Vec::new(),
5639            "stale-session",
5640        );
5641        adapter.capabilities.supports_session_load = false;
5642        assert!(adapter.reload().await.is_err());
5643        assert_eq!(adapter.session_id, None);
5644    }
5645
5646    #[tokio::test]
5647    async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
5648        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"}}'"#;
5649        let cwd = std::env::current_dir().expect("cwd");
5650        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
5651        adapter.start().await.expect("initialize");
5652        assert!(matches!(
5653            adapter.next_event().await,
5654            Some(Ok(AgentEvent::Ready { .. }))
5655        ));
5656        adapter.send_prompt("hello".into()).await.expect("prompt");
5657        assert!(matches!(
5658            adapter.next_event().await,
5659            Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
5660        ));
5661        assert!(matches!(
5662            adapter.next_event().await,
5663            Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
5664        ));
5665    }
5666
5667    #[test]
5668    fn acp_tool_patches_preserve_fields_and_honor_explicit_replacements() {
5669        let mut tools = std::collections::BTreeMap::new();
5670        let first = serde_json::json!({"sessionUpdate":"tool_call", "toolCallId":"read", "title":"Read config", "status":"in_progress",
5671            "content":[{"type":"content", "content":{"type":"text", "text":"old output"}}]});
5672        let initial = super::normalize_acp_tool(&first, &mut tools).unwrap();
5673        assert_eq!(initial.detail.as_deref(), Some("old output"));
5674        let completed = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call_update","toolCallId":"read","status":"completed"}), &mut tools).unwrap();
5675        assert_eq!(completed.title, "Read config");
5676        assert_eq!(completed.detail.as_deref(), Some("old output"));
5677        assert_eq!(completed.status, ToolStatus::Completed);
5678        let malformed = super::normalize_acp_tool(
5679            &serde_json::json!({"toolCallId":"read","title":3,"status":"unknown","content":null}),
5680            &mut tools,
5681        )
5682        .unwrap();
5683        assert_eq!(malformed, completed);
5684        let replaced = super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read","content":[false,{"type":"content","content":{"type":"text","text":"new output"}}]}), &mut tools).unwrap();
5685        assert_eq!(replaced.detail.as_deref(), Some("new output"));
5686        let cleared = super::normalize_acp_tool(
5687            &serde_json::json!({"toolCallId":"read","content":[]}),
5688            &mut tools,
5689        )
5690        .unwrap();
5691        assert_eq!(cleared.detail, None);
5692        let raw = super::normalize_acp_tool(
5693            &serde_json::json!({"toolCallId":"read","rawOutput":{"ok":true}}),
5694            &mut tools,
5695        )
5696        .unwrap();
5697        assert_eq!(raw.detail.as_deref(), Some("{\"ok\":true}"));
5698        let fresh = super::normalize_acp_tool(&serde_json::json!({"sessionUpdate":"tool_call","toolCallId":"read","title":"New call"}), &mut tools).unwrap();
5699        assert_eq!(fresh.status, ToolStatus::Pending);
5700        assert_eq!(fresh.detail, None);
5701        for invalid in [
5702            serde_json::json!({}),
5703            serde_json::json!({"toolCallId":7}),
5704            serde_json::json!({"toolCallId":" "}),
5705        ] {
5706            assert!(super::normalize_acp_tool(&invalid, &mut tools).is_none());
5707        }
5708        assert_eq!(tools.len(), 1);
5709        // IDs are opaque, not whitespace-normalized aliases of another tool.
5710        super::normalize_acp_tool(&serde_json::json!({"toolCallId":"read "}), &mut tools).unwrap();
5711        assert_eq!(tools.len(), 2);
5712    }
5713
5714    #[tokio::test]
5715    async fn acp_tool_status_only_notifications_retain_name_and_output() {
5716        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'
5717read _; echo '{"id":2,"result":{"sessionId":"s"}}'
5718read _
5719echo '{"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"}}]}}}'
5720echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"r","status":"completed"}}}'
5721echo '{"id":3,"result":{"stopReason":"end_turn"}}'"#;
5722        let mut adapter = AcpAdapter::new(
5723            0,
5724            std::env::current_dir().unwrap(),
5725            "sh",
5726            vec!["-c".into(), script.into()],
5727        );
5728        adapter.start().await.unwrap();
5729        adapter.next_event().await.unwrap().unwrap();
5730        adapter.send_prompt("read".into()).await.unwrap();
5731        for status in [ToolStatus::Running, ToolStatus::Completed] {
5732            let Some(Ok(AgentEvent::Tool { update, .. })) = adapter.next_event().await else {
5733                panic!("tool event");
5734            };
5735            assert_eq!(update.status, status);
5736            assert_eq!(update.title, "Read config");
5737            assert_eq!(update.detail.as_deref(), Some("file content"));
5738        }
5739        assert!(matches!(
5740            adapter.next_event().await,
5741            Some(Ok(AgentEvent::TurnComplete { .. }))
5742        ));
5743        adapter.stop().await.unwrap();
5744    }
5745
5746    #[tokio::test]
5747    async fn acp_reload_discards_old_queued_events_and_catalogs() {
5748        let script = r#"read _; echo '{"id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"id":2,"result":{"sessionId":"new"}}'"#;
5749        let mut adapter = AcpAdapter::new(
5750            0,
5751            std::env::current_dir().unwrap(),
5752            "sh",
5753            vec!["-c".into(), script.into()],
5754        );
5755        adapter.start().await.unwrap();
5756        adapter.queued_events.push_back(Ok(AgentEvent::Text {
5757            slot: 0,
5758            text: "stale".into(),
5759        }));
5760        adapter.modes = vec![Mode {
5761            id: "stale".into(),
5762            label: "Stale".into(),
5763        }];
5764        // Real peers echo each request's ID. Reset for this fixed-ID test script.
5765        adapter.next_request_id = 1;
5766        adapter.reload().await.unwrap();
5767        assert!(adapter.modes.is_empty());
5768        assert_eq!(adapter.queued_events.len(), 1);
5769        assert!(matches!(
5770            adapter.next_event().await,
5771            Some(Ok(AgentEvent::Ready { .. }))
5772        ));
5773        adapter.stop().await.unwrap();
5774        assert!(adapter.queued_events.is_empty());
5775    }
5776
5777    #[tokio::test]
5778    async fn acp_load_replays_history_without_starting_a_turn() {
5779        let script = r#"
5780read _
5781echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'
5782read request
5783case "$request" in *session/load*) ;; *) exit 2;; esac
5784echo '{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"text":"old question"}}}}'
5785echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"old answer"}}}}'
5786echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_thought_chunk","content":{"text":"old reasoning"}}}}'
5787echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call","toolCallId":"old-tool","title":"Read","status":"in_progress"}}}'
5788echo '{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"old-tool","title":"Read","status":"completed"}}}'
5789echo '{"jsonrpc":"2.0","id":2,"result":{}}'
5790read request
5791case "$request" in *session/prompt*) ;; *) exit 3;; esac
5792echo '{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"text":"new answer"}}}}'
5793echo '{"jsonrpc":"2.0","id":3,"result":{"stopReason":"end_turn"}}'
5794"#;
5795        let mut adapter = AcpAdapter::with_session_id(
5796            2,
5797            std::env::current_dir().unwrap(),
5798            "sh",
5799            vec!["-c".into(), script.into()],
5800            "saved",
5801        );
5802        adapter.start().await.unwrap();
5803        let mut state = crate::SessionState::new(3);
5804        for _ in 0..5 {
5805            let event = adapter.next_event().await.unwrap().unwrap();
5806            assert!(matches!(&event, AgentEvent::History { slot: 2, .. }));
5807            crate::reduce(&mut state, event);
5808            assert_eq!(state.active_slot, None);
5809        }
5810        assert!(matches!(
5811            adapter.next_event().await,
5812            Some(Ok(AgentEvent::Ready { slot: 2, .. }))
5813        ));
5814        adapter.send_prompt("new question".into()).await.unwrap();
5815        assert!(
5816            matches!(adapter.next_event().await, Some(Ok(AgentEvent::Text { text, .. })) if text == "new answer")
5817        );
5818        assert!(matches!(
5819            adapter.next_event().await,
5820            Some(Ok(AgentEvent::TurnComplete { slot: 2 }))
5821        ));
5822        adapter.stop().await.unwrap();
5823    }
5824
5825    #[tokio::test]
5826    async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
5827        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
5828        let cwd = std::env::current_dir().expect("cwd");
5829        let mut adapter = AcpAdapter::with_session_id(
5830            0,
5831            cwd,
5832            "sh",
5833            vec!["-c".into(), script.into()],
5834            "existing-session",
5835        );
5836        adapter.start().await.expect("load existing session");
5837        assert!(matches!(
5838            adapter.next_event().await,
5839            Some(Ok(AgentEvent::Ready { .. }))
5840        ));
5841    }
5842
5843    #[tokio::test]
5844    async fn acp_start_failure_reaps_transport_process() {
5845        // The child emits an invalid initialize response and exits. The
5846        // adapter must not retain a live child after protocol startup fails;
5847        // this is the path used when a configured ACP command is unavailable
5848        // or speaks a different protocol.
5849        let mut adapter = AcpAdapter::new(
5850            0,
5851            std::env::current_dir().expect("cwd"),
5852            "sh",
5853            vec!["-c".into(), "printf 'not-json\\n'".into()],
5854        );
5855        assert!(adapter.start().await.is_err());
5856        assert!(adapter.child.is_none());
5857        assert!(adapter.reader.is_none());
5858    }
5859
5860    #[tokio::test]
5861    async fn acp_transport_crash_is_reloaded_before_the_next_prompt() {
5862        let marker = unique_test_path("codeswarm-acp-retry", "count");
5863        let script = format!(
5864            r#"count=0
5865if [ -f '{0}' ]; then count=$(cat '{0}'); fi
5866count=$((count + 1))
5867printf '%s' "$count" > '{0}'
5868while IFS= read -r request; do
5869  id=$(printf '%s' "$request" | sed -n 's/.*"id":\([0-9][0-9]*\).*/\1/p')
5870  case "$request" in
5871    *initialize*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"agentCapabilities":{{"loadSession":true}}}}}}' ;;
5872    *session/new*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"sessionId":"saved-session"}}}}' ;;
5873    *session/load*) printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{}}}}' ;;
5874    *session/prompt*)
5875      if [ "$count" = 1 ]; then exit 0; fi
5876      printf '%s\n' '{{"jsonrpc":"2.0","method":"session/update","params":{{"update":{{"sessionUpdate":"agent_message_chunk","content":{{"text":"recovered"}}}}}}}}'
5877      printf '%s\n' '{{"jsonrpc":"2.0","id":'$id',"result":{{"stopReason":"end_turn"}}}}'
5878      ;;
5879  esac
5880done
5881"#,
5882            marker.display()
5883        );
5884        let mut adapter = AcpAdapter::new(
5885            0,
5886            std::env::current_dir().expect("cwd"),
5887            "sh",
5888            vec!["-c".into(), script],
5889        );
5890        adapter.start().await.expect("initial ACP startup");
5891        assert!(matches!(
5892            adapter.next_event().await,
5893            Some(Ok(AgentEvent::Ready { .. }))
5894        ));
5895        adapter
5896            .send_prompt("first".into())
5897            .await
5898            .expect("first prompt");
5899        assert!(matches!(
5900            adapter.next_event().await,
5901            Some(Err(AdapterError::Transport(_)))
5902        ));
5903        assert!(adapter.child.is_none());
5904        assert!(adapter.reader.is_none());
5905        assert_eq!(adapter.session_id(), Some("saved-session".into()));
5906
5907        adapter.reload().await.expect("reload ACP transport");
5908        assert!(matches!(
5909            adapter.next_event().await,
5910            Some(Ok(AgentEvent::Ready { .. }))
5911        ));
5912        adapter
5913            .send_prompt("retry".into())
5914            .await
5915            .expect("retry prompt");
5916        assert!(matches!(
5917            adapter.next_event().await,
5918            Some(Ok(AgentEvent::Text { text, .. })) if text == "recovered"
5919        ));
5920        assert!(matches!(
5921            adapter.next_event().await,
5922            Some(Ok(AgentEvent::TurnComplete { .. }))
5923        ));
5924        adapter.stop().await.expect("stop ACP");
5925        std::fs::remove_file(marker).expect("cleanup marker");
5926    }
5927
5928    #[tokio::test]
5929    async fn coordinator_reload_replays_context_and_reintroduces_a_crashed_slot() {
5930        let prompts = Arc::new(Mutex::new(Vec::new()));
5931        let healthy = ScriptedAdapter::new(
5932            0,
5933            AgentCapabilities::default(),
5934            [
5935                AgentEvent::Text {
5936                    slot: 0,
5937                    text: "peer context".into(),
5938                },
5939                AgentEvent::TurnComplete { slot: 0 },
5940            ],
5941        );
5942        let probe = ReloadProbeAdapter {
5943            slot: 1,
5944            crashed: false,
5945            reloaded: false,
5946            events: VecDeque::new(),
5947            prompts: Arc::clone(&prompts),
5948        };
5949        let mut relay = RelayHost::new(
5950            vec![
5951                AdapterHost::new(Box::new(healthy), None),
5952                AdapterHost::new(Box::new(probe), None),
5953            ],
5954            8,
5955        )
5956        .expect("relay");
5957        relay.start().await.expect("start");
5958        relay
5959            .run_turn("original task", 0)
5960            .await
5961            .expect("first turn");
5962        relay.run_turn("", 0).await.expect("crashed turn");
5963        relay.relay_mut().enqueue_human("retry", Some(1));
5964        relay.run_turn("", 0).await.expect("reloaded turn");
5965
5966        let prompts = prompts.lock().expect("prompts");
5967        let retry = prompts.last().expect("retry prompt");
5968        assert!(retry.contains("You are Reload probe"), "{retry}");
5969        assert!(retry.contains("original task"), "{retry}");
5970        assert!(retry.contains("peer context"), "{retry}");
5971        assert!(retry.contains("retry"), "{retry}");
5972    }
5973
5974    #[tokio::test]
5975    async fn acp_adapter_answers_permission_json_rpc_requests() {
5976        let path = std::env::temp_dir().join(format!(
5977            "codeswarm-permission-answer-{}",
5978            std::process::id()
5979        ));
5980        let script = format!(
5981            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"}}}}'"#,
5982            path.display()
5983        );
5984        let mut adapter = AcpAdapter::new(
5985            0,
5986            std::env::current_dir().expect("cwd"),
5987            "sh",
5988            vec!["-c".into(), script],
5989        );
5990        adapter.start().await.expect("start ACP");
5991        assert!(matches!(
5992            adapter.next_event().await,
5993            Some(Ok(AgentEvent::Ready { .. }))
5994        ));
5995        adapter.send_prompt("do it".into()).await.expect("prompt");
5996        assert!(matches!(
5997            adapter.next_event().await,
5998            Some(Ok(AgentEvent::Permission { request, .. }))
5999                if request.id == "9"
6000                    && request.options == ["Allow once"]
6001                    && request.option_ids == ["allow-once"]
6002        ));
6003        adapter
6004            .answer_permission(
6005                "9".into(),
6006                PermissionAnswer::Selected {
6007                    option_id: "allow-once".into(),
6008                },
6009            )
6010            .await
6011            .expect("permission answer");
6012        assert!(matches!(
6013            adapter.next_event().await,
6014            Some(Ok(AgentEvent::TurnComplete { .. }))
6015        ));
6016        let answer: Value = serde_json::from_str(
6017            &std::fs::read_to_string(&path).expect("captured permission answer"),
6018        )
6019        .expect("valid JSON-RPC answer");
6020        assert_eq!(answer["id"], 9);
6021        assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
6022        assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
6023        std::fs::remove_file(path).expect("cleanup");
6024    }
6025
6026    #[test]
6027    fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
6028        let event = parse_acp_notification(
6029            0,
6030            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
6031        )
6032        .expect("valid JSON-RPC request");
6033        assert!(event.is_none());
6034    }
6035
6036    #[tokio::test]
6037    async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
6038        let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
6039        std::fs::write(
6040            &script_path,
6041            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
6042        )
6043        .expect("write native test script");
6044        let mut adapter = AgyAdapter::new(
6045            0,
6046            std::env::current_dir().expect("cwd"),
6047            format!("sh {}", script_path.display()),
6048        );
6049        adapter.start().await.expect("start native adapter");
6050        assert!(adapter.next_event().await.is_some());
6051        assert!(adapter.next_event().await.is_some());
6052        adapter
6053            .send_prompt("continue".into())
6054            .await
6055            .expect("prompt");
6056        assert!(matches!(
6057            adapter.next_event().await,
6058            Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
6059        ));
6060        assert!(matches!(
6061            adapter.next_event().await,
6062            Some(Ok(AgentEvent::TurnComplete { .. }))
6063        ));
6064        adapter.stop().await.expect("stop native adapter");
6065        std::fs::remove_file(script_path).expect("cleanup native script");
6066    }
6067
6068    #[test]
6069    fn acp_workspace_file_access_is_root_bound_and_size_limited() {
6070        let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
6071        let _ = std::fs::remove_dir_all(&root);
6072        std::fs::create_dir_all(&root).expect("workspace");
6073        std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
6074        let outside =
6075            std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
6076        std::fs::write(&outside, "secret").expect("outside file");
6077        let link = root.join("outside-link");
6078        #[cfg(unix)]
6079        std::os::unix::fs::symlink(&outside, &link).expect("symlink");
6080        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6081
6082        assert_eq!(
6083            adapter
6084                .read_workspace_text("inside.txt", Some(2), Some(1))
6085                .expect("read inside"),
6086            "two"
6087        );
6088        std::fs::write(
6089            root.join("large.txt"),
6090            vec![b'x'; MAX_FILE_READ_BYTES + 1024],
6091        )
6092        .expect("large file");
6093        let bounded = adapter
6094            .read_workspace_text("large.txt", None, None)
6095            .expect("bounded read");
6096        assert!(bounded.len() <= MAX_FILE_READ_BYTES);
6097        #[cfg(unix)]
6098        {
6099            std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
6100                .expect("internal symlink");
6101            assert_eq!(
6102                adapter
6103                    .read_workspace_text("inside-link", None, None)
6104                    .expect("read internal symlink"),
6105                "one\ntwo\nthree\n"
6106            );
6107        }
6108        assert!(adapter.workspace_path("../codeswarm-outside").is_err());
6109        assert!(
6110            adapter
6111                .workspace_path(&outside.display().to_string())
6112                .is_err()
6113        );
6114        #[cfg(unix)]
6115        assert!(adapter.workspace_path("outside-link").is_err());
6116        #[cfg(unix)]
6117        std::fs::remove_file(link).expect("cleanup symlink");
6118        #[cfg(unix)]
6119        std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
6120        std::fs::remove_file(outside).expect("cleanup outside");
6121        std::fs::remove_dir_all(root).expect("cleanup workspace");
6122    }
6123
6124    #[tokio::test]
6125    async fn running_terminal_output_omits_exit_status_until_completion() {
6126        let root = unique_test_path("codeswarm-terminal-output", "dir");
6127        std::fs::create_dir_all(&root).expect("workspace");
6128        let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
6129        let result = adapter
6130            .terminal_create(&serde_json::json!({
6131                "command": "sh",
6132                "args": ["-c", "sleep 0.2; printf done"],
6133                "cwd": ".",
6134            }))
6135            .await
6136            .expect("terminal create");
6137        let id = result["terminalId"].as_str().expect("terminal id");
6138        let output = adapter.terminal_output(id).await.expect("terminal output");
6139        assert!(output.get("exitStatus").is_none());
6140        if let Some(terminal) = adapter.terminals.remove(id) {
6141            terminal.stop().await;
6142        }
6143        std::fs::remove_dir_all(root).expect("cleanup workspace");
6144    }
6145
6146    #[tokio::test]
6147    async fn acp_adapter_answers_workspace_read_requests() {
6148        let root =
6149            std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
6150        let _ = std::fs::remove_dir_all(&root);
6151        std::fs::create_dir_all(&root).expect("workspace");
6152        let source = root.join("inside.txt");
6153        let answer = root.join("answer.json");
6154        std::fs::write(&source, "workspace content").expect("source");
6155        let script = format!(
6156            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"}}}}'"#,
6157            source.display(),
6158            answer.display(),
6159        );
6160        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6161        adapter.start().await.expect("start ACP");
6162        assert!(matches!(
6163            adapter.next_event().await,
6164            Some(Ok(AgentEvent::Ready { .. }))
6165        ));
6166        adapter.send_prompt("read it".into()).await.expect("prompt");
6167        assert!(matches!(
6168            adapter.next_event().await,
6169            Some(Ok(AgentEvent::TurnComplete { .. }))
6170        ));
6171        let response: Value =
6172            serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
6173                .expect("response JSON");
6174        assert_eq!(response["id"], 9);
6175        assert_eq!(response["result"]["content"], "workspace content");
6176        adapter.stop().await.expect("stop ACP");
6177        std::fs::remove_dir_all(root).expect("cleanup workspace");
6178    }
6179
6180    #[tokio::test]
6181    async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
6182        let root =
6183            std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
6184        let _ = std::fs::remove_dir_all(&root);
6185        std::fs::create_dir_all(&root).expect("workspace");
6186        let create_request = serde_json::json!({
6187            "jsonrpc": "2.0",
6188            "id": 9,
6189            "method": "terminal/create",
6190            "params": {
6191                "sessionId": "s1",
6192                "command": "sh",
6193                "args": ["-c", "sleep 0.1; printf terminal-ok"],
6194                "cwd": ".",
6195            },
6196        });
6197        let wait_request = serde_json::json!({
6198            "jsonrpc": "2.0",
6199            "id": 10,
6200            "method": "terminal/wait_for_exit",
6201            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6202        });
6203        let output_request = serde_json::json!({
6204            "jsonrpc": "2.0",
6205            "id": 11,
6206            "method": "terminal/output",
6207            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
6208        });
6209        let create_answer = root.join("create-answer.json");
6210        let wait_answer = root.join("wait-answer.json");
6211        let output_answer = root.join("output-answer.json");
6212        let script = format!(
6213            "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\"}}}}'",
6214            create_request,
6215            create_answer.display(),
6216            wait_request,
6217            wait_answer.display(),
6218            output_request,
6219            output_answer.display(),
6220        );
6221        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
6222        adapter.start().await.expect("start ACP");
6223        assert!(matches!(
6224            adapter.next_event().await,
6225            Some(Ok(AgentEvent::Ready { .. }))
6226        ));
6227        adapter
6228            .send_prompt("run terminal".into())
6229            .await
6230            .expect("prompt");
6231        let mut saw_complete = false;
6232        for _ in 0..6 {
6233            match adapter.next_event().await {
6234                Some(Ok(AgentEvent::TurnComplete { .. })) => {
6235                    saw_complete = true;
6236                    break;
6237                }
6238                Some(_) => {}
6239                None => break,
6240            }
6241        }
6242        assert!(saw_complete, "terminal requests should not stall ACP");
6243        let create: Value = serde_json::from_str(
6244            &std::fs::read_to_string(&create_answer).expect("captured create response"),
6245        )
6246        .expect("create JSON");
6247        assert_eq!(create["result"]["terminalId"], "terminal-1");
6248        let output: Value = serde_json::from_str(
6249            &std::fs::read_to_string(&output_answer).expect("captured output response"),
6250        )
6251        .expect("output JSON");
6252        assert!(
6253            output["result"]["output"]
6254                .as_str()
6255                .unwrap_or_default()
6256                .contains("terminal-ok"),
6257            "output response: {output}"
6258        );
6259        adapter.stop().await.expect("stop ACP");
6260        std::fs::remove_dir_all(root).expect("cleanup workspace");
6261    }
6262
6263    #[tokio::test]
6264    async fn host_reduces_and_persists_adapter_events() {
6265        let path =
6266            std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
6267        let adapter = ScriptedAdapter::new(
6268            0,
6269            AgentCapabilities::default(),
6270            [AgentEvent::Text {
6271                slot: 0,
6272                text: "hello".into(),
6273            }],
6274        );
6275        let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
6276        host.start().await.expect("start");
6277        host.next_effects()
6278            .await
6279            .expect("event")
6280            .expect("valid event");
6281        assert_eq!(host.state.public_text[0].1, "hello");
6282        assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
6283        std::fs::remove_file(path).expect("cleanup");
6284    }
6285
6286    #[tokio::test]
6287    async fn relay_applies_default_policy_before_the_first_prompt() {
6288        let first_log = Arc::new(Mutex::new(Vec::new()));
6289        let second_log = Arc::new(Mutex::new(Vec::new()));
6290        let first = AdapterHost::new(
6291            Box::new(ModeOrderAdapter {
6292                slot: 0,
6293                log: Arc::clone(&first_log),
6294                phase: 0,
6295            }),
6296            None,
6297        );
6298        let second = AdapterHost::new(
6299            Box::new(ModeOrderAdapter {
6300                slot: 1,
6301                log: Arc::clone(&second_log),
6302                phase: 0,
6303            }),
6304            None,
6305        );
6306        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6307        relay.start().await.expect("start and synchronize policy");
6308        relay.run_turn("task", 0).await.expect("first turn");
6309        {
6310            let log = first_log.lock().expect("log");
6311            assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
6312        }
6313        assert_eq!(
6314            second_log.lock().expect("log").as_slice(),
6315            ["start", "mode:yolo"]
6316        );
6317        let added_log = Arc::new(Mutex::new(Vec::new()));
6318        relay
6319            .add_agent(
6320                AdapterHost::new(
6321                    Box::new(ModeOrderAdapter {
6322                        slot: 2,
6323                        log: Arc::clone(&added_log),
6324                        phase: 0,
6325                    }),
6326                    None,
6327                ),
6328                "Added",
6329                "added.example",
6330                "added-agent",
6331            )
6332            .await
6333            .expect("add with synchronized policy");
6334        assert_eq!(
6335            added_log.lock().expect("log").as_slice(),
6336            ["start", "mode:yolo"]
6337        );
6338        relay.drop_agent(2).await.expect("drop added agent");
6339        added_log.lock().expect("log").clear();
6340        relay
6341            .reload(2)
6342            .await
6343            .expect("reload with synchronized policy");
6344        assert_eq!(
6345            added_log.lock().expect("log").as_slice(),
6346            ["reload", "mode:yolo"]
6347        );
6348    }
6349
6350    #[tokio::test]
6351    async fn acp_roster_is_ready_before_any_prompt_is_sent() {
6352        let hosts = (0..2)
6353            .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
6354            .collect::<Vec<_>>();
6355        let startup_events = Arc::new(Mutex::new(Vec::new()));
6356        let captured = Arc::clone(&startup_events);
6357        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6358        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6359
6360        relay.start().await.expect("complete startup handshake");
6361
6362        assert!(relay.dispatches().is_empty());
6363        let ready_slots = startup_events
6364            .lock()
6365            .expect("events")
6366            .iter()
6367            .filter_map(|event| match event {
6368                AgentEvent::Ready { slot, .. } => Some(*slot),
6369                _ => None,
6370            })
6371            .collect::<Vec<_>>();
6372        assert_eq!(ready_slots, vec![0, 1]);
6373    }
6374
6375    #[tokio::test]
6376    async fn independent_roster_adapters_start_concurrently() {
6377        let barrier = Arc::new(tokio::sync::Barrier::new(2));
6378        let hosts = (0..2)
6379            .map(|slot| {
6380                AdapterHost::new(
6381                    Box::new(ConcurrentStartAdapter {
6382                        slot,
6383                        barrier: Arc::clone(&barrier),
6384                    }),
6385                    None,
6386                )
6387            })
6388            .collect::<Vec<_>>();
6389        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6390        tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
6391            .await
6392            .expect("startup should not serialize barrier participants")
6393            .expect("startup succeeds");
6394    }
6395
6396    #[tokio::test]
6397    async fn relay_host_dispatches_turns_sequentially() {
6398        let capabilities = AgentCapabilities {
6399            supports_cancel: true,
6400            ..AgentCapabilities::default()
6401        };
6402        let first = ScriptedAdapter::new(
6403            0,
6404            capabilities.clone(),
6405            [
6406                AgentEvent::Text {
6407                    slot: 0,
6408                    text: "first".into(),
6409                },
6410                AgentEvent::TurnComplete { slot: 0 },
6411            ],
6412        );
6413        let second = ScriptedAdapter::new(
6414            1,
6415            capabilities,
6416            [
6417                AgentEvent::Text {
6418                    slot: 1,
6419                    text: "review".into(),
6420                },
6421                AgentEvent::TurnComplete { slot: 1 },
6422            ],
6423        );
6424        let hosts = vec![
6425            AdapterHost::new(Box::new(first), None),
6426            AdapterHost::new(Box::new(second), None),
6427        ];
6428        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6429        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6430        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6431        let captured = std::sync::Arc::clone(&events);
6432        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6433        relay.start().await.expect("start");
6434        events.lock().expect("events").clear();
6435        assert!(matches!(
6436            relay.run_turn("task", 0).await.expect("first turn"),
6437            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6438        ));
6439        assert!(matches!(
6440            relay.run_turn("first", 0).await.expect("second turn"),
6441            crate::relay::RelayDecision::Dispatch {
6442                slot: 1,
6443                can_stop: true,
6444                ..
6445            }
6446        ));
6447        assert_eq!(
6448            relay
6449                .dispatches()
6450                .iter()
6451                .map(|(slot, _)| *slot)
6452                .collect::<Vec<_>>(),
6453            [0, 1]
6454        );
6455        assert!(relay.dispatches()[0].1.contains("You are Claude"));
6456        assert!(
6457            relay.dispatches()[0]
6458                .1
6459                .contains("CodeSwarm roster (ordered)")
6460        );
6461        assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
6462        assert!(relay.dispatches()[0].1.contains("2. Codex"));
6463        assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
6464        assert!(
6465            relay.dispatches()[1]
6466                .1
6467                .contains("stops all other agents and ends the entire automated relay")
6468        );
6469        assert!(relay.dispatches()[1].1.contains("Use it with extreme care"));
6470        assert!(
6471            relay.dispatches()[1]
6472                .1
6473                .contains("If there is any uncertainty, do not use it")
6474        );
6475        assert!(relay.dispatches()[0].1.contains("Do not use"));
6476        let lifecycle = events.lock().expect("events");
6477        let positions = lifecycle
6478            .iter()
6479            .filter_map(|event| match event {
6480                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
6481                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
6482                _ => None,
6483            })
6484            .collect::<Vec<_>>();
6485        assert_eq!(
6486            positions,
6487            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
6488        );
6489    }
6490
6491    #[tokio::test]
6492    async fn failed_resume_does_not_stop_or_dispatch_to_healthy_peer() {
6493        let stops = Arc::new(AtomicUsize::new(0));
6494        let failed = FailingStartAdapter {
6495            slot: 0,
6496            stopped: stops.clone(),
6497        };
6498        let healthy = ScriptedAdapter::new(
6499            1,
6500            AgentCapabilities::default(),
6501            [
6502                AgentEvent::Text {
6503                    slot: 1,
6504                    text: "healthy response".into(),
6505                },
6506                AgentEvent::TurnComplete { slot: 1 },
6507            ],
6508        );
6509        let events = Arc::new(std::sync::Mutex::new(Vec::new()));
6510        let captured = events.clone();
6511        let mut relay = RelayHost::new(
6512            vec![
6513                AdapterHost::new(Box::new(failed), None),
6514                AdapterHost::new(Box::new(healthy), None),
6515            ],
6516            4,
6517        )
6518        .unwrap();
6519        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
6520        relay.start_resuming().await.unwrap();
6521        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
6522        assert!(relay.dispatches().is_empty());
6523        assert_eq!(stops.load(Ordering::Relaxed), 1);
6524        assert!(
6525            events
6526                .lock()
6527                .unwrap()
6528                .iter()
6529                .any(|event| matches!(event, AgentEvent::Failed { slot: 0, .. }))
6530        );
6531        assert!(
6532            !events
6533                .lock()
6534                .unwrap()
6535                .iter()
6536                .any(|event| matches!(event, AgentEvent::Failed { slot: 1, .. }))
6537        );
6538        assert!(!relay.relay_mut().enqueue_human("do not retarget", Some(0)));
6539        assert!(
6540            relay
6541                .relay_mut()
6542                .enqueue_human("explicit healthy target", Some(1))
6543        );
6544        assert!(matches!(
6545            relay.run_turn("", 1).await.unwrap(),
6546            RelayDecision::Dispatch { slot: 1, .. }
6547        ));
6548        relay.stop().await.unwrap();
6549    }
6550
6551    #[tokio::test]
6552    async fn pair_strategy_wires_roles_into_non_direct_prompts() {
6553        let capabilities = AgentCapabilities::default();
6554        let first = ScriptedAdapter::new(
6555            0,
6556            capabilities.clone(),
6557            [
6558                AgentEvent::Text {
6559                    slot: 0,
6560                    text: "implemented".into(),
6561                },
6562                AgentEvent::TurnComplete { slot: 0 },
6563                AgentEvent::Text {
6564                    slot: 0,
6565                    text: format!("fixed review findings {STOP_TOKEN}"),
6566                },
6567                AgentEvent::TurnComplete { slot: 0 },
6568            ],
6569        );
6570        let second = ScriptedAdapter::new(
6571            1,
6572            capabilities,
6573            [
6574                AgentEvent::Text {
6575                    slot: 1,
6576                    text: "reviewed".into(),
6577                },
6578                AgentEvent::TurnComplete { slot: 1 },
6579                AgentEvent::Text {
6580                    slot: 1,
6581                    text: format!("approved {STOP_TOKEN}"),
6582                },
6583                AgentEvent::TurnComplete { slot: 1 },
6584                AgentEvent::Text {
6585                    slot: 1,
6586                    text: "new task".into(),
6587                },
6588                AgentEvent::TurnComplete { slot: 1 },
6589            ],
6590        );
6591        let hosts = vec![
6592            AdapterHost::new(Box::new(first), None),
6593            AdapterHost::new(Box::new(second), None),
6594        ];
6595        let mut relay = RelayHost::new(hosts, 4).expect("relay");
6596        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6597        relay.relay_mut().set_strategy(CollaborationStrategy::Pair);
6598        relay.start().await.expect("start");
6599        assert!(matches!(
6600            relay.run_turn("task", 0).await.expect("implementer turn"),
6601            RelayDecision::Dispatch {
6602                slot: 0,
6603                can_stop: false,
6604                ..
6605            }
6606        ));
6607        let implementer_prompt = &relay.dispatches()[0].1;
6608        assert!(implementer_prompt.contains("you are the implementer"));
6609        assert!(implementer_prompt.contains("pair reviewer will review the result next"));
6610        assert!(!implementer_prompt.contains("you are the reviewer"));
6611        assert!(implementer_prompt.contains("Do not use"));
6612        assert!(matches!(
6613            relay.run_turn("", 0).await.expect("reviewer turn"),
6614            RelayDecision::Dispatch {
6615                slot: 1,
6616                can_stop: true,
6617                ..
6618            }
6619        ));
6620        let reviewer_prompt = &relay.dispatches()[1].1;
6621        assert!(reviewer_prompt.contains("you are the reviewer"));
6622        assert!(reviewer_prompt.contains("Claude handed off"));
6623        assert!(reviewer_prompt.contains("concrete defects"));
6624        assert!(reviewer_prompt.contains("concise approval"));
6625        assert!(reviewer_prompt.contains(STOP_TOKEN));
6626        assert!(!reviewer_prompt.contains("you are the implementer"));
6627        relay.run_turn("", 0).await.unwrap();
6628        assert!(relay.dispatches()[2].1.contains("you are the implementer"));
6629        assert!(relay.dispatches()[2].1.contains("Do not use"));
6630        assert!(matches!(
6631            relay.run_turn("", 0).await.unwrap(),
6632            RelayDecision::Dispatch { slot: 1, .. }
6633        ));
6634        assert!(relay.dispatches()[3].1.contains("you are the reviewer"));
6635        assert!(relay.relay_mut().enqueue_human("new task", Some(1)));
6636        relay.run_turn("", 1).await.unwrap();
6637        assert!(relay.dispatches()[4].1.contains("you are the implementer"));
6638    }
6639
6640    #[tokio::test]
6641    async fn solo_roster_and_direct_prompts_omit_pair_roles() {
6642        let solo = ScriptedAdapter::new(
6643            0,
6644            AgentCapabilities::default(),
6645            [
6646                AgentEvent::Text {
6647                    slot: 0,
6648                    text: "solo".into(),
6649                },
6650                AgentEvent::TurnComplete { slot: 0 },
6651            ],
6652        );
6653        let mut solo_relay =
6654            RelayHost::new(vec![AdapterHost::new(Box::new(solo), None)], 4).expect("relay");
6655        solo_relay
6656            .relay_mut()
6657            .set_strategy(CollaborationStrategy::Pair);
6658        solo_relay.start().await.expect("start");
6659        solo_relay.run_turn("task", 0).await.expect("solo turn");
6660        assert!(!solo_relay.dispatches()[0].1.contains("Pair role"));
6661
6662        let roster_first = ScriptedAdapter::new(
6663            0,
6664            AgentCapabilities::default(),
6665            [AgentEvent::TurnComplete { slot: 0 }],
6666        );
6667        let roster_second = ScriptedAdapter::new(
6668            1,
6669            AgentCapabilities::default(),
6670            [AgentEvent::TurnComplete { slot: 1 }],
6671        );
6672        let mut roster = RelayHost::new(
6673            vec![
6674                AdapterHost::new(Box::new(roster_first), None),
6675                AdapterHost::new(Box::new(roster_second), None),
6676            ],
6677            4,
6678        )
6679        .expect("relay");
6680        roster.start().await.expect("start");
6681        roster.run_turn("task", 0).await.expect("first turn");
6682        roster.run_turn("", 0).await.expect("second turn");
6683        assert!(!roster.dispatches()[0].1.contains("Pair role"));
6684        assert!(!roster.dispatches()[1].1.contains("Pair role"));
6685
6686        let pair_first = ScriptedAdapter::new(
6687            0,
6688            AgentCapabilities::default(),
6689            [AgentEvent::TurnComplete { slot: 0 }],
6690        );
6691        let pair_second = ScriptedAdapter::new(
6692            1,
6693            AgentCapabilities::default(),
6694            [AgentEvent::TurnComplete { slot: 1 }],
6695        );
6696        let mut pair = RelayHost::new(
6697            vec![
6698                AdapterHost::new(Box::new(pair_first), None),
6699                AdapterHost::new(Box::new(pair_second), None),
6700            ],
6701            4,
6702        )
6703        .expect("relay");
6704        pair.relay_mut().set_strategy(CollaborationStrategy::Pair);
6705        assert_eq!(pair.relay_mut().enqueue_direct(1, "private"), Ok(true));
6706        pair.start().await.expect("start");
6707        assert!(matches!(
6708            pair.run_turn("ignored", 0).await.expect("direct turn"),
6709            RelayDecision::Dispatch {
6710                slot: 1,
6711                direct: true,
6712                ..
6713            }
6714        ));
6715        let direct_prompt = &pair.dispatches()[0].1;
6716        assert!(direct_prompt.contains("private"));
6717        assert!(!direct_prompt.contains("Pair role"));
6718    }
6719
6720    #[tokio::test]
6721    async fn relay_host_routes_around_a_usage_limited_agent() {
6722        let capabilities = AgentCapabilities::default();
6723        let first = ScriptedAdapter::new(
6724            0,
6725            capabilities.clone(),
6726            [
6727                AgentEvent::Text {
6728                    slot: 0,
6729                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
6730                           credits or try again later."
6731                        .into(),
6732                },
6733                AgentEvent::TurnComplete { slot: 0 },
6734            ],
6735        );
6736        let second = ScriptedAdapter::new(
6737            1,
6738            capabilities,
6739            [
6740                AgentEvent::Text {
6741                    slot: 1,
6742                    text: "review done".into(),
6743                },
6744                AgentEvent::TurnComplete { slot: 1 },
6745            ],
6746        );
6747        let hosts = vec![
6748            AdapterHost::new(Box::new(first), None),
6749            AdapterHost::new(Box::new(second), None),
6750        ];
6751        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
6752        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6753        let captured = std::sync::Arc::clone(&events);
6754        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6755        relay.start().await.expect("start");
6756        events.lock().expect("events").clear();
6757        assert!(matches!(
6758            relay.run_turn("task", 0).await.expect("limited turn"),
6759            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6760        ));
6761        assert!(
6762            events
6763                .lock()
6764                .expect("events")
6765                .iter()
6766                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6767        );
6768        // The next automatic turn skips the limited agent entirely.
6769        assert!(matches!(
6770            relay.run_turn("", 0).await.expect("next turn"),
6771            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6772        ));
6773        assert!(relay.relay().is_limited(0));
6774        // A reload restores the agent to the ring. (ScriptedAdapter cannot
6775        // feed further turns, so the restored routing itself is covered by
6776        // the relay unit tests.)
6777        relay.reload(0).await.expect("reload");
6778        assert!(!relay.relay().is_limited(0));
6779    }
6780
6781    #[tokio::test]
6782    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
6783        let limited = ScriptedAdapter::new(
6784            0,
6785            AgentCapabilities::default(),
6786            [AgentEvent::Failed {
6787                slot: 0,
6788                started: true,
6789                detail: "request failed: insufficient_quota".into(),
6790            }],
6791        );
6792        let healthy = ScriptedAdapter::new(
6793            1,
6794            AgentCapabilities::default(),
6795            [AgentEvent::TurnComplete { slot: 1 }],
6796        );
6797        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6798        let captured = std::sync::Arc::clone(&events);
6799        let mut relay = RelayHost::new(
6800            vec![
6801                AdapterHost::new(Box::new(limited), None),
6802                AdapterHost::new(Box::new(healthy), None),
6803            ],
6804            4,
6805        )
6806        .expect("relay");
6807        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
6808        relay.start().await.expect("start");
6809        events.lock().expect("events").clear();
6810
6811        assert!(matches!(
6812            relay.run_turn("task", 0).await.expect("limited failure"),
6813            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6814        ));
6815        assert!(relay.relay().is_limited(0));
6816        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
6817        {
6818            let events = events.lock().expect("events");
6819            assert!(
6820                events
6821                    .iter()
6822                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
6823            );
6824            assert!(
6825                !events
6826                    .iter()
6827                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
6828            );
6829        }
6830
6831        assert!(matches!(
6832            relay.run_turn("", 0).await.expect("healthy peer"),
6833            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6834        ));
6835    }
6836
6837    #[tokio::test]
6838    async fn relay_failure_is_skipped_for_one_batch_without_changing_the_roster() {
6839        let failed = ScriptedAdapter::new(
6840            0,
6841            AgentCapabilities::default(),
6842            [AgentEvent::Failed {
6843                slot: 0,
6844                started: true,
6845                detail: "connection lost".into(),
6846            }],
6847        );
6848        let healthy = ScriptedAdapter::new(
6849            1,
6850            AgentCapabilities::default(),
6851            [AgentEvent::TurnComplete { slot: 1 }],
6852        );
6853        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6854        let captured = std::sync::Arc::clone(&events);
6855        let mut relay = RelayHost::new(
6856            vec![
6857                AdapterHost::new(Box::new(failed), None),
6858                AdapterHost::new(Box::new(healthy), None),
6859            ],
6860            4,
6861        )
6862        .expect("relay");
6863        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6864        relay.start().await.expect("start");
6865
6866        assert!(matches!(
6867            relay.run_turn("task", 0).await.expect("handled failure"),
6868            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
6869        ));
6870        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![0, 1]);
6871        assert!(relay.relay().is_limited(0));
6872        assert!(events.lock().expect("lock").iter().any(|event| {
6873            matches!(
6874                event,
6875                AgentEvent::Failed {
6876                    slot: 0,
6877                    started: true,
6878                    ..
6879                }
6880            )
6881        }));
6882        assert!(matches!(
6883            relay.run_turn("", 0).await.expect("healthy peer"),
6884            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
6885        ));
6886    }
6887
6888    #[tokio::test]
6889    async fn codex_stop_does_not_skip_later_roster_reviewers() {
6890        let hosts = (0..3)
6891            .map(|slot| {
6892                AdapterHost::new(
6893                    Box::new(ScriptedAdapter::new(
6894                        slot,
6895                        AgentCapabilities::default(),
6896                        [
6897                            AgentEvent::Text {
6898                                slot,
6899                                text: STOP_TOKEN.into(),
6900                            },
6901                            AgentEvent::TurnComplete { slot },
6902                        ],
6903                    )),
6904                    None,
6905                )
6906            })
6907            .collect();
6908        let mut relay = RelayHost::new(hosts, 10).expect("relay");
6909        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
6910        relay.start().await.expect("start");
6911        for expected in 0..3 {
6912            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
6913                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
6914        }
6915        assert_eq!(
6916            relay.run_turn("", 0).await.expect("complete"),
6917            RelayDecision::Complete
6918        );
6919    }
6920
6921    #[tokio::test]
6922    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
6923        let first = ScriptedAdapter::new(
6924            0,
6925            AgentCapabilities::default(),
6926            [
6927                AgentEvent::Text {
6928                    slot: 0,
6929                    text: "done".into(),
6930                },
6931                AgentEvent::TurnComplete { slot: 0 },
6932            ],
6933        );
6934        let reviewer = ScriptedAdapter::new(
6935            1,
6936            AgentCapabilities::default(),
6937            [
6938                AgentEvent::Text {
6939                    slot: 1,
6940                    text: STOP_TOKEN.into(),
6941                },
6942                AgentEvent::TurnComplete { slot: 1 },
6943            ],
6944        );
6945        let mut relay = RelayHost::new(
6946            vec![
6947                AdapterHost::new(Box::new(first), None),
6948                AdapterHost::new(Box::new(reviewer), None),
6949            ],
6950            10,
6951        )
6952        .expect("relay");
6953        relay.start().await.expect("start");
6954        let first_decision = relay.run_turn("task", 0).await.expect("first");
6955        assert!(matches!(
6956            first_decision,
6957            RelayDecision::Dispatch { slot: 0, .. }
6958        ));
6959        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
6960        assert!(matches!(
6961            reviewer_decision,
6962            RelayDecision::Dispatch {
6963                slot: 1,
6964                can_stop: true,
6965                ..
6966            }
6967        ));
6968        assert_eq!(
6969            relay.run_turn("", 0).await.expect("complete"),
6970            RelayDecision::Complete
6971        );
6972    }
6973
6974    #[tokio::test]
6975    async fn relay_stream_emits_text_and_thought_endings_before_tools() {
6976        let tool = AgentEvent::Tool {
6977            slot: 0,
6978            update: crate::ToolUpdate {
6979                id: "read".into(),
6980                title: "Read file".into(),
6981                status: ToolStatus::Running,
6982                detail: None,
6983            },
6984        };
6985        let updates = vec![
6986            AgentEvent::Thought {
6987                slot: 0,
6988                text: "Check the buffer. ✈".into(),
6989            },
6990            AgentEvent::Text {
6991                slot: 0,
6992                text: "Let me check.".into(),
6993            },
6994            tool.clone(),
6995            AgentEvent::Text {
6996                slot: 0,
6997                text: "[CODE".into(),
6998            },
6999            AgentEvent::Text {
7000                slot: 0,
7001                text: " is ordinary.".into(),
7002            },
7003            AgentEvent::Text {
7004                slot: 0,
7005                text: "[CODESWARM:".into(),
7006            },
7007            AgentEvent::Text {
7008                slot: 0,
7009                text: "STOP] Done.".into(),
7010            },
7011            tool,
7012            AgentEvent::TurnComplete { slot: 0 },
7013        ];
7014        let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
7015        let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
7016        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7017        let captured = std::sync::Arc::clone(&events);
7018        let mut relay = RelayHost::new(
7019            vec![
7020                AdapterHost::new(Box::new(first), None),
7021                AdapterHost::new(Box::new(reviewer), None),
7022            ],
7023            2,
7024        )
7025        .expect("relay");
7026        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7027        relay.start().await.unwrap();
7028        relay.run_turn("task", 0).await.unwrap();
7029        let captured = events.lock().unwrap();
7030        let visible: Vec<_> = captured
7031            .iter()
7032            .filter(|event| {
7033                matches!(
7034                    event,
7035                    AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
7036                )
7037            })
7038            .cloned()
7039            .collect();
7040        assert_eq!(
7041            visible,
7042            vec![
7043                updates[0].clone(),
7044                updates[1].clone(),
7045                updates[2].clone(),
7046                AgentEvent::Text {
7047                    slot: 0,
7048                    text: "[CODE is ordinary.".into()
7049                },
7050                AgentEvent::Text {
7051                    slot: 0,
7052                    text: " Done.".into()
7053                },
7054                updates[7].clone(),
7055            ]
7056        );
7057    }
7058
7059    #[tokio::test]
7060    async fn roster_handoff_routes_only_terminal_message_markers_and_refreshes_targets() {
7061        let text = |value: &str| AgentEvent::Text {
7062            slot: 0,
7063            text: value.into(),
7064        };
7065        let thought = || AgentEvent::Thought {
7066            slot: 0,
7067            text: "still checking".into(),
7068        };
7069        let tool = || AgentEvent::Tool {
7070            slot: 0,
7071            update: crate::ToolUpdate {
7072                id: "read".into(),
7073                title: "Read file".into(),
7074                status: ToolStatus::Running,
7075                detail: None,
7076            },
7077        };
7078        let cases = vec![
7079            (vec![text("result [CODESWARM:NEXT:3]\n ")], 2),
7080            (
7081                vec![text("result [CODESWARM:"), text("NEXT:"), text("3]")],
7082                2,
7083            ),
7084            (vec![text("result [CODESWARM:NEXT:3]"), text(" more")], 1),
7085            (vec![text("result [CODESWARM:NEXT:3]"), thought()], 1),
7086            (
7087                vec![text("result [CODESWARM:NEXT:3]"), tool(), text(" ")],
7088                1,
7089            ),
7090            (
7091                vec![text("result [CODESWARM:NEXT:"), thought(), text("3]")],
7092                1,
7093            ),
7094            (
7095                vec![text("result"), thought(), text("[CODESWARM:NEXT:3]")],
7096                2,
7097            ),
7098            (vec![text("result [CODESWARM:NEXT:1]")], 1),
7099            (vec![text("result [CODESWARM:NEXT:0]")], 1),
7100            (vec![text("result [CODESWARM:NEXT:99]")], 1),
7101            (
7102                vec![
7103                    text("result [CODESWARM:NEXT:3]"),
7104                    AgentEvent::UsageUpdated {
7105                        slot: 0,
7106                        usage: crate::UsageUpdate { used: 1, size: 100 },
7107                    },
7108                ],
7109                2,
7110            ),
7111        ];
7112        for (mut updates, expected) in cases {
7113            updates.push(AgentEvent::TurnComplete { slot: 0 });
7114            let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates);
7115            let hosts = std::iter::once(AdapterHost::new(Box::new(first), None))
7116                .chain((1..3).map(|slot| {
7117                    AdapterHost::new(
7118                        Box::new(ScriptedAdapter::new(
7119                            slot,
7120                            AgentCapabilities::default(),
7121                            [AgentEvent::TurnComplete { slot }],
7122                        )),
7123                        None,
7124                    )
7125                }))
7126                .collect();
7127            let mut relay = RelayHost::new(hosts, 10).unwrap();
7128            relay.set_roster_names(vec!["Worker".into(), "Codex".into(), "Codex".into()]);
7129            let events = Arc::new(std::sync::Mutex::new(Vec::new()));
7130            let captured = events.clone();
7131            relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
7132            relay.start().await.unwrap();
7133            relay.run_turn("task", 0).await.unwrap();
7134            let prompt = &relay.dispatches()[0].1;
7135            assert!(prompt.contains("[CODESWARM:NEXT:2] → Codex"));
7136            assert!(prompt.contains("[CODESWARM:NEXT:3] → Codex"));
7137            assert!(!prompt.contains("[CODESWARM:NEXT:1]"));
7138            // The new names must appear even when introduction has already run.
7139            relay.introduced.fill(true);
7140            relay.set_roster_names(vec!["Replacement".into(), "Codex".into(), "Codex".into()]);
7141            let next = relay.run_turn("", 0).await.unwrap();
7142            assert!(
7143                matches!(next, RelayDecision::Dispatch { slot, can_stop: false, .. } if slot == expected),
7144                "{next:?}"
7145            );
7146            let prompt = &relay.dispatches()[1].1;
7147            assert!(prompt.contains("[CODESWARM:NEXT:1] → Replacement"));
7148            assert!(prompt.contains("result"));
7149            let public = prompt
7150                .split("Public updates:\n")
7151                .nth(1)
7152                .unwrap()
7153                .split("\n\nDo not use")
7154                .next()
7155                .unwrap();
7156            assert!(!public.contains("[CODESWARM:NEXT:"));
7157            let visible = events
7158                .lock()
7159                .unwrap()
7160                .iter()
7161                .filter_map(|event| match event {
7162                    AgentEvent::Text { text, .. } => Some(text.clone()),
7163                    _ => None,
7164                })
7165                .collect::<String>();
7166            assert!(visible.contains("result"));
7167            assert!(!visible.contains("[CODESWARM:"), "{visible}");
7168        }
7169    }
7170
7171    #[tokio::test]
7172    async fn reviewer_stop_requires_a_terminal_marker_after_all_activity() {
7173        let text = |value: &str| AgentEvent::Text {
7174            slot: 1,
7175            text: value.into(),
7176        };
7177        let thought = || AgentEvent::Thought {
7178            slot: 1,
7179            text: "still checking".into(),
7180        };
7181        let tool = || AgentEvent::Tool {
7182            slot: 1,
7183            update: crate::ToolUpdate {
7184                id: "read".into(),
7185                title: "Read file".into(),
7186                status: ToolStatus::Running,
7187                detail: None,
7188            },
7189        };
7190        let cases = vec![
7191            (vec![text(&format!("done {STOP_TOKEN}"))], true),
7192            (vec![text(STOP_TOKEN), text("\n  ")], true),
7193            (vec![text(STOP_TOKEN), text(" actually keep going")], false),
7194            (vec![text(STOP_TOKEN), thought()], false),
7195            (vec![text(STOP_TOKEN), tool()], false),
7196            (vec![text(STOP_TOKEN), tool(), text(" ")], false),
7197            (vec![text(STOP_TOKEN), tool(), text(STOP_TOKEN)], true),
7198            (vec![text("[CODESWARM:"), text("STOP]")], true),
7199            (vec![text("[CODESWARM:"), thought(), text("STOP]")], false),
7200            (
7201                vec![AgentEvent::Thought {
7202                    slot: 1,
7203                    text: STOP_TOKEN.into(),
7204                }],
7205                false,
7206            ),
7207            (
7208                vec![
7209                    text(STOP_TOKEN),
7210                    AgentEvent::UsageUpdated {
7211                        slot: 1,
7212                        usage: crate::UsageUpdate { used: 1, size: 100 },
7213                    },
7214                ],
7215                true,
7216            ),
7217        ];
7218        for (mut events, stop) in cases {
7219            let first = ScriptedAdapter::new(
7220                0,
7221                AgentCapabilities::default(),
7222                [
7223                    AgentEvent::Text {
7224                        slot: 0,
7225                        text: "initial response".into(),
7226                    },
7227                    AgentEvent::TurnComplete { slot: 0 },
7228                    AgentEvent::TurnComplete { slot: 0 },
7229                ],
7230            );
7231            events.push(AgentEvent::TurnComplete { slot: 1 });
7232            let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), events.clone());
7233            let mut relay = RelayHost::new(
7234                vec![
7235                    AdapterHost::new(Box::new(first), None),
7236                    AdapterHost::new(Box::new(reviewer), None),
7237                ],
7238                4,
7239            )
7240            .unwrap();
7241            relay.start().await.unwrap();
7242            relay.run_turn("task", 0).await.unwrap();
7243            relay.run_turn("", 0).await.unwrap();
7244            let next = relay.run_turn("", 0).await.unwrap();
7245            assert_eq!(
7246                matches!(next, RelayDecision::Complete),
7247                stop,
7248                "events={events:?}"
7249            );
7250            relay.stop().await.unwrap();
7251        }
7252    }
7253
7254    #[tokio::test]
7255    async fn stop_token_is_filtered_from_streamed_ui_events() {
7256        let first = ScriptedAdapter::new(
7257            0,
7258            AgentCapabilities::default(),
7259            [
7260                AgentEvent::Text {
7261                    slot: 0,
7262                    text: format!("visible {STOP_TOKEN} trailing"),
7263                },
7264                AgentEvent::TurnComplete { slot: 0 },
7265            ],
7266        );
7267        let reviewer = ScriptedAdapter::new(
7268            1,
7269            AgentCapabilities::default(),
7270            [AgentEvent::TurnComplete { slot: 1 }],
7271        );
7272        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7273        let captured = std::sync::Arc::clone(&events);
7274        let mut relay = RelayHost::new(
7275            vec![
7276                AdapterHost::new(Box::new(first), None),
7277                AdapterHost::new(Box::new(reviewer), None),
7278            ],
7279            2,
7280        )
7281        .expect("relay");
7282        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7283        relay.start().await.expect("start");
7284        relay.run_turn("task", 0).await.expect("turn");
7285        let captured = events.lock().expect("lock");
7286        assert!(captured.iter().all(|event| match event {
7287            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7288            _ => true,
7289        }));
7290        let visible = captured
7291            .iter()
7292            .filter_map(|event| match event {
7293                AgentEvent::Text { text, .. } => Some(text.as_str()),
7294                _ => None,
7295            })
7296            .collect::<String>();
7297        assert_eq!(visible, "visible  trailing");
7298    }
7299
7300    #[tokio::test]
7301    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
7302        let first = ScriptedAdapter::new(
7303            0,
7304            AgentCapabilities::default(),
7305            [
7306                AgentEvent::Text {
7307                    slot: 0,
7308                    text: "done".into(),
7309                },
7310                AgentEvent::TurnComplete { slot: 0 },
7311            ],
7312        );
7313        let reviewer = ScriptedAdapter::new(
7314            1,
7315            AgentCapabilities::default(),
7316            [
7317                AgentEvent::Text {
7318                    slot: 1,
7319                    text: STOP_TOKEN.into(),
7320                },
7321                AgentEvent::TurnComplete { slot: 1 },
7322            ],
7323        );
7324        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7325        let captured = std::sync::Arc::clone(&events);
7326        let mut relay = RelayHost::new(
7327            vec![
7328                AdapterHost::new(Box::new(first), None),
7329                AdapterHost::new(Box::new(reviewer), None),
7330            ],
7331            4,
7332        )
7333        .expect("relay");
7334        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7335        relay.start().await.expect("start");
7336        relay.run_turn("task", 0).await.expect("first turn");
7337        relay.run_turn("", 0).await.expect("review turn");
7338        let captured = events.lock().expect("lock");
7339        assert!(captured.iter().any(|event| {
7340            matches!(
7341                event,
7342                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7343            )
7344        }));
7345        assert!(captured.iter().all(|event| match event {
7346            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
7347            _ => true,
7348        }));
7349        let acknowledgment = captured
7350            .iter()
7351            .position(|event| {
7352                matches!(
7353                    event,
7354                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
7355                )
7356            })
7357            .expect("visible acknowledgment");
7358        let completion = captured
7359            .iter()
7360            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
7361            .expect("reviewer completion");
7362        assert!(acknowledgment < completion);
7363    }
7364
7365    #[tokio::test]
7366    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
7367        let first = ScriptedAdapter::new(
7368            0,
7369            AgentCapabilities::default(),
7370            [
7371                AgentEvent::Text {
7372                    slot: 0,
7373                    text: "done".into(),
7374                },
7375                AgentEvent::TurnComplete { slot: 0 },
7376            ],
7377        );
7378        let reviewer = ScriptedAdapter::new(
7379            1,
7380            AgentCapabilities::default(),
7381            [
7382                AgentEvent::Text {
7383                    slot: 1,
7384                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
7385                },
7386                AgentEvent::TurnComplete { slot: 1 },
7387            ],
7388        );
7389        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7390        let captured = std::sync::Arc::clone(&events);
7391        let mut relay = RelayHost::new(
7392            vec![
7393                AdapterHost::new(Box::new(first), None),
7394                AdapterHost::new(Box::new(reviewer), None),
7395            ],
7396            4,
7397        )
7398        .expect("relay");
7399        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7400        relay.start().await.expect("start");
7401        relay.run_turn("task", 0).await.expect("first turn");
7402        relay.run_turn("", 0).await.expect("review turn");
7403
7404        let visible = events
7405            .lock()
7406            .expect("lock")
7407            .iter()
7408            .filter_map(|event| match event {
7409                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
7410                _ => None,
7411            })
7412            .collect::<String>();
7413        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
7414        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
7415    }
7416
7417    #[tokio::test]
7418    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
7419        let first = AdapterHost::new(
7420            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
7421            None,
7422        );
7423        let second = AdapterHost::new(
7424            Box::new(ScriptedAdapter::new(
7425                1,
7426                AgentCapabilities::default(),
7427                [AgentEvent::TurnComplete { slot: 1 }],
7428            )),
7429            None,
7430        );
7431        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7432        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
7433        relay.set_event_sink(move |event| {
7434            if matches!(event, AgentEvent::Permission { .. }) {
7435                let _ = seen_sender.send(());
7436            }
7437        });
7438        relay.start().await.expect("start");
7439        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
7440        let answer = async move {
7441            seen_receiver.recv().await.expect("permission request");
7442            sender
7443                .send(super::RelayPermissionAnswer {
7444                    slot: 0,
7445                    request_id: "permission-1".into(),
7446                    answer: PermissionAnswer::Selected {
7447                        option_id: "allow".into(),
7448                    },
7449                })
7450                .expect("queue permission answer");
7451        };
7452        tokio::time::timeout(std::time::Duration::from_millis(100), async {
7453            let ((), result) = tokio::join!(
7454                answer,
7455                relay.run_turn_with_permissions("task", 0, &mut receiver)
7456            );
7457            result
7458        })
7459        .await
7460        .expect("permission-gated turn should not deadlock")
7461        .expect("turn completes");
7462    }
7463
7464    #[tokio::test]
7465    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
7466        let first = AdapterHost::new(
7467            Box::new(PendingAdapter {
7468                slot: 0,
7469                hang_on_cancel: false,
7470            }),
7471            None,
7472        );
7473        let second = AdapterHost::new(
7474            Box::new(ScriptedAdapter::new(
7475                1,
7476                AgentCapabilities::default(),
7477                [AgentEvent::TurnComplete { slot: 1 }],
7478            )),
7479            None,
7480        );
7481        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7482        relay.start().await.expect("start");
7483        let cancellation = relay.cancellation();
7484        let error = {
7485            let turn = relay.run_turn("task", 0);
7486            tokio::pin!(turn);
7487            cancellation.request();
7488            turn.await.expect_err("cancellation should stop turn")
7489        };
7490        assert!(error.to_string().contains("relay turn cancelled"));
7491
7492        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
7493        relay
7494            .run_turn("", 1)
7495            .await
7496            .expect("replacement job reaches the selected peer");
7497        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
7498        assert!(replacement.contains("replacement job"));
7499        assert!(replacement.contains("User "));
7500        assert!(replacement.contains(":\ntask"));
7501        let owner_updates = relay.relay_mut().unseen_context(0);
7502        assert!(owner_updates.contains("User "));
7503        assert!(owner_updates.contains(":\ntask"));
7504        assert!(owner_updates.contains(":\nreplacement job"));
7505    }
7506
7507    #[tokio::test]
7508    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
7509        let first = AdapterHost::new(
7510            Box::new(PendingAdapter {
7511                slot: 0,
7512                hang_on_cancel: true,
7513            }),
7514            None,
7515        );
7516        let second = AdapterHost::new(
7517            Box::new(PendingAdapter {
7518                slot: 1,
7519                hang_on_cancel: false,
7520            }),
7521            None,
7522        );
7523        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7524        relay.start().await.expect("start");
7525        let cancellation = relay.cancellation();
7526        let turn = relay.run_turn("task", 0);
7527        tokio::pin!(turn);
7528        cancellation.request();
7529        let error = turn.await.expect_err("cancellation should stop turn");
7530        assert!(error.to_string().contains("timed out"));
7531    }
7532
7533    #[tokio::test]
7534    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
7535        let event = [AgentEvent::TurnComplete { slot: 0 }];
7536        let first = AdapterHost::new(
7537            Box::new(ScriptedAdapter::new(
7538                0,
7539                AgentCapabilities::default(),
7540                event.clone(),
7541            )),
7542            None,
7543        );
7544        let second = AdapterHost::new(
7545            Box::new(ScriptedAdapter::new(
7546                1,
7547                AgentCapabilities::default(),
7548                [AgentEvent::TurnComplete { slot: 1 }],
7549            )),
7550            None,
7551        );
7552        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7553        relay.start().await.expect("start");
7554
7555        relay.pause();
7556        assert_eq!(
7557            relay.run_turn("paused", 0).await.expect("paused turn"),
7558            crate::relay::RelayDecision::Paused
7559        );
7560        assert!(relay.dispatches().is_empty());
7561
7562        relay.resume();
7563        relay.relay_mut().drop_agent(1).expect("drop reviewer");
7564        assert!(matches!(
7565            relay
7566                .run_turn("solo follow-up", 0)
7567                .await
7568                .expect("solo turn"),
7569            crate::relay::RelayDecision::Dispatch {
7570                slot: 0,
7571                can_stop: false,
7572                ..
7573            }
7574        ));
7575        assert_eq!(relay.dispatches().len(), 1);
7576    }
7577
7578    #[tokio::test]
7579    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
7580        let first = AdapterHost::new(
7581            Box::new(ScriptedAdapter::new(
7582                0,
7583                AgentCapabilities::default(),
7584                [AgentEvent::TurnComplete { slot: 0 }],
7585            )),
7586            None,
7587        );
7588        let second = AdapterHost::new(
7589            Box::new(ScriptedAdapter::new(
7590                1,
7591                AgentCapabilities::default(),
7592                [AgentEvent::TurnComplete { slot: 1 }],
7593            )),
7594            None,
7595        );
7596        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7597        relay.set_roster_names(vec!["First".into(), "Second".into()]);
7598        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
7599        relay.set_roster_launch_specs(vec![
7600            ("custom".into(), "owner".into()),
7601            ("custom".into(), "peer".into()),
7602        ]);
7603        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7604        let captured = std::sync::Arc::clone(&events);
7605        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7606        relay.start().await.expect("start");
7607        let slot = relay
7608            .add_agent(
7609                AdapterHost::new(
7610                    Box::new(ScriptedAdapter::new(
7611                        2,
7612                        AgentCapabilities::default(),
7613                        [AgentEvent::TurnComplete { slot: 2 }],
7614                    )),
7615                    None,
7616                ),
7617                "Reviewer",
7618                "reviewer.example",
7619                "reviewer --acp",
7620            )
7621            .await
7622            .expect("append agent");
7623        assert_eq!(slot, 2);
7624        assert_eq!(
7625            relay.relay().active_slots().collect::<Vec<_>>(),
7626            vec![0, 1, 2]
7627        );
7628        assert_eq!(
7629            relay
7630                .session_metadata()
7631                .get("agents")
7632                .and_then(|value| value.as_array())
7633                .map(Vec::len),
7634            Some(3)
7635        );
7636        relay.drop_agent(1).await.expect("drop middle peer");
7637        let metadata = relay.session_metadata();
7638        assert_eq!(
7639            metadata.get("agents"),
7640            Some(&serde_json::json!([
7641                {"slot": 0, "name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
7642                {"slot": 2, "name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
7643            ]))
7644        );
7645        assert!(
7646            events
7647                .lock()
7648                .expect("lock")
7649                .iter()
7650                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
7651        );
7652    }
7653
7654    #[tokio::test]
7655    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
7656        let path = unique_test_path("codeswarm-session-metadata", "json");
7657        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
7658        let writer = metadata_store.buffered().expect("metadata writer");
7659        let first = AdapterHost::new(
7660            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
7661            None,
7662        );
7663        let second = AdapterHost::new(
7664            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
7665            None,
7666        );
7667        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7668        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7669        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
7670        relay.set_roster_launch_specs(vec![
7671            ("custom".into(), "claude".into()),
7672            ("custom".into(), "codex".into()),
7673        ]);
7674        relay.set_session_metadata_writer(writer);
7675        relay.start().await.expect("start");
7676        relay.drop_agent(0).await.expect("drop first agent");
7677        relay.stop().await.expect("stop");
7678
7679        let loaded = metadata_store
7680            .read()
7681            .expect("read metadata")
7682            .expect("metadata snapshot");
7683        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
7684        assert_eq!(
7685            loaded.get("agents"),
7686            Some(&serde_json::json!([{
7687                "slot": 1, "name": "Codex", "identity": "openai.com", "protocol": "custom",
7688                "command": "codex", "supports_load_session": false
7689            }]))
7690        );
7691        assert!(loaded.get("owner").is_none());
7692        let _ = std::fs::remove_file(path);
7693    }
7694
7695    #[tokio::test]
7696    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
7697        let first = AdapterHost::new(
7698            Box::new(ScriptedAdapter::new(
7699                0,
7700                AgentCapabilities::default(),
7701                [
7702                    AgentEvent::Text {
7703                        slot: 0,
7704                        text: "owner stream".into(),
7705                    },
7706                    AgentEvent::TurnComplete { slot: 0 },
7707                ],
7708            )),
7709            None,
7710        );
7711        let second = AdapterHost::new(
7712            Box::new(ScriptedAdapter::new(
7713                1,
7714                AgentCapabilities::default(),
7715                [
7716                    AgentEvent::Text {
7717                        slot: 1,
7718                        text: "peer stream".into(),
7719                    },
7720                    AgentEvent::TurnComplete { slot: 1 },
7721                ],
7722            )),
7723            None,
7724        );
7725        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7726        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
7727        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
7728        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
7729        let captured = std::sync::Arc::clone(&events);
7730        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
7731        relay.start().await.expect("start");
7732
7733        relay.swap_agents(0, 1).expect("swap peers");
7734        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
7735        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
7736        relay.run_turn("task", 0).await.expect("swapped turn");
7737        let events = events.lock().expect("events");
7738        assert!(events.iter().any(|event| {
7739            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
7740        }));
7741        assert!(relay.dispatches()[0].1.contains("You are Peer"));
7742    }
7743
7744    #[tokio::test]
7745    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
7746        let path = unique_test_path("codeswarm-session-metadata", "json");
7747        let first = AdapterHost::new(
7748            Box::new(ScriptedAdapter::new(
7749                0,
7750                AgentCapabilities::default(),
7751                [AgentEvent::TurnComplete { slot: 0 }],
7752            )),
7753            None,
7754        );
7755        let second = AdapterHost::new(
7756            Box::new(ScriptedAdapter::new(
7757                1,
7758                AgentCapabilities::default(),
7759                [AgentEvent::TurnComplete { slot: 1 }],
7760            )),
7761            None,
7762        );
7763        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
7764        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
7765        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
7766        relay.set_roster_launch_specs(vec![
7767            ("custom".into(), "claude".into()),
7768            ("custom".into(), "codex".into()),
7769        ]);
7770        let writer = SessionMetadataStore::open(&path)
7771            .buffered()
7772            .expect("metadata writer");
7773        relay.set_session_metadata_writer(writer);
7774        relay.start().await.expect("start");
7775        relay.stop().await.expect("stop");
7776        let loaded = SessionMetadataStore::open(&path)
7777            .read()
7778            .expect("read metadata")
7779            .expect("metadata snapshot");
7780        let agents = loaded
7781            .get("agents")
7782            .and_then(|value| value.as_array())
7783            .expect("agents");
7784        assert_eq!(agents.len(), 2);
7785        assert_eq!(agents[0]["identity"], "claude.com");
7786        assert_eq!(agents[1]["identity"], "openai.com");
7787        let _ = std::fs::remove_file(path);
7788    }
7789
7790    #[tokio::test]
7791    async fn relay_host_routes_unseen_public_context_to_next_agent() {
7792        let first = AdapterHost::new(
7793            Box::new(ScriptedAdapter::new(
7794                0,
7795                AgentCapabilities::default(),
7796                [
7797                    AgentEvent::Text {
7798                        slot: 0,
7799                        text: "implemented the fix".into(),
7800                    },
7801                    AgentEvent::TurnComplete { slot: 0 },
7802                ],
7803            )),
7804            None,
7805        );
7806        let second = AdapterHost::new(
7807            Box::new(ScriptedAdapter::new(
7808                1,
7809                AgentCapabilities::default(),
7810                [AgentEvent::TurnComplete { slot: 1 }],
7811            )),
7812            None,
7813        );
7814        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
7815        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
7816        relay.start().await.expect("start");
7817        relay.run_turn("task", 0).await.expect("first turn");
7818        relay.run_turn("review this", 0).await.expect("review turn");
7819
7820        assert_eq!(relay.dispatches().len(), 2);
7821        assert_eq!(relay.dispatches()[0].0, 0);
7822        assert!(relay.dispatches()[0].1.contains("task"));
7823        assert!(relay.dispatches()[0].1.contains("You are Codex"));
7824        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
7825        assert_eq!(relay.dispatches()[1].0, 1);
7826        assert!(relay.dispatches()[1].1.contains("review this"));
7827        let public = relay.dispatches()[1]
7828            .1
7829            .split_once("Public updates:\n")
7830            .map(|(_, updates)| updates)
7831            .expect("review receives public context");
7832        let header = public
7833            .lines()
7834            .find(|line| line.starts_with("Codex "))
7835            .expect("named previous agent");
7836        let timestamp = header
7837            .strip_prefix("Codex ")
7838            .and_then(|value| value.strip_suffix(':'))
7839            .expect("timestamped header");
7840        assert_eq!(timestamp.len(), 5);
7841        assert_eq!(timestamp.as_bytes()[2], b':');
7842        assert!(
7843            timestamp
7844                .bytes()
7845                .enumerate()
7846                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
7847        );
7848        assert!(public.contains("implemented the fix"));
7849        assert!(!public.contains("Agent 0"));
7850    }
7851}