diff --git a/interfaces/http/weixin_ilink_api.py b/interfaces/http/weixin_ilink_api.py index 7b037637..9d427657 100644 --- a/interfaces/http/weixin_ilink_api.py +++ b/interfaces/http/weixin_ilink_api.py @@ -222,8 +222,10 @@ class _IlinkBridge: chat_id: str, text: str, context_token: str | None, + attachments: list[dict[str, Any]] | None = None, + media_path: str | None = None, ) -> None: - payload = { + payload: dict[str, Any] = { "channel": channel, "account_id": account_id, "chat_id": chat_id, @@ -234,6 +236,12 @@ class _IlinkBridge: "context_token": (context_token or "").strip(), "ts": _now_ms(), } + mp = str(media_path or "").strip() + if mp: + payload["media_path"] = mp + atts = [a for a in (attachments or []) if isinstance(a, dict)] + if atts: + payload["attachments"] = atts with self._lock: self._seq += 1 event = { @@ -381,7 +389,9 @@ async def _process_inbound_and_enqueue( ctx_token = _extract_context_token(payload) reply_text = str(item.get("text") or "").strip() reply_chat_id = str(item.get("chat_id") or chat_id or user_id).strip() or user_id - if not reply_text: + reply_attachments = item.get("attachments") if isinstance(item.get("attachments"), list) else [] + reply_media_path = str(item.get("media_path") or item.get("mediaPath") or "").strip() + if not reply_text and not reply_attachments and not reply_media_path: continue _BRIDGE.enqueue_reply( token=token, @@ -390,6 +400,8 @@ async def _process_inbound_and_enqueue( chat_id=reply_chat_id, text=reply_text, context_token=ctx_token, + attachments=[a for a in reply_attachments if isinstance(a, dict)], + media_path=reply_media_path or None, ) @@ -570,6 +582,8 @@ def enqueue_weixin_outbound_reply( text: str, context_token: str = "", token: str = "", + attachments: list[dict[str, Any]] | None = None, + media_path: str | None = None, ) -> str: """Enqueue a proactive Weixin/WeChat outbound message for the sidecar poll loop.""" tok = str( @@ -585,6 +599,8 @@ def enqueue_weixin_outbound_reply( chat_id=str(chat_id or "").strip(), text=str(text or ""), context_token=str(context_token or "").strip(), + attachments=attachments, + media_path=media_path, ) return str(_BRIDGE._seq) diff --git a/runtime/application/gateway/inbound_service.py b/runtime/application/gateway/inbound_service.py index 6fda6088..f46a9f62 100644 --- a/runtime/application/gateway/inbound_service.py +++ b/runtime/application/gateway/inbound_service.py @@ -566,28 +566,18 @@ def _rows_since_last_user_message(rows: list[Any]) -> list[Any]: return list(rows[last_user_idx + 1 :]) -_CHANNEL_DELIVERABLE_ATTACHMENT_TYPES = frozenset( - {"image_ref", "video_ref", "image", "input_image", "image_url"} -) -_CHANNEL_EXPLICIT_DELIVERABLE_TYPES = frozenset({"binary_ref", "text_ref"}) - - def _is_channel_deliverable_attachment(att: dict[str, Any]) -> bool: + """Channel outbound only sends attachments explicitly marked deliverable=true.""" if not isinstance(att, dict): return False - t = str(att.get("type") or "").strip().lower() - if t in _CHANNEL_DELIVERABLE_ATTACHMENT_TYPES: - return True - if t in _CHANNEL_EXPLICIT_DELIVERABLE_TYPES: - return att.get("deliverable") is True - return False + return att.get("deliverable") is True def _collect_recent_tool_attachments(*, store: Any, session_id: str) -> list[dict[str, Any]]: """Fallback for channel delivery: reuse tool media produced during the current user turn only. - Avoids re-sending images from earlier conversation turns when the latest assistant row has no attachments. - Visual media is always eligible; documents need explicit deliverable=true on the attachment ref. + Avoids re-sending attachments from earlier conversation turns when the latest assistant row has none. + All attachment types (image, video, document) require explicit deliverable=true. """ sid = str(session_id or "").strip() if not sid: @@ -1083,11 +1073,6 @@ def process_inbound_payload(payload: dict[str, Any]) -> dict[str, Any]: if dispatch_lang in {"zh", "en"} else resolve_runtime_lang(store=store, user_text=user_text) ) - if str(inbound.channel or "").strip().lower() in {"whatsapp", "wechat", "weixin"}: - from runtime.orchestration.group_ingest import build_channel_file_delivery_instruction - - ch_hint = build_channel_file_delivery_instruction(lang=lang) - user_text = f"{ch_hint}\n{user_text}" if user_text else ch_hint gw = OclawGateway(store=store) msg = StandardMessage( session_id=str(session_id), diff --git a/runtime/gateway.py b/runtime/gateway.py index 10177a35..d31ebefd 100644 --- a/runtime/gateway.py +++ b/runtime/gateway.py @@ -467,6 +467,17 @@ class OclawGateway: return True return False + @staticmethod + def _is_channel_delivery_channel(msg: StandardMessage) -> bool: + ch = str(getattr(msg, "channel", "") or "").strip().lower() + return ch in {"whatsapp", "wechat", "weixin"} + + @staticmethod + def _channel_file_delivery_system_hint(lang: str) -> str: + from runtime.orchestration.group_ingest import build_channel_file_delivery_instruction + + return str(build_channel_file_delivery_instruction(lang=lang) or "").strip() + @staticmethod def _tabular_query_system_hint(lang: str) -> str: limits = OclawGateway._tabular_limits_from_config() @@ -1159,6 +1170,10 @@ class OclawGateway: if model is None or tools is None: raise RuntimeError("executor missing model/tools") sys_prompt = system_prompt_override or str(getattr(selected_executor, "system_prompt", "") or "") + if self._is_channel_delivery_channel(msg): + ch_hint = self._channel_file_delivery_system_hint(lang) + if ch_hint: + sys_prompt = f"{sys_prompt}\n\n{ch_hint}".strip() if self._has_tabular_ref_attachments(msg): sys_prompt = f"{sys_prompt}\n\n{self._tabular_query_system_hint(lang)}".strip() if self._has_text_ref_attachments(msg): diff --git a/runtime/operations/weixin_bridge/official_runner.ts b/runtime/operations/weixin_bridge/official_runner.ts index e03710c1..089cb1b5 100644 --- a/runtime/operations/weixin_bridge/official_runner.ts +++ b/runtime/operations/weixin_bridge/official_runner.ts @@ -235,7 +235,10 @@ async function flushWeixinDbOutbound(args: { const toUser = String(item.chat_id || "").trim(); const text = String(item.text || "").trim(); const msgId = String(item.id || "").trim(); - if (!toUser || !text || !msgId) continue; + const mediaPath = String(item.media_path || item.mediaPath || "").trim(); + const attachments = Array.isArray(item.attachments) ? (item.attachments as Json[]) : []; + if (!toUser || !msgId) continue; + if (!text && !mediaPath && !attachments.length) continue; const contextToken = String( (item.context_token as string) || args.modules.getContextToken(args.accountId, toUser) || @@ -247,14 +250,20 @@ async function flushWeixinDbOutbound(args: { continue; } try { - await args.modules.sendMessageWeixin({ - to: toUser, - text, - opts: { - baseUrl: args.cloudBaseUrl, - token: args.token, - contextToken, + await deliverWeixinReply(args.modules, { + reply: { + chat_id: toUser, + text, + media_path: mediaPath, + attachments, + context_token: contextToken, }, + defaultTo: toUser, + cloudBaseUrl: args.cloudBaseUrl, + token: args.token, + accountId: args.accountId, + defaultContextToken: contextToken, + userContextTokens: args.userContextTokens, }); await postLocalJson(args.token, "weixin/outbound/ack", { id: msgId, ok: true }, 8000); log(`db proactive reply sent: id=${msgId} to=${toUser} textLen=${text.length}`); @@ -307,7 +316,10 @@ async function flushLocalProactiveReplies(args: { for (const r of msgs) { const toUser = String(r.chat_id || "").trim(); const text = String(r.text || "").trim(); - if (!toUser || !text) continue; + const mediaPath = String(r.media_path || r.mediaPath || "").trim(); + const attachments = Array.isArray(r.attachments) ? (r.attachments as Json[]) : []; + if (!toUser) continue; + if (!text && !mediaPath && !attachments.length) continue; const contextToken = String( (r.context_token as string) || args.modules.getContextToken(args.accountId, toUser) || @@ -320,14 +332,20 @@ async function flushLocalProactiveReplies(args: { continue; } try { - await args.modules.sendMessageWeixin({ - to: toUser, - text, - opts: { - baseUrl: args.cloudBaseUrl, - token: args.token, - contextToken, + await deliverWeixinReply(args.modules, { + reply: { + chat_id: toUser, + text, + media_path: mediaPath, + attachments, + context_token: contextToken, }, + defaultTo: toUser, + cloudBaseUrl: args.cloudBaseUrl, + token: args.token, + accountId: args.accountId, + defaultContextToken: contextToken, + userContextTokens: args.userContextTokens, }); log(`proactive reply sent: to=${toUser} textLen=${text.length}`); } catch (err) { @@ -450,16 +468,45 @@ function pickDownloadableMedia(msg: Json): Json | null { return null; } -function buildAttachmentsFromMedia(mediaOpts: Json): Json[] { +function attachmentNameFromMediaItem(item: Json | null, kind: string): string { + if (!item || typeof item !== "object") return ""; + if (kind === "file") { + const fileItem = item.file_item; + if (fileItem && typeof fileItem === "object") { + return String((fileItem as Json).file_name || (fileItem as Json).filename || (fileItem as Json).name || "").trim(); + } + } + if (kind === "image") { + const imageItem = item.image_item; + if (imageItem && typeof imageItem === "object") { + return String((imageItem as Json).file_name || (imageItem as Json).filename || "").trim(); + } + } + return ""; +} + +function buildAttachmentsFromMedia(mediaOpts: Json, mediaItem?: Json | null): Json[] { const out: Json[] = []; const pushIf = (key: string, mimeKey: string, kind: string): void => { const filePath = String(mediaOpts[key] || "").trim(); if (!filePath) return; - out.push({ + const name = String( + mediaOpts.fileOriginalName || + mediaOpts.originalFilename || + mediaOpts.file_name || + attachmentNameFromMediaItem(mediaItem || null, kind) || + "", + ).trim(); + const att: Json = { kind, local_path: filePath, media_type: String(mediaOpts[mimeKey] || "").trim(), - }); + }; + if (name) { + att.name = name; + att.filename = name; + } + out.push(att); }; pushIf("decryptedPicPath", "", "image"); pushIf("decryptedVideoPath", "", "video"); @@ -468,6 +515,70 @@ function buildAttachmentsFromMedia(mediaOpts: Json): Json[] { return out; } +async function deliverWeixinReply( + modules: OfficialModules, + params: { + reply: Json; + defaultTo: string; + cloudBaseUrl: string; + token: string; + accountId: string; + defaultContextToken?: string; + userContextTokens?: TokenMap; + }, +): Promise { + const reply = params.reply; + const text = String(reply.text || "").trim(); + const mediaPath = String(reply.media_path || reply.mediaPath || "").trim(); + const mediaUrl = String(reply.media_url || reply.mediaUrl || "").trim(); + const deliverTo = String(reply.chat_id || reply.to || params.defaultTo).trim() || params.defaultTo; + const replyContextToken = String( + reply.context_token || + modules.getContextToken(params.accountId, deliverTo) || + params.defaultContextToken || + params.userContextTokens?.[deliverTo] || + "", + ).trim(); + const replyAtts = Array.isArray(reply.attachments) ? (reply.attachments as Json[]) : []; + let inlineMediaPath = ""; + for (const att of replyAtts) { + if (!att || typeof att !== "object") continue; + const decoded = _decodeReplyBase64Attachment(att); + if (decoded?.filePath) { + inlineMediaPath = decoded.filePath; + break; + } + } + if (mediaPath || mediaUrl || inlineMediaPath) { + const filePath = mediaPath || mediaUrl || inlineMediaPath; + await modules.sendWeixinMediaFile({ + filePath, + to: deliverTo, + text, + opts: { + baseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: replyContextToken || undefined, + }, + cdnBaseUrl: DEFAULT_CDN_BASE_URL, + }); + log(`reply media sent to=${deliverTo} textLen=${text.length}`); + return true; + } + if (!text) return false; + await modules.sendMessageWeixin({ + to: deliverTo, + text, + opts: { + baseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: replyContextToken || undefined, + }, + }); + log(`reply text sent to=${deliverTo} textLen=${text.length}`); + return true; +} + function _decodeReplyBase64Attachment(att: Json): { filePath: string; mime: string } | null { const b64Raw = String( att.data_base64 || att.media_base64 || att.image_base64 || att.video_base64 || att.audio_base64 || att.data || "", @@ -551,7 +662,8 @@ async function handleInboundMessage( } const ctx = modules.weixinMessageToMsgContext(params.full, params.accountId, mediaOpts); const bodyText = String((ctx as Json).Body || "").trim(); - const attachCount = buildAttachmentsFromMedia(mediaOpts).length; + const inboundAttachments = buildAttachmentsFromMedia(mediaOpts, mediaItem); + const attachCount = inboundAttachments.length; log(`native call bodyLen=${bodyText.length} attachments=${attachCount}`); let replies: Json[] = []; const tNative0 = Date.now(); @@ -560,7 +672,7 @@ async function handleInboundMessage( channel: "wechat", account_id: params.accountId, ctx, - attachments: buildAttachmentsFromMedia(mediaOpts), + attachments: inboundAttachments, metadata: { source: "weixin_official_native", raw: { @@ -586,50 +698,15 @@ async function handleInboundMessage( return params.localCursor; } for (const reply of replies) { - const text = String(reply.text || "").trim(); - const mediaPath = String(reply.media_path || reply.mediaPath || "").trim(); - const mediaUrl = String(reply.media_url || reply.mediaUrl || "").trim(); - const deliverTo = String(reply.chat_id || reply.to || fromUser).trim() || fromUser; - const replyContextToken = String( - reply.context_token || modules.getContextToken(params.accountId, deliverTo) || contextToken || "", - ).trim(); - const replyAtts = Array.isArray(reply.attachments) ? (reply.attachments as Json[]) : []; - let inlineMediaPath = ""; - for (const att of replyAtts) { - if (!att || typeof att !== "object") continue; - const decoded = _decodeReplyBase64Attachment(att); - if (decoded?.filePath) { - inlineMediaPath = decoded.filePath; - break; - } - } - if (mediaPath || mediaUrl || inlineMediaPath) { - const filePath = mediaPath || mediaUrl || inlineMediaPath; - await modules.sendWeixinMediaFile({ - filePath, - to: deliverTo, - text, - opts: { - baseUrl: params.cloudBaseUrl, - token: params.token, - contextToken: replyContextToken || undefined, - }, - cdnBaseUrl: DEFAULT_CDN_BASE_URL, - }); - log(`reply media sent to=${deliverTo} textLen=${text.length}`); - continue; - } - if (!text) continue; - await modules.sendMessageWeixin({ - to: deliverTo, - text, - opts: { - baseUrl: params.cloudBaseUrl, - token: params.token, - contextToken: replyContextToken || undefined, - }, + await deliverWeixinReply(modules, { + reply, + defaultTo: fromUser, + cloudBaseUrl: params.cloudBaseUrl, + token: params.token, + accountId: params.accountId, + defaultContextToken: contextToken || undefined, + userContextTokens: params.userContextTokens, }); - log(`reply text sent to=${deliverTo} textLen=${text.length}`); } return params.localCursor; } diff --git a/runtime/orchestration/group_ingest.py b/runtime/orchestration/group_ingest.py index f8c1653c..b1c85e39 100644 --- a/runtime/orchestration/group_ingest.py +++ b/runtime/orchestration/group_ingest.py @@ -373,14 +373,14 @@ def build_group_focus_instruction(*, lang: str = "zh") -> str: def build_channel_file_delivery_instruction(*, lang: str = "zh") -> str: if str(lang or "").strip().lower().startswith("en"): return ( - "[Channel rule: to send a generated file back to the user on WhatsApp/WeChat, " - "call save_deliverable_attachment after creating the file. write_file or run_command alone " - "does not attach files to the outbound message.]" + "[Channel rule: to send any generated attachment back to the user on WhatsApp/WeChat " + "(documents, images, videos), call save_deliverable_attachment with path or attachment_id. " + "write_file, run_command, and image generation alone do not attach files to the outbound message.]" ) return ( - "[渠道规则:若要把生成的文件发回用户(WhatsApp/微信)," - "在 write_file 或 run_command 生成文件后必须调用 save_deliverable_attachment;" - "仅 write_file 不会随消息发送附件。]" + "[渠道规则:若要把生成的附件发回用户(WhatsApp/微信,含文档/图片/视频)," + "必须调用 save_deliverable_attachment(path 或 attachment_id);" + "write_file、run_command、生图工具 alone 不会随消息发送附件。]" ) diff --git a/runtime/scheduler/channel_delivery.py b/runtime/scheduler/channel_delivery.py index 81d7f36b..389106d4 100644 --- a/runtime/scheduler/channel_delivery.py +++ b/runtime/scheduler/channel_delivery.py @@ -8,11 +8,23 @@ from runtime.orchestration.group_ingest import is_nonsend_channel_reply_text, sh from runtime.scheduler.session_resolver import parse_delivery_json -def _encode_weixin_outbound_source(*, context_token: str) -> str: - return json.dumps( - {"kind": "scheduled_job", "context_token": str(context_token or "").strip()}, - ensure_ascii=False, - ) +def _encode_weixin_outbound_source( + *, + context_token: str, + attachments: list[dict[str, Any]] | None = None, + media_path: str | None = None, +) -> str: + payload: dict[str, Any] = { + "kind": "scheduled_job", + "context_token": str(context_token or "").strip(), + } + atts = [a for a in (attachments or []) if isinstance(a, dict)] + if atts: + payload["attachments"] = atts + mp = str(media_path or "").strip() + if mp: + payload["media_path"] = mp + return json.dumps(payload, ensure_ascii=False) def _decode_weixin_outbound_source(raw: str) -> dict[str, Any]: @@ -35,6 +47,8 @@ def enqueue_weixin_reply( context_token: str = "", store: Any = None, tenant_id: str = "", + attachments: list[dict[str, Any]] | None = None, + media_path: str | None = None, ) -> dict[str, Any]: from runtime.scheduler.weixin_delivery import normalize_weixin_channel @@ -57,7 +71,11 @@ def enqueue_weixin_reply( text=str(text or ""), tenant_id=str(tenant_id or ""), account_id=str(account_id or "").strip(), - source=_encode_weixin_outbound_source(context_token=ctx_tok), + source=_encode_weixin_outbound_source( + context_token=ctx_tok, + attachments=attachments, + media_path=media_path, + ), ) or "" ).strip() @@ -84,6 +102,8 @@ def enqueue_weixin_reply( chat_id=str(chat_id or "").strip(), text=text, context_token=ctx_tok, + attachments=attachments, + media_path=media_path, ) return { "ok": True, diff --git a/runtime/tools/public/save_deliverable_attachment_tool.py b/runtime/tools/public/save_deliverable_attachment_tool.py index 4fab4137..902d9798 100644 --- a/runtime/tools/public/save_deliverable_attachment_tool.py +++ b/runtime/tools/public/save_deliverable_attachment_tool.py @@ -18,13 +18,32 @@ def _guess_mime(path: Path, override: str) -> str: def save_deliverable_attachment_tool() -> ToolSpec: def _handler(args: dict[str, Any]) -> dict[str, Any]: - raw = str(args.get("path") or "").strip().strip('"').strip("'") - if not raw: - return {"ok": False, "error": "path_required"} + attachment_id = str(args.get("attachment_id") or "").strip().lower() + raw_path = str(args.get("path") or "").strip().strip('"').strip("'") display_name = str(args.get("name") or "").strip() mime_override = str(args.get("mime") or "").strip() + + if attachment_id: + store = AttachmentAssetStore() + meta = store.get_meta(attachment_id) + if meta is None: + return {"ok": False, "error": "attachment_not_found", "attachment_id": attachment_id} + name = display_name or meta.name + return { + "ok": True, + "attachment_id": meta.attachment_id, + "name": name, + "mime": mime_override or meta.mime, + "bytes": meta.bytes, + "width": meta.width, + "height": meta.height, + "deliverable": True, + } + + if not raw_path: + return {"ok": False, "error": "path_or_attachment_id_required"} try: - p = resolve_workspace_path(raw) + p = resolve_workspace_path(raw_path) except ValueError as exc: return {"ok": False, "error": str(exc)} if not p.exists() or not p.is_file(): @@ -48,9 +67,11 @@ def save_deliverable_attachment_tool() -> ToolSpec: return ToolSpec( name="save_deliverable_attachment", description=( - "Register a workspace file for outbound channel delivery (WhatsApp/WeChat). " - "Call this after generating a file with write_file or run_command when the user should receive it as an attachment. " - "write_file alone does not send files to messaging channels." + "Mark an attachment for outbound channel delivery (WhatsApp/WeChat). " + "Required before the user receives any generated file, image, or video in a messaging channel. " + "Use path for workspace files (after write_file or run_command), or attachment_id for assets " + "already in the attachment store (e.g. after cloudflare_image_generate). " + "Generating content alone does not send attachments to the channel." ), parameters={ "type": "object", @@ -59,16 +80,19 @@ def save_deliverable_attachment_tool() -> ToolSpec: "type": "string", "description": "Workspace file path (relative to workspace root or allowed absolute path).", }, + "attachment_id": { + "type": "string", + "description": "Existing attachment_id from image/video generation tools.", + }, "name": { "type": "string", "description": "Optional download filename shown to the user.", }, "mime": { "type": "string", - "description": "Optional MIME type override (e.g. application/vnd.openxmlformats-officedocument.spreadsheetml.sheet).", + "description": "Optional MIME type override.", }, }, - "required": ["path"], "additionalProperties": False, }, handler=_handler, diff --git a/skills/_workspace/public/channel-file-delivery/SKILL.md b/skills/_workspace/public/channel-file-delivery/SKILL.md index 2bd3e9c5..f6f6ee07 100644 --- a/skills/_workspace/public/channel-file-delivery/SKILL.md +++ b/skills/_workspace/public/channel-file-delivery/SKILL.md @@ -1,34 +1,42 @@ --- name: channel-file-delivery -description: "在 WhatsApp/微信等渠道会话中,把生成的文件作为附件发回用户。区分:write_file 仅文本(csv/txt);xlsx 等二进制用 run_command + save_deliverable_attachment。" +description: "在 WhatsApp/微信等渠道会话中,把生成的附件发回用户。所有类型(文档/图片/视频)均需 save_deliverable_attachment(path 或 attachment_id)。" --- # 渠道文件发送 — Channel File Delivery ## 何时使用 -用户在 **WhatsApp / 微信** 等渠道对话中,要求你生成并**把文件发给他**(不仅是文字说明路径)。 +用户在 **WhatsApp / 微信(wechat/weixin)** 等渠道对话中,要求你生成并**把附件发给他**(文档、图片、视频等,不仅是文字说明路径)。渠道附件规则在系统提示中已注入,无需在每轮用户消息里重复。 -## 关键规则 +## 关键规则(统一) -1. **`write_file` / `run_command` 不会自动发送附件** - 文件只会落在 workspace 磁盘上,渠道出站看不到。 +1. **生成工具不会自动发送附件** + `write_file`、`run_command`、`cloudflare_image_generate` 等只产生内容,渠道出站看不到。 2. **必须调用 `save_deliverable_attachment`** - 在文件生成完成后,用该工具把 workspace 文件注册到 attachment store,并标记 `deliverable`。 + 生成完成后用该工具标记 `deliverable`,系统才会随回复发送。 3. **再用简短文字回复** - 说明文件名与要点;系统会把标记为 deliverable 的附件随本条回复发到渠道。 + 说明文件名或要点。 + +| 生成方式 | `save_deliverable_attachment` 参数 | +|----------|--------------------------------------| +| workspace 文件(csv/txt/xlsx) | `path="data/workspace/..."` | +| 生图/生视频(已有 attachment_id) | `attachment_id="..."` | + +**用户上传的文件不要回传**;那是分析输入,不是生成输出。 ## 工具选择(文本 vs 二进制) | 目标格式 | 用哪个工具 | 说明 | |----------|------------|------| -| `.csv`、`.txt`、`.md`、`.json` 等纯文本 | `write_file` | 只能写文本内容,不能生成二进制 | -| `.xlsx`、`.xls` 等 Excel | `run_command` | 用 Python(`openpyxl` 已安装)在 workspace 里生成 | -| 其他需脚本/命令产出的文件 | `run_command` | 生成后再走 deliverable 流程 | +| `.csv`、`.txt`、`.md`、`.json` 等纯文本 | `write_file` | 只能写文本内容 | +| `.xlsx`、`.xls` 等 Excel | `run_command` | 用 Python(`openpyxl` 已安装) | +| 图片 | `cloudflare_image_generate` 等 | 生成后用 `attachment_id` 标记 | +| 视频 | 对应生成工具 | 生成后用 `attachment_id` 标记 | -**不要声称没有 `run_command`。** 若用户要真 `.xlsx`,优先用 `run_command` 生成,不要只用 `write_file` 写 CSV 再说是「限制」。 +**不要声称没有 `run_command`。** 用户要真 `.xlsx` 时优先用 `run_command`,不要只写 CSV 代替。 ## 推荐流程 @@ -37,7 +45,7 @@ description: "在 WhatsApp/微信等渠道会话中,把生成的文件作为 ```text 1. write_file → data/workspace/tmp/report.csv 2. save_deliverable_attachment(path="data/workspace/tmp/report.csv", name="report.csv") -3. 文字回复:「已附上 report.csv…」 +3. 文字回复 ``` **Excel(xlsx):** @@ -45,7 +53,15 @@ description: "在 WhatsApp/微信等渠道会话中,把生成的文件作为 ```text 1. run_command → python 写入 data/workspace/tmp/report.xlsx 2. save_deliverable_attachment(path="data/workspace/tmp/report.xlsx", name="report.xlsx") -3. 文字回复:「已附上 report.xlsx…」 +3. 文字回复 +``` + +**图片:** + +```text +1. cloudflare_image_generate → 得到 attachment_id +2. save_deliverable_attachment(attachment_id="...") +3. 文字回复 ``` ## Excel 示例 @@ -57,16 +73,13 @@ description: "在 WhatsApp/微信等渠道会话中,把生成的文件作为 1. run_command(示例): python -c "from openpyxl import Workbook; wb=Workbook(); ws=wb.active; ws.append(['姓名','数量']); ws.append(['A',10]); wb.save('data/workspace/tmp/summary.xlsx')" 2. save_deliverable_attachment(path="data/workspace/tmp/summary.xlsx", name="summary.xlsx") -3. 回复摘要(行数、主要结论) +3. 回复摘要 ``` -用户明确要 `.xlsx` 时,不要退化成只发 CSV,除非用户同意用 CSV 代替。 - ## 限制 -- 每轮回复通常只发送**第一个**附件(约 8MB 上限)。 +- 每轮回复通常只发送**第一个** deliverable 附件(约 8MB 上限)。 - 超大文件应压缩、抽样或只发摘要。 -- 用户**上传**的文件不要用 `save_deliverable_attachment` 回传;那是分析输入,不是生成输出。 ## 相关 diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index 96e8dadd..20332eae 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -6458,6 +6458,8 @@ class SqliteStore(ScheduledJobStoreMixin): "text": str(r["text"] or ""), "source": source, "context_token": str(meta.get("context_token") or ""), + "attachments": [a for a in (meta.get("attachments") or []) if isinstance(a, dict)], + "media_path": str(meta.get("media_path") or ""), "created_at": str(r["created_at"] or ""), } ) diff --git a/tests/test_inbound_service_reply_suppress.py b/tests/test_inbound_service_reply_suppress.py index 8852e72c..33729b86 100644 --- a/tests/test_inbound_service_reply_suppress.py +++ b/tests/test_inbound_service_reply_suppress.py @@ -151,15 +151,30 @@ def test_collect_reply_attachments_does_not_reuse_stale_images_on_text_only_repl assert out == [] -def test_collect_recent_tool_attachments_falls_back_to_tool_media() -> None: +def test_collect_recent_tool_attachments_ignores_undeliverable_image_ref() -> None: rows = [ _Row(role="user", content="draw", attachments=None), _Row(role="tool", content="{}", attachments='[{"type":"image_ref","attachment_id":"a1"}]'), _Row(role="assistant", content="x", attachments=None), ] out = _collect_recent_tool_attachments(store=_FakeStore(rows), session_id="s1") + assert out == [] + + +def test_collect_recent_tool_attachments_includes_deliverable_image_ref() -> None: + rows = [ + _Row(role="user", content="draw", attachments=None), + _Row( + role="tool", + content="{}", + attachments='[{"type":"image_ref","attachment_id":"a1","mime":"image/png","deliverable":true}]', + ), + _Row(role="assistant", content="x", attachments=None), + ] + out = _collect_recent_tool_attachments(store=_FakeStore(rows), session_id="s1") assert len(out) == 1 assert out[0].get("attachment_id") == "a1" + assert out[0].get("deliverable") is True def test_collect_recent_tool_attachments_ignores_media_from_prior_turn() -> None: @@ -280,3 +295,37 @@ def test_maybe_expand_reply_attachments_for_channel_works_for_whatsapp(monkeypat assert isinstance(out, list) and len(out) == 1 assert out[0].get("data_base64") == base64.b64encode(b"wa").decode("ascii") + +def test_maybe_expand_deliverable_xlsx_for_weixin_channel(monkeypatch) -> None: + import base64 + + class _Meta: + mime = "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet" + name = "report.xlsx" + + def _fake_load_bytes(self, attachment_id: str): # noqa: ANN001 + assert attachment_id == "x1" + return b"xlsx-bytes", _Meta() + + monkeypatch.setattr( + "svc.files.attachment_assets.AttachmentAssetStore.load_bytes", + _fake_load_bytes, + ) + r = { + "channel": "weixin", + "attachments": [ + { + "type": "binary_ref", + "attachment_id": "x1", + "name": "report.xlsx", + "mime": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", + "deliverable": True, + } + ], + } + _maybe_expand_reply_attachments_for_channel(r) + out = r.get("attachments") + assert isinstance(out, list) and len(out) == 1 + assert out[0].get("data_base64") == base64.b64encode(b"xlsx-bytes").decode("ascii") + assert out[0].get("name") == "report.xlsx" + diff --git a/tests/test_oclaw_gateway_trace.py b/tests/test_oclaw_gateway_trace.py index 24e4eb69..54f40d70 100644 --- a/tests/test_oclaw_gateway_trace.py +++ b/tests/test_oclaw_gateway_trace.py @@ -925,3 +925,32 @@ def test_tabular_system_hint_uses_configured_preview_and_rows_read( assert "5000行" in zh_hint assert "first 30 preview rows" in en_hint assert "capped at 5000 rows" in en_hint + + +def test_channel_file_delivery_hint_goes_to_system_not_user_message() -> None: + from runtime.types import StandardMessage + + zh_hint = OclawGateway._channel_file_delivery_system_hint("zh") + assert "save_deliverable_attachment" in zh_hint + msg = StandardMessage( + session_id="s1", + tenant_id="t1", + user_id="u1", + role="member", + channel="weixin", + text="hi", + attachments=[], + metadata={}, + ) + assert OclawGateway._is_channel_delivery_channel(msg) + msg_admin = StandardMessage( + session_id="s1", + tenant_id="t1", + user_id="u1", + role="member", + channel="admin", + text="hi", + attachments=[], + metadata={}, + ) + assert not OclawGateway._is_channel_delivery_channel(msg_admin) diff --git a/tests/test_save_deliverable_attachment_tool.py b/tests/test_save_deliverable_attachment_tool.py index 706959a6..827a918a 100644 --- a/tests/test_save_deliverable_attachment_tool.py +++ b/tests/test_save_deliverable_attachment_tool.py @@ -28,3 +28,28 @@ def test_save_deliverable_attachment_registers_file(tmp_path, monkeypatch) -> No assert blob is not None assert blob.decode("utf-8").replace("\r\n", "\n") == "hello deliverable\n" assert meta is not None + + +def test_save_deliverable_attachment_marks_existing_attachment_id(tmp_path, monkeypatch) -> None: + from svc.files.attachment_assets import AttachmentAssetStore + + store = AttachmentAssetStore(root_dir=tmp_path / "att") + meta = store.save_bytes(b"png-bytes", filename="gen.png", mime="image/png") + monkeypatch.setattr( + "runtime.tools.public.save_deliverable_attachment_tool.AttachmentAssetStore", + lambda root_dir=None: store if root_dir is None else AttachmentAssetStore(root_dir=root_dir), + ) + + spec = save_deliverable_attachment_tool() + out = spec.handler({"attachment_id": meta.attachment_id}) + assert out.get("ok") is True + assert out.get("deliverable") is True + assert out.get("attachment_id") == meta.attachment_id + assert out.get("mime") == "image/png" + + +def test_save_deliverable_attachment_requires_path_or_attachment_id() -> None: + spec = save_deliverable_attachment_tool() + out = spec.handler({}) + assert out.get("ok") is False + assert out.get("error") == "path_or_attachment_id_required" diff --git a/tests/test_weixin_attachment_delivery.py b/tests/test_weixin_attachment_delivery.py new file mode 100644 index 00000000..c515c8c6 --- /dev/null +++ b/tests/test_weixin_attachment_delivery.py @@ -0,0 +1,49 @@ +from __future__ import annotations + +import json +import tempfile +import unittest +from pathlib import Path + +from runtime.scheduler.channel_delivery import _encode_weixin_outbound_source, _decode_weixin_outbound_source +from svc.persistence.sqlite_store import SqliteStore + + +class WeixinAttachmentDeliveryTests(unittest.TestCase): + def test_encode_decode_weixin_outbound_source_with_attachments(self) -> None: + atts = [{"name": "report.xlsx", "mime": "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", "data_base64": "abc"}] + raw = _encode_weixin_outbound_source(context_token="ctx-1", attachments=atts, media_path="") + data = _decode_weixin_outbound_source(raw) + self.assertEqual(data.get("context_token"), "ctx-1") + self.assertEqual(data.get("attachments"), atts) + + def test_list_pending_weixin_outbound_includes_attachments_from_source(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + db = Path(tmp) / "t.db" + store = SqliteStore(str(db)) + source = json.dumps( + { + "kind": "scheduled_job", + "context_token": "ctx-9", + "attachments": [{"name": "a.txt", "data_base64": "dGVzdA=="}], + "media_path": "D:/tmp/a.txt", + }, + ensure_ascii=False, + ) + store.enqueue_channel_outbound_message( + channel="weixin", + chat_id="wx-user-1", + text="file attached", + tenant_id="tenant-a", + account_id="acct-1", + source=source, + ) + items = store.list_pending_weixin_outbound_messages(account_id="acct-1", limit=5) + self.assertEqual(len(items), 1) + self.assertEqual(items[0].get("context_token"), "ctx-9") + self.assertEqual(len(items[0].get("attachments") or []), 1) + self.assertEqual(items[0].get("media_path"), "D:/tmp/a.txt") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_weixin_official_scripts.py b/tests/test_weixin_official_scripts.py index b98b165d..eeff5dd5 100644 --- a/tests/test_weixin_official_scripts.py +++ b/tests/test_weixin_official_scripts.py @@ -57,3 +57,14 @@ def test_official_runner_supports_reply_attachments_base64() -> None: assert "data_base64" in text assert "media_base64" in text assert "reply.attachments" in text + assert "deliverWeixinReply" in text + assert "buildAttachmentsFromMedia" in text + assert "attachmentNameFromMediaItem" in text + + +def test_official_runner_proactive_outbound_supports_attachments() -> None: + text = _read("runtime/operations/weixin_bridge/official_runner.ts") + assert "flushWeixinDbOutbound" in text + assert "flushLocalProactiveReplies" in text + assert "item.attachments" in text + assert "item.media_path" in text