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 da33dd29c9..8c0e66ac2e 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 @@ -105,7 +105,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 0c5d236c78..16be0e2b91 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -200,7 +200,7 @@ conversation. 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 @@ -921,3 +921,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 08b423e61e..780937add2 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 @@ -143,7 +143,7 @@ Core request 在 provider 投递前保存带时间的观测。重复事件保留 较新的 Session。已选定注册 Agent 时固定观测其确切绑定 Session:即使该会话已 失败或关闭也如实读取,不退回更新的会话;执行证据缺失和未知状态明确显示不可判定。这两个命令不会打开 Session、调用模型或创建 Goal。`/help` 按角色列出命令、既有设置 → Lark 的工作区、 -执行器与解绑入口,以及目前仅支持文字的附件边界。 +执行器与解绑入口,以及图片支持和暂不可用的媒体/原宿主边界。 回归使用生产原生文件 store、持久队列、bound request 与 provider 受理/投递路径, provider 和协议执行器为合成 fixture。它验证排队、停止、读回和重复投递,不证明 @@ -562,3 +562,16 @@ Turn HTTP 预算包含既有附件额度的 base64 编码:最多四张图片 源码验收覆盖 HTTP 准入、会话持久化读回和合成 Codex 协议进程,以及打包后的桌面/窄屏 图片发送、历史操作卡与拒绝后的草稿恢复。该证据只证明传输和界面行为,不证明真实模型质量、 公开发布或已安装宿主验收。GQ06 的材料入口及 GQ07–09 的连续性仍须完成各自的交付与恢复验收。 + +### 飞书私聊默认复用原生 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/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index b22199fd27..33d776d74f 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") @@ -34,7 +36,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()): @@ -98,6 +101,7 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A self._record_steward_ingress(row, selected, current, client_id) 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) @@ -109,9 +113,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": @@ -182,6 +191,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 e92ce38056..611a110403 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -1267,6 +1267,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", @@ -1287,6 +1288,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, @@ -1303,6 +1306,7 @@ def enqueue_turn( session_id, client_turn_id=client_turn_id, message=message, + attachments=attachments, origin=origin, ) self.resume_session_queue( @@ -1418,7 +1422,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 d13527aebe..cf10102b26 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -1053,6 +1053,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", @@ -1061,6 +1062,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( @@ -1075,6 +1080,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, }, @@ -1111,6 +1117,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, @@ -1143,6 +1150,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 d9379e050f..c6f5339f61 100644 --- a/loopx/control_plane/collaboration/conversation_binding.ts +++ b/loopx/control_plane/collaboration/conversation_binding.ts @@ -266,6 +266,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"); } @@ -273,6 +277,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 49d93dcc70..5dba487066 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -22,6 +22,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_caption, private_message_images class LarkPrivateConversations: @@ -159,6 +160,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. @@ -180,25 +182,44 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: return {"status": "invalid_text"} if not text.strip(): return {"status": "empty_text"} - command = {"/status": "status", "/help": "help", "/new": "new", "/stop": "stop"}.get(text.strip()) - if text.strip() == "/agents": + 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_input = private_message_caption(record["source_content"]) if attachments else text.strip() + command = {"/status": "status", "/help": "help", "/new": "new", "/stop": "stop"}.get(command_input) + if command_input == "/agents": command = "agents" - elif text.strip() == "/project": + elif command_input == "/project": command = "select_project" - elif text.strip() == "/agent" or text.strip().startswith("/agent "): + elif command_input == "/agent" or command_input.startswith("/agent "): command = "select_agent" if binding["context_kind"] == "steward": for prefix, selected_command in [("/delegate", "commission"), ("/委托", "commission"), ("/confirm", "confirm_commission"), ("/cancel", "cancel_commission"), ("/stop-commission", "stop_commission"), ("/resume-commission", "resume_commission")]: - if text.strip() == prefix or text.strip().startswith(prefix + " "): + if command_input == prefix or command_input.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) @@ -392,7 +413,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): @@ -428,7 +450,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 返回普通项目对话。", @@ -453,7 +475,7 @@ def _status_text(snapshot: dict[str, Any], *, help_requested: bool) -> str: f"\n执行器:{markdown_scalar(snapshot['executor_endpoint_id'])}" f"\n观察时间:{markdown_scalar(snapshot['observed_at'])}" "\n工作区、执行器与解绑:本机 Chat → 设置 → Lark。变更或解绑会重新核验授权;已受理工作不会迁移到新会话。" - "\n图片/文件目前未交给模型,请改用文字。") + "\n可直接发送图片或图文消息(PNG/JPEG/GIF/WebP,最多 4 张,单张 5 MB、合计 12 MB)。文件与音视频暂不支持;选择原宿主 Agent 后仅支持文字。") 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..1589c9598d --- /dev/null +++ b/loopx/extensions/lark/private_images.py @@ -0,0 +1,87 @@ +"""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_caption(content: str) -> str: + """Control parsing reads the caption; model input keeps image placeholders.""" + return _IMAGE.sub("", content).strip() + + +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 4e1defbc17..176e9bb88c 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("owner steward scope reads the current registry, preserves its Session and never widens a project or legacy scope", () => { const steward = {...row, context_kind: "steward", grant: "portfolio_read", goal_ids: [], goal_scope: "all_registered"}; const next = planConversationBinding({...request, binding: steward}).state; diff --git a/tests/test_lark_private_conversations.py b/tests/test_lark_private_conversations.py index 0b6d4e319a..de8e7f9395 100644 --- a/tests/test_lark_private_conversations.py +++ b/tests/test_lark_private_conversations.py @@ -153,17 +153,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..5dd242893c --- /dev/null +++ b/tests/test_lark_private_images.py @@ -0,0 +1,164 @@ +"""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", "![Image](img_example)"), + ("post", "Inspect this diagram\n![Image](img_example)\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, "![Image](img_example)", "下载失败"), + (False, "![Image](img_example)\n", "文件或音视频"), + (False, "\n".join(f"![Image](img_{i})" for i in range(5)), "最多支持 4"), + (False, "", "资源信息")]) +def test_unavailable_or_partial_media_never_executes_text_alone(ordinary, failure, content, notice): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + transport.runner = image_runner(provider, fail=failure) + try: + event = provider.event("notes-app", "unavailable", content, kind="image" if not content else "post") + assert transport.admit("notes-app", event)["status"] == "command_recorded" + transport.reconcile() + assert store.list_sessions() == [] + assert any(notice in text and "未提交执行" in text for _, text in provider.writes) + finally: + runtime.close() + + +def test_revoked_during_image_download_cannot_create_native_work(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + binding = transport.bindings.read()["bindings"][0] + transport.runner = image_runner(provider, revoke=lambda: transport.bindings.disconnect( + binding["binding_id"], expected_revision=transport.bindings.read()["revision"])) + try: + event = provider.event("notes-app", "revoked-image", "![Image](img_example)", kind="image") + assert transport.admit("notes-app", event)["status"] == "command_rejected" + assert store.list_sessions() == [] and transport.core.pending() == [] + finally: + runtime.close() + + +def test_queued_image_survives_runtime_restart_and_changed_replay_is_rejected(ordinary, monkeypatch): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + transport.runner = image_runner(provider) + monkeypatch.setattr(runtime, "resume_session_queue", lambda **kw: None) + event = provider.event("notes-app", "persisted", "Caption\n![Image](img_example)", kind="post") + assert transport.admit("notes-app", event)["status"] == "durably_accepted" + row = transport.core.pending()[0] + sid, tid = row["session_id"], row["turn_id"] + thread = store.load_session(sid)["upstream_thread_id"] + assert store.load_turn(sid, tid)["status"] == "queued" + runtime.close() + restored = ChatSessionStore(store.root.parent) + original = restored.load_turn(sid, tid) + assert original["attachments"] == row["attachments"] + repeated, created = restored.create_queued_turn(sid, client_turn_id=original["client_turn_id"], + message=original["message"], attachments=original["attachments"], origin="lark") + assert not created and repeated["turn_id"] == tid + with pytest.raises(ValueError, match="different request"): + restored.create_queued_turn(sid, client_turn_id=original["client_turn_id"], + message=original["message"], attachments=[], origin="lark") + restarted = ChatRuntimeController(store=restored, codex_bin=str(ordinary[5]), + project_contexts=ordinary[2], registry_path=runtime.registry_path) + try: + restarted.resume_session_queue(session_id=sid, work_dir=ordinary[6], objective="") + assert restarted.wait_for_turn(session_id=sid, turn_id=tid, timeout_sec=10)["status"] == "completed" + assert restored.load_session(sid)["upstream_thread_id"] == thread + requests = [json.loads(line) for line in ordinary[4].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) + finally: + restarted.close() + + +@pytest.mark.parametrize("raw,notice", [(b"", "支持 PNG"), + (PNG_BYTES + b"x" * (5 * 1024 * 1024), "单张图片")], ids=["unsupported-type", "oversize"]) +def test_resource_bytes_are_checked_before_core_admission(raw, notice): + with pytest.raises(ValueError, match=notice): + private_message_images(content="![Image](img_example)", message_type="image", + message_id="om_example", profile="notes-app", cli_bin="lark-cli", + runner=image_runner(None, content=raw)) + + +@pytest.mark.parametrize("command", ["/status", "/stop", "/new", "/project", "/agents"]) +def test_image_control_caption_never_starts_or_controls_native_work(ordinary, command): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + transport.runner = image_runner(provider) + try: + event = provider.event("notes-app", "image-control", + command + "\n![Image](img_example)", kind="post") + assert transport.admit("notes-app", event)["status"] == "command_recorded" + transport.reconcile() + assert store.list_sessions() == [] + assert any("未提交执行" in text for _, text in provider.writes) + finally: + runtime.close() + + +def test_help_describes_default_images_without_creating_work(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + event = provider.event("notes-app", "image-help", "/help") + assert transport.admit("notes-app", event)["status"] == "command_recorded" + transport.reconcile() + answer = provider.writes[-1][1] + assert "PNG/JPEG/GIF/WebP" in answer + assert "文件与音视频暂不支持" in answer + assert "仅支持文字" in answer + assert store.list_sessions() == [] + finally: + runtime.close()