mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-06 06:52:20 +02:00
fix(runtime): terminate invalid background followups
This commit is contained in:
+23
-3
@@ -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)
|
||||
|
||||
+29
-16
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user