From ffedf616e39fcedf9046e364e56c0f61d2cdfed1 Mon Sep 17 00:00:00 2001 From: zenord Date: Tue, 18 Aug 2026 22:20:55 +0800 Subject: [PATCH] fix: use passive QQ event replies --- AGENTS.md | 6 +- README.md | 4 +- src/acp/assistant-manager.ts | 91 +++++++++++--- src/core/gateway.ts | 151 ++++++++++++++++++++--- src/platforms/qq/adapter.ts | 21 +++- src/roles/role-registry.ts | 3 +- test/assistant-manager.test.ts | 79 ++++++++++++ test/gateway.test.ts | 216 +++++++++++++++++++++++++++++++-- test/qq-adapter.test.ts | 37 ++++++ 9 files changed, 560 insertions(+), 48 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index b4cf7f3..487432c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -86,7 +86,7 @@ gori-agent list - 每个 ACP worker 使用独立进程组和随机 `GORI_AGENT_WORKER_TOKEN`;cancel、timeout、crash 或 assistant 隔离违约必须清理同进程组工具后代;bootstrap 禁止 `setsid`、`nohup`、detached/daemon/background 遗留进程,主动脱组属于无 cgroup/Bubblewrap 方案的边界。 - assistant 只接受 Kimi Code ACP,并依赖项目级 `tools: []`、`subagents: []` profile;permission deny 只是附加层。 - `src/acp/assistant-manager.ts` - - 每 conversation(chat + user)一个无工具 Assistant 会话,同群不同用户互相隔离(`GORI_ASSISTANT_ACTION_V1` envelope:create_proposal/confirm/start_next/cancel/stop,格式只修复一次)。 + - 每 conversation(chat + user)一个无工具 Assistant 会话,同群不同用户互相隔离(`GORI_ASSISTANT_ACTION_V1` envelope:create_proposal/confirm/start_next/cancel/stop,格式只修复一次);runtime 执行 action 后在 reply 末尾追加人话纠正(如 start_next 被阻塞、confirm/stop/cancel 无 owner 匹配),Assistant 不得自行宣称 action 已生效。 - 唯一活跃 Worker 执行已确认 Proposal(`GORI_WORKER_RESULT_V1`:SUCCESS/FAILED/NEEDS_CONFIRMATION);success 与 failure 都需用户确认才收尾,不自动开始下一个;执行中 dirty 是正常的 NEEDS_CONFIRMATION 确认点。 - capacity(maxAssistantSessions/maxProcesses)、idle sweep、cancel/confirm/stop、worker_lost 恢复、owner 事件通知。 - `src/core/durable-session-store.ts` @@ -99,7 +99,7 @@ gori-agent list - `src/core/gateway.ts` - 入站 allowlist、群 mention、命令和回复;不维护普通消息队列或 per-chat task lock,也没有固定时间处理中提醒。 - 每条异步入站使用独立串行 ReplyStream,`replySequence` 从 1 动态递增;发送失败只写安全日志且不阻止任务或后续发送。 - - Worker 落定事件经 `sendEvent` 作为新消息发送(QQ 同样是新消息)。 + - Worker 落定事件经 `sendEvent` 发送:优先引用该 chat 最近入站消息走被动回复窗口(群 4.5 分钟、C2C 55 分钟保守判定,`msg_id` + 递增 `msg_seq`);`msg_seq` 按 chat+messageId 共享计数(正常回复、事件、补发同一计数器,QQ 重复推送同一 msg_id 时沿用已用序号);无新鲜窗口或发送失败时不发主动消息,记录 lastEventDelivery/lastEventError/lastEventAt,补发失败保留队列并在下次入站时重试。 - `src/core/command-router.ts` - `/help`、`/status`、`/list`、`/confirm`、`/stop`、`/cancel`;未识别命令按普通消息处理。 - `src/platforms/*` @@ -212,7 +212,7 @@ Assistant 每条回复以隐藏 `GORI_ASSISTANT_ACTION_V1` envelope 结尾(`re - `SUCCESS` 与 `FAILED` 都进入 `awaiting_user_confirmation`,都需用户确认才落定为 `completed`/`failed`;failed 用 answer `retry` 确认可重试。 - 系统不自动开始下一个 Proposal;只有新确认或 `start_next` 推进队列。 - 执行中遇到 dirty/意外目标是正常的 `NEEDS_CONFIRMATION` 确认点,不是 blocked 终态;用户回答后同一 worker 继续。 -- Gateway 无固定时间提醒;worker 落定后由 Assistant 生成事件文案,作为新消息发给 owner chat。 +- Gateway 无固定时间提醒;worker 落定后由 Assistant 生成事件文案发给 owner chat,QQ 无主动消息权限时优先使用被动回复窗口(群 5 分钟/私聊 60 分钟,按 4.5/55 分钟保守判定),过期或失败则不主动发、记录 lastEventDelivery 并在下次入站时补发;`/status` 输出当前 chat 的事件投递状态与 schedulerState/blockedReason/nextAction。 Assistant 隔离:cwd 在实例私有 `state/assistant-workspaces/`,目录 key 使用进程随机 salt 的 HMAC,不得与项目 workspace 重叠;Kimi 项目级 agent override 必须设置 `tools: []`、`subagents: []`,ACP `mcpServers: []`,未显式配置时 `KIMI_CODE_HOME` 指向实例私有 `state/kimi/assistant/`;任何 tool update 或 permission request 都视为隔离违约,fail closed、终止进程组并删除 binding。每个 ACP worker 独立进程组 + 随机 token;runner 重启时 working Proposal 校验 token 清理旧进程组后标记 failed(worker_lost)。 diff --git a/README.md b/README.md index 2b1a38c..8905d13 100644 --- a/README.md +++ b/README.md @@ -228,7 +228,7 @@ Bot fingerprint 包含 bootstrap schema version、Bot ID、workspace、persona ``` - `/help`:显示自然语言使用说明和兜底命令列表。 -- `/status`:显示固定 Bot、agent、workspace、Assistant 会话数、各状态 Proposal 计数、Worker 是否在执行及当前用户是否 owner。 +- `/status`:显示固定 Bot、agent、workspace、Assistant 会话数、各状态 Proposal 计数、Worker 是否在执行及当前用户是否 owner;另含当前用户维度(`myQueuedProposals`、`myPendingConfirmations`、`schedulerState`、`blockedReason`、`nextAction`,他人 Proposal 只给脱敏原因,不暴露 title/id)与当前 chat 的事件投递状态(`lastEventDelivery`、`lastEventError`、`lastEventAt`)。 - `/list`:列出当前 chat 中当前用户拥有的 Proposal(id、状态、title)。 - `/confirm`:确认当前用户最近待确认的 Proposal(`proposed` 进入队列,`awaiting_user_confirmation` 落定完成/失败或带着回答继续)。 - `/stop`:停止当前用户正在执行的 Worker 并将其 Proposal 标记为 cancelled;对等待确认的 Proposal 按 pending 类型兜底落定(success→completed、failure→failed、step→cancelled);只有 Proposal 发起人可用。 @@ -236,7 +236,7 @@ Bot fingerprint 包含 bootstrap schema version、Bot ID、workspace、persona 日常操作以自然语言为主,Assistant 会自己生成 confirm/cancel/stop 等 action;这些命令是旁路 Assistant 的兜底入口。未识别的 `/` 命令按普通消息交给 Assistant。 -Gateway 不维护普通消息队列或 per-chat task lock,也不再有固定时间(15/60/180/480 秒)的处理中提醒;每条 allowlist 普通消息按 conversation(chat + user)串行交给 AssistantManager。Worker 落定结果(成功、失败或提问)时,AssistantManager 生成事件文案并经 `Gateway.sendEvent` 由平台 adapter 作为**新消息**发出(QQ 也是新消息,不引用原消息);发送失败只记录不含消息正文或 provider 错误详情的安全日志,不影响任务或后续发送。 +Gateway 不维护普通消息队列或 per-chat task lock,也不再有固定时间(15/60/180/480 秒)的处理中提醒;每条 allowlist 普通消息按 conversation(chat + user)串行交给 AssistantManager。Worker 落定结果(成功、失败或提问)时,AssistantManager 生成事件文案并经 `Gateway.sendEvent` 发出。QQ 无主动消息权限(HTTP 400 / code 40034105)时事件走被动回复窗口:`sendEvent` 引用该 chat 最近一次入站消息(`msg_id` + 递增 `msg_seq`),群聊窗口 5 分钟(按 4.5 分钟保守判定)、C2C 窗口 60 分钟(按 55 分钟保守判定);超过窗口或没有 messageId 时不发主动消息,记录为 skipped 并在该 chat 下次入站时补发。发送失败只记录不含消息正文或 provider 错误详情的安全日志(仅 HTTP status / QQ code),不影响任务或后续发送。 每条异步入站使用独立、串行的回复流,`replySequence` 从 1 动态递增;同步 webhook 直接返回 JSON 结果。 diff --git a/src/acp/assistant-manager.ts b/src/acp/assistant-manager.ts index 6823347..766961e 100644 --- a/src/acp/assistant-manager.ts +++ b/src/acp/assistant-manager.ts @@ -135,9 +135,30 @@ export class AssistantManager implements ConversationRuntime { } status(platform: string, chatId: string, userId?: string): Record { + const chatKey = chatKeyFor(platform, chatId); const proposals = this.proposals.list(); + const owned = (proposal: Proposal) => Boolean(userId) && proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId; + const mine = proposals.filter(owned); const active = this.active; const activeProposal = active ? this.proposals.get(active.proposalId) : undefined; + const awaiting = proposals.filter((proposal) => proposal.status === "awaiting_user_confirmation"); + const working = proposals.filter((proposal) => proposal.status === "working"); + const myQueued = mine.filter((proposal) => proposal.status === "queued").length; + const myPending = mine.filter((proposal) => proposal.status === "awaiting_user_confirmation").length; + const myProposed = mine.filter((proposal) => proposal.status === "proposed").length; + const schedulerState = active || working.length > 0 ? "working" : awaiting.length > 0 ? "awaiting_confirmation" : "idle"; + const blockedReason = awaiting.length > 0 + ? (awaiting.some(owned) ? "your proposal is awaiting your confirmation" : "another proposal is awaiting owner confirmation") + : active || working.length > 0 + ? (working.some(owned) || (activeProposal && owned(activeProposal)) ? "your proposal is working" : "another proposal is working") + : "none"; + const nextAction = myPending > 0 + ? "confirm your pending proposal" + : myProposed > 0 + ? "confirm your proposed proposal to queue it" + : myQueued > 0 + ? (schedulerState === "idle" ? "start_next" : "wait for the current proposal to settle") + : "none"; return { bot: this.bot.id, agent: this.bot.agent.id, @@ -145,10 +166,15 @@ export class AssistantManager implements ConversationRuntime { assistantSessions: this.assistantWorkers.size, proposals: proposals.length, queuedProposals: proposals.filter((proposal) => proposal.status === "queued").length, - workingProposals: proposals.filter((proposal) => proposal.status === "working").length, - awaitingConfirmation: proposals.filter((proposal) => proposal.status === "awaiting_user_confirmation").length, + workingProposals: working.length, + awaitingConfirmation: awaiting.length, workerRunning: Boolean(active?.worker.inFlight), - owner: Boolean(userId && activeProposal && activeProposal.ownerChatKey === chatKeyFor(platform, chatId) && activeProposal.requesterUserId === userId) + owner: Boolean(userId && activeProposal && owned(activeProposal)), + myQueuedProposals: myQueued, + myPendingConfirmations: myPending, + schedulerState, + blockedReason, + nextAction }; } @@ -188,8 +214,8 @@ export class AssistantManager implements ConversationRuntime { const parsed = await this.parseAssistantReply(worker, reply); if (worker.isolationViolated) throw new Error("Assistant session attempted forbidden tool activity"); await this.store.touchBinding(conversationKey); - await this.executeActions(chatKeyFor(request.platform, request.chatId), request, parsed.actions); - return parsed.reply; + const corrections = await this.executeActions(chatKeyFor(request.platform, request.chatId), request, parsed.actions); + return corrections.length > 0 ? `${parsed.reply}\n${corrections.join("\n")}` : parsed.reply; } catch (error) { if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined); throw error; @@ -203,7 +229,8 @@ export class AssistantManager implements ConversationRuntime { return parsed; } - private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise { + private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise { + const corrections: string[] = []; for (const action of actions) { if (action.type === "create_proposal") { await this.proposals.create({ @@ -214,15 +241,33 @@ export class AssistantManager implements ConversationRuntime { requesterUserId: request.userId }); } else if (action.type === "confirm") { - await this.scheduler(() => this.confirmLocked(chatKey, request.userId, action.id, action.answer)); + const ok = await this.scheduler(() => this.confirmLocked(chatKey, request.userId, action.id, action.answer)); + if (!ok) corrections.push("没有可确认的提案;你只能操作自己发起的任务。"); } else if (action.type === "start_next") { - await this.scheduler(() => this.tryStartLocked(chatKey, request.userId)); + const started = await this.scheduler(() => this.tryStartLocked(chatKey, request.userId)); + if (!started) corrections.push(this.startNextCorrection(chatKey, request.userId)); } else if (action.type === "cancel") { - await this.scheduler(() => this.cancelLocked(chatKey, request.userId, action.id)); + const ok = await this.scheduler(() => this.cancelLocked(chatKey, request.userId, action.id)); + if (!ok) corrections.push("没有可取消的提案;你只能操作自己发起的任务。"); } else if (action.type === "stop") { - await this.scheduler(() => this.stopActiveLocked(chatKey, request.userId)); + const ok = await this.scheduler(() => this.stopActiveLocked(chatKey, request.userId)); + if (!ok) corrections.push("没有可停止的任务;你只能操作自己发起的任务。"); } } + return corrections; + } + + private startNextCorrection(chatKey: string, userId: string): string { + const queuedMine = this.proposals.list({ status: "queued" }) + .some((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId); + const kept = queuedMine ? "你的任务已保留在队列里。" : ""; + if (this.proposals.list({ status: "awaiting_user_confirmation" }).length > 0) { + return `还不能开始:还有一个任务在等发起人确认。${kept}`.trim(); + } + if (this.active || this.proposals.list({ status: "working" }).length > 0) { + return `还不能开始:还有一个任务正在执行。${kept}`.trim(); + } + return "还不能开始:你没有已确认并排队的任务。"; } private async confirmLocked(chatKey: string, userId: string, id?: string, answer?: string): Promise { @@ -310,15 +355,16 @@ export class AssistantManager implements ConversationRuntime { await this.proposals.update(proposalId, { workerProcessGroup: undefined }); } - private async tryStartLocked(chatKey: string, userId: string): Promise { - if (this.active) return; - if (this.proposals.list({ status: "working" }).length > 0) return; - if (this.proposals.list({ status: "awaiting_user_confirmation" }).length > 0) return; + private async tryStartLocked(chatKey: string, userId: string): Promise { + if (this.active) return false; + if (this.proposals.list({ status: "working" }).length > 0) return false; + if (this.proposals.list({ status: "awaiting_user_confirmation" }).length > 0) return false; const next = this.proposals.list({ status: "queued" }) .find((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId); - if (!next) return; + if (!next) return false; await this.proposals.update(next.id, { status: "working", startedAt: Date.now() }); this.launchWorker(next.id); + return true; } private resolveTargetProposal(chatKey: string, userId: string, id: string | undefined, statuses: Proposal["status"][]): Proposal | undefined { @@ -600,11 +646,26 @@ export class AssistantManager implements ConversationRuntime { "[Proposal states]", ...this.proposalSummaries(chatKey, request.userId), "", + "[Scheduler state]", + this.schedulerStateSummary(chatKey, request.userId), + "", "[Worker state]", this.workerStateSummary(chatKey, request.userId) ].join("\n"); } + private schedulerStateSummary(chatKey: string, userId: string): string { + const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId; + const awaitingOther = this.proposals.list({ status: "awaiting_user_confirmation" }).some((proposal) => !owned(proposal)); + if (awaitingOther) return "blocked: another proposal is awaiting owner confirmation"; + const active = this.active; + const activeProposal = active ? this.proposals.get(active.proposalId) : undefined; + const workingOther = this.proposals.list({ status: "working" }).some((proposal) => !owned(proposal)) + || Boolean(activeProposal && !owned(activeProposal)); + if (workingOther) return "busy: another proposal is working"; + return "idle"; + } + private proposalSummaries(chatKey: string, userId: string): string[] { const proposals = this.proposals.list() .filter((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId); diff --git a/src/core/gateway.ts b/src/core/gateway.ts index 4e23919..bf61f8a 100644 --- a/src/core/gateway.ts +++ b/src/core/gateway.ts @@ -10,16 +10,33 @@ export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; interface ChatTarget { target: MessageTarget; adapter: PlatformAdapter; + lastInboundMessageId?: string; + lastInboundAt?: number; + inboundIsGroup: boolean; + lastEventDelivery?: "success" | "failed" | "skipped"; + lastEventError?: string; + lastEventAt?: number; } +export interface GatewayOptions { + now?: () => number; +} + +// QQ passive-reply windows (with safety margin): 5 minutes for group chats, 60 minutes for C2C. +const GROUP_PASSIVE_WINDOW_MS = 4.5 * 60_000; +const C2C_PASSIVE_WINDOW_MS = 55 * 60_000; + class ReplyStream { - private sequence = 0; private tail = Promise.resolve(); - constructor(private readonly message: IncomingMessage, private readonly adapter: PlatformAdapter) {} + constructor( + private readonly message: IncomingMessage, + private readonly adapter: PlatformAdapter, + private readonly nextSequence: () => number | undefined + ) {} - enqueue(text: string): Promise { - const replySequence = ++this.sequence; + enqueue(text: string, onError?: (error: unknown) => void): Promise { + const replySequence = this.nextSequence(); this.tail = this.tail.then(() => this.adapter.sendMessage({ target: { platform: this.message.platform, @@ -30,8 +47,9 @@ class ReplyStream { text, replyTo: this.message.messageId, replySequence - })).catch(() => { + })).catch((error: unknown) => { console.error(`Gateway delivery send failed (platform=${this.message.platform})`); + onError?.(error); }); return this.tail; } @@ -39,13 +57,21 @@ class ReplyStream { export class Gateway { private readonly chatTargets = new Map(); + private readonly pendingEvents = new Map(); + // QQ deduplicates passive replies by msg_id + msg_seq: all sends referencing the same inbound + // message (normal replies, events, redelivered pending events) share one counter per chat+messageId. + private readonly messageSequences = new Map>(); + private readonly now: () => number; private closed = false; readonly commandRouter = new CommandRouter(); constructor( private readonly policy: GatewayPolicy, - private readonly runtime: ConversationRuntime - ) {} + private readonly runtime: ConversationRuntime, + options: GatewayOptions = {} + ) { + this.now = options.now || Date.now; + } async receive(message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean } = {}): Promise { const policyError = this.checkPolicy(message); @@ -54,20 +80,29 @@ export class Gateway { return { ok: true, ignored: true, error: policyError }; } - this.chatTargets.set(chatKeyFor(message.platform, message.chatId), { + const chatKey = chatKeyFor(message.platform, message.chatId); + const previous = this.chatTargets.get(chatKey); + this.chatTargets.set(chatKey, { target: { platform: message.platform, chatId: message.chatId, userId: message.userId, raw: message.raw }, - adapter + adapter, + lastInboundMessageId: message.messageId, + lastInboundAt: this.now(), + inboundIsGroup: Boolean(message.isGroup), + lastEventDelivery: previous?.lastEventDelivery, + lastEventError: previous?.lastEventError, + lastEventAt: previous?.lastEventAt }); + await this.flushPendingEvents(chatKey, message, adapter, options); const command = this.commandRouter.parse(message.text); if (command) { const result = await this.executeCommandResult(command, message); - if (!options.synchronous) await new ReplyStream(message, adapter).enqueue(result.reply || ""); + if (!options.synchronous) await this.replyStream(chatKey, message, adapter).enqueue(result.reply || ""); return result; } @@ -81,12 +116,12 @@ export class Gateway { text: message.text, messageId: message.messageId })).text; - await new ReplyStream(message, adapter).enqueue(reply); + await this.replyStream(chatKey, message, adapter).enqueue(reply); return { ok: true, reply }; } catch (error) { const errorText = error instanceof Error ? error.message : String(error); const reply = `Agent error: ${errorText}`; - await new ReplyStream(message, adapter).enqueue(reply); + await this.replyStream(chatKey, message, adapter).enqueue(reply); return { ok: false, error: errorText, reply }; } } @@ -95,13 +130,89 @@ export class Gateway { if (this.closed) return; const entry = this.chatTargets.get(chatKey); if (!entry) return; + const windowMs = entry.inboundIsGroup ? GROUP_PASSIVE_WINDOW_MS : C2C_PASSIVE_WINDOW_MS; + const freshInbound = entry.lastInboundMessageId + && entry.lastInboundAt !== undefined + && this.now() - entry.lastInboundAt <= windowMs; + if (!freshInbound) { + // No fresh passive window and proactive messages may be unauthorized: skip and remind on the next inbound. + this.recordEventDelivery(chatKey, "skipped", "no fresh passive window"); + this.queuePendingEvent(chatKey, text); + console.log(`Gateway event skipped: no fresh passive window (platform=${entry.target.platform})`); + return; + } try { - await entry.adapter.sendMessage({ target: entry.target, text }); - } catch { + await entry.adapter.sendMessage({ + target: entry.target, + text, + replyTo: entry.lastInboundMessageId, + replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId) + }); + this.recordEventDelivery(chatKey, "success"); + } catch (error) { + this.recordEventDelivery(chatKey, "failed", safeEventError(error)); + this.queuePendingEvent(chatKey, text); console.error(`Gateway event send failed (platform=${entry.target.platform})`); } } + private recordEventDelivery(chatKey: string, delivery: "success" | "failed" | "skipped", error?: string): void { + const entry = this.chatTargets.get(chatKey); + if (!entry) return; + entry.lastEventDelivery = delivery; + entry.lastEventError = error; + entry.lastEventAt = this.now(); + } + + private queuePendingEvent(chatKey: string, text: string): void { + const pending = this.pendingEvents.get(chatKey) || []; + pending.push(text); + this.pendingEvents.set(chatKey, pending); + } + + private async flushPendingEvents(chatKey: string, message: IncomingMessage, adapter: PlatformAdapter, options: { synchronous?: boolean }): Promise { + if (options.synchronous) return; + const pending = this.pendingEvents.get(chatKey); + if (!pending || pending.length === 0) return; + this.pendingEvents.delete(chatKey); + const stream = this.replyStream(chatKey, message, adapter); + for (const text of pending) { + await stream.enqueue(text, (error) => { + // Redelivery failed: keep the event queued for the next inbound and record the failure. + this.recordEventDelivery(chatKey, "failed", safeEventError(error)); + this.queuePendingEvent(chatKey, text); + }); + } + } + + private replyStream(chatKey: string, message: IncomingMessage, adapter: PlatformAdapter): ReplyStream { + return new ReplyStream(message, adapter, () => this.nextReplySequence(chatKey, message.messageId)); + } + + private nextReplySequence(chatKey: string, messageId: string | undefined): number | undefined { + if (!messageId) return undefined; + let sequences = this.messageSequences.get(chatKey); + if (!sequences) { + sequences = new Map(); + this.messageSequences.set(chatKey, sequences); + } + const next = (sequences.get(messageId) || 0) + 1; + // Re-set to refresh insertion order, then bound memory to the most recent message ids per chat. + sequences.delete(messageId); + sequences.set(messageId, next); + while (sequences.size > 8) sequences.delete(sequences.keys().next().value!); + return next; + } + + private eventStatus(chatKey: string): Record { + const entry = this.chatTargets.get(chatKey); + return { + lastEventDelivery: entry?.lastEventDelivery || "none", + lastEventError: entry?.lastEventError || "none", + lastEventAt: entry?.lastEventAt ? new Date(entry.lastEventAt).toISOString() : "never" + }; + } + stats(): ReturnType { return this.runtime.stats(); } @@ -140,7 +251,8 @@ export class Gateway { case "help": return HELP_TEXT; case "status": { const status = this.runtime.status(message.platform, message.chatId, message.userId); - return `OK\n${Object.entries(status).map(([key, value]) => `${key}=${value}`).join("\n")}`; + const combined = { ...status, ...this.eventStatus(chatKeyFor(message.platform, message.chatId)) }; + return `OK\n${Object.entries(combined).map(([key, value]) => `${key}=${value}`).join("\n")}`; } case "list": { if (!this.runtime.listProposals) return "当前运行时不支持列出提案。"; @@ -170,6 +282,15 @@ export class Gateway { } +function safeEventError(error: unknown): string { + const message = error instanceof Error ? error.message : ""; + const http = /HTTP (\d{3})/.exec(message); + const code = /code=(\d+)/.exec(message); + const parts = [http ? `HTTP ${http[1]}` : "", code ? `QQ code ${code[1]}` : ""].filter(Boolean); + if (parts.length > 0) return parts.join(" "); + return error instanceof Error ? error.name : "error"; +} + const HELP_TEXT = [ "直接用自然语言告诉我你要做什么即可,我会自己判断并处理。", "兜底命令:", diff --git a/src/platforms/qq/adapter.ts b/src/platforms/qq/adapter.ts index 1012d0b..eb16145 100644 --- a/src/platforms/qq/adapter.ts +++ b/src/platforms/qq/adapter.ts @@ -66,17 +66,24 @@ export class QqAdapter implements PlatformAdapter { const path = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}/messages` : `/v2/users/${encodeURIComponent(targetId)}/messages`; console.log(`QQ send ${isGroup ? "group" : "user"} message length=${message.text.length}`); + // Passive replies reference the inbound message; proactive messages must send neither msg_id nor msg_seq. + const body: Record = { content: message.text }; + if (message.replyTo) { + body.msg_id = message.replyTo; + body.msg_seq = message.replySequence ?? 1; + } const response = await fetch(`https://api.sgroup.qq.com${path}`, { method: "POST", headers: { Authorization: `QQBot ${token}`, "Content-Type": "application/json; charset=utf-8" }, - body: JSON.stringify({ content: message.text, msg_id: message.replyTo, msg_seq: message.replySequence ?? 1 }) + body: JSON.stringify(body) }); - if (!response.ok) throw new Error(`QQ send failed: HTTP ${response.status}`); - const data = await response.json() as QqSendMessageResponse; - if (data.code && data.code !== 0) throw new Error(`QQ send failed: ${data.code} ${data.message || ""}`.trim()); + const data = await response.json().catch(() => undefined) as QqSendMessageResponse | undefined; + // Keep only safe error details (HTTP status / QQ code / QQ message); never log the request body. + if (!response.ok) throw new Error(`QQ send failed: HTTP ${response.status}${qqErrorDetail(data)}`); + if (data?.code && data.code !== 0) throw new Error(`QQ send failed:${qqErrorDetail(data)}`); } private handleValidation(data: QqWebhookEventData | undefined): WebhookResponse { @@ -141,6 +148,12 @@ export class QqAdapter implements PlatformAdapter { } } +function qqErrorDetail(data: QqSendMessageResponse | undefined): string { + if (!data) return ""; + const parts = [data.code ? `code=${data.code}` : "", data.message || ""].filter(Boolean); + return parts.length > 0 ? ` ${parts.join(" ")}` : ""; +} + function userIdFor(data: QqWebhookEventData, isGroup: boolean): string | undefined { if (isGroup) return data.author?.member_openid || data.member_openid || data.author?.user_openid || data.user_openid || data.author?.id; return data.author?.user_openid || data.user_openid || data.author?.id || data.member_openid || data.author?.member_openid; diff --git a/src/roles/role-registry.ts b/src/roles/role-registry.ts index 3ab745c..6b24640 100644 --- a/src/roles/role-registry.ts +++ b/src/roles/role-registry.ts @@ -2,7 +2,7 @@ import crypto from "node:crypto"; import type { AppConfig, BotConfig } from "../config.js"; import { SkillLoader, type LoadedSkill } from "./skill-loader.js"; -const BOOTSTRAP_SCHEMA_VERSION = 5; +const BOOTSTRAP_SCHEMA_VERSION = 6; export interface ResolvedBot extends BotConfig { loadedSkills: LoadedSkill[]; @@ -45,6 +45,7 @@ function buildAssistantBootstrap(bot: BotConfig): string { "Every reply must end with exactly one hidden action envelope: {\"reply\":\"...\",\"actions\":[...]}. Put the user-facing text in the JSON \"reply\" field, not outside the envelope.", "Supported actions: {\"type\":\"create_proposal\",\"title\":\"...\",\"goal\":\"...\",\"steps\":[\"...\"]}; {\"type\":\"confirm\",\"id\":\"optional\",\"answer\":\"optional\"}; {\"type\":\"start_next\"}; {\"type\":\"cancel\",\"id\":\"optional\"}; {\"type\":\"stop\"}. Use an empty actions array when no state change is needed.", "A proposal only starts after the user confirms it. Confirming a proposal whose worker reported SUCCESS marks it completed; confirming a FAILED proposal marks it failed; to retry a failed proposal, confirm with answer \"retry\". When a worker asks a question (NEEDS_CONFIRMATION), confirm with the user's answer to continue the same worker. \"stop\" terminates the active worker and cancels its proposal; \"start_next\" starts the oldest confirmed queued proposal. Never invent other actions or statuses.", + "Never claim an action has already taken effect. The runtime executes your actions after your reply and appends a correction to your message when something could not be done (for example when start_next is blocked by another proposal). Treat that correction as the truth and use the [Proposal states], [Scheduler state] and [Worker state] sections in each prompt as the only reliable state.", "In a group chat each proposal is owned by the user who requested it: only that user can confirm, stop, or cancel it. If another group member asks you to confirm, stop, or cancel a proposal they did not create, explain that only the proposal's initiator can do that and emit no action.", "When a worker reported SUCCESS and the user says something like \"够了\", \"停止\", \"不用继续\" or \"that's enough\", treat it as accepting the finished work: emit a \"confirm\" action so the proposal settles as completed. Use \"stop\" only to abort work that is still running or waiting on an answer." ].join("\n\n"); diff --git a/test/assistant-manager.test.ts b/test/assistant-manager.test.ts index 4d67eb2..98a85eb 100644 --- a/test/assistant-manager.test.ts +++ b/test/assistant-manager.test.ts @@ -437,6 +437,85 @@ test("envelope parsers accept valid tails and reject invalid ones", () => { assert.equal(parseWorkerResult(`x\n{"status":"DONE","summary":"s"}`), undefined); }); +test("start_next blocked by another user's pending confirmation appends a correction and exposes no details", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: succeed", "user-a")); + const first = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm", "user-a")); + await waitFor(() => harness.proposals.get(first.id)!.status === "awaiting_user_confirmation"); + + await harness.manager.prompt(request("create proposal: succeed", "user-b")); + const second = harness.proposals.list().find((proposal) => proposal.id !== first.id)!; + await harness.manager.prompt(request("confirm", "user-b")); + assert.equal(harness.proposals.get(second.id)!.status, "queued"); + + const blocked = await harness.manager.prompt(request("start next", "user-b")); + assert.match(blocked.text, /还不能开始/); + assert.match(blocked.text, /发起人确认/); + assert.match(blocked.text, /保留在队列里/); + assert.equal(harness.proposals.get(second.id)!.status, "queued"); + assert.equal(harness.proposals.get(first.id)!.status, "awaiting_user_confirmation"); + + const bPrompt = readLog(harness.logFile).find((entry) => entry.method === "session/prompt" && entry.text?.startsWith("[User message]\nstart next")); + assert.ok(bPrompt?.text); + assert.match(bPrompt.text!, /blocked: another proposal is awaiting owner confirmation/); + assert.doesNotMatch(bPrompt.text!, /goal done/); + + const bStatus = harness.manager.status("webhook", "chat-1", "user-b"); + assert.equal(bStatus.schedulerState, "awaiting_confirmation"); + assert.equal(bStatus.blockedReason, "another proposal is awaiting owner confirmation"); + assert.equal(bStatus.myQueuedProposals, 1); + assert.equal(bStatus.myPendingConfirmations, 0); + assert.equal(bStatus.nextAction, "wait for the current proposal to settle"); + + const aStatus = harness.manager.status("webhook", "chat-1", "user-a"); + assert.equal(aStatus.myPendingConfirmations, 1); + assert.equal(aStatus.blockedReason, "your proposal is awaiting your confirmation"); + assert.equal(aStatus.nextAction, "confirm your pending proposal"); + } finally { + await closeHarness(harness); + } +}); + +test("confirm and stop actions with no matching proposal append an ownership correction", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: hang", "user-a")); + + const confirmed = await harness.manager.prompt(request("confirm", "user-b")); + assert.match(confirmed.text, /没有可确认的提案/); + assert.match(confirmed.text, /自己发起/); + + const stopped = await harness.manager.prompt(request("stop", "user-b")); + assert.match(stopped.text, /没有可停止的任务/); + assert.match(stopped.text, /自己发起/); + + const cancelled = await harness.manager.prompt(request("cancel proposal", "user-b")); + assert.match(cancelled.text, /没有可取消的提案/); + assert.match(cancelled.text, /自己发起/); + + assert.equal(harness.proposals.list()[0]!.status, "proposed"); + } finally { + await closeHarness(harness); + } +}); + +test("start_next reports when the queue is empty instead of claiming a start", async () => { + const harness = await createHarness(); + try { + const reply = await harness.manager.prompt(request("start next", "user-a")); + assert.match(reply.text, /还不能开始/); + assert.match(reply.text, /没有已确认并排队/); + const status = harness.manager.status("webhook", "chat-1", "user-a"); + assert.equal(status.schedulerState, "idle"); + assert.equal(status.blockedReason, "none"); + assert.equal(status.nextAction, "none"); + } finally { + await closeHarness(harness); + } +}); + function processGroupExists(pgid: number): boolean { try { process.kill(-pgid, 0); return true; } catch (error) { return (error as NodeJS.ErrnoException).code === "EPERM"; } diff --git a/test/gateway.test.ts b/test/gateway.test.ts index f844172..87ecc11 100644 --- a/test/gateway.test.ts +++ b/test/gateway.test.ts @@ -68,6 +68,9 @@ const policy = { allowedUsers: [] as string[], allowedChats: [] as string[], req 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(); @@ -78,6 +81,25 @@ 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) }; } @@ -159,7 +181,7 @@ test("runtime errors remain runtime errors and are delivered through the same st assert.deepEqual(sent.map((outgoing) => [outgoing.text, outgoing.replySequence]), [["Agent error: failed", 1]]); }); -test("sendEvent pushes a fresh message to the latest chat target without replyTo", async () => { +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)); @@ -168,11 +190,14 @@ test("sendEvent pushes a fresh message to the latest chat target without replyTo await turn; await gateway.sendEvent("qq:chat", "任务完成了"); - const event = sent.at(-1)!; - assert.equal(event.text, "任务完成了"); - assert.equal(event.replyTo, undefined); - assert.equal(event.replySequence, undefined); - assert.deepEqual(event.target, { platform: "qq", chatId: "chat", userId: "user", raw: undefined }); + 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 () => { @@ -190,7 +215,8 @@ test("sendEvent uses the most recent target and adapter for the chat", async () 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, undefined); + 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 () => { @@ -198,13 +224,187 @@ test("sendEvent drops safely when the chat has no known target", async () => { 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.replyTo) throw new Error("sensitive provider response"); + if (outgoing.text.includes("event text")) throw new Error("sensitive provider response"); } }; const turn = gateway.receive(message("work", "event-failure"), adapter); diff --git a/test/qq-adapter.test.ts b/test/qq-adapter.test.ts index 877ff3c..c6f1198 100644 --- a/test/qq-adapter.test.ts +++ b/test/qq-adapter.test.ts @@ -51,3 +51,40 @@ test("sendMessage uses C2C endpoint and forwards reply sequences as msg_seq", as ]); } finally { globalThis.fetch = original; } }); + +test("sendMessage without replyTo omits msg_id and msg_seq from the body", async () => { + const original = globalThis.fetch; const bodies: Array> = []; + globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { + 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); + const target = { platform: "qq", chatId: "group:g1", raw: { group_openid: "g1" } }; + await adapter.sendMessage({ target, text: "proactive" }); + assert.deepEqual(bodies, [{ content: "proactive" }]); + } finally { globalThis.fetch = original; } +}); + +test("sendMessage surfaces safe QQ error details (HTTP status and code) without the request body", async () => { + const original = globalThis.fetch; + globalThis.fetch = (async (input: string | URL | Request) => { + if (String(input).includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + return new Response(JSON.stringify({ code: 40034105, message: "proactive message not allowed" }), { status: 400 }); + }) as typeof fetch; + try { + const adapter = new QqAdapter(config, { receive: async () => ({ ok: true }) } as never); + const target = { platform: "qq", chatId: "group:g1", raw: { group_openid: "g1" } }; + await assert.rejects( + adapter.sendMessage({ target, text: "secret message content" }), + (error: Error) => { + assert.match(error.message, /HTTP 400/); + assert.match(error.message, /code=40034105/); + assert.match(error.message, /proactive message not allowed/); + assert.doesNotMatch(error.message, /secret message content/); + return true; + } + ); + } finally { globalThis.fetch = original; } +});