Skip to main content

codeswarm_adapters/
adapters.rs

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