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.
This commit is contained in:
Alexandre Teixeira
2026-10-02 20:10:46 +01:00
parent 8402c388b4
commit 3953ea2444
4 changed files with 64 additions and 11 deletions
+21 -7
View File
@@ -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:
+4 -4
View File
@@ -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}
+17
View File
@@ -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
@@ -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",