Skip to main content

claude_wrapper/
duplex.rs

1//! Long-lived duplex stream-json sessions.
2//!
3//! [`DuplexSession`] holds a `claude` subprocess open in
4//! `--input-format stream-json --output-format stream-json` mode for
5//! the duration of a conversation. A single child is held open across
6//! many turns; user messages are written to its stdin, NDJSON events
7//! are read from its stdout and dispatched back to `send()` callers.
8//!
9//! # When to use
10//!
11//! [`DuplexSession`] is the recommended primitive for long-running
12//! hosts that drive multi-turn conversations: agent servers, IDE
13//! backends, daemons, chat UIs. Holding the child open across turns
14//! amortizes init cost and unlocks capabilities that are awkward or
15//! impossible from a transient subprocess: mid-turn permission
16//! decisions ([`PermissionHandler`]), clean
17//! [interrupts](DuplexSession::interrupt), and a typed
18//! [event subscriber stream](DuplexSession::subscribe) that fans out
19//! events to multiple consumers.
20//!
21//! For short-lived processes (CLIs, build scripts, batch jobs,
22//! lambdas) where each turn can stand on its own, prefer
23//! [`QueryCommand`] for one-off calls or [`Session`] for transient
24//! multi-turn with cumulative cost / history tracking.
25//!
26//! # Cost and budget bookkeeping
27//!
28//! [`DuplexSession`] itself keeps no cross-turn accounting: each
29//! [`TurnResult`] carries that turn's cost and nothing accumulates.
30//! For cumulative cost, turn history, and a
31//! [`BudgetTracker`]-enforced spend ceiling (send fails fast with
32//! [`Error::BudgetExceeded`] once the ceiling is hit), wrap the
33//! session in a [`Conversation`] --
34//! it is a thin bookkeeping layer over this module, not a different
35//! transport.
36//!
37//! [`QueryCommand`]: crate::QueryCommand
38//! [`Session`]: crate::session::Session
39//! [`Conversation`]: crate::conversation::Conversation
40//! [`BudgetTracker`]: crate::budget::BudgetTracker
41//!
42//! # Example
43//!
44//! ```no_run
45//! use claude_wrapper::Claude;
46//! use claude_wrapper::duplex::{DuplexOptions, DuplexSession};
47//!
48//! # async fn example() -> claude_wrapper::Result<()> {
49//! let claude = Claude::builder().build()?;
50//! let session = DuplexSession::spawn(
51//!     &claude,
52//!     DuplexOptions::default().model("haiku"),
53//! ).await?;
54//!
55//! let turn = session.send("hello").await?;
56//! if let Some(text) = turn.result_text() {
57//!     println!("{text}");
58//! }
59//!
60//! session.close().await?;
61//! # Ok(())
62//! # }
63//! ```
64//!
65//! # Subscribers
66//!
67//! For event-driven UIs that want to react to assistant tokens,
68//! tool-use blocks, or system events as they arrive, call
69//! [`DuplexSession::subscribe`] before issuing a [`DuplexSession::send`].
70//! Each receiver gets its own buffered view of the event stream;
71//! slow consumers see [`tokio::sync::broadcast::error::RecvError::Lagged`]
72//! rather than blocking the session task.
73//!
74//! ```no_run
75//! use claude_wrapper::Claude;
76//! use claude_wrapper::duplex::{DuplexOptions, DuplexSession, InboundEvent};
77//!
78//! # async fn example() -> claude_wrapper::Result<()> {
79//! let claude = Claude::builder().build()?;
80//! let session = DuplexSession::spawn(&claude, DuplexOptions::default()).await?;
81//!
82//! let mut rx = session.subscribe();
83//! let _turn = session.send("hello").await?;
84//!
85//! while let Ok(event) = rx.try_recv() {
86//!     match event {
87//!         InboundEvent::SystemInit { session_id } => {
88//!             println!("session id: {session_id}");
89//!         }
90//!         InboundEvent::Assistant(_) => {
91//!             // partial or complete assistant message
92//!         }
93//!         _ => {}
94//!     }
95//! }
96//!
97//! session.close().await?;
98//! # Ok(())
99//! # }
100//! ```
101//!
102//! For interleaved (concurrent) event handling while a turn is in
103//! flight, drive `rx.recv()` and the `send()` future together via
104//! `tokio::select!`. Pin the send future and use a block scope so
105//! its borrow of the session ends before [`DuplexSession::close`].
106//!
107//! # Mid-turn permission decisions
108//!
109//! Configure a [`PermissionHandler`] at spawn time to answer the
110//! CLI's permission prompts in-flight. The session writes
111//! `--permission-prompt-tool stdio` automatically when a handler is
112//! set, so the CLI emits `control_request` messages for tool use
113//! over the duplex channel rather than blocking on a TUI prompt.
114//!
115//! ```no_run
116//! use claude_wrapper::Claude;
117//! use claude_wrapper::duplex::{
118//!     DuplexOptions, DuplexSession, PermissionDecision, PermissionHandler,
119//! };
120//!
121//! # async fn example() -> claude_wrapper::Result<()> {
122//! let handler = PermissionHandler::new(|req| async move {
123//!     if req.tool_name == "Bash" {
124//!         PermissionDecision::Deny { message: "bash is denied".into() }
125//!     } else {
126//!         PermissionDecision::Allow { updated_input: None }
127//!     }
128//! });
129//!
130//! let claude = Claude::builder().build()?;
131//! let session = DuplexSession::spawn(
132//!     &claude,
133//!     DuplexOptions::default().on_permission(handler),
134//! ).await?;
135//! # Ok(())
136//! # }
137//! ```
138//!
139//! For human-in-the-loop UIs, return [`PermissionDecision::Defer`]
140//! from the handler, capture the [`PermissionRequest::request_id`],
141//! and answer later via [`DuplexSession::respond_to_permission`].
142//!
143//! **Known limitation:** as of claude CLI 2.1.x,
144//! `--permission-prompt-tool stdio` does not cause the CLI to emit
145//! `control_request {subtype: "can_use_tool"}` in
146//! `--print --output-format stream-json` mode. The permission handler
147//! registered here is wire-correct and unit-tested, but will not be
148//! invoked end-to-end until the upstream CLI bug is resolved. Tracked
149//! upstream at
150//! <https://github.com/anthropics/claude-agent-sdk-python/issues/469>.
151//!
152//! # Mid-turn interrupt
153//!
154//! [`DuplexSession::interrupt`] sends a clean
155//! `control_request {subtype: "interrupt"}` to the CLI. The CLI
156//! stops generating, closes the in-flight turn (`send().await`
157//! resolves with the truncated [`TurnResult`]), and answers our
158//! interrupt with a `control_response`. Use this instead of dropping
159//! the session or killing the child when you want to cancel one
160//! turn but keep the conversation going.
161//!
162//! ```no_run
163//! use std::time::Duration;
164//! use claude_wrapper::Claude;
165//! use claude_wrapper::duplex::{DuplexOptions, DuplexSession};
166//!
167//! # async fn example() -> claude_wrapper::Result<()> {
168//! let claude = Claude::builder().build()?;
169//! let session = DuplexSession::spawn(&claude, DuplexOptions::default()).await?;
170//!
171//! let send_fut = session.send("write a long essay about rust");
172//! let interrupt_fut = async {
173//!     tokio::time::sleep(Duration::from_millis(500)).await;
174//!     session.interrupt().await
175//! };
176//!
177//! let (turn, interrupt_result) = tokio::join!(send_fut, interrupt_fut);
178//! let _truncated = turn?;
179//! interrupt_result?;
180//! # Ok(())
181//! # }
182//! ```
183//!
184//! # Phased rollout
185//!
186//! This module rolled out in four PRs tracked in
187//! <https://github.com/joshrotenberg/claude-wrapper/issues/561>:
188//! `spawn`/`send`/`close` (PR 1), `subscribe` (PR 2), mid-turn
189//! permission handling (PR 3), and `interrupt` (PR 4, this one).
190
191use std::collections::HashMap;
192use std::future::Future;
193use std::pin::Pin;
194use std::process::Stdio;
195use std::sync::Arc;
196use std::time::Duration;
197
198use serde_json::Value;
199use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
200use tokio::process::{Child, ChildStdin, ChildStdout, Command};
201use tokio::sync::{broadcast, mpsc, oneshot, watch};
202use tokio::task::JoinHandle;
203use tracing::{debug, warn};
204
205use crate::Claude;
206use crate::command::spawn_args::SharedSpawnArgs;
207use crate::error::{Error, Result};
208use crate::tool_pattern::ToolPattern;
209use crate::types::{Effort, HermeticScope, PermissionMode};
210
211/// Default capacity of the per-session [`broadcast::Sender`] backing
212/// [`DuplexSession::subscribe`].
213///
214/// Override per-session via [`DuplexOptions::subscriber_capacity`].
215pub const DEFAULT_SUBSCRIBER_CAPACITY: usize = 256;
216
217/// A mid-turn permission prompt from the CLI for a single tool
218/// invocation.
219///
220/// Forwarded to the [`PermissionHandler`] registered via
221/// [`DuplexOptions::on_permission`]. Capture
222/// [`Self::request_id`] inside your handler if you intend to return
223/// [`PermissionDecision::Defer`] and answer later via
224/// [`DuplexSession::respond_to_permission`].
225#[derive(Debug, Clone)]
226pub struct PermissionRequest {
227    /// CLI-assigned correlation id. Pass this to
228    /// [`DuplexSession::respond_to_permission`] when deferring.
229    pub request_id: String,
230    /// The tool the model wants to use (e.g. `"Bash"`, `"Edit"`).
231    pub tool_name: String,
232    /// The tool's `input` payload as the model produced it.
233    pub input: Value,
234    /// The full `request` object as sent by the CLI, for fields not
235    /// promoted to typed accessors.
236    pub raw: Value,
237}
238
239/// The decision returned from a [`PermissionHandler`] (or passed to
240/// [`DuplexSession::respond_to_permission`] for deferred decisions).
241///
242/// `Allow` and `Deny` both write a control response to the CLI
243/// immediately. `Defer` causes the run loop to skip writing a
244/// response; the caller is then expected to invoke
245/// [`DuplexSession::respond_to_permission`] later. Passing `Defer`
246/// to `respond_to_permission` is a no-op.
247#[derive(Debug, Clone)]
248pub enum PermissionDecision {
249    /// Allow the tool to run, optionally with rewritten input.
250    Allow {
251        /// Replace the model's input with this object before running
252        /// the tool. `None` keeps the original input.
253        updated_input: Option<Value>,
254    },
255    /// Deny the tool. The `message` is surfaced to the model.
256    Deny {
257        /// Human-readable explanation given back to the model.
258        message: String,
259    },
260    /// Decision pending; the caller will supply it later via
261    /// [`DuplexSession::respond_to_permission`].
262    Defer,
263}
264
265type PermissionFuture = Pin<Box<dyn Future<Output = PermissionDecision> + Send + 'static>>;
266type PermissionFn = dyn Fn(PermissionRequest) -> PermissionFuture + Send + Sync + 'static;
267
268/// A user-supplied async callback invoked when the CLI requests
269/// permission to use a tool.
270///
271/// Construct with [`Self::new`], passing an `async fn` or
272/// async-block closure. Cheap to clone (`Arc` under the hood).
273///
274/// The handler runs inline on the duplex session's task. The CLI is
275/// blocked on the response while the handler runs, so awaiting an
276/// async policy check (DB lookup, remote call) is fine. If the
277/// decision needs human input on a different timescale, return
278/// [`PermissionDecision::Defer`] and answer via
279/// [`DuplexSession::respond_to_permission`] when ready.
280#[derive(Clone)]
281pub struct PermissionHandler {
282    inner: Arc<PermissionFn>,
283}
284
285impl PermissionHandler {
286    /// Wrap an async closure as a permission handler.
287    ///
288    /// # Example
289    ///
290    /// ```
291    /// use claude_wrapper::duplex::{PermissionDecision, PermissionHandler};
292    ///
293    /// let _handler = PermissionHandler::new(|req| async move {
294    ///     if req.tool_name == "Bash" {
295    ///         PermissionDecision::Deny { message: "no bash".into() }
296    ///     } else {
297    ///         PermissionDecision::Allow { updated_input: None }
298    ///     }
299    /// });
300    /// ```
301    pub fn new<F, Fut>(f: F) -> Self
302    where
303        F: Fn(PermissionRequest) -> Fut + Send + Sync + 'static,
304        Fut: Future<Output = PermissionDecision> + Send + 'static,
305    {
306        Self {
307            inner: Arc::new(move |req| Box::pin(f(req))),
308        }
309    }
310
311    fn invoke(&self, req: PermissionRequest) -> PermissionFuture {
312        (self.inner)(req)
313    }
314}
315
316impl std::fmt::Debug for PermissionHandler {
317    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
318        f.debug_struct("PermissionHandler").finish_non_exhaustive()
319    }
320}
321
322/// Configuration for [`DuplexSession::spawn`].
323///
324/// Builder methods cover the spawn-time options shared with
325/// [`QueryCommand`](crate::QueryCommand); the flag emission lives on a
326/// common internal `SharedSpawnArgs`, so the oneshot and duplex paths
327/// cannot drift on how a knob is rendered. The spawn call always
328/// includes
329/// `--print --verbose --input-format stream-json --output-format stream-json`
330/// regardless of these options.
331///
332/// A few `QueryCommand` knobs are intentionally not surfaced here
333/// because they only make sense for a oneshot run or are owned by the
334/// duplex transport itself:
335///
336/// - Transport is fixed: `output_format`, `input_format`,
337///   `include_partial_messages`, `verbose`, and `prompt_via_stdin` are
338///   pinned by the duplex spawn and not configurable.
339/// - `retry_policy` reruns a whole oneshot invocation; a duplex session
340///   holds one child open across turns, so there is nothing to retry at
341///   this layer.
342/// - `brief` and `from_pr` shape a single oneshot run (SendUserMessage
343///   for one-turn agent-to-user replies, resume-from-PR startup); a
344///   duplex host drives turns and session selection itself.
345/// - `prompt_suggestions` and `replay_user_messages` shape stdin/stream
346///   echoing; the duplex layer owns its own stream plumbing.
347///
348/// Use [`Self::arg`] if you need one of these on a duplex spawn anyway.
349#[derive(Debug, Default, Clone)]
350pub struct DuplexOptions {
351    // Spawn-time knobs shared with QueryCommand; the flag emission
352    // lives on SharedSpawnArgs so the two builders cannot drift.
353    shared: SharedSpawnArgs,
354    additional_args: Vec<String>,
355    subscriber_capacity: Option<usize>,
356    on_permission: Option<PermissionHandler>,
357}
358
359impl DuplexOptions {
360    /// Set the model for this session (`--model`).
361    #[must_use]
362    pub fn model(mut self, model: impl Into<String>) -> Self {
363        self.shared.model = Some(model.into());
364        self
365    }
366
367    /// Set the system prompt for this session (`--system-prompt`).
368    #[must_use]
369    pub fn system_prompt(mut self, prompt: impl Into<String>) -> Self {
370        self.shared.system_prompt = Some(prompt.into());
371        self
372    }
373
374    /// Append to the default system prompt (`--append-system-prompt`).
375    #[must_use]
376    pub fn append_system_prompt(mut self, prompt: impl Into<String>) -> Self {
377        self.shared.append_system_prompt = Some(prompt.into());
378        self
379    }
380
381    /// Resume a prior session by id (`--resume <session_id>`).
382    ///
383    /// Mirrors [`QueryCommand::resume`](crate::QueryCommand::resume)
384    /// for the duplex path. The spawned `claude` process picks up the
385    /// conversation that produced `session_id` and continues it; turns
386    /// sent through [`DuplexSession::send`] append to the existing
387    /// history rather than starting fresh.
388    ///
389    /// Use case: a host (IDE, MCP server, agent backend) wants to
390    /// upgrade a passive on-disk session to a live duplex one --
391    /// pulls the `session_id` out of the existing JSONL log, opens a
392    /// duplex session here, and the next turn extends the same
393    /// conversation.
394    ///
395    /// `resume` and [`Self::continue_session`] are mutually exclusive
396    /// at the CLI; passing both lets the CLI decide (it errors today).
397    #[must_use]
398    pub fn resume(mut self, session_id: impl Into<String>) -> Self {
399        self.shared.resume = Some(session_id.into());
400        self
401    }
402
403    /// Continue the most recent session in the current working
404    /// directory (`--continue`).
405    ///
406    /// Mirrors [`QueryCommand::continue_session`](crate::QueryCommand::continue_session)
407    /// for the duplex path. Use [`Self::resume`] to pick a specific
408    /// session id; use this when "the last one" is what you want.
409    #[must_use]
410    pub fn continue_session(mut self) -> Self {
411        self.shared.continue_session = true;
412        self
413    }
414
415    /// Run this session in a fresh git worktree (`--worktree [name]`).
416    ///
417    /// `name` is the optional worktree name (the CLI auto-generates
418    /// one if omitted). Calling this method always enables the
419    /// worktree flag, with or without a name.
420    ///
421    /// Use case: an agent host wants the chat's writes isolated from
422    /// the current working tree -- the chat opens with a fresh
423    /// worktree, mutations land there, and the host can inspect or
424    /// merge later.
425    #[must_use]
426    pub fn worktree(mut self, name: Option<impl Into<String>>) -> Self {
427        self.shared.worktree = true;
428        if let Some(n) = name {
429            self.shared.worktree_name = Some(n.into());
430        }
431        self
432    }
433
434    /// Pin the session to a named subagent (`--agent <name>`).
435    ///
436    /// `name` is resolved by the CLI in this order: inline
437    /// definitions from [`Self::agents_json`], then user-level
438    /// `~/.claude/agents/<name>.md` files, then project-level dirs
439    /// loaded by the active `--setting-sources`.
440    ///
441    /// **Caveat**: as of Claude Code 2.1.143, the CLI silently
442    /// ignores an unknown `name` and falls back to the default
443    /// behavior -- no warning, no error. Callers that want a hard
444    /// "agent must exist" semantics should validate the name out of
445    /// band (e.g. via [`crate::artifacts::AgentsRoot::get`]) before
446    /// passing it here.
447    #[must_use]
448    pub fn agent(mut self, name: impl Into<String>) -> Self {
449        self.shared.agent = Some(name.into());
450        self
451    }
452
453    /// Inline subagent definitions for this session
454    /// (`--agents <json>`).
455    ///
456    /// `json` is a JSON object keyed by agent name, with each value
457    /// carrying at least `description` and `prompt`. Inline
458    /// definitions take precedence over on-disk
459    /// `~/.claude/agents/*.md` of the same name. Pass [`Self::agent`]
460    /// to select which one to use as the session's persona.
461    ///
462    /// Example: `{"reviewer": {"description": "Reviews code",
463    /// "prompt": "You are a code reviewer"}}`.
464    #[must_use]
465    pub fn agents_json(mut self, json: impl Into<String>) -> Self {
466        self.shared.agents_json = Some(json.into());
467        self
468    }
469
470    /// Set the permission mode for this session
471    /// (`--permission-mode <mode>`).
472    ///
473    /// Mirrors [`QueryCommand::permission_mode`](crate::QueryCommand::permission_mode)
474    /// for the duplex path. The default mode (when this method isn't
475    /// called) drops to the CLI's interactive prompt for every
476    /// tool-use approval, which is broken for non-interactive duplex
477    /// sessions -- nothing answers the prompts and the session stalls
478    /// or fails. Call this with [`PermissionMode::AcceptEdits`] for
479    /// the "edit files autonomously" pattern, [`PermissionMode::Plan`]
480    /// for read-only planning, etc.
481    ///
482    /// Bypass mode is a footgun; reach for [`Self::dangerously_skip_permissions`]
483    /// (or, for stricter discipline, [`crate::dangerous::DangerousClient`])
484    /// when you really need it.
485    #[must_use]
486    pub fn permission_mode(mut self, mode: PermissionMode) -> Self {
487        self.shared.permission_mode = Some(mode);
488        self
489    }
490
491    /// Pass `--dangerously-skip-permissions` to the spawned session.
492    ///
493    /// Bypasses ALL permission checks -- file edits, bash, network,
494    /// the lot. Use only when you know the session runs in a trusted
495    /// sandbox (a fresh worktree, a container, etc.). For most "run
496    /// autonomously" cases you want [`Self::permission_mode`] with
497    /// [`PermissionMode::AcceptEdits`] instead.
498    #[must_use]
499    pub fn dangerously_skip_permissions(mut self) -> Self {
500        self.shared.dangerously_skip_permissions = true;
501        self
502    }
503
504    /// Start a new session under a caller-chosen id
505    /// (`--session-id <uuid>`).
506    ///
507    /// Mirrors [`QueryCommand::session_id`](crate::QueryCommand::session_id)
508    /// for the duplex path. Unlike [`Self::resume`] (pick up an
509    /// existing session) or [`Self::continue_session`] (pick up the
510    /// most recent one), this mints a fresh session whose id the host
511    /// knows up front -- useful when the host indexes sessions
512    /// externally before the first turn completes.
513    #[must_use]
514    pub fn session_id(mut self, id: impl Into<String>) -> Self {
515        self.shared.session_id = Some(id.into());
516        self
517    }
518
519    /// Set a JSON schema for structured output validation
520    /// (`--json-schema <schema>`).
521    ///
522    /// Mirrors [`QueryCommand::json_schema`](crate::QueryCommand::json_schema)
523    /// for the duplex path. `schema` is the inline JSON of the schema;
524    /// the turn's closing result message carries the validated
525    /// `structured_output`.
526    #[must_use]
527    pub fn json_schema(mut self, schema: impl Into<String>) -> Self {
528        self.shared.json_schema = Some(schema.into());
529        self
530    }
531
532    /// Add allowed tool patterns (`--allowed-tools`).
533    ///
534    /// Mirrors [`QueryCommand::allowed_tools`](crate::QueryCommand::allowed_tools)
535    /// for the duplex path: accepts anything convertible into
536    /// [`ToolPattern`], including bare strings (e.g. `"Bash"`,
537    /// `"Bash(git log:*)"`, `"mcp__my-server__*"`), and joins them
538    /// into the comma-separated form the CLI expects.
539    #[must_use]
540    pub fn allowed_tools<I, T>(mut self, tools: I) -> Self
541    where
542        I: IntoIterator<Item = T>,
543        T: Into<ToolPattern>,
544    {
545        self.shared
546            .allowed_tools
547            .extend(tools.into_iter().map(Into::into));
548        self
549    }
550
551    /// Add a single allowed tool pattern.
552    #[must_use]
553    pub fn allowed_tool(mut self, tool: impl Into<ToolPattern>) -> Self {
554        self.shared.allowed_tools.push(tool.into());
555        self
556    }
557
558    /// Add disallowed tool patterns (`--disallowed-tools`).
559    #[must_use]
560    pub fn disallowed_tools<I, T>(mut self, tools: I) -> Self
561    where
562        I: IntoIterator<Item = T>,
563        T: Into<ToolPattern>,
564    {
565        self.shared
566            .disallowed_tools
567            .extend(tools.into_iter().map(Into::into));
568        self
569    }
570
571    /// Add a single disallowed tool pattern.
572    #[must_use]
573    pub fn disallowed_tool(mut self, tool: impl Into<ToolPattern>) -> Self {
574        self.shared.disallowed_tools.push(tool.into());
575        self
576    }
577
578    /// Cap the number of agentic turns (`--max-turns <n>`).
579    ///
580    /// The cap applies per turn sent through [`DuplexSession::send`];
581    /// a turn that exhausts it closes with an `error_max_turns`
582    /// result rather than an assistant reply.
583    #[must_use]
584    pub fn max_turns(mut self, turns: u32) -> Self {
585        self.shared.max_turns = Some(turns);
586        self
587    }
588
589    /// Cap claude's own spend for the session
590    /// (`--max-budget-usd <usd>`).
591    ///
592    /// This is the CLI's cap, checked post-hoc after each API call,
593    /// so a session can overspend before tripping. It is distinct
594    /// from the wrapper's [`BudgetTracker`](crate::budget::BudgetTracker)
595    /// ceiling, which gates dispatch host-side -- attach one via
596    /// [`Conversation::with_budget`](crate::conversation::Conversation::with_budget)
597    /// to stop a duplex conversation before the next turn is sent.
598    #[must_use]
599    pub fn max_budget_usd(mut self, budget: f64) -> Self {
600        self.shared.max_budget_usd = Some(budget);
601        self
602    }
603
604    /// Set a fallback model for when the primary is overloaded
605    /// (`--fallback-model <model>`).
606    #[must_use]
607    pub fn fallback_model(mut self, model: impl Into<String>) -> Self {
608        self.shared.fallback_model = Some(model.into());
609        self
610    }
611
612    /// Set the reasoning effort level (`--effort <level>`).
613    #[must_use]
614    pub fn effort(mut self, effort: Effort) -> Self {
615        self.shared.effort = Some(effort);
616        self
617    }
618
619    /// Add an additional directory for tool access
620    /// (`--add-dir <dir>`, repeatable).
621    #[must_use]
622    pub fn add_dir(mut self, dir: impl Into<String>) -> Self {
623        self.shared.add_dir.push(dir.into());
624        self
625    }
626
627    /// Add an MCP config file path (`--mcp-config <path>`,
628    /// repeatable).
629    ///
630    /// Pair with [`crate::McpConfigBuilder`] to generate the file.
631    #[must_use]
632    pub fn mcp_config(mut self, path: impl Into<String>) -> Self {
633        self.shared.mcp_config.push(path.into());
634        self
635    }
636
637    /// Only use MCP servers from `--mcp-config` files, ignoring the
638    /// user- and project-level MCP configuration
639    /// (`--strict-mcp-config`).
640    #[must_use]
641    pub fn strict_mcp_config(mut self) -> Self {
642        self.shared.strict_mcp_config = true;
643        self
644    }
645
646    /// Comma-separated list of setting sources the CLI loads, for example
647    /// `"user,project,local"` (`--setting-sources`). Pass an empty string to
648    /// load none, sealing the session's promptspace against ambient project
649    /// config (agents, skills, `CLAUDE.md`). Mirrors
650    /// [`QueryCommand::setting_sources`](crate::QueryCommand::setting_sources).
651    #[must_use]
652    pub fn setting_sources(mut self, sources: impl Into<String>) -> Self {
653        self.shared.setting_sources = Some(sources.into());
654        self
655    }
656
657    /// Seal the ambient `~/.claude` config for a reproducible session
658    /// ([`HermeticScope::Full`]).
659    ///
660    /// Sets `--setting-sources ""`, `--strict-mcp-config`, and
661    /// `--exclude-dynamic-system-prompt-sections`. This gives a warm
662    /// duplex session a clean seal without the [`Self::arg`] escape
663    /// hatch. Mirrors
664    /// [`QueryCommand::hermetic`](crate::QueryCommand::hermetic).
665    ///
666    /// This is not [`Self::bare`]: a hermetic seal leaves OAuth and
667    /// keychain auth working, whereas `--bare` forces API-key billing.
668    /// A later [`Self::setting_sources`] call overrides the seal scope.
669    #[must_use]
670    pub fn hermetic(mut self) -> Self {
671        self.shared.apply_hermetic(HermeticScope::Full);
672        self
673    }
674
675    /// Seal the ambient `~/.claude` config at an explicit
676    /// [`HermeticScope`].
677    ///
678    /// See [`Self::hermetic`] for the flag set. Mirrors
679    /// [`QueryCommand::hermetic_scoped`](crate::QueryCommand::hermetic_scoped).
680    #[must_use]
681    pub fn hermetic_scoped(mut self, scope: HermeticScope) -> Self {
682        self.shared.apply_hermetic(scope);
683        self
684    }
685
686    /// Do not persist the session to on-disk history
687    /// (`--no-session-persistence`).
688    #[must_use]
689    pub fn no_session_persistence(mut self) -> Self {
690        self.shared.no_session_persistence = true;
691        self
692    }
693
694    /// Set the list of available built-in tools (`--tools`).
695    ///
696    /// Use `""` to disable all tools, `"default"` for all tools, or
697    /// specific tool names like `["Bash", "Edit", "Read"]`. This is
698    /// distinct from [`Self::allowed_tools`], which controls tool
699    /// permissions rather than which built-ins load. Mirrors
700    /// [`QueryCommand::tools`](crate::QueryCommand::tools).
701    #[must_use]
702    pub fn tools(mut self, tools: impl IntoIterator<Item = impl Into<String>>) -> Self {
703        self.shared.tools.extend(tools.into_iter().map(Into::into));
704        self
705    }
706
707    /// Add a file resource to download at startup (`--file`).
708    ///
709    /// Format: `file_id:relative_path` (e.g. `file_abc:doc.txt`).
710    /// Repeatable. Mirrors [`QueryCommand::file`](crate::QueryCommand::file).
711    #[must_use]
712    pub fn file(mut self, spec: impl Into<String>) -> Self {
713        self.shared.file.push(spec.into());
714        self
715    }
716
717    /// Path to a settings JSON file or a JSON string (`--settings`).
718    ///
719    /// Mirrors [`QueryCommand::settings`](crate::QueryCommand::settings).
720    #[must_use]
721    pub fn settings(mut self, settings: impl Into<String>) -> Self {
722        self.shared.settings = Some(settings.into());
723        self
724    }
725
726    /// When resuming, create a new session id instead of reusing the
727    /// original (`--fork-session`).
728    ///
729    /// Only meaningful alongside [`Self::resume`] or
730    /// [`Self::continue_session`]. Mirrors
731    /// [`QueryCommand::fork_session`](crate::QueryCommand::fork_session).
732    #[must_use]
733    pub fn fork_session(mut self) -> Self {
734        self.shared.fork_session = true;
735        self
736    }
737
738    /// Enable debug logging with an optional filter, e.g. `"api,hooks"`
739    /// (`--debug`). Mirrors
740    /// [`QueryCommand::debug_filter`](crate::QueryCommand::debug_filter).
741    #[must_use]
742    pub fn debug_filter(mut self, filter: impl Into<String>) -> Self {
743        self.shared.debug_filter = Some(filter.into());
744        self
745    }
746
747    /// Write debug logs to the given file path (`--debug-file`).
748    ///
749    /// Mirrors [`QueryCommand::debug_file`](crate::QueryCommand::debug_file).
750    #[must_use]
751    pub fn debug_file(mut self, path: impl Into<String>) -> Self {
752        self.shared.debug_file = Some(path.into());
753        self
754    }
755
756    /// Beta feature headers for API key authentication (`--betas`).
757    ///
758    /// Mirrors [`QueryCommand::betas`](crate::QueryCommand::betas).
759    #[must_use]
760    pub fn betas(mut self, betas: impl Into<String>) -> Self {
761        self.shared.betas = Some(betas.into());
762        self
763    }
764
765    /// Load plugins from the given directory for this session
766    /// (`--plugin-dir`). Repeatable. Mirrors
767    /// [`QueryCommand::plugin_dir`](crate::QueryCommand::plugin_dir).
768    #[must_use]
769    pub fn plugin_dir(mut self, dir: impl Into<String>) -> Self {
770        self.shared.plugin_dirs.push(dir.into());
771        self
772    }
773
774    /// Fetch a plugin `.zip` from a URL for this session only
775    /// (`--plugin-url`). Repeatable. Mirrors
776    /// [`QueryCommand::plugin_url`](crate::QueryCommand::plugin_url).
777    #[must_use]
778    pub fn plugin_url(mut self, url: impl Into<String>) -> Self {
779        self.shared.plugin_urls.push(url.into());
780        self
781    }
782
783    /// Create a tmux session for the worktree (`--tmux`).
784    ///
785    /// Mirrors [`QueryCommand::tmux`](crate::QueryCommand::tmux).
786    #[must_use]
787    pub fn tmux(mut self) -> Self {
788        self.shared.tmux = true;
789        self
790    }
791
792    /// Run in minimal mode (`--bare`).
793    ///
794    /// Skips hooks, LSP, plugin sync, attribution, auto-memory,
795    /// background prefetches, keychain reads, and `CLAUDE.md`
796    /// auto-discovery. Anthropic auth is restricted to
797    /// `ANTHROPIC_API_KEY` or `apiKeyHelper`; OAuth and keychain are
798    /// never read. Mirrors [`QueryCommand::bare`](crate::QueryCommand::bare).
799    #[must_use]
800    pub fn bare(mut self) -> Self {
801        self.shared.bare = true;
802        self
803    }
804
805    /// Start with all customizations disabled (`--safe-mode`).
806    ///
807    /// Disables `CLAUDE.md`, skills, plugins, hooks, MCP servers,
808    /// custom commands and agents, and output styles for
809    /// troubleshooting. Mirrors
810    /// [`QueryCommand::safe_mode`](crate::QueryCommand::safe_mode).
811    #[must_use]
812    pub fn safe_mode(mut self) -> Self {
813        self.shared.safe_mode = true;
814        self
815    }
816
817    /// Disable all slash-command skills (`--disable-slash-commands`).
818    ///
819    /// Mirrors
820    /// [`QueryCommand::disable_slash_commands`](crate::QueryCommand::disable_slash_commands).
821    #[must_use]
822    pub fn disable_slash_commands(mut self) -> Self {
823        self.shared.disable_slash_commands = true;
824        self
825    }
826
827    /// Include every hook lifecycle event in the stream-json output
828    /// (`--include-hook-events`).
829    ///
830    /// Duplex sessions always run in stream-json, so this takes effect
831    /// without extra configuration. Mirrors
832    /// [`QueryCommand::include_hook_events`](crate::QueryCommand::include_hook_events).
833    #[must_use]
834    pub fn include_hook_events(mut self) -> Self {
835        self.shared.include_hook_events = true;
836        self
837    }
838
839    /// Move per-machine sections (cwd, env info, memory paths, git
840    /// status) out of the system prompt and into the first user
841    /// message (`--exclude-dynamic-system-prompt-sections`).
842    ///
843    /// Improves cross-user prompt-cache reuse. Only applies with the
844    /// default system prompt; ignored with [`Self::system_prompt`].
845    /// Mirrors
846    /// [`QueryCommand::exclude_dynamic_system_prompt_sections`](crate::QueryCommand::exclude_dynamic_system_prompt_sections).
847    #[must_use]
848    pub fn exclude_dynamic_system_prompt_sections(mut self) -> Self {
849        self.shared.exclude_dynamic_system_prompt_sections = true;
850        self
851    }
852
853    /// Set a display name for this session (`--name`). Shown in the
854    /// `/resume` picker and terminal title. Mirrors
855    /// [`QueryCommand::name`](crate::QueryCommand::name).
856    #[must_use]
857    pub fn name(mut self, name: impl Into<String>) -> Self {
858        self.shared.name = Some(name.into());
859        self
860    }
861
862    /// Add a raw argument to the spawn command line.
863    ///
864    /// Escape hatch for flags not covered by the dedicated builder
865    /// methods.
866    #[must_use]
867    pub fn arg(mut self, arg: impl Into<String>) -> Self {
868        self.additional_args.push(arg.into());
869        self
870    }
871
872    /// Set the per-session [`broadcast::Sender`] capacity backing
873    /// [`DuplexSession::subscribe`].
874    ///
875    /// Defaults to [`DEFAULT_SUBSCRIBER_CAPACITY`] (256). Larger
876    /// values give slow subscribers more room before they
877    /// [`Lagged`](tokio::sync::broadcast::error::RecvError::Lagged);
878    /// smaller values reclaim memory if you do not subscribe.
879    #[must_use]
880    pub fn subscriber_capacity(mut self, capacity: usize) -> Self {
881        self.subscriber_capacity = Some(capacity);
882        self
883    }
884
885    /// Register a [`PermissionHandler`] to answer the CLI's tool-use
886    /// permission prompts in-flight.
887    ///
888    /// When set, the spawn command line includes
889    /// `--permission-prompt-tool stdio`, which configures the CLI to
890    /// emit `control_request` messages for tool use over the duplex
891    /// channel rather than blocking on a TUI prompt.
892    ///
893    /// Without a handler, the session does not pass
894    /// `--permission-prompt-tool` and the CLI applies its default
895    /// permission policy (driven by `--permission-mode`).
896    ///
897    /// **Known limitation:** as of claude CLI 2.1.x the CLI does not
898    /// emit `control_request {subtype: "can_use_tool"}` in stream-json
899    /// print mode, so this handler will not be invoked end-to-end until
900    /// an upstream fix lands. The wire handling is correct; see
901    /// <https://github.com/anthropics/claude-agent-sdk-python/issues/469>.
902    #[must_use]
903    pub fn on_permission(mut self, handler: PermissionHandler) -> Self {
904        self.on_permission = Some(handler);
905        self
906    }
907
908    fn into_args(self) -> Vec<String> {
909        let mut args = vec![
910            "--print".to_string(),
911            "--verbose".to_string(),
912            "--output-format".to_string(),
913            "stream-json".to_string(),
914            "--input-format".to_string(),
915            "stream-json".to_string(),
916        ];
917
918        self.shared.append_to(&mut args);
919
920        if self.on_permission.is_some() {
921            args.push("--permission-prompt-tool".to_string());
922            args.push("stdio".to_string());
923        }
924        args.extend(self.additional_args);
925
926        args
927    }
928}
929
930/// The result of one turn through a [`DuplexSession`].
931///
932/// `result` is the raw JSON of the `{"type": "result", ...}` message
933/// that closed the turn. `events` carries every other message
934/// received during the turn (system, assistant, stream_event, user)
935/// in arrival order, with the closing `result` excluded.
936#[derive(Debug, Clone)]
937pub struct TurnResult {
938    /// The raw `{"type": "result", ...}` message that ended the turn.
939    pub result: Value,
940    /// Every other message received during the turn, in order.
941    pub events: Vec<Value>,
942}
943
944impl TurnResult {
945    /// Extract `result.result` as a string, if present.
946    #[must_use]
947    pub fn result_text(&self) -> Option<&str> {
948        self.result.get("result").and_then(Value::as_str)
949    }
950
951    /// Extract `result.session_id`, if present.
952    #[must_use]
953    pub fn session_id(&self) -> Option<&str> {
954        self.result.get("session_id").and_then(Value::as_str)
955    }
956
957    /// Extract `total_cost_usd` (preferred) or the legacy `cost_usd`
958    /// field, if either is present.
959    ///
960    /// This is the cost of one turn. For the conversation-wide
961    /// running total, record turns through a
962    /// [`Conversation`](crate::conversation::Conversation).
963    #[must_use]
964    pub fn total_cost_usd(&self) -> Option<f64> {
965        self.result
966            .get("total_cost_usd")
967            .or_else(|| self.result.get("cost_usd"))
968            .and_then(Value::as_f64)
969    }
970
971    /// Extract `duration_ms`, if present.
972    #[must_use]
973    pub fn duration_ms(&self) -> Option<u64> {
974        self.result.get("duration_ms").and_then(Value::as_u64)
975    }
976}
977
978/// A classified inbound event broadcast to [`DuplexSession::subscribe`]
979/// receivers.
980///
981/// Every non-`result` message coming back from the CLI is broadcast as
982/// one of these variants. The closing `{"type": "result"}` message is
983/// not broadcast; it resolves the in-flight [`DuplexSession::send`]
984/// future and lands in [`TurnResult::result`].
985///
986/// Subscribers see the same set of events that accumulate in
987/// [`TurnResult::events`], in the same order, just classified. Adding
988/// a typed accessor for a new event type later (e.g. promoting a
989/// `system` subtype into its own variant) is non-breaking against the
990/// `Other` fallback.
991#[derive(Debug, Clone)]
992pub enum InboundEvent {
993    /// First `{"type": "system", "subtype": "init"}` event for the
994    /// session. Carries the CLI-assigned `session_id`.
995    SystemInit {
996        /// The CLI-assigned session id, useful for logging or
997        /// future resume support.
998        session_id: String,
999    },
1000    /// `{"type": "assistant", ...}` -- either a complete assistant
1001    /// message or, in stream-json mode, a partial chunk.
1002    Assistant(Value),
1003    /// `{"type": "stream_event", ...}` -- low-level streaming event
1004    /// emitted while a turn is in progress.
1005    StreamEvent(Value),
1006    /// `{"type": "user", ...}` -- typically a tool result echo from
1007    /// the CLI side.
1008    User(Value),
1009    /// Any other event type, including non-`init` `system` events
1010    /// and any message types not yet recognised by this enum.
1011    Other(Value),
1012}
1013
1014fn classify(msg: &Value) -> InboundEvent {
1015    match msg.get("type").and_then(Value::as_str) {
1016        Some("system") => {
1017            if msg.get("subtype").and_then(Value::as_str) == Some("init")
1018                && let Some(id) = msg.get("session_id").and_then(Value::as_str)
1019            {
1020                return InboundEvent::SystemInit {
1021                    session_id: id.to_string(),
1022                };
1023            }
1024            InboundEvent::Other(msg.clone())
1025        }
1026        Some("assistant") => InboundEvent::Assistant(msg.clone()),
1027        Some("stream_event") => InboundEvent::StreamEvent(msg.clone()),
1028        Some("user") => InboundEvent::User(msg.clone()),
1029        _ => InboundEvent::Other(msg.clone()),
1030    }
1031}
1032
1033/// Liveness state of a [`DuplexSession`]'s background task.
1034///
1035/// Surfaced through [`DuplexSession::is_alive`],
1036/// [`DuplexSession::exit_status`], and
1037/// [`DuplexSession::wait_for_exit`] for service-shaped hosts that
1038/// want non-consuming visibility into whether a session is still
1039/// usable. The closing [`DuplexSession::close`] still returns the
1040/// full [`Result`] for the one caller that consumes the session.
1041///
1042/// `Failed` carries a `String` rather than the full
1043/// [`Error`] because the underlying watch channel requires `Clone`
1044/// and `Error` is not `Clone` (its `Io` variant wraps a non-`Clone`
1045/// `std::io::Error`). The full error remains available via
1046/// [`DuplexSession::close`].
1047#[derive(Debug, Clone)]
1048pub enum SessionExitStatus {
1049    /// The session task is still running.
1050    Running,
1051    /// The session task completed normally (close, stdout EOF without
1052    /// error).
1053    Completed,
1054    /// The session task ended with an error. Carries the error's
1055    /// `Display` rendering.
1056    Failed(String),
1057}
1058
1059/// A long-lived `claude` subprocess in stream-json duplex mode.
1060///
1061/// Owns a background task that holds the child open, writes user
1062/// messages to its stdin, and reads NDJSON events from its stdout.
1063/// One turn at a time: calling [`Self::send`] while another turn is
1064/// in flight returns [`Error::DuplexTurnInFlight`].
1065///
1066/// See the [module docs](crate::duplex) for the full design.
1067#[derive(Debug)]
1068pub struct DuplexSession {
1069    outbound_tx: mpsc::UnboundedSender<OutboundMsg>,
1070    events_tx: broadcast::Sender<InboundEvent>,
1071    exit_rx: watch::Receiver<SessionExitStatus>,
1072    join: JoinHandle<Result<()>>,
1073}
1074
1075#[derive(Debug)]
1076enum OutboundMsg {
1077    Send {
1078        prompt: String,
1079        reply: oneshot::Sender<Result<TurnResult>>,
1080    },
1081    PermissionResponse {
1082        request_id: String,
1083        decision: PermissionDecision,
1084    },
1085    Interrupt {
1086        reply: oneshot::Sender<Result<()>>,
1087    },
1088}
1089
1090impl DuplexSession {
1091    /// Spawn a fresh `claude` subprocess in duplex mode.
1092    ///
1093    /// The child is started with
1094    /// `--print --verbose --input-format stream-json --output-format stream-json`
1095    /// plus any options applied via `opts`. The session task takes
1096    /// ownership of the child; dropping the returned handle (or
1097    /// calling [`Self::close`]) shuts the task down.
1098    pub async fn spawn(claude: &Claude, opts: DuplexOptions) -> Result<Self> {
1099        let capacity = opts
1100            .subscriber_capacity
1101            .unwrap_or(DEFAULT_SUBSCRIBER_CAPACITY);
1102        let permission_handler = opts.on_permission.clone();
1103
1104        let mut command_args = Vec::new();
1105        command_args.extend(claude.global_args.clone());
1106        command_args.extend(opts.into_args());
1107
1108        debug!(
1109            binary = %claude.binary.display(),
1110            args = ?command_args,
1111            "spawning duplex claude session"
1112        );
1113
1114        let mut cmd = Command::new(&claude.binary);
1115        cmd.args(&command_args)
1116            .env_remove("CLAUDECODE")
1117            .env_remove("CLAUDE_CODE_ENTRYPOINT")
1118            .envs(&claude.env)
1119            .stdin(Stdio::piped())
1120            .stdout(Stdio::piped())
1121            .stderr(Stdio::piped())
1122            .kill_on_drop(true);
1123
1124        if let Some(ref dir) = claude.working_dir {
1125            cmd.current_dir(dir);
1126        }
1127
1128        let mut child = cmd.spawn().map_err(|e| Error::Io {
1129            message: format!("failed to spawn claude: {e}"),
1130            source: e,
1131            working_dir: claude.working_dir.clone(),
1132        })?;
1133
1134        let stdin = child.stdin.take().expect("stdin was piped");
1135        let stdout = child.stdout.take().expect("stdout was piped");
1136
1137        let (outbound_tx, outbound_rx) = mpsc::unbounded_channel();
1138        let (events_tx, _initial_rx) = broadcast::channel(capacity);
1139        let (exit_tx, exit_rx) = watch::channel(SessionExitStatus::Running);
1140
1141        let join = tokio::spawn(run_session(
1142            child,
1143            stdin,
1144            stdout,
1145            outbound_rx,
1146            events_tx.clone(),
1147            permission_handler,
1148            exit_tx,
1149        ));
1150
1151        Ok(Self {
1152            outbound_tx,
1153            events_tx,
1154            exit_rx,
1155            join,
1156        })
1157    }
1158
1159    /// Send one user message and await the closing result event.
1160    ///
1161    /// Returns [`Error::DuplexTurnInFlight`] if another turn is
1162    /// already pending, and [`Error::DuplexClosed`] if the session
1163    /// task has already exited.
1164    ///
1165    /// The returned [`TurnResult`] is per-turn; nothing accumulates
1166    /// across calls. For cumulative cost/history and a budget
1167    /// ceiling, send through a
1168    /// [`Conversation`](crate::conversation::Conversation) instead.
1169    pub async fn send(&self, prompt: impl Into<String>) -> Result<TurnResult> {
1170        let (reply_tx, reply_rx) = oneshot::channel();
1171        self.outbound_tx
1172            .send(OutboundMsg::Send {
1173                prompt: prompt.into(),
1174                reply: reply_tx,
1175            })
1176            .map_err(|_| Error::DuplexClosed)?;
1177        reply_rx.await.map_err(|_| Error::DuplexClosed)?
1178    }
1179
1180    /// Subscribe to the session's classified inbound event stream.
1181    ///
1182    /// Returns a [`broadcast::Receiver<InboundEvent>`] that receives
1183    /// every non-`result` event as it arrives. Each subscriber gets
1184    /// its own buffered view; subscribers added later miss earlier
1185    /// events. Slow subscribers see
1186    /// [`RecvError::Lagged`](tokio::sync::broadcast::error::RecvError::Lagged)
1187    /// rather than blocking the session task.
1188    ///
1189    /// Subscribers see the same events that accumulate in
1190    /// [`TurnResult::events`], in the same order.
1191    ///
1192    /// # Example
1193    ///
1194    /// ```no_run
1195    /// use claude_wrapper::Claude;
1196    /// use claude_wrapper::duplex::{DuplexOptions, DuplexSession, InboundEvent};
1197    ///
1198    /// # async fn example() -> claude_wrapper::Result<()> {
1199    /// let claude = Claude::builder().build()?;
1200    /// let session = DuplexSession::spawn(&claude, DuplexOptions::default()).await?;
1201    /// let mut rx = session.subscribe();
1202    ///
1203    /// // Subscribe before send so we receive every event.
1204    /// let _turn = session.send("hello").await?;
1205    ///
1206    /// while let Ok(event) = rx.try_recv() {
1207    ///     if let InboundEvent::SystemInit { session_id } = event {
1208    ///         println!("session id: {session_id}");
1209    ///     }
1210    /// }
1211    /// # Ok(())
1212    /// # }
1213    /// ```
1214    #[must_use]
1215    pub fn subscribe(&self) -> broadcast::Receiver<InboundEvent> {
1216        self.events_tx.subscribe()
1217    }
1218
1219    /// Cheap, non-blocking liveness check.
1220    ///
1221    /// Returns `true` while the session task is running, `false` once
1222    /// it has exited (whether normally or with an error). Multiple
1223    /// concurrent callers are allowed, and the call does not consume
1224    /// the session: [`Self::close`] still works after polling.
1225    ///
1226    /// Reads the latest value from a `tokio::sync::watch` channel
1227    /// updated from inside the session task, so it never blocks and
1228    /// reflects state set just before the task returns.
1229    #[must_use]
1230    pub fn is_alive(&self) -> bool {
1231        matches!(*self.exit_rx.borrow(), SessionExitStatus::Running)
1232    }
1233
1234    /// Snapshot the session task's [`SessionExitStatus`].
1235    ///
1236    /// Returns [`SessionExitStatus::Running`] while the task is still
1237    /// alive, [`SessionExitStatus::Completed`] after a clean exit, or
1238    /// [`SessionExitStatus::Failed`] with the underlying error
1239    /// rendered to a string.
1240    ///
1241    /// Like [`Self::is_alive`], this is a cheap non-blocking read.
1242    #[must_use]
1243    pub fn exit_status(&self) -> SessionExitStatus {
1244        self.exit_rx.borrow().clone()
1245    }
1246
1247    /// Block until the session task transitions out of
1248    /// [`SessionExitStatus::Running`] and return the terminal status.
1249    ///
1250    /// Returns immediately if the task has already exited. Multiple
1251    /// concurrent callers are supported (each gets its own receiver
1252    /// clone), and the call does not consume the session.
1253    ///
1254    /// If the underlying watch sender is dropped without ever
1255    /// publishing a terminal state -- which should not happen in
1256    /// practice, but is treated defensively -- this returns the last
1257    /// observed value.
1258    pub async fn wait_for_exit(&self) -> SessionExitStatus {
1259        let mut rx = self.exit_rx.clone();
1260        loop {
1261            {
1262                let value = rx.borrow_and_update();
1263                if !matches!(*value, SessionExitStatus::Running) {
1264                    return value.clone();
1265                }
1266            }
1267            if rx.changed().await.is_err() {
1268                return rx.borrow().clone();
1269            }
1270        }
1271    }
1272
1273    /// Answer a deferred permission request from a different task.
1274    ///
1275    /// Use this after the [`PermissionHandler`] returned
1276    /// [`PermissionDecision::Defer`] for the matching `request_id`.
1277    /// Passing `decision = PermissionDecision::Defer` here is a
1278    /// no-op (logged at `warn`); pass `Allow` or `Deny`.
1279    ///
1280    /// Returns [`Error::DuplexClosed`] if the session task has
1281    /// already exited.
1282    ///
1283    /// # Example
1284    ///
1285    /// ```no_run
1286    /// use claude_wrapper::Claude;
1287    /// use claude_wrapper::duplex::{
1288    ///     DuplexOptions, DuplexSession, PermissionDecision, PermissionHandler,
1289    /// };
1290    /// use tokio::sync::mpsc;
1291    ///
1292    /// # async fn example() -> claude_wrapper::Result<()> {
1293    /// // Forward request_ids out to a UI thread; answer asynchronously.
1294    /// let (tx, _rx) = mpsc::unbounded_channel::<String>();
1295    /// let handler = PermissionHandler::new(move |req| {
1296    ///     let tx = tx.clone();
1297    ///     async move {
1298    ///         let _ = tx.send(req.request_id);
1299    ///         PermissionDecision::Defer
1300    ///     }
1301    /// });
1302    ///
1303    /// let claude = Claude::builder().build()?;
1304    /// let session = DuplexSession::spawn(
1305    ///     &claude,
1306    ///     DuplexOptions::default().on_permission(handler),
1307    /// ).await?;
1308    ///
1309    /// // ...later, from the UI thread:
1310    /// session.respond_to_permission(
1311    ///     "req-abc",
1312    ///     PermissionDecision::Allow { updated_input: None },
1313    /// )?;
1314    /// # Ok(())
1315    /// # }
1316    /// ```
1317    pub fn respond_to_permission(
1318        &self,
1319        request_id: impl Into<String>,
1320        decision: PermissionDecision,
1321    ) -> Result<()> {
1322        if matches!(decision, PermissionDecision::Defer) {
1323            warn!("respond_to_permission called with Defer; ignoring");
1324            return Ok(());
1325        }
1326        self.outbound_tx
1327            .send(OutboundMsg::PermissionResponse {
1328                request_id: request_id.into(),
1329                decision,
1330            })
1331            .map_err(|_| Error::DuplexClosed)?;
1332        Ok(())
1333    }
1334
1335    /// Send a clean interrupt to the CLI and wait for its
1336    /// acknowledgment.
1337    ///
1338    /// Writes a `control_request {subtype: "interrupt"}` and resolves
1339    /// when the matching `control_response` comes back. The
1340    /// in-flight turn (if any) closes shortly after with a truncated
1341    /// [`TurnResult`] -- the [`DuplexSession::send`] future for it
1342    /// resolves independently. Either ordering is possible; await
1343    /// both via `tokio::join!` if you care about both outcomes.
1344    ///
1345    /// Returns:
1346    /// - `Ok(())` when the CLI acknowledges with `subtype: "success"`.
1347    /// - [`Error::DuplexControlFailed`] when the CLI answers with an
1348    ///   error payload.
1349    /// - [`Error::DuplexClosed`] if the session task exited before
1350    ///   the response arrived.
1351    ///
1352    /// # Example
1353    ///
1354    /// ```no_run
1355    /// use std::time::Duration;
1356    /// use claude_wrapper::Claude;
1357    /// use claude_wrapper::duplex::{DuplexOptions, DuplexSession};
1358    ///
1359    /// # async fn example() -> claude_wrapper::Result<()> {
1360    /// let claude = Claude::builder().build()?;
1361    /// let session = DuplexSession::spawn(&claude, DuplexOptions::default()).await?;
1362    ///
1363    /// let send_fut = session.send("a question that triggers tool use");
1364    /// let interrupt_fut = async {
1365    ///     tokio::time::sleep(Duration::from_millis(250)).await;
1366    ///     session.interrupt().await
1367    /// };
1368    ///
1369    /// let (turn, interrupt) = tokio::join!(send_fut, interrupt_fut);
1370    /// let _truncated = turn?;
1371    /// interrupt?;
1372    /// # Ok(())
1373    /// # }
1374    /// ```
1375    pub async fn interrupt(&self) -> Result<()> {
1376        let (reply_tx, reply_rx) = oneshot::channel();
1377        self.outbound_tx
1378            .send(OutboundMsg::Interrupt { reply: reply_tx })
1379            .map_err(|_| Error::DuplexClosed)?;
1380        reply_rx.await.map_err(|_| Error::DuplexClosed)?
1381    }
1382
1383    /// Close the session and wait for the underlying task to exit.
1384    ///
1385    /// Drops the outbound channel sender, which the session task
1386    /// observes as `recv() -> None`, then closes stdin and reaps the
1387    /// child.
1388    pub async fn close(self) -> Result<()> {
1389        drop(self.outbound_tx);
1390        drop(self.events_tx);
1391        match self.join.await {
1392            Ok(result) => result,
1393            Err(e) if e.is_cancelled() => Ok(()),
1394            Err(e) => Err(Error::Io {
1395                message: format!("duplex session task panicked: {e}"),
1396                source: std::io::Error::other(e.to_string()),
1397                working_dir: None,
1398            }),
1399        }
1400    }
1401}
1402
1403/// Time budget for the graceful child shutdown after the run loop
1404/// exits. If the child is still alive after this deadline we SIGKILL
1405/// it so close() does not hang on a misbehaving subprocess.
1406const SHUTDOWN_BUDGET: Duration = Duration::from_secs(5);
1407
1408async fn run_session(
1409    mut child: Child,
1410    mut stdin: ChildStdin,
1411    stdout: ChildStdout,
1412    mut outbound_rx: mpsc::UnboundedReceiver<OutboundMsg>,
1413    events_tx: broadcast::Sender<InboundEvent>,
1414    permission_handler: Option<PermissionHandler>,
1415    exit_tx: watch::Sender<SessionExitStatus>,
1416) -> Result<()> {
1417    let mut lines = BufReader::new(stdout).lines();
1418    let mut pending: Option<(oneshot::Sender<Result<TurnResult>>, Vec<Value>)> = None;
1419    let mut pending_control: HashMap<String, oneshot::Sender<Result<()>>> = HashMap::new();
1420    let mut next_control_id: u64 = 0;
1421    let mut stream_err: Option<Error> = None;
1422
1423    loop {
1424        tokio::select! {
1425            biased;
1426
1427            line = lines.next_line() => match line {
1428                Ok(Some(l)) => {
1429                    if l.trim().is_empty() {
1430                        continue;
1431                    }
1432                    let parsed = match serde_json::from_str::<Value>(&l) {
1433                        Ok(v) => v,
1434                        Err(e) => {
1435                            debug!(line = %l, error = %e, "failed to parse duplex event, skipping");
1436                            continue;
1437                        }
1438                    };
1439                    match handle_inbound(parsed, &mut pending, &events_tx) {
1440                        InboundAction::None => {}
1441                        InboundAction::Permission(req) => {
1442                            let request_id = req.request_id.clone();
1443                            let decision = match permission_handler.as_ref() {
1444                                Some(h) => h.invoke(req).await,
1445                                None => {
1446                                    warn!(
1447                                        request_id = %request_id,
1448                                        "received can_use_tool with no permission handler; auto-denying"
1449                                    );
1450                                    PermissionDecision::Deny {
1451                                        message:
1452                                            "no permission handler configured on duplex session"
1453                                                .into(),
1454                                    }
1455                                }
1456                            };
1457                            if matches!(decision, PermissionDecision::Defer) {
1458                                debug!(
1459                                    request_id = %request_id,
1460                                    "permission handler deferred; waiting for respond_to_permission"
1461                                );
1462                            } else if let Err(e) =
1463                                write_permission_response(&mut stdin, &request_id, &decision).await
1464                            {
1465                                warn!(error = %e, "failed to write permission response");
1466                            }
1467                        }
1468                        InboundAction::ControlResponse { request_id, outcome } => {
1469                            if let Some(reply) = pending_control.remove(&request_id) {
1470                                let _ = reply.send(outcome);
1471                            } else {
1472                                debug!(
1473                                    request_id = %request_id,
1474                                    "received control_response with no pending request"
1475                                );
1476                            }
1477                        }
1478                    }
1479                }
1480                Ok(None) => break,
1481                Err(e) => {
1482                    stream_err = Some(Error::Io {
1483                        message: "failed to read duplex stdout".to_string(),
1484                        source: e,
1485                        working_dir: None,
1486                    });
1487                    break;
1488                }
1489            },
1490
1491            msg = outbound_rx.recv() => match msg {
1492                Some(OutboundMsg::Send { prompt, reply }) => {
1493                    if pending.is_some() {
1494                        let _ = reply.send(Err(Error::DuplexTurnInFlight));
1495                        continue;
1496                    }
1497                    if let Err(e) = write_user(&mut stdin, &prompt).await {
1498                        let _ = reply.send(Err(e));
1499                        continue;
1500                    }
1501                    pending = Some((reply, Vec::new()));
1502                }
1503                Some(OutboundMsg::PermissionResponse { request_id, decision }) => {
1504                    if let Err(e) =
1505                        write_permission_response(&mut stdin, &request_id, &decision).await
1506                    {
1507                        warn!(error = %e, "failed to write deferred permission response");
1508                    }
1509                }
1510                Some(OutboundMsg::Interrupt { reply }) => {
1511                    next_control_id += 1;
1512                    let request_id = format!("interrupt-{next_control_id}");
1513                    if let Err(e) =
1514                        write_control_request(&mut stdin, &request_id, "interrupt").await
1515                    {
1516                        let _ = reply.send(Err(e));
1517                        continue;
1518                    }
1519                    pending_control.insert(request_id, reply);
1520                }
1521                None => break,
1522            },
1523        }
1524    }
1525
1526    drop(stdin);
1527    match tokio::time::timeout(SHUTDOWN_BUDGET, child.wait()).await {
1528        Ok(Ok(_status)) => {}
1529        Ok(Err(e)) => {
1530            warn!(error = %e, "failed to wait for duplex child");
1531        }
1532        Err(_) => {
1533            warn!("duplex child did not exit within shutdown budget; killing");
1534            let _ = child.kill().await;
1535        }
1536    }
1537
1538    if let Some((reply, _)) = pending.take() {
1539        let _ = reply.send(Err(Error::DuplexClosed));
1540    }
1541    for (_, reply) in pending_control.drain() {
1542        let _ = reply.send(Err(Error::DuplexClosed));
1543    }
1544
1545    let result = match stream_err {
1546        Some(e) => Err(e),
1547        None => Ok(()),
1548    };
1549    let final_state = match &result {
1550        Ok(()) => SessionExitStatus::Completed,
1551        Err(e) => SessionExitStatus::Failed(e.to_string()),
1552    };
1553    let _ = exit_tx.send(final_state);
1554    result
1555}
1556
1557/// Action returned from [`handle_inbound`] for the run loop to act
1558/// on after the side-effects (broadcast, accumulate, resolve) are
1559/// done.
1560enum InboundAction {
1561    /// No further action -- side-effects were all handled inline.
1562    None,
1563    /// A `control_request {subtype: "can_use_tool"}` was received and
1564    /// needs the [`PermissionHandler`] invoked. The run loop awaits
1565    /// the handler and writes the response.
1566    Permission(PermissionRequest),
1567    /// A `control_response` matching one of our outbound
1568    /// `control_request`s arrived. The run loop matches `request_id`
1569    /// against its `pending_control` table and resolves the
1570    /// corresponding oneshot.
1571    ControlResponse {
1572        request_id: String,
1573        outcome: Result<()>,
1574    },
1575}
1576
1577fn handle_inbound(
1578    msg: Value,
1579    pending: &mut Option<(oneshot::Sender<Result<TurnResult>>, Vec<Value>)>,
1580    events_tx: &broadcast::Sender<InboundEvent>,
1581) -> InboundAction {
1582    match msg.get("type").and_then(Value::as_str) {
1583        Some("result") => {
1584            if let Some((reply, events)) = pending.take() {
1585                let _ = reply.send(Ok(TurnResult {
1586                    result: msg,
1587                    events,
1588                }));
1589            } else {
1590                debug!("dropping orphan result event with no pending turn");
1591            }
1592            InboundAction::None
1593        }
1594        Some("control_request") => {
1595            // can_use_tool flows through the permission handler;
1596            // anything else is logged + accumulated as Other for now.
1597            if msg
1598                .get("request")
1599                .and_then(|r| r.get("subtype"))
1600                .and_then(Value::as_str)
1601                == Some("can_use_tool")
1602                && let Some(req) = parse_permission_request(&msg)
1603            {
1604                if let Some((_, events)) = pending.as_mut() {
1605                    events.push(msg);
1606                }
1607                return InboundAction::Permission(req);
1608            }
1609            debug!(
1610                ?msg,
1611                "received unhandled control_request; treating as Other"
1612            );
1613            let _ = events_tx.send(InboundEvent::Other(msg.clone()));
1614            if let Some((_, events)) = pending.as_mut() {
1615                events.push(msg);
1616            }
1617            InboundAction::None
1618        }
1619        Some("control_response") => {
1620            if let Some((request_id, outcome)) = parse_control_response(&msg) {
1621                return InboundAction::ControlResponse {
1622                    request_id,
1623                    outcome,
1624                };
1625            }
1626            debug!(
1627                ?msg,
1628                "received malformed control_response; treating as Other"
1629            );
1630            let _ = events_tx.send(InboundEvent::Other(msg.clone()));
1631            if let Some((_, events)) = pending.as_mut() {
1632                events.push(msg);
1633            }
1634            InboundAction::None
1635        }
1636        _ => {
1637            // Broadcast a classified copy. Send error means no
1638            // subscribers, which is fine -- subscribers are optional.
1639            let _ = events_tx.send(classify(&msg));
1640
1641            if let Some((_, events)) = pending.as_mut() {
1642                events.push(msg);
1643            } else {
1644                debug!("dropping inbound event with no pending turn");
1645            }
1646            InboundAction::None
1647        }
1648    }
1649}
1650
1651fn parse_permission_request(msg: &Value) -> Option<PermissionRequest> {
1652    let request_id = msg.get("request_id").and_then(Value::as_str)?;
1653    let request = msg.get("request")?;
1654    let tool_name = request.get("tool_name").and_then(Value::as_str)?;
1655    let input = request.get("input").cloned().unwrap_or(Value::Null);
1656    Some(PermissionRequest {
1657        request_id: request_id.to_string(),
1658        tool_name: tool_name.to_string(),
1659        input,
1660        raw: request.clone(),
1661    })
1662}
1663
1664/// Pull `(request_id, outcome)` out of a `control_response` envelope.
1665///
1666/// Returns `None` if `request_id` is missing or the subtype is
1667/// unrecognised. `Some((id, Ok(())))` for `subtype: "success"`,
1668/// `Some((id, Err(DuplexControlFailed)))` for `subtype: "error"`.
1669fn parse_control_response(msg: &Value) -> Option<(String, Result<()>)> {
1670    let response = msg.get("response")?;
1671    let request_id = response.get("request_id").and_then(Value::as_str)?;
1672    let outcome = match response.get("subtype").and_then(Value::as_str) {
1673        Some("success") => Ok(()),
1674        Some("error") => {
1675            let message = response
1676                .get("error")
1677                .and_then(Value::as_str)
1678                .unwrap_or("unknown control_response error")
1679                .to_string();
1680            Err(Error::DuplexControlFailed { message })
1681        }
1682        _ => return None,
1683    };
1684    Some((request_id.to_string(), outcome))
1685}
1686
1687async fn write_user(stdin: &mut ChildStdin, prompt: &str) -> Result<()> {
1688    let user_msg = serde_json::json!({
1689        "type": "user",
1690        "message": {
1691            "role": "user",
1692            "content": prompt,
1693        },
1694        "parent_tool_use_id": null,
1695    });
1696    write_line(stdin, &user_msg, "user message").await
1697}
1698
1699async fn write_control_request(
1700    stdin: &mut ChildStdin,
1701    request_id: &str,
1702    subtype: &str,
1703) -> Result<()> {
1704    let envelope = serde_json::json!({
1705        "type": "control_request",
1706        "request_id": request_id,
1707        "request": { "subtype": subtype },
1708    });
1709    write_line(stdin, &envelope, "control_request").await
1710}
1711
1712async fn write_permission_response(
1713    stdin: &mut ChildStdin,
1714    request_id: &str,
1715    decision: &PermissionDecision,
1716) -> Result<()> {
1717    let inner = match decision {
1718        PermissionDecision::Allow { updated_input } => {
1719            let mut obj = serde_json::Map::new();
1720            obj.insert("behavior".to_string(), Value::String("allow".to_string()));
1721            if let Some(input) = updated_input {
1722                obj.insert("updatedInput".to_string(), input.clone());
1723            }
1724            Value::Object(obj)
1725        }
1726        PermissionDecision::Deny { message } => serde_json::json!({
1727            "behavior": "deny",
1728            "message": message,
1729        }),
1730        PermissionDecision::Defer => {
1731            // Caller path is supposed to filter this; defensive guard.
1732            return Ok(());
1733        }
1734    };
1735    let envelope = serde_json::json!({
1736        "type": "control_response",
1737        "response": {
1738            "request_id": request_id,
1739            "subtype": "success",
1740            "response": inner,
1741        },
1742    });
1743    write_line(stdin, &envelope, "control_response").await
1744}
1745
1746async fn write_line(stdin: &mut ChildStdin, value: &Value, what: &'static str) -> Result<()> {
1747    let mut line = serde_json::to_string(value).map_err(|e| Error::Json {
1748        message: format!("failed to serialize duplex {what}"),
1749        source: e,
1750    })?;
1751    line.push('\n');
1752    stdin
1753        .write_all(line.as_bytes())
1754        .await
1755        .map_err(|e| Error::Io {
1756            message: format!("failed to write {what} to duplex stdin"),
1757            source: e,
1758            working_dir: None,
1759        })?;
1760    stdin.flush().await.map_err(|e| Error::Io {
1761        message: "failed to flush duplex stdin".to_string(),
1762        source: e,
1763        working_dir: None,
1764    })?;
1765    Ok(())
1766}
1767
1768#[cfg(test)]
1769mod tests {
1770    use super::*;
1771    use serde_json::json;
1772
1773    #[test]
1774    fn into_args_default_includes_required_flags() {
1775        let args = DuplexOptions::default().into_args();
1776        assert!(args.contains(&"--print".to_string()));
1777        assert!(args.contains(&"--verbose".to_string()));
1778        assert!(
1779            args.windows(2)
1780                .any(|w| w == ["--output-format", "stream-json"])
1781        );
1782        assert!(
1783            args.windows(2)
1784                .any(|w| w == ["--input-format", "stream-json"])
1785        );
1786    }
1787
1788    #[test]
1789    fn into_args_includes_model() {
1790        let args = DuplexOptions::default().model("haiku").into_args();
1791        assert!(args.windows(2).any(|w| w == ["--model", "haiku"]));
1792    }
1793
1794    #[test]
1795    fn into_args_includes_system_prompts() {
1796        let args = DuplexOptions::default()
1797            .system_prompt("be concise")
1798            .append_system_prompt("also polite")
1799            .into_args();
1800        assert!(
1801            args.windows(2)
1802                .any(|w| w == ["--system-prompt", "be concise"])
1803        );
1804        assert!(
1805            args.windows(2)
1806                .any(|w| w == ["--append-system-prompt", "also polite"])
1807        );
1808    }
1809
1810    #[test]
1811    fn into_args_appends_raw_args_last() {
1812        let args = DuplexOptions::default()
1813            .arg("--add-dir")
1814            .arg("/tmp/foo")
1815            .into_args();
1816        // Last two entries should be the additional args, in order.
1817        assert_eq!(&args[args.len() - 2..], &["--add-dir", "/tmp/foo"]);
1818    }
1819
1820    #[test]
1821    fn into_args_includes_resume_when_set() {
1822        let args = DuplexOptions::default().resume("abc-123").into_args();
1823        assert!(args.windows(2).any(|w| w == ["--resume", "abc-123"]));
1824    }
1825
1826    #[test]
1827    fn into_args_omits_resume_by_default() {
1828        let args = DuplexOptions::default().into_args();
1829        assert!(
1830            !args.iter().any(|a| a == "--resume"),
1831            "--resume should not appear without an explicit resume(...) call; got {args:?}"
1832        );
1833    }
1834
1835    #[test]
1836    fn into_args_includes_continue_when_set() {
1837        let args = DuplexOptions::default().continue_session().into_args();
1838        assert!(args.iter().any(|a| a == "--continue"));
1839    }
1840
1841    #[test]
1842    fn into_args_omits_continue_by_default() {
1843        let args = DuplexOptions::default().into_args();
1844        assert!(!args.iter().any(|a| a == "--continue"));
1845    }
1846
1847    #[test]
1848    fn into_args_includes_worktree_flag_without_name() {
1849        let args = DuplexOptions::default().worktree(None::<&str>).into_args();
1850        assert!(args.iter().any(|a| a == "--worktree"));
1851        // No name means no positional follows --worktree.
1852        let pos = args.iter().position(|a| a == "--worktree").unwrap();
1853        assert!(
1854            args.get(pos + 1).is_none_or(|a| a.starts_with("--")),
1855            "--worktree without a name should not be followed by a positional; got {args:?}"
1856        );
1857    }
1858
1859    #[test]
1860    fn into_args_includes_worktree_flag_with_name() {
1861        let args = DuplexOptions::default()
1862            .worktree(Some("agent-xyz"))
1863            .into_args();
1864        let pos = args.iter().position(|a| a == "--worktree").unwrap();
1865        assert_eq!(args.get(pos + 1).map(String::as_str), Some("agent-xyz"));
1866    }
1867
1868    #[test]
1869    fn into_args_omits_worktree_by_default() {
1870        let args = DuplexOptions::default().into_args();
1871        assert!(
1872            !args.iter().any(|a| a == "--worktree"),
1873            "--worktree should not appear without an explicit worktree(...) call; got {args:?}"
1874        );
1875    }
1876
1877    #[test]
1878    fn worktree_lands_before_additional_args() {
1879        // Same `--` ordering bug class as resume.
1880        let args = DuplexOptions::default()
1881            .worktree(Some("foo"))
1882            .arg("--")
1883            .arg("trailing")
1884            .into_args();
1885        let wt_pos = args.iter().position(|a| a == "--worktree").unwrap();
1886        let dash_dash_pos = args.iter().position(|a| a == "--").unwrap();
1887        assert!(
1888            wt_pos < dash_dash_pos,
1889            "--worktree must precede `--` separator; got {args:?}"
1890        );
1891    }
1892
1893    #[test]
1894    fn into_args_includes_agent_when_set() {
1895        let args = DuplexOptions::default().agent("rust-qa").into_args();
1896        assert!(
1897            args.windows(2).any(|w| w == ["--agent", "rust-qa"]),
1898            "missing --agent rust-qa in {args:?}"
1899        );
1900    }
1901
1902    #[test]
1903    fn into_args_omits_agent_by_default() {
1904        let args = DuplexOptions::default().into_args();
1905        assert!(
1906            !args.iter().any(|a| a == "--agent"),
1907            "--agent should not appear without an explicit agent(...) call; got {args:?}"
1908        );
1909    }
1910
1911    #[test]
1912    fn into_args_includes_agents_json_when_set() {
1913        let json = r#"{"reviewer":{"description":"r","prompt":"p"}}"#;
1914        let args = DuplexOptions::default().agents_json(json).into_args();
1915        let pos = args.iter().position(|a| a == "--agents").unwrap();
1916        assert_eq!(args.get(pos + 1).map(String::as_str), Some(json));
1917    }
1918
1919    #[test]
1920    fn into_args_omits_agents_json_by_default() {
1921        let args = DuplexOptions::default().into_args();
1922        assert!(!args.iter().any(|a| a == "--agents"));
1923    }
1924
1925    #[test]
1926    fn agent_and_agents_json_compose() {
1927        let json = r#"{"reviewer":{"description":"r","prompt":"p"}}"#;
1928        let args = DuplexOptions::default()
1929            .agents_json(json)
1930            .agent("reviewer")
1931            .into_args();
1932        // Both flags present.
1933        assert!(args.iter().any(|a| a == "--agents"));
1934        assert!(args.iter().any(|a| a == "--agent"));
1935    }
1936
1937    #[test]
1938    fn agent_lands_before_additional_args() {
1939        let args = DuplexOptions::default()
1940            .agent("rust-qa")
1941            .arg("--")
1942            .arg("trailing")
1943            .into_args();
1944        let agent_pos = args.iter().position(|a| a == "--agent").unwrap();
1945        let dash_dash_pos = args.iter().position(|a| a == "--").unwrap();
1946        assert!(
1947            agent_pos < dash_dash_pos,
1948            "--agent must precede `--` separator; got {args:?}"
1949        );
1950    }
1951
1952    #[test]
1953    fn agents_json_lands_before_additional_args() {
1954        let args = DuplexOptions::default()
1955            .agents_json("{}")
1956            .arg("--")
1957            .arg("trailing")
1958            .into_args();
1959        let agents_pos = args.iter().position(|a| a == "--agents").unwrap();
1960        let dash_dash_pos = args.iter().position(|a| a == "--").unwrap();
1961        assert!(
1962            agents_pos < dash_dash_pos,
1963            "--agents must precede `--` separator; got {args:?}"
1964        );
1965    }
1966
1967    // -- QueryCommand knob-set parity (#672) -------------------------
1968
1969    #[test]
1970    fn into_args_includes_session_id() {
1971        let args = DuplexOptions::default().session_id("sid-9").into_args();
1972        assert!(args.windows(2).any(|w| w == ["--session-id", "sid-9"]));
1973    }
1974
1975    #[test]
1976    fn into_args_includes_setting_sources() {
1977        let args = DuplexOptions::default()
1978            .setting_sources("user,project")
1979            .into_args();
1980        assert!(
1981            args.windows(2)
1982                .any(|w| w == ["--setting-sources", "user,project"]),
1983            "got {args:?}"
1984        );
1985    }
1986
1987    #[test]
1988    fn into_args_omits_setting_sources_by_default() {
1989        let args = DuplexOptions::default().into_args();
1990        assert!(!args.iter().any(|a| a == "--setting-sources"));
1991    }
1992
1993    #[test]
1994    fn into_args_hermetic_emits_full_seal() {
1995        let args = DuplexOptions::default().hermetic().into_args();
1996        assert!(
1997            args.windows(2)
1998                .any(|w| w[0] == "--setting-sources" && w[1].is_empty()),
1999            "got {args:?}"
2000        );
2001        assert!(args.iter().any(|a| a == "--strict-mcp-config"));
2002        assert!(
2003            args.iter()
2004                .any(|a| a == "--exclude-dynamic-system-prompt-sections")
2005        );
2006        // A hermetic seal must never imply --bare.
2007        assert!(!args.iter().any(|a| a == "--bare"));
2008    }
2009
2010    #[test]
2011    fn into_args_hermetic_scoped_project_keeps_user() {
2012        let args = DuplexOptions::default()
2013            .hermetic_scoped(HermeticScope::Project)
2014            .into_args();
2015        assert!(args.windows(2).any(|w| w == ["--setting-sources", "user"]));
2016        assert!(args.iter().any(|a| a == "--strict-mcp-config"));
2017    }
2018
2019    #[test]
2020    fn into_args_includes_json_schema() {
2021        let schema = r#"{"type":"object"}"#;
2022        let args = DuplexOptions::default().json_schema(schema).into_args();
2023        assert!(args.windows(2).any(|w| w == ["--json-schema", schema]));
2024    }
2025
2026    #[test]
2027    fn into_args_joins_allowed_tools_comma_separated() {
2028        let args = DuplexOptions::default()
2029            .allowed_tools(["Read", "Bash(git log:*)"])
2030            .allowed_tool("Write")
2031            .into_args();
2032        assert!(
2033            args.windows(2)
2034                .any(|w| w == ["--allowed-tools", "Read,Bash(git log:*),Write"]),
2035            "missing joined --allowed-tools in {args:?}"
2036        );
2037    }
2038
2039    #[test]
2040    fn into_args_joins_disallowed_tools_comma_separated() {
2041        let args = DuplexOptions::default()
2042            .disallowed_tools(["WebSearch"])
2043            .disallowed_tool("WebFetch")
2044            .into_args();
2045        assert!(
2046            args.windows(2)
2047                .any(|w| w == ["--disallowed-tools", "WebSearch,WebFetch"]),
2048            "missing joined --disallowed-tools in {args:?}"
2049        );
2050    }
2051
2052    #[test]
2053    fn into_args_includes_caps() {
2054        let args = DuplexOptions::default()
2055            .max_turns(4)
2056            .max_budget_usd(0.25)
2057            .into_args();
2058        assert!(args.windows(2).any(|w| w == ["--max-turns", "4"]));
2059        assert!(args.windows(2).any(|w| w == ["--max-budget-usd", "0.25"]));
2060    }
2061
2062    #[test]
2063    fn into_args_includes_fallback_model_and_effort() {
2064        let args = DuplexOptions::default()
2065            .fallback_model("haiku")
2066            .effort(Effort::Low)
2067            .into_args();
2068        assert!(args.windows(2).any(|w| w == ["--fallback-model", "haiku"]));
2069        assert!(args.windows(2).any(|w| w == ["--effort", "low"]));
2070    }
2071
2072    #[test]
2073    fn into_args_repeats_add_dir_and_mcp_config() {
2074        let args = DuplexOptions::default()
2075            .add_dir("/a")
2076            .add_dir("/b")
2077            .mcp_config("x.json")
2078            .strict_mcp_config()
2079            .into_args();
2080        assert!(args.windows(2).any(|w| w == ["--add-dir", "/a"]));
2081        assert!(args.windows(2).any(|w| w == ["--add-dir", "/b"]));
2082        assert!(args.windows(2).any(|w| w == ["--mcp-config", "x.json"]));
2083        assert!(args.iter().any(|a| a == "--strict-mcp-config"));
2084    }
2085
2086    #[test]
2087    fn into_args_includes_no_session_persistence() {
2088        let args = DuplexOptions::default()
2089            .no_session_persistence()
2090            .into_args();
2091        assert!(args.iter().any(|a| a == "--no-session-persistence"));
2092    }
2093
2094    // ─── #690: parity builders promoted from QueryCommand ───
2095
2096    #[test]
2097    fn into_args_joins_tools_comma_separated() {
2098        let args = DuplexOptions::default()
2099            .tools(["Bash", "Read", "Edit"])
2100            .into_args();
2101        assert!(
2102            args.windows(2).any(|w| w == ["--tools", "Bash,Read,Edit"]),
2103            "missing joined --tools in {args:?}"
2104        );
2105    }
2106
2107    #[test]
2108    fn into_args_repeats_file_per_spec() {
2109        let args = DuplexOptions::default()
2110            .file("file_a:doc.txt")
2111            .file("file_b:notes.md")
2112            .into_args();
2113        assert_eq!(args.iter().filter(|a| *a == "--file").count(), 2);
2114        assert!(args.iter().any(|a| a == "file_a:doc.txt"));
2115        assert!(args.iter().any(|a| a == "file_b:notes.md"));
2116    }
2117
2118    #[test]
2119    fn into_args_includes_settings() {
2120        let args = DuplexOptions::default()
2121            .settings("/tmp/settings.json")
2122            .into_args();
2123        assert!(
2124            args.windows(2)
2125                .any(|w| w == ["--settings", "/tmp/settings.json"])
2126        );
2127    }
2128
2129    #[test]
2130    fn into_args_includes_fork_session() {
2131        let args = DuplexOptions::default().fork_session().into_args();
2132        assert!(args.iter().any(|a| a == "--fork-session"));
2133    }
2134
2135    #[test]
2136    fn into_args_includes_debug_filter_and_file() {
2137        let args = DuplexOptions::default()
2138            .debug_filter("api,hooks")
2139            .debug_file("/tmp/debug.log")
2140            .into_args();
2141        assert!(args.windows(2).any(|w| w == ["--debug", "api,hooks"]));
2142        assert!(
2143            args.windows(2)
2144                .any(|w| w == ["--debug-file", "/tmp/debug.log"])
2145        );
2146    }
2147
2148    #[test]
2149    fn into_args_includes_betas() {
2150        let args = DuplexOptions::default().betas("feature-x").into_args();
2151        assert!(args.windows(2).any(|w| w == ["--betas", "feature-x"]));
2152    }
2153
2154    #[test]
2155    fn into_args_repeats_plugin_dir_and_url() {
2156        let args = DuplexOptions::default()
2157            .plugin_dir("/plugins/a")
2158            .plugin_dir("/plugins/b")
2159            .plugin_url("https://example.com/p.zip")
2160            .into_args();
2161        assert_eq!(args.iter().filter(|a| *a == "--plugin-dir").count(), 2);
2162        assert!(
2163            args.windows(2)
2164                .any(|w| w == ["--plugin-url", "https://example.com/p.zip"])
2165        );
2166    }
2167
2168    #[test]
2169    fn into_args_includes_bare_family_bool_flags() {
2170        let args = DuplexOptions::default()
2171            .tmux()
2172            .bare()
2173            .safe_mode()
2174            .disable_slash_commands()
2175            .include_hook_events()
2176            .exclude_dynamic_system_prompt_sections()
2177            .into_args();
2178        for flag in [
2179            "--tmux",
2180            "--bare",
2181            "--safe-mode",
2182            "--disable-slash-commands",
2183            "--include-hook-events",
2184            "--exclude-dynamic-system-prompt-sections",
2185        ] {
2186            assert!(args.iter().any(|a| a == flag), "missing {flag} in {args:?}");
2187        }
2188    }
2189
2190    #[test]
2191    fn into_args_includes_name() {
2192        let args = DuplexOptions::default().name("my session").into_args();
2193        assert!(args.windows(2).any(|w| w == ["--name", "my session"]));
2194    }
2195
2196    #[test]
2197    fn into_args_omits_promoted_parity_flags_by_default() {
2198        let args = DuplexOptions::default().into_args();
2199        for flag in [
2200            "--tools",
2201            "--file",
2202            "--settings",
2203            "--fork-session",
2204            "--debug",
2205            "--debug-file",
2206            "--betas",
2207            "--plugin-dir",
2208            "--plugin-url",
2209            "--tmux",
2210            "--bare",
2211            "--safe-mode",
2212            "--disable-slash-commands",
2213            "--include-hook-events",
2214            "--exclude-dynamic-system-prompt-sections",
2215            "--name",
2216        ] {
2217            assert!(
2218                !args.iter().any(|a| a == flag),
2219                "{flag} should be absent by default; got {args:?}"
2220            );
2221        }
2222    }
2223
2224    #[test]
2225    fn parity_flags_land_before_additional_args() {
2226        // Same `--` ordering bug class as resume/agent.
2227        let args = DuplexOptions::default()
2228            .max_turns(2)
2229            .json_schema("{}")
2230            .arg("--")
2231            .arg("trailing")
2232            .into_args();
2233        let dash_dash_pos = args.iter().position(|a| a == "--").unwrap();
2234        for flag in ["--max-turns", "--json-schema"] {
2235            let pos = args.iter().position(|a| a == flag).unwrap();
2236            assert!(
2237                pos < dash_dash_pos,
2238                "{flag} must precede `--` separator; got {args:?}"
2239            );
2240        }
2241    }
2242
2243    #[test]
2244    fn into_args_omits_parity_flags_by_default() {
2245        let args = DuplexOptions::default().into_args();
2246        for flag in [
2247            "--session-id",
2248            "--json-schema",
2249            "--allowed-tools",
2250            "--disallowed-tools",
2251            "--max-turns",
2252            "--max-budget-usd",
2253            "--fallback-model",
2254            "--effort",
2255            "--add-dir",
2256            "--mcp-config",
2257            "--strict-mcp-config",
2258            "--no-session-persistence",
2259        ] {
2260            assert!(
2261                !args.iter().any(|a| a == flag),
2262                "{flag} should not appear by default; got {args:?}"
2263            );
2264        }
2265    }
2266
2267    #[test]
2268    fn resume_lands_before_additional_args() {
2269        // Catches the same class of bug as QueryCommand::execute_json
2270        // had: a flag appended after the user-supplied raw args (which
2271        // typically include `--`) gets eaten as a positional. Resume
2272        // must precede any caller-injected `arg(...)`.
2273        let args = DuplexOptions::default()
2274            .resume("xyz")
2275            .arg("--")
2276            .arg("trailing")
2277            .into_args();
2278        let resume_pos = args.iter().position(|a| a == "--resume").unwrap();
2279        let dash_dash_pos = args.iter().position(|a| a == "--").unwrap();
2280        assert!(
2281            resume_pos < dash_dash_pos,
2282            "--resume must precede `--` separator; got {args:?}"
2283        );
2284    }
2285
2286    #[test]
2287    fn turn_result_accessors_pull_from_result() {
2288        let r = TurnResult {
2289            result: json!({
2290                "type": "result",
2291                "result": "hello",
2292                "session_id": "sess-123",
2293                "total_cost_usd": 0.0042,
2294                "duration_ms": 1234_u64,
2295            }),
2296            events: vec![],
2297        };
2298        assert_eq!(r.result_text(), Some("hello"));
2299        assert_eq!(r.session_id(), Some("sess-123"));
2300        assert_eq!(r.total_cost_usd(), Some(0.0042));
2301        assert_eq!(r.duration_ms(), Some(1234));
2302    }
2303
2304    #[test]
2305    fn turn_result_total_cost_falls_back_to_legacy_field() {
2306        let r = TurnResult {
2307            result: json!({ "cost_usd": 0.5 }),
2308            events: vec![],
2309        };
2310        assert_eq!(r.total_cost_usd(), Some(0.5));
2311    }
2312
2313    #[test]
2314    fn turn_result_accessors_return_none_when_missing() {
2315        let r = TurnResult {
2316            result: json!({}),
2317            events: vec![],
2318        };
2319        assert_eq!(r.result_text(), None);
2320        assert_eq!(r.session_id(), None);
2321        assert_eq!(r.total_cost_usd(), None);
2322        assert_eq!(r.duration_ms(), None);
2323    }
2324
2325    #[test]
2326    fn handle_inbound_appends_non_result_to_pending_events() {
2327        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2328        let (events_tx, _events_rx) = broadcast::channel(16);
2329        let mut pending = Some((tx, Vec::new()));
2330        handle_inbound(
2331            json!({ "type": "assistant", "message": {} }),
2332            &mut pending,
2333            &events_tx,
2334        );
2335        let (_, events) = pending.as_ref().unwrap();
2336        assert_eq!(events.len(), 1);
2337        assert_eq!(
2338            events[0].get("type").and_then(Value::as_str),
2339            Some("assistant")
2340        );
2341    }
2342
2343    #[test]
2344    fn handle_inbound_resolves_pending_on_result() {
2345        let (tx, rx) = oneshot::channel::<Result<TurnResult>>();
2346        let (events_tx, _events_rx) = broadcast::channel(16);
2347        let mut pending = Some((tx, vec![json!({ "type": "assistant" })]));
2348        handle_inbound(
2349            json!({ "type": "result", "result": "ok" }),
2350            &mut pending,
2351            &events_tx,
2352        );
2353        assert!(pending.is_none());
2354        let received = rx.blocking_recv().unwrap().unwrap();
2355        assert_eq!(received.result_text(), Some("ok"));
2356        assert_eq!(received.events.len(), 1);
2357    }
2358
2359    #[test]
2360    fn handle_inbound_drops_orphans_without_pending_turn() {
2361        let (events_tx, _events_rx) = broadcast::channel(16);
2362        let mut pending: Option<(oneshot::Sender<Result<TurnResult>>, Vec<Value>)> = None;
2363        handle_inbound(json!({ "type": "assistant" }), &mut pending, &events_tx);
2364        handle_inbound(
2365            json!({ "type": "result", "result": "ok" }),
2366            &mut pending,
2367            &events_tx,
2368        );
2369        assert!(pending.is_none());
2370    }
2371
2372    #[test]
2373    fn handle_inbound_broadcasts_classified_event() {
2374        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2375        let (events_tx, mut events_rx) = broadcast::channel(16);
2376        let mut pending = Some((tx, Vec::new()));
2377        handle_inbound(
2378            json!({ "type": "assistant", "message": { "role": "assistant" } }),
2379            &mut pending,
2380            &events_tx,
2381        );
2382        let event = events_rx.try_recv().expect("classified event broadcast");
2383        assert!(matches!(event, InboundEvent::Assistant(_)));
2384    }
2385
2386    #[test]
2387    fn handle_inbound_does_not_broadcast_result() {
2388        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2389        let (events_tx, mut events_rx) = broadcast::channel(16);
2390        let mut pending = Some((tx, Vec::new()));
2391        handle_inbound(
2392            json!({ "type": "result", "result": "ok" }),
2393            &mut pending,
2394            &events_tx,
2395        );
2396        // Result is not broadcast -- it lands in TurnResult.result.
2397        assert!(events_rx.try_recv().is_err());
2398    }
2399
2400    #[test]
2401    fn classify_system_init_pulls_session_id() {
2402        let v = json!({
2403            "type": "system",
2404            "subtype": "init",
2405            "session_id": "sess-abc",
2406        });
2407        match classify(&v) {
2408            InboundEvent::SystemInit { session_id } => assert_eq!(session_id, "sess-abc"),
2409            other => panic!("expected SystemInit, got {other:?}"),
2410        }
2411    }
2412
2413    #[test]
2414    fn classify_system_without_init_subtype_is_other() {
2415        let v = json!({ "type": "system", "subtype": "compaction" });
2416        assert!(matches!(classify(&v), InboundEvent::Other(_)));
2417    }
2418
2419    #[test]
2420    fn classify_system_init_without_session_id_is_other() {
2421        let v = json!({ "type": "system", "subtype": "init" });
2422        assert!(matches!(classify(&v), InboundEvent::Other(_)));
2423    }
2424
2425    #[test]
2426    fn classify_assistant_stream_event_user() {
2427        assert!(matches!(
2428            classify(&json!({ "type": "assistant" })),
2429            InboundEvent::Assistant(_)
2430        ));
2431        assert!(matches!(
2432            classify(&json!({ "type": "stream_event" })),
2433            InboundEvent::StreamEvent(_)
2434        ));
2435        assert!(matches!(
2436            classify(&json!({ "type": "user" })),
2437            InboundEvent::User(_)
2438        ));
2439    }
2440
2441    #[test]
2442    fn classify_unknown_type_is_other() {
2443        assert!(matches!(
2444            classify(&json!({ "type": "control_request" })),
2445            InboundEvent::Other(_)
2446        ));
2447        assert!(matches!(
2448            classify(&json!({ "type": "future_thing" })),
2449            InboundEvent::Other(_)
2450        ));
2451        assert!(matches!(classify(&json!({})), InboundEvent::Other(_)));
2452    }
2453
2454    #[test]
2455    fn into_args_does_not_emit_subscriber_capacity_flag() {
2456        // subscriber_capacity is runtime config, not a CLI arg.
2457        let args = DuplexOptions::default().subscriber_capacity(64).into_args();
2458        assert!(!args.iter().any(|a| a.contains("subscriber")));
2459        assert!(!args.iter().any(|a| a.contains("capacity")));
2460    }
2461
2462    #[test]
2463    fn into_args_includes_permission_prompt_tool_when_handler_set() {
2464        let handler = PermissionHandler::new(|_req| async move {
2465            PermissionDecision::Allow {
2466                updated_input: None,
2467            }
2468        });
2469        let args = DuplexOptions::default().on_permission(handler).into_args();
2470        assert!(
2471            args.windows(2)
2472                .any(|w| w == ["--permission-prompt-tool", "stdio"])
2473        );
2474    }
2475
2476    #[test]
2477    fn into_args_omits_permission_prompt_tool_without_handler() {
2478        let args = DuplexOptions::default().into_args();
2479        assert!(!args.iter().any(|a| a == "--permission-prompt-tool"));
2480    }
2481
2482    #[test]
2483    fn into_args_emits_permission_mode_flag() {
2484        let args = DuplexOptions::default()
2485            .permission_mode(PermissionMode::AcceptEdits)
2486            .into_args();
2487        assert!(
2488            args.windows(2)
2489                .any(|w| w == ["--permission-mode", "acceptEdits"]),
2490            "missing --permission-mode acceptEdits in {args:?}"
2491        );
2492    }
2493
2494    #[test]
2495    fn into_args_emits_plan_mode() {
2496        let args = DuplexOptions::default()
2497            .permission_mode(PermissionMode::Plan)
2498            .into_args();
2499        assert!(args.windows(2).any(|w| w == ["--permission-mode", "plan"]));
2500    }
2501
2502    #[test]
2503    fn into_args_omits_permission_mode_by_default() {
2504        let args = DuplexOptions::default().into_args();
2505        assert!(!args.iter().any(|a| a == "--permission-mode"));
2506    }
2507
2508    #[test]
2509    fn into_args_emits_dangerously_skip_permissions_flag() {
2510        let args = DuplexOptions::default()
2511            .dangerously_skip_permissions()
2512            .into_args();
2513        assert!(args.iter().any(|a| a == "--dangerously-skip-permissions"));
2514    }
2515
2516    #[test]
2517    fn into_args_omits_dangerously_skip_by_default() {
2518        let args = DuplexOptions::default().into_args();
2519        assert!(!args.iter().any(|a| a == "--dangerously-skip-permissions"));
2520    }
2521
2522    #[test]
2523    fn parse_permission_request_extracts_fields() {
2524        let msg = json!({
2525            "type": "control_request",
2526            "request_id": "req-1",
2527            "request": {
2528                "subtype": "can_use_tool",
2529                "tool_name": "Bash",
2530                "input": { "command": "ls" }
2531            }
2532        });
2533        let req = parse_permission_request(&msg).expect("permission request");
2534        assert_eq!(req.request_id, "req-1");
2535        assert_eq!(req.tool_name, "Bash");
2536        assert_eq!(req.input, json!({ "command": "ls" }));
2537        assert_eq!(
2538            req.raw.get("subtype").and_then(Value::as_str),
2539            Some("can_use_tool")
2540        );
2541    }
2542
2543    #[test]
2544    fn parse_permission_request_returns_none_when_missing_request_id() {
2545        let msg = json!({
2546            "type": "control_request",
2547            "request": {
2548                "subtype": "can_use_tool",
2549                "tool_name": "Bash",
2550            }
2551        });
2552        assert!(parse_permission_request(&msg).is_none());
2553    }
2554
2555    #[test]
2556    fn parse_permission_request_returns_none_when_missing_tool_name() {
2557        let msg = json!({
2558            "type": "control_request",
2559            "request_id": "req-1",
2560            "request": { "subtype": "can_use_tool" }
2561        });
2562        assert!(parse_permission_request(&msg).is_none());
2563    }
2564
2565    #[test]
2566    fn parse_permission_request_handles_missing_input() {
2567        let msg = json!({
2568            "type": "control_request",
2569            "request_id": "req-1",
2570            "request": {
2571                "subtype": "can_use_tool",
2572                "tool_name": "Bash",
2573            }
2574        });
2575        let req = parse_permission_request(&msg).expect("request");
2576        assert_eq!(req.input, Value::Null);
2577    }
2578
2579    #[test]
2580    fn handle_inbound_returns_permission_for_can_use_tool() {
2581        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2582        let (events_tx, _events_rx) = broadcast::channel(16);
2583        let mut pending = Some((tx, Vec::new()));
2584        let action = handle_inbound(
2585            json!({
2586                "type": "control_request",
2587                "request_id": "req-1",
2588                "request": {
2589                    "subtype": "can_use_tool",
2590                    "tool_name": "Bash",
2591                    "input": { "command": "ls" }
2592                }
2593            }),
2594            &mut pending,
2595            &events_tx,
2596        );
2597        match action {
2598            InboundAction::Permission(req) => {
2599                assert_eq!(req.request_id, "req-1");
2600                assert_eq!(req.tool_name, "Bash");
2601            }
2602            InboundAction::None | InboundAction::ControlResponse { .. } => {
2603                panic!("expected Permission action");
2604            }
2605        }
2606        // Event should also be accumulated in the pending turn.
2607        let (_, events) = pending.as_ref().unwrap();
2608        assert_eq!(events.len(), 1);
2609    }
2610
2611    #[test]
2612    fn handle_inbound_treats_unknown_control_request_as_other() {
2613        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2614        let (events_tx, mut events_rx) = broadcast::channel(16);
2615        let mut pending = Some((tx, Vec::new()));
2616        let action = handle_inbound(
2617            json!({
2618                "type": "control_request",
2619                "request_id": "req-2",
2620                "request": { "subtype": "future_subtype" }
2621            }),
2622            &mut pending,
2623            &events_tx,
2624        );
2625        assert!(matches!(action, InboundAction::None));
2626        let event = events_rx.try_recv().expect("broadcast");
2627        assert!(matches!(event, InboundEvent::Other(_)));
2628    }
2629
2630    #[tokio::test]
2631    async fn permission_handler_invokes_closure_async() {
2632        let handler = PermissionHandler::new(|req| async move {
2633            if req.tool_name == "Bash" {
2634                PermissionDecision::Deny {
2635                    message: "no bash".into(),
2636                }
2637            } else {
2638                PermissionDecision::Allow {
2639                    updated_input: None,
2640                }
2641            }
2642        });
2643        let req = PermissionRequest {
2644            request_id: "r1".into(),
2645            tool_name: "Bash".into(),
2646            input: Value::Null,
2647            raw: Value::Null,
2648        };
2649        match handler.invoke(req).await {
2650            PermissionDecision::Deny { message } => assert_eq!(message, "no bash"),
2651            other => panic!("expected Deny, got {other:?}"),
2652        }
2653    }
2654
2655    #[test]
2656    fn parse_control_response_extracts_success() {
2657        let msg = json!({
2658            "type": "control_response",
2659            "response": {
2660                "request_id": "interrupt-1",
2661                "subtype": "success",
2662                "response": {}
2663            }
2664        });
2665        let (id, outcome) = parse_control_response(&msg).expect("parsed");
2666        assert_eq!(id, "interrupt-1");
2667        assert!(outcome.is_ok());
2668    }
2669
2670    #[test]
2671    fn parse_control_response_extracts_error_with_message() {
2672        let msg = json!({
2673            "type": "control_response",
2674            "response": {
2675                "request_id": "interrupt-2",
2676                "subtype": "error",
2677                "error": "no turn in flight"
2678            }
2679        });
2680        let (id, outcome) = parse_control_response(&msg).expect("parsed");
2681        assert_eq!(id, "interrupt-2");
2682        match outcome {
2683            Err(Error::DuplexControlFailed { message }) => {
2684                assert_eq!(message, "no turn in flight");
2685            }
2686            other => panic!("expected DuplexControlFailed, got {other:?}"),
2687        }
2688    }
2689
2690    #[test]
2691    fn parse_control_response_returns_none_on_missing_request_id() {
2692        let msg = json!({
2693            "type": "control_response",
2694            "response": { "subtype": "success" }
2695        });
2696        assert!(parse_control_response(&msg).is_none());
2697    }
2698
2699    #[test]
2700    fn parse_control_response_returns_none_on_unknown_subtype() {
2701        let msg = json!({
2702            "type": "control_response",
2703            "response": { "request_id": "x", "subtype": "future_subtype" }
2704        });
2705        assert!(parse_control_response(&msg).is_none());
2706    }
2707
2708    #[test]
2709    fn handle_inbound_returns_control_response_action() {
2710        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2711        let (events_tx, _events_rx) = broadcast::channel(16);
2712        let mut pending = Some((tx, Vec::new()));
2713        let action = handle_inbound(
2714            json!({
2715                "type": "control_response",
2716                "response": {
2717                    "request_id": "interrupt-1",
2718                    "subtype": "success",
2719                    "response": {}
2720                }
2721            }),
2722            &mut pending,
2723            &events_tx,
2724        );
2725        match action {
2726            InboundAction::ControlResponse {
2727                request_id,
2728                outcome,
2729            } => {
2730                assert_eq!(request_id, "interrupt-1");
2731                assert!(outcome.is_ok());
2732            }
2733            InboundAction::None | InboundAction::Permission(_) => {
2734                panic!("expected ControlResponse action");
2735            }
2736        }
2737    }
2738
2739    #[test]
2740    fn handle_inbound_treats_malformed_control_response_as_other() {
2741        let (tx, _reply_rx) = oneshot::channel::<Result<TurnResult>>();
2742        let (events_tx, mut events_rx) = broadcast::channel(16);
2743        let mut pending = Some((tx, Vec::new()));
2744        let action = handle_inbound(
2745            json!({
2746                "type": "control_response",
2747                "response": { "subtype": "success" }
2748            }),
2749            &mut pending,
2750            &events_tx,
2751        );
2752        assert!(matches!(action, InboundAction::None));
2753        let event = events_rx.try_recv().expect("broadcast");
2754        assert!(matches!(event, InboundEvent::Other(_)));
2755    }
2756
2757    #[tokio::test]
2758    async fn permission_handler_clones_arc() {
2759        let handler = PermissionHandler::new(|_req| async move {
2760            PermissionDecision::Allow {
2761                updated_input: None,
2762            }
2763        });
2764        let cloned = handler.clone();
2765        let req = PermissionRequest {
2766            request_id: "r1".into(),
2767            tool_name: "Read".into(),
2768            input: Value::Null,
2769            raw: Value::Null,
2770        };
2771        // Both handles invoke the same underlying closure.
2772        let _ = handler.invoke(req.clone()).await;
2773        let _ = cloned.invoke(req).await;
2774    }
2775
2776    /// Build a `DuplexSession` whose channels are wired up but whose
2777    /// background task is a no-op. Tests can drive the watch state
2778    /// machine via the returned `exit_tx` and observe the public
2779    /// accessors. The fake task idles on a oneshot so it stays alive
2780    /// for the life of the test (no JoinHandle::abort handshake
2781    /// needed).
2782    fn fake_session(
2783        initial: SessionExitStatus,
2784    ) -> (
2785        DuplexSession,
2786        watch::Sender<SessionExitStatus>,
2787        oneshot::Sender<()>,
2788    ) {
2789        let (outbound_tx, outbound_rx) = mpsc::unbounded_channel::<OutboundMsg>();
2790        let (events_tx, _events_rx) = broadcast::channel::<InboundEvent>(16);
2791        let (exit_tx, exit_rx) = watch::channel(initial);
2792        let (stop_tx, stop_rx) = oneshot::channel::<()>();
2793
2794        let join = tokio::spawn(async move {
2795            let _outbound_rx = outbound_rx;
2796            let _ = stop_rx.await;
2797            Ok::<(), Error>(())
2798        });
2799
2800        let session = DuplexSession {
2801            outbound_tx,
2802            events_tx,
2803            exit_rx,
2804            join,
2805        };
2806        (session, exit_tx, stop_tx)
2807    }
2808
2809    #[tokio::test]
2810    async fn is_alive_true_while_running() {
2811        let (session, _exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2812        assert!(session.is_alive());
2813    }
2814
2815    #[tokio::test]
2816    async fn is_alive_false_after_completed() {
2817        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2818        exit_tx.send(SessionExitStatus::Completed).unwrap();
2819        assert!(!session.is_alive());
2820    }
2821
2822    #[tokio::test]
2823    async fn is_alive_false_after_failed() {
2824        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2825        exit_tx
2826            .send(SessionExitStatus::Failed("boom".into()))
2827            .unwrap();
2828        assert!(!session.is_alive());
2829    }
2830
2831    #[tokio::test]
2832    async fn exit_status_reports_running_initially() {
2833        let (session, _exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2834        assert!(matches!(session.exit_status(), SessionExitStatus::Running));
2835    }
2836
2837    #[tokio::test]
2838    async fn exit_status_reflects_completed() {
2839        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2840        exit_tx.send(SessionExitStatus::Completed).unwrap();
2841        assert!(matches!(
2842            session.exit_status(),
2843            SessionExitStatus::Completed
2844        ));
2845    }
2846
2847    #[tokio::test]
2848    async fn exit_status_reflects_failed_with_message() {
2849        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2850        exit_tx
2851            .send(SessionExitStatus::Failed("oh no".into()))
2852            .unwrap();
2853        match session.exit_status() {
2854            SessionExitStatus::Failed(msg) => assert_eq!(msg, "oh no"),
2855            other => panic!("expected Failed, got {other:?}"),
2856        }
2857    }
2858
2859    #[tokio::test]
2860    async fn wait_for_exit_returns_immediately_when_already_terminal() {
2861        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2862        exit_tx.send(SessionExitStatus::Completed).unwrap();
2863        let status = tokio::time::timeout(Duration::from_secs(1), session.wait_for_exit())
2864            .await
2865            .expect("wait_for_exit should not block when already terminal");
2866        assert!(matches!(status, SessionExitStatus::Completed));
2867    }
2868
2869    #[tokio::test]
2870    async fn wait_for_exit_blocks_until_state_transitions() {
2871        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2872
2873        let waiter = async { session.wait_for_exit().await };
2874        let driver = async {
2875            tokio::time::sleep(Duration::from_millis(20)).await;
2876            exit_tx.send(SessionExitStatus::Completed).unwrap();
2877        };
2878        let (status, ()) = tokio::join!(waiter, driver);
2879        assert!(matches!(status, SessionExitStatus::Completed));
2880    }
2881
2882    #[tokio::test]
2883    async fn wait_for_exit_supports_multiple_observers() {
2884        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2885
2886        let waiter1 = async { session.wait_for_exit().await };
2887        let waiter2 = async { session.wait_for_exit().await };
2888        let driver = async {
2889            tokio::time::sleep(Duration::from_millis(20)).await;
2890            exit_tx
2891                .send(SessionExitStatus::Failed("crash".into()))
2892                .unwrap();
2893        };
2894        let (s1, s2, ()) = tokio::join!(waiter1, waiter2, driver);
2895        match s1 {
2896            SessionExitStatus::Failed(msg) => assert_eq!(msg, "crash"),
2897            other => panic!("waiter1 expected Failed, got {other:?}"),
2898        }
2899        match s2 {
2900            SessionExitStatus::Failed(msg) => assert_eq!(msg, "crash"),
2901            other => panic!("waiter2 expected Failed, got {other:?}"),
2902        }
2903    }
2904
2905    #[tokio::test]
2906    async fn wait_for_exit_returns_last_value_when_sender_dropped() {
2907        // Defensive: if exit_tx is dropped without ever publishing a
2908        // terminal value, wait_for_exit should fall back to the last
2909        // observed state rather than hang.
2910        let (session, exit_tx, _stop) = fake_session(SessionExitStatus::Running);
2911        let waiter = async { session.wait_for_exit().await };
2912        let driver = async {
2913            tokio::time::sleep(Duration::from_millis(20)).await;
2914            drop(exit_tx);
2915        };
2916        let (status, ()) = tokio::time::timeout(Duration::from_secs(1), async {
2917            tokio::join!(waiter, driver)
2918        })
2919        .await
2920        .expect("wait_for_exit must not hang when sender is dropped");
2921        assert!(matches!(status, SessionExitStatus::Running));
2922    }
2923}