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, stop_token_visible_end, 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                // Only a possible split marker needs to wait for another chunk.
1644                let visible_response = response.replace(STOP_TOKEN, "");
1645                let visible_end = stop_token_visible_end(&visible_response);
1646                if emitted_text < visible_end {
1647                    if let Some(sink) = &self.event_sink {
1648                        sink(AgentEvent::Text {
1649                            slot: *slot,
1650                            text: visible_response[emitted_text..visible_end].to_owned(),
1651                        });
1652                    }
1653                    emitted_text = visible_end;
1654                }
1655            } else if let Some(sink) = &self.event_sink {
1656                sink(update.event.clone());
1657            }
1658        };
1659        let (response, requested_stop) = strip_stop_token(&response);
1660        let response = response.replace(STOP_TOKEN, "");
1661        let accepted_stop = requested_stop && *can_stop;
1662        let needs_stop_acknowledgment = accepted_stop && response.is_empty();
1663        let response = if needs_stop_acknowledgment {
1664            DEFAULT_STOP_ACKNOWLEDGMENT.to_owned()
1665        } else {
1666            response
1667        };
1668        // A token-only reviewer response is intentionally hidden, but the
1669        // UI still needs the documented visible acknowledgement. Streamed
1670        // text is normally emitted above; this synthetic acknowledgement is
1671        // the one case where there was no visible adapter chunk to forward.
1672        if needs_stop_acknowledgment && let Some(sink) = &self.event_sink {
1673            sink(AgentEvent::Text {
1674                slot: *slot,
1675                text: response.clone(),
1676            });
1677        }
1678        if let Some(sink) = &self.event_sink {
1679            sink(completion_event);
1680        }
1681        // A provider plan that ran out mid-turn routes future turns around
1682        // the agent instead of back into the exhausted quota.
1683        if is_usage_limit_response(&response) {
1684            let detail = response.clone();
1685            let _ = self.relay.mark_limited(*slot);
1686            // A normal finish: the ring cursor advances past the limited
1687            // agent so the batch continues with a healthy peer.
1688            self.relay.finish(*slot, *direct, false);
1689            self.queue_session_metadata()?;
1690            if let Some(sink) = &self.event_sink {
1691                sink(AgentEvent::UsageLimitReached {
1692                    slot: *slot,
1693                    detail,
1694                });
1695            }
1696            return Ok(decision);
1697        }
1698        if !*direct && !response.is_empty() {
1699            self.relay
1700                .record_public(public_context_speaker(&speaker_name), response);
1701        }
1702        self.relay.mark_context_seen(*slot);
1703        self.relay.finish(*slot, *direct, accepted_stop);
1704        self.queue_session_metadata()?;
1705        Ok(decision)
1706    }
1707}
1708
1709fn report_relay_failure(
1710    relay: &mut Relay,
1711    event_sink: &Option<Arc<dyn Fn(AgentEvent) + Send + Sync>>,
1712    slot: RosterSlot,
1713    started: bool,
1714    detail: String,
1715) -> bool {
1716    // A quota rejection is not a crash: route around the agent without
1717    // tombstoning the slot so a recharge can restore it.
1718    if is_usage_limit_response(&detail) {
1719        let _ = relay.mark_limited(slot);
1720        if let Some(sink) = event_sink {
1721            sink(AgentEvent::UsageLimitReached { slot, detail });
1722        }
1723        return true;
1724    }
1725    let _ = relay.tombstone(slot);
1726    if let Some(sink) = event_sink {
1727        sink(AgentEvent::Failed {
1728            slot,
1729            started,
1730            detail,
1731        });
1732    }
1733    false
1734}
1735
1736/// Deterministic in-memory adapter used for contract and relay tests.
1737#[derive(Debug)]
1738pub struct ScriptedAdapter {
1739    slot: RosterSlot,
1740    capabilities: AgentCapabilities,
1741    events: VecDeque<AdapterResult<AgentEvent>>,
1742    prompts: Vec<String>,
1743}
1744
1745impl ScriptedAdapter {
1746    pub fn new(
1747        slot: RosterSlot,
1748        capabilities: AgentCapabilities,
1749        events: impl IntoIterator<Item = AgentEvent>,
1750    ) -> Self {
1751        Self {
1752            slot,
1753            capabilities,
1754            events: events.into_iter().map(Ok).collect(),
1755            prompts: Vec::new(),
1756        }
1757    }
1758
1759    pub fn prompts(&self) -> &[String] {
1760        &self.prompts
1761    }
1762}
1763
1764#[async_trait]
1765impl AgentAdapter for ScriptedAdapter {
1766    fn slot(&self) -> RosterSlot {
1767        self.slot
1768    }
1769
1770    fn capabilities(&self) -> AgentCapabilities {
1771        self.capabilities.clone()
1772    }
1773
1774    async fn start(&mut self) -> AdapterResult<()> {
1775        Ok(())
1776    }
1777
1778    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1779        self.prompts.push(prompt);
1780        Ok(())
1781    }
1782
1783    async fn cancel(&mut self) -> AdapterResult<bool> {
1784        Ok(self.capabilities.supports_cancel)
1785    }
1786
1787    async fn answer_permission(
1788        &mut self,
1789        _request_id: String,
1790        _answer: PermissionAnswer,
1791    ) -> AdapterResult<()> {
1792        if self.capabilities.supports_permissions {
1793            Ok(())
1794        } else {
1795            Err(AdapterError::Unsupported("permission answer"))
1796        }
1797    }
1798
1799    async fn set_mode(&mut self, _mode: String) -> AdapterResult<()> {
1800        if self.capabilities.supports_modes {
1801            Ok(())
1802        } else {
1803            Err(AdapterError::Unsupported("set_mode"))
1804        }
1805    }
1806
1807    async fn reload(&mut self) -> AdapterResult<()> {
1808        Ok(())
1809    }
1810
1811    async fn stop(&mut self) -> AdapterResult<()> {
1812        Ok(())
1813    }
1814
1815    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
1816        self.events.pop_front()
1817    }
1818}
1819
1820/// A direct stream-JSON adapter for Antigravity. It deliberately does not
1821/// pretend to be ACP; it translates its documented events into core events.
1822#[derive(Debug)]
1823pub struct AgyAdapter {
1824    slot: RosterSlot,
1825    cwd: PathBuf,
1826    command: String,
1827    mode: String,
1828    mode_policy: String,
1829    session_id: Option<String>,
1830    child: Option<Child>,
1831    sender: mpsc::Sender<AdapterResult<AgentEvent>>,
1832    receiver: mpsc::Receiver<AdapterResult<AgentEvent>>,
1833    /// Antigravity announces its conversation id in the `init` event. Keep it
1834    /// outside the stream task so the next prompt can resume the same native
1835    /// conversation without making the adapter reader/UI share a borrow.
1836    announced_session: Arc<Mutex<Option<String>>>,
1837    cancel_requested: Arc<AtomicBool>,
1838}
1839
1840impl AgyAdapter {
1841    pub fn new(slot: RosterSlot, cwd: PathBuf, command: impl Into<String>) -> Self {
1842        let (sender, receiver) = mpsc::channel(256);
1843        Self {
1844            slot,
1845            cwd,
1846            command: command.into(),
1847            mode: "default".into(),
1848            mode_policy: "agy:full-access".into(),
1849            session_id: None,
1850            child: None,
1851            sender,
1852            receiver,
1853            announced_session: Arc::new(Mutex::new(None)),
1854            cancel_requested: Arc::new(AtomicBool::new(false)),
1855        }
1856    }
1857
1858    pub fn with_session_id(
1859        slot: RosterSlot,
1860        cwd: PathBuf,
1861        command: impl Into<String>,
1862        session_id: impl Into<String>,
1863    ) -> Self {
1864        let mut adapter = Self::new(slot, cwd, command);
1865        adapter.session_id = Some(session_id.into());
1866        adapter
1867    }
1868
1869    fn modes() -> Vec<Mode> {
1870        vec![
1871            Mode {
1872                id: "agy:full-access".into(),
1873                label: "Auto pilot".into(),
1874            },
1875            Mode {
1876                id: "agy:manual".into(),
1877                label: "Manual".into(),
1878            },
1879            Mode {
1880                id: "accept-edits".into(),
1881                label: "Accept Edits".into(),
1882            },
1883            Mode {
1884                id: "plan".into(),
1885                label: "Plan".into(),
1886            },
1887        ]
1888    }
1889
1890    async fn emit(&self, event: AdapterResult<AgentEvent>) {
1891        let _ = self.sender.send(event).await;
1892    }
1893}
1894
1895#[async_trait]
1896impl AgentAdapter for AgyAdapter {
1897    fn slot(&self) -> RosterSlot {
1898        self.slot
1899    }
1900
1901    fn session_id(&self) -> Option<String> {
1902        self.session_id.clone()
1903    }
1904
1905    fn protocol(&self) -> &'static str {
1906        "native"
1907    }
1908
1909    fn capabilities(&self) -> AgentCapabilities {
1910        AgentCapabilities {
1911            supports_cancel: true,
1912            supports_modes: true,
1913            supports_permissions: false,
1914            supports_terminals: true,
1915            supports_session_load: true,
1916            supports_models: false,
1917        }
1918    }
1919
1920    async fn start(&mut self) -> AdapterResult<()> {
1921        self.cancel_requested.store(false, Ordering::Release);
1922        self.emit(Ok(AgentEvent::Ready {
1923            slot: self.slot,
1924            capabilities: self.capabilities(),
1925        }))
1926        .await;
1927        self.emit(Ok(AgentEvent::ModesReplaced {
1928            slot: self.slot,
1929            modes: Self::modes(),
1930            current_mode: Some(self.mode_policy.clone()),
1931        }))
1932        .await;
1933        Ok(())
1934    }
1935
1936    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
1937        if self.child.is_some() {
1938            return Err(AdapterError::Transport(
1939                "agent is already handling a turn".into(),
1940            ));
1941        }
1942        if self.session_id.is_none()
1943            && let Ok(session) = self.announced_session.lock()
1944        {
1945            self.session_id = session.clone();
1946        }
1947        self.cancel_requested.store(false, Ordering::Release);
1948        let (program, args) = parse_command_line(&self.command)
1949            .map_err(|error| AdapterError::Spawn(format!("invalid agent command: {error}")))?;
1950        let mut command = Command::new(program);
1951        isolate_process_group(&mut command);
1952        command
1953            .args(args)
1954            .arg("--print")
1955            .arg(prompt)
1956            .arg("--print-timeout")
1957            .arg("60m")
1958            .arg("--output-format")
1959            .arg("stream-json")
1960            .current_dir(&self.cwd)
1961            .env("CODESWARM_CWD", &self.cwd)
1962            .stdout(Stdio::piped())
1963            .stderr(Stdio::piped());
1964        if let Some(session_id) = &self.session_id {
1965            command.arg("--conversation").arg(session_id);
1966        }
1967        if self.mode != "default" {
1968            command.arg("--mode").arg(&self.mode);
1969        }
1970        let mut child = command
1971            .spawn()
1972            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
1973        let stdout = match child.stdout.take() {
1974            Some(stdout) => stdout,
1975            None => {
1976                let _ = terminate_child(&mut child).await;
1977                return Err(AdapterError::Transport("agent has no stdout".into()));
1978            }
1979        };
1980        let stderr = match child.stderr.take() {
1981            Some(stderr) => stderr,
1982            None => {
1983                let _ = terminate_child(&mut child).await;
1984                return Err(AdapterError::Transport("agent has no stderr".into()));
1985            }
1986        };
1987        let sender = self.sender.clone();
1988        let slot = self.slot;
1989        let announced_session = Arc::clone(&self.announced_session);
1990        let cancel_requested = Arc::clone(&self.cancel_requested);
1991        tokio::spawn(async move {
1992            let stderr_task = tokio::spawn(async move {
1993                const MAX_STDERR: usize = 32 * 1024;
1994                let mut stderr = BufReader::new(stderr);
1995                let mut bytes = Vec::new();
1996                let mut chunk = [0_u8; 4096];
1997                while let Ok(count) = stderr.read(&mut chunk).await {
1998                    if count == 0 {
1999                        break;
2000                    }
2001                    bytes.extend_from_slice(&chunk[..count]);
2002                    if bytes.len() > MAX_STDERR {
2003                        let keep_from = bytes.len() - MAX_STDERR;
2004                        bytes.drain(..keep_from);
2005                    }
2006                }
2007                String::from_utf8_lossy(&bytes).trim().to_owned()
2008            });
2009            let mut lines = BufReader::new(stdout).lines();
2010            let mut result: Option<Value> = None;
2011            let mut streamed_response = false;
2012            while let Ok(Some(line)) = lines.next_line().await {
2013                let value = match serde_json::from_str::<Value>(&line) {
2014                    Ok(value) => value,
2015                    Err(_) => {
2016                        // Native stream-json can contain diagnostic junk on
2017                        // stdout. Match the Python adapter's tolerant stream
2018                        // behavior and wait for the final result instead of
2019                        // turning one malformed line into a dead turn.
2020                        continue;
2021                    }
2022                };
2023                if value.get("event").and_then(Value::as_str) == Some("init")
2024                    && let Some(session_id) = value
2025                        .get("conversation_id")
2026                        .or_else(|| value.get("conversationId"))
2027                        .and_then(Value::as_str)
2028                        .filter(|id| !id.is_empty())
2029                    && let Ok(mut announced) = announced_session.lock()
2030                {
2031                    *announced = Some(session_id.to_owned());
2032                }
2033                if value.get("event").and_then(Value::as_str) == Some("result") {
2034                    result = value.get("result").cloned();
2035                }
2036                match parse_agy_value(slot, &value) {
2037                    Ok(Some(event)) => {
2038                        if matches!(event, AgentEvent::Text { .. }) {
2039                            streamed_response = true;
2040                        }
2041                        if sender.send(Ok(event)).await.is_err() {
2042                            break;
2043                        }
2044                    }
2045                    Ok(None) => {}
2046                    Err(error) => {
2047                        let _ = sender.send(Err(error)).await;
2048                    }
2049                }
2050            }
2051            let stderr = stderr_task.await.ok().unwrap_or_default();
2052            let succeeded = cancel_requested.load(Ordering::Acquire)
2053                || result
2054                    .as_ref()
2055                    .and_then(|result| result.get("status"))
2056                    .and_then(Value::as_str)
2057                    == Some("SUCCESS");
2058            if succeeded {
2059                // Some native stream-json wrappers emit only lifecycle events
2060                // and put the complete answer in the final result object. The
2061                // Python adapter surfaced that answer; do the same, while
2062                // avoiding duplication when token chunks were already sent.
2063                if !streamed_response
2064                    && let Some(response) = result
2065                        .as_ref()
2066                        .and_then(|result| result.get("response"))
2067                        .and_then(Value::as_str)
2068                        .filter(|response| !response.is_empty())
2069                {
2070                    let _ = sender
2071                        .send(Ok(AgentEvent::Text {
2072                            slot,
2073                            text: response.to_owned(),
2074                        }))
2075                        .await;
2076                }
2077                let _ = sender.send(Ok(AgentEvent::TurnComplete { slot })).await;
2078            } else {
2079                let detail = result
2080                    .as_ref()
2081                    .and_then(|result| result.get("error"))
2082                    .and_then(Value::as_str)
2083                    .filter(|detail| !detail.is_empty())
2084                    .map(str::to_owned)
2085                    .or_else(|| (!stderr.is_empty()).then_some(stderr))
2086                    .unwrap_or_else(|| "native stream ended before a successful result".into());
2087                let _ = sender
2088                    .send(Ok(AgentEvent::Failed {
2089                        slot,
2090                        started: true,
2091                        detail,
2092                    }))
2093                    .await;
2094            }
2095        });
2096        self.child = Some(child);
2097        Ok(())
2098    }
2099
2100    async fn cancel(&mut self) -> AdapterResult<bool> {
2101        self.cancel_requested.store(true, Ordering::Release);
2102        let Some(mut child) = self.child.take() else {
2103            return Ok(false);
2104        };
2105        // `Child` does not reap itself when dropped. Awaiting wait after the
2106        // kill keeps repeated prompts from accumulating zombies, especially
2107        // when cancellation happens before the stream reader observes EOF.
2108        terminate_child(&mut child).await?;
2109        let _ = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
2110            while let Some(event) = self.receiver.recv().await {
2111                if matches!(
2112                    event,
2113                    Ok(AgentEvent::TurnComplete { .. } | AgentEvent::Failed { .. })
2114                ) {
2115                    break;
2116                }
2117            }
2118        })
2119        .await;
2120        Ok(true)
2121    }
2122
2123    async fn answer_permission(
2124        &mut self,
2125        _request_id: String,
2126        _answer: PermissionAnswer,
2127    ) -> AdapterResult<()> {
2128        Err(AdapterError::Unsupported("permission answer"))
2129    }
2130
2131    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
2132        let (mode, mode_policy) = match mode.as_str() {
2133            "full-access" | "codeswarm:mode:full-access" | "auto" | "autopilot" => {
2134                ("default".to_owned(), "agy:full-access".to_owned())
2135            }
2136            "codeswarm:mode:plan" | "readonly" | "plan" => ("plan".to_owned(), "plan".to_owned()),
2137            "codeswarm:mode:accept-edits" | "acceptedits" | "accept-edits" => {
2138                ("accept-edits".to_owned(), "accept-edits".to_owned())
2139            }
2140            "codeswarm:mode:manual" | "manual" | "ask" | "default" => {
2141                ("default".to_owned(), "agy:manual".to_owned())
2142            }
2143            "agy:full-access" => ("default".to_owned(), "agy:full-access".to_owned()),
2144            "agy:manual" => ("default".to_owned(), "agy:manual".to_owned()),
2145            _ => return Err(AdapterError::Unsupported("requested Agy mode")),
2146        };
2147        self.mode = mode;
2148        self.mode_policy = mode_policy.clone();
2149        self.emit(Ok(AgentEvent::ModesReplaced {
2150            slot: self.slot,
2151            modes: Self::modes(),
2152            current_mode: Some(mode_policy),
2153        }))
2154        .await;
2155        Ok(())
2156    }
2157
2158    async fn reload(&mut self) -> AdapterResult<()> {
2159        self.stop().await?;
2160        self.start().await
2161    }
2162
2163    async fn stop(&mut self) -> AdapterResult<()> {
2164        let _ = self.cancel().await?;
2165        Ok(())
2166    }
2167
2168    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
2169        let event = self.receiver.recv().await;
2170        if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2171            && self.session_id.is_none()
2172            && let Ok(session) = self.announced_session.lock()
2173        {
2174            self.session_id = session.clone();
2175        }
2176        if matches!(event.as_ref(), Some(Ok(AgentEvent::TurnComplete { .. })))
2177            && let Some(mut child) = self.child.take()
2178        {
2179            let _ = child.wait().await;
2180        }
2181        event
2182    }
2183}
2184
2185#[cfg(test)]
2186#[cfg_attr(not(test), allow(dead_code))]
2187fn parse_agy_line(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
2188    let value: Value =
2189        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
2190    parse_agy_value(slot, &value)
2191}
2192
2193fn parse_agy_value(slot: RosterSlot, value: &Value) -> AdapterResult<Option<AgentEvent>> {
2194    let event = value.get("event").and_then(Value::as_str);
2195    if let Some(terminal) = parse_terminal_event(value, event) {
2196        return Ok(Some(AgentEvent::Terminal {
2197            slot,
2198            event: terminal,
2199        }));
2200    }
2201    match event {
2202        Some("step_update") => {
2203            let Some(update) = value.get("step_update") else {
2204                return Ok(None);
2205            };
2206            let is_response = update
2207                .get("step_type")
2208                .and_then(Value::as_str)
2209                .is_some_and(|kind| kind == "agent_response");
2210            let text = update.get("text_delta").and_then(Value::as_str);
2211            let response = is_response
2212                .then(|| text.map(str::to_owned))
2213                .flatten()
2214                .filter(|text| !text.is_empty())
2215                .map(|text| AgentEvent::Text { slot, text });
2216            Ok(response.or_else(|| parse_agy_tool(slot, value)))
2217        }
2218        _ => Ok(None),
2219    }
2220}
2221
2222fn parse_agy_tool(slot: RosterSlot, value: &Value) -> Option<AgentEvent> {
2223    let update = value.get("step_update")?;
2224    if update.get("step_type")?.as_str()? != "tool" {
2225        return None;
2226    }
2227    let step_index = update.get("step_index")?.as_i64()?;
2228    let title = update
2229        .get("tool_name")
2230        .and_then(Value::as_str)
2231        .unwrap_or("Tool call")
2232        .replace('_', " ");
2233    let status = match update.get("state").and_then(Value::as_str) {
2234        Some("DONE") => ToolStatus::Completed,
2235        Some("FAILED") => ToolStatus::Failed,
2236        Some("ACTIVE") => ToolStatus::Running,
2237        _ => ToolStatus::Pending,
2238    };
2239    let detail = update
2240        .get("tool_info")
2241        .and_then(|info| info.get("output"))
2242        .and_then(Value::as_str)
2243        .map(str::to_owned);
2244    Some(AgentEvent::Tool {
2245        slot,
2246        update: ToolUpdate {
2247            id: format!("agy-tool-{step_index}"),
2248            title,
2249            status,
2250            detail,
2251        },
2252    })
2253}
2254
2255/// Stdio ACP transport. Protocol-specific response handling belongs here,
2256/// keeping JSON-RPC framing outside the core and terminal renderer.
2257#[derive(Debug)]
2258pub struct AcpAdapter {
2259    slot: RosterSlot,
2260    program: String,
2261    args: Vec<String>,
2262    cwd: PathBuf,
2263    child: Option<Child>,
2264    reader: Option<BufReader<ChildStdout>>,
2265    capabilities: AgentCapabilities,
2266    modes: Vec<Mode>,
2267    models: Vec<Mode>,
2268    model_config_id: Option<String>,
2269    session_id: Option<String>,
2270    next_request_id: u64,
2271    prompt_request_id: Option<u64>,
2272    queued_events: VecDeque<AdapterResult<AgentEvent>>,
2273    stderr_task: Option<tokio::task::JoinHandle<String>>,
2274    terminals: BTreeMap<String, TerminalProcess>,
2275    next_terminal_id: u64,
2276}
2277
2278impl AcpAdapter {
2279    pub fn new(
2280        slot: RosterSlot,
2281        cwd: PathBuf,
2282        program: impl Into<String>,
2283        args: Vec<String>,
2284    ) -> Self {
2285        Self {
2286            slot,
2287            program: program.into(),
2288            args,
2289            cwd,
2290            child: None,
2291            reader: None,
2292            capabilities: AgentCapabilities::default(),
2293            modes: Vec::new(),
2294            models: Vec::new(),
2295            model_config_id: None,
2296            session_id: None,
2297            next_request_id: 1,
2298            prompt_request_id: None,
2299            queued_events: VecDeque::new(),
2300            stderr_task: None,
2301            terminals: BTreeMap::new(),
2302            next_terminal_id: 1,
2303        }
2304    }
2305
2306    pub fn with_session_id(
2307        slot: RosterSlot,
2308        cwd: PathBuf,
2309        program: impl Into<String>,
2310        args: Vec<String>,
2311        session_id: impl Into<String>,
2312    ) -> Self {
2313        let mut adapter = Self::new(slot, cwd, program, args);
2314        adapter.session_id = Some(session_id.into());
2315        adapter
2316    }
2317
2318    async fn request(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2319        self.request_with_timeout(method, params, std::time::Duration::from_secs(30))
2320            .await
2321    }
2322
2323    async fn request_with_timeout(
2324        &mut self,
2325        method: &str,
2326        params: Value,
2327        deadline: std::time::Duration,
2328    ) -> AdapterResult<Value> {
2329        tokio::time::timeout(deadline, self.request_inner(method, params))
2330            .await
2331            .map_err(|_| {
2332                AdapterError::Transport(format!(
2333                    "ACP {method} timed out; reload the agent to retry"
2334                ))
2335            })?
2336    }
2337
2338    async fn request_inner(&mut self, method: &str, params: Value) -> AdapterResult<Value> {
2339        let request_id = self.next_request_id;
2340        self.next_request_id += 1;
2341        self.write_json(serde_json::json!({
2342            "jsonrpc": "2.0",
2343            "id": request_id,
2344            "method": method,
2345            "params": params,
2346        }))
2347        .await?;
2348        loop {
2349            let line = self.read_line().await?;
2350            let value: Value = match serde_json::from_str(&line) {
2351                Ok(value) => value,
2352                Err(_) => {
2353                    // Keep the transport alive when a peer writes a stray
2354                    // diagnostic line. This is common with CLI wrappers and
2355                    // matches the baseline client's tolerant stream loop.
2356                    continue;
2357                }
2358            };
2359            if self.reject_empty_permission_request(&value).await? {
2360                continue;
2361            }
2362            if self.handle_client_request(&value).await? {
2363                continue;
2364            }
2365            if value
2366                .get("id")
2367                .is_some_and(|id| rpc_id_to_string(id) == request_id.to_string())
2368            {
2369                if let Some(error) = value.get("error") {
2370                    return Err(AdapterError::Protocol(error.to_string()));
2371                }
2372                return value
2373                    .get("result")
2374                    .cloned()
2375                    .ok_or_else(|| AdapterError::Protocol("response has no result".into()));
2376            }
2377            if let Some(event) = parse_acp_notification(self.slot, &line)? {
2378                self.queued_events.push_back(Ok(event));
2379            }
2380        }
2381    }
2382
2383    async fn write_json(&mut self, value: Value) -> AdapterResult<()> {
2384        let child = self
2385            .child
2386            .as_mut()
2387            .ok_or_else(|| AdapterError::Transport("ACP agent is not running".into()))?;
2388        let stdin = child
2389            .stdin
2390            .as_mut()
2391            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdin".into()))?;
2392        stdin
2393            .write_all(value.to_string().as_bytes())
2394            .await
2395            .map_err(|error| AdapterError::Transport(error.to_string()))?;
2396        stdin
2397            .write_all(b"\n")
2398            .await
2399            .map_err(|error| AdapterError::Transport(error.to_string()))
2400    }
2401
2402    /// ACP permission requests are JSON-RPC requests, not fire-and-forget
2403    /// notifications. An empty option list is invalid and must be answered
2404    /// with an error so the peer does not wait forever for a decision. This is
2405    /// the same validation performed by the Python ACP server.
2406    async fn reject_empty_permission_request(&mut self, value: &Value) -> AdapterResult<bool> {
2407        if value.get("method").and_then(Value::as_str) != Some("session/request_permission")
2408            || value.get("id").is_none()
2409        {
2410            return Ok(false);
2411        }
2412        let valid = value
2413            .get("params")
2414            .and_then(|params| params.get("options"))
2415            .and_then(Value::as_array)
2416            .is_some_and(|options| !options.is_empty());
2417        if valid {
2418            return Ok(false);
2419        }
2420        self.write_json(serde_json::json!({
2421            "jsonrpc": "2.0",
2422            "id": value.get("id").cloned().unwrap_or(Value::Null),
2423            "error": {
2424                "code": -32602,
2425                "message": "Permission request requires at least one option",
2426            },
2427        }))
2428        .await?;
2429        Ok(true)
2430    }
2431
2432    fn workspace_path(&self, path: &str) -> Result<PathBuf, String> {
2433        let root = self
2434            .cwd
2435            .canonicalize()
2436            .map_err(|error| format!("unable to resolve workspace: {error}"))?;
2437        let requested = Path::new(path);
2438        let candidate = if requested.is_absolute() {
2439            requested.to_path_buf()
2440        } else {
2441            root.join(requested)
2442        };
2443        let resolved = if !candidate.exists() {
2444            let parent = candidate
2445                .parent()
2446                .ok_or_else(|| "file path has no parent".to_owned())?
2447                .canonicalize()
2448                .map_err(|error| format!("unable to resolve parent directory: {error}"))?;
2449            parent.join(
2450                candidate
2451                    .file_name()
2452                    .ok_or_else(|| "file path has no filename".to_owned())?,
2453            )
2454        } else {
2455            candidate
2456                .canonicalize()
2457                .map_err(|error| format!("unable to resolve file path: {error}"))?
2458        };
2459        if !resolved.starts_with(&root) {
2460            return Err("file path is outside the project".into());
2461        }
2462        Ok(resolved)
2463    }
2464
2465    fn read_workspace_text(
2466        &self,
2467        path: &str,
2468        line: Option<i64>,
2469        limit: Option<i64>,
2470    ) -> Result<String, String> {
2471        if line.is_some_and(|line| line < 1) {
2472            return Err("line must be positive".into());
2473        }
2474        if limit.is_some_and(|limit| limit < 0) {
2475            return Err("limit must not be negative".into());
2476        }
2477        let path = self.workspace_path(path)?;
2478        let mut bytes = Vec::new();
2479        let mut source = match std::fs::File::open(path) {
2480            Ok(source) => source,
2481            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(String::new()),
2482            Err(error) => return Err(error.to_string()),
2483        };
2484        source
2485            .by_ref()
2486            .take((MAX_FILE_READ_BYTES as u64).saturating_add(1))
2487            .read_to_end(&mut bytes)
2488            .map_err(|error| error.to_string())?;
2489        bytes.truncate(MAX_FILE_READ_BYTES);
2490        let text = String::from_utf8_lossy(&bytes);
2491        if line.is_none() && limit.is_none() {
2492            return Ok(text.into_owned());
2493        }
2494        let start = line.map_or(0, |line| line as usize - 1);
2495        let limit = limit.unwrap_or(i64::MAX) as usize;
2496        let selected = text
2497            .split_inclusive('\n')
2498            .skip(start)
2499            .take(limit)
2500            .collect::<String>();
2501        if line.is_some() {
2502            Ok(selected.trim_end_matches('\n').to_owned())
2503        } else {
2504            Ok(selected)
2505        }
2506    }
2507
2508    fn write_workspace_text(&self, params: &Value) -> Result<(), String> {
2509        let path = params
2510            .get("path")
2511            .and_then(Value::as_str)
2512            .filter(|path| !path.is_empty())
2513            .ok_or("path must be a non-empty string")?;
2514        let content = params
2515            .get("content")
2516            .and_then(Value::as_str)
2517            .ok_or("content must be a string")?;
2518        let path = self.workspace_path(path)?;
2519        std::fs::write(path, content).map_err(|error| error.to_string())
2520    }
2521
2522    async fn terminal_create(&mut self, params: &Value) -> Result<Value, String> {
2523        let command = params
2524            .get("command")
2525            .and_then(Value::as_str)
2526            .filter(|command| !command.trim().is_empty())
2527            .ok_or_else(|| "terminal command is required".to_owned())?;
2528        let cwd = params.get("cwd").and_then(Value::as_str).unwrap_or(".");
2529        let cwd = self.workspace_path(cwd)?;
2530        if !cwd.is_dir() {
2531            return Err("terminal cwd is not a directory".into());
2532        }
2533        let mut process = Command::new(command);
2534        isolate_process_group(&mut process);
2535        if let Some(args) = params.get("args").and_then(Value::as_array) {
2536            process.args(args.iter().filter_map(Value::as_str));
2537        }
2538        process
2539            .current_dir(&cwd)
2540            .stdin(Stdio::null())
2541            .stdout(Stdio::piped())
2542            .stderr(Stdio::piped());
2543        if let Some(env) = params.get("env") {
2544            if let Some(entries) = env.as_array() {
2545                for entry in entries {
2546                    if let (Some(name), Some(value)) = (
2547                        entry.get("name").and_then(Value::as_str),
2548                        entry.get("value").and_then(Value::as_str),
2549                    ) {
2550                        process.env(name, value);
2551                    }
2552                }
2553            } else if let Some(entries) = env.as_object() {
2554                for (name, value) in entries {
2555                    if let Some(value) = value.as_str() {
2556                        process.env(name, value);
2557                    }
2558                }
2559            }
2560        }
2561        let mut child = process.spawn().map_err(|error| error.to_string())?;
2562        let stdout = child.stdout.take();
2563        let stderr = child.stderr.take();
2564        let (Some(stdout), Some(stderr)) = (stdout, stderr) else {
2565            let _ = terminate_child(&mut child).await;
2566            return Err("terminal has no output pipes".into());
2567        };
2568        let output = Arc::new(Mutex::new(Vec::new()));
2569        let truncated = Arc::new(AtomicBool::new(false));
2570        let output_readers = Arc::new(AtomicUsize::new(2));
2571        let output_limit = params
2572            .get("outputByteLimit")
2573            .and_then(Value::as_u64)
2574            .map_or(MAX_TERMINAL_OUTPUT_BYTES, |limit| {
2575                usize::try_from(limit)
2576                    .unwrap_or(MAX_TERMINAL_OUTPUT_BYTES)
2577                    .min(MAX_TERMINAL_OUTPUT_BYTES)
2578            });
2579        tokio::spawn(drain_terminal_output(
2580            stdout,
2581            Arc::clone(&output),
2582            Arc::clone(&truncated),
2583            Arc::clone(&output_readers),
2584            output_limit,
2585        ));
2586        tokio::spawn(drain_terminal_output(
2587            stderr,
2588            Arc::clone(&output),
2589            Arc::clone(&truncated),
2590            Arc::clone(&output_readers),
2591            output_limit,
2592        ));
2593        let id = format!("terminal-{}", self.next_terminal_id);
2594        self.next_terminal_id = self.next_terminal_id.saturating_add(1);
2595        let state = TerminalProcess {
2596            child: Arc::new(AsyncMutex::new(Some(child))),
2597            output,
2598            truncated,
2599            output_readers,
2600        };
2601        self.terminals.insert(id.clone(), state);
2602        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2603            slot: self.slot,
2604            event: TerminalEvent::Created {
2605                id: id.clone(),
2606                command: std::iter::once(command)
2607                    .chain(
2608                        params
2609                            .get("args")
2610                            .and_then(Value::as_array)
2611                            .into_iter()
2612                            .flatten()
2613                            .filter_map(Value::as_str),
2614                    )
2615                    .collect::<Vec<_>>()
2616                    .join(" "),
2617            },
2618        }));
2619        Ok(serde_json::json!({"terminalId": id}))
2620    }
2621
2622    async fn terminal_output(&mut self, id: &str) -> Result<Value, String> {
2623        let terminal = self
2624            .terminals
2625            .get(id)
2626            .ok_or_else(|| "terminal not found".to_owned())?
2627            .clone();
2628        let output = terminal
2629            .output
2630            .lock()
2631            .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
2632            .unwrap_or_default();
2633        let exit_code = terminal.exit_code().await;
2634        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2635            slot: self.slot,
2636            event: TerminalEvent::Output {
2637                id: id.to_owned(),
2638                text: output.clone(),
2639            },
2640        }));
2641        let mut response = serde_json::json!({
2642            "output": output,
2643            "truncated": terminal.truncated.load(Ordering::Acquire),
2644        });
2645        if let Some(code) = exit_code {
2646            response["exitStatus"] = serde_json::json!({"exitCode": code});
2647        }
2648        Ok(response)
2649    }
2650
2651    async fn terminal_wait(&mut self, id: &str) -> Result<Value, String> {
2652        let terminal = self
2653            .terminals
2654            .get(id)
2655            .ok_or_else(|| "terminal not found".to_owned())?
2656            .clone();
2657        let exit_code = terminal.wait().await;
2658        self.queued_events.push_back(Ok(AgentEvent::Terminal {
2659            slot: self.slot,
2660            event: TerminalEvent::Exited {
2661                id: id.to_owned(),
2662                code: exit_code.unwrap_or(-1),
2663            },
2664        }));
2665        Ok(serde_json::json!({"exitCode": exit_code, "signal": Value::Null}))
2666    }
2667
2668    /// Handle requests initiated by an ACP agent against the client. File
2669    /// access is mediated through the configured workspace root; unsupported
2670    /// requests receive a JSON-RPC error instead of hanging the agent.
2671    async fn handle_client_request(&mut self, value: &Value) -> AdapterResult<bool> {
2672        let Some(method) = value.get("method").and_then(Value::as_str) else {
2673            return Ok(false);
2674        };
2675        let Some(id) = value.get("id").cloned() else {
2676            return Ok(false);
2677        };
2678        // Permission requests are consumed by the normalized event parser and
2679        // answered later through the focused UI action, not here.
2680        if method == "session/request_permission" {
2681            return Ok(false);
2682        }
2683        let params = value.get("params").cloned().unwrap_or(Value::Null);
2684        let response = match method {
2685            "fs/read_text_file" => {
2686                let path = params.get("path").and_then(Value::as_str).unwrap_or("");
2687                let line = params.get("line").and_then(Value::as_i64);
2688                let limit = params.get("limit").and_then(Value::as_i64);
2689                match self.read_workspace_text(path, line, limit) {
2690                    Ok(content) => serde_json::json!({
2691                        "jsonrpc": "2.0",
2692                        "id": id,
2693                        "result": {"content": content},
2694                    }),
2695                    Err(message) => serde_json::json!({
2696                        "jsonrpc": "2.0",
2697                        "id": id,
2698                        "error": {"code": -32602, "message": message},
2699                    }),
2700                }
2701            }
2702            "fs/write_text_file" => {
2703                let result = self.write_workspace_text(&params);
2704                match result {
2705                    Ok(()) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}}),
2706                    Err(message) => serde_json::json!({
2707                        "jsonrpc": "2.0",
2708                        "id": id,
2709                        "error": {"code": -32602, "message": message},
2710                    }),
2711                }
2712            }
2713            "terminal/create" => match self.terminal_create(&params).await {
2714                Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2715                Err(message) => serde_json::json!({
2716                    "jsonrpc": "2.0",
2717                    "id": id,
2718                    "error": {"code": -32602, "message": message},
2719                }),
2720            },
2721            "terminal/output" => {
2722                let terminal_id = params
2723                    .get("terminalId")
2724                    .and_then(Value::as_str)
2725                    .unwrap_or("");
2726                match self.terminal_output(terminal_id).await {
2727                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2728                    Err(message) => serde_json::json!({
2729                        "jsonrpc": "2.0",
2730                        "id": id,
2731                        "error": {"code": -32602, "message": message},
2732                    }),
2733                }
2734            }
2735            "terminal/wait_for_exit" => {
2736                let terminal_id = params
2737                    .get("terminalId")
2738                    .and_then(Value::as_str)
2739                    .unwrap_or("");
2740                match self.terminal_wait(terminal_id).await {
2741                    Ok(result) => serde_json::json!({"jsonrpc": "2.0", "id": id, "result": result}),
2742                    Err(message) => serde_json::json!({
2743                        "jsonrpc": "2.0",
2744                        "id": id,
2745                        "error": {"code": -32602, "message": message},
2746                    }),
2747                }
2748            }
2749            "terminal/kill" => {
2750                let terminal_id = params
2751                    .get("terminalId")
2752                    .and_then(Value::as_str)
2753                    .unwrap_or("");
2754                if let Some(terminal) = self.terminals.get(terminal_id) {
2755                    terminal.kill().await;
2756                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2757                } else {
2758                    serde_json::json!({
2759                        "jsonrpc": "2.0",
2760                        "id": id,
2761                        "error": {"code": -32602, "message": "terminal not found"},
2762                    })
2763                }
2764            }
2765            "terminal/release" => {
2766                let terminal_id = params
2767                    .get("terminalId")
2768                    .and_then(Value::as_str)
2769                    .unwrap_or("");
2770                if let Some(terminal) = self.terminals.remove(terminal_id) {
2771                    terminal.stop().await;
2772                    self.queued_events.push_back(Ok(AgentEvent::Terminal {
2773                        slot: self.slot,
2774                        event: TerminalEvent::Released {
2775                            id: terminal_id.to_owned(),
2776                        },
2777                    }));
2778                    serde_json::json!({"jsonrpc": "2.0", "id": id, "result": {}})
2779                } else {
2780                    serde_json::json!({
2781                        "jsonrpc": "2.0",
2782                        "id": id,
2783                        "error": {"code": -32602, "message": "terminal not found"},
2784                    })
2785                }
2786            }
2787            _ => serde_json::json!({
2788                "jsonrpc": "2.0",
2789                "id": id,
2790                "error": {"code": -32601, "message": format!("unsupported client method: {method}")},
2791            }),
2792        };
2793        self.write_json(response).await?;
2794        Ok(true)
2795    }
2796
2797    async fn read_line(&mut self) -> AdapterResult<String> {
2798        let reader = self
2799            .reader
2800            .as_mut()
2801            .ok_or_else(|| AdapterError::Transport("ACP agent has no stdout".into()))?;
2802        read_bounded_line(reader).await
2803    }
2804
2805    async fn start(&mut self) -> AdapterResult<()> {
2806        // Starting an adapter twice must never orphan the first transport.
2807        if self.child.is_some() {
2808            self.stop().await?;
2809        }
2810        let mut command = Command::new(&self.program);
2811        isolate_process_group(&mut command);
2812        command
2813            .args(&self.args)
2814            .current_dir(&self.cwd)
2815            .stdin(Stdio::piped())
2816            .stdout(Stdio::piped())
2817            .stderr(Stdio::piped())
2818            .env("CODESWARM_CWD", &self.cwd);
2819        if self.program.to_ascii_lowercase().contains("gemini")
2820            || self
2821                .args
2822                .iter()
2823                .any(|arg| arg.to_ascii_lowercase().contains("gemini"))
2824        {
2825            command.env("GEMINI_TELEMETRY_ENABLED", "false");
2826        }
2827        let mut child = command
2828            .spawn()
2829            .map_err(|error| AdapterError::Spawn(error.to_string()))?;
2830        let stdout = match child.stdout.take() {
2831            Some(stdout) => stdout,
2832            None => {
2833                let _ = terminate_child(&mut child).await;
2834                return Err(AdapterError::Transport("ACP agent has no stdout".into()));
2835            }
2836        };
2837        let stderr = match child.stderr.take() {
2838            Some(stderr) => stderr,
2839            None => {
2840                let _ = terminate_child(&mut child).await;
2841                return Err(AdapterError::Transport("ACP agent has no stderr".into()));
2842            }
2843        };
2844        self.child = Some(child);
2845        self.reader = Some(BufReader::new(stdout));
2846        self.stderr_task = Some(tokio::spawn(drain_bounded(stderr, 32 * 1024)));
2847
2848        let initialize = match self
2849            .request(
2850                "initialize",
2851                serde_json::json!({
2852                    "protocolVersion": 1,
2853                    "clientCapabilities": {
2854                        "fs": {"readTextFile": true, "writeTextFile": true},
2855                        "terminal": true,
2856                    },
2857                    "clientInfo": {
2858                        "name": "CodeSwarm",
2859                        "title": "CodeSwarm",
2860                        "version": env!("CARGO_PKG_VERSION"),
2861                    },
2862                }),
2863            )
2864            .await
2865        {
2866            Ok(value) => value,
2867            Err(error) => {
2868                let _ = self.stop().await;
2869                return Err(error);
2870            }
2871        };
2872        let agent_capabilities = initialize
2873            .get("agentCapabilities")
2874            .cloned()
2875            .unwrap_or(Value::Null);
2876        self.capabilities = AgentCapabilities {
2877            supports_cancel: true,
2878            supports_modes: true,
2879            supports_permissions: true,
2880            supports_terminals: true,
2881            supports_session_load: agent_capabilities
2882                .get("loadSession")
2883                .and_then(Value::as_bool)
2884                .unwrap_or(false),
2885            supports_models: false,
2886        };
2887        let session = if let Some(session_id) = self.session_id.clone() {
2888            if !self.capabilities.supports_session_load {
2889                let _ = self.stop().await;
2890                return Err(AdapterError::Unsupported("session/load"));
2891            }
2892            match self
2893                .request(
2894                    "session/load",
2895                    serde_json::json!({
2896                        "cwd": self.cwd,
2897                        "mcpServers": [],
2898                        "sessionId": session_id,
2899                    }),
2900                )
2901                .await
2902            {
2903                Ok(value) => value,
2904                Err(error) => {
2905                    let _ = self.stop().await;
2906                    return Err(error);
2907                }
2908            }
2909        } else {
2910            let session = match self
2911                .request(
2912                    "session/new",
2913                    serde_json::json!({"cwd": self.cwd, "mcpServers": []}),
2914                )
2915                .await
2916            {
2917                Ok(value) => value,
2918                Err(error) => {
2919                    let _ = self.stop().await;
2920                    return Err(error);
2921                }
2922            };
2923            self.session_id = session
2924                .get("sessionId")
2925                .and_then(Value::as_str)
2926                .map(str::to_owned);
2927            if self.session_id.is_none() {
2928                let _ = self.stop().await;
2929                return Err(AdapterError::Protocol(
2930                    "session/new returned no sessionId".into(),
2931                ));
2932            }
2933            session
2934        };
2935        self.capabilities.supports_modes = false;
2936        if let Some(modes) = session.get("modes") {
2937            let available = modes
2938                .get("availableModes")
2939                .and_then(Value::as_array)
2940                .map(|modes| {
2941                    modes
2942                        .iter()
2943                        .filter_map(|mode| {
2944                            Some(Mode {
2945                                id: mode.get("id")?.as_str()?.to_owned(),
2946                                label: mode.get("name")?.as_str()?.to_owned(),
2947                            })
2948                        })
2949                        .collect::<Vec<_>>()
2950                })
2951                .unwrap_or_default();
2952            self.modes = available.clone();
2953            self.capabilities.supports_modes = !available.is_empty();
2954            self.queued_events.push_back(Ok(AgentEvent::ModesReplaced {
2955                slot: self.slot,
2956                modes: available,
2957                current_mode: modes
2958                    .get("currentModeId")
2959                    .and_then(Value::as_str)
2960                    .map(str::to_owned),
2961            }));
2962        }
2963        self.models.clear();
2964        self.model_config_id = None;
2965        let current_model =
2966            parse_model_config(&session).and_then(|(config_id, models, current)| {
2967                self.model_config_id = Some(config_id);
2968                self.models = models;
2969                current
2970            });
2971        self.capabilities.supports_models =
2972            self.model_config_id.is_some() && !self.models.is_empty();
2973        if let Some(config_id) = self.model_config_id.clone()
2974            && !self.models.is_empty()
2975        {
2976            self.queued_events.push_back(Ok(AgentEvent::ModelsReplaced {
2977                slot: self.slot,
2978                config_id,
2979                models: self.models.clone(),
2980                current_model,
2981            }));
2982        }
2983        self.queued_events.push_back(Ok(AgentEvent::Ready {
2984            slot: self.slot,
2985            capabilities: self.capabilities(),
2986        }));
2987        Ok(())
2988    }
2989}
2990
2991fn prompt_resource_paths(prompt: &str) -> Vec<String> {
2992    let characters = prompt.chars().collect::<Vec<_>>();
2993    let mut paths = Vec::new();
2994    let mut index = 0;
2995    while index < characters.len() {
2996        if characters[index] != '@' {
2997            index += 1;
2998            continue;
2999        }
3000        index += 1;
3001        let quoted = characters.get(index) == Some(&'"');
3002        if quoted {
3003            index += 1;
3004        }
3005        let start = index;
3006        while index < characters.len()
3007            && if quoted {
3008                characters[index] != '"'
3009            } else {
3010                !characters[index].is_whitespace()
3011            }
3012        {
3013            index += 1;
3014        }
3015        if index > start {
3016            paths.push(characters[start..index].iter().collect());
3017        }
3018        if quoted && index < characters.len() {
3019            index += 1;
3020        }
3021    }
3022    paths
3023}
3024
3025fn prompt_content_blocks(cwd: &Path, prompt: &str) -> Vec<Value> {
3026    let mut blocks = vec![serde_json::json!({"type": "text", "text": prompt})];
3027    for path in prompt_resource_paths(prompt) {
3028        if path.ends_with('/') {
3029            continue;
3030        }
3031        let Ok(resource) = resources::load(cwd, &path) else {
3032            continue;
3033        };
3034        let uri = format!("file://{}", resource.path.display());
3035        let resource_value = if let Some(text) = resource.text {
3036            serde_json::json!({
3037                "uri": uri,
3038                "text": text,
3039                "mimeType": resource.mime_type,
3040            })
3041        } else if let Some(data) = resource.data {
3042            serde_json::json!({
3043                "uri": uri,
3044                "blob": BASE64.encode(data),
3045                "mimeType": resource.mime_type,
3046            })
3047        } else {
3048            continue;
3049        };
3050        blocks.push(serde_json::json!({
3051            "type": "resource",
3052            "resource": resource_value,
3053        }));
3054    }
3055    blocks
3056}
3057
3058#[async_trait]
3059impl AgentAdapter for AcpAdapter {
3060    fn slot(&self) -> RosterSlot {
3061        self.slot
3062    }
3063
3064    fn session_id(&self) -> Option<String> {
3065        self.session_id.clone()
3066    }
3067
3068    fn protocol(&self) -> &'static str {
3069        "acp"
3070    }
3071
3072    fn capabilities(&self) -> AgentCapabilities {
3073        self.capabilities.clone()
3074    }
3075
3076    async fn start(&mut self) -> AdapterResult<()> {
3077        // Use the inherent implementation, which owns the cleanup boundary
3078        // around the multi-step ACP handshake.
3079        AcpAdapter::start(self).await
3080    }
3081
3082    async fn send_prompt(&mut self, prompt: String) -> AdapterResult<()> {
3083        let session_id = self
3084            .session_id
3085            .as_ref()
3086            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3087        let request_id = self.next_request_id;
3088        self.next_request_id += 1;
3089        let prompt_blocks = prompt_content_blocks(&self.cwd, &prompt);
3090        self.write_json(serde_json::json!({
3091            "jsonrpc": "2.0",
3092            "id": request_id,
3093            "method": "session/prompt",
3094            "params": {
3095                "sessionId": session_id,
3096                "prompt": prompt_blocks,
3097            },
3098        }))
3099        .await?;
3100        self.prompt_request_id = Some(request_id);
3101        Ok(())
3102    }
3103
3104    async fn cancel(&mut self) -> AdapterResult<bool> {
3105        let Some(session_id) = &self.session_id else {
3106            return Ok(false);
3107        };
3108        self.write_json(serde_json::json!({
3109            "jsonrpc": "2.0",
3110            "method": "session/cancel",
3111            "params": {"sessionId": session_id, "_meta": {}},
3112        }))
3113        .await?;
3114        let settled = tokio::time::timeout(CANCEL_SETTLE_TIMEOUT, async {
3115            loop {
3116                match <Self as AgentAdapter>::next_event(self).await {
3117                    Some(Ok(AgentEvent::TurnComplete { .. })) | None => break,
3118                    Some(Ok(_)) => {}
3119                    Some(Err(_)) => break,
3120                }
3121            }
3122        })
3123        .await
3124        .is_ok();
3125        if !settled {
3126            // A peer that never acknowledges cancellation cannot safely share
3127            // its stream with the next prompt. Restart the transport while
3128            // preserving a loadable provider session when supported.
3129            self.reload().await?;
3130        }
3131        Ok(true)
3132    }
3133
3134    async fn answer_permission(
3135        &mut self,
3136        request_id: String,
3137        answer: PermissionAnswer,
3138    ) -> AdapterResult<()> {
3139        let id = request_id
3140            .parse::<u64>()
3141            .map(Value::from)
3142            .unwrap_or_else(|_| Value::String(request_id));
3143        let outcome = match answer {
3144            PermissionAnswer::Selected { option_id } => {
3145                serde_json::json!({"outcome": "selected", "optionId": option_id})
3146            }
3147            PermissionAnswer::Cancelled => serde_json::json!({"outcome": "cancelled"}),
3148        };
3149        self.write_json(serde_json::json!({
3150            "jsonrpc": "2.0",
3151            "id": id,
3152            // RequestPermissionResponse wraps the selected/cancelled
3153            // discriminator in its `outcome` field. Keep this nested shape
3154            // compatible with the Python ACP server and ACP schema.
3155            "result": {"outcome": outcome},
3156        }))
3157        .await
3158    }
3159
3160    async fn set_mode(&mut self, mode: String) -> AdapterResult<()> {
3161        let session_id = self
3162            .session_id
3163            .as_ref()
3164            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3165        let policy = match mode.as_str() {
3166            "plan" => "codeswarm:mode:plan",
3167            "default" | "manual" => "codeswarm:mode:manual",
3168            "accept-edits" => "codeswarm:mode:accept-edits",
3169            "full-access" | "auto" | "autopilot" => "codeswarm:mode:full-access",
3170            other => other,
3171        };
3172        let native_mode = crate::policy::resolve(policy, &self.modes)
3173            .map(|mode| mode.id)
3174            .unwrap_or(mode);
3175        let _ = self
3176            .request(
3177                "session/set_mode",
3178                serde_json::json!({"sessionId": session_id, "modeId": native_mode.clone()}),
3179            )
3180            .await?;
3181        self.queued_events.push_back(Ok(AgentEvent::ModeUpdated {
3182            slot: self.slot,
3183            current_mode: native_mode,
3184        }));
3185        Ok(())
3186    }
3187
3188    async fn set_model(&mut self, model: String) -> AdapterResult<()> {
3189        let session_id = self
3190            .session_id
3191            .clone()
3192            .ok_or_else(|| AdapterError::Transport("ACP session is not initialized".into()))?;
3193        let config_id = self
3194            .model_config_id
3195            .clone()
3196            .ok_or(AdapterError::Unsupported("set_model"))?;
3197        if !self.models.iter().any(|candidate| candidate.id == model) {
3198            return Err(AdapterError::Protocol(
3199                "model is not advertised by the agent".into(),
3200            ));
3201        }
3202        let _ = self
3203            .request(
3204                "session/set_config_option",
3205                serde_json::json!({
3206                    "sessionId": session_id,
3207                    "configId": config_id,
3208                    "value": model,
3209                }),
3210            )
3211            .await?;
3212        Ok(())
3213    }
3214
3215    async fn reload(&mut self) -> AdapterResult<()> {
3216        // `stop` tears down the process and clears its transport-owned
3217        // session handle. A reload is different from a final shutdown: ACP
3218        // peers advertising `loadSession` must receive the prior ID so the
3219        // replacement process can resume the same conversation.
3220        let session_id = self
3221            .capabilities
3222            .supports_session_load
3223            .then(|| self.session_id.clone())
3224            .flatten();
3225        self.stop().await?;
3226        self.session_id = session_id.clone();
3227        let result = self.start().await;
3228        if result.is_err() {
3229            // `start` cleans up a partially initialized transport by calling
3230            // `stop`, which also clears the handle. Keep it available for a
3231            // subsequent retry after the coordinator reports the failure.
3232            self.session_id = session_id;
3233        }
3234        result
3235    }
3236
3237    async fn stop(&mut self) -> AdapterResult<()> {
3238        let terminals = std::mem::take(&mut self.terminals);
3239        for terminal in terminals.values() {
3240            terminal.stop().await;
3241        }
3242        if let Some(mut child) = self.child.take() {
3243            terminate_child(&mut child).await?;
3244        }
3245        self.reader = None;
3246        self.session_id = None;
3247        self.prompt_request_id = None;
3248        if let Some(task) = self.stderr_task.take() {
3249            task.abort();
3250            let _ = task.await;
3251        }
3252        Ok(())
3253    }
3254
3255    async fn next_event(&mut self) -> Option<AdapterResult<AgentEvent>> {
3256        if let Some(event) = self.queued_events.pop_front() {
3257            return Some(event);
3258        }
3259        loop {
3260            let line = match self.read_line().await {
3261                Ok(line) => line,
3262                Err(error) => return Some(Err(error)),
3263            };
3264            let value: Value = match serde_json::from_str(&line) {
3265                Ok(value) => value,
3266                Err(_) => continue,
3267            };
3268            match self.reject_empty_permission_request(&value).await {
3269                Ok(true) => continue,
3270                Ok(false) => {}
3271                Err(error) => return Some(Err(error)),
3272            }
3273            match self.handle_client_request(&value).await {
3274                Ok(true) => continue,
3275                Ok(false) => {}
3276                Err(error) => return Some(Err(error)),
3277            }
3278            match parse_acp_notification(self.slot, &line) {
3279                Ok(Some(event)) => {
3280                    if let AgentEvent::ModelsReplaced {
3281                        config_id, models, ..
3282                    } = &event
3283                    {
3284                        self.model_config_id = Some(config_id.clone());
3285                        self.models = models.clone();
3286                        self.capabilities.supports_models = !models.is_empty();
3287                    }
3288                    return Some(Ok(event));
3289                }
3290                Ok(None) => {}
3291                Err(error) => return Some(Err(error)),
3292            }
3293            if value.get("id").is_some_and(|id| {
3294                self.prompt_request_id
3295                    .is_some_and(|expected| rpc_id_to_string(id) == expected.to_string())
3296            }) {
3297                if let Some(error) = value.get("error") {
3298                    self.prompt_request_id = None;
3299                    return Some(Err(AdapterError::Protocol(error.to_string())));
3300                }
3301                self.prompt_request_id = None;
3302                return Some(Ok(AgentEvent::TurnComplete { slot: self.slot }));
3303            }
3304        }
3305    }
3306}
3307
3308fn parse_acp_notification(slot: RosterSlot, line: &str) -> AdapterResult<Option<AgentEvent>> {
3309    let value: Value =
3310        serde_json::from_str(line).map_err(|error| AdapterError::Protocol(error.to_string()))?;
3311    let method = value.get("method").and_then(Value::as_str);
3312    if method == Some("session/request_permission") {
3313        let params = value.get("params").cloned().unwrap_or(Value::Null);
3314        let request_id = value
3315            .get("id")
3316            .map(rpc_id_to_string)
3317            .unwrap_or_else(|| "permission".into());
3318        return Ok(parse_permission_event(
3319            slot,
3320            &params,
3321            &request_id,
3322            params.get("options"),
3323        ));
3324    }
3325    if method != Some("session/update") {
3326        return Ok(None);
3327    }
3328    let Some(update) = value.get("params").and_then(|params| params.get("update")) else {
3329        return Ok(None);
3330    };
3331    let kind = update.get("sessionUpdate").and_then(Value::as_str);
3332    if kind == Some("config_option_update")
3333        && let Some((config_id, models, current_model)) = parse_model_config(update)
3334    {
3335        return Ok(Some(AgentEvent::ModelsReplaced {
3336            slot,
3337            config_id,
3338            models,
3339            current_model,
3340        }));
3341    }
3342    if kind == Some("request_permission") {
3343        let request_id = update
3344            .get("toolCall")
3345            .and_then(|tool| tool.get("toolCallId"))
3346            .and_then(Value::as_str)
3347            .unwrap_or("permission");
3348        return Ok(parse_permission_event(
3349            slot,
3350            update,
3351            request_id,
3352            update.get("options"),
3353        ));
3354    }
3355    if kind == Some("available_commands_update") {
3356        let commands = update
3357            .get("availableCommands")
3358            .and_then(Value::as_array)
3359            .map(|commands| {
3360                commands
3361                    .iter()
3362                    .filter_map(|command| {
3363                        let name = command.get("name").and_then(Value::as_str)?.trim();
3364                        (!name.is_empty()).then(|| AgentCommand {
3365                            name: name.to_owned(),
3366                        })
3367                    })
3368                    .collect::<Vec<_>>()
3369            })
3370            .unwrap_or_default();
3371        return Ok(Some(AgentEvent::CommandsReplaced { slot, commands }));
3372    }
3373    if kind == Some("current_mode_update") {
3374        if let Some(mode) = update
3375            .get("currentModeId")
3376            .and_then(Value::as_str)
3377            .filter(|mode| !mode.trim().is_empty())
3378        {
3379            return Ok(Some(AgentEvent::ModeUpdated {
3380                slot,
3381                current_mode: mode.to_owned(),
3382            }));
3383        }
3384        return Ok(None);
3385    }
3386    if kind == Some("usage_update") {
3387        let Some(used) = update.get("used").and_then(Value::as_u64) else {
3388            return Ok(None);
3389        };
3390        let Some(size) = update.get("size").and_then(Value::as_u64) else {
3391            return Ok(None);
3392        };
3393        return Ok(Some(AgentEvent::UsageUpdated {
3394            slot,
3395            usage: UsageUpdate { used, size },
3396        }));
3397    }
3398    if let Some(terminal) = parse_terminal_event(update, kind) {
3399        return Ok(Some(AgentEvent::Terminal {
3400            slot,
3401            event: terminal,
3402        }));
3403    }
3404    let text = update
3405        .get("content")
3406        .and_then(|content| content.get("text"))
3407        .and_then(Value::as_str)
3408        .map(str::to_owned);
3409    if kind == Some("user_message_chunk") {
3410        return Ok(text
3411            .filter(|text| !text.is_empty())
3412            .map(|text| AgentEvent::UserText { slot, text }));
3413    }
3414    if kind == Some("agent_message_chunk")
3415        && let Some(mode) = text
3416            .as_deref()
3417            .and_then(|text| text.strip_prefix("[MODE_UPDATE]"))
3418            .map(str::trim)
3419            .filter(|mode| !mode.is_empty())
3420    {
3421        // Gemini's native ACP bridge historically encoded a mode change as a
3422        // control marker in the message stream. It is state, not transcript
3423        // content; expose it as a normalized catalog replacement instead of
3424        // leaking the marker into the conversation.
3425        return Ok(Some(AgentEvent::ModesReplaced {
3426            slot,
3427            modes: vec![Mode {
3428                id: mode.to_owned(),
3429                label: mode.to_owned(),
3430            }],
3431            current_mode: Some(mode.to_owned()),
3432        }));
3433    }
3434    match (kind, text) {
3435        (Some("agent_message_chunk"), Some(text)) if !text.is_empty() => {
3436            Ok(Some(AgentEvent::Text { slot, text }))
3437        }
3438        (Some("agent_thought_chunk"), Some(text)) if !text.is_empty() => {
3439            Ok(Some(AgentEvent::Thought { slot, text }))
3440        }
3441        (Some("tool_call"), _) | (Some("tool_call_update"), _) => {
3442            let id = update
3443                .get("toolCallId")
3444                .and_then(Value::as_str)
3445                .unwrap_or("unknown-tool")
3446                .to_owned();
3447            let title = update
3448                .get("title")
3449                .and_then(Value::as_str)
3450                .unwrap_or("Tool call")
3451                .to_owned();
3452            let status = match update.get("status").and_then(Value::as_str) {
3453                Some("completed") => ToolStatus::Completed,
3454                Some("failed") => ToolStatus::Failed,
3455                Some("in_progress") => ToolStatus::Running,
3456                _ => ToolStatus::Pending,
3457            };
3458            Ok(Some(AgentEvent::Tool {
3459                slot,
3460                update: ToolUpdate {
3461                    id,
3462                    title,
3463                    status,
3464                    detail: None,
3465                },
3466            }))
3467        }
3468        _ => Ok(None),
3469    }
3470}
3471
3472fn parse_model_config(value: &Value) -> Option<(String, Vec<Mode>, Option<String>)> {
3473    let config = value
3474        .get("configOptions")?
3475        .as_array()?
3476        .iter()
3477        .find(|option| {
3478            option.get("category").and_then(Value::as_str) == Some("model")
3479                && matches!(
3480                    option.get("type").and_then(Value::as_str),
3481                    Some("select" | "enum")
3482                )
3483        })?;
3484    let config_id = config.get("id")?.as_str()?.to_owned();
3485    let models = config
3486        .get("options")?
3487        .as_array()?
3488        .iter()
3489        .filter_map(|option| {
3490            let id = option.get("value")?.as_str()?.to_owned();
3491            let label = option
3492                .get("name")
3493                .or_else(|| option.get("label"))
3494                .and_then(Value::as_str)
3495                .unwrap_or(&id)
3496                .to_owned();
3497            Some(Mode { id, label })
3498        })
3499        .collect::<Vec<_>>();
3500    (!models.is_empty()).then(|| {
3501        let current = config
3502            .get("currentValue")
3503            .and_then(Value::as_str)
3504            .map(str::to_owned);
3505        (config_id, models, current)
3506    })
3507}
3508
3509/// Normalize terminal lifecycle updates emitted by ACP-compatible bridges and
3510/// native stream adapters. Protocols have used both snake_case update names
3511/// and a nested `terminal` object, so accept either without leaking that
3512/// shape beyond the adapter boundary.
3513fn parse_terminal_event(value: &Value, kind: Option<&str>) -> Option<TerminalEvent> {
3514    let nested = value.get("terminal").unwrap_or(value);
3515    let kind = kind.or_else(|| value.get("event").and_then(Value::as_str))?;
3516    let id = nested
3517        .get("terminalId")
3518        .or_else(|| nested.get("terminal_id"))
3519        .or_else(|| nested.get("id"))
3520        .and_then(Value::as_str)
3521        .unwrap_or("terminal")
3522        .to_owned();
3523    match kind {
3524        "terminal_created" | "terminal_create" | "terminal_started" => {
3525            let command = nested
3526                .get("command")
3527                .and_then(Value::as_str)
3528                .unwrap_or("")
3529                .to_owned();
3530            Some(TerminalEvent::Created { id, command })
3531        }
3532        "terminal_output" | "terminal_output_chunk" => {
3533            let text = nested
3534                .get("output")
3535                .or_else(|| nested.get("text"))
3536                .and_then(Value::as_str)
3537                .unwrap_or("")
3538                .to_owned();
3539            Some(TerminalEvent::Output { id, text })
3540        }
3541        "terminal_exited" | "terminal_exit" => {
3542            let code = nested
3543                .get("exitCode")
3544                .or_else(|| nested.get("exit_code"))
3545                .or_else(|| nested.get("code"))
3546                .and_then(Value::as_i64)
3547                .unwrap_or(0) as i32;
3548            Some(TerminalEvent::Exited { id, code })
3549        }
3550        "terminal_released" | "terminal_release" => Some(TerminalEvent::Released { id }),
3551        _ => None,
3552    }
3553}
3554
3555fn parse_permission_event(
3556    slot: RosterSlot,
3557    value: &Value,
3558    request_id: &str,
3559    options: Option<&Value>,
3560) -> Option<AgentEvent> {
3561    let tool = value.get("toolCall").unwrap_or(value);
3562    let title = tool
3563        .get("title")
3564        .and_then(Value::as_str)
3565        .unwrap_or("Agent requests permission")
3566        .to_owned();
3567    let (options, option_ids): (Vec<String>, Vec<String>) = options
3568        .and_then(Value::as_array)
3569        .map(|options| {
3570            options
3571                .iter()
3572                .filter_map(|option| {
3573                    let label = option
3574                        .get("name")
3575                        .or_else(|| option.get("optionId"))
3576                        .and_then(Value::as_str)?
3577                        .to_owned();
3578                    let option_id = option
3579                        .get("optionId")
3580                        .or_else(|| option.get("id"))
3581                        .and_then(Value::as_str)
3582                        .map(str::to_owned)
3583                        .unwrap_or_else(|| label.clone());
3584                    Some((label, option_id))
3585                })
3586                .unzip()
3587        })
3588        .unwrap_or_default();
3589    if options.is_empty() {
3590        return None;
3591    }
3592    Some(AgentEvent::Permission {
3593        slot,
3594        request: PermissionRequest {
3595            id: request_id.to_owned(),
3596            title,
3597            options,
3598            option_ids,
3599        },
3600    })
3601}
3602
3603fn rpc_id_to_string(value: &Value) -> String {
3604    value
3605        .as_str()
3606        .map(str::to_owned)
3607        .or_else(|| value.as_u64().map(|id| id.to_string()))
3608        .unwrap_or_else(|| value.to_string())
3609}
3610
3611#[cfg(test)]
3612mod tests {
3613    use super::{
3614        AcpAdapter, AdapterHost, AgentAdapter, AgyAdapter, MAX_ACP_LINE_BYTES, MAX_FILE_READ_BYTES,
3615        RelayHost, ScriptedAdapter, parse_acp_notification, parse_agy_line, parse_command_line,
3616        parse_model_config, prompt_content_blocks, read_bounded_line,
3617    };
3618    #[cfg(target_os = "linux")]
3619    use super::{isolate_process_group, terminate_child};
3620    use crate::TerminalEvent;
3621    use crate::{
3622        AgentCapabilities, AgentEvent, EventLog, Mode, PermissionAnswer, ToolStatus,
3623        persistence::SessionMetadataStore,
3624        relay::{DEFAULT_STOP_ACKNOWLEDGMENT, RelayDecision, STOP_TOKEN},
3625    };
3626    use async_trait::async_trait;
3627    use serde_json::Value;
3628    use std::sync::{
3629        Arc, Mutex,
3630        atomic::{AtomicUsize, Ordering},
3631    };
3632
3633    fn unique_test_path(stem: &str, extension: &str) -> std::path::PathBuf {
3634        let nonce = std::time::SystemTime::now()
3635            .duration_since(std::time::UNIX_EPOCH)
3636            .expect("clock")
3637            .as_nanos();
3638        std::env::temp_dir().join(format!("{stem}-{}-{nonce}.{extension}", std::process::id()))
3639    }
3640
3641    #[test]
3642    fn malformed_file_writes_preserve_existing_content() {
3643        let root = unique_test_path("codeswarm-write-validation", "dir");
3644        std::fs::create_dir_all(&root).unwrap();
3645        let file = root.join("keep.txt");
3646        std::fs::write(&file, "valuable content").unwrap();
3647        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
3648        for content in [
3649            Value::Null,
3650            serde_json::json!(false),
3651            serde_json::json!(42),
3652            serde_json::json!([]),
3653        ] {
3654            assert!(
3655                adapter
3656                    .write_workspace_text(
3657                        &serde_json::json!({"path":"keep.txt", "content": content})
3658                    )
3659                    .is_err()
3660            );
3661            assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3662        }
3663        assert!(
3664            adapter
3665                .write_workspace_text(&serde_json::json!({"path":"keep.txt"}))
3666                .is_err()
3667        );
3668        assert!(
3669            adapter
3670                .write_workspace_text(&serde_json::json!({"path": null, "content":"replacement"}))
3671                .is_err()
3672        );
3673        assert_eq!(std::fs::read_to_string(&file).unwrap(), "valuable content");
3674        adapter
3675            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":"replacement"}))
3676            .unwrap();
3677        assert_eq!(std::fs::read_to_string(&file).unwrap(), "replacement");
3678        adapter
3679            .write_workspace_text(&serde_json::json!({"path":"keep.txt", "content":""}))
3680            .unwrap();
3681        assert_eq!(std::fs::read_to_string(&file).unwrap(), "");
3682        std::fs::remove_dir_all(root).unwrap();
3683    }
3684
3685    #[tokio::test]
3686    async fn silent_acp_control_request_times_out_and_transport_can_be_stopped() {
3687        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{}}}'; read _; echo '{"jsonrpc":"2.0","id":"2","result":{"sessionId":"s"}}'; read _; read _"#;
3688        let mut adapter = AcpAdapter::new(
3689            0,
3690            std::env::current_dir().unwrap(),
3691            "sh",
3692            vec!["-c".into(), script.into()],
3693        );
3694        adapter.start().await.unwrap();
3695        let error = adapter
3696            .request_with_timeout(
3697                "session/set_mode",
3698                serde_json::json!({}),
3699                std::time::Duration::from_millis(10),
3700            )
3701            .await
3702            .unwrap_err();
3703        assert!(error.to_string().contains("session/set_mode timed out"));
3704        adapter.stop().await.unwrap();
3705        assert!(adapter.child.is_none());
3706    }
3707
3708    #[tokio::test]
3709    async fn goals_reach_every_roster_slot_without_native_goal_support() {
3710        use crate::goal::GoalCommand;
3711        let hosts = (0..3)
3712            .map(|slot| {
3713                AdapterHost::new(
3714                    Box::new(ScriptedAdapter::new(
3715                        slot,
3716                        AgentCapabilities::default(),
3717                        [
3718                            AgentEvent::TurnComplete { slot },
3719                            AgentEvent::TurnComplete { slot },
3720                        ],
3721                    )),
3722                    None,
3723                )
3724            })
3725            .collect();
3726        let mut relay = RelayHost::new(hosts, 10).unwrap();
3727        relay.start().await.unwrap();
3728        let task = relay
3729            .apply_goal(GoalCommand::Set("Ship the settings screen".into()))
3730            .unwrap()
3731            .unwrap();
3732        relay.relay_mut().enqueue_human(task, Some(0));
3733        for slot in 0..3 {
3734            relay.run_turn("", 0).await.unwrap();
3735            let (actual, prompt) = relay.dispatches().last().unwrap();
3736            assert_eq!(*actual, slot);
3737            assert!(prompt.contains("Active shared goal: Ship the settings screen"));
3738        }
3739        let snapshot = relay.session_metadata();
3740        let restored = crate::goal::Goal::from_metadata(snapshot.get("goal").unwrap());
3741        assert!(restored.is_some());
3742        relay.restore_goal(restored);
3743        relay.reload(0).await.unwrap();
3744        relay.run_turn("", 0).await.unwrap();
3745        assert!(
3746            relay
3747                .dispatches()
3748                .last()
3749                .unwrap()
3750                .1
3751                .contains("Active shared goal: Ship the settings screen")
3752        );
3753        relay.apply_goal(GoalCommand::Done).unwrap();
3754        relay.run_turn("", 0).await.unwrap();
3755        assert!(
3756            relay
3757                .dispatches()
3758                .last()
3759                .unwrap()
3760                .1
3761                .contains("No active shared goal")
3762        );
3763        relay.apply_goal(GoalCommand::Clear).unwrap();
3764        assert!(relay.session_metadata().get("goal").unwrap().is_null());
3765    }
3766
3767    #[tokio::test]
3768    async fn replacement_agent_receives_task_after_public_journal_pruning() {
3769        let hosts = (0..2)
3770            .map(|slot| {
3771                AdapterHost::new(
3772                    Box::new(ScriptedAdapter::new(
3773                        slot,
3774                        AgentCapabilities::default(),
3775                        [
3776                            AgentEvent::Text {
3777                                slot,
3778                                text: "progress".into(),
3779                            },
3780                            AgentEvent::TurnComplete { slot },
3781                            AgentEvent::Text {
3782                                slot,
3783                                text: "more progress".into(),
3784                            },
3785                            AgentEvent::TurnComplete { slot },
3786                        ],
3787                    )),
3788                    None,
3789                )
3790            })
3791            .collect();
3792        let mut relay = RelayHost::new(hosts, 10).unwrap();
3793        relay.start().await.unwrap();
3794        relay
3795            .relay_mut()
3796            .enqueue_human("Fix the login bug", Some(0));
3797        relay.run_turn("", 0).await.unwrap();
3798        relay.run_turn("", 0).await.unwrap();
3799        relay.run_turn("", 0).await.unwrap();
3800        relay.reload(1).await.unwrap();
3801        assert!(
3802            !relay
3803                .relay_mut()
3804                .unseen_context(1)
3805                .contains("Fix the login bug")
3806        );
3807        relay.run_turn("", 0).await.unwrap();
3808        assert!(
3809            relay
3810                .dispatches()
3811                .last()
3812                .unwrap()
3813                .1
3814                .contains("Shared task:\nFix the login bug")
3815        );
3816    }
3817
3818    #[cfg(target_os = "linux")]
3819    #[tokio::test]
3820    async fn termination_kills_only_the_verified_isolated_child_group() {
3821        use nix::unistd::{Pid, getpgid, getpgrp};
3822        use tokio::io::{AsyncBufReadExt, BufReader};
3823
3824        let own_group = getpgrp();
3825        let mut command = tokio::process::Command::new("sh");
3826        isolate_process_group(&mut command);
3827        command
3828            .arg("-c")
3829            .arg("sleep 60 & echo $!; wait")
3830            .stdout(std::process::Stdio::piped());
3831        let mut child = command.spawn().expect("spawn isolated shell");
3832        let leader = Pid::from_raw(child.id().expect("leader pid") as i32);
3833        assert_eq!(getpgid(Some(leader)).expect("leader group"), leader);
3834        assert_ne!(leader, own_group);
3835
3836        let stdout = child.stdout.take().expect("child stdout");
3837        let mut lines = BufReader::new(stdout).lines();
3838        let descendant = lines
3839            .next_line()
3840            .await
3841            .expect("read descendant pid")
3842            .expect("descendant pid")
3843            .parse::<i32>()
3844            .expect("numeric descendant pid");
3845        let descendant = Pid::from_raw(descendant);
3846        assert_eq!(getpgid(Some(descendant)).expect("descendant group"), leader);
3847
3848        terminate_child(&mut child).await.expect("terminate group");
3849        for _ in 0..100 {
3850            if !std::path::Path::new(&format!("/proc/{descendant}")).exists() {
3851                return;
3852            }
3853            tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3854        }
3855        panic!("descendant {descendant} survived isolated group termination");
3856    }
3857
3858    #[test]
3859    fn parses_configured_commands_with_shell_style_quotes_without_a_shell() {
3860        assert_eq!(
3861            parse_command_line(r#"npx -y "@agentclientprotocol/codex-acp" --flag 'two words'"#),
3862            Ok((
3863                "npx".into(),
3864                vec![
3865                    "-y".into(),
3866                    "@agentclientprotocol/codex-acp".into(),
3867                    "--flag".into(),
3868                    "two words".into(),
3869                ]
3870            ),)
3871        );
3872        assert_eq!(
3873            parse_command_line(r#"agent "" escaped\ argument"#),
3874            Ok(("agent".into(), vec!["".into(), "escaped argument".into()],))
3875        );
3876    }
3877
3878    #[test]
3879    fn acp_prompt_expands_safe_at_path_resources() {
3880        let root = unique_test_path("codeswarm-prompt-resource", "dir");
3881        std::fs::create_dir_all(&root).expect("workspace");
3882        std::fs::write(root.join("note.md"), "resource text").expect("resource");
3883        let blocks = prompt_content_blocks(&root, "inspect @note.md");
3884        assert_eq!(blocks[0]["type"], "text");
3885        assert_eq!(blocks[0]["text"], "inspect @note.md");
3886        assert_eq!(blocks[1]["type"], "resource");
3887        assert_eq!(blocks[1]["resource"]["text"], "resource text");
3888        assert_eq!(blocks[1]["resource"]["mimeType"], "text/markdown");
3889        std::fs::remove_dir_all(root).expect("cleanup workspace");
3890    }
3891
3892    #[tokio::test]
3893    async fn oversized_acp_frames_are_rejected_before_full_line_allocation() {
3894        let mut bytes = vec![b'x'; MAX_ACP_LINE_BYTES + 1];
3895        bytes.push(b'\n');
3896        let mut reader = tokio::io::BufReader::new(bytes.as_slice());
3897        assert!(matches!(
3898            read_bounded_line(&mut reader).await,
3899            Err(super::AdapterError::Protocol(detail)) if detail.contains("exceeds")
3900        ));
3901    }
3902
3903    #[test]
3904    fn rejects_malformed_configured_commands_before_spawn() {
3905        assert_eq!(
3906            parse_command_line("agent 'unfinished"),
3907            Err(super::CommandParseError::UnterminatedQuote)
3908        );
3909        assert_eq!(
3910            parse_command_line("agent\\"),
3911            Err(super::CommandParseError::TrailingEscape)
3912        );
3913        assert_eq!(
3914            parse_command_line("   \t"),
3915            Err(super::CommandParseError::Empty)
3916        );
3917    }
3918
3919    #[derive(Debug)]
3920    struct PendingAdapter {
3921        slot: usize,
3922        hang_on_cancel: bool,
3923    }
3924
3925    #[derive(Debug)]
3926    struct ConcurrentStartAdapter {
3927        slot: usize,
3928        barrier: Arc<tokio::sync::Barrier>,
3929    }
3930
3931    #[async_trait]
3932    impl AgentAdapter for ConcurrentStartAdapter {
3933        fn slot(&self) -> usize {
3934            self.slot
3935        }
3936
3937        fn capabilities(&self) -> AgentCapabilities {
3938            AgentCapabilities::default()
3939        }
3940
3941        async fn start(&mut self) -> super::AdapterResult<()> {
3942            self.barrier.wait().await;
3943            Ok(())
3944        }
3945
3946        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
3947            Ok(())
3948        }
3949
3950        async fn cancel(&mut self) -> super::AdapterResult<bool> {
3951            Ok(true)
3952        }
3953
3954        async fn answer_permission(
3955            &mut self,
3956            _request_id: String,
3957            _answer: PermissionAnswer,
3958        ) -> super::AdapterResult<()> {
3959            Ok(())
3960        }
3961
3962        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
3963            Ok(())
3964        }
3965
3966        async fn reload(&mut self) -> super::AdapterResult<()> {
3967            Ok(())
3968        }
3969
3970        async fn stop(&mut self) -> super::AdapterResult<()> {
3971            Ok(())
3972        }
3973
3974        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
3975            std::future::pending().await
3976        }
3977    }
3978
3979    #[derive(Debug)]
3980    struct PermissionBlockingAdapter {
3981        slot: usize,
3982        phase: u8,
3983    }
3984
3985    #[async_trait]
3986    impl AgentAdapter for PermissionBlockingAdapter {
3987        fn slot(&self) -> usize {
3988            self.slot
3989        }
3990
3991        fn capabilities(&self) -> AgentCapabilities {
3992            AgentCapabilities {
3993                supports_permissions: true,
3994                ..AgentCapabilities::default()
3995            }
3996        }
3997
3998        async fn start(&mut self) -> super::AdapterResult<()> {
3999            Ok(())
4000        }
4001
4002        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4003            Ok(())
4004        }
4005
4006        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4007            Ok(true)
4008        }
4009
4010        async fn answer_permission(
4011            &mut self,
4012            request_id: String,
4013            answer: PermissionAnswer,
4014        ) -> super::AdapterResult<()> {
4015            if self.phase != 1 || request_id != "permission-1" {
4016                return Err(super::AdapterError::Protocol(
4017                    "unexpected permission response".into(),
4018                ));
4019            }
4020            assert!(matches!(answer, PermissionAnswer::Selected { .. }));
4021            self.phase = 2;
4022            Ok(())
4023        }
4024
4025        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4026            Ok(())
4027        }
4028
4029        async fn reload(&mut self) -> super::AdapterResult<()> {
4030            Ok(())
4031        }
4032
4033        async fn stop(&mut self) -> super::AdapterResult<()> {
4034            Ok(())
4035        }
4036
4037        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4038            match self.phase {
4039                0 => {
4040                    self.phase = 1;
4041                    Some(Ok(AgentEvent::Permission {
4042                        slot: self.slot,
4043                        request: crate::PermissionRequest {
4044                            id: "permission-1".into(),
4045                            title: "Allow?".into(),
4046                            options: vec!["Allow".into()],
4047                            option_ids: vec!["allow".into()],
4048                        },
4049                    }))
4050                }
4051                1 => std::future::pending().await,
4052                _ => Some(Ok(AgentEvent::TurnComplete { slot: self.slot })),
4053            }
4054        }
4055    }
4056
4057    #[async_trait]
4058    impl AgentAdapter for PendingAdapter {
4059        fn slot(&self) -> usize {
4060            self.slot
4061        }
4062
4063        fn capabilities(&self) -> AgentCapabilities {
4064            AgentCapabilities {
4065                supports_cancel: true,
4066                ..AgentCapabilities::default()
4067            }
4068        }
4069
4070        async fn start(&mut self) -> super::AdapterResult<()> {
4071            Ok(())
4072        }
4073
4074        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4075            Ok(())
4076        }
4077
4078        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4079            if self.hang_on_cancel {
4080                return std::future::pending().await;
4081            }
4082            Ok(true)
4083        }
4084
4085        async fn answer_permission(
4086            &mut self,
4087            _request_id: String,
4088            _answer: PermissionAnswer,
4089        ) -> super::AdapterResult<()> {
4090            Ok(())
4091        }
4092
4093        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4094            Ok(())
4095        }
4096
4097        async fn reload(&mut self) -> super::AdapterResult<()> {
4098            Ok(())
4099        }
4100
4101        async fn stop(&mut self) -> super::AdapterResult<()> {
4102            Ok(())
4103        }
4104
4105        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4106            std::future::pending().await
4107        }
4108    }
4109
4110    #[derive(Debug)]
4111    struct StopTrackingAdapter {
4112        slot: usize,
4113        stopped: Arc<AtomicUsize>,
4114        fail_stop: bool,
4115    }
4116
4117    #[derive(Debug)]
4118    struct ModeOrderAdapter {
4119        slot: usize,
4120        log: Arc<Mutex<Vec<String>>>,
4121        phase: u8,
4122    }
4123
4124    #[derive(Debug)]
4125    struct StartupAcpAdapter {
4126        slot: usize,
4127        events: std::collections::VecDeque<AgentEvent>,
4128    }
4129
4130    impl StartupAcpAdapter {
4131        fn new(slot: usize) -> Self {
4132            Self {
4133                slot,
4134                events: [
4135                    AgentEvent::ModesReplaced {
4136                        slot,
4137                        modes: vec![Mode {
4138                            id: "full-access".into(),
4139                            label: "Auto pilot".into(),
4140                        }],
4141                        current_mode: Some("full-access".into()),
4142                    },
4143                    AgentEvent::Ready {
4144                        slot,
4145                        capabilities: AgentCapabilities {
4146                            supports_modes: true,
4147                            ..AgentCapabilities::default()
4148                        },
4149                    },
4150                ]
4151                .into(),
4152            }
4153        }
4154    }
4155
4156    #[async_trait]
4157    impl AgentAdapter for StartupAcpAdapter {
4158        fn slot(&self) -> usize {
4159            self.slot
4160        }
4161
4162        fn protocol(&self) -> &'static str {
4163            "acp"
4164        }
4165
4166        fn capabilities(&self) -> AgentCapabilities {
4167            AgentCapabilities {
4168                supports_modes: true,
4169                ..AgentCapabilities::default()
4170            }
4171        }
4172
4173        async fn start(&mut self) -> super::AdapterResult<()> {
4174            Ok(())
4175        }
4176
4177        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4178            Ok(())
4179        }
4180
4181        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4182            Ok(true)
4183        }
4184
4185        async fn answer_permission(
4186            &mut self,
4187            _request_id: String,
4188            _answer: PermissionAnswer,
4189        ) -> super::AdapterResult<()> {
4190            Ok(())
4191        }
4192
4193        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4194            Ok(())
4195        }
4196
4197        async fn reload(&mut self) -> super::AdapterResult<()> {
4198            self.events = Self::new(self.slot).events;
4199            Ok(())
4200        }
4201
4202        async fn stop(&mut self) -> super::AdapterResult<()> {
4203            Ok(())
4204        }
4205
4206        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4207            self.events.pop_front().map(Ok)
4208        }
4209    }
4210
4211    #[async_trait]
4212    impl AgentAdapter for ModeOrderAdapter {
4213        fn slot(&self) -> usize {
4214            self.slot
4215        }
4216
4217        fn capabilities(&self) -> AgentCapabilities {
4218            AgentCapabilities {
4219                supports_modes: true,
4220                ..AgentCapabilities::default()
4221            }
4222        }
4223
4224        async fn start(&mut self) -> super::AdapterResult<()> {
4225            self.log.lock().expect("log").push("start".into());
4226            Ok(())
4227        }
4228
4229        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4230            self.log.lock().expect("log").push("prompt".into());
4231            Ok(())
4232        }
4233
4234        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4235            Ok(true)
4236        }
4237
4238        async fn answer_permission(
4239            &mut self,
4240            _request_id: String,
4241            _answer: PermissionAnswer,
4242        ) -> super::AdapterResult<()> {
4243            Ok(())
4244        }
4245
4246        async fn set_mode(&mut self, mode: String) -> super::AdapterResult<()> {
4247            self.log.lock().expect("log").push(format!("mode:{mode}"));
4248            Ok(())
4249        }
4250
4251        async fn reload(&mut self) -> super::AdapterResult<()> {
4252            self.log.lock().expect("log").push("reload".into());
4253            self.phase = 0;
4254            Ok(())
4255        }
4256
4257        async fn stop(&mut self) -> super::AdapterResult<()> {
4258            Ok(())
4259        }
4260
4261        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4262            match self.phase {
4263                0 => {
4264                    self.phase = 1;
4265                    Some(Ok(AgentEvent::ModesReplaced {
4266                        slot: self.slot,
4267                        modes: vec![Mode {
4268                            id: "yolo".into(),
4269                            label: "YOLO".into(),
4270                        }],
4271                        current_mode: None,
4272                    }))
4273                }
4274                1 => {
4275                    self.phase = 2;
4276                    Some(Ok(AgentEvent::TurnComplete { slot: self.slot }))
4277                }
4278                _ => std::future::pending().await,
4279            }
4280        }
4281    }
4282
4283    #[async_trait]
4284    impl AgentAdapter for StopTrackingAdapter {
4285        fn slot(&self) -> usize {
4286            self.slot
4287        }
4288
4289        fn capabilities(&self) -> AgentCapabilities {
4290            AgentCapabilities::default()
4291        }
4292
4293        async fn start(&mut self) -> super::AdapterResult<()> {
4294            Ok(())
4295        }
4296
4297        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4298            Ok(())
4299        }
4300
4301        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4302            Ok(false)
4303        }
4304
4305        async fn answer_permission(
4306            &mut self,
4307            _request_id: String,
4308            _answer: PermissionAnswer,
4309        ) -> super::AdapterResult<()> {
4310            Ok(())
4311        }
4312
4313        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4314            Ok(())
4315        }
4316
4317        async fn reload(&mut self) -> super::AdapterResult<()> {
4318            Ok(())
4319        }
4320
4321        async fn stop(&mut self) -> super::AdapterResult<()> {
4322            self.stopped.fetch_add(1, Ordering::Relaxed);
4323            if self.fail_stop {
4324                Err(super::AdapterError::Transport("stop failed".into()))
4325            } else {
4326                Ok(())
4327            }
4328        }
4329
4330        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4331            None
4332        }
4333    }
4334
4335    /// Startup can fail after an adapter has allocated resources.  Keep a
4336    /// fixture that records whether the failing adapter itself receives the
4337    /// cleanup call, not just the already-started peers.
4338    #[derive(Debug)]
4339    struct FailingStartAdapter {
4340        slot: usize,
4341        stopped: Arc<AtomicUsize>,
4342    }
4343
4344    #[async_trait]
4345    impl AgentAdapter for FailingStartAdapter {
4346        fn slot(&self) -> usize {
4347            self.slot
4348        }
4349
4350        fn capabilities(&self) -> AgentCapabilities {
4351            AgentCapabilities::default()
4352        }
4353
4354        async fn start(&mut self) -> super::AdapterResult<()> {
4355            Err(super::AdapterError::Spawn("startup failed".into()))
4356        }
4357
4358        async fn send_prompt(&mut self, _prompt: String) -> super::AdapterResult<()> {
4359            Ok(())
4360        }
4361
4362        async fn cancel(&mut self) -> super::AdapterResult<bool> {
4363            Ok(false)
4364        }
4365
4366        async fn answer_permission(
4367            &mut self,
4368            _request_id: String,
4369            _answer: PermissionAnswer,
4370        ) -> super::AdapterResult<()> {
4371            Ok(())
4372        }
4373
4374        async fn set_mode(&mut self, _mode: String) -> super::AdapterResult<()> {
4375            Ok(())
4376        }
4377
4378        async fn reload(&mut self) -> super::AdapterResult<()> {
4379            Ok(())
4380        }
4381
4382        async fn stop(&mut self) -> super::AdapterResult<()> {
4383            self.stopped.fetch_add(1, Ordering::Relaxed);
4384            Ok(())
4385        }
4386
4387        async fn next_event(&mut self) -> Option<super::AdapterResult<AgentEvent>> {
4388            None
4389        }
4390    }
4391
4392    #[tokio::test]
4393    async fn relay_stop_attempts_every_adapter_after_one_shutdown_failure() {
4394        let stopped = Arc::new(AtomicUsize::new(0));
4395        let relay = RelayHost::new(
4396            vec![
4397                AdapterHost::new(
4398                    Box::new(StopTrackingAdapter {
4399                        slot: 0,
4400                        stopped: Arc::clone(&stopped),
4401                        fail_stop: true,
4402                    }),
4403                    None,
4404                ),
4405                AdapterHost::new(
4406                    Box::new(StopTrackingAdapter {
4407                        slot: 1,
4408                        stopped: Arc::clone(&stopped),
4409                        fail_stop: false,
4410                    }),
4411                    None,
4412                ),
4413            ],
4414            4,
4415        )
4416        .expect("relay");
4417        let mut relay = relay;
4418
4419        let error = relay.stop().await.expect_err("first stop failure");
4420        assert!(error.to_string().contains("stop failed"));
4421        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4422    }
4423
4424    #[tokio::test]
4425    async fn relay_start_cleans_up_the_adapter_that_failed_startup() {
4426        let stopped = Arc::new(AtomicUsize::new(0));
4427        let mut relay = RelayHost::new(
4428            vec![
4429                AdapterHost::new(
4430                    Box::new(StopTrackingAdapter {
4431                        slot: 0,
4432                        stopped: Arc::clone(&stopped),
4433                        fail_stop: false,
4434                    }),
4435                    None,
4436                ),
4437                AdapterHost::new(
4438                    Box::new(FailingStartAdapter {
4439                        slot: 1,
4440                        stopped: Arc::clone(&stopped),
4441                    }),
4442                    None,
4443                ),
4444            ],
4445            4,
4446        )
4447        .expect("relay");
4448
4449        assert!(relay.start().await.is_err());
4450        assert_eq!(stopped.load(Ordering::Relaxed), 2);
4451    }
4452
4453    #[test]
4454    fn parses_acp_text_without_ui_dependency() {
4455        let event = parse_acp_notification(
4456            2,
4457            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hello"}}}}"#,
4458        )
4459        .expect("valid ACP")
4460        .expect("text event");
4461        assert_eq!(
4462            event,
4463            AgentEvent::Text {
4464                slot: 2,
4465                text: "hello".into(),
4466            }
4467        );
4468    }
4469
4470    #[test]
4471    fn parses_acp_state_notifications_at_the_adapter_boundary() {
4472        let commands = parse_acp_notification(
4473            3,
4474            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"available_commands_update","availableCommands":[{"name":"review","description":"Review"},{"name":"","description":"bad"},{"name":7}]}}}"#,
4475        )
4476        .expect("valid ACP")
4477        .expect("commands event");
4478        assert_eq!(
4479            commands,
4480            AgentEvent::CommandsReplaced {
4481                slot: 3,
4482                commands: vec![crate::AgentCommand {
4483                    name: "review".into()
4484                }]
4485            }
4486        );
4487
4488        let mode = parse_acp_notification(
4489            3,
4490            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"current_mode_update","currentModeId":"review"}}}"#,
4491        )
4492        .expect("valid ACP")
4493        .expect("mode event");
4494        assert_eq!(
4495            mode,
4496            AgentEvent::ModeUpdated {
4497                slot: 3,
4498                current_mode: "review".into()
4499            }
4500        );
4501
4502        let usage = parse_acp_notification(
4503            3,
4504            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"usage_update","used":4200,"size":128000}}}"#,
4505        )
4506        .expect("valid ACP")
4507        .expect("usage event");
4508        assert_eq!(
4509            usage,
4510            AgentEvent::UsageUpdated {
4511                slot: 3,
4512                usage: crate::UsageUpdate {
4513                    used: 4200,
4514                    size: 128000
4515                }
4516            }
4517        );
4518
4519        let models = parse_acp_notification(
4520            3,
4521            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"}]}]}}}"#,
4522        )
4523        .expect("valid ACP")
4524        .expect("models event");
4525        assert!(matches!(
4526            models,
4527            AgentEvent::ModelsReplaced { slot: 3, models, current_model, .. }
4528                if models.len() == 2 && current_model.as_deref() == Some("smart")
4529        ));
4530        assert_eq!(
4531            parse_model_config(&serde_json::json!({
4532                "configOptions": [{"id": "model", "category": "model", "type": "select", "options": [{"name": "missing value"}]}]
4533            })),
4534            None
4535        );
4536
4537        let user = parse_acp_notification(
4538            3,
4539            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"context"}}}}"#,
4540        )
4541        .expect("valid ACP")
4542        .expect("user event");
4543        assert_eq!(
4544            user,
4545            AgentEvent::UserText {
4546                slot: 3,
4547                text: "context".into()
4548            }
4549        );
4550    }
4551
4552    #[test]
4553    fn parses_legacy_gemini_mode_marker_as_state_not_agent_text() {
4554        let event = parse_acp_notification(
4555            0,
4556            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"[MODE_UPDATE] yolo"}}}}"#,
4557        )
4558        .expect("valid ACP")
4559        .expect("mode event");
4560        assert!(matches!(
4561            event,
4562            AgentEvent::ModesReplaced { current_mode: Some(mode), modes, .. }
4563                if mode == "yolo" && modes[0].id == "yolo"
4564        ));
4565    }
4566
4567    #[test]
4568    fn parses_native_agy_text_without_acp_bridge() {
4569        let event = parse_agy_line(
4570            1,
4571            r#"{"event":"step_update","step_update":{"step_type":"agent_response","text_delta":"hello"}}"#,
4572        )
4573        .expect("valid stream-json")
4574        .expect("text event");
4575        assert_eq!(
4576            event,
4577            AgentEvent::Text {
4578                slot: 1,
4579                text: "hello".into(),
4580            }
4581        );
4582    }
4583
4584    #[test]
4585    fn parses_tool_lifecycle_from_each_protocol() {
4586        let agy = parse_agy_line(
4587            1,
4588            r#"{"event":"step_update","step_update":{"step_type":"tool","step_index":4,"tool_name":"run_command","state":"DONE","tool_info":{"output":"ok"}}}"#,
4589        )
4590        .expect("valid native tool")
4591        .expect("tool event");
4592        assert!(matches!(
4593            agy,
4594            AgentEvent::Tool {
4595                update: crate::ToolUpdate {
4596                    status: ToolStatus::Completed,
4597                    ..
4598                },
4599                ..
4600            }
4601        ));
4602
4603        let acp = parse_acp_notification(
4604            1,
4605            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"tool_call_update","toolCallId":"t1","title":"Run tests","status":"failed"}}}"#,
4606        )
4607        .expect("valid ACP tool")
4608        .expect("tool event");
4609        assert!(matches!(
4610            acp,
4611            AgentEvent::Tool {
4612                update: crate::ToolUpdate {
4613                    status: ToolStatus::Failed,
4614                    ..
4615                },
4616                ..
4617            }
4618        ));
4619    }
4620
4621    #[test]
4622    fn parses_terminal_lifecycle_from_acp_and_native_events() {
4623        let created = parse_acp_notification(
4624            0,
4625            r#"{"method":"session/update","params":{"update":{"sessionUpdate":"terminal_created","terminalId":"term-1","command":"cargo test"}}}"#,
4626        )
4627        .expect("valid ACP terminal")
4628        .expect("terminal event");
4629        assert_eq!(
4630            created,
4631            AgentEvent::Terminal {
4632                slot: 0,
4633                event: TerminalEvent::Created {
4634                    id: "term-1".into(),
4635                    command: "cargo test".into(),
4636                },
4637            }
4638        );
4639        let output = parse_agy_line(
4640            1,
4641            r#"{"event":"terminal_output","terminalId":"term-1","output":"ok\n"}"#,
4642        )
4643        .expect("valid native terminal")
4644        .expect("terminal event");
4645        assert_eq!(
4646            output,
4647            AgentEvent::Terminal {
4648                slot: 1,
4649                event: TerminalEvent::Output {
4650                    id: "term-1".into(),
4651                    text: "ok\n".into(),
4652                },
4653            }
4654        );
4655        let released = parse_agy_line(1, r#"{"event":"terminal_released","terminalId":"term-1"}"#)
4656            .expect("valid native release")
4657            .expect("terminal event");
4658        assert!(matches!(
4659            released,
4660            AgentEvent::Terminal {
4661                event: TerminalEvent::Released { id },
4662                ..
4663            } if id == "term-1"
4664        ));
4665    }
4666
4667    #[test]
4668    fn parses_acp_permission_requests() {
4669        let event = parse_acp_notification(
4670            0,
4671            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"}]}}}"#,
4672        )
4673        .expect("valid permission")
4674        .expect("permission event");
4675        assert!(matches!(
4676            event,
4677            AgentEvent::Permission { request, .. }
4678                if request.id == "t1"
4679                    && request.title == "Write file"
4680                    && request.options == ["Allow once", "Reject"]
4681                    && request.option_ids == ["allow-once", "reject"]
4682        ));
4683    }
4684
4685    #[test]
4686    fn parses_acp_permission_request_as_json_rpc_request() {
4687        let event = parse_acp_notification(
4688            2,
4689            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"sessionId":"s1","toolCall":{"title":"Write file"},"options":[{"optionId":"allow-once"},{"name":"reject"}]}}"#,
4690        )
4691        .expect("valid permission request")
4692        .expect("permission event");
4693        assert!(matches!(
4694            event,
4695            AgentEvent::Permission { request, .. }
4696                if request.id == "17"
4697                    && request.title == "Write file"
4698                    && request.options == ["allow-once", "reject"]
4699                    && request.option_ids == ["allow-once", "reject"]
4700        ));
4701    }
4702
4703    #[tokio::test]
4704    async fn native_adapter_explicitly_rejects_permission_answers() {
4705        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4706        assert_eq!(
4707            adapter
4708                .answer_permission(
4709                    "request".into(),
4710                    PermissionAnswer::Selected {
4711                        option_id: "allow".into()
4712                    },
4713                )
4714                .await,
4715            Err(super::AdapterError::Unsupported("permission answer"))
4716        );
4717    }
4718
4719    #[tokio::test]
4720    async fn native_mode_policy_aliases_resolve_to_its_supported_id() {
4721        let mut adapter = AgyAdapter::new(0, std::env::current_dir().expect("cwd"), "agy");
4722        adapter
4723            .set_mode("full-access".into())
4724            .await
4725            .expect("auto-pilot alias");
4726        assert!(matches!(
4727            adapter.next_event().await,
4728            Some(Ok(AgentEvent::ModesReplaced { current_mode: Some(mode), .. })) if mode == "agy:full-access"
4729        ));
4730    }
4731
4732    #[tokio::test]
4733    async fn native_stream_persists_announced_conversation_for_follow_up_turns() {
4734        let script_path = unique_test_path("codeswarm-agy-session", "sh");
4735        std::fs::write(
4736            &script_path,
4737            "#!/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",
4738        )
4739        .expect("write native test script");
4740        let mut adapter = AgyAdapter::new(
4741            0,
4742            std::env::current_dir().expect("cwd"),
4743            format!("sh {}", script_path.display()),
4744        );
4745        adapter.start().await.expect("start native adapter");
4746        // Startup emits its mode catalog and readiness before a turn.
4747        assert!(adapter.next_event().await.is_some());
4748        assert!(adapter.next_event().await.is_some());
4749        adapter
4750            .send_prompt("first".into())
4751            .await
4752            .expect("first prompt");
4753        while !matches!(
4754            adapter.next_event().await,
4755            Some(Ok(AgentEvent::TurnComplete { .. }))
4756        ) {}
4757        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4758        adapter
4759            .send_prompt("follow up".into())
4760            .await
4761            .expect("follow-up prompt");
4762        while !matches!(
4763            adapter.next_event().await,
4764            Some(Ok(AgentEvent::TurnComplete { .. }))
4765        ) {}
4766        assert_eq!(adapter.session_id.as_deref(), Some("native-session"));
4767        adapter.stop().await.expect("stop native adapter");
4768        std::fs::remove_file(script_path).expect("cleanup native script");
4769    }
4770
4771    #[tokio::test]
4772    async fn native_stream_reports_unsuccessful_result_as_crash_not_completion() {
4773        let script_path = unique_test_path("codeswarm-agy-failure", "sh");
4774        std::fs::write(
4775            &script_path,
4776            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"result\",\"result\":{\"status\":\"FAILURE\",\"error\":\"agent failed\"}}'\n",
4777        )
4778        .expect("write native test script");
4779        let mut adapter = AgyAdapter::new(
4780            0,
4781            std::env::current_dir().expect("cwd"),
4782            format!("sh {}", script_path.display()),
4783        );
4784        adapter.start().await.expect("start native adapter");
4785        assert!(adapter.next_event().await.is_some());
4786        assert!(adapter.next_event().await.is_some());
4787        adapter.send_prompt("fail".into()).await.expect("prompt");
4788        assert!(matches!(
4789            adapter.next_event().await,
4790            Some(Ok(AgentEvent::Failed { started: true, detail, .. }))
4791                if detail == "agent failed"
4792        ));
4793        adapter.stop().await.expect("stop native adapter");
4794        std::fs::remove_file(script_path).expect("cleanup native script");
4795    }
4796
4797    #[tokio::test]
4798    async fn acp_adapter_initializes_session_and_completes_a_prompt() {
4799        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"}}'"#;
4800        let cwd = std::env::current_dir().expect("cwd");
4801        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4802        adapter.start().await.expect("initialize");
4803        assert!(matches!(
4804            adapter.next_event().await,
4805            Some(Ok(AgentEvent::ModesReplaced { .. }))
4806        ));
4807        assert!(matches!(
4808            adapter.next_event().await,
4809            Some(Ok(AgentEvent::Ready { .. }))
4810        ));
4811        adapter.send_prompt("hello".into()).await.expect("prompt");
4812        assert!(matches!(
4813            adapter.next_event().await,
4814            Some(Ok(AgentEvent::Text { text, .. })) if text == "hello"
4815        ));
4816        assert!(matches!(
4817            adapter.next_event().await,
4818            Some(Ok(AgentEvent::TurnComplete { .. }))
4819        ));
4820    }
4821
4822    #[tokio::test]
4823    async fn acp_string_prompt_ids_complete_and_allow_a_follow_up_turn() {
4824        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"}}'"#;
4825        let cwd = std::env::current_dir().expect("cwd");
4826        let mut adapter = AcpAdapter::new(1, cwd, "sh", vec!["-c".into(), script.into()]);
4827        adapter.start().await.expect("initialize");
4828        assert!(adapter.next_event().await.is_some());
4829        assert!(adapter.next_event().await.is_some());
4830
4831        for (prompt, expected) in [("first prompt", "first"), ("follow up", "second")] {
4832            adapter.send_prompt(prompt.into()).await.expect("prompt");
4833            assert!(matches!(
4834                adapter.next_event().await,
4835                Some(Ok(AgentEvent::Text { text, .. })) if text == expected
4836            ));
4837            assert!(matches!(
4838                adapter.next_event().await,
4839                Some(Ok(AgentEvent::TurnComplete { slot: 1 }))
4840            ));
4841        }
4842        adapter.stop().await.expect("stop");
4843    }
4844
4845    #[tokio::test]
4846    async fn empty_acp_mode_catalog_disables_mode_control() {
4847        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":[]}}}'"#;
4848        let cwd = std::env::current_dir().expect("cwd");
4849        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4850        adapter.start().await.expect("initialize");
4851        assert!(!adapter.capabilities().supports_modes);
4852        assert!(matches!(
4853            adapter.next_event().await,
4854            Some(Ok(AgentEvent::ModesReplaced { modes, .. })) if modes.is_empty()
4855        ));
4856        assert!(matches!(
4857            adapter.next_event().await,
4858            Some(Ok(AgentEvent::Ready { capabilities, .. })) if !capabilities.supports_modes
4859        ));
4860        adapter.stop().await.expect("stop");
4861    }
4862
4863    #[tokio::test]
4864    async fn acp_models_are_discovered_live_and_changed_through_session_config() {
4865        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"#;
4866        let cwd = std::env::current_dir().expect("cwd");
4867        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4868        adapter.start().await.expect("initialize");
4869        assert!(adapter.capabilities().supports_models);
4870        assert!(matches!(
4871            adapter.next_event().await,
4872            Some(Ok(AgentEvent::ModelsReplaced { config_id, models, current_model, .. }))
4873                if config_id == "model"
4874                    && models == [Mode { id: "fast".into(), label: "Fast".into() }, Mode { id: "smart".into(), label: "Smart".into() }]
4875                    && current_model.as_deref() == Some("fast")
4876        ));
4877        assert!(matches!(
4878            adapter.next_event().await,
4879            Some(Ok(AgentEvent::Ready { capabilities, .. })) if capabilities.supports_models
4880        ));
4881        adapter.set_model("smart".into()).await.expect("set model");
4882        assert!(adapter.set_model("invented".into()).await.is_err());
4883        adapter.stop().await.expect("stop");
4884    }
4885
4886    #[tokio::test]
4887    async fn acp_mode_change_is_acknowledged_without_provider_notification() {
4888        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":{}}'"#;
4889        let cwd = std::env::current_dir().expect("cwd");
4890        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4891        adapter.start().await.expect("initialize");
4892        adapter
4893            .set_mode(crate::policy::DEFAULT_POLICY_ID.into())
4894            .await
4895            .expect("set mode");
4896        assert!(matches!(
4897            adapter.next_event().await,
4898            Some(Ok(AgentEvent::ModesReplaced { .. }))
4899        ));
4900        assert!(matches!(
4901            adapter.next_event().await,
4902            Some(Ok(AgentEvent::Ready { .. }))
4903        ));
4904        assert!(matches!(
4905            adapter.next_event().await,
4906            Some(Ok(AgentEvent::ModeUpdated { current_mode, .. })) if current_mode == "yolo"
4907        ));
4908        adapter.stop().await.expect("stop");
4909    }
4910
4911    #[tokio::test]
4912    async fn acp_reload_preserves_a_loadable_session_id() {
4913        let cwd = std::env::current_dir().expect("cwd");
4914        let mut adapter = AcpAdapter::with_session_id(
4915            0,
4916            cwd,
4917            "__codeswarm_missing_acp_for_reload_test__",
4918            Vec::new(),
4919            "saved-session",
4920        );
4921        adapter.capabilities.supports_session_load = true;
4922        // A failed replacement process still must not erase the session ID:
4923        // the coordinator can report the startup error and offer another
4924        // reload, preserving the only handle that can resume the conversation.
4925        assert!(adapter.reload().await.is_err());
4926        assert_eq!(adapter.session_id.as_deref(), Some("saved-session"));
4927    }
4928
4929    #[tokio::test]
4930    async fn acp_reload_starts_a_fresh_session_when_loading_is_not_supported() {
4931        let cwd = std::env::current_dir().expect("cwd");
4932        let mut adapter = AcpAdapter::with_session_id(
4933            0,
4934            cwd,
4935            "__codeswarm_missing_nonloadable_acp__",
4936            Vec::new(),
4937            "stale-session",
4938        );
4939        adapter.capabilities.supports_session_load = false;
4940        assert!(adapter.reload().await.is_err());
4941        assert_eq!(adapter.session_id, None);
4942    }
4943
4944    #[tokio::test]
4945    async fn acp_stream_ignores_diagnostic_junk_and_surfaces_prompt_errors() {
4946        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"}}'"#;
4947        let cwd = std::env::current_dir().expect("cwd");
4948        let mut adapter = AcpAdapter::new(0, cwd, "sh", vec!["-c".into(), script.into()]);
4949        adapter.start().await.expect("initialize");
4950        assert!(matches!(
4951            adapter.next_event().await,
4952            Some(Ok(AgentEvent::Ready { .. }))
4953        ));
4954        adapter.send_prompt("hello".into()).await.expect("prompt");
4955        assert!(matches!(
4956            adapter.next_event().await,
4957            Some(Ok(AgentEvent::Text { text, .. })) if text == "partial"
4958        ));
4959        assert!(matches!(
4960            adapter.next_event().await,
4961            Some(Err(super::AdapterError::Protocol(detail))) if detail.contains("capacity")
4962        ));
4963    }
4964
4965    #[tokio::test]
4966    async fn acp_adapter_loads_existing_session_when_capability_allows_it() {
4967        let script = r#"read _; echo '{"jsonrpc":"2.0","id":1,"result":{"agentCapabilities":{"loadSession":true}}}'; read _; echo '{"jsonrpc":"2.0","id":2,"result":{}}'"#;
4968        let cwd = std::env::current_dir().expect("cwd");
4969        let mut adapter = AcpAdapter::with_session_id(
4970            0,
4971            cwd,
4972            "sh",
4973            vec!["-c".into(), script.into()],
4974            "existing-session",
4975        );
4976        adapter.start().await.expect("load existing session");
4977        assert!(matches!(
4978            adapter.next_event().await,
4979            Some(Ok(AgentEvent::Ready { .. }))
4980        ));
4981    }
4982
4983    #[tokio::test]
4984    async fn acp_start_failure_reaps_transport_process() {
4985        // The child emits an invalid initialize response and exits. The
4986        // adapter must not retain a live child after protocol startup fails;
4987        // this is the path used when a configured ACP command is unavailable
4988        // or speaks a different protocol.
4989        let mut adapter = AcpAdapter::new(
4990            0,
4991            std::env::current_dir().expect("cwd"),
4992            "sh",
4993            vec!["-c".into(), "printf 'not-json\\n'".into()],
4994        );
4995        assert!(adapter.start().await.is_err());
4996        assert!(adapter.child.is_none());
4997        assert!(adapter.reader.is_none());
4998    }
4999
5000    #[tokio::test]
5001    async fn acp_adapter_answers_permission_json_rpc_requests() {
5002        let path = std::env::temp_dir().join(format!(
5003            "codeswarm-permission-answer-{}",
5004            std::process::id()
5005        ));
5006        let script = format!(
5007            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"}}}}'"#,
5008            path.display()
5009        );
5010        let mut adapter = AcpAdapter::new(
5011            0,
5012            std::env::current_dir().expect("cwd"),
5013            "sh",
5014            vec!["-c".into(), script],
5015        );
5016        adapter.start().await.expect("start ACP");
5017        assert!(matches!(
5018            adapter.next_event().await,
5019            Some(Ok(AgentEvent::Ready { .. }))
5020        ));
5021        adapter.send_prompt("do it".into()).await.expect("prompt");
5022        assert!(matches!(
5023            adapter.next_event().await,
5024            Some(Ok(AgentEvent::Permission { request, .. }))
5025                if request.id == "9"
5026                    && request.options == ["Allow once"]
5027                    && request.option_ids == ["allow-once"]
5028        ));
5029        adapter
5030            .answer_permission(
5031                "9".into(),
5032                PermissionAnswer::Selected {
5033                    option_id: "allow-once".into(),
5034                },
5035            )
5036            .await
5037            .expect("permission answer");
5038        assert!(matches!(
5039            adapter.next_event().await,
5040            Some(Ok(AgentEvent::TurnComplete { .. }))
5041        ));
5042        let answer: Value = serde_json::from_str(
5043            &std::fs::read_to_string(&path).expect("captured permission answer"),
5044        )
5045        .expect("valid JSON-RPC answer");
5046        assert_eq!(answer["id"], 9);
5047        assert_eq!(answer["result"]["outcome"]["outcome"], "selected");
5048        assert_eq!(answer["result"]["outcome"]["optionId"], "allow-once");
5049        std::fs::remove_file(path).expect("cleanup");
5050    }
5051
5052    #[test]
5053    fn empty_acp_permission_options_are_not_exposed_as_a_blank_prompt() {
5054        let event = parse_acp_notification(
5055            0,
5056            r#"{"jsonrpc":"2.0","id":17,"method":"session/request_permission","params":{"options":[]}}"#,
5057        )
5058        .expect("valid JSON-RPC request");
5059        assert!(event.is_none());
5060    }
5061
5062    #[tokio::test]
5063    async fn native_stream_uses_success_result_response_when_chunks_are_missing() {
5064        let script_path = unique_test_path("codeswarm-agy-result-response", "sh");
5065        std::fs::write(
5066            &script_path,
5067            "#!/bin/sh\nprintf '%s\\n' '{\"event\":\"step_update\",\"step_update\":\"malformed\"}' '{\"event\":\"result\",\"result\":{\"status\":\"SUCCESS\",\"response\":\"Recovered.\"}}'\n",
5068        )
5069        .expect("write native test script");
5070        let mut adapter = AgyAdapter::new(
5071            0,
5072            std::env::current_dir().expect("cwd"),
5073            format!("sh {}", script_path.display()),
5074        );
5075        adapter.start().await.expect("start native adapter");
5076        assert!(adapter.next_event().await.is_some());
5077        assert!(adapter.next_event().await.is_some());
5078        adapter
5079            .send_prompt("continue".into())
5080            .await
5081            .expect("prompt");
5082        assert!(matches!(
5083            adapter.next_event().await,
5084            Some(Ok(AgentEvent::Text { text, .. })) if text == "Recovered."
5085        ));
5086        assert!(matches!(
5087            adapter.next_event().await,
5088            Some(Ok(AgentEvent::TurnComplete { .. }))
5089        ));
5090        adapter.stop().await.expect("stop native adapter");
5091        std::fs::remove_file(script_path).expect("cleanup native script");
5092    }
5093
5094    #[test]
5095    fn acp_workspace_file_access_is_root_bound_and_size_limited() {
5096        let root = std::env::temp_dir().join(format!("codeswarm-fs-{}", std::process::id()));
5097        let _ = std::fs::remove_dir_all(&root);
5098        std::fs::create_dir_all(&root).expect("workspace");
5099        std::fs::write(root.join("inside.txt"), "one\ntwo\nthree\n").expect("inside file");
5100        let outside =
5101            std::env::temp_dir().join(format!("codeswarm-outside-{}", std::process::id()));
5102        std::fs::write(&outside, "secret").expect("outside file");
5103        let link = root.join("outside-link");
5104        #[cfg(unix)]
5105        std::os::unix::fs::symlink(&outside, &link).expect("symlink");
5106        let adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5107
5108        assert_eq!(
5109            adapter
5110                .read_workspace_text("inside.txt", Some(2), Some(1))
5111                .expect("read inside"),
5112            "two"
5113        );
5114        std::fs::write(
5115            root.join("large.txt"),
5116            vec![b'x'; MAX_FILE_READ_BYTES + 1024],
5117        )
5118        .expect("large file");
5119        let bounded = adapter
5120            .read_workspace_text("large.txt", None, None)
5121            .expect("bounded read");
5122        assert!(bounded.len() <= MAX_FILE_READ_BYTES);
5123        #[cfg(unix)]
5124        {
5125            std::os::unix::fs::symlink(root.join("inside.txt"), root.join("inside-link"))
5126                .expect("internal symlink");
5127            assert_eq!(
5128                adapter
5129                    .read_workspace_text("inside-link", None, None)
5130                    .expect("read internal symlink"),
5131                "one\ntwo\nthree\n"
5132            );
5133        }
5134        assert!(adapter.workspace_path("../codeswarm-outside").is_err());
5135        assert!(
5136            adapter
5137                .workspace_path(&outside.display().to_string())
5138                .is_err()
5139        );
5140        #[cfg(unix)]
5141        assert!(adapter.workspace_path("outside-link").is_err());
5142        #[cfg(unix)]
5143        std::fs::remove_file(link).expect("cleanup symlink");
5144        #[cfg(unix)]
5145        std::fs::remove_file(root.join("inside-link")).expect("internal link cleanup");
5146        std::fs::remove_file(outside).expect("cleanup outside");
5147        std::fs::remove_dir_all(root).expect("cleanup workspace");
5148    }
5149
5150    #[tokio::test]
5151    async fn running_terminal_output_omits_exit_status_until_completion() {
5152        let root = unique_test_path("codeswarm-terminal-output", "dir");
5153        std::fs::create_dir_all(&root).expect("workspace");
5154        let mut adapter = AcpAdapter::new(0, root.clone(), "unused", Vec::new());
5155        let result = adapter
5156            .terminal_create(&serde_json::json!({
5157                "command": "sh",
5158                "args": ["-c", "sleep 0.2; printf done"],
5159                "cwd": ".",
5160            }))
5161            .await
5162            .expect("terminal create");
5163        let id = result["terminalId"].as_str().expect("terminal id");
5164        let output = adapter.terminal_output(id).await.expect("terminal output");
5165        assert!(output.get("exitStatus").is_none());
5166        if let Some(terminal) = adapter.terminals.remove(id) {
5167            terminal.stop().await;
5168        }
5169        std::fs::remove_dir_all(root).expect("cleanup workspace");
5170    }
5171
5172    #[tokio::test]
5173    async fn acp_adapter_answers_workspace_read_requests() {
5174        let root =
5175            std::env::temp_dir().join(format!("codeswarm-fs-request-{}", std::process::id()));
5176        let _ = std::fs::remove_dir_all(&root);
5177        std::fs::create_dir_all(&root).expect("workspace");
5178        let source = root.join("inside.txt");
5179        let answer = root.join("answer.json");
5180        std::fs::write(&source, "workspace content").expect("source");
5181        let script = format!(
5182            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"}}}}'"#,
5183            source.display(),
5184            answer.display(),
5185        );
5186        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5187        adapter.start().await.expect("start ACP");
5188        assert!(matches!(
5189            adapter.next_event().await,
5190            Some(Ok(AgentEvent::Ready { .. }))
5191        ));
5192        adapter.send_prompt("read it".into()).await.expect("prompt");
5193        assert!(matches!(
5194            adapter.next_event().await,
5195            Some(Ok(AgentEvent::TurnComplete { .. }))
5196        ));
5197        let response: Value =
5198            serde_json::from_str(&std::fs::read_to_string(&answer).expect("captured fs response"))
5199                .expect("response JSON");
5200        assert_eq!(response["id"], 9);
5201        assert_eq!(response["result"]["content"], "workspace content");
5202        adapter.stop().await.expect("stop ACP");
5203        std::fs::remove_dir_all(root).expect("cleanup workspace");
5204    }
5205
5206    #[tokio::test]
5207    async fn acp_adapter_runs_and_reports_client_mediated_terminals() {
5208        let root =
5209            std::env::temp_dir().join(format!("codeswarm-terminal-request-{}", std::process::id()));
5210        let _ = std::fs::remove_dir_all(&root);
5211        std::fs::create_dir_all(&root).expect("workspace");
5212        let create_request = serde_json::json!({
5213            "jsonrpc": "2.0",
5214            "id": 9,
5215            "method": "terminal/create",
5216            "params": {
5217                "sessionId": "s1",
5218                "command": "sh",
5219                "args": ["-c", "sleep 0.1; printf terminal-ok"],
5220                "cwd": ".",
5221            },
5222        });
5223        let wait_request = serde_json::json!({
5224            "jsonrpc": "2.0",
5225            "id": 10,
5226            "method": "terminal/wait_for_exit",
5227            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5228        });
5229        let output_request = serde_json::json!({
5230            "jsonrpc": "2.0",
5231            "id": 11,
5232            "method": "terminal/output",
5233            "params": {"sessionId": "s1", "terminalId": "terminal-1"},
5234        });
5235        let create_answer = root.join("create-answer.json");
5236        let wait_answer = root.join("wait-answer.json");
5237        let output_answer = root.join("output-answer.json");
5238        let script = format!(
5239            "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\"}}}}'",
5240            create_request,
5241            create_answer.display(),
5242            wait_request,
5243            wait_answer.display(),
5244            output_request,
5245            output_answer.display(),
5246        );
5247        let mut adapter = AcpAdapter::new(0, root.clone(), "sh", vec!["-c".into(), script]);
5248        adapter.start().await.expect("start ACP");
5249        assert!(matches!(
5250            adapter.next_event().await,
5251            Some(Ok(AgentEvent::Ready { .. }))
5252        ));
5253        adapter
5254            .send_prompt("run terminal".into())
5255            .await
5256            .expect("prompt");
5257        let mut saw_complete = false;
5258        for _ in 0..6 {
5259            match adapter.next_event().await {
5260                Some(Ok(AgentEvent::TurnComplete { .. })) => {
5261                    saw_complete = true;
5262                    break;
5263                }
5264                Some(_) => {}
5265                None => break,
5266            }
5267        }
5268        assert!(saw_complete, "terminal requests should not stall ACP");
5269        let create: Value = serde_json::from_str(
5270            &std::fs::read_to_string(&create_answer).expect("captured create response"),
5271        )
5272        .expect("create JSON");
5273        assert_eq!(create["result"]["terminalId"], "terminal-1");
5274        let output: Value = serde_json::from_str(
5275            &std::fs::read_to_string(&output_answer).expect("captured output response"),
5276        )
5277        .expect("output JSON");
5278        assert!(
5279            output["result"]["output"]
5280                .as_str()
5281                .unwrap_or_default()
5282                .contains("terminal-ok"),
5283            "output response: {output}"
5284        );
5285        adapter.stop().await.expect("stop ACP");
5286        std::fs::remove_dir_all(root).expect("cleanup workspace");
5287    }
5288
5289    #[tokio::test]
5290    async fn host_reduces_and_persists_adapter_events() {
5291        let path =
5292            std::env::temp_dir().join(format!("codeswarm-host-{}.jsonl", std::process::id()));
5293        let adapter = ScriptedAdapter::new(
5294            0,
5295            AgentCapabilities::default(),
5296            [AgentEvent::Text {
5297                slot: 0,
5298                text: "hello".into(),
5299            }],
5300        );
5301        let mut host = AdapterHost::new(Box::new(adapter), Some(EventLog::open(&path)));
5302        host.start().await.expect("start");
5303        host.next_effects()
5304            .await
5305            .expect("event")
5306            .expect("valid event");
5307        assert_eq!(host.state.public_text[0].1, "hello");
5308        assert_eq!(EventLog::open(&path).read().expect("read").len(), 1);
5309        std::fs::remove_file(path).expect("cleanup");
5310    }
5311
5312    #[tokio::test]
5313    async fn relay_applies_default_policy_before_the_first_prompt() {
5314        let first_log = Arc::new(Mutex::new(Vec::new()));
5315        let second_log = Arc::new(Mutex::new(Vec::new()));
5316        let first = AdapterHost::new(
5317            Box::new(ModeOrderAdapter {
5318                slot: 0,
5319                log: Arc::clone(&first_log),
5320                phase: 0,
5321            }),
5322            None,
5323        );
5324        let second = AdapterHost::new(
5325            Box::new(ModeOrderAdapter {
5326                slot: 1,
5327                log: Arc::clone(&second_log),
5328                phase: 0,
5329            }),
5330            None,
5331        );
5332        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
5333        relay.start().await.expect("start and synchronize policy");
5334        relay.run_turn("task", 0).await.expect("first turn");
5335        {
5336            let log = first_log.lock().expect("log");
5337            assert_eq!(log.as_slice(), ["start", "mode:yolo", "prompt"]);
5338        }
5339        assert_eq!(
5340            second_log.lock().expect("log").as_slice(),
5341            ["start", "mode:yolo"]
5342        );
5343        let added_log = Arc::new(Mutex::new(Vec::new()));
5344        relay
5345            .add_agent(
5346                AdapterHost::new(
5347                    Box::new(ModeOrderAdapter {
5348                        slot: 2,
5349                        log: Arc::clone(&added_log),
5350                        phase: 0,
5351                    }),
5352                    None,
5353                ),
5354                "Added",
5355                "added.example",
5356                "added-agent",
5357            )
5358            .await
5359            .expect("add with synchronized policy");
5360        assert_eq!(
5361            added_log.lock().expect("log").as_slice(),
5362            ["start", "mode:yolo"]
5363        );
5364        relay.drop_agent(2).await.expect("drop added agent");
5365        added_log.lock().expect("log").clear();
5366        relay
5367            .reload(2)
5368            .await
5369            .expect("reload with synchronized policy");
5370        assert_eq!(
5371            added_log.lock().expect("log").as_slice(),
5372            ["reload", "mode:yolo"]
5373        );
5374    }
5375
5376    #[tokio::test]
5377    async fn acp_roster_is_ready_before_any_prompt_is_sent() {
5378        let hosts = (0..2)
5379            .map(|slot| AdapterHost::new(Box::new(StartupAcpAdapter::new(slot)), None))
5380            .collect::<Vec<_>>();
5381        let startup_events = Arc::new(Mutex::new(Vec::new()));
5382        let captured = Arc::clone(&startup_events);
5383        let mut relay = RelayHost::new(hosts, 4).expect("relay");
5384        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5385
5386        relay.start().await.expect("complete startup handshake");
5387
5388        assert!(relay.dispatches().is_empty());
5389        let ready_slots = startup_events
5390            .lock()
5391            .expect("events")
5392            .iter()
5393            .filter_map(|event| match event {
5394                AgentEvent::Ready { slot, .. } => Some(*slot),
5395                _ => None,
5396            })
5397            .collect::<Vec<_>>();
5398        assert_eq!(ready_slots, vec![0, 1]);
5399    }
5400
5401    #[tokio::test]
5402    async fn independent_roster_adapters_start_concurrently() {
5403        let barrier = Arc::new(tokio::sync::Barrier::new(2));
5404        let hosts = (0..2)
5405            .map(|slot| {
5406                AdapterHost::new(
5407                    Box::new(ConcurrentStartAdapter {
5408                        slot,
5409                        barrier: Arc::clone(&barrier),
5410                    }),
5411                    None,
5412                )
5413            })
5414            .collect::<Vec<_>>();
5415        let mut relay = RelayHost::new(hosts, 4).expect("relay");
5416        tokio::time::timeout(std::time::Duration::from_millis(100), relay.start())
5417            .await
5418            .expect("startup should not serialize barrier participants")
5419            .expect("startup succeeds");
5420    }
5421
5422    #[tokio::test]
5423    async fn relay_host_dispatches_turns_sequentially() {
5424        let capabilities = AgentCapabilities {
5425            supports_cancel: true,
5426            ..AgentCapabilities::default()
5427        };
5428        let first = ScriptedAdapter::new(
5429            0,
5430            capabilities.clone(),
5431            [
5432                AgentEvent::Text {
5433                    slot: 0,
5434                    text: "first".into(),
5435                },
5436                AgentEvent::TurnComplete { slot: 0 },
5437            ],
5438        );
5439        let second = ScriptedAdapter::new(
5440            1,
5441            capabilities,
5442            [
5443                AgentEvent::Text {
5444                    slot: 1,
5445                    text: "review".into(),
5446                },
5447                AgentEvent::TurnComplete { slot: 1 },
5448            ],
5449        );
5450        let hosts = vec![
5451            AdapterHost::new(Box::new(first), None),
5452            AdapterHost::new(Box::new(second), None),
5453        ];
5454        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5455        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
5456        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5457        let captured = std::sync::Arc::clone(&events);
5458        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5459        relay.start().await.expect("start");
5460        events.lock().expect("events").clear();
5461        assert!(matches!(
5462            relay.run_turn("task", 0).await.expect("first turn"),
5463            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5464        ));
5465        assert!(matches!(
5466            relay.run_turn("first", 0).await.expect("second turn"),
5467            crate::relay::RelayDecision::Dispatch {
5468                slot: 1,
5469                can_stop: true,
5470                ..
5471            }
5472        ));
5473        assert_eq!(
5474            relay
5475                .dispatches()
5476                .iter()
5477                .map(|(slot, _)| *slot)
5478                .collect::<Vec<_>>(),
5479            [0, 1]
5480        );
5481        assert!(relay.dispatches()[0].1.contains("You are Claude"));
5482        assert!(
5483            relay.dispatches()[0]
5484                .1
5485                .contains("CodeSwarm roster (ordered)")
5486        );
5487        assert!(relay.dispatches()[0].1.contains("1. Claude — you"));
5488        assert!(relay.dispatches()[0].1.contains("2. Codex"));
5489        assert!(relay.dispatches()[1].1.contains(STOP_TOKEN));
5490        assert!(relay.dispatches()[0].1.contains("Do not use"));
5491        let lifecycle = events.lock().expect("events");
5492        let positions = lifecycle
5493            .iter()
5494            .filter_map(|event| match event {
5495                AgentEvent::TurnStarted { slot } => Some(("start", *slot)),
5496                AgentEvent::TurnComplete { slot } => Some(("complete", *slot)),
5497                _ => None,
5498            })
5499            .collect::<Vec<_>>();
5500        assert_eq!(
5501            positions,
5502            [("start", 0), ("complete", 0), ("start", 1), ("complete", 1)]
5503        );
5504    }
5505
5506    #[tokio::test]
5507    async fn relay_host_routes_around_a_usage_limited_agent() {
5508        let capabilities = AgentCapabilities::default();
5509        let first = ScriptedAdapter::new(
5510            0,
5511            capabilities.clone(),
5512            [
5513                AgentEvent::Text {
5514                    slot: 0,
5515                    text: "You've hit your usage limit. Visit chatgpt.com to purchase more \
5516                           credits or try again later."
5517                        .into(),
5518                },
5519                AgentEvent::TurnComplete { slot: 0 },
5520            ],
5521        );
5522        let second = ScriptedAdapter::new(
5523            1,
5524            capabilities,
5525            [
5526                AgentEvent::Text {
5527                    slot: 1,
5528                    text: "review done".into(),
5529                },
5530                AgentEvent::TurnComplete { slot: 1 },
5531            ],
5532        );
5533        let hosts = vec![
5534            AdapterHost::new(Box::new(first), None),
5535            AdapterHost::new(Box::new(second), None),
5536        ];
5537        let mut relay = super::RelayHost::new(hosts, 4).expect("relay");
5538        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5539        let captured = std::sync::Arc::clone(&events);
5540        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5541        relay.start().await.expect("start");
5542        events.lock().expect("events").clear();
5543        assert!(matches!(
5544            relay.run_turn("task", 0).await.expect("limited turn"),
5545            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5546        ));
5547        assert!(
5548            events
5549                .lock()
5550                .expect("events")
5551                .iter()
5552                .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5553        );
5554        // The next automatic turn skips the limited agent entirely.
5555        assert!(matches!(
5556            relay.run_turn("", 0).await.expect("next turn"),
5557            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5558        ));
5559        assert!(relay.relay().is_limited(0));
5560        // A reload restores the agent to the ring. (ScriptedAdapter cannot
5561        // feed further turns, so the restored routing itself is covered by
5562        // the relay unit tests.)
5563        relay.reload(0).await.expect("reload");
5564        assert!(!relay.relay().is_limited(0));
5565    }
5566
5567    #[tokio::test]
5568    async fn relay_host_routes_around_usage_limit_failures_without_tombstoning() {
5569        let limited = ScriptedAdapter::new(
5570            0,
5571            AgentCapabilities::default(),
5572            [AgentEvent::Failed {
5573                slot: 0,
5574                started: true,
5575                detail: "request failed: insufficient_quota".into(),
5576            }],
5577        );
5578        let healthy = ScriptedAdapter::new(
5579            1,
5580            AgentCapabilities::default(),
5581            [AgentEvent::TurnComplete { slot: 1 }],
5582        );
5583        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5584        let captured = std::sync::Arc::clone(&events);
5585        let mut relay = RelayHost::new(
5586            vec![
5587                AdapterHost::new(Box::new(limited), None),
5588                AdapterHost::new(Box::new(healthy), None),
5589            ],
5590            4,
5591        )
5592        .expect("relay");
5593        relay.set_event_sink(move |event| captured.lock().expect("events").push(event));
5594        relay.start().await.expect("start");
5595        events.lock().expect("events").clear();
5596
5597        assert!(matches!(
5598            relay.run_turn("task", 0).await.expect("limited failure"),
5599            crate::relay::RelayDecision::Dispatch { slot: 0, .. }
5600        ));
5601        assert!(relay.relay().is_limited(0));
5602        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), [0, 1]);
5603        {
5604            let events = events.lock().expect("events");
5605            assert!(
5606                events
5607                    .iter()
5608                    .any(|event| matches!(event, AgentEvent::UsageLimitReached { slot: 0, .. }))
5609            );
5610            assert!(
5611                !events
5612                    .iter()
5613                    .any(|event| matches!(event, AgentEvent::Failed { .. }))
5614            );
5615        }
5616
5617        assert!(matches!(
5618            relay.run_turn("", 0).await.expect("healthy peer"),
5619            crate::relay::RelayDecision::Dispatch { slot: 1, .. }
5620        ));
5621    }
5622
5623    #[tokio::test]
5624    async fn relay_failure_is_tombstoned_and_reported_to_the_ui_sink() {
5625        let failed = ScriptedAdapter::new(
5626            0,
5627            AgentCapabilities::default(),
5628            [AgentEvent::Failed {
5629                slot: 0,
5630                started: true,
5631                detail: "connection lost".into(),
5632            }],
5633        );
5634        let healthy = ScriptedAdapter::new(
5635            1,
5636            AgentCapabilities::default(),
5637            [AgentEvent::TurnComplete { slot: 1 }],
5638        );
5639        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5640        let captured = std::sync::Arc::clone(&events);
5641        let mut relay = RelayHost::new(
5642            vec![
5643                AdapterHost::new(Box::new(failed), None),
5644                AdapterHost::new(Box::new(healthy), None),
5645            ],
5646            4,
5647        )
5648        .expect("relay");
5649        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5650        relay.start().await.expect("start");
5651
5652        let error = relay.run_turn("task", 0).await.expect_err("failure");
5653        assert!(error.to_string().contains("connection lost"));
5654        assert_eq!(relay.relay().active_slots().collect::<Vec<_>>(), vec![1]);
5655        assert!(events.lock().expect("lock").iter().any(|event| {
5656            matches!(
5657                event,
5658                AgentEvent::Failed {
5659                    slot: 0,
5660                    started: true,
5661                    ..
5662                }
5663            )
5664        }));
5665    }
5666
5667    #[tokio::test]
5668    async fn codex_stop_does_not_skip_later_roster_reviewers() {
5669        let hosts = (0..3)
5670            .map(|slot| {
5671                AdapterHost::new(
5672                    Box::new(ScriptedAdapter::new(
5673                        slot,
5674                        AgentCapabilities::default(),
5675                        [
5676                            AgentEvent::Text {
5677                                slot,
5678                                text: STOP_TOKEN.into(),
5679                            },
5680                            AgentEvent::TurnComplete { slot },
5681                        ],
5682                    )),
5683                    None,
5684                )
5685            })
5686            .collect();
5687        let mut relay = RelayHost::new(hosts, 10).expect("relay");
5688        relay.set_roster_names(vec!["Claude".into(), "Codex".into(), "Qwen".into()]);
5689        relay.start().await.expect("start");
5690        for expected in 0..3 {
5691            assert!(matches!(relay.run_turn("task", 0).await.expect("turn"),
5692                RelayDecision::Dispatch { slot, can_stop, .. } if slot == expected && can_stop == (expected == 2)));
5693        }
5694        assert_eq!(
5695            relay.run_turn("", 0).await.expect("complete"),
5696            RelayDecision::Complete
5697        );
5698    }
5699
5700    #[tokio::test]
5701    async fn reviewer_stop_token_ends_the_automatic_relay_sequence() {
5702        let first = ScriptedAdapter::new(
5703            0,
5704            AgentCapabilities::default(),
5705            [
5706                AgentEvent::Text {
5707                    slot: 0,
5708                    text: "done".into(),
5709                },
5710                AgentEvent::TurnComplete { slot: 0 },
5711            ],
5712        );
5713        let reviewer = ScriptedAdapter::new(
5714            1,
5715            AgentCapabilities::default(),
5716            [
5717                AgentEvent::Text {
5718                    slot: 1,
5719                    text: STOP_TOKEN.into(),
5720                },
5721                AgentEvent::TurnComplete { slot: 1 },
5722            ],
5723        );
5724        let mut relay = RelayHost::new(
5725            vec![
5726                AdapterHost::new(Box::new(first), None),
5727                AdapterHost::new(Box::new(reviewer), None),
5728            ],
5729            10,
5730        )
5731        .expect("relay");
5732        relay.start().await.expect("start");
5733        let first_decision = relay.run_turn("task", 0).await.expect("first");
5734        assert!(matches!(
5735            first_decision,
5736            RelayDecision::Dispatch { slot: 0, .. }
5737        ));
5738        let reviewer_decision = relay.run_turn("", 0).await.expect("reviewer");
5739        assert!(matches!(
5740            reviewer_decision,
5741            RelayDecision::Dispatch {
5742                slot: 1,
5743                can_stop: true,
5744                ..
5745            }
5746        ));
5747        assert_eq!(
5748            relay.run_turn("", 0).await.expect("complete"),
5749            RelayDecision::Complete
5750        );
5751    }
5752
5753    #[tokio::test]
5754    async fn relay_stream_emits_text_and_thought_endings_before_tools() {
5755        let tool = AgentEvent::Tool {
5756            slot: 0,
5757            update: crate::ToolUpdate {
5758                id: "read".into(),
5759                title: "Read file".into(),
5760                status: ToolStatus::Running,
5761                detail: None,
5762            },
5763        };
5764        let updates = vec![
5765            AgentEvent::Thought {
5766                slot: 0,
5767                text: "Check the buffer. ✈".into(),
5768            },
5769            AgentEvent::Text {
5770                slot: 0,
5771                text: "Let me check.".into(),
5772            },
5773            tool.clone(),
5774            AgentEvent::Text {
5775                slot: 0,
5776                text: "[CODE".into(),
5777            },
5778            AgentEvent::Text {
5779                slot: 0,
5780                text: " is ordinary.".into(),
5781            },
5782            AgentEvent::Text {
5783                slot: 0,
5784                text: "[CODESWARM:".into(),
5785            },
5786            AgentEvent::Text {
5787                slot: 0,
5788                text: "STOP] Done.".into(),
5789            },
5790            tool,
5791            AgentEvent::TurnComplete { slot: 0 },
5792        ];
5793        let first = ScriptedAdapter::new(0, AgentCapabilities::default(), updates.clone());
5794        let reviewer = ScriptedAdapter::new(1, AgentCapabilities::default(), []);
5795        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5796        let captured = std::sync::Arc::clone(&events);
5797        let mut relay = RelayHost::new(
5798            vec![
5799                AdapterHost::new(Box::new(first), None),
5800                AdapterHost::new(Box::new(reviewer), None),
5801            ],
5802            2,
5803        )
5804        .expect("relay");
5805        relay.set_event_sink(move |event| captured.lock().unwrap().push(event));
5806        relay.start().await.unwrap();
5807        relay.run_turn("task", 0).await.unwrap();
5808        let captured = events.lock().unwrap();
5809        let visible: Vec<_> = captured
5810            .iter()
5811            .filter(|event| {
5812                matches!(
5813                    event,
5814                    AgentEvent::Text { .. } | AgentEvent::Thought { .. } | AgentEvent::Tool { .. }
5815                )
5816            })
5817            .cloned()
5818            .collect();
5819        assert_eq!(
5820            visible,
5821            vec![
5822                updates[0].clone(),
5823                updates[1].clone(),
5824                updates[2].clone(),
5825                AgentEvent::Text {
5826                    slot: 0,
5827                    text: "[CODE is ordinary.".into()
5828                },
5829                AgentEvent::Text {
5830                    slot: 0,
5831                    text: " Done.".into()
5832                },
5833                updates[7].clone(),
5834            ]
5835        );
5836    }
5837
5838    #[tokio::test]
5839    async fn stop_token_is_filtered_from_streamed_ui_events() {
5840        let first = ScriptedAdapter::new(
5841            0,
5842            AgentCapabilities::default(),
5843            [
5844                AgentEvent::Text {
5845                    slot: 0,
5846                    text: format!("visible {STOP_TOKEN} trailing"),
5847                },
5848                AgentEvent::TurnComplete { slot: 0 },
5849            ],
5850        );
5851        let reviewer = ScriptedAdapter::new(
5852            1,
5853            AgentCapabilities::default(),
5854            [AgentEvent::TurnComplete { slot: 1 }],
5855        );
5856        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5857        let captured = std::sync::Arc::clone(&events);
5858        let mut relay = RelayHost::new(
5859            vec![
5860                AdapterHost::new(Box::new(first), None),
5861                AdapterHost::new(Box::new(reviewer), None),
5862            ],
5863            2,
5864        )
5865        .expect("relay");
5866        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5867        relay.start().await.expect("start");
5868        relay.run_turn("task", 0).await.expect("turn");
5869        let captured = events.lock().expect("lock");
5870        assert!(captured.iter().all(|event| match event {
5871            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5872            _ => true,
5873        }));
5874        let visible = captured
5875            .iter()
5876            .filter_map(|event| match event {
5877                AgentEvent::Text { text, .. } => Some(text.as_str()),
5878                _ => None,
5879            })
5880            .collect::<String>();
5881        assert_eq!(visible, "visible  trailing");
5882    }
5883
5884    #[tokio::test]
5885    async fn token_only_reviewer_response_emits_visible_acknowledgment() {
5886        let first = ScriptedAdapter::new(
5887            0,
5888            AgentCapabilities::default(),
5889            [
5890                AgentEvent::Text {
5891                    slot: 0,
5892                    text: "done".into(),
5893                },
5894                AgentEvent::TurnComplete { slot: 0 },
5895            ],
5896        );
5897        let reviewer = ScriptedAdapter::new(
5898            1,
5899            AgentCapabilities::default(),
5900            [
5901                AgentEvent::Text {
5902                    slot: 1,
5903                    text: STOP_TOKEN.into(),
5904                },
5905                AgentEvent::TurnComplete { slot: 1 },
5906            ],
5907        );
5908        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5909        let captured = std::sync::Arc::clone(&events);
5910        let mut relay = RelayHost::new(
5911            vec![
5912                AdapterHost::new(Box::new(first), None),
5913                AdapterHost::new(Box::new(reviewer), None),
5914            ],
5915            4,
5916        )
5917        .expect("relay");
5918        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5919        relay.start().await.expect("start");
5920        relay.run_turn("task", 0).await.expect("first turn");
5921        relay.run_turn("", 0).await.expect("review turn");
5922        let captured = events.lock().expect("lock");
5923        assert!(captured.iter().any(|event| {
5924            matches!(
5925                event,
5926                AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5927            )
5928        }));
5929        assert!(captured.iter().all(|event| match event {
5930            AgentEvent::Text { text, .. } => !text.contains(STOP_TOKEN),
5931            _ => true,
5932        }));
5933        let acknowledgment = captured
5934            .iter()
5935            .position(|event| {
5936                matches!(
5937                    event,
5938                    AgentEvent::Text { slot: 1, text } if text == DEFAULT_STOP_ACKNOWLEDGMENT
5939                )
5940            })
5941            .expect("visible acknowledgment");
5942        let completion = captured
5943            .iter()
5944            .position(|event| matches!(event, AgentEvent::TurnComplete { slot: 1 }))
5945            .expect("reviewer completion");
5946        assert!(acknowledgment < completion);
5947    }
5948
5949    #[tokio::test]
5950    async fn explicit_reviewer_acknowledgment_is_not_duplicated_at_stop() {
5951        let first = ScriptedAdapter::new(
5952            0,
5953            AgentCapabilities::default(),
5954            [
5955                AgentEvent::Text {
5956                    slot: 0,
5957                    text: "done".into(),
5958                },
5959                AgentEvent::TurnComplete { slot: 0 },
5960            ],
5961        );
5962        let reviewer = ScriptedAdapter::new(
5963            1,
5964            AgentCapabilities::default(),
5965            [
5966                AgentEvent::Text {
5967                    slot: 1,
5968                    text: format!("{DEFAULT_STOP_ACKNOWLEDGMENT}\n{STOP_TOKEN}"),
5969                },
5970                AgentEvent::TurnComplete { slot: 1 },
5971            ],
5972        );
5973        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
5974        let captured = std::sync::Arc::clone(&events);
5975        let mut relay = RelayHost::new(
5976            vec![
5977                AdapterHost::new(Box::new(first), None),
5978                AdapterHost::new(Box::new(reviewer), None),
5979            ],
5980            4,
5981        )
5982        .expect("relay");
5983        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
5984        relay.start().await.expect("start");
5985        relay.run_turn("task", 0).await.expect("first turn");
5986        relay.run_turn("", 0).await.expect("review turn");
5987
5988        let visible = events
5989            .lock()
5990            .expect("lock")
5991            .iter()
5992            .filter_map(|event| match event {
5993                AgentEvent::Text { slot: 1, text } => Some(text.as_str()),
5994                _ => None,
5995            })
5996            .collect::<String>();
5997        assert_eq!(visible.trim(), DEFAULT_STOP_ACKNOWLEDGMENT);
5998        assert_eq!(visible.matches(DEFAULT_STOP_ACKNOWLEDGMENT).count(), 1);
5999    }
6000
6001    #[tokio::test]
6002    async fn relay_permission_answer_is_consumed_before_the_turn_completes() {
6003        let first = AdapterHost::new(
6004            Box::new(PermissionBlockingAdapter { slot: 0, phase: 0 }),
6005            None,
6006        );
6007        let second = AdapterHost::new(
6008            Box::new(ScriptedAdapter::new(
6009                1,
6010                AgentCapabilities::default(),
6011                [AgentEvent::TurnComplete { slot: 1 }],
6012            )),
6013            None,
6014        );
6015        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6016        let (seen_sender, mut seen_receiver) = tokio::sync::mpsc::unbounded_channel();
6017        relay.set_event_sink(move |event| {
6018            if matches!(event, AgentEvent::Permission { .. }) {
6019                let _ = seen_sender.send(());
6020            }
6021        });
6022        relay.start().await.expect("start");
6023        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
6024        let answer = async move {
6025            seen_receiver.recv().await.expect("permission request");
6026            sender
6027                .send(super::RelayPermissionAnswer {
6028                    slot: 0,
6029                    request_id: "permission-1".into(),
6030                    answer: PermissionAnswer::Selected {
6031                        option_id: "allow".into(),
6032                    },
6033                })
6034                .expect("queue permission answer");
6035        };
6036        tokio::time::timeout(std::time::Duration::from_millis(100), async {
6037            let ((), result) = tokio::join!(
6038                answer,
6039                relay.run_turn_with_permissions("task", 0, &mut receiver)
6040            );
6041            result
6042        })
6043        .await
6044        .expect("permission-gated turn should not deadlock")
6045        .expect("turn completes");
6046    }
6047
6048    #[tokio::test]
6049    async fn relay_cancellation_interrupts_a_waiting_adapter_turn() {
6050        let first = AdapterHost::new(
6051            Box::new(PendingAdapter {
6052                slot: 0,
6053                hang_on_cancel: false,
6054            }),
6055            None,
6056        );
6057        let second = AdapterHost::new(
6058            Box::new(ScriptedAdapter::new(
6059                1,
6060                AgentCapabilities::default(),
6061                [AgentEvent::TurnComplete { slot: 1 }],
6062            )),
6063            None,
6064        );
6065        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6066        relay.start().await.expect("start");
6067        let cancellation = relay.cancellation();
6068        let error = {
6069            let turn = relay.run_turn("task", 0);
6070            tokio::pin!(turn);
6071            cancellation.request();
6072            turn.await.expect_err("cancellation should stop turn")
6073        };
6074        assert!(error.to_string().contains("relay turn cancelled"));
6075
6076        assert!(relay.relay_mut().enqueue_human("replacement job", Some(1)));
6077        relay
6078            .run_turn("", 1)
6079            .await
6080            .expect("replacement job reaches the selected peer");
6081        let replacement = &relay.dispatches().last().expect("replacement dispatch").1;
6082        assert!(replacement.contains("replacement job"));
6083        assert!(replacement.contains("User "));
6084        assert!(replacement.contains(":\ntask"));
6085        let owner_updates = relay.relay_mut().unseen_context(0);
6086        assert!(owner_updates.contains("User "));
6087        assert!(owner_updates.contains(":\ntask"));
6088        assert!(owner_updates.contains(":\nreplacement job"));
6089    }
6090
6091    #[tokio::test]
6092    async fn relay_cancellation_does_not_wait_forever_for_a_broken_adapter() {
6093        let first = AdapterHost::new(
6094            Box::new(PendingAdapter {
6095                slot: 0,
6096                hang_on_cancel: true,
6097            }),
6098            None,
6099        );
6100        let second = AdapterHost::new(
6101            Box::new(PendingAdapter {
6102                slot: 1,
6103                hang_on_cancel: false,
6104            }),
6105            None,
6106        );
6107        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6108        relay.start().await.expect("start");
6109        let cancellation = relay.cancellation();
6110        let turn = relay.run_turn("task", 0);
6111        tokio::pin!(turn);
6112        cancellation.request();
6113        let error = turn.await.expect_err("cancellation should stop turn");
6114        assert!(error.to_string().contains("timed out"));
6115    }
6116
6117    #[tokio::test]
6118    async fn relay_host_pause_and_single_healthy_agent_continues_without_peer_review() {
6119        let event = [AgentEvent::TurnComplete { slot: 0 }];
6120        let first = AdapterHost::new(
6121            Box::new(ScriptedAdapter::new(
6122                0,
6123                AgentCapabilities::default(),
6124                event.clone(),
6125            )),
6126            None,
6127        );
6128        let second = AdapterHost::new(
6129            Box::new(ScriptedAdapter::new(
6130                1,
6131                AgentCapabilities::default(),
6132                [AgentEvent::TurnComplete { slot: 1 }],
6133            )),
6134            None,
6135        );
6136        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6137        relay.start().await.expect("start");
6138
6139        relay.pause();
6140        assert_eq!(
6141            relay.run_turn("paused", 0).await.expect("paused turn"),
6142            crate::relay::RelayDecision::Paused
6143        );
6144        assert!(relay.dispatches().is_empty());
6145
6146        relay.resume();
6147        relay.relay_mut().drop_agent(1).expect("drop reviewer");
6148        assert!(matches!(
6149            relay
6150                .run_turn("solo follow-up", 0)
6151                .await
6152                .expect("solo turn"),
6153            crate::relay::RelayDecision::Dispatch {
6154                slot: 0,
6155                can_stop: false,
6156                ..
6157            }
6158        ));
6159        assert_eq!(relay.dispatches().len(), 1);
6160    }
6161
6162    #[tokio::test]
6163    async fn relay_host_can_append_a_started_adapter_in_a_new_slot() {
6164        let first = AdapterHost::new(
6165            Box::new(ScriptedAdapter::new(
6166                0,
6167                AgentCapabilities::default(),
6168                [AgentEvent::TurnComplete { slot: 0 }],
6169            )),
6170            None,
6171        );
6172        let second = AdapterHost::new(
6173            Box::new(ScriptedAdapter::new(
6174                1,
6175                AgentCapabilities::default(),
6176                [AgentEvent::TurnComplete { slot: 1 }],
6177            )),
6178            None,
6179        );
6180        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6181        relay.set_roster_names(vec!["First".into(), "Second".into()]);
6182        relay.set_roster_identities(vec!["owner.example".into(), "peer.example".into()]);
6183        relay.set_roster_launch_specs(vec![
6184            ("custom".into(), "owner".into()),
6185            ("custom".into(), "peer".into()),
6186        ]);
6187        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6188        let captured = std::sync::Arc::clone(&events);
6189        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6190        relay.start().await.expect("start");
6191        let slot = relay
6192            .add_agent(
6193                AdapterHost::new(
6194                    Box::new(ScriptedAdapter::new(
6195                        2,
6196                        AgentCapabilities::default(),
6197                        [AgentEvent::TurnComplete { slot: 2 }],
6198                    )),
6199                    None,
6200                ),
6201                "Reviewer",
6202                "reviewer.example",
6203                "reviewer --acp",
6204            )
6205            .await
6206            .expect("append agent");
6207        assert_eq!(slot, 2);
6208        assert_eq!(
6209            relay.relay().active_slots().collect::<Vec<_>>(),
6210            vec![0, 1, 2]
6211        );
6212        assert_eq!(
6213            relay
6214                .session_metadata()
6215                .get("agents")
6216                .and_then(|value| value.as_array())
6217                .map(Vec::len),
6218            Some(3)
6219        );
6220        relay.drop_agent(1).await.expect("drop middle peer");
6221        let metadata = relay.session_metadata();
6222        assert_eq!(
6223            metadata.get("agents"),
6224            Some(&serde_json::json!([
6225                {"name": "First", "identity": "owner.example", "protocol": "custom", "command": "owner", "supports_load_session": false},
6226                {"name": "Reviewer", "identity": "reviewer.example", "protocol": "custom", "command": "reviewer --acp", "supports_load_session": false}
6227            ]))
6228        );
6229        assert!(
6230            events
6231                .lock()
6232                .expect("lock")
6233                .iter()
6234                .any(|event| { matches!(event, AgentEvent::Ready { slot: 2, .. }) })
6235        );
6236    }
6237
6238    #[tokio::test]
6239    async fn relay_host_persists_coordinator_owned_runtime_metadata() {
6240        let path = unique_test_path("codeswarm-session-metadata", "json");
6241        let metadata_store = crate::persistence::SessionMetadataStore::open(&path);
6242        let writer = metadata_store.buffered().expect("metadata writer");
6243        let first = AdapterHost::new(
6244            Box::new(ScriptedAdapter::new(0, AgentCapabilities::default(), [])),
6245            None,
6246        );
6247        let second = AdapterHost::new(
6248            Box::new(ScriptedAdapter::new(1, AgentCapabilities::default(), [])),
6249            None,
6250        );
6251        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6252        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6253        relay.set_roster_identities(vec!["claude.ai".into(), "openai.com".into()]);
6254        relay.set_roster_launch_specs(vec![
6255            ("custom".into(), "claude".into()),
6256            ("custom".into(), "codex".into()),
6257        ]);
6258        relay.set_session_metadata_writer(writer);
6259        relay.start().await.expect("start");
6260        relay.drop_agent(0).await.expect("drop first agent");
6261        relay.stop().await.expect("stop");
6262
6263        let loaded = metadata_store
6264            .read()
6265            .expect("read metadata")
6266            .expect("metadata snapshot");
6267        assert_eq!(loaded.get("title"), Some(&serde_json::json!("CodeSwarm")));
6268        assert_eq!(
6269            loaded.get("agents"),
6270            Some(&serde_json::json!([{
6271                "name": "Codex", "identity": "openai.com", "protocol": "custom",
6272                "command": "codex", "supports_load_session": false
6273            }]))
6274        );
6275        assert!(loaded.get("owner").is_none());
6276        let _ = std::fs::remove_file(path);
6277    }
6278
6279    #[tokio::test]
6280    async fn relay_host_swaps_live_adapters_and_remaps_stream_events() {
6281        let first = AdapterHost::new(
6282            Box::new(ScriptedAdapter::new(
6283                0,
6284                AgentCapabilities::default(),
6285                [
6286                    AgentEvent::Text {
6287                        slot: 0,
6288                        text: "owner stream".into(),
6289                    },
6290                    AgentEvent::TurnComplete { slot: 0 },
6291                ],
6292            )),
6293            None,
6294        );
6295        let second = AdapterHost::new(
6296            Box::new(ScriptedAdapter::new(
6297                1,
6298                AgentCapabilities::default(),
6299                [
6300                    AgentEvent::Text {
6301                        slot: 1,
6302                        text: "peer stream".into(),
6303                    },
6304                    AgentEvent::TurnComplete { slot: 1 },
6305                ],
6306            )),
6307            None,
6308        );
6309        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6310        relay.set_roster_names(vec!["Owner".into(), "Peer".into()]);
6311        relay.set_roster_identities(vec!["first.example".into(), "second.example".into()]);
6312        let events = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
6313        let captured = std::sync::Arc::clone(&events);
6314        relay.set_event_sink(move |event| captured.lock().expect("lock").push(event));
6315        relay.start().await.expect("start");
6316
6317        relay.swap_agents(0, 1).expect("swap peers");
6318        assert_eq!(relay.active_slot_for_identity("first.example"), Some(1));
6319        assert_eq!(relay.active_slot_for_identity("second.example"), Some(0));
6320        relay.run_turn("task", 0).await.expect("swapped turn");
6321        let events = events.lock().expect("events");
6322        assert!(events.iter().any(|event| {
6323            matches!(event, AgentEvent::Text { slot: 0, text } if text == "peer stream")
6324        }));
6325        assert!(relay.dispatches()[0].1.contains("You are Peer"));
6326    }
6327
6328    #[tokio::test]
6329    async fn relay_host_persists_all_active_agent_metadata_off_thread() {
6330        let path = unique_test_path("codeswarm-session-metadata", "json");
6331        let first = AdapterHost::new(
6332            Box::new(ScriptedAdapter::new(
6333                0,
6334                AgentCapabilities::default(),
6335                [AgentEvent::TurnComplete { slot: 0 }],
6336            )),
6337            None,
6338        );
6339        let second = AdapterHost::new(
6340            Box::new(ScriptedAdapter::new(
6341                1,
6342                AgentCapabilities::default(),
6343                [AgentEvent::TurnComplete { slot: 1 }],
6344            )),
6345            None,
6346        );
6347        let mut relay = RelayHost::new(vec![first, second], 4).expect("relay");
6348        relay.set_roster_names(vec!["Claude".into(), "Codex".into()]);
6349        relay.set_roster_identities(vec!["claude.com".into(), "openai.com".into()]);
6350        relay.set_roster_launch_specs(vec![
6351            ("custom".into(), "claude".into()),
6352            ("custom".into(), "codex".into()),
6353        ]);
6354        let writer = SessionMetadataStore::open(&path)
6355            .buffered()
6356            .expect("metadata writer");
6357        relay.set_session_metadata_writer(writer);
6358        relay.start().await.expect("start");
6359        relay.stop().await.expect("stop");
6360        let loaded = SessionMetadataStore::open(&path)
6361            .read()
6362            .expect("read metadata")
6363            .expect("metadata snapshot");
6364        let agents = loaded
6365            .get("agents")
6366            .and_then(|value| value.as_array())
6367            .expect("agents");
6368        assert_eq!(agents.len(), 2);
6369        assert_eq!(agents[0]["identity"], "claude.com");
6370        assert_eq!(agents[1]["identity"], "openai.com");
6371        let _ = std::fs::remove_file(path);
6372    }
6373
6374    #[tokio::test]
6375    async fn relay_host_routes_unseen_public_context_to_next_agent() {
6376        let first = AdapterHost::new(
6377            Box::new(ScriptedAdapter::new(
6378                0,
6379                AgentCapabilities::default(),
6380                [
6381                    AgentEvent::Text {
6382                        slot: 0,
6383                        text: "implemented the fix".into(),
6384                    },
6385                    AgentEvent::TurnComplete { slot: 0 },
6386                ],
6387            )),
6388            None,
6389        );
6390        let second = AdapterHost::new(
6391            Box::new(ScriptedAdapter::new(
6392                1,
6393                AgentCapabilities::default(),
6394                [AgentEvent::TurnComplete { slot: 1 }],
6395            )),
6396            None,
6397        );
6398        let mut relay = super::RelayHost::new(vec![first, second], 4).expect("relay");
6399        relay.set_roster_names(vec!["Codex".into(), "Qwen".into()]);
6400        relay.start().await.expect("start");
6401        relay.run_turn("task", 0).await.expect("first turn");
6402        relay.run_turn("review this", 0).await.expect("review turn");
6403
6404        assert_eq!(relay.dispatches().len(), 2);
6405        assert_eq!(relay.dispatches()[0].0, 0);
6406        assert!(relay.dispatches()[0].1.contains("task"));
6407        assert!(relay.dispatches()[0].1.contains("You are Codex"));
6408        assert!(relay.dispatches()[0].1.contains("2. Qwen"));
6409        assert_eq!(relay.dispatches()[1].0, 1);
6410        assert!(relay.dispatches()[1].1.contains("review this"));
6411        let public = relay.dispatches()[1]
6412            .1
6413            .split_once("Public updates:\n")
6414            .map(|(_, updates)| updates)
6415            .expect("review receives public context");
6416        let header = public
6417            .lines()
6418            .find(|line| line.starts_with("Codex "))
6419            .expect("named previous agent");
6420        let timestamp = header
6421            .strip_prefix("Codex ")
6422            .and_then(|value| value.strip_suffix(':'))
6423            .expect("timestamped header");
6424        assert_eq!(timestamp.len(), 5);
6425        assert_eq!(timestamp.as_bytes()[2], b':');
6426        assert!(
6427            timestamp
6428                .bytes()
6429                .enumerate()
6430                .all(|(index, byte)| { index == 2 || byte.is_ascii_digit() })
6431        );
6432        assert!(public.contains("implemented the fix"));
6433        assert!(!public.contains("Agent 0"));
6434    }
6435}