4f7b4a9115
While a worker turn is running, the assistant prompt for the working proposal's owner includes a compact status card (elapsed time, last tool activity category with age, per-category counts) built from the ACP session/update stream the runtime already receives. Other users still see only the desensitized busy state, and raw update content never enters the assistant session.
224 lines
12 KiB
JavaScript
224 lines
12 KiB
JavaScript
#!/usr/bin/env node
|
|
import { spawn } from "node:child_process";
|
|
import fs from "node:fs";
|
|
import path from "node:path";
|
|
import { Readable, Writable } from "node:stream";
|
|
import * as acp from "@agentclientprotocol/sdk";
|
|
|
|
const PNG_HEADER = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
|
|
|
|
const pending = new Map();
|
|
const logFile = process.env.FAKE_ACP_LOG;
|
|
const log = (entry) => { if (logFile) fs.appendFileSync(logFile, `${JSON.stringify(entry)}\n`); };
|
|
const assistantEnvelope = (reply, actions) => `${reply}\n<GORI_ASSISTANT_ACTION_V2>${JSON.stringify({ reply, actions })}</GORI_ASSISTANT_ACTION_V2>`;
|
|
const workerEnvelope = (result) => `worker reply\n<GORI_WORKER_RESULT_V2>${JSON.stringify(result)}</GORI_WORKER_RESULT_V2>`;
|
|
const workerEnvelopeTruncated = (result) => `worker reply\n<GORI_WORKER_RESULT_V2>${JSON.stringify(result)}`;
|
|
|
|
const app = acp.agent({ name: "fake-acp-agent" })
|
|
.onRequest(acp.methods.agent.initialize, ({ params }) => {
|
|
log({ method: "initialize" });
|
|
return {
|
|
protocolVersion: params.protocolVersion,
|
|
agentCapabilities: {
|
|
loadSession: true,
|
|
promptCapabilities: { image: process.env.FAKE_ACP_IMAGE_CAP !== "0" },
|
|
sessionCapabilities: { resume: {}, close: {}, list: {} }
|
|
},
|
|
agentInfo: { name: "fake-acp-agent", version: "1" }
|
|
};
|
|
})
|
|
.onRequest(acp.methods.agent.session.new, async ({ params }) => {
|
|
const sessionId = `fake-${Date.now()}-${Math.random().toString(16).slice(2)}`;
|
|
log({ method: "session/new", sessionId, cwd: params.cwd });
|
|
if (process.env.FAKE_ACP_NEW_HANG === "1") await new Promise(() => undefined);
|
|
return { sessionId };
|
|
})
|
|
.onRequest(acp.methods.agent.session.load, ({ params }) => {
|
|
log({ method: "session/load", sessionId: params.sessionId });
|
|
if (process.env.FAKE_ACP_LOAD_FAIL === "1" || (process.env.FAKE_ACP_LOAD_FAIL_FILE && fs.existsSync(process.env.FAKE_ACP_LOAD_FAIL_FILE))) throw new Error("load failed");
|
|
return {};
|
|
})
|
|
.onRequest(acp.methods.agent.session.resume, ({ params }) => {
|
|
log({ method: "session/resume", sessionId: params.sessionId });
|
|
if (process.env.FAKE_ACP_RESUME_FAIL === "1" || (process.env.FAKE_ACP_RESUME_FAIL_FILE && fs.existsSync(process.env.FAKE_ACP_RESUME_FAIL_FILE))) throw new Error("resume failed");
|
|
return {};
|
|
})
|
|
.onRequest(acp.methods.agent.session.close, ({ params }) => { log({ method: "session/close", sessionId: params.sessionId }); return {}; })
|
|
.onRequest(acp.methods.agent.session.list, () => ({ sessions: [] }))
|
|
.onRequest(acp.methods.agent.session.prompt, async ({ params, client, signal }) => {
|
|
const text = params.prompt.filter((item) => item.type === "text").map((item) => item.text).join("");
|
|
const images = params.prompt.filter((item) => item.type === "image").length;
|
|
log({ method: "session/prompt", sessionId: params.sessionId, text, images });
|
|
if (text === "crash" || text.startsWith("crash\n\nWhen this turn")) process.exit(9);
|
|
if (text.includes("Initialize this ACP session")) {
|
|
if (process.env.FAKE_ACP_BOOTSTRAP_DELAY_MS) await new Promise((resolve) => setTimeout(resolve, Number(process.env.FAKE_ACP_BOOTSTRAP_DELAY_MS)));
|
|
await update(client, params.sessionId, "READY");
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
|
|
// Assistant protocol
|
|
if (text.startsWith("Your previous response did not end with a valid GORI_ASSISTANT_ACTION_V2")) {
|
|
await update(client, params.sessionId, assistantEnvelope("repaired", []));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (text.startsWith("Your previous response did not end with a valid GORI_WORKER_RESULT_V2")) {
|
|
await update(client, params.sessionId, process.env.FAKE_ACP_INVALID_REPAIR === "1"
|
|
? "still invalid"
|
|
: workerEnvelope({ status: "PENDING", summary: "repaired" }));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (text.includes("[Internal event")) {
|
|
await update(client, params.sessionId, assistantEnvelope("event received: worker update", []));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (text.includes("[User message]")) {
|
|
const userText = (text.split("[User message]\n")[1] || "").split("\n\n")[0].trim();
|
|
if (userText === "force tool") {
|
|
await client.notify(acp.methods.client.session.update, {
|
|
sessionId: params.sessionId,
|
|
update: { sessionUpdate: "tool_call", toolCallId: "read-1", title: "Read a file", kind: "read", status: "in_progress" }
|
|
});
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
let response;
|
|
if (userText.startsWith("create proposal:")) {
|
|
const goal = userText.slice("create proposal:".length).trim();
|
|
response = assistantEnvelope("proposal drafted", [{ type: "create_proposal", title: "Test proposal", goal, steps: ["step 1"] }]);
|
|
} else if (userText === "confirm") response = assistantEnvelope("confirmed", [{ type: "confirm" }]);
|
|
else if (userText === "finish") response = assistantEnvelope("finishing", [{ type: "finish" }]);
|
|
else if (userText.startsWith("follow up:")) response = assistantEnvelope("following up", [{ type: "follow_up", instruction: userText.slice("follow up:".length).trim() }]);
|
|
else if (userText.startsWith("adjust proposal:")) response = assistantEnvelope("adjusting", [{ type: "adjust_proposal", title: userText.slice("adjust proposal:".length).trim() }]);
|
|
else if (userText.startsWith("send image:")) response = assistantEnvelope("sending image", [{ type: "send_image", path: userText.slice("send image:".length).trim() }]);
|
|
else if (userText === "start next") response = assistantEnvelope("starting next", [{ type: "start_next" }]);
|
|
else if (userText === "stop") response = assistantEnvelope("stopping", [{ type: "stop" }]);
|
|
else if (userText === "cancel proposal") response = assistantEnvelope("cancelling", [{ type: "cancel" }]);
|
|
else if (userText === "ack pending") response = assistantEnvelope("任务「Test proposal」还在等你确认。", []);
|
|
else if (userText === "truncated assistant") response = `ok\n<GORI_ASSISTANT_ACTION_V2>${JSON.stringify({ reply: "ok", actions: [] })}`;
|
|
else if (userText === "invalid assistant") response = "invalid without envelope";
|
|
else response = assistantEnvelope("ok", []);
|
|
await update(client, params.sessionId, response);
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
|
|
// Worker protocol
|
|
if (text.startsWith("The user sent a follow-up instruction")) {
|
|
await update(client, params.sessionId, workerEnvelope({ status: "PENDING", summary: `continued ${params.sessionId}` }));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (text.startsWith("Execute this confirmed proposal")) {
|
|
const goal = ((text.split("Goal: ")[1] || "").split("\n")[0] || "").trim();
|
|
if (goal.includes("toolhang")) {
|
|
await client.notify(acp.methods.client.session.update, {
|
|
sessionId: params.sessionId,
|
|
update: { sessionUpdate: "tool_call", toolCallId: "read-1", title: "Read package.json", kind: "read", status: "in_progress" }
|
|
});
|
|
await client.notify(acp.methods.client.session.update, {
|
|
sessionId: params.sessionId,
|
|
update: { sessionUpdate: "tool_call", toolCallId: "read-2", title: "Read README.md", kind: "read", status: "in_progress" }
|
|
});
|
|
await client.notify(acp.methods.client.session.update, {
|
|
sessionId: params.sessionId,
|
|
update: { sessionUpdate: "tool_call", toolCallId: "exec-1", title: "bash npm test", kind: "execute", status: "in_progress" }
|
|
});
|
|
await new Promise((resolve) => {
|
|
const done = () => resolve(undefined);
|
|
pending.set(params.sessionId, done);
|
|
signal.addEventListener("abort", done, { once: true });
|
|
});
|
|
pending.delete(params.sessionId);
|
|
return { stopReason: "cancelled" };
|
|
}
|
|
if (goal.includes("hang")) {
|
|
await new Promise((resolve) => {
|
|
const done = () => resolve(undefined);
|
|
pending.set(params.sessionId, done);
|
|
signal.addEventListener("abort", done, { once: true });
|
|
});
|
|
pending.delete(params.sessionId);
|
|
return { stopReason: "cancelled" };
|
|
}
|
|
if (goal.includes("invalid")) {
|
|
await update(client, params.sessionId, "invalid without envelope");
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (process.env.FAKE_ACP_WORKER_GATE_FILE) {
|
|
while (!fs.existsSync(process.env.FAKE_ACP_WORKER_GATE_FILE)) await new Promise((resolve) => setTimeout(resolve, 5));
|
|
}
|
|
if (goal.includes("truncbad")) {
|
|
await update(client, params.sessionId, `worker reply\n<GORI_WORKER_RESULT_V2>{"status":"PENDING","summary":"cut off mid json"`);
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (goal.includes("truncattach")) {
|
|
const outbox = path.join(process.cwd(), ".gori-outbox");
|
|
fs.mkdirSync(outbox, { recursive: true });
|
|
fs.writeFileSync(path.join(outbox, "shot.png"), Buffer.concat([PNG_HEADER, Buffer.from("fake-png-payload-truncated")]));
|
|
await update(client, params.sessionId, workerEnvelopeTruncated({ status: "PENDING", summary: "made a pic", attachments: [{ path: ".gori-outbox/shot.png", mimeType: "image/png" }] }));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (goal.includes("attachbad")) {
|
|
await update(client, params.sessionId, workerEnvelope({ status: "PENDING", summary: "made a pic", attachments: [{ path: "/tmp/evil.png" }, { path: "not-an-image.txt" }] }));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
if (goal.includes("attach")) {
|
|
const outbox = path.join(process.cwd(), ".gori-outbox");
|
|
fs.mkdirSync(outbox, { recursive: true });
|
|
fs.writeFileSync(path.join(outbox, "shot.png"), Buffer.concat([PNG_HEADER, Buffer.from("fake-png-payload-v1")]));
|
|
fs.writeFileSync(path.join(process.cwd(), "not-an-image.txt"), "plain text");
|
|
await update(client, params.sessionId, workerEnvelope({ status: "PENDING", summary: "made a pic", attachments: [{ path: ".gori-outbox/shot.png", mimeType: "image/png" }] }));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
const result = goal.includes("ask")
|
|
? { status: "PENDING", summary: "hit a dirty target", question: "May I overwrite it?", workspaceDirty: true }
|
|
: goal.includes("fail")
|
|
? { status: "PENDING", summary: "something broke" }
|
|
: { status: "PENDING", summary: "goal done" };
|
|
await update(client, params.sessionId, workerEnvelope(result));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
|
|
// Legacy behaviors used by acp-worker/acp-client tests
|
|
const request = text.split("\n\nWhen this turn is finished")[0];
|
|
if (request === "spawn descendant") {
|
|
const descendant = spawn(process.execPath, ["-e", "process.on('SIGTERM', () => {}); setInterval(() => {}, 1000)"], { stdio: "ignore" });
|
|
if (process.env.FAKE_ACP_DESCENDANT_PID_FILE) fs.writeFileSync(process.env.FAKE_ACP_DESCENDANT_PID_FILE, `${descendant.pid}\n`);
|
|
await new Promise((resolve) => signal.addEventListener("abort", resolve, { once: true }));
|
|
return { stopReason: "cancelled" };
|
|
}
|
|
if (request === "hang") {
|
|
await new Promise((resolve) => {
|
|
const done = () => resolve(undefined);
|
|
pending.set(params.sessionId, done);
|
|
signal.addEventListener("abort", done, { once: true });
|
|
});
|
|
pending.delete(params.sessionId);
|
|
return { stopReason: "cancelled" };
|
|
}
|
|
if (request === "activity") {
|
|
await update(client, params.sessionId, "first");
|
|
await new Promise((resolve) => setTimeout(resolve, 600));
|
|
await update(client, params.sessionId, "second");
|
|
await new Promise((resolve) => setTimeout(resolve, 700));
|
|
return { stopReason: "end_turn" };
|
|
}
|
|
await update(client, params.sessionId, `reply:${params.sessionId}:${request}`);
|
|
return { stopReason: "end_turn" };
|
|
})
|
|
.onNotification(acp.methods.agent.session.cancel, ({ params }) => {
|
|
log({ method: "session/cancel", sessionId: params.sessionId });
|
|
pending.get(params.sessionId)?.();
|
|
});
|
|
|
|
function update(client, sessionId, text) {
|
|
return client.notify(acp.methods.client.session.update, {
|
|
sessionId,
|
|
update: { sessionUpdate: "agent_message_chunk", content: { type: "text", text } }
|
|
});
|
|
}
|
|
|
|
const stream = acp.ndJsonStream(
|
|
Writable.toWeb(process.stdout),
|
|
Readable.toWeb(process.stdin)
|
|
);
|
|
const connection = app.connect(stream);
|
|
await connection.closed;
|