mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-08 07:52:20 +02:00
merge(runtime): reconcile containment with frozen lab
This commit is contained in:
@@ -20375,6 +20375,10 @@ def _blocks_before_inference(turn_contract) -> bool:
|
||||
)
|
||||
|
||||
|
||||
from src.agent_runtime.authority import MISSING_AUTHORITY, active_request_authority, with_request_authority
|
||||
|
||||
|
||||
@with_request_authority
|
||||
@with_turn_contract
|
||||
@with_teacher_takeover
|
||||
@with_completion_gate
|
||||
@@ -20420,6 +20424,7 @@ async def stream_agent_loop(
|
||||
suppress_skills: bool = False,
|
||||
reasoning_effort: Optional[str] = None,
|
||||
_parent_run_id: Optional[str] = None,
|
||||
request_authority=MISSING_AUTHORITY,
|
||||
) -> AsyncGenerator[str, None]:
|
||||
"""Streaming agent loop generator.
|
||||
|
||||
@@ -32758,6 +32763,7 @@ async def stream_agent_loop(
|
||||
block.tool_type, block.content
|
||||
),
|
||||
request_text=_last_user,
|
||||
request_authority=active_request_authority(),
|
||||
)
|
||||
desc = f"{block.tool_type}: APPROVAL REQUIRED"
|
||||
result = {
|
||||
@@ -37372,6 +37378,7 @@ async def stream_agent_loop(
|
||||
active_document=active_document,
|
||||
active_email=active_email,
|
||||
turn_contract=turn_contract,
|
||||
request_authority=active_request_authority(),
|
||||
external_untrusted_context_seen=run_security.external_untrusted_context_seen,
|
||||
client_runtime_context=client_runtime_context,
|
||||
plan_mode=plan_mode,
|
||||
|
||||
@@ -0,0 +1,433 @@
|
||||
"""Server-owned request admission, independent of model tool availability."""
|
||||
from __future__ import annotations
|
||||
|
||||
from contextlib import aclosing, contextmanager
|
||||
from contextvars import ContextVar
|
||||
from dataclasses import dataclass, replace
|
||||
from functools import wraps
|
||||
from inspect import signature
|
||||
import json
|
||||
from pathlib import Path
|
||||
import re
|
||||
from uuid import uuid4
|
||||
|
||||
from src.tool_policy import ToolPolicy, build_effective_tool_policy
|
||||
from src.turn_contract import (
|
||||
FAMILY_TOOLS, canonical_tool, requested_capabilities,
|
||||
RequiredReadOperation, required_read_operation_for_request, selected_tools_for_request,
|
||||
)
|
||||
|
||||
|
||||
def _owner(value):
|
||||
return str(value or "").strip().casefold()
|
||||
|
||||
|
||||
def _pairs(pairs):
|
||||
result = {}
|
||||
for key, value in pairs:
|
||||
if key in result:
|
||||
raise ValueError("Duplicate operation argument")
|
||||
result[key] = value
|
||||
return result
|
||||
|
||||
|
||||
def _invalid_constant(value):
|
||||
raise ValueError("Non-finite operation argument")
|
||||
|
||||
|
||||
def _json(value):
|
||||
return json.dumps(value, sort_keys=True, separators=(",", ":"), allow_nan=False)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ExactOperation:
|
||||
tool: str
|
||||
input: str
|
||||
action: str | None = None
|
||||
transport_tool: str = ""
|
||||
|
||||
@classmethod
|
||||
def normalize(cls, tool, content):
|
||||
if not isinstance(tool, str) or not tool.strip() or not isinstance(content, str):
|
||||
raise ValueError("Operation requires a tool name and string input")
|
||||
transport_tool = tool.strip()
|
||||
tool = canonical_tool(transport_tool)
|
||||
normalized = content
|
||||
payload = None
|
||||
raw_input = tool in {"bash", "python"} or tool.startswith("scheduled__")
|
||||
if not raw_input and content.lstrip().startswith("{"):
|
||||
payload = json.loads(content, object_pairs_hook=_pairs, parse_constant=_invalid_constant)
|
||||
if not isinstance(payload, dict):
|
||||
raise ValueError("Structured tool input must be an object")
|
||||
normalized = _json(payload)
|
||||
# Reuse the runtime's existing multiplexed-action normalization; this
|
||||
# classifies input and never grants permission or changes the input.
|
||||
from src.tool_capabilities import _action_from_content
|
||||
action = _action_from_content(tool, content)
|
||||
if tool == "private_browser" and isinstance(payload, dict):
|
||||
action = payload.get("action")
|
||||
if action is not None and not isinstance(action, str):
|
||||
raise ValueError("Browser action must be a string")
|
||||
action = action.strip().casefold() if action else None
|
||||
return cls(tool, normalized, action, transport_tool)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class OperationGrant:
|
||||
tool: str
|
||||
actions: frozenset[str] | None = None
|
||||
inputs: frozenset[str] | None = None
|
||||
|
||||
def __post_init__(self):
|
||||
if not isinstance(self.tool, str) or not self.tool or canonical_tool(self.tool) != self.tool:
|
||||
raise ValueError("Grant requires a canonical tool identity")
|
||||
for values in (self.actions, self.inputs):
|
||||
if values is not None and (not isinstance(values, frozenset)
|
||||
or any(not isinstance(v, str) for v in values)):
|
||||
raise TypeError("Grant limits must be immutable string sets")
|
||||
|
||||
def permits(self, operation):
|
||||
return (self.tool == operation.tool
|
||||
and (self.actions is None or operation.action in self.actions)
|
||||
and (self.inputs is None or operation.input in self.inputs))
|
||||
|
||||
def intersect(self, other):
|
||||
if self.tool != other.tool:
|
||||
raise ValueError("Cannot intersect different operation classes")
|
||||
def limits(left, right):
|
||||
return right if left is None else left if right is None else left & right
|
||||
return OperationGrant(self.tool, limits(self.actions, other.actions),
|
||||
limits(self.inputs, other.inputs))
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RequestAuthority:
|
||||
request_id: str
|
||||
owner: str
|
||||
session_id: str
|
||||
workspace: str
|
||||
grants: tuple[OperationGrant, ...] = ()
|
||||
denied: frozenset[str] = frozenset()
|
||||
block_all: bool = False
|
||||
disable_mcp: bool = False
|
||||
inherited: bool = False
|
||||
|
||||
def __post_init__(self):
|
||||
if (not isinstance(self.request_id, str) or not self.request_id
|
||||
or any(not isinstance(v, str) for v in (self.owner, self.session_id, self.workspace))
|
||||
or not isinstance(self.grants, tuple)
|
||||
or any(not isinstance(g, OperationGrant) for g in self.grants)
|
||||
or len({g.tool for g in self.grants}) != len(self.grants)
|
||||
or not isinstance(self.denied, frozenset)
|
||||
or any(not isinstance(n, str) or canonical_tool(n) != n for n in self.denied)
|
||||
or any(type(v) is not bool for v in (self.block_all, self.disable_mcp, self.inherited))):
|
||||
raise ValueError("Malformed request authority")
|
||||
|
||||
@classmethod
|
||||
def empty(cls, *, owner=None, session_id=None, workspace=None):
|
||||
return cls(uuid4().hex, _owner(owner), str(session_id or ""), str(workspace or ""))
|
||||
|
||||
def bound_to(self, *, owner=None, session_id=None, workspace=None):
|
||||
return (self.owner == _owner(owner) and self.session_id == str(session_id or "")
|
||||
and self.workspace == str(workspace or ""))
|
||||
|
||||
def restricted(self, operation):
|
||||
return (self.block_all or operation.tool in self.denied
|
||||
or (self.disable_mcp and (operation.tool.startswith("mcp__")
|
||||
or operation.transport_tool.startswith("mcp__"))))
|
||||
|
||||
def permits(self, operation):
|
||||
return not self.restricted(operation) and any(g.permits(operation) for g in self.grants)
|
||||
|
||||
def restrict(self, policy=None, disabled_tools=()):
|
||||
policy = policy or ToolPolicy()
|
||||
return replace(self, denied=self.denied | frozenset(
|
||||
canonical_tool(n) for n in set(disabled_tools or ()) | policy.all_disabled_names()),
|
||||
block_all=self.block_all or policy.block_all_tool_calls,
|
||||
disable_mcp=self.disable_mcp or policy.disable_mcp)
|
||||
|
||||
def intersect(self, child):
|
||||
if not isinstance(child, RequestAuthority):
|
||||
raise TypeError("Child authority must be server-owned RequestAuthority")
|
||||
grants = []
|
||||
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]
|
||||
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)
|
||||
|
||||
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)
|
||||
|
||||
def to_dict(self):
|
||||
return {"version": 1, "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}
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, value):
|
||||
if (not isinstance(value, dict) or type(value.get("version")) is not int
|
||||
or value["version"] != 1):
|
||||
raise ValueError("Unsupported authority snapshot")
|
||||
def limits(value):
|
||||
if value is None:
|
||||
return None
|
||||
if not isinstance(value, list) or any(not isinstance(v, str) for v in value):
|
||||
raise ValueError("Malformed authority limits")
|
||||
return frozenset(value)
|
||||
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"])
|
||||
|
||||
|
||||
_BROWSER_READ_ACTIONS = frozenset({"open", "navigate", "snapshot", "text", "read", "find",
|
||||
"screenshot", "scroll", "back", "forward", "wait", "status", "close", "tabs"})
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SemanticIntent:
|
||||
"""Routing facts, with no execution permission or provider inventory."""
|
||||
capabilities: frozenset[str]
|
||||
selected_tools: frozenset[str] | None
|
||||
required_read: RequiredReadOperation | None
|
||||
|
||||
|
||||
def interpret_request(request_text, *, history=(), workspace=None, active_document=False,
|
||||
image_attachment=False):
|
||||
if not isinstance(request_text, str):
|
||||
raise TypeError("Intent requires request text")
|
||||
# Routing may use model/tool history. Admission may only inherit intent
|
||||
# from trusted user requests; a model's proposal or attempted tool call
|
||||
# cannot establish a new authorized operation class.
|
||||
history = tuple(history or ())
|
||||
trusted_history = []
|
||||
for row in history:
|
||||
get = row.get if isinstance(row, dict) else lambda key, default=None: getattr(row, key, default)
|
||||
metadata = get("metadata") or {}
|
||||
if isinstance(metadata, str):
|
||||
try:
|
||||
metadata = json.loads(metadata)
|
||||
except ValueError:
|
||||
metadata = {}
|
||||
if (get("role") == "user" and isinstance(metadata, dict)
|
||||
and metadata.get("trusted") is not False and not metadata.get("tool_gate_untrusted")):
|
||||
trusted_history.append({"role": "user", "content": get("content", "")})
|
||||
families = requested_capabilities(request_text, trusted_history,
|
||||
active_document=active_document, workspace=bool(workspace), image_attachment=image_attachment)
|
||||
selected = selected_tools_for_request(request_text)
|
||||
if (families <= {"unknown"} and selected is None
|
||||
and re.search(r"\b(?:lan|local\s+(?:network|ip)|tailscale|arp|ip\s+route|default\s+route|subnet|network\s+interface|neighbor\s+table|wifi|ethernet)\b", request_text, re.I)
|
||||
and re.search(r"\b(?:find|check|inspect|show|list|lookup|locate)\b", request_text, re.I)
|
||||
and not re.search(r"\b(?:web|internet|online)\b", request_text, re.I)):
|
||||
# An explicit local-network lookup is a host operation. The existing
|
||||
# router already chooses host_shell; neither its schema nor bridge
|
||||
# availability grants Bash/Python alongside this request.
|
||||
return SemanticIntent(frozenset({"shell_files"}), frozenset({"host_shell"}), None)
|
||||
return SemanticIntent(families, selected, required_read_operation_for_request(request_text, history))
|
||||
|
||||
|
||||
def create_request_authority(request_text, *, owner=None, session_id=None, workspace=None,
|
||||
history=(), policy=None, active_document=False,
|
||||
image_attachment=False, capabilities=None):
|
||||
"""Deterministic server policy over semantic facts, never schema inventory."""
|
||||
if not isinstance(request_text, str):
|
||||
raise TypeError("Authority requires trusted request text")
|
||||
intent = interpret_request(request_text, history=history, active_document=active_document,
|
||||
workspace=workspace, image_attachment=image_attachment)
|
||||
families = intent.capabilities
|
||||
if capabilities is not None:
|
||||
families |= frozenset(capabilities)
|
||||
tools = set().union(*(FAMILY_TOOLS.get(f, ()) for f in families))
|
||||
selected = intent.selected_tools
|
||||
if selected is not None:
|
||||
tools = tools & set(selected) if families else set(selected)
|
||||
if tools & {"web_search", "web_fetch"}:
|
||||
tools.add("private_browser")
|
||||
operation = intent.required_read
|
||||
if operation is not None and canonical_tool(operation.tool) in {canonical_tool(n) for n in tools}:
|
||||
tools = {canonical_tool(operation.tool)}
|
||||
else:
|
||||
operation = None
|
||||
grants = []
|
||||
for name in sorted(tools | {"ask_user", "update_plan"}):
|
||||
name = canonical_tool(name)
|
||||
actions = inputs = None
|
||||
if name == "private_browser":
|
||||
actions = _BROWSER_READ_ACTIONS
|
||||
# Explicit interaction intent admits its operation class. A
|
||||
# browser offered only as static-Web fallback gets no such grant.
|
||||
if re.search(r"\b(?:browser|browse|private_browser)\b", request_text, re.I):
|
||||
actions |= frozenset(action for action in ("click", "fill", "type", "press", "evaluate", "select")
|
||||
if re.search(r"\b" + action + r"\b", request_text, re.I))
|
||||
if operation is not None and name == canonical_tool(operation.tool):
|
||||
inputs = frozenset({ExactOperation.normalize(name, _json(dict(operation.args))).input})
|
||||
grants.append(OperationGrant(name, actions, inputs))
|
||||
authority = RequestAuthority(uuid4().hex, _owner(owner), str(session_id or ""),
|
||||
str(workspace or ""), tuple(grants))
|
||||
return authority.restrict(policy or build_effective_tool_policy(last_user_message=request_text))
|
||||
|
||||
|
||||
_ACTIVE: ContextVar[RequestAuthority | None] = ContextVar("request_authority", default=None)
|
||||
MISSING_AUTHORITY = object()
|
||||
|
||||
|
||||
def active_request_authority():
|
||||
return _ACTIVE.get()
|
||||
|
||||
|
||||
def is_internal_tool_request(request):
|
||||
"""HTTP authentication/owner attribution does not make a tool payload user intent."""
|
||||
from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN
|
||||
return (request.headers.get(INTERNAL_TOOL_HEADER) == INTERNAL_TOOL_TOKEN
|
||||
or getattr(request.state, "current_user", None) == "internal-tool")
|
||||
|
||||
|
||||
def require_user_approval_request(request):
|
||||
if is_internal_tool_request(request):
|
||||
from fastapi import HTTPException
|
||||
raise HTTPException(403, "Tool requests cannot submit user approval decisions.")
|
||||
|
||||
|
||||
def request_authority_for_http(request, request_text, **context):
|
||||
"""Known tool loopback is a continuation, never a fresh user grant source."""
|
||||
if is_internal_tool_request(request):
|
||||
return RequestAuthority.empty(owner=context.get("owner"),
|
||||
session_id=context.get("session_id"), workspace=context.get("workspace")).restrict(context.get("policy"))
|
||||
return create_request_authority(request_text, **context)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def bind_request_authority(authority):
|
||||
if not isinstance(authority, RequestAuthority):
|
||||
raise TypeError("Authority must be server-owned RequestAuthority")
|
||||
parent = _ACTIVE.get()
|
||||
authority = parent.intersect(authority) if parent is not None else authority
|
||||
token = _ACTIVE.set(authority)
|
||||
try:
|
||||
yield authority
|
||||
finally:
|
||||
_ACTIVE.reset(token)
|
||||
|
||||
|
||||
def _request_text(messages):
|
||||
for message in reversed(messages or ()):
|
||||
metadata = message.get("metadata") or {}
|
||||
if (message.get("role") != "user" or metadata.get("trusted") is False
|
||||
or metadata.get("tool_gate_untrusted")):
|
||||
continue
|
||||
content = message.get("content", "")
|
||||
if isinstance(content, str):
|
||||
return content
|
||||
if isinstance(content, list):
|
||||
return "\n".join(p.get("text", "") for p in content
|
||||
if isinstance(p, dict) and p.get("type") == "text")
|
||||
return ""
|
||||
|
||||
|
||||
def with_request_authority(func):
|
||||
"""Bind once per invocation; model rounds/fallbacks never recreate grants."""
|
||||
call_signature = signature(func)
|
||||
@wraps(func)
|
||||
async def wrapped(*args, **kwargs):
|
||||
bound = call_signature.bind(*args, **kwargs)
|
||||
bound.apply_defaults()
|
||||
parameters = bound.arguments
|
||||
parent = active_request_authority()
|
||||
authority = parameters.get("request_authority", MISSING_AUTHORITY)
|
||||
if authority is MISSING_AUTHORITY:
|
||||
approval = parameters.get("exact_approval")
|
||||
if parent is not None:
|
||||
authority = parent
|
||||
elif approval is not None:
|
||||
authority = approval.pending.request_authority or RequestAuthority.empty(
|
||||
owner=parameters.get("owner"), session_id=parameters.get("session_id"),
|
||||
workspace=parameters.get("workspace"))
|
||||
elif (parameters.get("_parent_run_id") or parameters.get("_is_teacher_run")
|
||||
or parameters.get("workload") == "background"):
|
||||
authority = RequestAuthority.empty(owner=parameters.get("owner"),
|
||||
session_id=parameters.get("session_id"), workspace=parameters.get("workspace"))
|
||||
else:
|
||||
authority = create_request_authority(_request_text(parameters.get("messages")),
|
||||
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")))
|
||||
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:
|
||||
authority = replace(authority, inherited=False)
|
||||
authority = authority.restrict(parameters.get("tool_policy"), parameters.get("disabled_tools"))
|
||||
with bind_request_authority(authority) as effective:
|
||||
if "request_authority" in parameters:
|
||||
parameters["request_authority"] = effective
|
||||
async with aclosing(func(*bound.args, **bound.kwargs)) as stream:
|
||||
async for chunk in stream:
|
||||
yield chunk
|
||||
return wrapped
|
||||
|
||||
|
||||
def task_operation(task_type, action, prompt):
|
||||
if task_type == "action":
|
||||
tool = ("bash" if action in {"run_local", "run_script", "ssh_command"}
|
||||
else "serve_model" if action == "cookbook_serve" else "scheduled__" + str(action))
|
||||
return ExactOperation.normalize(tool, str(prompt or ""))
|
||||
if task_type == "research":
|
||||
return ExactOperation.normalize("trigger_research", str(prompt or ""))
|
||||
return None
|
||||
|
||||
|
||||
def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authority=MISSING_AUTHORITY):
|
||||
"""Only direct ingress grants; a model-created task is capped by its parent."""
|
||||
operation = task_operation(task_type, action, prompt)
|
||||
authority = create_request_authority(str(prompt or ""), owner=owner)
|
||||
if operation is not None:
|
||||
authority = replace(authority, grants=(OperationGrant(operation.tool,
|
||||
inputs=frozenset({operation.input})),))
|
||||
parent = active_request_authority() if parent_authority is MISSING_AUTHORITY else parent_authority
|
||||
if parent_authority is None:
|
||||
parent = RequestAuthority.empty(owner=owner)
|
||||
if parent is not None:
|
||||
authority = parent.intersect(replace(authority, session_id=parent.session_id,
|
||||
workspace=parent.workspace))
|
||||
return _json({"task_input": [prompt, task_type, action], "authority": authority.to_dict()})
|
||||
|
||||
|
||||
def restore_task_authority(snapshot, prompt, task_type, action, *, owner=None, session_id=None):
|
||||
try:
|
||||
value = json.loads(snapshot)
|
||||
if value["task_input"] != [prompt, task_type, action]:
|
||||
raise ValueError("Scheduled request changed")
|
||||
return RequestAuthority.from_dict(value["authority"]).continuation(owner=owner, session_id=session_id)
|
||||
except (ValueError, TypeError, KeyError, AttributeError):
|
||||
return RequestAuthority.empty(owner=owner, session_id=session_id)
|
||||
|
||||
|
||||
def _background_path(job_id):
|
||||
if not isinstance(job_id, str) or not re.fullmatch(r"[A-Za-z0-9_-]+", job_id):
|
||||
raise ValueError("Invalid background authority identity")
|
||||
from src.constants import BG_JOBS_DIR
|
||||
return Path(BG_JOBS_DIR) / (job_id + ".authority.json")
|
||||
|
||||
|
||||
def save_background_authority(job_id, authority):
|
||||
from core.atomic_io import atomic_write_json
|
||||
atomic_write_json(_background_path(job_id), authority.to_dict())
|
||||
|
||||
|
||||
def restore_background_authority(job_id, *, owner=None, session_id=None):
|
||||
try:
|
||||
authority = RequestAuthority.from_dict(json.loads(_background_path(job_id).read_text()))
|
||||
if authority.session_id != str(session_id or ""):
|
||||
raise ValueError("Background session changed")
|
||||
return authority.continuation(owner=owner, session_id=session_id)
|
||||
except (OSError, ValueError, TypeError, KeyError, AttributeError):
|
||||
return RequestAuthority.empty(owner=owner, session_id=session_id)
|
||||
@@ -121,6 +121,10 @@ async def _create_bash_subprocess(
|
||||
bash,
|
||||
"-c",
|
||||
str(command or ""),
|
||||
stdin=asyncio.subprocess.DEVNULL,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
env=env,
|
||||
cwd=cwd,
|
||||
)
|
||||
kwargs = {"cwd": cwd} if cwd is not None else {}
|
||||
|
||||
+11
-2
@@ -36,11 +36,12 @@ def _background_result_message(rec):
|
||||
return untrusted_context_message("background job output", inject)
|
||||
|
||||
|
||||
async def _drain_agent(sess, messages):
|
||||
async def _drain_agent(sess, messages, request_authority=None):
|
||||
"""Run the agent loop headless against a session. Returns
|
||||
(final_prose, tool_events) — tool_events in the same shape the live chat
|
||||
saves, so the frontend rebuilds them as standard agent-thread tool cards."""
|
||||
from src.agent_loop import stream_agent_loop
|
||||
from src.agent_runtime.authority import RequestAuthority
|
||||
full = ""
|
||||
final_replaced = False
|
||||
tool_events = []
|
||||
@@ -52,6 +53,9 @@ async def _drain_agent(sess, messages):
|
||||
session_id=sess.id,
|
||||
max_rounds=_FOLLOWUP_MAX_ROUNDS,
|
||||
owner=getattr(sess, "owner", None),
|
||||
workspace=request_authority.workspace or None if request_authority is not None else None,
|
||||
request_authority=(request_authority or RequestAuthority.empty(
|
||||
owner=getattr(sess, "owner", None), session_id=sess.id)),
|
||||
):
|
||||
if not chunk.startswith("data: "):
|
||||
continue
|
||||
@@ -132,7 +136,12 @@ async def _run_followup(rec: dict) -> bool:
|
||||
context = sess.get_context_messages()
|
||||
context.append(_background_result_message(rec))
|
||||
|
||||
full, tool_events = await _drain_agent(sess, context)
|
||||
from src.agent_runtime.authority import restore_background_authority
|
||||
from src.settings import get_setting
|
||||
authority = restore_background_authority(
|
||||
rec["id"], owner=getattr(sess, "owner", None), session_id=sess.id)
|
||||
authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ())
|
||||
full, tool_events = await _drain_agent(sess, context, request_authority=authority)
|
||||
|
||||
# Persist ONLY the assistant continuation so it renders as a normal agent
|
||||
# turn — a standard chat bubble plus `tool_events` that the frontend
|
||||
|
||||
@@ -4984,6 +4984,10 @@ async def preview_lines_until_finish(response, finish_event=None):
|
||||
await asyncio.gather(finish_task, return_exceptions=True)
|
||||
|
||||
|
||||
from src.agent_runtime.authority import MISSING_AUTHORITY, with_request_authority
|
||||
|
||||
|
||||
@with_request_authority
|
||||
async def stream_preview(*, endpoint_url, model, messages, headers, turn_contract,
|
||||
session_id, owner, disabled_tools, tool_policy,
|
||||
history_session=None, external_untrusted_context_seen=False,
|
||||
@@ -4991,7 +4995,7 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
|
||||
client_runtime_context=None, max_tokens=768, max_rounds=8,
|
||||
max_tool_calls=0,
|
||||
external_tool_schemas=None, temperature=0.0,
|
||||
**ignored):
|
||||
request_authority=MISSING_AUTHORITY, **ignored):
|
||||
from src.generation_sampling import validate_temperature
|
||||
temperature = validate_temperature(temperature)
|
||||
# This path sends requests directly with httpx and therefore bypasses
|
||||
|
||||
+32
-2
@@ -1316,6 +1316,17 @@ class TaskScheduler:
|
||||
async def _execute_action(self, task, run_id: str | None = None) -> tuple:
|
||||
"""Execute a built-in action (no LLM needed)."""
|
||||
from src.builtin_actions import BUILTIN_ACTIONS
|
||||
from src.agent_runtime.authority import (
|
||||
bind_request_authority, restore_task_authority, task_operation,
|
||||
)
|
||||
authority = restore_task_authority(
|
||||
getattr(task, "request_authority_json", None), task.prompt, task.task_type,
|
||||
task.action, owner=task.owner)
|
||||
from src.settings import get_setting
|
||||
authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ())
|
||||
operation = task_operation(task.task_type, task.action, task.prompt)
|
||||
if operation is None or not authority.permits(operation):
|
||||
return "Scheduled action has no matching server request authority.", False
|
||||
|
||||
action_fn = BUILTIN_ACTIONS.get(task.action)
|
||||
if not action_fn:
|
||||
@@ -1342,7 +1353,8 @@ class TaskScheduler:
|
||||
if getattr(task, "model", None):
|
||||
kwargs["model"] = task.model
|
||||
kwargs["endpoint_url"] = getattr(task, "endpoint_url", None)
|
||||
result, success = await action_fn(**kwargs)
|
||||
with bind_request_authority(authority):
|
||||
result, success = await action_fn(**kwargs)
|
||||
if getattr(task, "model", None):
|
||||
self._last_run_model = task.model
|
||||
return result, success
|
||||
@@ -1945,6 +1957,7 @@ class TaskScheduler:
|
||||
datetime_context_msg: dict | None = None) -> str:
|
||||
"""Run the full agent loop with tool access, collecting the final text."""
|
||||
from src.agent_loop import stream_agent_loop
|
||||
from src.agent_runtime.authority import restore_task_authority
|
||||
|
||||
system_content = system_prompt or "You are a helpful assistant executing a scheduled task. Use available tools to complete the task thoroughly."
|
||||
user_content = override_user_message or task.prompt
|
||||
@@ -2000,6 +2013,10 @@ class TaskScheduler:
|
||||
_task_fallbacks = []
|
||||
# Close the stream in this task on every exit, including the
|
||||
# approval-pause break, so the agent run's context state unwinds here.
|
||||
request_authority = restore_task_authority(
|
||||
getattr(task, "request_authority_json", None), task.prompt,
|
||||
getattr(task, "task_type", "llm"), getattr(task, "action", None),
|
||||
owner=task.owner, session_id=session_id)
|
||||
async with contextlib.aclosing(stream_agent_loop(
|
||||
endpoint_url=endpoint_url,
|
||||
model=model,
|
||||
@@ -2007,11 +2024,13 @@ class TaskScheduler:
|
||||
max_rounds=_task_max_rounds,
|
||||
session_id=session_id,
|
||||
owner=task.owner,
|
||||
workspace=request_authority.workspace or None,
|
||||
headers=headers,
|
||||
disabled_tools=disabled_tools,
|
||||
relevant_tools=relevant_tools,
|
||||
fallbacks=_task_fallbacks,
|
||||
workload="background",
|
||||
request_authority=request_authority,
|
||||
)) as agent_stream:
|
||||
async for event_str in agent_stream:
|
||||
if event_str.startswith("data: ") and not event_str.startswith("data: [DONE]"):
|
||||
@@ -2109,6 +2128,14 @@ class TaskScheduler:
|
||||
|
||||
async def _execute_research_task(self, task, db) -> str:
|
||||
"""Execute a deep research task using DeepResearcher."""
|
||||
from src.agent_runtime.authority import bind_request_authority, restore_task_authority, task_operation
|
||||
from src.settings import get_setting
|
||||
authority = restore_task_authority(
|
||||
getattr(task, "request_authority_json", None), task.prompt, task.task_type,
|
||||
getattr(task, "action", None), owner=task.owner)
|
||||
authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ())
|
||||
if not authority.permits(task_operation(task.task_type, getattr(task, "action", None), task.prompt)):
|
||||
raise PermissionError("Scheduled research has no matching server request authority.")
|
||||
from core.database import Session as DbSession, ChatMessage
|
||||
from src.deep_research import DeepResearcher
|
||||
from src.research_handler import RESEARCH_DATA_DIR, ResearchHandler
|
||||
@@ -2180,7 +2207,8 @@ class TaskScheduler:
|
||||
)
|
||||
|
||||
started_ts = time.time()
|
||||
report = await researcher.research(task.prompt)
|
||||
with bind_request_authority(authority):
|
||||
report = await researcher.research(task.prompt)
|
||||
completed_ts = time.time()
|
||||
try:
|
||||
stats = researcher.get_stats() or {}
|
||||
@@ -2603,6 +2631,7 @@ class TaskScheduler:
|
||||
if (task.output_target or "session") == "session":
|
||||
task.output_target = defs.get("output_target", "none")
|
||||
seeded = []
|
||||
from src.agent_runtime.authority import seal_task_authority
|
||||
for action, defs in HOUSEKEEPING_DEFAULTS.items():
|
||||
if action in existing_actions:
|
||||
continue
|
||||
@@ -2620,6 +2649,7 @@ class TaskScheduler:
|
||||
name=defs["name"],
|
||||
task_type="action",
|
||||
action=action,
|
||||
request_authority_json=seal_task_authority(None, "action", action, owner=owner),
|
||||
trigger_type=trigger_type,
|
||||
trigger_event=defs.get("trigger_event"),
|
||||
trigger_count=defs.get("trigger_count"),
|
||||
|
||||
@@ -585,6 +585,7 @@ async def run_teacher_inline(
|
||||
external_untrusted_context_seen: bool = False,
|
||||
client_runtime_context: Optional[Dict[str, Any]] = None,
|
||||
plan_mode: bool = False,
|
||||
request_authority=None,
|
||||
):
|
||||
"""Async generator. Yields SSE event strings.
|
||||
|
||||
@@ -700,6 +701,7 @@ async def run_teacher_inline(
|
||||
active_email=active_email,
|
||||
turn_contract=turn_contract,
|
||||
_parent_run_id=parent_run_id,
|
||||
request_authority=request_authority,
|
||||
external_untrusted_context_seen=external_untrusted_context_seen,
|
||||
client_runtime_context=deepcopy(client_runtime_context),
|
||||
plan_mode=plan_mode,
|
||||
@@ -822,6 +824,7 @@ async def run_teacher_inline(
|
||||
workspace=workspace,
|
||||
external_untrusted_context_seen=True,
|
||||
capabilities=capabilities_for_action("manage_skills", skill_content),
|
||||
request_authority=request_authority,
|
||||
)
|
||||
approval = pending.public_payload(
|
||||
reason=(
|
||||
|
||||
@@ -25,6 +25,7 @@ from src.tool_approval_scopes import (
|
||||
scope_for_decision,
|
||||
)
|
||||
from src.tool_capabilities import ToolCapabilities, capabilities_for_action
|
||||
from src.agent_runtime.authority import RequestAuthority
|
||||
|
||||
|
||||
DEFAULT_APPROVAL_TTL_SECONDS = 10 * 60
|
||||
@@ -117,6 +118,7 @@ def _binding_payload(
|
||||
continuation_query: Any,
|
||||
effects: tuple[str, ...],
|
||||
result_integrity: str,
|
||||
request_authority: RequestAuthority | None = None,
|
||||
) -> dict[str, Any]:
|
||||
return {
|
||||
"owner": _normalized_owner(owner),
|
||||
@@ -137,6 +139,7 @@ def _binding_payload(
|
||||
"continuation_query": _normalized_continuation_query(continuation_query),
|
||||
"effects": list(effects),
|
||||
"result_integrity": str(result_integrity),
|
||||
"request_authority": request_authority.to_dict() if request_authority is not None else None,
|
||||
}
|
||||
|
||||
|
||||
@@ -165,6 +168,7 @@ class PendingToolApproval:
|
||||
# The originating user request is internal continuation context only; it
|
||||
# is never displayed or treated as authorization for the sealed action.
|
||||
request_text: str = ""
|
||||
request_authority: RequestAuthority | None = None
|
||||
|
||||
def public_payload(self, *, reason: str | None = None) -> dict[str, Any]:
|
||||
return {
|
||||
@@ -273,6 +277,7 @@ class ExactToolApproval:
|
||||
continuation_query=self.pending.continuation_query,
|
||||
effects=effects,
|
||||
result_integrity=result_integrity,
|
||||
request_authority=self.pending.request_authority,
|
||||
)
|
||||
return _canonical_digest(expected) == self.pending.digest
|
||||
|
||||
@@ -356,7 +361,10 @@ class ToolApprovalStore:
|
||||
external_untrusted_context_seen: bool,
|
||||
capabilities: ToolCapabilities,
|
||||
request_text: Any = "",
|
||||
request_authority: RequestAuthority | None = None,
|
||||
) -> PendingToolApproval:
|
||||
if request_authority is not None and not isinstance(request_authority, RequestAuthority):
|
||||
raise TypeError("Approval authority must be server-owned RequestAuthority")
|
||||
now = time.time()
|
||||
effects = tuple(sorted(effect.value for effect in capabilities.effects))
|
||||
result_integrity = capabilities.result_integrity.value
|
||||
@@ -375,6 +383,7 @@ class ToolApprovalStore:
|
||||
continuation_query=continuation_query,
|
||||
effects=effects,
|
||||
result_integrity=result_integrity,
|
||||
request_authority=request_authority,
|
||||
)
|
||||
pending = PendingToolApproval(
|
||||
approval_id=secrets.token_urlsafe(32),
|
||||
@@ -398,6 +407,7 @@ class ToolApprovalStore:
|
||||
selected_tools=tuple(payload["selected_tools"]),
|
||||
continuation_query=payload["continuation_query"],
|
||||
request_text=str(request_text or ""),
|
||||
request_authority=request_authority,
|
||||
)
|
||||
with self._lock:
|
||||
self._purge_expired_locked(now)
|
||||
|
||||
+69
-25
@@ -1214,6 +1214,7 @@ async def _direct_fallback(
|
||||
"client_runtime_context": client_runtime_context,
|
||||
"disabled_tools": frozenset(disabled_tools or ()),
|
||||
"tool_policy": tool_policy,
|
||||
"request_authority": active_request_authority(),
|
||||
}
|
||||
|
||||
from src.agent_tools import TOOL_HANDLERS
|
||||
@@ -1254,6 +1255,10 @@ async def _document_tool_dispatch(
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
from src.agent_runtime.journal import dispatched, mark_authorized, mark_dispatch, record_action
|
||||
from src.agent_runtime.authority import (
|
||||
MISSING_AUTHORITY, ExactOperation, RequestAuthority, active_request_authority,
|
||||
bind_request_authority, save_background_authority,
|
||||
)
|
||||
|
||||
|
||||
@record_action
|
||||
@@ -1273,6 +1278,7 @@ async def execute_tool_block(
|
||||
exact_approval: Optional[ExactToolApproval] = None,
|
||||
active_document_id: Optional[str] = None,
|
||||
client_runtime_context: Optional[Dict[str, Any]] = None,
|
||||
request_authority=MISSING_AUTHORITY,
|
||||
) -> Tuple[str, Dict]:
|
||||
"""Execute a single tool block. Returns (description, result_dict).
|
||||
|
||||
@@ -1294,6 +1300,40 @@ async def execute_tool_block(
|
||||
"NO_TOOL_SECURITY_CONTEXT"
|
||||
)
|
||||
|
||||
authority = active_request_authority() if request_authority is MISSING_AUTHORITY else request_authority
|
||||
parent = active_request_authority()
|
||||
if isinstance(authority, RequestAuthority) and parent is not None and authority is not parent:
|
||||
authority = parent.intersect(authority)
|
||||
try:
|
||||
operation = ExactOperation.normalize(getattr(block, "tool_type", None), getattr(block, "content", None))
|
||||
valid = isinstance(authority, RequestAuthority) and authority.bound_to(
|
||||
owner=owner, session_id=session_id, workspace=workspace)
|
||||
if valid:
|
||||
authority = authority.restrict(tool_policy, disabled_tools)
|
||||
exact_admission = bool(
|
||||
valid and not authority.inherited and not authority.restricted(operation)
|
||||
and exact_approval is not None and exact_approval.matches(
|
||||
owner=owner, session_id=session_id, workspace=workspace,
|
||||
tool_name=getattr(block, "tool_type", None), content=getattr(block, "content", None)))
|
||||
admitted = valid and (authority.permits(operation) or exact_admission)
|
||||
except (ValueError, TypeError, AttributeError) as error:
|
||||
return f"{getattr(block, 'tool_type', '')}: invalid arguments", {
|
||||
"error": (f"Tool arguments are not valid JSON: {error}"
|
||||
if isinstance(error, json.JSONDecodeError) else str(error)),
|
||||
"exit_code": 1, "blocked": True,
|
||||
"failure_kind": "request_authority_denied",
|
||||
}
|
||||
if not admitted:
|
||||
reason = "The exact operation is outside server request authority."
|
||||
if tool_policy and any(tool_policy.blocks(name) for name in email_tool_policy_names(getattr(block, "tool_type", ""))):
|
||||
reason = f"Execution of tool '{getattr(block, 'tool_type', '')}' is forbade by the active tool policy."
|
||||
elif isinstance(authority, RequestAuthority) and authority.restricted(operation):
|
||||
reason = "The exact operation is disabled by user or server request authority policy."
|
||||
return f"{getattr(block, 'tool_type', '')}: BLOCKED", {
|
||||
"error": reason,
|
||||
"exit_code": 1, "blocked": True, "failure_kind": "request_authority_denied",
|
||||
}
|
||||
|
||||
from src.turn_contract import active_turn_contract
|
||||
contract = active_turn_contract()
|
||||
if contract is not None and not contract.permits(getattr(block, "tool_type", "")):
|
||||
@@ -1388,31 +1428,32 @@ async def execute_tool_block(
|
||||
|
||||
token = _active_workspace.set(workspace or None)
|
||||
try:
|
||||
output = await _execute_tool_block_impl(
|
||||
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
|
||||
),
|
||||
active_document_id=active_document_id,
|
||||
client_runtime_context=client_runtime_context,
|
||||
)
|
||||
with bind_request_authority(authority):
|
||||
output = await _execute_tool_block_impl(
|
||||
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
|
||||
),
|
||||
active_document_id=active_document_id,
|
||||
client_runtime_context=client_runtime_context,
|
||||
)
|
||||
if isinstance(security_context, ToolRunSecurityContext):
|
||||
security_context.observe_tool_result(
|
||||
getattr(block, "tool_type", None),
|
||||
@@ -1623,6 +1664,9 @@ async def _execute_tool_block_impl(
|
||||
from src import bg_jobs
|
||||
mark_dispatch()
|
||||
rec = bg_jobs.launch(_bg_cmd, session_id=session_id, cwd=agent_cwd())
|
||||
# Only this server launch may seal detached-job authority; a
|
||||
# handler/bridge output carrying a job id is not a grant source.
|
||||
save_background_authority(rec["id"], active_request_authority())
|
||||
short = _bg_cmd.strip().split(chr(10))[0][:80]
|
||||
desc = f"bash (background): {short}"
|
||||
result = {
|
||||
|
||||
@@ -486,12 +486,15 @@ async def do_manage_tasks(content: str, owner: Optional[str] = None) -> Dict:
|
||||
# Guard each fallback with `or`: args.get("prompt", default) returns
|
||||
# None when the key is present but null, and None[:50] raises.
|
||||
name = args.get("name") or (args.get("prompt") or args.get("action_name") or "Task")[:50]
|
||||
from src.agent_runtime.authority import seal_task_authority
|
||||
|
||||
task = ScheduledTask(
|
||||
id=task_id,
|
||||
owner=owner,
|
||||
name=name,
|
||||
prompt=args.get("prompt"),
|
||||
request_authority_json=seal_task_authority(
|
||||
args.get("prompt"), task_type, args.get("action_name"), owner=owner),
|
||||
task_type=task_type,
|
||||
action=args.get("action_name"),
|
||||
schedule=args.get("schedule", "daily") if trigger_type == "schedule" else None,
|
||||
@@ -557,6 +560,10 @@ async def do_manage_tasks(content: str, owner: Optional[str] = None) -> Dict:
|
||||
if args.get("action_name") is not None:
|
||||
task.action = args["action_name"]
|
||||
changed.append("action")
|
||||
if any(args.get(field) is not None for field in ("prompt", "task_type", "action_name")):
|
||||
from src.agent_runtime.authority import seal_task_authority
|
||||
task.request_authority_json = seal_task_authority(
|
||||
task.prompt, task.task_type, task.action, owner=owner)
|
||||
if args.get("trigger_type") is not None:
|
||||
task.trigger_type = args["trigger_type"]
|
||||
changed.append("trigger_type")
|
||||
|
||||
Reference in New Issue
Block a user