5229059fe3
Proposal visibility is already global; this change makes finish, stop, and cancel shared queue-management actions so any user can unblock the single worker and shared queue. Confirm, adjust, follow_up, and start_next remain owner-only. Assistant/bootstrap copy, /help, /list, and runtime checks now align with that split.
733 lines
33 KiB
TypeScript
733 lines
33 KiB
TypeScript
import assert from "node:assert/strict";
|
||
import fs from "node:fs";
|
||
import os from "node:os";
|
||
import path from "node:path";
|
||
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";
|
||
|
||
const PNG_HEADER = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]);
|
||
|
||
async function tempPng(payload: string): Promise<{ dir: string; file: string }> {
|
||
const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-gateway-img-"));
|
||
const file = path.join(dir, "shot.png");
|
||
await fs.promises.writeFile(file, Buffer.concat([PNG_HEADER, Buffer.from(payload)]));
|
||
return { dir, file };
|
||
}
|
||
|
||
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;
|
||
resetResult = false;
|
||
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; }
|
||
async reset(_platform: string, _chatId: string, userId: string) { this.commandUsers.reset = userId; return this.resetResult; }
|
||
listProposals(_platform: string, _chatId: string, userId: string) { this.listed++; this.commandUsers.list = userId; return this.proposals; }
|
||
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) => {
|
||
const deadline = Date.now() + 5_000;
|
||
while (!condition()) {
|
||
if (Date.now() > deadline) break;
|
||
await flush();
|
||
await new Promise<void>((resolve) => setTimeout(resolve, 5));
|
||
}
|
||
assert.equal(condition(), true);
|
||
};
|
||
|
||
function recordingAdapter(sent: OutgoingMessage[] = [], supportsImages = false): PlatformAdapter {
|
||
return { name: "test", supportsImages, 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("sendEvent delivers the text first, then images, all sharing the msg_seq counter", async () => {
|
||
const { runtime, gateway } = createGateway();
|
||
const { dir, file } = await tempPng("payload-v1");
|
||
try {
|
||
const sent: OutgoingMessage[] = [];
|
||
const turn = gateway.receive(message("work", "img-event"), recordingAdapter(sent, true));
|
||
await waitFor(() => runtime.turns.length === 1);
|
||
runtime.turns[0].resolve("done");
|
||
await turn;
|
||
|
||
await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]);
|
||
const pieces = sent.map((outgoing) => [outgoing.text, outgoing.replySequence, outgoing.images?.length || 0]);
|
||
assert.deepEqual(pieces, [
|
||
["done", 1, 0],
|
||
["截图好了", 2, 0],
|
||
["", 3, 1]
|
||
]);
|
||
const image = sent[2]!.images![0]!;
|
||
assert.equal(image.mimeType, "image/png");
|
||
assert.equal(image.filename, "shot.png");
|
||
assert.equal(Buffer.from(image.data, "base64").toString("latin1"), Buffer.concat([PNG_HEADER, Buffer.from("payload-v1")]).toString("latin1"));
|
||
} finally {
|
||
await fs.promises.rm(dir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("a queued event re-reads its image at redelivery time and degrades when the file is gone", async () => {
|
||
let now = 1_000_000;
|
||
const runtime = new FakeRuntime();
|
||
const gateway = new Gateway(policy, runtime, { now: () => now });
|
||
const { dir, file } = await tempPng("payload-v1");
|
||
try {
|
||
const sent: OutgoingMessage[] = [];
|
||
const adapter = recordingAdapter(sent, true);
|
||
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", [{ path: file, filename: "shot.png" }]);
|
||
assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0);
|
||
|
||
// The worker updated the image after the event was queued: redelivery must send the new content.
|
||
await fs.promises.writeFile(file, Buffer.concat([PNG_HEADER, Buffer.from("payload-v2")]));
|
||
const second = gateway.receive(groupMessage("hello", "m2"), adapter);
|
||
await waitFor(() => runtime.turns.length === 2);
|
||
runtime.turns[1].resolve("ok");
|
||
await second;
|
||
const redelivered = sent.filter((outgoing) => outgoing.replyTo === "m2");
|
||
assert.deepEqual(redelivered.map((outgoing) => [outgoing.text, outgoing.replySequence, outgoing.images?.length || 0]), [
|
||
["stale event", 1, 0],
|
||
["", 2, 1],
|
||
["ok", 3, 0]
|
||
]);
|
||
assert.equal(Buffer.from(redelivered[1]!.images![0]!.data, "base64").toString("latin1"), Buffer.concat([PNG_HEADER, Buffer.from("payload-v2")]).toString("latin1"));
|
||
|
||
// Queue another event, then delete the file: the redelivery degrades to a text note.
|
||
now += 6 * 60_000;
|
||
await gateway.sendEvent("qq:group:g1", "another stale event", [{ path: file, filename: "shot.png" }]);
|
||
await fs.promises.rm(file);
|
||
const third = gateway.receive(groupMessage("hello again", "m3"), adapter);
|
||
await waitFor(() => runtime.turns.length === 3);
|
||
runtime.turns[2].resolve("ok again");
|
||
await third;
|
||
const degraded = sent.filter((outgoing) => outgoing.replyTo === "m3");
|
||
assert.equal(degraded.filter((outgoing) => outgoing.images?.length).length, 0);
|
||
assert.match(degraded[0]!.text, /another stale event/);
|
||
assert.match(degraded[0]!.text, /图片 shot\.png 发送失败/);
|
||
} finally {
|
||
await fs.promises.rm(dir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("adapters without image support receive images flattened into text lines", async () => {
|
||
const { runtime, gateway } = createGateway();
|
||
const { dir, file } = await tempPng("payload-v1");
|
||
try {
|
||
const sent: OutgoingMessage[] = [];
|
||
const turn = gateway.receive(message("work", "flat"), recordingAdapter(sent));
|
||
await waitFor(() => runtime.turns.length === 1);
|
||
runtime.turns[0].resolve("done");
|
||
await turn;
|
||
|
||
await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]);
|
||
const event = sent.find((outgoing) => outgoing.text.includes("截图好了"))!;
|
||
assert.equal(event.images, undefined);
|
||
assert.match(event.text, /截图好了/);
|
||
assert.match(event.text, /\[图片\] shot\.png/);
|
||
assert.equal(sent.filter((outgoing) => outgoing.images?.length).length, 0);
|
||
} finally {
|
||
await fs.promises.rm(dir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
test("an image send failure degrades to a text note without blocking the event", async () => {
|
||
const { runtime, gateway } = createGateway();
|
||
const { dir, file } = await tempPng("payload-v1");
|
||
try {
|
||
const sent: OutgoingMessage[] = [];
|
||
const errors: string[] = [];
|
||
const adapter: PlatformAdapter = {
|
||
name: "test", supportsImages: true, async handleWebhook() { return {}; },
|
||
async sendMessage(outgoing) {
|
||
if (outgoing.images?.length) throw new Error("sensitive provider response");
|
||
sent.push(outgoing);
|
||
}
|
||
};
|
||
const originalError = console.error;
|
||
console.error = (...args: unknown[]) => { errors.push(args.map(String).join(" ")); };
|
||
try {
|
||
const turn = gateway.receive(message("work", "img-fail"), adapter);
|
||
await waitFor(() => runtime.turns.length === 1);
|
||
runtime.turns[0].resolve("done");
|
||
await turn;
|
||
|
||
await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]);
|
||
} finally {
|
||
console.error = originalError;
|
||
}
|
||
const texts = sent.map((outgoing) => [outgoing.text, outgoing.replySequence]);
|
||
// The failed image attempt consumed msg_seq 3 (never reused in case QQ actually received it).
|
||
assert.deepEqual(texts, [
|
||
["done", 1],
|
||
["截图好了", 2],
|
||
["(图片 shot.png 发送失败)", 4]
|
||
]);
|
||
assert.deepEqual(errors, ["Gateway event image send failed (platform=qq)"]);
|
||
} finally {
|
||
await fs.promises.rm(dir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
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/);
|
||
assert.match(help.reply || "", /\/reset/);
|
||
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("/reset resets the current conversation session and reminds about unfinished proposals", async () => {
|
||
const { runtime, gateway } = createGateway();
|
||
const adapter = recordingAdapter();
|
||
|
||
const noSession = await gateway.receive(message("/reset"), adapter, { synchronous: true });
|
||
assert.equal(noSession.reply, "当前没有需要重置的会话。");
|
||
|
||
runtime.resetResult = true;
|
||
const reset = await gateway.receive(message("/reset"), adapter, { synchronous: true });
|
||
assert.match(reset.reply || "", /已重置当前对话/);
|
||
assert.doesNotMatch(reset.reply || "", /未完成任务/);
|
||
assert.equal(runtime.commandUsers.reset, "user");
|
||
|
||
runtime.proposals = [fakeProposal({ id: "p-1", status: "pending", pending: { summary: "s", receivedAt: 1 } })];
|
||
const reminded = await gateway.receive(message("/reset"), adapter, { synchronous: true });
|
||
assert.match(reminded.reply || "", /已重置当前对话/);
|
||
assert.match(reminded.reply || "", /你还有 1 个未完成任务,proposal 板不受影响/);
|
||
});
|
||
|
||
test("the list panel shows other members' proposals without exposing their identity", async () => {
|
||
const { runtime, gateway } = createGateway();
|
||
const adapter = recordingAdapter();
|
||
runtime.proposals = [
|
||
fakeProposal({ id: "p-mine", title: "我的任务", status: "working" }),
|
||
fakeProposal({ id: "p-other", title: "别人的任务", status: "pending", ownerChatKey: "qq:group:g-secret-9", requesterUserId: "someone-else", pending: { summary: "做到一半", receivedAt: 5 } })
|
||
];
|
||
|
||
const list = await gateway.receive(message("/list"), adapter, { synchronous: true });
|
||
const reply = list.reply || "";
|
||
assert.match(reply, /p-other「别人的任务」(其他成员):做到一半/);
|
||
assert.match(reply, /p-mine「我的任务」(你的)正在执行/);
|
||
assert.doesNotMatch(reply, /someone-else/);
|
||
assert.doesNotMatch(reply, /g-secret-9/);
|
||
});
|
||
|
||
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/);
|
||
});
|