Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/db/query_engine_reaper.py: 18%
107 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1"""Supervisor-side reaper for orphaned Prisma query-engine processes.
3Each proxy worker owns a Prisma query-engine subprocess whose only cleanup
4hook is an in-process ``atexit`` handler. When a multi-worker supervisor
5(uvicorn's multiprocess manager, the gunicorn arbiter) force-kills a hung or
6crashed worker, that handler never runs: the engine reparents to the nearest
7subreaper (PID 1 in a container, which is the supervisor itself under the
8standard docker entrypoint) and keeps its database connection pool
9established forever, while the replacement worker opens a fresh pool. Over
10repeated worker deaths the active database connections grow without bound.
12The reaper runs only in the supervisor process, where a query-engine process
13can never be a legitimate direct child: workers own their engines, and the
14supervisor never starts one. Any direct child whose command name begins with
15``query-engine`` is therefore an adopted orphan and is terminated
16(SIGTERM, bounded grace, SIGKILL) and reaped. On Linux the supervisor also
17marks itself a child subreaper so orphans reparent to it even when it is not
18PID 1.
20Linux-only by construction (``/proc`` scan, ``prctl``); a no-op elsewhere.
21"""
23import ctypes
24import os
25import signal
26import sys
27import threading
28import time
29from typing import Final
31from litellm._logging import verbose_proxy_logger
33QUERY_ENGINE_COMM_PREFIX: Final = "query-engine"
34REAPER_SCAN_INTERVAL_SECONDS: Final = 5.0
35SIGTERM_GRACE_SECONDS: Final = 10.0
36PR_SET_CHILD_SUBREAPER: Final = 36
39def set_child_subreaper() -> bool:
40 """Mark this process as a child subreaper so orphaned descendants
41 reparent to it instead of PID 1. Best-effort: when it fails (or on
42 non-Linux) the reaper still covers the containerized case where the
43 supervisor already is PID 1."""
44 if not sys.platform.startswith("linux"):
45 return False
46 try:
47 libc: Final = ctypes.CDLL(None, use_errno=True)
48 result: int = libc.prctl( # pyright: ignore[reportAny] # ctypes types foreign calls as Any; default restype is c_int
49 PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0
50 )
51 return result == 0
52 except (OSError, AttributeError):
53 return False
56def _read_comm_and_ppid(pid: int, proc_root: str) -> tuple[str, int] | None:
57 try:
58 with open(f"{proc_root}/{pid}/stat", encoding="ascii", errors="replace") as stat_file:
59 data: Final = stat_file.read()
60 except (FileNotFoundError, ProcessLookupError, PermissionError, OSError):
61 return None
62 lparen: Final = data.find("(")
63 rparen: Final = data.rfind(")")
64 if lparen == -1 or rparen == -1 or rparen < lparen:
65 return None
66 comm: Final = data[lparen + 1 : rparen]
67 fields: Final = data[rparen + 2 :].split()
68 if len(fields) < 2:
69 return None
70 try:
71 ppid: Final = int(fields[1])
72 except ValueError:
73 return None
74 return comm, ppid
77def list_orphaned_engine_pids(parent_pid: int, proc_root: str = "/proc") -> tuple[int, ...]:
78 """PIDs of direct children of ``parent_pid`` whose command name marks
79 them as Prisma query engines. In the supervisor these are always
80 adopted orphans: live engines are children of workers, not of the
81 supervisor."""
82 try:
83 entries: Final = os.listdir(proc_root)
84 except (FileNotFoundError, OSError):
85 return ()
86 candidate_pids: Final = (int(entry) for entry in entries if entry.isdigit())
87 return tuple(
88 pid
89 for pid in candidate_pids
90 if (info := _read_comm_and_ppid(pid, proc_root)) is not None
91 and info[1] == parent_pid
92 and info[0].startswith(QUERY_ENGINE_COMM_PREFIX)
93 )
96def _try_reap(pid: int) -> bool:
97 try:
98 reaped_pid, _ = os.waitpid(pid, os.WNOHANG)
99 except ChildProcessError:
100 return True
101 except OSError:
102 return True
103 return reaped_pid == pid
106def _send_signal(pid: int, signum: int) -> None:
107 try:
108 os.kill(pid, signum)
109 except (ProcessLookupError, PermissionError, OSError):
110 pass
113def _await_reaped(pids: tuple[int, ...], timeout_seconds: float) -> tuple[int, ...]:
114 """Poll until every PID is reaped or the shared deadline passes.
115 Returns the PIDs still alive at the deadline."""
116 deadline: Final = time.monotonic() + timeout_seconds
117 remaining = pids
118 while remaining and time.monotonic() < deadline:
119 remaining = tuple(pid for pid in remaining if not _try_reap(pid))
120 if remaining:
121 time.sleep(0.2)
122 return remaining
125def terminate_and_reap(pid: int, grace_seconds: float = SIGTERM_GRACE_SECONDS) -> None:
126 """SIGTERM the orphaned engine, escalate to SIGKILL after the grace
127 period, and reap it so it does not linger as a zombie."""
128 terminate_and_reap_all((pid,), grace_seconds=grace_seconds)
131def terminate_and_reap_all(
132 pids: tuple[int, ...],
133 grace_seconds: float = SIGTERM_GRACE_SECONDS,
134) -> None:
135 """Terminate a batch of orphaned engines concurrently: SIGTERM all of
136 them, share one grace period, SIGKILL the stragglers, and reap. The
137 shared deadline keeps cleanup time bounded when several workers die
138 at once instead of paying the grace period once per orphan."""
139 for pid in pids:
140 verbose_proxy_logger.warning(
141 "Reaping orphaned prisma query-engine PID %s (its worker process exited without cleanup).",
142 pid,
143 )
144 _send_signal(pid, signal.SIGTERM)
145 survivors: Final = _await_reaped(pids, grace_seconds)
146 if not survivors:
147 return
148 for pid in survivors:
149 verbose_proxy_logger.warning(
150 "Orphaned prisma query-engine PID %s did not exit within %.1fs of SIGTERM; sending SIGKILL.",
151 pid,
152 grace_seconds,
153 )
154 _send_signal(pid, signal.SIGKILL)
155 unkillable: Final = _await_reaped(survivors, 5.0)
156 for pid in unkillable:
157 verbose_proxy_logger.error(
158 "Orphaned prisma query-engine PID %s survived SIGKILL; will retry on the next scan.",
159 pid,
160 )
163def reap_orphaned_engines(parent_pid: int, proc_root: str = "/proc") -> tuple[int, ...]:
164 """One scan-and-reap pass. Returns the PIDs it acted on."""
165 orphaned_pids: Final = list_orphaned_engine_pids(parent_pid, proc_root=proc_root)
166 if orphaned_pids:
167 terminate_and_reap_all(orphaned_pids)
168 return orphaned_pids
171def _reaper_loop(parent_pid: int) -> None:
172 while True:
173 try:
174 reap_orphaned_engines(parent_pid)
175 except Exception as scan_error: # noqa: BLE001 # reaper thread must survive any scan failure
176 verbose_proxy_logger.debug("Orphaned query-engine scan failed: %s", scan_error)
177 time.sleep(REAPER_SCAN_INTERVAL_SECONDS)
180REAPER_THREAD_NAME: Final = "litellm-orphan-query-engine-reaper"
183def start_query_engine_reaper() -> threading.Thread | None:
184 """Start the reaper daemon thread in the supervisor process.
186 Must only be called from a process that never hosts the proxy app
187 itself (uvicorn with ``workers > 1``, the gunicorn arbiter): with a
188 single in-process uvicorn worker the query engine is a legitimate
189 direct child and must not be touched. Idempotent: a reaper already
190 running in this process is returned instead of starting a second one.
191 """
192 if not sys.platform.startswith("linux"):
193 return None
194 existing: Final = next(
195 (thread for thread in threading.enumerate() if thread.name == REAPER_THREAD_NAME),
196 None,
197 )
198 if existing is not None:
199 return existing
200 set_child_subreaper()
201 reaper_thread: Final = threading.Thread(
202 target=_reaper_loop,
203 args=(os.getpid(),),
204 daemon=True,
205 name=REAPER_THREAD_NAME,
206 )
207 reaper_thread.start()
208 verbose_proxy_logger.info(
209 "Started orphaned prisma query-engine reaper in supervisor process %s.",
210 os.getpid(),
211 )
212 return reaper_thread