mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-08 16:02:20 +02:00
feat(runtime): bind remote and owned resources to authority
This commit is contained in:
@@ -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 = {
|
||||
|
||||
@@ -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()})
|
||||
|
||||
|
||||
|
||||
@@ -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")
|
||||
@@ -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)
|
||||
+166
-12
@@ -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}))
|
||||
|
||||
@@ -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)]
|
||||
|
||||
@@ -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'])
|
||||
|
||||
@@ -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"}
|
||||
|
||||
@@ -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}
|
||||
|
||||
|
||||
+49
-1
@@ -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:
|
||||
|
||||
+45
-2
@@ -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)
|
||||
|
||||
+74
-31
@@ -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
|
||||
|
||||
+7
-3
@@ -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")
|
||||
|
||||
+8
-1
@@ -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}
|
||||
|
||||
+5
-1
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user