diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index b83e54bf7f..674967c1d9 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -2380,9 +2380,9 @@ export async function updateGoalOwnership(body: { goal_id: string; mode: Executi const privateConversationSchema = z.object({ - binding_id: z.string(), app_ref: z.string(), context_kind: z.literal("project"), + binding_id: z.string(), app_ref: z.string(), context_kind: z.enum(["project", "steward"]), project_ref: z.string(), project_title: z.string(), context_available: z.boolean(), executor_endpoint_id: z.string(), - grant: z.literal("workspace_read"), listener_status: z.string(), + grant: z.enum(["workspace_read", "portfolio_read"]), goal_count: z.number().int().default(0), listener_status: z.string(), pending_count: z.number().int(), recovery_count: z.number().int(), }); const privateConversationsSchema = z.object({ok: z.literal(true), revision: z.number().int(), @@ -2391,10 +2391,10 @@ export type PrivateConversation = z.infer; export async function fetchPrivateConversations() { return privateConversationsSchema.parse(await requestJson("/api/chat/lark/private-conversations")); } -export async function connectPrivateConversation(appRef: string, projectRef: string, executor: string) { +export async function connectPrivateConversation(appRef: string, projectRef: string, executor: string, contextKind: "project" | "steward" = "project") { return privateConversationsSchema.parse(await requestJson("/api/chat/lark/private-conversations", { method: "POST", headers: {"Content-Type": "application/json"}, - body: JSON.stringify({app_ref: appRef, project_ref: projectRef, executor_endpoint_id: executor}), + body: JSON.stringify({app_ref: appRef, project_ref: projectRef, executor_endpoint_id: executor, context_kind: contextKind}), })); } export async function disconnectPrivateConversation(bindingId: string, revision: number) { 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 066e903c50..1bdf2489cf 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 @@ -16,6 +16,7 @@ export function PrivateConversationPanel() { const [app, setApp] = useState(""); const [project, setProject] = useState(""); const [executor, setExecutor] = useState(""); + const [role, setRole] = useState<"project" | "steward">("project"); const [error, setError] = useState(""); const [busy, setBusy] = useState(false); @@ -54,10 +55,11 @@ export function PrivateConversationPanel() { finally {setBusy(false);} } return
-

{zh ? "本人私聊 · 项目对话" : "Owner private Chat · Project conversation"}

-

{zh ? "每个 App 单独核验登录本人,只读讨论所选工作区。普通私聊不会创建 Goal。" : "Verify the logged-in owner independently for each App. Discuss the selected workspace with a read grant; ordinary private Chat creates no Goal."}

+

{zh ? "本人私聊 · 项目助手与管家" : "Owner private Chat · Project assistant and steward"}

+

{zh ? "每个 App 单独核验本人。项目助手只读讨论工作区,不创建隐藏 Goal。管家从空 portfolio 开始,只管理在此入口明确确认的新委托。" : "Verify the owner independently for each App. Project Chat discusses the workspace without hidden Goals. A steward starts with an empty portfolio and manages only new commissions explicitly confirmed here."}

{rows.map(row =>
{row.app_ref} · {row.context_available ? row.project_title : (zh ? "工作区不可用" : "Workspace unavailable")} +

{row.context_kind === "steward" ? (zh ? `LoopX 管家 · ${row.goal_count === 0 ? "暂无已授权的新委托;没有继承旧目标。" : `${row.goal_count} 个已确认的新委托`}` : `LoopX steward · ${row.goal_count} new confirmed commissions; no inherited Goals.`) : (zh ? "普通项目助手 · 只读对话" : "Project assistant · Read-only Chat")}

{row.executor_endpoint_id} · {zh ? "监听状态" : "Listener"}: {listenerLabel(row.listener_status)}

{zh ? `待处理或回复:${row.pending_count}` : `Pending execution or reply: ${row.pending_count}`}

{row.recovery_count > 0 ?

{zh ? "存在尚未确认的发送回执。服务会读取原回执恢复;不要重新发送同一任务。检查 App 登录、权限和原会话后刷新状态。" : "A send receipt is unconfirmed. The service reads the original receipt to recover; avoid resending the same task. Check this App login, permissions and original Session, then refresh status."}

: null} @@ -69,6 +71,9 @@ export function PrivateConversationPanel() { {apps.map(app => )} + -
-

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

+

{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 ? "管家新委托:/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 34dbbeabd1..b7b3868374 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -60,7 +60,7 @@ browser coverage. A native Codex source canary separately qualifies distinct upstream threads with independently observed real App identities. Neither is a live Lark/model/mobile result. Installed service qualification, phone journeys, registered Agent selection, real incremental/media/permission interactions and -a separately granted long-running steward remain open acceptance. Runtime and +broader long-running coordination remain open acceptance. Runtime and permission-boundary changes require maintainer review before promotion. The private setup UI composes within Settings → Lark; the App workspace scope @@ -74,6 +74,57 @@ Synthetic product previews: [desktop](../../assets/personal-workspace/private-pr [narrow](../../assets/personal-workspace/private-project-conversations-narrow.png), [revoked workspace](../../assets/personal-workspace/private-project-workspace-revoked.png). +## Bound steward private Chat: explicit new commissions + +Settings → Lark can now select a steward role independently of ordinary project +Chat. The existing typed conversation binding owns its App, verified owner, +source, workspace and bounded portfolio. A verified empty scope is distinct from +an unavailable authorization; it contains no inherited Goals and does not certify +global inventory coverage. Ordinary Chat still has no Goal or manager identity. + +Only an explicit `/delegate --tokens N objective` prepares a `goal.create` +preview. Confirmation must arrive from that exact owner/App/source within fifteen +minutes. The immutable preview shows the read-only boundary, total native token +allowance and absence of automatic scheduling. Existing canonical Chat actions +create the Goal and return their receipt; Core adopts only that exact new creation +in the configured workspace. Existing Goal and single-workspace fallbacks cannot +redirect it. Neither normal conversation nor model prose creates a commission. + +The existing service worker advances the durable Core request without blocking +inbound admission. Codex native Goal continuation performs the read-only work; +its result returns through the original private-source delivery journal. +`/stop-commission` freezes the execution target at admission, while an explicit +`/resume-commission ... --tokens N` retains the original Session, native thread, +objective and cumulative usage. Native completion is host execution evidence, +not canonical Goal/Todo acceptance. An allowance includes previous usage and +context; an in-flight request can exceed it. No default heartbeat is enabled. + +Portfolio extension refreshes scoped evidence/tools in the same steward thread. +After creation commits, the existing request journal saves its exact resource +receipt before attempting portfolio adoption. An adoption or readback failure +keeps that operation queued for recovery. Recovery rechecks the original +binding and canonical receipt, adopts the same resources, and returns their +result without creating another Goal or model thread. Notification cannot +settle an operation whose adoption is pending. Fault journeys cover adoption +failure, lost adoption readback, an interrupted creation-owner call and revoked +binding; they qualify provider/Core IO recovery, not a full host restart. +App, owner, source and workspace identity remain frozen and rechecked. A silent +native event reader cannot block the control RPC receipt needed to pause work. +The provider-specific setup companion belongs to the Lark extension; the +conversation, request, scope and creation semantics remain in their typed owners. + +Regression journeys cover independently verified empty versus missing scope, +wrong App/source confirmation, expiry, creation/adoption, native result return, +stop, duplicate events and same-thread recovery. A local actual Codex source +canary separately exercised a budget-limited synthetic commission and resumed it +in the same native thread to return verified fixture findings. It did not prove +real Lark inbound or phone acceptance. Multi-Agent coordination, actual media / +permission interactions, installed login recovery and mobile journeys remain +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). + ## 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 5dccae4bd4..4d1a1f16ac 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 @@ -42,7 +42,7 @@ provider 读回确认回复;发生没有 receipt 的不确定写入时不盲 后续消息持久排队、exact stop,以及重启后原会话恢复和已确认回复不重复发送。 这些是合成 provider/协议验收;真实原生 Codex 另行验证独立线程与上下文隔离。 真实 Lark 收发、安装候选、手机旅程、注册 Agent 选择、媒体/增量/权限回调,以及 -新管家的明确长期委托仍未验收。当前私聊绑定仅支持普通只读项目会话。 +更广的长期协调仍未验收。这一普通项目绑定保持只读,新增管家入口见下一检查点。 私聊配置复用设置 → Lark;App 范围仍是本机普通对话的唯一入口。未存储 App 身份 的旧群聊 profile 保留原 profile-hash 监听锁键。没有私聊绑定时不增加鉴权;有绑定 @@ -53,6 +53,45 @@ extension,typed binding owner 继续保持 provider-neutral。 ![窄屏私聊设置](../../assets/personal-workspace/private-project-conversations-narrow.png) ![工作区撤权读回](../../assets/personal-workspace/private-project-workspace-revoked.png) +## 本人私聊管家:明确的新委托 + +设置 → Lark 可为独立 App 选择管家角色。既有 typed conversation binding +固定 App、独立核验的本人、来源、工作区和有界 portfolio。已核验的空范围与授权 +不可用分开:新管家没有继承旧 Goal;空范围也不证明全局 inventory 覆盖完整。 +普通项目聊天仍没有 Goal 或管家身份。 + +只有显式 `/delegate --tokens N 具体目标` 会准备既有 `goal.create` 预览。 +确认必须在十五分钟内从原本人、App 和私聊来源进入。预览固定只读边界、原生总 +token 上限和不启用默认调度的事实。既有 canonical Chat action 创建 Goal 并返回 +回执,Core 仅把该确切新创建加入此管家范围。已有 Goal 或单一工作区 fallback +不能重定向委托;普通聊天和模型文字不能创建委托。 + +既有服务 worker 推进已持久受理的 Core 操作,入站不等终态。只读工作由 Codex +原生 Goal continuation 执行,结果经原私聊 delivery journal 返回。 +`/stop-commission` 在受理时固定确切执行目标;显式 +`/resume-commission ... --tokens N` 保留原 Session、原生线程、目标和累计用量。 +原生完成是宿主执行证据,不是 canonical Goal/Todo 验收。总上限包含历史用量与 +上下文,运行中的请求可能超过上限;没有默认 heartbeat。 + +新委托扩展 portfolio 时,在原管家线程刷新有界证据和工具。创建提交后,既有 +请求日志先保存确切资源回执,再尝试加入管家范围。加入或读回失败时,操作保持 +排队以供恢复。恢复重新核验原 binding 与 canonical 回执,采用同一批资源并 +返回结果,不再创建 Goal 或模型线程。范围加入仍待恢复时,通知回执不能结算 +该操作。故障旅程覆盖加入失败、加入后读回丢失、创建 owner 调用中断及 binding +撤权;验证的是 provider/Core IO 恢复,不是完整宿主重启。App、本人、来源和 +工作区仍固定并重新核验。原生事件读取在没有输出时,也不能卡住停止所需的控制 +RPC 回执。Lark 专属设置 companion 归 extension;会话、请求、范围和创建语义 +仍由现有 typed owner 持有。 + +回归覆盖独立空态与无授权、另一 App/来源确认拒绝、过期、创建与 scope adoption、 +原生结果返回、停止、重复事件和同线程恢复。本地真实 Codex 源码 canary 另验证 +合成委托遇到额度限制后,沿用原生线程恢复并返回可核验的 fixture 结果。 +这不证明真实飞书入站或手机验收。多 Agent 协调、真实媒体/权限交互、正式安装与 +登录恢复、手机旅程仍开放;runtime/authority 发布继续等待维护者审核。 + +合成产品预览:[空管家与项目助手](../../assets/personal-workspace/private-steward-empty.png)、 +[窄视口](../../assets/personal-workspace/private-steward-empty-narrow.png)。 + ## 决策:让 App 成为工作会话持续进行的地方 用户应能在 LoopX 中说“接着做,结果给我” / “Keep going and bring me the result”, diff --git a/docs/assets/personal-workspace/private-steward-empty-narrow.png b/docs/assets/personal-workspace/private-steward-empty-narrow.png new file mode 100644 index 0000000000..16e0e17b62 Binary files /dev/null and b/docs/assets/personal-workspace/private-steward-empty-narrow.png differ diff --git a/docs/assets/personal-workspace/private-steward-empty.png b/docs/assets/personal-workspace/private-steward-empty.png new file mode 100644 index 0000000000..b2162ea7e4 Binary files /dev/null and b/docs/assets/personal-workspace/private-steward-empty.png differ diff --git a/loopx/capabilities/native_chat/conversation_bindings.py b/loopx/capabilities/native_chat/conversation_bindings.py index 6ff7abbcf2..615449f7a7 100644 --- a/loopx/capabilities/native_chat/conversation_bindings.py +++ b/loopx/capabilities/native_chat/conversation_bindings.py @@ -40,17 +40,23 @@ def _core(operation: str, params: dict[str, Any]) -> dict[str, Any]: return result def configure(self, *, transport_ref: str, project_ref: str, - executor_endpoint_id: str) -> dict[str, Any]: + executor_endpoint_id: str, context_kind: str = "project") -> dict[str, Any]: observation = self.observe(transport_ref) candidate = { "schema_version": "loopx_chat_conversation_binding_v0", "binding_id": uuid.uuid4().hex[:24], "transport_ref": transport_ref, "provider_ref": observation["provider_ref"], "operator_ref": observation["operator_ref"], - "context_kind": "project", "project_ref": project_ref, - "executor_endpoint_id": executor_endpoint_id, "grant": "workspace_read", "enabled": True, + "context_kind": context_kind, "project_ref": project_ref, + "executor_endpoint_id": executor_endpoint_id, + "grant": "workspace_read" if context_kind == "project" else "portfolio_read", "enabled": True, + **({"goal_ids": []} if context_kind == "steward" else {}), } with exclusive_file_lock(self.path, operation="configure_chat_conversation_binding"): current = self.read() + previous = next((row for row in current["bindings"] if row["transport_ref"] == transport_ref), None) + if (context_kind == "steward" and previous and all(previous.get(key) == candidate.get(key) + for key in ["context_kind", "project_ref", "executor_endpoint_id", "provider_ref", "operator_ref"])): + candidate["goal_ids"] = previous["goal_ids"] result = self._core("collaboration.conversation.binding", { "current": current, "expected_revision": current.get("revision"), "operation": "configure", "binding": candidate, "observation": observation, "available_projects": self.projects.available(), @@ -90,3 +96,28 @@ def resolve(self, *, binding_id: str, source_ref: str, sender_ref: str, def session_context(self, saved: dict[str, Any]) -> dict[str, Any]: return self.resolve(binding_id=saved["binding_id"], source_ref=saved["source_ref"], sender_ref=saved["operator_ref"], private_human_message=True, session_context=saved) + + def steward_scope(self, session: dict[str, Any]) -> list[str] | None: + saved = session.get("steward_context") + if not isinstance(saved, dict): + return None + selected = self.session_context(saved) + if session.get("goal_id") != "loopx-manager" or session.get("channel_id") != selected["channel_id"]: + raise ValueError("bound steward audience changed") + return list(selected["context"]["goal_ids"]) + + def adopt_created_goal(self, *, binding_id: str, source: dict[str, Any], proposal: dict[str, Any], + goal: dict[str, Any]) -> None: + selected = self.resolve(binding_id=binding_id, **source) + with exclusive_file_lock(self.path, operation="adopt_steward_created_goal"): + current = self.read() + result = self._core("collaboration.conversation.binding", { + "current": current, "expected_revision": current["revision"], "operation": "adopt_created_goal", + "binding_id": binding_id, "context": selected["context"], "proposal": proposal, + "goal": {"goal_id": goal["id"], "workspace_path": str(Path(goal["repo"]).resolve()), + "creation_operation_id": goal.get("creation_operation_id")}, + }) + if result["changed"]: + _atomic_write_json(self.path, result["state"]) + if self.read() != result["state"]: + raise OSError("steward scope publication did not verify") diff --git a/loopx/capabilities/native_chat/external_conversations.py b/loopx/capabilities/native_chat/external_conversations.py index dff21378f7..5b1d685062 100644 --- a/loopx/capabilities/native_chat/external_conversations.py +++ b/loopx/capabilities/native_chat/external_conversations.py @@ -7,6 +7,8 @@ from __future__ import annotations from pathlib import Path +from datetime import datetime, timedelta, timezone +import hashlib from typing import Any from ...chat_store import _atomic_write_json, _read_json @@ -18,13 +20,14 @@ def __init__(self, controller: Any) -> None: self.controller = controller self.bindings = controller.project_contexts.conversation_bindings self.root = controller.store.root / "external-requests" + 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]: 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"}: + if command not in {None, "status", "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" @@ -38,7 +41,8 @@ def admit(self, *, binding_id: str, source: dict[str, Any], request_ref: str, return row else: row = {"schema_version": "loopx_chat_external_request_v0", "request_ref": request_ref, - **expected, "status": "prepared", "session_id": None, "turn_id": None} + **expected, "status": "prepared", "session_id": None, "turn_id": None, + "created_at": datetime.now(timezone.utc).isoformat()} _atomic_write_json(path, row) return self._admit_prepared(path, row, selected) @@ -60,9 +64,27 @@ 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 - plan = effect_runtime_result("collaboration.conversation.request", {"request": row, "current_session": current}) + plan = effect_runtime_result("collaboration.conversation.request", { + "request": row, "current_session": current, "binding": selected["binding"]}) operation = plan["operation"] - if operation == "reply": + if operation == "steward_action": + parsed = self.bindings._core("collaboration.steward.command", {"command": row["command"], "message": row["message"]}) + if row["command"] in {"stop_commission", "resume_commission"} and not row.get("target_recorded"): + if self.actions is None: + raise ValueError("the native commission owner is unavailable") + proposal = self.actions.load(parsed["argument"]) + self.bindings._core("collaboration.steward.authorize_creation", {"context": selected["context"], + "proposal": proposal, "operation": "stop" if row["command"] == "stop_commission" else "resume"}) + resources = proposal["receipt"]["resource_ids"] + target = controller.store.load_session(resources["session_id"]) + if not target or target["goal_id"] != resources["goal_id"]: + raise ValueError("the exact commission Session is unavailable") + row.update(target_recorded=True, target_session_id=target["session_id"], + target_turn_id=target.get("active_turn_id"), target_goal_id=resources["goal_id"]) + _atomic_write_json(path, row) + row.update(status="command_queued", command_argument=parsed, session_id=plan["session_id"], + response="已持久受理此管家操作;正在核验原生操作与回执。") + elif operation == "reply": row.update(status="command_completed", session_id=plan["session_id"], response_code=plan["response_code"]) elif operation in {"new", "stop"}: row.update(target_recorded=True, session_id=plan["session_id"], turn_id=plan["turn_id"]) @@ -96,11 +118,155 @@ def _admit_prepared(self, path: Path, row: dict[str, Any], selected: dict[str, A def pending(self) -> list[dict[str, Any]]: return [_read_json(path) for path in sorted(self.root.glob("*.json"))] + def read_request(self, request_ref: str) -> dict[str, Any]: + return _read_json(self.root / f"{request_ref}.json") + + def commission_evidence(self, session: dict[str, Any]) -> list[dict[str, Any]]: + """Read original Core execution facts for this exact bound audience. + + Native terminal status/result delivery is separate from Goal/Todo + acceptance. A portfolio reader must not infer one from the other. + """ + saved = session["steward_context"] + allowed = set(self.bindings.steward_scope(session) or []) + latest: dict[str, dict[str, Any]] = {} + for row in sorted(self.pending(), key=lambda row: row.get("created_at", "")): + resources = row.get("commission_resources") or {} + if (row["binding_id"] == saved["binding_id"] and row["source"]["source_ref"] == saved["source_ref"] + and resources.get("goal_id") in allowed): + latest[resources["goal_id"]] = row + facts = [] + for goal_id, row in latest.items(): + resources = row["commission_resources"] + current = self.controller.store.load_session(resources["session_id"]) + turn = self.controller.store.load_turn(resources["session_id"], resources["turn_id"]) + if not current or not turn or current.get("goal_id") != goal_id: + continue + facts.append({"goal_id": goal_id, "turn_status": turn["status"], + "native_execution": current.get("native_goal"), + "result_delivery_verified": row.get("delivery_verified") is True, + "result_excerpt": str((turn.get("response") or {}).get("message") or "")[:500], + "canonical_acceptance_attested": False}) + return facts + + def _run_steward_action(self, row: dict[str, Any]) -> None: + """The existing service worker advances a journaled native operation. + + No provider owns Goal authority or a second scheduler. Confirmation is + persisted before the canonical action service performs any write. + """ + if self.actions is None: + return + path = self.root / f"{row['request_ref']}.json" + with exclusive_file_lock(path, operation="advance_steward_action"): + row = _read_json(path) + if row["status"] != "command_queued": + return + try: + selected = self.bindings.resolve(binding_id=row["binding_id"], **row["source"]) + context = selected["context"] + argument = row["command_argument"]["argument"] + if row.get("commission_adoption_pending"): + applied = self.actions.load(row["proposal_id"]) + self.bindings._core("collaboration.steward.authorize_creation", { + "context": context, "proposal": applied, "now": row["created_at"], "operation": "confirm"}) + if (applied.get("status") != "applied" + or applied["receipt"]["resource_ids"] != row["commission_resources"]): + raise ValueError("the applied commission receipt changed") + self._adopt_commission(row, applied) + elif row["command"] == "commission": + key = f"steward-commission-{row['request_ref']}" + existing = next((p for p in self.actions.store.list() if p["idempotency_key"] == key), None) + workspace = context["workspace_path"] + expiry = (datetime.fromisoformat(row["created_at"]) + timedelta(minutes=15)).isoformat() + audience = {key: context[key] for key in ["binding_id", "source_ref", "provider_ref", "operator_ref", "project_ref"]} + proposal = existing or self.actions.preview({"action_kind": "goal.create", + "summary": f"新管家委托:{argument[:180]}", "idempotency_key": key, + "context": {"kind": "manager", "goal_id": "loopx-manager", **audience, "expires_at": expiry}, + "normalized_parameters": {"goal_id": f"steward-{row['request_ref']}", "title": argument[:180], + "objective": argument, "completion_criteria": "按明确委托返回可核验结果;原生执行结束不代表 LoopX 验收。", + "execution_boundary": "只读所选工作区;不继承旧目标,不自动调度,不扩大宿主策略。", + "agent_id": selected["binding"]["executor_endpoint_id"], + "workspace_ref": "workspace-" + hashlib.sha256(workspace.encode()).hexdigest()[:12], + "heartbeat": {"enabled": False}, + "native_token_budget": row["command_argument"]["native_token_budget"]}}) + row.update(proposal_id=proposal["proposal_id"], response=( + f"待确认的新委托:{argument}\n只读执行;总 token 上限 {row['command_argument']['native_token_budget']},运行中请求可能超过该上限。" + "不自动设置长期调度;未继承旧目标。确认后创建 Goal 并启动原生持续执行;其完成状态不等于 LoopX 验收。" + f"\n15 分钟内发送 /confirm {proposal['proposal_id']};取消请发送 /cancel {proposal['proposal_id']}。")) + else: + proposal = self.actions.load(argument) + if proposal is None: + raise ValueError("the exact commission preview is unavailable") + from ...control_plane.effect_runtime import effect_runtime_result + effect_runtime_result("collaboration.steward.authorize_creation", { + "context": context, "proposal": proposal, "now": row["created_at"], + "operation": "stop" if row["command"] == "stop_commission" else + "resume" if row["command"] == "resume_commission" else "confirm"}) + # An authenticated, canonical source selected this exact + # preview; crash replay retains the same confirmation. + row.update(proposal_id=argument, confirmation_recorded=True) + _atomic_write_json(path, row) + if row["command"] == "stop_commission": + if row.get("target_turn_id"): + self.controller.interrupt_turn(session_id=row["target_session_id"], turn_id=row["target_turn_id"]) + row["response"] = "已对受理时记录的确切执行请求停止;没有停止其它委托或改变 Goal 验收状态。" if row.get("target_turn_id") else "受理此停止请求时,该委托没有正在执行的消息。" + elif row["command"] == "resume_commission": + goal = self.actions._goal(row["target_goal_id"]) + turn, _ = self.controller.submit_turn(session_id=row["target_session_id"], + client_turn_id=f"commission-resume-{row['request_ref']}", + message=f"/goal resume --tokens {row['command_argument']['native_token_budget']}", + work_dir=Path(goal["repo"]), objective=str(goal.get("objective") or "")) + row.update(commission_resources={"goal_id": row["target_goal_id"], "session_id": row["target_session_id"], + "turn_id": turn["turn_id"]}, response=f"已受理原委托的恢复,保留原生线程、目标及累计用量。总 token 上限 {row['command_argument']['native_token_budget']};结果会返回此私聊。\n停止:/stop-commission {argument}") + elif row["command"] == "cancel_commission": + self.actions.cancel(argument) + row["response"] = "已取消这份新委托预览;没有创建或启动 Goal。" + else: + result = self.actions.apply(argument, steward_context=context, steward_confirmed_at=row["created_at"]) + applied = result["proposal"] + if applied.get("status") != "applied": + raise ValueError("the creation preview became stale; prepare a new /delegate request") + resources = applied["receipt"]["resource_ids"] + # Canonical creation has committed. Retain its exact + # receipt before fallible adoption or notification. + row.update(commission_resources=resources, commission_adoption_pending=True, + commission_gate=bool(result.get("gate"))) + _atomic_write_json(path, row) + self._adopt_commission(row, applied) + row["status"] = "command_completed" + except (OSError, ValueError, KeyError, RuntimeError) as exc: + # No automatic wider permission, regenerated preview or new + # native thread is used to hide an unavailable operation. + if row["command"] == "confirm_commission" and row.get("commission_resources"): + row.update(status="command_queued", commission_adoption_pending=True, + failure_kind=type(exc).__name__, + response="已持久受理此管家操作;正在核验原生操作与回执。") + else: + row.update(status="command_completed", failure_kind=type(exc).__name__, + response="管家操作未完成:原工作区、授权、预览有效期或原生执行入口没有通过核验。原操作已保留;请在本机核对回执后再决定是否重试。") + _atomic_write_json(path, row) + + def _adopt_commission(self, row: dict[str, Any], applied: dict[str, Any]) -> None: + if self.actions is None: + raise ValueError("the native commission owner is unavailable") + resources = row["commission_resources"] + goal = self.actions._goal(resources["goal_id"]) + self.bindings.adopt_created_goal(binding_id=row["binding_id"], source=row["source"], proposal=applied, goal=goal) + row.update(commission_adoption_pending=False, response=( + f"原生委托已创建:{resources['goal_id']}。操作回执已读回;处理结果会回到此私聊。没有自动设置调度。" + f"\n停止此执行:/stop-commission {row['proposal_id']}")) + row.pop("failure_kind", None) + if row.get("commission_gate"): + row["response"] += "首轮执行有待处理 Gate;请在本机查看该 Goal 的操作回执。" + def record_delivery(self, request_ref: str, *, session_id: str | None, turn_id: str | None) -> None: # This transport receipt does not alter a canonical Turn or grant. path = self.root / f"{request_ref}.json" with exclusive_file_lock(path, operation="record_external_chat_delivery"): row = _read_json(path) + if row.get("commission_adoption_pending"): + raise ValueError("commission adoption is pending; notification cannot settle the effect") if (row.get("session_id"), row.get("turn_id")) != (session_id, turn_id): raise ValueError("delivery correlation changed") row["delivery_verified"] = True @@ -109,12 +275,14 @@ def record_delivery(self, request_ref: str, *, session_id: str | None, turn_id: def recover(self) -> None: for row in self.pending(): try: - if row.get("delivery_verified"): + if row.get("delivery_verified") and not row.get("commission_adoption_pending"): continue self.bindings.resolve(binding_id=row["binding_id"], **row["source"]) if row["status"] == "prepared": self.admit(binding_id=row["binding_id"], source=row["source"], request_ref=row["request_ref"], message=row["message"], command=row["command"]) + elif row["status"] == "command_queued": + self._run_steward_action(row) elif row["status"] == "accepted": session = self.controller.store.load_session(row["session_id"]) if session and session.get("status") != "closed": diff --git a/loopx/capabilities/native_chat/project_context.py b/loopx/capabilities/native_chat/project_context.py index 0d5eec44f1..ca9c352bcb 100644 --- a/loopx/capabilities/native_chat/project_context.py +++ b/loopx/capabilities/native_chat/project_context.py @@ -50,6 +50,16 @@ def resolve(self, project_ref: str, *, session_context: dict[str, Any] | None = raise ValueError(str(exc)) from exc def session_context(self, session: dict[str, Any]) -> dict[str, Any]: + steward = session.get("steward_context") + if isinstance(steward, dict): + if self.conversation_bindings is None: + raise ValueError("bound steward authority is unavailable") + selected = self.conversation_bindings.session_context(steward) + if session.get("goal_id") != "loopx-manager" or session.get("channel_id") != selected["channel_id"]: + raise ValueError("steward audience mismatch") + from ...chat_manager import MANAGER_AGENT_OBJECTIVE + return {"project": Path(selected["context"]["workspace_path"]), "objective": MANAGER_AGENT_OBJECTIVE, + "title": "Steward"} saved = session.get("project_context") if not isinstance(saved, dict) or session.get("goal_id") is not None: raise ValueError("invalid ordinary project Session") @@ -64,3 +74,26 @@ def session_context(self, session: dict[str, Any]) -> dict[str, Any]: return {"project": Path(selected["context"]["workspace_path"]), "objective": PROJECT_CONVERSATION_OBJECTIVE, "title": Path(selected["context"]["workspace_path"]).name} + + def open_bound(self, binding_id: str, source: dict[str, Any], *, executor: str, channel_id: str | None) -> dict[str, Any]: + if self.conversation_bindings is None: + raise ValueError("bound conversation authority is unavailable") + selected = self.conversation_bindings.resolve(binding_id=binding_id, **source) + if selected["binding"]["executor_endpoint_id"] != executor: + raise ValueError("executor does not match the conversation grant") + if channel_id is not None and channel_id != selected["channel_id"]: + raise ValueError("bound conversation channel mismatch") + steward = selected["binding"]["context_kind"] == "steward" + session = {"goal_id": "loopx-manager" if steward else None, "channel_id": selected["channel_id"], + "project_context": None if steward else selected["context"], + "steward_context": selected["context"] if steward else None} + return {**session, **self.session_context(session)} + + @staticmethod + def initialize_bound_scope(store, session): + steward = session.get("steward_context") + if not steward: + return session + from ...chat_manager_context import manager_authorization_scope_id + return store.update_session(session["session_id"], manager_authorization_scope_id=manager_authorization_scope_id( + steward["goal_ids"], runtime_root=store.root.parent, channel_id=session["channel_id"])) diff --git a/loopx/chat_action_normalization.py b/loopx/chat_action_normalization.py index 0a9dcdff85..2d2727d229 100644 --- a/loopx/chat_action_normalization.py +++ b/loopx/chat_action_normalization.py @@ -496,6 +496,7 @@ def _normalize( "heartbeat", "stop_condition", "initial_todos", + "native_token_budget", }, ) goal_id = _opaque(values.get("goal_id"), field="goal_id") @@ -508,6 +509,13 @@ def _normalize( "goal_id": goal_id, "title": _text(values.get("title"), field="title", limit=200), } + if values.get("native_token_budget") is not None: + budget = values["native_token_budget"] + if type(budget) is not int or not 1 <= budget <= 999_999_999: + raise ValueError("native_token_budget must be an explicit positive token allowance") + if values.get("agent_id") != "codex": + raise ValueError("native continuation requires the Codex endpoint") + result["native_token_budget"] = budget for field in ( "objective", "completion_criteria", diff --git a/loopx/chat_actions.py b/loopx/chat_actions.py index caa02ad3c9..1283a57bd6 100644 --- a/loopx/chat_actions.py +++ b/loopx/chat_actions.py @@ -440,6 +440,15 @@ def _project_for_goal_create( if not isinstance(parameters, Mapping) or not isinstance(context, Mapping): raise ValueError("typed Chat action proposal is malformed") workspace_ref = str(parameters.get("workspace_ref") or "current") + if context.get("binding_id"): + # The native steward freezes its configured workspace. Existing + # Goals and the single-Goal fallback cannot redirect a commission. + candidates = [root for root in self.workspace_roots if root.is_dir() and (root / ".git").exists() + and workspace_ref == f"workspace-{hashlib.sha256(str(root).encode('utf-8')).hexdigest()[:12]}"] + if len(candidates) != 1: + raise ValueError("the bound steward workspace is unavailable") + return candidates[0], {"id": "", "repo": str(candidates[0]), "domain": "project-goal-control-plane", + "adapter": {"kind": "generic_project_goal_v0"}} goals = registry_goals(self._registry()) context_goal_id = str(context.get("goal_id") or "").strip() source_goal = next( @@ -750,11 +759,12 @@ def _apply_goal_create( first_turn, created = self.runtime_controller.submit_turn( session_id=session_id, client_turn_id=f"goal-start-{proposal_id}", - message=( + message=(f"/goal start --tokens {parameters['native_token_budget']} {objective}" + if parameters.get("native_token_budget") else ( f"开始推进 Goal {goal_id}。先核对目标边界和现有 Todo," f"首个 Todo:{';'.join(str(item) for item in (parameters.get('initial_todos') or [])[:3]) or '按目标边界建立首个可验证进展'}。" "然后直接推进并报告可验证结果;遇到权限边界时停止并提出明确 Gate。" - ), + )), work_dir=project, objective=objective, ) @@ -1309,7 +1319,8 @@ def regenerate(self, proposal_id: str) -> dict[str, Any]: str(regenerated["proposal_id"]), regenerated_from=proposal_id ) - def apply(self, proposal_id: str) -> dict[str, Any]: + def apply(self, proposal_id: str, *, steward_context: dict[str, Any] | None = None, + steward_confirmed_at: str | None = None) -> dict[str, Any]: proposal = self.store.load(proposal_id) if proposal is None: raise KeyError("typed Chat action proposal was not found") @@ -1318,6 +1329,14 @@ def apply(self, proposal_id: str) -> dict[str, Any]: "proposal": proposal, "turn": self._turn_from_receipt(proposal.get("receipt")), } + if (proposal.get("context") or {}).get("binding_id"): + if steward_context is None: + raise ProtectedActionGate("goal.create", gate={"kind": "authenticated_steward_confirmation_required", + "summary": "Confirm this commission from its original owner private conversation.", + "next_action": "Use the exact /confirm command in the originating Bot before it expires."}) + from .control_plane.effect_runtime import effect_runtime_result + effect_runtime_result("collaboration.steward.authorize_creation", { + "context": steward_context, "proposal": proposal, "now": steward_confirmed_at or now_utc().isoformat()}) if proposal.get("action_kind") == "operation.execute": raise ProtectedActionGate( "operation.execute", diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index 096d33578c..40073b4bf8 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -835,7 +835,11 @@ def _request( try: message = waiter.get_nowait() except queue.Empty: - with self._message_dispatch_lock: + # A streaming reader can route this RPC response while holding + # the fence. Recheck our waiter instead of waiting for an event. + if not self._message_dispatch_lock.acquire(timeout=0.1): + continue + try: try: message = waiter.get_nowait() except queue.Empty: @@ -870,6 +874,8 @@ def _request( continue self._pending_events.put(message) continue + finally: + self._message_dispatch_lock.release() if message.get("id") == request_id: if message.get("error"): if method in { diff --git a/loopx/chat_coordination.py b/loopx/chat_coordination.py index 9c994256a7..909e688d65 100644 --- a/loopx/chat_coordination.py +++ b/loopx/chat_coordination.py @@ -50,6 +50,9 @@ def prepare_turn_context(controller, adapter, session, turn_id, event_sink, *, s # so it gets the bounded read inline. remote_evidence=not isinstance(adapter, CodexAppServerAdapter), ) + if scope.get("bound_steward") is True and isinstance(context.get("bound_steward"), dict): + from .capabilities.native_chat.external_conversations import ChatExternalConversations + context["bound_steward"]["executions"] = ChatExternalConversations(controller).commission_evidence(session) controller.store.append_event(session_id, turn_id, kind="manager.context", payload=context) if scope["kind"] == "external_audience": scope_id = str(context.get("authorization_scope_id") or "") @@ -63,7 +66,11 @@ def prepare_turn_context(controller, adapter, session, turn_id, event_sink, *, s "next_action": "Reconnect the manager to the intended Goal and retry the same message.", }, ) - if session.get("manager_authorization_scope_id") != scope_id: + if scope.get("bound_steward") is True: + # A new, explicitly confirmed commission extends this same owner's + # scope. Refresh tools/evidence without manufacturing a new thread. + controller.store.update_session(session_id, manager_authorization_scope_id=scope_id) + elif session.get("manager_authorization_scope_id") != scope_id: adapter.close_session() with controller.lock: if controller.adapters.get(session_id) is adapter: diff --git a/loopx/chat_manager_context.py b/loopx/chat_manager_context.py index 38e1595b40..78000bede5 100644 --- a/loopx/chat_manager_context.py +++ b/loopx/chat_manager_context.py @@ -425,7 +425,7 @@ def manager_turn_context( } owner_scope = conversation["private_conversation"] scope = conversation["goal_ids"] if owner_scope else authorized_goal_ids - if not owner_scope and not scope: + if not owner_scope and not scope and not (conversation.get("bound_steward") is True and scope == []): return unavailable_manager_context( "external_authorization_unavailable", evidence_window=_evidence_window( @@ -583,6 +583,8 @@ def manager_turn_context( ).hexdigest() if not owner_scope: result["authorization_scope_id"] = manager_authorization_scope_id(scope or [], runtime_root=runtime_root, channel_id=session.get("channel_id")) + if conversation.get("bound_steward") is True: + result["bound_steward"] = {"authorized": True, "goal_count": len(scope or []), "empty": scope == []} return result diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 9936fab716..6c07d4a643 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -570,21 +570,15 @@ def open_session( raise ValueError("mode must be resume_latest or new") selected_channel = channel_id or f"goal.{goal_id}" project_context = None + steward_context = None if conversation_binding_id is not None: if goal_id is not None or agent_goal_id is not None or project_ref is not None: raise ValueError("bound project Chat cannot substitute a Goal or workspace") - bindings = self.project_contexts.conversation_bindings - if bindings is None or source_context is None: - raise ValueError("bound project conversation authority is unavailable") - selected = bindings.resolve(binding_id=conversation_binding_id, **source_context) - if agent_id != selected["binding"]["executor_endpoint_id"]: - raise ValueError("executor does not match the conversation grant") - if channel_id is not None and channel_id != selected["channel_id"]: - raise ValueError("bound conversation channel mismatch") - project_context, selected_channel = selected["context"], selected["channel_id"] - context = self.project_contexts.session_context({ - "goal_id": None, "channel_id": selected_channel, "project_context": project_context, - }) + if source_context is None: + raise ValueError("bound conversation authority is unavailable") + context = self.project_contexts.open_bound(conversation_binding_id, source_context, executor=agent_id, channel_id=channel_id) + goal_id, selected_channel = context["goal_id"], context["channel_id"] + project_context, steward_context = context["project_context"], context["steward_context"] work_dir, objective = context["project"], context["objective"] elif project_ref is not None: if goal_id is not None or agent_goal_id is not None: @@ -679,7 +673,9 @@ def open_session( channel_id=selected_channel, codex_home=str(self.codex_home) if agent_id == "codex" else None, project_context=project_context, + steward_context=steward_context, ) + persisted = self.project_contexts.initialize_bound_scope(self.store, persisted) if is_manager_channel(selected_channel): assert manager_runtime is not None persisted = self.store.update_session( @@ -740,6 +736,8 @@ def _ensure_adapter_locked( if current_session is None or current_session.get("status") == "closed": raise KeyError("chat session was not found") session = current_session + if session.get("steward_context") is not None: + self.project_contexts.session_context(session) if session.get("project_context") is not None: context = self.project_contexts.session_context(session) work_dir, objective = context["project"], context["objective"] @@ -994,6 +992,8 @@ def submit_turn( session = self.store.load_session(session_id) if session is None: raise KeyError("chat session was not found") + if session.get("steward_context") is not None: + raise ValueError("bound steward Chat requires its external source admission") if session.get("project_context") is not None: if session["project_context"].get("audience") == "bound_owner": raise ValueError("bound project Chat requires its external source admission") @@ -1123,6 +1123,8 @@ def steer_active_turn( session = self.store.load_session(session_id) if session is None or session.get("status") == "closed": raise KeyError("chat session was not found") + if session.get("steward_context") is not None: + raise ValueError("bound steward Chat requires its external source admission") if session.get("project_context") is not None: if session["project_context"].get("audience") == "bound_owner": raise ValueError("bound project steering requires its external source admission") @@ -1260,6 +1262,10 @@ def enqueue_turn( session = self.store.load_session(session_id) if session is None or session.get("status") == "closed": raise KeyError("chat session was not found") + if session.get("steward_context") is not None: + if origin != "lark": + raise ValueError("bound steward Chat requires its external source admission") + self.project_contexts.session_context(session) if session.get("project_context") is not None: if conversation_scope(session, origin=origin)["kind"] != "project_workspace": raise ValueError("project grant does not authorize this external audience or source") diff --git a/loopx/chat_server.py b/loopx/chat_server.py index 5a03392c42..0a5551abfe 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -1595,11 +1595,13 @@ def serve_chat( store=server.chat_store, registry_path=resolved_registry_path, project_contexts=ChatProjectContexts(resolved_scan_roots), - manager_scope_resolver=lambda session: authorized_manager_goal_ids( + manager_scope_resolver=lambda session: ( + server.runtime_controller.project_contexts.conversation_bindings.steward_scope(session) + if isinstance(session.get("steward_context"), dict) else authorized_manager_goal_ids( build_lark_goal_topic_runtime_snapshot( registry_path=server.registry_path, runtime_root_override=server.runtime_root_override, ), session, runtime_root=runtime_root, - ), + )), codex_bin=codex_bin, claude_bin=claude_bin, kiro_cli_bin=kiro_cli_bin, @@ -1629,6 +1631,7 @@ def _lark_snapshot(): runtime_controller=server.runtime_controller, workspace_roots=resolved_scan_roots, ) + private_transport.core.actions = server.action_service # An admitted steward team preview is projected into the typed action store, # because that store is what the product surfaces list: the chat action # service owns it, so the channel hands the preview to that owner instead of diff --git a/loopx/chat_store.py b/loopx/chat_store.py index a87234f4da..93f5f01223 100644 --- a/loopx/chat_store.py +++ b/loopx/chat_store.py @@ -246,6 +246,7 @@ def create_session( attached_capabilities: dict[str, bool] | None = None, codex_home: str | None = None, project_context: dict[str, Any] | None = None, + steward_context: dict[str, Any] | None = None, ) -> dict[str, Any]: now = utc_now() token = _opaque_id(session_id or uuid.uuid4().hex, field="session_id") @@ -277,6 +278,13 @@ def create_session( if goal_id is not None or goal_instance_id is not None or channel_id != selected["channel_id"] or normalized_mode != CHAT_SESSION_MODE_MANAGED: raise ValueError("ordinary project Sessions require their exact channel and no Goal") project_context = selected["context"] + if steward_context is not None: + from .control_plane.effect_runtime import effect_runtime_result + selected = effect_runtime_result("collaboration.steward.session_identity", {"context": steward_context}) + if (project_context is not None or goal_id != "loopx-manager" or goal_instance_id is not None + or channel_id != selected["channel_id"] or normalized_mode != CHAT_SESSION_MODE_MANAGED): + raise ValueError("bound steward Sessions require their exact role and audience") + steward_context = selected["context"] normalized_goal_id = _opaque_id(goal_id, field="goal_id") if goal_id is not None else None if normalized_goal_id is None and project_context is None: raise ValueError("goal_id is required outside ordinary project Sessions") @@ -300,6 +308,7 @@ def create_session( "session_id": token, "goal_id": normalized_goal_id, **({"project_context": project_context} if project_context is not None else {}), + **({"steward_context": steward_context} if steward_context is not None else {}), **( { "goal_instance_id": _opaque_id( diff --git a/loopx/control_plane/collaboration/__init__.py b/loopx/control_plane/collaboration/__init__.py index 7c0dbed96f..b471b9d8ce 100644 --- a/loopx/control_plane/collaboration/__init__.py +++ b/loopx/control_plane/collaboration/__init__.py @@ -19,7 +19,7 @@ def conversation_trigger(mode: str | None = None, **evidence: bool) -> dict[str, def conversation_scope(session: dict[str, Any], *, origin: str | None = None) -> dict[str, Any]: return effect_runtime_result("collaboration.conversation.scope", { "channel_id": session.get("channel_id"), "goal_id": session.get("goal_id"), - "project_context": session.get("project_context"), + "project_context": session.get("project_context"), "steward_context": session.get("steward_context"), **({"origin": origin} if origin is not None else {}), }) diff --git a/loopx/control_plane/collaboration/conversation_binding.ts b/loopx/control_plane/collaboration/conversation_binding.ts index 60e7333439..bd6f11e2ce 100644 --- a/loopx/control_plane/collaboration/conversation_binding.ts +++ b/loopx/control_plane/collaboration/conversation_binding.ts @@ -1,11 +1,12 @@ import type {JsonObject} from "../effect_program.ts"; import {EffectRuntimeRequestError} from "../effect_runtime_errors.ts"; import {requireJsonObject, requireNonEmptyString} from "../runtime_decode.ts"; -import {normalizeProjectContext} from "./conversation_scope.ts"; +import {normalizeProjectContext, normalizeStewardContext, normalizeStewardGoalScope} from "./conversation_scope.ts"; /** Core owns the context and audience grant. A provider supplies verified, * opaque identity observations; neither a message nor a model selects them. - * This contract intentionally carries no Goal and no provider credentials. + * Ordinary project bindings carry no Goal. A steward adopts only newly + * confirmed creation receipts, never an existing or global portfolio. */ const BINDING_SCHEMA = "loopx_chat_conversation_binding_v0"; const SET_SCHEMA = "loopx_chat_conversation_bindings_v0"; @@ -24,8 +25,8 @@ function token(value: unknown, label: string): string { function binding(value: unknown): JsonObject { const row = requireJsonObject(value, "conversation binding"); - if (row.schema_version !== BINDING_SCHEMA || row.context_kind !== "project" - || row.grant !== "workspace_read" || row.enabled !== true) { + if (row.schema_version !== BINDING_SCHEMA || !["project", "steward"].includes(String(row.context_kind)) + || row.grant !== (row.context_kind === "project" ? "workspace_read" : "portfolio_read") || row.enabled !== true) { throw new EffectRuntimeRequestError("unsupported conversation binding"); } return { @@ -33,9 +34,10 @@ function binding(value: unknown): JsonObject { transport_ref: token(row.transport_ref, "transport reference"), provider_ref: ref(row.provider_ref, "provider identity"), operator_ref: ref(row.operator_ref, "verified operator identity"), - context_kind: "project", project_ref: ref(row.project_ref, "workspace reference"), + context_kind: row.context_kind, project_ref: ref(row.project_ref, "workspace reference"), executor_endpoint_id: token(row.executor_endpoint_id, "executor endpoint"), - grant: "workspace_read", enabled: true, + grant: row.grant, enabled: true, + ...(row.context_kind === "steward" ? {goal_ids: normalizeStewardGoalScope(row.goal_ids)} : {}), }; } @@ -81,6 +83,28 @@ export function planConversationBinding(params: JsonObject): JsonObject { } if (previous?.binding_id === candidate.binding_id) throw new EffectRuntimeRequestError("changed context requires a new binding identity"); rows = [...current.bindings.filter(row => row.transport_ref !== candidate.transport_ref), candidate]; + } else if (params.operation === "adopt_created_goal") { + const id = ref(params.binding_id, "binding identity"); + const previous = current.bindings.find(row => row.binding_id === id); + const context = normalizeStewardContext(params.context); + const proposal = requireJsonObject(params.proposal, "creation proposal"); + const principal = requireJsonObject(proposal.context, "creation audience"); + const receipt = requireJsonObject(proposal.receipt, "creation receipt"); + const resources = requireJsonObject(receipt.resource_ids, "creation resources"); + const goal = requireJsonObject(params.goal, "created Goal observation"); + if (!previous || previous.context_kind !== "steward" || previous.project_ref !== context.project_ref + || proposal.action_kind !== "goal.create" || proposal.status !== "applied" + || receipt.outcome !== "goal_created" || receipt.projection_verified !== true + || goal.goal_id !== resources.goal_id || goal.creation_operation_id !== proposal.proposal_id + || goal.workspace_path !== context.workspace_path) throw new EffectRuntimeRequestError("steward adoption requires its exact creation receipt"); + for (const key of ["binding_id", "source_ref", "provider_ref", "operator_ref"]) { + if (principal[key] !== context[key] || (key !== "source_ref" && previous[key] !== context[key])) { + throw new EffectRuntimeRequestError("creation receipt belongs to another audience"); + } + } + const goals = normalizeStewardGoalScope([...(previous.goal_ids as string[]), goal.goal_id]); + if (JSON.stringify(goals) === JSON.stringify(previous.goal_ids)) return {changed: false, state: current}; + rows = current.bindings.map(row => row.binding_id === id ? {...row, goal_ids: goals} : row); } else if (params.operation === "disconnect") { const id = ref(params.binding_id, "binding identity"); rows = current.bindings.filter(row => row.binding_id !== id); @@ -104,11 +128,20 @@ export function resolveBoundConversation(params: JsonObject): JsonObject { const projects = params.available_projects.map(normalizeProjectContext).filter(project => project.project_ref === row.project_ref); if (projects.length !== 1) throw new EffectRuntimeRequestError("workspace grant is unavailable or ambiguous"); const context = {...projects[0], audience: "bound_owner", binding_id: id, source_ref: source, - provider_ref: row.provider_ref, operator_ref: row.operator_ref}; - if (params.session_context !== undefined && JSON.stringify(params.session_context) !== JSON.stringify(context)) { - throw new EffectRuntimeRequestError("bound Session context changed"); + provider_ref: row.provider_ref, operator_ref: row.operator_ref, + ...(row.context_kind === "steward" ? {kind: "bound_steward", grant: "portfolio_read", goal_ids: row.goal_ids} : {})}; + if (params.session_context !== undefined) { + const saved = requireJsonObject(params.session_context, "bound Session context"); + // New commissions may extend the same owner's scope. Workspace, role and + // audience remain frozen; all evidence reads use the fresh binding scope. + const matches = row.context_kind === "steward" + ? JSON.stringify({...saved, goal_ids: []}) === JSON.stringify({...context, goal_ids: []}) + && normalizeStewardGoalScope(saved.goal_ids).every(g => (row.goal_ids as string[]).includes(g)) + : JSON.stringify(saved) === JSON.stringify(context); + if (!matches) throw new EffectRuntimeRequestError("bound Session context changed"); } - return {binding: row, context, channel_id: `project.external.${id}.${source}`}; + return {binding: row, context, channel_id: row.context_kind === "steward" + ? `manager.external.native.${id}.${source}` : `project.external.${id}.${source}`}; } /** Shared native command projection. Provider grammar carries the explicit @@ -118,13 +151,18 @@ 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"].includes(command as null | string)) { + if (![null, "status", "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"); 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 (["commission", "confirm_commission", "cancel_commission", "stop_commission", "resume_commission"].includes(String(command))) { + const selected = requireJsonObject(params.binding, "steward binding"); + 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") { return {operation: "reply", session_id: session, turn_id: null, response_code: command === "unsupported" ? "unsupported_attachment" : !current ? "no_session" @@ -136,3 +174,48 @@ export function planBoundConversationRequest(params: JsonObject): JsonObject { } return {operation: "admit_turn", client_turn_id: `external-${request}`, session_id: session, turn_id: null}; } + +/** 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 { + const command = params.command; + const message = requireNonEmptyString(params.message, "steward command").trim(); + const patterns: Record = {commission: /^\/(?:delegate|委托)\s+--tokens\s+([1-9][0-9]{0,8})\s+([\s\S]+)$/, + confirm_commission: /^\/confirm\s+(proposal-[a-f0-9]{32})$/, + cancel_commission: /^\/cancel\s+(proposal-[a-f0-9]{32})$/, + stop_commission: /^\/stop-commission\s+(proposal-[a-f0-9]{32})$/, + resume_commission: /^\/resume-commission\s+(proposal-[a-f0-9]{32})\s+--tokens\s+([1-9][0-9]{0,8})$/}; + const match = patterns[String(command)]?.exec(message); + const argument = command === "commission" ? match?.[2] : match?.[1]; + if (!match || !argument?.trim() || Array.from(argument).length > 1000) { + throw new EffectRuntimeRequestError("invalid explicit steward command"); + } + return {argument: argument.trim(), ...(command === "commission" ? {native_token_budget: Number(match[1])} + : command === "resume_commission" ? {native_token_budget: Number(match[2])} : {})}; +} + +export function authorizeStewardCreation(params: JsonObject): JsonObject { + const selected = normalizeStewardContext(params.context); + const proposal = requireJsonObject(params.proposal, "creation proposal"); + const audience = requireJsonObject(proposal.context, "creation audience"); + if (proposal.action_kind !== "goal.create" || audience.kind !== "manager" + || audience.goal_id !== "loopx-manager") throw new EffectRuntimeRequestError("not a steward creation preview"); + for (const field of ["binding_id", "source_ref", "provider_ref", "operator_ref", "project_ref"]) { + if (audience[field] !== selected[field]) throw new EffectRuntimeRequestError("confirmation audience changed"); + } + if (["stop", "resume"].includes(String(params.operation))) { + const receipt = requireJsonObject(proposal.receipt, "commission receipt"); + const resources = requireJsonObject(receipt.resource_ids, "commission resources"); + if (proposal.status !== "applied" || !(selected.goal_ids as string[]).includes(String(resources.goal_id))) { + throw new EffectRuntimeRequestError("execution control requires the exact adopted commission"); + } + return {authorized: true}; + } + const expiry = Date.parse(String(audience.expires_at)); + const now = Date.parse(String(params.now)); + if (!Number.isFinite(expiry) || !Number.isFinite(now) || now > expiry + || now < Date.parse(String(proposal.created_at)) - 300_000) { + throw new EffectRuntimeRequestError("steward confirmation expired or timestamp invalid"); + } + return {authorized: true}; +} diff --git a/loopx/control_plane/collaboration/conversation_scope.ts b/loopx/control_plane/collaboration/conversation_scope.ts index ccae68cd5a..00591ff04b 100644 --- a/loopx/control_plane/collaboration/conversation_scope.ts +++ b/loopx/control_plane/collaboration/conversation_scope.ts @@ -28,6 +28,30 @@ export function projectConversationIdentity(input: Record): Rec ? `project.${context.project_ref}` : `project.external.${context.binding_id}.${context.source_ref}`}; } +export function normalizeStewardGoalScope(value: unknown): string[] { + if (!Array.isArray(value) || value.length > 128 + || value.some(g => typeof g !== "string" || !/^[A-Za-z0-9._-]{1,160}$/.test(g) + || [".", "..", "loopx-manager"].includes(g))) throw new Error("invalid steward Goal scope"); + return [...new Set(value as string[])].sort(); +} + +/** An explicitly selected steward has a bounded portfolio, including an honest + * empty one. Identity is persisted by Core; message text cannot select it. + */ +export function normalizeStewardContext(value: unknown): Record { + if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("steward context unavailable"); + const row = value as Record; + const workspace = normalizeProjectContext({...row, kind: "project_workspace", grant: "workspace_read"}); + if (row.kind !== "bound_steward" || row.audience !== "bound_owner" || row.grant !== "portfolio_read" + ) throw new Error("invalid bounded steward context"); + return {...workspace, kind: "bound_steward", grant: "portfolio_read", goal_ids: normalizeStewardGoalScope(row.goal_ids)}; +} + +export function stewardConversationIdentity(input: Record): Record { + const context = normalizeStewardContext(input.context); + return {context, channel_id: `manager.external.native.${context.binding_id}.${context.source_ref}`}; +} + type ConversationScope = Record & ( | {kind: "owner_portfolio"; goal_ids: null; private_conversation: true} | {kind: "owner_goal"; goal_ids: [string]; private_conversation: true} @@ -54,6 +78,14 @@ export function resolveConversationScope(input: Record): Conver } } catch { /* Incomplete host identity grants no context. */ } } + if (goal === "loopx-manager" && (input.origin === undefined || input.origin === "lark")) { + try { + const selected = stewardConversationIdentity({context: input.steward_context}); + if (channel === selected.channel_id) { + return {kind: "external_audience", goal_ids: [], private_conversation: false, bound_steward: true}; + } + } catch { /* An incomplete steward identity cannot authorize a portfolio. */ } + } if (channel === "manager") { return {kind: "owner_portfolio", goal_ids: null, private_conversation: true}; } diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index f5ec3b1a4a..f15ed06755 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -578,6 +578,9 @@ export function createEffectRuntimeHandlers( ["collaboration.conversation.scope", lazyHandler(() => import("./collaboration/conversation_scope.ts"), ({resolveConversationScope}) => resolveConversationScope)], ["collaboration.project.context", lazyHandler(() => import("./collaboration/project_conversation.ts"), ({resolveProjectConversation}) => resolveProjectConversation)], ["collaboration.project.session_identity", lazyHandler(() => import("./collaboration/conversation_scope.ts"), ({projectConversationIdentity}) => projectConversationIdentity)], + ["collaboration.steward.session_identity", lazyHandler(() => import("./collaboration/conversation_scope.ts"), ({stewardConversationIdentity}) => stewardConversationIdentity)], + ["collaboration.steward.command", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({stewardCommand}) => stewardCommand)], + ["collaboration.steward.authorize_creation", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({authorizeStewardCreation}) => authorizeStewardCreation)], ["collaboration.conversation.binding", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({planConversationBinding}) => planConversationBinding)], ["collaboration.conversation.bound_context", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({resolveBoundConversation}) => resolveBoundConversation)], ["collaboration.conversation.request", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({planBoundConversationRequest}) => planBoundConversationRequest)], diff --git a/loopx/extensions/lark/private_conversation_api.py b/loopx/extensions/lark/private_conversation_api.py index 66f0b7c0b8..ed494884f1 100644 --- a/loopx/extensions/lark/private_conversation_api.py +++ b/loopx/extensions/lark/private_conversation_api.py @@ -20,21 +20,22 @@ def _private_conversations(self) -> None: {"binding_id": row["binding_id"], "app_ref": row["transport_ref"], "context_kind": row["context_kind"], "project_ref": row["project_ref"], "context_available": row["project_ref"] in titles, "project_title": titles.get(row["project_ref"], "Unavailable workspace"), "executor_endpoint_id": row["executor_endpoint_id"], "grant": row["grant"], + "goal_count": len(row.get("goal_ids", [])), "listener_status": health.get(row["transport_ref"], {}).get("status", "starting"), **deliveries.get(row["binding_id"], {"pending_count": 0, "recovery_count": 0})} for row in current["bindings"]]}) def _private_conversation_connect(self) -> None: - from ...extensions.lark.goal_topic_runtime import _active_profile_configs + from .goal_topic_runtime import _active_profile_configs from ...chat_lark_api import build_lark_goal_topic_runtime_snapshot try: body = self._read_json() - if set(body) != {"app_ref", "project_ref", "executor_endpoint_id"}: + if set(body) - {"context_kind"} != {"app_ref", "project_ref", "executor_endpoint_id"}: raise ValueError("select an App, authorized workspace and executor") profile = str(body["app_ref"]) existing = _active_profile_configs(build_lark_goal_topic_runtime_snapshot( registry_path=self.server.registry_path, runtime_root_override=self.server.runtime_root_override)) - from ...extensions.lark.conversation_identity import identity_ref + from .conversation_identity import identity_ref observed = self.server.runtime_controller.project_contexts.conversation_bindings.observe(profile) from ...chat_lark_api import _app_identity_for_private_guard group_apps = [str(config.get("bot_app_id") or _app_identity_for_private_guard( @@ -46,7 +47,8 @@ def _private_conversation_connect(self) -> None: if not any(row["agent_id"] == endpoint and row["available"] for row in self.server.runtime_controller.capabilities()): raise ValueError("the selected executor is unavailable") self.server.runtime_controller.project_contexts.conversation_bindings.configure( - transport_ref=profile, project_ref=str(body["project_ref"]), executor_endpoint_id=endpoint) + transport_ref=profile, project_ref=str(body["project_ref"]), executor_endpoint_id=endpoint, + context_kind=str(body.get("context_kind", "project"))) self.server.lark_goal_topic_runtime.refresh() except (ValueError, OSError, KeyError) as exc: self._send_error(str(exc), status=400) diff --git a/loopx/extensions/lark/private_conversations.py b/loopx/extensions/lark/private_conversations.py index ee3912553b..b51225569a 100644 --- a/loopx/extensions/lark/private_conversations.py +++ b/loopx/extensions/lark/private_conversations.py @@ -150,11 +150,22 @@ def admit(self, profile: str, event: dict[str, Any]) -> dict[str, Any]: if not text.strip(): return {"status": "empty_text"} command = {"/status": "status", "/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"), + ("/stop-commission", "stop_commission"), ("/resume-commission", "resume_commission")]: + if text.strip() == prefix or text.strip().startswith(prefix + " "): + command = selected_command + break if message_type != "text": command = "unsupported" try: admitted = self.core.admit(binding_id=binding["binding_id"], source=source, request_ref=request, message=text, command=command) + except ValueError: + record.update(status="rejected", response="操作格式不正确;新委托请使用 /delegate --tokens N 具体目标,确认或取消请使用原预览中的完整命令。") + _atomic_write_json(path, record) + return {"status": "command_rejected"} except RuntimeError as exc: if str(exc) != "session_queue_full": raise @@ -224,6 +235,13 @@ def reconcile(self) -> int: continue try: self.bindings.resolve(binding_id=record["binding_id"], **record["source"]) + if record["status"] in {"command_queued", "command_completed", "commission_running"}: + native = self.core.read_request(record["request_ref"]) + if native["status"] == "command_queued": + self._deliver(path, record, "admission", native["response"]) + continue + record.update(status=native["status"], response=native.get("response"), + commission_resources=native.get("commission_resources")) if record["status"] == "accepted": # Receipt follows persistent Core admission and is # independent of terminal execution and reply delivery. @@ -237,6 +255,20 @@ def reconcile(self) -> int: else: response = _command_text(str(record.get("response_code") or "")) or str(record.get("response") or "") if self._deliver(path, record, "terminal", response): + resources = record.get("commission_resources") or {} + if resources.get("session_id") and resources.get("turn_id"): + first_turn = self.core.controller.store.load_turn(resources["session_id"], resources["turn_id"]) + if not first_turn or first_turn["status"] not in {"completed", "failed", "interrupted", "expired"}: + record["status"] = "commission_running" + _atomic_write_json(path, record) + continue + result_text = str((first_turn.get("response") or {}).get("message") or "") + result_text = "委托执行结果:\n" + result_text if first_turn["status"] == "completed" else "委托首轮执行未完成;原 Goal 和回执已保留,请查看状态后决定恢复。" + proposal_id = native.get("proposal_id") + if proposal_id: + result_text += f"\n如需恢复暂停或额度受限的原执行:/resume-commission {proposal_id} --tokens N(N 为包含历史用量的总上限,须大于已用量;不会重开线程)。" + if not self._deliver(path, record, "commission_result", result_text): + continue config = self._inbox(record) acknowledge_lark_event_inbox(project=self.runtime_root, config_path=config, message_ids=[record["event"]["message_id"]], execute=True) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index a3fecd352d..35626f1f20 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -455,7 +455,7 @@ }, { "site": "loopx/chat_server.py::.serve_chat._wake_goal_context::codec_read:load_registry#1", - "line": 1658, + "line": 1661, "column": 20, "kind": "codec_read", "api": "load_registry", diff --git a/tests/control_plane_ts/conversation_binding.test.ts b/tests/control_plane_ts/conversation_binding.test.ts index fe8f522896..16a46f880f 100644 --- a/tests/control_plane_ts/conversation_binding.test.ts +++ b/tests/control_plane_ts/conversation_binding.test.ts @@ -1,6 +1,7 @@ import assert from "node:assert/strict"; import test from "node:test"; -import {planConversationBinding, resolveBoundConversation, planBoundConversationRequest} from "../../loopx/control_plane/collaboration/conversation_binding.ts"; +import {planConversationBinding, resolveBoundConversation, planBoundConversationRequest, + stewardCommand, authorizeStewardCreation} from "../../loopx/control_plane/collaboration/conversation_binding.ts"; import {resolveConversationScope} from "../../loopx/control_plane/collaboration/conversation_scope.ts"; const project = {kind: "project_workspace", project_ref: "a".repeat(24), workspace_path: "/authorized/notes", @@ -32,6 +33,49 @@ test("binding independently verifies the owner and does not create a Goal", () = assert.throws(() => planConversationBinding({...request, expected_revision: 1}), /revision/); }); +test("a selected steward has a verified empty portfolio and adopts only its own exact new creation", () => { + const steward = {...row, context_kind: "steward", grant: "portfolio_read", goal_ids: []}; + const next = planConversationBinding({...request, binding: steward}).state as typeof current; + const use = {current: next, binding_id: row.binding_id, source_ref: "e".repeat(24), + sender_ref: row.operator_ref, private_human_message: true, observation, available_projects: [project]}; + const selected = resolveBoundConversation(use); + assert.equal(resolveConversationScope({channel_id: selected.channel_id, goal_id: "loopx-manager", + steward_context: selected.context, origin: "lark"}).bound_steward, true); + assert.equal(resolveConversationScope({channel_id: selected.channel_id, goal_id: "loopx-manager", + steward_context: selected.context, origin: "web"}).bound_steward, undefined); + const context = selected.context as Record; + const proposal = {proposal_id: "proposal-" + "f".repeat(32), status: "applied", action_kind: "goal.create", + context, receipt: {outcome: "goal_created", projection_verified: true, resource_ids: {goal_id: "fresh-goal"}}}; + const adoption = {current: next, expected_revision: 1, operation: "adopt_created_goal", binding_id: row.binding_id, + context, proposal, goal: {goal_id: "fresh-goal", workspace_path: project.workspace_path, creation_operation_id: proposal.proposal_id}}; + const result = planConversationBinding(adoption).state as typeof current; + assert.deepEqual((result.bindings[0] as Record).goal_ids, ["fresh-goal"]); + assert.deepEqual((resolveBoundConversation({...use, current: result, session_context: context}).context as Record).goal_ids, ["fresh-goal"]); + for (const bad of [{...adoption, goal: {...adoption.goal, workspace_path: "/other/private"}}, + {...adoption, proposal: {...proposal, status: "preview_ready"}}, + {...adoption, proposal: {...proposal, context: {...context, source_ref: "f".repeat(24)}}}, + {...adoption, proposal: {...proposal, receipt: {...proposal.receipt, projection_verified: false}}}]) { + assert.throws(() => planConversationBinding(bad)); + } +}); + +test("a commission requires an explicit budget and confirmation stays within its audience and expiry", () => { + assert.deepEqual(stewardCommand({command: "commission", message: "/delegate --tokens 12000 Read README"}), + {argument: "Read README", native_token_budget: 12000}); + for (const message of ["delegate work", "/delegate Read README", "/delegate --tokens 0 work", "/delegate --tokens 1000"]) + assert.throws(() => stewardCommand({command: "commission", message})); + const context = {...project, kind: "bound_steward", audience: "bound_owner", grant: "portfolio_read", goal_ids: [], + binding_id: row.binding_id, source_ref: "e".repeat(24), provider_ref: row.provider_ref, operator_ref: row.operator_ref}; + const proposal = {action_kind: "goal.create", created_at: "2026-01-01T10:00:00Z", + context: {...context, kind: "manager", goal_id: "loopx-manager", expires_at: "2026-01-01T10:15:00Z"}}; + assert.equal(authorizeStewardCreation({context, proposal, now: "2026-01-01T10:01:00Z"}).authorized, true); + assert.throws(() => authorizeStewardCreation({context, proposal, now: "2026-01-01T10:16:00Z"}), /expired/); + assert.throws(() => authorizeStewardCreation({context: {...context, provider_ref: "f".repeat(24)}, proposal, + now: "2026-01-01T10:01:00Z"}), /audience/); + assert.throws(() => planBoundConversationRequest({binding: row, current_session: null, + request: {request_ref: "e".repeat(24), command: "commission"}}), /selected steward/); +}); + test("one App has one binding owner, and a context change cannot move an existing Session", () => { const next = planConversationBinding(request).state as typeof current; const second = {...row, transport_ref: "other-app", binding_id: "e".repeat(24)}; diff --git a/tests/test_chat_agent.py b/tests/test_chat_agent.py index b918a7182b..9121289d75 100644 --- a/tests/test_chat_agent.py +++ b/tests/test_chat_agent.py @@ -3,6 +3,9 @@ import io import json import queue +import threading +import time +from types import SimpleNamespace from pathlib import Path import loopx.chat_agent as chat_agent @@ -11,6 +14,34 @@ import pytest +def test_silent_event_reader_releases_dispatch_fence_for_control_receipt(tmp_path, monkeypatch): + session = chat_agent.CodexChatAgentSession(process=SimpleNamespace(poll=lambda: None), + messages=queue.Queue(), thread_id="synthetic-thread", work_dir=tmp_path) + arrived, errors = [], [] + + def events(): + try: + arrived.append(session._next_event(deadline=time.monotonic() + 3)) + except Exception as exc: + errors.append(exc) + + # The response arrives without an agent event, just like a silent native + # Goal stop/read. It must reach its waiter before the event stream ends. + monkeypatch.setattr(session, "_write", lambda packet: session.messages.put( + {"id": packet["id"], "result": {"goal": {"status": "active"}}})) + reader = threading.Thread(target=events) + reader.start() + try: + time.sleep(.05) + started = time.monotonic() + assert session._request("thread/goal/get", {"threadId": session.thread_id})["goal"]["status"] == "active" + assert time.monotonic() - started < 1 + finally: + session.messages.put({"method": "turn/completed", "params": {"turnId": "original"}}) + reader.join(timeout=4) + assert not errors and arrived == [{"method": "turn/completed", "params": {"turnId": "original"}}] + + class _FakeAppServerProcess: def __init__(self) -> None: responses = [ diff --git a/tests/test_native_steward_private.py b/tests/test_native_steward_private.py new file mode 100644 index 0000000000..87905f56ea --- /dev/null +++ b/tests/test_native_steward_private.py @@ -0,0 +1,315 @@ +"""Native steward journeys exercise the same Core, action store and provider. + +Synthetic transport/model fixtures never certify phone or real-model acceptance. +""" +import json +import time +from pathlib import Path +import runpy + +import pytest + +from test_lark_private_conversations import connect +from loopx.chat_action_store import ChatActionStore +from loopx.chat_actions import ChatActionService, ProtectedActionGate +from loopx.chat_manager_context import collect_manager_turn_context, manager_turn_context +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_store import ChatSessionStore +from loopx.capabilities.native_chat.project_context import ChatProjectContexts +from loopx.extensions.lark.conversation_identity import identity_ref + + +@pytest.fixture +def steward(tmp_path): + workspace = tmp_path / "fresh-workspace" + workspace.mkdir() + contexts = ChatProjectContexts([workspace]) + store = ChatSessionStore(tmp_path / "runtime") + capture, fake = tmp_path / "requests.jsonl", tmp_path / "codex" + source = runpy.run_path(str(Path(__file__).parents[1] / "examples/loopx-chat-runtime-smoke.py"))["FAKE_CODEX"] + source = source.replace(' method = request.get("method")', + f' with open({str(capture)!r}, "a") as output:\n output.write(json.dumps(request) + "\\n")\n' + ' method = request.get("method")') + fake.write_text(source) + fake.chmod(0o700) + runtime = ChatRuntimeController(store=store, codex_bin=str(fake), project_contexts=contexts, + registry_path=tmp_path / "registry.json") + runtime.registry_path.write_text(json.dumps({"schema_version": "0.1", "runtime_root": str(store.root.parent), "goals": []})) + (workspace / ".git").mkdir() + source = fake.read_text().replace('active_turn = None', 'active_turn = None\nnative_goal = None', 1) + begin = source.index(' elif method in {"thread/goal/set", "thread/goal/get"}:') + end = source.index(' elif method == "turn/start":', begin) + source = source[:begin] + source[end:] + source = source.replace(' elif method == "turn/start":', ''' elif method == "thread/goal/get": + result = {"goal": native_goal} + elif method == "thread/goal/set": + native_goal = {**(native_goal or {}), "threadId": "durable-thread", "tokensUsed": 120, "timeUsedSeconds": 1, **request["params"]} + print(json.dumps({"id": request_id, "result": {"goal": native_goal}}), flush=True) + if native_goal["status"] == "active": + turn = "native-first" + active_turn = turn + print(json.dumps({"method": "turn/started", "params": {"threadId": "durable-thread", "turn": {"id": turn}}}), flush=True) + if "wait for interrupt" in str(native_goal.get("objective")): + continue + print(json.dumps({"method": "item/agentMessage/delta", "params": {"threadId": "durable-thread", "turnId": turn, "delta": "Verified synthetic result: the workspace has one README and no inherited portfolio."}}), flush=True) + native_goal["status"] = "complete" + print(json.dumps({"method": "turn/completed", "params": {"threadId": "durable-thread", "turn": {"id": turn, "status": "completed"}}}), flush=True) + continue + elif method == "turn/start":''') + fake.write_text(source) + store, runtime, provider, transport = connect((store, runtime, contexts, None, capture, fake, workspace)) + ref = contexts.available()[0]["project_ref"] + binding = transport.bindings.configure(transport_ref="steward-app", project_ref=ref, + executor_endpoint_id="codex", context_kind="steward") + runtime.manager_scope_resolver = transport.bindings.steward_scope + transport.core.actions = ChatActionService(store=ChatActionStore(store.root / "actions"), + registry_path=runtime.registry_path, chat_store=store, runtime_controller=runtime, workspace_roots=[workspace]) + yield store, runtime, provider, transport, binding, capture, workspace + runtime.close() + + +def finish(runtime, row): + return runtime.wait_for_turn(session_id=row["session_id"], turn_id=row["turn_id"], timeout_sec=15) + + +def test_verified_empty_steward_is_not_missing_authorization_and_does_not_inherit(steward): + store, runtime, provider, transport, binding, _, _ = steward + transport.admit("steward-app", provider.event("steward-app", "empty", "What new work do you manage?")) + row = transport.core.pending()[0] + assert finish(runtime, row)["status"] == "completed" + session = store.load_session(row["session_id"]) + context = collect_manager_turn_context(runtime.registry_path, session, store.root.parent, runtime.manager_scope_resolver) + assert context["goals"] == [] and context["bound_steward"] == {"authorized": True, "goal_count": 0, "empty": True} + assert context["coverage"]["discovered"] == 0 and "external_authorization_unavailable" not in context["warnings"] + assert "empty_inventory" in context["warnings"] + assert json.loads(runtime.registry_path.read_text())["goals"] == [] + missing = manager_turn_context(runtime.registry_path, {"goal_id": "loopx-manager", "channel_id": "manager.external.unbound"}, store.root.parent, authorized_goal_ids=[]) + assert missing["warnings"] == ["external_authorization_unavailable"] + transport.admit("notes-app", provider.event("notes-app", "notes", "ordinary notes conversation")) + notes = next(r for r in transport.core.pending() if r["message"] == "ordinary notes conversation") + assert finish(runtime, notes)["status"] == "completed" + assert store.load_session(notes["session_id"])["goal_id"] is None + assert notes["session_id"] != row["session_id"] + assert "ordinary notes conversation" not in str(store.messages(row["session_id"])) + transport.bindings.disconnect(binding["binding_id"], expected_revision=transport.bindings.read()["revision"]) + assert collect_manager_turn_context(runtime.registry_path, session, store.root.parent, runtime.manager_scope_resolver)["warnings"] == ["external_authorization_unavailable"] + assert transport.admit("steward-app", provider.event("steward-app", "revoked", "more"))["status"] == "audience_rejected" + + +def test_confirmed_commission_runs_native_goal_returns_result_and_extends_same_session(steward): + store, runtime, provider, transport, binding, capture, _ = steward + transport.admit("steward-app", provider.event("steward-app", "before", "Check the empty portfolio")) + initial = transport.core.pending()[0] + assert finish(runtime, initial)["status"] == "completed" + original = store.load_session(initial["session_id"])["upstream_thread_id"] + event = provider.event("steward-app", "delegate", "/delegate --tokens 12000 Inspect the authorized README and report findings") + started = time.monotonic() + assert transport.admit("steward-app", event)["status"] == "command_recorded" + assert time.monotonic() - started < 3 + assert json.loads(runtime.registry_path.read_text())["goals"] == [] + transport.reconcile() + preview = transport.core.actions.store.list()[0] + assert preview["status"] == "preview_ready" and preview["normalized_parameters"]["heartbeat"] == {"enabled": False} + with pytest.raises(ProtectedActionGate): + transport.core.actions.apply(preview["proposal_id"]) + # The other App cannot consume even a known, exact proposal id. + text = "/confirm " + preview["proposal_id"] + transport.admit("notes-app", provider.event("notes-app", "wrong_app", text)) + wrong_app = next(r for r in transport.core.pending() if r["message"] == text) + assert finish(runtime, wrong_app)["status"] == "completed" + assert transport.core.actions.load(preview["proposal_id"])["status"] == "preview_ready" + confirm = provider.event("steward-app", "confirm", text) + transport.admit("steward-app", confirm) + transport.reconcile() + applied = transport.core.actions.load(preview["proposal_id"]) + assert applied["status"] == "applied", transport.core.pending() + resources = applied["receipt"]["resource_ids"] + turn = finish(runtime, resources) + assert turn["status"] == "completed", turn + assert "Verified synthetic result" in turn["response"]["message"] + assert "Codex Goal: complete" in turn["response"]["message"] + transport.reconcile() + assert any(profile == "steward-app" and "Verified synthetic result" in text for profile, text in provider.writes) + assert not any(profile == "notes-app" and "Verified synthetic result" in text for profile, text in provider.writes) + count = len(provider.writes) + transport.admit("steward-app", confirm) + transport.reconcile() + assert len(provider.writes) == count + assert len(json.loads(runtime.registry_path.read_text())["goals"]) == 1 + assert transport.bindings.read()["bindings"][1]["goal_ids"] == [resources["goal_id"]] + transport.admit("steward-app", provider.event("steward-app", "after", "Report the new commission status")) + after = next(r for r in transport.core.pending() if r["message"] == "Report the new commission status") + assert finish(runtime, after)["status"] == "completed" + assert after["session_id"] == initial["session_id"] + assert store.load_session(after["session_id"])["upstream_thread_id"] == original + context = collect_manager_turn_context(runtime.registry_path, store.load_session(after["session_id"]), store.root.parent, runtime.manager_scope_resolver) + assert [g["goal_id"] for g in context["goals"]] == [resources["goal_id"]] + facts = transport.core.commission_evidence(store.load_session(after["session_id"])) + assert facts[0]["native_execution"]["status"] == "complete" + assert facts[0]["result_delivery_verified"] is True + assert facts[0]["canonical_acceptance_attested"] is False + requests = [json.loads(line) for line in capture.read_text().splitlines()] + assert len([r for r in requests if r.get("method") == "thread/start"]) == 3 + assert any(r.get("method") == "thread/goal/set" for r in requests) + from loopx.chat_runtime import ChatRuntimeController + from loopx.chat_store import ChatSessionStore + runtime.close() + restarted = ChatRuntimeController(store=ChatSessionStore(store.root.parent), registry_path=runtime.registry_path, + project_contexts=runtime.project_contexts, codex_bin=runtime.codex_bin, + manager_scope_resolver=transport.bindings.steward_scope) + try: + current, resumed = restarted.open_session(goal_id=None, agent_id="codex", work_dir=runtime.registry_path.parent, + objective="ignored", mode="resume_latest", conversation_binding_id=binding["binding_id"], + source_context=initial["source"]) + assert resumed and current["session_id"] == initial["session_id"] and current["upstream_thread_id"] == original + finally: + restarted.close() + + +@pytest.mark.parametrize("fault", ["adoption", "adoption_readback", "creation_owner_crash"]) +def test_applied_commission_recovers_after_io_owner_restart_without_recreation(steward, monkeypatch, fault): + from loopx.extensions.lark.private_conversations import LarkPrivateConversations + + class InterruptedOwner(BaseException): + pass + + store, runtime, provider, transport, binding, capture, _ = steward + transport.admit("steward-app", provider.event("steward-app", "delegate", "/delegate --tokens 12000 Inspect README")) + transport.reconcile() + preview = transport.core.actions.store.list()[0] + text = "/confirm " + preview["proposal_id"] + transport.admit("steward-app", provider.event("steward-app", "confirm", text)) + request = next(r for r in transport.core.pending() if r["message"] == text) + original_adopt = transport.bindings.adopt_created_goal + original_apply = transport.core.actions.apply + failed = False + + def adopt(**kwargs): + nonlocal failed + if failed: + return original_adopt(**kwargs) + failed = True + if fault == "adoption_readback": + original_adopt(**kwargs) + raise OSError("temporary adoption IO failure") + + def apply(*args, **kwargs): + nonlocal failed + result = original_apply(*args, **kwargs) + if not failed: + failed = True + raise InterruptedOwner() + return result + + if fault == "creation_owner_crash": + monkeypatch.setattr(transport.core.actions, "apply", apply) + with pytest.raises(InterruptedOwner): + transport.reconcile() + else: + monkeypatch.setattr(transport.bindings, "adopt_created_goal", adopt) + transport.reconcile() + applied = transport.core.actions.load(preview["proposal_id"]) + resources = applied["receipt"]["resource_ids"] + assert applied["status"] == "applied" + assert finish(runtime, resources)["status"] == "completed" + pending = transport.core.read_request(request["request_ref"]) + assert pending["status"] == "command_queued" and not pending.get("delivery_verified") + assert not any("管家操作未完成" in message for _, message in provider.writes) + if fault != "creation_owner_crash": + assert pending["commission_resources"] == resources and pending["commission_adoption_pending"] + with pytest.raises(ValueError, match="adoption is pending"): + transport.core.record_delivery(request["request_ref"], session_id=pending["session_id"], turn_id=pending["turn_id"]) + monkeypatch.setattr(transport.core.actions, "apply", original_apply) + monkeypatch.setattr(transport.bindings, "adopt_created_goal", original_adopt) + threads_before = sum(json.loads(line).get("method") == "thread/start" for line in capture.read_text().splitlines()) + recovered = LarkPrivateConversations(controller=runtime, runtime_root=transport.runtime_root, + runner=provider, cli_bin=transport.cli_bin) + recovered.core.actions = transport.core.actions + recovered.reconcile() + result = recovered.core.read_request(request["request_ref"]) + assert result["commission_resources"] == resources and result["delivery_verified"] + assert not result.get("commission_adoption_pending") + assert transport.bindings.read()["bindings"][1]["goal_ids"] == [resources["goal_id"]] + assert len(json.loads(runtime.registry_path.read_text())["goals"]) == 1 + assert recovered.core.actions.load(preview["proposal_id"])["receipt"]["resource_ids"] == resources + assert sum(profile == "steward-app" and "Verified synthetic result" in message for profile, message in provider.writes) == 1 + writes = len(provider.writes) + recovered.reconcile() + assert len(provider.writes) == writes + assert sum(json.loads(line).get("method") == "thread/start" for line in capture.read_text().splitlines()) == threads_before + + +def test_post_commit_adoption_cannot_outlive_original_private_authority(steward, monkeypatch): + _, runtime, provider, transport, binding, _, _ = steward + transport.admit("steward-app", provider.event("steward-app", "delegate", "/delegate --tokens 12000 Inspect README")) + transport.reconcile() + preview = transport.core.actions.store.list()[0] + text = "/confirm " + preview["proposal_id"] + transport.admit("steward-app", provider.event("steward-app", "confirm", text)) + original = transport.bindings.adopt_created_goal + def unavailable(**kwargs): + raise OSError("temporary adoption IO failure") + monkeypatch.setattr(transport.bindings, "adopt_created_goal", unavailable) + transport.reconcile() + resources = transport.core.actions.load(preview["proposal_id"])["receipt"]["resource_ids"] + assert finish(runtime, resources)["status"] == "completed" + transport.bindings.disconnect(binding["binding_id"], expected_revision=transport.bindings.read()["revision"]) + monkeypatch.setattr(transport.bindings, "adopt_created_goal", original) + writes = len(provider.writes) + transport.reconcile() + row = next(r for r in transport.core.pending() if r["message"] == text) + assert row["commission_adoption_pending"] and not row.get("delivery_verified") + assert row["commission_resources"] == resources + assert len(provider.writes) == writes + assert len(json.loads(runtime.registry_path.read_text())["goals"]) == 1 + + +def test_expired_preview_and_wrong_source_cannot_create_a_goal(steward): + _, runtime, provider, transport, _, _, _ = steward + transport.admit("steward-app", provider.event("steward-app", "delegate", "/delegate --tokens 1000 Read only")) + transport.reconcile() + preview = transport.core.actions.store.list()[0] + text = "/confirm " + preview["proposal_id"] + wrong = provider.event("steward-app", "foreign", text) + wrong["chat_id"] = "oc_other_source" + provider.messages[wrong["message_id"]]["chat_id"] = wrong["chat_id"] + transport.admit("steward-app", wrong) + transport.reconcile() + assert json.loads(runtime.registry_path.read_text())["goals"] == [] + # Simulate a confirmation received beyond the immutable preview expiry. + event = provider.event("steward-app", "expired", text) + transport.admit("steward-app", event) + binding = transport.bindings.read()["bindings"][1] + request_ref = identity_ref(binding["provider_ref"], event["message_id"]) + row = transport.core.read_request(request_ref) + path = transport.core.root / f"{row['request_ref']}.json" + row["created_at"] = "2099-01-01T00:00:00+00:00" + path.write_text(json.dumps(row)) + transport.reconcile() + assert json.loads(runtime.registry_path.read_text())["goals"] == [] + assert transport.core.actions.load(preview["proposal_id"])["status"] == "preview_ready" + + +def test_exact_commission_stop_keeps_other_app_running_and_replay_keeps_target(steward): + store, runtime, provider, transport, _, _, _ = steward + transport.admit("steward-app", provider.event("steward-app", "delegate", "/delegate --tokens 12000 wait for interrupt")) + transport.reconcile() + preview = transport.core.actions.store.list()[0] + transport.admit("steward-app", provider.event("steward-app", "confirm", "/confirm " + preview["proposal_id"])) + transport.reconcile() + resources = transport.core.actions.load(preview["proposal_id"])["receipt"]["resource_ids"] + deadline = time.monotonic() + 10 + while not any(e["kind"] == "turn.started" for e in store.events_after(resources["session_id"], resources["turn_id"], None)): + assert time.monotonic() < deadline + time.sleep(.01) + transport.admit("notes-app", provider.event("notes-app", "independent", "ordinary independent answer")) + other = next(r for r in transport.core.pending() if r["message"] == "ordinary independent answer") + assert finish(runtime, other)["status"] == "completed" + stop = provider.event("steward-app", "stop_exact", "/stop-commission " + preview["proposal_id"]) + transport.admit("steward-app", stop) + transport.reconcile() + assert finish(runtime, resources)["status"] == "interrupted" + transport.admit("steward-app", stop) + transport.reconcile() + assert store.load_turn(other["session_id"], other["turn_id"])["status"] == "completed" + assert transport.core.actions.load(preview["proposal_id"])["receipt"]["resource_ids"] == resources