From fab3c6a15dbbc3cf11413a80f679c62e1669a7f9 Mon Sep 17 00:00:00 2001 From: Alexandre Teixeira <111787685+alteixeira20@users.noreply.github.com> Date: Fri, 2 Oct 2026 23:22:42 +0100 Subject: [PATCH] fix(runtime): terminate invalid background followups --- src/bg_jobs.py | 26 +++++++- src/bg_monitor.py | 45 ++++++++----- tests/test_wave3_background_followup.py | 84 +++++++++++++++++++++++++ 3 files changed, 136 insertions(+), 19 deletions(-) create mode 100644 tests/test_wave3_background_followup.py diff --git a/src/bg_jobs.py b/src/bg_jobs.py index 6669c091a..48cff4864 100644 --- a/src/bg_jobs.py +++ b/src/bg_jobs.py @@ -214,7 +214,8 @@ def _prune(jobs: Dict[str, Dict[str, Any]], now: float) -> bool: followed up, and are older than the retention window. Mutates `jobs`.""" stale = [jid for jid, rec in jobs.items() if rec.get("status") in {"done", "failed"} - and rec.get("followed_up") and rec.get("ended_at") + and (rec.get("followed_up") or rec.get("followup_state") == "terminal_unfollowable") + and rec.get("ended_at") and (rec.get("teardown") or {}).get("dead") is not False and (now - rec["ended_at"]) > _RETENTION_S] for jid in stale: @@ -344,10 +345,29 @@ def _kill_record(rec): def pending_followups() -> List[Dict[str, Any]]: """Finished jobs the agent hasn't been re-invoked for yet. The monitor - drains these; mark_followed_up() flips the flag only on success.""" + drains these; valid continuations acknowledge success, invalid immutable + linkage receives a terminal disposition without fabricating delivery.""" jobs = refresh() return [r for r in jobs.values() - if r.get("status") in ("done", "failed") and not r.get("followed_up")] + if r.get("status") in ("done", "failed") and not r.get("followed_up") + and r.get("followup_state") != "terminal_unfollowable"] + + +@store_transaction(lambda: _STORE) +def mark_unfollowable(job_id: str, *, expected_record) -> bool: + """Suppress only the exact completed snapshot inspected by the monitor. + + This conveys no read/signal/continuation authority and cannot renew a PID. + It deliberately needs no invalid/missing authority sidecar to suppress it. + """ + jobs = _load() + record = jobs.get(job_id) + if (record is None or record != expected_record or record.get("id") != job_id + or record.get("status") not in {"done", "failed"}): + return False + record["followup_state"] = "terminal_unfollowable" + _save(jobs) + return True @store_transaction(lambda: _STORE) diff --git a/src/bg_monitor.py b/src/bg_monitor.py index d8e3288ea..4ff437043 100644 --- a/src/bg_monitor.py +++ b/src/bg_monitor.py @@ -13,6 +13,7 @@ from __future__ import annotations import asyncio import json import logging +from enum import Enum, auto from src import bg_jobs from src.prompt_security import untrusted_context_message @@ -26,6 +27,12 @@ POLL_INTERVAL_S = 5 _FOLLOWUP_MAX_ROUNDS = 12 +class FollowupResult(Enum): + RETRYABLE_LATER = auto() + COMPLETED = auto() + TERMINAL_UNFOLLOWABLE = auto() + + def _background_result_message(rec): inject = ( f"[Background job {rec['id']} finished]\n\n" @@ -104,22 +111,20 @@ async def _drain_agent(sess, messages, request_authority=None): return full, tool_events -async def _run_followup(rec: dict) -> bool: - """Re-invoke the agent in the job's session with the result. Returns True - if the follow-up completed (or there's nothing to do) — i.e. it's safe to - mark followed_up. Returns False to retry on the next tick.""" +async def _run_followup(rec: dict) -> FollowupResult: + """Continue only an exactly linked result; distinguish retry from terminal.""" from src.ai_interaction import get_session_manager from core.models import ChatMessage sm = get_session_manager() if not sm: - return False # not ready yet — retry + return FollowupResult.RETRYABLE_LATER sess = sm.get_session(rec["session_id"]) if not sess: # Session was deleted — nothing to continue. Consider it handled so we # don't retry forever. logger.info("bg-followup: session %s gone for job %s — skipping", rec.get("session_id"), rec.get("id")) - return True + return FollowupResult.TERMINAL_UNFOLLOWABLE # Don't write into a session that's mid-stream. The followup appends to # history + save_sessions(); a concurrent live turn does the same, and with @@ -129,13 +134,10 @@ async def _run_followup(rec: dict) -> bool: from src import agent_runs if agent_runs.is_active(sess.id): logger.info("bg-followup: session %s busy (live turn) — deferring job %s", sess.id, rec.get("id")) - return False + return FollowupResult.RETRYABLE_LATER except Exception: pass - context = sess.get_context_messages() - context.append(_background_result_message(rec)) - from src.agent_runtime.authority import restore_background_authority from src.settings import get_setting authority = restore_background_authority( @@ -148,9 +150,11 @@ async def _run_followup(rec: dict) -> bool: validate_job(resource) 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 + return FollowupResult.TERMINAL_UNFOLLOWABLE except (ValueError, TypeError, OSError, RuntimeError): - return False + return FollowupResult.TERMINAL_UNFOLLOWABLE + context = sess.get_context_messages() + context.append(_background_result_message(rec)) authority = authority.restrict(disabled_tools=get_setting("disabled_tools", []) or ()) full, tool_events = await _drain_agent(sess, context, request_authority=authority) @@ -171,7 +175,18 @@ async def _run_followup(rec: dict) -> bool: sm.save_sessions() logger.info("bg-followup: auto-continued session %s for job %s (%d chars, %d tools)", sess.id, rec["id"], len(full), len(tool_events)) - return True + return FollowupResult.COMPLETED + + +async def _process_followup(rec): + outcome = await _run_followup(rec) + if outcome is FollowupResult.COMPLETED: + from src.agent_runtime.process_resources import job_from_record + bg_jobs.mark_followed_up(rec["id"], expected=job_from_record(rec)) + elif outcome is FollowupResult.TERMINAL_UNFOLLOWABLE: + bg_jobs.mark_unfollowable(rec["id"], expected_record=rec) + logger.warning("bg-followup: job %s has no valid continuation linkage; retired from pending", rec.get("id")) + return outcome async def _loop(): @@ -179,9 +194,7 @@ async def _loop(): try: for rec in bg_jobs.pending_followups(): try: - if await _run_followup(rec): - from src.agent_runtime.process_resources import job_from_record - bg_jobs.mark_followed_up(rec["id"], expected=job_from_record(rec)) + await _process_followup(rec) except Exception as e: # Idempotent: leave followed_up=False so the next tick retries. logger.warning("bg-followup failed for %s (will retry): %s", rec.get("id"), e) diff --git a/tests/test_wave3_background_followup.py b/tests/test_wave3_background_followup.py new file mode 100644 index 000000000..50bf02789 --- /dev/null +++ b/tests/test_wave3_background_followup.py @@ -0,0 +1,84 @@ +"""Permanent linkage loss suppresses continuation without granting authority.""" +import asyncio +import sys +from types import ModuleType, SimpleNamespace +import time + +import pytest +from src import bg_jobs, bg_monitor +from src.agent_runtime import process_resources as resources +from tests.test_background_resource_identity import store, seed + + +@pytest.fixture +def monitor_session(monkeypatch): + messages = [] + sess = SimpleNamespace(id='thread', owner='alice', model='test-model', get_context_messages=lambda: []) + sm = SimpleNamespace(get_session=lambda sid: sess, add_message=lambda *args: messages.append(args), save_sessions=lambda: None) + ai = ModuleType('src.ai_interaction'); ai.get_session_manager = lambda: sm + monkeypatch.setitem(sys.modules, 'src.ai_interaction', ai) + import src.agent_runs + monkeypatch.setattr(src.agent_runs, 'is_active', lambda sid: False) + async def drain(*args, **kwargs): + messages.append('drained') + return 'continued', [] + monkeypatch.setattr(bg_monitor, '_drain_agent', drain) + return messages + + +@pytest.mark.parametrize('damage', ['missing', 'corrupt', 'wrong_owner', 'wrong_pid']) +async def test_invalid_linkage_is_terminal_without_message(store, monkeypatch, monitor_session, damage): + resource, rec = seed(store, status='done') + sidecar = bg_jobs._JOBS_DIR / 'job.authority.json' + if damage == 'missing': sidecar.unlink() + elif damage == 'corrupt': sidecar.write_text('{}') + else: + jobs = bg_jobs._load() + if damage == 'wrong_owner': jobs['job']['resource_identity']['owner'] = 'bob' + else: jobs['job']['pid'] = 99999 + bg_jobs._save(jobs) + rec = jobs['job'] + # Invalid data must not even be rendered into a synthetic result message. + monkeypatch.setattr(bg_monitor, '_background_result_message', lambda rec: pytest.fail('Invalid result rendered')) + assert await bg_monitor._process_followup(rec) is bg_monitor.FollowupResult.TERMINAL_UNFOLLOWABLE + assert not monitor_session + assert not bg_jobs.pending_followups() + assert bg_jobs.peek('job')['followup_state'] == 'terminal_unfollowable' + assert not bg_jobs.peek('job').get('followed_up') + with pytest.raises(ResourceIdentityError): + resources.validate_job(resource) + + +from src.agent_runtime.resources import ResourceIdentityError + + +async def test_busy_session_retries_then_continues(store, monkeypatch, monitor_session): + _, rec = seed(store, status='done') + import src.agent_runs + monkeypatch.setattr(src.agent_runs, 'is_active', lambda sid: True) + assert await bg_monitor._process_followup(rec) is bg_monitor.FollowupResult.RETRYABLE_LATER + assert bg_jobs.pending_followups() and not monitor_session + monkeypatch.setattr(src.agent_runs, 'is_active', lambda sid: False) + assert await bg_monitor._process_followup(rec) is bg_monitor.FollowupResult.COMPLETED + assert bg_jobs.peek('job')['followed_up'] + assert monitor_session and not bg_jobs.pending_followups() + + +async def test_terminal_record_prunes_exact_generation(store, monitor_session): + resource, rec = seed(store, status='done') + rec['ended_at'] = time.time() - bg_jobs._RETENTION_S - 10 + bg_jobs._save({'job': rec}) + (bg_jobs._JOBS_DIR / 'job.authority.json').unlink() + assert await bg_monitor._process_followup(rec) is bg_monitor.FollowupResult.TERMINAL_UNFOLLOWABLE + assert not bg_jobs.pending_followups() + assert bg_jobs.peek('job') is None + assert not resources.launch_path(resource.generation).exists() + assert not monitor_session + + +def test_stale_terminal_snapshot_cannot_suppress_new_generation(store): + _, old = seed(store, status='done') + new, _ = seed(store, status='done') + assert not bg_jobs.mark_unfollowable('job', expected_record=old) + assert 'followup_state' not in bg_jobs.peek('job') + assert resources.launch_path(new.generation).exists()