feat(process): extract shared process io primitives

This commit is contained in:
Kayshen-X 2026-06-14 10:54:26 +08:00
parent 4d9b2d783d
commit 941d51f805
8 changed files with 393 additions and 77 deletions

9
Cargo.lock generated
View file

@ -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"

View file

@ -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 }

View file

@ -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<String, Stri
if let Some(path) = document_path {
command.arg(path);
}
let mut child = command
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
let mut child = spawn_null(&mut command)
.map_err(|e| format!("spawn {} --live-mcp: {e}", binary.display()))?;
// Only report a `documentPath` the editor will ACTUALLY open: its
@ -114,16 +110,20 @@ fn run_start_live(port: u16, document_path: Option<&str>) -> Result<String, Stri
// claim a document the editor never opened.
let opened = document_path.filter(|p| editor_will_open(p));
// Wait up to ~15s for the editor to bind + publish its port.
for _ in 0..150 {
if let Some((live_port, live_pid)) = reachable_live_port_file() {
match wait_for_child_or(&mut child, 150, Duration::from_millis(100), || {
reachable_live_port_file()
})
.map_err(|e| format!("wait for {} --live-mcp: {e}", binary.display()))?
{
WaitOutcome::Ready((live_port, live_pid)) => {
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<String,
// Per-instance token passed to the server (echoed in its `ping`) so a
// later `op stop` can prove the pid in our manager file owns the port.
let token = make_token();
let mut child = Command::new(&binary)
let mut command = Command::new(&binary);
command
.arg("--mcp-http")
.arg(port.to_string())
.arg(&document)
.env("OPENPENCIL_MCP_TOKEN", &token)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.env("OPENPENCIL_MCP_TOKEN", &token);
let mut child = spawn_null(&mut command)
.map_err(|e| format!("spawn {} --mcp-http: {e}", binary.display()))?;
let pid = child.id();
write_manager_files(pid, port, &token)?;
for _ in 0..30 {
// Confirm the MCP server actually answers WITH our token (not just
// an open port), so we never report success for a child that failed
// to bind or for an unrelated service on the port.
if crate::mcp_http_cli::mcp_ping_headless(port, &token) {
// Confirm the MCP server actually answers WITH our token (not just
// an open port), so we never report success for a child that failed
// to bind or for an unrelated service on the port.
match wait_for_child_or(&mut child, 30, Duration::from_millis(100), || {
crate::mcp_http_cli::mcp_ping_headless(port, &token).then_some(())
})
.map_err(|e| format!("wait for {} --mcp-http: {e}", binary.display()))?
{
WaitOutcome::Ready(()) => {
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<F: Fn() -> 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`.

View file

@ -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" }

View file

@ -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<std::sync::Mutex<String>> = 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;

View file

@ -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",
] }

View file

@ -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<T> {
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<Child> {
null_stdio(command).spawn()
}
/// Poll `probe` while also noticing if `child` exits first.
pub fn wait_for_child_or<T>(
child: &mut Child,
attempts: usize,
interval: Duration,
mut probe: impl FnMut() -> Option<T>,
) -> io::Result<WaitOutcome<T>> {
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<BufReader<ChildStdout>>;
/// Tokio child wrapper with piped stdin/stdout/stderr.
pub struct LineStreamChild {
child: tokio::process::Child,
stdin: Option<ChildStdin>,
lines: Option<LineStream>,
stderr: Option<ChildStderr>,
}
impl LineStreamChild {
/// Spawn `program` with piped stdio and the supplied args/envs.
pub fn spawn<P, A, E, K, V>(program: P, args: A, envs: E) -> io::Result<Self>
where
P: AsRef<OsStr>,
A: IntoIterator,
A::Item: AsRef<OsStr>,
E: IntoIterator<Item = (K, V)>,
K: AsRef<OsStr>,
V: AsRef<OsStr>,
{
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<Self> {
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<Option<String>> {
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<LineStream> {
self.lines.take()
}
/// Move stderr out so callers can drain or capture it separately.
pub fn take_stderr(&mut self) -> Option<ChildStderr> {
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<ExitStatus> {
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<ExitStatus> {
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())
}

View file

@ -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());
}