diff --git a/AGENTS.md b/AGENTS.md index d8635bd..d27b51c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -76,16 +76,21 @@ gori-agent list - 加载当前 Bot skills、计算 fingerprint、生成版本化 bootstrap。 - `src/acp/client.ts` - ACP initialize/new/resume/load/prompt/cancel 和 permission request。 + - 当前 session 的每个 `session/update` 只触发活动回调,不记录或泄露 update 内容。 - `src/acp/worker.ts` - 直接使用 resolved Bot 的 agent,在 `bot.workspace` 启动。 + - 记录 turn 开始/最近活动时间,并区分 `idle`、bootstrap `initializing`、普通 prompt `processing`。 - `src/acp/session-manager.ts` - 固定 Bot 的 binding、worker pool、恢复、idle sweep、capacity、cancel/reset。 + - status 暴露 `phase`、`runningSeconds`、`idleSeconds`。 - `src/core/durable-session-store.ts` - state v2,保存 bot/platform identity 和 chat binding。 - 单 writer lock、串行持久化、临时文件、fsync、原子 rename。 - `src/core/gateway.ts` - - 入站 allowlist、群 mention、命令、per-chat lock 和回复。 - - `/cancel` 必须绕过 chat lock。 + - 入站 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`。 - `src/core/command-router.ts` - `/help`、`/status`、`/cancel`、`/new`;旧 role 命令返回 retired 提示。 - `src/platforms/*` @@ -199,7 +204,7 @@ Binding 保存 agent ID、workspace、native session ID、Bot fingerprint 和时 首次消息创建 native session,发送隐藏 bootstrap,成功后保存 binding。重启/idle 后优先 resume,再 fallback load。prompt 失败不重放。 -`/new` 删除当前 chat binding;下一条消息创建新 session。`/cancel` 不等待 chat lock。 +`/new` 在 per-chat lock 内删除当前 chat binding;下一条消息创建新 session。`/status`、`/cancel`、`/help` 不等待 chat lock;status 合并 Gateway running/queued/started 时长与 ACP phase/running/idle 时长。 ## 6. 单平台装配 diff --git a/README.md b/README.md index d7dc325..1eb110d 100644 --- a/README.md +++ b/README.md @@ -206,12 +206,13 @@ Bot fingerprint 包含 Bot ID、workspace、persona、agent、permissions、skil /new ``` -- `/status` 显示固定 Bot、agent、workspace 和 session persisted/running 状态。 +- `/status` 绕过 per-chat lock,显示固定 Bot、agent、workspace、队列数、Gateway 当前任务时长,以及 ACP `phase`、`runningSeconds`、`idleSeconds`。`phase` 为 `idle`、bootstrap 的 `initializing` 或普通 prompt 的 `processing`;`idleSeconds` 按当前 session 的最近一次 `session/update` 活动计算,不记录或输出 update 内容。 - `/cancel` 绕过 per-chat lock,发送 ACP cancel。 -- `/new` 清除当前 chat binding;下一条普通消息创建新 native session。 +- `/help` 绕过 per-chat lock,可在长任务期间立即返回。 +- `/new` 保持串行,清除当前 chat binding;下一条普通消息创建新 native session。 - `/roles`、`/role`、`/agents`、`/agent` 会返回固定 retired 提示,不会转发给 ACP。 -同一 chat 串行执行,不同 chat 可并发。 +同一 chat 的普通消息和 `/new` 串行执行,不同 chat 可并发。异步平台的普通消息通过 policy 后会立即确认:空闲时发送“已收到,正在处理。可随时发送 /status 查看状态。”,已有同 chat 工作时发送“已收到,已排队。可随时发送 /status 查看状态。”。每条消息的 ACP prompt 会等待该确认完成或安全失败后再开始;排队消息的确认不等待前项。确认发送失败只写不含消息正文或 provider 错误详情的安全日志,不阻止已入队任务;同步 webhook 不额外发送确认。QQ 对同一入站 `msg_id` 的确认使用 `msg_seq=1`,最终或错误回复使用 `msg_seq=2`;单条命令回复使用 `msg_seq=1`。 ## HTTP 端点 diff --git a/src/acp/client.ts b/src/acp/client.ts index ad42bd2..1c2f42b 100644 --- a/src/acp/client.ts +++ b/src/acp/client.ts @@ -7,6 +7,7 @@ import type { PermissionPolicy } from "../config.js"; export interface AcpClientOptions { initializeTimeoutMs: number; policy: PermissionPolicy; + onSessionActivity?: () => void; } export class AcpClient { @@ -47,19 +48,24 @@ export class AcpClient { async resumeSession(sessionId: string, cwd: string): Promise { this.collecting = false; this.chunks = []; - if (this.capabilities.sessionCapabilities?.resume) { - try { - await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] }); - } catch (error) { - if (!this.capabilities.loadSession) throw error; - await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] }); - } - } else if (this.capabilities.loadSession) { - await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] }); - } else { - throw new Error("ACP backend cannot resume or load sessions"); - } this.activeSessionId = sessionId; + try { + if (this.capabilities.sessionCapabilities?.resume) { + try { + await this.connection.agent.request(acp.methods.agent.session.resume, { sessionId, cwd, mcpServers: [] }); + } catch (error) { + if (!this.capabilities.loadSession) throw error; + await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] }); + } + } else if (this.capabilities.loadSession) { + await this.connection.agent.request(acp.methods.agent.session.load, { sessionId, cwd, mcpServers: [] }); + } else { + throw new Error("ACP backend cannot resume or load sessions"); + } + } catch (error) { + this.activeSessionId = undefined; + throw error; + } this.chunks = []; } @@ -91,7 +97,9 @@ export class AcpClient { close(error?: unknown): void { this.connection.close(error); } private handleUpdate(notification: SessionNotification): void { - if (!this.collecting || notification.sessionId !== this.activeSessionId) return; + if (notification.sessionId !== this.activeSessionId) return; + this.options.onSessionActivity?.(); + if (!this.collecting) return; const update = notification.update; if (update.sessionUpdate === "agent_message_chunk" && update.content.type === "text") this.chunks.push(update.content.text); } diff --git a/src/acp/session-manager.ts b/src/acp/session-manager.ts index 16ac1d1..214b574 100644 --- a/src/acp/session-manager.ts +++ b/src/acp/session-manager.ts @@ -41,7 +41,7 @@ export class AcpSessionManager implements ConversationRuntime { try { if (!binding) { const now = Date.now(); - await worker.prompt(this.bot.bootstrap); + await worker.prompt(this.bot.bootstrap, "initializing"); binding = { chatKey, agentId: this.bot.agent.id, nativeSessionId: worker.nativeSessionId!, workspace: this.bot.workspace, botFingerprint: this.bot.fingerprint, createdAt: now, updatedAt: now @@ -75,12 +75,17 @@ export class AcpSessionManager implements ConversationRuntime { status(platform: string, chatId: string): Record { const chatKey = chatKeyFor(platform, chatId); const binding = this.store.getBinding(chatKey); + const worker = this.inFlight.get(chatKey); + const now = Date.now(); return { bot: this.bot.id, agent: this.bot.agent.id, workspace: this.bot.workspace, persisted: Boolean(binding), - running: this.inFlight.has(chatKey) + running: Boolean(worker), + phase: worker?.phase || "idle", + runningSeconds: worker?.turnStartedAt ? Math.max(0, Math.floor((now - worker.turnStartedAt) / 1000)) : 0, + idleSeconds: worker?.lastActivityAt ? Math.max(0, Math.floor((now - worker.lastActivityAt) / 1000)) : 0 }; } diff --git a/src/acp/worker.ts b/src/acp/worker.ts index ca4e003..7018063 100644 --- a/src/acp/worker.ts +++ b/src/acp/worker.ts @@ -5,6 +5,8 @@ import { AcpClient } from "./client.js"; export class AcpSessionRestoreError extends Error {} +export type AcpWorkerPhase = "idle" | "initializing" | "processing"; + export class AcpWorker { private child?: ChildProcessWithoutNullStreams; private client?: AcpClient; @@ -14,6 +16,9 @@ export class AcpWorker { private stderrBytes = 0; nativeSessionId?: string; lastUsedAt = Date.now(); + turnStartedAt?: number; + lastActivityAt?: number; + phase: AcpWorkerPhase = "idle"; inFlight = false; constructor( @@ -35,7 +40,11 @@ export class AcpWorker { this.exited = true; if (!this.stopping) this.crashed(new Error(`ACP worker exited (code=${code ?? "null"}, signal=${signal ?? "null"})`)); }); - this.client = new AcpClient(this.child, { initializeTimeoutMs: this.config.initializeTimeoutMs, policy: this.bot.permissions }); + this.client = new AcpClient(this.child, { + initializeTimeoutMs: this.config.initializeTimeoutMs, + policy: this.bot.permissions, + onSessionActivity: () => { this.lastActivityAt = Date.now(); } + }); try { await this.client.initialize(); if (nativeSessionId) { @@ -52,11 +61,15 @@ export class AcpWorker { } } - async prompt(text: string): Promise { + async prompt(text: string, phase: Exclude = "processing"): Promise { if (!this.client || !this.nativeSessionId || this.exited) throw new Error("ACP worker is not available"); if (this.inFlight) throw new Error("ACP worker already has an in-flight turn"); + const now = Date.now(); this.inFlight = true; - this.lastUsedAt = Date.now(); + this.phase = phase; + this.turnStartedAt = now; + this.lastActivityAt = now; + this.lastUsedAt = now; this.abort = new AbortController(); let timeout: NodeJS.Timeout | undefined; const timeoutPromise = new Promise((_resolve, reject) => { @@ -72,6 +85,9 @@ export class AcpWorker { } finally { if (timeout) clearTimeout(timeout); this.inFlight = false; + this.phase = "idle"; + this.turnStartedAt = undefined; + this.lastActivityAt = undefined; this.abort = undefined; this.lastUsedAt = Date.now(); } diff --git a/src/core/gateway.ts b/src/core/gateway.ts index 59be792..e37e985 100644 --- a/src/core/gateway.ts +++ b/src/core/gateway.ts @@ -7,8 +7,15 @@ import type { IncomingMessage } from "./types.js"; export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string } +interface ChatWorkState { + running: boolean; + queued: number; + startedAt?: number; +} + export class Gateway { private readonly locks = new Map>(); + private readonly chatWork = new Map(); readonly commandRouter = new CommandRouter(); constructor(private readonly policy: GatewayPolicy, private readonly runtime: ConversationRuntime) {} @@ -21,19 +28,28 @@ export class Gateway { } const command = this.commandRouter.parse(message.text); - if (command?.kind === "cancel") return this.reply(message, adapter, await this.cancelText(message), options); - if (command?.kind === "new") await this.runtime.cancel(message.platform, message.chatId); + if (command?.kind === "help" || command?.kind === "status" || command?.kind === "cancel") { + return this.reply(message, adapter, await this.executeCommand(command, message), 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; 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 this.reply(message, adapter, reply, options); + return this.reply(message, adapter, reply, options, replySequence); } 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); + if (!options.synchronous) await this.send(message, adapter, reply, replySequence); return { ok: false, error: errorText, reply }; } }); @@ -45,7 +61,14 @@ export class Gateway { switch (command.kind) { case "help": return ["Commands:", "/status", "/cancel", "/new", "/help"].join("\n"); case "status": { - const status = this.runtime.status(message.platform, message.chatId); + const chatKey = chatKeyFor(message.platform, message.chatId); + const work = this.chatWork.get(chatKey); + const status = { + ...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 + }; return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`; } case "new": @@ -60,13 +83,26 @@ 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 }): Promise { - if (!options.synchronous) await this.send(message, adapter, reply); + 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 send(message: IncomingMessage, adapter: PlatformAdapter, text: string): Promise { - return adapter.sendMessage({ target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, text, replyTo: message.messageId }); + 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 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 checkPolicy(message: IncomingMessage): string | undefined { @@ -78,14 +114,25 @@ export class Gateway { private async withChatLock(key: string, fn: () => Promise): Promise { const previous = this.locks.get(key) || Promise.resolve(); + const state = this.chatWork.get(key) || { running: false, queued: 0 }; + state.queued++; + this.chatWork.set(key, state); let release!: () => void; const current = new Promise((resolve) => { release = resolve; }); - const queued = previous.then(() => current); - this.locks.set(key, queued); + const lock = previous.then(() => current); + this.locks.set(key, lock); await previous; + state.queued--; + state.running = true; + state.startedAt = Date.now(); try { return await fn(); } finally { + state.running = false; + state.startedAt = undefined; release(); - if (this.locks.get(key) === queued) this.locks.delete(key); + if (this.locks.get(key) === lock) { + this.locks.delete(key); + this.chatWork.delete(key); + } } } } diff --git a/src/core/types.ts b/src/core/types.ts index b8304a4..1b31416 100644 --- a/src/core/types.ts +++ b/src/core/types.ts @@ -25,6 +25,7 @@ export interface OutgoingMessage { target: MessageTarget; text: string; replyTo?: string; + replySequence?: number; } export interface WebhookRequestContext { diff --git a/src/platforms/qq/adapter.ts b/src/platforms/qq/adapter.ts index 801e28b..a9cdf76 100644 --- a/src/platforms/qq/adapter.ts +++ b/src/platforms/qq/adapter.ts @@ -72,7 +72,7 @@ export class QqAdapter implements PlatformAdapter { Authorization: `QQBot ${token}`, "Content-Type": "application/json; charset=utf-8" }, - body: JSON.stringify({ content: message.text, msg_id: message.replyTo }) + body: JSON.stringify({ content: message.text, msg_id: message.replyTo, msg_seq: message.replySequence ?? 1 }) }); if (!response.ok) throw new Error(`QQ send failed: HTTP ${response.status}`); const data = await response.json() as QqSendMessageResponse; diff --git a/test/acp-client.test.ts b/test/acp-client.test.ts index df75788..5b0f5ac 100644 --- a/test/acp-client.test.ts +++ b/test/acp-client.test.ts @@ -1,7 +1,9 @@ import assert from "node:assert/strict"; +import { spawn } from "node:child_process"; +import path from "node:path"; import test from "node:test"; import type { RequestPermissionRequest } from "@agentclientprotocol/sdk"; -import { decidePermission } from "../src/acp/client.js"; +import { AcpClient, decidePermission } from "../src/acp/client.js"; const request = (rawInput: unknown, overrides: Partial = {}): RequestPermissionRequest => ({ sessionId: "session", @@ -30,3 +32,22 @@ test("permission deny and allowlist fail closed", () => { assert.deepEqual(decidePermission(request("git status", { name: undefined, title: "bash git status" }), unanchored).outcome, { outcome: "selected", optionId: "reject" }); assert.deepEqual(decidePermission(request({ command: "git status", env: { PATH: "/tmp" } }), unanchored).outcome, { outcome: "selected", optionId: "reject" }); }); + +test("reports activity for every update from the active session", async () => { + const child = spawn(process.execPath, [path.resolve("test/fixtures/fake-acp-agent.mjs")], { stdio: ["pipe", "pipe", "pipe"] }); + let activities = 0; + const client = new AcpClient(child, { + initializeTimeoutMs: 1_000, + policy: { mode: "deny", allowedTools: [], allowedCommandPatterns: [] }, + onSessionActivity: () => { activities++; } + }); + try { + await client.initialize(); + await client.newSession(path.resolve(".")); + assert.equal(await client.prompt("activity"), "firstsecond"); + assert.equal(activities, 2); + } finally { + client.close(); + child.kill("SIGTERM"); + } +}); diff --git a/test/acp-session-manager.test.ts b/test/acp-session-manager.test.ts index 4db54f9..8773264 100644 --- a/test/acp-session-manager.test.ts +++ b/test/acp-session-manager.test.ts @@ -95,3 +95,28 @@ test("cancel reaches hanging prompt, reset unbinds, and fingerprint changes sess assert.notEqual(second.text.split(":")[1], firstSession); await changed.manager.shutdown(); await changed.store.close(); }); + +test("status reports initializing, processing, running time, and update activity", async () => { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-acp-status-")); const state = path.join(dir, "state.json"); const log = path.join(dir, "fake.log"); + const current = await runtime(state, log, "test", { FAKE_ACP_BOOTSTRAP_DELAY_MS: "200" }); + const first = current.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "ready" }); + for (let count = 0; count < 50 && current.manager.status("qq", "chat").phase !== "initializing"; count++) { + await new Promise((resolve) => setTimeout(resolve, 10)); + } + assert.deepEqual(current.manager.status("qq", "chat"), { + bot: "test-bot", agent: "fake", workspace: path.resolve("."), persisted: false, running: true, + phase: "initializing", runningSeconds: 0, idleSeconds: 0 + }); + await first; + + const activity = current.manager.prompt({ platform: "qq", chatId: "chat", userId: "user", text: "activity" }); + await new Promise((resolve) => setTimeout(resolve, 1_100)); + const processing = current.manager.status("qq", "chat"); + assert.equal(processing.phase, "processing"); + assert.equal(processing.running, true); + assert.equal(processing.runningSeconds, 1); + assert.equal(processing.idleSeconds, 0); + await activity; + assert.equal(current.manager.status("qq", "chat").phase, "idle"); + await current.manager.shutdown(); await current.store.close(); +}); diff --git a/test/fixtures/fake-acp-agent.mjs b/test/fixtures/fake-acp-agent.mjs index 2783afe..d15fef1 100644 --- a/test/fixtures/fake-acp-agent.mjs +++ b/test/fixtures/fake-acp-agent.mjs @@ -58,6 +58,16 @@ const app = acp.agent({ name: "fake-acp-agent" }) pending.delete(params.sessionId); return { stopReason: "cancelled" }; } + if (text === "activity") { + await update(client, params.sessionId, "first"); + await new Promise((resolve) => setTimeout(resolve, 600)); + await update(client, params.sessionId, "second"); + await new Promise((resolve) => setTimeout(resolve, 700)); + return { stopReason: "end_turn" }; + } + if (text.includes("Initialize this ACP session") && process.env.FAKE_ACP_BOOTSTRAP_DELAY_MS) { + await new Promise((resolve) => setTimeout(resolve, Number(process.env.FAKE_ACP_BOOTSTRAP_DELAY_MS))); + } await update(client, params.sessionId, text.includes("Initialize this ACP session") ? "READY" : `reply:${params.sessionId}:${text}`); return { stopReason: "end_turn" }; }) diff --git a/test/gateway.test.ts b/test/gateway.test.ts index 2b27210..676672b 100644 --- a/test/gateway.test.ts +++ b/test/gateway.test.ts @@ -3,36 +3,157 @@ import test from "node:test"; import type { ConversationRuntime } from "../src/acp/types.js"; import type { PlatformAdapter } from "../src/core/adapter.js"; import { Gateway } from "../src/core/gateway.js"; -import type { IncomingMessage } from "../src/core/types.js"; +import type { IncomingMessage, OutgoingMessage } from "../src/core/types.js"; class FakeRuntime implements ConversationRuntime { - prompts = 0; cancelled = 0; resets = 0; release?: () => void; - async prompt() { this.prompts++; await new Promise((resolve) => { this.release = resolve; }); return { text: "done", botId: "test-bot", agentId: "kimi" }; } - async cancel() { this.cancelled++; this.release?.(); return true; } + 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" }; + } + async cancel() { this.cancelled++; this.releaseNext(); return true; } async reset() { this.resets++; } - status() { return { bot: "test-bot", agent: "kimi", workspace: "/tmp", running: Boolean(this.release) }; } + status() { return { bot: "test-bot", agent: "kimi", workspace: "/tmp", running: this.releases.length > 0, phase: "processing" }; } stats() { return { activeWorkers: 0, inFlight: 0, crashes: 0, persistedBindings: 0 }; } async shutdown() {} + releaseNext() { this.releases.shift()?.(); } } const policy = { allowedUsers: [] as string[], allowedChats: [] as string[], requireMentionInGroup: false }; -const adapter: PlatformAdapter = { name: "test", async handleWebhook() { return {}; }, async sendMessage() {} }; const message = (text: string): IncomingMessage => ({ platform: "qq", chatId: "chat", userId: "user", text }); +const waitFor = async (condition: () => boolean) => { + for (let count = 0; count < 50 && !condition(); count++) await new Promise((resolve) => setTimeout(resolve, 5)); + assert.equal(condition(), true); +}; -test("cancel bypasses the chat lock and new resets", async () => { - const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); +function recordingAdapter(sent: string[] = []): PlatformAdapter { + return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing.text); } }; +} + +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 turn = gateway.receive(message("work"), adapter, { synchronous: true }); - await new Promise((resolve) => setTimeout(resolve, 10)); + await waitFor(() => runtime.prompts === 1); + + const status = await gateway.receive(message("/status"), adapter, { synchronous: true }); + assert.match(status.reply || "", /gatewayRunning=true/); + assert.match(status.reply || "", /queued=0/); + 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."); await turn; - const reset = await gateway.receive(message("/new"), adapter, { synchronous: true }); - assert.match(reset.reply || "", /new native ACP session/); + + const second = gateway.receive(message("work"), adapter, { synchronous: true }); + await waitFor(() => runtime.prompts === 2); + const reset = gateway.receive(message("/new"), adapter, { synchronous: true }); + await new Promise((resolve) => setTimeout(resolve, 10)); + assert.equal(runtime.resets, 0); + runtime.releaseNext(); + 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); + assert.deepEqual(sent, []); + runtime.releaseNext(); + await turn; + 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 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/); } diff --git a/test/qq-adapter.test.ts b/test/qq-adapter.test.ts index 0980b6b..4c3f795 100644 --- a/test/qq-adapter.test.ts +++ b/test/qq-adapter.test.ts @@ -20,16 +20,23 @@ test("normalizes GROUP and C2C author openid while ACK remains immediate", async assert.equal(received[1].chatId, "user:u2"); assert.equal(received[1].userId, "u2"); }); -test("sendMessage uses nested author.user_openid for C2C endpoint", async () => { - const original = globalThis.fetch; const urls: string[] = []; - globalThis.fetch = (async (input: string | URL | Request) => { +test("sendMessage uses C2C endpoint and forwards reply sequences as msg_seq", async () => { + const original = globalThis.fetch; const urls: string[] = []; const bodies: Array> = []; + globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { urls.push(String(input)); if (String(input).includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + bodies.push(JSON.parse(String(init?.body)) as Record); return new Response(JSON.stringify({ id: "sent" }), { status: 200 }); }) as typeof fetch; try { const adapter = new QqAdapter(config, { receive: async () => ({ ok: true }) } as never); - await adapter.sendMessage({ target: { platform: "qq", chatId: "user:u2", raw: { author: { user_openid: "u2" } } }, text: "reply", replyTo: "m2" }); + 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 }); 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 } + ]); } finally { globalThis.fetch = original; } });