fix(runtime): restore authorized local control paths

This commit is contained in:
Alexandre Teixeira
2026-10-02 23:15:51 +01:00
parent 6094e2abe6
commit b89d178291
10 changed files with 389 additions and 20 deletions
+7
View File
@@ -523,6 +523,13 @@ def seal_task_authority(prompt, task_type, action, *, owner=None, parent_authori
if operation is not None:
authority = replace(authority, grants=(OperationGrant(operation.tool,
inputs=frozenset({operation.input})),))
if task_type == "action" and action == "cookbook_serve":
# The direct admin scheduling ingress selects the native Cookbook
# producer. Restore never infers this from task names/availability.
# Any model-created task still intersects with its parent's ceiling.
backend = NativeBackendResource("serve_model")
authority = replace(authority, backend_resources=tuple(dict.fromkeys(
(*authority.backend_resources, backend))))
parent = active_request_authority() if parent_authority is MISSING_AUTHORITY else parent_authority
if parent_authority is None:
parent = RequestAuthority.empty(owner=owner)
+105
View File
@@ -0,0 +1,105 @@
"""One-use transport capabilities for admitted local Cookbook producers.
The internal HTTP token authenticates transport only. A capability bridges one
server-owned request/operation/backend to one exact resolved local launch body.
It is never persisted, returned to the model, or usable for shell/job control.
"""
from contextlib import contextmanager
from dataclasses import dataclass
import hashlib
import json
import secrets
import threading
import time
from src.agent_runtime.resources import NativeBackendResource, ResourceIdentityError
CAPABILITY_HEADER = "X-Odysseus-Local-Model-Capability"
_ROUTES = {"download_model": "/api/model/download", "serve_model": "/api/model/serve",
"serve_preset": "/api/model/serve"}
_PENDING = {}
_LOCK = threading.Lock()
def _digest(payload):
return hashlib.sha256(json.dumps(payload, sort_keys=True, separators=(",", ":"),
allow_nan=False).encode()).hexdigest()
@dataclass(frozen=True)
class _Capability:
authority: object
operation: object
backend: NativeBackendResource
path: str
payload_digest: str
deadline: float
@contextmanager
def model_control_headers(tool, content, owner, payload, *, scheduled=False):
from src.tools._common import _internal_headers
headers = _internal_headers(owner)
if payload.get("remote_host"):
yield headers # Remote workload authority/transport is unchanged.
return
from src.agent_runtime.authority import active_request_authority, ExactOperation
from src.agent_runtime.remote_resources import active_backend_operation
from src.tool_security import owner_is_admin_or_single_user
authority = active_request_authority()
operation = ExactOperation.normalize(tool, content)
backend = active_backend_operation()
if (authority is None or authority.owner != str(owner or "").strip().casefold()
or tool not in _ROUTES or not owner_is_admin_or_single_user(owner)):
raise ResourceIdentityError("Local model producer has no matching server authority")
if scheduled:
# Called only by the server-owned scheduled action, after restoration of
# its immutable input ceiling. A task name or owner alone is not enough.
if tool != "serve_model" or not authority.permits(operation):
raise ResourceIdentityError("Scheduled local model input is outside authority")
resource = NativeBackendResource(tool)
if resource not in authority.backend_resources:
raise ResourceIdentityError("Scheduled local model backend is outside authority")
else:
# This binding exists only after dispatch admission (including one-use
# exact approval). A generic tool grant/header cannot create it over HTTP.
if (backend is None or backend.resource != NativeBackendResource(tool)
or (backend.request_id, backend.owner, backend.session_id,
backend.transport_tool, backend.exact_input) !=
(authority.request_id, authority.owner, authority.session_id, tool, operation.input)):
raise ResourceIdentityError("Local model producer operation or backend changed")
resource = backend.resource
capability = _Capability(authority, operation, resource, _ROUTES[tool], _digest(payload), time.monotonic() + 60)
token = secrets.token_urlsafe(32)
headers.update({CAPABILITY_HEADER: token, "X-Odysseus-Owner": authority.owner})
with _LOCK:
_PENDING[token] = capability
try:
yield headers
finally:
with _LOCK:
_PENDING.pop(token, None)
def consume_model_control(request, payload):
"""Claim exactly once at the local route, before any producer effect."""
from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN, INTERNAL_TOOL_USER
from src.auth_helpers import is_direct_loopback_request
token = request.headers.get(CAPABILITY_HEADER)
if not token:
return False
if (not is_direct_loopback_request(request)
or not secrets.compare_digest(request.headers.get(INTERNAL_TOOL_HEADER, ""), INTERNAL_TOOL_TOKEN)):
raise ResourceIdentityError("Local model transport is untrusted")
with _LOCK:
capability = _PENDING.get(token)
if (capability is None or capability.deadline < time.monotonic()
or request.method != "POST" or request.url.path != capability.path
or payload.get("remote_host") or _digest(payload) != capability.payload_digest
or request.headers.get("X-Odysseus-Owner", "") != capability.authority.owner
or getattr(request.state, "current_user", None) not in
(None, INTERNAL_TOOL_USER, capability.authority.owner)):
raise ResourceIdentityError("Local model capability binding changed or expired")
del _PENDING[token]
request.state.local_model_authority = capability.authority
return True
+21
View File
@@ -7,6 +7,27 @@ from fastapi import Request, HTTPException
from src.owner_identity import auth_disabled, effective_storage_owner
def is_direct_loopback_request(request: Request) -> bool:
"""Local operator transport, excluding reverse proxies and cross-site calls.
Locality supplies no model/tool authority. Native administration uses this
only in the operator's explicit auth-disabled single-user mode.
"""
client = getattr(request, "client", None)
if not client or client.host not in {"127.0.0.1", "::1"}:
return False
forwarding = ("cf-connecting-ip", "cf-ray", "cf-visitor", "x-forwarded-for",
"x-forwarded-host", "x-forwarded-proto", "x-real-ip", "forwarded")
if any(request.headers.get(name) for name in forwarding):
return False
if request.headers.get("sec-fetch-site") in {"cross-site", "same-site"}:
return False
origin = request.headers.get("origin")
if origin and origin != str(request.base_url).rstrip("/"):
return False
return True
def get_current_user(request: Request) -> Optional[str]:
"""Get current username from request state (set by auth middleware)."""
return getattr(request.state, 'current_user', None)
+6 -4
View File
@@ -3387,10 +3387,12 @@ async def action_cookbook_serve(
if srv.get("platform"): body["platform"] = srv["platform"]
try:
async with httpx.AsyncClient(timeout=30) as client:
r = await client.post(f"{internal_api_base()}/api/model/serve",
json=body, headers=headers)
data = r.json() if r.content else {}
from src.agent_runtime.local_model_control import model_control_headers
with model_control_headers("serve_model", command, owner, body, scheduled=True) as launch_headers:
async with httpx.AsyncClient(timeout=30) as client:
r = await client.post(f"{internal_api_base()}/api/model/serve",
json=body, headers=launch_headers)
data = r.json() if r.content else {}
except Exception as e:
return f"Launch HTTP failed: {e}", False
if not data.get("ok"):
+15 -9
View File
@@ -772,9 +772,11 @@ async def do_download_model(content: str, owner: Optional[str] = None) -> Dict:
if env_cfg.get("platform"): payload["platform"] = env_cfg["platform"]
if env_cfg.get("ssh_port"): payload["ssh_port"] = env_cfg["ssh_port"]
try:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/download",
json=payload, headers=_internal_headers())
from src.agent_runtime.local_model_control import model_control_headers
with model_control_headers("download_model", content, owner, payload) as launch_headers:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/download",
json=payload, headers=launch_headers)
data = resp.json()
if data.get("ok"):
sid = data.get("session_id", "?")
@@ -857,9 +859,11 @@ async def do_serve_model(content: str, owner: Optional[str] = None) -> Dict:
if env_cfg.get("platform"): payload["platform"] = env_cfg["platform"]
if env_cfg.get("ssh_port"): payload["ssh_port"] = env_cfg["ssh_port"]
try:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
json=payload, headers=_internal_headers())
from src.agent_runtime.local_model_control import model_control_headers
with model_control_headers("serve_model", content, owner, payload) as launch_headers:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
json=payload, headers=launch_headers)
data = resp.json()
if data.get("ok"):
sid = data.get("session_id", "?")
@@ -1908,9 +1912,11 @@ async def do_serve_preset(content: str, owner: Optional[str] = None) -> Dict:
payload["ssh_port"] = env_cfg["ssh_port"]
try:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
json=payload, headers=_internal_headers())
from src.agent_runtime.local_model_control import model_control_headers
with model_control_headers("serve_preset", content, owner, payload) as launch_headers:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(f"{_INTERNAL_BASE}/api/model/serve",
json=payload, headers=launch_headers)
data = resp.json()
if data.get("ok"):
sid = data.get("session_id", "?")