From b7e971032b41b94588e9610dbe98c2828499f238 Mon Sep 17 00:00:00 2001 From: Ishaan Jaffer Date: Wed, 6 May 2026 15:51:45 -0700 Subject: [PATCH] fix(cleanup): route dead-daemon sweep through _terminate_session_internal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _sweep_dead_daemons previously ran update_many directly on the session row, skipping provider.terminate. With NoopVMProvider this was harmless, but once Epic B swaps in a real VM provider it would orphan EC2 instances every time a daemon stopped heartbeating. Mirror the pattern from _sweep_expired_sessions: call _terminate_session_internal per-row so the provider is notified, then explicitly downgrade status from 'terminated' to 'error' (both are terminal — no further state transitions). Also drops the per-run update loop since _terminate_session_internal already cancels active runs and emits run_cancelled events. Greptile P1 (review #PRR_kwDOKALCgc78u9En). --- .../proxy/agent_session_endpoints/cleanup.py | 66 +++++++------------ 1 file changed, 25 insertions(+), 41 deletions(-) diff --git a/litellm/proxy/agent_session_endpoints/cleanup.py b/litellm/proxy/agent_session_endpoints/cleanup.py index 0211a281d2e..851c83ab6d2 100644 --- a/litellm/proxy/agent_session_endpoints/cleanup.py +++ b/litellm/proxy/agent_session_endpoints/cleanup.py @@ -22,7 +22,6 @@ from litellm.proxy.agent_session_endpoints.constants import ( CLEANUP_SWEEPER_INTERVAL_SECONDS, DAEMON_HEARTBEAT_DEAD_AFTER_SECONDS, EVENT_TYPE_RUN_ERROR, - RUN_ACTIVE_STATUSES, RUN_IDLE_TIMEOUT_SECONDS, RUN_STATUS_ERROR, SESSION_STATUS_ERROR, @@ -67,7 +66,15 @@ async def _sweep_dead_daemons(prisma_client) -> int: Skips ``provisioning`` sessions (they haven't registered yet — their own timeout is governed by `expires_at`). + + Routes through ``_terminate_session_internal`` so the VM provider + is notified — Greptile P1: bypassing this would orphan EC2 instances + once a real provider replaces ``NoopVMProvider``. """ + from litellm.proxy.agent_session_endpoints.session_endpoints import ( + _terminate_session_internal, + ) + threshold = _now() - timedelta(seconds=DAEMON_HEARTBEAT_DEAD_AFTER_SECONDS) rows = await prisma_client.db.litellm_agentsession.find_many( where={ @@ -82,57 +89,34 @@ async def _sweep_dead_daemons(prisma_client) -> int: if not rows: return 0 + # Terminate via the shared helper. It cancels active runs (with + # `run_cancelled` events), flips the session status, AND fires + # ``provider.terminate`` — exactly what we want here. + for row in rows: + try: + await _terminate_session_internal(row.id, reason="daemon_dead") + except Exception as exc: + verbose_proxy_logger.exception( + "sweeper: failed to terminate dead-daemon session=%s: %s", + row.id, + exc, + ) + + # Daemon-dead sessions land in ``error`` (not ``terminated``); the + # shared helper sets ``terminated`` so we explicitly downgrade to + # ``error`` here. (Both are terminal — no further state transitions.) now = _now() ids = [r.id for r in rows] await prisma_client.db.litellm_agentsession.update_many( where={"id": {"in": ids}}, data={ "status": SESSION_STATUS_ERROR, - "terminated_at": now, "updated_at": now, }, ) - # Also flip any active runs in those sessions to error. - runs = await prisma_client.db.litellm_agentrun.find_many( - where={ - "session_id": {"in": ids}, - "status": {"in": list(RUN_ACTIVE_STATUSES)}, - }, - take=500, - ) - for run in runs: - await prisma_client.db.litellm_agentrun.update( - where={"id": run.id}, - data={ - "status": RUN_STATUS_ERROR, - "terminated_at": now, - "updated_at": now, - }, - ) - last_evt = await prisma_client.db.litellm_agentrunevent.find_first( - where={"run_id": run.id}, order={"seq": "desc"} - ) - next_seq = (last_evt.seq + 1) if last_evt else 1 - try: - await prisma_client.db.litellm_agentrunevent.create( - data={ - "run_id": run.id, - "seq": next_seq, - "event_type": EVENT_TYPE_RUN_ERROR, - "payload": {"reason": "daemon_dead"}, - } - ) - except Exception as exc: - verbose_proxy_logger.warning( - "sweeper: skipped run_error event run=%s seq=%s: %s", - run.id, - next_seq, - exc, - ) - verbose_proxy_logger.info( - "sweeper: marked %d sessions error (dead daemon)", len(ids) + "sweeper: terminated %d sessions (dead daemon, via provider)", len(ids) ) return len(ids)