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