Skip to main content

codeswarm_adapters/
adapters.rs

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