diff --git a/scripts/validate_runtime_wave1.sh b/scripts/validate_runtime_wave1.sh new file mode 100644 index 000000000..6177bcfc3 --- /dev/null +++ b/scripts/validate_runtime_wave1.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash +# Focused runtime gate; no model inference or benchmark fixture access. +set -euo pipefail +cd "$(dirname "${BASH_SOURCE[0]}")/.." +export ODYSSEUS_TEST_STATIC_PORT=0 +export PYTHONDONTWRITEBYTECODE=1 +export PYTHON_DOTENV_DISABLED=1 +export ODYSSEUS_DATA_DIR="${ODYSSEUS_DATA_DIR:-/tmp/odysseus-runtime-decomposition-test-state}" +exec "${ODYSSEUS_TEST_PYTHON:-python3}" -m pytest -q -p no:cacheprovider \ + tests/test_runtime_evidence_contract.py tests/test_agent_evidence.py \ + tests/test_agent_evidence_loop.py tests/test_agent_render_ownership.py \ + tests/test_agent_runs_terminal_order.py tests/test_agent_loop.py \ + tests/test_tool_task_cancelled_on_disconnect.py tests/test_turn_contract.py \ + tests/test_agent_turn_contract_boundaries.py tests/test_loop_breaker_runaway.py \ + tests/test_chat_route_tool_policy.py tests/test_agent_runtime_context.py \ + tests/test_external_context_tool_gate.py tests/test_workspace_confine.py \ + tests/test_private_browser_tool.py tests/test_bg_jobs_store.py tests/test_bg_job_tools.py \ + tests/test_context_budget.py tests/test_context_compactor.py \ + tests/test_context_compactor_nonstring.py tests/test_generation_budget.py \ + tests/test_foreground_model_routing.py tests/test_tool_policy.py \ + tests/test_execution_bridge.py tests/test_tool_approvals.py \ + tests/test_tool_approval_single_action_scope.py tests/test_tool_approval_task_scope.py \ + tests/test_mcp_text_error_normalization.py tests/test_mcp_email_search_error_transport.py \ + tests/test_tool_path_confinement.py tests/test_workspace_artifact_tool_floor.py \ + tests/test_native_unattended_workspace_floor.py tests/test_misfenced_read_file_tool_call.py \ + tests/test_builtin_mcp_pythonpath.py tests/test_python_tool_import_paths.py "$@" diff --git a/src/agent_evidence.py b/src/agent_evidence.py index 5fc6165d8..d0f7560e6 100644 --- a/src/agent_evidence.py +++ b/src/agent_evidence.py @@ -9,6 +9,7 @@ from dataclasses import asdict, dataclass, field from enum import Enum from pathlib import Path from typing import Any, Iterable, Mapping, Sequence +from src.agent_runtime.identity import artifact_identity, artifact_version, executable_words, is_test_command, is_validation_command def workspace_artifact_is_usable(path: Path) -> bool: @@ -99,6 +100,10 @@ class EvidenceEvent: command_sha256: str = "" output_sha256: str = "" detail: str = "" + action_id: str = "" + execution_id: str = "" + artifact_id: str = "" + verification_id: str = "" def to_dict(self) -> dict[str, Any]: data = asdict(self) @@ -195,12 +200,12 @@ _VALIDATION_COMMAND_RE = re.compile( def command_is_validation(command: str) -> bool: """Return whether a shell command provides executable verification evidence.""" value = str(command or "") - return bool(_TEST_COMMAND_RE.search(value) or _VALIDATION_COMMAND_RE.search(value)) + return is_validation_command(value) def command_is_test(command: str) -> bool: """Return whether a shell command executes a recognized test runner.""" - return bool(_TEST_COMMAND_RE.search(str(command or ""))) + return is_test_command(str(command or "")) def _clean_path(value: str) -> str: @@ -441,18 +446,12 @@ def _path_is_mentioned(command: str, required_path: str) -> bool: return path in command or Path(path).name in command -def _artifact_path_matches_required(artifact_path: str, required_path: str) -> bool: - artifact = _clean_path(artifact_path) +def _artifact_path_matches_required(artifact_path: str, required_path: str, workspace: str = "") -> bool: + artifact = str(artifact_path or '').strip() required = _clean_path(required_path) if not artifact or not required: return False - if artifact == required: - return True - # Absolute requirements are exact output contracts; same basename in a - # different directory is not enough. - if artifact.startswith("/") or required.startswith("/"): - return False - return Path(artifact).name == Path(required).name + return artifact_identity(artifact, workspace) == artifact_identity(required, workspace) def _explicit_tool_paths(tool: str, command: str) -> list[str]: @@ -462,21 +461,21 @@ def _explicit_tool_paths(tool: str, command: str) -> list[str]: except (TypeError, json.JSONDecodeError): args = None if isinstance(args, Mapping): - path = _clean_path(str(args.get("path") or "")) + path = str(args.get("path") or "").strip() return [path] if path else [] # Keep compatibility with the legacy ``path\ncontent`` transport. - path = _clean_path(str(command or "").splitlines()[0] if command else "") + path = (str(command or "").splitlines()[0] if command else "").strip() return [path] if path else [] if tool == "edit_file": try: args = json.loads(command or "{}") except (TypeError, json.JSONDecodeError): return [] - path = _clean_path(str(args.get("path") or "")) if isinstance(args, dict) else "" + path = str(args.get("path") or "").strip() if isinstance(args, dict) else "" return [path] if path else [] if tool == "apply_patch": return [ - _clean_path(match.group(1)) + match.group(1).strip() for match in re.finditer(r"^\*\*\* (?:Add|Update|Delete) File:\s*(.+)$", command or "", re.MULTILINE) if _clean_path(match.group(1)) ] @@ -549,13 +548,13 @@ def _command_text(value: str) -> str: def _matches_declared_verifier(command: str, expected: Sequence[str]) -> bool: - actual = " ".join(_command_text(command).split()) + actual = executable_words(_command_text(command)) if not actual: return False return any( - normalized == actual or normalized in actual + normalized == actual for item in expected - if (normalized := " ".join(str(item or "").split())) + if (normalized := executable_words(str(item or ""))) ) @@ -574,6 +573,7 @@ class EvidenceLedger: def __init__(self, requirements: CompletionRequirements | None = None) -> None: self.requirements = requirements or CompletionRequirements() self.events: list[EvidenceEvent] = [] + self._verification_versions: dict[str, str] = {} @classmethod def from_tool_events( @@ -611,6 +611,10 @@ class EvidenceLedger: "command_sha256": _digest(command), "output_sha256": _digest(output), } + action_id = str(source.get('action_id') or '') + execution_id = str(source.get('execution_id') or '') + if action_id: + payload.update(action_id=action_id, execution_id=execution_id) evidence = EvidenceEvent( event_id=_event_id(payload, len(self.events)), kind=kind, @@ -623,6 +627,11 @@ class EvidenceLedger: command_sha256=payload["command_sha256"], output_sha256=payload["output_sha256"], detail=detail, + action_id=action_id, + execution_id=execution_id, + artifact_id=artifact_identity(artifact_path, self.requirements.workspace_root) if artifact_path else '', + verification_id=('verification-' + _event_id(payload, len(self.events))) + if kind in {EvidenceKind.VERIFIER_RESULT, EvidenceKind.ARTIFACT_VALIDATION} else '', ) self.events.append(evidence) return evidence @@ -631,8 +640,12 @@ class EvidenceLedger: tool = str(event.get("tool") or "") command = str(event.get("command") or "") exit_code = event.get("exit_code") - authoritative = isinstance(exit_code, int) and not isinstance(exit_code, bool) - success = authoritative and exit_code == 0 + authoritative = ( + isinstance(exit_code, int) and not isinstance(exit_code, bool) + and not event.get("blocked") and not event.get("approval_required") + and event.get("execution_attempted") is not False + ) + success = authoritative and exit_code == 0 and not event.get('error') if not authoritative: success = not bool(event.get("error")) self._append( @@ -644,7 +657,11 @@ class EvidenceLedger: explicit_paths = _explicit_tool_paths(tool, command) mutation_paths = list(explicit_paths) - if command_has_mutation_effect(command) and tool not in { + observed_changes = event.get('artifact_changes') + if isinstance(observed_changes, list) and tool in {'bash', 'python', 'host_shell'}: + mutation_paths.extend(path for path in self.requirements.required_artifacts + if artifact_identity(path, self.requirements.workspace_root) in observed_changes) + elif command_has_mutation_effect(command) and tool not in { "write_file", "edit_file", "apply_patch", @@ -657,7 +674,7 @@ class EvidenceLedger: ) seen_paths: set[str] = set() for path in mutation_paths: - path = _clean_path(path) + path = str(path or '').strip() if not path or path in seen_paths: continue seen_paths.add(path) @@ -675,12 +692,12 @@ class EvidenceLedger: except (TypeError, json.JSONDecodeError): read_args = None read_path = ( - _clean_path(str(read_args.get("path") or "")) + str(read_args.get("path") or "").strip() if isinstance(read_args, Mapping) - else "" + else command.strip() if read_args is None else "" ) if read_path and any( - _artifact_path_matches_required(read_path, required) + _artifact_path_matches_required(read_path, required, self.requirements.workspace_root) for required in self.requirements.required_artifacts ): self._append( @@ -692,10 +709,13 @@ class EvidenceLedger: detail="post-write artifact inspection", ) - if _TEST_COMMAND_RE.search(_command_text(command)) or _matches_declared_verifier( + if tool in {"bash", "host_shell"} and (command_is_test(_command_text(command)) or _matches_declared_verifier( command, self.requirements.verifier_commands, - ): + )): + if authoritative: + versions = event.get('artifact_versions') + self._verification_versions = dict(versions) if isinstance(versions, Mapping) else {} self._append( kind=EvidenceKind.VERIFIER_RESULT, success=success, @@ -703,7 +723,7 @@ class EvidenceLedger: source=event, detail="executable test/verifier command", ) - elif _VALIDATION_COMMAND_RE.search(command) and not mutation_paths: + elif tool in {"bash", "host_shell"} and is_validation_command(command) and not mutation_paths: for path in self.requirements.required_artifacts: if _path_is_mentioned(command, path): self._append( @@ -767,6 +787,19 @@ class EvidenceLedger: (latest_verifier.event_id,), ) + if latest_verifier and self.requirements.workspace_root: + for path in self.requirements.required_artifacts: + identity = artifact_identity(path, self.requirements.workspace_root) + expected = self._verification_versions.get(identity) + if expected in {'unobserved', 'missing-or-unreadable'}: + return CompletionDecision(CompletionStatus.BLOCKED, False, + 'artifact version could not be established for verification', + (latest_verifier.event_id,)) + if expected is not None and expected != artifact_version(path, self.requirements.workspace_root): + return CompletionDecision(CompletionStatus.BLOCKED, False, + 'artifact content changed after verification', + (latest_verifier.event_id,)) + satisfied_ids: list[str] = [] missing: list[str] = [] workspace_root = str(self.requirements.workspace_root or "").strip() @@ -774,7 +807,7 @@ class EvidenceLedger: matches = [ event for event in self.events if event.kind == EvidenceKind.ARTIFACT_MUTATION - and _artifact_path_matches_required(event.artifact_path, required) + and _artifact_path_matches_required(event.artifact_path, required, self.requirements.workspace_root) ] authoritative = [ event for event in matches @@ -793,10 +826,11 @@ class EvidenceLedger: and latest.tool in {"bash", "python"} ) filesystem_missing = False - if latest_success is not None and workspace_root and required.startswith("/workspace/"): + if latest_success is not None and workspace_root: try: root = Path(workspace_root).resolve() - candidate = (root / required.removeprefix("/workspace/")).resolve() + identity = artifact_identity(required, workspace_root) + candidate = (root / identity.removeprefix('workspace:')).resolve() if identity.startswith('workspace:') else Path(required).resolve() candidate.relative_to(root) filesystem_missing = not workspace_artifact_is_usable(candidate) except (OSError, RuntimeError, ValueError): @@ -852,14 +886,14 @@ class EvidenceLedger: if event.kind == EvidenceKind.ARTIFACT_MUTATION and event.authoritative and event.success - and _artifact_path_matches_required(event.artifact_path, required) + and _artifact_path_matches_required(event.artifact_path, required, self.requirements.workspace_root) ] matching_validations = [ (index, event) for index, event in enumerate(self.events) if event.kind == EvidenceKind.ARTIFACT_VALIDATION and event.authoritative - and _artifact_path_matches_required(event.artifact_path, required) + and _artifact_path_matches_required(event.artifact_path, required, self.requirements.workspace_root) ] if not matching_validations: continue @@ -882,6 +916,10 @@ class EvidenceLedger: current_validation_ids.append(latest_validation.event_id) if self.requirements.verifier_required and latest_verifier is None: + if self.requirements.executable_verifier_available: + return CompletionDecision(CompletionStatus.BLOCKED, False, + 'the request requires an executable verifier result', + tuple(satisfied_ids)) validation_ids: list[str] = [] for required in self.requirements.required_artifacts: matching_validation = [ @@ -890,7 +928,7 @@ class EvidenceLedger: if event.kind == EvidenceKind.ARTIFACT_VALIDATION and event.authoritative and event.success - and _artifact_path_matches_required(event.artifact_path, required) + and _artifact_path_matches_required(event.artifact_path, required, self.requirements.workspace_root) ] latest_validation = matching_validation[-1] if matching_validation else None if latest_validation is None or latest_validation[0] < latest_mutation_index: diff --git a/src/agent_loop.py b/src/agent_loop.py index b25ac0e1c..a8ca4e017 100644 --- a/src/agent_loop.py +++ b/src/agent_loop.py @@ -82,6 +82,8 @@ from src.tool_approvals import ( ) from src.tool_types import ToolBlock from src.turn_contract import selected_tools_for_request, with_turn_contract +from src.agent_runtime.journal import propose_action, execute_action +from src.agent_runtime.completion import with_completion_gate from src.tool_utils import _truncate, get_mcp_manager from src.agent_tools import ( parse_tool_blocks, @@ -20324,6 +20326,7 @@ def _blocks_before_inference(turn_contract) -> bool: @with_turn_contract +@with_completion_gate async def stream_agent_loop( endpoint_url: str, model: str, @@ -29694,7 +29697,10 @@ async def stream_agent_loop( _completion_requirements, ) _round_decision = _round_evidence.evaluate() - if not _round_decision.can_complete and _evidence_repair_rounds < 2: + # Missing evidence is an incomplete result, not a reason to + # manufacture additional provider rounds. Actual diagnostic + # failures can still enter the bounded recovery path. + if _round_decision.status.value == "failed" and _evidence_repair_rounds < 2: _evidence_repair_rounds += 1 _missing = ", ".join(_round_decision.missing_artifacts) _declared_verifiers = _completion_requirements.verifier_commands @@ -31423,6 +31429,15 @@ async def stream_agent_loop( local_network_budget_hit = False local_inspection_budget_hit = False for i, block in enumerate(tool_blocks): + native_call = converted_calls[i] if i < len(converted_calls) else None + tool_call_id = _resolved_tool_call_id( + native_call, + session_id=str(session_id or ""), + round_num=round_num, + tool_index=i, + tool_name=block.tool_type, + ) + _runtime_action = propose_action(block, tool_call_id, native_call) _call_signature = _tool_call_signature(block.tool_type, block.content) _previous_failure = _failed_call_history.get(_call_signature) _blocked_failed_retry = bool( @@ -31439,6 +31454,8 @@ async def stream_agent_loop( ) # --- Tool budget check --- if max_tool_calls > 0 and total_tool_calls >= max_tool_calls: + if _runtime_action is not None: + _runtime_action.finish({'blocked': True, 'exit_code': 1, 'error': 'tool budget exceeded'}) yield f'data: {json.dumps({"type": "budget_exceeded", "limit": max_tool_calls, "used": total_tool_calls})}\n\n' budget_hit = True break @@ -31450,26 +31467,22 @@ async def stream_agent_loop( ) ): local_network_budget_hit = True + if _runtime_action is not None: + _runtime_action.finish({'blocked': True, 'exit_code': 1, 'error': 'network action budget exceeded'}) break if ( _tui_local_inspection_turn and total_tool_calls >= _TUI_LOCAL_INSPECTION_TOOL_CALL_CAP ): local_inspection_budget_hit = True + if _runtime_action is not None: + _runtime_action.finish({'blocked': True, 'exit_code': 1, 'error': 'inspection budget exceeded'}) break if local_inspection_budget_hit: break if not (_blocked_failed_retry or _blocked_redundant_read): total_tool_calls += 1 - native_call = converted_calls[i] if i < len(converted_calls) else None - tool_call_id = _resolved_tool_call_id( - native_call, - session_id=str(session_id or ""), - round_num=round_num, - tool_index=i, - tool_name=block.tool_type, - ) normalized_native_block = _normalize_native_tool_shell_wrapper(block, _last_user) if normalized_native_block != block: logger.info( @@ -32722,7 +32735,8 @@ async def stream_agent_loop( "error": "Web recovery action is repeated or exceeds the execution budget.", "output": _web_execution_budget.instruction(), } - return await execute_tool_block( + return await execute_action( + execute_tool_block, _runtime_action, block, session_id=session_id, disabled_tools=disabled_tools, @@ -33256,6 +33270,10 @@ async def stream_agent_loop( # Emit tool_output (include ui_event data if present) tool_output_data = {"type": "tool_output", "tool": block.tool_type, "command": cmd_display, "output": output_text, "exit_code": result.get("exit_code"), "execution_attempted": _execution_attempted, "blocked": bool(result.get("blocked", False))} + if _runtime_action is not None: + _runtime_action.normalize(block, 'agent_loop compatibility adapters') + _runtime_action.finish(result) + tool_output_data['action_receipt'] = _runtime_action.to_dict() # Keep exact arguments on email mutation events. The frontend uses # these UIDs to reconcile an agent cleanup immediately, even when # a provider returns only human-readable MCP text. diff --git a/src/agent_runtime/__init__.py b/src/agent_runtime/__init__.py new file mode 100644 index 000000000..c0d70533b --- /dev/null +++ b/src/agent_runtime/__init__.py @@ -0,0 +1 @@ +"""Run-scoped contracts behind the public agent-loop compatibility facade.""" diff --git a/src/agent_runtime/completion.py b/src/agent_runtime/completion.py new file mode 100644 index 000000000..121c095fa --- /dev/null +++ b/src/agent_runtime/completion.py @@ -0,0 +1,207 @@ +"""One presentation gate between agent execution and externally visible prose. + +Tool, progress and interaction events stay live. Answer deltas are held until +the generator unwinds so a later replacement cannot conceal an earlier false +claim. This consumes no provider calls. Cancellation closes the inner generator +under the same journal/turn authority; it never emits a successful terminal event. +""" +from __future__ import annotations + +from contextlib import aclosing +from dataclasses import replace +from functools import wraps +from inspect import signature +import json +import re +from time import perf_counter + +from src.agent_evidence import ( + CompletionDecision, CompletionStatus, EvidenceKind, EvidenceLedger, + requirements_from_runtime_context, +) +from .journal import ActionJournal, bind_journal, current_journal + + +_TEST_CLAIM = re.compile( + r'\b(?:(?:all\s+)?(?:tests?|checks?|verification|suite)\s+(?:have\s+|has\s+|now\s+|are\s+|is\s+)*(?:passed|passing|successful|green)|' + r'(?:passed|passing)\s+(?:all\s+)?(?:the\s+)?tests?|\d+\s+passed)\b', re.I) +_TEST_STATUS_CLAIM = re.compile( + r'\b(?:tests?|pytest|unittest|test suite|checks?|verification)\s*[:—-]?\s*' + r'(?:all\s+|have\s+|has\s+|now\s+|are\s+|is\s+|ran\s+)*' + r'(?:pass(?:ed|ing)?|succeeded|successful(?:ly)?|green)\b|' + r'\b(?:zero|no|0)\s+(?:test\s+)?failures\b', re.I) +_TERMINAL_SUCCESS = re.compile(r'^\s*(?:done|completed|success|all done|all set|fixed)\b', re.I) +_EXECUTION_CLAIM = re.compile( + r'\b(?:(?:I|we|I\'ve|we\'ve)\s+(?:have\s+)?(?:successfully\s+)?(?:ran|executed|tested|verified|created|updated|modified|wrote|saved|fixed|completed)|' + r'(?:file|artifact|command|script|service|server)\s+(?:was\s+|has\s+been\s+|is\s+)?(?:successfully\s+)?(?:created|updated|written|saved|executed|started)|' + r'(?:successfully\s+)(?:ran|executed|created|updated|saved|completed))\b', re.I) + + +def completion_answer(text: str, ledger: EvidenceLedger, decision: CompletionDecision) -> tuple[str, str]: + """Return the answer and a reason if unsupported execution claims were removed.""" + if decision.status == CompletionStatus.AWAITING_USER: + # A question may still falsely assert that preceding work passed. + unsupported = '' + elif not decision.can_complete: + unsupported = decision.reason + else: + unsupported = '' + if (_TEST_CLAIM.search(text) or _TEST_STATUS_CLAIM.search(text)) and decision.status != CompletionStatus.VERIFIED: + unsupported = unsupported or 'no current passing executable verification supports the claim' + productive = [event for event in ledger.events + if event.authoritative and event.success + and event.tool not in {'update_plan', 'todowrite', 'ask_user'}] + if (_EXECUTION_CLAIM.search(text) or _TERMINAL_SUCCESS.search(text)) and not productive: + unsupported = unsupported or 'no successful operation supports the execution claim' + if not unsupported: + # For a declared execution contract, publish facts selected from the + # receipts rather than an unconstrained model claim (test counts, + # coverage and "everything fixed" cannot be inferred from exit status). + if decision.can_complete and (ledger.requirements.required_artifacts or ledger.requirements.verifier_required): + parts = [] + if ledger.requirements.required_artifacts: + parts.append('Output available: ' + ', '.join(ledger.requirements.required_artifacts) + '.') + if decision.status == CompletionStatus.VERIFIED: + parts.append('The latest executable verification passed.') + elif any(e.kind == EvidenceKind.ARTIFACT_VALIDATION and e.authoritative and e.success for e in ledger.events): + parts.append('Artifact readback verified. No passing executable test result was recorded.') + else: + parts.append('No passing executable test result was recorded.') + return ' '.join(parts), '' + return text, '' + missing = (" Missing artifacts: " + ", ".join(decision.missing_artifacts) + "." + if decision.missing_artifacts else '') + return "The task is incomplete: " + unsupported.rstrip('.') + '.' + missing, unsupported + + +def _event(data: dict) -> str: + return 'data: ' + json.dumps(data) + '\n\n' + + +def with_completion_gate(func): + call_signature = signature(func) + + @wraps(func) + async def wrapped(*args, **kwargs): + started = perf_counter() + first_answer_at = None + arguments = call_signature.bind(*args, **kwargs) + arguments.apply_defaults() + bound = arguments.arguments + messages = bound.get('messages') or [] + instruction = next((m.get('content', '') for m in reversed(messages) + if m.get('role') == 'user' and isinstance(m.get('content'), str)), '') + context = bound.get('client_runtime_context') or {} + requirements = requirements_from_runtime_context(context, instruction=instruction) + from src.tool_execution import vet_workspace + # A completion declaration is not a filesystem permission. Only the + # explicit, vetted runtime workspace may be read for artifact versions. + trusted_workspace = vet_workspace(bound.get('workspace')) if bound.get('workspace') else '' + requirements = replace(requirements, workspace_root=trusted_workspace or '') + parent = current_journal() + journal = parent if parent is not None and parent.workspace == requirements.workspace_root else ActionJournal( + workspace=requirements.workspace_root, observed_artifacts=requirements.required_artifacts) + answer_events: list[dict] = [] + metrics_events: list[dict] = [] + answer = '' + has_final = False + done = False + awaiting = False + exhausted = False + provider_error = False + with bind_journal(journal): + async with aclosing(func(*args, **kwargs)) as stream: + async for chunk in stream: + if chunk.strip() == 'data: [DONE]': + done = True + continue + try: + data = json.loads(chunk[6:]) if chunk.startswith('data: ') else None + except (ValueError, TypeError): + data = None + if not isinstance(data, dict): + if chunk.startswith('event: error'): + provider_error = True + yield chunk + continue + kind = data.get('type') + if kind == 'completion_decision': + existing = data.get('data') or {} + awaiting |= existing.get('status') == 'awaiting_user' + exhausted |= existing.get('status') == 'exhausted' + continue + if kind in {'metrics', 'agent_terminal'}: + metrics_events.append(data) + declared = (data.get('data') or {}).get('completion_requirements') + awaiting |= bool((data.get('data') or {}).get('missing_workspace')) + if isinstance(declared, dict): + requirements = requirements_from_runtime_context({'completion_requirements': declared}) + requirements = replace(requirements, workspace_root=trusted_workspace or '') + continue + if kind == 'ask_user': + awaiting = True + payload = data.get('data') or {} + if isinstance(payload.get('question'), str): + current = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements) + question, why = completion_answer(payload['question'], current, current.evaluate(awaiting_user=True)) + if why: + data = {**data, 'data': {**payload, 'question': question}} + chunk = _event(data) + if kind == 'final_response': + if first_answer_at is None: + first_answer_at = perf_counter() + answer = str(data.get('content') or '') + has_final = True + answer_events.append(data) + continue + if 'delta' in data and not data.get('thinking'): + if first_answer_at is None: + first_answer_at = perf_counter() + if has_final: + answer = '' + has_final = False + answer += str(data.get('delta') or '') + answer_events.append(data) + continue + yield chunk + if provider_error and not answer_events and not metrics_events: + return + ledger = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements) + decision = ledger.evaluate(exhausted=exhausted, awaiting_user=awaiting) + # Exhaustion limits execution; factual source synthesis can remain + # useful and must not be replaced merely because the budget ended. + presentation_decision = ledger.evaluate(awaiting_user=awaiting) if exhausted else decision + safe_answer, reason = completion_answer(answer, ledger, presentation_decision) + if reason and decision.can_complete: + decision = CompletionDecision(CompletionStatus.UNVERIFIED, False, reason, + decision.evidence_ids, decision.missing_artifacts) + released_at = perf_counter() + yield _event({'type': 'completion_decision', 'data': decision.to_dict()}) + # Evaluate each earlier draft as well as the final replacement. + # Never replay an unsupported intermediate success claim. + draft = ''.join(str(e.get('delta') or e.get('content') or '') for e in answer_events) + _, unsafe_draft = completion_answer(draft, ledger, presentation_decision) + replaced_answer = bool(reason or unsafe_draft or safe_answer != answer) + if replaced_answer: + yield _event({'type': 'final_response', 'content': safe_answer}) + else: + for event in answer_events: + yield _event(event) + for event in metrics_events: + metadata = event.setdefault('data', {}) + metadata.update(completion_decision=decision.to_dict(), evidence_events=ledger.to_list(), + action_receipts=journal.to_list(), completion_requirements=requirements.to_dict()) + metadata['completion_gate'] = { + 'buffer_seconds': released_at - first_answer_at if first_answer_at is not None else 0, + 'first_visible_answer_seconds': released_at - started, + 'additional_provider_calls': 0, + 'answer_replaced': replaced_answer, + } + if replaced_answer: + metadata['round_texts'] = [safe_answer] + metadata['completion_gate_reason'] = reason or unsafe_draft or 'receipt_summary' + yield _event(event) + if done: + yield 'data: [DONE]\n\n' + + return wrapped diff --git a/src/agent_runtime/identity.py b/src/agent_runtime/identity.py new file mode 100644 index 000000000..49381cde5 --- /dev/null +++ b/src/agent_runtime/identity.py @@ -0,0 +1,146 @@ +"""Canonical evidence identities; these helpers never grant filesystem access.""" +from __future__ import annotations + +import hashlib +import json +import os +from pathlib import Path, PurePosixPath +import re +import shlex +import stat +from .path_policy import _is_sensitive_path + + +def digest(value: object) -> str: + return hashlib.sha256(json.dumps(value, sort_keys=True, ensure_ascii=False, + separators=(",", ":"), default=str).encode()).hexdigest() + + +def artifact_version(value: str, workspace: str) -> str: + """Content version for a confined declared output; missing/unreadable is explicit.""" + identity = artifact_identity(value, workspace) + if not workspace or not identity.startswith('workspace:'): + return 'unobserved' + root = Path(workspace).resolve() + candidate = (root / identity.removeprefix('workspace:')).resolve() + if not candidate.is_relative_to(root) or _is_sensitive_path(str(candidate)): + return 'unobserved' + if not hasattr(os, 'O_NOFOLLOW') or not os.supports_dir_fd: + return 'unobserved' + directory = None + try: + directory = os.open(root, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW) + parts = candidate.relative_to(root).parts + if not parts: + return 'unobserved' + for part in parts[:-1]: + child = os.open(part, os.O_RDONLY | os.O_DIRECTORY | os.O_NOFOLLOW, dir_fd=directory) + os.close(directory) + directory = child + descriptor = os.open(parts[-1], os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK, dir_fd=directory) + with os.fdopen(descriptor, 'rb') as stream: + info = os.fstat(stream.fileno()) + if not stat.S_ISREG(info.st_mode) or info.st_size > 64 * 1024 * 1024: + return 'unobserved' + result = hashlib.sha256() + remaining = 64 * 1024 * 1024 + while block := stream.read(min(1024 * 1024, remaining + 1)): + remaining -= len(block) + if remaining < 0: + return 'unobserved' + result.update(block) + return result.hexdigest() + except (OSError, ValueError): + return 'missing-or-unreadable' + finally: + if directory is not None: + os.close(directory) + + +def artifact_identity(value: str, workspace: str = "") -> str: + """Unify relative, virtual and host aliases without basename matching. + + Resolving symlinks is evidence bookkeeping, never a confinement check. Paths + outside the workspace retain their absolute identity and cannot satisfy a + workspace obligation with the same basename. + """ + text = str(value or "").strip() + if not text: + return "" + path = PurePosixPath(text) + root = Path(workspace or "/workspace").resolve() + if path.parts[:2] == ('/', 'workspace'): + path = PurePosixPath(*path.parts[2:]) + candidate = Path(str(path)) + if not candidate.is_absolute(): + candidate = root / candidate + try: + resolved = candidate.resolve() + relative = resolved.relative_to(root) + return "workspace:" + relative.as_posix() + except ValueError: + return "absolute:" + str(candidate.resolve()) + except (OSError, RuntimeError): + return "unresolved:" + text + + +def executable_words(command: str) -> tuple[str, ...]: + """Recognize one foreground command, optionally after safe cd/set prefixes. + + This is deliberately conservative evidence parsing, not shell authorization. + Pipelines, control flow, substitutions and status-masking tails are not proof + that a verifier returned the recorded shell status. + """ + text = str(command or '').strip() + if any(marker in text for marker in ('`', '$(', '${', '\n', '\r')): + return () + try: + lexer = shlex.shlex(text, posix=True, punctuation_chars=';&|<>()') + lexer.whitespace_split = True + words = list(lexer) + except ValueError: + return () + while '&&' in words: + index = words.index('&&') + prefix = words[:index] + if not ((len(prefix) == 2 and prefix[0] == 'cd') or prefix == ['set', '-e']): + return () + words = words[index + 1:] + if any(word and all(c in ';&|<>()' for c in word) for word in words): + return () + while words and re.fullmatch(r'[A-Za-z_][A-Za-z0-9_]*=[^\n]*', words[0]): + words.pop(0) + return tuple(words) + + +def is_test_command(command: str) -> bool: + words = executable_words(command) + if not words: + return False + if any(word in {'--help', '-h', '--version', '--collect-only', '--co'} for word in words[1:]): + return False + binary = Path(words[0]).name + if binary in {'pytest', 'py.test'}: + return True + if re.fullmatch(r'python(?:\d+(?:\.\d+)?)?', binary): + args = list(words[1:]) + while args and args[0] in {'-I', '-S', '-s', '-E', '-B', '-u'}: + args.pop(0) + return len(args) >= 2 and args[:2] in (['-m', 'pytest'], ['-m', 'unittest']) + if binary in {'npm', 'pnpm', 'yarn', 'make', 'cargo', 'go'}: + args = words[1:] + return bool(args and (args[0] == 'test' or binary == 'npm' and args[:2] == ('run', 'test'))) + return bool(re.match(r'^/(?:tests?|verifier)/[^/]+', words[0])) + + +def is_validation_command(command: str) -> bool: + words = executable_words(command) + if not words: + return False + binary = Path(words[0]).name + return (is_test_command(command) + or binary in {'cat', 'head', 'tail', 'stat', 'wc', 'jq', 'cmp', 'diff', + 'coqc', 'gcc', 'g++', 'clang', 'clang++', 'javac', 'rustc'} + or binary == 'test' and len(words) > 1 and words[1] in {'-e', '-f', '-s', '-d'} + or binary in {'cargo', 'go', 'npm', 'pnpm', 'yarn'} and words[1:2] in {('build',), ('check',)} + or binary == 'npm' and words[1:3] == ('run', 'build')) diff --git a/src/agent_runtime/journal.py b/src/agent_runtime/journal.py new file mode 100644 index 000000000..8e8a118b5 --- /dev/null +++ b/src/agent_runtime/journal.py @@ -0,0 +1,205 @@ +"""Run-owned action history. Model text cannot insert authoritative receipts.""" +from __future__ import annotations + +from contextlib import contextmanager +from contextvars import ContextVar +from copy import deepcopy +from dataclasses import dataclass, field, asdict +from functools import wraps +from inspect import signature +from typing import Any +from uuid import uuid4 + +from .identity import artifact_identity, artifact_version, digest + + +@dataclass +class ActionReceipt: + action_id: str + call_id: str + proposed_tool: str + proposed_arguments: str + provider_arguments: Any = None + provider_tool: str = '' + tool: str = "" + arguments: str = "" + transitions: list[dict[str, Any]] = field(default_factory=list) + execution_id: str | None = None + operation_started: bool = False + outcome: dict[str, Any] | None = None + artifact_versions: dict[str, str] = field(default_factory=dict) + artifact_changes: list[str] | None = None + + def transition(self, stage: str, **details: Any) -> None: + self.transitions.append({'sequence': len(self.transitions), 'stage': stage, **details}) + + def normalize(self, block: Any, reason: str) -> None: + tool, arguments = str(block.tool_type), str(block.content) + if tool != self.tool or arguments != self.arguments or not any(t['stage'] == 'normalized' for t in self.transitions): + self.transition('normalized', reason=reason, tool=tool, arguments=arguments, + previous_sha256=digest((self.tool, self.arguments))) + self.tool, self.arguments = tool, arguments + + def finish(self, result: dict[str, Any]) -> None: + if self.outcome is not None: + return + code = result.get('exit_code') + valid_code = isinstance(code, int) and not isinstance(code, bool) + denied = bool(result.get('blocked') or result.get('approval_required') + or str(result.get('failure_kind', '')).endswith('_denied')) + self.outcome = { + 'exit_code': code if valid_code else None, + 'success': valid_code and code == 0 and not result.get('error') and not denied, + 'authoritative': self.execution_id is not None and valid_code and not denied, + 'blocked': denied, + 'output_sha256': digest(result.get('output') or result.get('error') or result.get('stdout') or ''), + } + self.transition('outcome', **self.outcome) + + def to_dict(self) -> dict[str, Any]: + return asdict(self) + + +@dataclass +class ActionJournal: + run_id: str = field(default_factory=lambda: uuid4().hex) + actions: list[ActionReceipt] = field(default_factory=list) + workspace: str = '' + observed_artifacts: tuple[str, ...] = () + + def capture_versions(self, action: ActionReceipt) -> None: + if self.workspace: + action.artifact_versions = { + artifact_identity(path, self.workspace): artifact_version(path, self.workspace) + for path in self.observed_artifacts + } + + def propose(self, block: Any, call_id: str = '', native_call: dict | None = None) -> ActionReceipt: + native = native_call or {} + function = native.get('function') or native + if not isinstance(function, dict): + function = {} + action = ActionReceipt( + action_id=f'{self.run_id}:action:{len(self.actions) + 1}', call_id=call_id, + proposed_tool=str(block.tool_type), proposed_arguments=str(block.content), + provider_arguments=deepcopy(function.get('arguments')), + provider_tool=str(function.get('name') or ''), + tool=str(block.tool_type), arguments=str(block.content), + ) + action.transition('proposed') + self.actions.append(action) + return action + + def to_list(self) -> list[dict[str, Any]]: + return [action.to_dict() for action in self.actions] + + def evidence_events(self) -> list[dict[str, Any]]: + return [dict(tool=a.tool, command=a.arguments, + exit_code=(a.outcome or {}).get('exit_code'), + error=not (a.outcome or {}).get('success'), + execution_attempted=bool((a.outcome or {}).get('authoritative')), + blocked=(a.outcome or {}).get('blocked', False), + action_id=a.action_id, execution_id=a.execution_id, + artifact_versions=a.artifact_versions, artifact_changes=a.artifact_changes) + for a in self.actions if a.outcome is not None] + + +_JOURNAL: ContextVar[ActionJournal | None] = ContextVar('runtime_action_journal', default=None) +_ACTION: ContextVar[ActionReceipt | None] = ContextVar('runtime_current_action', default=None) + + +@contextmanager +def bind_journal(journal: ActionJournal): + token = _JOURNAL.set(journal) + try: + yield journal + finally: + _JOURNAL.reset(token) + + +def current_journal() -> ActionJournal | None: + return _JOURNAL.get() + + +def propose_action(block: Any, call_id: str = '', native_call: dict | None = None) -> ActionReceipt | None: + journal = _JOURNAL.get() + return journal.propose(block, call_id, native_call) if journal else None + + +def mark_authorized() -> None: + action = _ACTION.get() + if action is not None and not any(t['stage'] == 'authorized' for t in action.transitions): + action.transition('authorized', authority='existing_dispatcher_policy') + + +def mark_dispatch() -> None: + action = _ACTION.get() + if action is not None and action.execution_id is None: + mark_authorized() + action.execution_id = action.action_id + ':execution:1' + action.transition('dispatched', execution_id=action.execution_id) + + +async def dispatched(operation): + """Record an actual backend invocation, distinct from router admission.""" + mark_dispatch() + return await operation + + +def mark_operation_started(backend: str, **details: Any) -> None: + action = _ACTION.get() + if action is not None: + action.operation_started = True + action.transition('operation_started', backend=backend, **details) + + +async def execute_action(executor, action: ActionReceipt | None, block: Any, **kwargs): + """Adapter binds the proposal across async tool-task execution and cleanup.""" + if action is not None: + action.normalize(block, 'agent_loop compatibility adapters') + token = _ACTION.set(action) + try: + return await executor(block, **kwargs) + finally: + _ACTION.reset(token) + + +def record_action(func): + call_signature = signature(func) + + @wraps(func) + async def wrapped(*args, **kwargs): + bound = call_signature.bind(*args, **kwargs) + block = bound.arguments['block'] + action = _ACTION.get() or propose_action(block) + token = _ACTION.set(action) + try: + journal = current_journal() + before = {} + if action is not None: + action.normalize(block, 'dispatcher input') + if journal is not None and journal.workspace: + journal.capture_versions(action) + before = dict(action.artifact_versions) + description, result = await func(*args, **kwargs) + if action is not None: + journal = current_journal() + if journal is not None: + journal.capture_versions(action) + if journal.workspace: + action.artifact_changes = [key for key, value in action.artifact_versions.items() + if before.get(key) != value] + if 'BLOCKED' in description and action.execution_id is None: + action.transition('authorization_denied', reason=str(result.get('error', ''))) + action.finish({**result, 'blocked': True}) + else: + action.finish(result) + return description, result + except BaseException as exc: + if action is not None: + action.transition('interrupted', category=type(exc).__name__) + raise + finally: + _ACTION.reset(token) + + return wrapped diff --git a/src/agent_runtime/path_policy.py b/src/agent_runtime/path_policy.py new file mode 100644 index 000000000..d03d2d83f --- /dev/null +++ b/src/agent_runtime/path_policy.py @@ -0,0 +1,26 @@ +"""Existing sensitive-path policy shared by tools and evidence observation. + +This is a deny predicate, not an authorization grant or a workspace scope. +""" +import os + +_SENSITIVE_BASENAMES: set[str] = { + ".ssh", ".gnupg", ".gitconfig", + ".bashrc", ".bash_profile", ".bash_logout", + ".zshrc", ".zprofile", ".zshenv", + ".profile", ".tcshrc", ".cshrc", ".env", ".netrc", +} +_SENSITIVE_FILE_PATTERNS: tuple[str, ...] = ( + "authorized_keys", "id_rsa", "id_ed25519", "id_ecdsa", + "known_hosts", "auth.json", "app.db", "settings.json", +) +_SENSITIVE_BASENAMES_CF = frozenset(b.casefold() for b in _SENSITIVE_BASENAMES) +_SENSITIVE_FILE_PATTERNS_CF = frozenset(p.casefold() for p in _SENSITIVE_FILE_PATTERNS) + + +def _is_sensitive_path(resolved: str) -> bool: + # Case folding is required even on POSIX: default macOS volumes are + # case insensitive but os.path.normcase there does not fold path names. + parts = [p.casefold() for p in resolved.split(os.sep)] + filename = parts[-1] if parts else "" + return any(part in _SENSITIVE_BASENAMES_CF for part in parts) or filename in _SENSITIVE_FILE_PATTERNS_CF diff --git a/src/agent_tools/subprocess_tools.py b/src/agent_tools/subprocess_tools.py index c8ffe9ddb..7537c82b5 100644 --- a/src/agent_tools/subprocess_tools.py +++ b/src/agent_tools/subprocess_tools.py @@ -16,6 +16,7 @@ from urllib.parse import urlparse import httpx from src.constants import MAX_OUTPUT_CHARS +from src.agent_runtime.journal import mark_operation_started # Agent shell calls must fail fast enough for the loop to recover and choose a # better tool. A one-hour default can pin an entire benchmark worker on an @@ -690,6 +691,7 @@ class BashTool: ) except RuntimeError as exc: return {"error": str(exc), "exit_code": 1} + mark_operation_started('subprocess', pid=proc.pid) stdout, stderr, rc, timed_out = await _run_subprocess_streaming( proc, timeout=DEFAULT_BASH_TIMEOUT, @@ -1015,6 +1017,7 @@ class PythonTool: env=_subproc_env, cwd=agent_cwd(), ) + mark_operation_started('subprocess', pid=proc.pid) stdout, stderr, rc, timed_out = await _run_subprocess_streaming( proc, timeout=DEFAULT_PYTHON_TIMEOUT, diff --git a/src/tool_execution.py b/src/tool_execution.py index 57707a693..451d5c7a8 100644 --- a/src/tool_execution.py +++ b/src/tool_execution.py @@ -727,49 +727,13 @@ async def _route_tool_via_bridge(tool: str, content: str, session_id: Optional[s # "tool_path_extra_roots" setting (list of path strings). # --------------------------------------------------------------------------- -_SENSITIVE_BASENAMES: set[str] = { - ".ssh", ".gnupg", ".gitconfig", - ".bashrc", ".bash_profile", ".bash_logout", - ".zshrc", ".zprofile", ".zshenv", - ".profile", ".tcshrc", ".cshrc", - ".env", ".netrc", -} - -_SENSITIVE_FILE_PATTERNS: tuple[str, ...] = ( - "authorized_keys", "id_rsa", "id_ed25519", "id_ecdsa", - "known_hosts", "auth.json", "app.db", "settings.json", +# Compatibility exports: the same deny predicate protects tool access and +# artifact observations, so hashing cannot become a sensitive-file side channel. +from src.agent_runtime.path_policy import ( + _SENSITIVE_BASENAMES, _SENSITIVE_FILE_PATTERNS, + _SENSITIVE_BASENAMES_CF, _SENSITIVE_FILE_PATTERNS_CF, _is_sensitive_path, ) -# Case-folded views used for matching. On a case-insensitive filesystem -# (Windows, default macOS) ".SSH/AUTHORIZED_KEYS" and ".env" resolve to the -# same protected files as their lowercase forms, so the deny-list has to fold -# case before comparing — the sibling resolver already normcases paths for the -# same reason. casefold (not os.path.normcase) because normcase is a no-op on -# POSIX, which is exactly where the macOS read-exfil path lives. -_SENSITIVE_BASENAMES_CF: frozenset[str] = frozenset(b.casefold() for b in _SENSITIVE_BASENAMES) -_SENSITIVE_FILE_PATTERNS_CF: frozenset[str] = frozenset(p.casefold() for p in _SENSITIVE_FILE_PATTERNS) - - -def _is_sensitive_path(resolved: str) -> bool: - """Return True if *resolved* falls under a sensitive directory or - matches a sensitive filename — regardless of what root it sits under. - - Matching is case-insensitive: on Windows / default macOS a case-variant - name (``.SSH``, ``AUTHORIZED_KEYS``, ``Id_Rsa``) points at the same file as - the lowercase form, so a case-sensitive check would let it slip past the - deny-list in every file tool that relies on it. - """ - parts = [p.casefold() for p in resolved.split(os.sep)] - filename = parts[-1] if parts else "" - - # Check if any path component is a sensitive directory. - for part in parts: - if part in _SENSITIVE_BASENAMES_CF: - return True - - # Check filename against known sensitive files. - return filename in _SENSITIVE_FILE_PATTERNS_CF - def _tool_path_roots() -> list[str]: """Return the list of directory roots that read_file / write_file @@ -1290,6 +1254,10 @@ async def _document_tool_dispatch( # Dispatcher # --------------------------------------------------------------------------- +from src.agent_runtime.journal import dispatched, mark_authorized, mark_dispatch, record_action + + +@record_action async def execute_tool_block( block: Any, session_id: Optional[str] = None, @@ -1615,14 +1583,15 @@ async def _execute_tool_block_impl( if rejected is not None: return rejected + mark_authorized() if bridge_owns_tool: try: - return await execution_bridge.route_tool( + return await dispatched(execution_bridge.route_tool( tool, content, session_id, client_runtime_context, - ) + )) except asyncio.CancelledError: raise except Exception as exc: @@ -1643,7 +1612,7 @@ async def _execute_tool_block_impl( ) if tool in _ROUTED_BRIDGE_TOOLS and _client_bridge(client_runtime_context) is not None: - return await _route_tool_via_bridge(tool, content, session_id, client_runtime_context) + return await dispatched(_route_tool_via_bridge(tool, content, session_id, client_runtime_context)) # Background execution: a `bash` block whose first line is the `#!bg` # marker runs DETACHED — returns a job id immediately so the chat stream @@ -1653,6 +1622,7 @@ async def _execute_tool_block_impl( _is_bg, _bg_cmd = _split_bg_marker(content) if _is_bg and _bg_cmd: from src import bg_jobs + mark_dispatch() rec = bg_jobs.launch(_bg_cmd, session_id=session_id, cwd=agent_cwd()) short = _bg_cmd.strip().split(chr(10))[0][:80] desc = f"bash (background): {short}" @@ -1678,37 +1648,37 @@ async def _execute_tool_block_impl( if tool in _MCP_TOOL_MAP: first_line = content.split(chr(10))[0][:80] desc = f"{tool}: {first_line}" - result = await _call_mcp_tool(tool, content, progress_cb=progress_cb) + result = await dispatched(_call_mcp_tool(tool, content, progress_cb=progress_cb)) elif tool in ("grep", "glob", "ls", "get_workspace", "host_shell"): # Code-navigation tools — no MCP server; run the direct implementation. first_line = content.split(chr(10))[0][:80] desc = f"{tool}: {first_line}" - result = await _direct_fallback( + result = await dispatched(_direct_fallback( tool, content, progress_cb=progress_cb, owner=owner, client_runtime_context=client_runtime_context, - ) \ + )) \ or {"error": f"{tool}: execution failed", "exit_code": 1} elif tool == "apply_patch" and _tui_host_bridge_patch_url(client_runtime_context): first_line = content.split(chr(10))[0][:80] desc = f"{tool}: {first_line}" if first_line else tool - result = await _apply_patch_via_tui_host_bridge(content, client_runtime_context) + result = await dispatched(_apply_patch_via_tui_host_bridge(content, client_runtime_context)) elif tool in ("apply_patch", "todowrite"): first_line = content.split(chr(10))[0][:80] desc = f"{tool}: {first_line}" if first_line else tool - result = await _direct_fallback(tool, content, session_id=session_id, owner=owner) \ + result = await dispatched(_direct_fallback(tool, content, session_id=session_id, owner=owner)) \ or {"error": f"{tool}: execution failed", "exit_code": 1} elif tool == "manage_bg_jobs": # Inspect/kill detached `bash` jobs; needs session_id to scope to chat. desc = f"manage_bg_jobs: {content.split(chr(10))[0][:80]}" - result = await _direct_fallback(tool, content, session_id=session_id, owner=owner) \ + result = await dispatched(_direct_fallback(tool, content, session_id=session_id, owner=owner)) \ or {"error": "manage_bg_jobs: execution failed", "exit_code": 1} elif tool in ("create_document", "update_document", "edit_document", "suggest_document", "manage_documents"): desc = f"{tool}: {content.split(chr(10))[0][:80]}" - result = await _document_tool_dispatch( + result = await dispatched(_document_tool_dispatch( tool, content, session_id, @@ -1716,14 +1686,14 @@ async def _execute_tool_block_impl( document_id=approved_document_id or active_document_id, document_version=approved_document_version, document_digest=approved_document_digest, - ) \ + )) \ or {"error": f"{tool}: execution failed", "exit_code": 1} if tool in ("edit_document", "suggest_document") and "title" in (result or {}): desc = f"{tool}: {result.get('title', '')}" elif tool == "search_chats": query = content.split("\n")[0].strip() desc = f"search_chats: {query[:80]}" - result = await do_search_chats(query, owner=owner) + result = await dispatched(do_search_chats(query, owner=owner)) elif tool in ("chat_with_model", "ask_teacher", "list_models"): # Migrated to the agent_tools registry (#3629): dispatched through # TOOL_HANDLERS with the owner/session ctx these tools need, instead @@ -1731,7 +1701,7 @@ async def _execute_tool_block_impl( # src/agent_tools/model_interaction_tools.py. first_line = content.split(chr(10))[0].strip()[:60] desc = f"{tool}: {first_line}" if first_line else tool - result = await _document_tool_dispatch(tool, content, session_id, owner) \ + result = await dispatched(_document_tool_dispatch(tool, content, session_id, owner)) \ or {"error": f"{tool}: execution failed", "exit_code": 1} elif tool in ("create_session", "list_sessions", "send_to_session", "manage_session"): # Migrated to the agent_tools registry (#3629): dispatched through @@ -1739,101 +1709,101 @@ async def _execute_tool_block_impl( # live in src/agent_tools/session_tools.py. first_line = content.split(chr(10))[0].strip()[:60] desc = f"{tool}: {first_line}" if first_line else tool - result = await _document_tool_dispatch(tool, content, session_id, owner) \ + result = await dispatched(_document_tool_dispatch(tool, content, session_id, owner)) \ or {"error": f"{tool}: execution failed", "exit_code": 1} elif tool in ("pipeline", "manage_memory", "ui_control"): from src.ai_interaction import dispatch_ai_tool - desc, result = await dispatch_ai_tool(tool, content, session_id, owner=owner) + desc, result = await dispatched(dispatch_ai_tool(tool, content, session_id, owner=owner)) elif tool == "manage_tasks": desc = "manage_tasks" - result = await do_manage_tasks(content, owner=owner) + result = await dispatched(do_manage_tasks(content, owner=owner)) elif tool == "manage_skills": desc = "manage_skills" - result = await do_manage_skills(content, owner=owner) + result = await dispatched(do_manage_skills(content, owner=owner)) elif tool == "api_call": first_line = content.split("\n")[0].strip()[:60] desc = f"api_call: {first_line}" - result = await do_api_call(content) + result = await dispatched(do_api_call(content)) elif tool in ("manage_endpoints", "manage_mcp", "manage_webhooks", "manage_tokens", "manage_settings"): # Registry-dispatched (agent_tools.admin_tools); owner threaded for ownership/admin checks. desc = tool - result = await _direct_fallback(tool, content, owner=owner) \ + result = await dispatched(_direct_fallback(tool, content, owner=owner)) \ or {"error": f"{tool}: execution failed", "exit_code": 1} elif tool == "manage_notes": desc = "manage_notes" - result = await do_manage_notes(content, owner=owner) + result = await dispatched(do_manage_notes(content, owner=owner)) elif tool == "manage_calendar": desc = "manage_calendar" - result = await do_manage_calendar(content, owner=owner) + result = await dispatched(do_manage_calendar(content, owner=owner)) elif tool == "download_model": desc = "download_model" - result = await do_download_model(content, owner=owner) + result = await dispatched(do_download_model(content, owner=owner)) elif tool == "serve_model": desc = "serve_model" - result = await do_serve_model(content, owner=owner) + result = await dispatched(do_serve_model(content, owner=owner)) elif tool == "list_served_models": desc = "list_served_models" - result = await do_list_served_models(content, owner=owner) + result = await dispatched(do_list_served_models(content, owner=owner)) elif tool == "stop_served_model": desc = "stop_served_model" - result = await do_stop_served_model(content, owner=owner) + result = await dispatched(do_stop_served_model(content, owner=owner)) elif tool == "tail_serve_output": desc = "tail_serve_output" - result = await do_tail_serve_output(content, owner=owner) + result = await dispatched(do_tail_serve_output(content, owner=owner)) elif tool == "list_downloads": desc = "list_downloads" - result = await do_list_downloads(content, owner=owner) + result = await dispatched(do_list_downloads(content, owner=owner)) elif tool == "cancel_download": desc = "cancel_download" - result = await do_cancel_download(content, owner=owner) + result = await dispatched(do_cancel_download(content, owner=owner)) elif tool == "search_hf_models": desc = "search_hf_models" - result = await do_search_hf_models(content, owner=owner) + result = await dispatched(do_search_hf_models(content, owner=owner)) elif tool == "list_cached_models": desc = "list_cached_models" - result = await do_list_cached_models(content, owner=owner) + result = await dispatched(do_list_cached_models(content, owner=owner)) elif tool == "app_api": desc = "app_api" - result = await do_app_api(content, owner=owner) + result = await dispatched(do_app_api(content, owner=owner)) elif tool == "list_serve_presets": desc = "list_serve_presets" - result = await do_list_serve_presets(content, owner=owner) + result = await dispatched(do_list_serve_presets(content, owner=owner)) elif tool == "serve_preset": desc = "serve_preset" - result = await do_serve_preset(content, owner=owner) + result = await dispatched(do_serve_preset(content, owner=owner)) elif tool == "adopt_served_model": desc = "adopt_served_model" - result = await do_adopt_served_model(content, owner=owner) + result = await dispatched(do_adopt_served_model(content, owner=owner)) elif tool == "list_cookbook_servers": desc = "list_cookbook_servers" - result = await do_list_cookbook_servers(content, owner=owner) + result = await dispatched(do_list_cookbook_servers(content, owner=owner)) elif tool == "edit_image": desc = "edit_image" - result = await do_edit_image(content, owner=owner) + result = await dispatched(do_edit_image(content, owner=owner)) elif tool == "edit_file": - result = await _direct_fallback(tool, content) or {"error": "edit failed", "exit_code": 1} + result = await dispatched(_direct_fallback(tool, content)) or {"error": "edit failed", "exit_code": 1} desc = result.get("output") or result.get("error") or "edit_file" elif tool == "trigger_research": desc = "trigger_research" - result = await do_trigger_research(content, owner=owner, chat_session_id=session_id) + result = await dispatched(do_trigger_research(content, owner=owner, chat_session_id=session_id)) elif tool == "manage_research": desc = "manage_research" - result = await do_manage_research(content, owner=owner) + result = await dispatched(do_manage_research(content, owner=owner)) elif tool == "resolve_contact": desc = "resolve_contact" - result = await do_resolve_contact(content, owner=owner) + result = await dispatched(do_resolve_contact(content, owner=owner)) elif tool == "manage_contact": desc = "manage_contact" - result = await do_manage_contact(content, owner=owner) + result = await dispatched(do_manage_contact(content, owner=owner)) elif tool == "vault_search": desc = "vault_search" - result = await do_vault_search(content, owner=owner) + result = await dispatched(do_vault_search(content, owner=owner)) elif tool == "vault_get": desc = "vault_get" - result = await do_vault_get(content, owner=owner) + result = await dispatched(do_vault_get(content, owner=owner)) elif tool == "vault_unlock": desc = "vault_unlock" - result = await do_vault_unlock(content, owner=owner) + result = await dispatched(do_vault_unlock(content, owner=owner)) elif tool in BUILTIN_EMAIL_TOOLS: # Bare email tool name from fenced-block models (e.g. Ollama) — route to MCP email server. # Non-admin owners never reach here: BUILTIN_EMAIL_TOOLS ⊆ NON_ADMIN_BLOCKED_TOOLS, @@ -1879,7 +1849,7 @@ async def _execute_tool_block_impl( if session_id: args = dict(args) args[_EMAIL_MCP_SESSION_ARG] = session_id - result = await mcp.call_tool(qualified, args) + result = await dispatched(mcp.call_tool(qualified, args)) else: result = {"error": "MCP manager not available", "exit_code": 1} elif tool.startswith("mcp__"): @@ -1898,7 +1868,7 @@ async def _execute_tool_block_impl( if session_id: args = dict(args) args[_EMAIL_MCP_SESSION_ARG] = session_id - result = _normalize_mcp_text_error(await mcp.call_tool(tool, args)) + result = _normalize_mcp_text_error(await dispatched(mcp.call_tool(tool, args))) else: desc = f"mcp: {tool}" result = {"error": "MCP manager not available", "exit_code": 1} @@ -1907,14 +1877,14 @@ async def _execute_tool_block_impl( elif tool in dynamic_handlers: first_line = content.split(chr(10))[0][:80] desc = f"registry: {tool} {first_line}".strip() - res = await _direct_fallback( + res = await dispatched(_direct_fallback( tool, content, progress_cb=progress_cb, session_id=session_id, owner=owner, client_runtime_context=client_runtime_context, - ) + )) if isinstance(res, tuple): desc, result = res diff --git a/tests/runtime_evidence_helpers.py b/tests/runtime_evidence_helpers.py new file mode 100644 index 000000000..f5aa483f3 --- /dev/null +++ b/tests/runtime_evidence_helpers.py @@ -0,0 +1,10 @@ +"""Dispatcher doubles must simulate the receipt boundary as well as the result.""" +from src.agent_runtime.journal import mark_dispatch, record_action + + +def authoritative_executor(function): + @record_action + async def execute(block, *args, **kwargs): + mark_dispatch() + return await function(block, *args, **kwargs) + return execute diff --git a/tests/test_agent_evidence_loop.py b/tests/test_agent_evidence_loop.py index 5b7acfde8..b9350edf6 100644 --- a/tests/test_agent_evidence_loop.py +++ b/tests/test_agent_evidence_loop.py @@ -4,6 +4,7 @@ import json import src.agent_loop as agent_loop from src.tool_parsing import ToolBlock from src.tool_capabilities import ToolGateDecision +from tests.runtime_evidence_helpers import authoritative_executor def _events(chunks): @@ -83,7 +84,7 @@ def _patch_loop(monkeypatch, responses, captured_kwargs=None): yield f'data: {json.dumps({"delta": response})}\n\n' yield "data: [DONE]\n\n" - monkeypatch.setattr(agent_loop, "execute_tool_block", execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(execute)) monkeypatch.setattr(agent_loop, "stream_llm_with_fallback", stream) return lambda: call_index @@ -117,14 +118,14 @@ def test_failed_workspace_mutation_attempts_are_not_hidden_by_successful_probe() assert agent_loop._failed_workspace_mutation_attempts([failed, probe], records) == 1 -def test_terminal_completion_repairs_missing_artifact_at_most_twice(monkeypatch): +def test_terminal_completion_missing_artifact_does_not_add_model_rounds(monkeypatch): calls = _patch_loop(monkeypatch, ["Done without writing anything."]) events = _run("Write answer.json", max_rounds=4) blocked = [event for event in events if event.get("type") == "completion_blocked"] - assert [event["attempt"] for event in blocked] == [1, 2] - assert calls() == 3 + assert blocked == [] + assert calls() == 1 decision = next(event["data"] for event in events if event.get("type") == "completion_decision") assert decision["status"] == "blocked" assert decision["missing_artifacts"] == ["answer.json"] @@ -145,7 +146,7 @@ def test_failed_trailing_tool_with_planning_prose_continues_artifact_task(monkey "exit_code": 1, } - monkeypatch.setattr(agent_loop, "execute_tool_block", fail_execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(fail_execute)) events = _run( "Create answer.json after inspecting the source", @@ -155,8 +156,8 @@ def test_failed_trailing_tool_with_planning_prose_continues_artifact_task(monkey assert calls() > 1 assert any( - event.get("type") == "completion_blocked" - and event.get("decision", {}).get("missing_artifacts") == ["answer.json"] + event.get("type") == "completion_decision" + and event.get("data", {}).get("missing_artifacts") == ["answer.json"] for event in events ) @@ -178,7 +179,7 @@ def test_exact_failed_call_is_blocked_across_planning_and_intervening_failure(mo executed.append(block.content) return block.tool_type, {"output": f"failed: {block.content}", "exit_code": 1} - monkeypatch.setattr(agent_loop, "execute_tool_block", fail_execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(fail_execute)) events = _run( "Create /tmp_workspace/results after classifying the files", @@ -218,7 +219,7 @@ def test_exact_failed_call_can_retry_after_successful_workspace_mutation(monkeyp return block.tool_type, {"output": "written", "exit_code": 0} return block.tool_type, {"output": "classifier failed", "exit_code": 1} - monkeypatch.setattr(agent_loop, "execute_tool_block", execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(execute)) events = _run( "Create /tmp_workspace/results after repairing and running the classifier", @@ -260,7 +261,7 @@ def test_terminal_artifact_task_repairs_after_consecutive_failed_batches(monkeyp "exit_code": 1, } - monkeypatch.setattr(agent_loop, "execute_tool_block", execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(execute)) events = _run( "Create answer.json and verify it", @@ -312,7 +313,7 @@ def test_varied_failed_artifact_mutations_have_cumulative_cap(monkeypatch): "exit_code": 1, } - monkeypatch.setattr(agent_loop, "execute_tool_block", execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(execute)) events = _run( "Create answer.json and verify it", @@ -355,7 +356,7 @@ def test_exact_successful_read_is_blocked_until_workspace_changes(monkeypatch): return block.tool_type, {"output": "written", "exit_code": 0} return block.tool_type, {"output": "7", "exit_code": 0} - monkeypatch.setattr(agent_loop, "execute_tool_block", execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(execute)) events = _run( "Create answer.json from the inspected workspace", @@ -378,7 +379,7 @@ def test_exact_successful_read_is_blocked_until_workspace_changes(monkeypatch): assert any(tool == "write_file" for tool, _ in executed) -def test_terminal_completion_recovers_fenced_body_after_two_repairs(monkeypatch): +def test_missing_evidence_does_not_generate_later_fenced_body_or_mutation(monkeypatch): executed = [] calls = _patch_loop( monkeypatch, @@ -395,20 +396,19 @@ def test_terminal_completion_recovers_fenced_body_after_two_repairs(monkeypatch) executed.append(block) return await original_execute(block, *args, **kwargs) - monkeypatch.setattr(agent_loop, "execute_tool_block", record_execute) + monkeypatch.setattr(agent_loop, "execute_tool_block", authoritative_executor(record_execute)) events = _run("Write answer.json", max_rounds=5) - assert [(block.tool_type, block.content) for block in executed] == [ - ("write_file", 'answer.json\n{"ok": true}'), - ] + assert executed == [] + assert calls() == 1 decision = next( event["data"] for event in events if event.get("type") == "completion_decision" ) - assert decision["status"] == "satisfied" - assert decision["can_complete"] is True + assert decision["status"] == "blocked" + assert decision["can_complete"] is False def test_successful_artifact_write_emits_satisfied_completion(monkeypatch): @@ -514,7 +514,8 @@ def test_verified_artifact_survives_provider_error_during_finish_round(monkeypat assert not any(event.get("type") == "agent_terminal" for event in events) final = next(event for event in events if event.get("type") == "final_response") assert "output.html" in final["content"] - assert "verified" in final["content"].lower() + assert "Output available" in final["content"] + assert "No passing executable test result" in final["content"] def test_uninspected_artifact_still_fails_on_provider_error(monkeypatch): diff --git a/tests/test_agent_runtime_context.py b/tests/test_agent_runtime_context.py index 4514bf187..b6fe40562 100644 --- a/tests/test_agent_runtime_context.py +++ b/tests/test_agent_runtime_context.py @@ -1122,7 +1122,8 @@ def test_native_host_shell_call_runs_through_bridge_and_threads_result(monkeypat # at import time, so a fresh `import src.tool_execution` here can bind a # different module object than the execute_tool_block agent_loop calls — # patching that fresh copy silently no-ops in full-suite runs. - _dispatch_globals = al.execute_tool_block.__globals__ + from inspect import unwrap + _dispatch_globals = unwrap(al.execute_tool_block).__globals__ monkeypatch.setitem( _dispatch_globals, "owner_is_admin_or_single_user", diff --git a/tests/test_foreground_model_routing.py b/tests/test_foreground_model_routing.py index 659c6761b..a683595e8 100644 --- a/tests/test_foreground_model_routing.py +++ b/tests/test_foreground_model_routing.py @@ -2331,6 +2331,9 @@ def test_multi_round_agent_uses_only_selected_model(monkeypatch): yield f'data: {json.dumps({"delta": "done"})}\n\n' yield "data: [DONE]\n\n" + from tests.runtime_evidence_helpers import authoritative_executor + + @authoritative_executor async def fake_execute(block, *args, **kwargs): return "bash", {"output": "ok", "exit_code": 0} diff --git a/tests/test_runtime_evidence_contract.py b/tests/test_runtime_evidence_contract.py new file mode 100644 index 000000000..3af727f98 --- /dev/null +++ b/tests/test_runtime_evidence_contract.py @@ -0,0 +1,293 @@ +"""Observable execution, stale evidence and completion-stream trust boundaries.""" +import asyncio +from contextlib import aclosing +from inspect import signature +import json +import os + +import pytest + +from src.agent_evidence import CompletionRequirements, EvidenceLedger, EvidenceKind +from src.agent_runtime.completion import completion_answer, with_completion_gate +from src.agent_runtime.identity import artifact_identity, artifact_version, is_test_command, is_validation_command +from src.agent_runtime.journal import ( + ActionJournal, bind_journal, current_journal, execute_action, mark_dispatch, + propose_action, record_action, +) +from src.tool_types import ToolBlock + + +@pytest.mark.parametrize('command', [ + 'python -m unittest discover -s tests -v', 'python3.12 -I -m unittest tests.test_app', + 'cd /workspace && python3 -m unittest', 'pytest -q tests/test_app.py', + '/usr/bin/python3 -m pytest', 'PYTHONPATH=. python -m unittest', 'npm run test', +]) +def test_actual_foreground_test_commands(command): + assert is_test_command(command) + + +@pytest.mark.parametrize('command', [ + 'echo python -m unittest', 'echo "pytest passed"', 'false && pytest', + 'pytest; true', 'pytest || true', 'pytest | cat', 'python -c "print(\'pytest\')"', + 'printf "python -m unittest"', 'pytest --help', 'pytest --collect-only', + 'python -m unittest --help', 'if false; then pytest; fi', 'echo $(pytest)', +]) +def test_non_execution_or_masked_status_is_not_verifier(command): + assert not is_test_command(command) + + +def test_echoed_readback_is_not_validation(): + assert not is_validation_command('echo cat answer.json') + assert is_validation_command('cat answer.json') + + +def test_workspace_path_aliases_and_unrelated_basenames(tmp_path): + (tmp_path / 'nested').mkdir() + (tmp_path / 'a.py').write_text('x') + (tmp_path / 'alias.py').symlink_to(tmp_path / 'a.py') + expected = artifact_identity('a.py', str(tmp_path)) + assert all(artifact_identity(path, str(tmp_path)) == expected for path in + ('./a.py', '/workspace/a.py', str(tmp_path / 'a.py'), 'nested/../a.py', 'alias.py')) + assert artifact_identity('nested/a.py', str(tmp_path)) != expected + assert artifact_identity('../a.py', str(tmp_path)) != expected + assert artifact_identity('/workspace-other/a.py', str(tmp_path)) != expected + assert artifact_identity('a.py.', str(tmp_path)) != expected + + +def test_literal_tool_path_punctuation_is_not_prose_to_strip(): + ledger = EvidenceLedger.from_tool_events([ + {'tool': 'write_file', 'command': '{"path":"app.py."}', 'exit_code': 0}, + ], CompletionRequirements(required_artifacts=('app.py',))) + assert ledger.evaluate().missing_artifacts == ('app.py',) + + +def test_artifact_observation_does_not_open_sensitive_or_outside_files(tmp_path, monkeypatch): + (tmp_path / '.SSH').mkdir() + (tmp_path / '.SSH' / 'id_rsa').write_text('sensitive fixture') + def forbidden(*args, **kwargs): + raise AssertionError('protected artifact must not be opened') + monkeypatch.setattr(os, 'open', forbidden) + assert artifact_version('.SSH/id_rsa', str(tmp_path)) == 'unobserved' + assert artifact_version('../outside', str(tmp_path)) == 'unobserved' + + +def test_fifo_artifact_observation_is_nonblocking(tmp_path): + os.mkfifo(tmp_path / 'pipe') + assert artifact_version('pipe', str(tmp_path)) == 'unobserved' + + +@pytest.mark.parametrize('nested', [True, False]) +def test_native_argument_shapes_are_preserved_without_mutable_aliases(nested): + function = {'name': 'provider_tool', 'arguments': {'value': 'original'}} + native = {'function': function} if nested else function + journal = ActionJournal() + action = journal.propose(ToolBlock('normalized_tool', '{}'), native_call=native) + function['arguments']['value'] = 'changed later' + assert action.provider_arguments == {'value': 'original'} + assert action.provider_tool == 'provider_tool' + + +def test_large_artifact_hashing_is_bounded(tmp_path): + with (tmp_path / 'large.bin').open('wb') as stream: + stream.truncate(64 * 1024 * 1024 + 1) + assert artifact_version('large.bin', str(tmp_path)) == 'unobserved' + + +@pytest.mark.asyncio +async def test_client_completion_declaration_cannot_grant_a_host_workspace(tmp_path): + seen = [] + @with_completion_gate + async def stream(messages, client_runtime_context=None): + seen.append(current_journal().workspace) + yield 'data: {"delta":"I cannot verify that."}\n\n' + yield 'data: [DONE]\n\n' + context = {'completion_requirements': {'workspace_root': str(tmp_path), 'required_artifacts': ['secret.txt']}} + _ = [chunk async for chunk in stream([], client_runtime_context=context)] + assert seen == [''] + + +def test_denied_and_never_dispatched_results_are_not_authoritative(): + for flags in ({'blocked': True}, {'execution_attempted': False}, {'approval_required': True}): + ledger = EvidenceLedger.from_tool_events([ + {'tool': 'bash', 'command': 'python -m unittest', 'exit_code': 0, **flags}], + CompletionRequirements(verifier_required=True, executable_verifier_available=True)) + assert not ledger.evaluate().can_complete + assert not any(e.authoritative for e in ledger.events) + + +def test_readback_does_not_substitute_for_required_executable_tests(): + ledger = EvidenceLedger.from_tool_events([ + {'tool': 'write_file', 'command': '{"path":"answer.json"}', 'exit_code': 0}, + {'tool': 'read_file', 'command': '/workspace/answer.json', 'exit_code': 0}, + ], CompletionRequirements(required_artifacts=('answer.json',), verifier_required=True, + executable_verifier_available=True)) + assert not ledger.evaluate().can_complete + + +@pytest.mark.parametrize('claim', ['All tests passed.', 'Tests: PASS', 'unittest succeeded', + 'Test suite ran successfully', 'No failures.', 'Done.', + 'I executed the command.', 'Successfully created the file.']) +def test_no_execution_receipts_cannot_support_adversarial_success_claims(claim): + ledger = EvidenceLedger() + answer, reason = completion_answer(claim, ledger, ledger.evaluate()) + assert reason + assert answer.startswith('The task is incomplete:') + + +def test_declared_execution_contract_does_not_publish_invented_test_counts(): + ledger = EvidenceLedger.from_tool_events([ + {'tool': 'write_file', 'command': '{"path":"app.py"}', 'exit_code': 0}, + {'tool': 'bash', 'command': 'python -m unittest', 'exit_code': 0}, + ], CompletionRequirements(required_artifacts=('app.py',))) + answer, _ = completion_answer('All 938 tests passed, 100% coverage, everything fixed.', ledger, ledger.evaluate()) + assert '938' not in answer and '100%' not in answer and 'everything' not in answer + assert 'executable verification passed' in answer + + +@record_action +async def successful_backend(block): + mark_dispatch() + return block.tool_type, {'exit_code': 0, 'output': 'OK'} + + +@pytest.mark.asyncio +async def test_normalization_preserves_provider_arguments_and_replay_identity(): + journal = ActionJournal(run_id='known') + original = ToolBlock('write_file', 'original arguments') + normalized = ToolBlock('bash', 'python -m unittest') + with bind_journal(journal): + action = propose_action(original, 'native-1', {'function': {'arguments': '{"original":true}'}}) + await execute_action(successful_backend, action, normalized) + receipt = action.to_dict() + assert receipt['proposed_arguments'] == 'original arguments' + assert receipt['provider_arguments'] == '{"original":true}' + assert receipt['arguments'] == normalized.content + assert [t['stage'] for t in receipt['transitions']] == ['proposed', 'normalized', 'authorized', 'dispatched', 'outcome'] + assert receipt['execution_id'] == 'known:action:1:execution:1' + first = EvidenceLedger.from_tool_events(journal.evidence_events()) + replay = EvidenceLedger.from_tool_events(json.loads(json.dumps(journal.evidence_events()))) + assert first.to_list() == replay.to_list() + assert first.evaluate().status.value == 'verified' + assert first.events[-1].verification_id + + +@pytest.mark.asyncio +async def test_changed_bytes_invalidate_a_passing_verifier(tmp_path): + path = tmp_path / 'app.py' + path.write_text('before') + journal = ActionJournal(workspace=str(tmp_path), observed_artifacts=('app.py',)) + with bind_journal(journal): + await successful_backend(ToolBlock('write_file', '{"path":"app.py"}')) + await successful_backend(ToolBlock('bash', 'python -m unittest')) + requirements = CompletionRequirements(required_artifacts=('app.py',), workspace_root=str(tmp_path)) + assert EvidenceLedger.from_tool_events(journal.evidence_events(), requirements).evaluate().can_complete + path.write_text('changed outside recorded call') + decision = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements).evaluate() + assert not decision.can_complete + assert 'changed after verification' in decision.reason + + +def decode(chunks): + return [json.loads(c[6:]) for c in chunks if c.strip() != 'data: [DONE]'] + + +@pytest.mark.asyncio +async def test_gate_holds_false_claim_until_decision_without_another_round(): + invocations = [] + @with_completion_gate + async def stream(messages, workspace=None, client_runtime_context=None): + invocations.append(1) + yield 'data: {"delta":"All tests "}\n\n' + yield 'data: {"type":"tool_start","tool":"bash"}\n\n' + yield 'data: {"delta":"passed."}\n\n' + yield 'data: {"type":"metrics","data":{}}\n\n' + yield 'data: [DONE]\n\n' + events = decode([c async for c in stream([{'role': 'user', 'content': 'Run the tests'}])]) + assert invocations == [1] + assert events[0]['type'] == 'tool_start' + assert events[1]['type'] == 'completion_decision' + assert not events[1]['data']['can_complete'] + assert all('All tests passed' not in str(e) for e in events) + assert events[2]['content'].startswith('The task is incomplete:') + assert events[3]['data']['round_texts'] == [events[2]['content']] + + +@pytest.mark.asyncio +async def test_gate_preserves_verified_answer_and_sse_shape(): + @with_completion_gate + async def stream(messages): + await successful_backend(ToolBlock('bash', 'python -m unittest')) + yield 'data: {"delta":"Tests passed."}\n\n' + yield 'data: [DONE]\n\n' + chunks = [c async for c in stream([])] + events = decode(chunks) + assert events[0]['data']['status'] == 'verified' + assert events[1] == {'delta': 'Tests passed.'} + assert chunks[-1] == 'data: [DONE]\n\n' + assert str(signature(stream)) == '(messages)' + + +@pytest.mark.asyncio +async def test_cancellation_unwinds_bound_journal_without_done_or_claims(): + closed = [] + @with_completion_gate + async def stream(messages): + try: + yield 'data: {"delta":"Tests passed."}\n\n' + yield 'data: {"type":"tool_start","tool":"bash"}\n\n' + await asyncio.Event().wait() + finally: + closed.append(current_journal() is not None) + async with aclosing(stream([])) as output: + assert json.loads((await anext(output))[6:])['type'] == 'tool_start' + assert closed == [True] + assert current_journal() is None + + +@pytest.mark.asyncio +async def test_real_unittest_dispatch_and_policy_denial_have_distinct_receipts(tmp_path, monkeypatch): + from src.tool_execution import execute_tool_block, NO_TOOL_SECURITY_CONTEXT + monkeypatch.setattr('src.tool_execution.owner_is_admin_or_single_user', lambda owner: True) + (tmp_path / 'test_sample.py').write_text('import unittest\nclass TestSample(unittest.TestCase):\n def test_ok(self): self.assertEqual(2+2,4)\n') + journal = ActionJournal() + with bind_journal(journal): + _, denied = await execute_tool_block(ToolBlock('bash', 'python3 -m unittest'), + workspace=str(tmp_path), disabled_tools={'bash'}, security_context=NO_TOOL_SECURITY_CONTEXT) + _, result = await execute_tool_block(ToolBlock('bash', 'python3 -m unittest -v'), + workspace=str(tmp_path), security_context=NO_TOOL_SECURITY_CONTEXT) + assert denied['exit_code'] != 0 + assert journal.actions[0].execution_id is None + assert not journal.actions[0].operation_started + assert result['exit_code'] == 0, result + assert 'Ran 1 test' in result['output'] + assert journal.actions[1].execution_id + assert journal.actions[1].operation_started + assert EvidenceLedger.from_tool_events(journal.evidence_events()).evaluate().status.value == 'verified' + + +@pytest.mark.asyncio +async def test_shell_writing_same_basename_elsewhere_is_not_required_mutation(tmp_path, monkeypatch): + from src.tool_execution import execute_tool_block, NO_TOOL_SECURITY_CONTEXT + monkeypatch.setattr('src.tool_execution.owner_is_admin_or_single_user', lambda owner: True) + (tmp_path / 'app.py').write_text('unchanged') + journal = ActionJournal(workspace=str(tmp_path), observed_artifacts=('app.py',)) + with bind_journal(journal): + _, result = await execute_tool_block(ToolBlock('bash', 'mkdir nested && printf changed > nested/app.py'), + workspace=str(tmp_path), security_context=NO_TOOL_SECURITY_CONTEXT) + assert result['exit_code'] == 0 + assert (tmp_path / 'nested' / 'app.py').read_text() == 'changed' + assert journal.actions[0].artifact_changes == [] + ledger = EvidenceLedger.from_tool_events(journal.evidence_events(), + CompletionRequirements(required_artifacts=('app.py',), workspace_root=str(tmp_path))) + assert not ledger.evaluate().can_complete + + +@pytest.mark.asyncio +async def test_unknown_tool_never_creates_dispatch_identity(monkeypatch): + from src.tool_execution import execute_tool_block, NO_TOOL_SECURITY_CONTEXT + monkeypatch.setattr('src.tool_execution.owner_is_admin_or_single_user', lambda owner: True) + journal = ActionJournal() + with bind_journal(journal): + await execute_tool_block(ToolBlock('unknown_nonexistent_tool', '{}'), security_context=NO_TOOL_SECURITY_CONTEXT) + assert journal.actions[0].execution_id is None + assert not journal.actions[0].outcome['authoritative'] diff --git a/tests/test_tool_policy.py b/tests/test_tool_policy.py index 952444f72..969cc978a 100644 --- a/tests/test_tool_policy.py +++ b/tests/test_tool_policy.py @@ -1681,6 +1681,9 @@ def test_calendar_create_response_includes_persistent_event_link(monkeypatch): set_user_timezone("Asia/Tokyo", 540) + from tests.runtime_evidence_helpers import authoritative_executor + + @authoritative_executor async def _fake_exec(block, *args, **kwargs): if '"list_calendars"' in (block.content or ""): return (