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 { 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; stopped = 0; listed = 0; cancelResult = true; confirmResult = true; stopResult = true; proposals: Proposal[] = []; commandUsers: Record = {}; 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 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((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(); 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 || "", /\/stop/); assert.match(help.reply || "", /\/cancel/); runtime.turns[0].resolve("done"); await turn; }); test("list, confirm, 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 || "", /proposal-1 \[proposed\] 修复登录页/); 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 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.stop, "user"); assert.equal(runtime.commandUsers.cancel, "user"); runtime.confirmResult = false; runtime.stopResult = false; runtime.cancelResult = false; assert.equal((await gateway.receive(message("/confirm"), 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("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/); });