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        is_usage_limit_response, strip_stop_token,
24    },
25    resources,
26};
27use async_trait::async_trait;
28use base64::{Engine, engine::general_purpose::STANDARD as BASE64};
29use serde_json::Value;
30use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncWriteExt, BufReader};
31use tokio::process::{Child, ChildStdout, Command};
32use tokio::sync::{Mutex as AsyncMutex, Notify, mpsc};
33
34pub type AdapterResult<T> = Result<T, AdapterError>;
35
36/// Keep a peer from allocating unbounded memory for one newline-delimited
37/// protocol frame. The Python client applied the same boundary (10 MiB) to
38/// its asyncio reader; ACP messages are normally much smaller than this.
39const MAX_ACP_LINE_BYTES: usize = 10 * 1024 * 1024;
40const MAX_FILE_READ_BYTES: usize = 4 * 1024 * 1024;
41const MAX_TERMINAL_OUTPUT_BYTES: usize = 1024 * 1024;
42
43#[derive(Clone, Debug)]
44struct TerminalProcess {
45    child: Arc<AsyncMutex<Option<Child>>>,
46    output: Arc<Mutex<Vec<u8>>>,
47    truncated: Arc<AtomicBool>,
48    output_readers: Arc<AtomicUsize>,
49}
50
51impl TerminalProcess {
52    async fn kill(&self) {
53        if let Some(child) = self.child.lock().await.as_mut() {
54            #[cfg(unix)]
55            if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
56                tokio::time::sleep(std::time::Duration::from_millis(100)).await;
57                signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
58            }
59            let _ = child.start_kill();
60        }
61    }
62
63    async fn stop(&self) {
64        if let Some(mut child) = self.child.lock().await.take() {
65            let _ = terminate_child(&mut child).await;
66        }
67    }
68
69    async fn wait(&self) -> Option<i32> {
70        loop {
71            let code = {
72                let mut child = self.child.lock().await;
73                match child.as_mut() {
74                    None => Some(-1),
75                    Some(child) => match child.try_wait() {
76                        Ok(Some(status)) => Some(status.code().unwrap_or(-1)),
77                        Ok(None) => None,
78                        // A terminal whose child handle can no longer be
79                        // polled must not leave `wait_for_exit` spinning
80                        // forever. Preserve the protocol's integer exit
81                        // contract with an unknown/failure sentinel.
82                        Err(_) => Some(-1),
83                    },
84                }
85            };
86            if code.is_some() {
87                while self.output_readers.load(Ordering::Acquire) != 0 {
88                    tokio::task::yield_now().await;
89                }
90                return code;
91            }
92            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
93        }
94    }
95
96    async fn exit_code(&self) -> Option<i32> {
97        let mut child = self.child.lock().await;
98        child
99            .as_mut()
100            .and_then(|child| child.try_wait().ok().flatten())
101            .map(|status| status.code().unwrap_or(-1))
102    }
103}
104
105#[derive(Clone, Debug)]
106pub struct HostUpdate {
107    pub event: AgentEvent,
108    pub effects: Vec<Effect>,
109}
110
111#[derive(Clone, Debug, Eq, PartialEq)]
112pub enum AdapterError {
113    Unsupported(&'static str),
114    Spawn(String),
115    Transport(String),
116    Protocol(String),
117}
118
119impl std::fmt::Display for AdapterError {
120    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
121        match self {
122            Self::Unsupported(operation) => write!(formatter, "unsupported operation: {operation}"),
123            Self::Spawn(error) => write!(formatter, "unable to launch agent: {error}"),
124            Self::Transport(error) => write!(formatter, "agent transport error: {error}"),
125            Self::Protocol(error) => write!(formatter, "agent protocol error: {error}"),
126        }
127    }
128}
129
130impl std::error::Error for AdapterError {}
131
132fn floor_char_boundary(text: &str, index: usize) -> usize {
133    let mut index = index.min(text.len());
134    while index > 0 && !text.is_char_boundary(index) {
135        index -= 1;
136    }
137    index
138}
139
140/// A shell-free argv parser for commands stored in the agent catalog.
141///
142/// Agent commands are configuration data, not shell snippets: expansion and
143/// pipelines are intentionally unsupported. We still accept the quoting users
144/// expect when entering a command (`'...'`, `"..."`, and backslash escapes),
145/// so a configured executable or argument containing whitespace is passed to
146/// `Command` as one argument and cannot be re-split accidentally.
147#[derive(Clone, Debug, Eq, PartialEq)]
148pub enum CommandParseError {
149    Empty,
150    UnterminatedQuote,
151    TrailingEscape,
152}
153
154impl std::fmt::Display for CommandParseError {
155    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
156        let message = match self {
157            Self::Empty => "command is empty",
158            Self::UnterminatedQuote => "command contains an unterminated quote",
159            Self::TrailingEscape => "command ends with an incomplete escape",
160        };
161        formatter.write_str(message)
162    }
163}
164
165impl std::error::Error for CommandParseError {}
166
167/// Parse a configured command into an executable and argv without invoking a
168/// shell. This is shared by native and ACP adapters.
169pub fn parse_command_line(command: &str) -> Result<(String, Vec<String>), CommandParseError> {
170    let mut argv = Vec::new();
171    let mut argument = String::new();
172    let mut quoted = None;
173    let mut escaped = false;
174    let mut started = false;
175
176    for character in command.chars() {
177        if escaped {
178            argument.push(character);
179            escaped = false;
180            started = true;
181            continue;
182        }
183        match (quoted, character) {
184            (_, '\\') if quoted != Some('\'') => {
185                escaped = true;
186                started = true;
187            }
188            (None, '\'' | '"') => {
189                quoted = Some(character);
190                started = true;
191            }
192            (Some(quote), character) if character == quote => quoted = None,
193            (None, character) if character.is_whitespace() => {
194                if started {
195                    argv.push(std::mem::take(&mut argument));
196                    started = false;
197                }
198            }
199            (_, character) => {
200                argument.push(character);
201                started = true;
202            }
203        }
204    }
205
206    if escaped {
207        return Err(CommandParseError::TrailingEscape);
208    }
209    if quoted.is_some() {
210        return Err(CommandParseError::UnterminatedQuote);
211    }
212    if started {
213        argv.push(argument);
214    }
215    let Some((program, args)) = argv.split_first() else {
216        return Err(CommandParseError::Empty);
217    };
218    Ok((program.clone(), args.to_vec()))
219}
220
221/// Kill and reap a child process. Tokio intentionally does not reap a child
222/// when its handle is dropped, so every adapter shutdown path must await this
223/// helper before releasing the handle.
224async fn terminate_child(child: &mut Child) -> AdapterResult<()> {
225    #[cfg(unix)]
226    if signal_isolated_process_group(child, nix::sys::signal::Signal::SIGTERM) {
227        tokio::time::sleep(std::time::Duration::from_millis(100)).await;
228        signal_isolated_process_group(child, nix::sys::signal::Signal::SIGKILL);
229    }
230    let kill_error = child.start_kill().err();
231    let wait_error = child.wait().await.err();
232    if let Some(error) = kill_error.or(wait_error) {
233        return Err(AdapterError::Transport(error.to_string()));
234    }
235    Ok(())
236}
237
238#[cfg(unix)]
239fn signal_isolated_process_group(child: &Child, signal: nix::sys::signal::Signal) -> bool {
240    use nix::{
241        sys::signal::killpg,
242        unistd::{Pid, getpgid, getpgrp},
243    };
244
245    let Some(raw_pid) = child.id().and_then(|pid| i32::try_from(pid).ok()) else {
246        return false;
247    };
248    let pid = Pid::from_raw(raw_pid);
249    // `process_group(0)` makes the child its own group leader. Verify that
250    // invariant at signal time and also reject CodeSwarm's own group. This
251    // keeps descendant cleanup without any possibility of reaching the
252    // containing shell or terminal multiplexer.
253    if getpgid(Some(pid)).ok() == Some(pid) && pid != getpgrp() {
254        let _ = killpg(pid, signal);
255        true
256    } else {
257        false
258    }
259}
260
261fn isolate_process_group(command: &mut Command) {
262    #[cfg(unix)]
263    command.process_group(0);
264}
265
266/// Drain diagnostics without allowing a noisy peer to block on stderr or
267/// retain an unbounded failure report.
268async fn drain_bounded<R>(mut reader: R, limit: usize) -> String
269where
270    R: AsyncRead + Unpin,
271{
272    let mut bytes = Vec::new();
273    let mut chunk = [0_u8; 4096];
274    while let Ok(count) = reader.read(&mut chunk).await {
275        if count == 0 {
276            break;
277        }
278        bytes.extend_from_slice(&chunk[..count]);
279        if bytes.len() > limit {
280            let keep_from = bytes.len() - limit;
281            bytes.drain(..keep_from);
282        }
283    }
284    String::from_utf8_lossy(&bytes).trim().to_owned()
285}
286
287/// Read one JSONL frame without allowing a peer to allocate an unbounded
288/// string before the size check runs. `AsyncBufReadExt::read_line` checks only
289/// after it has appended the complete line, so the bounded fill/consume loop
290/// below is intentional.
291async fn read_bounded_line<R>(reader: &mut R) -> AdapterResult<String>
292where
293    R: AsyncBufRead + Unpin,
294{
295    let mut bytes = Vec::with_capacity(4096);
296    loop {
297        let buffer = reader
298            .fill_buf()
299            .await
300            .map_err(|error| AdapterError::Transport(error.to_string()))?;
301        if buffer.is_empty() {
302            if bytes.is_empty() {
303                return Err(AdapterError::Transport("ACP stream closed".into()));
304            }
305            break;
306        }
307        let newline = buffer.iter().position(|byte| *byte == b'\n');
308        let available = newline.map_or(buffer.len(), |index| index + 1);
309        let remaining = MAX_ACP_LINE_BYTES
310            .saturating_add(1)
311            .saturating_sub(bytes.len());
312        if available > remaining {
313            reader.consume(remaining);
314            return Err(AdapterError::Protocol(format!(
315                "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
316            )));
317        }
318        bytes.extend_from_slice(&buffer[..available]);
319        reader.consume(available);
320        if newline.is_some() {
321            break;
322        }
323    }
324    if bytes.len() > MAX_ACP_LINE_BYTES {
325        return Err(AdapterError::Protocol(format!(
326            "ACP protocol line exceeds {MAX_ACP_LINE_BYTES} bytes"
327        )));
328    }
329    String::from_utf8(bytes).map_err(|error| AdapterError::Protocol(error.to_string()))
330}
331
332async fn drain_terminal_output<R>(
333    mut reader: R,
334    output: Arc<Mutex<Vec<u8>>>,
335    truncated: Arc<AtomicBool>,
336    output_readers: Arc<AtomicUsize>,
337    limit: usize,
338) where
339    R: AsyncRead + Unpin,
340{
341    let mut chunk = [0_u8; 4096];
342    while let Ok(count) = reader.read(&mut chunk).await {
343        if count == 0 {
344            break;
345        }
346        if let Ok(mut bytes) = output.lock() {
347            let remaining = limit.saturating_sub(bytes.len());
348            if count > remaining {
349                bytes.extend_from_slice(&chunk[..remaining]);
350                truncated.store(true, Ordering::Release);
351            } else {
352                bytes.extend_from_slice(&chunk[..count]);
353            }
354        }
355    }
356    output_readers.fetch_sub(1, Ordering::AcqRel);
357}
358
359/// Uniform control plane for ACP and custom command-line adapters.
360#[async_trait]
361pub trait AgentAdapter: Send {
362    fn slot(&self) -> RosterSlot;
363    /// Human-readable adapter identity used in the first-turn roster context.
364    /// Catalog-backed launchers can override this with their display name;
365    /// direct command adapters still get a deterministic fallback.
366    fn display_name(&self) -> String {
367        format!("Agent {}", self.slot().saturating_add(1))
368    }
369    /// Return the protocol session handle when this adapter has one. Custom
370    /// adapters may omit it; runtime metadata then still preserves roster
371    /// identity without claiming the session is resumable.
372    fn session_id(&self) -> Option<String> {
373        None
374    }
375    /// Stable protocol label recorded in the runtime session snapshot.
376    /// Custom adapters may override this when they expose a resumable
377    /// protocol distinct from the built-in native and ACP bridges.
378    fn protocol(&self) -> &'static str {
379        "custom"
380    }
381    fn capabilities(&self) -> AgentCapabilities;
382    async fn start(&mut self) -> AdapterResult<()>;
383    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()>;
384    async fn cancel(&mut self) -> AdapterResult<bool>;
385    async fn answer_permission(
386        &mut self,
387        request_id: String,
388        answer: PermissionAnswer,
389    ) -> AdapterResult<()>;
390    async fn set_mode(&mut self, mode: String) -> AdapterResult<()>;
391    async fn set_model(&mut self, _model: String) -> AdapterResult<()> {
392        Err(AdapterError::Unsupported("set_model"))
393    }
394    async fn reload(&mut self) -> AdapterResult<()>;
395    async fn stop(&mut self) -> AdapterResult<()>;
396    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>>;
397}
398
399/// Preserve a stable CodeSwarm roster slot when an already-running adapter is
400/// moved to another logical position (for example, a roster reorder). Native
401/// protocol implementations keep their original slot internally, while this
402/// boundary rewrites the normalized event identity seen by the reducer.
403struct SlotMappedAdapter {
404    logical_slot: RosterSlot,
405    inner: Box<dyn AgentAdapter>,
406}
407
408impl std::fmt::Debug for SlotMappedAdapter {
409    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
410        formatter
411            .debug_struct("SlotMappedAdapter")
412            .field("logical_slot", &self.logical_slot)
413            .field("inner_slot", &self.inner.slot())
414            .finish_non_exhaustive()
415    }
416}
417
418fn map_event_slot(event: AgentEvent, slot: RosterSlot) -> AgentEvent {
419    match event {
420        AgentEvent::GoalUpdated { goal } => AgentEvent::GoalUpdated { goal },
421        AgentEvent::RosterUpdated { update } => AgentEvent::RosterUpdated { update },
422        AgentEvent::Ready { capabilities, .. } => AgentEvent::Ready { slot, capabilities },
423        AgentEvent::TurnStarted { .. } => AgentEvent::TurnStarted { slot },
424        AgentEvent::ModesReplaced {
425            modes,
426            current_mode,
427            ..
428        } => AgentEvent::ModesReplaced {
429            slot,
430            modes,
431            current_mode,
432        },
433        AgentEvent::ModeUpdated { current_mode, .. } => {
434            AgentEvent::ModeUpdated { slot, current_mode }
435        }
436        AgentEvent::ModelsReplaced {
437            config_id,
438            models,
439            current_model,
440            ..
441        } => AgentEvent::ModelsReplaced {
442            slot,
443            config_id,
444            models,
445            current_model,
446        },
447        AgentEvent::ModelUpdated { current_model, .. } => AgentEvent::ModelUpdated {
448            slot,
449            current_model,
450        },
451        AgentEvent::UserText { text, .. } => AgentEvent::UserText { slot, text },
452        AgentEvent::CommandsReplaced { commands, .. } => {
453            AgentEvent::CommandsReplaced { slot, commands }
454        }
455        AgentEvent::UsageUpdated { usage, .. } => AgentEvent::UsageUpdated { slot, usage },
456        AgentEvent::Text { text, .. } => AgentEvent::Text { slot, text },
457        AgentEvent::Thought { text, .. } => AgentEvent::Thought { slot, text },
458        AgentEvent::Tool { update, .. } => AgentEvent::Tool { slot, update },
459        AgentEvent::Permission { request, .. } => AgentEvent::Permission { slot, request },
460        AgentEvent::Terminal { event, .. } => AgentEvent::Terminal { slot, event },
461        AgentEvent::TurnComplete { .. } => AgentEvent::TurnComplete { slot },
462        AgentEvent::UsageLimitReached { detail, .. } => {
463            AgentEvent::UsageLimitReached { slot, detail }
464        }
465        AgentEvent::Failed {
466            started, detail, ..
467        } => AgentEvent::Failed {
468            slot,
469            started,
470            detail,
471        },
472    }
473}
474
475#[async_trait]
476impl AgentAdapter for SlotMappedAdapter {
477    fn slot(&self) -> RosterSlot {
478        self.logical_slot
479    }
480
481    fn display_name(&self) -> String {
482        self.inner.display_name()
483    }
484
485    fn session_id(&self) -> Option<String> {
486        self.inner.session_id()
487    }
488
489    fn protocol(&self) -> &'static str {
490        self.inner.protocol()
491    }
492
493    fn capabilities(&self) -> AgentCapabilities {
494        self.inner.capabilities()
495    }
496
497    async fn start(&mut self) -> AdapterResult<()> {
498        self.inner.start().await
499    }
500
501    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
502        self.inner.send_prompt(prompt).await
503    }
504
505    async fn cancel(&mut self) -> AdapterResult<bool> {
506        self.inner.cancel().await
507    }
508
509    async fn answer_permission(
510        &mut self,
511        request_id: String,
512        answer: PermissionAnswer,
513    ) -> AdapterResult<()> {
514        self.inner.answer_permission(request_id, answer).await
515    }
516
517    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
518        self.inner.set_mode(mode).await
519    }
520
521    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
522        self.inner.set_model(model).await
523    }
524
525    async fn reload(&mut self) -> AdapterResult<()> {
526        self.inner.reload().await
527    }
528
529    async fn stop(&mut self) -> AdapterResult<()> {
530        self.inner.stop().await
531    }
532
533    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
534        let slot = self.logical_slot;
535        self.inner
536            .next_event()
537            .await
538            .map(|result| result.map(|event| map_event_slot(event, slot)))
539    }
540}
541
542/// Owns one adapter and feeds normalized events through the deterministic core
543/// reducer. The UI consumes effects and state snapshots.
544pub struct AdapterHost {
545    adapter: Box<dyn AgentAdapter>,
546    pub state: SessionState,
547    pub last_error: Option<String>,
548    event_log: Option<EventLog>,
549}
550
551impl std::fmt::Debug for AdapterHost {
552    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
553        formatter
554            .debug_struct("AdapterHost")
555            .field("state", &self.state)
556            .field("last_error", &self.last_error)
557            .field("event_log", &self.event_log)
558            .finish_non_exhaustive()
559    }
560}
561
562impl AdapterHost {
563    pub fn new(adapter: Box<dyn AgentAdapter>, event_log: Option<EventLog>) -> Self {
564        let slot = adapter.slot();
565        Self {
566            adapter,
567            state: SessionState::new(slot.saturating_add(1)),
568            last_error: None,
569            event_log,
570        }
571    }
572
573    pub async fn start(&mut self) -> AdapterResult<()> {
574        self.adapter.start().await
575    }
576
577    pub async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
578        self.adapter.send_prompt(prompt).await
579    }
580
581    pub async fn cancel(&mut self) -> AdapterResult<bool> {
582        self.adapter.cancel().await
583    }
584
585    pub async fn answer_permission(
586        &mut self,
587        request_id: String,
588        answer: PermissionAnswer,
589    ) -> AdapterResult<()> {
590        self.adapter.answer_permission(request_id, answer).await
591    }
592
593    pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
594        self.adapter.set_mode(mode).await
595    }
596
597    pub async fn set_model(&mut self, model: String) -> AdapterResult<()> {
598        self.adapter.set_model(model).await
599    }
600
601    pub async fn reload(&mut self) -> AdapterResult<()> {
602        self.adapter.reload().await?;
603        let slot = self.adapter.slot();
604        if let Some(agent) = self.state.slots.get_mut(slot) {
605            agent.active = true;
606            agent.capabilities = self.adapter.capabilities();
607        }
608        self.last_error = None;
609        Ok(())
610    }
611
612    pub async fn stop(&mut self) -> AdapterResult<()> {
613        self.adapter.stop().await
614    }
615
616    pub async fn next_effects(&mut self) -> Option<AdapterResult<Vec<Effect>>> {
617        Some(self.next_update().await?.map(|update| update.effects))
618    }
619
620    pub async fn next_update(&mut self) -> Option<AdapterResult<HostUpdate>> {
621        let event = match self.adapter.next_event().await {
622            None => return None,
623            Some(Err(error)) => {
624                self.last_error = Some(error.to_string());
625                let slot = self.adapter.slot();
626                let failure = AgentEvent::Failed {
627                    slot,
628                    started: true,
629                    detail: error.to_string(),
630                };
631                let effects = reduce(&mut self.state, failure.clone());
632                return Some(Ok(HostUpdate {
633                    event: failure,
634                    effects,
635                }));
636            }
637            Some(Ok(event)) => event,
638        };
639        if let Some(log) = &self.event_log
640            && let Err(error) = log.append(&event)
641        {
642            return Some(Err(AdapterError::Transport(error.to_string())));
643        }
644        let effects = reduce(&mut self.state, event.clone());
645        Some(Ok(HostUpdate { event, effects }))
646    }
647
648    pub fn adapter(&self) -> &dyn AgentAdapter {
649        &*self.adapter
650    }
651
652    pub fn session_id(&self) -> Option<String> {
653        self.adapter.session_id()
654    }
655
656    /// Move this host to a new logical roster slot. The adapter process is not
657    /// restarted; only normalized event identities and reducer state move.
658    fn remap(self, logical_slot: RosterSlot) -> Self {
659        let old_slot = self.adapter.slot();
660        if old_slot == logical_slot {
661            return self;
662        }
663
664        let mut state = self.state;
665        if state.slots.len() <= logical_slot {
666            state
667                .slots
668                .resize(logical_slot.saturating_add(1), Default::default());
669        }
670        if let Some(agent) = state.slots.get(old_slot).cloned() {
671            state.slots[logical_slot] = agent;
672        }
673        if state.active_slot == Some(old_slot) {
674            state.active_slot = Some(logical_slot);
675        }
676        for (slot, _) in &mut state.queued_prompts {
677            if *slot == old_slot {
678                *slot = logical_slot;
679            }
680        }
681        for (slot, _) in &mut state.public_text {
682            if *slot == old_slot {
683                *slot = logical_slot;
684            }
685        }
686
687        Self {
688            adapter: Box::new(SlotMappedAdapter {
689                logical_slot,
690                inner: self.adapter,
691            }),
692            state,
693            last_error: self.last_error,
694            event_log: self.event_log,
695        }
696    }
697}
698
699/// Sequential multi-adapter runner. It intentionally never polls two
700/// adapters concurrently: the next prompt depends on the prior response.
701pub struct RelayHost {
702    goal: Option<crate::goal::Goal>,
703    hosts: Vec<AdapterHost>,
704    relay: Relay,
705    introduced: Vec<bool>,
706    roster_names: Vec<String>,
707    roster_identities: Vec<String>,
708    roster_launch_specs: Vec<(String, String)>,
709    desired_policy: String,
710    metadata_writer: Option<BufferedSessionMetadataStore>,
711    metadata_workspace: Option<String>,
712    dispatches: Vec<(RosterSlot, String)>,
713    event_sink: Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
714    cancel_requested: Arc<AtomicBool>,
715    cancel_notify: Arc<Notify>,
716}
717
718/// A clonable signal used by a terminal control loop to interrupt the active
719/// relay turn without borrowing the relay while its adapter is being polled.
720#[derive(Clone, Debug)]
721pub struct RelayCancellation {
722    requested: Arc<AtomicBool>,
723    notify: Arc<Notify>,
724}
725
726/// A permission response delivered while a relay turn is still streaming.
727///
728/// Relay turns own the active adapter mutably, so permission answers must be
729/// consumed inside that turn rather than deferred by the outer coordinator.
730#[derive(Debug)]
731pub struct RelayPermissionAnswer {
732    pub slot: RosterSlot,
733    pub request_id: String,
734    pub answer: PermissionAnswer,
735}
736
737#[cfg(test)]
738const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
739#[cfg(not(test))]
740const CANCEL_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
741
742#[cfg(test)]
743const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(20);
744#[cfg(not(test))]
745const CANCEL_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
746
747/// Bound third-party cancellation hooks so a broken adapter cannot freeze the
748/// terminal control loop.  Native and ACP adapters normally return promptly;
749/// the timeout is for custom adapters whose cancellation future may never
750/// resolve.
751async fn cancel_with_timeout(host: &mut AdapterHost) -> AdapterResult<bool> {
752    tokio::time::timeout(CANCEL_TIMEOUT, host.cancel())
753        .await
754        .map_err(|_| AdapterError::Transport("adapter cancellation timed out".into()))?
755}
756
757fn canonical_policy_id(policy: &str) -> &str {
758    match policy {
759        "plan" => "codeswarm:mode:plan",
760        "default" | "manual" => "codeswarm:mode:manual",
761        "accept-edits" => "codeswarm:mode:accept-edits",
762        "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
763        other => other,
764    }
765}
766
767async fn apply_policy_to_host(host: &mut AdapterHost, policy: &str) -> AdapterResult<()> {
768    if !host.adapter().capabilities().supports_modes {
769        return Ok(());
770    }
771    let policy_id = canonical_policy_id(policy);
772    let slot = host.adapter().slot();
773    let advertised = host
774        .state
775        .slots
776        .get(slot)
777        .map(|agent| agent.modes.as_slice())
778        .unwrap_or_default();
779    let native = if advertised.is_empty() {
780        match policy_id {
781            "codeswarm:mode:plan" => "plan".into(),
782            "codeswarm:mode:manual" => "default".into(),
783            "codeswarm:mode:accept-edits" => "accept-edits".into(),
784            "codeswarm:mode:full-access" => "full-access".into(),
785            other => other.into(),
786        }
787    } else {
788        crate::policy::resolve(policy_id, advertised)
789            .map(|mode| mode.id)
790            .ok_or(AdapterError::Unsupported(
791                "desired policy is unavailable for adapter",
792            ))?
793    };
794    host.set_mode(native).await
795}
796
797async fn refresh_mode_catalog(
798    host: &mut AdapterHost,
799    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
800) -> AdapterResult<bool> {
801    if !host.adapter().capabilities().supports_modes {
802        return Ok(false);
803    }
804    tokio::time::timeout(std::time::Duration::from_secs(2), async {
805        let mut ready_seen = false;
806        loop {
807            let update = host.next_update().await.ok_or_else(|| {
808                AdapterError::Transport("adapter ended before advertising modes".into())
809            })??;
810            let catalog_ready = matches!(update.event, AgentEvent::ModesReplaced { .. });
811            ready_seen |= matches!(update.event, AgentEvent::Ready { .. });
812            if let Some(sink) = event_sink {
813                sink(update.event);
814            }
815            if catalog_ready {
816                return Ok(ready_seen);
817            }
818        }
819    })
820    .await
821    .map_err(|_| AdapterError::Transport("adapter mode catalog timed out".into()))?
822}
823
824async fn refresh_adapter_startup(
825    host: &mut AdapterHost,
826    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
827) -> AdapterResult<()> {
828    let ready_seen = refresh_mode_catalog(host, event_sink).await?;
829    if ready_seen || !matches!(host.adapter().protocol(), "native" | "acp") {
830        return Ok(());
831    }
832    tokio::time::timeout(std::time::Duration::from_secs(2), async {
833        loop {
834            let update = host.next_update().await.ok_or_else(|| {
835                AdapterError::Transport("adapter ended before becoming ready".into())
836            })??;
837            let ready = matches!(update.event, AgentEvent::Ready { .. });
838            if let Some(sink) = event_sink {
839                sink(update.event);
840            }
841            if ready {
842                return Ok(());
843            }
844        }
845    })
846    .await
847    .map_err(|_| AdapterError::Transport("adapter ready handshake timed out".into()))?
848}
849
850fn public_context_speaker(name: &str) -> String {
851    let now = time::OffsetDateTime::now_local().unwrap_or_else(|_| time::OffsetDateTime::now_utc());
852    format!("{name} {:02}:{:02}", now.hour(), now.minute())
853}
854
855impl RelayCancellation {
856    pub fn request(&self) {
857        self.requested.store(true, Ordering::Release);
858        self.notify.notify_one();
859    }
860}
861
862impl std::fmt::Debug for RelayHost {
863    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
864        formatter
865            .debug_struct("RelayHost")
866            .field("hosts", &self.hosts)
867            .field("relay", &self.relay)
868            .field("dispatches", &self.dispatches)
869            .field("event_sink", &self.event_sink.is_some())
870            .field(
871                "cancel_requested",
872                &self.cancel_requested.load(Ordering::Acquire),
873            )
874            .finish()
875    }
876}
877
878impl RelayHost {
879    pub fn new(hosts: Vec<AdapterHost>, max_rounds: usize) -> Result<Self, AdapterError> {
880        if hosts.is_empty() {
881            return Err(AdapterError::Unsupported("relay requires an adapter"));
882        }
883        Ok(Self {
884            relay: Relay::new(hosts.len(), max_rounds),
885            introduced: vec![false; hosts.len()],
886            roster_names: hosts
887                .iter()
888                .map(|host| host.adapter().display_name())
889                .collect(),
890            roster_identities: hosts
891                .iter()
892                .map(|host| host.adapter().display_name())
893                .collect(),
894            roster_launch_specs: Vec::new(),
895            desired_policy: crate::policy::DEFAULT_POLICY_ID.into(),
896            metadata_writer: None,
897            metadata_workspace: None,
898            hosts,
899            dispatches: Vec::new(),
900            goal: None,
901            event_sink: None,
902            cancel_requested: Arc::new(AtomicBool::new(false)),
903            cancel_notify: Arc::new(Notify::new()),
904        })
905    }
906
907    /// Send each normalized event to a client while a turn is being drained.
908    /// The callback runs synchronously on the relay task and should only
909    /// enqueue the event; expensive rendering must happen outside the callback.
910    pub fn set_event_sink<F>(&mut self, sink: F)
911    where
912        F: Fn(AgentEvent) + Send + Sync + 'static,
913    {
914        self.event_sink = Some(Arc::new(sink));
915    }
916
917    /// Install catalog-backed names for first-turn introductions. Direct
918    /// scripted/custom adapters retain deterministic `Agent N` fallbacks.
919    pub fn set_roster_names(&mut self, names: Vec<String>) {
920        if names.len() == self.hosts.len() {
921            self.roster_names = names;
922        }
923    }
924
925    /// Install catalog identities used by the launcher and persisted session
926    /// metadata. Names remain a separate display concern, so custom adapters
927    /// can retain friendly labels while still restoring by stable identity.
928    pub fn set_roster_identities(&mut self, identities: Vec<String>) {
929        if identities.len() == self.hosts.len() {
930            self.roster_identities = identities;
931        }
932    }
933
934    pub fn set_roster_launch_specs(&mut self, specs: Vec<(String, String)>) {
935        if specs.len() == self.hosts.len() {
936            self.roster_launch_specs = specs;
937        }
938    }
939
940    /// Attach a background metadata writer. Runtime changes enqueue complete
941    /// snapshots; no metadata filesystem work runs on the terminal thread.
942    pub fn set_session_metadata_writer(&mut self, writer: BufferedSessionMetadataStore) {
943        self.metadata_writer = Some(writer);
944    }
945
946    pub fn set_session_metadata_workspace(&mut self, workspace: impl Into<String>) {
947        self.metadata_workspace = Some(workspace.into());
948    }
949
950    /// Build the coordinator-owned session snapshot for the active roster.
951    pub fn session_metadata(&self) -> SessionMetadata {
952        let active = self.relay.active_slots().collect::<Vec<_>>();
953        let mut data = serde_json::Map::new();
954        data.insert(
955            "goal".into(),
956            serde_json::to_value(&self.goal).expect("goal serializes"),
957        );
958        if let Some(workspace) = &self.metadata_workspace {
959            data.insert("cwd".into(), serde_json::Value::String(workspace.clone()));
960        }
961        data.insert(
962            "title".into(),
963            serde_json::Value::String("CodeSwarm".into()),
964        );
965        data.insert(
966            "agents".into(),
967            serde_json::Value::Array(
968                active
969                    .into_iter()
970                    .filter_map(|slot| {
971                        let host = self.hosts.get(slot)?;
972                        let (protocol, command) = self.roster_launch_specs.get(slot)?;
973                        let mut agent = serde_json::Map::new();
974                        agent.insert(
975                            "name".into(),
976                            serde_json::Value::String(
977                                self.roster_names
978                                    .get(slot)
979                                    .cloned()
980                                    .unwrap_or_else(|| host.adapter().display_name()),
981                            ),
982                        );
983                        agent.insert(
984                            "identity".into(),
985                            serde_json::Value::String(
986                                self.roster_identities
987                                    .get(slot)
988                                    .cloned()
989                                    .unwrap_or_else(|| host.adapter().display_name()),
990                            ),
991                        );
992                        agent.insert(
993                            "protocol".into(),
994                            serde_json::Value::String(protocol.clone()),
995                        );
996                        agent.insert("command".into(), serde_json::Value::String(command.clone()));
997                        agent.insert(
998                            "supports_load_session".into(),
999                            serde_json::Value::Bool(
1000                                host.adapter().capabilities().supports_session_load,
1001                            ),
1002                        );
1003                        if let Some(session_id) = host.session_id() {
1004                            agent
1005                                .insert("session_id".into(), serde_json::Value::String(session_id));
1006                        }
1007                        Some(serde_json::Value::Object(agent))
1008                    })
1009                    .collect(),
1010            ),
1011        );
1012        SessionMetadata::new(data)
1013    }
1014
1015    fn queue_session_metadata(&self) -> AdapterResult<()> {
1016        if let Some(writer) = &self.metadata_writer {
1017            writer
1018                .write(self.session_metadata())
1019                .map_err(|error| AdapterError::Transport(error.to_string()))?;
1020        }
1021        Ok(())
1022    }
1023
1024    pub fn restore_goal(&mut self, goal: Option<crate::goal::Goal>) {
1025        self.goal = goal;
1026        if let Some(sink) = &self.event_sink {
1027            sink(AgentEvent::GoalUpdated {
1028                goal: self.goal.clone(),
1029            });
1030        }
1031    }
1032
1033    pub fn apply_goal(
1034        &mut self,
1035        command: crate::goal::GoalCommand,
1036    ) -> Result<Option<String>, String> {
1037        let task = crate::goal::apply(&mut self.goal, command)?;
1038        if let Some(sink) = &self.event_sink {
1039            sink(AgentEvent::GoalUpdated {
1040                goal: self.goal.clone(),
1041            });
1042        }
1043        self.queue_session_metadata()
1044            .map_err(|error| error.to_string())?;
1045        Ok(task)
1046    }
1047
1048    pub fn roster_names(&self) -> &[String] {
1049        &self.roster_names
1050    }
1051
1052    pub fn session_ids(&self) -> Vec<Option<String>> {
1053        self.hosts.iter().map(AdapterHost::session_id).collect()
1054    }
1055
1056    pub async fn start(&mut self) -> AdapterResult<()> {
1057        let event_sink = self.event_sink.clone();
1058        let startups = self.hosts.iter_mut().map(|host| {
1059            let event_sink = event_sink.clone();
1060            async move {
1061                host.start().await?;
1062                refresh_adapter_startup(host, &event_sink).await
1063            }
1064        });
1065        let results = futures::future::join_all(startups).await;
1066        if let Some(error) = results.into_iter().find_map(Result::err) {
1067            // Startup is transactional even though independent adapters are
1068            // warmed concurrently. Every host gets a cleanup attempt.
1069            for host in &mut self.hosts {
1070                let _ = host.stop().await;
1071            }
1072            return Err(error);
1073        }
1074        if let Err(error) = self
1075            .set_policy(crate::policy::DEFAULT_POLICY_ID.into())
1076            .await
1077        {
1078            for host in &mut self.hosts {
1079                let _ = host.stop().await;
1080            }
1081            return Err(error);
1082        }
1083        let _ = self.queue_session_metadata();
1084        Ok(())
1085    }
1086
1087    pub async fn stop(&mut self) -> AdapterResult<()> {
1088        // A third-party adapter can fail during shutdown (for example after
1089        // its transport has already disappeared). Always give every roster
1090        // member a chance to clean up, then return the first error so callers
1091        // still get an actionable failure without leaking later processes.
1092        let mut first_error = None;
1093        for host in &mut self.hosts {
1094            if let Err(error) = host.stop().await
1095                && first_error.is_none()
1096            {
1097                first_error = Some(error);
1098            }
1099        }
1100        if let Some(writer) = &self.metadata_writer
1101            && let Err(error) = writer.flush()
1102            && first_error.is_none()
1103        {
1104            first_error = Some(AdapterError::Transport(error.to_string()));
1105        }
1106        first_error.map_or(Ok(()), Err)
1107    }
1108
1109    /// Forward a normalized permission answer to the adapter owning `slot`.
1110    /// This keeps protocol-specific response framing out of the relay and UI.
1111    pub async fn answer_permission(
1112        &mut self,
1113        slot: RosterSlot,
1114        request_id: String,
1115        answer: PermissionAnswer,
1116    ) -> AdapterResult<()> {
1117        let host = self
1118            .hosts
1119            .get_mut(slot)
1120            .ok_or_else(|| AdapterError::Transport("permission target is missing".into()))?;
1121        host.answer_permission(request_id, answer).await
1122    }
1123
1124    pub fn pause(&mut self) {
1125        self.relay.pause();
1126    }
1127
1128    pub fn resume(&mut self) {
1129        self.relay.resume();
1130    }
1131
1132    /// Select how future non-direct turns are routed. Queued prompts and the
1133    /// shared context remain intact when a user changes this setting.
1134    pub fn set_strategy(&mut self, strategy: CollaborationStrategy) {
1135        self.relay.set_strategy(strategy);
1136    }
1137
1138    pub fn strategy(&self) -> CollaborationStrategy {
1139        self.relay.strategy()
1140    }
1141
1142    pub fn roster_identity(&self, slot: RosterSlot) -> Option<&str> {
1143        self.roster_identities.get(slot).map(String::as_str)
1144    }
1145
1146    pub fn active_slot_for_identity(&self, identity: &str) -> Option<RosterSlot> {
1147        self.relay.active_slots().find(|slot| {
1148            self.roster_identity(*slot)
1149                .is_some_and(|candidate| candidate.eq_ignore_ascii_case(identity))
1150        })
1151    }
1152
1153    /// Apply one semantic mode to every active adapter that advertises mode
1154    /// support. Native adapters may translate the policy to their own IDs.
1155    pub async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
1156        let active = self.relay.active_slots().collect::<Vec<_>>();
1157        for slot in active {
1158            let Some(host) = self.hosts.get_mut(slot) else {
1159                continue;
1160            };
1161            if host.adapter().capabilities().supports_modes {
1162                host.set_mode(mode.clone()).await?;
1163            }
1164        }
1165        Ok(())
1166    }
1167
1168    pub async fn set_model(&mut self, slot: RosterSlot, model: String) -> AdapterResult<()> {
1169        let host = self
1170            .hosts
1171            .get_mut(slot)
1172            .ok_or_else(|| AdapterError::Transport("model target is missing".into()))?;
1173        if !host.adapter().capabilities().supports_models {
1174            return Err(AdapterError::Unsupported("set_model"));
1175        }
1176        host.set_model(model.clone()).await?;
1177        if let Some(sink) = &self.event_sink {
1178            sink(AgentEvent::ModelUpdated {
1179                slot,
1180                current_model: model,
1181            });
1182        }
1183        Ok(())
1184    }
1185
1186    /// Translate a CodeSwarm semantic policy to each adapter's currently
1187    /// advertised native mode, falling back to conventional IDs before the
1188    /// first catalog update arrives.
1189    pub async fn set_policy(&mut self, policy: String) -> AdapterResult<()> {
1190        self.desired_policy = canonical_policy_id(&policy).to_owned();
1191        let desired_policy = self.desired_policy.clone();
1192        let mut first_error = None;
1193        let active = self.relay.active_slots().collect::<Vec<_>>();
1194        for active_slot in active {
1195            let Some(host) = self.hosts.get_mut(active_slot) else {
1196                continue;
1197            };
1198            if let Err(error) = apply_policy_to_host(host, &desired_policy).await
1199                && first_error.is_none()
1200            {
1201                first_error = Some(error);
1202            }
1203        }
1204        first_error.map_or(Ok(()), Err)
1205    }
1206
1207    pub async fn reload(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1208        let desired_policy = self.desired_policy.clone();
1209        let event_sink = self.event_sink.clone();
1210        let host = self
1211            .hosts
1212            .get_mut(slot)
1213            .ok_or_else(|| AdapterError::Transport("reload target is missing".into()))?;
1214        host.reload().await?;
1215        refresh_adapter_startup(host, &event_sink).await?;
1216        apply_policy_to_host(host, &desired_policy).await?;
1217        if let Some(introduced) = self.introduced.get_mut(slot) {
1218            *introduced = false;
1219        }
1220        self.relay
1221            .reactivate(slot)
1222            .map_err(|error| AdapterError::Transport(error.into()))?;
1223        let _ = self.queue_session_metadata();
1224        if let Some(sink) = &self.event_sink {
1225            sink(AgentEvent::RosterUpdated {
1226                update: RosterUpdate::Reloaded { slot },
1227            });
1228        }
1229        let _ = self.relay.clear_limited(slot);
1230        Ok(())
1231    }
1232
1233    /// Stop and tombstone an agent while preserving its stable roster slot.
1234    pub async fn drop_agent(&mut self, slot: RosterSlot) -> AdapterResult<()> {
1235        self.relay
1236            .drop_agent(slot)
1237            .map_err(|error| AdapterError::Transport(error.into()))?;
1238        let _stop_result = if let Some(host) = self.hosts.get_mut(slot) {
1239            host.stop().await
1240        } else {
1241            Ok(())
1242        };
1243        // Persist the tombstone even when a third-party process reports a
1244        // shutdown error; otherwise the next launch can resurrect a slot the
1245        // user explicitly removed.
1246        let _ = self.queue_session_metadata();
1247        if let Some(sink) = &self.event_sink {
1248            sink(AgentEvent::RosterUpdated {
1249                update: RosterUpdate::Dropped { slot },
1250            });
1251        }
1252        Ok(())
1253    }
1254
1255    /// Start and append a new adapter in the next stable roster slot. The
1256    /// adapter is started before it becomes visible to the relay, so a failed
1257    /// add cannot leave a half-active slot behind.
1258    pub async fn add_agent(
1259        &mut self,
1260        mut host: AdapterHost,
1261        name: impl Into<String>,
1262        identity: impl Into<String>,
1263        command: impl Into<String>,
1264    ) -> AdapterResult<RosterSlot> {
1265        let slot = self.hosts.len();
1266        if host.adapter().slot() != slot {
1267            return Err(AdapterError::Transport(
1268                "new adapter slot must append after the existing roster".into(),
1269            ));
1270        }
1271        if let Err(error) = host.start().await {
1272            let _ = host.stop().await;
1273            return Err(error);
1274        }
1275        if let Err(error) = refresh_adapter_startup(&mut host, &self.event_sink).await {
1276            let _ = host.stop().await;
1277            return Err(error);
1278        }
1279        if let Err(error) = apply_policy_to_host(&mut host, &self.desired_policy).await {
1280            let _ = host.stop().await;
1281            return Err(error);
1282        }
1283        let capabilities = host.adapter().capabilities();
1284        self.hosts.push(host);
1285        self.relay.add_agent();
1286        self.introduced.push(false);
1287        let name = name.into();
1288        let identity = identity.into();
1289        self.roster_names.push(name.clone());
1290        self.roster_identities.push(identity.clone());
1291        self.roster_launch_specs
1292            .push((self.hosts[slot].adapter().protocol().into(), command.into()));
1293        let _ = self.queue_session_metadata();
1294        if let Some(sink) = &self.event_sink {
1295            sink(AgentEvent::RosterUpdated {
1296                update: RosterUpdate::Added {
1297                    slot,
1298                    name,
1299                    identity,
1300                },
1301            });
1302            sink(AgentEvent::Ready { slot, capabilities });
1303        }
1304        Ok(slot)
1305    }
1306
1307    /// Reorder two live adapters without restarting either process. The
1308    /// logical slot wrapper keeps normalized events, reducer state, and
1309    /// permission routing aligned with the new roster order.
1310    pub fn swap_agents(&mut self, first: RosterSlot, second: RosterSlot) -> AdapterResult<()> {
1311        if first == second {
1312            return Ok(());
1313        }
1314        if first >= self.hosts.len() || second >= self.hosts.len() {
1315            return Err(AdapterError::Transport("roster slot out of range".into()));
1316        }
1317        self.relay
1318            .swap_agents(first, second)
1319            .map_err(|error| AdapterError::Transport(error.into()))?;
1320        let low = first.min(second);
1321        let high = first.max(second);
1322        let high_host = self.hosts.remove(high);
1323        let low_host = self.hosts.remove(low);
1324        self.hosts.insert(low, high_host.remap(low));
1325        self.hosts.insert(high, low_host.remap(high));
1326        self.roster_names.swap(first, second);
1327        self.roster_identities.swap(first, second);
1328        if self.roster_launch_specs.len() == self.hosts.len() {
1329            self.roster_launch_specs.swap(first, second);
1330        }
1331        self.introduced.swap(first, second);
1332        if let Some(sink) = &self.event_sink {
1333            sink(AgentEvent::RosterUpdated {
1334                update: RosterUpdate::Swapped { first, second },
1335            });
1336            sink(AgentEvent::Ready {
1337                slot: first,
1338                capabilities: self.hosts[first].adapter().capabilities(),
1339            });
1340            sink(AgentEvent::Ready {
1341                slot: second,
1342                capabilities: self.hosts[second].adapter().capabilities(),
1343            });
1344        }
1345        let _ = self.queue_session_metadata();
1346        Ok(())
1347    }
1348
1349    pub fn relay(&self) -> &Relay {
1350        &self.relay
1351    }
1352
1353    pub fn next_slot(&self) -> RosterSlot {
1354        self.hosts.len()
1355    }
1356
1357    pub fn relay_mut(&mut self) -> &mut Relay {
1358        &mut self.relay
1359    }
1360
1361    pub fn cancellation(&self) -> RelayCancellation {
1362        RelayCancellation {
1363            requested: Arc::clone(&self.cancel_requested),
1364            notify: Arc::clone(&self.cancel_notify),
1365        }
1366    }
1367
1368    /// Prompts sent to adapters, in causal dispatch order. This is useful to
1369    /// diagnostics and makes the context-routing boundary observable without
1370    /// exposing protocol-specific adapter internals.
1371    pub fn dispatches(&self) -> &[(RosterSlot, String)] {
1372        &self.dispatches
1373    }
1374
1375    pub async fn run_turn(
1376        &mut self,
1377        task: impl Into<String>,
1378        first_slot: RosterSlot,
1379    ) -> AdapterResult<RelayDecision> {
1380        self.run_turn_inner(task.into(), first_slot, None).await
1381    }
1382
1383    pub async fn run_turn_with_permissions(
1384        &mut self,
1385        task: impl Into<String>,
1386        first_slot: RosterSlot,
1387        permissions: &mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>,
1388    ) -> AdapterResult<RelayDecision> {
1389        self.run_turn_inner(task.into(), first_slot, Some(permissions))
1390            .await
1391    }
1392
1393    async fn run_turn_inner(
1394        &mut self,
1395        task: String,
1396        first_slot: RosterSlot,
1397        mut permissions: Option<&mut tokio::sync::mpsc::UnboundedReceiver<RelayPermissionAnswer>>,
1398    ) -> AdapterResult<RelayDecision> {
1399        let decision = self.relay.begin(task, first_slot);
1400        let RelayDecision::Dispatch {
1401            slot,
1402            prompt,
1403            direct,
1404            can_stop,
1405        } = &decision
1406        else {
1407            return Ok(decision);
1408        };
1409        let speaker_name = self
1410            .roster_names
1411            .get(*slot)
1412            .cloned()
1413            .unwrap_or_else(|| self.hosts[*slot].adapter().display_name());
1414        let unseen = self.relay.unseen_context(*slot);
1415        // A non-direct prompt with text is a new public human turn. Record it
1416        // after collecting the recipient's existing unseen context so that
1417        // the selected agent receives the prompt once, while every peer gets
1418        // it through the shared journal on its next turn. Recording before
1419        // adapter I/O also preserves the job when the selected turn is
1420        // cancelled.
1421        if !*direct && !prompt.trim().is_empty() {
1422            if self.relay.shared_task().is_none() {
1423                self.relay.set_shared_task(prompt.clone());
1424            }
1425            self.relay
1426                .record_public(public_context_speaker("User"), prompt.clone());
1427        }
1428        let prompt = if unseen.is_empty() {
1429            prompt.clone()
1430        } else {
1431            format!("{prompt}\n\nPublic updates:\n{unseen}")
1432        };
1433        let introduction = if !self.introduced.get(*slot).copied().unwrap_or(false) {
1434            let self_name = speaker_name.clone();
1435            let roster = self
1436                .relay
1437                .active_slots()
1438                .map(|candidate| {
1439                    let name = self
1440                        .roster_names
1441                        .get(candidate)
1442                        .cloned()
1443                        .unwrap_or_else(|| self.hosts[candidate].adapter().display_name());
1444                    let role = if candidate == *slot { " — you" } else { "" };
1445                    format!("{}. {name}{role}", candidate.saturating_add(1))
1446                })
1447                .collect::<Vec<_>>();
1448            let shared_task = self
1449                .relay
1450                .shared_task()
1451                .filter(|task| *task != prompt)
1452                .map(|task| format!("\n\nShared task:\n{task}"))
1453                .unwrap_or_default();
1454            format!(
1455                "You are {self_name}.\nCodeSwarm roster (ordered):\n{}\n\
1456                 Turns relay sequentially through this roster. Treat the user request as the shared task; use timestamped public updates as conversation context.{shared_task}",
1457                roster.join("\n")
1458            )
1459        } else {
1460            String::new()
1461        };
1462        let prompt = format!(
1463            "{introduction}{separator}{prompt}\n\n{}",
1464            if *can_stop {
1465                format!(
1466                    "You are reviewing another agent. If no meaningful correction is needed,\nreply with an optional emoji followed by {STOP_TOKEN} on the final line.\nCodeSwarm hides the token and ends this review batch."
1467                )
1468            } else {
1469                format!(
1470                    "Do not use {STOP_TOKEN} on this turn. Your response must be offered to another agent for review."
1471                )
1472            },
1473            separator = if introduction.is_empty() { "" } else { "\n\n" },
1474        );
1475        let event_sink = self.event_sink.clone();
1476        let host = self
1477            .hosts
1478            .get_mut(*slot)
1479            .ok_or_else(|| AdapterError::Transport("relay selected missing adapter".into()))?;
1480        let prompt = crate::goal::prompt(self.goal.as_ref(), &prompt);
1481        if let Err(error) = host.send_prompt(prompt.clone()).await {
1482            let limited =
1483                report_relay_failure(&mut self.relay, &event_sink, *slot, true, error.to_string());
1484            if limited {
1485                self.relay.finish(*slot, *direct, false);
1486            }
1487            let _ = self.queue_session_metadata();
1488            if limited {
1489                return Ok(decision);
1490            }
1491            return Err(error);
1492        }
1493        if let Some(sink) = &self.event_sink {
1494            sink(AgentEvent::TurnStarted { slot: *slot });
1495        }
1496        if let Some(introduced) = self.introduced.get_mut(*slot) {
1497            *introduced = true;
1498        }
1499        self.dispatches.push((*slot, prompt));
1500        let mut response = String::new();
1501        let mut emitted_text = 0usize;
1502        let completion_event = loop {
1503            if self.cancel_requested.swap(false, Ordering::AcqRel) {
1504                if let Err(error) = cancel_with_timeout(host).await {
1505                    let limited = report_relay_failure(
1506                        &mut self.relay,
1507                        &event_sink,
1508                        *slot,
1509                        true,
1510                        error.to_string(),
1511                    );
1512                    if limited {
1513                        self.relay.finish(*slot, *direct, false);
1514                    }
1515                    let _ = self.queue_session_metadata();
1516                    if limited {
1517                        return Ok(decision);
1518                    }
1519                    return Err(error);
1520                }
1521                return Err(AdapterError::Transport("relay turn cancelled".into()));
1522            }
1523            let update = tokio::select! {
1524                update = host.next_update() => match update {
1525                    Some(Ok(update)) => update,
1526                    Some(Err(error)) => {
1527                        let limited = report_relay_failure(
1528                            &mut self.relay,
1529                            &event_sink,
1530                            *slot,
1531                            true,
1532                            error.to_string(),
1533                        );
1534                        if limited {
1535                            self.relay.finish(*slot, *direct, false);
1536                        }
1537                        let _ = self.queue_session_metadata();
1538                        if limited {
1539                            return Ok(decision);
1540                        }
1541                        return Err(error);
1542                    }
1543                    None => {
1544                        let error = AdapterError::Transport("adapter ended during turn".into());
1545                        let limited = report_relay_failure(
1546                            &mut self.relay,
1547                            &event_sink,
1548                            *slot,
1549                            true,
1550                            error.to_string(),
1551                        );
1552                        if limited {
1553                            self.relay.finish(*slot, *direct, false);
1554                        }
1555                        let _ = self.queue_session_metadata();
1556                        if limited {
1557                            return Ok(decision);
1558                        }
1559                        return Err(error);
1560                    }
1561                },
1562                _ = self.cancel_notify.notified() => {
1563                    if !self.cancel_requested.swap(false, Ordering::AcqRel) {
1564                        continue;
1565                    }
1566                    if let Err(error) = cancel_with_timeout(host).await {
1567                        let limited = report_relay_failure(
1568                            &mut self.relay,
1569                            &event_sink,
1570                            *slot,
1571                            true,
1572                            error.to_string(),
1573                        );
1574                        if limited {
1575                            self.relay.finish(*slot, *direct, false);
1576                        }
1577                        let _ = self.queue_session_metadata();
1578                        if limited {
1579                            return Ok(decision);
1580                        }
1581                        return Err(error);
1582                    }
1583                    return Err(AdapterError::Transport("relay turn cancelled".into()));
1584                },
1585                permission = async {
1586                    match permissions.as_mut() {
1587                        Some(receiver) => receiver.recv().await,
1588                        None => std::future::pending().await,
1589                    }
1590                } => {
1591                    let Some(permission) = permission else {
1592                        permissions = None;
1593                        continue;
1594                    };
1595                    if permission.slot != *slot {
1596                        return Err(AdapterError::Transport(
1597                            "permission response targets an inactive relay slot".into(),
1598                        ));
1599                    }
1600                    host.answer_permission(permission.request_id, permission.answer).await?;
1601                    continue;
1602                },
1603            };
1604            match &update.event {
1605                AgentEvent::Text { text, .. } => response.push_str(text),
1606                AgentEvent::TurnComplete { .. } => {
1607                    let visible_response = response.replace(STOP_TOKEN, "");
1608                    let visible_start = emitted_text.min(visible_response.len());
1609                    let visible_start = floor_char_boundary(&visible_response, visible_start);
1610                    if visible_start < visible_response.len()
1611                        && let Some(sink) = &self.event_sink
1612                    {
1613                        sink(AgentEvent::Text {
1614                            slot: *slot,
1615                            text: visible_response[visible_start..].to_owned(),
1616                        });
1617                    }
1618                    self.cancel_requested.store(false, Ordering::Release);
1619                    break update.event.clone();
1620                }
1621                AgentEvent::Failed {
1622                    started, detail, ..
1623                } => {
1624                    let limited = report_relay_failure(
1625                        &mut self.relay,
1626                        &event_sink,
1627                        *slot,
1628                        *started,
1629                        detail.clone(),
1630                    );
1631                    if limited {
1632                        self.relay.finish(*slot, *direct, false);
1633                    }
1634                    let _ = self.queue_session_metadata();
1635                    if limited {
1636                        return Ok(decision);
1637                    }
1638                    return Err(AdapterError::Transport(detail.clone()));
1639                }
1640                _ => {}
1641            }
1642            if let AgentEvent::Text { .. } = &update.event {
1643                // Hold a short suffix until the next chunk so a stop token
1644                // split across stream messages never leaks into the UI.
1645                let visible_end = response
1646                    .find(STOP_TOKEN)
1647                    .unwrap_or_else(|| response.len().saturating_sub(STOP_TOKEN.len() - 1));
1648                let visible_end = floor_char_boundary(&response, visible_end);
1649                if emitted_text < visible_end {
1650                    if let Some(sink) = &self.event_sink {
1651                        sink(AgentEvent::Text {
1652                            slot: *slot,
1653                            text: response[emitted_text..visible_end].to_owned(),
1654                        });
1655                    }
1656                    emitted_text = visible_end;
1657                }
1658            } else if let Some(sink) = &self.event_sink {
1659                sink(update.event.clone());
1660            }
1661        };
1662        let (response, requested_stop) = strip_stop_token(&response);
1663        let response = response.replace(STOP_TOKEN, "");
1664        let accepted_stop = requested_stop && *can_stop;
1665        let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1666        let response = if needs_stop_acknowledgment {
1667            DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1668        } else {
1669            response
1670        };
1671        // A token-only reviewer response is intentionally hidden, but the
1672        // UI still needs the documented visible acknowledgement. Streamed
1673        // text is normally emitted above; this synthetic acknowledgement is
1674        // the one case where there was no visible adapter chunk to forward.
1675        if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1676            sink(AgentEvent::Text {
1677                slot: *slot,
1678                text: response.clone(),
1679            });
1680        }
1681        if let Some(sink) = &self.event_sink {
1682            sink(completion_event);
1683        }
1684        // A provider plan that ran out mid-turn routes future turns around
1685        // the agent instead of back into the exhausted quota.
1686        if is_usage_limit_response(&response) {
1687            let detail = response.clone();
1688            let _ = self.relay.mark_limited(*slot);
1689            // A normal finish: the ring cursor advances past the limited
1690            // agent so the batch continues with a healthy peer.
1691            self.relay.finish(*slot, *direct, false);
1692            self.queue_session_metadata()?;
1693            if let Some(sink) = &self.event_sink {
1694                sink(AgentEvent::UsageLimitReached {
1695                    slot: *slot,
1696                    detail,
1697                });
1698            }
1699            return Ok(decision);
1700        }
1701        if !*direct && !response.is_empty() {
1702            self.relay
1703                .record_public(public_context_speaker(&speaker_name), response);
1704        }
1705        self.relay.mark_context_seen(*slot);
1706        self.relay.finish(*slot, *direct, accepted_stop);
1707        self.queue_session_metadata()?;
1708        Ok(decision)
1709    }
1710}
1711
1712fn report_relay_failure(
1713    relay: &mut Relay,
1714    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1715    slot: RosterSlot,
1716    started: bool,
1717    detail: String,
1718) -> bool {
1719    // A quota rejection is not a crash: route around the agent without
1720    // tombstoning the slot so a recharge can restore it.
1721    if is_usage_limit_response(&detail) {
1722        let _ = relay.mark_limited(slot);
1723        if let Some(sink) = event_sink {
1724            sink(AgentEvent::UsageLimitReached { slot, detail });
1725        }
1726        return true;
1727    }
1728    let _ = relay.tombstone(slot);
1729    if let Some(sink) = event_sink {
1730        sink(AgentEvent::Failed {
1731            slot,
1732            started,
1733            detail,
1734        });
1735    }
1736    false
1737}
1738
1739/// Deterministic in-memory adapter used for contract and relay tests.
1740#[derive(Debug)]
1741pub struct ScriptedAdapter {
1742    slot: RosterSlot,
1743    capabilities: AgentCapabilities,
1744    events: VecDeque<AdapterResult<AgentEvent>>,
1745    prompts: Vec<String>,
1746}
1747
1748impl ScriptedAdapter {
1749    pub fn new(
1750        slot: RosterSlot,
1751        capabilities: AgentCapabilities,
1752        events: impl IntoIterator<Item = AgentEvent>,
1753    ) -> Self {
1754        Self {
1755            slot,
1756            capabilities,
1757            events: events.into_iter().map(Ok).collect(),
1758            prompts: Vec::new(),
1759        }
1760    }
1761
1762    pub fn prompts(&self) -> &[String] {
1763        &self.prompts
1764    }
1765}
1766
1767#[async_trait]
1768impl AgentAdapter for ScriptedAdapter {
1769    fn slot(&self) -> RosterSlot {
1770        self.slot
1771    }
1772
1773    fn capabilities(&self) -> AgentCapabilities {
1774        self.capabilities.clone()
1775    }
1776
1777    async fn start(&mut self) -> AdapterResult<()> {
1778        Ok(())
1779    }
1780
1781    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1782        self.prompts.push(prompt);
1783        Ok(())
1784    }
1785
1786    async fn cancel(&mut self) -> AdapterResult<bool> {
1787        Ok(self.capabilities.supports_cancel)
1788    }
1789
1790    async fn answer_permission(
1791        &mut self,
1792        _request_id: String,
1793        _answer: PermissionAnswer,
1794    ) -> AdapterResult<()> {
1795        if self.capabilities.supports_permissions {
1796            Ok(())
1797        } else {
1798            Err(AdapterError::Unsupported("permission answer"))
1799        }
1800    }
1801
1802    async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
1803        if self.capabilities.supports_modes {
1804            Ok(())
1805        } else {
1806            Err(AdapterError::Unsupported("set_mode"))
1807        }
1808    }
1809
1810    async fn reload(&mut self) -> AdapterResult<()> {
1811        Ok(())
1812    }
1813
1814    async fn stop(&mut self) -> AdapterResult<()> {
1815        Ok(())
1816    }
1817
1818    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
1819        self.events.pop_front()
1820    }
1821}
1822
1823/// A direct stream-JSON adapter for Antigravity. It deliberately does not
1824/// pretend to be ACP; it translates its documented events into core events.
1825#[derive(Debug)]
1826pub struct AgyAdapter {
1827    slot: RosterSlot,
1828    cwd: PathBuf,
1829    command: String,
1830    mode: String,
1831    mode_policy: String,
1832    session_id: Option<String>,
1833    child: Option<Child>,
1834    sender: mpsc::Sender<AdapterResult<AgentEvent>>,
1835    receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
1836    /// Antigravity announces its conversation id in the `init` event. Keep it
1837    /// outside the stream task so the next prompt can resume the same native
1838    /// conversation without making the adapter reader/UI share a borrow.
1839    announced_session: Arc<Mutex<Option<String>>>,
1840    cancel_requested: Arc<AtomicBool>,
1841}
1842
1843impl AgyAdapter {
1844    pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
1845        let (sender, receiver) = mpsc::channel(256);
1846        Self {
1847            slot,
1848            cwd,
1849            command: command.into(),
1850            mode: "default".into(),
1851            mode_policy: "agy:full-access".into(),
1852            session_id: None,
1853            child: None,
1854            sender,
1855            receiver,
1856            announced_session: Arc::new(Mutex::new(None)),
1857            cancel_requested: Arc::new(AtomicBool::new(false)),
1858        }
1859    }
1860
1861    pub fn with_session_id(
1862        slot: RosterSlot,
1863        cwd: PathBuf,
1864        command: impl Into<String>,
1865        session_id: impl Into<String>,
1866    ) -> Self {
1867        let mut adapter = Self::new(slot, cwd, command);
1868        adapter.session_id = Some(session_id.into());
1869        adapter
1870    }
1871
1872    fn modes() -> Vec<Mode> {
1873        vec![
1874            Mode {
1875                id: "agy:full-access".into(),
1876                label: "Auto pilot".into(),
1877            },
1878            Mode {
1879                id: "agy:manual".into(),
1880                label: "Manual".into(),
1881            },
1882            Mode {
1883                id: "accept-edits".into(),
1884                label: "Accept Edits".into(),
1885            },
1886            Mode {
1887                id: "plan".into(),
1888                label: "Plan".into(),
1889            },
1890        ]
1891    }
1892
1893    async fn emit(&self, event: AdapterResult<AgentEvent>) {
1894        let _ = self.sender.send(event).await;
1895    }
1896}
1897
1898#[async_trait]
1899impl AgentAdapter for AgyAdapter {
1900    fn slot(&self) -> RosterSlot {
1901        self.slot
1902    }
1903
1904    fn session_id(&self) -> Option<String> {
1905        self.session_id.clone()
1906    }
1907
1908    fn protocol(&self) -> &'static str {
1909        "native"
1910    }
1911
1912    fn capabilities(&self) -> AgentCapabilities {
1913        AgentCapabilities {
1914            supports_cancel: true,
1915            supports_modes: true,
1916            supports_permissions: false,
1917            supports_terminals: true,
1918            supports_session_load: true,
1919            supports_models: false,
1920        }
1921    }
1922
1923    async fn start(&mut self) -> AdapterResult<()> {
1924        self.cancel_requested.store(false, Ordering::Release);
1925        self.emit(Ok(AgentEvent::Ready {
1926            slot: self.slot,
1927            capabilities: self.capabilities(),
1928        }))
1929        .await;
1930        self.emit(Ok(AgentEvent::ModesReplaced {
1931            slot: self.slot,
1932            modes: Self::modes(),
1933            current_mode: Some(self.mode_policy.clone()),
1934        }))
1935        .await;
1936        Ok(())
1937    }
1938
1939    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1940        if self.child.is_some() {
1941            return Err(AdapterError::Transport(
1942                "agent is already handling a turn".into(),
1943            ));
1944        }
1945        if self.session_id.is_none()
1946            && let Ok(session) = self.announced_session.lock()
1947        {
1948            self.session_id = session.clone();
1949        }
1950        self.cancel_requested.store(false, Ordering::Release);
1951        let (program, args) = parse_command_line(&self.command)
1952            .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
1953        let mut command = Command::new(program);
1954        isolate_process_group(&mut command);
1955        command
1956            .args(args)
1957            .arg("--print")
1958            .arg(prompt)
1959            .arg("--print-timeout")
1960            .arg("60m")
1961            .arg("--output-format")
1962            .arg("stream-json")
1963            .current_dir(&self.cwd)
1964            .env("CODESWARM_CWD", &self.cwd)
1965            .stdout(Stdio::piped())
1966            .stderr(Stdio::piped());
1967        if let Some(session_id) = &self.session_id {
1968            command.arg("--conversation").arg(session_id);
1969        }
1970        if self.mode != "default" {
1971            command.arg("--mode").arg(&self.mode);
1972        }
1973        let mut child = command
1974            .spawn()
1975            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
1976        let stdout = match child.stdout.take() {
1977            Some(stdout) => stdout,
1978            None => {
1979                let _ = terminate_child(&mut child).await;
1980                return Err(AdapterError::Transport("agent has no stdout".into()));
1981            }
1982        };
1983        let stderr = match child.stderr.take() {
1984            Some(stderr) => stderr,
1985            None => {
1986                let _ = terminate_child(&mut child).await;
1987                return Err(AdapterError::Transport("agent has no stderr".into()));
1988            }
1989        };
1990        let sender = self.sender.clone();
1991        let slot = self.slot;
1992        let announced_session = Arc::clone(&self.announced_session);
1993        let cancel_requested = Arc::clone(&self.cancel_requested);
1994        tokio::spawn(async move {
1995            let stderr_task = tokio::spawn(async move {
1996                const MAX_STDERR: usize = 32 * 1024;
1997                let mut stderr = BufReader::new(stderr);
1998                let mut bytes = Vec::new();
1999                let mut chunk = [0_u8; 4096];
2000                while let Ok(count) = stderr.read(&mut chunk).await {
2001                    if count == 0 {
2002                        break;
2003                    }
2004                    bytes.extend_from_slice(&chunk[..count]);
2005                    if bytes.len() > MAX_STDERR {
2006                        let keep_from = bytes.len() - MAX_STDERR;
2007                        bytes.drain(..keep_from);
2008                    }
2009                }
2010                String::from_utf8_lossy(&bytes).trim().to_owned()
2011            });
2012            let mut lines = BufReader::new(stdout).lines();
2013            let mut result: Option<Value> = None;
2014            let mut streamed_response = false;
2015            while let Ok(Some(line)) = lines.next_line().await {
2016                let value = match serde_json::from_str::<Value>(&line) {
2017                    Ok(value) => value,
2018                    Err(_) => {
2019                        // Native stream-json can contain diagnostic junk on
2020                        // stdout. Match the Python adapter's tolerant stream
2021                        // behavior and wait for the final result instead of
2022                        // turning one malformed line into a dead turn.
2023                        continue;
2024                    }
2025                };
2026                if value.get("event").and_then(Value::as_str) == Some("init")
2027                    && let Some(session_id) = value
2028                        .get("conversation_id")
2029                        .or_else(|| value.get("conversationId"))
2030                        .and_then(Value::as_str)
2031                        .filter(|id| !id.is_empty())
2032                    && let Ok(mut announced) = announced_session.lock()
2033                {
2034                    *announced = Some(session_id.to_owned());
2035                }
2036                if value.get("event").and_then(Value::as_str) == Some("result") {
2037                    result = value.get("result").cloned();
2038                }
2039                match parse_agy_value(slot, &value) {
2040                    Ok(Some(event)) => {
2041                        if matches!(event, AgentEvent::Text { .. }) {
2042                            streamed_response = true;
2043                        }
2044                        if sender.send(Ok(event)).await.is_err() {
2045                            break;
2046                        }
2047                    }
2048                    Ok(None) => {}
2049                    Err(error) => {
2050                        let _ = sender.send(Err(error)).await;
2051                    }
2052                }
2053            }
2054            let stderr = stderr_task.await.ok().unwrap_or_default();
2055            let succeeded = cancel_requested.load(Ordering::Acquire)
2056                || result
2057                    .as_ref()
2058                    .and_then(|result| result.get("status"))
2059                    .and_then(Value::as_str)
2060                    == Some("SUCCESS");
2061            if succeeded {
2062                // Some native stream-json wrappers emit only lifecycle events
2063                // and put the complete answer in the final result object. The
2064                // Python adapter surfaced that answer; do the same, while
2065                // avoiding duplication when token chunks were already sent.
2066                if !streamed_response
2067                    && let Some(response) = result
2068                        .as_ref()
2069                        .and_then(|result| result.get("response"))
2070                        .and_then(Value::as_str)
2071                        .filter(|response| !response.is_empty())
2072                {
2073                    let _ = sender
2074                        .send(Ok(AgentEvent::Text {
2075                            slot,
2076                            text: response.to_owned(),
2077                        }))
2078                        .await;
2079                }
2080                let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2081            } else {
2082                let detail = result
2083                    .as_ref()
2084                    .and_then(|result| result.get("error"))
2085                    .and_then(Value::as_str)
2086                    .filter(|detail| !detail.is_empty())
2087                    .map(str::to_owned)
2088                    .or_else(|| (!stderr.is_empty()).then_some(stderr))
2089                    .unwrap_or_else(|| "native stream ended before a successful result".into());
2090                let _ = sender
2091                    .send(Ok(AgentEvent::Failed {
2092                        slot,
2093                        started: true,
2094                        detail,
2095                    }))
2096                    .await;
2097            }
2098        });
2099        self.child = Some(child);
2100        Ok(())
2101    }
2102
2103    async fn cancel(&mut self) -> AdapterResult<bool> {
2104        self.cancel_requested.store(true, Ordering::Release);
2105        let Some(mut child) = self.child.take() else {
2106            return Ok(false);
2107        };
2108        // `Child` does not reap itself when dropped. Awaiting wait after the
2109        // kill keeps repeated prompts from accumulating zombies, especially
2110        // when cancellation happens before the stream reader observes EOF.
2111        terminate_child(&mut child).await?;
2112        let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2113            while let Some(event) = self.receiver.recv().await {
2114                if matches!(
2115                    event,
2116                    Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2117                ) {
2118                    break;
2119                }
2120            }
2121        })
2122        .await;
2123        Ok(true)
2124    }
2125
2126    async fn answer_permission(
2127        &mut self,
2128        _request_id: String,
2129        _answer: PermissionAnswer,
2130    ) -> AdapterResult<()> {
2131        Err(AdapterError::Unsupported("permission answer"))
2132    }
2133
2134    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2135        let (mode, mode_policy) = match mode.as_str() {
2136            "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2137                ("default".to_owned(), "agy:full-access".to_owned())
2138            }
2139            "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2140            "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2141                ("accept-edits".to_owned(), "accept-edits".to_owned())
2142            }
2143            "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2144                ("default".to_owned(), "agy:manual".to_owned())
2145            }
2146            "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2147            "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2148            _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2149        };
2150        self.mode = mode;
2151        self.mode_policy = mode_policy.clone();
2152        self.emit(Ok(AgentEvent::ModesReplaced {
2153            slot: self.slot,
2154            modes: Self::modes(),
2155            current_mode: Some(mode_policy),
2156        }))
2157        .await;
2158        Ok(())
2159    }
2160
2161    async fn reload(&mut self) -> AdapterResult<()> {
2162        self.stop().await?;
2163        self.start().await
2164    }
2165
2166    async fn stop(&mut self) -> AdapterResult<()> {
2167        let _ = self.cancel().await?;
2168        Ok(())
2169    }
2170
2171    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2172        let event = self.receiver.recv().await;
2173        if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2174            && self.session_id.is_none()
2175            && let Ok(session) = self.announced_session.lock()
2176        {
2177            self.session_id = session.clone();
2178        }
2179        if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2180            && let Some(mut child) = self.child.take()
2181        {
2182            let _ = child.wait().await;
2183        }
2184        event
2185    }
2186}
2187
2188#[cfg(test)]
2189#[cfg_attr(not(test), allow(dead_code))]
2190fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2191    let value: Value =
2192        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2193    parse_agy_value(slot, &value)
2194}
2195
2196fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2197    let event = value.get("event").and_then(Value::as_str);
2198    if let Some(terminal) = parse_terminal_event(value, event) {
2199        return Ok(Some(AgentEvent::Terminal {
2200            slot,
2201            event: terminal,
2202        }));
2203    }
2204    match event {
2205        Some("step_update") => {
2206            let Some(update) = value.get("step_update") else {
2207                return Ok(None);
2208            };
2209            let is_response = update
2210                .get("step_type")
2211                .and_then(Value::as_str)
2212                .is_some_and(|kind| kind == "agent_response");
2213            let text = update.get("text_delta").and_then(Value::as_str);
2214            let response = is_response
2215                .then(|| text.map(str::to_owned))
2216                .flatten()
2217                .filter(|text| !text.is_empty())
2218                .map(|text| AgentEvent::Text { slot, text });
2219            Ok(response.or_else(|| parse_agy_tool(slot, value)))
2220        }
2221        _ => Ok(None),
2222    }
2223}
2224
2225fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2226    let update = value.get("step_update")?;
2227    if update.get("step_type")?.as_str()? != "tool" {
2228        return None;
2229    }
2230    let step_index = update.get("step_index")?.as_i64()?;
2231    let title = update
2232        .get("tool_name")
2233        .and_then(Value::as_str)
2234        .unwrap_or("Tool call")
2235        .replace('_', " ");
2236    let status = match update.get("state").and_then(Value::as_str) {
2237        Some("DONE") => ToolStatus::Completed,
2238        Some("FAILED") => ToolStatus::Failed,
2239        Some("ACTIVE") => ToolStatus::Running,
2240        _ => ToolStatus::Pending,
2241    };
2242    let detail = update
2243        .get("tool_info")
2244        .and_then(|info| info.get("output"))
2245        .and_then(Value::as_str)
2246        .map(str::to_owned);
2247    Some(AgentEvent::Tool {
2248        slot,
2249        update: ToolUpdate {
2250            id: format!("agy-tool-{step_index}"),
2251            title,
2252            status,
2253            detail,
2254        },
2255    })
2256}
2257
2258/// Stdio ACP transport. Protocol-specific response handling belongs here,
2259/// keeping JSON-RPC framing outside the core and terminal renderer.
2260#[derive(Debug)]
2261pub struct AcpAdapter {
2262    slot: RosterSlot,
2263    program: String,
2264    args: Vec<String>,
2265    cwd: PathBuf,
2266    child: Option<Child>,
2267    reader: Option<BufReader<ChildStdout>>,
2268    capabilities: AgentCapabilities,
2269    modes: Vec<Mode>,
2270    models: Vec<Mode>,
2271    model_config_id: Option<String>,
2272    session_id: Option<String>,
2273    next_request_id: u64,
2274    prompt_request_id: Option<u64>,
2275    queued_events: VecDeque<AdapterResult<AgentEvent>>,
2276    stderr_task: Option<tokio::task::JoinHandle<String>>,
2277    terminals: BTreeMap<String, TerminalProcess>,
2278    next_terminal_id: u64,
2279}
2280
2281impl AcpAdapter {
2282    pub fn new(
2283        slot: RosterSlot,
2284        cwd: PathBuf,
2285        program: impl Into<String>,
2286        args: Vec<String>,
2287    ) -> Self {
2288        Self {
2289            slot,
2290            program: program.into(),
2291            args,
2292            cwd,
2293            child: None,
2294            reader: None,
2295            capabilities: AgentCapabilities::default(),
2296            modes: Vec::new(),
2297            models: Vec::new(),
2298            model_config_id: None,
2299            session_id: None,
2300            next_request_id: 1,
2301            prompt_request_id: None,
2302            queued_events: VecDeque::new(),
2303            stderr_task: None,
2304            terminals: BTreeMap::new(),
2305            next_terminal_id: 1,
2306        }
2307    }
2308
2309    pub fn with_session_id(
2310        slot: RosterSlot,
2311        cwd: PathBuf,
2312        program: impl Into<String>,
2313        args: Vec<String>,
2314        session_id: impl Into<String>,
2315    ) -> Self {
2316        let mut adapter = Self::new(slot, cwd, program, args);
2317        adapter.session_id = Some(session_id.into());
2318        adapter
2319    }
2320
2321    async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2322        self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2323            .await
2324    }
2325
2326    async fn request_with_timeout(
2327        &mut self,
2328        method: &str,
2329        params: Value,
2330        deadline: std::time::Duration,
2331    ) -> AdapterResult<Value> {
2332        tokio::time::timeout(deadline, self.request_inner(method, params))
2333            .await
2334            .map_err(|_| {
2335                AdapterError::Transport(format!(
2336                    "ACP {method} timed out; reload the agent to retry"
2337                ))
2338            })?
2339    }
2340
2341    async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2342        let request_id = self.next_request_id;
2343        self.next_request_id += 1;
2344        self.write_json(serde_json::json!({
2345            "jsonrpc": "2.0",
2346            "id": request_id,
2347            "method": method,
2348            "params": params,
2349        }))
2350        .await?;
2351        loop {
2352            let line = self.read_line().await?;
2353            let value: Value = match serde_json::from_str(&line) {
2354                Ok(value) => value,
2355                Err(_) => {
2356                    // Keep the transport alive when a peer writes a stray
2357                    // diagnostic line. This is common with CLI wrappers and
2358                    // matches the baseline client's tolerant stream loop.
2359                    continue;
2360                }
2361            };
2362            if self.reject_empty_permission_request(&value).await? {
2363                continue;
2364            }
2365            if self.handle_client_request(&value).await? {
2366                continue;
2367            }
2368            if value
2369                .get("id")
2370                .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2371            {
2372                if let Some(error) = value.get("error") {
2373                    return Err(AdapterError::Protocol(error.to_string()));
2374                }
2375                return value
2376                    .get("result")
2377                    .cloned()
2378                    .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2379            }
2380            if let Some(event) = parse_acp_notification(self.slot, &line)? {
2381                self.queued_events.push_back(Ok(event));
2382            }
2383        }
2384    }
2385
2386    async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2387        let child = self
2388            .child
2389            .as_mut()
2390            .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2391        let stdin = child
2392            .stdin
2393            .as_mut()
2394            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2395        stdin
2396            .write_all(value.to_string().as_bytes())
2397            .await
2398            .map_err(|error| AdapterError::Transport(error.to_string()))?;
2399        stdin
2400            .write_all(b"\n")
2401            .await
2402            .map_err(|error| AdapterError::Transport(error.to_string()))
2403    }
2404
2405    /// ACP permission requests are JSON-RPC requests, not fire-and-forget
2406    /// notifications. An empty option list is invalid and must be answered
2407    /// with an error so the peer does not wait forever for a decision. This is
2408    /// the same validation performed by the Python ACP server.
2409    async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2410        if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2411            || value.get("id").is_none()
2412        {
2413            return Ok(false);
2414        }
2415        let valid = value
2416            .get("params")
2417            .and_then(|params| params.get("options"))
2418            .and_then(Value::as_array)
2419            .is_some_and(|options| !options.is_empty());
2420        if valid {
2421            return Ok(false);
2422        }
2423        self.write_json(serde_json::json!({
2424            "jsonrpc": "2.0",
2425            "id": value.get("id").cloned().unwrap_or(Value::Null),
2426            "error": {
2427                "code": -32602,
2428                "message": "Permission request requires at least one option",
2429            },
2430        }))
2431        .await?;
2432        Ok(true)
2433    }
2434
2435    fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2436        let root = self
2437            .cwd
2438            .canonicalize()
2439            .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2440        let requested = Path::new(path);
2441        let candidate = if requested.is_absolute() {
2442            requested.to_path_buf()
2443        } else {
2444            root.join(requested)
2445        };
2446        let resolved = if !candidate.exists() {
2447            let parent = candidate
2448                .parent()
2449                .ok_or_else(|| "file path has no parent".to_owned())?
2450                .canonicalize()
2451                .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2452            parent.join(
2453                candidate
2454                    .file_name()
2455                    .ok_or_else(|| "file path has no filename".to_owned())?,
2456            )
2457        } else {
2458            candidate
2459                .canonicalize()
2460                .map_err(|error| format!("unable to resolve file path: {error}"))?
2461        };
2462        if !resolved.starts_with(&root) {
2463            return Err("file path is outside the project".into());
2464        }
2465        Ok(resolved)
2466    }
2467
2468    fn read_workspace_text(
2469        &self,
2470        path: &str,
2471        line: Option<i64>,
2472        limit: Option<i64>,
2473    ) -> Result<String, String> {
2474        if line.is_some_and(|line| line < 1) {
2475            return Err("line must be positive".into());
2476        }
2477        if limit.is_some_and(|limit| limit < 0) {
2478            return Err("limit must not be negative".into());
2479        }
2480        let path = self.workspace_path(path)?;
2481        let mut bytes = Vec::new();
2482        let mut source = match std::fs::File::open(path) {
2483            Ok(source) => source,
2484            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2485            Err(error) => return Err(error.to_string()),
2486        };
2487        source
2488            .by_ref()
2489            .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2490            .read_to_end(&mut bytes)
2491            .map_err(|error| error.to_string())?;
2492        bytes.truncate(MAX_FILE_READ_BYTES);
2493        let text = String::from_utf8_lossy(&bytes);
2494        if line.is_none() && limit.is_none() {
2495            return Ok(text.into_owned());
2496        }
2497        let start = line.map_or(0, |line| line as usize - 1);
2498        let limit = limit.unwrap_or(i64::MAX) as usize;
2499        let selected = text
2500            .split_inclusive('\n')
2501            .skip(start)
2502            .take(limit)
2503            .collect::<String>();
2504        if line.is_some() {
2505            Ok(selected.trim_end_matches('\n').to_owned())
2506        } else {
2507            Ok(selected)
2508        }
2509    }
2510
2511    fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2512        let path = params
2513            .get("path")
2514            .and_then(Value::as_str)
2515            .filter(|path| !path.is_empty())
2516            .ok_or("path must be a non-empty string")?;
2517        let content = params
2518            .get("content")
2519            .and_then(Value::as_str)
2520            .ok_or("content must be a string")?;
2521        let path = self.workspace_path(path)?;
2522        std::fs::write(path, content).map_err(|error| error.to_string())
2523    }
2524
2525    async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2526        let command = params
2527            .get("command")
2528            .and_then(Value::as_str)
2529            .filter(|command| !command.trim().is_empty())
2530            .ok_or_else(|| "terminal command is required".to_owned())?;
2531        let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2532        let cwd = self.workspace_path(cwd)?;
2533        if !cwd.is_dir() {
2534            return Err("terminal cwd is not a directory".into());
2535        }
2536        let mut process = Command::new(command);
2537        isolate_process_group(&mut process);
2538        if let Some(args) = params.get("args").and_then(Value::as_array) {
2539            process.args(args.iter().filter_map(Value::as_str));
2540        }
2541        process
2542            .current_dir(&cwd)
2543            .stdin(Stdio::null())
2544            .stdout(Stdio::piped())
2545            .stderr(Stdio::piped());
2546        if let Some(env) = params.get("env") {
2547            if let Some(entries) = env.as_array() {
2548                for entry in entries {
2549                    if let (Some(name), Some(value)) = (
2550                        entry.get("name").and_then(Value::as_str),
2551                        entry.get("value").and_then(Value::as_str),
2552                    ) {
2553                        process.env(name, value);
2554                    }
2555                }
2556            } else if let Some(entries) = env.as_object() {
2557                for (name, value) in entries {
2558                    if let Some(value) = value.as_str() {
2559                        process.env(name, value);
2560                    }
2561                }
2562            }
2563        }
2564        let mut child = process.spawn().map_err(|error| error.to_string())?;
2565        let stdout = child.stdout.take();
2566        let stderr = child.stderr.take();
2567        let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2568            let _ = terminate_child(&mut child).await;
2569            return Err("terminal has no output pipes".into());
2570        };
2571        let output = Arc::new(Mutex::new(Vec::new()));
2572        let truncated = Arc::new(AtomicBool::new(false));
2573        let output_readers = Arc::new(AtomicUsize::new(2));
2574        let output_limit = params
2575            .get("outputByteLimit")
2576            .and_then(Value::as_u64)
2577            .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2578                usize::try_from(limit)
2579                    .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2580                    .min(MAX_TERMINAL_OUTPUT_BYTES)
2581            });
2582        tokio::spawn(drain_terminal_output(
2583            stdout,
2584            Arc::clone(&output),
2585            Arc::clone(&truncated),
2586            Arc::clone(&output_readers),
2587            output_limit,
2588        ));
2589        tokio::spawn(drain_terminal_output(
2590            stderr,
2591            Arc::clone(&output),
2592            Arc::clone(&truncated),
2593            Arc::clone(&output_readers),
2594            output_limit,
2595        ));
2596        let id = format!("terminal-{}", self.next_terminal_id);
2597        self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2598        let state = TerminalProcess {
2599            child: Arc::new(AsyncMutex::new(Some(child))),
2600            output,
2601            truncated,
2602            output_readers,
2603        };
2604        self.terminals.insert(id.clone(), state);
2605        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2606            slot: self.slot,
2607            event: TerminalEvent::Created {
2608                id: id.clone(),
2609                command: std::iter::once(command)
2610                    .chain(
2611                        params
2612                            .get("args")
2613                            .and_then(Value::as_array)
2614                            .into_iter()
2615                            .flatten()
2616                            .filter_map(Value::as_str),
2617                    )
2618                    .collect::<Vec<_>>()
2619                    .join(" "),
2620            },
2621        }));
2622        Ok(serde_json::json!({"terminalId": id}))
2623    }
2624
2625    async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2626        let terminal = self
2627            .terminals
2628            .get(id)
2629            .ok_or_else(|| "terminal not found".to_owned())?
2630            .clone();
2631        let output = terminal
2632            .output
2633            .lock()
2634            .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2635            .unwrap_or_default();
2636        let exit_code = terminal.exit_code().await;
2637        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2638            slot: self.slot,
2639            event: TerminalEvent::Output {
2640                id: id.to_owned(),
2641                text: output.clone(),
2642            },
2643        }));
2644        let mut response = serde_json::json!({
2645            "output": output,
2646            "truncated": terminal.truncated.load(Ordering::Acquire),
2647        });
2648        if let Some(code) = exit_code {
2649            response["exitStatus"] = serde_json::json!({"exitCode": code});
2650        }
2651        Ok(response)
2652    }
2653
2654    async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2655        let terminal = self
2656            .terminals
2657            .get(id)
2658            .ok_or_else(|| "terminal not found".to_owned())?
2659            .clone();
2660        let exit_code = terminal.wait().await;
2661        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2662            slot: self.slot,
2663            event: TerminalEvent::Exited {
2664                id: id.to_owned(),
2665                code: exit_code.unwrap_or(-1),
2666            },
2667        }));
2668        Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2669    }
2670
2671    /// Handle requests initiated by an ACP agent against the client. File
2672    /// access is mediated through the configured workspace root; unsupported
2673    /// requests receive a JSON-RPC error instead of hanging the agent.
2674    async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2675        let Some(method) = value.get("method").and_then(Value::as_str) else {
2676            return Ok(false);
2677        };
2678        let Some(id) = value.get("id").cloned() else {
2679            return Ok(false);
2680        };
2681        // Permission requests are consumed by the normalized event parser and
2682        // answered later through the focused UI action, not here.
2683        if method == "session/request_permission" {
2684            return Ok(false);
2685        }
2686        let params = value.get("params").cloned().unwrap_or(Value::Null);
2687        let response = match method {
2688            "fs/read_text_file" => {
2689                let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2690                let line = params.get("line").and_then(Value::as_i64);
2691                let limit = params.get("limit").and_then(Value::as_i64);
2692                match self.read_workspace_text(path, line, limit) {
2693                    Ok(content) => serde_json::json!({
2694                        "jsonrpc": "2.0",
2695                        "id": id,
2696                        "result": {"content": content},
2697                    }),
2698                    Err(message) => serde_json::json!({
2699                        "jsonrpc": "2.0",
2700                        "id": id,
2701                        "error": {"code": -32602, "message": message},
2702                    }),
2703                }
2704            }
2705            "fs/write_text_file" => {
2706                let result = self.write_workspace_text(&params);
2707                match result {
2708                    Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2709                    Err(message) => serde_json::json!({
2710                        "jsonrpc": "2.0",
2711                        "id": id,
2712                        "error": {"code": -32602, "message": message},
2713                    }),
2714                }
2715            }
2716            "terminal/create" => match self.terminal_create(&params).await {
2717                Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2718                Err(message) => serde_json::json!({
2719                    "jsonrpc": "2.0",
2720                    "id": id,
2721                    "error": {"code": -32602, "message": message},
2722                }),
2723            },
2724            "terminal/output" => {
2725                let terminal_id = params
2726                    .get("terminalId")
2727                    .and_then(Value::as_str)
2728                    .unwrap_or("");
2729                match self.terminal_output(terminal_id).await {
2730                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2731                    Err(message) => serde_json::json!({
2732                        "jsonrpc": "2.0",
2733                        "id": id,
2734                        "error": {"code": -32602, "message": message},
2735                    }),
2736                }
2737            }
2738            "terminal/wait_for_exit" => {
2739                let terminal_id = params
2740                    .get("terminalId")
2741                    .and_then(Value::as_str)
2742                    .unwrap_or("");
2743                match self.terminal_wait(terminal_id).await {
2744                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2745                    Err(message) => serde_json::json!({
2746                        "jsonrpc": "2.0",
2747                        "id": id,
2748                        "error": {"code": -32602, "message": message},
2749                    }),
2750                }
2751            }
2752            "terminal/kill" => {
2753                let terminal_id = params
2754                    .get("terminalId")
2755                    .and_then(Value::as_str)
2756                    .unwrap_or("");
2757                if let Some(terminal) = self.terminals.get(terminal_id) {
2758                    terminal.kill().await;
2759                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2760                } else {
2761                    serde_json::json!({
2762                        "jsonrpc": "2.0",
2763                        "id": id,
2764                        "error": {"code": -32602, "message": "terminal not found"},
2765                    })
2766                }
2767            }
2768            "terminal/release" => {
2769                let terminal_id = params
2770                    .get("terminalId")
2771                    .and_then(Value::as_str)
2772                    .unwrap_or("");
2773                if let Some(terminal) = self.terminals.remove(terminal_id) {
2774                    terminal.stop().await;
2775                    self.queued_events.push_back(Ok(AgentEvent::Terminal {
2776                        slot: self.slot,
2777                        event: TerminalEvent::Released {
2778                            id: terminal_id.to_owned(),
2779                        },
2780                    }));
2781                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2782                } else {
2783                    serde_json::json!({
2784                        "jsonrpc": "2.0",
2785                        "id": id,
2786                        "error": {"code": -32602, "message": "terminal not found"},
2787                    })
2788                }
2789            }
2790            _ => serde_json::json!({
2791                "jsonrpc": "2.0",
2792                "id": id,
2793                "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
2794            }),
2795        };
2796        self.write_json(response).await?;
2797        Ok(true)
2798    }
2799
2800    async fn read_line(&mut self) -> AdapterResult<String> {
2801        let reader = self
2802            .reader
2803            .as_mut()
2804            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
2805        read_bounded_line(reader).await
2806    }
2807
2808    async fn start(&mut self) -> AdapterResult<()> {
2809        // Starting an adapter twice must never orphan the first transport.
2810        if self.child.is_some() {
2811            self.stop().await?;
2812        }
2813        let mut command = Command::new(&self.program);
2814        isolate_process_group(&mut command);
2815        command
2816            .args(&self.args)
2817            .current_dir(&self.cwd)
2818            .stdin(Stdio::piped())
2819            .stdout(Stdio::piped())
2820            .stderr(Stdio::piped())
2821            .env("CODESWARM_CWD", &self.cwd);
2822        if self.program.to_ascii_lowercase().contains("gemini")
2823            || self
2824                .args
2825                .iter()
2826                .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
2827        {
2828            command.env("GEMINI_TELEMETRY_ENABLED", "false");
2829        }
2830        let mut child = command
2831            .spawn()
2832            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2833        let stdout = match child.stdout.take() {
2834            Some(stdout) => stdout,
2835            None => {
2836                let _ = terminate_child(&mut child).await;
2837                return Err(AdapterError::Transport("ACP agent has no stdout".into()));
2838            }
2839        };
2840        let stderr = match child.stderr.take() {
2841            Some(stderr) => stderr,
2842            None => {
2843                let _ = terminate_child(&mut child).await;
2844                return Err(AdapterError::Transport("ACP agent has no stderr".into()));
2845            }
2846        };
2847        self.child = Some(child);
2848        self.reader = Some(BufReader::new(stdout));
2849        self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
2850
2851        let initialize = match self
2852            .request(
2853                "initialize",
2854                serde_json::json!({
2855                    "protocolVersion": 1,
2856                    "clientCapabilities": {
2857                        "fs": {"readTextFile": true, "writeTextFile": true},
2858                        "terminal": true,
2859                    },
2860                    "clientInfo": {
2861                        "name": "CodeSwarm",
2862                        "title": "CodeSwarm",
2863                        "version": env!("CARGO_PKG_VERSION"),
2864                    },
2865                }),
2866            )
2867            .await
2868        {
2869            Ok(value) => value,
2870            Err(error) => {
2871                let _ = self.stop().await;
2872                return Err(error);
2873            }
2874        };
2875        let agent_capabilities = initialize
2876            .get("agentCapabilities")
2877            .cloned()
2878            .unwrap_or(Value::Null);
2879        self.capabilities = AgentCapabilities {
2880            supports_cancel: true,
2881            supports_modes: true,
2882            supports_permissions: true,
2883            supports_terminals: true,
2884            supports_session_load: agent_capabilities
2885                .get("loadSession")
2886                .and_then(Value::as_bool)
2887                .unwrap_or(false),
2888            supports_models: false,
2889        };
2890        let session = if let Some(session_id) = self.session_id.clone() {
2891            if !self.capabilities.supports_session_load {
2892                let _ = self.stop().await;
2893                return Err(AdapterError::Unsupported("session/load"));
2894            }
2895            match self
2896                .request(
2897                    "session/load",
2898                    serde_json::json!({
2899                        "cwd": self.cwd,
2900                        "mcpServers": [],
2901                        "sessionId": session_id,
2902                    }),
2903                )
2904                .await
2905            {
2906                Ok(value) => value,
2907                Err(error) => {
2908                    let _ = self.stop().await;
2909                    return Err(error);
2910                }
2911            }
2912        } else {
2913            let session = match self
2914                .request(
2915                    "session/new",
2916                    serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
2917                )
2918                .await
2919            {
2920                Ok(value) => value,
2921                Err(error) => {
2922                    let _ = self.stop().await;
2923                    return Err(error);
2924                }
2925            };
2926            self.session_id = session
2927                .get("sessionId")
2928                .and_then(Value::as_str)
2929                .map(str::to_owned);
2930            if self.session_id.is_none() {
2931                let _ = self.stop().await;
2932                return Err(AdapterError::Protocol(
2933                    "session/new returned no sessionId".into(),
2934                ));
2935            }
2936            session
2937        };
2938        self.capabilities.supports_modes = false;
2939        if let Some(modes) = session.get("modes") {
2940            let available = modes
2941                .get("availableModes")
2942                .and_then(Value::as_array)
2943                .map(|modes| {
2944                    modes
2945                        .iter()
2946                        .filter_map(|mode| {
2947                            Some(Mode {
2948                                id: mode.get("id")?.as_str()?.to_owned(),
2949                                label: mode.get("name")?.as_str()?.to_owned(),
2950                            })
2951                        })
2952                        .collect::<Vec<_>>()
2953                })
2954                .unwrap_or_default();
2955            self.modes = available.clone();
2956            self.capabilities.supports_modes = !available.is_empty();
2957            self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
2958                slot: self.slot,
2959                modes: available,
2960                current_mode: modes
2961                    .get("currentModeId")
2962                    .and_then(Value::as_str)
2963                    .map(str::to_owned),
2964            }));
2965        }
2966        self.models.clear();
2967        self.model_config_id = None;
2968        let current_model =
2969            parse_model_config(&session).and_then(|(config_id, models, current)| {
2970                self.model_config_id = Some(config_id);
2971                self.models = models;
2972                current
2973            });
2974        self.capabilities.supports_models =
2975            self.model_config_id.is_some() && !self.models.is_empty();
2976        if let Some(config_id) = self.model_config_id.clone()
2977            && !self.models.is_empty()
2978        {
2979            self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
2980                slot: self.slot,
2981                config_id,
2982                models: self.models.clone(),
2983                current_model,
2984            }));
2985        }
2986        self.queued_events.push_back(Ok(AgentEvent::Ready {
2987            slot: self.slot,
2988            capabilities: self.capabilities(),
2989        }));
2990        Ok(())
2991    }
2992}
2993
2994fn prompt_resource_paths(prompt: &str) -> Vec<String> {
2995    let characters = prompt.chars().collect::<Vec<_>>();
2996    let mut paths = Vec::new();
2997    let mut index = 0;
2998    while index < characters.len() {
2999        if characters[index] != '@' {
3000            index += 1;
3001            continue;
3002        }
3003        index += 1;
3004        let quoted = characters.get(index) == Some(&'"');
3005        if quoted {
3006            index += 1;
3007        }
3008        let start = index;
3009        while index < characters.len()
3010            && if quoted {
3011                characters[index] != '"'
3012            } else {
3013                !characters[index].is_whitespace()
3014            }
3015        {
3016            index += 1;
3017        }
3018        if index > start {
3019            paths.push(characters[start..index].iter().collect());
3020        }
3021        if quoted && index < characters.len() {
3022            index += 1;
3023        }
3024    }
3025    paths
3026}
3027
3028fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3029    let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3030    for path in prompt_resource_paths(prompt) {
3031        if path.ends_with('/') {
3032            continue;
3033        }
3034        let Ok(resource) = resources::load(cwd, &path) else {
3035            continue;
3036        };
3037        let uri = format!("file://{}", resource.path.display());
3038        let resource_value = if let Some(text) = resource.text {
3039            serde_json::json!({
3040                "uri": uri,
3041                "text": text,
3042                "mimeType": resource.mime_type,
3043            })
3044        } else if let Some(data) = resource.data {
3045            serde_json::json!({
3046                "uri": uri,
3047                "blob": BASE64.encode(data),
3048                "mimeType": resource.mime_type,
3049            })
3050        } else {
3051            continue;
3052        };
3053        blocks.push(serde_json::json!({
3054            "type": "resource",
3055            "resource": resource_value,
3056        }));
3057    }
3058    blocks
3059}
3060
3061#[async_trait]
3062impl AgentAdapter for AcpAdapter {
3063    fn slot(&self) -> RosterSlot {
3064        self.slot
3065    }
3066
3067    fn session_id(&self) -> Option<String> {
3068        self.session_id.clone()
3069    }
3070
3071    fn protocol(&self) -> &'static str {
3072        "acp"
3073    }
3074
3075    fn capabilities(&self) -> AgentCapabilities {
3076        self.capabilities.clone()
3077    }
3078
3079    async fn start(&mut self) -> AdapterResult<()> {
3080        // Use the inherent implementation, which owns the cleanup boundary
3081        // around the multi-step ACP handshake.
3082        AcpAdapter::start(self).await
3083    }
3084
3085    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3086        let session_id = self
3087            .session_id
3088            .as_ref()
3089            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3090        let request_id = self.next_request_id;
3091        self.next_request_id += 1;
3092        let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3093        self.write_json(serde_json::json!({
3094            "jsonrpc": "2.0",
3095            "id": request_id,
3096            "method": "session/prompt",
3097            "params": {
3098                "sessionId": session_id,
3099                "prompt": prompt_blocks,
3100            },
3101        }))
3102        .await?;
3103        self.prompt_request_id = Some(request_id);
3104        Ok(())
3105    }
3106
3107    async fn cancel(&mut self) -> AdapterResult<bool> {
3108        let Some(session_id) = &self.session_id else {
3109            return Ok(false);
3110        };
3111        self.write_json(serde_json::json!({
3112            "jsonrpc": "2.0",
3113            "method": "session/cancel",
3114            "params": {"sessionId": session_id, "_meta": {}},
3115        }))
3116        .await?;
3117        let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3118            loop {
3119                match <Self as AgentAdapter>::next_event(self).await {
3120                    Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3121                    Some(Ok(_)) => {}
3122                    Some(Err(_)) => break,
3123                }
3124            }
3125        })
3126        .await
3127        .is_ok();
3128        if !settled {
3129            // A peer that never acknowledges cancellation cannot safely share
3130            // its stream with the next prompt. Restart the transport while
3131            // preserving a loadable provider session when supported.
3132            self.reload().await?;
3133        }
3134        Ok(true)
3135    }
3136
3137    async fn answer_permission(
3138        &mut self,
3139        request_id: String,
3140        answer: PermissionAnswer,
3141    ) -> AdapterResult<()> {
3142        let id = request_id
3143            .parse::<u64>()
3144            .map(Value::from)
3145            .unwrap_or_else(|_| Value::String(request_id));
3146        let outcome = match answer {
3147            PermissionAnswer::Selected { option_id } => {
3148                serde_json::json!({"outcome": "selected", "optionId": option_id})
3149            }
3150            PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3151        };
3152        self.write_json(serde_json::json!({
3153            "jsonrpc": "2.0",
3154            "id": id,
3155            // RequestPermissionResponse wraps the selected/cancelled
3156            // discriminator in its `outcome` field. Keep this nested shape
3157            // compatible with the Python ACP server and ACP schema.
3158            "result": {"outcome": outcome},
3159        }))
3160        .await
3161    }
3162
3163    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3164        let session_id = self
3165            .session_id
3166            .as_ref()
3167            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3168        let policy = match mode.as_str() {
3169            "plan" => "codeswarm:mode:plan",
3170            "default" | "manual" => "codeswarm:mode:manual",
3171            "accept-edits" => "codeswarm:mode:accept-edits",
3172            "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3173            other => other,
3174        };
3175        let native_mode = crate::policy::resolve(policy, &self.modes)
3176            .map(|mode| mode.id)
3177            .unwrap_or(mode);
3178        let _ = self
3179            .request(
3180                "session/set_mode",
3181                serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3182            )
3183            .await?;
3184        self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3185            slot: self.slot,
3186            current_mode: native_mode,
3187        }));
3188        Ok(())
3189    }
3190
3191    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3192        let session_id = self
3193            .session_id
3194            .clone()
3195            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3196        let config_id = self
3197            .model_config_id
3198            .clone()
3199            .ok_or(AdapterError::Unsupported("set_model"))?;
3200        if !self.models.iter().any(|candidate| candidate.id == model) {
3201            return Err(AdapterError::Protocol(
3202                "model is not advertised by the agent".into(),
3203            ));
3204        }
3205        let _ = self
3206            .request(
3207                "session/set_config_option",
3208                serde_json::json!({
3209                    "sessionId": session_id,
3210                    "configId": config_id,
3211                    "value": model,
3212                }),
3213            )
3214            .await?;
3215        Ok(())
3216    }
3217
3218    async fn reload(&mut self) -> AdapterResult<()> {
3219        // `stop` tears down the process and clears its transport-owned
3220        // session handle. A reload is different from a final shutdown: ACP
3221        // peers advertising `loadSession` must receive the prior ID so the
3222        // replacement process can resume the same conversation.
3223        let session_id = self
3224            .capabilities
3225            .supports_session_load
3226            .then(|| self.session_id.clone())
3227            .flatten();
3228        self.stop().await?;
3229        self.session_id = session_id.clone();
3230        let result = self.start().await;
3231        if result.is_err() {
3232            // `start` cleans up a partially initialized transport by calling
3233            // `stop`, which also clears the handle. Keep it available for a
3234            // subsequent retry after the coordinator reports the failure.
3235            self.session_id = session_id;
3236        }
3237        result
3238    }
3239
3240    async fn stop(&mut self) -> AdapterResult<()> {
3241        let terminals = std::mem::take(&mut self.terminals);
3242        for terminal in terminals.values() {
3243            terminal.stop().await;
3244        }
3245        if let Some(mut child) = self.child.take() {
3246            terminate_child(&mut child).await?;
3247        }
3248        self.reader = None;
3249        self.session_id = None;
3250        self.prompt_request_id = None;
3251        if let Some(task) = self.stderr_task.take() {
3252            task.abort();
3253            let _ = task.await;
3254        }
3255        Ok(())
3256    }
3257
3258    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3259        if let Some(event) = self.queued_events.pop_front() {
3260            return Some(event);
3261        }
3262        loop {
3263            let line = match self.read_line().await {
3264                Ok(line) => line,
3265                Err(error) => return Some(Err(error)),
3266            };
3267            let value: Value = match serde_json::from_str(&line) {
3268                Ok(value) => value,
3269                Err(_) => continue,
3270            };
3271            match self.reject_empty_permission_request(&value).await {
3272                Ok(true) => continue,
3273                Ok(false) => {}
3274                Err(error) => return Some(Err(error)),
3275            }
3276            match self.handle_client_request(&value).await {
3277                Ok(true) => continue,
3278                Ok(false) => {}
3279                Err(error) => return Some(Err(error)),
3280            }
3281            match parse_acp_notification(self.slot, &line) {
3282                Ok(Some(event)) => {
3283                    if let AgentEvent::ModelsReplaced {
3284                        config_id, models, ..
3285                    } = &event
3286                    {
3287                        self.model_config_id = Some(config_id.clone());
3288                        self.models = models.clone();
3289                        self.capabilities.supports_models = !models.is_empty();
3290                    }
3291                    return Some(Ok(event));
3292                }
3293                Ok(None) => {}
3294                Err(error) => return Some(Err(error)),
3295            }
3296            if value.get("id").is_some_and(|id| {
3297                self.prompt_request_id
3298                    .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3299            }) {
3300                if let Some(error) = value.get("error") {
3301                    self.prompt_request_id = None;
3302                    return Some(Err(AdapterError::Protocol(error.to_string())));
3303                }
3304                self.prompt_request_id = None;
3305                return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3306            }
3307        }
3308    }
3309}
3310
3311fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3312    let value: Value =
3313        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3314    let method = value.get("method").and_then(Value::as_str);
3315    if method == Some("session/request_permission") {
3316        let params = value.get("params").cloned().unwrap_or(Value::Null);
3317        let request_id = value
3318            .get("id")
3319            .map(rpc_id_to_string)
3320            .unwrap_or_else(|| "permission".into());
3321        return Ok(parse_permission_event(
3322            slot,
3323            &params,
3324            &request_id,
3325            params.get("options"),
3326        ));
3327    }
3328    if method != Some("session/update") {
3329        return Ok(None);
3330    }
3331    let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3332        return Ok(None);
3333    };
3334    let kind = update.get("sessionUpdate").and_then(Value::as_str);
3335    if kind == Some("config_option_update")
3336        && let Some((config_id, models, current_model)) = parse_model_config(update)
3337    {
3338        return Ok(Some(AgentEvent::ModelsReplaced {
3339            slot,
3340            config_id,
3341            models,
3342            current_model,
3343        }));
3344    }
3345    if kind == Some("request_permission") {
3346        let request_id = update
3347            .get("toolCall")
3348            .and_then(|tool| tool.get("toolCallId"))
3349            .and_then(Value::as_str)
3350            .unwrap_or("permission");
3351        return Ok(parse_permission_event(
3352            slot,
3353            update,
3354            request_id,
3355            update.get("options"),
3356        ));
3357    }
3358    if kind == Some("available_commands_update") {
3359        let commands = update
3360            .get("availableCommands")
3361            .and_then(Value::as_array)
3362            .map(|commands| {
3363                commands
3364                    .iter()
3365                    .filter_map(|command| {
3366                        let name = command.get("name").and_then(Value::as_str)?.trim();
3367                        (!name.is_empty()).then(|| AgentCommand {
3368                            name: name.to_owned(),
3369                        })
3370                    })
3371                    .collect::<Vec<_>>()
3372            })
3373            .unwrap_or_default();
3374        return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3375    }
3376    if kind == Some("current_mode_update") {
3377        if let Some(mode) = update
3378            .get("currentModeId")
3379            .and_then(Value::as_str)
3380            .filter(|mode| !mode.trim().is_empty())
3381        {
3382            return Ok(Some(AgentEvent::ModeUpdated {
3383                slot,
3384                current_mode: mode.to_owned(),
3385            }));
3386        }
3387        return Ok(None);
3388    }
3389    if kind == Some("usage_update") {
3390        let Some(used) = update.get("used").and_then(Value::as_u64) else {
3391            return Ok(None);
3392        };
3393        let Some(size) = update.get("size").and_then(Value::as_u64) else {
3394            return Ok(None);
3395        };
3396        return Ok(Some(AgentEvent::UsageUpdated {
3397            slot,
3398            usage: UsageUpdate { used, size },
3399        }));
3400    }
3401    if let Some(terminal) = parse_terminal_event(update, kind) {
3402        return Ok(Some(AgentEvent::Terminal {
3403            slot,
3404            event: terminal,
3405        }));
3406    }
3407    let text = update
3408        .get("content")
3409        .and_then(|content| content.get("text"))
3410        .and_then(Value::as_str)
3411        .map(str::to_owned);
3412    if kind == Some("user_message_chunk") {
3413        return Ok(text
3414            .filter(|text| !text.is_empty())
3415            .map(|text| AgentEvent::UserText { slot, text }));
3416    }
3417    if kind == Some("agent_message_chunk")
3418        && let Some(mode) = text
3419            .as_deref()
3420            .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3421            .map(str::trim)
3422            .filter(|mode| !mode.is_empty())
3423    {
3424        // Gemini's native ACP bridge historically encoded a mode change as a
3425        // control marker in the message stream. It is state, not transcript
3426        // content; expose it as a normalized catalog replacement instead of
3427        // leaking the marker into the conversation.
3428        return Ok(Some(AgentEvent::ModesReplaced {
3429            slot,
3430            modes: vec![Mode {
3431                id: mode.to_owned(),
3432                label: mode.to_owned(),
3433            }],
3434            current_mode: Some(mode.to_owned()),
3435        }));
3436    }
3437    match (kind, text) {
3438        (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3439            Ok(Some(AgentEvent::Text { slot, text }))
3440        }
3441        (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3442            Ok(Some(AgentEvent::Thought { slot, text }))
3443        }
3444        (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3445            let id = update
3446                .get("toolCallId")
3447                .and_then(Value::as_str)
3448                .unwrap_or("unknown-tool")
3449                .to_owned();
3450            let title = update
3451                .get("title")
3452                .and_then(Value::as_str)
3453                .unwrap_or("Tool call")
3454                .to_owned();
3455            let status = match update.get("status").and_then(Value::as_str) {
3456                Some("completed") => ToolStatus::Completed,
3457                Some("failed") => ToolStatus::Failed,
3458                Some("in_progress") => ToolStatus::Running,
3459                _ => ToolStatus::Pending,
3460            };
3461            Ok(Some(AgentEvent::Tool {
3462                slot,
3463                update: ToolUpdate {
3464                    id,
3465                    title,
3466                    status,
3467                    detail: None,
3468                },
3469            }))
3470        }
3471        _ => Ok(None),
3472    }
3473}
3474
3475fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3476    let config = value
3477        .get("configOptions")?
3478        .as_array()?
3479        .iter()
3480        .find(|option| {
3481            option.get("category").and_then(Value::as_str) == Some("model")
3482                && matches!(
3483                    option.get("type").and_then(Value::as_str),
3484                    Some("select" | "enum")
3485                )
3486        })?;
3487    let config_id = config.get("id")?.as_str()?.to_owned();
3488    let models = config
3489        .get("options")?
3490        .as_array()?
3491        .iter()
3492        .filter_map(|option| {
3493            let id = option.get("value")?.as_str()?.to_owned();
3494            let label = option
3495                .get("name")
3496                .or_else(|| option.get("label"))
3497                .and_then(Value::as_str)
3498                .unwrap_or(&id)
3499                .to_owned();
3500            Some(Mode { id, label })
3501        })
3502        .collect::<Vec<_>>();
3503    (!models.is_empty()).then(|| {
3504        let current = config
3505            .get("currentValue")
3506            .and_then(Value::as_str)
3507            .map(str::to_owned);
3508        (config_id, models, current)
3509    })
3510}
3511
3512/// Normalize terminal lifecycle updates emitted by ACP-compatible bridges and
3513/// native stream adapters. Protocols have used both snake_case update names
3514/// and a nested `terminal` object, so accept either without leaking that
3515/// shape beyond the adapter boundary.
3516fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3517    let nested = value.get("terminal").unwrap_or(value);
3518    let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3519    let id = nested
3520        .get("terminalId")
3521        .or_else(|| nested.get("terminal_id"))
3522        .or_else(|| nested.get("id"))
3523        .and_then(Value::as_str)
3524        .unwrap_or("terminal")
3525        .to_owned();
3526    match kind {
3527        "terminal_created" | "terminal_create" | "terminal_started" => {
3528            let command = nested
3529                .get("command")
3530                .and_then(Value::as_str)
3531                .unwrap_or("")
3532                .to_owned();
3533            Some(TerminalEvent::Created { id, command })
3534        }
3535        "terminal_output" | "terminal_output_chunk" => {
3536            let text = nested
3537                .get("output")
3538                .or_else(|| nested.get("text"))
3539                .and_then(Value::as_str)
3540                .unwrap_or("")
3541                .to_owned();
3542            Some(TerminalEvent::Output { id, text })
3543        }
3544        "terminal_exited" | "terminal_exit" => {
3545            let code = nested
3546                .get("exitCode")
3547                .or_else(|| nested.get("exit_code"))
3548                .or_else(|| nested.get("code"))
3549                .and_then(Value::as_i64)
3550                .unwrap_or(0) as i32;
3551            Some(TerminalEvent::Exited { id, code })
3552        }
3553        "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3554        _ => None,
3555    }
3556}
3557
3558fn parse_permission_event(
3559    slot: RosterSlot,
3560    value: &Value,
3561    request_id: &str,
3562    options: Option<&Value>,
3563) -> Option<AgentEvent> {
3564    let tool = value.get("toolCall").unwrap_or(value);
3565    let title = tool
3566        .get("title")
3567        .and_then(Value::as_str)
3568        .unwrap_or("Agent requests permission")
3569        .to_owned();
3570    let (options, option_ids): (Vec<String>, Vec<String>) = options
3571        .and_then(Value::as_array)
3572        .map(|options| {
3573            options
3574                .iter()
3575                .filter_map(|option| {
3576                    let label = option
3577                        .get("name")
3578                        .or_else(|| option.get("optionId"))
3579                        .and_then(Value::as_str)?
3580                        .to_owned();
3581                    let option_id = option
3582                        .get("optionId")
3583                        .or_else(|| option.get("id"))
3584                        .and_then(Value::as_str)
3585                        .map(str::to_owned)
3586                        .unwrap_or_else(|| label.clone());
3587                    Some((label, option_id))
3588                })
3589                .unzip()
3590        })
3591        .unwrap_or_default();
3592    if options.is_empty() {
3593        return None;
3594    }
3595    Some(AgentEvent::Permission {
3596        slot,
3597        request: PermissionRequest {
3598            id: request_id.to_owned(),
3599            title,
3600            options,
3601            option_ids,
3602        },
3603    })
3604}
3605
3606fn rpc_id_to_string(value: &Value) -> String {
3607    value
3608        .as_str()
3609        .map(str::to_owned)
3610        .or_else(|| value.as_u64().map(|id| id.to_string()))
3611        .unwrap_or_else(|| value.to_string())
3612}
3613
3614#[cfg(test)]
3615mod tests {
3616    use super::{
3617        AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3618        RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3619        parse_model_config, prompt_content_blocks, read_bounded_line,
3620    };
3621    #[cfg(target_os = "linux")]
3622    use super::{isolate_process_group, terminate_child};
3623    use crate::TerminalEvent;
3624    use crate::{
3625        AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3626        persistence::SessionMetadataStore,
3627        relay::{DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3628    };
3629    use async_trait::async_trait;
3630    use serde_json::Value;
3631    use std::sync::{
3632        Arc, Mutex,
3633        atomic::{AtomicUsize, Ordering},
3634    };
3635
3636    fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3637        let nonce = std::time::SystemTime::now()
3638            .duration_since(std::time::UNIX_EPOCH)
3639            .expect("clock")
3640            .as_nanos();
3641        std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3642    }
3643
3644    #[test]
3645    fn malformed_file_writes_preserve_existing_content() {
3646        let root = unique_test_path("codeswarm-write-validation", "dir");
3647        std::fs::create_dir_all(&root).unwrap();
3648        let file = root.join("keep.txt");
3649        std::fs::write(&file, "valuable content").unwrap();
3650        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3651        for content in [
3652            Value::Null,
3653            serde_json::json!(false),
3654            serde_json::json!(42),
3655            serde_json::json!([]),
3656        ] {
3657            assert!(
3658                adapter
3659                    .write_workspace_text(
3660                        &serde_json::json!({"path":"keep.txt", "content": content})
3661                    )
3662                    .is_err()
3663            );
3664            assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3665        }
3666        assert!(
3667            adapter
3668                .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
3669                .is_err()
3670        );
3671        assert!(
3672            adapter
3673                .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
3674                .is_err()
3675        );
3676        assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3677        adapter
3678            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
3679            .unwrap();
3680        assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
3681        adapter
3682            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
3683            .unwrap();
3684        assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
3685        std::fs::remove_dir_all(root).unwrap();
3686    }
3687
3688    #[tokio::test]
3689    async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
3690        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
3691        let mut adapter = AcpAdapter::new(
3692            0,
3693            std::env::current_dir().unwrap(),
3694            "sh",
3695            vec!["-c".into(), script.into()],
3696        );
3697        adapter.start().await.unwrap();
3698        let error = adapter
3699            .request_with_timeout(
3700                "session/set_mode",
3701                serde_json::json!({}),
3702                std::time::Duration::from_millis(10),
3703            )
3704            .await
3705            .unwrap_err();
3706        assert!(error.to_string().contains("session/set_mode timed out"));
3707        adapter.stop().await.unwrap();
3708        assert!(adapter.child.is_none());
3709    }
3710
3711    #[tokio::test]
3712    async fn goals_reach_every_roster_slot_without_native_goal_support() {
3713        use crate::goal::GoalCommand;
3714        let hosts = (0..3)
3715            .map(|slot| {
3716                AdapterHost::new(
3717                    Box::new(ScriptedAdapter::new(
3718                        slot,
3719                        AgentCapabilities::default(),
3720                        [
3721                            AgentEvent::TurnComplete { slot },
3722                            AgentEvent::TurnComplete { slot },
3723                        ],
3724                    )),
3725                    None,
3726                )
3727            })
3728            .collect();
3729        let mut relay = RelayHost::new(hosts, 10).unwrap();
3730        relay.start().await.unwrap();
3731        let task = relay
3732            .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
3733            .unwrap()
3734            .unwrap();
3735        relay.relay_mut().enqueue_human(task, Some(0));
3736        for slot in 0..3 {
3737            relay.run_turn("", 0).await.unwrap();
3738            let (actual, prompt) = relay.dispatches().last().unwrap();
3739            assert_eq!(*actual, slot);
3740            assert!(prompt.contains("Active shared goal: Ship the settings screen"));
3741        }
3742        let snapshot = relay.session_metadata();
3743        let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
3744        assert!(restored.is_some());
3745        relay.restore_goal(restored);
3746        relay.reload(0).await.unwrap();
3747        relay.run_turn("", 0).await.unwrap();
3748        assert!(
3749            relay
3750                .dispatches()
3751                .last()
3752                .unwrap()
3753                .1
3754                .contains("Active shared goal: Ship the settings screen")
3755        );
3756        relay.apply_goal(GoalCommand::Done).unwrap();
3757        relay.run_turn("", 0).await.unwrap();
3758        assert!(
3759            relay
3760                .dispatches()
3761                .last()
3762                .unwrap()
3763                .1
3764                .contains("No active shared goal")
3765        );
3766        relay.apply_goal(GoalCommand::Clear).unwrap();
3767        assert!(relay.session_metadata().get("goal").unwrap().is_null());
3768    }
3769
3770    #[tokio::test]
3771    async fn replacement_agent_receives_task_after_public_journal_pruning() {
3772        let hosts = (0..2)
3773            .map(|slot| {
3774                AdapterHost::new(
3775                    Box::new(ScriptedAdapter::new(
3776                        slot,
3777                        AgentCapabilities::default(),
3778                        [
3779                            AgentEvent::Text {
3780                                slot,
3781                                text: "progress".into(),
3782                            },
3783                            AgentEvent::TurnComplete { slot },
3784                            AgentEvent::Text {
3785                                slot,
3786                                text: "more progress".into(),
3787                            },
3788                            AgentEvent::TurnComplete { slot },
3789                        ],
3790                    )),
3791                    None,
3792                )
3793            })
3794            .collect();
3795        let mut relay = RelayHost::new(hosts, 10).unwrap();
3796        relay.start().await.unwrap();
3797        relay
3798            .relay_mut()
3799            .enqueue_human("Fix the login bug", Some(0));
3800        relay.run_turn("", 0).await.unwrap();
3801        relay.run_turn("", 0).await.unwrap();
3802        relay.run_turn("", 0).await.unwrap();
3803        relay.reload(1).await.unwrap();
3804        assert!(
3805            !relay
3806                .relay_mut()
3807                .unseen_context(1)
3808                .contains("Fix the login bug")
3809        );
3810        relay.run_turn("", 0).await.unwrap();
3811        assert!(
3812            relay
3813                .dispatches()
3814                .last()
3815                .unwrap()
3816                .1
3817                .contains("Shared task:\nFix the login bug")
3818        );
3819    }
3820
3821    #[cfg(target_os = "linux")]
3822    #[tokio::test]
3823    async fn termination_kills_only_the_verified_isolated_child_group() {
3824        use nix::unistd::{Pid, getpgid, getpgrp};
3825        use tokio::io::{AsyncBufReadExt, BufReader};
3826
3827        let own_group = getpgrp();
3828        let mut command = tokio::process::Command::new("sh");
3829        isolate_process_group(&mut command);
3830        command
3831            .arg("-c")
3832            .arg("sleep 60 & echo $!; wait")
3833            .stdout(std::process::Stdio::piped());
3834        let mut child = command.spawn().expect("spawn isolated shell");
3835        let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
3836        assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
3837        assert_ne!(leader, own_group);
3838
3839        let stdout = child.stdout.take().expect("child stdout");
3840        let mut lines = BufReader::new(stdout).lines();
3841        let descendant = lines
3842            .next_line()
3843            .await
3844            .expect("read descendant pid")
3845            .expect("descendant pid")
3846            .parse::<i32>()
3847            .expect("numeric descendant pid");
3848        let descendant = Pid::from_raw(descendant);
3849        assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
3850
3851        terminate_child(&mut child).await.expect("terminate group");
3852        for _ in 0..100 {
3853            if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
3854                return;
3855            }
3856            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3857        }
3858        panic!("descendant {descendant} survived isolated group termination");
3859    }
3860
3861    #[test]
3862    fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
3863        assert_eq!(
3864            parse_command_line(r#"npx -y "@agentclientprotocol/codex-acp" --flag 'two words'"#),
3865            Ok((
3866                "npx".into(),
3867                vec![
3868                    "-y".into(),
3869                    "@agentclientprotocol/codex-acp".into(),
3870                    "--flag".into(),
3871                    "two words".into(),
3872                ]
3873            ),)
3874        );
3875        assert_eq!(
3876            parse_command_line(r#"agent "" escaped\ argument"#),
3877            Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
3878        );
3879    }
3880
3881    #[test]
3882    fn acp_prompt_expands_safe_at_path_resources() {
3883        let root = unique_test_path("codeswarm-prompt-resource", "dir");
3884        std::fs::create_dir_all(&root).expect("workspace");
3885        std::fs::write(root.join("note.md"), "resource text").expect("resource");
3886        let blocks = prompt_content_blocks(&root, "inspect @note.md");
3887        assert_eq!(blocks[0]["type"], "text");
3888        assert_eq!(blocks[0]["text"], "inspect @note.md");
3889        assert_eq!(blocks[1]["type"], "resource");
3890        assert_eq!(blocks[1]["resource"]["text"], "resource text");
3891        assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
3892        std::fs::remove_dir_all(root).expect("cleanup workspace");
3893    }
3894
3895    #[tokio::test]
3896    async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
3897        let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
3898        bytes.push(b'\n');
3899        let mut reader = tokio::io::BufReader::new(bytes.as_slice());
3900        assert!(matches!(
3901            read_bounded_line(&mut reader).await,
3902            Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
3903        ));
3904    }
3905
3906    #[test]
3907    fn rejects_malformed_configured_commands_before_spawn() {
3908        assert_eq!(
3909            parse_command_line("agent 'unfinished"),
3910            Err(super::CommandParseError::UnterminatedQuote)
3911        );
3912        assert_eq!(
3913            parse_command_line("agent\\"),
3914            Err(super::CommandParseError::TrailingEscape)
3915        );
3916        assert_eq!(
3917            parse_command_line("   \t"),
3918            Err(super::CommandParseError::Empty)
3919        );
3920    }
3921
3922    #[derive(Debug)]
3923    struct PendingAdapter {
3924        slot: usize,
3925        hang_on_cancel: bool,
3926    }
3927
3928    #[derive(Debug)]
3929    struct ConcurrentStartAdapter {
3930        slot: usize,
3931        barrier: Arc<tokio::sync::Barrier>,
3932    }
3933
3934    #[async_trait]
3935    impl AgentAdapter for ConcurrentStartAdapter {
3936        fn slot(&self) -> usize {
3937            self.slot
3938        }
3939
3940        fn capabilities(&self) -> AgentCapabilities {
3941            AgentCapabilities::default()
3942        }
3943
3944        async fn start(&mut self) -> super::AdapterResult<()> {
3945            self.barrier.wait().await;
3946            Ok(())
3947        }
3948
3949        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3950            Ok(())
3951        }
3952
3953        async fn cancel(&mut self) -> super::AdapterResult<bool> {
3954            Ok(true)
3955        }
3956
3957        async fn answer_permission(
3958            &mut self,
3959            _request_id: String,
3960            _answer: PermissionAnswer,
3961        ) -> super::AdapterResult<()> {
3962            Ok(())
3963        }
3964
3965        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3966            Ok(())
3967        }
3968
3969        async fn reload(&mut self) -> super::AdapterResult<()> {
3970            Ok(())
3971        }
3972
3973        async fn stop(&mut self) -> super::AdapterResult<()> {
3974            Ok(())
3975        }
3976
3977        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3978            std::future::pending().await
3979        }
3980    }
3981
3982    #[derive(Debug)]
3983    struct PermissionBlockingAdapter {
3984        slot: usize,
3985        phase: u8,
3986    }
3987
3988    #[async_trait]
3989    impl AgentAdapter for PermissionBlockingAdapter {
3990        fn slot(&self) -> usize {
3991            self.slot
3992        }
3993
3994        fn capabilities(&self) -> AgentCapabilities {
3995            AgentCapabilities {
3996                supports_permissions: true,
3997                ..AgentCapabilities::default()
3998            }
3999        }
4000
4001        async fn start(&mut self) -> super::AdapterResult<()> {
4002            Ok(())
4003        }
4004
4005        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4006            Ok(())
4007        }
4008
4009        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4010            Ok(true)
4011        }
4012
4013        async fn answer_permission(
4014            &mut self,
4015            request_id: String,
4016            answer: PermissionAnswer,
4017        ) -> super::AdapterResult<()> {
4018            if self.phase != 1 || request_id != "permission-1" {
4019                return Err(super::AdapterError::Protocol(
4020                    "unexpected permission response".into(),
4021                ));
4022            }
4023            assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4024            self.phase = 2;
4025            Ok(())
4026        }
4027
4028        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4029            Ok(())
4030        }
4031
4032        async fn reload(&mut self) -> super::AdapterResult<()> {
4033            Ok(())
4034        }
4035
4036        async fn stop(&mut self) -> super::AdapterResult<()> {
4037            Ok(())
4038        }
4039
4040        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4041            match self.phase {
4042                0 => {
4043                    self.phase = 1;
4044                    Some(Ok(AgentEvent::Permission {
4045                        slot: self.slot,
4046                        request: crate::PermissionRequest {
4047                            id: "permission-1".into(),
4048                            title: "Allow?".into(),
4049                            options: vec!["Allow".into()],
4050                            option_ids: vec!["allow".into()],
4051                        },
4052                    }))
4053                }
4054                1 => std::future::pending().await,
4055                _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4056            }
4057        }
4058    }
4059
4060    #[async_trait]
4061    impl AgentAdapter for PendingAdapter {
4062        fn slot(&self) -> usize {
4063            self.slot
4064        }
4065
4066        fn capabilities(&self) -> AgentCapabilities {
4067            AgentCapabilities {
4068                supports_cancel: true,
4069                ..AgentCapabilities::default()
4070            }
4071        }
4072
4073        async fn start(&mut self) -> super::AdapterResult<()> {
4074            Ok(())
4075        }
4076
4077        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4078            Ok(())
4079        }
4080
4081        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4082            if self.hang_on_cancel {
4083                return std::future::pending().await;
4084            }
4085            Ok(true)
4086        }
4087
4088        async fn answer_permission(
4089            &mut self,
4090            _request_id: String,
4091            _answer: PermissionAnswer,
4092        ) -> super::AdapterResult<()> {
4093            Ok(())
4094        }
4095
4096        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4097            Ok(())
4098        }
4099
4100        async fn reload(&mut self) -> super::AdapterResult<()> {
4101            Ok(())
4102        }
4103
4104        async fn stop(&mut self) -> super::AdapterResult<()> {
4105            Ok(())
4106        }
4107
4108        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4109            std::future::pending().await
4110        }
4111    }
4112
4113    #[derive(Debug)]
4114    struct StopTrackingAdapter {
4115        slot: usize,
4116        stopped: Arc<AtomicUsize>,
4117        fail_stop: bool,
4118    }
4119
4120    #[derive(Debug)]
4121    struct ModeOrderAdapter {
4122        slot: usize,
4123        log: Arc<Mutex<Vec<String>>>,
4124        phase: u8,
4125    }
4126
4127    #[derive(Debug)]
4128    struct StartupAcpAdapter {
4129        slot: usize,
4130        events: std::collections::VecDeque<AgentEvent>,
4131    }
4132
4133    impl StartupAcpAdapter {
4134        fn new(slot: usize) -> Self {
4135            Self {
4136                slot,
4137                events: [
4138                    AgentEvent::ModesReplaced {
4139                        slot,
4140                        modes: vec![Mode {
4141                            id: "full-access".into(),
4142                            label: "Auto pilot".into(),
4143                        }],
4144                        current_mode: Some("full-access".into()),
4145                    },
4146                    AgentEvent::Ready {
4147                        slot,
4148                        capabilities: AgentCapabilities {
4149                            supports_modes: true,
4150                            ..AgentCapabilities::default()
4151                        },
4152                    },
4153                ]
4154                .into(),
4155            }
4156        }
4157    }
4158
4159    #[async_trait]
4160    impl AgentAdapter for StartupAcpAdapter {
4161        fn slot(&self) -> usize {
4162            self.slot
4163        }
4164
4165        fn protocol(&self) -> &'static str {
4166            "acp"
4167        }
4168
4169        fn capabilities(&self) -> AgentCapabilities {
4170            AgentCapabilities {
4171                supports_modes: true,
4172                ..AgentCapabilities::default()
4173            }
4174        }
4175
4176        async fn start(&mut self) -> super::AdapterResult<()> {
4177            Ok(())
4178        }
4179
4180        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4181            Ok(())
4182        }
4183
4184        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4185            Ok(true)
4186        }
4187
4188        async fn answer_permission(
4189            &mut self,
4190            _request_id: String,
4191            _answer: PermissionAnswer,
4192        ) -> super::AdapterResult<()> {
4193            Ok(())
4194        }
4195
4196        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4197            Ok(())
4198        }
4199
4200        async fn reload(&mut self) -> super::AdapterResult<()> {
4201            self.events = Self::new(self.slot).events;
4202            Ok(())
4203        }
4204
4205        async fn stop(&mut self) -> super::AdapterResult<()> {
4206            Ok(())
4207        }
4208
4209        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4210            self.events.pop_front().map(Ok)
4211        }
4212    }
4213
4214    #[async_trait]
4215    impl AgentAdapter for ModeOrderAdapter {
4216        fn slot(&self) -> usize {
4217            self.slot
4218        }
4219
4220        fn capabilities(&self) -> AgentCapabilities {
4221            AgentCapabilities {
4222                supports_modes: true,
4223                ..AgentCapabilities::default()
4224            }
4225        }
4226
4227        async fn start(&mut self) -> super::AdapterResult<()> {
4228            self.log.lock().expect("log").push("start".into());
4229            Ok(())
4230        }
4231
4232        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4233            self.log.lock().expect("log").push("prompt".into());
4234            Ok(())
4235        }
4236
4237        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4238            Ok(true)
4239        }
4240
4241        async fn answer_permission(
4242            &mut self,
4243            _request_id: String,
4244            _answer: PermissionAnswer,
4245        ) -> super::AdapterResult<()> {
4246            Ok(())
4247        }
4248
4249        async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4250            self.log.lock().expect("log").push(format!("mode:{mode}"));
4251            Ok(())
4252        }
4253
4254        async fn reload(&mut self) -> super::AdapterResult<()> {
4255            self.log.lock().expect("log").push("reload".into());
4256            self.phase = 0;
4257            Ok(())
4258        }
4259
4260        async fn stop(&mut self) -> super::AdapterResult<()> {
4261            Ok(())
4262        }
4263
4264        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4265            match self.phase {
4266                0 => {
4267                    self.phase = 1;
4268                    Some(Ok(AgentEvent::ModesReplaced {
4269                        slot: self.slot,
4270                        modes: vec![Mode {
4271                            id: "yolo".into(),
4272                            label: "YOLO".into(),
4273                        }],
4274                        current_mode: None,
4275                    }))
4276                }
4277                1 => {
4278                    self.phase = 2;
4279                    Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4280                }
4281                _ => std::future::pending().await,
4282            }
4283        }
4284    }
4285
4286    #[async_trait]
4287    impl AgentAdapter for StopTrackingAdapter {
4288        fn slot(&self) -> usize {
4289            self.slot
4290        }
4291
4292        fn capabilities(&self) -> AgentCapabilities {
4293            AgentCapabilities::default()
4294        }
4295
4296        async fn start(&mut self) -> super::AdapterResult<()> {
4297            Ok(())
4298        }
4299
4300        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4301            Ok(())
4302        }
4303
4304        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4305            Ok(false)
4306        }
4307
4308        async fn answer_permission(
4309            &mut self,
4310            _request_id: String,
4311            _answer: PermissionAnswer,
4312        ) -> super::AdapterResult<()> {
4313            Ok(())
4314        }
4315
4316        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4317            Ok(())
4318        }
4319
4320        async fn reload(&mut self) -> super::AdapterResult<()> {
4321            Ok(())
4322        }
4323
4324        async fn stop(&mut self) -> super::AdapterResult<()> {
4325            self.stopped.fetch_add(1, Ordering::Relaxed);
4326            if self.fail_stop {
4327                Err(super::AdapterError::Transport("stop failed".into()))
4328            } else {
4329                Ok(())
4330            }
4331        }
4332
4333        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4334            None
4335        }
4336    }
4337
4338    /// Startup can fail after an adapter has allocated resources.  Keep a
4339    /// fixture that records whether the failing adapter itself receives the
4340    /// cleanup call, not just the already-started peers.
4341    #[derive(Debug)]
4342    struct FailingStartAdapter {
4343        slot: usize,
4344        stopped: Arc<AtomicUsize>,
4345    }
4346
4347    #[async_trait]
4348    impl AgentAdapter for FailingStartAdapter {
4349        fn slot(&self) -> usize {
4350            self.slot
4351        }
4352
4353        fn capabilities(&self) -> AgentCapabilities {
4354            AgentCapabilities::default()
4355        }
4356
4357        async fn start(&mut self) -> super::AdapterResult<()> {
4358            Err(super::AdapterError::Spawn("startup failed".into()))
4359        }
4360
4361        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4362            Ok(())
4363        }
4364
4365        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4366            Ok(false)
4367        }
4368
4369        async fn answer_permission(
4370            &mut self,
4371            _request_id: String,
4372            _answer: PermissionAnswer,
4373        ) -> super::AdapterResult<()> {
4374            Ok(())
4375        }
4376
4377        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4378            Ok(())
4379        }
4380
4381        async fn reload(&mut self) -> super::AdapterResult<()> {
4382            Ok(())
4383        }
4384
4385        async fn stop(&mut self) -> super::AdapterResult<()> {
4386            self.stopped.fetch_add(1, Ordering::Relaxed);
4387            Ok(())
4388        }
4389
4390        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4391            None
4392        }
4393    }
4394
4395    #[tokio::test]
4396    async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4397        let stopped = Arc::new(AtomicUsize::new(0));
4398        let relay = RelayHost::new(
4399            vec![
4400                AdapterHost::new(
4401                    Box::new(StopTrackingAdapter {
4402                        slot: 0,
4403                        stopped: Arc::clone(&stopped),
4404                        fail_stop: true,
4405                    }),
4406                    None,
4407                ),
4408                AdapterHost::new(
4409                    Box::new(StopTrackingAdapter {
4410                        slot: 1,
4411                        stopped: Arc::clone(&stopped),
4412                        fail_stop: false,
4413                    }),
4414                    None,
4415                ),
4416            ],
4417            4,
4418        )
4419        .expect("relay");
4420        let mut relay = relay;
4421
4422        let error = relay.stop().await.expect_err("first stop failure");
4423        assert!(error.to_string().contains("stop failed"));
4424        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4425    }
4426
4427    #[tokio::test]
4428    async fn relay_start_cleans_up_the_adapter_that_failed_startup() {
4429        let stopped = Arc::new(AtomicUsize::new(0));
4430        let mut relay = RelayHost::new(
4431            vec![
4432                AdapterHost::new(
4433                    Box::new(StopTrackingAdapter {
4434                        slot: 0,
4435                        stopped: Arc::clone(&stopped),
4436                        fail_stop: false,
4437                    }),
4438                    None,
4439                ),
4440                AdapterHost::new(
4441                    Box::new(FailingStartAdapter {
4442                        slot: 1,
4443                        stopped: Arc::clone(&stopped),
4444                    }),
4445                    None,
4446                ),
4447            ],
4448            4,
4449        )
4450        .expect("relay");
4451
4452        assert!(relay.start().await.is_err());
4453        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4454    }
4455
4456    #[test]
4457    fn parses_acp_text_without_ui_dependency() {
4458        let event = parse_acp_notification(
4459            2,
4460            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4461        )
4462        .expect("valid ACP")
4463        .expect("text event");
4464        assert_eq!(
4465            event,
4466            AgentEvent::Text {
4467                slot: 2,
4468                text: "hello".into(),
4469            }
4470        );
4471    }
4472
4473    #[test]
4474    fn parses_acp_state_notifications_at_the_adapter_boundary() {
4475        let commands = parse_acp_notification(
4476            3,
4477            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4478        )
4479        .expect("valid ACP")
4480        .expect("commands event");
4481        assert_eq!(
4482            commands,
4483            AgentEvent::CommandsReplaced {
4484                slot: 3,
4485                commands: vec![crate::AgentCommand {
4486                    name: "review".into()
4487                }]
4488            }
4489        );
4490
4491        let mode = parse_acp_notification(
4492            3,
4493            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4494        )
4495        .expect("valid ACP")
4496        .expect("mode event");
4497        assert_eq!(
4498            mode,
4499            AgentEvent::ModeUpdated {
4500                slot: 3,
4501                current_mode: "review".into()
4502            }
4503        );
4504
4505        let usage = parse_acp_notification(
4506            3,
4507            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4508        )
4509        .expect("valid ACP")
4510        .expect("usage event");
4511        assert_eq!(
4512            usage,
4513            AgentEvent::UsageUpdated {
4514                slot: 3,
4515                usage: crate::UsageUpdate {
4516                    used: 4200,
4517                    size: 128000
4518                }
4519            }
4520        );
4521
4522        let models = parse_acp_notification(
4523            3,
4524            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"}]}]}}}"#,
4525        )
4526        .expect("valid ACP")
4527        .expect("models event");
4528        assert!(matches!(
4529            models,
4530            AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4531                if models.len() == 2 && current_model.as_deref() == Some("smart")
4532        ));
4533        assert_eq!(
4534            parse_model_config(&serde_json::json!({
4535                "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4536            })),
4537            None
4538        );
4539
4540        let user = parse_acp_notification(
4541            3,
4542            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4543        )
4544        .expect("valid ACP")
4545        .expect("user event");
4546        assert_eq!(
4547            user,
4548            AgentEvent::UserText {
4549                slot: 3,
4550                text: "context".into()
4551            }
4552        );
4553    }
4554
4555    #[test]
4556    fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4557        let event = parse_acp_notification(
4558            0,
4559            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4560        )
4561        .expect("valid ACP")
4562        .expect("mode event");
4563        assert!(matches!(
4564            event,
4565            AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4566                if mode == "yolo" && modes[0].id == "yolo"
4567        ));
4568    }
4569
4570    #[test]
4571    fn parses_native_agy_text_without_acp_bridge() {
4572        let event = parse_agy_line(
4573            1,
4574            r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4575        )
4576        .expect("valid stream-json")
4577        .expect("text event");
4578        assert_eq!(
4579            event,
4580            AgentEvent::Text {
4581                slot: 1,
4582                text: "hello".into(),
4583            }
4584        );
4585    }
4586
4587    #[test]
4588    fn parses_tool_lifecycle_from_each_protocol() {
4589        let agy = parse_agy_line(
4590            1,
4591            r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
4592        )
4593        .expect("valid native tool")
4594        .expect("tool event");
4595        assert!(matches!(
4596            agy,
4597            AgentEvent::Tool {
4598                update: crate::ToolUpdate {
4599                    status: ToolStatus::Completed,
4600                    ..
4601                },
4602                ..
4603            }
4604        ));
4605
4606        let acp = parse_acp_notification(
4607            1,
4608            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
4609        )
4610        .expect("valid ACP tool")
4611        .expect("tool event");
4612        assert!(matches!(
4613            acp,
4614            AgentEvent::Tool {
4615                update: crate::ToolUpdate {
4616                    status: ToolStatus::Failed,
4617                    ..
4618                },
4619                ..
4620            }
4621        ));
4622    }
4623
4624    #[test]
4625    fn parses_terminal_lifecycle_from_acp_and_native_events() {
4626        let created = parse_acp_notification(
4627            0,
4628            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
4629        )
4630        .expect("valid ACP terminal")
4631        .expect("terminal event");
4632        assert_eq!(
4633            created,
4634            AgentEvent::Terminal {
4635                slot: 0,
4636                event: TerminalEvent::Created {
4637                    id: "term-1".into(),
4638                    command: "cargo test".into(),
4639                },
4640            }
4641        );
4642        let output = parse_agy_line(
4643            1,
4644            r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
4645        )
4646        .expect("valid native terminal")
4647        .expect("terminal event");
4648        assert_eq!(
4649            output,
4650            AgentEvent::Terminal {
4651                slot: 1,
4652                event: TerminalEvent::Output {
4653                    id: "term-1".into(),
4654                    text: "ok\n".into(),
4655                },
4656            }
4657        );
4658        let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
4659            .expect("valid native release")
4660            .expect("terminal event");
4661        assert!(matches!(
4662            released,
4663            AgentEvent::Terminal {
4664                event: TerminalEvent::Released { id },
4665                ..
4666            } if id == "term-1"
4667        ));
4668    }
4669
4670    #[test]
4671    fn parses_acp_permission_requests() {
4672        let event = parse_acp_notification(
4673            0,
4674            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"}]}}}"#,
4675        )
4676        .expect("valid permission")
4677        .expect("permission event");
4678        assert!(matches!(
4679            event,
4680            AgentEvent::Permission { request, .. }
4681                if request.id == "t1"
4682                    && request.title == "Write file"
4683                    && request.options == ["Allow once", "Reject"]
4684                    && request.option_ids == ["allow-once", "reject"]
4685        ));
4686    }
4687
4688    #[test]
4689    fn parses_acp_permission_request_as_json_rpc_request() {
4690        let event = parse_acp_notification(
4691            2,
4692            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
4693        )
4694        .expect("valid permission request")
4695        .expect("permission event");
4696        assert!(matches!(
4697            event,
4698            AgentEvent::Permission { request, .. }
4699                if request.id == "17"
4700                    && request.title == "Write file"
4701                    && request.options == ["allow-once", "reject"]
4702                    && request.option_ids == ["allow-once", "reject"]
4703        ));
4704    }
4705
4706    #[tokio::test]
4707    async fn native_adapter_explicitly_rejects_permission_answers() {
4708        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4709        assert_eq!(
4710            adapter
4711                .answer_permission(
4712                    "request".into(),
4713                    PermissionAnswer::Selected {
4714                        option_id: "allow".into()
4715                    },
4716                )
4717                .await,
4718            Err(super::AdapterError::Unsupported("permission answer"))
4719        );
4720    }
4721
4722    #[tokio::test]
4723    async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
4724        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4725        adapter
4726            .set_mode("full-access".into())
4727            .await
4728            .expect("auto-pilot alias");
4729        assert!(matches!(
4730            adapter.next_event().await,
4731            Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
4732        ));
4733    }
4734
4735    #[tokio::test]
4736    async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
4737        let script_path = unique_test_path("codeswarm-agy-session", "sh");
4738        std::fs::write(
4739            &script_path,
4740            "#!/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",
4741        )
4742        .expect("write native test script");
4743        let mut adapter = AgyAdapter::new(
4744            0,
4745            std::env::current_dir().expect("cwd"),
4746            format!("sh {}", script_path.display()),
4747        );
4748        adapter.start().await.expect("start native adapter");
4749        // Startup emits its mode catalog and readiness before a turn.
4750        assert!(adapter.next_event().await.is_some());
4751        assert!(adapter.next_event().await.is_some());
4752        adapter
4753            .send_prompt("first".into())
4754            .await
4755            .expect("first prompt");
4756        while !matches!(
4757            adapter.next_event().await,
4758            Some(Ok(AgentEvent::TurnComplete { .. }))
4759        ) {}
4760        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4761        adapter
4762            .send_prompt("follow up".into())
4763            .await
4764            .expect("follow-up prompt");
4765        while !matches!(
4766            adapter.next_event().await,
4767            Some(Ok(AgentEvent::TurnComplete { .. }))
4768        ) {}
4769        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4770        adapter.stop().await.expect("stop native adapter");
4771        std::fs::remove_file(script_path).expect("cleanup native script");
4772    }
4773
4774    #[tokio::test]
4775    async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
4776        let script_path = unique_test_path("codeswarm-agy-failure", "sh");
4777        std::fs::write(
4778            &script_path,
4779            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
4780        )
4781        .expect("write native test script");
4782        let mut adapter = AgyAdapter::new(
4783            0,
4784            std::env::current_dir().expect("cwd"),
4785            format!("sh {}", script_path.display()),
4786        );
4787        adapter.start().await.expect("start native adapter");
4788        assert!(adapter.next_event().await.is_some());
4789        assert!(adapter.next_event().await.is_some());
4790        adapter.send_prompt("fail".into()).await.expect("prompt");
4791        assert!(matches!(
4792            adapter.next_event().await,
4793            Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
4794                if detail == "agent failed"
4795        ));
4796        adapter.stop().await.expect("stop native adapter");
4797        std::fs::remove_file(script_path).expect("cleanup native script");
4798    }
4799
4800    #[tokio::test]
4801    async fn acp_adapter_initializes_session_and_completes_a_prompt() {
4802        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"}}'"#;
4803        let cwd = std::env::current_dir().expect("cwd");
4804        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4805        adapter.start().await.expect("initialize");
4806        assert!(matches!(
4807            adapter.next_event().await,
4808            Some(Ok(AgentEvent::ModesReplaced { .. }))
4809        ));
4810        assert!(matches!(
4811            adapter.next_event().await,
4812            Some(Ok(AgentEvent::Ready { .. }))
4813        ));
4814        adapter.send_prompt("hello".into()).await.expect("prompt");
4815        assert!(matches!(
4816            adapter.next_event().await,
4817            Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
4818        ));
4819        assert!(matches!(
4820            adapter.next_event().await,
4821            Some(Ok(AgentEvent::TurnComplete { .. }))
4822        ));
4823    }
4824
4825    #[tokio::test]
4826    async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
4827        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"}}'"#;
4828        let cwd = std::env::current_dir().expect("cwd");
4829        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
4830        adapter.start().await.expect("initialize");
4831        assert!(adapter.next_event().await.is_some());
4832        assert!(adapter.next_event().await.is_some());
4833
4834        for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
4835            adapter.send_prompt(prompt.into()).await.expect("prompt");
4836            assert!(matches!(
4837                adapter.next_event().await,
4838                Some(Ok(AgentEvent::Text { text, .. })) if text == expected
4839            ));
4840            assert!(matches!(
4841                adapter.next_event().await,
4842                Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
4843            ));
4844        }
4845        adapter.stop().await.expect("stop");
4846    }
4847
4848    #[tokio::test]
4849    async fn empty_acp_mode_catalog_disables_mode_control() {
4850        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":[]}}}'"#;
4851        let cwd = std::env::current_dir().expect("cwd");
4852        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4853        adapter.start().await.expect("initialize");
4854        assert!(!adapter.capabilities().supports_modes);
4855        assert!(matches!(
4856            adapter.next_event().await,
4857            Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
4858        ));
4859        assert!(matches!(
4860            adapter.next_event().await,
4861            Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
4862        ));
4863        adapter.stop().await.expect("stop");
4864    }
4865
4866    #[tokio::test]
4867    async fn acp_models_are_discovered_live_and_changed_through_session_config() {
4868        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"#;
4869        let cwd = std::env::current_dir().expect("cwd");
4870        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4871        adapter.start().await.expect("initialize");
4872        assert!(adapter.capabilities().supports_models);
4873        assert!(matches!(
4874            adapter.next_event().await,
4875            Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
4876                if config_id == "model"
4877                    && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
4878                    && current_model.as_deref() == Some("fast")
4879        ));
4880        assert!(matches!(
4881            adapter.next_event().await,
4882            Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
4883        ));
4884        adapter.set_model("smart".into()).await.expect("set model");
4885        assert!(adapter.set_model("invented".into()).await.is_err());
4886        adapter.stop().await.expect("stop");
4887    }
4888
4889    #[tokio::test]
4890    async fn acp_mode_change_is_acknowledged_without_provider_notification() {
4891        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":{}}'"#;
4892        let cwd = std::env::current_dir().expect("cwd");
4893        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4894        adapter.start().await.expect("initialize");
4895        adapter
4896            .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
4897            .await
4898            .expect("set mode");
4899        assert!(matches!(
4900            adapter.next_event().await,
4901            Some(Ok(AgentEvent::ModesReplaced { .. }))
4902        ));
4903        assert!(matches!(
4904            adapter.next_event().await,
4905            Some(Ok(AgentEvent::Ready { .. }))
4906        ));
4907        assert!(matches!(
4908            adapter.next_event().await,
4909            Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
4910        ));
4911        adapter.stop().await.expect("stop");
4912    }
4913
4914    #[tokio::test]
4915    async fn acp_reload_preserves_a_loadable_session_id() {
4916        let cwd = std::env::current_dir().expect("cwd");
4917        let mut adapter = AcpAdapter::with_session_id(
4918            0,
4919            cwd,
4920            "__codeswarm_missing_acp_for_reload_test__",
4921            Vec::new(),
4922            "saved-session",
4923        );
4924        adapter.capabilities.supports_session_load = true;
4925        // A failed replacement process still must not erase the session ID:
4926        // the coordinator can report the startup error and offer another
4927        // reload, preserving the only handle that can resume the conversation.
4928        assert!(adapter.reload().await.is_err());
4929        assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
4930    }
4931
4932    #[tokio::test]
4933    async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
4934        let cwd = std::env::current_dir().expect("cwd");
4935        let mut adapter = AcpAdapter::with_session_id(
4936            0,
4937            cwd,
4938            "__codeswarm_missing_nonloadable_acp__",
4939            Vec::new(),
4940            "stale-session",
4941        );
4942        adapter.capabilities.supports_session_load = false;
4943        assert!(adapter.reload().await.is_err());
4944        assert_eq!(adapter.session_id, None);
4945    }
4946
4947    #[tokio::test]
4948    async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
4949        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"}}'"#;
4950        let cwd = std::env::current_dir().expect("cwd");
4951        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4952        adapter.start().await.expect("initialize");
4953        assert!(matches!(
4954            adapter.next_event().await,
4955            Some(Ok(AgentEvent::Ready { .. }))
4956        ));
4957        adapter.send_prompt("hello".into()).await.expect("prompt");
4958        assert!(matches!(
4959            adapter.next_event().await,
4960            Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
4961        ));
4962        assert!(matches!(
4963            adapter.next_event().await,
4964            Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
4965        ));
4966    }
4967
4968    #[tokio::test]
4969    async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
4970        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
4971        let cwd = std::env::current_dir().expect("cwd");
4972        let mut adapter = AcpAdapter::with_session_id(
4973            0,
4974            cwd,
4975            "sh",
4976            vec!["-c".into(), script.into()],
4977            "existing-session",
4978        );
4979        adapter.start().await.expect("load existing session");
4980        assert!(matches!(
4981            adapter.next_event().await,
4982            Some(Ok(AgentEvent::Ready { .. }))
4983        ));
4984    }
4985
4986    #[tokio::test]
4987    async fn acp_start_failure_reaps_transport_process() {
4988        // The child emits an invalid initialize response and exits. The
4989        // adapter must not retain a live child after protocol startup fails;
4990        // this is the path used when a configured ACP command is unavailable
4991        // or speaks a different protocol.
4992        let mut adapter = AcpAdapter::new(
4993            0,
4994            std::env::current_dir().expect("cwd"),
4995            "sh",
4996            vec!["-c".into(), "printf 'not-json\\n'".into()],
4997        );
4998        assert!(adapter.start().await.is_err());
4999        assert!(adapter.child.is_none());
5000        assert!(adapter.reader.is_none());
5001    }
5002
5003    #[tokio::test]
5004    async fn acp_adapter_answers_permission_json_rpc_requests() {
5005        let path = std::env::temp_dir().join(format!(
5006            "codeswarm-permission-answer-{}",
5007            std::process::id()
5008        ));
5009        let script = format!(
5010            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"}}}}'"#,
5011            path.display()
5012        );
5013        let mut adapter = AcpAdapter::new(
5014            0,
5015            std::env::current_dir().expect("cwd"),
5016            "sh",
5017            vec!["-c".into(), script],
5018        );
5019        adapter.start().await.expect("start ACP");
5020        assert!(matches!(
5021            adapter.next_event().await,
5022            Some(Ok(AgentEvent::Ready { .. }))
5023        ));
5024        adapter.send_prompt("do it".into()).await.expect("prompt");
5025        assert!(matches!(
5026            adapter.next_event().await,
5027            Some(Ok(AgentEvent::Permission { request, .. }))
5028                if request.id == "9"
5029                    && request.options == ["Allow once"]
5030                    && request.option_ids == ["allow-once"]
5031        ));
5032        adapter
5033            .answer_permission(
5034                "9".into(),
5035                PermissionAnswer::Selected {
5036                    option_id: "allow-once".into(),
5037                },
5038            )
5039            .await
5040            .expect("permission answer");
5041        assert!(matches!(
5042            adapter.next_event().await,
5043            Some(Ok(AgentEvent::TurnComplete { .. }))
5044        ));
5045        let answer: Value = serde_json::from_str(
5046            &std::fs::read_to_string(&path).expect("captured permission answer"),
5047        )
5048        .expect("valid JSON-RPC answer");
5049        assert_eq!(answer["id"], 9);
5050        assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5051        assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5052        std::fs::remove_file(path).expect("cleanup");
5053    }
5054
5055    #[test]
5056    fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5057        let event = parse_acp_notification(
5058            0,
5059            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5060        )
5061        .expect("valid JSON-RPC request");
5062        assert!(event.is_none());
5063    }
5064
5065    #[tokio::test]
5066    async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5067        let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5068        std::fs::write(
5069            &script_path,
5070            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5071        )
5072        .expect("write native test script");
5073        let mut adapter = AgyAdapter::new(
5074            0,
5075            std::env::current_dir().expect("cwd"),
5076            format!("sh {}", script_path.display()),
5077        );
5078        adapter.start().await.expect("start native adapter");
5079        assert!(adapter.next_event().await.is_some());
5080        assert!(adapter.next_event().await.is_some());
5081        adapter
5082            .send_prompt("continue".into())
5083            .await
5084            .expect("prompt");
5085        assert!(matches!(
5086            adapter.next_event().await,
5087            Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5088        ));
5089        assert!(matches!(
5090            adapter.next_event().await,
5091            Some(Ok(AgentEvent::TurnComplete { .. }))
5092        ));
5093        adapter.stop().await.expect("stop native adapter");
5094        std::fs::remove_file(script_path).expect("cleanup native script");
5095    }
5096
5097    #[test]
5098    fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5099        let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5100        let _ = std::fs::remove_dir_all(&root);
5101        std::fs::create_dir_all(&root).expect("workspace");
5102        std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5103        let outside =
5104            std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5105        std::fs::write(&outside, "secret").expect("outside file");
5106        let link = root.join("outside-link");
5107        #[cfg(unix)]
5108        std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5109        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5110
5111        assert_eq!(
5112            adapter
5113                .read_workspace_text("inside.txt", Some(2), Some(1))
5114                .expect("read inside"),
5115            "two"
5116        );
5117        std::fs::write(
5118            root.join("large.txt"),
5119            vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5120        )
5121        .expect("large file");
5122        let bounded = adapter
5123            .read_workspace_text("large.txt", None, None)
5124            .expect("bounded read");
5125        assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5126        #[cfg(unix)]
5127        {
5128            std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5129                .expect("internal symlink");
5130            assert_eq!(
5131                adapter
5132                    .read_workspace_text("inside-link", None, None)
5133                    .expect("read internal symlink"),
5134                "one\ntwo\nthree\n"
5135            );
5136        }
5137        assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5138        assert!(
5139            adapter
5140                .workspace_path(&outside.display().to_string())
5141                .is_err()
5142        );
5143        #[cfg(unix)]
5144        assert!(adapter.workspace_path("outside-link").is_err());
5145        #[cfg(unix)]
5146        std::fs::remove_file(link).expect("cleanup symlink");
5147        #[cfg(unix)]
5148        std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5149        std::fs::remove_file(outside).expect("cleanup outside");
5150        std::fs::remove_dir_all(root).expect("cleanup workspace");
5151    }
5152
5153    #[tokio::test]
5154    async fn running_terminal_output_omits_exit_status_until_completion() {
5155        let root = unique_test_path("codeswarm-terminal-output", "dir");
5156        std::fs::create_dir_all(&root).expect("workspace");
5157        let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5158        let result = adapter
5159            .terminal_create(&serde_json::json!({
5160                "command": "sh",
5161                "args": ["-c", "sleep 0.2; printf done"],
5162                "cwd": ".",
5163            }))
5164            .await
5165            .expect("terminal create");
5166        let id = result["terminalId"].as_str().expect("terminal id");
5167        let output = adapter.terminal_output(id).await.expect("terminal output");
5168        assert!(output.get("exitStatus").is_none());
5169        if let Some(terminal) = adapter.terminals.remove(id) {
5170            terminal.stop().await;
5171        }
5172        std::fs::remove_dir_all(root).expect("cleanup workspace");
5173    }
5174
5175    #[tokio::test]
5176    async fn acp_adapter_answers_workspace_read_requests() {
5177        let root =
5178            std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
5179        let _ = std::fs::remove_dir_all(&root);
5180        std::fs::create_dir_all(&root).expect("workspace");
5181        let source = root.join("inside.txt");
5182        let answer = root.join("answer.json");
5183        std::fs::write(&source, "workspace content").expect("source");
5184        let script = format!(
5185            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"}}}}'"#,
5186            source.display(),
5187            answer.display(),
5188        );
5189        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5190        adapter.start().await.expect("start ACP");
5191        assert!(matches!(
5192            adapter.next_event().await,
5193            Some(Ok(AgentEvent::Ready { .. }))
5194        ));
5195        adapter.send_prompt("read it".into()).await.expect("prompt");
5196        assert!(matches!(
5197            adapter.next_event().await,
5198            Some(Ok(AgentEvent::TurnComplete { .. }))
5199        ));
5200        let response: Value =
5201            serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
5202                .expect("response JSON");
5203        assert_eq!(response["id"], 9);
5204        assert_eq!(response["result"]["content"], "workspace content");
5205        adapter.stop().await.expect("stop ACP");
5206        std::fs::remove_dir_all(root).expect("cleanup workspace");
5207    }
5208
5209    #[tokio::test]
5210    async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
5211        let root =
5212            std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
5213        let _ = std::fs::remove_dir_all(&root);
5214        std::fs::create_dir_all(&root).expect("workspace");
5215        let create_request = serde_json::json!({
5216            "jsonrpc": "2.0",
5217            "id": 9,
5218            "method": "terminal/create",
5219            "params": {
5220                "sessionId": "s1",
5221                "command": "sh",
5222                "args": ["-c", "sleep 0.1; printf terminal-ok"],
5223                "cwd": ".",
5224            },
5225        });
5226        let wait_request = serde_json::json!({
5227            "jsonrpc": "2.0",
5228            "id": 10,
5229            "method": "terminal/wait_for_exit",
5230            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5231        });
5232        let output_request = serde_json::json!({
5233            "jsonrpc": "2.0",
5234            "id": 11,
5235            "method": "terminal/output",
5236            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5237        });
5238        let create_answer = root.join("create-answer.json");
5239        let wait_answer = root.join("wait-answer.json");
5240        let output_answer = root.join("output-answer.json");
5241        let script = format!(
5242            "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\"}}}}'",
5243            create_request,
5244            create_answer.display(),
5245            wait_request,
5246            wait_answer.display(),
5247            output_request,
5248            output_answer.display(),
5249        );
5250        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5251        adapter.start().await.expect("start ACP");
5252        assert!(matches!(
5253            adapter.next_event().await,
5254            Some(Ok(AgentEvent::Ready { .. }))
5255        ));
5256        adapter
5257            .send_prompt("run terminal".into())
5258            .await
5259            .expect("prompt");
5260        let mut saw_complete = false;
5261        for _ in 0..6 {
5262            match adapter.next_event().await {
5263                Some(Ok(AgentEvent::TurnComplete { .. })) => {
5264                    saw_complete = true;
5265                    break;
5266                }
5267                Some(_) => {}
5268                None => break,
5269            }
5270        }
5271        assert!(saw_complete, "terminal requests should not stall ACP");
5272        let create: Value = serde_json::from_str(
5273            &std::fs::read_to_string(&create_answer).expect("captured create response"),
5274        )
5275        .expect("create JSON");
5276        assert_eq!(create["result"]["terminalId"], "terminal-1");
5277        let output: Value = serde_json::from_str(
5278            &std::fs::read_to_string(&output_answer).expect("captured output response"),
5279        )
5280        .expect("output JSON");
5281        assert!(
5282            output["result"]["output"]
5283                .as_str()
5284                .unwrap_or_default()
5285                .contains("terminal-ok"),
5286            "output response: {output}"
5287        );
5288        adapter.stop().await.expect("stop ACP");
5289        std::fs::remove_dir_all(root).expect("cleanup workspace");
5290    }
5291
5292    #[tokio::test]
5293    async fn host_reduces_and_persists_adapter_events() {
5294        let path =
5295            std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
5296        let adapter = ScriptedAdapter::new(
5297            0,
5298            AgentCapabilities::default(),
5299            [AgentEvent::Text {
5300                slot: 0,
5301                text: "hello".into(),
5302            }],
5303        );
5304        let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
5305        host.start().await.expect("start");
5306        host.next_effects()
5307            .await
5308            .expect("event")
5309            .expect("valid event");
5310        assert_eq!(host.state.public_text[0].1, "hello");
5311        assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
5312        std::fs::remove_file(path).expect("cleanup");
5313    }
5314
5315    #[tokio::test]
5316    async fn relay_applies_default_policy_before_the_first_prompt() {
5317        let first_log = Arc::new(Mutex::new(Vec::new()));
5318        let second_log = Arc::new(Mutex::new(Vec::new()));
5319        let first = AdapterHost::new(
5320            Box::new(ModeOrderAdapter {
5321                slot: 0,
5322                log: Arc::clone(&first_log),
5323                phase: 0,
5324            }),
5325            None,
5326        );
5327        let second = AdapterHost::new(
5328            Box::new(ModeOrderAdapter {
5329                slot: 1,
5330                log: Arc::clone(&second_log),
5331                phase: 0,
5332            }),
5333            None,
5334        );
5335        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5336        relay.start().await.expect("start and synchronize policy");
5337        relay.run_turn("task", 0).await.expect("first turn");
5338        {
5339            let log = first_log.lock().expect("log");
5340            assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
5341        }
5342        assert_eq!(
5343            second_log.lock().expect("log").as_slice(),
5344            ["start", "mode:yolo"]
5345        );
5346        let added_log = Arc::new(Mutex::new(Vec::new()));
5347        relay
5348            .add_agent(
5349                AdapterHost::new(
5350                    Box::new(ModeOrderAdapter {
5351                        slot: 2,
5352                        log: Arc::clone(&added_log),
5353                        phase: 0,
5354                    }),
5355                    None,
5356                ),
5357                "Added",
5358                "added.example",
5359                "added-agent",
5360            )
5361            .await
5362            .expect("add with synchronized policy");
5363        assert_eq!(
5364            added_log.lock().expect("log").as_slice(),
5365            ["start", "mode:yolo"]
5366        );
5367        relay.drop_agent(2).await.expect("drop added agent");
5368        added_log.lock().expect("log").clear();
5369        relay
5370            .reload(2)
5371            .await
5372            .expect("reload with synchronized policy");
5373        assert_eq!(
5374            added_log.lock().expect("log").as_slice(),
5375            ["reload", "mode:yolo"]
5376        );
5377    }
5378
5379    #[tokio::test]
5380    async fn acp_roster_is_ready_before_any_prompt_is_sent() {
5381        let hosts = (0..2)
5382            .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
5383            .collect::<Vec<_>>();
5384        let startup_events = Arc::new(Mutex::new(Vec::new()));
5385        let captured = Arc::clone(&startup_events);
5386        let mut relay = RelayHost::new(hosts, 4).expect("relay");
5387        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5388
5389        relay.start().await.expect("complete startup handshake");
5390
5391        assert!(relay.dispatches().is_empty());
5392        let ready_slots = startup_events
5393            .lock()
5394            .expect("events")
5395            .iter()
5396            .filter_map(|event| match event {
5397                AgentEvent::Ready { slot, .. } => Some(*slot),
5398                _ => None,
5399            })
5400            .collect::<Vec<_>>();
5401        assert_eq!(ready_slots, vec![0, 1]);
5402    }
5403
5404    #[tokio::test]
5405    async fn independent_roster_adapters_start_concurrently() {
5406        let barrier = Arc::new(tokio::sync::Barrier::new(2));
5407        let hosts = (0..2)
5408            .map(|slot| {
5409                AdapterHost::new(
5410                    Box::new(ConcurrentStartAdapter {
5411                        slot,
5412                        barrier: Arc::clone(&barrier),
5413                    }),
5414                    None,
5415                )
5416            })
5417            .collect::<Vec<_>>();
5418        let mut relay = RelayHost::new(hosts, 4).expect("relay");
5419        tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
5420            .await
5421            .expect("startup should not serialize barrier participants")
5422            .expect("startup succeeds");
5423    }
5424
5425    #[tokio::test]
5426    async fn relay_host_dispatches_turns_sequentially() {
5427        let capabilities = AgentCapabilities {
5428            supports_cancel: true,
5429            ..AgentCapabilities::default()
5430        };
5431        let first = ScriptedAdapter::new(
5432            0,
5433            capabilities.clone(),
5434            [
5435                AgentEvent::Text {
5436                    slot: 0,
5437                    text: "first".into(),
5438                },
5439                AgentEvent::TurnComplete { slot: 0 },
5440            ],
5441        );
5442        let second = ScriptedAdapter::new(
5443            1,
5444            capabilities,
5445            [
5446                AgentEvent::Text {
5447                    slot: 1,
5448                    text: "review".into(),
5449                },
5450                AgentEvent::TurnComplete { slot: 1 },
5451            ],
5452        );
5453        let hosts = vec![
5454            AdapterHost::new(Box::new(first), None),
5455            AdapterHost::new(Box::new(second), None),
5456        ];
5457        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5458        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5459        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5460        let captured = std::sync::Arc::clone(&events);
5461        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5462        relay.start().await.expect("start");
5463        events.lock().expect("events").clear();
5464        assert!(matches!(
5465            relay.run_turn("task", 0).await.expect("first turn"),
5466            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5467        ));
5468        assert!(matches!(
5469            relay.run_turn("first", 0).await.expect("second turn"),
5470            crate::relay::RelayDecision::Dispatch {
5471                slot: 1,
5472                can_stop: true,
5473                ..
5474            }
5475        ));
5476        assert_eq!(
5477            relay
5478                .dispatches()
5479                .iter()
5480                .map(|(slot, _)| *slot)
5481                .collect::<Vec<_>>(),
5482            [0, 1]
5483        );
5484        assert!(relay.dispatches()[0].1.contains("You are Claude"));
5485        assert!(
5486            relay.dispatches()[0]
5487                .1
5488                .contains("CodeSwarm roster (ordered)")
5489        );
5490        assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
5491        assert!(relay.dispatches()[0].1.contains("2. Codex"));
5492        assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
5493        assert!(relay.dispatches()[0].1.contains("Do not use"));
5494        let lifecycle = events.lock().expect("events");
5495        let positions = lifecycle
5496            .iter()
5497            .filter_map(|event| match event {
5498                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
5499                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
5500                _ => None,
5501            })
5502            .collect::<Vec<_>>();
5503        assert_eq!(
5504            positions,
5505            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
5506        );
5507    }
5508
5509    #[tokio::test]
5510    async fn relay_host_routes_around_a_usage_limited_agent() {
5511        let capabilities = AgentCapabilities::default();
5512        let first = ScriptedAdapter::new(
5513            0,
5514            capabilities.clone(),
5515            [
5516                AgentEvent::Text {
5517                    slot: 0,
5518                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
5519                           credits or try again later."
5520                        .into(),
5521                },
5522                AgentEvent::TurnComplete { slot: 0 },
5523            ],
5524        );
5525        let second = ScriptedAdapter::new(
5526            1,
5527            capabilities,
5528            [
5529                AgentEvent::Text {
5530                    slot: 1,
5531                    text: "review done".into(),
5532                },
5533                AgentEvent::TurnComplete { slot: 1 },
5534            ],
5535        );
5536        let hosts = vec![
5537            AdapterHost::new(Box::new(first), None),
5538            AdapterHost::new(Box::new(second), None),
5539        ];
5540        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5541        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5542        let captured = std::sync::Arc::clone(&events);
5543        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5544        relay.start().await.expect("start");
5545        events.lock().expect("events").clear();
5546        assert!(matches!(
5547            relay.run_turn("task", 0).await.expect("limited turn"),
5548            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5549        ));
5550        assert!(
5551            events
5552                .lock()
5553                .expect("events")
5554                .iter()
5555                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5556        );
5557        // The next automatic turn skips the limited agent entirely.
5558        assert!(matches!(
5559            relay.run_turn("", 0).await.expect("next turn"),
5560            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5561        ));
5562        assert!(relay.relay().is_limited(0));
5563        // A reload restores the agent to the ring. (ScriptedAdapter cannot
5564        // feed further turns, so the restored routing itself is covered by
5565        // the relay unit tests.)
5566        relay.reload(0).await.expect("reload");
5567        assert!(!relay.relay().is_limited(0));
5568    }
5569
5570    #[tokio::test]
5571    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
5572        let limited = ScriptedAdapter::new(
5573            0,
5574            AgentCapabilities::default(),
5575            [AgentEvent::Failed {
5576                slot: 0,
5577                started: true,
5578                detail: "request failed: insufficient_quota".into(),
5579            }],
5580        );
5581        let healthy = ScriptedAdapter::new(
5582            1,
5583            AgentCapabilities::default(),
5584            [AgentEvent::TurnComplete { slot: 1 }],
5585        );
5586        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5587        let captured = std::sync::Arc::clone(&events);
5588        let mut relay = RelayHost::new(
5589            vec![
5590                AdapterHost::new(Box::new(limited), None),
5591                AdapterHost::new(Box::new(healthy), None),
5592            ],
5593            4,
5594        )
5595        .expect("relay");
5596        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5597        relay.start().await.expect("start");
5598        events.lock().expect("events").clear();
5599
5600        assert!(matches!(
5601            relay.run_turn("task", 0).await.expect("limited failure"),
5602            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5603        ));
5604        assert!(relay.relay().is_limited(0));
5605        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
5606        {
5607            let events = events.lock().expect("events");
5608            assert!(
5609                events
5610                    .iter()
5611                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5612            );
5613            assert!(
5614                !events
5615                    .iter()
5616                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
5617            );
5618        }
5619
5620        assert!(matches!(
5621            relay.run_turn("", 0).await.expect("healthy peer"),
5622            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5623        ));
5624    }
5625
5626    #[tokio::test]
5627    async fn relay_failure_is_tombstoned_and_reported_to_the_ui_sink() {
5628        let failed = ScriptedAdapter::new(
5629            0,
5630            AgentCapabilities::default(),
5631            [AgentEvent::Failed {
5632                slot: 0,
5633                started: true,
5634                detail: "connection lost".into(),
5635            }],
5636        );
5637        let healthy = ScriptedAdapter::new(
5638            1,
5639            AgentCapabilities::default(),
5640            [AgentEvent::TurnComplete { slot: 1 }],
5641        );
5642        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5643        let captured = std::sync::Arc::clone(&events);
5644        let mut relay = RelayHost::new(
5645            vec![
5646                AdapterHost::new(Box::new(failed), None),
5647                AdapterHost::new(Box::new(healthy), None),
5648            ],
5649            4,
5650        )
5651        .expect("relay");
5652        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5653        relay.start().await.expect("start");
5654
5655        let error = relay.run_turn("task", 0).await.expect_err("failure");
5656        assert!(error.to_string().contains("connection lost"));
5657        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
5658        assert!(events.lock().expect("lock").iter().any(|event| {
5659            matches!(
5660                event,
5661                AgentEvent::Failed {
5662                    slot: 0,
5663                    started: true,
5664                    ..
5665                }
5666            )
5667        }));
5668    }
5669
5670    #[tokio::test]
5671    async fn codex_stop_does_not_skip_later_roster_reviewers() {
5672        let hosts = (0..3)
5673            .map(|slot| {
5674                AdapterHost::new(
5675                    Box::new(ScriptedAdapter::new(
5676                        slot,
5677                        AgentCapabilities::default(),
5678                        [
5679                            AgentEvent::Text {
5680                                slot,
5681                                text: STOP_TOKEN.into(),
5682                            },
5683                            AgentEvent::TurnComplete { slot },
5684                        ],
5685                    )),
5686                    None,
5687                )
5688            })
5689            .collect();
5690        let mut relay = RelayHost::new(hosts, 10).expect("relay");
5691        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
5692        relay.start().await.expect("start");
5693        for expected in 0..3 {
5694            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
5695                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
5696        }
5697        assert_eq!(
5698            relay.run_turn("", 0).await.expect("complete"),
5699            RelayDecision::Complete
5700        );
5701    }
5702
5703    #[tokio::test]
5704    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
5705        let first = ScriptedAdapter::new(
5706            0,
5707            AgentCapabilities::default(),
5708            [
5709                AgentEvent::Text {
5710                    slot: 0,
5711                    text: "done".into(),
5712                },
5713                AgentEvent::TurnComplete { slot: 0 },
5714            ],
5715        );
5716        let reviewer = ScriptedAdapter::new(
5717            1,
5718            AgentCapabilities::default(),
5719            [
5720                AgentEvent::Text {
5721                    slot: 1,
5722                    text: STOP_TOKEN.into(),
5723                },
5724                AgentEvent::TurnComplete { slot: 1 },
5725            ],
5726        );
5727        let mut relay = RelayHost::new(
5728            vec![
5729                AdapterHost::new(Box::new(first), None),
5730                AdapterHost::new(Box::new(reviewer), None),
5731            ],
5732            10,
5733        )
5734        .expect("relay");
5735        relay.start().await.expect("start");
5736        let first_decision = relay.run_turn("task", 0).await.expect("first");
5737        assert!(matches!(
5738            first_decision,
5739            RelayDecision::Dispatch { slot: 0, .. }
5740        ));
5741        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
5742        assert!(matches!(
5743            reviewer_decision,
5744            RelayDecision::Dispatch {
5745                slot: 1,
5746                can_stop: true,
5747                ..
5748            }
5749        ));
5750        assert_eq!(
5751            relay.run_turn("", 0).await.expect("complete"),
5752            RelayDecision::Complete
5753        );
5754    }
5755
5756    #[tokio::test]
5757    async fn stop_token_is_filtered_from_streamed_ui_events() {
5758        let first = ScriptedAdapter::new(
5759            0,
5760            AgentCapabilities::default(),
5761            [
5762                AgentEvent::Text {
5763                    slot: 0,
5764                    text: format!("visible {STOP_TOKEN} trailing"),
5765                },
5766                AgentEvent::TurnComplete { slot: 0 },
5767            ],
5768        );
5769        let reviewer = ScriptedAdapter::new(
5770            1,
5771            AgentCapabilities::default(),
5772            [AgentEvent::TurnComplete { slot: 1 }],
5773        );
5774        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5775        let captured = std::sync::Arc::clone(&events);
5776        let mut relay = RelayHost::new(
5777            vec![
5778                AdapterHost::new(Box::new(first), None),
5779                AdapterHost::new(Box::new(reviewer), None),
5780            ],
5781            2,
5782        )
5783        .expect("relay");
5784        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5785        relay.start().await.expect("start");
5786        relay.run_turn("task", 0).await.expect("turn");
5787        let captured = events.lock().expect("lock");
5788        assert!(captured.iter().all(|event| match event {
5789            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5790            _ => true,
5791        }));
5792        let visible = captured
5793            .iter()
5794            .filter_map(|event| match event {
5795                AgentEvent::Text { text, .. } => Some(text.as_str()),
5796                _ => None,
5797            })
5798            .collect::<String>();
5799        assert_eq!(visible, "visible  trailing");
5800    }
5801
5802    #[tokio::test]
5803    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
5804        let first = ScriptedAdapter::new(
5805            0,
5806            AgentCapabilities::default(),
5807            [
5808                AgentEvent::Text {
5809                    slot: 0,
5810                    text: "done".into(),
5811                },
5812                AgentEvent::TurnComplete { slot: 0 },
5813            ],
5814        );
5815        let reviewer = ScriptedAdapter::new(
5816            1,
5817            AgentCapabilities::default(),
5818            [
5819                AgentEvent::Text {
5820                    slot: 1,
5821                    text: STOP_TOKEN.into(),
5822                },
5823                AgentEvent::TurnComplete { slot: 1 },
5824            ],
5825        );
5826        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5827        let captured = std::sync::Arc::clone(&events);
5828        let mut relay = RelayHost::new(
5829            vec![
5830                AdapterHost::new(Box::new(first), None),
5831                AdapterHost::new(Box::new(reviewer), None),
5832            ],
5833            4,
5834        )
5835        .expect("relay");
5836        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5837        relay.start().await.expect("start");
5838        relay.run_turn("task", 0).await.expect("first turn");
5839        relay.run_turn("", 0).await.expect("review turn");
5840        let captured = events.lock().expect("lock");
5841        assert!(captured.iter().any(|event| {
5842            matches!(
5843                event,
5844                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5845            )
5846        }));
5847        assert!(captured.iter().all(|event| match event {
5848            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5849            _ => true,
5850        }));
5851        let acknowledgment = captured
5852            .iter()
5853            .position(|event| {
5854                matches!(
5855                    event,
5856                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5857                )
5858            })
5859            .expect("visible acknowledgment");
5860        let completion = captured
5861            .iter()
5862            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
5863            .expect("reviewer completion");
5864        assert!(acknowledgment < completion);
5865    }
5866
5867    #[tokio::test]
5868    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
5869        let first = ScriptedAdapter::new(
5870            0,
5871            AgentCapabilities::default(),
5872            [
5873                AgentEvent::Text {
5874                    slot: 0,
5875                    text: "done".into(),
5876                },
5877                AgentEvent::TurnComplete { slot: 0 },
5878            ],
5879        );
5880        let reviewer = ScriptedAdapter::new(
5881            1,
5882            AgentCapabilities::default(),
5883            [
5884                AgentEvent::Text {
5885                    slot: 1,
5886                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
5887                },
5888                AgentEvent::TurnComplete { slot: 1 },
5889            ],
5890        );
5891        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5892        let captured = std::sync::Arc::clone(&events);
5893        let mut relay = RelayHost::new(
5894            vec![
5895                AdapterHost::new(Box::new(first), None),
5896                AdapterHost::new(Box::new(reviewer), None),
5897            ],
5898            4,
5899        )
5900        .expect("relay");
5901        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5902        relay.start().await.expect("start");
5903        relay.run_turn("task", 0).await.expect("first turn");
5904        relay.run_turn("", 0).await.expect("review turn");
5905
5906        let visible = events
5907            .lock()
5908            .expect("lock")
5909            .iter()
5910            .filter_map(|event| match event {
5911                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
5912                _ => None,
5913            })
5914            .collect::<String>();
5915        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
5916        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
5917    }
5918
5919    #[tokio::test]
5920    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
5921        let first = AdapterHost::new(
5922            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
5923            None,
5924        );
5925        let second = AdapterHost::new(
5926            Box::new(ScriptedAdapter::new(
5927                1,
5928                AgentCapabilities::default(),
5929                [AgentEvent::TurnComplete { slot: 1 }],
5930            )),
5931            None,
5932        );
5933        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5934        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
5935        relay.set_event_sink(move |event| {
5936            if matches!(event, AgentEvent::Permission { .. }) {
5937                let _ = seen_sender.send(());
5938            }
5939        });
5940        relay.start().await.expect("start");
5941        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
5942        let answer = async move {
5943            seen_receiver.recv().await.expect("permission request");
5944            sender
5945                .send(super::RelayPermissionAnswer {
5946                    slot: 0,
5947                    request_id: "permission-1".into(),
5948                    answer: PermissionAnswer::Selected {
5949                        option_id: "allow".into(),
5950                    },
5951                })
5952                .expect("queue permission answer");
5953        };
5954        tokio::time::timeout(std::time::Duration::from_millis(100), async {
5955            let ((), result) = tokio::join!(
5956                answer,
5957                relay.run_turn_with_permissions("task", 0, &mut receiver)
5958            );
5959            result
5960        })
5961        .await
5962        .expect("permission-gated turn should not deadlock")
5963        .expect("turn completes");
5964    }
5965
5966    #[tokio::test]
5967    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
5968        let first = AdapterHost::new(
5969            Box::new(PendingAdapter {
5970                slot: 0,
5971                hang_on_cancel: false,
5972            }),
5973            None,
5974        );
5975        let second = AdapterHost::new(
5976            Box::new(ScriptedAdapter::new(
5977                1,
5978                AgentCapabilities::default(),
5979                [AgentEvent::TurnComplete { slot: 1 }],
5980            )),
5981            None,
5982        );
5983        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
5984        relay.start().await.expect("start");
5985        let cancellation = relay.cancellation();
5986        let error = {
5987            let turn = relay.run_turn("task", 0);
5988            tokio::pin!(turn);
5989            cancellation.request();
5990            turn.await.expect_err("cancellation should stop turn")
5991        };
5992        assert!(error.to_string().contains("relay turn cancelled"));
5993
5994        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
5995        relay
5996            .run_turn("", 1)
5997            .await
5998            .expect("replacement job reaches the selected peer");
5999        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
6000        assert!(replacement.contains("replacement job"));
6001        assert!(replacement.contains("User "));
6002        assert!(replacement.contains(":\ntask"));
6003        let owner_updates = relay.relay_mut().unseen_context(0);
6004        assert!(owner_updates.contains("User "));
6005        assert!(owner_updates.contains(":\ntask"));
6006        assert!(owner_updates.contains(":\nreplacement job"));
6007    }
6008
6009    #[tokio::test]
6010    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
6011        let first = AdapterHost::new(
6012            Box::new(PendingAdapter {
6013                slot: 0,
6014                hang_on_cancel: true,
6015            }),
6016            None,
6017        );
6018        let second = AdapterHost::new(
6019            Box::new(PendingAdapter {
6020                slot: 1,
6021                hang_on_cancel: false,
6022            }),
6023            None,
6024        );
6025        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6026        relay.start().await.expect("start");
6027        let cancellation = relay.cancellation();
6028        let turn = relay.run_turn("task", 0);
6029        tokio::pin!(turn);
6030        cancellation.request();
6031        let error = turn.await.expect_err("cancellation should stop turn");
6032        assert!(error.to_string().contains("timed out"));
6033    }
6034
6035    #[tokio::test]
6036    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
6037        let event = [AgentEvent::TurnComplete { slot: 0 }];
6038        let first = AdapterHost::new(
6039            Box::new(ScriptedAdapter::new(
6040                0,
6041                AgentCapabilities::default(),
6042                event.clone(),
6043            )),
6044            None,
6045        );
6046        let second = AdapterHost::new(
6047            Box::new(ScriptedAdapter::new(
6048                1,
6049                AgentCapabilities::default(),
6050                [AgentEvent::TurnComplete { slot: 1 }],
6051            )),
6052            None,
6053        );
6054        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6055        relay.start().await.expect("start");
6056
6057        relay.pause();
6058        assert_eq!(
6059            relay.run_turn("paused", 0).await.expect("paused turn"),
6060            crate::relay::RelayDecision::Paused
6061        );
6062        assert!(relay.dispatches().is_empty());
6063
6064        relay.resume();
6065        relay.relay_mut().drop_agent(1).expect("drop reviewer");
6066        assert!(matches!(
6067            relay
6068                .run_turn("solo follow-up", 0)
6069                .await
6070                .expect("solo turn"),
6071            crate::relay::RelayDecision::Dispatch {
6072                slot: 0,
6073                can_stop: false,
6074                ..
6075            }
6076        ));
6077        assert_eq!(relay.dispatches().len(), 1);
6078    }
6079
6080    #[tokio::test]
6081    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
6082        let first = AdapterHost::new(
6083            Box::new(ScriptedAdapter::new(
6084                0,
6085                AgentCapabilities::default(),
6086                [AgentEvent::TurnComplete { slot: 0 }],
6087            )),
6088            None,
6089        );
6090        let second = AdapterHost::new(
6091            Box::new(ScriptedAdapter::new(
6092                1,
6093                AgentCapabilities::default(),
6094                [AgentEvent::TurnComplete { slot: 1 }],
6095            )),
6096            None,
6097        );
6098        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6099        relay.set_roster_names(vec!["First".into(), "Second".into()]);
6100        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
6101        relay.set_roster_launch_specs(vec![
6102            ("custom".into(), "owner".into()),
6103            ("custom".into(), "peer".into()),
6104        ]);
6105        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6106        let captured = std::sync::Arc::clone(&events);
6107        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6108        relay.start().await.expect("start");
6109        let slot = relay
6110            .add_agent(
6111                AdapterHost::new(
6112                    Box::new(ScriptedAdapter::new(
6113                        2,
6114                        AgentCapabilities::default(),
6115                        [AgentEvent::TurnComplete { slot: 2 }],
6116                    )),
6117                    None,
6118                ),
6119                "Reviewer",
6120                "reviewer.example",
6121                "reviewer --acp",
6122            )
6123            .await
6124            .expect("append agent");
6125        assert_eq!(slot, 2);
6126        assert_eq!(
6127            relay.relay().active_slots().collect::<Vec<_>>(),
6128            vec![0, 1, 2]
6129        );
6130        assert_eq!(
6131            relay
6132                .session_metadata()
6133                .get("agents")
6134                .and_then(|value| value.as_array())
6135                .map(Vec::len),
6136            Some(3)
6137        );
6138        relay.drop_agent(1).await.expect("drop middle peer");
6139        let metadata = relay.session_metadata();
6140        assert_eq!(
6141            metadata.get("agents"),
6142            Some(&serde_json::json!([
6143                {"name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
6144                {"name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
6145            ]))
6146        );
6147        assert!(
6148            events
6149                .lock()
6150                .expect("lock")
6151                .iter()
6152                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
6153        );
6154    }
6155
6156    #[tokio::test]
6157    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
6158        let path = unique_test_path("codeswarm-session-metadata", "json");
6159        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
6160        let writer = metadata_store.buffered().expect("metadata writer");
6161        let first = AdapterHost::new(
6162            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
6163            None,
6164        );
6165        let second = AdapterHost::new(
6166            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
6167            None,
6168        );
6169        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6170        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6171        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
6172        relay.set_roster_launch_specs(vec![
6173            ("custom".into(), "claude".into()),
6174            ("custom".into(), "codex".into()),
6175        ]);
6176        relay.set_session_metadata_writer(writer);
6177        relay.start().await.expect("start");
6178        relay.drop_agent(0).await.expect("drop first agent");
6179        relay.stop().await.expect("stop");
6180
6181        let loaded = metadata_store
6182            .read()
6183            .expect("read metadata")
6184            .expect("metadata snapshot");
6185        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
6186        assert_eq!(
6187            loaded.get("agents"),
6188            Some(&serde_json::json!([{
6189                "name": "Codex", "identity": "openai.com", "protocol": "custom",
6190                "command": "codex", "supports_load_session": false
6191            }]))
6192        );
6193        assert!(loaded.get("owner").is_none());
6194        let _ = std::fs::remove_file(path);
6195    }
6196
6197    #[tokio::test]
6198    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
6199        let first = AdapterHost::new(
6200            Box::new(ScriptedAdapter::new(
6201                0,
6202                AgentCapabilities::default(),
6203                [
6204                    AgentEvent::Text {
6205                        slot: 0,
6206                        text: "owner stream".into(),
6207                    },
6208                    AgentEvent::TurnComplete { slot: 0 },
6209                ],
6210            )),
6211            None,
6212        );
6213        let second = AdapterHost::new(
6214            Box::new(ScriptedAdapter::new(
6215                1,
6216                AgentCapabilities::default(),
6217                [
6218                    AgentEvent::Text {
6219                        slot: 1,
6220                        text: "peer stream".into(),
6221                    },
6222                    AgentEvent::TurnComplete { slot: 1 },
6223                ],
6224            )),
6225            None,
6226        );
6227        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6228        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
6229        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
6230        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6231        let captured = std::sync::Arc::clone(&events);
6232        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6233        relay.start().await.expect("start");
6234
6235        relay.swap_agents(0, 1).expect("swap peers");
6236        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
6237        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
6238        relay.run_turn("task", 0).await.expect("swapped turn");
6239        let events = events.lock().expect("events");
6240        assert!(events.iter().any(|event| {
6241            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
6242        }));
6243        assert!(relay.dispatches()[0].1.contains("You are Peer"));
6244    }
6245
6246    #[tokio::test]
6247    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
6248        let path = unique_test_path("codeswarm-session-metadata", "json");
6249        let first = AdapterHost::new(
6250            Box::new(ScriptedAdapter::new(
6251                0,
6252                AgentCapabilities::default(),
6253                [AgentEvent::TurnComplete { slot: 0 }],
6254            )),
6255            None,
6256        );
6257        let second = AdapterHost::new(
6258            Box::new(ScriptedAdapter::new(
6259                1,
6260                AgentCapabilities::default(),
6261                [AgentEvent::TurnComplete { slot: 1 }],
6262            )),
6263            None,
6264        );
6265        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6266        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6267        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
6268        relay.set_roster_launch_specs(vec![
6269            ("custom".into(), "claude".into()),
6270            ("custom".into(), "codex".into()),
6271        ]);
6272        let writer = SessionMetadataStore::open(&path)
6273            .buffered()
6274            .expect("metadata writer");
6275        relay.set_session_metadata_writer(writer);
6276        relay.start().await.expect("start");
6277        relay.stop().await.expect("stop");
6278        let loaded = SessionMetadataStore::open(&path)
6279            .read()
6280            .expect("read metadata")
6281            .expect("metadata snapshot");
6282        let agents = loaded
6283            .get("agents")
6284            .and_then(|value| value.as_array())
6285            .expect("agents");
6286        assert_eq!(agents.len(), 2);
6287        assert_eq!(agents[0]["identity"], "claude.com");
6288        assert_eq!(agents[1]["identity"], "openai.com");
6289        let _ = std::fs::remove_file(path);
6290    }
6291
6292    #[tokio::test]
6293    async fn relay_host_routes_unseen_public_context_to_next_agent() {
6294        let first = AdapterHost::new(
6295            Box::new(ScriptedAdapter::new(
6296                0,
6297                AgentCapabilities::default(),
6298                [
6299                    AgentEvent::Text {
6300                        slot: 0,
6301                        text: "implemented the fix".into(),
6302                    },
6303                    AgentEvent::TurnComplete { slot: 0 },
6304                ],
6305            )),
6306            None,
6307        );
6308        let second = AdapterHost::new(
6309            Box::new(ScriptedAdapter::new(
6310                1,
6311                AgentCapabilities::default(),
6312                [AgentEvent::TurnComplete { slot: 1 }],
6313            )),
6314            None,
6315        );
6316        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6317        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
6318        relay.start().await.expect("start");
6319        relay.run_turn("task", 0).await.expect("first turn");
6320        relay.run_turn("review this", 0).await.expect("review turn");
6321
6322        assert_eq!(relay.dispatches().len(), 2);
6323        assert_eq!(relay.dispatches()[0].0, 0);
6324        assert!(relay.dispatches()[0].1.contains("task"));
6325        assert!(relay.dispatches()[0].1.contains("You are Codex"));
6326        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
6327        assert_eq!(relay.dispatches()[1].0, 1);
6328        assert!(relay.dispatches()[1].1.contains("review this"));
6329        let public = relay.dispatches()[1]
6330            .1
6331            .split_once("Public updates:\n")
6332            .map(|(_, updates)| updates)
6333            .expect("review receives public context");
6334        let header = public
6335            .lines()
6336            .find(|line| line.starts_with("Codex "))
6337            .expect("named previous agent");
6338        let timestamp = header
6339            .strip_prefix("Codex ")
6340            .and_then(|value| value.strip_suffix(':'))
6341            .expect("timestamped header");
6342        assert_eq!(timestamp.len(), 5);
6343        assert_eq!(timestamp.as_bytes()[2], b':');
6344        assert!(
6345            timestamp
6346                .bytes()
6347                .enumerate()
6348                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
6349        );
6350        assert!(public.contains("implemented the fix"));
6351        assert!(!public.contains("Agent 0"));
6352    }
6353}