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