From d7547c994f61af8d82c19c449c10c82e11c6eda7 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:13:55 +0800 Subject: [PATCH 1/9] fix(chat): execute explicitly bound inbox work through governed turns Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- .../capabilities/manager_context/execution.py | 138 ++++++++++++ .../manager_context/inspection.py | 1 + loopx/chat_agent.py | 2 + loopx/chat_coordination.py | 5 + loopx/chat_runtime.py | 10 +- .../collaboration/semantic_request.ts | 10 +- .../collaboration/source_grant_observation.py | 39 ++-- .../collaboration/source_grants.ts | 39 ++++ .../control_plane/effect_runtime_handlers.ts | 1 + .../project_registry_io_manifest_v1.json | 2 +- .../control_plane_ts/semantic_request.test.ts | 4 + tests/control_plane_ts/source_grants.test.ts | 17 +- tests/test_manager_context_execution.py | 208 ++++++++++++++++++ 13 files changed, 451 insertions(+), 25 deletions(-) create mode 100644 loopx/capabilities/manager_context/execution.py create mode 100644 tests/test_manager_context_execution.py diff --git a/loopx/capabilities/manager_context/execution.py b/loopx/capabilities/manager_context/execution.py new file mode 100644 index 0000000000..f5e7a0084f --- /dev/null +++ b/loopx/capabilities/manager_context/execution.py @@ -0,0 +1,138 @@ +"""Chat selection over explicit operator bindings and the existing Turn owner. + +This adapter provisions no work, identity, grants, host profile or scheduler. +Context receipt, launch observation, receiver adoption and return stay separate. +""" + +from pathlib import Path +from collections.abc import Callable +from typing import Any + +from ...agent_registry import load_goal_from_registry +from ...collaboration_mcp import Delegations +from ...control_plane.collaboration import conversation_scope +from ...control_plane.collaboration.source_grant_observation import ( + external_source_policy, + registered_context_recipients, + source_context_authority, +) +from ...control_plane.effect_runtime import EffectRuntimeRejected, effect_runtime_result +from ...control_plane.projects.registry_codec import load_project_registry +from ...orchestration import compact_orchestration_policy, normalize_subagent_execution_config + + +def _grants(root: Path, registry: Path, session: dict[str, Any], turn: dict[str, Any]) -> list[dict[str, Any]]: + if conversation_scope(session, origin=turn.get("origin", "unknown"))["kind"] != "external_audience": + return [] + # Also reject runtime-incompatible registries through the existing owner. + authority = source_context_authority(root, registry, session, turn) + if authority["mode"] != "context_only": + raise ValueError("current source authorization unavailable") + ingress, source = external_source_policy(root, session, turn) + observed = registered_context_recipients(load_project_registry(registry)) + result = effect_runtime_result("collaboration.source.execution_bindings", { + "source": source, "sender_id": ingress["sender_id"], "available": observed["available"], + }) + return list(result["bindings"]) + + +def _service(root: Path, registry: Path, grant: dict[str, Any]) -> tuple[Delegations, dict[str, Any]]: + goal = load_goal_from_registry(registry, grant["goal_id"]) + if not goal: + raise ValueError("execution Goal unavailable") + configured = compact_orchestration_policy(goal.get("spawn_policy")).get("execution_config") + if not configured: + raise ValueError("Goal execution configuration unavailable") + relative = Path(normalize_subagent_execution_config(configured)) + workspace = Path(goal["repo"]).resolve() + config_root = workspace / ".loopx" / "config" + config = (workspace / relative).resolve() + if not config.is_relative_to(config_root) or not config.is_file() or config.stat().st_size > 1_000_000: + raise ValueError("Goal execution configuration unavailable") + service = Delegations(root, registry, grant["goal_id"], grant["requester_agent_id"], config) + binding = service.binding(grant["binding_id"], require_active=True) + if binding["agent_id"] != grant["agent_id"]: + raise ValueError("execution binding recipient changed") + return service, binding + + +def catalog(root: Path, registry: Path, session: dict[str, Any], turn: dict[str, Any]) -> dict[str, Any]: + """Expose only exact task choices, never commands, paths or host credentials.""" + try: + rows = [] + for grant in _grants(root, registry, session, turn): + _, binding = _service(root, registry, grant) + rows.append({key: grant[key] for key in ("goal_id", "agent_id", "binding_id")} + | {"todo_id": binding["todo_id"]}) + return {"available": True, "bindings": rows} + except (OSError, ValueError, KeyError, TypeError, EffectRuntimeRejected): + return {"available": False, "bindings": [], "reason": "execution_bindings_unavailable"} + + +def dispatch(root: Path, registry: Path, *, session: dict[str, Any], turn: dict[str, Any], + request: dict[str, Any], receipt: dict[str, Any], + execution_allowed: Callable[[], bool]) -> dict[str, Any]: + """Submit this delivered brief to exactly one separately granted binding. + + Stable operation identity uses the original inbox receipt. No retry invents + another operation, and no completed/stopped binding is reset or replaced. + The independent receiver must adopt and report through the existing inbox. + """ + selected = request.get("execution_binding_id") + if selected is None: + return {"submitted": False} + try: + if not execution_allowed(): + raise ValueError("source Turn is no longer active") + target = {key: request[key] for key in ("goal_id", "agent_id")} + if receipt.get("status") != "delivered" or any(receipt.get(key) != value for key, value in target.items()): + raise ValueError("original handoff receipt mismatch") + grant = next((row for row in _grants(root, registry, session, turn) + if all(row[key] == value for key, value in target.items()) + and row["binding_id"] == selected), None) + if grant is None: + raise ValueError("execution binding not granted to this source") + service, binding = _service(root, registry, grant) + operation = "context-" + receipt["request_id"] + # Recover an existing operation instead of probing a now-completed Todo + # and misreporting an accepted result as a refused new launch. + if service.path(operation).exists(): + row = service.read(operation) + return {"submitted": True, "operation_id": operation, "todo_id": binding["todo_id"], + "status": row["status"], "replayed": True} + preflight = service.inspect(selected) + # Match the existing start owner: an unprobed configured runtime may + # attempt bounded execution, but is never presented as ready/running. + # Known refusal, missing canonical acceptance or an unavailable runtime + # must stop before dispatch. Actual launch and acceptance remain owned + # by the governed Turn, not this point-in-time preview. + if (preflight["state"] not in {"launchable", "runtime_unverified"} + or not all(preflight.get(key) is True for key in + ("turn_eligible", "acceptance_ready", "authority_ready"))): + return {"submitted": False, "preflight": preflight, "reason": "execution_not_launchable"} + # Re-read the source and operator selection after the potentially slow + # preview. Delegations.start itself rechecks the exact binding and Turn. + if not execution_allowed() or grant not in _grants(root, registry, session, turn): + raise ValueError("source execution grant changed before launch") + requirement = ( + "\nOriginal owner inbox request " + receipt["request_id"] + ": read its original context, " + "independently acknowledge adopt/defer/reject, and return an audience-safe conclusion " + "for that exact request with manager-inbox report. A peer return alone does not reply " + "to the original conversation. Final prose is not a receipt." + ) + brief = {**request["brief"], "return_requirement": request["brief"]["return_requirement"] + requirement} + result = service.start(selected, operation, brief, + conversation={"session_id": session["session_id"], "turn_id": turn["turn_id"]}) + return {"submitted": True, "operation_id": operation, "todo_id": binding["todo_id"], + "status": result["status"], "runtime_readiness": preflight["state"], "replayed": False} + except (OSError, ValueError, KeyError, TypeError, EffectRuntimeRejected): + return {"submitted": False, "reason": "execution_binding_or_admission_unavailable"} + + +def handoff_message(receipt: dict[str, Any], execution: dict[str, Any]) -> str: + prefix = "已将原消息交给 " + str(receipt["agent_id"]) + "。" + if execution.get("submitted"): + return prefix + "已提交受控执行;受理不代表完成,接收方的处理结论将回到本次对话。" + if execution.get("reason"): + return prefix + "交接已保存,但执行绑定或任务准入未通过,尚未启动执行。" + return prefix + "材料已进入收件箱,尚未启动执行;收到接收方的处理结论后会回到这里。" diff --git a/loopx/capabilities/manager_context/inspection.py b/loopx/capabilities/manager_context/inspection.py index b39d4fc250..ffb16a1c14 100644 --- a/loopx/capabilities/manager_context/inspection.py +++ b/loopx/capabilities/manager_context/inspection.py @@ -191,6 +191,7 @@ def manager_index(context: dict[str, Any]) -> dict[str, Any]: if row.get("activation_state") != "stopped" ], "context_delegation": context.get("context_delegation"), + "context_execution": context.get("context_execution"), "evidence_sources": context.get("evidence_sources", [])[:12], "evidence_source_count": len(context.get("evidence_sources", [])), "agent_discovery": {"tool": read_tool, "view": "agents", "scope": "permitted_registry", diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index bf8e5ce92f..62f34ca0d6 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -383,6 +383,8 @@ def _turn_prompt( "A continuation, correction or status question belongs to the established Goal/owner. Preserve its constraints; do not restart, create a duplicate Goal or ask for permission already granted. " "For requested work, inspect the supplied Goal directory and relevant work/Agent evidence (using the declared read tool when incomplete). An empty delivery-grant list does not prove there is no existing work. " "Use context_handoff for a uniquely relevant, active and currently granted existing owner when the user asks for that work, even without the word delegate. " + "When context_execution.bindings supplies an exact existing Todo binding for the requested work, read that Todo and select its binding_id as context_handoff.execution_binding_id to submit governed execution. " + "Only select an explicitly cataloged binding that covers this request; registration and context delivery do not authorize execution. For consultation or unrelated/missing task bindings omit execution_binding_id. Never create a hidden Todo, change host settings or reuse a completed/stopped task to obtain launch. " "A correction to requested work is authorized context for its existing owner: send the corrected constraints in context_handoff, proposals=[], without asking to approve a Todo edit. Only direct control-plane record/configuration edits use that separate preview path. " "Do not redirect a Goal Chat back to its own owner: handle its follow-up in the current conversation. Registration alone is not delivery authority or execution readiness. " "Compare ALL plausible existing work items before selecting. A Goal ID, row order, or word overlap is not evidence of user intent. If two active items cover the requested subject and history does not distinguish them, context_handoff MUST be null; ask which in message, with goal_draft=null. " diff --git a/loopx/chat_coordination.py b/loopx/chat_coordination.py index 909e688d65..97fdc973f9 100644 --- a/loopx/chat_coordination.py +++ b/loopx/chat_coordination.py @@ -107,6 +107,11 @@ def prepare_turn_context(controller, adapter, session, turn_id, event_sink, *, s controller.store.root.parent, controller.registry_path, session, controller.store.load_turn(session_id, turn_id) or {}, ) + from .capabilities.manager_context.execution import catalog + context["context_execution"] = catalog( + controller.store.root.parent, controller.registry_path, session, + controller.store.load_turn(session_id, turn_id) or {}, + ) if isinstance(adapter, CodexAppServerAdapter): from .capabilities.manager_context.inspection import ManagerInspection, manager_index from .chat_manager_context import manager_authorization_scope_id diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index 2694a636dd..717cf8daae 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -1646,11 +1646,15 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None: receipt = deliver(self.store.root.parent, self.registry_path, session=session, turn=self.store.load_turn(session_id, turn_id) or {}, request=response["context_handoff"]) + from .capabilities.manager_context.execution import dispatch, handoff_message + execution = dispatch(self.store.root.parent, self.registry_path, + session=session, turn=self.store.load_turn(session_id, turn_id) or {}, + request=response["context_handoff"], receipt=receipt, + execution_allowed=lambda: not execution_ended()) response = {**response, "proposals": [], "gate": None, "context_handoff_receipt": receipt, - "message": ("已将交办说明和原消息交给 " if response["context_handoff"].get("brief") else "已将原消息交给 ") + receipt["agent_id"] + - "。材料已进入收件箱,后续处理结论会自动回到这里。" - "(委托 " + receipt["request_id"][:8] + ")"} + "context_execution": execution, + "message": handoff_message(receipt, execution)} except (OSError, ValueError): response = {**response, "proposals": [], "gate": None, "message": "材料尚未转交:目标绑定、来源授权或持久收件回读未通过。需要修复交接链路;没有改动任务或优先级。"} diff --git a/loopx/control_plane/collaboration/semantic_request.ts b/loopx/control_plane/collaboration/semantic_request.ts index 1fbbdb6dc0..1eede1e2a2 100644 --- a/loopx/control_plane/collaboration/semantic_request.ts +++ b/loopx/control_plane/collaboration/semantic_request.ts @@ -86,7 +86,7 @@ export function normalizeCollaborationBrief(value: unknown): JsonObject { export function normalizeCollaborationRequest(value: unknown): JsonObject { const request = requireJsonObject(value, "request"); - exactKeys(request, ["goal_id", "agent_id", "brief"], ["goal_id", "agent_id"]); + exactKeys(request, ["goal_id", "agent_id", "brief", "execution_binding_id"], ["goal_id", "agent_id"]); const result: JsonObject = {}; for (const key of ["goal_id", "agent_id"]) { const id = requireNonEmptyString(request[key], key); @@ -94,6 +94,14 @@ export function normalizeCollaborationRequest(value: unknown): JsonObject { result[key] = id; } if (request.brief !== undefined) result.brief = normalizeCollaborationBrief(request.brief); + // Selection is data, not a grant. The host separately rechecks the source, + // operator binding, canonical task, acceptance and governed Turn admission. + if (request.execution_binding_id !== undefined) { + if (!result.brief || typeof request.execution_binding_id !== "string" || !ID.test(request.execution_binding_id)) { + throw new EffectRuntimeRequestError("execution selection requires an exact binding and semantic brief"); + } + result.execution_binding_id = request.execution_binding_id; + } return result; } diff --git a/loopx/control_plane/collaboration/source_grant_observation.py b/loopx/control_plane/collaboration/source_grant_observation.py index fb88066738..61f73d1841 100644 --- a/loopx/control_plane/collaboration/source_grant_observation.py +++ b/loopx/control_plane/collaboration/source_grant_observation.py @@ -1,6 +1,7 @@ """Provider/store observation for the typed context-recipient policy owner.""" from pathlib import Path +from typing import Any from ...agent_registry import registered_agent_ids_for_goal from ..goals.activation import goal_is_stopped @@ -15,10 +16,24 @@ POLICY_SCHEMA = "loopx_manager_context_policy_v1" -def registered_context_recipients(registry: dict) -> dict: +def external_source_policy(runtime_root: Path, session: dict[str, Any], turn: dict[str, Any]) -> tuple[dict[str, Any], dict[str, Any]]: + """One provider provenance check for context and explicit execution grants.""" + ingress = _read(_root(runtime_root) / "ingress" / + (_hash([session["session_id"], turn["client_turn_id"]]) + ".json")) + if (ingress["channel"] != session.get("channel_id") + or ingress["message_digest"] != _hash(turn.get("message")) + or turn.get("origin") != "lark"): + raise ValueError("source mismatch") + policy = _read(_root(runtime_root) / "policy.json") + if policy.get("schema_version") != POLICY_SCHEMA: + raise ValueError("invalid policy") + return ingress, policy.get("sources", {}).get(ingress["channel"], {}) + + +def registered_context_recipients(registry: dict[str, Any]) -> dict[str, Any]: """Observe active Goal membership, ignoring unreadable activation rows.""" active_goals = [] - available = [] + available: list[dict[str, str]] = [] for goal in registry.get("goals", []): if not isinstance(goal, dict) or not goal.get("id"): continue @@ -34,8 +49,8 @@ def registered_context_recipients(registry: dict) -> dict: def source_context_authority( - runtime_root: Path, registry_path: Path, session: dict, turn: dict -) -> dict: + runtime_root: Path, registry_path: Path, session: dict[str, Any], turn: dict[str, Any] +) -> dict[str, Any]: """Return only a write-only recipient catalog; no cross-audience Goal evidence.""" if registry_path is None: return {"mode": "unavailable", "targets": []} @@ -59,21 +74,7 @@ def source_context_authority( if scope["kind"] != "external_audience": return {"mode": "unavailable", "targets": []} try: - ingress = _read( - _root(runtime_root) - / "ingress" - / (_hash([session["session_id"], turn["client_turn_id"]]) + ".json") - ) - if ( - ingress["channel"] != session.get("channel_id") - or ingress["message_digest"] != _hash(turn.get("message")) - or turn.get("origin") != "lark" - ): - raise ValueError("source mismatch") - policy = _read(_root(runtime_root) / "policy.json") - if policy.get("schema_version") != POLICY_SCHEMA: - raise ValueError("invalid policy") - grants = policy.get("sources", {}).get(ingress["channel"], {}) + ingress, grants = external_source_policy(runtime_root, session, turn) selected = effect_runtime_result("collaboration.source.recipients", { "source": grants, "sender_id": ingress["sender_id"], "available": observed["available"], diff --git a/loopx/control_plane/collaboration/source_grants.ts b/loopx/control_plane/collaboration/source_grants.ts index f9b1c601ec..81ce0b6f70 100644 --- a/loopx/control_plane/collaboration/source_grants.ts +++ b/loopx/control_plane/collaboration/source_grants.ts @@ -60,6 +60,45 @@ export function resolveSourceRecipients(params: JsonObject): JsonObject { a.goal_id.localeCompare(b.goal_id) || a.agent_id!.localeCompare(b.agent_id!)) }; } +/** A separate, exact operator grant for existing governed work. Context delivery + * alone never authorizes launch; neither models nor sources choose host profiles. + */ +export function sourceExecutionBindings(params: JsonObject): JsonObject { + const source = requireJsonObject(params.source, "source policy"); + const authorized = resolveSourceRecipients(params).targets as Recipient[]; + const available = recipients(params.available, "registered recipients", true); + const rows = Object.hasOwn(source, "execution_bindings") ? source.execution_bindings : []; + if (!Array.isArray(rows) || rows.length > 100) { + throw new EffectRuntimeRequestError("bounded source execution bindings required"); + } + const bindings = rows.map(raw => { + const row = requireJsonObject(raw, "source execution binding"); + const keys = ["goal_id", "agent_id", "requester_agent_id", "binding_id"]; + if (Object.keys(row).some(key => !keys.includes(key)) || keys.some(key => !Object.hasOwn(row, key))) { + throw new EffectRuntimeRequestError("exact source execution binding fields required"); + } + const binding: JsonObject = {}; + for (const key of keys) { + const id = requireNonEmptyString(row[key], key); + if (!/^[A-Za-z0-9][A-Za-z0-9._-]{0,159}$/.test(id)) { + throw new EffectRuntimeRequestError("invalid source execution binding identity"); + } + binding[key] = id; + } + if (binding.requester_agent_id === binding.agent_id) { + throw new EffectRuntimeRequestError("execution requester must be an independent registered Agent"); + } + return binding; + }); + const identities = bindings.map(row => JSON.stringify([row.goal_id, row.agent_id, row.binding_id])); + if (new Set(identities).size !== identities.length) { + throw new EffectRuntimeRequestError("ambiguous source execution binding"); + } + return {bindings: bindings.filter(row => + authorized.some(target => target.goal_id === row.goal_id && target.agent_id === row.agent_id) + && available.some(target => target.goal_id === row.goal_id && target.agent_id === row.requester_agent_id))}; +} + /** Plan one trusted-local configuration change; the adapter owns locking and IO. * Missing agent_id chooses the managed Goal, explicit Agent chooses an exception. */ diff --git a/loopx/control_plane/effect_runtime_handlers.ts b/loopx/control_plane/effect_runtime_handlers.ts index 2689961f25..66911145ca 100644 --- a/loopx/control_plane/effect_runtime_handlers.ts +++ b/loopx/control_plane/effect_runtime_handlers.ts @@ -589,6 +589,7 @@ export function createEffectRuntimeHandlers( ["collaboration.conversation.agent_target", lazyHandler(() => import("./collaboration/conversation_binding.ts"), ({resolveConversationAgentTarget}) => resolveConversationAgentTarget)], ["collaboration.peer.context_access", lazyHandler(() => import("./collaboration/peer_context.ts"), ({requirePeerContextAccess}) => requirePeerContextAccess)], ["collaboration.source.recipients", lazyHandler(() => import("./collaboration/source_grants.ts"), ({resolveSourceRecipients}) => resolveSourceRecipients)], + ["collaboration.source.execution_bindings", lazyHandler(() => import("./collaboration/source_grants.ts"), ({sourceExecutionBindings}) => sourceExecutionBindings)], ["collaboration.source.configure_recipient", lazyHandler(() => import("./collaboration/source_grants.ts"), ({configureSourceRecipient}) => configureSourceRecipient)], ["collaboration.conversation.reply_context", lazyHandler(() => import("./collaboration/conversation_reply_context.ts"), ({projectConversationReplyContext}) => projectConversationReplyContext)], ["chat.turn.accept", lazyHandler(() => import("./turn_driver/chat_turn_acceptance.ts"), ({planChatTurnAcceptance}) => planChatTurnAcceptance)], diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 56243372b1..5af92abb3f 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -1039,7 +1039,7 @@ }, { "site": "loopx/control_plane/collaboration/source_grant_observation.py::.source_context_authority::codec_read:load_project_registry#1", - "line": 43, + "line": 58, "column": 20, "kind": "codec_read", "api": "load_project_registry", diff --git a/tests/control_plane_ts/semantic_request.test.ts b/tests/control_plane_ts/semantic_request.test.ts index f32701cc86..5563db9b1e 100644 --- a/tests/control_plane_ts/semantic_request.test.ts +++ b/tests/control_plane_ts/semantic_request.test.ts @@ -10,6 +10,10 @@ const brief = { }; assert.deepEqual(normalizeCollaborationRequest(target), target); assert.deepEqual(normalizeCollaborationRequest({ ...target, brief }), { ...target, brief }); +assert.deepEqual(normalizeCollaborationRequest({ ...target, brief, execution_binding_id: "review" }), + {...target, brief, execution_binding_id: "review"}); +assert.throws(() => normalizeCollaborationRequest({...target, execution_binding_id: "review"})); +assert.throws(() => normalizeCollaborationRequest({...target, brief, execution_binding_id: "../escape"})); for (const request of [ { ...target, priority: "P0" }, { ...target, agent_id: "../reviewer" }, { ...target, brief: { ...brief, acceptance: [] } }, diff --git a/tests/control_plane_ts/source_grants.test.ts b/tests/control_plane_ts/source_grants.test.ts index 39e0833bf6..7dfce2fc18 100644 --- a/tests/control_plane_ts/source_grants.test.ts +++ b/tests/control_plane_ts/source_grants.test.ts @@ -1,6 +1,6 @@ import assert from "node:assert/strict"; import test from "node:test"; -import { configureSourceRecipient, resolveSourceRecipients } from "../../loopx/control_plane/collaboration/source_grants.ts"; +import { configureSourceRecipient, resolveSourceRecipients, sourceExecutionBindings } from "../../loopx/control_plane/collaboration/source_grants.ts"; const worker = { goal_id: "research", agent_id: "worker" }; const peer = { goal_id: "research", agent_id: "peer" }; @@ -8,6 +8,21 @@ const other = { goal_id: "other", agent_id: "worker" }; const available = [worker, peer, other]; const source = { local_delivery_scope: "selected", sender_ids: ["owner"], targets: [{ goal_id: "research" }] }; +test("execution needs an independent exact grant, current sender and both registered identities", () => { + const binding = {...worker, binding_id: "review", requester_agent_id: "peer"}; + const params = {source, sender_id: "owner", available}; + assert.deepEqual(sourceExecutionBindings(params), {bindings: []}); + const granted = {...source, execution_bindings: [binding]}; + assert.deepEqual(sourceExecutionBindings({...params, source: granted}), {bindings: [binding]}); + assert.deepEqual(sourceExecutionBindings({...params, source: granted, available: [worker]}), {bindings: []}); + assert.deepEqual(sourceExecutionBindings({...params, source: {...granted, blocked_targets: [worker]}}), {bindings: []}); + assert.throws(() => sourceExecutionBindings({...params, source: granted, sender_id: "other"})); + for (const bad of [null, {}, [binding, binding], [{...binding, requester_agent_id: "worker"}], + [{...binding, workspace: "/injected"}], [{...binding, binding_id: "../escape"}]]) { + assert.throws(() => sourceExecutionBindings({...params, source: {...source, execution_bindings: bad}})); + } +}); + test("a selected managed Goal includes current and future registered Agents, never another Goal", () => { assert.deepEqual(resolveSourceRecipients({ sender_id: "owner", source, available: [worker, other] }), { targets: [worker] }); assert.deepEqual(resolveSourceRecipients({ sender_id: "owner", source, available }), { targets: [peer, worker] }); diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py new file mode 100644 index 0000000000..b1d7654064 --- /dev/null +++ b/tests/test_manager_context_execution.py @@ -0,0 +1,208 @@ +"""Exact source grants, task-bound dispatch, revocation and truthful fallback. + +The governed integration uses the existing explicit fixture host. Live model, +provider delivery and routing quality require separate product qualification. +""" + +import json +from pathlib import Path + +import pytest + +from loopx.capabilities.manager_context import ( + POLICY_SCHEMA, _root, _write, deliver, register_ingress, +) +from loopx.capabilities.manager_context import execution +from loopx.chat import normalize_agent_response +from loopx.chat_agent import _turn_prompt +from loopx.chat_store import ChatSessionStore +from loopx.capabilities.manager_context.inspection import manager_index +from test_local_delegation import brief, wait, service as delegation_service # noqa: F401 +from test_independent_delegation_validation import independent_binding + + +def source(root, registry, *, goal_id, agent_id, requester, binding): + store = ChatSessionStore(root) + session = store.create_session(goal_id="loopx-manager", agent_id="codex", + adapter_kind="codex_app_server", upstream_thread_id="original", + channel_id="manager.external.test") + turn, _ = store.create_turn(session["session_id"], client_turn_id="current-request", + message="Do the authorized bounded review and return here.", origin="lark") + grant = {"goal_id": goal_id, "agent_id": agent_id, "requester_agent_id": requester, "binding_id": binding} + policy = {"schema_version": POLICY_SCHEMA, "sources": {session["channel_id"]: { + "sender_ids": ["owner"], "execution_bindings": [grant], + }}} + _write(_root(root) / "policy.json", policy) + register_ingress(root, session_id=session["session_id"], client_turn_id=turn["client_turn_id"], + channel=session["channel_id"], sender_id="owner", message=turn["message"], + source_id="lark:exact-message") + request = {"goal_id": goal_id, "agent_id": agent_id, "execution_binding_id": binding, "brief": brief()} + receipt = deliver(root, registry, session=session, turn=turn, request=request) + return store, session, turn, request, receipt, policy + + +@pytest.fixture +def flow(tmp_path, monkeypatch): + registry = tmp_path / "registry.json" + registry.write_text(json.dumps({"goals": [{"id": "research", "repo": str(tmp_path), + "spawn_policy": {"execution_config": ".loopx/config/delegations.json"}, + "coordination": {"registered_agents": ["lead", "worker"]}}]})) + config = tmp_path / ".loopx/config/delegations.json" + config.parent.mkdir(parents=True) + config.write_text("{}") + started = [] + + class BoundService: + def __init__(self, *args): + pass + + def binding(self, binding_id, **kwargs): + assert binding_id == "review" + return {"id": binding_id, "agent_id": "worker", "todo_id": "todo_current"} + + def path(self, operation): + return tmp_path / operation + + def inspect(self, binding_id): + return {"state": "launchable", "turn_eligible": True, "acceptance_ready": True, "authority_ready": True} + + def start(self, binding_id, operation, semantic_brief, **kwargs): + started.append((binding_id, operation, semantic_brief, kwargs)) + self.path(operation).write_text("prepared") + return {"status": "prepared"} + + def read(self, operation): + return {"status": "accepted"} + + monkeypatch.setattr(execution, "Delegations", BoundService) + values = source(tmp_path, registry, goal_id="research", agent_id="worker", requester="lead", binding="review") + return tmp_path, registry, values, started + + +def dispatch(flow, **changes): + root, registry, (_, session, turn, request, receipt, _), _ = flow + return execution.dispatch(root, registry, session=session, turn=turn, request=request, receipt=receipt, + execution_allowed=lambda: True, **changes) + + +def test_exact_catalog_and_launch_keep_original_conversation_and_operation(flow): + root, registry, (_, session, turn, request, receipt, _), started = flow + assert execution.catalog(root, registry, session, turn) == {"available": True, "bindings": [ + {"goal_id": "research", "agent_id": "worker", "binding_id": "review", "todo_id": "todo_current"}]} + normalized = normalize_agent_response({"context_handoff": request}) + assert normalized["context_handoff"]["execution_binding_id"] == "review" + assert dispatch(flow)["status"] == "prepared" + assert started[0][1] == "context-" + receipt["request_id"] + assert started[0][3]["conversation"] == {"session_id": session["session_id"], "turn_id": turn["turn_id"]} + assert receipt["request_id"] in started[0][2]["return_requirement"] + assert dispatch(flow)["replayed"] + assert len(started) == 1 + + +def test_model_receives_authorized_task_choices_after_manager_context_compaction(flow): + root, registry, (_, session, turn, _, _, _), _ = flow + choices = execution.catalog(root, registry, session, turn) + prompt = _turn_prompt(turn["message"], context_summary=json.dumps(manager_index({ + "context_execution": choices, + }))) + # The production prompt projection must retain the operator's exact task + # choice. Registration/route discovery cannot substitute for this grant. + assert '"binding_id": "review"' in prompt and '"todo_id": "todo_current"' in prompt + assert str(root) not in prompt and "delegations.json" not in prompt + + +@pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped"]) +def test_no_launch_after_source_or_registration_changes(flow, change): + root, registry, (_, session, turn, _, _, policy), started = flow + if change == "sender": + policy["sources"][session["channel_id"]]["sender_ids"] = ["someone-else"] + elif change == "body": + turn["message"] = "Different input" + elif change == "channel": + session["channel_id"] = "manager.external.other-app" + elif change == "revoke": + policy["sources"][session["channel_id"]].pop("execution_bindings") + elif change == "blocked": + policy["sources"][session["channel_id"]]["blocked_targets"] = [{"goal_id": "research"}] + else: + data = json.loads(registry.read_text()) + if change == "requester": + data["goals"][0]["coordination"]["registered_agents"] = ["worker"] + else: + data["goals"][0]["activation_state"] = "stopped" + registry.write_text(json.dumps(data)) + _write(_root(root) / "policy.json", policy) + assert not dispatch(flow)["submitted"] + assert not started + + +@pytest.mark.parametrize("state", ["turn_blocked", "acceptance_unavailable", "runtime_unavailable"]) +def test_preflight_never_becomes_a_launch_permit(flow, monkeypatch, state): + monkeypatch.setattr(execution.Delegations, "inspect", lambda *_: {"state": state}) + result = dispatch(flow) + assert result["preflight"]["state"] == state and not result["submitted"] + assert not flow[3] + + +def test_stop_during_preview_and_context_only_selection_do_not_launch(flow, monkeypatch): + root, registry, (_, session, turn, request, receipt, _), started = flow + active = [True] + def preview(*_): + active[0] = False + return {"state": "launchable", "turn_eligible": True, "acceptance_ready": True, "authority_ready": True} + monkeypatch.setattr(execution.Delegations, "inspect", preview) + result = execution.dispatch(root, registry, session=session, turn=turn, request=request, receipt=receipt, + execution_allowed=lambda: active[0]) + assert not result["submitted"] and not started + request.pop("execution_binding_id") + assert dispatch(flow) == {"submitted": False} + assert "尚未启动执行" in execution.handoff_message(receipt, {"submitted": False}) + + +def test_governed_worker_adopts_original_request_and_returns_without_another_model_turn(delegation_service): # noqa: F811 + root, service = delegation_service + independent_binding(delegation_service) + registry = service.registry + data = json.loads(registry.read_text()) + project = Path(data["goals"][0]["repo"]) + config = project / ".loopx/config/delegations.json" + config.parent.mkdir(parents=True, exist_ok=True) + config.write_bytes(service.config.read_bytes()) + data["goals"][0]["spawn_policy"] = {"execution_config": ".loopx/config/delegations.json"} + registry.write_text(json.dumps(data)) + store, session, turn, request, receipt, _ = source(service.root, registry, goal_id=service.goal_id, + agent_id="analyst", requester="lead", binding="analysis") + # The fixture receiver (not the Chat caller) reads, decides and returns the + # original owner request as well as its separately accepted peer result. + host = root / "fixture-host.py" + content = host.read_text() + anchor = "print(json.dumps(build_result" + statement = ( + "from loopx.control_plane.collaboration.peers import read_inbox\n" + "from loopx.capabilities.manager_context.roundtrip import report\n" + "for item in read_inbox(root / 'runtime', root / 'registry.json', envelope['goal_id'], actor)['items']:\n" + " if item.get('source_kind') != 'peer':\n" + " acknowledge(root / 'runtime', envelope['goal_id'], actor, item['request_id'], 'adopt', 'Receiver independently read the original scope.')\n" + " report(root / 'runtime', envelope['goal_id'], actor, item['request_id'], 'conclusion', 'Independent fixture result returned to the original request.')\n" + ) + host.write_text(content.replace(anchor, statement + anchor)) + launched = execution.dispatch(service.root, registry, session=session, turn=turn, request=request, receipt=receipt, + execution_allowed=lambda: True) + assert launched["submitted"], launched + assert launched["runtime_readiness"] == "runtime_unverified" + worker = execution._service(service.root, registry, {"goal_id": service.goal_id, "agent_id": "analyst", + "requester_agent_id": "lead", "binding_id": "analysis"})[0] + result = wait(worker, launched["operation_id"]) + assert result["status"] == "accepted", result + store.update_turn(session["session_id"], turn["turn_id"], status="completing", + response={"context_handoff_receipt": receipt}) + store.finalize_managed_turn_completion(session["session_id"], turn["turn_id"]) + from loopx.capabilities.manager_context.roundtrip import drain + sends = [] + def transport(route, *_args, **_kwargs): + sends.append(route["session_id"]) + return {"message_id": "provider-return", "reply_verified": True} + drain(service.root, registry, store, transport) + assert sends == [session["session_id"]] + assert worker.read(launched["operation_id"])["status"] == "accepted" + assert (Path(worker.binding("analysis")["workspace"]) / "host-invocations").read_text() == "1" From 9b0982bc90bbe2926cc5cfa854305fac7bb84e26 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:13:55 +0800 Subject: [PATCH 2/9] docs(chat): explain task-bound execution and remaining steward gates Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- .../rfcs/loopx-overall-roadmap-v0.md | 10 ++++ loopx/capabilities/manager_context/README.md | 46 +++++++++++++++++++ 2 files changed, 56 insertions(+) diff --git a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md index 02899a4dbc..0c55115463 100644 --- a/docs/architecture/rfcs/loopx-overall-roadmap-v0.md +++ b/docs/architecture/rfcs/loopx-overall-roadmap-v0.md @@ -777,6 +777,16 @@ deduplication and return. App settings select and read back the trigger per connection. External host-tool permission and sender-bound delegation remain separate gaps; receiving a request does not establish execution authority. +An external steward can now select an operator-granted existing task binding +from its exact source catalog and submit the handoff through the same governed +delegation/Turn owner. Durable context receipt, bounded launch, receiver +adoption, task acceptance and original-audience return remain distinct. No +identity, task, runtime policy or scheduler is provisioned by the handoff. Local +file/SQLite fixtures qualify dispatch, independent acceptance and return; +provider/model routing, new-task allocation, operator UI discovery and complete +stop/result presentation remain separate R3 gates. See +[bound inbox execution](../../../loopx/capabilities/manager_context/README.md#executing-already-bound-work). + Resolve source context before routing: a short reply retains the exact same-conversation parent, with missing and truncated material explicit. TS owns the bounded context projection; Lark transport and the private inbox preserve provider ancestry without granting diff --git a/loopx/capabilities/manager_context/README.md b/loopx/capabilities/manager_context/README.md index 39d7a46f61..ad52c9129f 100644 --- a/loopx/capabilities/manager_context/README.md +++ b/loopx/capabilities/manager_context/README.md @@ -23,6 +23,52 @@ Keep the policy private (0600); do not commit it. A read grant without a sender grant is insufficient. This does not grant remote delivery, evidence reads, worker launch, Todo/lease changes or protected operations. +### Executing already bound work + +An operator may separately authorize an external audience to select existing +governed task bindings. Add an `execution_bindings` list to that exact private +`sources[channel]` row (retain its independently verified `sender_ids`): + +```json +{"execution_bindings":[{"goal_id":"research","agent_id":"worker", + "requester_agent_id":"lead","binding_id":"review"}]} +``` + +The Goal must already point at its operator-owned `.loopx/config/` binding file +through `loopx configure-goal --goal-id research --subagent-execution-config +.loopx/config/delegations.json --execute`. Provision its real task, validation, +registered requester, workspace and host using the existing +[local delegation interface](../../../docs/reference/local-delegation.md). +Neither registration, a sender grant, read access nor a context brief creates +this execution grant. No task, role, profile or schedule is provisioned here. + +Chat exposes only the authorized Goal/Agent/binding/Todo identities in +`context_execution.bindings`. For requested work covered by that current task, +the model may select `execution_binding_id` in its semantic `context_handoff`. +Assessment-only handoffs omit it. The host rechecks provider provenance, source +revocation, active membership, Goal configuration, the original Turn and +canonical preflight, then uses `Delegations.start`. Commands, roots, requester +identity and host policy cannot come from the model. An unprobed runtime retains +`runtime_unverified`; it may attempt the existing bounded execution, and is never +reported as ready, running or complete. Known unavailability or refused task +admission does not launch. Retry preserves `context-` and +does not resume, reset or replace a stopped/completed operation. + +The receiver independently reads/adopts the original inbox request and returns +an audience-safe conclusion with the existing `manager-inbox report` path, in +addition to its peer result. Launch submission, receiver conclusion, canonical +acceptance and provider delivery remain separate evidence. The existing return +service replies to the original conversation; it does not start another model +thread. Plain inbox delivery now says that execution has not started. + +Remove that source's exact execution grant to prevent later launches/replays. +Already launched work keeps its original lifecycle: inspect and stop its exact +operation with `loopx delegation`, rather than assuming `/stop` of the manager +also stops an independently governed worker. Existing private configuration and +journals remain local. Automatic allocation of new tasks, operator UI discovery, +general inbox activation, cross-host readiness and live result/stop presentation +are still separate product work; this path only selects already configured work. + The existing local operator commands preview, apply and verify exceptions: ```sh From 0741bb9f212aa80e68c088285a7af169359fb523 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:17:35 +0800 Subject: [PATCH 3/9] chore(semantics): register the bound execution registry read Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- loopx/semantics/project_registry_io_manifest_v1.json | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 5af92abb3f..b3c5508ee5 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -261,6 +261,14 @@ "api": "load_registry", "classification": "codec_api" }, + { + "site": "loopx/capabilities/manager_context/execution.py::._grants::codec_read:load_project_registry#1", + "line": 32, + "column": 46, + "kind": "codec_read", + "api": "load_project_registry", + "classification": "codec_api" + }, { "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", "line": 879, From eabc77df44ba5405a252072f5664cfa4aa30d31c Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:28:26 +0800 Subject: [PATCH 4/9] fix(chat): scope execution guidance to granted task choices Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- loopx/capabilities/manager_context/inspection.py | 2 +- loopx/chat_agent.py | 14 ++++++++++++-- loopx/chat_coordination.py | 4 +++- tests/test_manager_context_execution.py | 3 +++ 4 files changed, 19 insertions(+), 4 deletions(-) diff --git a/loopx/capabilities/manager_context/inspection.py b/loopx/capabilities/manager_context/inspection.py index ffb16a1c14..677c5a41a0 100644 --- a/loopx/capabilities/manager_context/inspection.py +++ b/loopx/capabilities/manager_context/inspection.py @@ -191,7 +191,7 @@ def manager_index(context: dict[str, Any]) -> dict[str, Any]: if row.get("activation_state") != "stopped" ], "context_delegation": context.get("context_delegation"), - "context_execution": context.get("context_execution"), + **({"context_execution": context["context_execution"]} if "context_execution" in context else {}), "evidence_sources": context.get("evidence_sources", [])[:12], "evidence_source_count": len(context.get("evidence_sources", [])), "agent_discovery": {"tool": read_tool, "view": "agents", "scope": "permitted_registry", diff --git a/loopx/chat_agent.py b/loopx/chat_agent.py index 62f34ca0d6..edc966b8f3 100644 --- a/loopx/chat_agent.py +++ b/loopx/chat_agent.py @@ -326,6 +326,16 @@ def _turn_prompt( execution_mode: bool = False, runtime_profile: str = "restricted", ) -> str: + try: + supplied = json.loads(context_summary) + choices = supplied.get("context_execution") if isinstance(supplied, dict) else None + except (ValueError, TypeError): + choices = None + execution_guidance = ( + "When context_execution.bindings supplies an exact existing Todo binding for the requested work, read that Todo and select its binding_id as context_handoff.execution_binding_id to submit governed execution. " + "Only select an explicitly cataloged binding that covers this request; registration and context delivery do not authorize execution. For consultation or unrelated/missing task bindings omit execution_binding_id. Never create a hidden Todo, change host settings or reuse a completed/stopped task to obtain launch. " + if isinstance(choices, dict) and choices.get("bindings") else "" + ) envelope = { "schema_version": CHAT_AGENT_RESPONSE_SCHEMA_VERSION, "message": "Complete answer for the operator, at the depth this task needs.", @@ -383,8 +393,8 @@ def _turn_prompt( "A continuation, correction or status question belongs to the established Goal/owner. Preserve its constraints; do not restart, create a duplicate Goal or ask for permission already granted. " "For requested work, inspect the supplied Goal directory and relevant work/Agent evidence (using the declared read tool when incomplete). An empty delivery-grant list does not prove there is no existing work. " "Use context_handoff for a uniquely relevant, active and currently granted existing owner when the user asks for that work, even without the word delegate. " - "When context_execution.bindings supplies an exact existing Todo binding for the requested work, read that Todo and select its binding_id as context_handoff.execution_binding_id to submit governed execution. " - "Only select an explicitly cataloged binding that covers this request; registration and context delivery do not authorize execution. For consultation or unrelated/missing task bindings omit execution_binding_id. Never create a hidden Todo, change host settings or reuse a completed/stopped task to obtain launch. " + + execution_guidance + + "A correction to requested work is authorized context for its existing owner: send the corrected constraints in context_handoff, proposals=[], without asking to approve a Todo edit. Only direct control-plane record/configuration edits use that separate preview path. " "Do not redirect a Goal Chat back to its own owner: handle its follow-up in the current conversation. Registration alone is not delivery authority or execution readiness. " "Compare ALL plausible existing work items before selecting. A Goal ID, row order, or word overlap is not evidence of user intent. If two active items cover the requested subject and history does not distinguish them, context_handoff MUST be null; ask which in message, with goal_draft=null. " diff --git a/loopx/chat_coordination.py b/loopx/chat_coordination.py index 97fdc973f9..46d581a790 100644 --- a/loopx/chat_coordination.py +++ b/loopx/chat_coordination.py @@ -108,10 +108,12 @@ def prepare_turn_context(controller, adapter, session, turn_id, event_sink, *, s controller.store.load_turn(session_id, turn_id) or {}, ) from .capabilities.manager_context.execution import catalog - context["context_execution"] = catalog( + execution_catalog = catalog( controller.store.root.parent, controller.registry_path, session, controller.store.load_turn(session_id, turn_id) or {}, ) + if execution_catalog["bindings"] or not execution_catalog["available"]: + context["context_execution"] = execution_catalog if isinstance(adapter, CodexAppServerAdapter): from .capabilities.manager_context.inspection import ManagerInspection, manager_index from .chat_manager_context import manager_authorization_scope_id diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index b1d7654064..fdd116cb72 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -109,6 +109,9 @@ def test_model_receives_authorized_task_choices_after_manager_context_compaction # choice. Registration/route discovery cannot substitute for this grant. assert '"binding_id": "review"' in prompt and '"todo_id": "todo_current"' in prompt assert str(root) not in prompt and "delegations.json" not in prompt + disabled = manager_index({}) + assert "context_execution" not in disabled + assert "execution_binding_id" not in _turn_prompt(turn["message"], context_summary=json.dumps(disabled)) @pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped"]) From fb80f69e9eba92f85196869252fe5d45bd365ad6 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 09:54:05 +0800 Subject: [PATCH 5/9] fix(collaboration): settle native validation and external wake boundaries Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- loopx/capabilities/manager_context/README.md | 7 ++++ loopx/collaboration_mcp.py | 9 +++- .../control_plane/collaboration/chat_mode.ts | 14 +++++-- tests/control_plane_ts/chat_mode.test.ts | 6 ++- .../test_independent_delegation_validation.py | 42 +++++++++++++++++++ tests/test_manager_context_execution.py | 18 ++++++-- 6 files changed, 87 insertions(+), 9 deletions(-) diff --git a/loopx/capabilities/manager_context/README.md b/loopx/capabilities/manager_context/README.md index ad52c9129f..54b96e54ca 100644 --- a/loopx/capabilities/manager_context/README.md +++ b/loopx/capabilities/manager_context/README.md @@ -60,6 +60,13 @@ addition to its peer result. Launch submission, receiver conclusion, canonical acceptance and provider delivery remain separate evidence. The existing return service replies to the original conversation; it does not start another model thread. Plain inbox delivery now says that execution has not started. +The native postcondition entry retires only its exact operation-owned temporary +host input before checking a clean delivery worktree; unrelated files and actual +artifact changes still fail canonical validation. Validation recovery resumes the +original Turn after its retained host result, without invoking the model again. +An external conversation is not a native Goal wake owner. The typed wake owner +settles that separate intent as `no_wake_owner`; the exact original inbox return +still carries the receiver's conclusion, without a hidden Goal or retry loop. Remove that source's exact execution grant to prevent later launches/replays. Already launched work keeps its original lifecycle: inspect and stop its exact diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index 8659773ded..0855ff0e24 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -1935,7 +1935,14 @@ def main(): service = Delegations(args.runtime_root, args.registry, args.goal_id, args.agent_id, args.execution_config) if args.delegation_action == "validate": - service._validate(service._bound(_read(service.path(args.operation_id)))) + row = _read(service.path(args.operation_id)) + binding = service._bound(row) + # The native Turn invokes this after its Host returns, before the + # worker's outer finally. Retire only the exact operation-owned + # input before the canonical clean-worktree check; user files and + # actual delivery changes must still be rejected by that check. + service._clear_delegation_bootstrap(row, binding) + service._validate(binding) else: service._stop_signal = install_worker_stop_signal( service._stop_path(service.path(args.operation_id))) diff --git a/loopx/control_plane/collaboration/chat_mode.ts b/loopx/control_plane/collaboration/chat_mode.ts index 03cce4cf7d..89197f2890 100644 --- a/loopx/control_plane/collaboration/chat_mode.ts +++ b/loopx/control_plane/collaboration/chat_mode.ts @@ -15,10 +15,16 @@ export function planChatMode(input: JsonObject): JsonObject { const session = requireJsonObject(input.session, "conversation session"); const operation = input.operation; requireThat(["configure", "start", "resume", "pause", "exit", "message", "wake"].includes(String(operation)), "unsupported conversation operation"); - requireThat(resolveConversationScope(session).kind === "owner_goal" - && input.origin === (operation === "wake" ? "host" : "web") - && session.session_mode !== "attached_host" - && session.agent_id === "codex", "LoopX mode requires a local managed Codex Goal conversation"); + const localOwner = resolveConversationScope(session).kind === "owner_goal" + && session.session_mode !== "attached_host" && session.agent_id === "codex"; + // External inbox returns have their own exact audience owner. A recorded + // conversation is provenance, not authority to resume a native Goal. Settle + // this host intent instead of throwing and retrying it on every pump tick. + if (operation === "wake" && input.origin === "host" && !localOwner) { + return {operation, state: "refused", reason: "no_wake_owner"}; + } + requireThat(localOwner && input.origin === (operation === "wake" ? "host" : "web"), + "LoopX mode requires a local managed Codex Goal conversation"); const settings = requireJsonObject(input.settings, "conversation settings"); const native = requireJsonObject(input.native ?? {}, "native Goal observation"); if (operation === "wake") return planDelegationWake(input, session, settings, native); diff --git a/tests/control_plane_ts/chat_mode.test.ts b/tests/control_plane_ts/chat_mode.test.ts index 5892700680..48d556c775 100644 --- a/tests/control_plane_ts/chat_mode.test.ts +++ b/tests/control_plane_ts/chat_mode.test.ts @@ -95,7 +95,11 @@ test("a host wake reuses the resume facts and returns a typed outcome, never an assert.throws(() => planChatMode({...input, origin: "host"}), /local managed/); assert.throws(() => planChatMode({...wake, origin: "web"}), /local managed/); assert.throws(() => planChatMode({...wake, origin: "external"}), /local managed/); - assert.throws(() => planChatMode({...wake, session: {...enabled, channel_id: "manager"}})); + for (const changes of [{channel_id: "manager"}, {channel_id: "manager.external.test"}, + {session_mode: "attached_host"}, {agent_id: "other"}]) { + assert.deepEqual(planChatMode({...wake, session: {...enabled, ...changes}}), + {operation: "wake", state: "refused", reason: "no_wake_owner"}); + } }); test("only a provider-dispatched wake Turn is dispatch evidence; a queued or merely starting one is replayed", () => { diff --git a/tests/test_independent_delegation_validation.py b/tests/test_independent_delegation_validation.py index abedf641ec..2d58dfcdf5 100644 --- a/tests/test_independent_delegation_validation.py +++ b/tests/test_independent_delegation_validation.py @@ -113,6 +113,48 @@ def git(*args: str) -> None: runner._validate(binding) +@pytest.mark.parametrize("extra_file", [None, "DELEGATION.json", "unrelated.txt"]) +def test_native_validator_retires_only_its_exact_host_input(service, monkeypatch, extra_file): + root, runner = service + repository = root / "validation-repository" + repository.mkdir() + def git(*args): + subprocess.run(["git", "-C", str(repository), "-c", "user.name=Fixture", + "-c", "user.email=fixture@example.invalid", *args], + check=True, capture_output=True) + git("init", "-b", "main") + (repository / "marker").write_text("ok") + git("add", "marker") + git("commit", "-m", "Validation fixture") + remote = "https://example.invalid/synthetic/validation.git" + git("remote", "add", "origin", remote) + workspace = root / "validation-worker" + git("worktree", "add", "--detach", str(workspace), "HEAD") + independent_binding(service, task_repository=remote, + validation_argv=[sys.executable, "-c", + "from pathlib import Path; assert Path('marker').read_text() == 'ok'"]) + config = json.loads(runner.config.read_text()) + config["bindings"][0]["workspace"] = str(workspace) + runner.config.write_text(json.dumps(config)) + monkeypatch.setattr(runner, "_spawn", lambda _: None) + runner.start("analysis", "validator-host-input", brief()) + row = json.loads(runner.path("validator-host-input").read_text()) + binding = runner.binding("analysis") + runner._write_delegation_bootstrap(row, binding) + if extra_file: + (workspace / extra_file).write_text("Caller-owned content") + arguments = runner._execution_arguments(binding, "validator-host-input") + argv = json.loads(arguments[arguments.index("--validation-command-json") + 1]) + result = subprocess.run(argv, cwd=workspace, capture_output=True, text=True, timeout=30) + assert (result.returncode == 0) is (extra_file is None), result.stderr + assert (workspace / "marker").read_text() == "ok" + if extra_file: + assert (workspace / extra_file).read_text() == "Caller-owned content" + else: + assert not (workspace / "DELEGATION.json").exists() + assert not demo.canonical_tasks(root)[binding["todo_id"]]["done"] + + @pytest.mark.parametrize("handoff_mode", ["soft_claim", "hard_lease"]) def test_independent_result_reconnects_and_revalidates_without_goal_binding(service, monkeypatch, handoff_mode): root, runner = service diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index fdd116cb72..d86c52f217 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -21,7 +21,7 @@ from test_independent_delegation_validation import independent_binding -def source(root, registry, *, goal_id, agent_id, requester, binding): +def source(root, registry, *, goal_id, agent_id, requester, binding, source_id="lark:exact-message"): store = ChatSessionStore(root) session = store.create_session(goal_id="loopx-manager", agent_id="codex", adapter_kind="codex_app_server", upstream_thread_id="original", @@ -35,7 +35,7 @@ def source(root, registry, *, goal_id, agent_id, requester, binding): _write(_root(root) / "policy.json", policy) register_ingress(root, session_id=session["session_id"], client_turn_id=turn["client_turn_id"], channel=session["channel_id"], sender_id="owner", message=turn["message"], - source_id="lark:exact-message") + source_id=source_id) request = {"goal_id": goal_id, "agent_id": agent_id, "execution_binding_id": binding, "brief": brief()} receipt = deliver(root, registry, session=session, turn=turn, request=request) return store, session, turn, request, receipt, policy @@ -114,7 +114,7 @@ def test_model_receives_authorized_task_choices_after_manager_context_compaction assert "execution_binding_id" not in _turn_prompt(turn["message"], context_summary=json.dumps(disabled)) -@pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped"]) +@pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped", "binding"]) def test_no_launch_after_source_or_registration_changes(flow, change): root, registry, (_, session, turn, _, _, policy), started = flow if change == "sender": @@ -123,6 +123,10 @@ def test_no_launch_after_source_or_registration_changes(flow, change): turn["message"] = "Different input" elif change == "channel": session["channel_id"] = "manager.external.other-app" + elif change == "binding": + # A newly configured choice in the same Goal/Agent is not covered by + # the existing source's exact binding consent. + flow[2][3]["execution_binding_id"] = "new-task-choice" elif change == "revoke": policy["sources"][session["channel_id"]].pop("execution_bindings") elif change == "blocked": @@ -209,3 +213,11 @@ def transport(route, *_args, **_kwargs): assert sends == [session["session_id"]] assert worker.read(launched["operation_id"])["status"] == "accepted" assert (Path(worker.binding("analysis")["workspace"]) / "host-invocations").read_text() == "1" + # A new user request cannot reactivate the completed task merely because + # its source grant and registered Agent still exist. + _, new_session, new_turn, new_request, new_receipt, _ = source(service.root, registry, + goal_id=service.goal_id, agent_id="analyst", requester="lead", binding="analysis", source_id="lark:second-message") + refused = execution.dispatch(service.root, registry, session=new_session, turn=new_turn, + request=new_request, receipt=new_receipt, execution_allowed=lambda: True) + assert not refused["submitted"] and refused["reason"] == "execution_not_launchable", refused + assert (Path(worker.binding("analysis")["workspace"]) / "host-invocations").read_text() == "1" From f076d9fbd1cea8be372bbc3451adfd1d15162380 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:26:05 +0800 Subject: [PATCH 6/9] refactor(chat): keep bound handoff presentation in its owning capability Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- .../capabilities/manager_context/execution.py | 23 +++++++++++++++ loopx/chat_runtime.py | 28 +++++-------------- tests/test_manager_context_execution.py | 26 +++++++++++++++++ 3 files changed, 56 insertions(+), 21 deletions(-) diff --git a/loopx/capabilities/manager_context/execution.py b/loopx/capabilities/manager_context/execution.py index f5e7a0084f..7d7745502b 100644 --- a/loopx/capabilities/manager_context/execution.py +++ b/loopx/capabilities/manager_context/execution.py @@ -136,3 +136,26 @@ def handoff_message(receipt: dict[str, Any], execution: dict[str, Any]) -> str: if execution.get("reason"): return prefix + "交接已保存,但执行绑定或任务准入未通过,尚未启动执行。" return prefix + "材料已进入收件箱,尚未启动执行;收到接收方的处理结论后会回到这里。" + + +def handoff_response(root: Path, registry: Path, *, session: dict[str, Any], + turn: dict[str, Any], response: dict[str, Any], + source_authorized: Callable[[], bool], + execution_allowed: Callable[[], bool]) -> dict[str, Any]: + """Present context delivery and separately admitted execution on one path.""" + from . import deliver + + try: + if not source_authorized(): + raise ValueError("manager connection authority is no longer available") + receipt = deliver(root, registry, session=session, turn=turn, + request=response["context_handoff"]) + execution = dispatch(root, registry, session=session, turn=turn, + request=response["context_handoff"], receipt=receipt, + execution_allowed=execution_allowed) + return {**response, "proposals": [], "gate": None, + "context_handoff_receipt": receipt, "context_execution": execution, + "message": handoff_message(receipt, execution)} + except (OSError, ValueError): + return {**response, "proposals": [], "gate": None, + "message": "材料尚未转交:目标绑定、来源授权或持久收件回读未通过。需要修复交接链路;没有改动任务或优先级。"} diff --git a/loopx/chat_runtime.py b/loopx/chat_runtime.py index eb6112edb9..bb3cba654b 100644 --- a/loopx/chat_runtime.py +++ b/loopx/chat_runtime.py @@ -1647,29 +1647,15 @@ def event_sink(kind: str, payload: dict[str, Any]) -> None: event_buffer.close() return if response.get("context_handoff") is not None: - from .capabilities.manager_context import deliver + from .capabilities.manager_context.execution import handoff_response if scope["kind"] == "unavailable": raise ValueError("context handoff requires a scoped conversation") - try: - if scope["kind"] == "external_audience" and ( - self.manager_scope_resolver is None or not self.manager_scope_resolver(session) - ): - raise ValueError("manager connection authority is no longer available") - receipt = deliver(self.store.root.parent, self.registry_path, - session=session, turn=self.store.load_turn(session_id, turn_id) or {}, - request=response["context_handoff"]) - from .capabilities.manager_context.execution import dispatch, handoff_message - execution = dispatch(self.store.root.parent, self.registry_path, - session=session, turn=self.store.load_turn(session_id, turn_id) or {}, - request=response["context_handoff"], receipt=receipt, - execution_allowed=lambda: not execution_ended()) - response = {**response, "proposals": [], "gate": None, - "context_handoff_receipt": receipt, - "context_execution": execution, - "message": handoff_message(receipt, execution)} - except (OSError, ValueError): - response = {**response, "proposals": [], "gate": None, - "message": "材料尚未转交:目标绑定、来源授权或持久收件回读未通过。需要修复交接链路;没有改动任务或优先级。"} + response = handoff_response( + self.store.root.parent, self.registry_path, session=session, + turn=self.store.load_turn(session_id, turn_id) or {}, response=response, + source_authorized=lambda: scope["kind"] != "external_audience" or bool( + self.manager_scope_resolver and self.manager_scope_resolver(session)), + execution_allowed=lambda: not execution_ended()) response = offer_team_plan_confirmation( store=self.store, session=session, diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index d86c52f217..7a501c7ffd 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -114,6 +114,32 @@ def test_model_receives_authorized_task_choices_after_manager_context_compaction assert "execution_binding_id" not in _turn_prompt(turn["message"], context_summary=json.dumps(disabled)) +def test_handoff_response_preserves_receipt_and_separate_execution_status(flow): + root, registry, (_, session, turn, request, _, _), started = flow + response = execution.handoff_response(root, registry, session=session, turn=turn, + response={"context_handoff": request, "message": "Unverified model completion claim.", + "proposals": [{"kind": "unused"}], "gate": {"kind": "unused"}}, + source_authorized=lambda: True, execution_allowed=lambda: True) + assert response["context_handoff_receipt"]["status"] == "delivered" + assert response["context_execution"]["status"] == "prepared" + assert "受理不代表完成" in response["message"] + assert response["proposals"] == [] and response["gate"] is None + assert len(started) == 1 + + +def test_handoff_scope_revocation_stops_before_inbox_delivery(flow, monkeypatch): + root, registry, (_, session, turn, request, _, _), started = flow + from loopx.capabilities import manager_context + def forbidden_delivery(*args, **kwargs): + pytest.fail("revoked manager scope must not publish an inbox request") + monkeypatch.setattr(manager_context, "deliver", forbidden_delivery) + response = execution.handoff_response(root, registry, session=session, turn=turn, + response={"context_handoff": request}, source_authorized=lambda: False, + execution_allowed=lambda: True) + assert "尚未转交" in response["message"] + assert "context_handoff_receipt" not in response and not started + + @pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped", "binding"]) def test_no_launch_after_source_or_registration_changes(flow, change): root, registry, (_, session, turn, _, _, policy), started = flow From e470ce28e6716ad212c7ce04146d527035bc3f34 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 17:36:16 +0800 Subject: [PATCH 7/9] fix(chat): reuse guarded source observation for execution grants Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- .../capabilities/manager_context/execution.py | 17 ++-------- .../collaboration/source_grant_observation.py | 31 ++++++++++++++----- .../project_registry_io_manifest_v1.json | 14 ++------- tests/test_manager_context_execution.py | 11 +++++++ 4 files changed, 41 insertions(+), 32 deletions(-) diff --git a/loopx/capabilities/manager_context/execution.py b/loopx/capabilities/manager_context/execution.py index 7d7745502b..1ae5c090c2 100644 --- a/loopx/capabilities/manager_context/execution.py +++ b/loopx/capabilities/manager_context/execution.py @@ -12,27 +12,16 @@ from ...collaboration_mcp import Delegations from ...control_plane.collaboration import conversation_scope from ...control_plane.collaboration.source_grant_observation import ( - external_source_policy, - registered_context_recipients, - source_context_authority, + source_execution_bindings, ) -from ...control_plane.effect_runtime import EffectRuntimeRejected, effect_runtime_result -from ...control_plane.projects.registry_codec import load_project_registry +from ...control_plane.effect_runtime import EffectRuntimeRejected from ...orchestration import compact_orchestration_policy, normalize_subagent_execution_config def _grants(root: Path, registry: Path, session: dict[str, Any], turn: dict[str, Any]) -> list[dict[str, Any]]: if conversation_scope(session, origin=turn.get("origin", "unknown"))["kind"] != "external_audience": return [] - # Also reject runtime-incompatible registries through the existing owner. - authority = source_context_authority(root, registry, session, turn) - if authority["mode"] != "context_only": - raise ValueError("current source authorization unavailable") - ingress, source = external_source_policy(root, session, turn) - observed = registered_context_recipients(load_project_registry(registry)) - result = effect_runtime_result("collaboration.source.execution_bindings", { - "source": source, "sender_id": ingress["sender_id"], "available": observed["available"], - }) + result = source_execution_bindings(root, registry, session, turn) return list(result["bindings"]) diff --git a/loopx/control_plane/collaboration/source_grant_observation.py b/loopx/control_plane/collaboration/source_grant_observation.py index cbd7d94476..8659cf2cb7 100644 --- a/loopx/control_plane/collaboration/source_grant_observation.py +++ b/loopx/control_plane/collaboration/source_grant_observation.py @@ -117,6 +117,29 @@ def source_context_target_authority( return _source_context_grant(runtime_root, session, turn, [target]) +def _source_registry_recipients(registry_path: Path) -> dict: + registry = load_project_registry(registry_path) + if not isinstance(registry, dict): + raise ValueError("invalid registry") + require_runtime_compatible_project_registry( + registry, operation="context source recipient observation" + ) + return registered_context_recipients(registry) + + +def source_execution_bindings( + runtime_root: Path, registry_path: Path, session: dict, turn: dict +) -> dict: + """Observe one compatible registry; TS owns sender and exact binding grants.""" + if conversation_scope(session, origin=turn.get("origin", "unknown"))["kind"] != "external_audience": + return {"bindings": []} + observed = _source_registry_recipients(registry_path) + ingress, source = external_source_policy(runtime_root, session, turn) + return effect_runtime_result("collaboration.source.execution_bindings", { + "source": source, "sender_id": ingress["sender_id"], "available": observed["available"], + }) + + def source_context_authority( runtime_root: Path, registry_path: Path, session: dict, turn: dict ) -> dict: @@ -124,15 +147,9 @@ def source_context_authority( if registry_path is None: return {"mode": "unavailable", "targets": []} try: - registry = load_project_registry(registry_path) - if not isinstance(registry, dict): - raise ValueError("invalid registry") - require_runtime_compatible_project_registry( - registry, operation="context source recipient observation" - ) + observed = _source_registry_recipients(registry_path) except (OSError, ValueError, TypeError): return {"mode": "unavailable", "targets": []} - observed = registered_context_recipients(registry) return _source_context_grant( runtime_root, session, diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 806108059b..9c889cfc7a 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -261,14 +261,6 @@ "api": "load_registry", "classification": "codec_api" }, - { - "site": "loopx/capabilities/manager_context/execution.py::._grants::codec_read:load_project_registry#1", - "line": 32, - "column": 46, - "kind": "codec_read", - "api": "load_project_registry", - "classification": "codec_api" - }, { "site": "loopx/capabilities/manager_context/roundtrip.py::.drain::codec_read:load_project_registry#1", "line": 884, @@ -1046,9 +1038,9 @@ "classification": "codec_api" }, { - "site": "loopx/control_plane/collaboration/source_grant_observation.py::.source_context_authority::codec_read:load_project_registry#1", - "line": 127, - "column": 20, + "site": "loopx/control_plane/collaboration/source_grant_observation.py::._source_registry_recipients::codec_read:load_project_registry#1", + "line": 121, + "column": 16, "kind": "codec_read", "api": "load_project_registry", "classification": "codec_api" diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index 7a501c7ffd..1fe18f9010 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -140,6 +140,17 @@ def forbidden_delivery(*args, **kwargs): assert "context_handoff_receipt" not in response and not started +def test_lifecycle_only_registry_cannot_authorize_host_execution(flow): + from loopx.control_plane.projects.registry_codec import SOURCE_SESSION_PROFILE_ID + root, registry, (_, session, turn, _, _, _), started = flow + data = json.loads(registry.read_text()) + data["profile_id"] = SOURCE_SESSION_PROFILE_ID + registry.write_text(json.dumps(data)) + assert execution.catalog(root, registry, session, turn)["available"] is False + assert dispatch(flow)["submitted"] is False + assert not started + + @pytest.mark.parametrize("change", ["sender", "body", "channel", "revoke", "blocked", "requester", "stopped", "binding"]) def test_no_launch_after_source_or_registration_changes(flow, change): root, registry, (_, session, turn, _, _, policy), started = flow From b8b6cd2ff80d22979110168583d0c53abb217aa2 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 22:56:08 +0800 Subject: [PATCH 8/9] fix(manager-context): keep host return routing outside brief budgets Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- loopx/capabilities/manager_context/README.md | 5 + .../capabilities/manager_context/execution.py | 12 +- loopx/collaboration_mcp.py | 25 +++- .../collaboration/source_grant_observation.py | 9 +- tests/test_manager_context_execution.py | 108 ++++++++++++++++-- 5 files changed, 136 insertions(+), 23 deletions(-) diff --git a/loopx/capabilities/manager_context/README.md b/loopx/capabilities/manager_context/README.md index 54b96e54ca..5dc97164d4 100644 --- a/loopx/capabilities/manager_context/README.md +++ b/loopx/capabilities/manager_context/README.md @@ -54,6 +54,11 @@ reported as ready, running or complete. Known unavailability or refused task admission does not launch. Retry preserves `context-` and does not resume, reset or replace a stopped/completed operation. +The original brief and its field/encoded-byte limits remain unchanged. The trusted +host binds the original inbox request in the existing operation identity and +receiver bootstrap instructions, outside user-authored fields. Replay preserves +that binding; it cannot substitute another request. + The receiver independently reads/adopts the original inbox request and returns an audience-safe conclusion with the existing `manager-inbox report` path, in addition to its peer result. Launch submission, receiver conclusion, canonical diff --git a/loopx/capabilities/manager_context/execution.py b/loopx/capabilities/manager_context/execution.py index 1ae5c090c2..32a355adc3 100644 --- a/loopx/capabilities/manager_context/execution.py +++ b/loopx/capabilities/manager_context/execution.py @@ -103,15 +103,9 @@ def dispatch(root: Path, registry: Path, *, session: dict[str, Any], turn: dict[ # preview. Delegations.start itself rechecks the exact binding and Turn. if not execution_allowed() or grant not in _grants(root, registry, session, turn): raise ValueError("source execution grant changed before launch") - requirement = ( - "\nOriginal owner inbox request " + receipt["request_id"] + ": read its original context, " - "independently acknowledge adopt/defer/reject, and return an audience-safe conclusion " - "for that exact request with manager-inbox report. A peer return alone does not reply " - "to the original conversation. Final prose is not a receipt." - ) - brief = {**request["brief"], "return_requirement": request["brief"]["return_requirement"] + requirement} - result = service.start(selected, operation, brief, - conversation={"session_id": session["session_id"], "turn_id": turn["turn_id"]}) + result = service.start(selected, operation, request["brief"], + conversation={"session_id": session["session_id"], "turn_id": turn["turn_id"]}, + source_request_id=receipt["request_id"]) return {"submitted": True, "operation_id": operation, "todo_id": binding["todo_id"], "status": result["status"], "runtime_readiness": preflight["state"], "replayed": False} except (OSError, ValueError, KeyError, TypeError, EffectRuntimeRejected): diff --git a/loopx/collaboration_mcp.py b/loopx/collaboration_mcp.py index 9ae53da6c9..6ea8612cf8 100644 --- a/loopx/collaboration_mcp.py +++ b/loopx/collaboration_mcp.py @@ -643,19 +643,28 @@ def authority_fault(exc: OSError | ValueError) -> dict[str, object]: def start(self, binding_id: str, operation_id: str, brief: dict, parent_request_id: str | None = None, *, conversation: dict | None = None, - confirmed_operation_id: str | None = None) -> dict: + confirmed_operation_id: str | None = None, + source_request_id: str | None = None) -> dict: """Start or replay one bound operation. ``conversation`` is supplied only by the trusted Chat host, never by the model: the session and Turn that started the operation. It is kept on first creation and never replaced, so a later wake returns to - that conversation and no other. + that conversation and no other. ``source_request_id`` is also host-only: + the receiver's existing inbox request, separate from the user brief. """ binding = self.binding(binding_id, require_active=True) require_operation_id(operation_id) brief = normalize_request({"goal_id": self.goal_id, "agent_id": binding["agent_id"], "brief": brief})["brief"] if any(item.get("delegation", {}).get("operation_id") == operation_id for item in brief["inputs"]): raise ValueError("delegation cannot depend on itself") + if source_request_id is not None: + with collaboration_goal_scope(self.registry, goal_id=self.goal_id, agents=(), + caller_goal_ref=self._caller_goal_ref()) as scope: + original = _entry(self.root, self.goal_id, binding["agent_id"], source_request_id, scope=scope) + decide_collaboration_lifecycle(scope, operation="history_inspect", record=original) + if original.get("brief") != brief: + raise ValueError("source request brief does not match delegated work") path = self.path(operation_id) with exclusive_file_lock(path.with_suffix(".dispatch")): exists = path.exists() @@ -665,6 +674,8 @@ def start(self, binding_id: str, operation_id: str, brief: dict, binding["agent_id"], operation_id, brief, parent_request_id, caller_goal_ref=self._caller_goal_ref()) identity = {"binding": binding, "request_id": delivered["request_id"], "operation_id": operation_id} + if source_request_id is not None: + identity["source_request_id"] = source_request_id if confirmed_operation_id is not None: # Internal callback adapter only: a canonical locator/CAS fence, # not an executor identity or domain execution permission. @@ -1504,16 +1515,24 @@ def _delegation_bootstrap(self, row: dict, binding: dict) -> dict: operation="history_inspect", record=entry, ) + source_id = row["identity"].get("source_request_id") + source_instruction = ( + " Original owner inbox request " + source_id + ": read its original context, " + "independently acknowledge adopt/defer/reject, and return an audience-safe conclusion " + "for that exact request with manager-inbox report. A peer return alone does not reply " + "to the original conversation. Final prose is not a receipt." + ) if source_id else "" return { "request_id": request_id, "brief": entry["brief"], + **({"source_request_id": source_id} if source_id else {}), "instruction": ( "Use the loopx_delegation tools to read_context and call " "assess_request for this request before working. If you adopt " "it, call return_result with the evidence-backed conclusion " "after validation. Final-answer prose alone is not an adoption " "or return receipt." - ), + ) + source_instruction, } def _write_delegation_bootstrap(self, row: dict, binding: dict) -> None: diff --git a/loopx/control_plane/collaboration/source_grant_observation.py b/loopx/control_plane/collaboration/source_grant_observation.py index 042100e4ef..138c1646d2 100644 --- a/loopx/control_plane/collaboration/source_grant_observation.py +++ b/loopx/control_plane/collaboration/source_grant_observation.py @@ -119,7 +119,7 @@ def source_context_target_authority( return _source_context_grant(runtime_root, session, turn, [target]) -def _source_registry_recipients(registry_path: Path | None, *, context_only: bool = False) -> dict: +def _source_registry_recipients(registry_path: Path | None, *, context_only: bool = False) -> dict[str, Any]: if registry_path is None: raise ValueError("context source registry unavailable") registry = load_project_registry(registry_path) @@ -156,16 +156,17 @@ def _source_registry_recipients(registry_path: Path | None, *, context_only: boo def source_execution_bindings( - runtime_root: Path, registry_path: Path | None, session: dict, turn: dict -) -> dict: + runtime_root: Path, registry_path: Path | None, session: dict[str, Any], turn: dict[str, Any] +) -> dict[str, Any]: """Observe one compatible registry; TS owns sender and exact binding grants.""" if conversation_scope(session, origin=turn.get("origin", "unknown"))["kind"] != "external_audience": return {"bindings": []} observed = _source_registry_recipients(registry_path) ingress, source = external_source_policy(runtime_root, session, turn) - return effect_runtime_result("collaboration.source.execution_bindings", { + result: dict[str, Any] = effect_runtime_result("collaboration.source.execution_bindings", { "source": source, "sender_id": ingress["sender_id"], "available": observed["available"], }) + return result def source_context_authority( diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index 4301230502..f5c33e795f 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -21,7 +21,7 @@ from test_independent_delegation_validation import independent_binding -def source(root, registry, *, goal_id, agent_id, requester, binding, source_id="lark:exact-message"): +def source(root, registry, *, goal_id, agent_id, requester, binding, source_id="lark:exact-message", semantic_brief=None): store = ChatSessionStore(root) session = store.create_session(goal_id="loopx-manager", agent_id="codex", adapter_kind="codex_app_server", upstream_thread_id="original", @@ -36,7 +36,7 @@ def source(root, registry, *, goal_id, agent_id, requester, binding, source_id=" register_ingress(root, session_id=session["session_id"], client_turn_id=turn["client_turn_id"], channel=session["channel_id"], sender_id="owner", message=turn["message"], source_id=source_id) - request = {"goal_id": goal_id, "agent_id": agent_id, "execution_binding_id": binding, "brief": brief()} + request = {"goal_id": goal_id, "agent_id": agent_id, "execution_binding_id": binding, "brief": brief() if semantic_brief is None else semantic_brief} receipt = deliver(root, registry, session=session, turn=turn, request=request) return store, session, turn, request, receipt, policy @@ -94,7 +94,8 @@ def test_exact_catalog_and_launch_keep_original_conversation_and_operation(flow) assert dispatch(flow)["status"] == "prepared" assert started[0][1] == "context-" + receipt["request_id"] assert started[0][3]["conversation"] == {"session_id": session["session_id"], "turn_id": turn["turn_id"]} - assert receipt["request_id"] in started[0][2]["return_requirement"] + assert started[0][2] == request["brief"] + assert started[0][3]["source_request_id"] == receipt["request_id"] assert dispatch(flow)["replayed"] assert len(started) == 1 @@ -204,7 +205,22 @@ def preview(*_): assert "尚未启动执行" in execution.handoff_message(receipt, {"submitted": False}) -def test_governed_worker_adopts_original_request_and_returns_without_another_model_turn(delegation_service): # noqa: F811 +def budget_brief(boundary): + semantic = brief() + if boundary == "return_field": + semantic["return_requirement"] = "r" * 2000 + else: + semantic["context"] = "c" * 6000 + semantic["constraints"] = ["x" * 1000] * 8 + semantic["return_requirement"] = "r" * 1000 + current = len(json.dumps(semantic, ensure_ascii=False, separators=(",", ":")).encode()) + semantic["constraints"].append("x" * (15800 - current - 3)) + assert len(json.dumps(semantic, ensure_ascii=False, separators=(",", ":")).encode()) == 15800 + return semantic + + +@pytest.mark.parametrize("boundary", ["short", "return_field", "encoded_total"]) +def test_governed_worker_adopts_original_request_and_returns_without_another_model_turn(delegation_service, boundary): # noqa: F811 root, service = delegation_service independent_binding(delegation_service) registry = service.registry @@ -216,7 +232,7 @@ def test_governed_worker_adopts_original_request_and_returns_without_another_mod data["goals"][0]["spawn_policy"] = {"execution_config": ".loopx/config/delegations.json"} registry.write_text(json.dumps(data)) store, session, turn, request, receipt, _ = source(service.root, registry, goal_id=service.goal_id, - agent_id="analyst", requester="lead", binding="analysis") + agent_id="analyst", requester="lead", binding="analysis", semantic_brief=budget_brief(boundary)) # The fixture receiver (not the Chat caller) reads, decides and returns the # original owner request as well as its separately accepted peer result. host = root / "fixture-host.py" @@ -225,8 +241,11 @@ def test_governed_worker_adopts_original_request_and_returns_without_another_mod statement = ( "from loopx.control_plane.collaboration.peers import read_inbox\n" "from loopx.capabilities.manager_context.roundtrip import report\n" - "for item in read_inbox(root / 'runtime', root / 'registry.json', envelope['goal_id'], actor)['items']:\n" - " if item.get('source_kind') != 'peer':\n" + "source_id = delegation['source_request_id']\n" + "items = read_inbox(root / 'runtime', root / 'registry.json', envelope['goal_id'], actor)['items']\n" + "assert any(item['request_id'] == source_id for item in items)\n" + "for item in items:\n" + " if item['request_id'] == source_id:\n" " acknowledge(root / 'runtime', envelope['goal_id'], actor, item['request_id'], 'adopt', 'Receiver independently read the original scope.')\n" " report(root / 'runtime', envelope['goal_id'], actor, item['request_id'], 'conclusion', 'Independent fixture result returned to the original request.')\n" ) @@ -259,3 +278,78 @@ def transport(route, *_args, **_kwargs): request=new_request, receipt=new_receipt, execution_allowed=lambda: True) assert not refused["submitted"] and refused["reason"] == "execution_not_launchable", refused assert (Path(worker.binding("analysis")["workspace"]) / "host-invocations").read_text() == "1" + + +@pytest.mark.parametrize("boundary", ["return_field", "encoded_total"]) +def test_legal_brief_budget_survives_real_chat_dispatch(delegation_service, monkeypatch, boundary): # noqa: F811 + """Internal return routing must not consume a caller's semantic budget.""" + from loopx.control_plane.collaboration.inbox import normalize_request, _read + + root, service = delegation_service + independent_binding(delegation_service) + data = json.loads(service.registry.read_text()) + project = Path(data["goals"][0]["repo"]) + config = project / ".loopx/config/delegations.json" + config.parent.mkdir(parents=True, exist_ok=True) + config.write_bytes(service.config.read_bytes()) + data["goals"][0]["spawn_policy"] = {"execution_config": ".loopx/config/delegations.json"} + service.registry.write_text(json.dumps(data)) + semantic = budget_brief(boundary) + semantic = normalize_request({"goal_id": service.goal_id, "agent_id": "analyst", "brief": semantic})["brief"] + store, session, turn, request, _, _ = source(service.root, service.registry, + goal_id=service.goal_id, agent_id="analyst", requester="lead", binding="analysis", semantic_brief=semantic) + # Let the production Chat adapter deliver and call the real start owner, + # but keep this admission oracle separate from the actual worker test. + started = [] + monkeypatch.setattr(execution.Delegations, "_spawn", lambda _, operation: started.append(operation)) + from loopx.chat_runtime import ChatRuntimeController + import loopx.chat_manager_context as manager_context + controller = ChatRuntimeController(store=store, codex_bin="codex", registry_path=service.registry, + manager_scope_resolver=lambda _: [service.goal_id]) + store.update_session(session["session_id"], manager_authorization_scope_id="fixture-scope") + monkeypatch.setattr(manager_context, "collect_manager_turn_context", lambda *_, **__: { + "coverage": {}, "goals": [], "authorization_scope_id": "fixture-scope"}) + model_calls = [] + class Adapter: + upstream_thread_id = "fixture-upstream" + def start_turn(self, message, sink): + assert "context_execution" in message and '"binding_id": "analysis"' in message + model_calls.append(message) + return {"context_handoff": request, "proposals": [], "gate": None, "message": "Preparing work"} + def close_session(self): + pass + try: + controller._run_turn(session_id=session["session_id"], turn_id=turn["turn_id"], + message=turn["message"], attachments=[], adapter=Adapter()) + completed = store.load_turn(session["session_id"], turn["turn_id"]) + assert completed["status"] == "completed", completed + response = completed["response"] + assert len(model_calls) == 1 + finally: + controller.close() + assert response["context_handoff_receipt"]["status"] == "delivered" + assert response["context_execution"]["submitted"], response + worker, binding = execution._service(service.root, service.registry, {"goal_id": service.goal_id, + "agent_id": "analyst", "requester_agent_id": "lead", "binding_id": "analysis"}) + operation = response["context_execution"]["operation_id"] + row = _read(worker.path(operation)) + bootstrap = worker._delegation_bootstrap(row, binding) + assert bootstrap["brief"] == request["brief"] + assert response["context_handoff_receipt"]["request_id"] in bootstrap["instruction"] + assert started == [operation] + assert bootstrap["source_request_id"] == response["context_handoff_receipt"]["request_id"] + before = worker.path(operation).read_bytes() + # Another legitimate request with the same brief cannot replace the cause + # of this operation; invalid/mismatched references fail before dispatch. + _, _, _, _, other, _ = source(service.root, service.registry, goal_id=service.goal_id, + agent_id="analyst", requester="lead", binding="analysis", source_id="lark:other", semantic_brief=semantic) + with pytest.raises(ValueError, match="identity conflict"): + worker.start("analysis", operation, semantic, source_request_id=other["request_id"]) + with pytest.raises(ValueError, match="invalid context request id"): + worker.start("analysis", "invalid-source", semantic, source_request_id="../outside") + with pytest.raises(ValueError, match="source request brief"): + worker.start("analysis", "mismatched-source", brief(), source_request_id=other["request_id"]) + assert worker.path(operation).read_bytes() == before + assert started == [operation] + assert not worker.path("invalid-source").exists() and not worker.path("mismatched-source").exists() + assert not (Path(binding["workspace"]) / "host-invocations").exists() From f0cac9a372bbc23e4d16752be9fed3ba5c923021 Mon Sep 17 00:00:00 2001 From: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> Date: Mon, 5 Oct 2026 23:04:48 +0800 Subject: [PATCH 9/9] fix(manager-context): omit empty execution prompt projection Signed-off-by: LoopX Agent <337587101+loopx-agent@users.noreply.github.com> --- loopx/capabilities/manager_context/inspection.py | 4 +++- tests/test_manager_context_execution.py | 5 +++++ 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/loopx/capabilities/manager_context/inspection.py b/loopx/capabilities/manager_context/inspection.py index 677c5a41a0..b6ac034f29 100644 --- a/loopx/capabilities/manager_context/inspection.py +++ b/loopx/capabilities/manager_context/inspection.py @@ -164,6 +164,7 @@ def rejected_read_arguments(arguments: dict[str, Any]) -> list[str]: def manager_index(context: dict[str, Any]) -> dict[str, Any]: """A small directory, never a second mutable progress store.""" read_tool = CONTEXT_TOOL_NAME if context.get("scope") == "owner_goal" else TOOL_NAME + execution = context.get("context_execution") return { "schema_version": "manager_evidence_index_v1", "snapshot_id": context.get("snapshot_id"), @@ -191,7 +192,8 @@ def manager_index(context: dict[str, Any]) -> dict[str, Any]: if row.get("activation_state") != "stopped" ], "context_delegation": context.get("context_delegation"), - **({"context_execution": context["context_execution"]} if "context_execution" in context else {}), + **({"context_execution": execution} if isinstance(execution, dict) + and (execution.get("bindings") or execution.get("available") is False) else {}), "evidence_sources": context.get("evidence_sources", [])[:12], "evidence_source_count": len(context.get("evidence_sources", [])), "agent_discovery": {"tool": read_tool, "view": "agents", "scope": "permitted_registry", diff --git a/tests/test_manager_context_execution.py b/tests/test_manager_context_execution.py index f5c33e795f..1857136ad8 100644 --- a/tests/test_manager_context_execution.py +++ b/tests/test_manager_context_execution.py @@ -113,6 +113,11 @@ def test_model_receives_authorized_task_choices_after_manager_context_compaction disabled = manager_index({}) assert "context_execution" not in disabled assert "execution_binding_id" not in _turn_prompt(turn["message"], context_summary=json.dumps(disabled)) + assert manager_index({"context_execution": {"available": True, "bindings": []}}) == disabled + unavailable = manager_index({"context_execution": {"available": False, "bindings": [], + "reason": "execution_bindings_unavailable"}}) + assert unavailable["context_execution"]["available"] is False + assert "execution_binding_id" not in _turn_prompt(turn["message"], context_summary=json.dumps(unavailable)) def test_handoff_response_preserves_receipt_and_separate_execution_status(flow):