Files
gori-agent/test/gateway.test.ts
T
zenord 7aa610294c feat: pending-finish proposals and QQ image input
Proposal semantics: worker results no longer distinguish success/fail;
only an explicit finish settles a pending proposal. Pending owner input
follows up by resuming the original worker session. Dirty pending blocks
start_next globally; clean pending only blocks its owner.

QQ adapter now downloads image attachments and passes them as ACP image
content blocks; video/file attachments degrade to link text.
2026-08-19 13:36:23 +08:00

546 lines
24 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import assert from "node:assert/strict";
import test from "node:test";
import type { ConversationRequest, ConversationRuntime } from "../src/acp/types.js";
import type { PlatformAdapter } from "../src/core/adapter.js";
import { Gateway } from "../src/core/gateway.js";
import type { Proposal } from "../src/core/proposal-store.js";
import type { IncomingMessage, OutgoingMessage } from "../src/core/types.js";
interface PendingTurn {
request: ConversationRequest;
resolve(text: string): void;
reject(error: Error): void;
}
function fakeProposal(overrides: Partial<Proposal> = {}): Proposal {
return {
id: "proposal-1",
title: "修复登录页",
goal: "goal",
steps: ["step"],
ownerChatKey: "qq:chat",
requesterUserId: "user",
status: "proposed",
createdAt: 1,
updatedAt: 1,
...overrides
};
}
class FakeRuntime implements ConversationRuntime {
readonly turns: PendingTurn[] = [];
cancelled = 0;
confirmed = 0;
finished = 0;
stopped = 0;
listed = 0;
cancelResult = true;
confirmResult = true;
finishResult = true;
stopResult = true;
proposals: Proposal[] = [];
commandUsers: Record<string, string | undefined> = {};
async prompt(request: ConversationRequest) {
return new Promise<{ text: string; botId: string; agentId: string }>((resolve, reject) => {
this.turns.push({
request,
resolve: (text) => resolve({ text, botId: "test-bot", agentId: "kimi" }),
reject
});
});
}
async cancel(_platform: string, _chatId: string, userId: string) { this.cancelled++; this.commandUsers.cancel = userId; return this.cancelResult; }
async confirm(_platform: string, _chatId: string, userId: string) { this.confirmed++; this.commandUsers.confirm = userId; return this.confirmResult; }
async finish(_platform: string, _chatId: string, userId: string) { this.finished++; this.commandUsers.finish = userId; return this.finishResult; }
async stop(_platform: string, _chatId: string, userId: string) { this.stopped++; this.commandUsers.stop = userId; return this.stopResult; }
listProposals(_platform: string, _chatId: string, userId: string) { this.listed++; this.commandUsers.list = userId; return this.proposals; }
async reset() {}
status(_platform?: string, _chatId?: string, userId?: string) {
this.commandUsers.status = userId;
return {
bot: "test-bot", agent: "kimi", workspace: "/tmp",
running: this.turns.length > 0, phase: "processing"
};
}
stats() { return { activeWorkers: 0, inFlight: 0, crashes: 0, persistedBindings: 0 }; }
async shutdown() {}
}
const policy = { allowedUsers: [] as string[], allowedChats: [] as string[], requireMentionInGroup: false };
const message = (text: string, messageId = text, chatId = "chat"): IncomingMessage => ({
platform: "qq", chatId, userId: "user", text, messageId
});
const groupMessage = (text: string, messageId: string): IncomingMessage => ({
platform: "qq", chatId: "group:g1", userId: "user", text, messageId, isGroup: true
});
const flush = async () => { await Promise.resolve(); await Promise.resolve(); await new Promise<void>((resolve) => setImmediate(resolve)); };
const waitFor = async (condition: () => boolean) => {
for (let count = 0; count < 50 && !condition(); count++) await flush();
assert.equal(condition(), true);
};
function recordingAdapter(sent: OutgoingMessage[] = []): PlatformAdapter {
return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } };
}
// Simulates QQ passive-reply dedup: a repeated msg_id + msg_seq pair is rejected.
function qqDedupAdapter(sent: OutgoingMessage[] = [], fail?: (outgoing: OutgoingMessage) => string | undefined): PlatformAdapter {
const used = new Set<string>();
return {
name: "test",
async handleWebhook() { return {}; },
async sendMessage(outgoing) {
const failure = fail?.(outgoing);
if (failure) throw new Error(failure);
if (outgoing.replyTo) {
const pair = `${outgoing.replyTo}|${outgoing.replySequence}`;
if (used.has(pair)) throw new Error("QQ send failed: HTTP 400 code=400304018 duplicate msg_id+msg_seq");
used.add(pair);
}
sent.push(outgoing);
}
};
}
function createGateway(runtime = new FakeRuntime()) {
return { runtime, gateway: new Gateway(policy, runtime) };
}
function messagesFor(sent: OutgoingMessage[], replyTo: string) {
return sent.filter((outgoing) => outgoing.replyTo === replyTo);
}
test("short asynchronous work sends only the real result with sequence one", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(message("work", "short"), recordingAdapter(sent));
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
assert.deepEqual(await turn, { ok: true, reply: "done" });
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [["done", 1]]);
});
test("ordinary asynchronous messages send nothing before the runtime resolves", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(message("work", "no-timer"), recordingAdapter(sent));
await waitFor(() => runtime.turns.length === 1);
await flush();
assert.deepEqual(sent, []);
runtime.turns[0].resolve("done");
await turn;
assert.deepEqual(sent.map((outgoing) => outgoing.text), ["done"]);
});
test("ordinary messages reach the runtime immediately without a Gateway queue", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const adapter = recordingAdapter(sent);
const first = gateway.receive(message("one", "first"), adapter);
const second = gateway.receive(message("two", "second"), adapter);
await waitFor(() => runtime.turns.length === 2);
assert.deepEqual(sent, []);
runtime.turns[0].resolve("first done");
runtime.turns[1].resolve("second done");
await Promise.all([first, second]);
assert.deepEqual(messagesFor(sent, "first").map((outgoing) => [outgoing.text, outgoing.replySequence]), [["first done", 1]]);
assert.deepEqual(messagesFor(sent, "second").map((outgoing) => [outgoing.text, outgoing.replySequence]), [["second done", 1]]);
});
test("delivery failures are logged safely and do not block the result", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const errors: string[] = [];
const adapter: PlatformAdapter = {
name: "test", async handleWebhook() { return {}; },
async sendMessage(outgoing) {
sent.push(outgoing);
throw new Error("sensitive provider response");
}
};
const originalError = console.error;
console.error = (...args: unknown[]) => { errors.push(args.map(String).join(" ")); };
try {
const turn = gateway.receive(message("work", "failure"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
assert.deepEqual(await turn, { ok: true, reply: "done" });
} finally {
console.error = originalError;
}
assert.deepEqual(sent.map((outgoing) => outgoing.text), ["done"]);
assert.deepEqual(errors, ["Gateway delivery send failed (platform=qq)"]);
});
test("runtime errors remain runtime errors and are delivered through the same stream", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(message("work", "runtime-error"), recordingAdapter(sent));
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].reject(new Error("failed"));
assert.deepEqual(await turn, { ok: false, error: "failed", reply: "Agent error: failed" });
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [["Agent error: failed", 1]]);
});
test("sendEvent uses a fresh passive window with replyTo and an incrementing sequence", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(message("work", "event-source"), recordingAdapter(sent));
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await turn;
await gateway.sendEvent("qq:chat", "任务完成了");
await gateway.sendEvent("qq:chat", "又完成了一步");
// The normal reply already consumed msg_seq 1 for "event-source"; events share the same counter.
const events = sent.filter((outgoing) => outgoing.text !== "done");
assert.deepEqual(events.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
["任务完成了", "event-source", 2],
["又完成了一步", "event-source", 3]
]);
assert.deepEqual(events[0]!.target, { platform: "qq", chatId: "chat", userId: "user", raw: undefined });
});
test("sendEvent uses the most recent target and adapter for the chat", async () => {
const { runtime, gateway } = createGateway();
const firstSent: OutgoingMessage[] = [];
const secondSent: OutgoingMessage[] = [];
const first = gateway.receive(message("one", "one", "chat"), recordingAdapter(firstSent));
await waitFor(() => runtime.turns.length === 1);
const second = gateway.receive(message("two", "two", "chat"), recordingAdapter(secondSent));
await waitFor(() => runtime.turns.length === 2);
runtime.turns[0].resolve("one done");
runtime.turns[1].resolve("two done");
await Promise.all([first, second]);
await gateway.sendEvent("qq:chat", "event");
assert.equal(firstSent.filter((outgoing) => outgoing.text === "event").length, 0);
assert.equal(secondSent.at(-1)?.text, "event");
assert.equal(secondSent.at(-1)?.replyTo, "two");
assert.equal(secondSent.at(-1)?.replySequence, 2); // "two done" consumed seq 1 for message "two"
});
test("sendEvent drops safely when the chat has no known target", async () => {
const { gateway } = createGateway();
await gateway.sendEvent("qq:unknown", "event");
});
test("expired group passive window skips the event and reminds on the next inbound", async () => {
let now = 1_000_000;
const runtime = new FakeRuntime();
const gateway = new Gateway(policy, runtime, { now: () => now });
const sent: OutgoingMessage[] = [];
const adapter = recordingAdapter(sent);
const turn = gateway.receive(groupMessage("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await turn;
now += 4 * 60_000; // within the 4.5 minute group window
await gateway.sendEvent("qq:group:g1", "fresh event");
assert.deepEqual(sent.at(-1), {
target: { platform: "qq", chatId: "group:g1", userId: "user", raw: undefined },
text: "fresh event", replyTo: "m1", replySequence: 2 // "done" consumed seq 1 for m1
});
now += 2 * 60_000; // past the group window
await gateway.sendEvent("qq:group:g1", "stale event");
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
const status = await gateway.receive(groupMessage("/status", "m2"), adapter, { synchronous: true });
assert.match(status.reply || "", /lastEventDelivery=skipped/);
assert.match(status.reply || "", /lastEventError=no fresh passive window/);
assert.match(status.reply || "", /lastEventAt=\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}/);
const reminder = gateway.receive(groupMessage("hello", "m3"), adapter);
await waitFor(() => runtime.turns.length === 2);
runtime.turns[1].resolve("ok");
await reminder;
const reminded = sent.find((outgoing) => outgoing.text === "stale event");
assert.equal(reminded?.replyTo, "m3");
assert.equal(reminded?.replySequence, 1);
// The stale event was already flushed once; a later inbound must not repeat it.
const again = gateway.receive(groupMessage("hello again", "m4"), adapter);
await waitFor(() => runtime.turns.length === 3);
runtime.turns[2].resolve("ok");
await again;
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 1);
});
test("status reports the last successful event delivery for the current chat", async () => {
const { runtime, gateway } = createGateway();
const adapter = recordingAdapter();
const turn = gateway.receive(message("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await turn;
await gateway.sendEvent("qq:chat", "event");
const status = await gateway.receive(message("/status", "m2"), adapter, { synchronous: true });
assert.match(status.reply || "", /lastEventDelivery=success/);
assert.match(status.reply || "", /lastEventError=none/);
assert.match(status.reply || "", /lastEventAt=\d{4}-\d{2}-\d{2}T/);
});
test("normal replies and fresh events share one msg_seq counter per inbound message", async () => {
const runtime = new FakeRuntime();
const gateway = new Gateway(policy, runtime);
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(groupMessage("work", "m1"), qqDedupAdapter(sent));
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await turn;
// With independent counters these events would reuse (m1, seq 1) and QQ would reject them.
await gateway.sendEvent("qq:group:g1", "event one");
await gateway.sendEvent("qq:group:g1", "event two");
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
["done", "m1", 1],
["event one", "m1", 2],
["event two", "m1", 3]
]);
});
test("a redelivered inbound message id continues its msg_seq counter instead of restarting", async () => {
const runtime = new FakeRuntime();
const gateway = new Gateway(policy, runtime);
const sent: OutgoingMessage[] = [];
const adapter = qqDedupAdapter(sent);
const first = gateway.receive(groupMessage("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("first done");
await first;
// QQ pushed the same msg_id again: the reply must not reuse (m1, seq 1).
const second = gateway.receive(groupMessage("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 2);
runtime.turns[1].resolve("second done");
await second;
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
["first done", "m1", 1],
["second done", "m1", 2]
]);
});
test("flushed pending events and the current reply share the msg_seq counter", async () => {
let now = 1_000_000;
const runtime = new FakeRuntime();
const gateway = new Gateway(policy, runtime, { now: () => now });
const sent: OutgoingMessage[] = [];
const adapter = qqDedupAdapter(sent);
const first = gateway.receive(groupMessage("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await first;
now += 6 * 60_000; // past the group window: the event is queued instead of sent
await gateway.sendEvent("qq:group:g1", "stale event");
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
const second = gateway.receive(groupMessage("hello", "m2"), adapter);
await waitFor(() => runtime.turns.length === 2);
runtime.turns[1].resolve("ok");
await second;
// The redelivery takes (m2, seq 1) and the current turn's reply takes (m2, seq 2); no pair collides.
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
["done", "m1", 1],
["stale event", "m2", 1],
["ok", "m2", 2]
]);
});
test("failed pending event redelivery stays queued, records failed, and retries on the next inbound", async () => {
let now = 1_000_000;
const runtime = new FakeRuntime();
const gateway = new Gateway(policy, runtime, { now: () => now });
const sent: OutgoingMessage[] = [];
let failEvent = true;
const adapter = qqDedupAdapter(sent, (outgoing) => failEvent && outgoing.text === "stale event"
? "QQ send failed: HTTP 400 code=40034105 proactive message not allowed"
: undefined);
const first = gateway.receive(groupMessage("work", "m1"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await first;
now += 6 * 60_000;
await gateway.sendEvent("qq:group:g1", "stale event"); // skipped, queued
const second = gateway.receive(groupMessage("hello", "m2"), adapter);
await waitFor(() => runtime.turns.length === 2);
runtime.turns[1].resolve("ok");
await second;
// The redelivery failed: nothing was delivered, but the normal reply still went out.
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [
["done", "m1", 1],
["ok", "m2", 2]
]);
const status = await gateway.receive(groupMessage("/status", "ms"), adapter, { synchronous: true });
assert.match(status.reply || "", /lastEventDelivery=failed/);
assert.match(status.reply || "", /lastEventError=HTTP 400 QQ code 40034105/);
failEvent = false;
const third = gateway.receive(groupMessage("hello again", "m3"), adapter);
await waitFor(() => runtime.turns.length === 3);
runtime.turns[2].resolve("ok again");
await third;
const redelivered = sent.filter((outgoing) => outgoing.text === "stale event");
assert.equal(redelivered.length, 1);
assert.equal(redelivered[0]!.replyTo, "m3");
assert.equal(redelivered[0]!.replySequence, 1);
});
test("sendEvent delivery failures are logged safely without message content", async () => {
const { runtime, gateway } = createGateway();
const errors: string[] = [];
const adapter: PlatformAdapter = {
name: "test", async handleWebhook() { return {}; },
async sendMessage(outgoing) {
if (outgoing.text.includes("event text")) throw new Error("sensitive provider response");
}
};
const turn = gateway.receive(message("work", "event-failure"), adapter);
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
await turn;
const originalError = console.error;
console.error = (...args: unknown[]) => { errors.push(args.map(String).join(" ")); };
try {
await gateway.sendEvent("qq:chat", "event text must not leak");
} finally {
console.error = originalError;
}
assert.deepEqual(errors, ["Gateway event send failed (platform=qq)"]);
});
test("status and help reach the runtime or reply directly", async () => {
const { runtime, gateway } = createGateway();
const adapter = recordingAdapter();
const turn = gateway.receive(message("work"), adapter, { synchronous: true });
await waitFor(() => runtime.turns.length === 1);
const status = await gateway.receive(message("/status"), adapter, { synchronous: true });
assert.match(status.reply || "", /running=true/);
assert.equal(runtime.commandUsers.status, "user");
const help = await gateway.receive(message("/help"), adapter, { synchronous: true });
assert.match(help.reply || "", /自然语言/);
assert.match(help.reply || "", /\/confirm/);
assert.match(help.reply || "", /\/finish/);
assert.match(help.reply || "", /\/stop/);
assert.match(help.reply || "", /\/cancel/);
runtime.turns[0].resolve("done");
await turn;
});
test("list, confirm, finish, stop, and cancel forward to the runtime with the chat identity", async () => {
const { runtime, gateway } = createGateway();
const adapter = recordingAdapter();
runtime.proposals = [fakeProposal()];
const list = await gateway.receive(message("/list"), adapter, { synchronous: true });
assert.equal(runtime.listed, 1);
assert.match(list.reply || "", /未确认(proposed)/);
assert.match(list.reply || "", /proposal-1「修复登录页」/);
runtime.proposals = [];
const emptyList = await gateway.receive(message("/list"), adapter, { synchronous: true });
assert.equal(emptyList.reply, "当前没有提案。");
const confirm = await gateway.receive(message("/confirm"), adapter, { synchronous: true });
assert.equal(runtime.confirmed, 1);
assert.equal(confirm.reply, "已确认,加入队列。");
const finish = await gateway.receive(message("/finish"), adapter, { synchronous: true });
assert.equal(runtime.finished, 1);
assert.equal(finish.reply, "已结束该任务。");
const stop = await gateway.receive(message("/stop"), adapter, { synchronous: true });
assert.equal(runtime.stopped, 1);
assert.equal(stop.reply, "已停止当前任务,任务转为待确认。");
const cancel = await gateway.receive(message("/cancel"), adapter, { synchronous: true });
assert.equal(runtime.cancelled, 1);
assert.equal(cancel.reply, "已取消最近的提案。");
assert.equal(runtime.commandUsers.list, "user");
assert.equal(runtime.commandUsers.confirm, "user");
assert.equal(runtime.commandUsers.finish, "user");
assert.equal(runtime.commandUsers.stop, "user");
assert.equal(runtime.commandUsers.cancel, "user");
runtime.confirmResult = false;
runtime.finishResult = false;
runtime.stopResult = false;
runtime.cancelResult = false;
assert.equal((await gateway.receive(message("/confirm"), adapter, { synchronous: true })).reply, "当前没有待确认的提案。");
assert.equal((await gateway.receive(message("/finish"), adapter, { synchronous: true })).reply, "当前没有待结束的任务。");
assert.equal((await gateway.receive(message("/stop"), adapter, { synchronous: true })).reply, "当前没有正在执行的任务。");
assert.equal((await gateway.receive(message("/cancel"), adapter, { synchronous: true })).reply, "没有可取消的提案。");
});
test("the list panel orders pending first, then working, queued, proposed, and recent finished", async () => {
const { runtime, gateway } = createGateway();
const adapter = recordingAdapter();
runtime.proposals = [
fakeProposal({ id: "p-finished", title: "旧任务", status: "finished", finishKind: "done", finishedAt: 10 }),
fakeProposal({ id: "p-proposed", title: "新想法", status: "proposed" }),
fakeProposal({ id: "p-working", title: "干活中", status: "working" }),
fakeProposal({ id: "p-queued", title: "排队任务", status: "queued" }),
fakeProposal({ id: "p-pending", title: "待确认任务", status: "pending", pending: { summary: "做了一半", question: "继续吗?", receivedAt: 5 } }),
fakeProposal({ id: "p-cancelled", title: "取消任务", status: "finished", finishKind: "cancelled", finishedAt: 20 })
];
const list = await gateway.receive(message("/list"), adapter, { synchronous: true });
const reply = list.reply || "";
const order = ["待确认(等你处理)", "进行中", "排队中", "未确认(proposed)", "最近结束"]
.map((section) => reply.indexOf(section));
assert.ok(order.every((index) => index >= 0), reply);
assert.deepEqual([...order].sort((a, b) => a - b), order);
assert.match(reply, /p-pending「待确认任务」:做了一半(问你:继续吗?)/);
assert.match(reply, /\/finish 结束/);
assert.match(reply, /p-working「干活中」正在执行/);
assert.match(reply, /p-cancelled「取消任务」已取消/);
assert.match(reply, /p-finished「旧任务」已完成/);
});
test("synchronous work sends no delivery messages", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
const turn = gateway.receive(message("work"), recordingAdapter(sent), { synchronous: true });
await waitFor(() => runtime.turns.length === 1);
runtime.turns[0].resolve("done");
assert.deepEqual(await turn, { ok: true, reply: "done" });
assert.deepEqual(sent, []);
});
test("asynchronous help replies with sequence one without reaching the runtime", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
await gateway.receive(message("/help", "help"), recordingAdapter(sent));
assert.deepEqual(messagesFor(sent, "help").map((outgoing) => outgoing.replySequence), [1]);
assert.equal(runtime.turns.length, 0);
});
test("asynchronous commands reply with sequence one without reaching prompt", async () => {
const { runtime, gateway } = createGateway();
const sent: OutgoingMessage[] = [];
runtime.proposals = [fakeProposal()];
const result = await gateway.receive(message("/list", "list"), recordingAdapter(sent));
assert.equal(result.ok, true);
assert.equal(runtime.turns.length, 0);
assert.deepEqual(messagesFor(sent, "list").map((outgoing) => outgoing.replySequence), [1]);
assert.match(messagesFor(sent, "list")[0].text, /proposal-1/);
});