diff --git a/src/agent_runtime/effects.py b/src/agent_runtime/effects.py new file mode 100644 index 000000000..e69572eec --- /dev/null +++ b/src/agent_runtime/effects.py @@ -0,0 +1,800 @@ +"""Wave 4 effect claims, outcomes, observations and verification. + +This module consumes exact Wave 3 resource identities. It never resolves a +selector, discovers an alias, grants an operation or performs I/O. A +``ResourceRef`` can only be built from an already-admitted typed Wave 3 resource +object; names, paths, PIDs, URLs, labels and dictionaries are not accepted. + +Facts are kept separate: + +* a claim records intent and scope before backend invocation, not dispatch; +* an outcome records what the executor reported, not the resulting state; +* an observation records state seen through an admitted mechanism; +* verification is derived from fresh, relevant, complete observations made + after the effect settled, and never from receipts or acknowledgements. + +History is append-only. Invalidation and freshness are computed from the +ordered record history; earlier records are never rewritten. Refresh is a new +observation. Unknown scope is conservative, never "no impact". +""" +from __future__ import annotations + +from dataclasses import dataclass +from enum import Enum +from pathlib import PurePosixPath +import hashlib +import json +import re +from typing import Any, Iterable, Mapping + + +def _sha(value: Any) -> str: + return hashlib.sha256(json.dumps(value, sort_keys=True, separators=(",", ":"), + ensure_ascii=False, default=str).encode()).hexdigest() + + +def _text(value: Any, label: str, *, optional: bool = False) -> None: + if (not isinstance(value, str) or (not value and not optional) + or any(c in value for c in ("\0", "\n", "\r"))): + raise ValueError(f"Invalid effect {label}") + + +def _position(value: Any) -> None: + if type(value) is not int or value < 0: + raise ValueError("Effect history position must be a nonnegative integer") + + +_SHA256 = re.compile(r"[a-f0-9]{64}") + + +# --------------------------------------------------------------------------- +# Exact resource references (Wave 3 consumption only) +# --------------------------------------------------------------------------- + +class ResourceKind(str, Enum): + FILESYSTEM = "filesystem" + PROCESS = "process" + BACKGROUND_JOB = "background_job" + OWNED = "owned" + EXTERNAL = "external" + BROWSER_SESSION = "browser_session" + + +@dataclass(frozen=True) +class ResourceRef: + """Historical reference to one exact admitted Wave 3 resource. + + ``location`` identifies where the resource lives (including the identity of + its sealed root/namespace); ``incarnation`` identifies the object observed + there when the reference was taken. Replacement keeps the location and + changes the incarnation, so evidence never transfers to a replacement. + ``snapshot_sha256`` digests the full Wave 3 snapshot for audit. A ref is not + authority: it is not accepted by any dispatcher, resolver or grant. + """ + + kind: ResourceKind + role: str + location: tuple[str, ...] + incarnation: str + snapshot_sha256: str + + def __post_init__(self) -> None: + if not isinstance(self.kind, ResourceKind): + raise ValueError("Unsupported effect resource kind") + _text(self.role, "resource role") + _text(self.incarnation, "resource incarnation", optional=True) + if (not isinstance(self.location, tuple) or len(self.location) < 2 + or any(not isinstance(part, str) or any(c in part for c in ("\0", "\n", "\r")) + for part in self.location) + or self.location[0] != self.kind.value): + raise ValueError("Malformed effect resource location") + if not _SHA256.fullmatch(self.snapshot_sha256 or ""): + raise ValueError("Malformed effect resource snapshot digest") + + @property + def location_key(self) -> str: + return _sha(list(self.location)) + + def same_location(self, other: "ResourceRef") -> bool: + return self.kind is other.kind and self.location == other.location + + def overlaps(self, other: "ResourceRef") -> bool: + """Conservative relevance between two exact references. + + Filesystem relevance is ancestor-or-self within one sealed root + identity: a mutation of ``d/x`` invalidates a listing of ``d`` and a + replacement of ``d`` invalidates observations of ``d/x``. Other kinds + only overlap at the same exact location. No alias discovery is done. + """ + if self.kind is not other.kind: + return False + if self.kind is not ResourceKind.FILESYSTEM: + return self.location == other.location + if self.location[:-1] != other.location[:-1]: + return False + left, right = PurePosixPath(self.location[-1]), PurePosixPath(other.location[-1]) + return left == right or left.is_relative_to(right) or right.is_relative_to(left) + + def to_dict(self) -> dict[str, Any]: + return {"kind": self.kind.value, "role": self.role, "location": list(self.location), + "incarnation": self.incarnation, "snapshot_sha256": self.snapshot_sha256} + + @classmethod + def from_dict(cls, value: Any) -> "ResourceRef": + """Reload a persisted historical reference. This creates no authority.""" + if (not isinstance(value, dict) + or set(value) != {"kind", "role", "location", "incarnation", "snapshot_sha256"} + or not isinstance(value["location"], list)): + raise ValueError("Malformed persisted effect resource reference") + return cls(ResourceKind(value["kind"]), value["role"], tuple(value["location"]), + value["incarnation"], value["snapshot_sha256"]) + + +def resource_ref(resource: Any, role: str) -> ResourceRef: + """Reference an exact typed Wave 3 resource; anything else is refused. + + Browser page/document resources are refused: Wave 3 fails closed for page + authority and Wave 4 must not promote page observations into identity. + """ + from src.agent_runtime import resources as wave3 + if isinstance(resource, wave3.BrowserPageResource): + raise TypeError("Browser page resources are not effect-bindable") + if isinstance(resource, wave3.FilesystemResource): + root = resource.root + location = ("filesystem", root.scope.value, root.owner, root.path, + str(root.identity.device), str(root.identity.inode), resource.path) + chain = [[a.path, a.identity.device, a.identity.inode] for a in resource.ancestors] + identity = resource.identity + incarnation = ("absent:" + _sha(chain) if identity is None else + f"{identity.kind}:{identity.device}:{identity.inode}:" + _sha(chain)) + return ResourceRef(ResourceKind.FILESYSTEM, role, location, incarnation, _sha(resource.to_dict())) + if isinstance(resource, wave3.ProcessResource): + ident = resource.identity + 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.BackgroundJobResource): + location = ("background_job", resource.namespace, resource.owner, resource.request_id, + resource.thread_id, resource.job_id, resource.generation) + return ResourceRef(ResourceKind.BACKGROUND_JOB, role, location, resource.generation, + _sha(resource.to_dict())) + if isinstance(resource, wave3.OwnedResource): + location = ("owned", resource.namespace, resource.owner, resource.thread_id, + resource.collection, resource.record_id) + return ResourceRef(ResourceKind.OWNED, role, location, resource.revision, _sha(resource.to_dict())) + if isinstance(resource, wave3.ExternalResource): + location = ("external", resource.namespace, resource.owner, resource.endpoint_id, + resource.server_id, resource.tool_id) + return ResourceRef(ResourceKind.EXTERNAL, role, location, resource.incarnation, _sha(resource.to_dict())) + if isinstance(resource, wave3.BrowserSessionResource): + observation = resource.observation + location = ("browser_session", resource.owner, resource.thread_id, observation.session_key) + return ResourceRef(ResourceKind.BROWSER_SESSION, role, location, observation.session_incarnation, + _sha(resource.to_dict())) + raise TypeError("Effect scope requires an exact Wave 3 resource identity") + + +def bound_filesystem_refs(bound: Any) -> tuple[ResourceRef, ...]: + """References for an admitted ``BoundFilesystemOperation``'s exact bindings.""" + from src.agent_runtime.resource_binding import BoundFilesystemOperation + if not isinstance(bound, BoundFilesystemOperation): + raise TypeError("Filesystem effect scope requires a server-owned bound operation") + return tuple(resource_ref(binding.resource, binding.role) for binding in bound.bindings) + + +# --------------------------------------------------------------------------- +# Claims +# --------------------------------------------------------------------------- + +@dataclass(frozen=True) +class OperationRef: + """Final normalized operation reference; not a second normalization API.""" + + tool: str + action: str + input_sha256: str + request_id: str = "" + + def __post_init__(self) -> None: + _text(self.tool, "operation tool") + _text(self.action, "operation action", optional=True) + _text(self.request_id, "operation request", optional=True) + if not _SHA256.fullmatch(self.input_sha256 or ""): + raise ValueError("Malformed operation input digest") + + @classmethod + def from_exact(cls, operation: Any, execution_input: str | None = None, request_id: str = "") -> "OperationRef": + from src.agent_runtime.authority import ExactOperation + if not isinstance(operation, ExactOperation): + raise TypeError("Effect claims require the admitted exact operation") + body = operation.input if execution_input is None else execution_input + return cls(str(operation.tool), str(operation.action or ""), _sha(body), request_id or "") + + def to_dict(self) -> dict[str, Any]: + return {"tool": self.tool, "action": self.action, "input_sha256": self.input_sha256, + "request_id": self.request_id} + + @classmethod + def from_dict(cls, value: Any) -> "OperationRef": + if not isinstance(value, dict) or set(value) != {"tool", "action", "input_sha256", "request_id"}: + raise ValueError("Malformed persisted operation reference") + return cls(**value) + + +class Predicate(str, Enum): + EXISTS = "exists" + ABSENT = "absent" + CONTENT_SHA256 = "content_sha256" + # The observed content digest differs from ``expected`` (the pre-state). + CONTENT_CHANGED = "content_changed" + + +_PREDICATE_KINDS = { + Predicate.EXISTS: {ResourceKind.FILESYSTEM, ResourceKind.OWNED, ResourceKind.EXTERNAL}, + Predicate.ABSENT: {ResourceKind.FILESYSTEM, ResourceKind.OWNED, ResourceKind.EXTERNAL}, + Predicate.CONTENT_SHA256: {ResourceKind.FILESYSTEM, ResourceKind.OWNED, ResourceKind.EXTERNAL}, + Predicate.CONTENT_CHANGED: {ResourceKind.FILESYSTEM, ResourceKind.OWNED, ResourceKind.EXTERNAL}, +} + + +@dataclass(frozen=True) +class Postcondition: + """An explicit requested post-state predicate on one exact claimed target.""" + + target: ResourceRef + predicate: Predicate + expected: str = "" + + def __post_init__(self) -> None: + if not isinstance(self.target, ResourceRef) or not isinstance(self.predicate, Predicate): + raise ValueError("Malformed postcondition") + if self.target.kind not in _PREDICATE_KINDS[self.predicate]: + raise ValueError("Predicate is not supported for this resource kind") + needs_digest = self.predicate in {Predicate.CONTENT_SHA256, Predicate.CONTENT_CHANGED} + if needs_digest != bool(_SHA256.fullmatch(self.expected or "")) or (not needs_digest and self.expected): + raise ValueError("Malformed postcondition expectation") + + def to_dict(self) -> dict[str, Any]: + return {"target": self.target.to_dict(), "predicate": self.predicate.value, "expected": self.expected} + + @classmethod + def from_dict(cls, value: Any) -> "Postcondition": + if not isinstance(value, dict) or set(value) != {"target", "predicate", "expected"}: + raise ValueError("Malformed persisted postcondition") + return cls(ResourceRef.from_dict(value["target"]), Predicate(value["predicate"]), value["expected"]) + + +@dataclass(frozen=True) +class EffectClaim: + """Server-owned claim, persisted before backend invocation. + + The claim states intent and scope; it is not evidence that dispatch, the + backend operation, or any mutation happened. ``impact_scope`` holds the + exact admitted bindings the operation may change; empty means unknown + scope, never no impact. ``dependencies`` are resources the predicate + relies on without being mutation targets. + """ + + effect_id: str + run_id: str + action_id: str + sequence: int + operation: OperationRef + impact_scope: tuple[ResourceRef, ...] = () + dependencies: tuple[ResourceRef, ...] = () + obligations: tuple[Postcondition, ...] = () + parent_run_id: str = "" + external: bool = False + + def __post_init__(self) -> None: + for name in ("effect_id", "run_id", "action_id"): + _text(getattr(self, name), name) + _text(self.parent_run_id, "parent run", optional=True) + _position(self.sequence) + if not isinstance(self.operation, OperationRef) or type(self.external) is not bool: + raise ValueError("Malformed effect claim") + for name in ("impact_scope", "dependencies"): + refs = getattr(self, name) + if not isinstance(refs, tuple) or any(not isinstance(r, ResourceRef) for r in refs): + raise ValueError("Effect scope must be exact resource references") + if (not isinstance(self.obligations, tuple) + or any(not isinstance(o, Postcondition) for o in self.obligations)): + raise ValueError("Malformed effect obligations") + for obligation in self.obligations: + if not any(obligation.target == ref for ref in self.impact_scope): + raise ValueError("Postcondition target must be a claimed impact binding") + + @property + def unknown_scope(self) -> bool: + return not self.impact_scope + + def to_dict(self) -> dict[str, Any]: + return {"effect_id": self.effect_id, "run_id": self.run_id, "action_id": self.action_id, + "sequence": self.sequence, "operation": self.operation.to_dict(), + "impact_scope": [r.to_dict() for r in self.impact_scope], + "dependencies": [r.to_dict() for r in self.dependencies], + "obligations": [o.to_dict() for o in self.obligations], + "parent_run_id": self.parent_run_id, "external": self.external} + + @classmethod + def from_dict(cls, value: Any) -> "EffectClaim": + keys = {"effect_id", "run_id", "action_id", "sequence", "operation", "impact_scope", + "dependencies", "obligations", "parent_run_id", "external"} + if not isinstance(value, dict) or set(value) != keys or any( + not isinstance(value[k], list) for k in ("impact_scope", "dependencies", "obligations")): + raise ValueError("Malformed persisted effect claim") + return cls(value["effect_id"], value["run_id"], value["action_id"], value["sequence"], + OperationRef.from_dict(value["operation"]), + tuple(ResourceRef.from_dict(r) for r in value["impact_scope"]), + tuple(ResourceRef.from_dict(r) for r in value["dependencies"]), + tuple(Postcondition.from_dict(o) for o in value["obligations"]), + value["parent_run_id"], value["external"]) + + +# --------------------------------------------------------------------------- +# Outcomes +# --------------------------------------------------------------------------- + +class ExecutionOutcome(str, Enum): + NOT_EXECUTED = "not_executed" # refused before backend invocation + ATTEMPTED = "attempted" # claimed; no settled outcome yet + REPORTED_SUCCESS = "reported_success" # executor reported success; not post-state + FAILED = "failed" + TIMED_OUT = "timed_out" + CANCELLED = "cancelled" + RUNNING = "running" # admitted/background; not completed work + INTERRUPTED = "interrupted" # unknown: lost, crashed or replayed + + +class Impact(str, Enum): + NONE = "none" # known no-op: the backend was never invoked + POSSIBLE = "possible" # may have changed state, including partially + CHANGED = "changed" # a trusted before/after capture differs + + +class CleanupState(str, Enum): + NOT_APPLICABLE = "not_applicable" + VERIFIED = "verified" + FAILED = "failed" + UNKNOWN = "unknown" + + +_SETTLED = {ExecutionOutcome.NOT_EXECUTED, ExecutionOutcome.REPORTED_SUCCESS, ExecutionOutcome.FAILED, + ExecutionOutcome.TIMED_OUT, ExecutionOutcome.CANCELLED, ExecutionOutcome.INTERRUPTED} + + +@dataclass(frozen=True) +class ProducerFacts: + """Bounded typed producer facts; arbitrary returned data is never kept. + + These are execution/lifecycle facts reported by a server producer. None of + them is a post-state observation. + """ + + exit_code: int | None = None + timed_out: bool = False + output_truncated: bool = False + failure_kind: str = "" + job_state: str = "" + remote_acknowledged: bool = False + external: 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"): + if type(getattr(self, name)) is not bool: + raise ValueError("Malformed producer flag") + for name in ("failure_kind", "job_state"): + value = getattr(self, name) + _text(value, name, optional=True) + if len(value) > 64 or (value and not re.fullmatch(r"[a-z0-9_.:-]+", value)): + raise ValueError("Malformed producer label") + + def to_dict(self) -> dict[str, Any]: + 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} + + @classmethod + def from_dict(cls, value: Any) -> "ProducerFacts": + if not isinstance(value, dict) or set(value) != set(cls.__dataclass_fields__): + raise ValueError("Malformed persisted producer facts") + return cls(**value) + + +def _label(value: Any) -> str: + text = value.strip().lower() if isinstance(value, str) else "" + return text if len(text) <= 64 and re.fullmatch(r"[a-z0-9_.:-]+", text) else "" + + +def producer_facts(result: Any) -> ProducerFacts: + """Project a dispatcher result into typed facts without trusting its shape. + + Only exact scalar types are copied. Anything else becomes the default, so a + forged or malformed dictionary can only lose information, not add trust. + """ + if not isinstance(result, Mapping): + return ProducerFacts() + code = result.get("exit_code") + containment = result.get("containment") + external = isinstance(containment, Mapping) and containment.get("external") is True + job = result.get("status") if isinstance(result.get("job_id"), str) else "" + return ProducerFacts( + exit_code=code if type(code) is int else None, + timed_out=result.get("timed_out") is True or _label(result.get("failure_kind")) == "timeout", + output_truncated=result.get("output_truncated") is True or result.get("truncated") is True, + failure_kind=_label(result.get("failure_kind")), + job_state=_label(job), + external=external, + ) + + +@dataclass(frozen=True) +class EffectOutcome: + """Append-only execution outcome for one claim. + + ``impact`` must not claim no change for anything that reached a backend. + ``cleanup`` is recorded separately: cleanup success is not business-effect + success and cleanup failure does not erase an achieved effect. + """ + + effect_id: str + sequence: int + execution: ExecutionOutcome + impact: Impact + facts: ProducerFacts = ProducerFacts() + cleanup: CleanupState = CleanupState.NOT_APPLICABLE + execution_id: str = "" + replayed: bool = False + + def __post_init__(self) -> None: + _text(self.effect_id, "effect identifier") + _text(self.execution_id, "execution identifier", optional=True) + _position(self.sequence) + if (not isinstance(self.execution, ExecutionOutcome) or not isinstance(self.impact, Impact) + or not isinstance(self.facts, ProducerFacts) or not isinstance(self.cleanup, CleanupState) + or type(self.replayed) is not bool): + raise ValueError("Malformed effect outcome") + if self.execution is ExecutionOutcome.ATTEMPTED: + raise ValueError("ATTEMPTED is derived from a claim without an outcome") + if (self.impact is Impact.NONE) != (self.execution is ExecutionOutcome.NOT_EXECUTED): + raise ValueError("Only a refused, never-invoked operation is a known no-op") + if self.execution is ExecutionOutcome.NOT_EXECUTED and self.execution_id: + raise ValueError("A refused operation has no execution identity") + + def to_dict(self) -> dict[str, Any]: + return {"effect_id": self.effect_id, "sequence": self.sequence, "execution": self.execution.value, + "impact": self.impact.value, "facts": self.facts.to_dict(), "cleanup": self.cleanup.value, + "execution_id": self.execution_id, "replayed": self.replayed} + + @classmethod + def from_dict(cls, value: Any) -> "EffectOutcome": + if not isinstance(value, dict) or set(value) != set(cls.__dataclass_fields__): + raise ValueError("Malformed persisted effect outcome") + return cls(value["effect_id"], value["sequence"], ExecutionOutcome(value["execution"]), + Impact(value["impact"]), ProducerFacts.from_dict(value["facts"]), + CleanupState(value["cleanup"]), value["execution_id"], value["replayed"]) + + +# --------------------------------------------------------------------------- +# Observations +# --------------------------------------------------------------------------- + +class ObservationMechanism(str, Enum): + FILESYSTEM_READ = "filesystem_read" # admitted read of the exact binding + OWNED_RECORD_READ = "owned_record_read" # admitted owner-scoped readback + REMOTE_READBACK = "remote_readback" # admitted independent remote query + PROCESS_OWNERSHIP = "process_ownership" # lifecycle owner's verdict + JOB_STATE = "job_state" # background job record transition + BROWSER_SESSION = "browser_session" # session lifecycle metadata only + # The following are never post-state verification. + EXECUTION_RECEIPT = "execution_receipt" + REMOTE_ACKNOWLEDGEMENT = "remote_acknowledgement" + + +class Coverage(str, Enum): + COMPLETE = "complete" + PARTIAL = "partial" + + +# Mechanisms able to decide a postcondition for each resource kind. Process, +# job and browser-session observations are lifecycle facts: they can make +# earlier evidence stale but cannot verify a file/record/remote predicate. +_VERIFYING = { + ResourceKind.FILESYSTEM: {ObservationMechanism.FILESYSTEM_READ}, + ResourceKind.OWNED: {ObservationMechanism.OWNED_RECORD_READ}, + ResourceKind.EXTERNAL: {ObservationMechanism.REMOTE_READBACK}, +} +_ADMITTED_READS = {ObservationMechanism.FILESYSTEM_READ, ObservationMechanism.OWNED_RECORD_READ, + ObservationMechanism.REMOTE_READBACK} + + +@dataclass(frozen=True) +class Observation: + """State seen through one mechanism for one exact resource. + + ``exists``/``content_sha256`` are what the mechanism saw; ``None``/empty + means not observed. A PARTIAL observation (offset/limit/truncated read, + listing, existence-only probe) never decides a whole-content predicate. + Admitted reads must name the journal action that performed them. + """ + + observation_id: str + sequence: int + resource: ResourceRef + mechanism: ObservationMechanism + coverage: Coverage + source_action_id: str = "" + source_execution_id: str = "" + exists: bool | None = None + content_sha256: str = "" + evidence_event_id: str = "" + + def __post_init__(self) -> None: + _text(self.observation_id, "observation identifier") + for name in ("source_action_id", "source_execution_id", "evidence_event_id"): + _text(getattr(self, name), name, optional=True) + _position(self.sequence) + if (not isinstance(self.resource, ResourceRef) or not isinstance(self.mechanism, ObservationMechanism) + or not isinstance(self.coverage, Coverage) + or (self.exists is not None and type(self.exists) is not bool)): + raise ValueError("Malformed observation") + if self.content_sha256 and (not _SHA256.fullmatch(self.content_sha256) or self.exists is not True): + raise ValueError("Malformed observed content digest") + if self.mechanism in _ADMITTED_READS and not self.source_action_id: + raise ValueError("Readback observations require the admitted action that performed them") + + def to_dict(self) -> dict[str, Any]: + return {"observation_id": self.observation_id, "sequence": self.sequence, + "resource": self.resource.to_dict(), "mechanism": self.mechanism.value, + "coverage": self.coverage.value, "source_action_id": self.source_action_id, + "source_execution_id": self.source_execution_id, "exists": self.exists, + "content_sha256": self.content_sha256, "evidence_event_id": self.evidence_event_id} + + @classmethod + def from_dict(cls, value: Any) -> "Observation": + if not isinstance(value, dict) or set(value) != set(cls.__dataclass_fields__): + raise ValueError("Malformed persisted observation") + return cls(**{**value, "resource": ResourceRef.from_dict(value["resource"]), + "mechanism": ObservationMechanism(value["mechanism"]), + "coverage": Coverage(value["coverage"])}) + + +def predicate_holds(postcondition: Postcondition, observation: Observation) -> bool | None: + """Decide one predicate from one observation; ``None`` means undecidable. + + The check is performed here from the observed state, so no adapter can + attest verification by labelling an unrelated read. + """ + target = postcondition.target + if (not observation.resource.same_location(target) + or observation.mechanism not in _VERIFYING.get(target.kind, set())): + return None + predicate = postcondition.predicate + if predicate is Predicate.ABSENT: + return None if observation.exists is None else not observation.exists + if predicate is Predicate.EXISTS: + return observation.exists + if observation.exists is False: + return False + if observation.coverage is not Coverage.COMPLETE or not observation.content_sha256: + return None + if predicate is Predicate.CONTENT_SHA256: + return observation.content_sha256 == postcondition.expected + return observation.content_sha256 != postcondition.expected + + +# --------------------------------------------------------------------------- +# History, invalidation and freshness +# --------------------------------------------------------------------------- + +class Freshness(str, Enum): + FRESH = "fresh" + STALE = "stale" # a later possible mutation or replacement overlaps + UNSETTLED = "unsettled" # an overlapping effect was still in flight + + +@dataclass(frozen=True) +class EffectHistory: + """An immutable, totally ordered view of one effect log. + + Sequences are unique positions in one log. Duplicate positions are rejected + rather than ordered arbitrarily. + """ + + claims: tuple[EffectClaim, ...] = () + outcomes: tuple[EffectOutcome, ...] = () + observations: tuple[Observation, ...] = () + + def __post_init__(self) -> None: + positions = [r.sequence for r in (*self.claims, *self.outcomes, *self.observations)] + if len(positions) != len(set(positions)): + raise ValueError("Effect history positions must be unique") + ids = [c.effect_id for c in self.claims] + if len(ids) != len(set(ids)): + raise ValueError("Effect claims must have unique identifiers") + claim_at = {c.effect_id: c.sequence for c in self.claims} + settled: set[str] = set() + for outcome in sorted(self.outcomes, key=lambda o: o.sequence): + if outcome.effect_id not in claim_at or outcome.sequence <= claim_at[outcome.effect_id]: + raise ValueError("Outcome must follow its claim in one history") + # A RUNNING effect may later settle (background continuation or + # replay interruption); a settled outcome is never replaced. + if outcome.effect_id in settled: + raise ValueError("A settled effect outcome cannot be replaced") + if outcome.execution is not ExecutionOutcome.RUNNING: + settled.add(outcome.effect_id) + + def claim(self, effect_id: str) -> EffectClaim | None: + return next((c for c in self.claims if c.effect_id == effect_id), None) + + def latest_outcome(self, effect_id: str, before: int | None = None) -> EffectOutcome | None: + matching = [o for o in self.outcomes if o.effect_id == effect_id + and (before is None or o.sequence < before)] + return max(matching, key=lambda o: o.sequence) if matching else None + + def execution(self, effect_id: str, before: int | None = None) -> ExecutionOutcome: + outcome = self.latest_outcome(effect_id, before) + return ExecutionOutcome.ATTEMPTED if outcome is None else outcome.execution + + +def _claim_touches(claim: EffectClaim, resource: ResourceRef) -> bool: + return claim.unknown_scope or any(ref.overlaps(resource) for ref in claim.impact_scope) + + +def invalidated_by(observation: Observation, history: EffectHistory) -> tuple[str, ...]: + """Identifiers of later records that make ``observation`` stale. + + Any later claim that may touch the resource invalidates it once the claim + exists (it may already be executing), unless it settled as a known no-op. + A later observation of the same location with a different incarnation + reveals replacement. Execution receipts are never invalidated: they remain + historical execution facts. + """ + if observation.mechanism in {ObservationMechanism.EXECUTION_RECEIPT, + ObservationMechanism.REMOTE_ACKNOWLEDGEMENT}: + return () + reasons: list[str] = [] + for claim in history.claims: + if claim.sequence <= observation.sequence or not _claim_touches(claim, observation.resource): + continue + outcome = history.latest_outcome(claim.effect_id) + if outcome is not None and outcome.impact is Impact.NONE: + continue + reasons.append(claim.effect_id) + for later in history.observations: + if (later.sequence > observation.sequence and later.resource.same_location(observation.resource) + and later.resource.incarnation != observation.resource.incarnation): + reasons.append(later.observation_id) + return tuple(dict.fromkeys(reasons)) + + +def freshness(observation: Observation, history: EffectHistory) -> Freshness: + if invalidated_by(observation, history): + return Freshness.STALE + for claim in history.claims: + if claim.sequence < observation.sequence and _claim_touches(claim, observation.resource): + state = history.execution(claim.effect_id, before=observation.sequence) + if state in {ExecutionOutcome.ATTEMPTED, ExecutionOutcome.RUNNING}: + return Freshness.UNSETTLED + return Freshness.FRESH + + +# --------------------------------------------------------------------------- +# Verification +# --------------------------------------------------------------------------- + +class EffectVerdict(str, Enum): + NOT_EXECUTED = "not_executed" + PENDING = "pending" # attempted/running; not settled + VERIFIED = "verified" # reported success + fresh matching post-state + STATE_OBSERVED = "state_observed" # matching post-state; causality unknown + UNVERIFIED = "unverified" # no adequate fresh evidence + CONTRADICTED = "contradicted" # latest fresh check shows the predicate false + FAILED = "failed" # execution failed; never effect success + + +_VERDICT_RANK = {EffectVerdict.FAILED: 0, EffectVerdict.CONTRADICTED: 1, EffectVerdict.PENDING: 2, + EffectVerdict.UNVERIFIED: 3, EffectVerdict.NOT_EXECUTED: 4, + EffectVerdict.STATE_OBSERVED: 5, EffectVerdict.VERIFIED: 6} + + +@dataclass(frozen=True) +class EffectAssessment: + effect_id: str + action_id: str + execution: ExecutionOutcome + impact: Impact | None + verdict: EffectVerdict + reason: str + cleanup: CleanupState = CleanupState.NOT_APPLICABLE + observation_ids: tuple[str, ...] = () + targets: tuple[ResourceRef, ...] = () + + @property + def unresolved_impact(self) -> bool: + """Resources may have changed in a way no fresh evidence has settled.""" + return (self.impact is not Impact.NONE + and self.execution is not ExecutionOutcome.REPORTED_SUCCESS + and self.verdict not in {EffectVerdict.STATE_OBSERVED, EffectVerdict.CONTRADICTED}) + + def to_dict(self) -> dict[str, Any]: + return {"effect_id": self.effect_id, "action_id": self.action_id, "execution": self.execution.value, + "impact": None if self.impact is None else self.impact.value, "verdict": self.verdict.value, + "reason": self.reason, "cleanup": self.cleanup.value, + "observation_ids": list(self.observation_ids), + "targets": [t.to_dict() for t in self.targets]} + + +def _assess_obligation(claim: EffectClaim, settled: EffectOutcome, obligation: Postcondition, + history: EffectHistory) -> tuple[EffectVerdict, str, str]: + candidates = [o for o in history.observations + if o.sequence > settled.sequence and o.resource.same_location(obligation.target) + and o.mechanism in _VERIFYING.get(obligation.target.kind, set())] + if not candidates: + return EffectVerdict.UNVERIFIED, "no authorized post-settlement observation of the target", "" + # The newest check wins. A newer partial or failed check never falls back + # to an earlier complete one. + latest = max(candidates, key=lambda o: o.sequence) + state = freshness(latest, history) + if state is not Freshness.FRESH: + return EffectVerdict.UNVERIFIED, f"the latest target observation is {state.value}", latest.observation_id + holds = predicate_holds(obligation, latest) + if holds is None: + return EffectVerdict.UNVERIFIED, "the latest observation does not decide the postcondition", latest.observation_id + if not holds: + return EffectVerdict.CONTRADICTED, "the latest fresh observation contradicts the postcondition", latest.observation_id + execution = settled.execution + if execution is ExecutionOutcome.FAILED: + return EffectVerdict.FAILED, "execution failed; matching state is not attributed to it", latest.observation_id + if execution is ExecutionOutcome.REPORTED_SUCCESS: + return EffectVerdict.VERIFIED, "fresh authorized observation matches the postcondition", latest.observation_id + return (EffectVerdict.STATE_OBSERVED, + "state matches, but this execution's outcome is unknown; causality is not established", + latest.observation_id) + + +def assess(claim: EffectClaim, history: EffectHistory) -> EffectAssessment: + """Derive a claim's verdict from the append-only history.""" + settled = history.latest_outcome(claim.effect_id) + targets = tuple(o.target for o in claim.obligations) + if settled is None: + return EffectAssessment(claim.effect_id, claim.action_id, ExecutionOutcome.ATTEMPTED, None, + EffectVerdict.PENDING, "no settled execution outcome", targets=targets) + base = dict(effect_id=claim.effect_id, action_id=claim.action_id, execution=settled.execution, + impact=settled.impact, cleanup=settled.cleanup, targets=targets) + if settled.execution is ExecutionOutcome.NOT_EXECUTED: + return EffectAssessment(**base, verdict=EffectVerdict.NOT_EXECUTED, reason="refused before invocation") + if settled.execution is ExecutionOutcome.RUNNING: + return EffectAssessment(**base, verdict=EffectVerdict.PENDING, + reason="background execution has not settled") + if not claim.obligations: + verdict = EffectVerdict.FAILED if settled.execution is ExecutionOutcome.FAILED else EffectVerdict.UNVERIFIED + return EffectAssessment(**base, verdict=verdict, reason="no explicit postcondition obligation") + results = [_assess_obligation(claim, settled, o, history) for o in claim.obligations] + worst = min(results, key=lambda r: _VERDICT_RANK[r[0]]) + if settled.execution is ExecutionOutcome.FAILED and worst[0] is not EffectVerdict.CONTRADICTED: + worst = (EffectVerdict.FAILED, worst[1] if worst[0] is EffectVerdict.FAILED else + "execution failed and may have partially changed the target", worst[2]) + return EffectAssessment(**base, verdict=worst[0], reason=worst[1], + observation_ids=tuple(dict.fromkeys(r[2] for r in results if r[2]))) + + +def assess_all(history: EffectHistory) -> tuple[EffectAssessment, ...]: + return tuple(assess(claim, history) for claim in sorted(history.claims, key=lambda c: c.sequence)) + + +def replay_interrupted(history: EffectHistory, next_sequence: int) -> tuple[EffectOutcome, ...]: + """Outcomes to append for claims that never settled before a reload. + + Unknown remains unknown: the backend may or may not have been invoked, so + impact is POSSIBLE. Running background effects are left to their own + lifecycle owner and are not converted here. + """ + _position(next_sequence) + pending = [c for c in sorted(history.claims, key=lambda c: c.sequence) + if history.latest_outcome(c.effect_id) is None] + return tuple(EffectOutcome(c.effect_id, next_sequence + i, ExecutionOutcome.INTERRUPTED, + Impact.POSSIBLE, replayed=True) for i, c in enumerate(pending)) diff --git a/tests/test_effects_foundation.py b/tests/test_effects_foundation.py new file mode 100644 index 000000000..7d0976243 --- /dev/null +++ b/tests/test_effects_foundation.py @@ -0,0 +1,466 @@ +"""Wave 4 effect semantics against exact Wave 3 resource identities.""" +from __future__ import annotations + +import hashlib +import os + +import pytest + +from src.agent_runtime import effects as fx +from src.agent_runtime.authority import ExactOperation +from src.agent_runtime.resources import ( + BrowserPageResource, ExternalResource, FilesystemResource, FilesystemRoot, OwnedResource, ProcessResource, +) +from src.process_lifecycle import ProcessIdentity + + +def sha(data: bytes) -> str: + return hashlib.sha256(data).hexdigest() + + +@pytest.fixture +def root(tmp_path): + workspace = tmp_path / "ws" + workspace.mkdir() + return FilesystemRoot.seal(str(workspace)) + + +def fs_ref(root, name, role="target", *, missing=False): + resource = FilesystemResource.resolve(root, os.path.join(root.path, name), allow_missing=missing) + return fx.resource_ref(resource, role) + + +def op(tool="write_file"): + return fx.OperationRef("write_file" if tool == "write_file" else tool, "", "0" * 64) + + +class Log: + """Test helper assigning one total order, like the runtime effect log.""" + + def __init__(self): + self.claims, self.outcomes, self.observations, self.seq = [], [], [], 0 + + def _next(self): + self.seq += 1 + return self.seq + + def claim(self, scope=(), obligations=(), effect_id=None, dependencies=()): + claim = fx.EffectClaim(effect_id or f"e{len(self.claims) + 1}", "run", f"a{len(self.claims) + 1}", + self._next(), op(), tuple(scope), tuple(dependencies), tuple(obligations)) + self.claims.append(claim) + return claim + + def outcome(self, claim, execution=fx.ExecutionOutcome.REPORTED_SUCCESS, impact=fx.Impact.POSSIBLE, **kw): + outcome = fx.EffectOutcome(claim.effect_id, self._next(), execution, impact, + execution_id="" if impact is fx.Impact.NONE else claim.action_id + ":x", **kw) + self.outcomes.append(outcome) + return outcome + + def observe(self, resource, *, exists=True, digest="", coverage=fx.Coverage.COMPLETE, + mechanism=fx.ObservationMechanism.FILESYSTEM_READ): + observation = fx.Observation(f"o{len(self.observations) + 1}", self._next(), resource, mechanism, coverage, + source_action_id="read-action", exists=exists, content_sha256=digest) + self.observations.append(observation) + return observation + + @property + def history(self): + return fx.EffectHistory(tuple(self.claims), tuple(self.outcomes), tuple(self.observations)) + + +def content(target, body=b"hello"): + return fx.Postcondition(target, fx.Predicate.CONTENT_SHA256, sha(body)) + + +# -- exact resource references -------------------------------------------- + +def test_refs_only_accept_typed_wave3_resources(root): + for forged in ({"kind": "filesystem", "path": "/etc/passwd"}, "/workspace/a.txt", 1234, + ("filesystem", "x")): + with pytest.raises(TypeError): + fx.resource_ref(forged, "target") + page = object.__new__(BrowserPageResource) + with pytest.raises(TypeError, match="page"): + fx.resource_ref(page, "target") + + +def test_filesystem_replacement_changes_incarnation_not_location(root): + path = os.path.join(root.path, "a.txt") + with open(path, "w") as handle: + handle.write("one") + before = fs_ref(root, "a.txt") + os.replace(_write(root, "tmp", "two"), path) + after = fs_ref(root, "a.txt") + assert before.same_location(after) + assert before.incarnation != after.incarnation + assert fs_ref(root, "missing.txt", missing=True).incarnation.startswith("absent:") + + +def _write(root, name, text): + path = os.path.join(root.path, name) + with open(path, "w") as handle: + handle.write(text) + return path + + +def test_filesystem_overlap_is_ancestor_or_self_within_one_sealed_root(root, tmp_path): + os.mkdir(os.path.join(root.path, "d")) + _write(root, "d/x.txt", "x") + _write(root, "dx.txt", "x") + directory = fs_ref(root, "d", "search_root") + child = fs_ref(root, "d/x.txt") + sibling = fs_ref(root, "dx.txt") + assert directory.overlaps(child) and child.overlaps(directory) + assert not sibling.overlaps(directory) + other_dir = tmp_path / "other" + other_dir.mkdir() + (other_dir / "d").mkdir() + other = FilesystemRoot.seal(str(other_dir)) + assert not fx.resource_ref(FilesystemResource.resolve(other, str(other_dir / "d")), "target").overlaps(directory) + + +def test_process_pid_reuse_is_a_different_location(): + first = ProcessResource("native:containment", "u", "r", "t", ProcessIdentity(4242, "boot:1:100", None), "leader") + reused = ProcessResource("native:containment", "u", "r", "t", ProcessIdentity(4242, "boot:1:999", None), "leader") + a, b = fx.resource_ref(first, "subject"), fx.resource_ref(reused, "subject") + assert not a.same_location(b) and not a.overlaps(b) + + +def test_owned_revision_is_incarnation_and_external_never_contained(): + v1 = fx.resource_ref(OwnedResource("notes", "u", "t", "notes", "n1", "rev-1"), "target") + v2 = fx.resource_ref(OwnedResource("notes", "u", "t", "notes", "n1", "rev-2"), "target") + assert v1.same_location(v2) and v1.incarnation != v2.incarnation + remote = fx.resource_ref(ExternalResource("mcp", "ep", "srv", "tool", "inc-1"), "target") + assert remote.kind is fx.ResourceKind.EXTERNAL + + +def test_ref_round_trip_is_historical_and_strict(root): + ref = fs_ref(root, "a.txt", missing=True) + assert fx.ResourceRef.from_dict(ref.to_dict()) == ref + with pytest.raises(ValueError): + fx.ResourceRef.from_dict({**ref.to_dict(), "extra": 1}) + with pytest.raises(ValueError): + fx.ResourceRef.from_dict({**ref.to_dict(), "location": ["owned", "x"]}) + + +def test_operation_ref_requires_admitted_exact_operation(): + with pytest.raises(TypeError): + fx.OperationRef.from_exact({"tool": "write_file"}) + exact = ExactOperation.normalize("write_file", '{"path": "a.txt", "content": "x"}') + assert fx.OperationRef.from_exact(exact).tool == "write_file" + + +# -- claims and outcomes --------------------------------------------------- + +def test_claim_obligations_must_target_claimed_scope(root): + target, other = fs_ref(root, "a.txt", missing=True), fs_ref(root, "b.txt", missing=True) + with pytest.raises(ValueError, match="claimed impact"): + fx.EffectClaim("e", "run", "a", 0, op(), (target,), (), (content(other),)) + claim = fx.EffectClaim("e", "run", "a", 0, op(), (target,), (), (content(target),)) + assert fx.EffectClaim.from_dict(claim.to_dict()) == claim + assert fx.EffectClaim("e", "run", "a", 0, op()).unknown_scope + + +def test_known_noop_only_for_refusal_before_invocation(): + with pytest.raises(ValueError): + fx.EffectOutcome("e", 1, fx.ExecutionOutcome.FAILED, fx.Impact.NONE) + with pytest.raises(ValueError): + fx.EffectOutcome("e", 1, fx.ExecutionOutcome.REPORTED_SUCCESS, fx.Impact.NONE) + with pytest.raises(ValueError): + fx.EffectOutcome("e", 1, fx.ExecutionOutcome.NOT_EXECUTED, fx.Impact.POSSIBLE) + with pytest.raises(ValueError, match="derived"): + fx.EffectOutcome("e", 1, fx.ExecutionOutcome.ATTEMPTED, fx.Impact.POSSIBLE) + assert fx.EffectOutcome("e", 1, fx.ExecutionOutcome.NOT_EXECUTED, fx.Impact.NONE).impact is fx.Impact.NONE + + +def test_forged_producer_dictionaries_cannot_add_trust(): + forged = {"exit_code": True, "timed_out": "yes", "failure_kind": "x\ny", "status": "finished", + "containment": {"external": "true"}, "verified": True, "postcondition": "ok"} + facts = fx.producer_facts(forged) + assert facts == fx.ProducerFacts() + assert fx.producer_facts(["not", "a", "mapping"]) == fx.ProducerFacts() + real = fx.producer_facts({"exit_code": 0, "job_id": "j1", "status": "running", + "containment": {"external": True}, "output_truncated": True}) + assert (real.exit_code, real.job_state, real.external, real.output_truncated) == (0, "running", True, True) + + +def test_history_rejects_ambiguous_order_and_replaced_outcomes(root): + log = Log() + claim = log.claim((fs_ref(root, "a.txt", missing=True),)) + log.outcome(claim) + with pytest.raises(ValueError, match="cannot be replaced"): + log.outcome(claim, fx.ExecutionOutcome.FAILED) + log.history + clash = fx.Observation("o", claim.sequence, claim.impact_scope[0], fx.ObservationMechanism.FILESYSTEM_READ, + fx.Coverage.COMPLETE, source_action_id="r") + with pytest.raises(ValueError, match="unique"): + fx.EffectHistory((claim,), (), (clash,)) + + +# -- verification ---------------------------------------------------------- + +def test_execution_success_is_not_verification(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.UNVERIFIED + + +def test_fresh_complete_readback_verifies_reported_success(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + observed = log.observe(target, digest=sha(b"hello")) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.VERIFIED + assert assessment.observation_ids == (observed.observation_id,) + + +def test_receipts_and_acknowledgements_never_verify(root): + for mechanism in (fx.ObservationMechanism.EXECUTION_RECEIPT, fx.ObservationMechanism.REMOTE_ACKNOWLEDGEMENT, + fx.ObservationMechanism.PROCESS_OWNERSHIP, fx.ObservationMechanism.JOB_STATE): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + log.observe(target, digest=sha(b"hello"), mechanism=mechanism) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + + +def test_independent_remote_readback_differs_from_acknowledgement(): + remote = fx.resource_ref(ExternalResource("mcp", "ep", "srv", "tool", "inc"), "target") + log = Log() + claim = log.claim((remote,), (fx.Postcondition(remote, fx.Predicate.EXISTS),)) + log.outcome(claim, facts=fx.ProducerFacts(exit_code=0, remote_acknowledged=True, external=True)) + log.observe(remote, mechanism=fx.ObservationMechanism.REMOTE_ACKNOWLEDGEMENT) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + log.observe(remote, mechanism=fx.ObservationMechanism.REMOTE_READBACK) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.VERIFIED + + +def test_verifier_before_mutation_does_not_count(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + log.observe(target, digest=sha(b"hello")) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + + +def test_observation_while_effect_in_flight_is_unsettled(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + early = log.observe(target, digest=sha(b"hello")) + log.outcome(claim) + assert fx.freshness(early, log.history) is fx.Freshness.UNSETTLED + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + + +def test_stale_evidence_after_later_mutation_is_preserved_but_stale(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + observed = log.observe(target, digest=sha(b"hello")) + later = log.claim((target,)) + log.outcome(later, fx.ExecutionOutcome.FAILED) + history = log.history + assert observed in history.observations # history is never rewritten + assert fx.invalidated_by(observed, history) == (later.effect_id,) + assert fx.assess(claim, history).verdict is fx.EffectVerdict.UNVERIFIED + + +def test_concurrent_mutation_between_effect_and_verification(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + racing = log.claim((target,)) + log.observe(target, digest=sha(b"hello")) + log.outcome(racing) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + + +def test_unknown_mutation_scope_invalidates_everything_earlier(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + observed = log.observe(target, digest=sha(b"hello")) + unknown = log.claim(()) + log.outcome(unknown, fx.ExecutionOutcome.INTERRUPTED) + assert fx.invalidated_by(observed, log.history) == (unknown.effect_id,) + + +def test_refused_operation_is_a_known_noop_and_preserves_freshness(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + log.observe(target, digest=sha(b"hello")) + refused = log.claim((target,)) + log.outcome(refused, fx.ExecutionOutcome.NOT_EXECUTED, fx.Impact.NONE) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.VERIFIED + assert fx.assess(refused, log.history).verdict is fx.EffectVerdict.NOT_EXECUTED + + +def test_unrelated_resource_mutation_does_not_invalidate(root): + log = Log() + target, other = fs_ref(root, "a.txt", missing=True), fs_ref(root, "b.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + log.observe(target, digest=sha(b"hello")) + log.outcome(log.claim((other,))) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.VERIFIED + + +def test_replacement_revealed_by_later_observation_makes_earlier_stale(root): + _write(root, "a.txt", "hello") + original = fs_ref(root, "a.txt") + log = Log() + claim = log.claim((original,), (content(original),)) + log.outcome(claim) + first = log.observe(original, digest=sha(b"hello")) + os.replace(_write(root, "tmp", "hello"), os.path.join(root.path, "a.txt")) + replacement = fs_ref(root, "a.txt") + second = log.observe(replacement, digest=sha(b"hello")) + assert fx.invalidated_by(first, log.history) == (second.observation_id,) + # The latest check is of the replacement: same bytes, still fresh. + assert fx.freshness(second, log.history) is fx.Freshness.FRESH + + +def test_partial_read_cannot_verify_whole_content_and_blocks_fallback(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + log.observe(target, digest=sha(b"hello")) + log.observe(target, coverage=fx.Coverage.PARTIAL) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + # A partial read can still decide existence. + exists = fx.Postcondition(target, fx.Predicate.EXISTS) + log2 = Log() + claim2 = log2.claim((target,), (exists,)) + log2.outcome(claim2) + log2.observe(target, coverage=fx.Coverage.PARTIAL) + assert fx.assess(claim2, log2.history).verdict is fx.EffectVerdict.VERIFIED + + +def test_partial_verifier_coverage_of_multiple_obligations(root): + log = Log() + a, b = fs_ref(root, "a.txt", missing=True), fs_ref(root, "b.txt", missing=True) + claim = log.claim((a, b), (content(a), content(b))) + log.outcome(claim) + log.observe(a, digest=sha(b"hello")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + log.observe(b, digest=sha(b"hello")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.VERIFIED + + +def test_contradicting_fresh_observation(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim) + log.observe(target, digest=sha(b"other")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.CONTRADICTED + + +def test_unknown_execution_matching_state_is_not_causation(root): + for execution in (fx.ExecutionOutcome.INTERRUPTED, fx.ExecutionOutcome.TIMED_OUT, + fx.ExecutionOutcome.CANCELLED): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim, execution) + log.observe(target, digest=sha(b"hello")) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.STATE_OBSERVED + assert "causality" in assessment.reason + + +def test_failed_execution_never_becomes_success(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim, fx.ExecutionOutcome.FAILED) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.FAILED and assessment.unresolved_impact + log.observe(target, digest=sha(b"hello")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.FAILED + + +def test_unknown_partial_effect_has_unresolved_impact(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim, fx.ExecutionOutcome.TIMED_OUT, facts=fx.ProducerFacts(timed_out=True)) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.UNVERIFIED and assessment.unresolved_impact + + +def test_cleanup_failure_is_preserved_separately_from_effect(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim, cleanup=fx.CleanupState.FAILED) + log.observe(target, digest=sha(b"hello")) + assessment = fx.assess(claim, log.history) + assert assessment.verdict is fx.EffectVerdict.VERIFIED + assert assessment.cleanup is fx.CleanupState.FAILED + + +def test_background_running_is_pending_until_settled(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + log.outcome(claim, fx.ExecutionOutcome.RUNNING, facts=fx.ProducerFacts(job_state="running")) + log.observe(target, digest=sha(b"hello")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.PENDING + log.outcome(claim) + # The observation predates settlement; a new one is required. + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.UNVERIFIED + log.observe(target, digest=sha(b"hello")) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.VERIFIED + + +def test_interrupted_claim_replays_as_unknown_never_success(root): + log = Log() + target = fs_ref(root, "a.txt", missing=True) + claim = log.claim((target,), (content(target),)) + assert fx.assess(claim, log.history).verdict is fx.EffectVerdict.PENDING + appended = fx.replay_interrupted(log.history, log.seq + 1) + assert [o.execution for o in appended] == [fx.ExecutionOutcome.INTERRUPTED] + assert appended[0].impact is fx.Impact.POSSIBLE and appended[0].replayed + history = fx.EffectHistory((claim,), appended, ()) + assert fx.assess(claim, history).unresolved_impact + assert fx.replay_interrupted(history, 99) == () + + +def test_stale_owned_revision_does_not_verify(): + v1 = fx.resource_ref(OwnedResource("notes", "u", "t", "notes", "n1", "rev-1"), "target") + v2 = fx.resource_ref(OwnedResource("notes", "u", "t", "notes", "n1", "rev-2"), "target") + log = Log() + claim = log.claim((v1,), (fx.Postcondition(v1, fx.Predicate.EXISTS),)) + log.outcome(claim) + first = log.observe(v1, mechanism=fx.ObservationMechanism.OWNED_RECORD_READ) + second = log.observe(v2, mechanism=fx.ObservationMechanism.OWNED_RECORD_READ) + # The old revision's readback never inherits freshness once a different + # revision of the same display ID is observed. + assert fx.invalidated_by(first, log.history) == (second.observation_id,) + assert fx.freshness(first, log.history) is fx.Freshness.STALE + # Verification rests only on the newest readback, of the current revision. + assert fx.assess(claim, log.history).observation_ids == (second.observation_id,) + + +def test_readbacks_require_admitted_source_action(root): + target = fs_ref(root, "a.txt", missing=True) + with pytest.raises(ValueError, match="admitted action"): + fx.Observation("o", 1, target, fx.ObservationMechanism.FILESYSTEM_READ, fx.Coverage.COMPLETE) + with pytest.raises(ValueError): + fx.Observation("o", 1, target, fx.ObservationMechanism.FILESYSTEM_READ, fx.Coverage.COMPLETE, + source_action_id="a", exists=False, content_sha256=sha(b"x"))