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