fix(session): finalize owned SQLite session rows on AIAgent.close()
Funnel session finalization through AIAgent.close() — the single terminal path every agent (CLI, gateway, subagent, cron) funnels through — so finished agents stop leaving rows with ended_at IS NULL. The biggest leak source was delegate_task subagent + background-review forks whose close() never ended their row. end_session() is first-reason-wins and no-ops on an already-ended row, so a 'compression'/'cron_complete'/'cli_close' reason set by an earlier terminal path is never clobbered. /resume already calls reopen_session(), so finalizing-on-close does not break resumability. Temporary helper agents that rotate/share the session forward (manual compression, gateway session-hygiene) opt out via _end_session_on_close=False. Also stop the long-running gateway heartbeat once the executor is done or the session slot is rebound to a different agent, preventing a stale 'running: delegate_task' bubble from outliving its run. Closes #12029.
This commit is contained in:
@@ -1100,6 +1100,12 @@ def init_agent(
|
||||
agent._parent_session_id = parent_session_id
|
||||
agent._last_flushed_db_idx = 0 # tracks DB-write cursor to prevent duplicate writes
|
||||
agent._session_db_created = False # DB row deferred to run_conversation()
|
||||
# Most agents own their session row and should finalize it on close().
|
||||
# Some temporary helper agents (manual compression / session-hygiene /
|
||||
# background-review forks) rotate or share the session forward to a
|
||||
# continuation row that must remain open after the helper is torn down;
|
||||
# those callers explicitly set this flag to False.
|
||||
agent._end_session_on_close = True
|
||||
agent._session_init_model_config = {
|
||||
"max_iterations": agent.max_iterations,
|
||||
"reasoning_config": reasoning_config,
|
||||
|
||||
@@ -4686,6 +4686,27 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
pass
|
||||
self._cleanup_agent_resources(agent)
|
||||
|
||||
def _should_emit_long_running_notification(
|
||||
self,
|
||||
session_key: Optional[str],
|
||||
agent: Any,
|
||||
executor_task: Optional[Any],
|
||||
) -> bool:
|
||||
"""Only emit the heartbeat while this task still owns the live run.
|
||||
|
||||
Guards against a stale ``running: delegate_task`` heartbeat outliving the
|
||||
run that started it: stop once the executor finishes, the agent is gone,
|
||||
or the session key has been rebound to a different live agent (e.g. the
|
||||
user sent ``/new`` and a fresh agent took the slot mid-run, #12029).
|
||||
"""
|
||||
if agent is None:
|
||||
return False
|
||||
if executor_task is not None and executor_task.done():
|
||||
return False
|
||||
if session_key and self._running_agents.get(session_key) is not agent:
|
||||
return False
|
||||
return True
|
||||
|
||||
def _cleanup_agent_resources(self, agent: Any) -> None:
|
||||
"""Best-effort cleanup for temporary or cached agent instances."""
|
||||
if agent is None:
|
||||
@@ -9194,6 +9215,13 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
session_id=session_entry.session_id,
|
||||
)
|
||||
try:
|
||||
# The hygiene agent rotates the session
|
||||
# forward to a continuation id that becomes
|
||||
# the gateway session's live row. It must
|
||||
# never finalize on close() (today it has no
|
||||
# session_db so close() no-ops, but this
|
||||
# guards a future where one is wired in).
|
||||
_hyg_agent._end_session_on_close = False
|
||||
_hyg_agent._print_fn = lambda *a, **kw: None
|
||||
|
||||
loop = asyncio.get_running_loop()
|
||||
@@ -16274,6 +16302,20 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
_heartbeat_msg_id: Optional[str] = None
|
||||
while True:
|
||||
await asyncio.sleep(_NOTIFY_INTERVAL)
|
||||
# Stop heartbeating once this run no longer owns the session
|
||||
# slot or the executor has finished — otherwise a stale
|
||||
# "running: delegate_task" bubble can outlive the run that
|
||||
# spawned it (#12029). _executor_task is a closure var bound
|
||||
# just after this task is scheduled; tolerate the brief window
|
||||
# before then (the first wake is _NOTIFY_INTERVAL away anyway).
|
||||
try:
|
||||
_exec_ref = _executor_task
|
||||
except NameError:
|
||||
_exec_ref = None
|
||||
if not self._should_emit_long_running_notification(
|
||||
session_key, agent_holder[0], _exec_ref
|
||||
):
|
||||
break
|
||||
_elapsed_mins = int((time.time() - _notify_start) // 60)
|
||||
# Include agent activity context if available. Default
|
||||
# heartbeat is terse: elapsed + current tool. Verbose
|
||||
|
||||
@@ -3250,6 +3250,22 @@ class AIAgent:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# 7. Finalize the owned SQLite session row unless this agent is only a
|
||||
# temporary helper that deliberately handed session ownership forward
|
||||
# (manual compression helpers that rotate to a continuation session_id,
|
||||
# or background-review forks that share the live parent's session_id and
|
||||
# must leave it open). end_session() is first-reason-wins and no-ops on
|
||||
# an already-ended row, so this never clobbers a 'compression' /
|
||||
# 'cron_complete' / 'cli_close' reason set by an earlier terminal path.
|
||||
try:
|
||||
if getattr(self, "_end_session_on_close", True):
|
||||
session_db = getattr(self, "_session_db", None)
|
||||
session_id = getattr(self, "session_id", None)
|
||||
if session_db and session_id:
|
||||
session_db.end_session(session_id, "agent_close")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _hydrate_todo_store(self, history: List[Dict[str, Any]]) -> None:
|
||||
"""
|
||||
Recover todo state from conversation history.
|
||||
|
||||
@@ -715,3 +715,62 @@ class TestBusySessionOnboardingHint:
|
||||
assert "/busy interrupt" in content
|
||||
# Must NOT tell the user to /busy queue when they're already on queue.
|
||||
assert "/busy queue" not in content
|
||||
|
||||
|
||||
class TestLongRunningNotificationOwnership:
|
||||
"""The long-running heartbeat must stop once its run no longer owns the
|
||||
session slot or the executor finished — otherwise a stale
|
||||
'running: delegate_task' bubble outlives the run that spawned it (#12029).
|
||||
"""
|
||||
|
||||
def test_notification_stops_after_session_ownership_moves(self):
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
runner._running_agents = {}
|
||||
|
||||
original_agent = MagicMock()
|
||||
replacement_agent = MagicMock()
|
||||
runner._running_agents["sess"] = replacement_agent
|
||||
|
||||
assert runner._should_emit_long_running_notification(
|
||||
"sess", original_agent, executor_task=None
|
||||
) is False
|
||||
|
||||
def test_notification_stops_after_executor_finishes(self):
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
agent = MagicMock()
|
||||
runner._running_agents = {"sess": agent}
|
||||
|
||||
done_task = MagicMock()
|
||||
done_task.done.return_value = True
|
||||
|
||||
assert runner._should_emit_long_running_notification(
|
||||
"sess", agent, executor_task=done_task
|
||||
) is False
|
||||
|
||||
def test_notification_stops_when_agent_is_gone(self):
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
runner._running_agents = {}
|
||||
|
||||
assert runner._should_emit_long_running_notification(
|
||||
"sess", None, executor_task=None
|
||||
) is False
|
||||
|
||||
def test_notification_continues_for_live_active_run(self):
|
||||
from gateway.run import GatewayRunner
|
||||
|
||||
runner = object.__new__(GatewayRunner)
|
||||
agent = MagicMock()
|
||||
runner._running_agents = {"sess": agent}
|
||||
|
||||
live_task = MagicMock()
|
||||
live_task.done.return_value = False
|
||||
|
||||
assert runner._should_emit_long_running_notification(
|
||||
"sess", agent, executor_task=live_task
|
||||
) is True
|
||||
|
||||
@@ -155,6 +155,59 @@ class TestAgentCloseMethod:
|
||||
child_2.close.assert_called_once()
|
||||
assert agent._active_children == []
|
||||
|
||||
def test_close_ends_owned_session_row(self):
|
||||
"""close() finalizes the agent's owned SQLite session row."""
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
with patch("run_agent.AIAgent.__init__", return_value=None):
|
||||
from run_agent import AIAgent
|
||||
agent = AIAgent.__new__(AIAgent)
|
||||
agent.session_id = "test-close-session-row"
|
||||
agent._active_children = []
|
||||
agent._active_children_lock = threading.Lock()
|
||||
agent.client = None
|
||||
agent._end_session_on_close = True
|
||||
agent._session_db = MagicMock()
|
||||
|
||||
agent.close()
|
||||
|
||||
agent._session_db.end_session.assert_called_once_with(
|
||||
"test-close-session-row", "agent_close"
|
||||
)
|
||||
|
||||
def test_close_skips_session_end_for_forwarded_continuation_agents(self):
|
||||
"""Helper agents that handed session ownership forward opt out."""
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
with patch("run_agent.AIAgent.__init__", return_value=None):
|
||||
from run_agent import AIAgent
|
||||
agent = AIAgent.__new__(AIAgent)
|
||||
agent.session_id = "test-close-forwarded-session"
|
||||
agent._active_children = []
|
||||
agent._active_children_lock = threading.Lock()
|
||||
agent.client = None
|
||||
agent._end_session_on_close = False
|
||||
agent._session_db = MagicMock()
|
||||
|
||||
agent.close()
|
||||
|
||||
agent._session_db.end_session.assert_not_called()
|
||||
|
||||
def test_close_session_end_noops_without_session_db(self):
|
||||
"""close() is a no-op for session finalization when no DB is wired in."""
|
||||
from unittest.mock import patch
|
||||
|
||||
with patch("run_agent.AIAgent.__init__", return_value=None):
|
||||
from run_agent import AIAgent
|
||||
agent = AIAgent.__new__(AIAgent)
|
||||
agent.session_id = "test-close-no-db"
|
||||
agent._active_children = []
|
||||
agent._active_children_lock = threading.Lock()
|
||||
agent.client = None
|
||||
# No _session_db / _end_session_on_close attributes at all —
|
||||
# getattr defaults must keep close() from raising.
|
||||
agent.close() # must not raise
|
||||
|
||||
def test_close_survives_partial_failures(self):
|
||||
"""close() continues cleanup even if one step fails."""
|
||||
from unittest.mock import patch
|
||||
|
||||
Reference in New Issue
Block a user