From 3953ea244463775e68d2f1eb795c5df18a21baaa Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Fri, 2 Oct 2026 20:10:46 +0100 Subject: [PATCH] feat(runtime): settle background launch effects from validated job lifecycle The background monitor already validates the exact Wave 3 job linkage (job_from_record + validate_job) before continuing a session. At that point it now records the job's settlement against the durable launch claim through the launch-generation index, using typed lifecycle facts from the server-owned record. Settlement is idempotent across deferred retries, is execution evidence only, and never reads the delivered output: the injected report stays attributed content. Failure to record leaves the claim running/unknown and never blocks the follow-up. --- src/agent_runtime/effect_adapters.py | 28 ++++++++++++++++------ src/agent_tools/bg_job_tools.py | 8 +++---- src/bg_monitor.py | 17 +++++++++++++ tests/test_effect_verification_adapters.py | 22 +++++++++++++++++ 4 files changed, 64 insertions(+), 11 deletions(-) diff --git a/src/agent_runtime/effect_adapters.py b/src/agent_runtime/effect_adapters.py index bd8ba3d50..34f1fbfc7 100644 --- a/src/agent_runtime/effect_adapters.py +++ b/src/agent_runtime/effect_adapters.py @@ -330,20 +330,34 @@ def _settle_background(log: Any, capture: DispatchCapture, result: dict) -> None job_facts = result.get("job") if not isinstance(job_facts, dict) or len(capture.process.jobs) != 1: return + settle_background_job(capture.process.jobs[0], job_facts, log=log) + + +def settle_background_job(job: Any, job_facts: Any, *, log: Any = None) -> None: + """Settle the RUNNING launch claim of one exact, Wave 3-validated job. + + ``job`` must be a ``BackgroundJobResource`` the caller obtained through + Wave 3 validation (an admitted job read, or the monitor's + ``job_from_record``/``validate_job``). ``job_facts`` are typed lifecycle + facts from that server-owned record; delivered output is never consulted. + """ + from src.agent_runtime.effect_log import EffectLog, EffectPersistenceError, effects_dir + from src.agent_runtime.resources import BackgroundJobResource + if not isinstance(job, BackgroundJobResource) or not isinstance(job_facts, dict): + 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: + owner = log if log is not None and 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: # 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) + directory = log.path.parent if log is not None and log.path is not None else effects_dir() + indexed = EffectLog.launch_owner(job.generation, directory=directory) if indexed is not None: try: - owner = EffectLog.open(indexed[0], directory=log.path.parent) + owner = EffectLog.open(indexed[0], directory=directory) except (EffectPersistenceError, ValueError): owner = None if owner is None: diff --git a/src/agent_tools/bg_job_tools.py b/src/agent_tools/bg_job_tools.py index 0f7260025..abab4d7b5 100644 --- a/src/agent_tools/bg_job_tools.py +++ b/src/agent_tools/bg_job_tools.py @@ -44,7 +44,7 @@ def _status_label(rec: Dict[str, Any]) -> str: return status -def _job_facts(rec: Dict[str, Any]) -> Dict[str, Any]: +def job_lifecycle_facts(rec: Dict[str, Any]) -> Dict[str, Any]: """Typed lifecycle facts from the exact admitted job record. Execution evidence only: completion of a job is not verification of any @@ -115,19 +115,19 @@ class ManageBgJobsTool: if action in _KILL_ACTIONS: if rec.get("status") != "running": return {"output": f"Job `{job_id}` already {_status_label(rec)}; nothing to kill.", "exit_code": 0, - "job": _job_facts(rec)} + "job": job_lifecycle_facts(rec)} killed = bg_jobs.kill(job_id, expected=resource) if not killed or not killed.get("killed"): return {"error": f"Could not verify termination of background job `{job_id}`.", "exit_code": 1, "teardown": (killed or {}).get("teardown")} return {"output": f"Killed background job `{job_id}` ({(killed or {}).get('command', '').splitlines()[0][:80]}).", "exit_code": 0, - "job": _job_facts(killed)} + "job": job_lifecycle_facts(killed)} out = rec.get("output") or "(no output yet)" return { "output": f"Job `{job_id}` [{_status_label(rec)}, {_age(rec)}]\nCommand: {rec.get('command')}\n\nOutput:\n{out}", "exit_code": 0, - "job": _job_facts(rec), + "job": job_lifecycle_facts(rec), } return {"error": f"manage_bg_jobs: unknown action '{action}'. Use list, output, or kill.", "exit_code": 1} diff --git a/src/bg_monitor.py b/src/bg_monitor.py index d8e3288ea..790783bc9 100644 --- a/src/bg_monitor.py +++ b/src/bg_monitor.py @@ -36,6 +36,22 @@ def _background_result_message(rec): return untrusted_context_message("background job output", inject) +def _settle_launch_effect(resource, rec): + """Record the exact job's settlement against its durable launch claim. + + Uses only the Wave 3-validated job identity and typed lifecycle facts from + the server-owned record. Settlement is execution evidence; the delivered + output remains attributed content and verifies nothing. Best-effort: a + failure leaves the claim running/unknown and never blocks the follow-up. + """ + try: + from src.agent_runtime.effect_adapters import settle_background_job + from src.agent_tools.bg_job_tools import job_lifecycle_facts + settle_background_job(resource, job_lifecycle_facts(rec)) + except Exception as error: # noqa: BLE001 + logger.warning("bg-followup: effect settlement for %s was not recorded: %s", rec.get("id"), error) + + async def _drain_agent(sess, messages, request_authority=None): """Run the agent loop headless against a session. Returns (final_prose, tool_events) — tool_events in the same shape the live chat @@ -146,6 +162,7 @@ async def _run_followup(rec: dict) -> bool: try: resource = job_from_record(rec) validate_job(resource) + _settle_launch_effect(resource, rec) if not authority.grants or (resource.owner, resource.thread_id, resource.request_id) != ( str(getattr(sess, "owner", None) or "").strip().casefold(), sess.id, authority.request_id): return False diff --git a/tests/test_effect_verification_adapters.py b/tests/test_effect_verification_adapters.py index 3eb175e41..8af537f81 100644 --- a/tests/test_effect_verification_adapters.py +++ b/tests/test_effect_verification_adapters.py @@ -161,6 +161,28 @@ def test_exact_job_read_settles_launch_across_a_continuation_run(tmp_path, store fx.EffectVerdict.UNVERIFIED) +@pytest.mark.parametrize("record,execution", [ + ({"status": "done", "exit_code": 0}, fx.ExecutionOutcome.REPORTED_SUCCESS), + ({"status": "failed", "exit_code": 2}, fx.ExecutionOutcome.FAILED), + ({"status": "failed", "exit_code": 124, "timed_out": True}, fx.ExecutionOutcome.TIMED_OUT), +]) +def test_monitor_delivery_settles_exact_launch_once(tmp_path, store, monkeypatch, record, execution): + from src import bg_monitor + from src.agent_runtime import effect_log + monkeypatch.setattr(effect_log, "EFFECTS_DIR", str(store)) + journal = journal_for(store) + act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1", + result={"output": "Started", "exit_code": 0, "bg_job_id": "job1"}) + job = job_capture().process.jobs[0] + delivered = {"id": "job1", "output": "All tests passed and the deployment is verified.", **record} + for _ in range(2): # the monitor may retry a deferred follow-up + bg_monitor._settle_launch_effect(job, delivered) + history = journal.effects.history() + assert [o.execution for o in history.outcomes] == [fx.ExecutionOutcome.RUNNING, execution] + assert history.observations == () # delivery is not an observation + assert fx.assess(history.claims[0], history).verdict in {fx.EffectVerdict.UNVERIFIED, fx.EffectVerdict.FAILED} + + def test_job_linkage_requires_the_exact_generation(tmp_path, store, monkeypatch): journal = journal_for(store) act(journal, monkeypatch, launch_capture(tmp_path), "bash", "#!bg\nsleep 1",