diff --git a/crates/openpencil-desktop/src/chat_runtime.rs b/crates/openpencil-desktop/src/chat_runtime.rs index b0740963c..9aa4ac689 100644 --- a/crates/openpencil-desktop/src/chat_runtime.rs +++ b/crates/openpencil-desktop/src/chat_runtime.rs @@ -20,7 +20,7 @@ //! that pumps `Event`s into a `std::sync::mpsc::channel`; the returned //! iterator drains the receiver until it goes idle. -use std::sync::{mpsc, Arc, OnceLock}; +use std::sync::{Arc, OnceLock}; use agent::abort::AbortController; use agent::provider::Provider; @@ -31,6 +31,7 @@ use openpencil_shell_core::chat_provider::{ ChatDelta, ChatProvider, ChatRequest, StopReason, }; use tokio::runtime::{Builder, Runtime}; +use tokio::sync::mpsc; /// Process-wide tokio runtime used for every BuiltIn chat turn. We /// own a single multi-thread runtime instead of spinning one up per @@ -50,12 +51,18 @@ pub(crate) fn shared_runtime() -> &'static Runtime { /// `ChatProvider` impl that drives `agent::QueryEngine` directly. The /// engine carries its own message store, so repeated `send` calls -/// retain conversation history across turns (multi-turn chat). The -/// `system_prompt` field on `ChatRequest` is applied once at construct -/// time via `with_system`; subsequent requests override it per turn by -/// rebuilding the engine — cheap because the `Provider` is `Arc`'d. +/// retain conversation history across turns (multi-turn chat). +/// +/// Per-`ChatRequest` `system_prompt` + `max_output_tokens` are honored +/// by rebuilding a turn-local engine that shares the provider + +/// message-store handle. `with_system` clones cheaply (provider is +/// `Arc`), so each `send` pays at most a few small +/// allocations. pub struct BuiltInProvider { - engine: Arc, + provider: Arc, + model: String, + default_system: Option, + default_max_output_tokens: u32, label: String, } @@ -71,12 +78,11 @@ impl BuiltInProvider { max_output_tokens: u32, label: impl Into, ) -> Self { - let mut engine = QueryEngine::new(provider, model).with_max_output_tokens(max_output_tokens); - if let Some(sys) = system_prompt { - engine = engine.with_system(sys); - } Self { - engine: Arc::new(engine), + provider, + model: model.into(), + default_system: system_prompt, + default_max_output_tokens: max_output_tokens.max(1), label: label.into(), } } @@ -91,65 +97,133 @@ impl ChatProvider for BuiltInProvider { &self, request: ChatRequest, ) -> Box + Send> { - let engine = self.engine.clone(); + // Build a per-turn engine so the request's `system_prompt` + // and `max_output_tokens` actually take effect (codex CONCERN + // 1). Provider is `Arc`'d so this is cheap. + let system = if request.system_prompt.is_empty() { + self.default_system.clone() + } else { + Some(request.system_prompt.clone()) + }; + let max_tokens = if request.max_output_tokens == 0 { + self.default_max_output_tokens + } else { + request.max_output_tokens + }; + let mut engine = + QueryEngine::new(self.provider.clone(), self.model.clone()) + .with_max_output_tokens(max_tokens); + if let Some(sys) = system { + engine = engine.with_system(sys); + } + let engine = Arc::new(engine); let abort = AbortController::new(); - let (tx, rx) = mpsc::channel::(); + let (tx, rx) = mpsc::channel::(64); shared_runtime().spawn(async move { - // Per-turn engine settings would land via a per-turn - // builder; for now we honor `max_output_tokens` set at - // construct time and pass the user message straight in. - // `system_prompt` on the request is ignored because the - // engine's system prompt is fixed at construct time — - // changing it per-turn would need a fresh `QueryEngine`. - let _ = request.system_prompt; // see note above - let _ = request.max_output_tokens; // honored at construct let stream = match engine.run(request.user_message, abort).await { Ok(s) => s, Err(e) => { - let _ = tx.send(ChatDelta::Error(e.to_string())); - let _ = tx.send(ChatDelta::Done { - stop_reason: StopReason::Aborted, - }); + let _ = tx.send(ChatDelta::Error(e.to_string())).await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; return; } }; let mut stream = stream; + let mut emitted_done = false; while let Some(item) = stream.next().await { match item { Ok(Event::TextDelta { delta }) => { - if tx.send(ChatDelta::TextDelta(delta)).is_err() { - break; + if tx.send(ChatDelta::TextDelta(delta)).await.is_err() { + return; } } Ok(Event::Thinking { delta }) => { - if tx.send(ChatDelta::Thinking(delta)).is_err() { - break; + if tx.send(ChatDelta::Thinking(delta)).await.is_err() { + return; } } Ok(Event::ToolUse { name, input, .. }) => { let args = input.to_string(); - if tx.send(ChatDelta::ToolUse { name, args }).is_err() { - break; + if tx + .send(ChatDelta::ToolUse { name, args }) + .await + .is_err() + { + return; } } Ok(Event::Result { data }) => { let reason = map_stop_reason(data.stop_reason.as_deref()); - let _ = tx.send(ChatDelta::Done { stop_reason: reason }); + let _ = tx.send(ChatDelta::Done { stop_reason: reason }).await; + emitted_done = true; break; } Ok(Event::Error { code, message }) => { - let _ = tx.send(ChatDelta::Error(format!("{code}: {message}"))); + let _ = tx + .send(ChatDelta::Error(format!("{code}: {message}"))) + .await; + // Always send a terminal `Done` after `Error` + // so consumers can distinguish a completed + // error path from a silently-dropped channel + // (codex CONCERN 2). + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; + emitted_done = true; break; } Ok(_) => {} // ToolResult / Usage / Notice / Unknown — silent Err(e) => { - let _ = tx.send(ChatDelta::Error(e.to_string())); + let _ = tx.send(ChatDelta::Error(e.to_string())).await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; + emitted_done = true; break; } } } + if !emitted_done { + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::EndTurn, + }) + .await; + } }); - Box::new(rx.into_iter()) + Box::new(BlockingRecvIter::new(rx)) + } +} + +/// Adapter that turns a `tokio::sync::mpsc::Receiver` into +/// a sync `Iterator`. The receiver's +/// `blocking_recv` blocks the calling thread until a value arrives or +/// the channel closes — exactly the contract `ChatProvider::send`'s +/// iterator return needs. Sharing this helper across both BuiltIn + +/// Subprocess (and the future HttpServer / Acp) keeps the async ↔ +/// sync bridge in one place. +pub(crate) struct BlockingRecvIter { + rx: mpsc::Receiver, +} + +impl BlockingRecvIter { + pub(crate) fn new(rx: mpsc::Receiver) -> Self { + Self { rx } + } +} + +impl Iterator for BlockingRecvIter { + type Item = T; + fn next(&mut self) -> Option { + self.rx.blocking_recv() } } diff --git a/crates/openpencil-desktop/src/chat_subprocess.rs b/crates/openpencil-desktop/src/chat_subprocess.rs index 146a7f74e..d2a73b2e6 100644 --- a/crates/openpencil-desktop/src/chat_subprocess.rs +++ b/crates/openpencil-desktop/src/chat_subprocess.rs @@ -26,15 +26,16 @@ //! for the duration of stdout pumping; killed by drop on early channel //! close (user navigated away mid-stream). -use std::sync::{mpsc, Arc}; +use std::sync::Arc; use openpencil_shell_core::chat_provider::{ ChatDelta, ChatProvider, ChatRequest, CliName, StopReason, }; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::process::Command; +use tokio::sync::mpsc; -use crate::chat_runtime::shared_runtime; +use crate::chat_runtime::{shared_runtime, BlockingRecvIter}; /// `ChatProvider` impl that bridges to a CLI binary via stdio. /// Construct via [`SubprocessProvider::for_cli`] or @@ -51,7 +52,14 @@ impl SubprocessProvider { /// `--print` / `--stream-json` flags (best-effort defaults — each /// CLI's real flag set lands as user-tunable in the settings /// modal). The returned provider's label is `cli.label()`. - pub fn for_cli(cli: CliName) -> Self { + /// + /// Returns `None` when `cli` is in the HttpServer category + /// (Codex / OpenCode) — those go through the dedicated + /// HttpServerProvider, not this subprocess bridge. Callers who + /// truly want a generic stdio pipe to a `codex` / `opencode` + /// binary can use [`SubprocessProvider::with_binary`] instead + /// (codex CONCERN 5: don't silently accept the wrong backend). + pub fn for_cli(cli: CliName) -> Option { let args: Vec = match cli { CliName::ClaudeCode => vec![ "--print".into(), @@ -60,17 +68,13 @@ impl SubprocessProvider { ], CliName::Gemini => vec!["--quiet".into()], CliName::Copilot => vec!["suggest".into()], - // The Subprocess kind only applies to subprocess-IPC CLIs. - // Codex / OpenCode use HTTP-server mode; if a caller picks - // them here we still spawn the binary but pass through - // generic stdin/stdout (the http_server bridge is correct). - CliName::Codex | CliName::OpenCode => Vec::new(), + CliName::Codex | CliName::OpenCode => return None, }; - Self { + Some(Self { binary: cli.default_binary().into(), args, label: cli.label().into(), - } + }) } /// Build a subprocess provider with a user-supplied binary path @@ -101,7 +105,7 @@ impl ChatProvider for SubprocessProvider { ) -> Box + Send> { let binary = self.binary.clone(); let args = Arc::new(self.args.clone()); - let (tx, rx) = mpsc::channel::(); + let (tx, rx) = mpsc::channel::(64); shared_runtime().spawn(async move { let mut child = match Command::new(&binary) .args(args.iter()) @@ -112,67 +116,145 @@ impl ChatProvider for SubprocessProvider { { Ok(c) => c, Err(e) => { - let _ = tx.send(ChatDelta::Error(format!("spawn {binary}: {e}"))); - let _ = tx.send(ChatDelta::Done { - stop_reason: StopReason::Aborted, - }); + let _ = tx + .send(ChatDelta::Error(format!("spawn {binary}: {e}"))) + .await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; return; } }; + // Drain stderr to /dev/null on a sibling task so the CLI + // never deadlocks on a full stderr pipe (codex BLOCK 1). + // Future work could route it into a status-bar notice + // channel; today we just keep the pipe drained. + if let Some(stderr) = child.stderr.take() { + tokio::spawn(async move { + let mut lines = BufReader::new(stderr).lines(); + while let Ok(Some(_)) = lines.next_line().await {} + }); + } + if let Some(mut stdin) = child.stdin.take() { // Feed the user message + close stdin so the CLI // sees EOF and starts responding instead of waiting - // for more input. - let _ = stdin.write_all(request.user_message.as_bytes()).await; - let _ = stdin.shutdown().await; + // for more input. Stdin write errors surface as a + // chat error so the user sees the broken-pipe instead + // of silent normal completion (codex CONCERN 3). + if let Err(e) = stdin.write_all(request.user_message.as_bytes()).await { + let _ = tx + .send(ChatDelta::Error(format!("stdin write: {e}"))) + .await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; + let _ = child.start_kill(); + let _ = child.wait().await; + return; + } + let _ = stdin.shutdown().await; // EOF; ignore close error } let stdout = match child.stdout.take() { Some(s) => s, None => { - let _ = tx.send(ChatDelta::Error("no stdout from CLI".into())); - let _ = tx.send(ChatDelta::Done { - stop_reason: StopReason::Aborted, - }); + let _ = tx + .send(ChatDelta::Error("no stdout from CLI".into())) + .await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; return; } }; let mut lines = BufReader::new(stdout).lines(); let mut emitted_done = false; + let mut terminal_error = false; loop { - match lines.next_line().await { - Ok(Some(line)) => { - let delta = parse_line(&line); - if matches!(delta, ChatDelta::Done { .. }) { - emitted_done = true; - } - if tx.send(delta).is_err() { - // Receiver dropped — kill the child so - // we don't leak a hung CLI process. - let _ = child.start_kill(); - return; - } - } - Ok(None) => break, // stdout EOF - Err(e) => { - let _ = tx.send(ChatDelta::Error(e.to_string())); + // Race the line read against channel closure so the + // bridge notices an idle receiver-drop without + // waiting for the next CLI output (codex BLOCK 2). + // `tokio::sync::mpsc::Sender::closed()` returns a + // future that resolves when every receiver has been + // dropped — no polling required. + tokio::select! { + biased; + _ = tx.closed() => { + let _ = child.start_kill(); break; } + result = lines.next_line() => match result { + Ok(Some(line)) => { + let delta = parse_line(&line); + let is_done = matches!(delta, ChatDelta::Done { .. }); + if tx.send(delta).await.is_err() { + let _ = child.start_kill(); + break; + } + if is_done { + // CLI signaled turn end — stop reading + // even if stdout stays open (codex + // BLOCK 3). + emitted_done = true; + break; + } + } + Ok(None) => break, // stdout EOF + Err(e) => { + let _ = tx.send(ChatDelta::Error(e.to_string())).await; + // Terminal: don't paper over an I/O error + // with `EndTurn` (codex BLOCK 5). + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; + terminal_error = true; + break; + } + }, } } - if !emitted_done { - let _ = tx.send(ChatDelta::Done { - stop_reason: StopReason::EndTurn, - }); + // Reap the child + interpret exit status (codex BLOCK 4): + // non-zero exit with no prior Done surfaces as Error + + // Aborted instead of an unrelated `EndTurn`. + let status = child.wait().await.ok(); + if !emitted_done && !terminal_error { + let nonzero = status.as_ref().map(|s| !s.success()).unwrap_or(false); + if nonzero { + let code = status + .as_ref() + .and_then(|s| s.code()) + .map(|c| c.to_string()) + .unwrap_or_else(|| "?".into()); + let _ = tx + .send(ChatDelta::Error(format!("CLI exited with status {code}"))) + .await; + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::Aborted, + }) + .await; + } else { + let _ = tx + .send(ChatDelta::Done { + stop_reason: StopReason::EndTurn, + }) + .await; + } } - // Reap the child so we don't leave a zombie. Ignored - // error: child may have already exited or been killed. - let _ = child.wait().await; }); - Box::new(rx.into_iter()) + Box::new(BlockingRecvIter::new(rx)) } } @@ -198,49 +280,39 @@ fn parse_line(line: &str) -> ChatDelta { } }; let ty = val.get("type").and_then(|v| v.as_str()).unwrap_or(""); + // Strict shape per type — missing or wrong-type required fields + // become `Error` deltas instead of silent empty deltas (codex + // BLOCK 6). Better to surface a parse problem than to feed empty + // strings into the chat panel. match ty { - "text" => { - let delta = val - .get("delta") - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - ChatDelta::TextDelta(delta) - } - "thinking" => { - let delta = val - .get("delta") - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - ChatDelta::Thinking(delta) - } - "tool_use" => { - let name = val - .get("name") - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - let args = val - .get("args") - .map(|v| v.to_string()) - .unwrap_or_else(|| "{}".into()); - ChatDelta::ToolUse { name, args } - } + "text" => match val.get("delta").and_then(|v| v.as_str()) { + Some(s) => ChatDelta::TextDelta(s.to_string()), + None => ChatDelta::Error(format!("malformed text event: {trimmed}")), + }, + "thinking" => match val.get("delta").and_then(|v| v.as_str()) { + Some(s) => ChatDelta::Thinking(s.to_string()), + None => ChatDelta::Error(format!("malformed thinking event: {trimmed}")), + }, + "tool_use" => match ( + val.get("name").and_then(|v| v.as_str()), + val.get("args"), + ) { + (Some(name), Some(args)) => ChatDelta::ToolUse { + name: name.to_string(), + args: args.to_string(), + }, + _ => ChatDelta::Error(format!("malformed tool_use event: {trimmed}")), + }, "done" => { let reason = val.get("stop_reason").and_then(|v| v.as_str()); ChatDelta::Done { stop_reason: map_stop_reason(reason), } } - "error" => { - let msg = val - .get("message") - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); - ChatDelta::Error(msg) - } + "error" => match val.get("message").and_then(|v| v.as_str()) { + Some(msg) => ChatDelta::Error(msg.to_string()), + None => ChatDelta::Error(format!("malformed error event: {trimmed}")), + }, _ => { // Unknown structured event — surface the raw line so the // user can debug what their CLI is emitting. @@ -336,7 +408,7 @@ mod tests { #[test] fn for_cli_claude_code_seeds_print_stream_json_flags() { - let p = SubprocessProvider::for_cli(CliName::ClaudeCode); + let p = SubprocessProvider::for_cli(CliName::ClaudeCode).unwrap(); assert_eq!(p.binary, "claude"); assert!(p.args.iter().any(|a| a == "--print")); assert!(p.args.iter().any(|a| a == "stream-json")); @@ -346,15 +418,40 @@ mod tests { #[test] fn for_cli_uses_default_binary_per_cli_name() { assert_eq!( - SubprocessProvider::for_cli(CliName::Gemini).binary, + SubprocessProvider::for_cli(CliName::Gemini).unwrap().binary, "gemini" ); assert_eq!( - SubprocessProvider::for_cli(CliName::Copilot).binary, + SubprocessProvider::for_cli(CliName::Copilot).unwrap().binary, "gh-copilot" ); } + #[test] + fn for_cli_rejects_http_server_kinds() { + assert!(SubprocessProvider::for_cli(CliName::Codex).is_none()); + assert!(SubprocessProvider::for_cli(CliName::OpenCode).is_none()); + } + + #[test] + fn parse_line_malformed_structured_event_is_error() { + // "text" with no "delta" → Error + match parse_line(r#"{"type":"text"}"#) { + ChatDelta::Error(s) => assert!(s.contains("malformed text")), + other => panic!("expected Error, got {other:?}"), + } + // "tool_use" missing "name" → Error + match parse_line(r#"{"type":"tool_use","args":{}}"#) { + ChatDelta::Error(s) => assert!(s.contains("malformed tool_use")), + other => panic!("expected Error, got {other:?}"), + } + // "error" missing "message" → Error + match parse_line(r#"{"type":"error"}"#) { + ChatDelta::Error(s) => assert!(s.contains("malformed error")), + other => panic!("expected Error, got {other:?}"), + } + } + #[test] fn spawn_failure_surfaces_error_and_done() { // Use a binary path that's guaranteed not to exist on PATH. diff --git a/crates/openpencil-shell-core/src/chat_provider.rs b/crates/openpencil-shell-core/src/chat_provider.rs index aa19e9a71..91f82f6a3 100644 --- a/crates/openpencil-shell-core/src/chat_provider.rs +++ b/crates/openpencil-shell-core/src/chat_provider.rs @@ -1,4 +1,4 @@ -//! Chat / agent provider abstraction. Three backend categories +//! Chat / agent provider abstraction. Four backend categories //! mirror the architecture decision in //! [[project_agent_runtime]]: //!