File size: 5,577 Bytes
ece5049
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
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);
}