Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions agent/review_idle_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -95,8 +95,15 @@ def note_turn_finished(self) -> None:
def enqueue(self, agent: Any, session_key: str, kwargs: Dict[str, Any]) -> None:
"""Add (or replace — newest snapshot wins) a session's pending review, keeping the ORIGINAL
enqueue time on coalesce so a busy session cannot push its age-out forever."""
kwargs = dict(kwargs)
kwargs.setdefault("_review_snapshot_created_at", self._now())
with self._lock:
existing = self._pending.get(session_key)
# A cancelled worker can unwind after the next turn queued a newer snapshot.
# Its retry keeps the original creation time; arrival order is not freshness.
if (existing is not None and kwargs["_review_snapshot_created_at"]
< existing.kwargs["_review_snapshot_created_at"]):
return
enqueued_at = existing.enqueued_at if existing is not None else self._now()
self._pending[session_key] = _PendingReview(agent, session_key, kwargs, enqueued_at)
self._ensure_thread()
Expand Down
7 changes: 5 additions & 2 deletions run_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -772,7 +772,8 @@ def _spawn_background_review(self, messages_snapshot: List[Dict], review_memory:
def _spawn_background_review_now(self, messages_snapshot: List[Dict], review_memory: bool = False,
review_skills: bool = False, focus: Optional[str] = None,
task_cfg: Optional[Dict[str, Any]] = None, _requeue_attempts: int = 0,
explicit: bool = False) -> None:
explicit: bool = False,
_review_snapshot_created_at: Optional[float] = None) -> None:
"""Spawn the background memory/skill review thread.

``threading.Thread`` is constructed here so tests patching ``run_agent.threading.Thread`` keep working.
Expand All @@ -786,6 +787,8 @@ def _spawn_background_review_now(self, messages_snapshot: List[Dict], review_mem
)
from tools.thread_context import propagate_context_to_thread

if _review_snapshot_created_at is None:
_review_snapshot_created_at = time.monotonic()
review_run = prepare_background_review_run(self)
if review_run is None:
return
Expand All @@ -800,7 +803,7 @@ def _target_with_requeue() -> None:
self._maybe_requeue_preempted_review(review_run, dict(
messages_snapshot=messages_snapshot, review_memory=review_memory, review_skills=review_skills,
focus=focus, task_cfg=task_cfg, _requeue_attempts=_requeue_attempts + 1,
explicit=explicit))
explicit=explicit, _review_snapshot_created_at=_review_snapshot_created_at))

# Carry the active profile into the review thread so MEMORY.md / skill review writes land in the
# right profile.
Expand Down
67 changes: 67 additions & 0 deletions tests/agent/test_review_requeue_order.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
"""A preempted review retains its original position relative to newer snapshots."""

import threading
from types import SimpleNamespace

import pytest

from agent import review_idle_queue
from run_agent import AIAgent


@pytest.mark.parametrize("retry_first", [False, True])
def test_preempted_retry_cannot_replace_a_newer_snapshot(monkeypatch, retry_first):
queue = review_idle_queue.ReviewIdleQueue()
monkeypatch.setattr(queue, "_ensure_thread", lambda: None)
monkeypatch.setattr(review_idle_queue, "QUEUE", queue)
monkeypatch.setattr(review_idle_queue, "review_targets_managed_local", lambda *_: True)
now = [10.0]
queue._now = lambda: now[0]
parent = SimpleNamespace(session_id="session", _REVIEW_REQUEUE_MAX_ATTEMPTS=3)
old = [{"role": "user", "content": "Original task"}]
new = old + [{"role": "user", "content": "The corrected requirement"}]
queue.enqueue(parent, "session", {"messages_snapshot": old, "task_cfg": {}})
now[0] += 2000
running = queue._pop_dispatchable()
assert running is not None
retry = dict(running.kwargs, _requeue_attempts=1)
cancelled = threading.Event()
cancelled.set()

def requeue():
AIAgent._maybe_requeue_preempted_review(
parent, SimpleNamespace(cancel_requested=cancelled), retry)

if retry_first:
requeue()
now[0] += 1
queue.enqueue(parent, "session", {"messages_snapshot": new, "task_cfg": {}})
admitted_at = queue._pending["session"].enqueued_at
if not retry_first:
requeue()
pending = queue._pending["session"]
assert pending.kwargs["messages_snapshot"] == new
assert pending.kwargs.get("_requeue_attempts", 0) == 0
assert pending.enqueued_at == admitted_at


def test_spawn_preserves_snapshot_age_through_the_worker(monkeypatch):
from agent import background_review

done = threading.Event()
captured = []
run = SimpleNamespace()
monkeypatch.setattr(background_review, "prepare_background_review_run", lambda _: run)
monkeypatch.setattr(background_review, "spawn_background_review_thread", lambda *a, **k: (lambda: None, ""))

def capture(finished, kwargs):
captured.append((finished, kwargs))
done.set()

parent = SimpleNamespace(_maybe_requeue_preempted_review=capture)
AIAgent._spawn_background_review_now(
parent, [{"role": "user", "content": "Task"}], _review_snapshot_created_at=12.5)
assert done.wait(10)
assert captured[0][0] is run
assert captured[0][1]["_review_snapshot_created_at"] == 12.5
assert captured[0][1]["_requeue_attempts"] == 1