diff --git a/src/agent_evidence.py b/src/agent_evidence.py index 54d14367c..29daebebe 100644 --- a/src/agent_evidence.py +++ b/src/agent_evidence.py @@ -632,6 +632,84 @@ class EvidenceLedger: # Retain receipt command identity privately for presentation matching; # model prose and client dictionaries never populate this evidence. self._verifier_commands: dict[str, tuple[str, ...]] = {} + # Wave 4 effect assessments from the run's journal, plus the journal + # order of actions so receipt evidence and effects share one ordering. + self.effects: list[dict[str, Any]] = [] + self._action_order: dict[str, int] = {} + + def record_effects(self, entries: Iterable[Mapping[str, Any]], action_order: Mapping[str, int], + partial_reads: Iterable[str] = ()) -> None: + """Consume server-derived effect assessments (never model/client data). + + ``partial_reads`` names read actions whose admitted observation was + partial (offset/limit, truncation or extraction): such a read cannot + validate omitted content, so its validation event is not authoritative. + """ + from dataclasses import replace + from src.agent_runtime.effects import EffectAssessment + self.effects = [dict(entry) for entry in entries + if isinstance(entry, Mapping) and isinstance(entry.get("assessment"), EffectAssessment)] + self._action_order = {str(k): v for k, v in action_order.items() if type(v) is int} + partial = set(partial_reads) + self.events = [replace(event, authoritative=False, detail="partial read; omitted content is unvalidated") + if event.kind == EvidenceKind.ARTIFACT_VALIDATION and event.tool == "read_file" + and event.action_id in partial else event for event in self.events] + + def _last_success_ordinal(self, required: str) -> int: + return max((self._action_order.get(event.action_id, 0) for event in self.events + if event.kind == EvidenceKind.ARTIFACT_MUTATION and event.authoritative and event.success + and _artifact_path_matches_required(event.artifact_path, required, self.requirements.workspace_root)), + default=0) + + def _later_effects(self, required: str) -> list[tuple[dict[str, Any], bool]]: + """Effects after the artifact's last successful mutation, with targeting.""" + floor = self._last_success_ordinal(required) + later = [] + for entry in self.effects: + ordinal = entry.get("ordinal") + if type(ordinal) is not int or ordinal <= floor: + continue + explicit = any(_artifact_path_matches_required(path, required, self.requirements.workspace_root) + for path in entry.get("paths") or ()) + if explicit or entry.get("unknown_scope"): + later.append((entry, explicit)) + return later + + def _effect_unsettled(self, required: str) -> bool: + """A later operation may have partially changed this artifact. + + Explicit targets are unsettled by unknown/timed-out/cancelled outcomes + and by failures after the producer reached its mutation stage (atomic + refusals keep the earlier artifact). Unknown-scope effects are + unsettled when nothing captured their settlement: cancellation, + interruption, or failed process teardown; settled shell/Python changes + are already tracked through artifact version capture. + """ + from src.agent_runtime.effects import CleanupState, ExecutionOutcome + unknown = {ExecutionOutcome.ATTEMPTED, ExecutionOutcome.INTERRUPTED, ExecutionOutcome.CANCELLED} + for entry, explicit in self._later_effects(required): + assessment = entry["assessment"] + if not assessment.unresolved_impact: + continue + if assessment.execution in unknown or assessment.execution is ExecutionOutcome.RUNNING: + return True + if explicit and (assessment.execution is ExecutionOutcome.TIMED_OUT + or (assessment.execution is ExecutionOutcome.FAILED and entry.get("mutation_attempted"))): + return True + if not explicit and assessment.cleanup is CleanupState.FAILED: + return True + return False + + def _effect_contradicted(self, required: str) -> str: + """The latest effect targeting the artifact, if fresh readback contradicts it.""" + from src.agent_runtime.effects import EffectVerdict + targeting = [entry for entry in self.effects if type(entry.get("ordinal")) is int + and any(_artifact_path_matches_required(path, required, self.requirements.workspace_root) + for path in entry.get("paths") or ())] + if not targeting: + return "" + latest = max(targeting, key=lambda entry: entry["ordinal"])["assessment"] + return latest.effect_id if latest.verdict is EffectVerdict.CONTRADICTED else "" @classmethod def from_tool_events( @@ -827,6 +905,8 @@ class EvidenceLedger: and matching[-1].tool in {'bash', 'python'}) if not successful or destructive_failure: return False + if self.effects and (self._effect_unsettled(path) or self._effect_contradicted(path)): + return False return True def record_media_ingress(self, metadata: Mapping[str, Any]) -> None: @@ -897,8 +977,16 @@ class EvidenceLedger: 'artifact content changed after verification', (latest_verifier.event_id,)) + for required in self.requirements.required_artifacts if self.effects else (): + contradicted = self._effect_contradicted(required) + if contradicted: + return CompletionDecision(CompletionStatus.FAILED, False, + "fresh readback contradicts the requested artifact content", + (), (required,)) + satisfied_ids: list[str] = [] missing: list[str] = [] + unsettled: list[str] = [] workspace_root = str(self.requirements.workspace_root or "").strip() for required in self.requirements.required_artifacts: matches = [ @@ -934,15 +1022,20 @@ class EvidenceLedger: filesystem_missing = True if latest_success is None or destructive_failure or filesystem_missing: missing.append(required) + elif self.effects and self._effect_unsettled(required): + # Earlier success is historical; a later possible change to + # this artifact has no settled evidence. + unsettled.append(required) else: satisfied_ids.append(latest_success.event_id) - if missing: + if missing or unsettled: return CompletionDecision( CompletionStatus.BLOCKED, False, - "required artifacts lack successful mutation evidence", + "required artifacts lack successful mutation evidence" if missing else + "a later operation may have changed a required artifact without settled evidence", tuple(satisfied_ids), - tuple(missing), + tuple([*missing, *unsettled]), ) latest_mutation_index = max( diff --git a/src/agent_runtime/completion.py b/src/agent_runtime/completion.py index b8b80700c..943074df3 100644 --- a/src/agent_runtime/completion.py +++ b/src/agent_runtime/completion.py @@ -20,9 +20,19 @@ from src.agent_evidence import ( requirements_from_runtime_context, _execution_obligation, _unquoted_statements, _ARTIFACT_PATH, ) +from .effect_log import EffectLog from .journal import ActionJournal, bind_journal, current_journal +def _ledger(journal: ActionJournal, requirements) -> EvidenceLedger: + """The single evidence view used for the decision and the prose filter.""" + ledger = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements) + ledger.record_effects(journal.effect_entries(), + {action.action_id: index for index, action in enumerate(journal.actions, 1)}, + journal.partial_reads()) + return ledger + + _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) @@ -192,6 +202,10 @@ def with_completion_gate(func): journal = ActionJournal( workspace=requirements.workspace_root, observed_artifacts=requirements.required_artifacts, parent_run_id=bound.get('_parent_run_id') or (parent.run_id if parent is not None else None)) + # One durable effect log per run lineage gives child effects and parent + # observations a single total order for invalidation. + journal.effects = (parent.effects if parent is not None and parent.effects is not None + else EffectLog(journal.run_id)) answer_events: list[dict] = [] metrics_events: list[dict] = [] answer = '' @@ -241,7 +255,7 @@ def with_completion_gate(func): awaiting = True payload = data.get('data') or {} if isinstance(payload.get('question'), str): - current = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements) + current = _ledger(journal, requirements) question, why = completion_answer(payload['question'], current, current.evaluate(awaiting_user=True)) if why: data = {**data, 'data': {**payload, 'question': question}} @@ -290,7 +304,7 @@ def with_completion_gate(func): presentation_replaced = True answer = terminal_answer answer_events = [event for event in answer_events if event.get('thinking') is True] - ledger = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements) + ledger = _ledger(journal, requirements) decision = ledger.evaluate(exhausted=exhausted, awaiting_user=awaiting) if provider_error: decision = replace(decision, status=CompletionStatus.FAILED, @@ -334,6 +348,8 @@ def with_completion_gate(func): metadata.update(completion_decision=decision.to_dict(), evidence_events=ledger.to_list(), action_receipts=journal.to_list(), completion_requirements=requirements.to_dict(), run_id=journal.run_id, parent_run_id=journal.parent_run_id) + if ledger.effects: + metadata['effect_assessments'] = [entry['assessment'].to_dict() for entry in ledger.effects] 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, diff --git a/src/agent_runtime/effect_adapters.py b/src/agent_runtime/effect_adapters.py new file mode 100644 index 000000000..bd8ba3d50 --- /dev/null +++ b/src/agent_runtime/effect_adapters.py @@ -0,0 +1,370 @@ +"""Server-boundary adapters from admitted Wave 3 bindings to effect records. + +Runs only inside the dispatcher's existing admission scope: the bindings read +here are the contextvars the dispatcher bound after authority, resource and +approval checks. Nothing here admits, resolves, broadens or re-derives a +resource. Observations are recorded only for operations that were themselves +admitted reads of the exact bound resource; evidence bookkeeping never performs +a read that the operation was not already admitted to perform. +""" +from __future__ import annotations + +import asyncio +from dataclasses import dataclass, field +import hashlib +import json +import logging +import os +import stat +from typing import Any + +from src.agent_runtime.effects import ( + CleanupState, Coverage, EffectClaim, ExecutionOutcome, Impact, ObservationMechanism, OperationRef, + Postcondition, Predicate, ProducerFacts, ResourceKind, ResourceRef, producer_facts, resource_ref, +) + + +_FILESYSTEM_READS = frozenset({"read_file", "ls", "glob", "grep"}) +_JOB_READS = frozenset({"list", "ls", "jobs", "output", "get", "read", "tail", "status", "show"}) +_OWNED_READS = frozenset({"vault_get", "vault_search", "list_sessions", "search_chats"}) +_JOB_SETTLED = {"done", "failed"} +logger = logging.getLogger(__name__) + + +@dataclass +class DispatchCapture: + """The admitted bindings that were live when the backend was invoked.""" + + filesystem: Any = None + owned: Any = None + process: Any = None + backend: Any = None + browser: Any = None + claim: EffectClaim | None = None + read_only: bool = False + paths: tuple[str, ...] = field(default_factory=tuple) + + +def capture_dispatch() -> DispatchCapture: + from src.agent_runtime.owned_resources import active_owned_operation + from src.agent_runtime.process_resources import active_process_operation + from src.agent_runtime.remote_resources import active_backend_operation + from src.agent_runtime.resource_binding import active_resource_operation + import sys + browser_module = sys.modules.get("src.browser_identity") + browser = browser_module._ACTIVE.get() if browser_module is not None else None + return DispatchCapture(active_resource_operation(), active_owned_operation(), active_process_operation(), + active_backend_operation(), browser) + + +def _exact_operation(capture: DispatchCapture): + for bound in (capture.filesystem, capture.owned, capture.process, capture.browser): + if bound is not None: + return bound.operation, getattr(bound, "execution_input", None), getattr(bound, "request_id", "") + return None, None, "" + + +def _operation(capture: DispatchCapture, action: Any) -> OperationRef: + operation, execution_input, request_id = _exact_operation(capture) + if operation is not None: + return OperationRef.from_exact(operation, execution_input, request_id) + backend = capture.backend + # Unbound tools still name their final normalized dispatcher input. + digest = hashlib.sha256(str(action.arguments).encode("utf-8", errors="replace")).hexdigest() + return OperationRef(str(action.tool) or "unknown", "", digest, + getattr(backend, "request_id", "") if backend is not None else "") + + +def _write_file_digest(execution_input: str, path: str) -> str: + """The exact bytes WriteFileTool commits for this admitted input, or ''.""" + from src.agent_tools.filesystem_tools import _unwrap_fenced_source_body + try: + args = json.loads(execution_input) + except (TypeError, ValueError): + return "" + body = args.get("content") if isinstance(args, dict) else None + if not isinstance(body, str) or os.linesep != "\n": + return "" + return hashlib.sha256(_unwrap_fenced_source_body(body, path).encode("utf-8")).hexdigest() + + +def _filesystem_scope(bound: Any) -> tuple[tuple[ResourceRef, ...], tuple[Postcondition, ...]]: + from src.agent_tools.filesystem_tools import _parse_agent_patch + tool = bound.operation.tool + refs = tuple(resource_ref(b.resource, b.role) for b in bound.bindings) + obligations: list[Postcondition] = [] + if tool == "write_file": + target = refs[0] + expected = _write_file_digest(bound.execution_input, bound.bindings[0].resource.path) + obligations.append(Postcondition(target, Predicate.CONTENT_SHA256, expected) if expected + else Postcondition(target, Predicate.EXISTS)) + elif tool == "edit_file": + obligations.append(Postcondition(refs[0], Predicate.EXISTS)) + elif tool == "apply_patch": + ops = _parse_agent_patch(json.loads(bound.execution_input)["patch_text"]) + for op, ref in zip(ops, refs): + if op["kind"] == "add": + digest = hashlib.sha256(op["content"].encode("utf-8")).hexdigest() + obligations.append(Postcondition(ref, Predicate.CONTENT_SHA256, digest)) + elif op["kind"] == "delete": + obligations.append(Postcondition(ref, Predicate.ABSENT)) + else: + obligations.append(Postcondition(ref, Predicate.EXISTS)) + return refs, tuple(obligations) + + +def classify(capture: DispatchCapture) -> dict[str, Any] | None: + """Claim scope for the captured bindings, or None for an admitted read. + + Unbound operations get an unknown-scope claim: they may change anything. + """ + impact: tuple[ResourceRef, ...] = () + dependencies: tuple[ResourceRef, ...] = () + obligations: tuple[Postcondition, ...] = () + external = False + if capture.browser is not None: + # Wave 3 admits only session metadata. A page binding is never + # effect-bindable; leave its scope unknown rather than infer it. + if capture.browser.page is None: + return None + elif capture.filesystem is not None: + if capture.filesystem.operation.tool in _FILESYSTEM_READS: + return None + impact, obligations = _filesystem_scope(capture.filesystem) + elif capture.process is not None: + bound = capture.process + if bound.launch is not None: + # An arbitrary command has unknown impact scope; the exact launch + # reservation is kept only as lineage for background settlement. + dependencies = (resource_ref(bound.launch, "launch"),) + else: + action = str(json.loads(bound.operation.input or "{}").get("action", "list")).strip().lower() + if action in _JOB_READS: + return None + impact = tuple(resource_ref(job, "job") for job in bound.jobs) + tuple( + resource_ref(process, "process") for job in bound.jobs for process in job.processes) + tuple( + resource_ref(process, "process") for process in bound.processes) + elif capture.owned is not None: + if capture.owned.operation.tool in _OWNED_READS: + return None + impact = tuple(resource_ref(r, "record") for r in capture.owned.resources) + dependencies = tuple(resource_ref(a.file, "attachment") for a in capture.owned.attachments) + if capture.backend is not None: + from src.agent_runtime.resources import ExternalResource + if isinstance(capture.backend.resource, ExternalResource): + external = True + impact = (*impact, resource_ref(capture.backend.resource, "backend")) + return {"impact_scope": impact, "dependencies": dependencies, "obligations": obligations, "external": external} + + +def begin_effect(journal: Any, action: Any) -> DispatchCapture: + """Capture bindings and durably claim a possible effect before invocation.""" + capture = capture_dispatch() + log = journal.effects + try: + scope = classify(capture) + except Exception: # noqa: BLE001 - classification never blocks dispatch + # An unclassifiable admitted operation may change anything. + logger.warning("Effect scope classification failed; claiming unknown scope", exc_info=True) + scope = {"impact_scope": (), "dependencies": (), "obligations": (), "external": False} + if scope is None: + capture.read_only = True + else: + capture.claim = log.claim(effect_id=action.action_id + ":effect", run_id=journal.run_id, + action_id=action.action_id, operation=_operation(capture, action), + parent_run_id=journal.parent_run_id or "", **scope) + capture.paths = tuple(ref.location[-1] for ref in capture.claim.impact_scope + if ref.kind is ResourceKind.FILESYSTEM) + for ref in capture.claim.dependencies: + if ref.kind is ResourceKind.PROCESS_LAUNCH: + try: + log.index_launch(ref.incarnation, capture.claim.effect_id) + except (OSError, ValueError): + # Without the index a later turn cannot settle this + # launch: it stays running/unknown, never successful. + logger.warning("Background launch lineage was not indexed", exc_info=True) + return capture + + +def _execution(result: Any, facts: ProducerFacts) -> ExecutionOutcome: + if not isinstance(result, dict): + return ExecutionOutcome.INTERRUPTED + if facts.timed_out: + return ExecutionOutcome.TIMED_OUT + if isinstance(result.get("bg_job_id"), str) and facts.exit_code == 0: + return ExecutionOutcome.RUNNING + if result.get("detached") is True or result.get("status") == "running" or result.get("running") is True: + return ExecutionOutcome.RUNNING + denied = bool(result.get("blocked") or result.get("approval_required") + or facts.failure_kind.endswith("_denied")) + if facts.exit_code == 0 and not result.get("error") and not denied: + return ExecutionOutcome.REPORTED_SUCCESS + return ExecutionOutcome.FAILED + + +def _cleanup(result: Any, facts: ProducerFacts) -> CleanupState: + if not isinstance(result, dict): + return CleanupState.UNKNOWN + if facts.failure_kind == "process_teardown_failed": + return CleanupState.FAILED + teardown = result.get("teardown") + if isinstance(teardown, dict) and type(teardown.get("dead")) is bool: + return CleanupState.VERIFIED if teardown["dead"] else CleanupState.FAILED + if facts.external: + # External execution reports no locally observed teardown. + return CleanupState.UNKNOWN + return CleanupState.NOT_APPLICABLE + + +def settle_effect(journal: Any, action: Any, capture: DispatchCapture | None, *, + result: Any = None, error: BaseException | None = None) -> None: + """Append the outcome and any admitted-read observations for one action.""" + if capture is None: + return + log = journal.effects + if capture.claim is not None: + if error is not None: + execution = (ExecutionOutcome.CANCELLED if isinstance(error, asyncio.CancelledError) + else ExecutionOutcome.INTERRUPTED) + facts, cleanup = ProducerFacts(), CleanupState.UNKNOWN + else: + facts = producer_facts(result) + if capture.backend is not None and capture.claim.external: + facts = ProducerFacts(**{**facts.to_dict(), "external": True, + "remote_acknowledged": facts.exit_code == 0}) + execution, cleanup = _execution(result, facts), _cleanup(result, facts) + log.outcome(effect_id=capture.claim.effect_id, execution=execution, impact=Impact.POSSIBLE, + facts=facts, cleanup=cleanup, execution_id=action.execution_id or "") + if (execution is ExecutionOutcome.REPORTED_SUCCESS and capture.process is not None + and capture.process.launch is None): + _settle_background(log, capture, result) # e.g. an exact kill + return + if error is not None or not isinstance(result, dict) or result.get("exit_code") != 0 or result.get("error"): + return + for fields in _observations(capture, action, result): + log.observe(**fields) + if capture.process is not None and capture.process.launch is None: + _settle_background(log, capture, result) + + +# -- observations ------------------------------------------------------------ + +def _read_whole(resource: Any, limit: int) -> bytes | None: + """Re-read the exact admitted source binding; None if it is not stable.""" + flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0) | getattr(os, "O_CLOEXEC", 0) + try: + resource.validate() + descriptor = os.open(resource.path, flags) + with os.fdopen(descriptor, "rb") as stream: + info = os.fstat(stream.fileno()) + identity = resource.identity + if (not stat.S_ISREG(info.st_mode) or identity is None + or (info.st_dev, info.st_ino) != (identity.device, identity.inode)): + return None + data = stream.read(limit + 1) + resource.validate() + except (OSError, ValueError): + return None + return data + + +def _file_observation(capture: DispatchCapture, action: Any) -> dict[str, Any] | None: + from src.agent_tools import filesystem_tools as producer + bound = capture.filesystem + binding = bound.bindings[0] + resource = binding.resource + args = json.loads(bound.execution_input) + partial = bool(args.get("offset") or args.get("limit")) or ( + os.path.splitext(resource.path)[1].lower() in producer._STRUCTURED_DOCUMENT_SUFFIXES) + data = _read_whole(resource, producer.MAX_READ_CHARS * 4) + if data is None: + return None + if len(data) > producer.MAX_READ_CHARS * 4 or len(data.decode("utf-8", errors="replace")) > producer.MAX_READ_CHARS: + partial = True # the producer truncated what it read + complete = not partial + return dict(observation_id=action.action_id + ":observation", resource=resource_ref(resource, binding.role), + mechanism=ObservationMechanism.FILESYSTEM_READ, + coverage=Coverage.COMPLETE if complete else Coverage.PARTIAL, + source_action_id=action.action_id, source_execution_id=action.execution_id or "", + exists=True, content_sha256=hashlib.sha256(data).hexdigest() if complete else "") + + +def _observations(capture: DispatchCapture, action: Any, result: dict) -> list[dict[str, Any]]: + base = dict(source_action_id=action.action_id, source_execution_id=action.execution_id or "") + if capture.browser is not None and capture.browser.page is None: + # Session lifecycle metadata only; never page/document state. + return [dict(observation_id=action.action_id + ":observation", + resource=resource_ref(capture.browser.session, "session"), + mechanism=ObservationMechanism.BROWSER_SESSION, coverage=Coverage.PARTIAL, + exists=True, **base)] + if capture.filesystem is not None: + tool = capture.filesystem.operation.tool + if tool == "read_file": + observation = _file_observation(capture, action) + return [observation] if observation else [] + if tool in _FILESYSTEM_READS: + # Listings/searches are partial: they cannot decide content. + return [dict(observation_id=f"{action.action_id}:observation:{i}", resource=resource_ref(b.resource, b.role), + mechanism=ObservationMechanism.FILESYSTEM_READ, coverage=Coverage.PARTIAL, exists=True, **base) + for i, b in enumerate(capture.filesystem.bindings)] + if capture.owned is not None and capture.owned.operation.tool in _OWNED_READS: + return [dict(observation_id=f"{action.action_id}:observation:{i}", resource=resource_ref(r, "record"), + mechanism=ObservationMechanism.OWNED_RECORD_READ, coverage=Coverage.PARTIAL, exists=True, **base) + for i, r in enumerate(capture.owned.resources) if r.record_id != "*"] + if capture.process is not None and capture.process.launch is None: + job = result.get("job") + if isinstance(job, dict) and len(capture.process.jobs) == 1: + return [dict(observation_id=action.action_id + ":observation", + resource=resource_ref(capture.process.jobs[0], "job"), + mechanism=ObservationMechanism.JOB_STATE, coverage=Coverage.PARTIAL, exists=True, **base)] + return [] + + +def _settle_background(log: Any, capture: DispatchCapture, result: dict) -> None: + """Settle a RUNNING launch claim from an admitted read of its exact job. + + Linkage is the Wave 3 launch generation plus owner/request/thread, already + validated by ``job_from_record`` at admission. Job completion is execution + evidence for that claim; it verifies no postcondition. + """ + job_facts = result.get("job") + if not isinstance(job_facts, dict) or len(capture.process.jobs) != 1: + return + status = job_facts.get("status") + if status not in _JOB_SETTLED: + return + job = capture.process.jobs[0] + lineage = ("process_launch", "native:containment", job.owner, job.request_id, job.thread_id, job.generation) + owner = log if any(any(ref.kind is ResourceKind.PROCESS_LAUNCH and ref.location == lineage + for ref in c.dependencies) for c in log.history().claims) else None + if owner is None and log.path is not None: + # Background continuation: the launch was claimed by an earlier run. + from src.agent_runtime.effect_log import EffectLog, EffectPersistenceError + indexed = EffectLog.launch_owner(job.generation, directory=log.path.parent) + if indexed is not None: + try: + owner = EffectLog.open(indexed[0], directory=log.path.parent) + except (EffectPersistenceError, ValueError): + owner = None + if owner is None: + return + history = owner.history() + for claim in history.claims: + if not any(ref.kind is ResourceKind.PROCESS_LAUNCH and ref.location == lineage for ref in claim.dependencies): + continue + latest = history.latest_outcome(claim.effect_id) + if latest is None or latest.execution is not ExecutionOutcome.RUNNING: + continue + code = job_facts.get("exit_code") + code = code if type(code) is int else None + if job_facts.get("timed_out") is True: + execution = ExecutionOutcome.TIMED_OUT + elif job_facts.get("killed") is True: + execution = ExecutionOutcome.CANCELLED + elif status == "done" and code == 0 and job_facts.get("died") is not True: + execution = ExecutionOutcome.REPORTED_SUCCESS + else: + execution = ExecutionOutcome.FAILED + facts = ProducerFacts(exit_code=code, timed_out=job_facts.get("timed_out") is True, job_state=status) + owner.outcome(effect_id=claim.effect_id, execution=execution, impact=Impact.POSSIBLE, facts=facts, + cleanup=CleanupState.UNKNOWN, execution_id=latest.execution_id) diff --git a/src/agent_runtime/effect_log.py b/src/agent_runtime/effect_log.py index 34ead8372..7addcb0b1 100644 --- a/src/agent_runtime/effect_log.py +++ b/src/agent_runtime/effect_log.py @@ -17,6 +17,7 @@ from pathlib import Path import re import stat import threading +import weakref from typing import Any from src.constants import DATA_DIR @@ -41,11 +42,17 @@ def effects_dir() -> Path: class EffectLog: + # Logs still owned by a live run in this process. A later turn appends to + # the same object rather than a second copy with its own sequence counter. + _LIVE: "weakref.WeakValueDictionary[tuple[str, str], EffectLog]" = weakref.WeakValueDictionary() + def __init__(self, run_id: str, *, durable: bool = True, directory: str | os.PathLike | None = None) -> None: if not isinstance(run_id, str) or not _RUN_ID.fullmatch(run_id): raise ValueError("Effect log requires a server-generated run identifier") self.run_id = run_id self.path = (Path(directory) if directory is not None else effects_dir()) / f"{run_id}.jsonl" if durable else None + if self.path is not None: + self._LIVE[(str(self.path.parent), run_id)] = self self._claims: list[EffectClaim] = [] self._outcomes: list[EffectOutcome] = [] self._observations: list[Observation] = [] @@ -167,6 +174,60 @@ class EffectLog: default=0) return log + @classmethod + def open(cls, run_id: str, *, directory: str | os.PathLike | None = None) -> "EffectLog": + """The live log for a run, or its replayed durable history. + + A log that is not live belongs to a finished or crashed run, so its + unsettled claims are recovered as interrupted before any append. + """ + base = Path(directory) if directory is not None else effects_dir() + live = cls._LIVE.get((str(base), run_id)) + if live is not None: + return live + log = cls.load(run_id, directory=base) + log.recover_interrupted() + return log + + # -- background launch lineage ------------------------------------------ + + def index_launch(self, generation: str, effect_id: str) -> None: + """Durably map an exact Wave 3 launch generation to its claim.""" + if self.path is None: + return + if not _RUN_ID.fullmatch(generation or ""): + raise ValueError("Malformed launch generation") + target = self.path.parent / f"launch-{generation}.json" + temporary = target.with_suffix(".tmp") + data = json.dumps({"run_id": self.run_id, "effect_id": effect_id}, sort_keys=True).encode() + flags = os.O_WRONLY | os.O_CREAT | os.O_TRUNC | getattr(os, "O_NOFOLLOW", 0) + descriptor = os.open(temporary, flags, 0o600) + try: + os.write(descriptor, data) + os.fsync(descriptor) + finally: + os.close(descriptor) + os.replace(temporary, target) + + @staticmethod + def launch_owner(generation: str, *, directory: str | os.PathLike | None = None) -> tuple[str, str] | None: + if not _RUN_ID.fullmatch(generation or ""): + return None + base = Path(directory) if directory is not None else effects_dir() + try: + descriptor = os.open(base / f"launch-{generation}.json", os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0)) + with os.fdopen(descriptor, "rb") as stream: + if os.fstat(stream.fileno()).st_nlink != 1: + return None + value = json.loads(stream.read(4096)) + except (OSError, ValueError): + return None + if (not isinstance(value, dict) or set(value) != {"run_id", "effect_id"} + or not isinstance(value["run_id"], str) or not _RUN_ID.fullmatch(value["run_id"]) + or not isinstance(value["effect_id"], str)): + return None + return value["run_id"], value["effect_id"] + def recover_interrupted(self) -> tuple[EffectOutcome, ...]: """Append INTERRUPTED outcomes for claims that never settled.""" with self._lock: diff --git a/src/agent_runtime/effects.py b/src/agent_runtime/effects.py index e69572eec..994557645 100644 --- a/src/agent_runtime/effects.py +++ b/src/agent_runtime/effects.py @@ -54,6 +54,7 @@ _SHA256 = re.compile(r"[a-f0-9]{64}") class ResourceKind(str, Enum): FILESYSTEM = "filesystem" PROCESS = "process" + PROCESS_LAUNCH = "process_launch" BACKGROUND_JOB = "background_job" OWNED = "owned" EXTERNAL = "external" @@ -108,6 +109,11 @@ class ResourceRef: """ if self.kind is not other.kind: return False + if self.kind is ResourceKind.OWNED: + # A collection binding ("*") covers every record it can create, + # list or change; specific records only overlap themselves. + return self.location[:-1] == other.location[:-1] and ( + self.location[-1] == other.location[-1] or "*" in (self.location[-1], other.location[-1])) if self.kind is not ResourceKind.FILESYSTEM: return self.location == other.location if self.location[:-1] != other.location[:-1]: @@ -153,6 +159,13 @@ def resource_ref(resource: Any, role: str) -> ResourceRef: location = ("process", resource.namespace, resource.owner, resource.request_id, resource.thread_id, str(ident.pid), ident.start_token, resource.role) return ResourceRef(ResourceKind.PROCESS, role, location, ident.start_token, _sha(resource.to_dict())) + if isinstance(resource, wave3.ProcessLaunchResource): + # The reservation generation is the exact launch -> job linkage that + # Wave 3 validates in ``job_from_record``. + location = ("process_launch", resource.namespace, resource.owner, resource.request_id, + resource.thread_id, resource.generation) + return ResourceRef(ResourceKind.PROCESS_LAUNCH, role, location, resource.generation, + _sha(resource.to_dict())) if isinstance(resource, wave3.BackgroundJobResource): location = ("background_job", resource.namespace, resource.owner, resource.request_id, resource.thread_id, resource.job_id, resource.generation) @@ -378,11 +391,13 @@ class ProducerFacts: job_state: str = "" remote_acknowledged: bool = False external: bool = False + # The producer reached its mutation stage before reporting failure. + mutation_attempted: bool = False def __post_init__(self) -> None: if self.exit_code is not None and type(self.exit_code) is not int: raise ValueError("Malformed producer exit code") - for name in ("timed_out", "output_truncated", "remote_acknowledged", "external"): + for name in ("timed_out", "output_truncated", "remote_acknowledged", "external", "mutation_attempted"): if type(getattr(self, name)) is not bool: raise ValueError("Malformed producer flag") for name in ("failure_kind", "job_state"): @@ -395,7 +410,7 @@ class ProducerFacts: return {"exit_code": self.exit_code, "timed_out": self.timed_out, "output_truncated": self.output_truncated, "failure_kind": self.failure_kind, "job_state": self.job_state, "remote_acknowledged": self.remote_acknowledged, - "external": self.external} + "external": self.external, "mutation_attempted": self.mutation_attempted} @classmethod def from_dict(cls, value: Any) -> "ProducerFacts": @@ -428,6 +443,7 @@ def producer_facts(result: Any) -> ProducerFacts: failure_kind=_label(result.get("failure_kind")), job_state=_label(job), external=external, + mutation_attempted=result.get("mutation_attempted") is True, ) diff --git a/src/agent_runtime/journal.py b/src/agent_runtime/journal.py index 1d2c8456b..200718bf5 100644 --- a/src/agent_runtime/journal.py +++ b/src/agent_runtime/journal.py @@ -7,6 +7,7 @@ from copy import deepcopy from dataclasses import dataclass, field, asdict from functools import wraps from inspect import signature +import logging from typing import Any from uuid import uuid4 @@ -67,6 +68,43 @@ class ActionJournal: workspace: str = '' observed_artifacts: tuple[str, ...] = () parent_run_id: str | None = None + # Durable Wave 4 effect log, shared across one run lineage; None disables. + effects: Any = field(default=None, repr=False, compare=False) + _dispatches: dict[str, Any] = field(default_factory=dict, repr=False, compare=False) + + def effect_entries(self) -> list[dict[str, Any]]: + """Effect assessments ordered against this journal's actions. + + Ordinal is the 1-based position of the action in this journal, so the + ledger can compare effects with receipt-derived evidence. Effects from + other journals in the lineage carry no ordinal here. + """ + if self.effects is None: + return [] + order = {action.action_id: index for index, action in enumerate(self.actions, 1)} + changes = {action.action_id: action.artifact_changes for action in self.actions} + history = self.effects.history() + entries = [] + for assessment in self.effects.assessments(): + claim = history.claim(assessment.effect_id) + outcome = history.latest_outcome(assessment.effect_id) + entries.append({ + 'ordinal': order.get(assessment.action_id), 'assessment': assessment, + 'tool': claim.operation.tool, 'unknown_scope': claim.unknown_scope, + 'paths': tuple(ref.location[-1] for ref in claim.impact_scope if ref.kind.value == 'filesystem'), + 'mutation_attempted': bool(outcome and outcome.facts.mutation_attempted), + 'artifact_changes': changes.get(assessment.action_id), + }) + return entries + + def partial_reads(self) -> tuple[str, ...]: + """Read actions in this journal whose admitted observation was partial.""" + if self.effects is None: + return () + mine = {action.action_id for action in self.actions} + return tuple(o.source_action_id for o in self.effects.history().observations + if o.source_action_id in mine and o.mechanism.value == 'filesystem_read' + and o.coverage.value == 'partial') def capture_versions(self, action: ActionReceipt) -> None: if self.workspace: @@ -138,6 +176,16 @@ def mark_authorized() -> None: def mark_dispatch() -> None: action = _ACTION.get() if action is not None and action.execution_id is None: + journal = _JOURNAL.get() + if journal is not None and journal.effects is not None: + # Durable claim first. If it cannot be persisted this raises and + # the action stays undispatched: the backend is never invoked. + from .effect_adapters import begin_effect + capture = begin_effect(journal, action) + journal._dispatches[action.action_id] = capture + if capture.claim is not None: + action.transition('effect_claimed', effect_id=capture.claim.effect_id, + sequence=capture.claim.sequence) mark_authorized() action.execution_id = action.action_id + ':execution:1' action.transition('dispatched', execution_id=action.execution_id) @@ -145,10 +193,32 @@ def mark_dispatch() -> None: async def dispatched(operation): """Record an actual backend invocation, distinct from router admission.""" - mark_dispatch() + try: + mark_dispatch() + except BaseException: + close = getattr(operation, 'close', None) + if close is not None: + close() # never invoked; do not leave an un-awaited coroutine + raise return await operation +def _settle(journal: ActionJournal | None, action: ActionReceipt, **outcome: Any) -> None: + if journal is None or journal.effects is None: + return + capture = journal._dispatches.pop(action.action_id, None) + if capture is None: + return + from .effect_adapters import settle_effect + try: + settle_effect(journal, action, capture, **outcome) + except Exception: # noqa: BLE001 - bookkeeping must not alter the tool result + # The claim stays unsettled (ATTEMPTED), which assesses as pending + # with possible impact: conservative, never a manufactured success. + journal.effects.degraded = True + logging.getLogger(__name__).warning('Effect outcome could not be recorded', exc_info=True) + + def mark_operation_started(backend: str, **details: Any) -> None: action = _ACTION.get() if action is not None: @@ -197,10 +267,14 @@ def record_action(func): action.finish({**result, 'blocked': True}) else: action.finish(result) + # Structured producer facts are projected here, before the + # receipt reduction drops them. + _settle(journal, action, result=result) return description, result except BaseException as exc: if action is not None: action.transition('interrupted', category=type(exc).__name__) + _settle(current_journal(), action, error=exc) raise finally: _ACTION.reset(token) diff --git a/src/agent_tools/bg_job_tools.py b/src/agent_tools/bg_job_tools.py index 8d3fd5c1a..0f7260025 100644 --- a/src/agent_tools/bg_job_tools.py +++ b/src/agent_tools/bg_job_tools.py @@ -44,6 +44,19 @@ def _status_label(rec: Dict[str, Any]) -> str: return status +def _job_facts(rec: Dict[str, Any]) -> Dict[str, Any]: + """Typed lifecycle facts from the exact admitted job record. + + Execution evidence only: completion of a job is not verification of any + filesystem, service or external state its command was meant to change. + """ + code = rec.get("exit_code") + return {"status": rec.get("status") if rec.get("status") in {"running", "done", "failed"} else "unknown", + "exit_code": code if type(code) is int else None, + "timed_out": rec.get("timed_out") is True, "killed": rec.get("killed") is True, + "died": rec.get("died") is True} + + def _row(rec: Dict[str, Any]) -> str: cmd = (rec.get("command") or "").strip().splitlines()[0][:80] return f"[{rec.get('id')}] {_status_label(rec)} | {_age(rec)} | {cmd}" @@ -101,17 +114,20 @@ class ManageBgJobsTool: if action in _KILL_ACTIONS: if rec.get("status") != "running": - return {"output": f"Job `{job_id}` already {_status_label(rec)}; nothing to kill.", "exit_code": 0} + return {"output": f"Job `{job_id}` already {_status_label(rec)}; nothing to kill.", "exit_code": 0, + "job": _job_facts(rec)} killed = bg_jobs.kill(job_id, expected=resource) if not killed or not killed.get("killed"): return {"error": f"Could not verify termination of background job `{job_id}`.", "exit_code": 1, "teardown": (killed or {}).get("teardown")} - return {"output": f"Killed background job `{job_id}` ({(killed or {}).get('command', '').splitlines()[0][:80]}).", "exit_code": 0} + return {"output": f"Killed background job `{job_id}` ({(killed or {}).get('command', '').splitlines()[0][:80]}).", "exit_code": 0, + "job": _job_facts(killed)} out = rec.get("output") or "(no output yet)" return { "output": f"Job `{job_id}` [{_status_label(rec)}, {_age(rec)}]\nCommand: {rec.get('command')}\n\nOutput:\n{out}", "exit_code": 0, + "job": _job_facts(rec), } return {"error": f"manage_bg_jobs: unknown action '{action}'. Use list, output, or kill.", "exit_code": 1} diff --git a/src/agent_tools/filesystem_tools.py b/src/agent_tools/filesystem_tools.py index 89d161d43..53854da61 100644 --- a/src/agent_tools/filesystem_tools.py +++ b/src/agent_tools/filesystem_tools.py @@ -156,20 +156,24 @@ class EditFileTool: if count > 1 and not replace_all: return original, None, f"not_unique:{count}" updated = original.replace(old, new) if replace_all else original.replace(old, new, 1) + attempted.append(True) with open(path, "w", encoding="utf-8", newline="") as f: f.write(updated) return original, updated, "ok" + # In-place rewrite: a failure after truncation may leave partial bytes. + attempted = [] + partial = lambda: {"mutation_attempted": True} if attempted else {} try: original, updated, status = await asyncio.to_thread(_apply) except FileNotFoundError: - return {"error": f"edit_file: {path}: not found (use write_file to create it)", "exit_code": 1} + return {"error": f"edit_file: {path}: not found (use write_file to create it)", "exit_code": 1, **partial()} except (IsADirectoryError, UnicodeDecodeError): - return {"error": f"edit_file: {path}: not an editable text file", "exit_code": 1} + return {"error": f"edit_file: {path}: not an editable text file", "exit_code": 1, **partial()} except PermissionError: - return {"error": f"edit_file: {path}: permission denied", "exit_code": 1} + return {"error": f"edit_file: {path}: permission denied", "exit_code": 1, **partial()} except OSError as e: - return {"error": f"edit_file: {path}: {e}", "exit_code": 1} + return {"error": f"edit_file: {path}: {e}", "exit_code": 1, **partial()} if status == "not_found": return {"error": f"edit_file: old_string not found in {path}. Read the file and match it exactly.", "exit_code": 1} @@ -332,6 +336,9 @@ class WriteFileTool: "exit_code": 1, "binary_artifact_preserved": target_existed, } + # This writer truncates in place. Once that stage is reached, a failure + # may leave a partial file; report it so effect evidence stays honest. + attempted = [] try: def _write(): old = "" @@ -343,14 +350,17 @@ class WriteFileTool: d = os.path.dirname(path) if d: os.makedirs(d, exist_ok=True) + attempted.append(True) with open(path, "w", encoding="utf-8") as f: f.write(body) return old, len(body) old_content, size = await asyncio.to_thread(_write) except PermissionError: - return {"error": f"write_file: {path}: permission denied", "exit_code": 1} + return {"error": f"write_file: {path}: permission denied", "exit_code": 1, + **({"mutation_attempted": True} if attempted else {})} except OSError as e: - return {"error": f"write_file: {path}: {e}", "exit_code": 1} + return {"error": f"write_file: {path}: {e}", "exit_code": 1, + **({"mutation_attempted": True} if attempted else {})} diff = _unified_diff(old_content, body, path) result = { "output": f"Wrote {size} bytes to {_display_tool_path(path)}", diff --git a/src/agent_tools/subprocess_tools.py b/src/agent_tools/subprocess_tools.py index 399755caf..fc34eb046 100644 --- a/src/agent_tools/subprocess_tools.py +++ b/src/agent_tools/subprocess_tools.py @@ -571,7 +571,7 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg "stderr": _truncate(result.stderr, MAX_OUTPUT_CHARS)} if result.timed_out: return {**common, "error": f"{tool}: timed out after {timeout}s; process tree terminated.{capture_note}", - "exit_code": 124, "stdout": _truncate(result.stdout, MAX_OUTPUT_CHARS), + "exit_code": 124, "timed_out": True, "stdout": _truncate(result.stdout, MAX_OUTPUT_CHARS), "stderr": _truncate(result.stderr, MAX_OUTPUT_CHARS)} if tool == "python": child_failure = _python_child_runtime_failure(result.stdout, result.stderr, result.exit_code) diff --git a/tests/test_effect_resource_bindings.py b/tests/test_effect_resource_bindings.py new file mode 100644 index 000000000..23a378577 --- /dev/null +++ b/tests/test_effect_resource_bindings.py @@ -0,0 +1,227 @@ +"""Wave 4 effects through the real dispatcher and exact Wave 3 filesystem bindings.""" +from __future__ import annotations + +import asyncio +import hashlib +import json +import os + +import pytest + +from src import tool_execution +from src.agent_evidence import CompletionRequirements, CompletionStatus, EvidenceKind +from src.agent_runtime import effects as fx +from src.agent_runtime.authority import OperationGrant, RequestAuthority +from src.agent_runtime.completion import _ledger, completion_answer +from src.agent_runtime.effect_log import EffectLog +from src.agent_runtime.journal import ActionJournal, bind_journal +from src.agent_tools import TOOL_HANDLERS +from src.tool_capabilities import ToolRunSecurityContext +from src.tool_types import ToolBlock + + +@pytest.fixture +def ws(tmp_path, monkeypatch): + work = tmp_path / "ws" + work.mkdir() + monkeypatch.setattr(tool_execution, "_owner_is_admin", lambda owner: True) + return work + + +@pytest.fixture +def run(ws, tmp_path): + journal = ActionJournal(workspace=str(ws), observed_artifacts=("a.txt",)) + journal.effects = EffectLog(journal.run_id, directory=tmp_path / "fx") + authority = RequestAuthority("request", "alice", "thread", str(ws), tuple( + OperationGrant(tool) for tool in ("write_file", "read_file", "edit_file", "apply_patch", "ls"))) + + async def call(tool, args): + content = args if isinstance(args, str) else json.dumps(args) + with bind_journal(journal): + return await tool_execution.execute_tool_block( + ToolBlock(tool, content), owner="alice", session_id="thread", workspace=str(ws), + security_context=ToolRunSecurityContext(external_untrusted_context_seen=False), + request_authority=authority) + + def go(tool, args): + return asyncio.run(call(tool, args)) + + go.journal = journal + return go + + +def sha(text): + return hashlib.sha256(text.encode()).hexdigest() + + +def verdicts(journal): + return [a.verdict for a in journal.effects.assessments()] + + +def ledger(journal, ws): + return _ledger(journal, CompletionRequirements(required_artifacts=("a.txt",), workspace_root=str(ws))) + + +def records(journal): + return [json.loads(line)["type"] for line in journal.effects.path.read_text().splitlines()] + + +def test_claim_is_durable_before_the_producer_runs(run, monkeypatch): + seen = [] + original = TOOL_HANDLERS["write_file"] + + async def spy(content, ctx): + seen.append(records(run.journal)) + return await original(content, ctx) + + monkeypatch.setitem(TOOL_HANDLERS, "write_file", spy) + _, result = run("write_file", {"path": "a.txt", "content": "hello\n"}) + assert result["exit_code"] == 0 + assert seen == [["claim"]], "the claim must be on disk, with no outcome, at backend invocation" + claim = run.journal.effects.history().claims[0] + assert [ref.role for ref in claim.impact_scope] == ["destination"] + assert claim.obligations[0].predicate is fx.Predicate.CONTENT_SHA256 + assert claim.obligations[0].expected == sha("hello\n") + receipt = run.journal.actions[0] + stages = [t["stage"] for t in receipt.transitions] + assert stages.index("effect_claimed") < stages.index("dispatched") + + +def test_persistence_failure_refuses_invocation(run, monkeypatch, tmp_path): + called = [] + monkeypatch.setitem(TOOL_HANDLERS, "write_file", lambda content, ctx: called.append(1)) + blocker = tmp_path / "blocker" + blocker.write_text("x") + run.journal.effects = EffectLog(run.journal.run_id, directory=blocker) + description, result = run("write_file", {"path": "a.txt", "content": "hello\n"}) + assert not called and "BLOCKED" in description and result["blocked"] is True + assert run.journal.actions[0].execution_id is None + assert not (tmp_path / "ws" / "a.txt").exists() + + +def test_execution_success_then_complete_readback_verifies(run): + run("write_file", {"path": "a.txt", "content": "hello\n"}) + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED] + run("read_file", {"path": "a.txt"}) + assert verdicts(run.journal) == [fx.EffectVerdict.VERIFIED] + observation = run.journal.effects.history().observations[0] + assert observation.source_action_id == run.journal.actions[1].action_id + assert observation.content_sha256 == sha("hello\n") + + +def test_partial_read_neither_verifies_nor_validates(run, ws): + run("write_file", {"path": "a.txt", "content": "one\ntwo\n"}) + run("read_file", {"path": "a.txt", "offset": 1, "limit": 1}) + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED] + current = ledger(run.journal, ws) + validations = [e for e in current.events if e.kind == EvidenceKind.ARTIFACT_VALIDATION] + assert validations and not any(e.authoritative for e in validations) + + +def test_later_mutation_makes_earlier_verification_stale(run): + run("write_file", {"path": "a.txt", "content": "one\n"}) + run("read_file", {"path": "a.txt"}) + run("write_file", {"path": "a.txt", "content": "two\n"}) + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED, fx.EffectVerdict.UNVERIFIED] + run("read_file", {"path": "a.txt"}) + assert verdicts(run.journal) == [fx.EffectVerdict.CONTRADICTED, fx.EffectVerdict.VERIFIED] + + +def test_unrecorded_change_is_contradicted_and_fails_completion(run, ws): + run("write_file", {"path": "a.txt", "content": "hello\n"}) + (ws / "a.txt").write_text("tampered\n") + run("read_file", {"path": "a.txt"}) + assert verdicts(run.journal) == [fx.EffectVerdict.CONTRADICTED] + decision = ledger(run.journal, ws).evaluate() + assert decision.status == CompletionStatus.FAILED and decision.missing_artifacts == ("a.txt",) + + +def test_cancelled_write_unsettles_an_earlier_success(run, ws, monkeypatch): + run("write_file", {"path": "a.txt", "content": "hello\n"}) + assert ledger(run.journal, ws).evaluate().can_complete + + async def cancelled(content, ctx): + raise asyncio.CancelledError + + monkeypatch.setitem(TOOL_HANDLERS, "write_file", cancelled) + with pytest.raises(asyncio.CancelledError): + run("write_file", {"path": "a.txt", "content": "again\n"}) + assessment = run.journal.effects.assessments()[-1] + assert assessment.execution is fx.ExecutionOutcome.CANCELLED and assessment.unresolved_impact + current = ledger(run.journal, ws) + decision = current.evaluate() + assert decision.status == CompletionStatus.BLOCKED and "settled" in decision.reason + prose, why = completion_answer("I wrote a.txt.", current, decision) + assert prose.startswith("The task is incomplete: a later operation may have changed") + assert "I wrote a.txt" not in prose and why + # The same artifact claim is unsupported by the shared ledger view. + assert not current._supports_artifact_claim(EvidenceKind.ARTIFACT_MUTATION, ("a.txt",)) + + +def test_mid_write_failure_unsettles_but_refusal_preserves(run, ws, monkeypatch): + run("write_file", {"path": "a.txt", "content": "hello\n"}) + # A deterministic refusal before the mutation stage keeps the artifact. + run("write_file", {"path": "a.txt", "content": ""}) + assert run.journal.effects.assessments()[-1].execution is fx.ExecutionOutcome.FAILED + assert ledger(run.journal, ws).evaluate().can_complete + + real_open = open + + def failing_open(path, mode="r", *args, **kwargs): + if "w" in mode and str(path).endswith("a.txt"): + handle = real_open(path, mode, *args, **kwargs) # truncates + handle.close() + raise OSError("disk full") + return real_open(path, mode, *args, **kwargs) + + monkeypatch.setattr("builtins.open", failing_open) + _, result = run("write_file", {"path": "a.txt", "content": "hello again\n"}) + monkeypatch.setattr("builtins.open", real_open) + assert result.get("mutation_attempted") is True + decision = ledger(run.journal, ws).evaluate() + assert decision.status == CompletionStatus.BLOCKED and decision.missing_artifacts == ("a.txt",) + + +def test_refused_operation_creates_no_claim(run, tmp_path): + outside = tmp_path / "outside.txt" + description, _ = run("write_file", {"path": str(outside), "content": "x"}) + assert "BLOCKED" in description + assert run.journal.effects.history().claims == () + assert not outside.exists() + + +def test_forged_producer_fields_do_not_verify(run, monkeypatch): + async def forged(content, ctx): + return {"output": "verified", "exit_code": 0, "verified": True, "content_sha256": sha("hello\n"), + "observation": {"coverage": "complete"}} + + monkeypatch.setitem(TOOL_HANDLERS, "write_file", forged) + run("write_file", {"path": "a.txt", "content": "hello\n"}) + assert run.journal.effects.history().observations == () + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED] + + +def test_patch_obligations_follow_exact_bindings(run, ws): + (ws / "old.txt").write_text("x\n") + patch = "*** Begin Patch\n*** Add File: new.txt\n+hello\n*** Delete File: old.txt\n*** End Patch" + _, result = run("apply_patch", {"patch_text": patch}) + assert result["exit_code"] == 0, result + claim = run.journal.effects.history().claims[0] + assert {o.predicate for o in claim.obligations} == {fx.Predicate.CONTENT_SHA256, fx.Predicate.ABSENT} + + +def test_listing_is_partial_and_does_not_verify_content(run): + run("write_file", {"path": "a.txt", "content": "hello\n"}) + run("ls", {"path": "."}) + observation = run.journal.effects.history().observations[0] + assert observation.coverage is fx.Coverage.PARTIAL + assert verdicts(run.journal) == [fx.EffectVerdict.UNVERIFIED] + + +def test_ordinary_read_only_turn_completes_normally(run, ws): + (ws / "a.txt").write_text("existing\n") + run("read_file", {"path": "a.txt"}) + assert run.journal.effects.history().claims == () + assert not run.journal.effects.path.exists() + current = _ledger(run.journal, CompletionRequirements(workspace_root=str(ws))) + assert current.evaluate().can_complete diff --git a/tests/test_effect_verification_adapters.py b/tests/test_effect_verification_adapters.py new file mode 100644 index 000000000..3eb175e41 --- /dev/null +++ b/tests/test_effect_verification_adapters.py @@ -0,0 +1,331 @@ +"""Wave 4 adapters for process/background, owned, external and browser bindings. + +These drive ``begin_effect``/``settle_effect`` with real Wave 3 bound-operation +objects. The dispatcher's contextvar capture is replaced by the same objects so +that each producer family can be exercised without its live backend. +""" +from __future__ import annotations + +import asyncio +import gc +import json + +import pytest + +from src import browser_identity +from src.agent_evidence import CompletionRequirements, CompletionStatus +from src.agent_runtime import effect_adapters as adapters +from src.agent_runtime import effects as fx +from src.agent_runtime.authority import ExactOperation +from src.agent_runtime.completion import _ledger +from src.agent_runtime.effect_log import EffectLog +from src.agent_runtime.journal import ActionJournal +from src.agent_runtime.owned_resources import BoundOwnedOperation +from src.agent_runtime.process_resources import BoundProcessOperation, digest as process_digest +from src.agent_runtime.remote_resources import BoundBackendOperation +from src.agent_runtime.resources import ( + BackgroundJobResource, BrowserPageResource, BrowserSessionObservation, BrowserSessionResource, ExternalResource, + FilesystemRoot, NativeBackendResource, OwnedResource, ProcessLaunchResource, ProcessLaunchScope, ProcessResource, +) +from src.process_lifecycle import ProcessIdentity +from src.tool_types import ToolBlock + + +GENERATION = "c" * 32 + + +@pytest.fixture +def store(tmp_path): + return tmp_path / "fx" + + +def journal_for(store, parent=None): + journal = ActionJournal(parent_run_id=parent.run_id if parent else None) + journal.effects = parent.effects if parent else EffectLog(journal.run_id, directory=store) + return journal + + +def act(journal, monkeypatch, capture, tool="bash", content="{}", *, result=None, error=None): + """One admitted action: claim at dispatch, then settle with a producer result.""" + action = journal.propose(ToolBlock(tool, content)) + monkeypatch.setattr(adapters, "capture_dispatch", lambda: capture) + captured = adapters.begin_effect(journal, action) + action.execution_id = action.action_id + ":execution:1" + if result is not None: + action.finish(result) + adapters.settle_effect(journal, action, captured, result=result, error=error) + return action, captured + + +def launch_capture(tmp_path, generation=GENERATION): + workspace = tmp_path / "ws" + workspace.mkdir(exist_ok=True) + operation = ExactOperation.normalize("bash", "#!bg\nsleep 1") + scope = ProcessLaunchScope(NativeBackendResource("bash"), FilesystemRoot.seal(str(workspace)), frozenset({"filesystem"})) + launch = ProcessLaunchResource("native:containment", "alice", "request", "thread", generation, "bash", + process_digest(operation.input), scope, "b" * 64) + return adapters.DispatchCapture(process=BoundProcessOperation(operation, "request", "alice", "thread", launch)) + + +def job_capture(action="status", generation=GENERATION, job_id="job1"): + supervisor = ProcessResource("native:bg_jobs", "alice", "request", "thread", + ProcessIdentity(4242, "boot:1:100", None), "supervisor", job_id, "cont-1") + job = BackgroundJobResource("native:bg_jobs", job_id, generation, "alice", "request", "thread", "cont-1", + (supervisor,)) + operation = ExactOperation.normalize("manage_bg_jobs", json.dumps({"action": action, "job_id": job_id})) + return adapters.DispatchCapture(process=BoundProcessOperation(operation, "request", "alice", "thread", jobs=(job,))) + + +def job_result(status, exit_code=None, **flags): + return {"output": "Job report says everything succeeded and was verified.", "exit_code": 0, + "job": {"status": status, "exit_code": exit_code, "timed_out": False, "killed": False, + "died": False, **flags}} + + +# -- process / background ---------------------------------------------------- + +def test_process_exit_is_execution_evidence_not_a_postcondition(tmp_path, store, monkeypatch): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "ls", + result={"output": "ok", "exit_code": 0, "teardown": {"dead": True}}) + claim = journal.effects.history().claims[0] + assert claim.unknown_scope, "an arbitrary command has unknown impact scope" + assert [ref.kind for ref in claim.dependencies] == [fx.ResourceKind.PROCESS_LAUNCH] + assessment = journal.effects.assessments()[0] + assert (assessment.execution, assessment.verdict, assessment.cleanup) == ( + fx.ExecutionOutcome.REPORTED_SUCCESS, fx.EffectVerdict.UNVERIFIED, fx.CleanupState.VERIFIED) + + +@pytest.mark.parametrize("result,execution,cleanup", [ + ({"error": "timed out", "exit_code": 124, "timed_out": True, "teardown": {"dead": True}}, + fx.ExecutionOutcome.TIMED_OUT, fx.CleanupState.VERIFIED), + ({"error": "teardown", "exit_code": 1, "failure_kind": "process_teardown_failed", "teardown": {"dead": False}}, + fx.ExecutionOutcome.FAILED, fx.CleanupState.FAILED), + ({"output": "", "exit_code": 0, "status": "running", "detached": True, "containment": {"external": True}}, + fx.ExecutionOutcome.RUNNING, fx.CleanupState.UNKNOWN), +]) +def test_process_outcomes_are_preserved_separately(tmp_path, store, monkeypatch, result, execution, cleanup): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "x", result=result) + assessment = journal.effects.assessments()[0] + assert (assessment.execution, assessment.cleanup) == (execution, cleanup) + assert assessment.unresolved_impact + + +def test_cleanup_failure_after_command_unsettles_required_artifact(tmp_path, store, monkeypatch): + workspace = tmp_path / "ws" + journal = journal_for(store) + journal.workspace, journal.observed_artifacts = str(workspace), ("out.txt",) + write = journal.propose(ToolBlock("write_file", json.dumps({"path": "out.txt", "content": "x"}))) + write.execution_id = write.action_id + ":execution:1" + write.finish({"output": "Wrote", "exit_code": 0}) + (workspace).mkdir(exist_ok=True) + (workspace / "out.txt").write_text("x") + requirements = CompletionRequirements(required_artifacts=("out.txt",), workspace_root=str(workspace)) + assert _ledger(journal, requirements).evaluate().can_complete + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "x", + result={"error": "teardown", "exit_code": 1, "failure_kind": "process_teardown_failed", + "teardown": {"dead": False}}) + decision = _ledger(journal, requirements).evaluate() + assert decision.status == CompletionStatus.BLOCKED and decision.missing_artifacts == ("out.txt",) + + +def test_background_launch_is_running_not_completed_work(tmp_path, store, monkeypatch): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started background job `job1`.", "exit_code": 0, "bg_job_id": "job1"}) + assessment = journal.effects.assessments()[0] + assert (assessment.execution, assessment.verdict) == (fx.ExecutionOutcome.RUNNING, fx.EffectVerdict.PENDING) + + +def test_exact_job_read_settles_launch_across_a_continuation_run(tmp_path, store, monkeypatch): + first = journal_for(store) + act(first, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started", "exit_code": 0, "bg_job_id": "job1"}) + launch_effect = first.effects.history().claims[0] + # A still-running job does not settle anything. + second = journal_for(store) + act(second, monkeypatch, job_capture(), "manage_bg_jobs", "{}", result=job_result("running")) + assert fx.assess(launch_effect, first.effects.history()).verdict is fx.EffectVerdict.PENDING + assert second.effects.history().observations[0].mechanism is fx.ObservationMechanism.JOB_STATE + + del first + gc.collect() # the launching run is gone: settle through its durable log + third = journal_for(store) + act(third, monkeypatch, job_capture(), "manage_bg_jobs", "{}", result=job_result("done", 0)) + reloaded = EffectLog.load(launch_effect.run_id, directory=store) + assessment = fx.assess(launch_effect, reloaded.history()) + # Delivered completion is execution evidence; the job's report prose is + # attributed content and verifies nothing. + assert (assessment.execution, assessment.verdict) == (fx.ExecutionOutcome.REPORTED_SUCCESS, + fx.EffectVerdict.UNVERIFIED) + + +def test_job_linkage_requires_the_exact_generation(tmp_path, store, monkeypatch): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started", "exit_code": 0, "bg_job_id": "job1"}) + # Same display job id, different launch generation: a replacement job. + act(journal, monkeypatch, job_capture(generation="d" * 32), "manage_bg_jobs", "{}", + result=job_result("done", 0)) + assert journal.effects.assessments()[0].execution is fx.ExecutionOutcome.RUNNING + + +def test_killed_job_settles_as_cancelled(tmp_path, store, monkeypatch): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started", "exit_code": 0, "bg_job_id": "job1"}) + act(journal, monkeypatch, job_capture("kill"), "manage_bg_jobs", "{}", + result={**job_result("failed", -9, killed=True), "output": "Killed"}) + kill_claim = journal.effects.history().claims[1] + assert [ref.kind for ref in kill_claim.impact_scope] == [fx.ResourceKind.BACKGROUND_JOB, fx.ResourceKind.PROCESS] + assert journal.effects.assessments()[0].execution is fx.ExecutionOutcome.CANCELLED + + +def test_running_background_work_unsettles_later_required_artifact(tmp_path, store, monkeypatch): + workspace = tmp_path / "ws" + workspace.mkdir() + (workspace / "out.txt").write_text("x") + journal = journal_for(store) + journal.workspace = str(workspace) + write = journal.propose(ToolBlock("write_file", json.dumps({"path": "out.txt", "content": "x"}))) + write.execution_id = write.action_id + ":execution:1" + write.finish({"output": "Wrote", "exit_code": 0}) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started", "exit_code": 0, "bg_job_id": "job1"}) + requirements = CompletionRequirements(required_artifacts=("out.txt",), workspace_root=str(workspace)) + assert _ledger(journal, requirements).evaluate().status == CompletionStatus.BLOCKED + + +# -- owned records ------------------------------------------------------------ + +def owned_capture(tool, payload, record): + operation = ExactOperation.normalize(tool, json.dumps(payload)) + return adapters.DispatchCapture(owned=BoundOwnedOperation(operation, operation.input, "request", "alice", + "thread", (record,))) + + +def test_owned_mutation_claims_exact_record_and_stays_unverified(store, monkeypatch): + record = OwnedResource("notes", "alice", "thread", "notes", "n1", "rev-1") + journal = journal_for(store) + act(journal, monkeypatch, owned_capture("manage_notes", {"action": "update", "id": "n1"}, record), + "manage_notes", result={"output": "Note updated and verified.", "exit_code": 0}) + claim = journal.effects.history().claims[0] + assert claim.impact_scope == (fx.resource_ref(record, "record"),) + assert journal.effects.assessments()[0].verdict is fx.EffectVerdict.UNVERIFIED + + +def test_same_display_id_new_revision_does_not_inherit_freshness(store, monkeypatch): + journal = journal_for(store) + old = OwnedResource("vault", "alice", "thread", "vault", "rec", "rev-1") + new = OwnedResource("vault", "alice", "thread", "vault", "rec", "rev-2") + act(journal, monkeypatch, owned_capture("vault_get", {"id": "rec"}, old), "vault_get", + result={"output": "secret", "exit_code": 0}) + act(journal, monkeypatch, owned_capture("vault_get", {"id": "rec"}, new), "vault_get", + result={"output": "secret", "exit_code": 0}) + history = journal.effects.history() + first, second = history.observations + assert history.claims == () + assert fx.freshness(first, history) is fx.Freshness.STALE + assert fx.freshness(second, history) is fx.Freshness.FRESH + + +# -- external / MCP ------------------------------------------------------------- + +def test_remote_success_is_acknowledgement_not_state(store, monkeypatch): + remote = ExternalResource("mcp", "endpoint", "server", "tool", "inc-1") + bound = BoundBackendOperation(remote, "request", "alice", "thread", "mcp__server__tool", "{}") + journal = journal_for(store) + act(journal, monkeypatch, adapters.DispatchCapture(backend=bound), "mcp__server__tool", + result={"output": "Successfully created and verified the record.", "exit_code": 0}) + claim = journal.effects.history().claims[0] + assert claim.external and claim.impact_scope == (fx.resource_ref(remote, "backend"),) + outcome = journal.effects.history().outcomes[0] + assert outcome.facts.remote_acknowledged and outcome.facts.external + assert outcome.cleanup is fx.CleanupState.UNKNOWN + assert journal.effects.history().observations == () + assert journal.effects.assessments()[0].verdict is fx.EffectVerdict.UNVERIFIED + + +def test_remote_failure_after_send_may_have_changed_state(store, monkeypatch): + remote = ExternalResource("mcp", "endpoint", "server", "tool", "inc-1") + bound = BoundBackendOperation(remote, "request", "alice", "thread", "mcp__server__tool", "{}") + journal = journal_for(store) + act(journal, monkeypatch, adapters.DispatchCapture(backend=bound), "mcp__server__tool", + error=TimeoutError("transport closed after send")) + assessment = journal.effects.assessments()[0] + assert assessment.execution is fx.ExecutionOutcome.INTERRUPTED and assessment.unresolved_impact + + +# -- browser session metadata only -------------------------------------------- + +def session(monkeypatch, incarnation_seed="1"): + monkeypatch.setattr(browser_identity, "PRODUCER_HASHES", {"linux-x64": "e" * 64}) + values = {"producer_namespace": "native:agent-browser", "producer_version": "0.35.0", "platform": "linux-x64", + "binary_sha256": "e" * 64, "configuration_digest": "1" * 64, "session_key": "ody-" + "a" * 24, + "daemon": {"pid": 4321, "start_token": "boot:" + incarnation_seed, "pgid": 4321}, + "browser_instance_digest": incarnation_seed * 64} + observation = BrowserSessionObservation(**{**values, "daemon": ProcessIdentity(4321, "boot:" + incarnation_seed, 4321), + "session_incarnation": browser_identity.incarnation(values)}) + return BrowserSessionResource("alice", "thread", observation) + + +def test_browser_session_info_is_lifecycle_observation_only(store, monkeypatch): + journal = journal_for(store) + operation = ExactOperation.normalize("private_browser", json.dumps({"action": "session_info"})) + first = browser_identity.BoundBrowserOperation(operation, "request", "alice", "thread", session(monkeypatch, "1")) + replaced = browser_identity.BoundBrowserOperation(operation, "request", "alice", "thread", session(monkeypatch, "2")) + for bound in (first, replaced): + act(journal, monkeypatch, adapters.DispatchCapture(browser=bound), "private_browser", + result={"output": "{}", "exit_code": 0, "executed": True, "browser_page_operations_supported": False}) + history = journal.effects.history() + assert history.claims == () + assert {o.mechanism for o in history.observations} == {fx.ObservationMechanism.BROWSER_SESSION} + # Session replacement never transfers freshness to the new session. + assert fx.freshness(history.observations[0], history) is fx.Freshness.STALE + # A session observation decides no file/record/remote postcondition. + assert all(o.coverage is fx.Coverage.PARTIAL for o in history.observations) + + +def test_browser_page_binding_never_becomes_effect_scope(store, monkeypatch): + journal = journal_for(store) + owner_session = session(monkeypatch) + page = BrowserPageResource(owner_session, "A" * 32, "loader") + operation = ExactOperation.normalize("private_browser", json.dumps({"action": "session_info"})) + bound = browser_identity.BoundBrowserOperation(operation, "request", "alice", "thread", owner_session, page) + act(journal, monkeypatch, adapters.DispatchCapture(browser=bound), "private_browser", + result={"output": "{}", "exit_code": 0}) + claim = journal.effects.history().claims[0] + assert claim.unknown_scope and journal.effects.history().observations == () + with pytest.raises(TypeError): + fx.resource_ref(page, "target") + + +# -- lineage -------------------------------------------------------------------- + +def test_child_effects_share_lineage_order_and_invalidate_parent_evidence(tmp_path, store, monkeypatch): + parent = journal_for(store) + old = OwnedResource("vault", "alice", "thread", "vault", "rec", "rev-1") + act(parent, monkeypatch, owned_capture("vault_get", {"id": "rec"}, old), "vault_get", + result={"output": "x", "exit_code": 0}) + child = journal_for(store, parent=parent) + act(child, monkeypatch, launch_capture(tmp_path), "bash", "x", result={"output": "", "exit_code": 0}) + history = parent.effects.history() + claim = history.claims[0] + assert (claim.run_id, claim.parent_run_id) == (child.run_id, parent.run_id) + # The child's unknown-scope command may have changed the parent's record. + assert fx.freshness(history.observations[0], history) is fx.Freshness.STALE + + +def test_classification_failure_claims_unknown_scope(store, monkeypatch): + journal = journal_for(store) + monkeypatch.setattr(adapters, "classify", lambda capture: (_ for _ in ()).throw(KeyError("bug"))) + act(journal, monkeypatch, adapters.DispatchCapture(), "anything", result={"output": "", "exit_code": 0}) + assert journal.effects.history().claims[0].unknown_scope + + +def test_cancellation_is_recorded_without_inventing_a_result(tmp_path, store, monkeypatch): + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "x", error=asyncio.CancelledError()) + assessment = journal.effects.assessments()[0] + assert (assessment.execution, assessment.cleanup) == (fx.ExecutionOutcome.CANCELLED, fx.CleanupState.UNKNOWN)