fix(desktop/chat): address codex review BLOCKs + CONCERNs on chat backend trio

Codex review of 3d754fdc / 85d93e7c / 1b168a84 flagged six BLOCKs +
five CONCERNs + one NIT. This commit closes them.

BLOCKs (all six fixed):

1. stderr pipe not drained → CLI deadlocks on full pipe. Now spawned
   sibling task drains stderr to /dev/null for the lifetime of the
   child.
2. Receiver-drop only detected on next stdout line → idle CLI keeps
   running forever. Switched both BuiltIn + Subprocess channels from
   `std::sync::mpsc` to `tokio::sync::mpsc` so `tx.closed()` is a
   future we can race against `lines.next_line()` in a `select!`.
   Sync iterator wrapper `BlockingRecvIter` in chat_runtime.rs uses
   `Receiver::blocking_recv`.
3. `Done` event marked the flag but didn't break the read loop. Now
   breaks immediately on any structured `done` so stdout staying
   open past turn-end doesn't hang the iterator.
4. EOF path always emitted `Done { EndTurn }`, ignoring child exit
   status. Now reaps the child and surfaces non-zero exit as
   `Error("CLI exited with status N") + Done { Aborted }`.
5. stdout read error fell through to `Done { EndTurn }`. Now emits
   `Error + Done { Aborted }` so I/O failures aren't reported as
   normal completion.
6. Malformed structured events silently produced empty deltas /
   empty errors / nameless tool calls. `parse_line` now requires the
   shape's mandatory fields and emits `Error("malformed X event:
   ...")` when they're missing. Done's `stop_reason` stays optional
   (already had a sensible `EndTurn` fallback).

CONCERNs (4 of 5 fixed):

1. `ChatRequest.system_prompt` + `max_output_tokens` ignored by
   BuiltInProvider. Now honored per-turn: each `send` builds a
   turn-local `QueryEngine` cloning the provider Arc (cheap) +
   applying request's system/max-tokens with constructor defaults
   as fallback. `BuiltInProvider` struct now holds the pieces
   instead of a pre-built engine.
2. BuiltIn `Error` path emitted no terminal `Done`. Now every
   error-terminal branch sends `Done { Aborted }` so consumers can
   distinguish "stream errored" from "channel silently closed."
3. Subprocess stdin write errors silently ignored. Now surfaces as
   `Error("stdin write: ...") + Done { Aborted }` with child kill +
   wait.
5. `SubprocessProvider::for_cli(Codex | OpenCode)` accepted the
   wrong backend silently. Signature now returns `Option<Self>`;
   HttpServer-category CLIs return `None`. Direct stdio bridging to
   a `codex` binary still possible via `with_binary`.

CONCERN 4 (Subprocess silently drops `system_prompt` +
`max_output_tokens`) intentionally deferred — those fields have no
universal CLI mapping; the planned settings-modal flow lets users
encode them in argv directly via `with_binary`. Documented in the
module header as the contract.

NIT 1 (chat_provider.rs doc said "three" but listed four backends)
fixed.

Tests: 14 → 16 native chat tests + 250 shell-core tests pass. Two
new tests cover for_cli's `None` return for HttpServer kinds + the
malformed structured-event paths.
This commit is contained in:
Kayshen-X 2026-05-14 16:16:24 +08:00
parent b788c088ca
commit b39d69eb70
3 changed files with 292 additions and 121 deletions

View file

@ -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<dyn Provider>`), so each `send` pays at most a few small
/// allocations.
pub struct BuiltInProvider {
engine: Arc<QueryEngine>,
provider: Arc<dyn Provider>,
model: String,
default_system: Option<String>,
default_max_output_tokens: u32,
label: String,
}
@ -71,12 +78,11 @@ impl BuiltInProvider {
max_output_tokens: u32,
label: impl Into<String>,
) -> 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<dyn Iterator<Item = ChatDelta> + 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::<ChatDelta>();
let (tx, rx) = mpsc::channel::<ChatDelta>(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<ChatDelta>` into
/// a sync `Iterator<Item = ChatDelta>`. 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<T> {
rx: mpsc::Receiver<T>,
}
impl<T> BlockingRecvIter<T> {
pub(crate) fn new(rx: mpsc::Receiver<T>) -> Self {
Self { rx }
}
}
impl<T> Iterator for BlockingRecvIter<T> {
type Item = T;
fn next(&mut self) -> Option<T> {
self.rx.blocking_recv()
}
}

View file

@ -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<Self> {
let args: Vec<String> = 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<dyn Iterator<Item = ChatDelta> + Send> {
let binary = self.binary.clone();
let args = Arc::new(self.args.clone());
let (tx, rx) = mpsc::channel::<ChatDelta>();
let (tx, rx) = mpsc::channel::<ChatDelta>(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.

View file

@ -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]]:
//!