diff --git a/docs/runtime-decomposition/wave-3-independent-adapters.md b/docs/runtime-decomposition/wave-3-independent-adapters.md new file mode 100644 index 000000000..1125adba8 --- /dev/null +++ b/docs/runtime-decomposition/wave-3-independent-adapters.md @@ -0,0 +1,282 @@ +# Wave 3 independent adapters + +Continuation base: `8ae6ee43936bdc5fe1da1297f87fb7b56be4a6cc`, directly +above canonical `a80c164dbe3e8bde4fb29b45c5d1c61404f2fede`. +The read-only continuation audit reviewed that checkpoint, its callers and tests, +then used the following design for this slice. The original A–I inventory remains +in `wave-3-resource-identity.md`; this supplement specifies the independent +adapters and the adversarial corrections. Process/browser adapters are deferred. + +## A. Re-audit and implicit-resource inventory + +| Site | Observation and decision | +| --- | --- | +| `resources.intersect_roots`, `RequestAuthority.intersect`, `bind_request_authority`, `seal_task_authority` | Descendant intersection already checks the parent observation. The equal-root shortcut did not revalidate it. Validate both observations before any intersection result; a fresh descendant never renews a replaced parent. | +| Dispatcher empty-root exact-approval fallback | Proposal roots serve only to re-resolve and compare one captured operation. Never install them into request authority. Test restored versions 1/2, sibling/parent access, replay, aliases and request/owner/session changes. | +| Native read/write/edit/patch, navigation and media workspace paths | Canonical control-path denial omitted hardlinked control objects. Also deny observed device/inode aliases, private configuration/DB/index paths and background control files, including configured paths from loaded producers. Directory grep's ripgrep branch scans descendants without bound checks: use the existing per-file resolver before reading. Filter bound ls/glob results through the same resolver. Media source/destination resolution uses the same control-state denial. This does not introduce a media filesystem adapter. | +| `McpManager.connect_server`, successful connection registration, `call_tool` | Server ID and qualified tool are mutable connection selectors. Seal the actual connection, configured endpoint origin and opaque epoch; revalidate at transport. A bound call cannot reconnect/retry into another producer. No transport redesign. | +| `_MCP_TOOL_MAP`, qualified/bare email dispatch | Availability previously selected backend/fallback. Preserve native filesystem semantics; snapshot other configured backends at trusted admission and pin dispatch. Discovery never creates operation grants. | +| Scoped `AgentExecutionBridge`, TUI bridge, HTTP request bridge | Callback objects or validated endpoint configuration determine execution. Capture object/configuration identity and exact tool, not a local filesystem observation. HTTP bridge factory and admission must produce the same configuration identity. | +| `do_api_call`, registered integrations | Names/IDs resolve through mutable configuration. Resolve aliases uniquely, bind integration ID, origin and configuration epoch; use the ID during execution and compare the loaded configuration before HTTP work. Generic API grants do not authorize the configured integration inventory: explicit trusted backend scope or one exact approval is required. Paths may contain tokens, so serialize origins and opaque epochs, not URL paths. | +| Document handlers / active document | Context/global active ID or most-recent lookup occurred during execution. Resolve server context or owner-scoped latest once; pass exact ID/version/digest and normalized selector. Global active changes cannot select another record. | +| Attachment OCR / upload index | URI resolves through mutable owner/path/hash index. Capture owner-checked row identity and confined file observation; consume the captured path. Keep the upload producer's owner check, without administrator override. | +| Thread management / send / history searches | `current`, line/JSON ID aliases and history target must bind caller owner and invocation thread. Capture exact selected thread row; collection searches bind the owner namespace. Existing owner-filtered search/cache boundaries remain. | +| Notes / native memories | Prefix and title selection can choose the first row later. Resolve uniquely within owner scope and normalize full ID; exact lookup in bound execution. Capture DB revision or opaque private memory revision. | +| Vault configuration / CLI | Global config had no owner producer binding. Legacy unowned config refuses runtime access. Authenticated settings save establishes owner and drops legacy session material; subsequent runtime reads require that owner, endpoint/configuration observation and an item observed by the server search producer. Names/prefixes resolve uniquely in that owner/configuration catalog to an exact UUID. Unknown UUIDs cannot manufacture a record observation. No credential appears in identity. | +| Builtin memory / RAG MCP stores | Memory producer has a fixed configured owner. Bind that owner and reject another caller or an ownerless producer. Legacy builtin RAG has no owner contract and cannot acquire private scope from discovery; refuse its runtime identity. | +| Generic `app_api` loopback | Internal-token calls could bypass migrated record domains. Refuse those namespace paths, including encoded/relative path aliases; callers use dedicated resource-bound operations. This is a migration guard, not an expanded internal API capability. | + +Other owner domains (calendar/contact/research/task/dynamic-tool stores), opaque +native script semantics and unrelated internal API paths remain separate adapter +work. Their existing permission gates are not described as typed enforcement. +This slice does not make a whole-runtime containment or private-data claim. + +## B. Typed model + +`resources.py` owns the additive immutable contracts: + +* `NativeBackendResource`: fixed native namespace and exact tool. Availability + cannot replace it with an MCP filesystem. +* `ExternalResource`: backend namespace, configured server ID, credential-free + endpoint origin, exact tool ID, connection/configuration epoch and optional + producer owner. Always `external=true`, `contained=false`. +* `OwnedScope`: namespace, owner, invocation thread and either an explicit record + ID set or a server-granted owner collection. The collection is a typed scope, + not a wildcard model selector or a capability floor. +* `OwnedResource`: namespace/collection, owner, invocation thread, exact record + ID, observed revision and storage-thread linkage where applicable. + +Attachment bindings additionally carry the existing typed filesystem observation +under the owner's private upload root. Context adapters are in +`remote_resources.py` and `owned_resources.py`; they grant no operation names. + +## C. Normalized operation/resource binding + +`ExactOperation` retains the original normalized proposal. Backend bindings +capture that exact input, caller and request alongside the backend identity. +Owned bindings carry original operation plus server-normalized execution input, +record observations and document execution context. Approval serialization seals +normalized input digests without copying credential-bearing arguments into the +identity. Existing approval content/digest and one-use claim remain mandatory. + +Collection creation/search/list operations bind owner collection identity; +specific reads/mutations bind exact records. A restricted record set cannot admit +a collection operation. Native filesystem bindings keep all existing source and +destination rules; patch moves remain unsupported and fail before execution. + +## D. Validation flow + +1. Server semantic admission grants operations independently of the tool inventory. +2. Trusted authority construction snapshots backend resources for those grants + and admits relevant owner/thread scopes. Restored snapshots never run this + constructor's implicit sealing path. +3. Request binding, parent intersection, policy and TurnContract gates run first. +4. Resolve backend and record selectors centrally, or consume the proposal's + exact sealed identities. Compare ownership, request/thread and resource scope. +5. Revalidate observations before consuming the existing one-use approval and + again at dispatch/producer entry. Bind contexts with `finally` reset. +6. Execute normalized input on the pinned backend/record. MCP and integration + producers compare their actual connection/configuration at the call boundary. + +Filesystem checks remain pathname observations, not descriptor-relative atomic +execution. Inode reuse, concurrent path replacement after validation and DB +changes between observation and mutation remain limitations. Record revisions +identify selected state; they are not new Wave 4 evidence or effect claims. + +## E. Alias, rename and ownership rules + +Backend aliases must resolve uniquely to the approved server/configuration. A +changed endpoint, connection or alias fails before claim/effect. Document +active/latest and thread current selectors resolve once on the server; an +approval consumes the captured ID even when the current UI alias changes. Missing, +stale, conflicting or ambiguous records fail closed. Notes/memory prefixes cannot +fall through to another title/record during bound execution. + +Child scopes intersect exact backend identities and owned record sets. Session +continuations may rebind the invocation namespace under the existing trusted +continuation rules, retaining owner, record limits and backend observations; +they do not synthesize a record from copied history. Exact approvals may admit +only their captured operation for a non-inherited legacy authority; they never +install a general resource scope or widen a parent's record/backend scope. +Inherited proposals themselves must fit their originating operation, backend, +filesystem and record scopes. A later approval resumption that resets the existing +inherited marker cannot reconstruct an identity excluded at proposal time. +Private read identity grants no additional send/egress operation. + +## F. Integration points + +Authority construction/persistence/intersection; central dispatch; approval +proposal/digest; HTTP request bridge admission; MCP successful connection/call +boundary; integration alias/configuration lookup; document dispatch context; +attachment OCR; notes/native memory exact lookup; authenticated vault settings +and owner-bound vault search producers. `agent_loop` changes only forward existing runtime context to proposal +capture. No loop decomposition, containment redesign or lifecycle change. + +## G. Migration + +Authority snapshots become version 3. Versions 1/2 restore empty backend/owned +scope fields. Fixed local dispatch compatibility retains existing operation gates; +no legacy snapshot reconstructs an external backend or owned collection. Exact +proposal snapshots can admit one operation without renewing general authority. + +Remote connection identities expire on reconnect/restart; private configuration +epochs use an in-process keyed opaque identifier. Restored stale epochs refuse +execution and require fresh trusted admission. Legacy unowned vault/RAG and +unresolved MCP connections fail closed. No remote owner, resource containment or +semantic page claim is inferred from successful transport. +Vault record observations describe the last server search response. Configuration +changes or refreshed record observations invalidate sealed operations; this is +not fresh remote semantic verification or a CLI process/account lifecycle claim. + +## H. Required verification + +New regressions cover equal/subtree stale parent intersection through direct, +context and task callers; restored empty-root exact approvals; control-state +direct/relative/symlink/hardlink reads/writes/search; backend availability, exact +tool/selectors, reconnect/endpoint/alias changes, legacy restoration, child +intersection, credentials and external flags; owned record aliases, revisions, +owner/thread changes, narrow scopes, attachments, vault/native memory identities, +generic loopback bypasses and context cleanup on success/error/cancel/nesting. +Focused existing suites cover RequestAuthority, TurnContract transcription/OCR/ +tasks, approvals, nested invocation, filesystem confinement, MCP/bridge routing, +documents/uploads/history and owner-scoped stores. Validation results are recorded +below; no full repository suite is run. + +Final validation on the checkpoint tree: **2,435 passed, 2 skipped, 4 warnings** +across the 88 focused files below (56.60 seconds). The skips are the existing +`/tmp`-symlink platform case and a containment shortfall case when `RLIMIT_AS` +can be lowered. The full repository suite was not run. + +Tests used `/tmp/odysseus-wave3-validation/bin/python`, an isolated venv with +system site packages plus `bcrypt`, `pyotp`, `mcp<2` and `pypdfium2`. The command +was that interpreter followed by `-m pytest -q -rs --disable-warnings +--maxfail=10` and the exact file arguments below. Earlier overlapping targeted +runs are not added to the final count. + +Static gates passed with empty output: + +```sh +python3 -m compileall -q app.py core routes services src tests scripts +git diff --check +git grep -n -E '^(<<<<<<< |=======$|>>>>>>> )' || true +git ls-files -u +``` + +
+Exact focused test file arguments + +```text +tests/test_resource_identity.py +tests/test_owned_resource_identity.py +tests/test_remote_resource_identity.py +tests/test_request_authority.py +tests/test_tool_approvals.py +tests/test_tool_approval_single_action_scope.py +tests/test_tool_approval_task_scope.py +tests/test_workspace_confine.py +tests/test_tool_path_confinement.py +tests/test_path_confinement_boundary.py +tests/test_filesystem_tool_argument_validation.py +tests/test_code_nav_tools.py +tests/test_apply_patch_transaction.py +tests/test_execution_bridge.py +tests/test_production_external_bridge.py +tests/test_turn_contract.py +tests/test_turn_contract_read_operations.py +tests/test_turn_contract_integration.py +tests/test_agent_turn_contract_boundaries.py +tests/test_explicit_personal_turn_contract.py +tests/test_nested_invocation_ownership.py +tests/test_containment_contract.py +tests/test_containment_enforcement.py +tests/test_containment_process_tree.py +tests/test_native_execution_containment.py +tests/test_background_containment.py +tests/test_process_ownership.py +tests/test_bg_jobs_store.py +tests/test_bg_job_tools.py +tests/test_execution_filesystem_boundary.py +tests/test_mcp_manager.py +tests/test_mcp_reconnect_args.py +tests/test_mcp_text_error_normalization.py +tests/test_mcp_param_hint_hardening.py +tests/test_mcp_tool_params_in_prompt.py +tests/test_mcp_memory_owner_scope.py +tests/test_mcp_cache_invalidation.py +tests/test_multiple_mcp_servers_timeout.py +tests/test_mcp_dependency_compatibility.py +tests/test_builtin_mcp_bg_tasks.py +tests/test_builtin_mcp_pythonpath.py +tests/test_builtin_mcp_npx_cache.py +tests/test_mcp_add_server_args_validation.py +tests/test_manage_mcp_command_allowlist.py +tests/test_document_tool_owner_scope.py +tests/test_owned_document_query.py +tests/test_document_session_owner_scope.py +tests/test_active_document_mutation_guard.py +tests/test_native_document_stream.py +tests/test_document_followup_integrity.py +tests/test_document_active_restore.py +tests/test_attachment_refs.py +tests/test_upload_handler_atomicity.py +tests/test_upload_handler_cleanup.py +tests/test_upload_handler_rename_owner.py +tests/test_upload_routes_owner_scope.py +tests/test_resolve_upload_path_nondict.py +tests/test_personal_upload_isolation.py +tests/test_personal_upload_privilege.py +tests/test_extract_text_tool.py +tests/test_media_ingress.py +tests/test_session_tools_registry.py +tests/test_session_owner_attribution.py +tests/test_session_list_owner_scope.py +tests/test_session_endpoint_owner_scope.py +tests/test_session_search.py +tests/test_session_search_batch_fetch.py +tests/test_history_topics_owner_scope.py +tests/test_history_order_by_timestamp_regression.py +tests/test_history_db_fallback_hidden.py +tests/test_memory_owner_isolation.py +tests/test_memory_routes_session_owner.py +tests/test_manage_memory_json_contract.py +tests/test_manage_memory_list.py +tests/test_memory_store_unreadable_no_wipe.py +tests/test_manage_notes_search_contract.py +tests/test_notes_fail_closed_auth.py +tests/test_notes_checklist_state.py +tests/test_vault_password_not_in_argv.py +tests/test_vault_routes_shim.py +tests/test_external_context_tool_gate.py +tests/test_chat_route_tool_policy.py +tests/test_product_turn_contract_route.py +tests/test_native_tool_result_threading.py +tests/test_host_shell_polling.py +tests/test_integrations_url_join.py +tests/test_integration_api_call_ssrf.py +tests/test_integrations_api_call_truncation.py +``` + +
+ +## I. Wave 4 / Wave 5B collision boundaries + +Wave 4 retains durable claim, effects, evidence freshness, provenance and egress +policy. No private content is licensed for transfer by a resource identity. +Existing containment/browser receipts are not authority or semantic verification. + +Wave 5B must freeze the shared `ProcessIdentity` and lifecycle API before these +seams are implemented: + +* Native `_run_owned_command` and process ownership checks: consume the producer's + verified process identity and lifecycle namespace/incarnation, linking the + admitted execution backend/root and containment receipt without granting scope. +* `bg_jobs.launch/get/kill`, monitor continuations and authority sidecars: link + the durable owner/thread/job identity to that same verified lifecycle identity + and receipt. A model job ID or restored PID never reconstructs it. +* Browser lifecycle `session_for`/receipt and private/MCP browser producers: + consume the frozen producer/process lifecycle identity, then bind owner/thread, + browser session incarnation and page/navigation observations separately. + Producer liveness is not verification of remote page meaning. + +This continuation implements none of those adapters and creates no parallel +`ProcessIdentity`. Existing inert process/browser types are unchanged. diff --git a/routes/chat_routes.py b/routes/chat_routes.py index 4ef072afa..85b89eba9 100644 --- a/routes/chat_routes.py +++ b/routes/chat_routes.py @@ -814,10 +814,13 @@ def _external_execution_bridge( raise ValueError("external execution bridge returned an invalid payload") return str(payload.get("description") or tool), payload["result"] + from src.agent_runtime.remote_resources import configuration_incarnation return AgentExecutionBridge( route_tool=route_tool, supported_tools=supported, name="request_local_http", + endpoint_id=url, + configuration_id=configuration_incarnation((url, token, tuple(sorted(supported)))), ) @@ -3418,6 +3421,7 @@ def setup_chat_routes( _request_authority = request_authority_for_http( request, message, owner=_user, session_id=session, workspace=workspace, history=_turn_history, policy=tool_policy, + client_runtime_context=client_runtime_context, active_document=bool(active_doc), image_attachment=any(str(a.get('mime') or '').startswith('image/') for a in (ctx.preprocessed.attachment_meta or [])), diff --git a/routes/vault/vault_routes.py b/routes/vault/vault_routes.py index 7e97500f0..88cd625d9 100644 --- a/routes/vault/vault_routes.py +++ b/routes/vault/vault_routes.py @@ -13,6 +13,7 @@ import asyncio from pathlib import Path from datetime import datetime from fastapi import APIRouter, Request +from fastapi import HTTPException from pydantic import BaseModel from core.middleware import require_admin @@ -77,6 +78,19 @@ def _save_config(cfg: dict): safe_chmod(str(VAULT_FILE), 0o600) +def _bind_config_owner(cfg: dict, request: Request): + from src.auth_helpers import effective_user + from src.owner_identity import effective_storage_owner + owner = effective_storage_owner(effective_user(request)) + if not owner or (cfg.get("owner") and cfg["owner"] != owner): + raise HTTPException(403, "Vault configuration requires its explicit owner") + if not cfg.get("owner"): + # Legacy credentials cannot silently acquire a new ownership binding. + cfg.pop("session", None) + cfg.pop("unlocked_at", None) + cfg["owner"] = owner + + async def _run_bw(args: list, session: str = None, input_text: str = None, bw_password: str = None) -> tuple: env = {} @@ -144,6 +158,7 @@ def setup_vault_routes(): """Save vault URL + email. Runs 'bw config server' to point at Vaultwarden.""" require_admin(request) cfg = _load_config() + _bind_config_owner(cfg, request) cfg["server_url"] = req.server_url.strip().rstrip("/") cfg["email"] = req.email.strip() diff --git a/src/agent_loop.py b/src/agent_loop.py index f9130a1e6..e87f789c0 100644 --- a/src/agent_loop.py +++ b/src/agent_loop.py @@ -32769,6 +32769,7 @@ async def stream_agent_loop( ), request_text=_last_user, request_authority=active_request_authority(), + client_runtime_context=client_runtime_context, ) desc = f"{block.tool_type}: APPROVAL REQUIRED" result = { diff --git a/src/agent_runtime/authority.py b/src/agent_runtime/authority.py index b86fd161f..1ffa9adad 100644 --- a/src/agent_runtime/authority.py +++ b/src/agent_runtime/authority.py @@ -11,7 +11,10 @@ from pathlib import Path import re from uuid import uuid4 -from src.agent_runtime.resources import FilesystemRoot, intersect_roots +from src.agent_runtime.resources import ( + FilesystemRoot, ExternalResource, NativeBackendResource, OwnedScope, + backend_from_dict, intersect_roots, seal_owned_scopes, +) from src.tool_policy import ToolPolicy, build_effective_tool_policy from src.turn_contract import ( FAMILY_TOOLS, canonical_tool, requested_capabilities, @@ -115,6 +118,8 @@ class RequestAuthority: # None is only the trusted constructor's instruction to seal a workspace. # Persisted/child authorities always carry an explicit tuple, including (). resource_roots: tuple[FilesystemRoot, ...] | None = None + backend_resources: tuple[ExternalResource | NativeBackendResource, ...] | None = None + owned_scopes: tuple[OwnedScope, ...] | None = None def __post_init__(self): if (not isinstance(self.request_id, str) or not self.request_id @@ -138,11 +143,24 @@ class RequestAuthority: or any(not isinstance(r, FilesystemRoot) or (r.owner and r.owner != self.owner) for r in self.resource_roots)): raise ValueError("Malformed request resource roots") + if self.backend_resources is None: + from src.agent_runtime.remote_resources import seal_backends + object.__setattr__(self, "backend_resources", seal_backends((g.tool for g in self.grants), owner=self.owner)) + if self.owned_scopes is None: + object.__setattr__(self, "owned_scopes", seal_owned_scopes( + self.owner, self.session_id, (g.tool for g in self.grants))) + if (not isinstance(self.backend_resources, tuple) + or any(not isinstance(r, (ExternalResource, NativeBackendResource)) + or (isinstance(r, ExternalResource) and r.owner and r.owner != self.owner) for r in self.backend_resources) + or not isinstance(self.owned_scopes, tuple) + or any(not isinstance(s, OwnedScope) or (s.owner, s.thread_id) != (self.owner, self.session_id) + for s in self.owned_scopes)): + raise ValueError("Malformed backend or owned resource scope") @classmethod def empty(cls, *, owner=None, session_id=None, workspace=None): return cls(uuid4().hex, _owner(owner), str(session_id or ""), str(workspace or ""), - resource_roots=()) + resource_roots=(), backend_resources=(), owned_scopes=()) def bound_to(self, *, owner=None, session_id=None, workspace=None): return (self.owner == _owner(owner) and self.session_id == str(session_id or "") @@ -168,35 +186,44 @@ class RequestAuthority: raise TypeError("Child authority must be server-owned RequestAuthority") grants = [] roots = () + backends = () + owned = () if (self.owner, self.session_id, self.workspace) == (child.owner, child.session_id, child.workspace): theirs = {g.tool: g for g in child.grants} grants = [g.intersect(theirs[g.tool]) for g in self.grants if g.tool in theirs] roots = intersect_roots(self.resource_roots, child.resource_roots) + backends = tuple(r for r in self.backend_resources if r in child.backend_resources) + owned = tuple(s for left in self.owned_scopes for right in child.owned_scopes + if (s := left.intersect(right)) is not None) return replace(self, grants=tuple(grants), denied=self.denied | child.denied, block_all=self.block_all or child.block_all, disable_mcp=self.disable_mcp or child.disable_mcp, inherited=True, - resource_roots=roots) + resource_roots=roots, backend_resources=backends, owned_scopes=owned) def continuation(self, *, owner=None, session_id=None): """A server continuation may rebind a session, never change owner/grants.""" if self.owner != _owner(owner): return RequestAuthority.empty(owner=owner, session_id=session_id) - return replace(self, session_id=str(session_id or ""), inherited=True) + rebound = str(session_id or "") + return replace(self, session_id=rebound, inherited=True, + owned_scopes=tuple(replace(s, thread_id=rebound) for s in self.owned_scopes) if rebound else ()) def to_dict(self): - return {"version": 2, "request_id": self.request_id, "owner": self.owner, + return {"version": 3, "request_id": self.request_id, "owner": self.owner, "session_id": self.session_id, "workspace": self.workspace, "grants": [{"tool": g.tool, "actions": None if g.actions is None else sorted(g.actions), "inputs": None if g.inputs is None else sorted(g.inputs)} for g in self.grants], "denied": sorted(self.denied), "block_all": self.block_all, "disable_mcp": self.disable_mcp, "inherited": self.inherited, - "resource_roots": [r.to_dict() for r in self.resource_roots]} + "resource_roots": [r.to_dict() for r in self.resource_roots], + "backend_resources": [r.to_dict() for r in self.backend_resources], + "owned_scopes": [s.to_dict() for s in self.owned_scopes]} @classmethod def from_dict(cls, value): if (not isinstance(value, dict) or type(value.get("version")) is not int - or value["version"] not in {1, 2}): + or value["version"] not in {1, 2, 3}): raise ValueError("Unsupported authority snapshot") def limits(value): if value is None: @@ -204,14 +231,19 @@ class RequestAuthority: if not isinstance(value, list) or any(not isinstance(v, str) for v in value): raise ValueError("Malformed authority limits") return frozenset(value) - roots = value["resource_roots"] if value["version"] == 2 else [] + roots = value["resource_roots"] if value["version"] >= 2 else [] if not isinstance(roots, list): raise ValueError("Malformed request resource snapshot") + backends = value["backend_resources"] if value["version"] == 3 else [] + owned = value["owned_scopes"] if value["version"] == 3 else [] + if not isinstance(backends, list) or not isinstance(owned, list): + raise ValueError("Malformed request resource scope snapshot") return cls(value["request_id"], value["owner"], value["session_id"], value["workspace"], tuple(OperationGrant(g["tool"], limits(g["actions"]), limits(g["inputs"])) for g in value["grants"]), limits(value["denied"]), value["block_all"], value["disable_mcp"], value["inherited"], - tuple(FilesystemRoot.from_dict(r) for r in roots)) + tuple(FilesystemRoot.from_dict(r) for r in roots), + tuple(backend_from_dict(r) for r in backends), tuple(OwnedScope.from_dict(s) for s in owned)) _BROWSER_READ_ACTIONS = frozenset({"open", "navigate", "snapshot", "text", "read", "find", @@ -262,7 +294,7 @@ def interpret_request(request_text, *, history=(), workspace=None, active_docume def create_request_authority(request_text, *, owner=None, session_id=None, workspace=None, history=(), policy=None, active_document=False, - image_attachment=False, capabilities=None): + image_attachment=False, capabilities=None, client_runtime_context=None): """Deterministic server policy over semantic facts, never schema inventory.""" if not isinstance(request_text, str): raise TypeError("Authority requires trusted request text") @@ -298,6 +330,10 @@ def create_request_authority(request_text, *, owner=None, session_id=None, works grants.append(OperationGrant(name, actions, inputs)) authority = RequestAuthority(uuid4().hex, _owner(owner), str(session_id or ""), str(workspace or ""), tuple(grants)) + if client_runtime_context is not None: + from src.agent_runtime.remote_resources import seal_backends + authority = replace(authority, backend_resources=seal_backends( + (g.tool for g in authority.grants), context=client_runtime_context, owner=authority.owner)) return authority.restrict(policy or build_effective_tool_policy(last_user_message=request_text)) @@ -385,7 +421,8 @@ def with_request_authority(func): owner=parameters.get("owner"), session_id=parameters.get("session_id"), workspace=parameters.get("workspace"), history=getattr(parameters.get("history_session"), "history", ()) or (), - active_document=bool(parameters.get("active_document"))) + active_document=bool(parameters.get("active_document")), + client_runtime_context=parameters.get("client_runtime_context")) if not isinstance(authority, RequestAuthority): raise TypeError("Missing or malformed server request authority") if parent is None and parameters.get("exact_approval") is not None: @@ -423,7 +460,9 @@ def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authori if parent is not None: authority = parent.intersect(replace(authority, session_id=parent.session_id, workspace=parent.workspace, - resource_roots=parent.resource_roots)) + resource_roots=parent.resource_roots, + backend_resources=parent.backend_resources, + owned_scopes=parent.owned_scopes)) return _json({"task_input": [prompt, task_type, action], "authority": authority.to_dict()}) diff --git a/src/agent_runtime/owned_resources.py b/src/agent_runtime/owned_resources.py new file mode 100644 index 000000000..676870381 --- /dev/null +++ b/src/agent_runtime/owned_resources.py @@ -0,0 +1,451 @@ +"""Resolve owned selectors before execution and consume exact server identities.""" +from contextlib import contextmanager +from contextvars import ContextVar +from dataclasses import dataclass +import json +import re +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from src.agent_runtime.authority import ExactOperation + +from src.agent_runtime.resources import ( + FilesystemResource, FilesystemRoot, FilesystemScope, OwnedResource, + OWNED_TOOL_NAMESPACES, ResourceIdentityError, +) + + +def _args(content): + if not isinstance(json.loads(content or "{}"), dict): + raise ResourceIdentityError("Owned resource arguments must be an object") + from src.tools._common import _parse_tool_args + value = _parse_tool_args(content) + if not isinstance(value, dict): + raise ResourceIdentityError("Owned resource arguments must be an object") + return dict(value) + + +def _selector(args, keys): + values = [args[k] for k in keys if k in args and args[k] not in (None, "")] + if any(not isinstance(v, str) or not v.strip() for v in values): + raise ResourceIdentityError("Record selectors must be strings") + values = [v.strip() for v in values] + if len(set(values)) > 1: + raise ResourceIdentityError("Conflicting record aliases") + return values[0] if values else "" + + +def _revision(row, namespace): + created = getattr(row, "created_at", None) + updated = getattr(row, "updated_at", None) + if created is None or not hasattr(created, "isoformat") or updated is None or not hasattr(updated, "isoformat"): + raise ResourceIdentityError("Record has no observable revision") + version = getattr(row, "version_count", "") if namespace == "documents" else "" + if namespace == "documents" and type(version) is not int: + raise ResourceIdentityError("Document version is unresolved") + return f"{created.isoformat()}:{updated.isoformat()}:{version}" + + +def _record(namespace, owner, thread, row): + if (getattr(row, "owner", None) != owner or not isinstance(getattr(row, "id", None), str) + or row.id in {"", "*"}): + raise ResourceIdentityError("Record ownership is unresolved") + linked = str(getattr(row, "session_id", "") or "") if namespace == "documents" else row.id if namespace == "threads" else "" + return OwnedResource(namespace, owner, thread, namespace, row.id, _revision(row, namespace), linked) + + +def _row(namespace, identifier, owner): + from core.database import SessionLocal, Document, Session, Note + model = {"documents": Document, "threads": Session, "notes": Note}[namespace] + db = SessionLocal() + try: + row = db.query(model).filter(model.id == identifier, model.owner == owner).first() + if row is None or (namespace == "documents" and not row.is_active): + raise ResourceIdentityError("Owned record is missing or inaccessible") + db.expunge(row) + return row + finally: + db.close() + + +@dataclass(frozen=True) +class AttachmentResource: + record: OwnedResource + file: FilesystemResource + + def __post_init__(self): + if not isinstance(self.record, OwnedResource) or not isinstance(self.file, FilesystemResource) or self.file.root.owner != self.record.owner: + raise ValueError("Malformed attachment identity") + + def to_dict(self): + return {"record": self.record.to_dict(), "file": self.file.to_dict()} + + +def _attachment(identifier, owner, thread): + from src.tool_utils import get_upload_handler + handler = get_upload_handler() + if handler is None: + raise ResourceIdentityError("Attachment store is unavailable") + info = handler.resolve_upload(identifier, owner=owner, allow_admin=False) + if not isinstance(info, dict) or info.get("id") != identifier or info.get("owner") != owner: + raise ResourceIdentityError("Attachment ownership is unresolved") + root = FilesystemRoot.seal(handler.upload_dir, scope=FilesystemScope.PRIVATE, owner=owner) + file = FilesystemResource.resolve(root, info.get("path")) + if file.identity.kind != "file": + raise ResourceIdentityError("Attachment must identify a file") + revision = str(info.get("checksum_sha256") or info.get("hash") or info.get("uploaded_at") or "") + if not revision: + raise ResourceIdentityError("Attachment has no observable revision") + return AttachmentResource(OwnedResource("attachments", owner, thread, "attachments", identifier, revision), file) + + +_VAULT_RECORDS = {} + + +def _vault_revision(cfg, owner): + from src.agent_runtime.remote_resources import endpoint_identity, configuration_incarnation + if not isinstance(cfg, dict) or cfg.get("owner") != owner: + raise ResourceIdentityError("Vault configuration has no matching explicit owner") + endpoint = endpoint_identity(cfg.get("server_url") or cfg.get("url") or "") + return endpoint + ":" + configuration_incarnation((cfg.get("server_url") or cfg.get("url"), cfg.get("email"), cfg.get("unlocked_at"), cfg.get("session"))) + + +def observe_vault_records(owner, cfg, records): + """Only a server search response produces record observations, not grants.""" + from src.tools.vault import _load_vault_config + from src.agent_runtime.remote_resources import configuration_incarnation + from uuid import UUID + revision = _vault_revision(cfg, owner) + if _vault_revision(_load_vault_config(), owner) != revision or not isinstance(records, list): + raise ResourceIdentityError("Vault producer configuration changed") + observed = {} + for row in records: + if not isinstance(row, dict): + raise ResourceIdentityError("Malformed vault producer record") + try: + identifier = str(UUID(row.get("id", ""))) + except (ValueError, TypeError, AttributeError) as error: + raise ResourceIdentityError("Vault producer record has no exact UUID") from error + if identifier in observed or not isinstance(row.get("name", ""), str): + raise ResourceIdentityError("Ambiguous vault producer identity") + observed[identifier] = (row.get("name", ""), configuration_incarnation(json.dumps(row, sort_keys=True, allow_nan=False))) + catalog = _VAULT_RECORDS.setdefault((owner, revision), {}) + catalog.update(observed) + + +def _vault_resource(owner, thread, identifier): + from src.tools.vault import _load_vault_config + revision = _vault_revision(_load_vault_config(), owner) + if identifier != "*": + record = _VAULT_RECORDS.get((owner, revision), {}).get(identifier) + if record is None: + raise ResourceIdentityError("Vault record has no server observation; search the owner vault first") + revision += ":" + record[1] + return OwnedResource("vault", owner, thread, "vault", identifier, revision) + + +def _vault_selector(owner, selector): + from src.tools.vault import _load_vault_config + revision = _vault_revision(_load_vault_config(), owner) + rows = _VAULT_RECORDS.get((owner, revision), {}) + if selector in rows: + return selector + matches = [identifier for identifier, (name, _) in rows.items() + if identifier.startswith(selector) or name == selector] + if not selector or len(matches) != 1: + raise ResourceIdentityError("Vault selector is missing or ambiguous") + return matches[0] + + +def _memory_record(identifier, owner, thread, *, prefix=False): + from src.ai_interaction import _memory_manager + if _memory_manager is None: + raise ResourceIdentityError("Memory store is unavailable") + rows = [row for row in _memory_manager.load(owner=owner) if isinstance(row, dict) + and row.get("owner") == owner and isinstance(row.get("id"), str) + and (row["id"].startswith(identifier) if prefix else row["id"] == identifier)] + if len(rows) != 1 or not identifier or rows[0].get("timestamp") is None or rows[0]["id"] in {"", "*"}: + raise ResourceIdentityError("Memory selector is missing or ambiguous") + from src.agent_runtime.remote_resources import configuration_incarnation + row = rows[0] + # A same-second edit still changes the private revision without serializing content. + revision = configuration_incarnation(json.dumps(row, sort_keys=True, allow_nan=False)) + return OwnedResource("memory", owner, thread, "memory", row["id"], revision) + + +@dataclass(frozen=True) +class BoundOwnedOperation: + operation: "ExactOperation" + execution_input: str + request_id: str + owner: str + thread_id: str + resources: tuple[OwnedResource, ...] + attachments: tuple[AttachmentResource, ...] = () + document_id: str = "" + document_version: int | None = None + document_digest: str = "" + + def __post_init__(self): + from src.agent_runtime.authority import ExactOperation + if (not isinstance(self.operation, ExactOperation) or not self.owner or not self.thread_id + or any(not isinstance(v, str) for v in (self.execution_input, self.request_id, self.owner, self.thread_id, self.document_id, self.document_digest)) + or not isinstance(self.resources, tuple) or not self.resources + or any(not isinstance(r, OwnedResource) or (r.owner, r.thread_id) != (self.owner, self.thread_id) for r in self.resources) + or not isinstance(self.attachments, tuple) or any(not isinstance(a, AttachmentResource) for a in self.attachments)): + raise ValueError("Malformed owned resource operation") + namespace = OWNED_TOOL_NAMESPACES.get(self.operation.tool) + if (any(r.namespace != namespace or r.collection != namespace or (r.record_id != "*" and not r.revision) for r in self.resources) + or tuple(a.record for a in self.attachments) != tuple(r for r in self.resources if r.namespace == "attachments")): + raise ValueError("Malformed owned resource identity") + if self.document_id: + if (type(self.document_version) is not int or self.document_version < 1 + or not re.fullmatch(r"[0-9a-f]{64}", self.document_digest) + or not any(r.namespace == "documents" and r.record_id == self.document_id for r in self.resources)): + raise ValueError("Malformed document binding") + elif any(r.namespace == "documents" and r.record_id != "*" for r in self.resources): + raise ValueError("Missing document binding") + + def to_dict(self): + from src.agent_runtime.remote_resources import configuration_incarnation + return {"request_id": self.request_id, "owner": self.owner, "thread_id": self.thread_id, + "tool": self.operation.transport_tool, + "execution_input_digest": configuration_incarnation(self.execution_input), + "resources": [r.to_dict() for r in self.resources], + "attachments": [a.to_dict() for a in self.attachments], + "document_id": self.document_id, "document_version": self.document_version, + "document_digest": self.document_digest} + + def validate(self): + for resource in self.resources: + if resource.record_id == "*": + if resource.namespace == "vault" and _vault_resource(self.owner, self.thread_id, "*") != resource: + raise ResourceIdentityError("Vault identity changed") + continue + if resource.namespace == "attachments": + expected = next((a for a in self.attachments if a.record == resource), None) + if expected is None or _attachment(resource.record_id, self.owner, self.thread_id) != expected: + raise ResourceIdentityError("Attachment identity changed") + expected.file.validate() + elif resource.namespace == "vault": + if _vault_resource(self.owner, self.thread_id, resource.record_id) != resource: + raise ResourceIdentityError("Vault identity changed") + elif resource.namespace == "memory": + if _memory_record(resource.record_id, self.owner, self.thread_id) != resource: + raise ResourceIdentityError("Memory identity changed") + elif _record(resource.namespace, self.owner, self.thread_id, + _row(resource.namespace, resource.record_id, self.owner)) != resource: + raise ResourceIdentityError("Owned record identity changed") + + +def needs_owned_binding(operation): + if operation.tool == "app_api": + # The generic internal-token bridge must not bypass migrated owner + # namespaces. Dedicated tools carry their typed record operations. + from urllib.parse import unquote, urlsplit + import posixpath + args = _args(operation.input) + path = args.get("path", "") + if not isinstance(path, str): + raise ResourceIdentityError("Malformed internal resource selector") + for _ in range(4): + decoded = unquote(path) + if decoded == path: + break + path = decoded + if "%" in path or "\\" in path: + raise ResourceIdentityError("Unresolved internal resource selector") + path = posixpath.normpath(urlsplit(path).path) + private = {"document", "documents", "session", "sessions", "history", "chat", "chats", + "notes", "memory", "vault", "upload", "uploads", "attachments"} + segments = path.strip("/").split("/") + if len(segments) >= 2 and segments[0] == "api" and segments[1].casefold() in private: + raise ResourceIdentityError("Owned records require a dedicated resource-bound tool") + return False + if operation.tool not in OWNED_TOOL_NAMESPACES: + return False + if operation.tool in {"extract_text", "inspect_media", "transcribe_media"}: + return "odysseus://attachment/" in operation.input + return True + + +def resolve_owned_operation(operation, *, owner, thread_id, request_id="", document_id=None): + if not owner or not thread_id: + raise ResourceIdentityError("Owned operations require an owner and invocation thread") + if document_id is not None and (not isinstance(document_id, str) or not document_id.strip()): + raise ResourceIdentityError("Malformed server document selector") + namespace = OWNED_TOOL_NAMESPACES[operation.tool] + args = _args(operation.input) if operation.tool not in {"create_document", "edit_document", "update_document", "suggest_document", "send_to_session", "create_session", "list_sessions", "search_chats", "manage_session", "manage_memory"} else {} + execution_input = operation.input + resources = [] + attachments = [] + doc_id = "" + doc_version = None + doc_digest = "" + collection = lambda: OwnedResource(namespace, owner, thread_id, namespace, "*") + if namespace == "documents": + action = str(args.get("action") or "list").strip().lower() + if operation.tool == "create_document" or (operation.tool == "manage_documents" and action in {"list", "search", "find", "tidy"}): + resources.append(collection()) + else: + identifier = _selector(args, ("document_id", "id", "uid")) or document_id or "" + if identifier in {"active", "current"}: + if not document_id or document_id in {"active", "current", "latest"}: + raise ResourceIdentityError("Active document selector is unresolved") + identifier = document_id + if not identifier and operation.tool == "manage_documents" and action != "delete": + raise ResourceIdentityError("Document selector is required") + if not identifier or identifier == "latest": + from core.database import SessionLocal, Document + db = SessionLocal() + try: + row = db.query(Document).filter(Document.owner == owner, Document.is_active == True).order_by(Document.updated_at.desc(), Document.id).first() + identifier = row.id if row is not None else "" + finally: + db.close() + if not identifier: + raise ResourceIdentityError("Document selector is unresolved") + row = _row(namespace, identifier, owner) + resources.append(_record(namespace, owner, thread_id, row)) + doc_id, doc_version = row.id, row.version_count + from src.tool_approvals import document_content_digest + doc_digest = document_content_digest(row.current_content) + if operation.tool == "manage_documents": + for key in ("id", "uid"): + args.pop(key, None) + args["document_id"] = doc_id + execution_input = json.dumps(args, sort_keys=True) + elif namespace == "threads": + if operation.tool in {"list_sessions", "search_chats", "create_session"}: + resources.append(collection()) + else: + if operation.tool == "send_to_session": + identifier, _, message = operation.input.partition("\n") + identifier = identifier.strip() + else: + if operation.input.lstrip().startswith("{"): + args = _args(operation.input) + else: + lines = operation.input.strip().split("\n", 2) + args = {"action": lines[0], "session_id": lines[1] if len(lines) > 1 else ""} + if len(lines) > 2: + args["value"] = lines[2] + if args.get("action") == "list": + resources.append(collection()) + identifier = _selector(args, ("session_id", "session", "id")) + if not resources: + identifier = thread_id if identifier == "current" else identifier + row = _row(namespace, identifier, owner) + resources.append(_record(namespace, owner, thread_id, row)) + if operation.tool == "send_to_session": + execution_input = row.id + "\n" + message + else: + args.pop("id", None) + args.pop("session", None) + args["session_id"] = row.id + execution_input = json.dumps(args, sort_keys=True) + elif namespace == "notes": + action = str(args.get("action") or "").strip().lower().replace("-", "_") + if action in {"list", "search", "find", "add", "create", "new", "save", "remind"}: + resources.append(collection()) + else: + identifier = _selector(args, ("id", "note_id", "noteId")) + from core.database import SessionLocal, Note + db = SessionLocal() + try: + q = db.query(Note).filter(Note.owner == owner) + if identifier: + rows = q.filter(Note.id.startswith(identifier, autoescape=True)).limit(2).all() + else: + title = _selector(args, ("title", "query", "text")) + rows = q.filter(Note.title == title).limit(2).all() if title else [] + if len(rows) != 1: + raise ResourceIdentityError("Note selector is missing or ambiguous") + identifier = rows[0].id + finally: + db.close() + row = _row(namespace, identifier, owner) + resources.append(_record(namespace, owner, thread_id, row)) + args.pop("note_id", None) + args.pop("noteId", None) + args["id"] = identifier + execution_input = json.dumps(args, sort_keys=True) + elif namespace == "attachments": + selector = args.get("path") + match = re.fullmatch(r"odysseus://attachment/([A-Za-z0-9_-]+(?:\.[A-Za-z0-9]+)?)", selector or "") + if match is None: + raise ResourceIdentityError("Malformed attachment selector") + attachment = _attachment(match[1], owner, thread_id) + resources.append(attachment.record) + attachments.append(attachment) + elif namespace == "memory": + from src.ai_interaction import _manage_memory_lines + lines = _manage_memory_lines(operation.input) + if not lines: + raise ResourceIdentityError("Memory action is unresolved") + action = lines[0].strip().lower() + if action in {"list", "search", "add"}: + resources.append(collection()) + elif action in {"edit", "delete"} and len(lines) >= 2: + resource = _memory_record(lines[1].strip(), owner, thread_id, prefix=True) + resources.append(resource) + lines[1] = resource.record_id + execution_input = "\n".join(lines) + else: + raise ResourceIdentityError("Memory operation is unresolved") + elif namespace == "vault": + identifier = "*" + if operation.tool == "vault_get": + identifier = _vault_selector(owner, _selector(args, ("item_id",))) + args["item_id"] = identifier + execution_input = json.dumps(args, sort_keys=True) + resources.append(_vault_resource(owner, thread_id, identifier)) + bound = BoundOwnedOperation(operation, execution_input, request_id, owner, thread_id, + tuple(resources), tuple(attachments), doc_id, doc_version, doc_digest) + bound.validate() + return bound + + +def admit_owned_operation(authority, operation, *, document_id=None, approved=None, exact_admission=False): + bound = (approved if approved is not None else resolve_owned_operation(operation, owner=authority.owner, + thread_id=authority.session_id, request_id=authority.request_id, document_id=document_id)) + if (not isinstance(bound, BoundOwnedOperation) or bound.operation != operation + or (bound.owner, bound.thread_id) != (authority.owner, authority.session_id) + or (bound.request_id and bound.request_id != authority.request_id)): + raise ResourceIdentityError("Owned operation approval binding changed") + if not all(any(scope.permits(r) for scope in authority.owned_scopes) for r in bound.resources): + if not (approved is not None and exact_admission and not authority.inherited and not authority.owned_scopes): + raise ResourceIdentityError("Owned resource exceeds parent/request scope") + bound.validate() + return bound + + +_ACTIVE = ContextVar("owned_resource_operation", default=None) + + +def active_owned_operation(): + return _ACTIVE.get() + + +@contextmanager +def bind_owned_operation(operation): + if operation is not None: + if not isinstance(operation, BoundOwnedOperation): + raise TypeError("Owned operation must be server-owned") + operation.validate() + token = _ACTIVE.set(operation) + try: + yield operation + finally: + _ACTIVE.reset(token) + + +def bound_attachment_path(owner, selector): + operation = active_owned_operation() + if operation is None: + return None + operation.validate() + for attachment in operation.attachments: + if owner == operation.owner and selector == "odysseus://attachment/" + attachment.record.record_id: + return attachment.file.path + raise ResourceIdentityError("Attachment is not declared by this operation") diff --git a/src/agent_runtime/remote_resources.py b/src/agent_runtime/remote_resources.py new file mode 100644 index 000000000..6074fdb6e --- /dev/null +++ b/src/agent_runtime/remote_resources.py @@ -0,0 +1,238 @@ +"""Backend resolution and pinning, independent of transport and lifecycle. + +Connection/configuration incarnations here are not process identities. Backend +snapshots are captured by trusted admission; discovery never supplies a grant. +""" +from contextlib import contextmanager +from contextvars import ContextVar +from dataclasses import dataclass +from urllib.parse import urlsplit, urlunsplit +from uuid import uuid4 +import hashlib +import hmac +import secrets +import json + +from src.agent_runtime.resources import ExternalResource, NativeBackendResource, ResourceIdentityError + + +def endpoint_identity(url): + """Credential-free origin. Paths may themselves contain access tokens.""" + if not isinstance(url, str) or any(c in url for c in ("\0", "\n", "\r")): + raise ValueError("Malformed resource endpoint") + parsed = urlsplit(url) + if parsed.scheme not in {"http", "https"} or not parsed.hostname: + raise ValueError("Resource endpoint requires an HTTP origin") + host = parsed.hostname.lower() + if ":" in host: + host = "[" + host + "]" + port = parsed.port + if port and port != (443 if parsed.scheme == "https" else 80): + host += f":{port}" + return urlunsplit((parsed.scheme, host, "", "", "")) + + +_CLIENT_ENDPOINTS = {} +_CONFIG_KEY = secrets.token_bytes(32) + + +def configuration_incarnation(value): + """Opaque in-process configuration identity, including secret URL changes.""" + return hmac.new(_CONFIG_KEY, str(value).encode(), hashlib.sha256).hexdigest() + + +def _client_resource(tool, context, *, admission): + from src.tool_execution import _client_bridge, _tui_host_bridge_patch_url, _ROUTED_BRIDGE_TOOLS + bridge = _client_bridge(context) + target = _tui_host_bridge_patch_url(context) if tool == "apply_patch" else None + if target is not None: + url = target[0] + elif bridge is not None and (tool in _ROUTED_BRIDGE_TOOLS or tool == "host_shell"): + url = bridge["url"] + else: + return None + endpoint = endpoint_identity(url) + # Never cache credentials. A request cannot create a registry entry during + # dispatch; only trusted server admission may register an endpoint. + key = configuration_incarnation((url, bridge.get("token") if bridge else None)) + if admission: + _CLIENT_ENDPOINTS.setdefault(key, uuid4().hex) + incarnation = _CLIENT_ENDPOINTS.get(key) + if incarnation is None: + raise ResourceIdentityError("External bridge endpoint is not sealed") + return ExternalResource("client_bridge", endpoint, "tui", tool, incarnation) + + +def http_bridge_resource(tool, context, *, admission=False): + config = context.get("external_execution_bridge") if isinstance(context, dict) else None + if not isinstance(config, dict) or tool not in (config.get("supported_tools") or ()): + return None + url, token = config.get("url"), config.get("token") + if not isinstance(token, str) or not token: + raise ResourceIdentityError("External HTTP bridge has no server configuration") + epoch = configuration_incarnation((url, token, tuple(sorted(config["supported_tools"])))) + if admission: + _CLIENT_ENDPOINTS.setdefault(epoch, epoch) + if epoch not in _CLIENT_ENDPOINTS: + raise ResourceIdentityError("External HTTP bridge configuration is not sealed") + return ExternalResource("execution_bridge", endpoint_identity(url), "request_local_http", tool, epoch) + + +def integration_resource(config): + if not isinstance(config, dict) or not config.get("enabled", True) or not isinstance(config.get("id"), str) or not config["id"]: + raise ResourceIdentityError("Integration identity is unresolved") + endpoint = endpoint_identity(config.get("base_url")) + epoch = configuration_incarnation(json.dumps(config, sort_keys=True, allow_nan=False)) + return ExternalResource("integration", endpoint, config["id"], "api_call", epoch) + + +def api_arguments(content): + if content.lstrip().startswith("{"): + args = json.loads(content) + else: + lines = content.strip().split("\n", 2) + args = {"integration": lines[0].strip()} + if len(lines) > 1: + method, _, path = lines[1].strip().partition(" ") + args.update(method=method, path=path or "/") + if len(lines) > 2: + args["body"] = json.loads(lines[2]) + selector = args.get("integration") + if not isinstance(selector, str) or not selector.strip(): + raise ResourceIdentityError("Integration selector is unresolved") + return args + + +def resolve_backend(tool, *, context=None, admission=False, content="", owner=None): + from src.tool_execution import get_active_execution_bridge, get_mcp_manager, _MCP_TOOL_MAP + from src.tool_security import BUILTIN_EMAIL_TOOLS + bridge = get_active_execution_bridge() + if bridge is not None and tool in bridge.supported_tools: + return bridge.resource_identity(tool) + configured_bridge = http_bridge_resource(tool, context, admission=admission) + if configured_bridge is not None: + return configured_bridge + client = _client_resource(tool, context, admission=admission) + if client is not None: + return client + if tool == "api_call": + from src.integrations import load_integrations + selector = api_arguments(content)["integration"] + rows = [row for row in load_integrations() if row.get("id") == selector + or str(row.get("name", "")).casefold() == selector.casefold()] + if len(rows) != 1: + raise ResourceIdentityError("Integration alias is missing or ambiguous") + return integration_resource(rows[0]) + qualified = tool + required = tool.startswith("mcp__") or tool in BUILTIN_EMAIL_TOOLS + if tool in BUILTIN_EMAIL_TOOLS: + qualified = "mcp__email__" + tool + elif tool in _MCP_TOOL_MAP and tool not in {"read_file", "write_file", "generate_image"}: + server, name = _MCP_TOOL_MAP[tool] + qualified = f"mcp__{server}__{name}" + if qualified.startswith("mcp__"): + manager = get_mcp_manager() + identity = manager.resource_identity(qualified) if manager is not None else None + if isinstance(identity, ExternalResource): + if identity.owner and owner != identity.owner: + raise ResourceIdentityError("MCP backend belongs to another owner") + return identity + if required: + raise ResourceIdentityError("MCP backend/tool identity is unresolved") + if tool == "host_shell": + raise ResourceIdentityError("Host-shell backend identity is unresolved") + return NativeBackendResource(tool) + + +def seal_backends(tools, *, context=None, owner=None): + result = [] + for tool in tools: + try: + if tool == "api_call": + # A generic API operation grant does not select an integration. + # Trusted admission must supply its explicit backend identity, + # or a user can approve one fully sealed exact operation. + continue + result.append(resolve_backend(tool, context=context, admission=True, owner=owner)) + except (ValueError, TypeError, AttributeError): + continue + return tuple(dict.fromkeys(result)) + + +@dataclass(frozen=True) +class BoundBackendOperation: + resource: ExternalResource | NativeBackendResource + request_id: str + owner: str + session_id: str + transport_tool: str + exact_input: str + + def __post_init__(self): + if not isinstance(self.resource, (ExternalResource, NativeBackendResource)): + raise ValueError("Malformed bound backend operation") + if any(not isinstance(v, str) for v in (self.request_id, self.owner, self.session_id, self.transport_tool, self.exact_input)): + raise ValueError("Malformed backend operation binding") + + def to_dict(self): + # Exact arguments/selectors are already digest-bound by the approval's + # original content. Keep credentials out of the identity serializer. + return {"resource": self.resource.to_dict(), "request_id": self.request_id, + "owner": self.owner, "session_id": self.session_id, "tool": self.transport_tool, + "input_digest": configuration_incarnation(self.exact_input)} + + def validate(self, context=None): + current = resolve_backend(self.transport_tool, context=context, content=self.exact_input, owner=self.owner) + if current != self.resource: + # A pinned native backend remains native when MCP availability + # changes. It cannot be upgraded to an external backend. + if isinstance(self.resource, NativeBackendResource) and isinstance(current, ExternalResource) and current.namespace == "mcp": + return + raise ResourceIdentityError("Backend resource identity changed") + + +def bind_backend_for_operation(authority, operation, *, context=None, approved=None, exact_admission=False): + current = resolve_backend(operation.transport_tool, context=context, content=operation.input, owner=authority.owner) + native = NativeBackendResource(operation.transport_tool) + if approved is not None: + if (not isinstance(approved, BoundBackendOperation) + or (approved.request_id and approved.request_id != authority.request_id) + or (approved.owner, approved.session_id) != (authority.owner, authority.session_id) + or (approved.transport_tool, approved.exact_input) != (operation.transport_tool, operation.input)): + raise ResourceIdentityError("Approved backend binding changed") + selected = approved.resource + elif current in authority.backend_resources: + selected = current + elif native in authority.backend_resources: + selected = native + elif isinstance(current, NativeBackendResource) and not authority.inherited: + # Legacy operation authority can only retain the fixed local backend; + # it cannot reconstruct any external backend from current availability. + selected = current + else: + raise ResourceIdentityError("External backend is outside sealed request scope") + if isinstance(selected, ExternalResource) and selected not in authority.backend_resources: + if not (exact_admission and approved is not None and not authority.inherited): + raise ResourceIdentityError("External backend exceeds parent/request scope") + bound = BoundBackendOperation(selected, authority.request_id, authority.owner, authority.session_id, + operation.transport_tool, operation.input) + bound.validate(context) + return bound + + +_ACTIVE = ContextVar("backend_resource_operation", default=None) + + +def active_backend_operation(): + return _ACTIVE.get() + + +@contextmanager +def bind_backend_operation(operation): + if operation is not None and not isinstance(operation, BoundBackendOperation): + raise TypeError("Backend operation must be server-owned") + token = _ACTIVE.set(operation) + try: + yield operation + finally: + _ACTIVE.reset(token) diff --git a/src/agent_runtime/resources.py b/src/agent_runtime/resources.py index c8b76d474..d63a9ed00 100644 --- a/src/agent_runtime/resources.py +++ b/src/agent_runtime/resources.py @@ -10,6 +10,7 @@ from enum import Enum import os from pathlib import Path import stat +import sys from src.agent_runtime.path_policy import _is_sensitive_path from src.path_confinement import canonical_root, confine @@ -30,11 +31,65 @@ def _absolute(value): def _control_plane_path(path): # Execution snapshots/receipts are server state, even if a workspace root # contains the data directory. A writable user file cannot mint authority. - from src.constants import BG_JOBS_DIR, BG_JOBS_FILE, CONTAINMENT_STATE_FILE - if path in {canonical_root(BG_JOBS_FILE), canonical_root(CONTAINMENT_STATE_FILE)}: + from src import constants + protected = {canonical_root(getattr(constants, name)) for name in ( + "BG_JOBS_FILE", "CONTAINMENT_STATE_FILE", "APP_DB", "AUTH_FILE", + "SETTINGS_FILE", "SESSIONS_FILE", "USER_PREFS_FILE", "VAULT_FILE", + "SCHEDULED_EMAILS_DB", "EMAIL_CACHE_DB", "MEMORY_FILE", "INTEGRATIONS_FILE", + )} + job_dirs = {canonical_root(constants.BG_JOBS_DIR)} + # Producers may have configured paths different from the default constants. + # Inspect already-loaded server metadata without initializing a store here. + bg = sys.modules.get("src.bg_jobs") + if bg is not None: + for name, targets in (("_STORE", protected), ("_JOBS_DIR", job_dirs)): + value = getattr(bg, name, None) + if isinstance(value, (str, os.PathLike)): + targets.add(canonical_root(value)) + containment = sys.modules.get("src.containment") + if containment is not None: + value = containment._store_path() + if isinstance(value, (str, os.PathLike)): + protected.add(canonical_root(value)) + database = sys.modules.get("core.database") + url = getattr(getattr(database, "engine", None), "url", None) + if url is not None and url.get_backend_name() == "sqlite": + location = url.database + if isinstance(location, str) and location not in {"", ":memory:"}: + from urllib.parse import unquote + if location.startswith("file:"): + location = unquote(location[5:].split("?", 1)[0]) + protected.update(canonical_root(location + suffix) for suffix in ("", "-wal", "-shm", "-journal")) + from src.tool_utils import get_upload_handler + uploader = get_upload_handler() + if uploader is not None and isinstance(getattr(uploader, "upload_dir", None), (str, os.PathLike)): + protected.add(canonical_root(Path(uploader.upload_dir) / "uploads.json")) + for directory in job_dirs: + jobs = Path(directory) + if Path(path).is_relative_to(jobs): + return True + if jobs.exists(): + # Uninspectable state fails closed; hardlinks retain object identity. + protected.update(canonical_root(p) for p in jobs.iterdir()) + protected.update(canonical_root(getattr(constants, name) + suffix) + for name in ("APP_DB", "SCHEDULED_EMAILS_DB", "EMAIL_CACHE_DB") + for suffix in ("-wal", "-shm", "-journal")) + protected.add(canonical_root(Path(constants.DATA_DIR) / ".app_key")) + protected.add(canonical_root(Path(constants.UPLOAD_DIR) / "uploads.json")) + if path in protected: return True - return (Path(path).is_relative_to(canonical_root(BG_JOBS_DIR)) - and path.endswith(".authority.json")) + try: + candidate = os.stat(path) + except FileNotFoundError: + return False + for control in protected: + try: + observed = os.stat(control) + except FileNotFoundError: + continue + if (candidate.st_dev, candidate.st_ino) == (observed.st_dev, observed.st_ino): + return True + return False class FilesystemScope(str, Enum): @@ -201,18 +256,16 @@ def intersect_roots(parent, child): for right in child: if (left.scope, left.owner) != (right.scope, right.owner): continue - if left == right: - result.append(left) - continue try: + left.validate() + right.validate() + if left == right: + result.append(left) + continue if Path(right.path).is_relative_to(left.path): # A newly sealed child may not renew a replaced parent root. - left.validate() - right.validate() result.append(right) elif Path(left.path).is_relative_to(right.path): - left.validate() - right.validate() result.append(left) except (OSError, ValueError, RuntimeError): continue @@ -279,12 +332,44 @@ class ExternalResource: tool_id: str incarnation: str external: bool = True + contained: bool = False + owner: str = "" def __post_init__(self): for name in ("namespace", "endpoint_id", "server_id", "tool_id", "incarnation"): _text(getattr(self, name), name) - if self.external is not True: + if self.external is not True or self.contained is not False: raise ValueError("External resource cannot attest local containment") + _text(self.owner, "external owner", optional=True) + + def to_dict(self): + return {"kind": "external", **asdict(self)} + + +@dataclass(frozen=True) +class NativeBackendResource: + tool_id: str + namespace: str = "native" + external: bool = False + contained: bool = False + + def __post_init__(self): + _text(self.tool_id, "native tool") + if self.namespace != "native" or self.external is not False or self.contained is not False: + raise ValueError("Malformed native backend identity") + + def to_dict(self): + return {"kind": "native", **asdict(self)} + + +def backend_from_dict(value): + if not isinstance(value, dict): + raise ValueError("Malformed backend snapshot") + fields = dict(value) + kind = fields.pop("kind", None) + if kind not in {"native", "external"}: + raise ValueError("Malformed backend kind") + return (NativeBackendResource if kind == "native" else ExternalResource)(**fields) @dataclass(frozen=True) @@ -295,8 +380,77 @@ class OwnedResource: collection: str record_id: str revision: str = "" + record_thread_id: str = "" def __post_init__(self): for name in ("namespace", "owner", "thread_id", "collection", "record_id"): _text(getattr(self, name), name) _text(self.revision, "revision", optional=True) + _text(self.record_thread_id, "record thread", optional=True) + + def to_dict(self): + return asdict(self) + + +@dataclass(frozen=True) +class OwnedScope: + namespace: str + owner: str + thread_id: str + record_ids: frozenset[str] | None = None + + def __post_init__(self): + for name in ("namespace", "owner", "thread_id"): + _text(getattr(self, name), name) + if self.record_ids is not None: + if not isinstance(self.record_ids, frozenset): + raise ValueError("Owned scope must be immutable") + for identifier in self.record_ids: + _text(identifier, "record identifier") + if identifier == "*": + raise ValueError("Collection authority must be explicit") + + def permits(self, resource): + return (isinstance(resource, OwnedResource) + and (self.namespace, self.owner, self.thread_id) == + (resource.namespace, resource.owner, resource.thread_id) + and resource.collection == self.namespace + and (self.record_ids is None or resource.record_id in self.record_ids)) + + def intersect(self, other): + if (self.namespace, self.owner, self.thread_id) != (other.namespace, other.owner, other.thread_id): + return None + ids = (other.record_ids if self.record_ids is None else self.record_ids if other.record_ids is None + else self.record_ids & other.record_ids) + return OwnedScope(self.namespace, self.owner, self.thread_id, ids) + + def to_dict(self): + return {"namespace": self.namespace, "owner": self.owner, "thread_id": self.thread_id, + "record_ids": None if self.record_ids is None else sorted(self.record_ids)} + + @classmethod + def from_dict(cls, value): + if not isinstance(value, dict) or set(value) != {"namespace", "owner", "thread_id", "record_ids"}: + raise ValueError("Malformed owned scope snapshot") + ids = value["record_ids"] + if ids is not None and (not isinstance(ids, list) or any(not isinstance(v, str) for v in ids)): + raise ValueError("Malformed owned record limits") + return cls(value["namespace"], value["owner"], value["thread_id"], + None if ids is None else frozenset(ids)) + + +OWNED_TOOL_NAMESPACES = { + **{name: "documents" for name in ("create_document", "edit_document", "update_document", "suggest_document", "manage_documents")}, + **{name: "threads" for name in ("create_session", "list_sessions", "manage_session", "send_to_session", "search_chats")}, + **{name: "attachments" for name in ("extract_text", "inspect_media", "transcribe_media")}, + "manage_notes": "notes", + "manage_memory": "memory", + **{name: "vault" for name in ("vault_get", "vault_search", "vault_unlock")}, +} + + +def seal_owned_scopes(owner, thread_id, tools): + if not owner or not thread_id: + return () + return tuple(OwnedScope(namespace, owner, thread_id) + for namespace in sorted({OWNED_TOOL_NAMESPACES[t] for t in tools if t in OWNED_TOOL_NAMESPACES})) diff --git a/src/agent_tools/filesystem_tools.py b/src/agent_tools/filesystem_tools.py index 89fd30975..89d161d43 100644 --- a/src/agent_tools/filesystem_tools.py +++ b/src/agent_tools/filesystem_tools.py @@ -26,6 +26,18 @@ _BINARY_ARTIFACT_SUFFIXES = _STRUCTURED_DOCUMENT_SUFFIXES | frozenset({ ".png", ".wav", ".webm", ".webp", ".zip", }) + +def _visible_bound_resource(path): + from src.agent_runtime.resource_binding import active_resource_operation + bound = active_resource_operation() + if bound is None: + return True + try: + bound.resolve_path(path) + return True + except (ValueError, OSError, RuntimeError): + return False + # Models frequently put source artifacts in a Markdown code fence even when a # tool schema asks for the raw file body. Persisting that fence makes HTML, # CSS, JavaScript, and source files invalid. Restrict normalization to @@ -610,6 +622,8 @@ class LsTool: for entry in it: if entry.name.startswith("."): continue + if not _visible_bound_resource(entry.path): + continue try: is_dir = entry.is_dir(follow_symlinks=False) size = entry.stat(follow_symlinks=False).st_size if not is_dir else 0 @@ -681,7 +695,7 @@ class GlobTool: # .ssh/id_rsa, …) falls through to the walk, which skips it — # otherwise glob would surface secret paths that read_file / # grep already refuse to touch. - if inside and os.path.exists(cand) and not _is_sensitive_path(cand): + if inside and os.path.exists(cand) and not _is_sensitive_path(cand) and _visible_bound_resource(cand): return [cand], None # Literal not at exact path — fall through to walk so # e.g. "foo.py" still matches at any depth (like rglob). @@ -705,7 +719,7 @@ class GlobTool: if regex.fullmatch(rel) or regex.fullmatch(name): # Skip deny-listed sensitive files (.env, id_rsa, # known_hosts, …) the same way grep does. - if _is_sensitive_path(os.path.realpath(full)): + if _is_sensitive_path(os.path.realpath(full)) or not _visible_bound_resource(full): continue try: mtime = os.stat(full).st_mtime @@ -766,9 +780,12 @@ class GrepTool: def _grep(): import re as _re import shutil + from src.agent_runtime.resource_binding import active_resource_operation if not os.path.exists(root): return None, f"grep: search target not found: {_display_tool_path(root)}" - rg = shutil.which("rg") + # The pathname-only fast path scans before individual resources can + # be checked. Bound searches must validate every file before read. + rg = None if active_resource_operation() is not None else shutil.which("rg") if rg: cmd = [rg, "--line-number", "--with-filename", "--no-heading", "--color=never", "--max-count", str(max_hits)] diff --git a/src/agent_tools/media_tools.py b/src/agent_tools/media_tools.py index 493b161df..e5915e4f5 100644 --- a/src/agent_tools/media_tools.py +++ b/src/agent_tools/media_tools.py @@ -184,6 +184,9 @@ def _resolve_workspace_path( raise ValueError( f"{tool_name} {field_name} must stay inside the active workspace" ) from exc + from src.agent_runtime.resources import _control_plane_path + if _control_plane_path(str(resolved)): + raise ValueError(f"{tool_name} {field_name} addresses execution-control state") if must_exist and not resolved.is_file(): raise FileNotFoundError(f"media file not found: {raw}") return resolved @@ -445,8 +448,10 @@ class ExtractTextTool: from src.tool_utils import get_upload_handler ref = re.fullmatch(r'odysseus://attachment/([A-Za-z0-9_-]+(?:\.[A-Za-z0-9]+)?)', raw_path) owner = (_ctx or {}).get('owner') + from src.agent_runtime.owned_resources import bound_attachment_path + bound_path = bound_attachment_path(owner, raw_path) handler = get_upload_handler() - info = handler.resolve_upload(ref[1], owner=owner, allow_admin=False) if ref and owner and handler else None + info = {"path": bound_path} if bound_path else handler.resolve_upload(ref[1], owner=owner, allow_admin=False) if ref and owner and handler else None if not info or not info.get('path'): raise ValueError('Uploaded image not found or not accessible to this user') path = Path(info['path']) diff --git a/src/ai_interaction.py b/src/ai_interaction.py index ed767e678..416b0c2f4 100644 --- a/src/ai_interaction.py +++ b/src/ai_interaction.py @@ -394,6 +394,11 @@ async def do_manage_memory(content: str, session_id: Optional[str] = None, owner if not _memory_manager: return {"error": "Memory manager not available"} + from src.agent_runtime.owned_resources import active_owned_operation + bound = active_owned_operation() + if bound is not None: + bound.validate() + lines = _manage_memory_lines(content) if not lines: return {"error": "Need at least 1 line: action"} @@ -465,7 +470,7 @@ async def do_manage_memory(content: str, session_id: Optional[str] = None, owner memories = _memory_manager.load_all() found = False for m in memories: - if m.get("id", "").startswith(memory_id): + if (m.get("id", "") == memory_id if bound is not None else m.get("id", "").startswith(memory_id)): # Verify ownership if owner and m.get("owner") != owner: return {"error": f"Memory '{memory_id}' not found"} @@ -498,7 +503,7 @@ async def do_manage_memory(content: str, session_id: Optional[str] = None, owner full_id = None delete_id = None for m in memories: - if m.get("id", "").startswith(memory_id): + if (m.get("id", "") == memory_id if bound is not None else m.get("id", "").startswith(memory_id)): # Verify ownership if owner and m.get("owner") != owner: return {"error": f"Memory '{memory_id}' not found"} diff --git a/src/integrations.py b/src/integrations.py index 82806a24a..763b930c1 100644 --- a/src/integrations.py +++ b/src/integrations.py @@ -517,6 +517,12 @@ async def execute_api_call( if not integration: return {"error": f"Integration not found: {integration_id}", "exit_code": 1} + from src.agent_runtime.remote_resources import active_backend_operation, integration_resource + bound = active_backend_operation() + if bound is not None and bound.resource != integration_resource(integration): + return {"error": "Integration resource identity changed", "exit_code": 1, + "failure_kind": "resource_identity_denied"} + if not integration.get("enabled", True): return {"error": f"Integration '{integration.get('name')}' is disabled", "exit_code": 1} diff --git a/src/mcp_manager.py b/src/mcp_manager.py index 21d9e9cad..53112009d 100644 --- a/src/mcp_manager.py +++ b/src/mcp_manager.py @@ -164,6 +164,10 @@ class McpManager: self._owner_tasks: Dict[str, asyncio.Task] = {} # Tracking updates to tools/connections for RAG indexing / prompt cache self._generation = 0 + # Identity of the actual connection, not a PID or lifecycle contract. + self._resource_connections = {} + self._resource_endpoints = {} + self._resource_owners = {} async def connect_server( self, @@ -177,6 +181,13 @@ class McpManager: ) -> bool: """Connect to an MCP server via stdio, SSE, or Streamable HTTP transport.""" try: + from src.agent_runtime.remote_resources import endpoint_identity, configuration_incarnation + self._resource_endpoints[server_id] = ( + endpoint_identity(url) if transport in {"sse", "http"} else f"stdio:{server_id}", + configuration_incarnation((transport, url, command, args, env))) + if server_id == "memory": + effective_env = {**os.environ, **(env or {})} + self._resource_owners[server_id] = str(effective_env.get("ODYSSEUS_MCP_MEMORY_OWNER") or effective_env.get("ODYSSEUS_MEMORY_OWNER") or "").strip() if transport == "stdio": res = await self._connect_stdio(server_id, name, command, args or [], env or {}) elif transport == "sse": @@ -243,6 +254,7 @@ class McpManager: identity = ", ".join(identity_hints) if identity_hints else "" self._sessions[server_id] = session + self._register_resource_connection(server_id, session) self._stacks[server_id] = stack self._tools[server_id] = tools self._connections[server_id] = { @@ -302,6 +314,7 @@ class McpManager: }) self._sessions[server_id] = session + self._register_resource_connection(server_id, session) self._stacks[server_id] = stack self._tools[server_id] = tools self._connections[server_id] = { @@ -385,6 +398,7 @@ class McpManager: }) self._sessions[server_id] = session + self._register_resource_connection(server_id, session) self._stacks[server_id] = stack self._tools[server_id] = tools self._connections[server_id] = { @@ -445,6 +459,7 @@ class McpManager: logger.warning(f"Error closing MCP server {server_id}: {e}") self._sessions.pop(server_id, None) + self._resource_connections.pop(server_id, None) self._tools.pop(server_id, None) self._connections.pop(server_id, None) self._generation += 1 @@ -516,6 +531,33 @@ class McpManager: "name": srv.name, } + def _register_resource_connection(self, server_id, session): + from uuid import uuid4 + endpoint = self._resource_endpoints.get(server_id) + if endpoint: + self._resource_connections[server_id] = (endpoint, uuid4().hex, session, + self._resource_owners.get(server_id, "")) + + def resource_identity(self, qualified_name): + from src.agent_runtime.resources import ExternalResource + parts = qualified_name.split("__", 2) + if len(parts) != 3 or parts[0] != "mcp" or not parts[1] or not parts[2]: + return None + _, server, tool = parts + # The builtin memory producer uses a fixed owner, not model arguments. + # The builtin RAG producer has no owner contract; its legacy global + # store cannot acquire private read scope through discovery. + if server == "rag" or (server == "memory" and not self._resource_owners.get(server)): + return None + record = self._resource_connections.get(server) + if (not record or self._sessions.get(server) is not record[2] + or self._resource_endpoints.get(server) != record[0] + or self._resource_owners.get(server, "") != record[3] + or not any(row.get("name") == tool for row in self._tools.get(server, []))): + return None + return ExternalResource("mcp", record[0][0], server, qualified_name, record[1], + owner=record[3]) + async def call_tool(self, qualified_name: str, arguments: Dict) -> Dict: """Call an MCP tool by its qualified name (mcp__{server_id}__{tool_name}). @@ -532,6 +574,12 @@ class McpManager: if not session: return {"error": f"MCP server not connected: {server_id}", "exit_code": 1} + from src.agent_runtime.remote_resources import active_backend_operation + bound_backend = active_backend_operation() + if bound_backend is not None and self.resource_identity(qualified_name) != bound_backend.resource: + return {"error": "MCP resource binding changed", "exit_code": 1, + "failure_kind": "resource_identity_denied"} + try: if server_id == BROWSER_MCP_SERVER_ID: # The shared Playwright browser must not hold a turn forever. @@ -554,7 +602,7 @@ class McpManager: result = await self._do_call(session, tool_name, arguments) except Exception as e: # Auto-reconnect for builtin servers whose subprocess may have died - if self.is_builtin(server_id): + if bound_backend is None and self.is_builtin(server_id): logger.warning(f"MCP call failed for {qualified_name}, attempting reconnect: {e}") reconnected = await self._reconnect_builtin(server_id) if reconnected: diff --git a/src/tool_approvals.py b/src/tool_approvals.py index d392bfce4..416fa9e90 100644 --- a/src/tool_approvals.py +++ b/src/tool_approvals.py @@ -29,6 +29,8 @@ from src.agent_runtime.authority import RequestAuthority if TYPE_CHECKING: from src.agent_runtime.resource_binding import BoundFilesystemOperation + from src.agent_runtime.remote_resources import BoundBackendOperation + from src.agent_runtime.owned_resources import BoundOwnedOperation DEFAULT_APPROVAL_TTL_SECONDS = 10 * 60 @@ -123,6 +125,8 @@ def _binding_payload( result_integrity: str, request_authority: RequestAuthority | None = None, resource_operation=None, + backend_operation=None, + owned_operation=None, ) -> dict[str, Any]: return { "owner": _normalized_owner(owner), @@ -145,6 +149,8 @@ def _binding_payload( "result_integrity": str(result_integrity), "request_authority": request_authority.to_dict() if request_authority is not None else None, "resource_operation": resource_operation.to_dict() if resource_operation is not None else None, + "backend_operation": backend_operation.to_dict() if backend_operation is not None else None, + "owned_operation": owned_operation.to_dict() if owned_operation is not None else None, } @@ -176,6 +182,8 @@ class PendingToolApproval: request_authority: RequestAuthority | None = None # Server-resolved targets at proposal time; never read from the approval UI. resource_operation: BoundFilesystemOperation | None = None + backend_operation: BoundBackendOperation | None = None + owned_operation: BoundOwnedOperation | None = None def public_payload(self, *, reason: str | None = None) -> dict[str, Any]: return { @@ -286,6 +294,8 @@ class ExactToolApproval: result_integrity=result_integrity, request_authority=self.pending.request_authority, resource_operation=self.pending.resource_operation, + backend_operation=self.pending.backend_operation, + owned_operation=self.pending.owned_operation, ) return _canonical_digest(expected) == self.pending.digest @@ -370,6 +380,7 @@ class ToolApprovalStore: capabilities: ToolCapabilities, request_text: Any = "", request_authority: RequestAuthority | None = None, + client_runtime_context: dict | None = None, ) -> PendingToolApproval: if request_authority is not None and not isinstance(request_authority, RequestAuthority): raise TypeError("Approval authority must be server-owned RequestAuthority") @@ -378,10 +389,38 @@ class ToolApprovalStore: from src.agent_runtime.resource_binding import NATIVE_FILESYSTEM_TOOLS, resolve_filesystem_operation from src.agent_runtime.resources import FilesystemRoot resource_operation = None - if tool_name in NATIVE_FILESYSTEM_TOOLS: + backend_operation = None + owned_operation = None + from src.agent_runtime.remote_resources import BoundBackendOperation, resolve_backend + from src.agent_runtime.owned_resources import needs_owned_binding, resolve_owned_operation + from src.agent_runtime.resources import NativeBackendResource + try: + operation = ExactOperation.normalize(tool_name, content) + backend = resolve_backend(operation.transport_tool, context=client_runtime_context, content=operation.input, owner=_normalized_owner(owner)) + if request_authority is not None and request_authority.inherited: + if not request_authority.permits(operation) or backend not in request_authority.backend_resources: + raise ValueError("Child approval exceeds originating authority") + backend_operation = BoundBackendOperation(backend, + request_authority.request_id if request_authority is not None else "", + _normalized_owner(owner), str(session_id or ""), operation.transport_tool, operation.input) + if isinstance(backend, NativeBackendResource) and needs_owned_binding(operation): + resolved_owned = resolve_owned_operation(operation, owner=_normalized_owner(owner), + thread_id=str(session_id or ""), request_id=backend_operation.request_id, + document_id=document_id) + if request_authority is not None and request_authority.inherited: + if not all(any(scope.permits(r) for scope in request_authority.owned_scopes) for r in resolved_owned.resources): + raise ValueError("Child approval exceeds originating record scope") + owned_operation = resolved_owned + if owned_operation.document_id: + document_id = owned_operation.document_id + document_version = owned_operation.document_version + document_digest = owned_operation.document_digest + except (ValueError, TypeError, OSError, RuntimeError, AttributeError): + pass # Unresolved proposals are never reconstructed at execution. + if tool_name in NATIVE_FILESYSTEM_TOOLS and backend_operation is not None and isinstance(backend_operation.resource, NativeBackendResource): try: roots = request_authority.resource_roots if request_authority is not None else () - if not roots and workspace: + if not roots and workspace and (request_authority is None or not request_authority.inherited): roots = (FilesystemRoot.seal(workspace, owner=_normalized_owner(owner)),) resource_operation = resolve_filesystem_operation( ExactOperation.normalize(tool_name, content), roots=roots, workspace=workspace or "", @@ -409,6 +448,8 @@ class ToolApprovalStore: result_integrity=result_integrity, request_authority=request_authority, resource_operation=resource_operation, + backend_operation=backend_operation, + owned_operation=owned_operation, ) pending = PendingToolApproval( approval_id=secrets.token_urlsafe(32), @@ -434,6 +475,8 @@ class ToolApprovalStore: request_text=str(request_text or ""), request_authority=request_authority, resource_operation=resource_operation, + backend_operation=backend_operation, + owned_operation=owned_operation, ) with self._lock: self._purge_expired_locked(now) diff --git a/src/tool_execution.py b/src/tool_execution.py index 7e77c4408..9b7a41a8a 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -19,7 +19,7 @@ import secrets import sys import time from contextlib import contextmanager -from dataclasses import dataclass, replace +from dataclasses import dataclass, field, replace from typing import Any, Awaitable, Callable, Dict, Iterator, Optional, Tuple @@ -48,7 +48,13 @@ from src.agent_runtime.resource_binding import ( NATIVE_FILESYSTEM_TOOLS, active_resource_operation, bind_resource_operation, resolve_filesystem_operation, ) -from src.agent_runtime.resources import ResourceIdentityError +from src.agent_runtime.resources import ExternalResource, NativeBackendResource, ResourceIdentityError +from src.agent_runtime.remote_resources import ( + active_backend_operation, bind_backend_operation, bind_backend_for_operation, +) +from src.agent_runtime.owned_resources import ( + active_owned_operation, bind_owned_operation, admit_owned_operation, needs_owned_binding, +) class _MissingToolSecurityContext: @@ -88,6 +94,9 @@ class AgentExecutionBridge: route_tool: ExecutionBridgeHandler supported_tools: frozenset[str] name: str = "external_environment" + endpoint_id: str = "" + incarnation: str = field(default_factory=lambda: secrets.token_hex(16), init=False) + configuration_id: str = "" def __post_init__(self) -> None: if not callable(self.route_tool): @@ -95,6 +104,12 @@ class AgentExecutionBridge: if not self.supported_tools: raise ValueError("execution bridge supported_tools cannot be empty") + def resource_identity(self, tool): + from src.agent_runtime.resources import ExternalResource + from src.agent_runtime.remote_resources import endpoint_identity + endpoint = endpoint_identity(self.endpoint_id) if self.endpoint_id else "bridge:" + self.name + return ExternalResource("execution_bridge", endpoint, self.name, tool, self.configuration_id or self.incarnation) + _active_execution_bridge: contextvars.ContextVar[AgentExecutionBridge | None] = ( contextvars.ContextVar("agent_execution_bridge", default=None) @@ -303,7 +318,7 @@ def _client_bridge(client_runtime_context: Optional[Dict]) -> Optional[Dict]: context = client_runtime_context if isinstance(client_runtime_context, dict) else {} if str(context.get("surface") or "").strip() != "odysseus-tui": return None - bridge = context.get("host_shell_bridge") + bridge = context.get("host_shell_bridge") or context.get("hostShellBridge") if not isinstance(bridge, dict): return None url = str(bridge.get("url") or "").strip() @@ -1153,8 +1168,13 @@ async def _call_mcp_tool( progress_cb: Optional[Callable[[Dict], Awaitable[None]]] = None, ) -> Dict: """Route a legacy tool call through the MCP manager, with direct fallbacks.""" + bound = active_backend_operation() + if bound is not None and isinstance(bound.resource, NativeBackendResource): + return await _direct_fallback(tool, content, progress_cb=progress_cb) or {"error": f"Native tool '{tool}' unavailable", "exit_code": 1} mcp = get_mcp_manager() if not mcp: + if bound is not None: + raise ResourceIdentityError("Pinned MCP backend is unavailable") return await _direct_fallback(tool, content, progress_cb=progress_cb) or {"error": f"MCP manager not available for tool '{tool}'", "exit_code": 1} server_id, tool_name = _MCP_TOOL_MAP[tool] @@ -1168,7 +1188,7 @@ async def _call_mcp_tool( result = _normalize_mcp_text_error(result) # If MCP server not connected, try direct fallback - if isinstance(result, dict) and result.get("exit_code") == 1 and "not connected" in result.get("error", ""): + if bound is None and isinstance(result, dict) and result.get("exit_code") == 1 and "not connected" in result.get("error", ""): fallback = await _direct_fallback(tool, content, progress_cb=progress_cb) if fallback: return fallback @@ -1240,6 +1260,9 @@ async def _direct_fallback( _subproc_env = _agent_subprocess_env() try: + owned = active_owned_operation() + if owned is not None: + owned.validate() ctx = { "progress_cb": progress_cb, "subproc_env": _subproc_env, @@ -1250,6 +1273,7 @@ async def _direct_fallback( "tool_policy": tool_policy, "request_authority": active_request_authority(), "resource_operation": active_resource_operation(), + "owned_operation": active_owned_operation(), } from src.agent_tools import TOOL_HANDLERS @@ -1273,6 +1297,9 @@ async def _document_tool_dispatch( ) -> Optional[Dict]: """Route a document tool through TOOL_HANDLERS with the right ctx shape.""" from src.agent_tools import TOOL_HANDLERS + owned = active_owned_operation() + if owned is not None: + owned.validate() ctx = { "session_id": session_id, "owner": owner, @@ -1377,15 +1404,30 @@ async def execute_tool_block( "exit_code": 1, "failure_kind": "turn_contract_denied", } - # External executors require their own adapters. Local observations must - # never stand in for remote resource or containment identities. - execution_bridge = get_active_execution_bridge() transport = operation.transport_tool - external_resource_call = ( - (execution_bridge is not None and transport in execution_bridge.supported_tools) - or (transport in _ROUTED_BRIDGE_TOOLS and _client_bridge(client_runtime_context) is not None) - or (transport == "apply_patch" and _tui_host_bridge_patch_url(client_runtime_context)) - ) + try: + pending = exact_approval.pending if exact_approval is not None else None + if pending is not None and pending.backend_operation is None: + raise ResourceIdentityError("Approved action has no sealed backend identity") + backend_operation = bind_backend_for_operation( + authority, operation, context=client_runtime_context, + approved=pending.backend_operation if pending is not None else None, + exact_admission=exact_admission) + external_resource_call = isinstance(backend_operation.resource, ExternalResource) + owned_operation = None + if needs_owned_binding(operation) and not external_resource_call: + if pending is not None and pending.owned_operation is None: + raise ResourceIdentityError("Approved action has no sealed owned resource identity") + owned_operation = admit_owned_operation( + authority, operation, document_id=active_document_id, + approved=pending.owned_operation if pending is not None else None, + exact_admission=exact_admission) + except (ValueError, TypeError, OSError, RuntimeError, AttributeError) as error: + return f"{transport}: BLOCKED", { + "error": str(error), "exit_code": 1, "blocked": True, + "failure_kind": "resource_identity_denied", + **({"policy": "exact_tool_approval"} if exact_approval is not None else {}), + } resource_operation = None if operation.tool in NATIVE_FILESYSTEM_TOOLS and not external_resource_call: try: @@ -1498,29 +1540,21 @@ async def execute_tool_block( token = _active_workspace.set(workspace or None) try: - with bind_request_authority(authority), bind_resource_operation(resource_operation): + backend_operation.validate(client_runtime_context) + normalized = resource_operation or owned_operation + sealed_document = owned_operation or (exact_approval.pending if approval_claimed else None) + with (bind_request_authority(authority), bind_resource_operation(resource_operation), + bind_backend_operation(backend_operation), bind_owned_operation(owned_operation)): output = await _execute_tool_block_impl( - ToolBlock(transport, resource_operation.execution_input) if resource_operation is not None else block, + ToolBlock(transport, normalized.execution_input) if normalized is not None else block, session_id=session_id, disabled_tools=disabled_tools, owner=owner, progress_cb=progress_cb, tool_policy=tool_policy, - approved_document_id=( - exact_approval.pending.document_id - if approval_claimed - else None - ), - approved_document_version=( - exact_approval.pending.document_version - if approval_claimed - else None - ), - approved_document_digest=( - exact_approval.pending.document_digest - if approval_claimed - else None - ), + approved_document_id=sealed_document.document_id if sealed_document is not None else None, + approved_document_version=sealed_document.document_version if sealed_document is not None else None, + approved_document_digest=sealed_document.document_digest if sealed_document is not None else None, active_document_id=active_document_id, client_runtime_context=client_runtime_context, ) @@ -1664,11 +1698,20 @@ async def _execute_tool_block_impl( return desc, result execution_bridge = get_active_execution_bridge() + backend = active_backend_operation() + if backend is not None: + backend.validate(client_runtime_context) + owned = active_owned_operation() + if owned is not None: + owned.validate() bridge_owns_tool = ( execution_bridge is not None and tool in execution_bridge.supported_tools and active_resource_operation() is None + and (backend is None or backend.resource.namespace == "execution_bridge") ) + if backend is not None and backend.resource.namespace == "execution_bridge" and not bridge_owns_tool: + raise ResourceIdentityError("Pinned external execution bridge is unavailable") # Public-owner restrictions protect tools executed by this deployment. # A request-scoped execution bridge is a separate, explicit authority for @@ -1727,7 +1770,7 @@ async def _execute_tool_block_impl( }, ) - if (active_resource_operation() is None and tool in _ROUTED_BRIDGE_TOOLS + if (active_resource_operation() is None and (backend is None or backend.resource.namespace == "client_bridge") and tool in _ROUTED_BRIDGE_TOOLS and _client_bridge(client_runtime_context) is not None): return await dispatched(_route_tool_via_bridge(tool, content, session_id, client_runtime_context)) @@ -1735,7 +1778,7 @@ async def _execute_tool_block_impl( # marker runs DETACHED — returns a job id immediately so the chat stream # isn't held open for a multi-minute install/ffmpeg/download. The always-on # monitor re-invokes the agent with the full output when the job finishes. - if tool == "bash" and session_id: + if tool == "bash" and session_id and (backend is None or isinstance(backend.resource, NativeBackendResource)): _is_bg, _bg_cmd = _split_bg_marker(content) if _is_bg and _bg_cmd: from src import bg_jobs diff --git a/src/tools/notes.py b/src/tools/notes.py index fb6a812d3..f464e0dca 100644 --- a/src/tools/notes.py +++ b/src/tools/notes.py @@ -77,6 +77,10 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict: if action == "list" and list_search_query: action = "search" args.setdefault("query", list_search_query) + from src.agent_runtime.owned_resources import active_owned_operation + bound = active_owned_operation() + if bound is not None: + bound.validate() db = SessionLocal() def _norm_note_title(value: str) -> str: @@ -98,7 +102,7 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict: def _note_by_prefix(note_id: str): if not note_id: return None - q = db.query(Note).filter(Note.id.startswith(note_id)) + q = db.query(Note).filter(Note.id == note_id if bound is not None else Note.id.startswith(note_id)) if owner: q = q.filter(Note.owner == owner) return q.first() @@ -415,7 +419,7 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict: elif action == "update": note_id = _note_id_arg() note = _note_by_prefix(note_id) - if not note: + if not note and bound is None: title_query = str( args.get("title") or args.get("query") @@ -489,7 +493,7 @@ async def do_manage_notes(content: str, owner: Optional[str] = None) -> Dict: elif action == "delete": note_id = _note_id_arg() note = _note_by_prefix(note_id) - if not note: + if not note and bound is None: title_query = str( args.get("title") or args.get("query") diff --git a/src/tools/system.py b/src/tools/system.py index d60d8f585..a5e9d340b 100644 --- a/src/tools/system.py +++ b/src/tools/system.py @@ -689,9 +689,16 @@ async def do_api_call(content: str) -> Dict: pass integration_name = args.get("integration", "") + from src.agent_runtime.remote_resources import active_backend_operation + bound = active_backend_operation() + if bound is not None: + bound.validate() + if bound.resource.namespace != "integration": + return {"error": "API call has no integration resource binding", "exit_code": 1} + integration_name = bound.resource.server_id integrations = load_integrations() intg = next((i for i in integrations if i["id"] == integration_name - or i["name"].lower() == integration_name.lower()), None) + or (bound is None and i["name"].lower() == integration_name.lower())), None) if not intg: available = ", ".join(i["name"] for i in integrations if i.get("enabled", True)) return {"error": f"No integration matching '{integration_name}'. Available: {available or 'none configured'}", "exit_code": 1} diff --git a/src/tools/vault.py b/src/tools/vault.py index fbb3bfcf9..2f1065b81 100644 --- a/src/tools/vault.py +++ b/src/tools/vault.py @@ -68,6 +68,10 @@ async def do_vault_search(content: str, owner: Optional[str] = None) -> Dict: except json.JSONDecodeError: return {"error": "Failed to parse bw output", "exit_code": 1} + from src.agent_runtime.owned_resources import active_owned_operation, observe_vault_records + if active_owned_operation() is not None: + observe_vault_records(owner, cfg, items) + if not items: return {"output": f"No vault items match '{query}'.", "exit_code": 0} @@ -79,7 +83,7 @@ async def do_vault_search(content: str, owner: Optional[str] = None) -> Dict: username = login.get("username", "") uris = login.get("uris") or [] url = uris[0].get("uri", "") if uris else "" - parts = [f"[{item_id[:8]}] {name}"] + parts = [f"[{item_id}] {name}"] if username: parts.append(f"user: {username}") if url: diff --git a/tests/runtime_evidence_helpers.py b/tests/runtime_evidence_helpers.py index 1ab245a15..e12c633bd 100644 --- a/tests/runtime_evidence_helpers.py +++ b/tests/runtime_evidence_helpers.py @@ -14,16 +14,20 @@ def server_authorized_executor(executor): from src.agent_runtime.authority import OperationGrant, RequestAuthority from src.tool_policy import known_tool_names from src.turn_contract import canonical_tool + from src.agent_runtime.remote_resources import seal_backends call_signature = signature(executor) @wraps(executor) async def execute(*args, **kwargs): bound = call_signature.bind(*args, **kwargs) parameters = bound.arguments + grants = tuple(OperationGrant(name) for name in sorted( + {canonical_tool(n) for n in known_tool_names()} | {"list_dir", "find_files"})) kwargs.setdefault("request_authority", RequestAuthority( "standalone-test-request", str(parameters.get("owner") or "").strip().casefold(), str(parameters.get("session_id") or ""), str(parameters.get("workspace") or ""), - tuple(OperationGrant(name) for name in sorted( - {canonical_tool(n) for n in known_tool_names()} | {"list_dir", "find_files"})), + grants, + backend_resources=seal_backends((g.tool for g in grants), context=parameters.get("client_runtime_context"), + owner=str(parameters.get("owner") or "").strip().casefold()), )) return await executor(*args, **kwargs) return execute diff --git a/tests/test_owned_resource_identity.py b/tests/test_owned_resource_identity.py new file mode 100644 index 000000000..79eb2e39b --- /dev/null +++ b/tests/test_owned_resource_identity.py @@ -0,0 +1,429 @@ +"""Server resolution pins record aliases before approval and dispatch.""" +import asyncio +from dataclasses import replace +from datetime import datetime, timedelta +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker + +from src.agent_runtime.authority import ExactOperation, OperationGrant, RequestAuthority, bind_request_authority +from src.agent_runtime.owned_resources import ( + active_owned_operation, admit_owned_operation, bind_owned_operation, + bound_attachment_path, resolve_owned_operation, + observe_vault_records, +) +from src.agent_runtime.remote_resources import active_backend_operation +from src.agent_runtime.resources import OwnedScope, ResourceIdentityError +from src.tool_approvals import ToolApprovalStore +from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action +from src.tool_types import ToolBlock + + +@pytest.fixture(autouse=True) +def fresh_vault_observations(monkeypatch): + from src.agent_runtime import owned_resources + monkeypatch.setattr(owned_resources, "_VAULT_RECORDS", {}) + + +def grant(*tools, scopes=None, owner="alice", thread="s"): + return RequestAuthority("owned-request", owner, thread, "", + tuple(OperationGrant(t) for t in tools), owned_scopes=scopes) + + +async def dispatch(authority, tool, content, **kwargs): + from src import tool_execution as execution + return await execution.execute_tool_block(ToolBlock(tool, content), owner=authority.owner, + session_id=authority.session_id, request_authority=authority, + security_context=kwargs.pop("security_context", execution.NO_TOOL_SECURITY_CONTEXT), **kwargs) + + +def approval(authority, tool, content, **kwargs): + store = ToolApprovalStore() + pending = store.create(owner=authority.owner, session_id=authority.session_id, origin_run_id="run", + tool_name=tool, content=content, workspace=None, request_authority=authority, + external_untrusted_context_seen=True, capabilities=capabilities_for_action(tool, content), **kwargs) + return store.consume(pending.approval_id, decision="approve", owner=authority.owner, session_id=authority.session_id) + + +@pytest.fixture +def records(monkeypatch): + import core.database as db + import src.database as compatibility + from src.agent_tools import document_tools + engine = create_engine("sqlite:///:memory:") + db.Base.metadata.create_all(engine) + factory = sessionmaker(bind=engine) + monkeypatch.setattr(db, "SessionLocal", factory) + monkeypatch.setattr(compatibility, "SessionLocal", factory) + from src import tool_execution as execution + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + now = datetime(2026, 1, 1) + with factory() as connection: + for identifier, owner in (("s", "alice"), ("other", "alice"), ("foreign", "bob")): + connection.add(db.Session(id=identifier, owner=owner, name=identifier, endpoint_url="https://model.test", model="test")) + for identifier, owner, thread, offset in (("d1", "alice", "s", 1), ("d2", "alice", "other", 2), ("private", "bob", "foreign", 3)): + connection.add(db.Document(id=identifier, owner=owner, session_id=thread, title=identifier, + current_content=identifier + " original", language="text", version_count=1, + created_at=now, updated_at=now + timedelta(days=offset))) + connection.add(db.Note(id="note-one", owner="alice", title="first", content="original")) + connection.add(db.Note(id="note-two", owner="alice", title="second", content="original")) + connection.add(db.Note(id="note-foreign", owner="bob", title="private", content="private")) + connection.commit() + monkeypatch.setattr(document_tools, "_active_document_id", None) + yield factory + engine.dispose() + + +@pytest.mark.parametrize("selector", ["active", "current", "latest"]) +async def test_document_alias_resolves_once_and_does_not_follow_new_active_or_latest(records, monkeypatch, selector): + from src import tool_execution as execution + authority = grant("manage_documents") + content = json.dumps({"action": "read", "document_id": selector}) + exact = approval(authority, "manage_documents", content, document_id="d1") + bound = exact.pending.owned_operation + expected = "d2" if selector == "latest" else "d1" + assert bound.document_id == expected + import core.database as db + with records() as connection: + connection.add(db.Document(id="newest", owner="alice", session_id="s", title="newest", + current_content="newest content", version_count=1, updated_at=datetime(2030, 1, 1))) + connection.commit() + seen = [] + async def implementation(block, **kwargs): + seen.append((json.loads(block.content)["document_id"], kwargs["approved_document_id"], active_owned_operation())) + return "read", {"exit_code": 0} + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + _, result = await dispatch(authority, "manage_documents", content, active_document_id="newest", + exact_approval=exact, security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["exit_code"] == 0 + assert seen[0][:2] == (expected, expected) + assert seen[0][2] is bound + assert active_owned_operation() is None + + +async def test_document_runtime_executes_captured_id_not_process_global_alias(records): + from src.agent_tools import document_tools + authority = grant("update_document") + exact = approval(authority, "update_document", "replacement", document_id="d1") + document_tools.set_active_document("d2") + _, result = await dispatch(authority, "update_document", "replacement", active_document_id="d2", + exact_approval=exact, security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result.get("exit_code", 0) == 0 and not result.get("error") + import core.database as db + with records() as connection: + assert connection.get(db.Document, "d1").current_content == "replacement" + assert connection.get(db.Document, "d2").current_content == "d2 original" + + +@pytest.mark.parametrize("change", ["revision", "owner", "thread", "deleted", "request", "invocation_thread"]) +async def test_stale_or_rebound_document_approval_fails_before_effect(records, monkeypatch, change): + import core.database as db + from src import tool_execution as execution + authority = grant("update_document") + exact = approval(authority, "update_document", "replacement", document_id="d1") + if change in {"request", "invocation_thread"}: + authority = replace(authority, **({"request_id": "other"} if change == "request" else + {"session_id": "other", "owned_scopes": (OwnedScope("documents", "alice", "other"),)})) + else: + with records() as connection: + row = connection.get(db.Document, "d1") + if change == "revision": + row.current_content = "changed" + row.version_count += 1 + elif change == "owner": + row.owner = "bob" + elif change == "thread": + row.session_id = "other" + else: + connection.delete(row) + connection.commit() + implementation = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + _, result = await dispatch(authority, "update_document", "replacement", exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["failure_kind"] == "resource_identity_denied" + implementation.assert_not_awaited() + assert not exact._claimed + + +@pytest.mark.parametrize("owner,thread", [("", "s"), ("alice", ""), ("bob", "s")]) +async def test_owner_and_invocation_thread_are_mandatory(records, owner, thread): + _, result = await dispatch(grant("manage_documents", owner=owner, thread=thread), + "manage_documents", '{"action":"read","id":"d1"}') + assert result["failure_kind"] == "resource_identity_denied" + + +async def test_child_record_scope_is_intersection_and_exact_approval_cannot_widen(records, monkeypatch): + parent = grant("manage_documents", scopes=(OwnedScope("documents", "alice", "s", frozenset({"d1"})),)) + child = grant("manage_documents") + assert parent.intersect(child).owned_scopes == parent.owned_scopes + exact = approval(child, "manage_documents", '{"action":"read","id":"d2"}') + from src import tool_execution as execution + handler = AsyncMock(return_value=("read", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + with bind_request_authority(parent): + _, denied = await dispatch(child, "manage_documents", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + _, allowed = await dispatch(child, "manage_documents", '{"action":"read","id":"d1"}') + _, collection = await dispatch(child, "manage_documents", '{"action":"list"}') + assert denied["failure_kind"] == collection["failure_kind"] == "resource_identity_denied" + assert allowed["exit_code"] == 0 and handler.await_count == 1 + + +async def test_legacy_owned_approval_is_exact_one_use_not_reconstructed_scope(records, monkeypatch): + snapshot = grant("manage_documents").to_dict() + snapshot["version"] = 2 + authority = RequestAuthority.from_dict(snapshot) + exact = approval(authority, "manage_documents", '{"action":"read","id":"d1"}') + from src import tool_execution as execution + handler = AsyncMock(return_value=("read", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + _, denied = await dispatch(authority, "manage_documents", exact.pending.content) + assert denied["failure_kind"] == "resource_identity_denied" + security = ToolRunSecurityContext(external_untrusted_context_seen=True) + _, allowed = await dispatch(authority, "manage_documents", exact.pending.content, exact_approval=exact, security_context=security) + _, replay = await dispatch(authority, "manage_documents", exact.pending.content, exact_approval=exact, security_context=security) + assert allowed["exit_code"] == 0 and replay["exit_code"] == 1 + assert authority.owned_scopes == () and handler.await_count == 1 + + +@pytest.mark.parametrize("tool,content", [("manage_session", '{"action":"rename","session":"current","value":"new"}'), + ("manage_session", "rename\ncurrent\nnew"), ("send_to_session", "current\nhello")]) +def test_thread_current_alias_becomes_exact_owned_identity(records, tool, content): + bound = resolve_owned_operation(ExactOperation.normalize(tool, content), owner="alice", thread_id="s") + assert bound.resources[0].record_id == bound.resources[0].record_thread_id == "s" + assert "current" not in bound.execution_input + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize(tool, content.replace("current", "foreign")), owner="alice", thread_id="s") + + +@pytest.mark.parametrize("selector", ["note-o", "first"]) +def test_note_alias_resolves_once_and_prefix_ambiguity_fails_closed(records, selector): + content = {"action": "update", "content": "replacement", "id" if selector == "note-o" else "title": selector} + bound = resolve_owned_operation(ExactOperation.normalize("manage_notes", json.dumps(content)), owner="alice", thread_id="s") + assert bound.resources[0].record_id == json.loads(bound.execution_input)["id"] == "note-one" + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("manage_notes", '{"action":"view","id":"note-"}'), owner="alice", thread_id="s") + + +@pytest.fixture +def attachment_store(tmp_path, monkeypatch): + from src import tool_utils + file = tmp_path / "image.png" + file.write_bytes(b"test image") + row = {"id": "upload.png", "owner": "alice", "path": str(file), "hash": "observed-hash"} + handler = SimpleNamespace(upload_dir=str(tmp_path), resolve_upload=lambda identifier, *, owner, allow_admin: + dict(row) if identifier == row.get("id") and owner == row.get("owner") and not allow_admin else None) + monkeypatch.setattr(tool_utils, "get_upload_handler", lambda: handler) + return row, file + + +@pytest.mark.parametrize("change", ["owner", "path", "file", "missing"]) +def test_attachment_ownership_index_and_file_identity_are_pinned(attachment_store, change): + row, file = attachment_store + bound = resolve_owned_operation(ExactOperation.normalize("extract_text", '{"path":"odysseus://attachment/upload.png"}'), owner="alice", thread_id="s") + with bind_owned_operation(bound): + assert bound_attachment_path("alice", "odysseus://attachment/upload.png") == str(file) + with pytest.raises(ResourceIdentityError): + bound_attachment_path("bob", "odysseus://attachment/upload.png") + if change == "owner": + row["owner"] = "bob" + elif change == "path": + other = file.with_name("other.png") + other.write_bytes(b"other") + row["path"] = str(other) + elif change == "file": + file.rename(file.with_name("old.png")) + file.write_bytes(b"replacement") + else: + row.clear() + with pytest.raises(ResourceIdentityError): + bound.validate() + + +def test_memory_prefix_is_owner_scoped_exact_and_revision_sensitive(monkeypatch): + from src import ai_interaction + rows = [{"id": "memory-one", "owner": "alice", "text": "secret", "timestamp": 1}, + {"id": "memory-other", "owner": "bob", "text": "private", "timestamp": 1}] + monkeypatch.setattr(ai_interaction, "_memory_manager", SimpleNamespace(load=lambda owner: rows)) + bound = resolve_owned_operation(ExactOperation.normalize("manage_memory", "edit\nmemory-o\nreplacement"), owner="alice", thread_id="s") + assert bound.resources[0].record_id == "memory-one" + assert bound.execution_input == "edit\nmemory-one\nreplacement" + assert "secret" not in json.dumps(bound.to_dict()) + rows[0]["text"] = "changed in same second" + with pytest.raises(ResourceIdentityError): + bound.validate() + + +@pytest.mark.parametrize("change", ["owner", "endpoint", "session"]) +def test_private_vault_identity_binds_owner_endpoint_and_exact_uuid_without_credentials(monkeypatch, change): + from src.tools import vault + cfg = {"owner": "alice", "server_url": "https://vault.test", "session": "SECRET_SESSION", "unlocked_at": "observed"} + monkeypatch.setattr(vault, "_load_vault_config", lambda: cfg) + observe_vault_records("alice", cfg, [{"id": "12345678-1234-1234-1234-123456789abc", "name": "bank", "login": {"password": "PRIVATE_PASSWORD"}}]) + tool = ExactOperation.normalize("vault_get", '{"item_id":"12345678-1234-1234-1234-123456789abc","reason":"requested"}') + bound = resolve_owned_operation(tool, owner="alice", thread_id="s") + assert "SECRET_SESSION" not in json.dumps(bound.to_dict()) and "PRIVATE_PASSWORD" not in json.dumps(bound.to_dict()) + cfg.update({"owner": "bob"} if change == "owner" else {"server_url": "https://other.test"} if change == "endpoint" else {"session": "OTHER_SECRET"}) + with pytest.raises(ResourceIdentityError): + bound.validate() + + +def test_legacy_vault_and_model_name_alias_do_not_create_private_identity(monkeypatch): + from src.tools import vault + monkeypatch.setattr(vault, "_load_vault_config", lambda: {"server_url": "https://vault.test", "session": "secret"}) + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("vault_search", '{"query":"bank"}'), owner="alice", thread_id="s") + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("vault_get", '{"item_id":"latest","reason":"requested"}'), owner="alice", thread_id="s") + + +@pytest.mark.parametrize("error", [None, RuntimeError, asyncio.CancelledError]) +async def test_owned_and_backend_context_restore_after_success_error_cancel_and_nested_call(records, monkeypatch, error): + from src import tool_execution as execution + authority = grant("manage_documents") + parent = admit_owned_operation(authority, ExactOperation.normalize("manage_documents", '{"action":"read","id":"d1"}')) + async def implementation(block, **kwargs): + assert active_owned_operation().resources[0].record_id == "d2" + assert active_backend_operation() is not None + if error: + raise error("stop") + return "read", {"exit_code": 0} + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + with bind_owned_operation(parent): + if error: + with pytest.raises(error): + await dispatch(authority, "manage_documents", '{"action":"read","id":"d2"}') + else: + await dispatch(authority, "manage_documents", '{"action":"read","id":"d2"}') + assert active_owned_operation() is parent + assert active_backend_operation() is None + assert active_owned_operation() is None + + +@pytest.mark.parametrize("path", ["/api/document/d1", "/api/history/s", "/api/vault/config", "/api/memory", "/api/notes", + "/api/upload/upload.png", "/api/%64ocument/d1", "/api/cookbook/../document/d1"]) +async def test_generic_internal_bridge_cannot_bypass_owned_resource_adapter(monkeypatch, path): + from src import tool_execution as execution + handler = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + _, result = await dispatch(grant("app_api"), "app_api", json.dumps({"path": path})) + assert result["failure_kind"] == "resource_identity_denied" + handler.assert_not_awaited() + + +def test_owned_snapshot_roundtrip_and_malformed_scopes_fail_closed(records): + authority = grant("manage_documents", scopes=(OwnedScope("documents", "alice", "s", frozenset({"d1"})),)) + snapshot = json.loads(json.dumps(authority.to_dict())) + assert RequestAuthority.from_dict(snapshot) == authority + for mutation in ({"owner": "bob"}, {"thread_id": "other"}, {"record_ids": ["*"]}, {"record_ids": [1]}): + changed = json.loads(json.dumps(snapshot)) + changed["owned_scopes"][0].update(mutation) + with pytest.raises((ValueError, TypeError)): + RequestAuthority.from_dict(changed) + + +def test_vault_config_owner_is_produced_by_authenticated_request_and_drops_legacy_session(): + from routes.vault.vault_routes import _bind_config_owner + from fastapi import HTTPException + request = SimpleNamespace(state=SimpleNamespace(current_user="alice", api_token=False)) + cfg = {"session": "legacy-secret", "unlocked_at": "legacy"} + _bind_config_owner(cfg, request) + assert cfg == {"owner": "alice"} + request.state.current_user = "bob" + with pytest.raises(HTTPException): + _bind_config_owner(cfg, request) + + +def test_missing_proposal_record_is_not_reconstructed_after_it_appears(records): + authority = grant("manage_documents") + exact = approval(authority, "manage_documents", '{"action":"read","id":"not-yet"}') + assert exact.pending.owned_operation is None + import core.database as db + with records() as connection: + connection.add(db.Document(id="not-yet", owner="alice", title="appeared", current_content="content", version_count=1)) + connection.commit() + _, result = asyncio.run(dispatch(authority, "manage_documents", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True))) + assert result["failure_kind"] == "resource_identity_denied" and not exact._claimed + + +def test_malformed_record_identity_and_normalized_approval_tampering_fail_closed(records): + authority = grant("manage_documents") + exact = approval(authority, "manage_documents", '{"action":"read","id":"d1"}') + bound = exact.pending.owned_operation + with pytest.raises(ValueError): + replace(bound, resources=(replace(bound.resources[0], revision=""),)) + with pytest.raises(ValueError): + replace(bound, document_id="d2") + exact.pending = replace(exact.pending, owned_operation=replace(bound, execution_input='{"action":"read","document_id":"d2"}')) + assert not exact.matches(owner="alice", session_id="s", workspace=None, + tool_name="manage_documents", content=exact.pending.content) + + +def test_note_prefix_wildcards_cannot_create_selector_authority(records): + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("manage_notes", '{"action":"view","id":"note-o%"}'), owner="alice", thread_id="s") + + +@pytest.mark.parametrize("content", ['{"action":"read"}', '{"action":"read","id":"active"}', '{"action":"read","id":"current"}', '{"action":"read","id":42}', '{"action":"read","id":"d1","uid":"d2"}']) +def test_missing_malformed_and_conflicting_document_selectors_fail_closed(records, content): + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("manage_documents", content), owner="alice", thread_id="s") + + +async def test_resumed_child_approval_cannot_restore_excluded_record(records): + parent = grant("manage_documents", scopes=(OwnedScope("documents", "alice", "s", frozenset({"d1"})),)) + child = parent.intersect(grant("manage_documents")) + exact = approval(child, "manage_documents", '{"action":"read","id":"d2"}') + assert exact.pending.owned_operation is None + _, result = await dispatch(replace(child, inherited=False), "manage_documents", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["failure_kind"] == "resource_identity_denied" and not exact._claimed + + +async def test_vault_search_producer_supplies_exact_owned_item_identity_and_alias_binding(monkeypatch): + from src.tools import vault + from src import tool_execution as execution + cfg = {"owner": "alice", "server_url": "https://vault.test", "session": "SECRET_SESSION", "unlocked_at": "observed"} + item = {"id": "12345678-1234-1234-1234-123456789abc", "name": "bank", "login": {"password": "PRIVATE_PASSWORD"}} + monkeypatch.setattr(vault, "_load_vault_config", lambda: cfg) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + cli = AsyncMock(side_effect=[(json.dumps([item]), "", 0), (json.dumps(item), "", 0)]) + monkeypatch.setattr(vault, "_run_bw", cli) + authority = grant("vault_search", "vault_get") + _, missing = await dispatch(authority, "vault_get", json.dumps({"item_id": item["id"], "reason": "requested"})) + assert missing["failure_kind"] == "resource_identity_denied" + cli.assert_not_awaited() + _, search = await dispatch(authority, "vault_search", '{"query":"bank"}') + assert search["exit_code"] == 0 and item["id"] in search["output"] + exact = approval(authority, "vault_get", '{"item_id":"bank","reason":"requested"}') + assert exact.pending.owned_operation.resources[0].record_id == item["id"] + # A later producer result with the same alias cannot change the approved ID. + other = {"id": "87654321-1234-1234-1234-123456789abc", "name": "bank", "login": {"password": "OTHER_PASSWORD"}} + observe_vault_records("alice", cfg, [other]) + _, result = await dispatch(authority, "vault_get", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["exit_code"] == 0 and "PRIVATE_PASSWORD" in result["output"] and "OTHER_PASSWORD" not in result["output"] + assert cli.await_args.args[0] == ["get", "item", item["id"]] + with pytest.raises(ResourceIdentityError): + resolve_owned_operation(ExactOperation.normalize("vault_get", exact.pending.content), owner="alice", thread_id="s") + + +def test_vault_producer_revision_change_invalidates_sealed_item_and_cannot_cross_owner(monkeypatch): + from src.tools import vault + cfg = {"owner": "alice", "server_url": "https://vault.test", "session": "SECRET"} + item = {"id": "12345678-1234-1234-1234-123456789abc", "name": "bank", "revisionDate": "one"} + monkeypatch.setattr(vault, "_load_vault_config", lambda: cfg) + observe_vault_records("alice", cfg, [item]) + bound = resolve_owned_operation(ExactOperation.normalize("vault_get", '{"item_id":"12345678","reason":"requested"}'), owner="alice", thread_id="s") + assert json.loads(bound.execution_input)["item_id"] == item["id"] + with pytest.raises(ResourceIdentityError): + observe_vault_records("bob", cfg, [item]) + observe_vault_records("alice", cfg, [{**item, "revisionDate": "two"}]) + with pytest.raises(ResourceIdentityError): + bound.validate() diff --git a/tests/test_remote_resource_identity.py b/tests/test_remote_resource_identity.py new file mode 100644 index 000000000..885059ad5 --- /dev/null +++ b/tests/test_remote_resource_identity.py @@ -0,0 +1,397 @@ +"""Backend selection is resolution, never an operation or resource grant.""" +import asyncio +from dataclasses import replace +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock + +import pytest + +from src.agent_runtime.authority import ExactOperation, OperationGrant, RequestAuthority, bind_request_authority +from src.agent_runtime.remote_resources import ( + active_backend_operation, bind_backend_operation, bind_backend_for_operation, + configuration_incarnation, endpoint_identity, integration_resource, seal_backends, +) +from src.agent_runtime.resources import ExternalResource, NativeBackendResource, ResourceIdentityError +from src.mcp_manager import McpManager +from src.tool_approvals import ToolApprovalStore +from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action +from src.tool_types import ToolBlock + + +def grant(*tools, resources=None): + return RequestAuthority("remote-request", "alice", "s", "", + tuple(OperationGrant(t) for t in tools), backend_resources=resources) + + +async def dispatch(authority, tool, content="{}", **kwargs): + from src import tool_execution as execution + return await execution.execute_tool_block(ToolBlock(tool, content), owner=authority.owner, + session_id=authority.session_id, request_authority=authority, + security_context=kwargs.pop("security_context", execution.NO_TOOL_SECURITY_CONTEXT), **kwargs) + + +def approval(authority, tool, content="{}", **kwargs): + store = ToolApprovalStore() + pending = store.create(owner=authority.owner, session_id=authority.session_id, origin_run_id="run", + tool_name=tool, content=content, workspace=None, request_authority=authority, + external_untrusted_context_seen=True, capabilities=capabilities_for_action(tool, content), **kwargs) + return store.consume(pending.approval_id, decision="approve", owner=authority.owner, session_id=authority.session_id) + + +@pytest.fixture +def manager(monkeypatch): + from src import tool_execution as execution + value = McpManager() + monkeypatch.setattr(execution, "get_mcp_manager", lambda: value) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + return value + + +def connect(manager, server="alpha", tools=("read", "write"), url="https://example.test/mcp?token=SECRET"): + session = SimpleNamespace(call_tool=AsyncMock(return_value=SimpleNamespace( + content=[SimpleNamespace(text="remote result")], isError=False))) + manager._sessions[server] = session + manager._tools[server] = [{"name": tool} for tool in tools] + manager._resource_endpoints[server] = (endpoint_identity(url), configuration_incarnation(url)) + manager._register_resource_connection(server, session) + return session + + +@pytest.mark.parametrize("kind", ["availability", "selection", "model_name", "legacy"]) +async def test_remote_availability_does_not_create_resource_authority(manager, kind): + session = connect(manager) + authority = grant("mcp__alpha__read", resources=()) if kind != "model_name" else grant() + if kind == "legacy": + snapshot = grant("mcp__alpha__read").to_dict() + snapshot["version"] = 2 + authority = RequestAuthority.from_dict(snapshot) + _, result = await dispatch(authority, "mcp__alpha__read") + assert result["exit_code"] == 1 + session.call_tool.assert_not_awaited() + + +async def test_qualified_mcp_binds_exact_tool_and_backend(manager): + session = connect(manager) + authority = grant("mcp__alpha__read") + _, allowed = await dispatch(authority, "mcp__alpha__read", '{"record":"one"}') + assert allowed["exit_code"] == 0 + session.call_tool.assert_awaited_once_with("read", {"record": "one"}) + _, denied = await dispatch(replace(authority, grants=(OperationGrant("mcp__alpha__write"),)), "mcp__alpha__write") + assert denied["failure_kind"] == "resource_identity_denied" + assert active_backend_operation() is None + + +@pytest.mark.parametrize("change", ["session", "endpoint", "path", "query", "discovery"]) +async def test_remote_identity_changes_invalidate_admission_and_exact_approval(manager, change): + old = connect(manager) + authority = grant("mcp__alpha__read") + exact = approval(authority, "mcp__alpha__read", '{"resource":"one"}') + assert exact.pending.backend_operation is not None + if change == "session": + new = connect(manager) + elif change == "discovery": + manager._tools["alpha"] = [{"name": "write"}] + else: + url = {"endpoint": "https://other.test/mcp", "path": "https://example.test/other", + "query": "https://example.test/mcp?token=OTHER"}[change] + manager._resource_endpoints["alpha"] = (endpoint_identity(url), configuration_incarnation(url)) + for extra in ({}, {"exact_approval": exact, "security_context": ToolRunSecurityContext(external_untrusted_context_seen=True)}): + _, result = await dispatch(authority, "mcp__alpha__read", '{"resource":"one"}', **extra) + assert result["failure_kind"] == "resource_identity_denied" + old.call_tool.assert_not_awaited() + if change == "session": + new.call_tool.assert_not_awaited() + assert not exact._claimed + + +@pytest.mark.parametrize("change", ["tool", "selector", "request", "owner", "session"]) +async def test_remote_approval_is_bound_to_operation_and_request(manager, change): + session = connect(manager) + authority = grant("mcp__alpha__read", "mcp__alpha__write") + exact = approval(authority, "mcp__alpha__read", '{"record":"one"}') + tool, content = "mcp__alpha__read", '{"record":"one"}' + if change == "tool": + tool = "mcp__alpha__write" + elif change == "selector": + content = '{"record":"two"}' + else: + authority = replace(authority, **{"request": {"request_id": "other"}, "owner": {"owner": "bob"}, + "session": {"session_id": "other"}}[change]) + _, result = await dispatch(authority, tool, content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["exit_code"] == 1 + session.call_tool.assert_not_awaited() + assert not exact._claimed + + +async def test_legacy_exact_remote_approval_is_one_use_and_does_not_mint_backend_scope(manager): + session = connect(manager) + authority = grant(resources=()) + exact = approval(authority, "mcp__alpha__read") + security = ToolRunSecurityContext(external_untrusted_context_seen=True) + _, result = await dispatch(authority, "mcp__alpha__read", exact_approval=exact, security_context=security) + assert result["exit_code"] == 0 + _, replay = await dispatch(authority, "mcp__alpha__read", exact_approval=exact, security_context=security) + assert replay["exit_code"] == 1 + assert session.call_tool.await_count == 1 and authority.backend_resources == () + assert "backend_operation" not in exact.pending.public_payload() + + +async def test_child_cannot_use_parent_ungranted_backend_or_exact_approval(manager): + session = connect(manager) + parent = grant("mcp__alpha__read", resources=()) + child = grant("mcp__alpha__read") + exact = approval(child, "mcp__alpha__read") + with bind_request_authority(parent): + _, denied = await dispatch(child, "mcp__alpha__read", exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert denied["failure_kind"] == "resource_identity_denied" + session.call_tool.assert_not_awaited() + + +def test_remote_snapshots_exclude_credentials_and_cannot_claim_containment(manager): + connect(manager, url="https://user:PASSWORD@example.test/SECRET_PATH?token=TOKEN") + authority = grant("mcp__alpha__read") + snapshot = json.dumps(authority.to_dict()) + assert all(secret not in snapshot for secret in ("PASSWORD", "SECRET_PATH", "TOKEN", "user:")) + resource = authority.backend_resources[0] + assert resource.endpoint_id == "https://example.test" + assert resource.external is True and resource.contained is False + assert RequestAuthority.from_dict(json.loads(snapshot)) == authority + with pytest.raises(ValueError): + replace(resource, contained=True) + + +async def test_mcp_revalidates_at_transport_and_never_retries_bound_calls(manager): + manager._resource_owners["memory"] = "alice" + session = connect(manager, server="memory") + authority = grant("mcp__memory__read") + operation = ExactOperation.normalize("mcp__memory__read", "{}") + bound = bind_backend_for_operation(authority, operation) + reconnect = AsyncMock() + manager._reconnect_builtin = reconnect + session.call_tool.side_effect = RuntimeError("disconnected") + with bind_backend_operation(bound): + result = await manager.call_tool(operation.tool, {}) + assert result["exit_code"] == 1 + replacement = connect(manager, server="memory") + result = await manager.call_tool(operation.tool, {}) + assert result["failure_kind"] == "resource_identity_denied" + replacement.call_tool.assert_not_awaited() + reconnect.assert_not_awaited() + + +@pytest.mark.parametrize("owner", ["", "bob"]) +async def test_builtin_memory_backend_requires_its_configured_owner(manager, owner): + manager._resource_owners["memory"] = owner + session = connect(manager, server="memory") + _, result = await dispatch(grant("mcp__memory__read"), "mcp__memory__read") + assert result["failure_kind"] == "resource_identity_denied" + session.call_tool.assert_not_awaited() + + +async def test_native_filesystem_cannot_be_redirected_through_mcp(manager, tmp_path): + session = connect(manager, server="filesystem", tools=("read_file",)) + (tmp_path / "a").write_text("native contents") + authority = RequestAuthority("request", "alice", "s", str(tmp_path), (OperationGrant("read_file"),)) + from src import tool_execution as execution + _, result = await execution.execute_tool_block(ToolBlock("read_file", "a"), owner="alice", session_id="s", + workspace=str(tmp_path), request_authority=authority, security_context=execution.NO_TOOL_SECURITY_CONTEXT) + assert result["output"] == "native contents" + assert isinstance(authority.backend_resources[0], NativeBackendResource) + session.call_tool.assert_not_awaited() + + +@pytest.mark.parametrize("change", ["alias", "endpoint", "secret_path"]) +async def test_integration_alias_and_configuration_cannot_retarget_approval(monkeypatch, change): + from src import integrations + rows = [{"id": "one", "name": "service", "base_url": "https://service.test/SECRET", "enabled": True}] + monkeypatch.setattr(integrations, "load_integrations", lambda: rows) + authority = grant("api_call", resources=(integration_resource(rows[0]),)) + exact = approval(authority, "api_call", '{"integration":"service","path":"/record/one"}') + assert exact.pending.backend_operation.resource.server_id == "one" + assert "SECRET" not in json.dumps(authority.to_dict()) + if change == "alias": + rows[:] = [{**rows[0], "id": "two"}] + else: + rows[0]["base_url"] = "https://other.test/SECRET" if change == "endpoint" else "https://service.test/OTHER" + from src import tool_execution as execution + handler = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + _, result = await dispatch(authority, "api_call", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["failure_kind"] == "resource_identity_denied" + handler.assert_not_awaited() + + +@pytest.mark.parametrize("error", [None, RuntimeError, asyncio.CancelledError]) +async def test_scoped_bridge_context_restores_and_replacement_is_ungranted(monkeypatch, error): + from src import tool_execution as execution + seen = [] + async def route(*args): + seen.append(active_backend_operation().resource) + if error: + raise error("stop") + return "bridge", {"exit_code": 0} + bridge = execution.AgentExecutionBridge(route, frozenset({"host_shell"}), name="test") + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + with execution.bind_execution_bridge(bridge): + authority = grant("host_shell") + parent_bound = bind_backend_for_operation(authority, ExactOperation.normalize("host_shell", "parent")) + with bind_backend_operation(parent_bound): + if error is asyncio.CancelledError: + with pytest.raises(error): + await dispatch(authority, "host_shell", "pwd") + else: + await dispatch(authority, "host_shell", "pwd") + assert active_backend_operation() is parent_bound + assert active_backend_operation() is None + assert seen[0].external and not seen[0].contained + with execution.bind_execution_bridge(replace(bridge)): + _, denied = await dispatch(authority, "host_shell", "pwd") + assert denied["failure_kind"] == "resource_identity_denied" + + +async def test_tui_endpoint_is_registered_only_at_trusted_admission(monkeypatch): + from src import tool_execution as execution + context = {"surface": "odysseus-tui", "host_shell_bridge": {"url": "http://127.0.0.1:17654/run", "token": "TOKEN"}} + authority = grant("bash", resources=(NativeBackendResource("bash"),)) + _, denied = await dispatch(authority, "bash", "pwd", client_runtime_context=context) + assert denied["failure_kind"] == "resource_identity_denied" + authority = replace(authority, backend_resources=seal_backends(["bash"], context=context, owner="alice")) + monkeypatch.setattr(execution, "_bridge_post", AsyncMock(return_value={"exit_code": 0, "stdout": "external", "stderr": ""})) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + _, allowed = await dispatch(authority, "bash", "pwd", client_runtime_context=context) + assert allowed["exit_code"] == 0 + context["host_shell_bridge"]["url"] = "http://127.0.0.1:17655/run" + _, denied = await dispatch(authority, "bash", "pwd", client_runtime_context=context) + assert denied["failure_kind"] == "resource_identity_denied" + + +@pytest.mark.parametrize("change", [None, "url", "token"]) +async def test_http_bridge_factory_and_admission_share_config_identity(monkeypatch, change): + from src import tool_execution as execution + from routes.chat_routes import _external_execution_bridge + context = {"external_execution_bridge": {"url": "http://127.0.0.1:17654/execute", + "token": "SECRET_TOKEN", "supported_tools": ["host_shell"]}} + authority = grant("host_shell", resources=seal_backends(["host_shell"], context=context, owner="alice")) + if change: + context["external_execution_bridge"][change] = "http://127.0.0.1:17655/execute" if change == "url" else "OTHER_TOKEN" + bridge = _external_execution_bridge(context) + route = AsyncMock(return_value=("bridge", {"exit_code": 0})) + bridge = replace(bridge, route_tool=route) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + with execution.bind_execution_bridge(bridge): + _, result = await dispatch(authority, "host_shell", "pwd", client_runtime_context=context) + if change: + assert result["failure_kind"] == "resource_identity_denied" + route.assert_not_awaited() + else: + assert result["exit_code"] == 0 + route.assert_awaited_once() + assert "SECRET_TOKEN" not in json.dumps(authority.to_dict()) + + +async def test_backend_alias_cannot_retarget_a_legacy_tool_after_approval(manager, monkeypatch): + from src import tool_execution as execution + first = connect(manager, server="web_fetch", tools=("web_fetch",)) + second = connect(manager, server="other", tools=("fetch",)) + authority = grant("web_fetch") + exact = approval(authority, "web_fetch", "https://page.test/one") + monkeypatch.setitem(execution._MCP_TOOL_MAP, "web_fetch", ("other", "fetch")) + _, result = await dispatch(authority, "web_fetch", exact.pending.content, exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["failure_kind"] == "resource_identity_denied" + first.call_tool.assert_not_awaited() + second.call_tool.assert_not_awaited() + + +async def test_native_backend_is_pinned_when_mcp_becomes_available(manager, monkeypatch): + from src import tool_execution as execution + authority = grant("web_fetch") + session = connect(manager, server="web_fetch", tools=("web_fetch",)) + fallback = AsyncMock(return_value={"output": "native", "exit_code": 0}) + monkeypatch.setattr(execution, "_direct_fallback", fallback) + _, result = await dispatch(authority, "web_fetch", "https://page.test/one") + assert result["output"] == "native" + session.call_tool.assert_not_awaited() + + +async def test_integration_inventory_does_not_supply_backend_scope(monkeypatch): + from src import integrations + rows = [{"id": "one", "name": "service", "base_url": "https://service.test", "enabled": True}] + monkeypatch.setattr(integrations, "load_integrations", lambda: rows) + authority = grant("api_call") + assert authority.backend_resources == () + _, denied = await dispatch(authority, "api_call", '{"integration":"service"}') + assert denied["failure_kind"] == "resource_identity_denied" + + +async def test_external_bash_marker_cannot_switch_to_local_background_execution(manager, monkeypatch): + from src import bg_jobs + session = connect(manager, server="bash", tools=("bash",)) + launch = AsyncMock() + monkeypatch.setattr(bg_jobs, "launch", launch) + _, result = await dispatch(grant("bash"), "bash", "#!bg\npwd") + assert result["exit_code"] == 0 + assert session.call_tool.await_count == 1 + launch.assert_not_called() + + +@pytest.mark.parametrize("alias", ["host_shell_bridge", "hostShellBridge"]) +def test_host_shell_bridge_aliases_resolve_to_one_external_identity(alias): + context = {"surface": "odysseus-tui", alias: {"url": "http://127.0.0.1:17654/run", "token": "TOKEN"}} + resources = seal_backends(["host_shell"], context=context, owner="alice") + assert len(resources) == 1 and isinstance(resources[0], ExternalResource) + assert resources[0].external and not resources[0].contained + + +async def test_host_shell_cannot_reconstruct_an_unsealed_external_backend(monkeypatch): + from src import tool_execution as execution + handler = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + _, result = await dispatch(grant("host_shell"), "host_shell", "pwd") + assert result["failure_kind"] == "resource_identity_denied" + handler.assert_not_awaited() + + +async def test_http_backend_without_bound_producer_cannot_fall_back_to_native(monkeypatch): + from src import tool_execution as execution + context = {"external_execution_bridge": {"url": "http://127.0.0.1:17654/execute", + "token": "TOKEN", "supported_tools": ["bash"]}} + authority = grant("bash", resources=seal_backends(["bash"], context=context, owner="alice")) + fallback = AsyncMock() + monkeypatch.setattr(execution, "_direct_fallback", fallback) + _, denied = await dispatch(authority, "bash", "pwd", client_runtime_context=context) + assert denied["failure_kind"] == "resource_identity_denied" + fallback.assert_not_awaited() + + +async def test_resumed_child_approval_cannot_restore_excluded_backend(manager): + session = connect(manager) + child = grant("mcp__alpha__read", resources=()).intersect(grant("mcp__alpha__read")) + exact = approval(child, "mcp__alpha__read") + assert exact.pending.backend_operation is None + _, result = await dispatch(replace(child, inherited=False), "mcp__alpha__read", exact_approval=exact, + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True)) + assert result["failure_kind"] == "resource_identity_denied" + session.call_tool.assert_not_awaited() + + +async def test_integration_executes_server_resolved_id_and_revalidates_loaded_configuration(monkeypatch): + from src import integrations, tool_execution as execution + row = {"id": "one", "name": "service", "base_url": "https://service.test", "enabled": True} + monkeypatch.setattr(integrations, "load_integrations", lambda: [dict(row)]) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + authority = grant("api_call", resources=(integration_resource(row),)) + producer = AsyncMock(return_value={"output": "remote", "exit_code": 0}) + original = integrations.execute_api_call + monkeypatch.setattr(integrations, "execute_api_call", producer) + _, allowed = await dispatch(authority, "api_call", '{"integration":"service","method":"GET","path":"/record/one"}') + assert allowed["exit_code"] == 0 + assert producer.await_args.args == ("one", "GET", "/record/one") + monkeypatch.setattr(integrations, "execute_api_call", original) + monkeypatch.setattr(integrations, "_find_integration", lambda identifier: {**row, "base_url": "https://other.test"}) + _, denied = await dispatch(authority, "api_call", '{"integration":"service","path":"/record/one"}') + assert denied["failure_kind"] == "resource_identity_denied" diff --git a/tests/test_resource_identity.py b/tests/test_resource_identity.py index 8c9daf46b..bc61c15de 100644 --- a/tests/test_resource_identity.py +++ b/tests/test_resource_identity.py @@ -81,6 +81,31 @@ def test_symlink_escape_is_not_a_resource(tmp_path): resolve(authority(workspace, "read_file"), "read_file", "alias") +@pytest.mark.parametrize("alias", ["direct", "relative", "symlink", "hardlink"]) +@pytest.mark.parametrize("must_exist", [True, False]) +def test_media_workspace_paths_cannot_address_control_state(tmp_path, monkeypatch, alias, must_exist): + from src import constants, tool_execution + from src.agent_tools.media_tools import _resolve_workspace_path + control = tmp_path / "receipts.json" + control.write_text("private execution state") + monkeypatch.setattr(constants, "CONTAINMENT_STATE_FILE", str(control)) + monkeypatch.setattr(tool_execution, "get_active_workspace", lambda: str(tmp_path)) + if alias == "direct": + selector = str(control) + elif alias == "relative": + selector = "./receipts.json" + else: + target = tmp_path / "image.png" + if alias == "symlink": + target.symlink_to(control) + else: + os.link(control, target) + selector = "/workspace/image.png" + with pytest.raises(ValueError, match="execution-control"): + _resolve_workspace_path(selector, must_exist=must_exist) + assert control.read_text() == "private execution state" + + def test_destination_binds_absence_and_existing_ancestors(tmp_path): parent = tmp_path / "existing" parent.mkdir() @@ -143,6 +168,93 @@ async def test_user_filesystem_scope_cannot_write_server_execution_state(tmp_pat assert not (tmp_path / target).exists() +@pytest.mark.parametrize("state", ["authority", "jobs", "containment", "result", "exit", "database", "vault", "uploads"]) +@pytest.mark.parametrize("alias", ["direct", "relative", "symlink", "hardlink"]) +async def test_control_files_cannot_be_read_or_written_through_aliases(tmp_path, monkeypatch, state, alias): + import src.constants as constants + jobs = tmp_path / "jobs" + jobs.mkdir() + monkeypatch.setattr(constants, "BG_JOBS_DIR", str(jobs)) + monkeypatch.setattr(constants, "DATA_DIR", str(tmp_path)) + monkeypatch.setattr(constants, "UPLOAD_DIR", str(tmp_path / "uploads")) + for name, filename in (("BG_JOBS_FILE", "jobs.json"), ("CONTAINMENT_STATE_FILE", "receipts.json"), + ("APP_DB", "private.db"), ("VAULT_FILE", "vault.json")): + monkeypatch.setattr(constants, name, str(tmp_path / filename)) + filename = {"authority": "jobs/job.authority.json", "jobs": "jobs.json", "containment": "receipts.json", + "result": "jobs/job.result.json", "exit": "jobs/job.exit", "database": "private.db", + "vault": "vault.json", "uploads": "uploads/uploads.json"}[state] + target = tmp_path / filename + target.parent.mkdir(exist_ok=True) + target.write_text("control-secret") + selector = str(target) + if alias == "relative": + selector = "./" + filename + elif alias in {"symlink", "hardlink"}: + link = tmp_path / "ordinary.txt" + try: + link.symlink_to(target) if alias == "symlink" else os.link(target, link) + except OSError as error: + pytest.skip(f"Platform cannot create {alias}: {error}") + selector = str(link) + grant = authority(tmp_path, "read_file", "write_file") + for tool, content in (("read_file", selector), ("write_file", selector + "\nforged")): + _, result = await dispatch(grant, tool, content) + assert result["failure_kind"] == "resource_identity_denied" + assert target.read_text() == "control-secret" + + +async def test_directory_grep_does_not_scan_control_state_or_hardlinks(tmp_path, monkeypatch): + import src.constants as constants + control = tmp_path / "jobs.json" + control.write_text("UNIQUE_CONTROL_SECRET") + (tmp_path / "ordinary").write_text("visible text") + os.link(control, tmp_path / "innocent.txt") + monkeypatch.setattr(constants, "BG_JOBS_FILE", str(control)) + _, result = await dispatch(authority(tmp_path, "grep"), "grep", '{"pattern":"UNIQUE_CONTROL_SECRET","path":"."}') + assert result["exit_code"] == 0 + assert "No matches" in result["output"] + + +@pytest.mark.parametrize("tool,content", [("glob", '{"pattern":"*.json","path":"."}'), ("ls", ".")]) +async def test_directory_enumeration_does_not_address_control_files(tmp_path, monkeypatch, tool, content): + import src.constants as constants + control = tmp_path / "jobs.json" + control.write_text("control") + monkeypatch.setattr(constants, "BG_JOBS_FILE", str(control)) + _, result = await dispatch(authority(tmp_path, tool), tool, content) + assert result["exit_code"] == 0 + assert "jobs.json" not in result["output"] + + +@pytest.mark.parametrize("producer", ["database", "containment", "jobs", "uploads"]) +@pytest.mark.parametrize("alias", ["direct", "hardlink"]) +async def test_configured_control_producer_paths_are_protected(tmp_path, monkeypatch, producer, alias): + target = tmp_path / "custom" / "state" + target.parent.mkdir() + if producer == "database": + import core.database as database + monkeypatch.setattr(database, "engine", SimpleNamespace(url=SimpleNamespace( + get_backend_name=lambda: "sqlite", database=str(target)))) + elif producer == "containment": + from src import containment + monkeypatch.setattr(containment, "_store_path", lambda: target) + elif producer == "jobs": + from src import bg_jobs + monkeypatch.setattr(bg_jobs, "_STORE", target) + else: + from src import tool_utils + target = target.parent / "uploads.json" + monkeypatch.setattr(tool_utils, "get_upload_handler", lambda: SimpleNamespace(upload_dir=str(target.parent))) + target.write_text("server state") + selector = str(target) + if alias == "hardlink": + link = tmp_path / "ordinary" + os.link(target, link) + selector = str(link) + _, result = await dispatch(authority(tmp_path, "read_file"), "read_file", selector) + assert result["failure_kind"] == "resource_identity_denied" + + @pytest.mark.parametrize("roots", [(), None]) async def test_nonworkspace_allowlist_and_operation_do_not_grant_resources(tmp_path, monkeypatch, roots): from src import tool_execution as execution @@ -217,6 +329,43 @@ def test_child_cannot_renew_replaced_parent_root(tmp_path): assert parent.intersect(child).resource_roots == () +@pytest.mark.parametrize("caller", ["intersection", "context", "task"]) +@pytest.mark.parametrize("child_location", ["root", "subtree"]) +def test_replaced_parent_cannot_be_renewed_by_new_child_observation(tmp_path, caller, child_location): + root = tmp_path / "root" + root.mkdir() + parent = authority(root, "read_file") + root.rename(tmp_path / "old") + root.mkdir() + sub = root / "sub" + sub.mkdir() + (sub / "a").write_text("replacement") + child_root = FilesystemRoot.seal(root if child_location == "root" else sub, owner="alice") + child = authority(root, "read_file", roots=(child_root,)) + if caller == "intersection": + effective = parent.intersect(child) + elif caller == "context": + with bind_request_authority(parent), bind_request_authority(child) as effective: + assert effective.resource_roots == () + else: + with bind_request_authority(parent): + sealed = seal_task_authority("Read files in the workspace", "llm", None, owner="alice") + effective = restore_task_authority(sealed, "Read files in the workspace", "llm", None, + owner="alice", session_id="continuation") + assert effective.resource_roots == () + with pytest.raises(ValueError): + resolve(effective, "read_file", "sub/a") + + +def test_equal_stale_roots_are_revalidated(tmp_path): + root = tmp_path / "root" + root.mkdir() + parent = authority(root, "read_file") + root.rename(tmp_path / "old") + root.mkdir() + assert parent.intersect(parent).resource_roots == () + + @pytest.mark.parametrize("legacy", [False, True]) async def test_snapshot_preserves_incarnation_and_never_reconstructs_legacy(tmp_path, legacy): root = tmp_path / "root" @@ -336,7 +485,7 @@ async def test_last_dispatch_validation_refuses_replacement_and_resets_context(t assert execution.get_active_workspace() is None -@pytest.mark.parametrize("error_type", [RuntimeError, asyncio.CancelledError]) +@pytest.mark.parametrize("error_type", [None, RuntimeError, asyncio.CancelledError]) async def test_nested_resource_context_restores_on_failure_or_cancellation(tmp_path, monkeypatch, error_type): from src import tool_execution as execution for name in ("parent", "child"): @@ -345,10 +494,15 @@ async def test_nested_resource_context_restores_on_failure_or_cancellation(tmp_p parent = resolve(grant, "read_file", "parent") async def implementation(block, **kwargs): assert active_resource_operation().bindings[0].resource.path == str(tmp_path / "child") - raise error_type("stop") + if error_type: + raise error_type("stop") + return "read", {"exit_code": 0} monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) with bind_resource_operation(parent): - with pytest.raises(error_type): + if error_type: + with pytest.raises(error_type): + await dispatch(grant, "read_file", "child") + else: await dispatch(grant, "read_file", "child") assert active_resource_operation() is parent assert active_resource_operation() is None @@ -480,6 +634,63 @@ async def test_exact_user_approval_binds_only_one_missing_destination(tmp_path): assert not (tmp_path / "other.txt").exists() +@pytest.mark.parametrize("version", [1, 2]) +async def test_restored_empty_roots_approval_is_exact_and_never_restores_generic_scope(tmp_path, version): + (tmp_path / "approved").write_text("approved content") + (tmp_path / "sibling").write_text("private sibling") + snapshot = authority(tmp_path, "read_file", "write_file", "ls").to_dict() + snapshot["version"] = version + snapshot["resource_roots"] = [] + restored = RequestAuthority.from_dict(snapshot) + exact, security = approval(restored, "read_file", "approved") + assert exact.pending.resource_operation is not None + for tool, content in (("read_file", "sibling"), ("ls", "."), ("write_file", "sibling\nx")): + _, blocked = await dispatch(restored, tool, content, exact_approval=exact, security_context=security) + assert blocked["exit_code"] == 1 + _, unapproved = await dispatch(restored, tool, content) + assert unapproved["failure_kind"] == "resource_identity_denied" + _, allowed = await dispatch(restored, "read_file", "approved", exact_approval=exact, security_context=security) + assert allowed["output"] == "approved content" + _, replay = await dispatch(restored, "read_file", "approved", exact_approval=exact, security_context=security) + assert replay["exit_code"] == 1 + assert restored.resource_roots == () and restored.backend_resources == () + assert (tmp_path / "sibling").read_text() == "private sibling" + + +@pytest.mark.parametrize("change", ["alias", "request", "session", "owner"]) +async def test_restored_exact_filesystem_binding_rejects_retarget_and_rebinding(tmp_path, change): + (tmp_path / "a").write_text("a") + (tmp_path / "b").write_text("b") + (tmp_path / "alias").symlink_to(tmp_path / "a") + snapshot = authority(tmp_path, "read_file").to_dict() + snapshot["version"] = 1 + restored = RequestAuthority.from_dict(snapshot) + exact, security = approval(restored, "read_file", "alias") + if change == "alias": + (tmp_path / "alias").unlink() + (tmp_path / "alias").symlink_to(tmp_path / "b") + else: + restored = replace(restored, **{"request": {"request_id": "other"}, + "session": {"session_id": "other"}, "owner": {"owner": "bob"}}[change]) + _, result = await dispatch(restored, "read_file", "alias", exact_approval=exact, security_context=security) + assert result["exit_code"] == 1 + assert exact.matches(owner="alice", session_id="s", workspace=str(tmp_path), tool_name="read_file", content="alias") + + +async def test_resumed_child_approval_cannot_renew_replaced_parent_root(tmp_path): + root = tmp_path / "root" + root.mkdir() + parent = authority(root, "read_file") + root.rename(tmp_path / "old") + root.mkdir() + (root / "new").write_text("replacement") + child = parent.intersect(authority(root, "read_file")) + exact, security = approval(child, "read_file", "new") + assert child.resource_roots == () and exact.pending.resource_operation is None + _, result = await dispatch(replace(child, inherited=False), "read_file", "new", exact_approval=exact, security_context=security) + assert result["failure_kind"] == "resource_identity_denied" and not exact._claimed + + @pytest.mark.parametrize("request_text,denied", [ ("Transcribe /workspace/audio.wav", "read_file"), ("OCR extract exact text from /workspace/image.png", "write_file"), diff --git a/tests/test_tool_approvals.py b/tests/test_tool_approvals.py index be88b0d87..e5f793683 100644 --- a/tests/test_tool_approvals.py +++ b/tests/test_tool_approvals.py @@ -204,6 +204,16 @@ async def test_dispatcher_claims_approval_immediately_before_execution(monkeypat @pytest.mark.asyncio async def test_dispatcher_uses_sealed_document_target(monkeypatch): import src.tool_execution as tool_execution + from datetime import datetime + from types import SimpleNamespace + from src.agent_runtime import owned_resources + # This dispatcher fixture seals an observed owned row, as production does; + # model/document text alone cannot stand in for a resource identity. + row = SimpleNamespace(id="document-7", owner="alice", session_id="session-1", + version_count=4, current_content="original", created_at=datetime(2026, 1, 1), + updated_at=datetime(2026, 1, 2)) + monkeypatch.setattr(owned_resources, "_row", lambda namespace, identifier, owner: row + if (namespace, identifier, owner) == ("documents", "document-7", "alice") else None) store = ToolApprovalStore() content = '{"content":"replacement"}'