From 680a449cd8804c3c81edbc855c8297a98c9ec346 Mon Sep 17 00:00:00 2001 From: zenord Date: Mon, 17 Aug 2026 15:35:18 +0800 Subject: [PATCH] feat: delay long-running task notices --- AGENTS.md | 5 +- README.md | 4 +- src/core/gateway.ts | 238 ++++++++++++++++--- src/server.ts | 1 + test/gateway.test.ts | 495 ++++++++++++++++++++++++++++++---------- test/qq-adapter.test.ts | 10 +- 6 files changed, 589 insertions(+), 164 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index d27b51c..3149292 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -89,8 +89,9 @@ gori-agent list - `src/core/gateway.ts` - 入站 allowlist、群 mention、命令、per-chat lock、队列状态和回复。 - policy 后 `/status`、`/cancel`、`/help` 必须绕过 chat lock,`/new` 必须保持串行。 - - 异步平台普通消息入队即发送中文处理/排队确认;对应 prompt 等确认完成或安全失败后才开始,排队确认不等待前项;确认失败只写安全日志且不阻止任务,同步 webhook 不发送确认。 - - QQ 同一入站 `msg_id` 的确认使用 `msg_seq=1`,最终/错误回复使用 `msg_seq=2`;单条命令回复使用 `msg_seq=1`。 + - 异步平台普通消息轮到后立即 prompt,不等待提示发送;15 秒内完成只发结果,否则按 queued/running 提示,queued 转 running 时补开始提示,并在入站后 60 秒、180 秒及之后每 300 秒结合 ACP `idleSeconds` 低频提醒;同步 webhook 不提示。 + - 每条异步入站使用独立串行 ReplyStream,`replySequence` 从 1 动态递增;发送失败只写安全日志且不阻止任务或后续发送,runtime/delivery 错误分离,最终消息入流后释放 chat lock。 + - 每条入站只维护一个基于绝对 deadline 的可取消 timer;完成先标记 finished 并清 timer,shutdown 清理全部 timer,生产 timer 必须 `unref`,Gateway 可注入 fake clock/schedule 供测试。 - `src/core/command-router.ts` - `/help`、`/status`、`/cancel`、`/new`;旧 role 命令返回 retired 提示。 - `src/platforms/*` diff --git a/README.md b/README.md index 1eb110d..ea335fd 100644 --- a/README.md +++ b/README.md @@ -212,7 +212,9 @@ Bot fingerprint 包含 Bot ID、workspace、persona、agent、permissions、skil - `/new` 保持串行,清除当前 chat binding;下一条普通消息创建新 native session。 - `/roles`、`/role`、`/agents`、`/agent` 会返回固定 retired 提示,不会转发给 ACP。 -同一 chat 的普通消息和 `/new` 串行执行,不同 chat 可并发。异步平台的普通消息通过 policy 后会立即确认:空闲时发送“已收到,正在处理。可随时发送 /status 查看状态。”,已有同 chat 工作时发送“已收到,已排队。可随时发送 /status 查看状态。”。每条消息的 ACP prompt 会等待该确认完成或安全失败后再开始;排队消息的确认不等待前项。确认发送失败只写不含消息正文或 provider 错误详情的安全日志,不阻止已入队任务;同步 webhook 不额外发送确认。QQ 对同一入站 `msg_id` 的确认使用 `msg_seq=1`,最终或错误回复使用 `msg_seq=2`;单条命令回复使用 `msg_seq=1`。 +同一 chat 的普通消息和 `/new` 串行执行,不同 chat 可并发。异步平台的普通消息轮到后立即开始 ACP prompt,不等待状态提示发送:15 秒内完成只发送真实结果;15 秒仍未完成时,running 发送“这条还在处理,完成后我会直接回复。”,queued 发送“前面还有任务,这条还在排队,轮到后我马上处理。”;发送过 queued 提示的消息真正开始时再发送“轮到这条了,我开始处理。”。入站后 60 秒、180 秒及之后每 300 秒按 queued/running 和 ACP `idleSeconds` 发送低频口语化提醒;延迟执行的 timer 不补发历史提醒。同步 webhook 不发送这些提示。 + +每条异步入站有独立、串行的回复流,`replySequence` 从 1 动态递增;因此短任务最终回复使用 `msg_seq=1`,已发送提示时后续消息使用下一序号。平台发送失败只记录不含消息正文或 provider 错误详情的安全日志,不影响 ACP 任务或回复流中的后续发送,也不会把成功的 ACP 结果误报为 `Agent error`。最终结果加入回复流后即释放 per-chat lock,平台发送延迟不会阻塞下一项任务;完成与 shutdown 都会清理状态提醒 timer。 ## HTTP 端点 diff --git a/src/core/gateway.ts b/src/core/gateway.ts index e37e985..de36260 100644 --- a/src/core/gateway.ts +++ b/src/core/gateway.ts @@ -7,18 +7,75 @@ import type { IncomingMessage } from "./types.js"; export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string } +export interface GatewayTimer { cancel(): void } + +export interface GatewayOptions { + now?: () => number; + schedule?: (callback: () => void, delayMs: number) => GatewayTimer; +} + interface ChatWorkState { running: boolean; queued: number; startedAt?: number; } +interface AsyncInbound { + phase: "queued" | "running" | "finished"; + receivedAt: number; + initialNoticeSent: boolean; + queuedNoticeSent: boolean; + timer?: GatewayTimer; + stream: ReplyStream; +} + +interface LockedResult { + result: GatewayResult; + delivery: Promise; +} + +class ReplyStream { + private sequence = 0; + private tail = Promise.resolve(); + + constructor(private readonly message: IncomingMessage, private readonly adapter: PlatformAdapter) {} + + enqueue(text: string): Promise { + const replySequence = ++this.sequence; + this.tail = this.tail.then(() => this.adapter.sendMessage({ + target: { + platform: this.message.platform, + chatId: this.message.chatId, + userId: this.message.userId, + raw: this.message.raw + }, + text, + replyTo: this.message.messageId, + replySequence + })).catch(() => { + console.error(`Gateway delivery send failed (platform=${this.message.platform})`); + }); + return this.tail; + } +} + export class Gateway { private readonly locks = new Map>(); private readonly chatWork = new Map(); + private readonly asyncInbounds = new Set(); + private readonly now: () => number; + private readonly schedule: (callback: () => void, delayMs: number) => GatewayTimer; + private closed = false; readonly commandRouter = new CommandRouter(); - constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime) {} + constructor( + private readonly policy: GatewayPolicy, + private readonly runtime: ConversationRuntime, + options: GatewayOptions = {} + ) { + this.now = options.now || Date.now; + this.schedule = options.schedule || defaultSchedule; + } async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise { const policyError = this.checkPolicy(message); @@ -29,33 +86,92 @@ export class Gateway { const command = this.commandRouter.parse(message.text); if (command?.kind === "help" || command?.kind === "status" || command?.kind === "cancel") { - return this.reply(message, adapter, await this.executeCommand(command, message), options); + return this.executeAndReply(command, message, adapter, options); } const chatKey = chatKeyFor(message.platform, message.chatId); - const existing = this.chatWork.get(chatKey); - const wasBusy = Boolean(existing?.running || existing?.queued); - const acknowledgement = !command && !options.synchronous - ? this.sendAcknowledgement(message, adapter, wasBusy) - : Promise.resolve(); - return this.withChatLock(chatKey, async () => { - await acknowledgement; - const replySequence = !command && !options.synchronous ? 2 : 1; + if (command || options.synchronous) { + const result = await this.withChatLock(chatKey, () => this.executeSynchronous(command, message)); + if (!options.synchronous) await new ReplyStream(message, adapter).enqueue(result.reply || ""); + return result; + } + + const inbound: AsyncInbound = { + phase: "queued", + receivedAt: this.now(), + initialNoticeSent: false, + queuedNoticeSent: false, + stream: new ReplyStream(message, adapter) + }; + this.asyncInbounds.add(inbound); + this.armTimer(inbound, inbound.receivedAt + 15_000, message); + + const locked = await this.withChatLock(chatKey, async () => { + inbound.phase = "running"; + if (inbound.queuedNoticeSent) void inbound.stream.enqueue("轮到这条了,我开始处理。"); try { - const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({ - platform: message.platform, chatId: message.chatId, userId: message.userId, text: message.text, messageId: message.messageId + const reply = (await this.runtime.prompt({ + platform: message.platform, + chatId: message.chatId, + userId: message.userId, + text: message.text, + messageId: message.messageId })).text; - return this.reply(message, adapter, reply, options, replySequence); + this.finishInbound(inbound); + return { result: { ok: true, reply }, delivery: inbound.stream.enqueue(reply) }; } catch (error) { const errorText = error instanceof Error ? error.message : String(error); const reply = `Agent error: ${errorText}`; - if (!options.synchronous) await this.send(message, adapter, reply, replySequence); - return { ok: false, error: errorText, reply }; + this.finishInbound(inbound); + return { result: { ok: false, error: errorText, reply }, delivery: inbound.stream.enqueue(reply) }; } }); + await locked.delivery; + return locked.result; } - stats(): ReturnType & { lockedChats: number } { return { ...this.runtime.stats(), lockedChats: this.locks.size }; } + stats(): ReturnType & { lockedChats: number } { + return { ...this.runtime.stats(), lockedChats: this.locks.size }; + } + + shutdown(): void { + if (this.closed) return; + this.closed = true; + for (const inbound of this.asyncInbounds) { + inbound.phase = "finished"; + inbound.timer?.cancel(); + inbound.timer = undefined; + } + this.asyncInbounds.clear(); + } + + private async executeAndReply( + command: ParsedCommand, + message: IncomingMessage, + adapter: PlatformAdapter, + options: { synchronous?: boolean } + ): Promise { + const result = await this.executeSynchronous(command, message); + if (options.synchronous) return result; + await new ReplyStream(message, adapter).enqueue(result.reply || ""); + return result; + } + + private async executeSynchronous(command: ParsedCommand | undefined, message: IncomingMessage): Promise { + try { + const reply = command ? await this.executeCommand(command, message) : (await this.runtime.prompt({ + platform: message.platform, + chatId: message.chatId, + userId: message.userId, + text: message.text, + messageId: message.messageId + })).text; + return { ok: true, reply }; + } catch (error) { + const errorText = error instanceof Error ? error.message : String(error); + return { ok: false, error: errorText, reply: `Agent error: ${errorText}` }; + } + } private async executeCommand(command: ParsedCommand, message: IncomingMessage): Promise { switch (command.kind) { @@ -67,7 +183,8 @@ export class Gateway { ...this.runtime.status(message.platform, message.chatId), gatewayRunning: Boolean(work?.running), queued: work?.queued || 0, - gatewayRunningSeconds: work?.startedAt ? Math.max(0, Math.floor((Date.now() - work.startedAt) / 1000)) : 0 + gatewayRunningSeconds: work?.startedAt !== undefined + ? Math.max(0, Math.floor((this.now() - work.startedAt) / 1000)) : 0 }; return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`; } @@ -83,26 +200,55 @@ export class Gateway { return await this.runtime.cancel(message.platform, message.chatId) ? "Cancellation requested." : "No active turn to cancel."; } - private async reply(message: IncomingMessage, adapter: PlatformAdapter, reply: string, options: { synchronous?: boolean }, replySequence = 1): Promise { - if (!options.synchronous) await this.send(message, adapter, reply, replySequence); - return { ok: true, reply }; + private armTimer(inbound: AsyncInbound, deadline: number, message: IncomingMessage): void { + if (this.closed || inbound.phase === "finished") return; + inbound.timer?.cancel(); + inbound.timer = this.schedule(() => { + inbound.timer = undefined; + if (this.closed || inbound.phase === "finished") return; + const now = this.now(); + if (now < deadline) { + this.armTimer(inbound, deadline, message); + return; + } + if (!inbound.initialNoticeSent) { + inbound.initialNoticeSent = true; + inbound.queuedNoticeSent = inbound.phase === "queued"; + void inbound.stream.enqueue(inbound.phase === "queued" + ? "前面还有任务,这条还在排队,轮到后我马上处理。" + : "这条还在处理,完成后我会直接回复。"); + } else { + void inbound.stream.enqueue(this.reminderText(inbound, message)); + } + this.armTimer(inbound, nextReminderDeadline(inbound.receivedAt, now), message); + }, Math.max(0, deadline - this.now())); } - private send(message: IncomingMessage, adapter: PlatformAdapter, text: string, replySequence: number): Promise { - return adapter.sendMessage({ - target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, - text, - replyTo: message.messageId, - replySequence - }); + private reminderText(inbound: AsyncInbound, message: IncomingMessage): string { + let idleSeconds: number | undefined; + try { + const value = this.runtime.status(message.platform, message.chatId).idleSeconds; + if (typeof value === "number" && Number.isFinite(value) && value >= 0) idleSeconds = Math.floor(value); + } catch { + // A status sampling failure must not affect the active turn. + } + const activity = activityText(idleSeconds); + if (inbound.phase === "queued") { + return activity + ? `前面的任务还在处理,${activity};这条仍在排队,轮到后我马上处理。` + : "前面的任务还在处理,这条仍在排队,轮到后我马上处理。"; + } + return activity + ? `这条还在处理中,${activity};完成后我会直接回复。` + : "这条还在处理中,完成后我会直接回复。"; } - private async sendAcknowledgement(message: IncomingMessage, adapter: PlatformAdapter, queued: boolean): Promise { - const text = queued - ? "已收到,已排队。可随时发送 /status 查看状态。" - : "已收到,正在处理。可随时发送 /status 查看状态。"; - try { await this.send(message, adapter, text, 1); } - catch { console.error(`Gateway acknowledgement send failed (platform=${message.platform})`); } + private finishInbound(inbound: AsyncInbound): void { + if (inbound.phase === "finished") return; + inbound.phase = "finished"; + inbound.timer?.cancel(); + inbound.timer = undefined; + this.asyncInbounds.delete(inbound); } private checkPolicy(message: IncomingMessage): string | undefined { @@ -124,8 +270,10 @@ export class Gateway { await previous; state.queued--; state.running = true; - state.startedAt = Date.now(); - try { return await fn(); } finally { + state.startedAt = this.now(); + try { + return await fn(); + } finally { state.running = false; state.startedAt = undefined; release(); @@ -136,3 +284,23 @@ export class Gateway { } } } + +function defaultSchedule(callback: () => void, delayMs: number): GatewayTimer { + const timer = setTimeout(callback, delayMs); + timer.unref(); + return { cancel: () => clearTimeout(timer) }; +} + +function nextReminderDeadline(receivedAt: number, now: number): number { + const elapsed = now - receivedAt; + if (elapsed < 60_000) return receivedAt + 60_000; + if (elapsed < 180_000) return receivedAt + 180_000; + return receivedAt + 180_000 + (Math.floor((elapsed - 180_000) / 300_000) + 1) * 300_000; +} + +function activityText(idleSeconds: number | undefined): string | undefined { + if (idleSeconds === undefined) return undefined; + if (idleSeconds < 30) return "刚刚还有进展"; + if (idleSeconds < 120) return `最近 ${idleSeconds} 秒没新进展`; + return `最近 ${Math.floor(idleSeconds / 60)} 分钟没新进展`; +} diff --git a/src/server.ts b/src/server.ts index 66792cb..7de77a9 100644 --- a/src/server.ts +++ b/src/server.ts @@ -70,6 +70,7 @@ export async function createGatewayRuntime(config: AppConfig): Promise closing ||= (async () => { qqGatewayClient?.stop(); + gateway.shutdown(); const results = await Promise.allSettled([sessionManager.shutdown(), store.close()]); const failed = results.find((result): result is PromiseRejectedResult => result.status === "rejected"); if (failed) throw failed.reason; diff --git a/test/gateway.test.ts b/test/gateway.test.ts index 676672b..a80d2fb 100644 --- a/test/gateway.test.ts +++ b/test/gateway.test.ts @@ -1,162 +1,415 @@ import assert from "node:assert/strict"; import test from "node:test"; -import type { ConversationRuntime } from "../src/acp/types.js"; +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 { Gateway, type GatewayTimer } from "../src/core/gateway.js"; import type { IncomingMessage, OutgoingMessage } from "../src/core/types.js"; +interface PendingTurn { + request: ConversationRequest; + resolve(text: string): void; + reject(error: Error): void; +} + class FakeRuntime implements ConversationRuntime { - prompts = 0; cancelled = 0; resets = 0; promptError?: Error; - private readonly releases: Array<() => void> = []; - async prompt() { - this.prompts++; - if (this.promptError) throw this.promptError; - await new Promise((resolve) => { this.releases.push(resolve); }); - return { text: "done", botId: "test-bot", agentId: "kimi" }; + readonly turns: PendingTurn[] = []; + cancelled = 0; + resets = 0; + idleSeconds = 0; + + 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() { this.cancelled++; this.releaseNext(); return true; } + async cancel() { this.cancelled++; this.turns.at(-1)?.resolve("cancelled"); return true; } async reset() { this.resets++; } - status() { return { bot: "test-bot", agent: "kimi", workspace: "/tmp", running: this.releases.length > 0, phase: "processing" }; } + status() { + return { + bot: "test-bot", agent: "kimi", workspace: "/tmp", + running: this.turns.length > 0, phase: "processing", idleSeconds: this.idleSeconds + }; + } stats() { return { activeWorkers: 0, inFlight: 0, crashes: 0, persistedBindings: 0 }; } async shutdown() {} - releaseNext() { this.releases.shift()?.(); } +} + +interface ScheduledTimer extends GatewayTimer { + deadline: number; + callback: () => void; + cancelled: boolean; + order: number; +} + +class FakeClock { + nowMs = 0; + private order = 0; + readonly timers: ScheduledTimer[] = []; + + readonly now = () => this.nowMs; + readonly schedule = (callback: () => void, delayMs: number): GatewayTimer => { + const timer: ScheduledTimer = { + deadline: this.nowMs + delayMs, + callback, + cancelled: false, + order: this.order++, + cancel() { timer.cancelled = true; } + }; + this.timers.push(timer); + return timer; + }; + + set(timeMs: number) { this.nowMs = timeMs; } + + async advanceTo(timeMs: number) { + while (true) { + const next = this.timers + .filter((timer) => !timer.cancelled && timer.deadline <= timeMs) + .sort((left, right) => left.deadline - right.deadline || left.order - right.order)[0]; + if (!next) break; + next.cancelled = true; + this.nowMs = next.deadline; + next.callback(); + await flush(); + } + this.nowMs = timeMs; + } + + async jumpTo(timeMs: number) { + this.nowMs = timeMs; + const due = this.timers + .filter((timer) => !timer.cancelled && timer.deadline <= timeMs) + .sort((left, right) => left.deadline - right.deadline || left.order - right.order); + for (const timer of due) { + if (timer.cancelled) continue; + timer.cancelled = true; + timer.callback(); + await flush(); + } + } + + activeTimers() { return this.timers.filter((timer) => !timer.cancelled).length; } } const policy = { allowedUsers: [] as string[], allowedChats: [] as string[], requireMentionInGroup: false }; -const message = (text: string): IncomingMessage => ({ platform: "qq", chatId: "chat", userId: "user", text }); +const message = (text: string, messageId = text, chatId = "chat"): IncomingMessage => ({ + platform: "qq", chatId, userId: "user", text, messageId +}); +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 new Promise((resolve) => setTimeout(resolve, 5)); + for (let count = 0; count < 50 && !condition(); count++) await flush(); assert.equal(condition(), true); }; -function recordingAdapter(sent: string[] = []): PlatformAdapter { - return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing.text); } }; +function recordingAdapter(sent: OutgoingMessage[] = []): PlatformAdapter { + return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } }; } +function createGateway(runtime = new FakeRuntime(), clock = new FakeClock()) { + return { runtime, clock, gateway: new Gateway(policy, runtime, { now: clock.now, schedule: clock.schedule }) }; +} + +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("completion and timer callbacks are safe in both orders at 15 seconds", async () => { + { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "completion-first"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + clock.set(15_000); + runtime.turns[0].resolve("done"); + await turn; + await clock.advanceTo(15_000); + assert.deepEqual(sent.map((outgoing) => outgoing.text), ["done"]); + } + { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "timer-first"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(15_000); + runtime.turns[0].resolve("done"); + await turn; + assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [ + ["这条还在处理,完成后我会直接回复。", 1], + ["done", 2] + ]); + } +}); + +test("running and queued notices are independent and queued work announces start", async () => { + const { runtime, clock, 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 === 1); + await clock.advanceTo(15_000); + assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), ["这条还在处理,完成后我会直接回复。"]); + assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.text), ["前面还有任务,这条还在排队,轮到后我马上处理。"]); + + runtime.turns[0].resolve("first done"); + await waitFor(() => runtime.turns.length === 2); + assert.deepEqual(messagesFor(sent, "second").map((outgoing) => [outgoing.text, outgoing.replySequence]), [ + ["前面还有任务,这条还在排队,轮到后我马上处理。", 1], + ["轮到这条了,我开始处理。", 2] + ]); + runtime.turns[1].resolve("second done"); + await Promise.all([first, second]); + assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.replySequence), [1, 2]); + assert.deepEqual(messagesFor(sent, "second").map((outgoing) => outgoing.replySequence), [1, 2, 3]); +}); + +test("running reminders fire at 60, 180, and 480 seconds with dynamic activity text", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "reminders"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(15_000); + runtime.idleSeconds = 10; + await clock.advanceTo(60_000); + runtime.idleSeconds = 60; + await clock.advanceTo(180_000); + runtime.idleSeconds = 180; + await clock.advanceTo(480_000); + assert.deepEqual(sent.map((outgoing) => outgoing.text), [ + "这条还在处理,完成后我会直接回复。", + "这条还在处理中,刚刚还有进展;完成后我会直接回复。", + "这条还在处理中,最近 60 秒没新进展;完成后我会直接回复。", + "这条还在处理中,最近 3 分钟没新进展;完成后我会直接回复。" + ]); + runtime.turns[0].resolve("done"); + await turn; +}); + +test("queued reminders use queue wording and sample the active ACP status", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const adapter = recordingAdapter(sent); + const first = gateway.receive(message("one", "active"), adapter); + const second = gateway.receive(message("two", "queued"), adapter); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(15_000); + runtime.idleSeconds = 45; + await clock.advanceTo(60_000); + assert.equal(messagesFor(sent, "queued").at(-1)?.text, + "前面的任务还在处理,最近 45 秒没新进展;这条仍在排队,轮到后我马上处理。"); + runtime.turns[0].resolve("one done"); + await waitFor(() => runtime.turns.length === 2); + runtime.turns[1].resolve("two done"); + await Promise.all([first, second]); +}); + +test("a delayed timer jump emits one reminder instead of backfilling every deadline", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "jump"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(15_000); + await clock.jumpTo(1_000_000); + assert.equal(sent.length, 2); + assert.equal(clock.activeTimers(), 1); + runtime.turns[0].resolve("done"); + await turn; +}); + +test("completion and shutdown cancel all inbound timers", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const first = gateway.receive(message("one", "complete"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + assert.equal(clock.activeTimers(), 1); + runtime.turns[0].resolve("done"); + await first; + assert.equal(clock.activeTimers(), 0); + + const second = gateway.receive(message("two", "shutdown"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 2); + assert.equal(clock.activeTimers(), 1); + gateway.shutdown(); + assert.equal(clock.activeTimers(), 0); + await clock.advanceTo(600_000); + assert.deepEqual(messagesFor(sent, "shutdown"), []); + runtime.turns[1].resolve("done after shutdown"); + await second; +}); + +test("prompt does not wait for a notice send and final delivery delay releases the chat lock", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + let releaseFirstSend!: () => void; + const firstSend = new Promise((resolve) => { releaseFirstSend = resolve; }); + const adapter: PlatformAdapter = { + name: "test", async handleWebhook() { return {}; }, + async sendMessage(outgoing) { + sent.push(outgoing); + if (outgoing.replyTo === "first" && outgoing.replySequence === 1) await firstSend; + } + }; + const first = gateway.receive(message("one", "first"), adapter); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(15_000); + assert.equal(runtime.turns.length, 1); + runtime.turns[0].resolve("first done"); + const second = gateway.receive(message("two", "second"), adapter); + await waitFor(() => runtime.turns.length === 2); + assert.equal(messagesFor(sent, "first").length, 1); + releaseFirstSend(); + runtime.turns[1].resolve("second done"); + await Promise.all([first, second]); + assert.deepEqual(messagesFor(sent, "first").map((outgoing) => outgoing.text), [ + "这条还在处理,完成后我会直接回复。", "first done" + ]); +}); + +test("delivery failures are logged safely and do not block later sends", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const errors: string[] = []; + let sends = 0; + const adapter: PlatformAdapter = { + name: "test", async handleWebhook() { return {}; }, + async sendMessage(outgoing) { + sent.push(outgoing); + if (++sends === 1) 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); + await clock.advanceTo(15_000); + 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("a failed final send does not turn a successful runtime result into Agent error", 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("provider detail"); } + }; + const originalError = console.error; + console.error = (...args: unknown[]) => { errors.push(args.map(String).join(" ")); }; + try { + const turn = gateway.receive(message("work", "final-failure"), adapter); + await waitFor(() => runtime.turns.length === 1); + runtime.turns[0].resolve("real result"); + assert.deepEqual(await turn, { ok: true, reply: "real result" }); + } finally { + console.error = originalError; + } + assert.deepEqual(sent.map((outgoing) => outgoing.text), ["real result"]); + 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("status, cancel, and help bypass the chat lock while new remains serial", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const adapter = recordingAdapter(); + const { runtime, clock, gateway } = createGateway(); + const adapter = recordingAdapter(); const turn = gateway.receive(message("work"), adapter, { synchronous: true }); - await waitFor(() => runtime.prompts === 1); + await waitFor(() => runtime.turns.length === 1); + clock.set(5_000); const status = await gateway.receive(message("/status"), adapter, { synchronous: true }); assert.match(status.reply || "", /gatewayRunning=true/); assert.match(status.reply || "", /queued=0/); + assert.match(status.reply || "", /gatewayRunningSeconds=5/); assert.match((await gateway.receive(message("/help"), adapter, { synchronous: true })).reply || "", /\/status/); - const cancelled = await gateway.receive(message("/cancel"), adapter, { synchronous: true }); - assert.equal(cancelled.reply, "Cancellation requested."); + assert.equal((await gateway.receive(message("/cancel"), adapter, { synchronous: true })).reply, "Cancellation requested."); await turn; const second = gateway.receive(message("work"), adapter, { synchronous: true }); - await waitFor(() => runtime.prompts === 2); + await waitFor(() => runtime.turns.length === 2); const reset = gateway.receive(message("/new"), adapter, { synchronous: true }); - await new Promise((resolve) => setTimeout(resolve, 10)); + await flush(); assert.equal(runtime.resets, 0); - runtime.releaseNext(); + runtime.turns[1].resolve("done"); await second; await reset; assert.equal(runtime.resets, 1); }); -test("asynchronous messages acknowledge processing and queue state", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const sent: string[] = []; const adapter = recordingAdapter(sent); - const first = gateway.receive(message("one"), adapter); - await waitFor(() => runtime.prompts === 1); - const second = gateway.receive(message("two"), adapter); - await waitFor(() => sent.length >= 2); - assert.deepEqual(sent.slice(0, 2), [ - "已收到,正在处理。可随时发送 /status 查看状态。", - "已收到,已排队。可随时发送 /status 查看状态。" - ]); - const status = await gateway.receive(message("/status"), adapter, { synchronous: true }); - assert.match(status.reply || "", /queued=1/); - assert.match(status.reply || "", /gatewayRunningSeconds=\d+/); - runtime.releaseNext(); - await waitFor(() => runtime.prompts === 2); - runtime.releaseNext(); - await Promise.all([first, second]); -}); - -test("prompt waits for acknowledgement and final reply uses sequence two", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const sent: OutgoingMessage[] = []; - let releaseAcknowledgement!: () => void; - const acknowledgement = new Promise((resolve) => { releaseAcknowledgement = resolve; }); - const adapter: PlatformAdapter = { - name: "test", async handleWebhook() { return {}; }, - async sendMessage(outgoing) { - sent.push(outgoing); - if (outgoing.replySequence === 1) await acknowledgement; - } - }; - const turn = gateway.receive({ ...message("work"), messageId: "message-1" }, adapter); - await waitFor(() => sent.length === 1); - assert.equal(runtime.prompts, 0); - assert.equal(sent[0].replySequence, 1); - releaseAcknowledgement(); - await waitFor(() => runtime.prompts === 1); - runtime.releaseNext(); - await turn; - assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replyTo, outgoing.replySequence]), [ - ["已收到,正在处理。可随时发送 /status 查看状态。", "message-1", 1], - ["done", "message-1", 2] - ]); -}); - -test("acknowledgement failure is logged safely and does not block work", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const errors: string[] = []; - let sends = 0; - const adapter: PlatformAdapter = { - name: "test", async handleWebhook() { return {}; }, - async sendMessage() { if (++sends === 1) 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"), adapter); - await waitFor(() => runtime.prompts === 1); - runtime.releaseNext(); - assert.equal((await turn).reply, "done"); - } finally { - console.error = originalError; - } - assert.deepEqual(errors, ["Gateway acknowledgement send failed (platform=qq)"]); -}); - -test("asynchronous prompt errors reply with sequence two", async () => { - const runtime = new FakeRuntime(); runtime.promptError = new Error("failed"); - const gateway = new Gateway(policy, runtime); const sent: OutgoingMessage[] = []; - const adapter: PlatformAdapter = { - name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } - }; - const result = await gateway.receive(message("work"), adapter); - assert.equal(result.ok, false); - assert.deepEqual(sent.map((outgoing) => outgoing.replySequence), [1, 2]); -}); - -test("single asynchronous command replies with sequence one", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const sent: OutgoingMessage[] = []; - const adapter: PlatformAdapter = { - name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } - }; - await gateway.receive(message("/help"), adapter); - assert.equal(sent.length, 1); - assert.equal(sent[0].replySequence, 1); -}); - -test("synchronous messages do not send acknowledgements", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const sent: string[] = []; const adapter = recordingAdapter(sent); - const turn = gateway.receive(message("work"), adapter, { synchronous: true }); - await waitFor(() => runtime.prompts === 1); +test("synchronous work sends no notices or delivery messages", async () => { + const { runtime, clock, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work"), recordingAdapter(sent), { synchronous: true }); + await waitFor(() => runtime.turns.length === 1); + await clock.advanceTo(600_000); assert.deepEqual(sent, []); - runtime.releaseNext(); - await turn; + runtime.turns[0].resolve("done"); + assert.deepEqual(await turn, { ok: true, reply: "done" }); assert.deepEqual(sent, []); }); -test("fixed Bot status and retired role commands never reach ACP", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); const adapter = recordingAdapter(); - for (const command of ["/roles", "/role ops", "/agents", "/agent ops"]) { - assert.match((await gateway.receive(message(command), adapter, { synchronous: true })).reply || "", /removed in Config v3/); - } - assert.match((await gateway.receive(message("/status"), adapter, { synchronous: true })).reply || "", /bot=test-bot/); - assert.equal(runtime.prompts, 0); +test("asynchronous help replies with sequence one", 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 new waits for prior work and then replies with sequence one", async () => { + const { runtime, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const adapter = recordingAdapter(sent); + const work = gateway.receive(message("work", "work"), adapter); + await waitFor(() => runtime.turns.length === 1); + + const reset = gateway.receive(message("/new", "new"), adapter); + await flush(); + assert.equal(runtime.resets, 0); + assert.deepEqual(messagesFor(sent, "new"), []); + + runtime.turns[0].resolve("done"); + await work; + assert.deepEqual(await reset, { ok: true, reply: "Started a new native ACP session for this chat." }); + assert.equal(runtime.resets, 1); + assert.deepEqual(messagesFor(sent, "new").map((outgoing) => [outgoing.text, outgoing.replySequence]), [ + ["Started a new native ACP session for this chat.", 1] + ]); +}); + +test("asynchronous retired-role command replies with sequence one without reaching ACP", async () => { + const { runtime, gateway } = createGateway(); + const sent: OutgoingMessage[] = []; + const result = await gateway.receive(message("/roles", "retired"), recordingAdapter(sent)); + assert.match(result.reply || "", /removed in Config v3/); + assert.equal(runtime.turns.length, 0); + assert.deepEqual(messagesFor(sent, "retired").map((outgoing) => outgoing.replySequence), [1]); }); diff --git a/test/qq-adapter.test.ts b/test/qq-adapter.test.ts index 4c3f795..322a1bc 100644 --- a/test/qq-adapter.test.ts +++ b/test/qq-adapter.test.ts @@ -8,7 +8,7 @@ const config: QqConfig = { botNames: ["Bot"], intents: 33_554_432, shard: [0, 1] }; -test("normalizes GROUP and C2C author openid while ACK remains immediate", async () => { +test("normalizes GROUP and C2C author openid while webhook ACK remains immediate", async () => { const received: any[] = []; const gateway = { receive: async (message: unknown) => { received.push(message); throw new Error("ACP failed"); } }; const adapter = new QqAdapter(config, gateway as never); @@ -31,12 +31,12 @@ test("sendMessage uses C2C endpoint and forwards reply sequences as msg_seq", as try { const adapter = new QqAdapter(config, { receive: async () => ({ ok: true }) } as never); const target = { platform: "qq", chatId: "user:u2", raw: { author: { user_openid: "u2" } } }; - await adapter.sendMessage({ target, text: "ack", replyTo: "m2", replySequence: 1 }); - await adapter.sendMessage({ target, text: "reply", replyTo: "m2", replySequence: 2 }); + await adapter.sendMessage({ target, text: "first", replyTo: "m2", replySequence: 1 }); + await adapter.sendMessage({ target, text: "later", replyTo: "m2", replySequence: 7 }); assert.ok(urls.some((url) => url.endsWith("/v2/users/u2/messages"))); assert.deepEqual(bodies, [ - { content: "ack", msg_id: "m2", msg_seq: 1 }, - { content: "reply", msg_id: "m2", msg_seq: 2 } + { content: "first", msg_id: "m2", msg_seq: 1 }, + { content: "later", msg_id: "m2", msg_seq: 7 } ]); } finally { globalThis.fetch = original; } });