mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-07 07:22:21 +02:00
feat(runtime): claim effects before dispatch and gate completion on them
Dispatcher seam: mark_dispatch, which runs inside the live Wave 3 binding scope immediately before backend invocation, now durably claims a possible effect before execution_id is assigned. If the claim cannot be persisted the action stays undispatched and the dispatcher returns BLOCKED; dispatched() closes the never-awaited coroutine. record_action appends the outcome (including cancellation/interruption) and admitted-read observations before the receipt reduction drops producer facts. Adapters consume only the bound operations the dispatcher admitted: filesystem bindings give exact scope and predicates (write_file content digest after fence unwrapping, apply_patch add/delete, edit existence); bash/python launches have unknown scope with the launch generation as lineage; job kills scope the exact job and its processes; owned operations scope their exact revisioned records; external backends are claimed as external and never verified by acknowledgement; browser session_info yields session lifecycle observations only, and a page binding is never effect scope. Complete read_file re-reads the exact bound source to digest it; offset/limit, truncation, extraction and listings are partial. Background launches stay RUNNING until an admitted read of the exact job generation (via a durable launch index, across continuation runs) reports settlement. Producer seams: typed job lifecycle facts on manage_bg_jobs reads/kills, a structured timed_out flag on containment timeouts, and mutation_attempted on in-place write_file/edit_file failures after truncation. Completion: the existing EvidenceLedger consumes effect assessments through a single helper used for the decision, ask_user and prose filtering. A required artifact is unsettled by a later unresolved effect that may have touched it, a fresh contradicting readback fails the decision, and partial reads no longer count as artifact validation. Ordinary conversation and read-only turns are unchanged; no second completion policy is introduced.
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
@@ -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:
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user