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