Skip to main content

codeswarm_adapters/
adapters.rs

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