From 941d51f805260abc5797de6bae9f98e75942e2df Mon Sep 17 00:00:00 2001 From: Kayshen-X Date: Sun, 14 Jun 2026 10:54:26 +0800 Subject: [PATCH] feat(process): extract shared process io primitives --- Cargo.lock | 9 + crates/op-cli/Cargo.toml | 1 + crates/op-cli/src/app_control_cli.rs | 81 ++++---- crates/op-host-desktop/Cargo.toml | 2 + crates/op-host-desktop/src/chat_subprocess.rs | 61 +++--- crates/op-process-io/Cargo.toml | 23 +++ crates/op-process-io/src/lib.rs | 187 ++++++++++++++++++ crates/op-process-io/tests/process_io.rs | 106 ++++++++++ 8 files changed, 393 insertions(+), 77 deletions(-) create mode 100644 crates/op-process-io/Cargo.toml create mode 100644 crates/op-process-io/src/lib.rs create mode 100644 crates/op-process-io/tests/process_io.rs diff --git a/Cargo.lock b/Cargo.lock index 7029d5625..013b230bb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2998,6 +2998,7 @@ version = "0.8.0" dependencies = [ "op-config-store", "op-figma", + "op-process-io", "op-rpc-transport", "serde_json", ] @@ -3121,6 +3122,7 @@ dependencies = [ "op-opmerge", "op-orchestrator", "op-pen-loader", + "op-process-io", "reqwest 0.12.28", "rfd", "serde", @@ -3233,6 +3235,13 @@ dependencies = [ "serde_stacker", ] +[[package]] +name = "op-process-io" +version = "0.8.0" +dependencies = [ + "tokio", +] + [[package]] name = "op-rpc-transport" version = "0.8.0" diff --git a/crates/op-cli/Cargo.toml b/crates/op-cli/Cargo.toml index 5d1b1b1d3..e45dee0f3 100644 --- a/crates/op-cli/Cargo.toml +++ b/crates/op-cli/Cargo.toml @@ -14,5 +14,6 @@ path = "src/main.rs" [dependencies] op-config-store = { path = "../op-config-store" } op-figma = { path = "../op-figma" } +op-process-io = { path = "../op-process-io" } op-rpc-transport = { path = "../op-rpc-transport" } serde_json = { workspace = true } diff --git a/crates/op-cli/src/app_control_cli.rs b/crates/op-cli/src/app_control_cli.rs index a72900f40..b3abb1c90 100644 --- a/crates/op-cli/src/app_control_cli.rs +++ b/crates/op-cli/src/app_control_cli.rs @@ -1,9 +1,9 @@ +use op_process_io::{spawn_null, wait_for_child_or, wait_until_false, WaitOutcome}; use serde_json::json; use std::env; use std::fs; use std::path::{Path, PathBuf}; use std::process::{Command, Stdio}; -use std::thread; use std::time::{Duration, SystemTime, UNIX_EPOCH}; const PID_FILE_NAME: &str = "openpencil-mcp-server.pid"; @@ -101,11 +101,7 @@ fn run_start_live(port: u16, document_path: Option<&str>) -> Result) -> Result { return Ok(start_json(live_pid, live_port, opened.map(Path::new))); } - if let Ok(Some(status)) = child.try_wait() { + WaitOutcome::Exited(status) => { return Err(format!( "OpenPencil editor exited before serving the live MCP server: {status}" )); } - thread::sleep(Duration::from_millis(100)); + WaitOutcome::TimedOut => {} } // Timed out without a verified live server. Report failure honestly // rather than a fabricated success on the requested port — the editor @@ -182,12 +182,8 @@ fn run_start_web( if let Some(host) = host { command.arg("--host").arg(host); } - let mut child = command - .env("OPENPENCIL_MCP_TOKEN", &token) - .stdin(Stdio::null()) - .stdout(Stdio::null()) - .stderr(Stdio::null()) - .spawn() + command.env("OPENPENCIL_MCP_TOKEN", &token); + let mut child = spawn_null(&mut command) .map_err(|e| format!("spawn {} --serve-web: {e}", binary.display()))?; let pid = child.id(); @@ -196,19 +192,23 @@ fn run_start_web( // Wait up to ~5s for the daemon's HTTP server to answer. The token-authed // JSON-RPC ping doubles as the identity check (never trust a foreign // listener on the port), exactly like the headless start path. - for _ in 0..50 { - if crate::mcp_http_cli::mcp_ping_headless(port, &token) { + match wait_for_child_or(&mut child, 50, Duration::from_millis(100), || { + crate::mcp_http_cli::mcp_ping_headless(port, &token).then_some(()) + }) + .map_err(|e| format!("wait for {} --serve-web: {e}", binary.display()))? + { + WaitOutcome::Ready(()) => { let url = format!("http://127.0.0.1:{port}"); open_in_browser(&url); return Ok(start_web_json(pid, port, document.as_deref(), host)); } - if let Ok(Some(status)) = child.try_wait() { + WaitOutcome::Exited(status) => { remove_manager_files(); return Err(format!( "OpenPencil web daemon exited before accepting connections: {status}" )); } - thread::sleep(Duration::from_millis(100)); + WaitOutcome::TimedOut => {} } Err(format!( "OpenPencil web daemon did not respond on 127.0.0.1:{port} within 5s" @@ -240,11 +240,7 @@ fn open_in_browser(url: &str) { c.arg(url); c }; - let _ = command - .stdin(Stdio::null()) - .stdout(Stdio::null()) - .stderr(Stdio::null()) - .spawn(); + let _ = spawn_null(&mut command); } /// `op start --web` success JSON: the headless `start_json` shape plus @@ -282,34 +278,36 @@ fn run_start_headless(port: u16, document_path: Option<&str>) -> Result { return Ok(start_json(pid, port, Some(&document))); } - if let Ok(Some(status)) = child.try_wait() { + WaitOutcome::Exited(status) => { remove_manager_files(); return Err(format!( "OpenPencil MCP server exited before accepting connections: {status}" )); } - thread::sleep(Duration::from_millis(100)); + WaitOutcome::TimedOut => {} } Err(format!( @@ -374,12 +372,7 @@ fn remove_live_port_file() { /// false (the server stopped responding to its token) or the budget /// elapses — used to confirm a graceful shutdown actually took effect. fn wait_until bool>(still_up: F) { - for _ in 0..30 { - if !still_up() { - return; - } - thread::sleep(Duration::from_millis(100)); - } + let _ = wait_until_false(30, Duration::from_millis(100), still_up); } /// Path to the Rust live discovery file `~/.openpencil/.op-mcp-port`. diff --git a/crates/op-host-desktop/Cargo.toml b/crates/op-host-desktop/Cargo.toml index 811ee99c8..9274e50dd 100644 --- a/crates/op-host-desktop/Cargo.toml +++ b/crates/op-host-desktop/Cargo.toml @@ -88,6 +88,8 @@ op-orchestrator = { path = "../op-orchestrator" } # extracted into op-ai; the `src/chat_*.rs` real transports + the # model-discovery path import them through `op_ai::*`. op-ai = { path = "../op-ai" } +# Shared subprocess IO primitives for chat CLI bridges. +op-process-io = { path = "../op-process-io" } # AI skill engine — the BuiltIn agent provider resolves # generation-phase prompt skills through it (`chat_runtime.rs`). op-ai-skills = { path = "../op-ai-skills" } diff --git a/crates/op-host-desktop/src/chat_subprocess.rs b/crates/op-host-desktop/src/chat_subprocess.rs index 39d27963b..65fa80b78 100644 --- a/crates/op-host-desktop/src/chat_subprocess.rs +++ b/crates/op-host-desktop/src/chat_subprocess.rs @@ -58,7 +58,8 @@ use std::time::Duration; use op_ai::chat_provider::{ ChatDelta, ChatProvider, ChatRequest, CliName, EffortLevel, StopReason, ThinkingMode, }; -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; +use op_process_io::LineStreamChild; +use tokio::io::{AsyncBufReadExt, BufReader}; use tokio::sync::mpsc; use crate::chat_runtime::{prompt_with_system_prompt, shared_runtime, BlockingRecvIter}; @@ -387,15 +388,12 @@ impl ChatProvider for SubprocessProvider { // Keep staged attachment temp files alive for the turn. let _guard = guard; let mut cmd = build_command(&binary, &args); - cmd.stdin(std::process::Stdio::piped()) - .stdout(std::process::Stdio::piped()) - .stderr(std::process::Stdio::piped()); // Set the child's env from the per-CLI policy. We // env_clear first because tokio::process Command // otherwise inherits the parent env verbatim. cmd.env_clear(); cmd.envs(env_pairs); - let mut child = match cmd.spawn() { + let mut child = match LineStreamChild::spawn_command(cmd) { Ok(c) => c, Err(e) => { let _ = tx @@ -416,7 +414,7 @@ impl ChatProvider for SubprocessProvider { // (`extractCodexCliError`); other CLIs discard it (TS // parity: gemini stream path discards stderr). let stderr_tail: Arc> = Arc::default(); - if let Some(stderr) = child.stderr.take() { + if let Some(stderr) = child.take_stderr() { let capture = (cli == Some(CliName::Codex)).then(|| Arc::clone(&stderr_tail)); tokio::spawn(async move { let mut lines = BufReader::new(stderr).lines(); @@ -435,36 +433,34 @@ impl ChatProvider for SubprocessProvider { }); } - if let Some(mut stdin) = child.stdin.take() { - match prompt_mode { - PromptMode::Stdin => { - // Feed the user message + close stdin so the - // CLI sees EOF and starts responding. Stdin - // write errors surface as a chat error. - if let Err(e) = stdin.write_all(prompt.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; - } - } - PromptMode::PositionalArg => { - // No stdin write — prompt is in argv. Close - // stdin immediately so the CLI doesn't sit - // waiting on it (Claude Code's `--print` mode - // exits if stdin stays open with no input). + match prompt_mode { + PromptMode::Stdin => { + // Feed the user message + close stdin so the CLI + // sees EOF and starts responding. Stdin write + // errors surface as a chat error. + if let Err(e) = child.feed(prompt.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 + PromptMode::PositionalArg => { + // No stdin write — prompt is in argv. Close stdin + // immediately so the CLI doesn't sit waiting on it + // (Claude Code's `--print` mode exits if stdin + // stays open with no input). + } } + let _ = child.close_stdin().await; // EOF; ignore close error - let stdout = match child.stdout.take() { - Some(s) => s, + let mut lines = match child.take_lines() { + Some(lines) => lines, None => { let _ = tx.send(ChatDelta::Error("no stdout from CLI".into())).await; let _ = tx @@ -478,7 +474,6 @@ impl ChatProvider for SubprocessProvider { let deadline = tokio::time::Instant::now() + turn_timeout.unwrap_or(Duration::from_secs(60 * 60 * 24 * 365)); - let mut lines = BufReader::new(stdout).lines(); let mut emitted_done = false; let mut terminal_error = false; let mut emitted_text = false; diff --git a/crates/op-process-io/Cargo.toml b/crates/op-process-io/Cargo.toml new file mode 100644 index 000000000..81cbd4d50 --- /dev/null +++ b/crates/op-process-io/Cargo.toml @@ -0,0 +1,23 @@ +[package] +name = "op-process-io" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +description = "Shared OpenPencil process spawning, line-stream, and shutdown primitives" + +[dependencies] +tokio = { version = "1", default-features = false, features = [ + "io-util", + "process", + "time", +] } + +[dev-dependencies] +tokio = { version = "1", default-features = false, features = [ + "io-util", + "macros", + "process", + "rt", + "time", +] } diff --git a/crates/op-process-io/src/lib.rs b/crates/op-process-io/src/lib.rs new file mode 100644 index 000000000..293c7a56b --- /dev/null +++ b/crates/op-process-io/src/lib.rs @@ -0,0 +1,187 @@ +//! Shared process IO primitives for OpenPencil native crates. + +use std::ffi::OsStr; +use std::io; +use std::process::{Child, Command as StdCommand, ExitStatus, Stdio}; +use std::thread; +use std::time::Duration; + +use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, Lines}; +use tokio::process::{ChildStderr, ChildStdin, ChildStdout, Command as TokioCommand}; + +/// Result of polling a spawned child while waiting for an external +/// readiness signal. +#[derive(Debug, PartialEq, Eq)] +pub enum WaitOutcome { + Ready(T), + Exited(ExitStatus), + TimedOut, +} + +/// Apply the detached daemon stdio policy used by CLI-launched +/// OpenPencil processes. +pub fn null_stdio(command: &mut StdCommand) -> &mut StdCommand { + command + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) +} + +/// Spawn a std child with stdin/stdout/stderr connected to null. +pub fn spawn_null(command: &mut StdCommand) -> io::Result { + null_stdio(command).spawn() +} + +/// Poll `probe` while also noticing if `child` exits first. +pub fn wait_for_child_or( + child: &mut Child, + attempts: usize, + interval: Duration, + mut probe: impl FnMut() -> Option, +) -> io::Result> { + for _ in 0..attempts { + if let Some(value) = probe() { + return Ok(WaitOutcome::Ready(value)); + } + if let Some(status) = child.try_wait()? { + return Ok(WaitOutcome::Exited(status)); + } + thread::sleep(interval); + } + Ok(WaitOutcome::TimedOut) +} + +/// Poll until `still_up` becomes false, returning whether it stopped +/// within the allotted attempts. +pub fn wait_until_false( + attempts: usize, + interval: Duration, + mut still_up: impl FnMut() -> bool, +) -> bool { + for _ in 0..attempts { + if !still_up() { + return true; + } + thread::sleep(interval); + } + false +} + +/// Async stdout line stream for a piped child process. +pub type LineStream = Lines>; + +/// Tokio child wrapper with piped stdin/stdout/stderr. +pub struct LineStreamChild { + child: tokio::process::Child, + stdin: Option, + lines: Option, + stderr: Option, +} + +impl LineStreamChild { + /// Spawn `program` with piped stdio and the supplied args/envs. + pub fn spawn(program: P, args: A, envs: E) -> io::Result + where + P: AsRef, + A: IntoIterator, + A::Item: AsRef, + E: IntoIterator, + K: AsRef, + V: AsRef, + { + let mut command = TokioCommand::new(program); + command.args(args); + command.envs(envs); + Self::spawn_command(command) + } + + /// Spawn a preconfigured tokio command after forcing piped stdio. + pub fn spawn_command(mut command: TokioCommand) -> io::Result { + pipe_stdio(&mut command); + let mut child = command.spawn()?; + let stdout = child + .stdout + .take() + .ok_or_else(|| io::Error::new(io::ErrorKind::Other, "child stdout was not piped"))?; + Ok(Self { + stdin: child.stdin.take(), + lines: Some(BufReader::new(stdout).lines()), + stderr: child.stderr.take(), + child, + }) + } + + /// Write bytes to stdin without adding a newline. + pub async fn feed(&mut self, text: impl AsRef<[u8]>) -> io::Result<()> { + let Some(stdin) = &mut self.stdin else { + return Err(io::Error::new( + io::ErrorKind::BrokenPipe, + "child stdin is closed", + )); + }; + stdin.write_all(text.as_ref()).await + } + + /// Close stdin, signaling EOF to children that read from it. + pub async fn close_stdin(&mut self) -> io::Result<()> { + if let Some(mut stdin) = self.stdin.take() { + stdin.shutdown().await?; + } + Ok(()) + } + + /// Read the next stdout line. + pub async fn next_line(&mut self) -> io::Result> { + match &mut self.lines { + Some(lines) => lines.next_line().await, + None => Ok(None), + } + } + + /// Borrow the active stdout line stream. + pub fn lines(&mut self) -> Option<&mut LineStream> { + self.lines.as_mut() + } + + /// Move the stdout line stream out for select loops that must + /// still operate on the child concurrently. + pub fn take_lines(&mut self) -> Option { + self.lines.take() + } + + /// Move stderr out so callers can drain or capture it separately. + pub fn take_stderr(&mut self) -> Option { + self.stderr.take() + } + + /// Start platform termination without waiting for reaping. + pub fn start_kill(&mut self) -> io::Result<()> { + self.child.start_kill() + } + + /// Wait for process exit. + pub async fn wait(&mut self) -> io::Result { + self.child.wait().await + } + + /// Close stdin and wait for the process to exit, killing it if it + /// ignores EOF beyond `budget`. + pub async fn kill_graceful(&mut self, budget: Duration) -> io::Result { + let _ = self.close_stdin().await; + match tokio::time::timeout(budget, self.child.wait()).await { + Ok(status) => status, + Err(_) => { + self.child.start_kill()?; + self.child.wait().await + } + } + } +} + +/// Apply the piped stdio policy used by line-stream subprocesses. +pub fn pipe_stdio(command: &mut TokioCommand) -> &mut TokioCommand { + command + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) +} diff --git a/crates/op-process-io/tests/process_io.rs b/crates/op-process-io/tests/process_io.rs new file mode 100644 index 000000000..47b6bdf49 --- /dev/null +++ b/crates/op-process-io/tests/process_io.rs @@ -0,0 +1,106 @@ +#![cfg(unix)] + +use std::process::Command; +use std::time::Duration; + +use op_process_io::{spawn_null, wait_for_child_or, LineStreamChild, WaitOutcome}; + +#[test] +fn wait_for_child_or_reports_ready_before_process_exit() { + let mut command = Command::new("sh"); + command.args(["-c", "sleep 5"]); + let mut child = spawn_null(&mut command).expect("spawn child"); + + let outcome = wait_for_child_or(&mut child, 3, Duration::from_millis(10), || Some("ready")) + .expect("wait for readiness"); + + assert_eq!(outcome, WaitOutcome::Ready("ready")); + let _ = child.kill(); + let _ = child.wait(); +} + +#[test] +fn wait_for_child_or_reports_early_exit_status() { + let mut command = Command::new("sh"); + command.args(["-c", "exit 7"]); + let mut child = spawn_null(&mut command).expect("spawn child"); + + let outcome = wait_for_child_or::<()>(&mut child, 10, Duration::from_millis(10), || None) + .expect("wait for child exit"); + + let WaitOutcome::Exited(status) = outcome else { + panic!("expected exited outcome"); + }; + assert_eq!(status.code(), Some(7)); +} + +#[tokio::test] +async fn line_stream_child_reads_stdout_lines() { + let mut child = LineStreamChild::spawn( + "sh", + ["-c", "printf 'one\\ntwo\\n'"], + std::iter::empty::<(&str, &str)>(), + ) + .expect("spawn child"); + + assert_eq!( + child.next_line().await.expect("read line"), + Some("one".into()) + ); + assert_eq!( + child.next_line().await.expect("read line"), + Some("two".into()) + ); + assert!(child.wait().await.expect("wait").success()); +} + +#[tokio::test] +async fn line_stream_child_feeds_stdin_and_closes_it() { + let mut child = LineStreamChild::spawn( + "sh", + ["-c", "IFS= read -r line; printf '%s\\n' \"$line\""], + std::iter::empty::<(&str, &str)>(), + ) + .expect("spawn child"); + + child.feed("hello\n").await.expect("feed stdin"); + child.close_stdin().await.expect("close stdin"); + + assert_eq!( + child.next_line().await.expect("read echoed line"), + Some("hello".into()) + ); + assert!(child.wait().await.expect("wait").success()); +} + +#[tokio::test] +async fn kill_graceful_closes_stdin_before_forcing_exit() { + let mut child = LineStreamChild::spawn( + "sh", + ["-c", "cat >/dev/null"], + std::iter::empty::<(&str, &str)>(), + ) + .expect("spawn child"); + + let status = child + .kill_graceful(Duration::from_secs(1)) + .await + .expect("graceful kill"); + assert!(status.success()); +} + +#[tokio::test] +async fn kill_graceful_forces_exit_after_timeout() { + let mut child = LineStreamChild::spawn( + "sh", + ["-c", "while true; do sleep 1; done"], + std::iter::empty::<(&str, &str)>(), + ) + .expect("spawn child"); + + let status = child + .kill_graceful(Duration::from_millis(20)) + .await + .expect("forced kill"); + assert!(!status.success()); +}