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