@@ -5,6 +5,7 @@
|
|||||||
### Fixed
|
### Fixed
|
||||||
|
|
||||||
- Fixed managed npm extension updates to avoid package managers installing or resolving pi host packages as peer dependencies ([#4907](https://github.com/earendil-works/pi/issues/4907)).
|
- Fixed managed npm extension updates to avoid package managers installing or resolving pi host packages as peer dependencies ([#4907](https://github.com/earendil-works/pi/issues/4907)).
|
||||||
|
- Fixed RPC mode raw stdout writes to retry transient backpressure errors and flush queued protocol output during shutdown ([#4897](https://github.com/earendil-works/pi/issues/4897)).
|
||||||
|
|
||||||
## [0.75.5] - 2026-05-23
|
## [0.75.5] - 2026-05-23
|
||||||
|
|
||||||
|
|||||||
@@ -6,6 +6,42 @@ interface StdoutTakeoverState {
|
|||||||
|
|
||||||
let stdoutTakeoverState: StdoutTakeoverState | undefined;
|
let stdoutTakeoverState: StdoutTakeoverState | undefined;
|
||||||
|
|
||||||
|
const RAW_STDOUT_RETRY_DELAY_MS = 10;
|
||||||
|
|
||||||
|
let rawStdoutWriteTail: Promise<void> = Promise.resolve();
|
||||||
|
|
||||||
|
function getRawStdoutWrite(): StdoutTakeoverState["rawStdoutWrite"] {
|
||||||
|
if (stdoutTakeoverState) {
|
||||||
|
return stdoutTakeoverState.rawStdoutWrite;
|
||||||
|
}
|
||||||
|
return process.stdout.write.bind(process.stdout) as StdoutTakeoverState["rawStdoutWrite"];
|
||||||
|
}
|
||||||
|
|
||||||
|
async function writeRawStdoutChunk(text: string): Promise<void> {
|
||||||
|
while (true) {
|
||||||
|
try {
|
||||||
|
await new Promise<void>((resolve, reject) => {
|
||||||
|
try {
|
||||||
|
getRawStdoutWrite()(text, (error) => {
|
||||||
|
if (error) reject(error);
|
||||||
|
else resolve();
|
||||||
|
});
|
||||||
|
} catch (error) {
|
||||||
|
reject(error instanceof Error ? error : new Error(String(error)));
|
||||||
|
}
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
} catch (error) {
|
||||||
|
const writeError = error instanceof Error ? error : new Error(String(error));
|
||||||
|
const code = (writeError as Error & { code?: unknown }).code;
|
||||||
|
if (code !== "ENOBUFS" && code !== "EAGAIN" && code !== "EWOULDBLOCK") {
|
||||||
|
throw writeError;
|
||||||
|
}
|
||||||
|
await new Promise<void>((resolve) => setTimeout(resolve, RAW_STDOUT_RETRY_DELAY_MS));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export function takeOverStdout(): void {
|
export function takeOverStdout(): void {
|
||||||
if (stdoutTakeoverState) {
|
if (stdoutTakeoverState) {
|
||||||
return;
|
return;
|
||||||
@@ -47,28 +83,26 @@ export function isStdoutTakenOver(): boolean {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export function writeRawStdout(text: string): void {
|
export function writeRawStdout(text: string): void {
|
||||||
if (stdoutTakeoverState) {
|
if (text.length === 0) {
|
||||||
stdoutTakeoverState.rawStdoutWrite(text);
|
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
process.stdout.write(text);
|
rawStdoutWriteTail = rawStdoutWriteTail.then(() => writeRawStdoutChunk(text));
|
||||||
|
void rawStdoutWriteTail.catch(() => {
|
||||||
|
process.exit(1);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function waitForRawStdoutBackpressure(): Promise<void> {
|
||||||
|
while (true) {
|
||||||
|
const tail = rawStdoutWriteTail;
|
||||||
|
await tail;
|
||||||
|
if (tail === rawStdoutWriteTail) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function flushRawStdout(): Promise<void> {
|
export async function flushRawStdout(): Promise<void> {
|
||||||
if (stdoutTakeoverState) {
|
await waitForRawStdoutBackpressure();
|
||||||
await new Promise<void>((resolve, reject) => {
|
await writeRawStdoutChunk("");
|
||||||
stdoutTakeoverState?.rawStdoutWrite("", (err) => {
|
|
||||||
if (err) reject(err);
|
|
||||||
else resolve();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
await new Promise<void>((resolve, reject) => {
|
|
||||||
process.stdout.write("", (err) => {
|
|
||||||
if (err) reject(err);
|
|
||||||
else resolve();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -19,7 +19,12 @@ import type {
|
|||||||
ExtensionWidgetOptions,
|
ExtensionWidgetOptions,
|
||||||
WorkingIndicatorOptions,
|
WorkingIndicatorOptions,
|
||||||
} from "../../core/extensions/index.ts";
|
} from "../../core/extensions/index.ts";
|
||||||
import { takeOverStdout, writeRawStdout } from "../../core/output-guard.ts";
|
import {
|
||||||
|
flushRawStdout,
|
||||||
|
takeOverStdout,
|
||||||
|
waitForRawStdoutBackpressure,
|
||||||
|
writeRawStdout,
|
||||||
|
} from "../../core/output-guard.ts";
|
||||||
import { killTrackedDetachedChildren } from "../../utils/shell.ts";
|
import { killTrackedDetachedChildren } from "../../utils/shell.ts";
|
||||||
import { type Theme, theme } from "../interactive/theme/theme.ts";
|
import { type Theme, theme } from "../interactive/theme/theme.ts";
|
||||||
import { attachJsonlLineReader, serializeJsonLine } from "./jsonl.ts";
|
import { attachJsonlLineReader, serializeJsonLine } from "./jsonl.ts";
|
||||||
@@ -49,6 +54,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
takeOverStdout();
|
takeOverStdout();
|
||||||
let session = runtimeHost.session;
|
let session = runtimeHost.session;
|
||||||
let unsubscribe: (() => void) | undefined;
|
let unsubscribe: (() => void) | undefined;
|
||||||
|
let unsubscribeBackpressure: (() => void) | undefined;
|
||||||
|
|
||||||
const output = (obj: RpcResponse | RpcExtensionUIRequest | object) => {
|
const output = (obj: RpcResponse | RpcExtensionUIRequest | object) => {
|
||||||
writeRawStdout(serializeJsonLine(obj));
|
writeRawStdout(serializeJsonLine(obj));
|
||||||
@@ -343,9 +349,13 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
});
|
});
|
||||||
|
|
||||||
unsubscribe?.();
|
unsubscribe?.();
|
||||||
|
unsubscribeBackpressure?.();
|
||||||
unsubscribe = session.subscribe((event) => {
|
unsubscribe = session.subscribe((event) => {
|
||||||
output(event);
|
output(event);
|
||||||
});
|
});
|
||||||
|
unsubscribeBackpressure = session.agent.subscribe(async () => {
|
||||||
|
await waitForRawStdoutBackpressure();
|
||||||
|
});
|
||||||
};
|
};
|
||||||
|
|
||||||
const registerSignalHandlers = (): void => {
|
const registerSignalHandlers = (): void => {
|
||||||
@@ -357,7 +367,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
for (const signal of signals) {
|
for (const signal of signals) {
|
||||||
const handler = () => {
|
const handler = () => {
|
||||||
killTrackedDetachedChildren();
|
killTrackedDetachedChildren();
|
||||||
void shutdown(signal === "SIGHUP" ? 129 : 143);
|
void shutdown(signal === "SIGHUP" ? 129 : 143, signal);
|
||||||
};
|
};
|
||||||
process.on(signal, handler);
|
process.on(signal, handler);
|
||||||
signalCleanupHandlers.push(() => process.off(signal, handler));
|
signalCleanupHandlers.push(() => process.off(signal, handler));
|
||||||
@@ -665,7 +675,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
*/
|
*/
|
||||||
let detachInput = () => {};
|
let detachInput = () => {};
|
||||||
|
|
||||||
async function shutdown(exitCode = 0): Promise<never> {
|
async function shutdown(exitCode = 0, signal?: NodeJS.Signals): Promise<never> {
|
||||||
if (shuttingDown) {
|
if (shuttingDown) {
|
||||||
process.exit(exitCode);
|
process.exit(exitCode);
|
||||||
}
|
}
|
||||||
@@ -674,9 +684,13 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
cleanup();
|
cleanup();
|
||||||
}
|
}
|
||||||
unsubscribe?.();
|
unsubscribe?.();
|
||||||
|
unsubscribeBackpressure?.();
|
||||||
await runtimeHost.dispose();
|
await runtimeHost.dispose();
|
||||||
detachInput();
|
detachInput();
|
||||||
process.stdin.pause();
|
process.stdin.pause();
|
||||||
|
if (signal !== "SIGTERM") {
|
||||||
|
await flushRawStdout();
|
||||||
|
}
|
||||||
process.exit(exitCode);
|
process.exit(exitCode);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -697,6 +711,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
`Failed to parse command: ${parseError instanceof Error ? parseError.message : String(parseError)}`,
|
`Failed to parse command: ${parseError instanceof Error ? parseError.message : String(parseError)}`,
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
|
await waitForRawStdoutBackpressure();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -721,6 +736,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
const response = await handleCommand(command);
|
const response = await handleCommand(command);
|
||||||
if (response) {
|
if (response) {
|
||||||
output(response);
|
output(response);
|
||||||
|
await waitForRawStdoutBackpressure();
|
||||||
}
|
}
|
||||||
await checkShutdownRequested();
|
await checkShutdownRequested();
|
||||||
} catch (commandError: unknown) {
|
} catch (commandError: unknown) {
|
||||||
@@ -731,6 +747,7 @@ export async function runRpcMode(runtimeHost: AgentSessionRuntime): Promise<neve
|
|||||||
commandError instanceof Error ? commandError.message : String(commandError),
|
commandError instanceof Error ? commandError.message : String(commandError),
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
|
await waitForRawStdoutBackpressure();
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|||||||
@@ -25,7 +25,9 @@ const rpcIo = vi.hoisted(() => ({
|
|||||||
}));
|
}));
|
||||||
|
|
||||||
vi.mock("../src/core/output-guard.js", () => ({
|
vi.mock("../src/core/output-guard.js", () => ({
|
||||||
|
flushRawStdout: vi.fn(async () => {}),
|
||||||
takeOverStdout: vi.fn(),
|
takeOverStdout: vi.fn(),
|
||||||
|
waitForRawStdoutBackpressure: vi.fn(async () => {}),
|
||||||
writeRawStdout: (line: string) => {
|
writeRawStdout: (line: string) => {
|
||||||
rpcIo.outputLines.push(line);
|
rpcIo.outputLines.push(line);
|
||||||
},
|
},
|
||||||
|
|||||||
Reference in New Issue
Block a user