e9f1d595fa
When an assistant turn emits both confirm and start_next, confirm may already start the newly queued proposal. The later start_next then sees the just-started task as an active blocker and appends a false running-task correction. Suppress that redundant start_next correction when confirm in the same turn already started the current user proposal. Also tighten the worker bootstrap wording to run confirmed low-risk work continuously until a real decision point.
225 lines
12 KiB
JavaScript
225 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 === "confirm and start next") response = assistantEnvelope("confirmed", [{ type: "confirm" }, { type: "start_next" }]);
|
|
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;
|