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