chat-ui-agent-runtime / patch-pi-rpc.mjs
Mike0021's picture
Harden reconnect-safe Agent runtime
ece5049 verified
Raw History Blame Contribute Delete
5.58 kB
import { readFile, writeFile } from "node:fs/promises";
import { fileURLToPath } from "node:url";
import { resolve } from "node:path";
const QUEUE_PATCH_MARKER = 'case "replace_queue"';
const CORRELATED_UI_PATCH_MARKER = "const requestId = response.requestId";
const FOLLOW_UP_HANDLER = ` case "follow_up": {
await session.followUp(command.message, command.images);
return success(id, "follow_up");
}`;
const REPLACE_QUEUE_HANDLER = `${FOLLOW_UP_HANDLER}
case "replace_queue": {
// chat-ui extension for editable pending chips. This is one
// synchronous compare-and-swap so a stale browser snapshot
// cannot overwrite a queue that Pi already consumed or another
// tab changed.
const validQueue = (value) => Array.isArray(value)
&& value.length <= 50
&& value.every((item) => typeof item === "string" && item.length > 0 && item.length <= 1024);
const expected = command.expected;
const replacement = command.replacement;
if (!expected || !replacement
|| !validQueue(expected.steering) || !validQueue(expected.followUp)
|| !validQueue(replacement.steering) || !validQueue(replacement.followUp)) {
return error(id, "replace_queue", "Invalid queue snapshot");
}
const currentSteering = [...session.getSteeringMessages()];
const currentFollowUp = [...session.getFollowUpMessages()];
const unchanged = JSON.stringify(currentSteering) === JSON.stringify(expected.steering)
&& JSON.stringify(currentFollowUp) === JSON.stringify(expected.followUp);
if (!unchanged) {
return error(id, "replace_queue", "Pending queue changed; refresh and try again");
}
session.clearQueue();
// These pinned 0.84.2 methods have synchronous bodies. Calling
// them without awaiting keeps clear + refill atomic on Node's
// event loop while updating both AgentSession and core queues.
for (const message of replacement.steering) void session._queueSteer(message);
for (const message of replacement.followUp) void session._queueFollowUp(message);
return success(id, "replace_queue", replacement);
}`;
const EXTENSION_UI_RESPONSE_HANDLER = ` // Handle extension UI responses
if (typeof parsed === "object" &&
parsed !== null &&
"type" in parsed &&
parsed.type === "extension_ui_response") {
const response = parsed;
const pending = pendingExtensionRequests.get(response.id);
if (pending) {
pendingExtensionRequests.delete(response.id);
pending.resolve(response);
}
return;
}`;
const CORRELATED_EXTENSION_UI_RESPONSE_HANDLER = ` // Handle extension UI responses. chat-ui supplies a stable command id
// plus the native request id so retries receive a semantic success/error.
if (typeof parsed === "object" &&
parsed !== null &&
"type" in parsed &&
parsed.type === "extension_ui_response") {
const response = parsed;
const requestId = response.requestId ?? response.id;
const correlationId = response.requestId === undefined ? undefined : response.id;
const pending = pendingExtensionRequests.get(requestId);
if (pending) {
pendingExtensionRequests.delete(requestId);
pending.resolve({ ...response, id: requestId });
if (correlationId !== undefined) {
output(success(correlationId, "extension_ui_response"));
await waitForRawStdoutBackpressure();
}
}
else if (correlationId !== undefined) {
output(error(correlationId, "extension_ui_response", "Extension UI request is no longer pending"));
await waitForRawStdoutBackpressure();
}
return;
}`;
export function patchPiRpcSource(source) {
let patched = source;
if (!patched.includes(QUEUE_PATCH_MARKER)) {
const occurrences = patched.split(FOLLOW_UP_HANDLER).length - 1;
if (occurrences !== 1) {
throw new Error(`Expected one pinned Pi follow_up handler, found ${occurrences}`);
}
patched = patched.replace(FOLLOW_UP_HANDLER, REPLACE_QUEUE_HANDLER);
}
if (!patched.includes(CORRELATED_UI_PATCH_MARKER)) {
const occurrences = patched.split(EXTENSION_UI_RESPONSE_HANDLER).length - 1;
if (occurrences !== 1) {
throw new Error(`Expected one pinned Pi extension UI response handler, found ${occurrences}`);
}
patched = patched.replace(
EXTENSION_UI_RESPONSE_HANDLER,
CORRELATED_EXTENSION_UI_RESPONSE_HANDLER
);
}
return patched;
}
export async function patchPiRpcFile(path) {
const source = await readFile(path, "utf8");
const patched = patchPiRpcSource(source);
if (patched !== source) await writeFile(path, patched, "utf8");
}
const invokedPath = process.argv[1] ? resolve(process.argv[1]) : undefined;
if (invokedPath === fileURLToPath(import.meta.url)) {
const target = process.argv[2];
if (!target) throw new Error("Usage: node patch-pi-rpc.mjs <rpc-mode.js>");
await patchPiRpcFile(target);
}