diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index c990d2c8a..2c89988e3 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -108,6 +108,7 @@ ### Fixed +- Fixed RPC and RPC-UI startup crashes when an in-process extension claimed Bun's singleton stdin stream before the protocol reader ([#5898](https://github.com/can1357/oh-my-pi/issues/5898)). - Fixed loading issues for linked legacy extensions importing `DefaultPackageManager` or `linkedom`. - Fixed the advisor retrying terminal, non-retriable provider failures (e.g., blocked prompts), ensuring they fail immediately while transient failures still retry. - Fixed an issue where reassigning the `plan` role model mid-planning did not take effect until the next plan-mode entry; it now applies at the next turn boundary. diff --git a/packages/coding-agent/src/main.ts b/packages/coding-agent/src/main.ts index 7c0a066a8..c40784f9f 100644 --- a/packages/coding-agent/src/main.ts +++ b/packages/coding-agent/src/main.ts @@ -57,6 +57,7 @@ import { registerDaemonProjectPresence } from "./launch/presence"; import type { MCPManager } from "./mcp"; import { InteractiveMode } from "./modes/interactive-mode"; import type { PrintModeOptions } from "./modes/print-mode"; +import { claimRpcInput } from "./modes/rpc/rpc-input"; import { CURRENT_SETUP_VERSION } from "./modes/setup-version"; import { initTheme, stopThemeWatcher } from "./modes/theme/theme"; import type { SubmittedUserInput } from "./modes/types"; @@ -98,6 +99,7 @@ type RunRpcMode = ( session: AgentSession, setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void, eventBus?: EventBus, + input?: ReadableStream, ) => Promise; export function writeStartupNotice(parsedArgs: Pick, text: string): void { @@ -1111,6 +1113,9 @@ export async function runRootCommand( process.stderr.write(`${chalk.red("Error: @file arguments are not supported in RPC mode")}\n`); process.exit(1); } + const mode = parsedArgs.mode || "text"; + // RPC owns stdin. Claim its singleton stream before plugin/extension discovery can load an in-process consumer. + const rpcInput = mode === "rpc" || mode === "rpc-ui" ? claimRpcInput() : undefined; // Kick off plugin-root preload in parallel with the remaining startup work. // Awaited later (before extension/skill discovery in createAgentSession needs it). @@ -1157,7 +1162,6 @@ export async function runRootCommand( if (parsedArgs.noTitle || parsedArgs.mode === "rpc" || parsedArgs.mode === "rpc-ui" || parsedArgs.mode === "acp") { Bun.env.PI_NO_TITLE = "1"; } - const mode = parsedArgs.mode || "text"; const isProtocolMode = mode === "rpc" || mode === "rpc-ui" || mode === "acp"; // Protocol modes own stdin; treating it as prompt text would consume JSON-RPC frames before their transports start. const pipedInput = isProtocolMode ? undefined : await logger.time("readPipedInput", readPipedInput); @@ -1518,7 +1522,7 @@ export async function runRootCommand( // Branch-only protocol runner: keep RPC host code out of normal interactive startup. const runRpcMode: RunRpcMode = (await import("./modes/rpc/rpc-mode")).runRpcMode; stopStartupWatchdog(); - await runRpcMode(session, mode === "rpc-ui" ? setToolUIContext : undefined, eventBus); + await runRpcMode(session, mode === "rpc-ui" ? setToolUIContext : undefined, eventBus, rpcInput); } else if (isInteractive) { const versionCheckPromise = checkForNewVersion(VERSION).catch(() => undefined); const changelogMarkdown = await logger.time("main:getChangelogForDisplay", getChangelogForDisplay, parsedArgs); diff --git a/packages/coding-agent/src/modes/rpc/rpc-input.ts b/packages/coding-agent/src/modes/rpc/rpc-input.ts new file mode 100644 index 000000000..c99a92cd2 --- /dev/null +++ b/packages/coding-agent/src/modes/rpc/rpc-input.ts @@ -0,0 +1,38 @@ +/** + * Claims Bun's singleton stdin reader immediately and exposes a separately readable stream. + * RPC startup uses this before extension discovery so in-process modules cannot steal protocol input. + */ +export function claimRpcInput(): ReadableStream { + const reader = Bun.stdin.stream().getReader(); + let released = false; + const release = () => { + if (released) return; + released = true; + try { + reader.releaseLock(); + } catch {} + }; + return new ReadableStream({ + async pull(controller) { + try { + const result = await reader.read(); + if (result.done) { + release(); + controller.close(); + } else { + controller.enqueue(result.value); + } + } catch (error) { + release(); + controller.error(error); + } + }, + async cancel() { + try { + await reader.cancel(); + } finally { + release(); + } + }, + }); +} diff --git a/packages/coding-agent/src/modes/rpc/rpc-mode.ts b/packages/coding-agent/src/modes/rpc/rpc-mode.ts index 06cfd0acd..8f966c946 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-mode.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-mode.ts @@ -34,6 +34,7 @@ import type { EventBus } from "../../utils/event-bus"; import { initializeExtensions } from "../runtime-init"; import { isRpcHostToolResult, isRpcHostToolUpdate, RpcHostToolBridge } from "./host-tools"; import { isRpcHostUriResult, RpcHostUriBridge } from "./host-uris"; +import { claimRpcInput } from "./rpc-input"; import { RpcSubagentRegistry, readRpcSubagentTranscript } from "./rpc-subagents"; import type { RpcCommand, @@ -607,6 +608,7 @@ export async function runRpcMode( session: AgentSession, setToolUIContext?: (uiContext: ExtensionUIContext, hasUI: boolean) => void, eventBus?: EventBus, + input: ReadableStream = claimRpcInput(), ): Promise { // Signal to RPC clients that the server is ready to accept commands // Suppress terminal notifications: they write \x07 (BEL) or OSC sequences directly to @@ -1381,7 +1383,7 @@ export async function runRpcMode( // line is reported as an error frame and the loop keeps running instead of // throwing out of the generator and killing the whole process (issue #5194). const decoder = new TextDecoder(); - for await (const line of readLines(Bun.stdin.stream())) { + for await (const line of readLines(input ?? Bun.stdin.stream())) { const text = decoder.decode(line).trim(); if (!text) continue; let parsed: unknown; diff --git a/packages/coding-agent/test/fixtures/locked-stdin-reader.ts b/packages/coding-agent/test/fixtures/locked-stdin-reader.ts new file mode 100644 index 000000000..65e5ccecb --- /dev/null +++ b/packages/coding-agent/test/fixtures/locked-stdin-reader.ts @@ -0,0 +1,4 @@ +const lockedStdinReader = Bun.stdin.stream().getReader(); +void lockedStdinReader; + +export default function lockedStdinReaderExtension(): void {} diff --git a/packages/coding-agent/test/rpc-stdin-lock.test.ts b/packages/coding-agent/test/rpc-stdin-lock.test.ts new file mode 100644 index 000000000..46b610532 --- /dev/null +++ b/packages/coding-agent/test/rpc-stdin-lock.test.ts @@ -0,0 +1,60 @@ +import { describe, expect, test } from "bun:test"; +import * as path from "node:path"; +import { isRecord, readJsonl } from "@oh-my-pi/pi-utils"; + +async function expectRpcModeOwnsStdin(mode: "rpc" | "rpc-ui"): Promise { + const cliPath = path.join(import.meta.dir, "..", "src", "cli.ts"); + const extensionPath = path.join(import.meta.dir, "fixtures", "locked-stdin-reader.ts"); + const child = Bun.spawn( + [ + "bun", + cliPath, + "--extension", + extensionPath, + "--mode", + mode, + "--provider", + "anthropic", + "--model", + "claude-sonnet-4-5", + ], + { + cwd: path.join(import.meta.dir, ".."), + env: { ...Bun.env, PI_NO_TITLE: "1" }, + stdin: "pipe", + stdout: "pipe", + stderr: "pipe", + }, + ); + const stderrPromise = new Response(child.stderr).text(); + + child.stdin.write(`${JSON.stringify({ type: "get_state", id: "probe" })}\n`); + await child.stdin.flush(); + + let stateResponse: Record | undefined; + try { + for await (const frame of readJsonl(child.stdout as ReadableStream)) { + if (isRecord(frame) && frame.type === "response" && frame.id === "probe") { + stateResponse = frame; + break; + } + } + } finally { + child.stdin.end(); + child.kill(); + await child.exited.catch(() => {}); + } + + const stderr = await stderrPromise; + expect(stderr).not.toContain("ReadableStream is locked"); + expect(stateResponse?.success).toBe(true); +} + +describe("RPC mode stdin ownership", () => { + test("rpc claims stdin before extensions can lock its singleton stream", () => expectRpcModeOwnsStdin("rpc"), 30000); + test( + "rpc-ui claims stdin before extensions can lock its singleton stream", + () => expectRpcModeOwnsStdin("rpc-ui"), + 30000, + ); +});