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 1bdf2489cf..7496818de9 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 @@ -84,7 +84,7 @@ export function PrivateConversationPanel() {
-

{zh ? "从手机发送文字开始;后续消息进入原会话队列。/status 查看聊天状态,/stop 停止当前聊天执行,/new 开启新会话。图片、文件会明确提示暂不支持。" : "Send text from your phone to begin; follow-ups queue in the same Session. /status checks Chat state, /stop stops the current Chat Turn, /new starts a new conversation. Images and files receive an explicit unsupported response."}

+

{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 ? "管家新委托:/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 b7b3868374..d2e500124b 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -125,6 +125,27 @@ open; runtime/authority promotion still requires maintainer review. Synthetic product previews: [empty steward and project assistant](../../assets/personal-workspace/private-steward-empty.png), [narrow](../../assets/personal-workspace/private-steward-empty-narrow.png). +## Private Chat status and help: scoped observations + +`/status` and `/help` use the existing typed bound-request owner to report the +authorized role, workspace, executor endpoint, native Session/active Turn and +durable queued-Turn count. Steward status counts only that binding's freshly +authorized commissions. An executor endpoint is not a selected registered Agent; +native execution ending is not Goal acceptance or proof of result delivery. + +The Core request persists a timestamped observation before provider delivery. +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. + +Regression coverage uses the production native filesystem store, durable queue, +bound request and provider admission/reconciliation paths with a synthetic +provider and protocol executor. It qualifies queue/stop/readback and duplicate +delivery, not live provider, mobile or installed-service acceptance of this delta. +Registered Agent selection and broader coordination remain open. + ## Decision: make the App the place where work conversations continue Users should be able to say “接着做,结果给我” / “Keep going and bring me the result” 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 4d1a1f16ac..33b033cd28 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 @@ -92,6 +92,22 @@ RPC 回执。Lark 专属设置 companion 归 extension;会话、请求、范 合成产品预览:[空管家与项目助手](../../assets/personal-workspace/private-steward-empty.png)、 [窄视口](../../assets/personal-workspace/private-steward-empty-narrow.png)。 +## 私聊状态与帮助:授权范围内的观测 + +`/status` 与 `/help` 复用既有 typed bound-request owner,展示已授权角色、工作区、 +executor endpoint、原生 Session/active Turn 与已持久排队数量。管家只统计此 binding +当前获授权的新委托。executor endpoint 不代表已经选择注册 Agent;原生执行结束 +不代表 Goal 验收,也不证明结果已经投递。 + +Core request 在 provider 投递前保存带时间的观测。重复事件保留原快照,不切换到 +较新的 Session;执行证据缺失和未知状态明确显示不可判定。这两个命令不会打开 +Session、调用模型或创建 Goal。`/help` 按角色列出命令、既有设置 → Lark 的工作区、 +执行器与解绑入口,以及目前仅支持文字的附件边界。 + +回归使用生产原生文件 store、持久队列、bound request 与 provider 受理/投递路径, +provider 和协议执行器为合成 fixture。它验证排队、停止、读回和重复投递,不证明 +本增量的真实 provider、手机或正式安装验收。注册 Agent 选择及更广协调仍开放。 + ## 决策:让 App 成为工作会话持续进行的地方 用户应能在 LoopX 中说“接着做,结果给我” / “Keep going and bring me the result”, diff --git a/loopx/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index 5b1d685062..7e94ee634a 100644 --- a/loopx/capabilities/native_chat/external_conversations.py +++ b/loopx/capabilities/native_chat/external_conversations.py @@ -27,7 +27,7 @@ def admit(self, *, binding_id: str, source: dict[str, Any], request_ref: str, import re if not re.fullmatch(r"[a-f0-9]{24}", request_ref): raise ValueError("invalid external request reference") - if command not in {None, "status", "new", "stop", "unsupported", "commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"}: + if command not in {None, "status", "help", "new", "stop", "unsupported", "commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"}: raise ValueError("unsupported external conversation command") selected = self.bindings.resolve(binding_id=binding_id, **source) path = self.root / f"{request_ref}.json" @@ -50,6 +50,12 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A controller = self.controller current = controller.store.latest_session(goal_id=None, agent_id=selected["binding"]["executor_endpoint_id"], channel_id=selected["channel_id"]) + if row["command"] in {"status", "help"}: + # Observation must retain failed/closed originals. Admission still + # uses the resumable selector and never resumes from this snapshot. + current = max(controller.store.session_candidates(goal_id=None, + agent_id=selected["binding"]["executor_endpoint_id"], channel_id=selected["channel_id"]), + key=lambda candidate: str(candidate.get("updated_at") or ""), default=None) if row.get("session_id") and row["command"] is None: current = controller.store.load_session(row["session_id"]) if current is None: @@ -64,8 +70,15 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A _atomic_write_json(path, row) return row from ...control_plane.effect_runtime import effect_runtime_result + observations = {} + if row["command"] in {"status", "help"}: + active_id = current.get("active_turn_id") if current else None + 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} plan = effect_runtime_result("collaboration.conversation.request", { - "request": row, "current_session": current, "binding": selected["binding"]}) + "request": row, "current_session": current, "binding": selected["binding"], **observations}) operation = plan["operation"] if operation == "steward_action": parsed = self.bindings._core("collaboration.steward.command", {"command": row["command"], "message": row["message"]}) @@ -86,6 +99,8 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A response="已持久受理此管家操作;正在核验原生操作与回执。") elif operation == "reply": row.update(status="command_completed", session_id=plan["session_id"], response_code=plan["response_code"]) + if "status_snapshot" in plan: + row["status_snapshot"] = plan["status_snapshot"] elif operation in {"new", "stop"}: row.update(target_recorded=True, session_id=plan["session_id"], turn_id=plan["turn_id"]) _atomic_write_json(path, row) diff --git a/loopx/control_plane/collaboration/conversation_binding.ts b/loopx/control_plane/collaboration/conversation_binding.ts index bd6f11e2ce..e6d30f2a41 100644 --- a/loopx/control_plane/collaboration/conversation_binding.ts +++ b/loopx/control_plane/collaboration/conversation_binding.ts @@ -151,7 +151,7 @@ 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; - if (![null, "status", "new", "stop", "unsupported", "commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"].includes(command as null | string)) { + if (![null, "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"); } const current = params.current_session === null ? null : requireJsonObject(params.current_session, "current Session"); @@ -163,10 +163,15 @@ export function planBoundConversationRequest(params: JsonObject): JsonObject { if (selected.context_kind !== "steward") throw new EffectRuntimeRequestError("explicit commissions require a selected steward"); return {operation: "steward_action", session_id: session, turn_id: null}; } - if (command === "status" || command === "unsupported") { + if (command === "status" || command === "help") { return {operation: "reply", session_id: session, turn_id: null, - response_code: command === "unsupported" ? "unsupported_attachment" : !current ? "no_session" - : current.active_turn_id ? "active_session" : "ready_session"}; + response_code: command === "help" ? "conversation_help" : !current ? "no_session" + : current.active_turn_id ? "active_session" : "ready_session", + status_snapshot: boundConversationStatus(params, current)}; + } + if (command === "unsupported") { + return {operation: "reply", session_id: session, turn_id: null, + response_code: "unsupported_attachment"}; } if (command === "stop" || command === "new") { return {operation: command, session_id: session, turn_id: turn, @@ -175,6 +180,50 @@ export function planBoundConversationRequest(params: JsonObject): JsonObject { return {operation: "admit_turn", client_turn_id: `external-${request}`, session_id: session, turn_id: null}; } +/** A labelled observation of the same authorized context and canonical queue. + * It neither creates a Session nor certifies execution or result delivery. + * The external request persists this snapshot so duplicate delivery cannot + * silently substitute a later Session or another binding's current state. + */ +function boundConversationStatus(params: JsonObject, current: JsonObject | null): JsonObject { + const selected = binding(params.binding); + const steward = selected.context_kind === "steward"; + const context = steward ? normalizeStewardContext(params.context) : normalizeProjectContext(params.context); + for (const key of ["binding_id", "project_ref", "provider_ref", "operator_ref"]) { + if (context[key] !== selected[key]) throw new EffectRuntimeRequestError("status context belongs to another binding"); + } + const source = ref(context.source_ref, "status source identity"); + const channel = steward ? `manager.external.native.${selected.binding_id}.${source}` + : `project.external.${selected.binding_id}.${source}`; + if (current) { + const saved = steward ? normalizeStewardContext(current.steward_context) : normalizeProjectContext(current.project_context); + if (current.channel_id !== channel || current.goal_id !== (steward ? "loopx-manager" : null) + || JSON.stringify({...saved, goal_ids: []}) !== JSON.stringify({...context, goal_ids: []})) { + throw new EffectRuntimeRequestError("status Session context changed"); + } + } + if (!Number.isSafeInteger(params.queued_count) || Number(params.queued_count) < 0 + || (!current && params.queued_count !== 0)) throw new EffectRuntimeRequestError("invalid canonical queue observation"); + const instant = requireNonEmptyString(params.observed_at, "status observation time"); + if (instant.length > 64 || !Number.isFinite(Date.parse(instant))) throw new EffectRuntimeRequestError("invalid status observation time"); + const turn = params.active_turn === null ? null : requireJsonObject(params.active_turn, "observed active Turn"); + if (turn && (!current || turn.session_id !== current.session_id || turn.turn_id !== current.active_turn_id)) { + throw new EffectRuntimeRequestError("status Turn belongs to another Session"); + } + for (const item of [current, turn]) { + if (item && (typeof item.status !== "string" || !/^[a-z][a-z0-9_]{0,63}$/.test(item.status))) { + throw new EffectRuntimeRequestError("invalid canonical status observation"); + } + } + return {schema_version: "loopx_chat_bound_status_v0", observed_at: instant, + context_kind: selected.context_kind, workspace_path: context.workspace_path, + executor_endpoint_id: selected.executor_endpoint_id, grant: selected.grant, + authorized_commission_count: steward ? (selected.goal_ids as string[]).length : 0, + session_status: current?.status ?? null, active_turn_status: turn?.status ?? null, + active_turn_observation_available: !current?.active_turn_id || turn !== null, + queued_count: params.queued_count}; +} + /** Explicit text grammar selects one existing typed Goal operation. Normal * conversation and model prose never grant creation or scheduling authority. */ export function stewardCommand(params: JsonObject): JsonObject { diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index b51225569a..d0e0bfe9b9 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -149,7 +149,7 @@ 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", "/new": "new", "/stop": "stop"}.get(text.strip()) + command = {"/status": "status", "/help": "help", "/new": "new", "/stop": "stop"}.get(text.strip()) if binding["context_kind"] == "steward": for prefix, selected_command in [("/delegate", "commission"), ("/委托", "commission"), ("/confirm", "confirm_commission"), ("/cancel", "cancel_commission"), @@ -241,7 +241,8 @@ def reconcile(self) -> int: self._deliver(path, record, "admission", native["response"]) continue record.update(status=native["status"], response=native.get("response"), - commission_resources=native.get("commission_resources")) + commission_resources=native.get("commission_resources"), + status_snapshot=native.get("status_snapshot")) if record["status"] == "accepted": # Receipt follows persistent Core admission and is # independent of terminal execution and reply delivery. @@ -254,6 +255,8 @@ def reconcile(self) -> int: "本次执行失败或已过期,原会话已保留;请发送 /status 后再决定是否重试。") else: response = _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): resources = record.get("commission_resources") or {} if resources.get("session_id") and resources.get("turn_id"): @@ -287,3 +290,34 @@ def _command_text(code: str) -> str: "no_session": "尚无会话;发送文字即可开始。", "active_session": "正在执行;后续文字会进入同一会话队列。", "ready_session": "会话已就绪,可继续发送文字。", "new_session": "已关闭此前会话;下一条文字将开启新会话。", "stop_requested": "已请求停止这条消息对应的执行。", "no_active_turn": "当前没有正在执行的消息。"}.get(code, "") + + +def _status_text(snapshot: dict[str, Any], *, help_requested: bool) -> str: + """Localize Core facts; never infer an Agent, grant or model completion.""" + steward = snapshot["context_kind"] == "steward" + phases = {"queued": "已受理等待执行", "starting": "正在启动", "running": "正在执行", + "completing": "正在收尾", "interrupting": "正在停止", "completed": "原生执行结束", + "interrupted": "已停止", "timed_out": "执行超时", "failed": "执行失败"} + if snapshot["session_status"] is None: + state = "尚无会话;发送文字即可开始。" + elif not snapshot["active_turn_observation_available"]: + state = "执行证据暂不可读;请在本机检查原会话。" + elif snapshot["active_turn_status"]: + state = phases.get(snapshot["active_turn_status"], "执行状态暂不可判定;请在本机检查原会话。") + else: + state = {"failed": "会话恢复失败;请在本机检查原会话。", + "resume_failed": "会话恢复失败;请在本机检查原会话。", "stale": "会话需要恢复。", + "starting": "会话正在启动。", "resuming": "会话正在恢复。", + "ready": "会话可继续。", "closed": "会话已关闭。"}.get( + snapshot["session_status"], "会话状态暂不可判定;请在本机检查原会话。") + text = (f"状态快照({snapshot['observed_at']})\n角色:{'长期管家' if steward else '普通项目对话'}" + f"\n工作区:{snapshot['workspace_path']}\n执行器:{snapshot['executor_endpoint_id']}" + f"\n{state}\n已持久排队:{snapshot['queued_count']} 条。") + text += (f"\n已授权新委托:{snapshot['authorized_commission_count']};执行结束不代表委托验收。" if steward else + "\n当前仅有工作区只读授权;没有自动选用注册 Agent 或创建 Goal。") + text += "\n/status 查看状态;/stop 停止当前聊天执行;/new 关闭当前聊天并开启下次新会话;/help 查看用法。" + if help_requested: + text += "\n工作区、执行器与解绑:本机 Chat → 设置 → Lark。变更或解绑会重新核验授权;已受理工作不会迁移到新会话。图片/文件目前未交给模型,请改用文字。" + if steward: + text += "\n新委托:/delegate --tokens N 具体目标;读完预览后从原私聊发送完整 /confirm。/cancel 取消预览;/stop-commission 和 /resume-commission 使用原回执中的完整命令。" + return text diff --git a/tests/control_plane_ts/conversation_binding.test.ts b/tests/control_plane_ts/conversation_binding.test.ts index 16a46f880f..195c596a9e 100644 --- a/tests/control_plane_ts/conversation_binding.test.ts +++ b/tests/control_plane_ts/conversation_binding.test.ts @@ -118,3 +118,60 @@ test("native commands retain their recorded target across a new Session and rede assert.throws(() => planBoundConversationRequest({request: {...request, command: "grant"}, current_session: current}), /unsupported/); }); + +test("status and help project only the verified binding without opening a Session", () => { + const context = resolveBoundConversation({current: planConversationBinding(request).state, + binding_id: row.binding_id, source_ref: "e".repeat(24), sender_ref: row.operator_ref, + private_human_message: true, observation, available_projects: [project]}).context; + const input = {binding: row, context, current_session: null, queued_count: 0, active_turn: null, + observed_at: "2026-01-01T10:00:00Z", request: {request_ref: "f".repeat(24), command: "status"}}; + const plan = planBoundConversationRequest(input); + assert.equal(plan.operation, "reply"); + assert.equal(plan.session_id, null); + assert.equal(plan.response_code, "no_session"); + assert.deepEqual(plan.status_snapshot, {schema_version: "loopx_chat_bound_status_v0", + observed_at: input.observed_at, context_kind: "project", workspace_path: project.workspace_path, + executor_endpoint_id: "codex", grant: "workspace_read", authorized_commission_count: 0, + session_status: null, active_turn_status: null, active_turn_observation_available: true, queued_count: 0}); + assert.equal(planBoundConversationRequest({...input, request: {...input.request, command: "help"}}).response_code, + "conversation_help"); + for (const bad of [{...input, binding: {...row, provider_ref: "a".repeat(24)}}, + {...input, queued_count: -1}, {...input, queued_count: .5}, {...input, queued_count: 1}, + {...input, observed_at: "unknown"}]) assert.throws(() => planBoundConversationRequest(bad)); +}); + +test("status reads the canonical queue and exact active Turn, retaining unavailable evidence", () => { + const selected = resolveBoundConversation({current: planConversationBinding(request).state, + binding_id: row.binding_id, source_ref: "e".repeat(24), sender_ref: row.operator_ref, + private_human_message: true, observation, available_projects: [project]}); + const session = {session_id: "original", channel_id: selected.channel_id, goal_id: null, + project_context: selected.context, status: "busy", active_turn_id: "original-turn"}; + const turn = {session_id: session.session_id, turn_id: session.active_turn_id, status: "running"}; + const input = {binding: row, context: selected.context, current_session: session, queued_count: 2, + active_turn: turn, observed_at: "2026-01-01T10:00:00Z", + request: {request_ref: "f".repeat(24), command: "status"}}; + const snapshot = planBoundConversationRequest(input).status_snapshot as Record; + assert.equal(snapshot.queued_count, 2); + assert.equal(snapshot.active_turn_status, "running"); + assert.equal((planBoundConversationRequest({...input, active_turn: null}).status_snapshot as Record) + .active_turn_observation_available, false); + for (const bad of [{...input, current_session: {...session, channel_id: "another-owner"}}, + {...input, active_turn: {...turn, session_id: "another-session"}}, + {...input, active_turn: {...turn, turn_id: "later-turn"}}, + {...input, active_turn: {...turn, status: {pretend: "completed"}}}]) { + assert.throws(() => planBoundConversationRequest(bad)); + } +}); + +test("steward status counts its fresh grant, not a global portfolio or old Session scope", () => { + const steward = {...row, context_kind: "steward", grant: "portfolio_read", goal_ids: ["fresh-goal"]}; + const context = {...project, kind: "bound_steward", audience: "bound_owner", grant: "portfolio_read", + goal_ids: ["fresh-goal"], binding_id: row.binding_id, source_ref: "e".repeat(24), + provider_ref: row.provider_ref, operator_ref: row.operator_ref}; + const plan = planBoundConversationRequest({binding: steward, context, queued_count: 0, active_turn: null, + observed_at: "2026-01-01T10:00:00Z", request: {request_ref: "f".repeat(24), command: "help"}, + current_session: {session_id: "steward", channel_id: `manager.external.native.${row.binding_id}.${context.source_ref}`, + goal_id: "loopx-manager", steward_context: {...context, goal_ids: []}, status: "ready", active_turn_id: null}}); + assert.equal((plan.status_snapshot as Record).authorized_commission_count, 1); + assert.equal(JSON.stringify(plan).includes("fresh-goal"), false); +}); diff --git a/tests/test_lark_private_status.py b/tests/test_lark_private_status.py new file mode 100644 index 0000000000..18c12f2d9d --- /dev/null +++ b/tests/test_lark_private_status.py @@ -0,0 +1,179 @@ +"""Status observes real native files/queue without creating model work.""" +import time + +import pytest + +from test_chat_ordinary_project import ordinary # noqa: F401 +from test_lark_private_conversations import connect +from loopx.extensions.lark.private_conversations import LarkPrivateConversations, _status_text +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_agent import CodexChatAgentError + + +def test_help_and_empty_steward_status_do_not_open_sessions_or_borrow_scope(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + transport.admit("notes-app", provider.event("notes-app", "help", "/help")) + assert transport.reconcile() == 1 + message = provider.writes[-1][1] + assert "角色:普通项目对话" in message and "只读授权" in message + assert "/new" in message and "设置 → Lark" in message + assert "/delegate" not in message and "没有自动选用注册 Agent" in message + assert store.list_sessions() == [] + project = runtime.project_contexts.available()[0] + transport.bindings.configure(transport_ref="steward-app", project_ref=project["project_ref"], + executor_endpoint_id="codex", context_kind="steward") + transport.admit("steward-app", provider.event("steward-app", "empty", "/status")) + assert transport.reconcile() == 1 + assert "角色:长期管家" in provider.writes[-1][1] + assert "已授权新委托:0" in provider.writes[-1][1] + assert "执行结束不代表委托验收" in provider.writes[-1][1] + assert store.list_sessions() == [] + assert all(row["turn_id"] is None for row in transport.core.pending()) + finally: + runtime.close() + + +def test_pending_status_readback_recovers_from_files_without_resending(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + provider.verify_replies = False + transport.admit("notes-app", provider.event("notes-app", "status", "/status")) + snapshot = transport.core.pending()[0]["status_snapshot"] + assert transport.reconcile() == 0 + assert len(provider.writes) == 1 + assert transport.health()[transport.bindings.read()["bindings"][0]["binding_id"]]["recovery_count"] == 1 + # A new production transport instance observes the original durable + # intent and native request; it neither opens a Session nor resends. + restarted = LarkPrivateConversations(controller=runtime, runtime_root=store.root.parent, + runner=provider, cli_bin="lark-cli") + provider.verify_replies = True + assert restarted.reconcile() == 1 + assert len(provider.writes) == 1 + assert restarted.core.pending()[0]["status_snapshot"] == snapshot + assert store.list_sessions() == [] + finally: + runtime.close() + + +@pytest.mark.parametrize("resume_error", [False, True]) +def test_status_and_help_observe_original_session_after_actual_upstream_resume(ordinary, resume_error): # noqa: F811 + import json + store, original_runtime, provider, transport = connect(ordinary) + _, _, contexts, _, capture, fake, workspace = ordinary + restarted = None + try: + transport.admit("notes-app", provider.event("notes-app", "initial", "initial")) + original = transport.core.pending()[0] + sid = original["session_id"] + original_runtime.wait_for_turn(session_id=sid, turn_id=original["turn_id"], timeout_sec=10) + assert transport.reconcile() == 1 + upstream = store.load_session(sid)["upstream_thread_id"] + original_runtime.close() + if resume_error: + fake.write_text(fake.read_text().replace( + ' elif method in {"thread/start", "thread/resume"}:', + ' elif method == "thread/resume":\n' + ' print(json.dumps({"id": request_id, "error": {"code": -32000, "message": "Original thread unavailable"}}), flush=True)\n' + ' continue\n' + ' elif method in {"thread/start", "thread/resume"}:')) + restarted = ChatRuntimeController(store=store, codex_bin=str(fake), project_contexts=contexts, + registry_path=original_runtime.registry_path) + def resume(): + return restarted.open_session(goal_id=None, agent_id="codex", work_dir=workspace, objective="", + mode="resume_latest", conversation_binding_id=original["binding_id"], source_context=original["source"]) + if resume_error: + with pytest.raises(CodexChatAgentError, match="could not be restored"): + resume() + else: + assert resume()[0]["session_id"] == sid + before_observation = capture.read_text().splitlines() + assert any(json.loads(line).get("method") == "thread/resume" for line in before_observation) + transport = LarkPrivateConversations(controller=restarted, runtime_root=store.root.parent, + runner=provider, cli_bin="lark-cli") + for command in ["status", "help"]: + transport.admit("notes-app", provider.event("notes-app", command, f"/{command}")) + row = next(row for row in transport.core.pending() if row["command"] == command) + assert row["session_id"] == sid + assert row["status_snapshot"]["session_status"] == ("resume_failed" if resume_error else "ready") + transport.reconcile() + assert ("会话恢复失败" if resume_error else "会话可继续") in provider.writes[-1][1] + assert "尚无会话" not in provider.writes[-1][1] + assert len(store.list_sessions()) == 1 + assert store.load_session(sid)["upstream_thread_id"] == upstream + requests = [json.loads(line) for line in capture.read_text().splitlines()] + assert len([row for row in requests if row.get("method") == "thread/start"]) == 1 + assert capture.read_text().splitlines() == before_observation + finally: + if restarted: + restarted.close() + original_runtime.close() + + +def test_status_after_new_observes_closed_original_without_reopening_it(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + transport.admit("notes-app", provider.event("notes-app", "initial", "initial")) + first = transport.core.pending()[0] + runtime.wait_for_turn(session_id=first["session_id"], turn_id=first["turn_id"], timeout_sec=10) + transport.admit("notes-app", provider.event("notes-app", "new", "/new")) + transport.admit("notes-app", provider.event("notes-app", "status", "/status")) + snapshot = next(row for row in transport.core.pending() if row["command"] == "status")["status_snapshot"] + assert snapshot["session_status"] == "closed" + transport.reconcile() + assert any("会话已关闭" in text for _, text in provider.writes) + assert len(store.list_sessions()) == 1 + transport.admit("notes-app", provider.event("notes-app", "next", "next")) + second = next(row for row in transport.core.pending() if row["message"] == "next") + assert second["session_id"] != first["session_id"] + runtime.wait_for_turn(session_id=second["session_id"], turn_id=second["turn_id"], timeout_sec=10) + finally: + runtime.close() + + +def test_busy_status_reads_durable_queue_and_duplicate_retains_original_snapshot(ordinary): # noqa: F811 + store, runtime, provider, transport = connect(ordinary) + try: + transport.admit("notes-app", provider.event("notes-app", "slow", "wait for interrupt")) + first = transport.core.pending()[0] + sid = first["session_id"] + deadline = time.monotonic() + 10 + while store.load_session(sid).get("active_turn_id") != first["turn_id"]: + assert time.monotonic() < deadline + time.sleep(.01) + transport.admit("notes-app", provider.event("notes-app", "queued", "follow-up")) + status = provider.event("notes-app", "status", "/status") + transport.admit("notes-app", status) + original = next(row for row in transport.core.pending() if row["command"] == "status") + assert original["status_snapshot"]["queued_count"] == 1 + assert original["status_snapshot"]["active_turn_status"] in {"starting", "running"} + transport.reconcile() + assert any("已持久排队:1 条" in text for _, text in provider.writes) + transport.admit("notes-app", provider.event("notes-app", "stop", "/stop")) + queued = next(row for row in transport.core.pending() if row["message"] == "follow-up") + runtime.wait_for_turn(session_id=sid, turn_id=queued["turn_id"], timeout_sec=10) + transport.reconcile() + before = len(provider.writes) + transport.admit("notes-app", {**status, "event_id": "redelivery"}) + transport.reconcile() + assert len(provider.writes) == before + assert transport.core.read_request(original["request_ref"])["status_snapshot"] == original["status_snapshot"] + transport.admit("notes-app", provider.event("notes-app", "current", "/status")) + transport.reconcile() + assert "已持久排队:0 条" in provider.writes[-1][1] + assert len(store.list_sessions()) == 1 + assert store.load_session(sid)["goal_id"] is None + assert not any("source_ref" in text or "operator_ref" in text for _, text in provider.writes) + finally: + runtime.close() + + +@pytest.mark.parametrize("unknown", ["session", "turn", "missing"]) +def test_unavailable_execution_evidence_is_not_presented_as_ready(unknown): + snapshot = {"context_kind": "project", "observed_at": "2026-01-01T10:00:00Z", + "workspace_path": "/authorized/notes", "executor_endpoint_id": "codex", "grant": "workspace_read", + "queued_count": 0, "authorized_commission_count": 0, "session_status": "future_state" if unknown == "session" else "busy", + "active_turn_status": "future_state" if unknown == "turn" else None, + "active_turn_observation_available": unknown != "missing"} + text = _status_text(snapshot, help_requested=False) + assert "暂不可" in text and "会话可继续" not in text