diff --git a/apps/presentation/dashboard/src/features/personal-workspace/private-conversation-panel.tsx b/apps/presentation/dashboard/src/features/personal-workspace/private-conversation-panel.tsx index 7b60b191e8..1efda91c8a 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/private-conversation-panel.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/private-conversation-panel.tsx @@ -99,7 +99,7 @@ export function PrivateConversationPanel() {
{zh ? "从手机发送文字开始;后续消息进入原会话队列。/status 查看工作区、角色与持久排队状态,/help 查看用法与解绑入口,/stop 停止当前聊天执行,/new 开启新会话。图片、文件会明确提示暂不支持。" : "Send text from your phone to begin; follow-ups queue in the same Session. /status shows the workspace, role and durable queue, /help explains commands and where to unbind, /stop stops the current Chat Turn, /new starts a new conversation. Images and files receive an explicit unsupported response."}
+{zh ? "发送文字、图片或图文消息开始;后续消息进入原会话队列。/status 查看工作区、角色与持久排队状态,/help 查看用法与解绑入口,/stop 停止当前聊天执行,/new 开启新会话。文件、音视频、附在控制命令或已选择 Agent 上的图片会明确提示暂不支持。" : "Send text, images or image/text posts to begin; follow-ups queue in the same Session. /status shows the workspace, role and durable queue, /help explains commands and where to unbind, /stop stops the current Chat Turn, /new starts a new conversation. Files, audio/video, and images sent with control commands or to a selected attached Agent receive an explicit unsupported response."}
{zh ? "管家新委托:/delegate --tokens N 具体目标。先读预览,再用原私聊的完整 /confirm 命令确认;15 分钟过期。原生执行保持只读,总 token 上限可能被运行中的请求超过;没有默认定时调度。回执提供 /stop-commission 停止和 /resume-commission 恢复命令;恢复保留原线程及累计用量。" : "Steward commission: /delegate --tokens N objective. Read the preview, then use its full /confirm command in the original private Chat within 15 minutes. Native execution remains read-only; in-flight requests can exceed the total token allowance. No default schedule. Receipts provide /stop-commission and /resume-commission commands; recovery retains the original thread and cumulative usage."}
{error ?{error}
: null} ; diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index e6f6181508..bd643be837 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -208,7 +208,7 @@ Duplicate events retain that snapshot instead of switching to a newer Session. Missing execution evidence and unknown states stay explicitly unavailable. These commands open no Session, invoke no model and create no Goal. `/help` shows role-specific commands and the existing Settings → Lark entry for workspace, -executor and revocation, including the text-only attachment boundary. +executor and revocation, including supported images and unavailable media/attached-host boundaries. Regression coverage uses the production native filesystem store, durable queue, bound request and provider admission/reconciliation paths with a synthetic @@ -939,3 +939,23 @@ historical cards and rejected-draft recovery. These fixtures establish transport and interface behavior, not live model quality, public posting or installed-host acceptance. GQ06's material entry and GQ07–09's continuity remain subject to their full delivery and recovery acceptance. + + +### Default Lark private images reuse native Turn attachments + +Ordinary project and steward private conversations accept images and image/text +posts by default. The provider verifies the canonical message under its receiving +App, downloads only that message's resources as that App, and passes bounded +PNG/JPEG/GIF/WebP data into the existing Core request and durable Session queue. +Limits remain four images, 5 MiB each and 12 MiB total. Captions survive; resource +keys and private image bytes do not enter typed routing observations. Duplicate +events reuse downloaded input and the original Turn; restart drains that same +Turn and upstream thread. Grants are checked again after download and on return. + +Failed downloads, unsupported files/audio/video, and images sent with control +commands or to an attached host receive an explicit non-execution notice. The +provider does not execute only the text of a partially supported post. Attached +host media and file delivery remain separate gaps. Regression covers model-wire +image input, unchanged Session, replay, durable restart and download-time +revocation; live provider/model acceptance is reported separately. No new +Session authority, queue, worker or feature toggle is introduced. diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md index 4cf1553fde..945ce1d789 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.zh-CN.md @@ -139,7 +139,7 @@ executor endpoint、原生 Session/active Turn 与已持久排队数量。管家 Core request 在 provider 投递前保存带时间的观测。重复事件保留原快照,不切换到 较新的 Session;执行证据缺失和未知状态明确显示不可判定。这两个命令不会打开 Session、调用模型或创建 Goal。`/help` 按角色列出命令、既有设置 → Lark 的工作区、 -执行器与解绑入口,以及目前仅支持文字的附件边界。 +执行器与解绑入口,以及图片支持和暂不可用的媒体/原宿主边界。 回归使用生产原生文件 store、持久队列、bound request 与 provider 受理/投递路径, provider 和协议执行器为合成 fixture。它验证排队、停止、读回和重复投递,不证明 @@ -578,3 +578,17 @@ typed Core、HTTP 与原生宿主回归覆盖默认读写、明确只读、工 旧会话拒绝和原线程恢复。真实 Codex canary 按项目规则编辑并读回合成笔记,保留原文、 不创建 Goal;这只是宿主/文件系统结果,不代表真实 Lark 写入、素材 intake 或发布完成。 维护者 review、安装与 Lark 旅程、更多 IM 交互仍未关闭。 + + +### 飞书私聊默认复用原生 Turn 图片附件 + +普通项目与管家私聊默认接收图片和图文消息。provider 在接收 App 下核验 canonical +message,仅以该 App 身份下载属于这条消息的资源,再将 PNG/JPEG/GIF/WebP 交给 +既有 Core request 与持久 Session queue。沿用四张、单张 5 MiB、合计 12 MiB 上限。 +保留配文;资源 key 和私有图片字节不进入 typed routing 观测。重复事件复用原输入 +和 Turn;重启后仍由原 Turn、原 upstream thread 执行。下载后及回复前重新核验授权。 + +下载失败、文件/音视频、携图控制命令或原宿主 Agent 图片请求均明确告知未提交执行, +不会只执行混合消息的文字部分。原宿主媒体与文件交付仍待补齐。回归覆盖图片模型输入、 +原 Session、重复投递、持久重启与下载中撤权;真实 provider/model 验收另行记录。 +本增量不新增 Session authority、queue、worker 或默认关闭的功能开关。 diff --git a/loopx/canary/module_metric_baseline.json b/loopx/canary/module_metric_baseline.json index bcaeb848a7..f1e70669e0 100644 --- a/loopx/canary/module_metric_baseline.json +++ b/loopx/canary/module_metric_baseline.json @@ -32,7 +32,8 @@ }, "loopx/chat_runtime.py": { "any_count": 35, - "dict_any_count": 0 + "dict_any_count": 0, + "lines": 2004 }, "loopx/chat_server.py": { "any_count": 25, diff --git a/loopx/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index 4dbdb07870..7d56802c65 100644 --- a/loopx/capabilities/native_chat/external_conversations.py +++ b/loopx/capabilities/native_chat/external_conversations.py @@ -12,6 +12,7 @@ from typing import Any from ...chat_store import _atomic_write_json, _read_json +from ...chat_attachments import normalize_chat_image_attachments from ...file_lock import exclusive_file_lock @@ -24,7 +25,8 @@ def __init__(self, controller: Any) -> None: self.actions: Any | None = None def admit(self, *, binding_id: str, source: dict[str, Any], request_ref: str, - message: str, command: str | None = None) -> dict[str, Any]: + message: str, command: str | None = None, + attachments: list[dict[str, Any]] | None = None) -> dict[str, Any]: import re if not re.fullmatch(r"[a-f0-9]{24}", request_ref): raise ValueError("invalid external request reference") @@ -33,7 +35,8 @@ def admit(self, *, binding_id: str, source: dict[str, Any], request_ref: str, path = self.root / f"{request_ref}.json" with exclusive_file_lock(self.root / "source-fences" / f"{binding_id}.{source['source_ref']}.json", operation="route_external_chat_request"), exclusive_file_lock(path, operation="admit_external_chat_request"): selected = self.bindings.resolve(binding_id=binding_id, **source) - expected = {"binding_id": binding_id, "source": source, "message": message, "command": command} + expected = {"binding_id": binding_id, "source": source, "message": message, "command": command, + "attachments": normalize_chat_image_attachments(attachments) or None} if path.exists(): row = _read_json(path) if any(row.get(key) != value for key, value in expected.items()): @@ -83,6 +86,7 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A # Session check. Never move an accepted request to a new Session. turn, _ = controller.store.create_queued_turn(current["session_id"], client_turn_id=client_id, message=row["message"], origin="lark", + attachments=row.get("attachments"), external_agent_target={"target": target, "context": selected["context"]} if target else None) row.update(status="accepted", turn_id=turn["turn_id"]) _atomic_write_json(path, row) @@ -94,9 +98,14 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A observations = {"context": selected["context"], "observed_at": datetime.now(timezone.utc).isoformat(), "queued_count": len(controller.store.queued_turns(current["session_id"])) if current else 0, - "active_turn": controller.store.load_turn(current["session_id"], active_id) if active_id else None} + "active_turn": ({key: value for key, value in controller.store.load_turn(current["session_id"], active_id).items() + if key != "attachments"} if active_id else None)} + # Routing needs presence, not private image bytes. Persisted attachments + # remain in the native request/Turn and never enter the effect bridge. plan = effect_runtime_result("collaboration.conversation.request", { - "request": row, "current_session": current, "binding": selected["binding"], + "request": {key: value for key, value in row.items() if key != "attachments"}, + "attachment_count": len(row.get("attachments") or []), + "current_session": current, "binding": selected["binding"], "agent_target": target, **observations}) operation = plan["operation"] if operation == "select_recipient": @@ -166,6 +175,7 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A try: turn, _ = controller.enqueue_turn(session_id=current["session_id"], client_turn_id=plan["client_turn_id"], message=row["message"], + attachments=row.get("attachments"), work_dir=Path("."), objective="", origin="lark", external_agent_target={"target": target, "context": selected["context"]} if target else None) row.update(status="accepted", turn_id=turn["turn_id"]) diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 276ca4d5e6..0435747aff 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -1259,6 +1259,7 @@ def enqueue_turn( session_id: str, client_turn_id: str, message: str, + attachments: list[AttachmentPayload] | None = None, work_dir: Path, objective: str, origin: str = "external", @@ -1279,6 +1280,8 @@ def enqueue_turn( context = self.project_contexts.session_context(session) work_dir, objective = context["project"], context["objective"] if session.get("session_mode") == CHAT_SESSION_MODE_ATTACHED: + if attachments: + raise ValueError("attached host session queue does not yet accept attachments") turn, created = enqueue_attached_agent_turn( store=self.store, registry_path=self.registry_path, @@ -1295,6 +1298,7 @@ def enqueue_turn( session_id, client_turn_id=client_turn_id, message=message, + attachments=attachments, origin=origin, ) self.resume_session_queue( @@ -1410,7 +1414,7 @@ def _drain_session_queue( session_id=session_id, turn_id=turn_id, message=str(turn.get("message") or ""), - attachments=[], + attachments=turn.get("attachments") or [], adapter=adapter, done_event=done_event, ) diff --git a/loopx/chat_store.py b/loopx/chat_store.py index 4aadc882e2..097a04b40d 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -1029,6 +1029,7 @@ def create_queued_turn( *, client_turn_id: str, message: str, + attachments: list[dict[str, Any]] | None = None, goal_instance_id: str | None = None, ttl_seconds: int = SESSION_QUEUE_TTL_SECONDS, origin: str = "external", @@ -1037,6 +1038,10 @@ def create_queued_turn( """Persist one bounded follow-up without replacing the active Turn.""" client_id = _opaque_id(client_turn_id, field="client_turn_id") + from .chat_attachments import normalize_chat_image_attachments, validate_chat_turn_envelope + normalized_attachments = normalize_chat_image_attachments(attachments) or None + if normalized_attachments: + validate_chat_turn_envelope({"message": message, "attachments": normalized_attachments}) session_path = self._session_path(session_id) with self._session_lock(session_id): with exclusive_file_lock( @@ -1051,6 +1056,7 @@ def create_queued_turn( identity="client_turn_id", request={ "message": str(message), + "attachments": normalized_attachments, "origin": _opaque_id(origin, field="origin"), "external_agent_target": external_agent_target, }, @@ -1087,6 +1093,7 @@ def create_queued_turn( "status": "queued", **({"external_agent_target": external_agent_target} if external_agent_target is not None else {}), "message": str(message), + **({"attachments": normalized_attachments} if normalized_attachments else {}), "origin": _opaque_id(origin, field="origin"), "upstream_turn_id": None, "response": None, @@ -1119,6 +1126,7 @@ def create_queued_turn( text=message, turn_id=turn_id, origin=origin, + attachments=normalized_attachments, ) self.append_event( session_id, diff --git a/loopx/control_plane/collaboration/conversation_binding.ts b/loopx/control_plane/collaboration/conversation_binding.ts index 6743e9d380..e24d057d2f 100644 --- a/loopx/control_plane/collaboration/conversation_binding.ts +++ b/loopx/control_plane/collaboration/conversation_binding.ts @@ -257,6 +257,10 @@ export function planBoundConversationRequest(params: JsonObject): JsonObject { const row = requireJsonObject(params.request, "external request"); const request = ref(row.request_ref, "external request identity"); const command = row.command; + const imageCount = params.attachment_count ?? 0; + if (!Number.isSafeInteger(imageCount) || Number(imageCount) < 0 || Number(imageCount) > 4) { + throw new EffectRuntimeRequestError("invalid external image attachment count"); + } if (![null, "agents", "select_agent", "select_project", "status", "help", "new", "stop", "unsupported", "commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"].includes(command as null | string)) { throw new EffectRuntimeRequestError("unsupported external conversation command"); } @@ -264,6 +268,11 @@ export function planBoundConversationRequest(params: JsonObject): JsonObject { const target = row.target_recorded === true ? row : current; const session = target?.session_id ?? null; const turn = row.target_recorded === true ? row.turn_id ?? null : current?.active_turn_id ?? null; + if (Number(imageCount) > 0 && (command !== null || params.agent_target != null)) { + // Preserve the existing attached-host capability boundary and never drop + // images while executing a control command or handing off to that host. + return {operation: "reply", session_id: session, turn_id: null, response_code: "unsupported_attachment"}; + } if (["agents", "select_agent", "select_project"].includes(String(command))) { return {operation: "select_recipient", session_id: null, turn_id: null}; } diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index c6e4b71dd5..d5e642e7b9 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -19,6 +19,7 @@ from .goal_channel_transport import APP_ID_PATTERN, call, json_payload, lark_args from .inbox_reply import _message, reply_lark_event_inbox, verify_lark_inbox_reply from .inbox_reactions import mark_lark_event_inbox_processing, mark_lark_event_inbox_received +from .private_images import private_message_images class LarkPrivateConversations: @@ -156,6 +157,7 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: return {"status": "source_verification_failed"} message_type = str(event.get("message_type") or "") text = "" + attachments: list[dict[str, Any]] = [] if message_type == "text": # lark-cli renders event and mget content as plain text. Read # the full canonical message; do not decode a rendered event. @@ -177,6 +179,24 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: return {"status": "invalid_text"} if not text.strip(): return {"status": "empty_text"} + elif message_type in {"image", "post"}: + if source_message.get("msg_type", source_message.get("message_type")) != message_type: + return {"status": "source_conflict"} + content = source_message.get("content") + if not isinstance(content, str): + return {"status": "source_verification_failed"} + if record.get("source_content", content) != content: + return {"status": "source_conflict"} + record["source_content"] = content + if "attachments" not in record and not record.get("attachment_notice"): + try: + text, attachments = private_message_images(content=content, message_type=message_type, + message_id=event["message_id"], profile=profile, cli_bin=self.cli_bin, runner=self.runner) + record.update(message=text, attachments=attachments) + except ValueError as exc: + record["attachment_notice"] = str(exc) + _atomic_write_json(path, record) + text, attachments = record.get("message", ""), record.get("attachments", []) command = {"/status": "status", "/help": "help", "/new": "new", "/stop": "stop"}.get(text.strip()) if text.strip() == "/agents": command = "agents" @@ -191,11 +211,11 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: if text.strip() == prefix or text.strip().startswith(prefix + " "): command = selected_command break - if message_type != "text": + if message_type not in {"text", "image", "post"} or record.get("attachment_notice"): command = "unsupported" try: admitted = self.core.admit(binding_id=binding["binding_id"], source=source, - request_ref=request, message=text, command=command) + request_ref=request, message=text, command=command, attachments=attachments) except ValueError: record.update(status="rejected", response="操作或原授权不可用;Agent 请先用 /agents 查看确切命令,/project 返回项目对话。新委托请使用 /delegate --tokens N 具体目标,确认或取消请使用原预览中的完整命令。") _atomic_write_json(path, record) @@ -350,7 +370,8 @@ def inbox() -> Path: else: if record["status"] != "rejected": self._feedback(path, record, inbox=inbox) - response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") + response = (str(record.get("attachment_notice") or "") + or _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "")) if record.get("status_snapshot"): response = _status_text(record["status_snapshot"], help_requested=native.get("command") == "help") if self._deliver(path, record, "terminal", response, inbox=inbox): @@ -384,7 +405,7 @@ def inbox() -> Path: def _command_text(code: str) -> str: - return {"unsupported_attachment": "此入口目前只支持文字;图片或文件没有交给模型。请发送文字描述。", + return {"unsupported_attachment": "这条消息未提交执行。普通项目或管家对话支持文字与图片;文件、音视频及所选 Agent 的原宿主暂不支持图片。请将文字与图片单独发送,或用 /project 返回项目对话。", "no_session": "尚无会话;发送文字即可开始。", "active_session": "正在执行;后续文字会进入同一会话队列。", "ready_session": "会话已就绪,可继续发送文字。", "new_session": "已关闭此前会话;下一条文字将开启新会话。", "attached_control_unavailable": "原 Agent 宿主尚不支持此处的实时停止或新建会话;原执行没有被停止或替换。请在原宿主处理,/project 返回普通项目对话。", @@ -421,7 +442,8 @@ def _status_text(snapshot: dict[str, Any], *, help_requested: bool) -> str: if not steward: text += "\n/agents 查看本 App 已授权的 Agent;使用列表中的完整 /agent 命令选择,/project 返回此项目会话。" if help_requested: - text += "\n工作区、执行器与解绑:本机 Chat → 设置 → Lark。变更或解绑会重新核验授权;已受理工作不会迁移到新会话。图片/文件目前未交给模型,请改用文字。" + text += "\n可直接发送图片或图文消息(PNG/JPEG/GIF/WebP,最多 4 张,单张 5 MB、合计 12 MB)。文件与音视频暂不支持;选择原宿主 Agent 后仅支持文字。" + text += "\n工作区、执行器与解绑:本机 Chat → 设置 → Lark。变更或解绑会重新核验授权;已受理工作不会迁移到新会话。" if steward: text += "\n新委托:/delegate --tokens N 具体目标;读完预览后从原私聊发送完整 /confirm。/cancel 取消预览;/stop-commission 和 /resume-commission 使用原回执中的完整命令。" return text diff --git a/loopx/extensions/lark/private_images.py b/loopx/extensions/lark/private_images.py new file mode 100644 index 0000000000..5b58fe4fcf --- /dev/null +++ b/loopx/extensions/lark/private_images.py @@ -0,0 +1,82 @@ +"""Bounded image IO for a message already verified under its receiving App. + +Resource keys come from lark-cli's canonical message rendering; the provider +download endpoint checks that each key belongs to this exact message. Core's +existing image normalization, Session and Turn remain the only model boundary. +""" +from __future__ import annotations + +import base64 +import re +from pathlib import Path +from tempfile import TemporaryDirectory +from typing import Any + +from ...chat_attachments import ( + CHAT_IMAGE_MAX_BYTES, CHAT_IMAGE_MAX_COUNT, CHAT_IMAGE_MAX_TOTAL_BYTES, normalize_chat_image_attachments, +) +from .goal_channel_transport import call, json_payload, lark_args + +_IMAGE = re.compile(r"(?:!\[[^\]]*\]\(|\[Image: *)(img_[A-Za-z0-9_-]+)[)\]]") +_OTHER_RESOURCE = re.compile(r"<(?:file|folder|audio|video|media)\b") + + +def private_message_images(*, content: str, message_type: str, message_id: str, + profile: str, cli_bin: str, runner: Any) -> tuple[str, list[dict[str, Any]]]: + """Read all images or reject the whole message; never execute a partial post.""" + if _OTHER_RESOURCE.search(content): + raise ValueError("这条消息包含暂不支持的文件或音视频,尚未提交执行。请将文字与图片单独发送。") + keys = list(dict.fromkeys(_IMAGE.findall(content))) + if message_type == "image" and not keys: + raise ValueError("未能读取这张图片的资源信息,尚未提交执行。请重新发送图片。") + if len(keys) > CHAT_IMAGE_MAX_COUNT: + raise ValueError("一次最多支持 4 张图片,尚未提交执行。请分开发送。") + attachments = [] + total = 0 + with TemporaryDirectory(prefix="loopx-lark-images-") as temporary: + root = Path(temporary).resolve() + for index, key in enumerate(keys, 1): + result = call(lambda args, _cwd, timeout: runner(args, root, timeout), + lark_args(cli_bin=cli_bin, profile=profile, tail=["im", "+messages-resources-download", + "--message-id", message_id, "--file-key", key, "--type", "image", "--as", "bot", + "--output", f"./image-{index}", "--format", "json"])) + payload = json_payload(result) + if result.get("returncode") != 0 or payload.get("ok") is not True: + raise ValueError("图片下载失败,尚未提交执行。请重发;若仍失败,请检查此 App 的消息读取权限。") + data = payload.get("data") or {} + if not isinstance(data, dict): + raise ValueError("图片下载结果不可用,尚未提交执行。请重新发送。") + path = Path(str(data.get("saved_path") or "")) + path = path if path.is_absolute() else root / path + if path.is_symlink() or not path.resolve().is_relative_to(root) or not path.is_file(): + raise ValueError("图片下载结果不可用,尚未提交执行。请重新发送。") + try: + if path.stat().st_size > CHAT_IMAGE_MAX_BYTES: + raise ValueError("单张图片最多支持 5 MB,尚未提交执行。请压缩后重发。") + with path.open("rb") as stream: + raw = stream.read(CHAT_IMAGE_MAX_BYTES + 1) + except OSError as exc: + raise ValueError("图片下载结果不可读取,尚未提交执行。请重新发送。") from exc + if len(raw) > CHAT_IMAGE_MAX_BYTES: + raise ValueError("单张图片最多支持 5 MB,尚未提交执行。请压缩后重发。") + total += len(raw) + if total > CHAT_IMAGE_MAX_TOTAL_BYTES: + raise ValueError("图片总量超过限制(最多 12 MB),尚未提交执行。请分开发送。") + if raw.startswith(b"\x89PNG\r\n\x1a\n"): + mime = "image/png" + elif raw.startswith(b"\xff\xd8\xff"): + mime = "image/jpeg" + elif raw.startswith((b"GIF87a", b"GIF89a")): + mime = "image/gif" + elif raw.startswith(b"RIFF") and raw[8:12] == b"WEBP": + mime = "image/webp" + else: + raise ValueError("支持 PNG、JPEG、GIF 和 WebP 图片;这份资源尚未提交执行。") + attachments.append({"id": f"lark-image-{index}", "name": f"image-{index}", "mime_type": mime, + "data_url": f"data:{mime};base64," + base64.b64encode(raw).decode("ascii"), "size": len(raw)}) + try: + attachments = normalize_chat_image_attachments(attachments) + except ValueError as exc: + raise ValueError("图片总量超过限制(最多 12 MB),尚未提交执行。请分开发送。") from exc + text = _IMAGE.sub(lambda match: f"[图片 {keys.index(match[1]) + 1}]", content).strip() + return text or "请查看这张图片。", attachments diff --git a/tests/control_plane_ts/conversation_binding.test.ts b/tests/control_plane_ts/conversation_binding.test.ts index dc5fc99862..e9d52d083a 100644 --- a/tests/control_plane_ts/conversation_binding.test.ts +++ b/tests/control_plane_ts/conversation_binding.test.ts @@ -16,6 +16,20 @@ const current = {schema_version: "loopx_chat_conversation_bindings_v0", revision const request = {current, expected_revision: 0, operation: "configure", binding: row, observation, available_projects: [project]}; +test("images use the managed conversation and never disappear into commands or an attached host", () => { + const request = {request_ref: "e".repeat(24), command: null}; + assert.equal(planBoundConversationRequest({request, current_session: null, attachment_count: 1}).operation, "admit_turn"); + for (const command of ["stop", "new", "select_agent", "commission"]) { + assert.equal(planBoundConversationRequest({request: {...request, command}, current_session: null, + attachment_count: 1}).response_code, "unsupported_attachment"); + } + assert.equal(planBoundConversationRequest({request, current_session: null, attachment_count: 1, + agent_target: {session_id: "existing-host"}}).response_code, "unsupported_attachment"); + for (const count of [-1, 5, 0.5, "1"]) { + assert.throws(() => planBoundConversationRequest({request, current_session: null, attachment_count: count}), /attachment count/); + } +}); + test("explicit project writes remain App-bound and cannot exceed the host grant, another executor or an existing Session", () => { const original = planConversationBinding(request).state as typeof current; const use = {current: original, binding_id: row.binding_id, source_ref: "e".repeat(24), diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index 77e3228ae0..4192db16d2 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -140,17 +140,17 @@ def test_native_private_admission_queue_other_app_stop_and_verified_delivery(ord runtime.close() -def test_source_rejection_attachment_notice_and_ambiguous_reply_readback(ordinary): # noqa: F811 +def test_source_rejection_unsupported_file_notice_and_ambiguous_reply_readback(ordinary): # noqa: F811 _, runtime, provider, transport = connect(ordinary) try: - image = provider.event("notes-app", "image", "", kind="image") + image = provider.event("notes-app", "file", "", kind="file") assert transport.admit("notes-app", {**image, "sender_type": "app"})["status"] == "audience_rejected" assert transport.admit("steward-app", image)["status"] == "audience_rejected" assert transport.core.pending() == [] assert transport.admit("notes-app", image)["status"] == "command_recorded" provider.verify_replies = False assert transport.reconcile() == 0 - assert len(provider.writes) == 1 and "图片或文件没有交给模型" in provider.writes[0][1] + assert len(provider.writes) == 1 and "未提交执行" in provider.writes[0][1] assert transport.reconcile() == 0 and len(provider.writes) == 1 provider.verify_replies = True assert transport.reconcile() == 1 and len(provider.writes) == 1 diff --git a/tests/test_lark_private_images.py b/tests/test_lark_private_images.py new file mode 100644 index 0000000000..f0ca4a786f --- /dev/null +++ b/tests/test_lark_private_images.py @@ -0,0 +1,134 @@ +"""Default media journeys and failed downloads must never silently lose input.""" +import json +from pathlib import Path + +import pytest +from test_chat_ordinary_project import ordinary # noqa: F401 +from test_chat_image_attachments import PNG_BYTES, PNG_DATA_URL +from test_lark_private_conversations import connect + +from loopx.extensions.lark.private_images import private_message_images +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_store import ChatSessionStore + + +def image_runner(provider, *, fail=False, content=PNG_BYTES, revoke=None): + def run(args, cwd=None, timeout=None): + if "+messages-resources-download" not in args: + return provider(args, cwd, timeout) + if provider is not None: + provider.calls.append(list(args)) + assert args[args.index("--as") + 1] == "bot" + if revoke: + revoke() + if fail: + return {"returncode": 1, "stdout": '{"ok":false}', "stderr": ""} + path = Path(cwd) / args[args.index("--output") + 1] + path.write_bytes(content) + return {"returncode": 0, "stdout": json.dumps({"ok": True, + "data": {"saved_path": str(path), "size_bytes": len(content)}})} + return run + + +@pytest.mark.parametrize("kind,content", [("image", "[Image: img_example]"), ("image", ""), + ("post", "Inspect this diagram\n\nKeep the caption")]) +def test_default_images_reach_codex_in_original_session_and_replay_once(ordinary, kind, content): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + capture = ordinary[4] + transport.runner = image_runner(provider) + try: + transport.admit("notes-app", provider.event("notes-app", "first", "Remember our context")) + first = transport.core.pending()[0] + runtime.wait_for_turn(session_id=first["session_id"], turn_id=first["turn_id"], timeout_sec=10) + event = provider.event("notes-app", "diagram", content, kind=kind) + assert transport.admit("notes-app", event)["status"] == "durably_accepted" + row = next(r for r in transport.core.pending() if r.get("attachments")) + assert row["session_id"] == first["session_id"] + result = runtime.wait_for_turn(session_id=row["session_id"], turn_id=row["turn_id"], timeout_sec=10) + assert result["status"] == "completed" + assert result["attachments"][0]["data_url"] == PNG_DATA_URL + assert transport.admit("notes-app", {**event, "event_id": "redelivery"})["status"] == "durably_accepted" + assert len([c for c in provider.calls if "+messages-resources-download" in c]) == 1 + transcript = store.messages(row["session_id"]) + assert len([m for m in transcript if m.get("turn_id") == row["turn_id"] and m["role"] == "user"]) == 1 + requests = [json.loads(line) for line in capture.read_text().splitlines()] + wire = [r["params"]["input"] for r in requests if r.get("method") == "turn/start"][-1] + assert any(part.get("type") == "image" and part["url"] == PNG_DATA_URL for part in wire) + assert any("[图片 1]" in part.get("text", "") for part in wire) + if kind == "post": + assert "Keep the caption" in row["message"] + assert all(s["goal_id"] is None for s in store.list_sessions()) + finally: + runtime.close() + + +@pytest.mark.parametrize("failure,content,notice", [(True, "", "下载失败"), + (False, "\n