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 "); await patchPiRpcFile(target); }