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
28 changes: 18 additions & 10 deletions agent/review_idle_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@

from __future__ import annotations

import contextvars
import json
import logging
import threading
Expand Down Expand Up @@ -64,6 +65,7 @@ class _PendingReview:
session_key: str
kwargs: Dict[str, Any]
enqueued_at: float
context: contextvars.Context


class ReviewIdleQueue:
Expand Down Expand Up @@ -98,7 +100,8 @@ def enqueue(self, agent: Any, session_key: str, kwargs: Dict[str, Any]) -> None:
with self._lock:
existing = self._pending.get(session_key)
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._pending[session_key] = _PendingReview(
agent, session_key, kwargs, enqueued_at, contextvars.copy_context())
self._ensure_thread()
self._wake.set()
logger.info("Background review deferred (session=%s, queued=%d)", session_key[-12:], len(self._pending))
Expand Down Expand Up @@ -149,20 +152,25 @@ def _run(self) -> None:
try:
item = self._pop_dispatchable()
if item is not None:
if not self._still_enabled(item):
logger.info(
"Deferred background review dropped: reviews were disabled while it was queued (session=%s)",
item.session_key[-12:])
continue
logger.info(
"Dispatching deferred background review (session=%s, waited=%.0fs, queued=%d)",
item.session_key[-12:], self._now() - item.enqueued_at, self.pending_count())
item.agent._spawn_background_review_now(**item.kwargs)
# The shared dispatcher has no caller profile. Enter the item's context
# before reading config AND spawning the context-propagating review worker.
item.context.run(self._dispatch, item)
except Exception: # noqa: BLE001 — dispatcher must survive anything
logger.warning("Deferred review dispatch failed", exc_info=True)
if item is None:
time.sleep(_POLL_INTERVAL_S)

def _dispatch(self, item: _PendingReview) -> None:
if not self._still_enabled(item):
logger.info(
"Deferred background review dropped: reviews were disabled while it was queued (session=%s)",
item.session_key[-12:])
return
logger.info(
"Dispatching deferred background review (session=%s, waited=%.0fs, queued=%d)",
item.session_key[-12:], self._now() - item.enqueued_at, self.pending_count())
item.agent._spawn_background_review_now(**item.kwargs)

@staticmethod
def _still_enabled(item: _PendingReview) -> bool:
"""Re-check the enabled gate at DISPATCH time (disabling reviews while queued must stick). Fail-open."""
Expand Down
77 changes: 77 additions & 0 deletions tests/agent/test_review_idle_profile_scope.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
"""Deferred review dispatch must use the profile that admitted each item."""

import threading

import pytest

from agent.review_idle_queue import ReviewIdleQueue
from hermes_constants import reset_hermes_home_override, set_hermes_home_override
from tools.memory_tool_store import MemoryStore


@pytest.mark.parametrize("ambient_enabled", [True, False])
def test_deferred_review_keeps_profile_config_and_memory(tmp_path, monkeypatch, ambient_enabled):
ambient = tmp_path / "ambient"
profiles = [tmp_path / "first", tmp_path / "second"]
enabled = not ambient_enabled
for home, gate in [(ambient, ambient_enabled), *[(p, enabled) for p in profiles]]:
home.mkdir()
(home / "config.yaml").write_text(
f"auxiliary:\n background_review:\n enabled: {str(gate).lower()}\n",
encoding="utf-8",
)
monkeypatch.setenv("HERMES_HOME", str(ambient))
queue = ReviewIdleQueue()
monkeypatch.setattr(queue, "_ensure_thread", lambda: None)
clock = [0.0]
queue._now = lambda: clock[0]
writes = []

class Parent:
def _spawn_background_review_now(self, **kwargs):
# Real profile-aware memory I/O at the spawn boundary; no model call.
writes.append(MemoryStore().add("memory", kwargs["fact"]))

parent = Parent()
for index, home in enumerate(profiles):
token = set_hermes_home_override(home)
try:
queue.enqueue(parent, str(index), {"fact": f"Project {index} uses Python", "task_cfg": {}})
finally:
reset_hermes_home_override(token)
clock[0] = 4000.0 # age-out avoids contacting a model server

class StopWorker(BaseException):
pass

class DrainWake(threading.Event):
def wait(self, timeout=None):
if queue.pending_count() == 0:
raise StopWorker
return True

queue._wake = DrainWake()
failures = []

def drain():
try:
queue._run()
except StopWorker:
pass
except BaseException as exc:
failures.append(exc)

worker = threading.Thread(target=drain, daemon=True)
worker.start()
worker.join(timeout=10)
assert not worker.is_alive()
assert not failures
assert queue.pending_count() == 0
assert len(writes) == (2 if enabled else 0)
assert all(result["success"] for result in writes)
assert not (ambient / "memories" / "MEMORY.md").exists()
for index, home in enumerate(profiles):
path = home / "memories" / "MEMORY.md"
assert path.exists() == enabled
if enabled:
assert path.read_text(encoding="utf-8") == f"Project {index} uses Python"