From 6b506b8c55de74f0ad04d817e471c0ff52481ecf Mon Sep 17 00:00:00 2001 From: zenord Date: Wed, 19 Aug 2026 14:39:18 +0800 Subject: [PATCH] feat: send workspace images to QQ chats Worker results may report image attachments stored inside the workspace; the runtime validates them (containment, png/jpg magic, size) and delivers them with the pending event. Assistants gain a send_image action so users can ask for an image later. QQ uploads via /files with base64 file_data and sends msg_type 7 rich media, sharing the same msg_id/msg_seq counter as text replies; non-image adapters flatten images to text. --- .kimi-code/skills/gori-agent-deploy/SKILL.md | 2 + AGENTS.md | 6 +- README.md | 6 +- src/acp/assistant-manager.ts | 83 ++++++++-- src/acp/types.ts | 5 +- src/core/adapter.ts | 3 + src/core/gateway.ts | 128 +++++++++++---- src/core/proposal-store.ts | 18 +++ src/core/types.ts | 14 ++ src/core/workspace-images.ts | 66 ++++++++ src/platforms/qq/adapter.ts | 44 +++++- src/roles/role-registry.ts | 7 +- src/server.ts | 2 +- test/assistant-manager.test.ts | 94 ++++++++++- test/fixtures/fake-acp-agent.mjs | 16 ++ test/gateway.test.ts | 156 ++++++++++++++++++- test/qq-adapter.test.ts | 71 +++++++++ 17 files changed, 659 insertions(+), 62 deletions(-) create mode 100644 src/core/workspace-images.ts diff --git a/.kimi-code/skills/gori-agent-deploy/SKILL.md b/.kimi-code/skills/gori-agent-deploy/SKILL.md index 7a3fed4..2f3c898 100644 --- a/.kimi-code/skills/gori-agent-deploy/SKILL.md +++ b/.kimi-code/skills/gori-agent-deploy/SKILL.md @@ -172,6 +172,8 @@ Config v3 支持: - 同一 `msg_id` 的 `msg_seq` 必须由 Gateway 共享递增,正常回复、事件、补发不能各自从 1 开始。 - `/status` 应能看到 `schedulerState`、`blockedReason`、`lastEventDelivery`、`lastEventError`、`lastEventAt`。 +出站图片同样没有主动权限,走同一被动窗口与共享 `msg_seq`:Worker 把给用户看的 png/jpg 存到 workspace 内(建议 `.gori-outbox/`)并在 result envelope 的 `attachments` 上报,runtime 校验 workspace containment、magic bytes、单张 ≤10MB、最多 3 张;发送时先 `POST /v2/{groups|users}/{id}/files` 上传(`file_type: 1` + base64 `file_data` + `srv_send_msg: false`)再 `msg_type: 7` + `media.file_info` 发送,群/私聊上传的 file_info 不通用。落定事件先发文案再发图;窗口过期补发时按路径重读文件,文件没了降级为文本说明;上传/发送失败不阻断文本。用户后来说「把图发给我」时 Assistant 用 `send_image` action 发同一张图。 + ## 验收清单 部署完成后至少检查: diff --git a/AGENTS.md b/AGENTS.md index 900a836..f619090 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -97,6 +97,8 @@ gori-agent list - proposals.json(version 2)保存 Proposal:title/goal/steps、owner chat、发起用户 requesterUserId、状态流(proposed/queued/working/pending/finished)、pending(summary/question?/workspaceDirty?/receivedAt)、finished(finishKind done|cancelled、finishNote?)、worker session 与进程组 PGID/token;同样的 lock 与原子写入纪律。打开 v1 文件时先写 `proposals.json.v1-.bak`(0600)备份再按固定映射迁移。 - `src/core/workspace-scope.ts` - canonical workspace + 跨实例 Config v3 扫描;相同或父子 workspace 在 doctor/start/runner fail closed。lease 已退役。 +- `src/core/workspace-images.ts` + - 出站图片校验与读取:路径必须 resolved 在 canonical workspace 内、png/jpg magic bytes、单张 ≤10MB;发送时按路径重读。 - `src/core/gateway.ts` - 入站 allowlist、群 mention、命令和回复;不维护普通消息队列或 per-chat task lock,也没有固定时间处理中提醒。 - 每条异步入站使用独立串行 ReplyStream,`replySequence` 从 1 动态递增;发送失败只写安全日志且不阻止任务或后续发送。 @@ -207,7 +209,7 @@ Kimi Code agent `args` 必须严格为 `["acp"]`,不能添加可能绕过 assi Proposal 状态流:`proposed → queued → working → pending → finished`。`proposed` 必须用户确认才进入 `queued`;`pending` 就是「等用户决定」,不区分 success/failure;只有 `finish` 把 pending 落定为 `finished(done)`,`cancel`(proposed/queued/pending)落定为 `finished(cancelled)`;没有 working、owner 自己没有 pending、全局没有 workspaceDirty pending 时,最早确认的 queued Proposal 才可通过 confirm 或 start_next 开始。 -Assistant 每条回复以隐藏 `GORI_ASSISTANT_ACTION_V2` envelope 结尾(`reply` + `actions`),action 仅 `create_proposal`、`confirm`、`adjust_proposal`(仅 proposed/queued)、`follow_up`(pending → working,优先 resume 原 worker native session,失败则带 Proposal 上下文新起 session,用户图片附件随 prompt 给 Worker)、`finish`、`start_next`、`cancel`、`stop`,格式错误只修复一次。Worker 每轮以隐藏 `GORI_WORKER_RESULT_V2` envelope 收尾:仅 `PENDING`(`summary` 必填,可选 `question`、`workspaceDirty`),同样只修复一次。 +Assistant 每条回复以隐藏 `GORI_ASSISTANT_ACTION_V2` envelope 结尾(`reply` + `actions`),action 仅 `create_proposal`、`confirm`、`adjust_proposal`(仅 proposed/queued)、`follow_up`(pending → working,优先 resume 原 worker native session,失败则带 Proposal 上下文新起 session,用户图片附件随 prompt 给 Worker)、`finish`、`send_image`(把 Worker 报告过的 workspace 内图片发给用户)、`start_next`、`cancel`、`stop`,格式错误只修复一次。Worker 每轮以隐藏 `GORI_WORKER_RESULT_V2` envelope 收尾:仅 `PENDING`(`summary` 必填,可选 `question`、`workspaceDirty`、`attachments`),同样只修复一次。 确认语义: @@ -224,6 +226,8 @@ state v3(`acp-sessions.json`)header 保存 `botId` 和 platform,只存 ass 命令 `/help`、`/status`、`/list`、`/confirm`、`/finish`、`/stop`、`/cancel` 旁路 Assistant;全部 owner-only(chat + user):`/confirm` 对 proposed,`/finish` 对 pending,`/stop` 仅对 working(→pending),`/cancel` 对 proposed/queued/pending;`/list` 面板按 pending → working → queued → proposed → 最近 finished 排列。QQ 入站图片附件(image/*,单张 ≤5MB、每条最多 3 张、10 秒下载超时)下载为 base64 经 `IncomingMessage.attachments` 透传;agent 声明 `promptCapabilities.image` 时作为 ACP image content block 发给 Assistant/Worker,否则降级为文本说明;视频/文件附件不下载,仅以 `[视频] ` / `[文件] ` 文本拼接。 +出站图片(QQ):Worker 在 `GORI_WORKER_RESULT_V2` 的 `attachments`(最多 3 个 `{path, mimeType?}`)上报 workspace 内 png/jpg(建议 `.gori-outbox/`),Assistant 也可用 `send_image { path }` 主动发图;runtime 校验 resolved realpath 必须在 canonical workspace 内、png/jpg magic bytes、单张 ≤10MB,违规丢弃并在事件文本说明。发送走 `POST /v2/{groups|users}/{id}/files` 上传(`file_type: 1` + base64 `file_data` + `srv_send_msg: false`)后 `msg_type: 7` + `media.file_info` 发送;群/私聊上传的 file_info 不通用,按目标分别上传。图片与文本共用 Gateway 的 `msg_id` + 递增 `msg_seq` 计数器;落定事件先文案后图;被动窗口过期时照旧记录并下次入站补发,补发按路径重读文件、文件缺失降级为文本说明;上传/发送失败不阻断文本,降级文本说明 + 安全日志。其他平台 adapter 无 `supportsImages` 标记时图片降级为 `[图片] <文件名>` 文本。 + ## 6. 单平台装配 Server 只构造 `gateway.platform.type` 对应 adapter,只挂载对应 webhook: diff --git a/README.md b/README.md index fd8052e..5c6098d 100644 --- a/README.md +++ b/README.md @@ -175,6 +175,8 @@ QQ websocket 示例: QQ 入站附件:`image/*` 附件会被下载(每条消息最多 3 张、单张超过 5MB 跳过、10 秒下载超时、缺 scheme 的 URL 自动补 `https:`)并转成 base64 随消息传给 Agent;agent 声明 `promptCapabilities.image` 时作为 ACP image content block 发送,否则降级为文本说明。视频/文件等非图片附件不下载,仅以 `[视频] ` / `[文件] ` 文本拼进消息,交给模型自由使用;纯图片消息使用占位文本「(发来一张图片)」。日志只记录附件类型/大小/数量,不记录 URL 全文或 base64。 +QQ 出站图片:Worker/Assistant 报告的 workspace 内图片(png/jpg、≤10MB、最多 3 张)经 `POST /v2/{groups|users}/{id}/files` 上传(`file_type: 1`、`file_data` base64、`srv_send_msg: false`)后以 `msg_type: 7` + `media.file_info` 发送;群上传的 file_info 只能发群、私聊上传的只能发私聊,按目标分别上传。无主动消息权限(40034105),图片与文本一样只能走被动回复窗口,且与文本共用同一 `msg_id` + 递增 `msg_seq` 计数器;Worker 落定事件先发文案再发图,窗口过期时按现有规则记录并下次入站补发,补发时按路径重读文件,文件不存在则降级为文本说明。图片上传/发送失败不阻断文本,降级为文本说明 + 安全日志。其他平台 adapter 不支持图片时把图片降级为 `[图片] <文件名>` 文本行。 + ### ACP 生命周期 默认 `runtime.acp.promptTimeoutMs` 是 `14400000`(4 小时),适合长构建/部署任务。其他默认值: @@ -192,9 +194,9 @@ QQ 入站附件:`image/*` 附件会被下载(每条消息最多 3 张、单 运行链路是三层: -- **Assistant**:每个 conversation(chat + user)一个无工具 ACP 会话,只与用户对话;同群不同用户的会话互相隔离。它把用户意图整理成 Proposal(title、goal、steps),并解释 Worker 的反馈。Assistant 每条回复必须以隐藏 `GORI_ASSISTANT_ACTION_V2` envelope 结尾(`reply` + `actions`),action 只有 `create_proposal`、`confirm`、`adjust_proposal`、`follow_up`、`finish`、`start_next`、`cancel`、`stop`;格式错误只修复一次。 +- **Assistant**:每个 conversation(chat + user)一个无工具 ACP 会话,只与用户对话;同群不同用户的会话互相隔离。它把用户意图整理成 Proposal(title、goal、steps),并解释 Worker 的反馈。Assistant 每条回复必须以隐藏 `GORI_ASSISTANT_ACTION_V2` envelope 结尾(`reply` + `actions`),action 只有 `create_proposal`、`confirm`、`adjust_proposal`、`follow_up`、`finish`、`send_image`、`start_next`、`cancel`、`stop`;格式错误只修复一次。`send_image { path }` 把 Worker 报告过的 workspace 内图片发给用户,与 Worker 附件同样的路径/类型/大小校验。 - **Proposal**:一份工作单,owner 是 chat + 发起用户;只有发起人本人可以 confirm/adjust/follow_up/finish/stop/cancel/list 它,`start_next` 也只启动发起人自己的 queued Proposal。状态流为 `proposed → queued → working → pending → finished`。`proposed` 只有用户确认后才进入 `queued`;`pending` 就是「等用户决定」,不再区分 success/failure;只有 `finish` 把 pending 落定为 `finished(done)`,`cancel` 落定为 `finished(cancelled)`。 -- **Worker**:同一时刻全实例只有一个,在 `bot.workspace` 用 `bot.permissions` policy 执行一个已确认 Proposal。每轮必须以隐藏 `GORI_WORKER_RESULT_V2` envelope 收尾:`PENDING`(`summary` 必填,可带 `question`、`workspaceDirty`),不区分成功/失败,只把结果交给用户。 +- **Worker**:同一时刻全实例只有一个,在 `bot.workspace` 用 `bot.permissions` policy 执行一个已确认 Proposal。每轮必须以隐藏 `GORI_WORKER_RESULT_V2` envelope 收尾:`PENDING`(`summary` 必填,可带 `question`、`workspaceDirty`),不区分成功/失败,只把结果交给用户。Worker 给用户看的图片(png/jpg)必须保存在 workspace 内(建议 `.gori-outbox/`),并通过 `attachments: [{ path, mimeType? }]`(最多 3 个)上报;runtime 校验路径必须在 workspace 内、magic bytes 为 png/jpg、单张 ≤10MB,违规的丢弃并在事件文本里说明。 确认语义是刻意的: diff --git a/src/acp/assistant-manager.ts b/src/acp/assistant-manager.ts index a0aae4b..bcd161f 100644 --- a/src/acp/assistant-manager.ts +++ b/src/acp/assistant-manager.ts @@ -4,7 +4,8 @@ import path from "node:path"; import type { AcpConfig, PermissionPolicy } from "../config.js"; import { chatKeyFor, type AssistantBinding, type DurableSessionStore } from "../core/durable-session-store.js"; import type { Proposal, ProposalPending, ProposalStore } from "../core/proposal-store.js"; -import type { IncomingAttachment } from "../core/types.js"; +import type { IncomingAttachment, OutgoingImageRef } from "../core/types.js"; +import { MAX_OUTBOUND_IMAGES, validateWorkspaceImage } from "../core/workspace-images.js"; import type { ResolvedBot } from "../roles/role-registry.js"; import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeEventSink, RuntimeStats } from "./types.js"; import type { AcpPromptContent } from "./client.js"; @@ -20,6 +21,7 @@ export type AssistantAction = | { type: "adjust_proposal"; id?: string; title?: string; goal?: string; steps?: string[] } | { type: "follow_up"; id?: string; instruction: string } | { type: "finish"; id?: string; note?: string } + | { type: "send_image"; path: string } | { type: "start_next" } | { type: "cancel"; id?: string } | { type: "stop" }; @@ -36,6 +38,7 @@ export interface ParsedWorkerResult { summary: string; question?: string; workspaceDirty?: boolean; + attachments?: { path: string; mimeType?: string }[]; } interface ActiveWorker { @@ -111,8 +114,8 @@ export class AssistantManager implements ConversationRuntime { async prompt(request: ConversationRequest): Promise { if (this.shuttingDown) throw new Error("ACP runtime is shutting down"); const conversationKey = assistantKeyFor(chatKeyFor(request.platform, request.chatId), request.userId); - const reply = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request)); - return { text: reply, botId: this.bot.id, agentId: this.bot.agent.id, mode: "assistant" }; + const result = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request)); + return { text: result.text, images: result.images, botId: this.bot.id, agentId: this.bot.agent.id, mode: "assistant" }; } async cancel(platform: string, chatId: string, userId: string): Promise { @@ -226,7 +229,7 @@ export class AssistantManager implements ConversationRuntime { await Promise.allSettled([...this.chatChains.values()]); } - private async runAssistantTurn(conversationKey: string, request: ConversationRequest): Promise { + private async runAssistantTurn(conversationKey: string, request: ConversationRequest): Promise<{ text: string; images?: OutgoingImageRef[] }> { const worker = await this.acquireAssistantWorker(conversationKey); try { const reply = await worker.prompt(this.assistantPromptContent(worker, request)); @@ -234,10 +237,10 @@ export class AssistantManager implements ConversationRuntime { if (worker.isolationViolated) throw new Error("Assistant session attempted forbidden tool activity"); await this.store.touchBinding(conversationKey); const chatKey = chatKeyFor(request.platform, request.chatId); - const corrections = await this.executeActions(chatKey, request, parsed.actions); + const { corrections, images } = await this.executeActions(chatKey, request, parsed.actions); const reminders = this.pendingReminders(chatKey, request.userId, parsed.reply); const extras = [...corrections, ...reminders]; - return extras.length > 0 ? `${parsed.reply}\n${extras.join("\n")}` : parsed.reply; + return { text: extras.length > 0 ? `${parsed.reply}\n${extras.join("\n")}` : parsed.reply, images }; } catch (error) { if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined); throw error; @@ -269,8 +272,9 @@ 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<{ corrections: string[]; images?: OutgoingImageRef[] }> { const corrections: string[] = []; + const images: OutgoingImageRef[] = []; for (const action of actions) { if (action.type === "create_proposal") { await this.proposals.create({ @@ -292,6 +296,10 @@ export class AssistantManager implements ConversationRuntime { } else if (action.type === "finish") { const ok = await this.scheduler(() => this.finishLocked(chatKey, request.userId, action.id, action.note)); if (!ok) corrections.push("没有可结束的任务;只有待确认的任务可以 finish,且你只能操作自己发起的任务。"); + } else if (action.type === "send_image") { + const image = await validateWorkspaceImage(this.bot.workspace, action.path); + if (image) images.push({ path: image.path, mimeType: image.mimeType, filename: image.filename }); + else corrections.push("这张图发不出去:路径不在工作区内,或者不是有效的 png/jpg 图片。"); } else if (action.type === "start_next") { const started = await this.scheduler(() => this.tryStartLocked(chatKey, request.userId)); if (!started) corrections.push(this.startNextCorrection(chatKey, request.userId)); @@ -303,7 +311,7 @@ export class AssistantManager implements ConversationRuntime { if (!ok) corrections.push("没有可停止的任务;你只能操作自己发起的任务。"); } } - return corrections; + return { corrections, images: images.length > 0 ? images.slice(0, MAX_OUTBOUND_IMAGES) : undefined }; } private startNextCorrection(chatKey: string, userId: string): string { @@ -517,11 +525,23 @@ export class AssistantManager implements ConversationRuntime { if (!proposal || proposal.status !== "working") { active.settled = true; return false; } active.settled = true; this.active = undefined; + const attachments: { path: string; mimeType?: string }[] = []; + const droppedAttachments: string[] = []; + for (const ref of result.attachments || []) { + const image = await validateWorkspaceImage(this.bot.workspace, ref.path); + if (image) attachments.push({ path: image.path, mimeType: image.mimeType }); + else droppedAttachments.push(ref.path); + } + if (droppedAttachments.length > 0) { + console.error(`Worker reported ${droppedAttachments.length} invalid attachment(s) for proposal ${active.proposalId}; dropped (outside workspace or not a valid png/jpg)`); + } const pending: ProposalPending = { summary: result.summary, question: result.question, workspaceDirty: result.workspaceDirty, - receivedAt: Date.now() + receivedAt: Date.now(), + attachments: attachments.length > 0 ? attachments : undefined, + droppedAttachments: droppedAttachments.length > 0 ? droppedAttachments : undefined }; await this.proposals.update(active.proposalId, { status: "pending", @@ -578,11 +598,11 @@ export class AssistantManager implements ConversationRuntime { if (this.shuttingDown) return; try { const reply = await this.runAssistantEvent(conversationKey, proposal.id); - if (reply) await sink(chatKey, reply); + if (reply) await sink(chatKey, reply, pendingImageRefs(this.proposals.get(proposalId))); } catch (error) { console.error(`Assistant event notification failed (${error instanceof Error ? error.name : "unknown error"})`); const latest = this.proposals.get(proposalId); - if (latest?.pending) await sink(chatKey, fallbackEventText(latest)).catch(() => undefined); + if (latest?.pending) await sink(chatKey, fallbackEventText(latest), pendingImageRefs(latest)).catch(() => undefined); } }).catch(() => undefined); } @@ -844,6 +864,10 @@ function parseAssistantAction(value: unknown): AssistantAction | undefined { if ((value.id !== undefined && typeof value.id !== "string") || (value.note !== undefined && typeof value.note !== "string")) return undefined; return { type: "finish", id: value.id as string | undefined, note: value.note as string | undefined }; } + if (value.type === "send_image") { + if (typeof value.path !== "string" || value.path.length === 0) return undefined; + return { type: "send_image", path: value.path }; + } if (value.type === "start_next") return { type: "start_next" }; if (value.type === "cancel") { if (value.id !== undefined && typeof value.id !== "string") return undefined; @@ -862,16 +886,27 @@ export function parseWorkerResult(text: string): ParsedWorkerResult | undefined if (value.status !== "PENDING") return undefined; if ((value.question !== undefined && typeof value.question !== "string") || (value.workspaceDirty !== undefined && typeof value.workspaceDirty !== "boolean")) return undefined; + let attachments: { path: string; mimeType?: string }[] | undefined; + if (value.attachments !== undefined) { + if (!Array.isArray(value.attachments) || value.attachments.length > MAX_OUTBOUND_IMAGES) return undefined; + attachments = []; + for (const item of value.attachments) { + if (!isRecord(item) || typeof item.path !== "string" || item.path.length === 0 + || (item.mimeType !== undefined && typeof item.mimeType !== "string")) return undefined; + attachments.push({ path: item.path, mimeType: item.mimeType as string | undefined }); + } + } return { status: "PENDING", summary: value.summary, question: value.question as string | undefined, - workspaceDirty: value.workspaceDirty as boolean | undefined + workspaceDirty: value.workspaceDirty as boolean | undefined, + attachments }; } function withWorkerResultProtocol(input: string | AcpPromptContent[]): string | AcpPromptContent[] { - const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: {\"status\":\"PENDING\",\"summary\":\"...\"}. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. PENDING is the only status; do not emit any other status value or any text after the envelope."; + const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: {\"status\":\"PENDING\",\"summary\":\"...\"}. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 of them as \"attachments\": [{\"path\":\"relative/or/absolute/path\"}]; never report paths outside the workspace. PENDING is the only status; do not emit any other status value or any text after the envelope."; if (typeof input === "string") return `${input}\n\n${instruction}`; return [...input, { type: "text" as const, text: `\n\n${instruction}` }]; } @@ -908,18 +943,38 @@ function workerFollowUpPrompt(instruction: string, question?: string): string { function assistantEventPrompt(proposal: Proposal): string { const pending = proposal.pending!; + const attachments = pending.attachments?.length + ? `Image attachments produced by the worker (already sent to the user after your message): ${pending.attachments.map((attachment) => attachment.path).join(", ")}` + : "Image attachments: none"; + const dropped = pending.droppedAttachments?.length + ? `Note: ${pending.droppedAttachments.length} reported attachment(s) were invalid (outside the workspace or not png/jpg) and were dropped; mention this to the user.` + : ""; return [ "[Internal event from gori-agent. This is not a user message.]", `The worker for proposal "${proposal.title}" (id ${proposal.id}) reported its result and the proposal is now waiting for the user's decision.`, `Summary: ${pending.summary}`, pending.question ? `Question for the user: ${pending.question}` : "Question for the user: none", + attachments, + ...(dropped ? [dropped] : []), "Tell the user, in your own words, what happened and what they can do next (finish to close the task, follow up with new instructions to continue it, or stop). End with the usual GORI_ASSISTANT_ACTION_V2 envelope; its actions array MUST be empty because actions are ignored for internal events." ].join("\n"); } function fallbackEventText(proposal: Proposal): string { const pending = proposal.pending!; - return `任务「${proposal.title}」等你确认。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""} 说 finish 结束,或直接说要求继续改。`; + const attachments = pending.attachments?.length ? `(附 ${pending.attachments.length} 张图片)` : ""; + const dropped = pending.droppedAttachments?.length ? `(${pending.droppedAttachments.length} 个附件校验失败被丢弃)` : ""; + return `任务「${proposal.title}」等你确认。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""}${attachments}${dropped} 说 finish 结束,或直接说要求继续改。`; +} + +function pendingImageRefs(proposal: Proposal | undefined): OutgoingImageRef[] | undefined { + const attachments = proposal?.pending?.attachments; + if (!attachments || attachments.length === 0) return undefined; + return attachments.map((attachment) => ({ + path: attachment.path, + mimeType: attachment.mimeType, + filename: attachment.path.split("/").pop() + })); } function isRecord(value: unknown): value is Record { diff --git a/src/acp/types.ts b/src/acp/types.ts index c639257..8abc5a7 100644 --- a/src/acp/types.ts +++ b/src/acp/types.ts @@ -1,6 +1,6 @@ import type { PermissionPolicy } from "../config.js"; import type { Proposal } from "../core/proposal-store.js"; -import type { IncomingAttachment } from "../core/types.js"; +import type { IncomingAttachment, OutgoingImageRef } from "../core/types.js"; export interface AcpBackendSpec { id: string; @@ -23,6 +23,7 @@ export interface ConversationResponse { botId: string; agentId: string; mode?: "assistant"; + images?: OutgoingImageRef[]; } export interface RuntimeStats { @@ -32,7 +33,7 @@ export interface RuntimeStats { persistedBindings: number; } -export type RuntimeEventSink = (chatKey: string, text: string) => Promise; +export type RuntimeEventSink = (chatKey: string, text: string, images?: OutgoingImageRef[]) => Promise; export interface ConversationRuntime { prompt(request: ConversationRequest): Promise; diff --git a/src/core/adapter.ts b/src/core/adapter.ts index 8664220..b0dc1fd 100644 --- a/src/core/adapter.ts +++ b/src/core/adapter.ts @@ -2,6 +2,9 @@ import type { OutgoingMessage, WebhookRequestContext, WebhookResponse } from "./ export interface PlatformAdapter { name: string; + // True when the adapter can deliver OutgoingMessage.images (e.g. QQ msg_type 7); + // gateways degrade images to text notes for adapters without this flag. + supportsImages?: boolean; handleWebhook(context: WebhookRequestContext): Promise; sendMessage(message: OutgoingMessage): Promise; } diff --git a/src/core/gateway.ts b/src/core/gateway.ts index cf93f82..59cb581 100644 --- a/src/core/gateway.ts +++ b/src/core/gateway.ts @@ -4,7 +4,8 @@ import type { PlatformAdapter } from "./adapter.js"; import { CommandRouter, type ParsedCommand } from "./command-router.js"; import { chatKeyFor } from "./durable-session-store.js"; import type { Proposal } from "./proposal-store.js"; -import type { IncomingMessage, MessageTarget } from "./types.js"; +import type { IncomingMessage, MessageTarget, OutgoingImage, OutgoingImageRef } from "./types.js"; +import { readOutboundImage } from "./workspace-images.js"; export interface GatewayResult { ok: boolean; reply?: string; ignored?: boolean; error?: string } @@ -19,6 +20,11 @@ interface ChatTarget { lastEventAt?: number; } +interface PendingEvent { + text: string; + images?: OutgoingImageRef[]; +} + export interface GatewayOptions { now?: () => number; } @@ -27,6 +33,28 @@ export interface GatewayOptions { const GROUP_PASSIVE_WINDOW_MS = 4.5 * 60_000; const C2C_PASSIVE_WINDOW_MS = 55 * 60_000; +// Reads image files at send time (the worker may have updated them); unreadable files degrade +// to a text note, and adapters without image support get a plain [图片] text line instead. +async function resolveOutboundImages(text: string, refs: OutgoingImageRef[] | undefined, imageCapable: boolean): Promise<{ text: string; images: OutgoingImage[] }> { + if (!refs || refs.length === 0) return { text, images: [] }; + const images: OutgoingImage[] = []; + const failures: string[] = []; + for (const ref of refs) { + const image = await readOutboundImage(ref); + if (image) images.push(image); + else failures.push(ref.filename || "image"); + } + let finalText = text; + if (failures.length > 0) { + finalText = [finalText, ...failures.map((name) => `(图片 ${name} 发送失败:文件不存在或已失效)`)].filter((part) => part.length > 0).join("\n"); + } + if (!imageCapable && images.length > 0) { + finalText = [finalText, ...images.map((image) => `[图片] ${image.filename || "image"}`)].filter((part) => part.length > 0).join("\n"); + return { text: finalText, images: [] }; + } + return { text: finalText, images }; +} + class ReplyStream { private tail = Promise.resolve(); @@ -36,29 +64,47 @@ class ReplyStream { private readonly nextSequence: () => number | undefined ) {} - enqueue(text: string, onError?: (error: unknown) => void): Promise { - const replySequence = this.nextSequence(); - this.tail = this.tail.then(() => this.adapter.sendMessage({ + enqueue(entry: string | PendingEvent, onError?: (error: unknown) => void): Promise { + const request = typeof entry === "string" ? { text: entry } : entry; + this.tail = this.tail.then(() => this.deliver(request)).catch((error: unknown) => { + console.error(`Gateway delivery send failed (platform=${this.message.platform})`); + onError?.(error); + }); + return this.tail; + } + + private async deliver(request: PendingEvent): Promise { + const { text, images } = await resolveOutboundImages(request.text, request.images, this.adapter.supportsImages === true); + if (text || images.length === 0) await this.send({ text }); + for (const image of images) { + try { + await this.send({ text: "", images: [image] }); + } catch (error) { + console.error(`Gateway image delivery failed (platform=${this.message.platform})`); + await this.send({ text: `(图片 ${image.filename || "image"} 发送失败)` }).catch(() => undefined); + } + } + } + + private send(piece: { text: string; images?: OutgoingImage[] }): Promise { + return this.adapter.sendMessage({ target: { platform: this.message.platform, chatId: this.message.chatId, userId: this.message.userId, raw: this.message.raw }, - text, + text: piece.text, replyTo: this.message.messageId, - replySequence - })).catch((error: unknown) => { - console.error(`Gateway delivery send failed (platform=${this.message.platform})`); - onError?.(error); + replySequence: this.nextSequence(), + images: piece.images }); - return this.tail; } } export class Gateway { private readonly chatTargets = new Map(); - private readonly pendingEvents = 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>(); @@ -110,16 +156,16 @@ export class Gateway { if (options.synchronous) return this.promptResult(message); try { - const reply = (await this.runtime.prompt({ + const response = await this.runtime.prompt({ platform: message.platform, chatId: message.chatId, userId: message.userId, text: message.text, messageId: message.messageId, attachments: message.attachments - })).text; - await this.replyStream(chatKey, message, adapter).enqueue(reply); - return { ok: true, reply }; + }); + await this.replyStream(chatKey, message, adapter).enqueue({ text: response.text, images: response.images }); + return { ok: true, reply: response.text }; } catch (error) { const errorText = error instanceof Error ? error.message : String(error); const reply = `Agent error: ${errorText}`; @@ -128,7 +174,7 @@ export class Gateway { } } - async sendEvent(chatKey: string, text: string): Promise { + async sendEvent(chatKey: string, text: string, images?: OutgoingImageRef[]): Promise { if (this.closed) return; const entry = this.chatTargets.get(chatKey); if (!entry) return; @@ -139,21 +185,44 @@ export class Gateway { 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); + this.queuePendingEvent(chatKey, { text, images }); console.log(`Gateway event skipped: no fresh passive window (platform=${entry.target.platform})`); return; } try { - await entry.adapter.sendMessage({ - target: entry.target, - text, - replyTo: entry.lastInboundMessageId, - replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId) - }); + const resolved = await resolveOutboundImages(text, images, entry.adapter.supportsImages === true); + if (resolved.text || resolved.images.length === 0) { + await entry.adapter.sendMessage({ + target: entry.target, + text: resolved.text, + replyTo: entry.lastInboundMessageId, + replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId) + }); + } + for (const image of resolved.images) { + try { + await entry.adapter.sendMessage({ + target: entry.target, + text: "", + replyTo: entry.lastInboundMessageId, + replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId), + images: [image] + }); + } catch (imageError) { + // Image upload/send failures never block the event text; degrade to a text note. + console.error(`Gateway event image send failed (platform=${entry.target.platform})`); + await entry.adapter.sendMessage({ + target: entry.target, + text: `(图片 ${image.filename || "image"} 发送失败)`, + replyTo: entry.lastInboundMessageId, + replySequence: this.nextReplySequence(chatKey, entry.lastInboundMessageId) + }).catch(() => undefined); + } + } this.recordEventDelivery(chatKey, "success"); } catch (error) { this.recordEventDelivery(chatKey, "failed", safeEventError(error)); - this.queuePendingEvent(chatKey, text); + this.queuePendingEvent(chatKey, { text, images }); console.error(`Gateway event send failed (platform=${entry.target.platform})`); } } @@ -166,9 +235,9 @@ export class Gateway { entry.lastEventAt = this.now(); } - private queuePendingEvent(chatKey: string, text: string): void { + private queuePendingEvent(chatKey: string, event: PendingEvent): void { const pending = this.pendingEvents.get(chatKey) || []; - pending.push(text); + pending.push(event); this.pendingEvents.set(chatKey, pending); } @@ -178,11 +247,12 @@ export class Gateway { 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) => { + for (const event of pending) { + // Images are re-read from disk at redelivery time (the worker may have updated them). + await stream.enqueue(event, (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); + this.queuePendingEvent(chatKey, event); }); } } diff --git a/src/core/proposal-store.ts b/src/core/proposal-store.ts index 99866d9..41cb6b2 100644 --- a/src/core/proposal-store.ts +++ b/src/core/proposal-store.ts @@ -19,11 +19,18 @@ export type ProposalStatus = | "pending" | "finished"; +export interface ProposalPendingAttachment { + path: string; + mimeType?: string; +} + export interface ProposalPending { summary: string; question?: string; workspaceDirty?: boolean; receivedAt: number; + attachments?: ProposalPendingAttachment[]; + droppedAttachments?: string[]; } export type ProposalFinishKind = "done" | "cancelled"; @@ -325,6 +332,17 @@ function validateProposal(proposal: Proposal, id: string): void { || typeof pending.receivedAt !== "number") { invalid("pending"); } + if (pending.attachments !== undefined + && (!Array.isArray(pending.attachments) || pending.attachments.length > 3 + || pending.attachments.some((attachment) => !isRecord(attachment) + || typeof attachment.path !== "string" || attachment.path.length === 0 + || (attachment.mimeType !== undefined && typeof attachment.mimeType !== "string")))) { + invalid("pending.attachments"); + } + if (pending.droppedAttachments !== undefined + && (!Array.isArray(pending.droppedAttachments) || pending.droppedAttachments.some((entry) => typeof entry !== "string"))) { + invalid("pending.droppedAttachments"); + } } if (proposal.finishKind !== undefined && !FINISH_KINDS.includes(proposal.finishKind)) invalid("finishKind"); if (proposal.finishNote !== undefined && typeof proposal.finishNote !== "string") invalid("finishNote"); diff --git a/src/core/types.ts b/src/core/types.ts index 4c2579d..3ae07ce 100644 --- a/src/core/types.ts +++ b/src/core/types.ts @@ -28,11 +28,25 @@ export interface IncomingMessage { raw?: unknown; } +export interface OutgoingImage { + mimeType: string; + data: string; // base64 + filename?: string; +} + +// A reference to an image file on disk (inside bot.workspace); resolved to OutgoingImage at send time. +export interface OutgoingImageRef { + path: string; + mimeType?: string; + filename?: string; +} + export interface OutgoingMessage { target: MessageTarget; text: string; replyTo?: string; replySequence?: number; + images?: OutgoingImage[]; } export interface WebhookRequestContext { diff --git a/src/core/workspace-images.ts b/src/core/workspace-images.ts new file mode 100644 index 0000000..f84bb73 --- /dev/null +++ b/src/core/workspace-images.ts @@ -0,0 +1,66 @@ +import fs from "node:fs"; +import path from "node:path"; +import type { OutgoingImage, OutgoingImageRef } from "./types.js"; + +// Outbound images (Worker screenshots etc.) must live inside bot.workspace and be png or jpg. +export const MAX_OUTBOUND_IMAGE_BYTES = 10 * 1024 * 1024; +export const MAX_OUTBOUND_IMAGES = 3; + +const PNG_MAGIC = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]); +const JPG_MAGIC = Buffer.from([0xff, 0xd8, 0xff]); + +export interface ValidatedWorkspaceImage { + path: string; // resolved real path inside the workspace + mimeType: "image/png" | "image/jpeg"; + filename: string; +} + +// Resolves a worker/assistant-reported image path against the workspace and validates +// containment, size, and magic bytes. Returns undefined (never throws) when invalid. +export async function validateWorkspaceImage(workspace: string, refPath: string): Promise { + if (typeof refPath !== "string" || refPath.length === 0) return undefined; + try { + const root = fs.realpathSync.native(workspace); + const candidate = path.resolve(root, refPath); + const resolved = await fs.promises.realpath(candidate); + const relative = path.relative(root, resolved); + if (relative === "" || relative.startsWith("..") || path.isAbsolute(relative)) return undefined; + const stat = await fs.promises.lstat(resolved); + if (!stat.isFile() || stat.isSymbolicLink() || stat.size === 0 || stat.size > MAX_OUTBOUND_IMAGE_BYTES) return undefined; + const handle = await fs.promises.open(resolved, "r"); + let header: Buffer; + try { + header = Buffer.alloc(PNG_MAGIC.length); + const { bytesRead } = await handle.read(header, 0, PNG_MAGIC.length, 0); + header = header.subarray(0, bytesRead); + } finally { + await handle.close(); + } + const mimeType = magicMimeType(header); + if (!mimeType) return undefined; + return { path: resolved, mimeType, filename: path.basename(resolved) }; + } catch { + return undefined; + } +} + +// Re-reads a previously validated image at send time (the worker may have updated it); +// re-checks type and size and returns undefined when the file is gone or no longer valid. +export async function readOutboundImage(ref: OutgoingImageRef): Promise { + try { + const stat = await fs.promises.lstat(ref.path); + if (!stat.isFile() || stat.isSymbolicLink() || stat.size === 0 || stat.size > MAX_OUTBOUND_IMAGE_BYTES) return undefined; + const data = await fs.promises.readFile(ref.path); + const mimeType = magicMimeType(data.subarray(0, PNG_MAGIC.length)); + if (!mimeType) return undefined; + return { mimeType, data: data.toString("base64"), filename: ref.filename || path.basename(ref.path) }; + } catch { + return undefined; + } +} + +function magicMimeType(header: Buffer): "image/png" | "image/jpeg" | undefined { + if (header.length >= PNG_MAGIC.length && header.subarray(0, PNG_MAGIC.length).equals(PNG_MAGIC)) return "image/png"; + if (header.length >= JPG_MAGIC.length && header.subarray(0, JPG_MAGIC.length).equals(JPG_MAGIC)) return "image/jpeg"; + return undefined; +} diff --git a/src/platforms/qq/adapter.ts b/src/platforms/qq/adapter.ts index c9471ec..9821178 100644 --- a/src/platforms/qq/adapter.ts +++ b/src/platforms/qq/adapter.ts @@ -20,6 +20,7 @@ const ATTACHMENT_DOWNLOAD_TIMEOUT_MS = 10_000; export class QqAdapter implements PlatformAdapter { readonly name = "qq"; + readonly supportsImages = true; private accessToken?: { token: string; expiresAt: number }; constructor( @@ -68,7 +69,15 @@ export class QqAdapter implements PlatformAdapter { const userOpenId = stringField(raw.user_openid) || stringField(author.user_openid); const isGroup = Boolean(groupOpenId) || message.target.chatId.startsWith("group:"); const targetId = groupOpenId || userOpenId || message.target.chatId.replace(/^group:/, "").replace(/^user:/, ""); - const path = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}/messages` : `/v2/users/${encodeURIComponent(targetId)}/messages`; + const base = isGroup ? `/v2/groups/${encodeURIComponent(targetId)}` : `/v2/users/${encodeURIComponent(targetId)}`; + + // Outbound images: upload via /files (group uploads can only be sent to groups, user + // uploads only to C2C — the base path already matches the target), then send msg_type 7. + for (const image of message.images || []) { + await this.sendImageMessage(token, base, isGroup, image, message); + } + if (message.images?.length && !message.text) return; + 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. @@ -77,7 +86,7 @@ export class QqAdapter implements PlatformAdapter { body.msg_id = message.replyTo; body.msg_seq = message.replySequence ?? 1; } - const response = await fetch(`https://api.sgroup.qq.com${path}`, { + const response = await fetch(`https://api.sgroup.qq.com${base}/messages`, { method: "POST", headers: { Authorization: `QQBot ${token}`, @@ -91,6 +100,37 @@ export class QqAdapter implements PlatformAdapter { if (data?.code && data.code !== 0) throw new Error(`QQ send failed:${qqErrorDetail(data)}`); } + private async sendImageMessage(token: string, base: string, isGroup: boolean, image: { mimeType: string; data: string; filename?: string }, message: OutgoingMessage): Promise { + console.log(`QQ send ${isGroup ? "group" : "user"} image type=${image.mimeType} bytes=${Math.floor(image.data.length * 3 / 4)}`); + const headers = { + Authorization: `QQBot ${token}`, + "Content-Type": "application/json; charset=utf-8" + }; + const upload = await fetch(`https://api.sgroup.qq.com${base}/files`, { + method: "POST", + headers, + body: JSON.stringify({ file_type: 1, file_data: image.data, srv_send_msg: false }) + }); + const uploaded = await upload.json().catch(() => undefined) as { file_info?: string; code?: number; message?: string } | undefined; + if (!upload.ok) throw new Error(`QQ image upload failed: HTTP ${upload.status}${qqErrorDetail(uploaded)}`); + if (uploaded?.code && uploaded.code !== 0) throw new Error(`QQ image upload failed:${qqErrorDetail(uploaded)}`); + if (!uploaded?.file_info) throw new Error("QQ image upload failed: missing file_info"); + + const body: Record = { msg_type: 7, media: { file_info: uploaded.file_info }, content: "" }; + if (message.replyTo) { + body.msg_id = message.replyTo; + body.msg_seq = message.replySequence ?? 1; + } + const response = await fetch(`https://api.sgroup.qq.com${base}/messages`, { + method: "POST", + headers, + body: JSON.stringify(body) + }); + const data = await response.json().catch(() => undefined) as QqSendMessageResponse | undefined; + if (!response.ok) throw new Error(`QQ image send failed: HTTP ${response.status}${qqErrorDetail(data)}`); + if (data?.code && data.code !== 0) throw new Error(`QQ image send failed:${qqErrorDetail(data)}`); + } + private handleValidation(data: QqWebhookEventData | undefined): WebhookResponse { const secret = this.callbackSecret(); if (!secret) return jsonResponse({ ok: false, error: "QQ botSecret is required for callback validation" }, 400); diff --git a/src/roles/role-registry.ts b/src/roles/role-registry.ts index 5af9a16..627fe08 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 = 8; +const BOOTSTRAP_SCHEMA_VERSION = 9; export interface ResolvedBot extends BotConfig { loadedSkills: LoadedSkill[]; @@ -46,8 +46,9 @@ function buildAssistantBootstrap(bot: BotConfig): string { "You are the user-facing Assistant. You have no tools and cannot inspect files, run commands, call Skills or MCP. You chat with the user, propose work, and explain worker feedback. Never claim to have executed anything yourself.", "Speaking style: talk like a reliable colleague, not a console. Lead with the conclusion, then the reason, then the next step. Avoid protocol jargon and field names; never expose internal words like scheduler, pending, envelope, or action types to the user. Default to 2-4 sentences. Do not repeat proposal IDs unless the user asks. If you are unsure, say so plainly. When something is blocked, always give the user an actionable next step.", "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\"}; {\"type\":\"adjust_proposal\",\"id\":\"optional\",\"title\":\"optional\",\"goal\":\"optional\",\"steps\":\"optional\"}; {\"type\":\"follow_up\",\"id\":\"optional\",\"instruction\":\"...\"}; {\"type\":\"finish\",\"id\":\"optional\",\"note\":\"optional\"}; {\"type\":\"start_next\"}; {\"type\":\"cancel\",\"id\":\"optional\"}; {\"type\":\"stop\"}. Use an empty actions array when no state change is needed.", + "Supported actions: {\"type\":\"create_proposal\",\"title\":\"...\",\"goal\":\"...\",\"steps\":[\"...\"]}; {\"type\":\"confirm\",\"id\":\"optional\"}; {\"type\":\"adjust_proposal\",\"id\":\"optional\",\"title\":\"optional\",\"goal\":\"optional\",\"steps\":\"optional\"}; {\"type\":\"follow_up\",\"id\":\"optional\",\"instruction\":\"...\"}; {\"type\":\"finish\",\"id\":\"optional\",\"note\":\"optional\"}; {\"type\":\"send_image\",\"path\":\"...\"}; {\"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. When a worker finishes a turn, the proposal becomes pending: it waits for the user's decision with a summary (and maybe a question). The worker never reports success or failure; treat every result as information for the user. \"finish\" closes a pending proposal as done; \"follow_up\" sends the user's new instruction to the same proposal and resumes its worker; \"cancel\" drops a proposed, queued, or pending proposal; \"stop\" aborts the running worker and leaves the proposal pending; \"start_next\" starts the oldest confirmed queued proposal. \"adjust_proposal\" edits a proposal that has not started yet. Never invent other actions or statuses.", + "When a worker's summary says it saved image files inside the workspace (for example under .gori-outbox/) and the user asks to see one, emit \"send_image\" with that exact path. Only send paths a worker actually reported; never invent paths.", "Whenever the [Pending proposals] section in a prompt lists one of the user's proposals, your reply MUST acknowledge it: remind the user what is waiting and that they can say finish to close it or just keep talking to continue it.", "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], [Pending proposals], [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, adjust, follow up, finish, stop, or cancel it. If another group member asks you to operate a proposal they did not create, explain that only the proposal's initiator can do that and emit no action.", @@ -63,7 +64,7 @@ function buildWorkerBootstrap(bot: BotConfig, skills: LoadedSkill[]): string { bot.persona ? `Persona:\n${bot.persona}` : "Persona: general coding assistant", `Permission policy enforced by the ACP client: ${JSON.stringify(bot.permissions)}`, "You are the Worker. You execute exactly one confirmed proposal in the configured workspace and never talk to the user directly.", - "For every turn, end your response with exactly one hidden result envelope: {\"status\":\"PENDING\",\"summary\":\"...\"}. PENDING is the only status: it hands the result back to the user. Always include a short user-readable \"summary\" of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing. Set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. Never emit any other status or text after the envelope.", + "For every turn, end your response with exactly one hidden result envelope: {\"status\":\"PENDING\",\"summary\":\"...\"}. PENDING is the only status: it hands the result back to the user. Always include a short user-readable \"summary\" of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing. Set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 as \"attachments\": [{\"path\":\"relative/or/absolute/path\",\"mimeType\":\"optional\"}]; never report paths outside the workspace such as /tmp. Never emit any other status or text after the envelope.", "Keep every command and tool process attached to this ACP worker. Never daemonize, call setsid, use nohup, create a detached process, or leave a background process running after the turn.", ...skills.map((skill) => `Skill ${skill.id} (${skill.file}):\n${skill.content}`) ].join("\n\n"); diff --git a/src/server.ts b/src/server.ts index 6edd5a5..20be285 100644 --- a/src/server.ts +++ b/src/server.ts @@ -47,7 +47,7 @@ export async function createGatewayRuntime(config: AppConfig): Promise gateway.sendEvent(chatKey, text)); + manager.setEventSink((chatKey, text, images) => gateway.sendEvent(chatKey, text, images)); const platformAdapter = createPlatformAdapter(config, gateway); const qqGatewayClient = config.gateway.platform.type === "qq" && config.gateway.platform.connectionMode === "websocket" ? new QqGatewayClient(config.gateway.platform, platformAdapter as QqAdapter) : undefined; diff --git a/test/assistant-manager.test.ts b/test/assistant-manager.test.ts index 54bc3d8..89d39da 100644 --- a/test/assistant-manager.test.ts +++ b/test/assistant-manager.test.ts @@ -22,7 +22,7 @@ interface Harness { store: DurableSessionStore; proposals: ProposalStore; manager: AssistantManager; - events: { chatKey: string; text: string }[]; + events: { chatKey: string; text: string; images?: { path: string; mimeType?: string; filename?: string }[] }[]; logFile: string; } @@ -41,12 +41,12 @@ async function createHarness(options: { initialize?: boolean; agentEnv?: Record< const proposals = new ProposalStore(path.join(home, "state", "proposals.json"), identity); await store.open(); await proposals.open(); - const events: { chatKey: string; text: string }[] = []; + const events: { chatKey: string; text: string; images?: { path: string; mimeType?: string; filename?: string }[] }[] = []; const manager = new AssistantManager(config.runtime.acp, new BotProfileResolver(config).bot, store, proposals, { assistantWorkspaceHome: home, allowUnverifiedAssistantAgent: true }); - manager.setEventSink(async (chatKey, text) => { events.push({ chatKey, text }); }); + manager.setEventSink(async (chatKey, text, images) => { events.push({ chatKey, text, images }); }); if (options.initialize !== false) await manager.initialize(); return { home, workspace, store, proposals, manager, events, logFile }; } @@ -70,7 +70,7 @@ async function reopenManager(harness: Harness): Promise { assistantWorkspaceHome: harness.home, allowUnverifiedAssistantAgent: true }); - harness.manager.setEventSink(async (chatKey, text) => { harness.events.push({ chatKey, text }); }); + harness.manager.setEventSink(async (chatKey, text, images) => { harness.events.push({ chatKey, text, images }); }); await harness.manager.initialize(); } @@ -485,10 +485,21 @@ test("envelope parsers accept valid tails and reject invalid ones", () => { assert.equal(parseAssistantActions(`x\n{"reply":"hi","actions":[]}`), undefined); assert.deepEqual( parseWorkerResult(`text\n{"status":"PENDING","summary":"s","question":"q","workspaceDirty":true}`), - { status: "PENDING", summary: "s", question: "q", workspaceDirty: true } + { status: "PENDING", summary: "s", question: "q", workspaceDirty: true, attachments: undefined } ); + assert.deepEqual( + parseWorkerResult(`text\n{"status":"PENDING","summary":"s","attachments":[{"path":".gori-outbox/a.png","mimeType":"image/png"}]}`), + { status: "PENDING", summary: "s", question: undefined, workspaceDirty: undefined, attachments: [{ path: ".gori-outbox/a.png", mimeType: "image/png" }] } + ); + assert.equal(parseWorkerResult(`x\n{"status":"PENDING","summary":"s","attachments":[{"path":"a.png"},{"path":"b.png"},{"path":"c.png"},{"path":"d.png"}]}`), undefined); + assert.equal(parseWorkerResult(`x\n{"status":"PENDING","summary":"s","attachments":[{"mimeType":"image/png"}]}`), undefined); assert.equal(parseWorkerResult(`x\n{"status":"SUCCESS","summary":"s"}`), undefined); assert.equal(parseWorkerResult(`x\n{"status":"PENDING","summary":"s"}`), undefined); + assert.deepEqual( + parseAssistantActions(`text\n{"reply":"hi","actions":[{"type":"send_image","path":".gori-outbox/a.png"}]}`), + { reply: "hi", actions: [{ type: "send_image", path: ".gori-outbox/a.png" }] } + ); + assert.equal(parseAssistantActions(`x\n{"reply":"hi","actions":[{"type":"send_image"}]}`), undefined); }); test("an owner's own pending proposal blocks their start_next with an actionable correction", async () => { @@ -678,6 +689,79 @@ test("image attachments degrade to a text note when the agent has no image capab } }); +test("worker attachments are validated against the workspace and ride along with the owner event", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: attach")); + const proposal = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm")); + await waitFor(() => harness.proposals.get(proposal.id)!.status === "pending"); + + const pending = harness.proposals.get(proposal.id)!.pending!; + assert.equal(pending.attachments?.length, 1); + const stored = pending.attachments![0]!; + assert.equal(stored.mimeType, "image/png"); + const workspaceRoot = fs.realpathSync(harness.workspace); + assert.ok(stored.path.startsWith(`${workspaceRoot}${path.sep}`), stored.path); + assert.ok(stored.path.endsWith(path.join(".gori-outbox", "shot.png")), stored.path); + assert.equal(pending.droppedAttachments, undefined); + + await waitFor(() => harness.events.length > 0); + assert.equal(harness.events[0]!.images?.length, 1); + assert.equal(harness.events[0]!.images![0]!.path, stored.path); + assert.equal(harness.events[0]!.images![0]!.filename, "shot.png"); + } finally { + await closeHarness(harness); + } +}); + +test("worker attachments outside the workspace or not png/jpg are dropped without failing the result", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: attachbad")); + const proposal = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm")); + await waitFor(() => harness.proposals.get(proposal.id)!.status === "pending"); + + const pending = harness.proposals.get(proposal.id)!.pending!; + assert.equal(pending.summary, "made a pic"); + assert.equal(pending.attachments, undefined); + assert.deepEqual(pending.droppedAttachments, ["/tmp/evil.png", "not-an-image.txt"]); + + await waitFor(() => harness.events.length > 0); + assert.equal(harness.events[0]!.images, undefined); + } finally { + await closeHarness(harness); + } +}); + +test("send_image validates the path and attaches the image to the reply", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: attach")); + const proposal = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm")); + await waitFor(() => harness.proposals.get(proposal.id)!.status === "pending"); + + const sent = await harness.manager.prompt(request("send image: .gori-outbox/shot.png")); + assert.equal(sent.images?.length, 1); + assert.equal(sent.images![0]!.filename, "shot.png"); + assert.equal(sent.images![0]!.mimeType, "image/png"); + assert.ok(sent.images![0]!.path.startsWith(fs.realpathSync(harness.workspace))); + assert.doesNotMatch(sent.text, /发不出去/); + + const outside = await harness.manager.prompt(request("send image: /tmp/evil.png")); + assert.equal(outside.images, undefined); + assert.match(outside.text, /发不出去/); + + const notImage = await harness.manager.prompt(request("send image: not-an-image.txt")); + assert.equal(notImage.images, undefined); + assert.match(notImage.text, /发不出去/); + } 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/fixtures/fake-acp-agent.mjs b/test/fixtures/fake-acp-agent.mjs index 362aae8..324cdc3 100644 --- a/test/fixtures/fake-acp-agent.mjs +++ b/test/fixtures/fake-acp-agent.mjs @@ -1,9 +1,12 @@ #!/usr/bin/env node import { spawn } from "node:child_process"; import fs from "node:fs"; +import path from "node:path"; import { Readable, Writable } from "node:stream"; import * as acp from "@agentclientprotocol/sdk"; +const PNG_HEADER = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]); + const pending = new Map(); const logFile = process.env.FAKE_ACP_LOG; const log = (entry) => { if (logFile) fs.appendFileSync(logFile, `${JSON.stringify(entry)}\n`); }; @@ -84,6 +87,7 @@ const app = acp.agent({ name: "fake-acp-agent" }) else if (userText === "finish") response = assistantEnvelope("finishing", [{ type: "finish" }]); else if (userText.startsWith("follow up:")) response = assistantEnvelope("following up", [{ type: "follow_up", instruction: userText.slice("follow up:".length).trim() }]); else if (userText.startsWith("adjust proposal:")) response = assistantEnvelope("adjusting", [{ type: "adjust_proposal", title: userText.slice("adjust proposal:".length).trim() }]); + else if (userText.startsWith("send image:")) response = assistantEnvelope("sending image", [{ type: "send_image", path: userText.slice("send image:".length).trim() }]); else if (userText === "start next") response = assistantEnvelope("starting next", [{ type: "start_next" }]); else if (userText === "stop") response = assistantEnvelope("stopping", [{ type: "stop" }]); else if (userText === "cancel proposal") response = assistantEnvelope("cancelling", [{ type: "cancel" }]); @@ -117,6 +121,18 @@ const app = acp.agent({ name: "fake-acp-agent" }) if (process.env.FAKE_ACP_WORKER_GATE_FILE) { while (!fs.existsSync(process.env.FAKE_ACP_WORKER_GATE_FILE)) await new Promise((resolve) => setTimeout(resolve, 5)); } + if (goal.includes("attachbad")) { + await update(client, params.sessionId, workerEnvelope({ status: "PENDING", summary: "made a pic", attachments: [{ path: "/tmp/evil.png" }, { path: "not-an-image.txt" }] })); + return { stopReason: "end_turn" }; + } + if (goal.includes("attach")) { + const outbox = path.join(process.cwd(), ".gori-outbox"); + fs.mkdirSync(outbox, { recursive: true }); + fs.writeFileSync(path.join(outbox, "shot.png"), Buffer.concat([PNG_HEADER, Buffer.from("fake-png-payload-v1")])); + fs.writeFileSync(path.join(process.cwd(), "not-an-image.txt"), "plain text"); + await update(client, params.sessionId, workerEnvelope({ status: "PENDING", summary: "made a pic", attachments: [{ path: ".gori-outbox/shot.png", mimeType: "image/png" }] })); + return { stopReason: "end_turn" }; + } const result = goal.includes("ask") ? { status: "PENDING", summary: "hit a dirty target", question: "May I overwrite it?", workspaceDirty: true } : goal.includes("fail") diff --git a/test/gateway.test.ts b/test/gateway.test.ts index e410893..4bf04e9 100644 --- a/test/gateway.test.ts +++ b/test/gateway.test.ts @@ -1,4 +1,7 @@ import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; import test from "node:test"; import type { ConversationRequest, ConversationRuntime } from "../src/acp/types.js"; import type { PlatformAdapter } from "../src/core/adapter.js"; @@ -6,6 +9,15 @@ import { Gateway } from "../src/core/gateway.js"; import type { Proposal } from "../src/core/proposal-store.js"; import type { IncomingMessage, OutgoingMessage } from "../src/core/types.js"; +const PNG_HEADER = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]); + +async function tempPng(payload: string): Promise<{ dir: string; file: string }> { + const dir = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-gateway-img-")); + const file = path.join(dir, "shot.png"); + await fs.promises.writeFile(file, Buffer.concat([PNG_HEADER, Buffer.from(payload)])); + return { dir, file }; +} + interface PendingTurn { request: ConversationRequest; resolve(text: string): void; @@ -76,12 +88,17 @@ const groupMessage = (text: string, messageId: string): IncomingMessage => ({ }); 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(); + const deadline = Date.now() + 5_000; + while (!condition()) { + if (Date.now() > deadline) break; + await flush(); + await new Promise((resolve) => setTimeout(resolve, 5)); + } assert.equal(condition(), true); }; -function recordingAdapter(sent: OutgoingMessage[] = []): PlatformAdapter { - return { name: "test", async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } }; +function recordingAdapter(sent: OutgoingMessage[] = [], supportsImages = false): PlatformAdapter { + return { name: "test", supportsImages, async handleWebhook() { return {}; }, async sendMessage(outgoing) { sent.push(outgoing); } }; } // Simulates QQ passive-reply dedup: a repeated msg_id + msg_seq pair is rejected. @@ -286,6 +303,139 @@ test("status reports the last successful event delivery for the current chat", a assert.match(status.reply || "", /lastEventAt=\d{4}-\d{2}-\d{2}T/); }); +test("sendEvent delivers the text first, then images, all sharing the msg_seq counter", async () => { + const { runtime, gateway } = createGateway(); + const { dir, file } = await tempPng("payload-v1"); + try { + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "img-event"), recordingAdapter(sent, true)); + await waitFor(() => runtime.turns.length === 1); + runtime.turns[0].resolve("done"); + await turn; + + await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]); + const pieces = sent.map((outgoing) => [outgoing.text, outgoing.replySequence, outgoing.images?.length || 0]); + assert.deepEqual(pieces, [ + ["done", 1, 0], + ["截图好了", 2, 0], + ["", 3, 1] + ]); + const image = sent[2]!.images![0]!; + assert.equal(image.mimeType, "image/png"); + assert.equal(image.filename, "shot.png"); + assert.equal(Buffer.from(image.data, "base64").toString("latin1"), Buffer.concat([PNG_HEADER, Buffer.from("payload-v1")]).toString("latin1")); + } finally { + await fs.promises.rm(dir, { recursive: true, force: true }); + } +}); + +test("a queued event re-reads its image at redelivery time and degrades when the file is gone", async () => { + let now = 1_000_000; + const runtime = new FakeRuntime(); + const gateway = new Gateway(policy, runtime, { now: () => now }); + const { dir, file } = await tempPng("payload-v1"); + try { + const sent: OutgoingMessage[] = []; + const adapter = recordingAdapter(sent, true); + const first = gateway.receive(groupMessage("work", "m1"), adapter); + await waitFor(() => runtime.turns.length === 1); + runtime.turns[0].resolve("done"); + await first; + + now += 6 * 60_000; // past the group window: the event is queued instead of sent + await gateway.sendEvent("qq:group:g1", "stale event", [{ path: file, filename: "shot.png" }]); + assert.equal(sent.filter((outgoing) => outgoing.text === "stale event").length, 0); + + // The worker updated the image after the event was queued: redelivery must send the new content. + await fs.promises.writeFile(file, Buffer.concat([PNG_HEADER, Buffer.from("payload-v2")])); + const second = gateway.receive(groupMessage("hello", "m2"), adapter); + await waitFor(() => runtime.turns.length === 2); + runtime.turns[1].resolve("ok"); + await second; + const redelivered = sent.filter((outgoing) => outgoing.replyTo === "m2"); + assert.deepEqual(redelivered.map((outgoing) => [outgoing.text, outgoing.replySequence, outgoing.images?.length || 0]), [ + ["stale event", 1, 0], + ["", 2, 1], + ["ok", 3, 0] + ]); + assert.equal(Buffer.from(redelivered[1]!.images![0]!.data, "base64").toString("latin1"), Buffer.concat([PNG_HEADER, Buffer.from("payload-v2")]).toString("latin1")); + + // Queue another event, then delete the file: the redelivery degrades to a text note. + now += 6 * 60_000; + await gateway.sendEvent("qq:group:g1", "another stale event", [{ path: file, filename: "shot.png" }]); + await fs.promises.rm(file); + const third = gateway.receive(groupMessage("hello again", "m3"), adapter); + await waitFor(() => runtime.turns.length === 3); + runtime.turns[2].resolve("ok again"); + await third; + const degraded = sent.filter((outgoing) => outgoing.replyTo === "m3"); + assert.equal(degraded.filter((outgoing) => outgoing.images?.length).length, 0); + assert.match(degraded[0]!.text, /another stale event/); + assert.match(degraded[0]!.text, /图片 shot\.png 发送失败/); + } finally { + await fs.promises.rm(dir, { recursive: true, force: true }); + } +}); + +test("adapters without image support receive images flattened into text lines", async () => { + const { runtime, gateway } = createGateway(); + const { dir, file } = await tempPng("payload-v1"); + try { + const sent: OutgoingMessage[] = []; + const turn = gateway.receive(message("work", "flat"), recordingAdapter(sent)); + await waitFor(() => runtime.turns.length === 1); + runtime.turns[0].resolve("done"); + await turn; + + await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]); + const event = sent.find((outgoing) => outgoing.text.includes("截图好了"))!; + assert.equal(event.images, undefined); + assert.match(event.text, /截图好了/); + assert.match(event.text, /\[图片\] shot\.png/); + assert.equal(sent.filter((outgoing) => outgoing.images?.length).length, 0); + } finally { + await fs.promises.rm(dir, { recursive: true, force: true }); + } +}); + +test("an image send failure degrades to a text note without blocking the event", async () => { + const { runtime, gateway } = createGateway(); + const { dir, file } = await tempPng("payload-v1"); + try { + const sent: OutgoingMessage[] = []; + const errors: string[] = []; + const adapter: PlatformAdapter = { + name: "test", supportsImages: true, async handleWebhook() { return {}; }, + async sendMessage(outgoing) { + if (outgoing.images?.length) throw new Error("sensitive provider response"); + sent.push(outgoing); + } + }; + const originalError = console.error; + console.error = (...args: unknown[]) => { errors.push(args.map(String).join(" ")); }; + try { + const turn = gateway.receive(message("work", "img-fail"), adapter); + await waitFor(() => runtime.turns.length === 1); + runtime.turns[0].resolve("done"); + await turn; + + await gateway.sendEvent("qq:chat", "截图好了", [{ path: file, filename: "shot.png" }]); + } finally { + console.error = originalError; + } + const texts = sent.map((outgoing) => [outgoing.text, outgoing.replySequence]); + // The failed image attempt consumed msg_seq 3 (never reused in case QQ actually received it). + assert.deepEqual(texts, [ + ["done", 1], + ["截图好了", 2], + ["(图片 shot.png 发送失败)", 4] + ]); + assert.deepEqual(errors, ["Gateway event image send failed (platform=qq)"]); + } finally { + await fs.promises.rm(dir, { recursive: true, force: true }); + } +}); + test("normal replies and fresh events share one msg_seq counter per inbound message", async () => { const runtime = new FakeRuntime(); const gateway = new Gateway(policy, runtime); diff --git a/test/qq-adapter.test.ts b/test/qq-adapter.test.ts index cb3edd0..469256a 100644 --- a/test/qq-adapter.test.ts +++ b/test/qq-adapter.test.ts @@ -171,6 +171,77 @@ test("sendMessage uses C2C endpoint and forwards reply sequences as msg_seq", as } finally { globalThis.fetch = original; } }); +test("sendMessage uploads images via /files and sends msg_type 7 media with the shared msg_seq", async () => { + const original = globalThis.fetch; const urls: string[] = []; const bodies: Array> = []; + globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { + const url = String(input); + urls.push(url); + if (url.includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + bodies.push(JSON.parse(String(init?.body)) as Record); + if (url.endsWith("/files")) return new Response(JSON.stringify({ file_info: "FILEINFO", ttl: 60 }), { status: 200 }); + 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" } }; + const image = { mimeType: "image/png", data: Buffer.from("fake-png").toString("base64"), filename: "shot.png" }; + await adapter.sendMessage({ target, text: "", images: [image], replyTo: "m1", replySequence: 3 }); + const apiCalls = urls.filter((url) => !url.includes("getAppAccessToken")).map((url) => url.replace("https://api.sgroup.qq.com", "")); + assert.deepEqual(apiCalls, ["/v2/groups/g1/files", "/v2/groups/g1/messages"]); + assert.deepEqual(bodies, [ + { file_type: 1, file_data: image.data, srv_send_msg: false }, + { msg_type: 7, media: { file_info: "FILEINFO" }, content: "", msg_id: "m1", msg_seq: 3 } + ]); + } finally { globalThis.fetch = original; } +}); + +test("sendMessage delivers text and images to C2C with independent upload per target", async () => { + const original = globalThis.fetch; const urls: string[] = []; const bodies: Array> = []; + globalThis.fetch = (async (input: string | URL | Request, init?: RequestInit) => { + const url = String(input); + urls.push(url); + if (url.includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + bodies.push(JSON.parse(String(init?.body)) as Record); + if (url.endsWith("/files")) return new Response(JSON.stringify({ file_info: "FILEINFO" }), { status: 200 }); + 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: "user:u2", raw: { author: { user_openid: "u2" } } }; + const image = { mimeType: "image/jpeg", data: Buffer.from("fake-jpg").toString("base64") }; + await adapter.sendMessage({ target, text: "看图", images: [image], replyTo: "m2", replySequence: 5 }); + assert.ok(urls.some((url) => url.endsWith("/v2/users/u2/files"))); + assert.equal(urls.filter((url) => url.endsWith("/v2/users/u2/messages")).length, 2); + assert.deepEqual(bodies[0], { file_type: 1, file_data: image.data, srv_send_msg: false }); + assert.deepEqual(bodies[1], { msg_type: 7, media: { file_info: "FILEINFO" }, content: "", msg_id: "m2", msg_seq: 5 }); + assert.deepEqual(bodies[2], { content: "看图", msg_id: "m2", msg_seq: 5 }); + } finally { globalThis.fetch = original; } +}); + +test("image upload failures surface safe errors without the base64 payload", async () => { + const original = globalThis.fetch; + globalThis.fetch = (async (input: string | URL | Request) => { + const url = String(input); + if (url.includes("getAppAccessToken")) return new Response(JSON.stringify({ access_token: "token", expires_in: 7200 }), { status: 200 }); + return new Response(JSON.stringify({ code: 40034001, message: "invalid file_data" }), { 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" } }; + const image = { mimeType: "image/png", data: Buffer.from("secret-image-bytes").toString("base64") }; + await assert.rejects( + adapter.sendMessage({ target, text: "", images: [image], replyTo: "m1", replySequence: 1 }), + (error: Error) => { + assert.match(error.message, /QQ image upload failed: HTTP 400/); + assert.match(error.message, /code=40034001/); + assert.doesNotMatch(error.message, /secret-image-bytes/); + assert.doesNotMatch(error.message, /base64/); + return true; + } + ); + } 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) => {