diff --git a/core/database.py b/core/database.py index 6addc95c4..1621492c7 100644 --- a/core/database.py +++ b/core/database.py @@ -777,6 +777,7 @@ class ScheduledTask(TimestampMixin, Base): owner = Column(String, nullable=True, index=True) name = Column(String, nullable=False, default="Untitled Task") prompt = Column(Text, nullable=True) # LLM prompt (for task_type="llm") + request_authority_json = Column(Text, nullable=True) # server-only admitted request snapshot task_type = Column(String, default="llm") # "llm" | "action" action = Column(String, nullable=True) # builtin action name (for task_type="action") schedule = Column(String, nullable=True) # "once", "daily", "weekly", "monthly" @@ -2335,6 +2336,15 @@ def _migrate_seed_email_account(): # Any future migrations or schema changes that temporarily violate foreign-key # constraints will fail. To perform such operations, foreign_keys must be # temporarily disabled around the migration workflow. +def _migrate_add_task_authority_column(): + """Retain snapshots after legacy task-table rebuilds; support all DBs.""" + from sqlalchemy import inspect + with engine.begin() as conn: + columns = {column["name"] for column in inspect(conn).get_columns("scheduled_tasks")} + if "request_authority_json" not in columns: + conn.execute(text("ALTER TABLE scheduled_tasks ADD COLUMN request_authority_json TEXT")) + + def init_db(): """ Initialize the database by creating all tables. @@ -2412,6 +2422,7 @@ def init_db(): _migrate_add_oauth_config() _migrate_add_email_oauth_columns() _migrate_add_task_automation_columns() + _migrate_add_task_authority_column() _migrate_add_disabled_tools() _migrate_add_mcp_oauth_tokens_column() _migrate_add_task_v2_columns() diff --git a/docs/runtime-decomposition/wave-2-request-authority.md b/docs/runtime-decomposition/wave-2-request-authority.md new file mode 100644 index 000000000..44a72c350 --- /dev/null +++ b/docs/runtime-decomposition/wave-2-request-authority.md @@ -0,0 +1,248 @@ +# Wave 2: request authority + +Base: `d6c3c98c75e03f70c05ebe4058c6fa12e0395f62`, branch +`feature/runtime-request-authority`. Discovery and this plan precede production +changes. No later runtime waves are included. + +## Discovered call paths + +`routes/chat_routes.py` parses mode, toggles, workspace, approval decisions and +runtime context. User intent can promote Chat to Agent. Owner privileges, +global disabled tools, compare/incognito and plan restrictions produce +`ToolPolicy`. Compact/native routes resolve `TurnContract`; regular/full models +can receive the full enabled schema inventory. The route calls +`_stream_agent_with_execution_bridge` and `stream_agent_loop`. Detached runs +retain this generator; reconnecting subscribes to it rather than creating a new +invocation. Their stream IDs are distinct from journal IDs. + +`src/turn_contract.py` classifies request families and selected tools, resolves +exact safe reads, and filters schema availability. Empty-family routing has a +legacy core inventory. Warm tools and editor availability may enlarge offers. +Transcription, OCR and tasks have narrow selection; static web retrieval may +offer private_browser for fallback. These routing choices are not grants. + +`src/agent_loop.py` selects provider/profile transports, parses native or textual +tool blocks, repairs calls, performs deterministic preflights and retries, and +calls `src/tool_execution.py:execute_tool_block`. Compact preview uses +`src/clean_agent_preview.py` but reaches the same dispatcher. The dispatcher +checks run security, exact approval, contract membership, disabled tools, +ToolPolicy, owner restrictions and bridges before MCP/dynamic/built-in handlers. +It forwards policy to dynamic handlers. Legacy loop reconciliation removes +disabled names found in a contract's offered inventory. This must not erase a +request-authority denial. + +Approvals use `src/tool_approvals.py`. A server record binds tool/content, owner, +session, workspace, document id/version/digest, origin run and continuation +state. Consume is destructive; claim is one-use. Task/chat scopes bypass an +existing run-security gate; they do not define the requested operation classes. +Approval continuation executes the sealed action in round zero. Denial exits +the route without execution. + +Generic app_api forwards both the internal token and the caller's owner to +loopback HTTP. Its blocklist does not exclude Chat/skill approval ingress. +Matching owner/session/input bindings alone therefore cannot distinguish a +model-produced HTTP decision from a user approval. Those existing ingress +points need an explicit internal-tool rejection before consuming approval. +Internal HTTP skill-test task bodies likewise cannot mint fresh authority. +The same origin rule applies to generic Chat HTTP entry: a loopback generated +message is not a new trusted user request, even with correct owner attribution. +Both Chat entry points use the existing non-persistence switch for these +messages and append explicitly untrusted transient context instead. Later +referential turns cannot inherit their operation class as prior user intent. + +Teacher takeover is queued by the student, then owned by the outer adapter in +`src/teacher_escalation.py`. It invokes a child loop after the student gate closes +and forwards policy, contract, workspace and runtime context. The teacher's +synthetic user message is model context, not a new authority source. + +`src/task_scheduler.py:_execute_assistant` composes crew/global restrictions and +RAG/default shell availability. `_run_agent_loop` supplies task.prompt or a +synthetic override as a user message, with background provider fallback. Exact +approval pauses are retired because there is no interactive approver. +`_execute_action` invokes BUILTIN_ACTIONS directly, with a separate admin gate. +`src/tools/system.py:do_manage_tasks` and `routes/task/task_routes.py` create/edit +persisted tasks. No authority snapshot currently survives scheduling. + +Detached Bash dispatch launches `bg_jobs.launch` and returns bg_job_id. +`src/bg_monitor.py:_run_followup` appends an explicitly untrusted result to session +context and re-enters the loop. It currently forwards neither the originating +authority nor its request restrictions. Skill tests/audits in +`routes/skills_routes.py` also invoke the loop with task/user messages; generated +audit context must not manufacture grants. + +| Question | Current source | +| --- | --- | +| Requested operation | User intent classifiers, exact safe-read resolver; ultimately parsed/repaired model tool block | +| Available capabilities | Registry/MCP inventory, profiles, RAG, TurnContract and request-specific schema filters | +| Authorized capabilities | Fragmented policy, privileges, run security and approval checks; no independent envelope | +| Restrictions | Route toggles, owner/global policy, plan/compare/incognito, dispatcher owner/workspace checks | +| Approval required | Deterministic run-security decision; model output can propose the action but cannot consume approval | +| Approval input scope | Server-sealed exact tool/content and owner/session/workspace/document binding | +| Nested state | Explicit policy/contract/workspace/context forwarding and journal lineage; no authority snapshot | +| Model influence | Tool/input proposals, repairs, recovery choices, generated task/audit prompts; availability currently participates in execution gating | + +## Implementation plan and contract + +1. Add immutable `ExactOperation`, `OperationGrant` and `RequestAuthority` in + `src/agent_runtime/authority.py`. Normalize canonical tool identity and JSON + inputs (reject duplicate keys/non-finite values); retain exact raw text for + Bash/Python, built-in scheduled actions and non-JSON inputs. Grants contain an operation class/tool identity, optional + action limits and exact input limits. Authority has its own request id, + owner/session/workspace binding, immutable grants and hard denials. It is + independent of schema presence, model/profile, stream/journal/receipt IDs. +2. Create authority from trusted request text/history and deterministic policy + at the chat route before availability reconciliation. The general loop + boundary creates it for other trusted direct callers, without consulting + schemas, relevant_tools, forced_tools or model output. Authority family + inheritance reads only trusted user history. Tool-history exact reads may + narrow an already admitted class, never create a class. Unknown intent grants + no execution floor. Neutral interaction/planning controls remain explicit. +3. Keep semantic classification and availability in TurnContract. Resolve + authority grants separately from those semantic facts and hard policy. + Exact safe reads restrict action/identifiers. Static web fallback authorizes + browser reading/navigation, not arbitrary click/evaluate/form operations. + Media/task families do not inherit the shell inventory. +4. Bind authority around the whole logical stream, including teacher takeover; + forward it explicitly to teacher children and approval records. Children + inherit the parent or intersect explicit authority with it. Policy denials + union; grants intersect; a child cannot replace the parent scope. Restore + the parent on close/error/cancellation. Capture restrictions before legacy + offered-tool reconciliation can erase them. +5. Enforce at `execute_tool_block`, before approvals are claimed or handlers, + bridges/MCP/process dispatch begin. Current policy/disabled gates still win. + Missing/malformed dispatcher state fails closed. Standalone callers/tests + must supply explicit server authority. Journal ownership remains unchanged; + denied calls produce no authoritative execution receipt. +6. Existing approvals remain one-use exact claims. Seal the originating + authority in the approval digest. Resumption keeps original class limits and + current hard restrictions. The approved exact operation may cross its + original class boundary only through the consumed, matching server record + at that call; it does not mutate authority for subsequent calls. Nested + execution cannot use an approval to exceed its parent ceiling. Existing + task/chat UI and run-security scope semantics are unchanged. + Chat/skill approval ingress rejects validated internal-tool requests before + consumption; identity impersonation is not a user approval decision. + A shared HTTP factory admits trusted user requests and produces an empty, + policy-restricted envelope for known internal-tool Chat/skill requests. +7. Persist a server-only authority snapshot and task-input binding on scheduled + records. Direct authenticated task ingress can admit its user-supplied task; + task creation inside model execution intersects with parent authority. + Scheduler overrides, retries and provider fallbacks reuse that snapshot. + Missing/stale snapshots grant no tool authority. Newly seeded server-owned + housekeeping jobs receive exact snapshots at their static creation point; + existing rows are not retrospectively authorized by their names/actions. + Internal tool HTTP task payloads cannot become fresh user requests across an + ASGI context boundary. Built-in actions receive + an exact admission check. Persist detached-job authority in a separate + authority sidecar at dispatch; monitor continuations reuse it and current + denials. Do not edit bg_jobs/process containment implementation. +8. Production files: new authority module; routes/chat_routes.py; + src/agent_loop.py; src/tool_execution.py; src/teacher_escalation.py; + src/tool_approvals.py; core/database.py; routes/task/task_routes.py; + src/tools/system.py; src/task_scheduler.py; src/bg_monitor.py; + routes/skills_routes.py. Change preview only if direct-entry binding is + required by validation. No TurnContract/profile/schema redesign. +9. Shared hotspots: route/loop/dispatch, approvals and task/database integration. + One coordinator writes all production files. Keep changes confined to + authority creation, forwarding, persistence and admission. Do not modify + containment, provenance/effect classification or egress implementation. +10. Focused regressions: available schema/bridge/dynamic handler without grants; + model-selected unrelated tool/action; explicit class admission; exact read + arguments; narrow transcription/OCR/tasks/browser fallback; hard denials + despite offered-tool reconciliation; retry/fallback stability; child and + teacher non-widening and restoration; malformed/missing state; exact + approval mismatch/replay and continuation scope; scheduled snapshot/input + binding and synthetic override; detached followup inheritance; journal + denial evidence. Preserve existing policy-forwarding and Ajax assertions. + +Validation: new focused tests; existing contract/policy/capability/profile +tests; scripts/validate_runtime_wave1.sh; broad affected runtime tests; full +pytest; compileall; JS/MJS syntax; diff check and conflict-marker scan. Any +production edit after full pytest requires affected tests and full pytest again. + +## Implemented boundaries and remaining limits + +The preview entry also binds authority because it supports direct callers. +Research task admission binds the snapshot around the researcher, so nested +execution cannot infer grants from generated research context. LAN lookup +intent has a narrow host_shell-only admission rule; it adds neither Bash nor +Python and does not alter Ajax schemas or profiles. + +Scheduled loop entry explicitly forwards the restored workspace as well as +the envelope; rebinding the continuation session never drops confinement to +the original workspace. Only the actual server Bash launch seals a detached +job sidecar. A handler/bridge result claiming a job id cannot create one. + +Snapshots are trusted server state, stored in the task database and detached +job authority sidecars. Missing, malformed, changed-input, wrong-owner or +wrong-session snapshots fail closed. Legacy tasks need a trusted task-input +save to obtain a snapshot; legacy detached jobs have no execution grants on +followup. No broad backfill, authority-mode UI, containment, effect/egress or +receipt/journal redesign is included. Sidecars follow the detached job's server +storage trust assumptions; retention/integrity hardening is outside this slice. + +Class admission deliberately reuses the deterministic semantic classifiers. +Unrecognized intent has only explicit ask_user/update_plan controls. This can +deny unsupported phrasing and generated default skill tests/audits; model +prompts and tool inventory cannot repair that denial. Existing exact approvals +can admit one sealed root operation, never widen subsequent calls or nested +authority. They still require the existing armed security context, matching +bindings, one-use claim, document checks and current hard restrictions. + +Standalone dispatcher test fixtures now supply explicit registry grants to +continue exercising their original handler/policy/confinement assertions. +New authority tests use the raw dispatcher and prove denial before dispatch. + +## File ownership and reasons + +| Production file | Wave 2 change | +| --- | --- | +| src/agent_runtime/authority.py | Immutable intent/admission/operation API, trusted factory, intersection/context binding, task/job snapshots | +| routes/chat_routes.py | Capture authority before availability reconciliation; pass it into execution; guard approval ingress | +| src/agent_loop.py | Bind logical-invocation authority; capture it in approvals and teacher takeover | +| src/tool_execution.py | Normalize/check operations before dispatch and approval claims; bind handler context; seal actual detached launch | +| src/teacher_escalation.py | Explicit child/approval inheritance without synthetic-prompt grants | +| src/tool_approvals.py | Bind immutable originating authority into exact approval digest | +| src/clean_agent_preview.py | Bind authority at the supported direct preview entry | +| core/database.py | Add nullable server-only scheduled snapshot column and additive migration | +| routes/task/task_routes.py | Seal direct user task inputs; deny fresh grants to internal-tool HTTP payloads | +| src/tools/system.py | Cap model-created/edited task snapshots by active authority | +| src/task_scheduler.py | Restore original scope/workspace for loops, admit exact built-ins/research, seal new static defaults | +| src/bg_monitor.py | Restore original detached-job scope and current hard restrictions | +| routes/skills_routes.py | Separate explicit user task authority from generated/internal skill prompts; guard approval ingress | + +Shared hotspots touched: chat routes, agent loop, central dispatcher, preview, +teacher escalation, approvals, task CRUD/scheduler/system handlers, database, +background monitor and skill entry routes. All production edits have one writer. +TurnContract, tool schemas, model profiles, journal/completion foundations, +bg_jobs/process containment and effect/egress implementations are untouched. + +`tests/test_request_authority.py` adds the focused authority regressions. +`tests/runtime_evidence_helpers.py` adds explicit standalone server fixture +grants. Original assertions are preserved in these adapted fixture suites: + +- tests/test_agent_external_tool_schemas.py +- tests/test_ask_user_tool.py +- tests/test_client_tool_routing.py +- tests/test_edit_file.py +- tests/test_execution_bridge.py +- tests/test_external_context_tool_gate.py +- tests/test_image_creation_routing.py +- tests/test_review_regressions.py +- tests/test_runtime_evidence_contract.py +- tests/test_task_cookbook_admin_gate.py +- tests/test_task_scheduler_cancel.py +- tests/test_tool_approvals.py +- tests/test_tool_path_confinement.py +- tests/test_tool_policy.py +- tests/test_turn_contract.py +- tests/test_turn_contract_integration.py +- tests/test_update_plan_tool.py +- tests/test_weather_search_recovery.py +- tests/test_workspace_confine.py + +`website/configuration-reference.md` is regenerated solely to update the +chat-route environment-read line number. This document records discovery, +the pre-edit plan, implementation boundaries and file ownership. The validation +report records final commands/results. No production files in parallel lanes +are claimed. diff --git a/docs/runtime-decomposition/wave-2-validation-full.md b/docs/runtime-decomposition/wave-2-validation-full.md new file mode 100644 index 000000000..579b12c5a --- /dev/null +++ b/docs/runtime-decomposition/wave-2-validation-full.md @@ -0,0 +1,80 @@ +# Wave 2 final validation + +Worktree: `odysseus-runtime-request-authority`; branch: +`feature/runtime-request-authority`. +Starting SHA: `d6c3c98c75e03f70c05ebe4058c6fa12e0395f62`. +The final SHA is the local commit containing this report, returned in the final +implementation report. No rebase, merge, push or PR was performed. + +All results below apply to the final production code. The last production +changes addressed internal HTTP request/approval origin and transient untrusted +Chat context. Focused, Wave 1.1, broad runtime and full pytest were rerun after +those changes. Subsequent edits only recorded results and removed temporary +validation logs. + +| Gate | Final result | +| --- | --- | +| New Wave 2 authority tests | 58 passed, 1 warning; 1.23s | +| Relevant contract/policy/approval/capability/Ajax/task/background tests | 1500 passed, 28 skipped, 1 warning; 30.55s | +| Wave 1.1 validation script | 2292 passed, 1 warning; 65.75s | +| Broad affected runtime suite | 3079 passed, 28 skipped, 1 warning; 92.83s | +| Full pytest | 11644 passed, 54 skipped, 2 xfailed, 182 warnings, 6 subtests passed; 444.40s | +| Python compileall | Passed | +| JS/MJS syntax | Passed for all 361 tracked files | +| Git whitespace gate | Passed | +| Conflict-marker scan | Passed | + +The existing release smoke hook skipped because `APP_PORT` was unset; no live +instance was driven. Full pytest includes its existing skips and expected +failures. Warnings are retained in the local raw log. Missing development test +dependencies and Playwright Chromium were installed locally, without changing +project dependency declarations. No global dotenv-disable override was used. + +## Commands + +```sh +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py tests/test_turn_contract*.py tests/test_tool_policy.py tests/test_tool_approval*.py tests/test_execution_capabilities.py tests/test_ajax*.py tests/test_task_*.py tests/test_bg_*.py + +ODYSSEUS_TEST_PYTHON="$PWD/.venv/bin/python" bash scripts/validate_runtime_wave1.sh + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q tests/test_request_authority.py tests/test_agent_*.py tests/test_turn_contract*.py tests/test_tool_policy.py tests/test_tool_approval*.py tests/test_task_*.py tests/test_bg_*.py tests/test_*completion*.py tests/test_foreground_model_routing.py tests/test_client_tool_routing.py tests/test_workspace_confine.py tests/test_product_turn_contract_route.py tests/test_execution_bridge.py tests/test_execution_capabilities.py tests/test_ajax*.py tests/test_external_context_tool_gate.py tests/test_tool_path_confinement.py tests/test_edit_file.py tests/test_runtime_evidence_contract.py tests/test_review_regressions.py tests/test_image_creation_routing.py tests/test_ask_user_tool.py tests/test_update_plan_tool.py tests/test_weather_search_recovery.py tests/test_clean_agent_preview.py tests/test_skill_audit*.py tests/test_preview_execution_evidence.py + +ODYSSEUS_TEST_STATIC_PORT=0 .venv/bin/python -m pytest -q + +.venv/bin/python -m compileall -q -x '(^|/)(\.venv|\.git|node_modules|data|logs|uploads)/' . +git ls-files -z '*.js' '*.mjs' | xargs -0 -n 1 node --check +git diff --check +# Staged whitespace check used --cached --check with all 37 changed paths explicit. +git grep --cached -l -E '^(<<<<<<< |=======$|>>>>>>> )' -- '*.py' '*.js' '*.mjs' '*.html' '*.css' '*.json' '*.md' '*.sh' +``` + +Conflict-marker grep returns exit 1 with no matches on success. +The context firewall rejected the unbounded staged whitespace command before +execution; the exact-path check passed. No admitted source inspection was +blocked by staging. +Local raw validation outputs are archived under the ignored +`.venv/wave2-validation/` directory; they are not committed. + +## Regression scope and limits + +The 58 authority tests cover schema/handler/model-selection non-authority, +narrow media/tasks/browser behavior, exact reads, deterministic grants, hard +denials, malformed/missing state, retry and nested inheritance, teacher +forwarding, exact approval scope/replay/digest, scheduled input sealing and +workspace restoration, detached followups and actual-launch-only sealing, +internal HTTP origin, untrusted Chat persistence, and denied-call journal +completion evidence. Existing fixture assertions remain intact; standalone +dispatch fixtures now provide explicit server authority. + +Remaining limits: class admission uses deterministic request classifiers and +can reject unsupported phrasing; legacy task/job snapshots fail closed until +trusted resealing; snapshots assume trusted server database/job storage; +sidecar retention hardening is deferred. Existing approvals can admit one exact +root operation without granting subsequent or nested operations. + +No Wave 3, 3-S, 4, 5 or 6 work was started. No containment, effect/egress, +provenance, authority-mode UI, journal or completion-foundation redesign is +included. File ownership and the discovery/implementation contract are recorded +in [wave-2-request-authority.md](wave-2-request-authority.md). diff --git a/routes/chat_routes.py b/routes/chat_routes.py index 967e36b16..09596ccbc 100644 --- a/routes/chat_routes.py +++ b/routes/chat_routes.py @@ -86,6 +86,7 @@ from src.model_profiles import ( tool_schema_profile, ) from src.tool_execution import AgentExecutionBridge, bind_execution_bridge +from src.agent_runtime.authority import is_internal_tool_request, request_authority_for_http from src.turn_contract import ( FAMILY_TOOLS, bind_turn_contract, preserve_bound_editor_selected_tools, requested_capabilities, resolve_turn_contract, @@ -95,6 +96,15 @@ from src.turn_contract import ( logger = logging.getLogger(__name__) + +def _append_internal_chat_context(ctx, message): + tagged = untrusted_context_message("internal tool request", message) + ctx.messages.append(tagged) + routed = getattr(ctx, "route_messages", None) + if routed is not None and routed is not ctx.messages: + routed.append(tagged) + + # Track active streams for partial-save safety net _active_streams: Dict[str, dict] = {} @@ -2196,7 +2206,10 @@ def setup_chat_routes( webhook_manager=webhook_manager, allow_tool_preprocessing=allow_tool_preprocessing, defer_context_shaping=foreground_policy.enabled, + persist_user_message=not is_internal_tool_request(request), ) + if is_internal_tool_request(request): + _append_internal_chat_context(ctx, message) # Research injection research_blocked_by_policy = ( @@ -2648,6 +2661,8 @@ def setup_chat_routes( ) owner = effective_user(request) if tool_approval_id: + from src.agent_runtime.authority import require_user_approval_request + require_user_approval_request(request) pending_tool_approval = tool_approval_store.peek(tool_approval_id) normalized_owner = str(owner or "").strip().casefold() if ( @@ -2907,10 +2922,12 @@ def setup_chat_routes( and pending_tool_approval.continuation_query else None ), - persist_user_message=not tool_approval_continuation, + persist_user_message=not tool_approval_continuation and not is_internal_tool_request(request), interaction_mode=chat_mode, auto_escalated=auto_escalated, ) + if is_internal_tool_request(request): + _append_internal_chat_context(ctx, message) _research_flags = {"do": do_research} # Mutable container for generator scope @@ -3349,6 +3366,24 @@ def setup_chat_routes( "manage_documents", "create_document", "edit_document", "update_document", }.issubset(disabled_tools), } + # Capture permission state before schema selection/reconciliation. + # Only deterministic request intent supplies grants, never inventory. + _request_authority = request_authority_for_http( + request, message, owner=_user, session_id=session, workspace=workspace, + history=_turn_history, policy=tool_policy, + active_document=bool(active_doc), + image_attachment=any(str(a.get('mime') or '').startswith('image/') + for a in (ctx.preprocessed.attachment_meta or [])), + capabilities=({'search_browser'} if ( + _explicit_browser_intent or _external_discovery_intent + ) else ()), + ) + if exact_tool_approval is not None: + from src.agent_runtime.authority import RequestAuthority + _request_authority = ( + exact_tool_approval.pending.request_authority + or RequestAuthority.empty(owner=_user, session_id=session, workspace=workspace) + ).restrict(tool_policy) _turn_contract = None # Image models execute directly, not through the text-agent inventory. # Keep the permission policy above, but do not apply routing omissions @@ -3477,6 +3512,7 @@ def setup_chat_routes( warm_tools=_warm_tools, message=message, history=getattr(sess, "history", []) or [], ) + _request_authority = _request_authority.restrict(_contract_policy) # Resolution already applies user, owner, and global policy. An # admitted tool must not later be rejected by the stale # pre-contract disabled snapshot during execution. @@ -4415,6 +4451,7 @@ def setup_chat_routes( cwd=_agent_turn_cwd(sess, client_runtime_context), forced_tools=_forced_tools, turn_contract=_turn_contract, + request_authority=_request_authority, uploaded_files=ctx.uploaded_files, defer_context_shaping=_foreground_policy.enabled, external_untrusted_context_seen=external_untrusted_context_seen, diff --git a/routes/skills_routes.py b/routes/skills_routes.py index ef5f65047..ec5a49876 100644 --- a/routes/skills_routes.py +++ b/routes/skills_routes.py @@ -534,11 +534,13 @@ async def _run_skill_test_job( messages=None, transcript=None, exact_approval=None, + request_authority=None, ): """Background coroutine: run the skill in an agent loop, capture a condensed log + transcript, then have the judge grade it. Writes into _skill_test_jobs.""" import json as _json from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import RequestAuthority job = _skill_test_jobs.get(key) if job is None: @@ -559,6 +561,9 @@ async def _run_skill_test_job( url, model, messages, headers=headers, temperature=0.3, max_tokens=0, max_rounds=8, owner=owner, exact_approval=exact_approval, + request_authority=(request_authority or ( + exact_approval.pending.request_authority if exact_approval is not None else None + ) or RequestAuthority.empty(owner=owner)), ): if not chunk.startswith("data: ") or chunk.strip() == "data: [DONE]": continue @@ -1039,6 +1044,7 @@ async def _run_skill_audit_arm(messages: list[dict], url, model, headers, owner, """Run one audit arm in the agent loop; return transcript, stats, approval.""" import json as _json from src.agent_loop import stream_agent_loop + from src.agent_runtime.authority import RequestAuthority, active_request_authority transcript = [] approval_required = None stats = {"turns": 0, "tool_calls": 0} @@ -1051,6 +1057,7 @@ async def _run_skill_audit_arm(messages: list[dict], url, model, headers, owner, url, model, messages, headers=headers, temperature=0.3, max_tokens=4096, max_rounds=8, owner=owner, workload=workload, suppress_skills=True, + request_authority=(active_request_authority() or RequestAuthority.empty(owner=owner)), ): # Streams can include an SSE event line before the data line, # notably `event: error`. Do not silently discard those failures. @@ -2001,6 +2008,9 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: user = _owner(request) body = await request.json() task = (body.get("task") or "").strip() + from src.agent_runtime.authority import RequestAuthority, request_authority_for_http + request_authority = (request_authority_for_http(request, task, owner=user) + if task else RequestAuthority.empty(owner=user)) skills = skills_manager.load(owner=user) match = next((s for s in skills if s.get("name") == skill_id or s.get("id") == skill_id), None) @@ -2064,9 +2074,12 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: "model": model, "headers": headers, "owner": user, + "request_authority": request_authority, }, } - _asyncio.create_task(_run_skill_test_job(key, name, md, task, url, model, headers, user, skills_manager)) + _asyncio.create_task(_run_skill_test_job( + key, name, md, task, url, model, headers, user, skills_manager, + request_authority=request_authority)) return {"ok": True, "status": "running", "skill": name, "model": model} @router.post("/{skill_id}/test-approval") @@ -2074,6 +2087,8 @@ def setup_skills_routes(skills_manager: SkillsManager) -> APIRouter: """Resume a manual skill test with one exact server-sealed action.""" import asyncio as _asyncio from src.tool_approvals import tool_approval_store + from src.agent_runtime.authority import require_user_approval_request + require_user_approval_request(request) user = _owner(request) skills = skills_manager.load(owner=user) diff --git a/routes/task/task_routes.py b/routes/task/task_routes.py index c19e73ac9..0749cb548 100644 --- a/routes/task/task_routes.py +++ b/routes/task/task_routes.py @@ -25,6 +25,14 @@ from routes.prefs_routes import _load_for_user, _save_for_user logger = logging.getLogger(__name__) +def _seal_request_task_authority(request, prompt, task_type, action, owner): + from src.agent_runtime.authority import MISSING_AUTHORITY, is_internal_tool_request, seal_task_authority + # A tool HTTP call starts another ASGI context. Its model-produced body is + # not a fresh user request, even though the internal token authenticates it. + parent = None if is_internal_tool_request(request) else MISSING_AUTHORITY + return seal_task_authority(prompt, task_type, action, owner=owner, parent_authority=parent) + + def _maybe_cascade_calendar_event(task) -> None: """Delete the linked calendar event when a cookbook_serve task is removed. Two lookup strategies: @@ -530,6 +538,8 @@ def setup_task_routes(task_scheduler) -> APIRouter: owner=user, name=name, prompt=req.prompt, + request_authority_json=_seal_request_task_authority( + request, req.prompt, req.task_type, req.action, user), task_type=req.task_type, action=req.action, schedule=req.schedule, @@ -737,6 +747,9 @@ def setup_task_routes(task_scheduler) -> APIRouter: task.task_type = req.task_type if req.action is not None: task.action = req.action + if any(value is not None for value in (req.prompt, req.task_type, req.action)): + task.request_authority_json = _seal_request_task_authority( + request, task.prompt, task.task_type, task.action, user) if req.output_target is not None: task.output_target = req.output_target if req.model is not None: diff --git a/src/agent_loop.py b/src/agent_loop.py index 988b454f0..5de35898a 100644 --- a/src/agent_loop.py +++ b/src/agent_loop.py @@ -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, diff --git a/src/agent_runtime/authority.py b/src/agent_runtime/authority.py new file mode 100644 index 000000000..ef8552903 --- /dev/null +++ b/src/agent_runtime/authority.py @@ -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) diff --git a/src/bg_monitor.py b/src/bg_monitor.py index 2c17c3a1b..086faae19 100644 --- a/src/bg_monitor.py +++ b/src/bg_monitor.py @@ -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 diff --git a/src/clean_agent_preview.py b/src/clean_agent_preview.py index 83a448db1..4c0302eb7 100644 --- a/src/clean_agent_preview.py +++ b/src/clean_agent_preview.py @@ -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 diff --git a/src/task_scheduler.py b/src/task_scheduler.py index ce0103f48..02fa9970f 100644 --- a/src/task_scheduler.py +++ b/src/task_scheduler.py @@ -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"), diff --git a/src/teacher_escalation.py b/src/teacher_escalation.py index 1f646b583..fda3f4cec 100644 --- a/src/teacher_escalation.py +++ b/src/teacher_escalation.py @@ -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=( diff --git a/src/tool_approvals.py b/src/tool_approvals.py index 5e688172b..7144fc3cf 100644 --- a/src/tool_approvals.py +++ b/src/tool_approvals.py @@ -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) diff --git a/src/tool_execution.py b/src/tool_execution.py index e72594cba..b9cd65c61 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -1219,6 +1219,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 @@ -1259,6 +1260,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 @@ -1278,6 +1283,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). @@ -1299,6 +1305,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", "")): @@ -1393,31 +1433,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), @@ -1628,6 +1669,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 = { diff --git a/src/tools/system.py b/src/tools/system.py index 76783badc..d60d8f585 100644 --- a/src/tools/system.py +++ b/src/tools/system.py @@ -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") diff --git a/tests/runtime_evidence_helpers.py b/tests/runtime_evidence_helpers.py index f5aa483f3..1ab245a15 100644 --- a/tests/runtime_evidence_helpers.py +++ b/tests/runtime_evidence_helpers.py @@ -2,6 +2,33 @@ from src.agent_runtime.journal import mark_dispatch, record_action +def server_authorized_executor(executor): + """Give standalone dispatcher fixtures their explicit server grants. + + These existing suites exercise handlers, policy, confinement and approvals. + Their fixture grants cover the declared native tool registry, independently + of the proposed block. Request-authority denial tests use the raw dispatcher. + """ + from functools import wraps + from inspect import signature + from src.agent_runtime.authority import OperationGrant, RequestAuthority + from src.tool_policy import known_tool_names + from src.turn_contract import canonical_tool + call_signature = signature(executor) + @wraps(executor) + async def execute(*args, **kwargs): + bound = call_signature.bind(*args, **kwargs) + parameters = bound.arguments + kwargs.setdefault("request_authority", RequestAuthority( + "standalone-test-request", str(parameters.get("owner") or "").strip().casefold(), + str(parameters.get("session_id") or ""), str(parameters.get("workspace") or ""), + tuple(OperationGrant(name) for name in sorted( + {canonical_tool(n) for n in known_tool_names()} | {"list_dir", "find_files"})), + )) + return await executor(*args, **kwargs) + return execute + + def authoritative_executor(function): @record_action async def execute(block, *args, **kwargs): diff --git a/tests/test_agent_external_tool_schemas.py b/tests/test_agent_external_tool_schemas.py index f6aa8c46e..f6688a870 100644 --- a/tests/test_agent_external_tool_schemas.py +++ b/tests/test_agent_external_tool_schemas.py @@ -333,6 +333,7 @@ def test_external_tool_images_are_threaded_as_multimodal_evidence(): def test_declared_external_call_reaches_scoped_bridge(monkeypatch): + from src.agent_runtime.authority import OperationGrant, RequestAuthority bridge_calls = [] round_no = 0 monkeypatch.setattr(agent_loop, "get_setting", lambda key, default=None: default) @@ -367,6 +368,8 @@ def test_declared_external_call_reaches_scoped_bridge(monkeypatch): [{"role": "user", "content": "Perform the declared operation."}], max_rounds=2, owner="pewds", + request_authority=RequestAuthority("declared-fixture", "pewds", "", "", + (OperationGrant("inspect_state"),)), relevant_tools={"inspect_state"}, forced_tools={"inspect_state"}, fallbacks=[], @@ -389,6 +392,7 @@ def test_declared_external_call_reaches_scoped_bridge(monkeypatch): def test_known_native_tool_reaches_scoped_bridge_without_redeclared_schema(monkeypatch): + from src.agent_runtime.authority import create_request_authority bridge_calls = [] round_no = 0 monkeypatch.setattr(agent_loop, "get_setting", lambda key, default=None: default) @@ -435,6 +439,7 @@ def test_known_native_tool_reaches_scoped_bridge_without_redeclared_schema(monke [{"role": "user", "content": "Search email for Project Alpha."}], max_rounds=2, owner="public-user", + request_authority=create_request_authority("Search email for Project Alpha.", owner="public-user"), relevant_tools={"search_emails"}, forced_tools={"search_emails"}, fallbacks=[], diff --git a/tests/test_ask_user_tool.py b/tests/test_ask_user_tool.py index 4094701bb..facaf2d01 100644 --- a/tests/test_ask_user_tool.py +++ b/tests/test_ask_user_tool.py @@ -11,6 +11,9 @@ from src.agent_tools import ToolBlock, TOOL_TAGS # noqa: E402 (import first to from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block from src.tool_index import ALWAYS_AVAILABLE, BUILTIN_TOOL_DESCRIPTIONS from src.tool_security import is_public_blocked_tool +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _run(content): diff --git a/tests/test_client_tool_routing.py b/tests/test_client_tool_routing.py index 83fe271d1..fe050d7b7 100644 --- a/tests/test_client_tool_routing.py +++ b/tests/test_client_tool_routing.py @@ -14,6 +14,7 @@ from routes.chat_routes import _agent_turn_cwd import routes.chat_routes as chat_routes from src import tool_execution as _te from src.agent_loop import _is_explicit_local_network_request +from tests.runtime_evidence_helpers import server_authorized_executor # Hold module-object references (not just from-imported names): other test # modules re-import src.tool_execution via sys.modules pops, so string-target @@ -28,7 +29,7 @@ _ROUTED_BRIDGE_TOOLS = _te._ROUTED_BRIDGE_TOOLS async def _execute_tool_block_for_unit_tests(*args, **kwargs): """Use the explicit non-security-context test mode for dispatch tests.""" kwargs.setdefault("security_context", _te.NO_TOOL_SECURITY_CONTEXT) - return await _te.execute_tool_block(*args, **kwargs) + return await server_authorized_executor(_te.execute_tool_block)(*args, **kwargs) execute_tool_block = _execute_tool_block_for_unit_tests diff --git a/tests/test_edit_file.py b/tests/test_edit_file.py index b3fb202e1..6f94a3961 100644 --- a/tests/test_edit_file.py +++ b/tests/test_edit_file.py @@ -4,6 +4,14 @@ import os import tempfile import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src import tool_security from src.tool_security import ( diff --git a/tests/test_execution_bridge.py b/tests/test_execution_bridge.py index 9c72e0569..bcfc33b96 100644 --- a/tests/test_execution_bridge.py +++ b/tests/test_execution_bridge.py @@ -1,6 +1,7 @@ import asyncio import logging import src.tool_execution as tool_execution +from tests.runtime_evidence_helpers import server_authorized_executor from src.tool_execution import ( AgentExecutionBridge, @@ -11,6 +12,9 @@ from src.tool_execution import ( ) +execute_tool_block = server_authorized_executor(execute_tool_block) + + class Block: tool_type = "host_shell" diff --git a/tests/test_external_context_tool_gate.py b/tests/test_external_context_tool_gate.py index cc8d7f4f0..473cba897 100644 --- a/tests/test_external_context_tool_gate.py +++ b/tests/test_external_context_tool_gate.py @@ -7,6 +7,14 @@ from collections import namedtuple from pathlib import Path import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from tests.helpers.document_source import document_source from tests.helpers.js_modules import email_library_paths diff --git a/tests/test_image_creation_routing.py b/tests/test_image_creation_routing.py index 9a2189775..a571e960b 100644 --- a/tests/test_image_creation_routing.py +++ b/tests/test_image_creation_routing.py @@ -125,6 +125,8 @@ async def test_generation_dispatch_uses_owner_aware_backend(monkeypatch, exit_co from types import SimpleNamespace from src import ai_interaction, tool_execution from src.agent_runtime.journal import ActionJournal, bind_journal + from tests.runtime_evidence_helpers import server_authorized_executor + execute = server_authorized_executor(tool_execution.execute_tool_block) calls = [] async def generate(content, **kwargs): calls.append((content, kwargs)) @@ -140,9 +142,9 @@ async def test_generation_dispatch_uses_owner_aware_backend(monkeypatch, exit_co block = SimpleNamespace(tool_type='generate_image', content='{"prompt":"A city"}') journal = ActionJournal() with bind_journal(journal): - _, denied = await tool_execution.execute_tool_block(block, owner='pewds', session_id='fixture', + _, denied = await execute(block, owner='pewds', session_id='fixture', disabled_tools={'generate_image'}, security_context=tool_execution.NO_TOOL_SECURITY_CONTEXT) - _, result = await tool_execution.execute_tool_block(block, owner='pewds', session_id='fixture', + _, result = await execute(block, owner='pewds', session_id='fixture', security_context=tool_execution.NO_TOOL_SECURITY_CONTEXT) assert denied['exit_code'] != 0 assert journal.actions[0].execution_id is None diff --git a/tests/test_request_authority.py b/tests/test_request_authority.py new file mode 100644 index 000000000..2f25449ca --- /dev/null +++ b/tests/test_request_authority.py @@ -0,0 +1,584 @@ +"""Request grants are independent of tool offerings and model proposals.""" +import asyncio +from contextlib import nullcontext +from dataclasses import replace +import json +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock + +import pytest + +from src.agent_runtime.authority import ( + MISSING_AUTHORITY, ExactOperation, OperationGrant, RequestAuthority, + active_request_authority, bind_request_authority, create_request_authority, + restore_background_authority, restore_task_authority, save_background_authority, + seal_task_authority, task_operation, with_request_authority, + is_internal_tool_request, require_user_approval_request, + request_authority_for_http, +) +from src.tool_policy import ToolPolicy +from src.tool_types import ToolBlock + + +def authority(*tools, owner="alice", session_id="s", workspace=""): + return RequestAuthority("request-test", owner, session_id, workspace, + tuple(OperationGrant(tool) for tool in tools)) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("offering", ["schema", "bridge", "dynamic"]) +async def test_availability_and_model_selection_do_not_grant_execution(monkeypatch, offering): + from src import tool_execution as execution + from src.turn_contract import bind_turn_contract, resolve_turn_contract + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + offered_handler = AsyncMock() + if offering == "dynamic": + import src.agent_tools + monkeypatch.setitem(src.agent_tools.TOOL_HANDLERS, "bash", offered_handler) + bridge_context = execution.bind_execution_bridge(execution.AgentExecutionBridge( + route_tool=offered_handler, supported_tools=frozenset({"bash"}))) if offering == "bridge" else nullcontext() + contract = resolve_turn_contract(capabilities={"shell_files"}, policy=ToolPolicy(), + schemas=[{"function": {"name": "bash"}}]) + with bind_turn_contract(contract), bridge_context: + description, result = await execution.execute_tool_block( + ToolBlock("bash", "echo 'authorized by model'"), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert "BLOCKED" in description + assert result["failure_kind"] == "request_authority_denied" + implementation.assert_not_awaited() + offered_handler.assert_not_awaited() + + +@pytest.mark.parametrize("user_text,allowed,denied", [ + ("Transcribe /workspace/input/audio.wav", "transcribe_media", "bash"), + ("OCR extract exact text from /workspace/input/image.png", "extract_text", "python"), + ("List my tasks", "manage_tasks", "web_fetch"), + ("Search the web for current weather", "web_search", "bash"), +]) +def test_explicit_request_classes_remain_narrow(user_text, allowed, denied): + grant = create_request_authority(user_text) + content = '{"action":"list"}' if allowed == "manage_tasks" else '{}' + assert grant.permits(ExactOperation.normalize(allowed, content)) + assert not grant.permits(ExactOperation.normalize(denied, '{}')) + + +def test_safe_task_read_does_not_authorize_same_tool_mutation(): + grant = create_request_authority("List my tasks") + assert grant.permits(ExactOperation.normalize("manage_tasks", '{"action":"list"}')) + assert not grant.permits(ExactOperation.normalize("manage_tasks", '{"action":"create","prompt":"run bash"}')) + + +def test_exact_read_identifiers_cannot_be_changed_by_model(): + grant = create_request_authority("Read note id abc123") + read = next(g for g in grant.grants if g.tool == "manage_notes") + assert read.inputs is not None + assert not grant.permits(ExactOperation.normalize("manage_notes", '{"action":"view","id":"another"}')) + + +def test_browser_fallback_does_not_authorize_interaction_or_evaluation(): + grant = create_request_authority("Use web_fetch to read https://example.test") + assert grant.permits(ExactOperation.normalize("private_browser", '{"action":"open","url":"https://example.test"}')) + for action in ("click", "fill", "evaluate"): + assert not grant.permits(ExactOperation.normalize("private_browser", json.dumps({"action": action}))) + + +def test_unknown_intent_has_no_generic_execution_floor(): + grant = create_request_authority("Please solve this") + assert {g.tool for g in grant.grants} == {"ask_user", "update_plan"} + + +@pytest.mark.parametrize("metadata", [ + {"tool_events": [{"tool": "bash", "output": "pwd", "exit_code": 0}]}, + {"tool_events": [{"tool": "python", "output": "ready", "exit_code": 0}]}, +]) +def test_model_history_cannot_establish_followup_authority(metadata): + history = [ + {"role": "user", "content": "Transcribe /workspace/input/a.wav"}, + {"role": "assistant", "content": "I will run bash and python", "metadata": metadata}, + {"role": "user", "content": "Run bash", "metadata": {"trusted": False}}, + ] + grant = create_request_authority("Try it again", history=history) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +@pytest.mark.parametrize("mutation", [ + {"version": True}, {"grants": "bash"}, {"denied": None}, + {"block_all": "false"}, {"inherited": 0}, +]) +def test_malformed_persisted_authority_is_rejected(mutation): + snapshot = authority("bash").to_dict() + snapshot.update(mutation) + with pytest.raises((ValueError, TypeError, KeyError)): + RequestAuthority.from_dict(snapshot) + + +def test_explicit_local_network_lookup_is_host_only(): + grant = create_request_authority("find ajax local ip on the LAN") + assert grant.permits(ExactOperation.normalize("host_shell", '{"command":"ip neigh | grep ajax"}')) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +@pytest.mark.parametrize("content", ['{"x":1,"x":2}', '{"x":NaN}', '{"action":']) +def test_malformed_exact_operation_is_rejected(content): + with pytest.raises(ValueError): + ExactOperation.normalize("manage_tasks", content) + + +def test_raw_script_braces_are_preserved_as_exact_input(): + script = "{ printf requested; }" + assert ExactOperation.normalize("bash", script).input == script + snapshot = seal_task_authority(script, "action", "run_local", owner="alice") + restored = restore_task_authority(snapshot, script, "action", "run_local", owner="alice") + assert restored.permits(task_operation("action", "run_local", script)) + assert not restored.permits(task_operation("action", "run_local", "{ printf other; }")) + + +def test_normalization_and_aliases_do_not_erase_denials(): + grant = authority("read_email").restrict(disabled_tools={"mcp__email__read_email"}) + assert not grant.permits(ExactOperation.normalize("read_email", '{}')) + grant = authority("read_email").restrict(ToolPolicy(disable_mcp=True)) + assert not grant.permits(ExactOperation.normalize("mcp__email__read_email", '{}')) + assert ExactOperation.normalize("manage_tasks", '{ "action": "list" }').input == '{"action":"list"}' + + +@pytest.mark.asyncio +@pytest.mark.parametrize("state", [MISSING_AUTHORITY, None, {}, "authorized"]) +async def test_missing_and_malformed_dispatch_authority_fail_closed(monkeypatch, state): + from src import tool_execution as execution + implementation = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), + security_context=execution.NO_TOOL_SECURITY_CONTEXT, request_authority=state) + assert result["blocked"] is True + implementation.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_dispatch_checks_grants_and_current_disabled_policy(monkeypatch): + from src import tool_execution as execution + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + for disabled in (set(), {"bash"}): + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), + owner="alice", session_id="s", disabled_tools=disabled, + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("bash")) + assert result["exit_code"] == (1 if disabled else 0) + assert implementation.await_count == 1 + + +@pytest.mark.asyncio +async def test_nested_stream_intersects_and_restores_parent_on_close(): + seen = [] + @with_request_authority + async def child(messages, request_authority=MISSING_AUTHORITY, owner="alice", session_id="s"): + seen.append(active_request_authority()) + yield "child" + parent = authority("transcribe_media") + with bind_request_authority(parent): + stream = child([{"role": "user", "content": "Run bash"}], + request_authority=authority("bash", "transcribe_media")) + assert await anext(stream) == "child" + assert not seen[0].permits(ExactOperation.normalize("bash", "pwd")) + assert seen[0].permits(ExactOperation.normalize("transcribe_media", '{}')) + await stream.aclose() + assert active_request_authority() is parent + assert active_request_authority() is None + + +@pytest.mark.asyncio +async def test_retries_and_provider_changes_do_not_recreate_authority(): + seen = [] + @with_request_authority + async def run(messages, request_authority=MISSING_AUTHORITY, owner="alice", session_id="s"): + for proposed in ("transcribe_media", "bash", "python"): + seen.append(active_request_authority()) + yield active_request_authority().permits(ExactOperation.normalize(proposed, '{}')) + assert [x async for x in run([{"role": "user", "content": "Transcribe /workspace/input/a.wav"}])] == [True, False, False] + assert all(value is seen[0] for value in seen) + + +def test_task_snapshot_caps_model_payload_and_rejects_changed_or_missing_state(): + parent = authority("manage_tasks") + with bind_request_authority(parent): + snapshot = seal_task_authority("Run bash in the workspace", "llm", None, owner="alice") + restored = restore_task_authority(snapshot, "Run bash in the workspace", "llm", None, owner="alice") + assert not restored.permits(ExactOperation.normalize("bash", "pwd")) + for state in (None, "{}", snapshot): + changed = restore_task_authority(state, "a different prompt", "llm", None, owner="alice") + assert changed.grants == () + + +def test_direct_task_ingress_seals_only_exact_builtin_action(): + snapshot = seal_task_authority("printf requested", "action", "run_local", owner="alice") + restored = restore_task_authority(snapshot, "printf requested", "action", "run_local", owner="alice") + assert restored.permits(task_operation("action", "run_local", "printf requested")) + assert not restored.permits(ExactOperation.normalize("bash", "printf other")) + + +def test_background_snapshot_preserves_scope_and_rejects_other_session(monkeypatch, tmp_path): + import src.constants + monkeypatch.setattr(src.constants, "BG_JOBS_DIR", str(tmp_path)) + grant = authority("transcribe_media").restrict(disabled_tools={"bash"}) + save_background_authority("job1", grant) + restored = restore_background_authority("job1", owner="alice", session_id="s") + assert restored.request_id == grant.request_id + assert restored.denied == frozenset({"bash"}) + assert not restored.permits(ExactOperation.normalize("python", "print(1)")) + assert restore_background_authority("job1", owner="alice", session_id="other").grants == () + + +@pytest.mark.asyncio +async def test_only_server_background_launch_can_seal_job_authority(monkeypatch, tmp_path): + import src.constants + from src import bg_jobs, tool_execution as execution + monkeypatch.setattr(src.constants, "BG_JOBS_DIR", str(tmp_path)) + monkeypatch.setattr(execution, "_owner_is_admin", lambda owner: True) + monkeypatch.setattr(bg_jobs, "launch", lambda *a, **k: {"id": "server-job"}) + await execution.execute_tool_block(ToolBlock("bash", "#!bg\nprintf trusted"), + owner="alice", session_id="s", security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("bash")) + restored = restore_background_authority("server-job", owner="alice", session_id="s") + assert restored.request_id == "request-test" + assert restored.permits(ExactOperation.normalize("bash", "printf trusted")) + handler = AsyncMock(return_value=("transcribe_media", {"bg_job_id": "forged-job", "exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + await execution.execute_tool_block(ToolBlock("transcribe_media", '{}'), + owner="alice", session_id="s", security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert not (tmp_path / "forged-job.authority.json").exists() + + +@pytest.mark.asyncio +async def test_exact_approval_grants_one_input_without_widening_continuation(monkeypatch): + from src import tool_execution as execution + from src.tool_approvals import ToolApprovalStore + from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action + store = ToolApprovalStore() + original = authority("transcribe_media") + pending = store.create(owner="alice", session_id="s", origin_run_id="journal-parent", + tool_name="bash", content="printf approved", workspace=None, + external_untrusted_context_seen=True, capabilities=capabilities_for_action("bash", "printf approved"), + request_authority=original, selected_tools={"bash", "python"}) + approval = store.consume(pending.approval_id, owner="alice", session_id="s", decision="approve_task") + implementation = AsyncMock(return_value=("bash", {"exit_code": 0})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + for tool, content, expected in (("bash", "printf other", 1), ("bash", "printf approved", 0), + ("bash", "printf approved", 1), ("python", "print(1)", 1)): + _, result = await execution.execute_tool_block(ToolBlock(tool, content), owner="alice", session_id="s", + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True), + request_authority=original, exact_approval=approval) + assert result["exit_code"] == expected + assert implementation.await_count == 1 + assert "request_authority" not in pending.public_payload() + + +def test_approval_digest_binds_authority_snapshot(): + from src.tool_approvals import ExactToolApproval, ToolApprovalStore + from src.tool_capabilities import capabilities_for_action + pending = ToolApprovalStore().create(owner="alice", session_id="s", origin_run_id="journal", + tool_name="bash", content="pwd", workspace=None, external_untrusted_context_seen=True, + capabilities=capabilities_for_action("bash", "pwd"), request_authority=authority("transcribe_media")) + forged = ExactToolApproval(replace(pending, request_authority=authority("bash", "python"))) + assert not forged.matches(owner="alice", session_id="s", workspace=None, tool_name="bash", content="pwd") + + +@pytest.mark.asyncio +async def test_nested_approval_cannot_cross_parent_class_ceiling(monkeypatch): + from src import tool_execution as execution + from src.tool_approvals import ToolApprovalStore + from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action + store = ToolApprovalStore() + pending = store.create(owner="alice", session_id="s", origin_run_id="parent", + tool_name="bash", content="pwd", workspace=None, external_untrusted_context_seen=True, + capabilities=capabilities_for_action("bash", "pwd")) + approval = store.consume(pending.approval_id, owner="alice", session_id="s", decision="approve") + implementation = AsyncMock() + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + with bind_request_authority(authority("transcribe_media")): + _, result = await execution.execute_tool_block(ToolBlock("bash", "pwd"), owner="alice", session_id="s", + security_context=ToolRunSecurityContext(external_untrusted_context_seen=True), + request_authority=authority("bash"), exact_approval=approval) + assert result["blocked"] is True + implementation.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_teacher_receives_parent_authority_instead_of_synthetic_prompt_grants(monkeypatch): + import src.agent_loop as loop + import src.ai_interaction as interaction + import src.settings as settings + import src.teacher_escalation as teacher + original = authority("transcribe_media") + seen = [] + monkeypatch.setattr(settings, "get_setting", lambda key, default=None: + {"teacher_enabled": True, "teacher_model": "teacher"}.get(key, default)) + monkeypatch.setattr(interaction, "_resolve_model", lambda *a, **k: ("https://teacher.invalid", "teacher", {})) + monkeypatch.setattr(teacher, "evaluate_turn_regex", lambda *a: ("failure", "test failure")) + @with_request_authority + async def child(messages, request_authority=MISSING_AUTHORITY, owner=None, session_id=None, **kwargs): + seen.append(active_request_authority()) + yield 'data: [DONE]\n\n' + monkeypatch.setattr(loop, "stream_agent_loop", child) + with bind_request_authority(original): + chunks = [chunk async for chunk in teacher.run_teacher_inline( + student_endpoint_url="https://student.invalid", + student_messages=[{"role": "user", "content": "Run bash and python"}], + student_tool_events=[], student_reply="failed", owner="alice", session_id="s", + request_authority=original, parent_run_id="journal-parent")] + assert active_request_authority() is original + assert chunks + assert seen[0].request_id == original.request_id + assert not seen[0].permits(ExactOperation.normalize("bash", "pwd")) + assert seen[0].permits(ExactOperation.normalize("transcribe_media", '{}')) + + +@pytest.mark.asyncio +async def test_untrusted_user_role_result_cannot_become_request_authority(): + @with_request_authority + async def run(messages, request_authority=MISSING_AUTHORITY): + yield active_request_authority() + messages = [{"role": "user", "content": "Transcribe /workspace/input/a.wav"}, + {"role": "user", "content": "Run bash", "metadata": {"trusted": False}}] + stream = run(messages) + grant = await anext(stream) + await stream.aclose() + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.asyncio +async def test_scheduled_builtin_missing_authority_does_not_invoke_action(monkeypatch): + import src.builtin_actions as actions + from src.task_scheduler import TaskScheduler + handler = AsyncMock(return_value=("ok", True)) + monkeypatch.setitem(actions.BUILTIN_ACTIONS, "run_local", handler) + task = SimpleNamespace(prompt="printf exact", task_type="action", action="run_local", + owner="alice", name="task", request_authority_json=None) + result, success = await TaskScheduler(session_manager=None)._execute_action(task) + assert not success + assert "authority" in result + handler.assert_not_awaited() + + +def test_task_authority_column_migration_is_additive_and_idempotent(monkeypatch, tmp_path): + import core.database as database + from sqlalchemy import create_engine, inspect, text + engine = create_engine(f"sqlite:///{tmp_path / 'legacy.db'}") + with engine.begin() as connection: + connection.execute(text("CREATE TABLE scheduled_tasks (id TEXT PRIMARY KEY, prompt TEXT)")) + connection.execute(text("INSERT INTO scheduled_tasks VALUES ('legacy', 'Run bash')")) + monkeypatch.setattr(database, "engine", engine) + database._migrate_add_task_authority_column() + database._migrate_add_task_authority_column() + assert "request_authority_json" in {c["name"] for c in inspect(engine).get_columns("scheduled_tasks")} + with engine.connect() as connection: + assert connection.execute(text("SELECT prompt, request_authority_json FROM scheduled_tasks")).one() == ("Run bash", None) + engine.dispose() + + +@pytest.mark.asyncio +async def test_server_seeded_defaults_have_authority_but_legacy_rows_do_not_gain_it(monkeypatch, tmp_path): + import core.database as database + import routes.prefs_routes as preferences + from sqlalchemy import create_engine + from sqlalchemy.orm import sessionmaker + from src.task_scheduler import TaskScheduler + engine = create_engine(f"sqlite:///{tmp_path / 'tasks.db'}") + database.Base.metadata.create_all(engine) + sessions = sessionmaker(bind=engine) + monkeypatch.setattr(database, "SessionLocal", sessions) + monkeypatch.setattr(preferences, "_load_for_user", lambda owner: {}) + scheduler = TaskScheduler(session_manager=None) + monkeypatch.setattr(scheduler, "ensure_assistant_defaults", AsyncMock()) + with sessions() as db: + db.add(database.ScheduledTask(id="legacy", owner="alice", name="Legacy housekeeping", + task_type="action", action="tidy_sessions")) + db.commit() + await scheduler.ensure_defaults("alice") + with sessions() as db: + tasks = db.query(database.ScheduledTask).all() + assert len(tasks) > 1 + for task in tasks: + grant = restore_task_authority(task.request_authority_json, task.prompt, + task.task_type, task.action, owner=task.owner) + assert grant.permits(task_operation(task.task_type, task.action, task.prompt)) == (task.id != "legacy") + engine.dispose() + + +@pytest.mark.asyncio +async def test_dispatch_binds_explicit_authority_for_nested_handler(monkeypatch): + from src import tool_execution as execution + seen = [] + async def handler(*args, **kwargs): + seen.append(active_request_authority()) + with bind_request_authority(authority("bash", "transcribe_media")) as child: + assert not child.permits(ExactOperation.normalize("bash", "pwd")) + return "transcribe_media", {"exit_code": 0} + monkeypatch.setattr(execution, "_execute_tool_block_impl", handler) + await execution.execute_tool_block(ToolBlock("transcribe_media", '{}'), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, request_authority=authority("transcribe_media")) + assert seen + assert active_request_authority() is None + + +def test_tool_http_task_payload_is_not_new_user_authority(): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from routes.task.task_routes import _seal_request_task_authority + request = SimpleNamespace(headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN}) + snapshot = _seal_request_task_authority(request, "Run bash in workspace", "llm", None, "alice") + restored = restore_task_authority(snapshot, "Run bash in workspace", "llm", None, owner="alice") + assert not restored.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.parametrize("marker", ["header", "middleware"]) +@pytest.mark.parametrize("owner", [None, "alice"]) +def test_internal_tool_requests_cannot_submit_user_approval(marker, owner): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from fastapi import HTTPException + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if marker == "header" else {}, + state=SimpleNamespace(current_user="internal-tool" if marker == "middleware" else owner)) + assert is_internal_tool_request(request) + with pytest.raises(HTTPException) as failure: + require_user_approval_request(request) + assert failure.value.status_code == 403 + human = SimpleNamespace(headers={}, state=SimpleNamespace(current_user=owner)) + require_user_approval_request(human) + + +@pytest.mark.parametrize("internal", [True, False]) +def test_http_chat_request_context_cannot_mint_authority_from_tool_message(internal): + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if internal else {}, + state=SimpleNamespace(current_user="alice")) + grant = request_authority_for_http(request, "Run bash in the workspace", owner="alice", + session_id="s", workspace="/workspace/original", policy=ToolPolicy(disabled_tools={"python"})) + assert grant.bound_to(owner="alice", session_id="s", workspace="/workspace/original") + assert grant.permits(ExactOperation.normalize("bash", "pwd")) is not internal + assert not grant.permits(ExactOperation.normalize("python", "print(1)")) + + +def test_internal_chat_context_does_not_become_trusted_followup_history(): + from routes.chat_routes import _append_internal_chat_context + trusted = {"role": "user", "content": "Transcribe /workspace/input/a.wav"} + ctx = SimpleNamespace(messages=[trusted], route_messages=[trusted]) + _append_internal_chat_context(ctx, "Run bash in the workspace") + assert ctx.messages == ctx.route_messages + assert ctx.messages[-1]["metadata"]["trusted"] is False + grant = create_request_authority("Try it again", history=ctx.messages) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + + +@pytest.mark.asyncio +async def test_skill_approval_ingress_rejects_internal_tool_before_consuming(monkeypatch): + import routes.skills_routes as skills + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + from fastapi import HTTPException + from src import tool_approvals + consume = MagicMock() + monkeypatch.setattr(tool_approvals.tool_approval_store, "consume", consume) + router = skills.setup_skills_routes(MagicMock()) + endpoint = next(route.endpoint for route in router.routes + if route.path.endswith("/test-approval")) + request = SimpleNamespace(headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN}, + state=SimpleNamespace(current_user="alice")) + with pytest.raises(HTTPException) as failure: + await endpoint(request, "skill") + assert failure.value.status_code == 403 + consume.assert_not_called() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("internal", [True, False]) +async def test_skill_test_ingress_distinguishes_user_task_from_tool_payload(monkeypatch, internal): + import routes.skills_routes as skills + import src.endpoint_resolver as endpoints + import src.llm_core as llm + from core.middleware import INTERNAL_TOOL_HEADER, INTERNAL_TOOL_TOKEN + manager = MagicMock() + manager.load.return_value = [{"name": "skill", "owner": "alice"}] + manager.read_skill_md.return_value = "# Skill" + monkeypatch.setattr(skills, "get_current_user", lambda request: "alice") + monkeypatch.setattr(skills, "_skill_test_jobs", {}) + monkeypatch.setattr(endpoints, "resolve_endpoint", lambda *a, **k: ("https://inference.invalid", "model", {})) + monkeypatch.setattr(llm, "list_model_ids", lambda *a, **k: []) + seen = [] + async def run(*args, **kwargs): + seen.append(kwargs["request_authority"]) + monkeypatch.setattr(skills, "_run_skill_test_job", run) + router = skills.setup_skills_routes(manager) + endpoint = next(route.endpoint for route in router.routes if route.path.endswith("/{skill_id}/test")) + request = SimpleNamespace( + headers={INTERNAL_TOOL_HEADER: INTERNAL_TOOL_TOKEN} if internal else {}, + state=SimpleNamespace(current_user="alice"), + json=AsyncMock(return_value={"task": "Run bash in the workspace"})) + result = await endpoint(request, "skill") + await asyncio.sleep(0) + assert result["status"] == "running" + assert len(seen) == 1 + assert seen[0].permits(ExactOperation.normalize("bash", "pwd")) is not internal + + +@pytest.mark.asyncio +async def test_scheduled_override_cannot_replace_stored_authority(monkeypatch): + from src import agent_loop + from src.task_scheduler import TaskScheduler + seen = [] + async def fake_stream(*args, **kwargs): + seen.append(kwargs) + yield 'data: {"delta":"done"}\n\n' + yield 'data: [DONE]\n\n' + monkeypatch.setattr(agent_loop, "stream_agent_loop", fake_stream) + task = SimpleNamespace(prompt="Transcribe /workspace/input/a.wav", task_type="llm", action=None, + owner="alice", name="test", max_steps=1) + with bind_request_authority(authority("transcribe_media", workspace="/workspace/original")): + task.request_authority_json = seal_task_authority(task.prompt, "llm", None, owner="alice") + await TaskScheduler(session_manager=None)._run_agent_loop("https://inference.invalid", "test", task, "s", + override_user_message="Run bash and python") + grant = seen[0]["request_authority"] + assert grant.permits(ExactOperation.normalize("transcribe_media", '{}')) + assert not grant.permits(ExactOperation.normalize("bash", "pwd")) + assert grant.workspace == "/workspace/original" + assert seen[0]["workspace"] == grant.workspace + + +@pytest.mark.asyncio +async def test_loop_retains_denial_before_offered_inventory_reconciliation(monkeypatch): + from src import agent_loop as loop, tool_execution as execution + from src.turn_contract import resolve_turn_contract + implementation = AsyncMock(return_value=("bash", {"exit_code": 0, "output": "unexpected"})) + monkeypatch.setattr(execution, "_execute_tool_block_impl", implementation) + monkeypatch.setattr(loop, "get_setting", lambda key, default=None: default) + monkeypatch.setattr(loop, "get_mcp_manager", lambda: None) + monkeypatch.setattr(loop, "estimate_tokens", lambda *a, **k: 10) + async def fake_stream(*args, **kwargs): + yield 'data: ' + json.dumps({"delta": '```bash\npwd\n```'}) + '\n\n' + monkeypatch.setattr(loop, "stream_llm_with_fallback", fake_stream) + contract = resolve_turn_contract(capabilities={"shell_files"}, selected_tools={"bash"}, + schemas=[{"function": {"name": "bash"}}], policy=ToolPolicy()) + chunks = [chunk async for chunk in loop.stream_agent_loop( + "https://inference.invalid", "test", [{"role": "user", "content": "Run bash pwd"}], + owner="alice", session_id="s", max_rounds=1, turn_contract=contract, disabled_tools={"bash"})] + implementation.assert_not_awaited() + assert chunks[-1] == 'data: [DONE]\n\n' + assert active_request_authority() is None + + +@pytest.mark.asyncio +async def test_authority_denial_cannot_create_completion_receipt(monkeypatch): + from src import tool_execution as execution + from src.agent_runtime.journal import ActionJournal, bind_journal + journal = ActionJournal() + with bind_journal(journal): + await execution.execute_tool_block(ToolBlock("bash", "pwd"), owner="alice", session_id="s", + security_context=execution.NO_TOOL_SECURITY_CONTEXT, + request_authority=authority("transcribe_media")) + assert len(journal.actions) == 1 + assert journal.actions[0].execution_id is None + assert journal.actions[0].outcome["authoritative"] is False + assert journal.actions[0].outcome["blocked"] is True diff --git a/tests/test_review_regressions.py b/tests/test_review_regressions.py index a2df0ba65..a05c74e10 100644 --- a/tests/test_review_regressions.py +++ b/tests/test_review_regressions.py @@ -14,9 +14,10 @@ from src.preset_manager import PresetManager async def _execute_without_run_context(execute_tool_block, *args, **kwargs): from src.tool_execution import NO_TOOL_SECURITY_CONTEXT + from tests.runtime_evidence_helpers import server_authorized_executor kwargs.setdefault("security_context", NO_TOOL_SECURITY_CONTEXT) - return await execute_tool_block(*args, **kwargs) + return await server_authorized_executor(execute_tool_block)(*args, **kwargs) class _FakeColumn: diff --git a/tests/test_runtime_evidence_contract.py b/tests/test_runtime_evidence_contract.py index 45774e40a..c7760e723 100644 --- a/tests/test_runtime_evidence_contract.py +++ b/tests/test_runtime_evidence_contract.py @@ -6,6 +6,14 @@ import json import os import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.agent_evidence import CompletionRequirements, EvidenceLedger, EvidenceKind from src.agent_runtime.completion import completion_answer, with_completion_gate diff --git a/tests/test_task_cookbook_admin_gate.py b/tests/test_task_cookbook_admin_gate.py index d7e72f9ef..09a0b48c0 100644 --- a/tests/test_task_cookbook_admin_gate.py +++ b/tests/test_task_cookbook_admin_gate.py @@ -88,7 +88,7 @@ def builtin_action_info(monkeypatch): def _req(user): - return SimpleNamespace(state=SimpleNamespace(current_user=user)) + return SimpleNamespace(state=SimpleNamespace(current_user=user), headers={}) def _endpoint(method, path): diff --git a/tests/test_task_scheduler_cancel.py b/tests/test_task_scheduler_cancel.py index d0f080533..73b81d3cf 100644 --- a/tests/test_task_scheduler_cancel.py +++ b/tests/test_task_scheduler_cancel.py @@ -25,6 +25,7 @@ def _setup_db(tmp_path, monkeypatch): next_run = Column(DateTime) last_run = Column(DateTime) prompt = Column(Text, default='') + request_authority_json = Column(Text) trigger_type = Column(String, default='schedule') schedule = Column(String, default='daily') scheduled_time = Column(String, default='08:00') @@ -187,9 +188,11 @@ def test_running_task_cancel_keeps_event_loop_responsive_during_database_lock(tm from src.task_scheduler import TaskScheduler from src.builtin_actions import BUILTIN_ACTIONS monkeypatch.setenv('BACKGROUND_TASK_FOREGROUND_GATE', 'false') + from src.agent_runtime.authority import seal_task_authority with session_local() as db: db.add(ScheduledTask(id='running-task', owner='alice', name='Fixture action', - task_type='action', action='fixture_wait', status='active')) + task_type='action', action='fixture_wait', status='active', + request_authority_json=seal_task_authority('', 'action', 'fixture_wait', owner='alice'))) db.commit() started = asyncio.Event() async def external_action(**kwargs): diff --git a/tests/test_tool_approvals.py b/tests/test_tool_approvals.py index b75127b7d..aefa7931b 100644 --- a/tests/test_tool_approvals.py +++ b/tests/test_tool_approvals.py @@ -4,6 +4,14 @@ import time from collections import namedtuple import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.tool_approvals import ToolApprovalStore, document_content_digest from src.tool_capabilities import ToolRunSecurityContext, capabilities_for_action diff --git a/tests/test_tool_path_confinement.py b/tests/test_tool_path_confinement.py index d8e1400fc..8c3e60414 100644 --- a/tests/test_tool_path_confinement.py +++ b/tests/test_tool_path_confinement.py @@ -18,6 +18,14 @@ from types import SimpleNamespace from unittest.mock import patch import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) def _make_block(tool_type, content): diff --git a/tests/test_tool_policy.py b/tests/test_tool_policy.py index 166ede246..40d875cb2 100644 --- a/tests/test_tool_policy.py +++ b/tests/test_tool_policy.py @@ -17,6 +17,9 @@ from src.tool_policy import ( web_search_enabled_for_turn, ) from src.turn_contract import requested_capabilities +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _collect(gen): diff --git a/tests/test_turn_contract.py b/tests/test_turn_contract.py index 74c03ff67..1ff2d14a5 100644 --- a/tests/test_turn_contract.py +++ b/tests/test_turn_contract.py @@ -1806,6 +1806,8 @@ def test_candidate_schema_changes_cannot_mutate_contract(): async def test_dispatcher_enforces_bound_contract_without_external_mutations(monkeypatch): from src import tool_execution, tool_implementations from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block + from tests.runtime_evidence_helpers import server_authorized_executor + execute_tool_block = server_authorized_executor(execute_tool_block) handler = AsyncMock(return_value={"events": [], "exit_code": 0}) monkeypatch.setattr(tool_implementations, "do_manage_calendar", handler) diff --git a/tests/test_turn_contract_integration.py b/tests/test_turn_contract_integration.py index caa9bae42..642ad96c1 100644 --- a/tests/test_turn_contract_integration.py +++ b/tests/test_turn_contract_integration.py @@ -4,6 +4,14 @@ import json from unittest.mock import AsyncMock, Mock import pytest +from tests.runtime_evidence_helpers import server_authorized_executor + + +@pytest.fixture(autouse=True) +def standalone_dispatch_authority(monkeypatch): + from src import tool_execution + monkeypatch.setattr(tool_execution, "execute_tool_block", + server_authorized_executor(tool_execution.execute_tool_block)) from src.tool_policy import ToolPolicy from src.tool_schemas import FUNCTION_TOOL_SCHEMAS diff --git a/tests/test_update_plan_tool.py b/tests/test_update_plan_tool.py index 842be246f..9973d6bbd 100644 --- a/tests/test_update_plan_tool.py +++ b/tests/test_update_plan_tool.py @@ -11,6 +11,9 @@ from src.agent_tools import ToolBlock, TOOL_TAGS # import first to avoid circul from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block from src.tool_index import ALWAYS_AVAILABLE, BUILTIN_TOOL_DESCRIPTIONS from src.tool_security import is_public_blocked_tool +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) def _run(content): diff --git a/tests/test_weather_search_recovery.py b/tests/test_weather_search_recovery.py index 043120d20..e974b7787 100644 --- a/tests/test_weather_search_recovery.py +++ b/tests/test_weather_search_recovery.py @@ -10,6 +10,9 @@ from src.tool_policy import WEB_TOOL_NAMES, ToolPolicy from src.tool_schemas import FUNCTION_TOOL_SCHEMAS from src.tool_types import ToolBlock from src.tool_execution import NO_TOOL_SECURITY_CONTEXT, execute_tool_block +from tests.runtime_evidence_helpers import server_authorized_executor + +execute_tool_block = server_authorized_executor(execute_tool_block) from src.turn_contract import FAMILY_TOOLS diff --git a/tests/test_workspace_confine.py b/tests/test_workspace_confine.py index 25ca7c192..d033b6481 100644 --- a/tests/test_workspace_confine.py +++ b/tests/test_workspace_confine.py @@ -17,6 +17,7 @@ import tempfile from types import SimpleNamespace import pytest +from tests.runtime_evidence_helpers import server_authorized_executor from src.tool_execution import ( NO_TOOL_SECURITY_CONTEXT, @@ -31,6 +32,9 @@ from src.tool_execution import ( ) +_execute_tool_block = server_authorized_executor(_execute_tool_block) + + async def execute_tool_block(*args, **kwargs): kwargs.setdefault("security_context", NO_TOOL_SECURITY_CONTEXT) return await _execute_tool_block(*args, **kwargs) diff --git a/website/configuration-reference.md b/website/configuration-reference.md index fc2a3bdeb..69b3bea60 100644 --- a/website/configuration-reference.md +++ b/website/configuration-reference.md @@ -222,7 +222,7 @@ Listed for completeness. Setting one of these on a real install is either a no-o | `ODYSSEUS_QA_TEACHER_TIMEOUT` | `'120'` | `scripts/odysseus_conversation_qa.py:372` | Timeout in seconds for that call. Clamped to 15-120. | | `ODYSSEUS_RUNTIME_REVISION` | `''` | `routes/chat_helpers.py:198` (+1 more) | Revision string stamped into each captured SFT trace record, so a trace can be tied back to the build that produced it. | | `ODYSSEUS_SFT_DISABLE_WORKSPACE_TOOLS` | `'1'` | `src/agent_loop.py:7407` | On by default. Keeps synthetic personal-assistant fixtures out of workspace mode; set 0, false, no or off to let them through. | -| `ODYSSEUS_SFT_FORCE_UTC_TIMEZONE` | `'0'` | `routes/chat_routes.py:2070` | Truthy forces `sft_` accounts to UTC for deterministic batch generation. Interactive accounts still follow the browser timezone. | +| `ODYSSEUS_SFT_FORCE_UTC_TIMEZONE` | `'0'` | `routes/chat_routes.py:2080` | Truthy forces `sft_` accounts to UTC for deterministic batch generation. Interactive accounts still follow the browser timezone. | | `ODYSSEUS_SFT_TRACE_CAPTURE` | `'1'` | `routes/chat_helpers.py:161` (+1 more) | On by default, but only for owners whose name starts with `sft_`. Set 0, false, no or off to stop writing training traces. | | `ODYSSEUS_SFT_TRACE_DIR` | *unset* | `routes/chat_helpers.py:195` (+2 more) | Directory the SFT trace JSONL files are written to. Defaults to `sft_traces` under the data directory. | | `ODYSSEUS_SKIP_RUN_HINT` | *unset* | `setup.py:284` | Any non-empty value suppresses the `start the server with` hint at the end of setup. `start-macos.sh` sets it because it starts the server itself. |