diff --git a/agent/agent_init.py b/agent/agent_init.py index 6f0edf4fb4..ffefcee5eb 100644 --- a/agent/agent_init.py +++ b/agent/agent_init.py @@ -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, diff --git a/gateway/run.py b/gateway/run.py index e84b5feee8..f105d27a25 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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 diff --git a/run_agent.py b/run_agent.py index b086400b6c..3d295caf27 100644 --- a/run_agent.py +++ b/run_agent.py @@ -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. diff --git a/tests/gateway/test_busy_session_ack.py b/tests/gateway/test_busy_session_ack.py index c58031fdb5..a77c527d2e 100644 --- a/tests/gateway/test_busy_session_ack.py +++ b/tests/gateway/test_busy_session_ack.py @@ -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 diff --git a/tests/tools/test_zombie_process_cleanup.py b/tests/tools/test_zombie_process_cleanup.py index e31e042fb2..a8b745f541 100644 --- a/tests/tools/test_zombie_process_cleanup.py +++ b/tests/tools/test_zombie_process_cleanup.py @@ -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