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

1"""Supervisor-side reaper for orphaned Prisma query-engine processes. 

2 

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. 

11 

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. 

19 

20Linux-only by construction (``/proc`` scan, ``prctl``); a no-op elsewhere. 

21""" 

22 

23import ctypes 

24import os 

25import signal 

26import sys 

27import threading 

28import time 

29from typing import Final 

30 

31from litellm._logging import verbose_proxy_logger 

32 

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 

37 

38 

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 

54 

55 

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 

75 

76 

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 ) 

94 

95 

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 

104 

105 

106def _send_signal(pid: int, signum: int) -> None: 

107 try: 

108 os.kill(pid, signum) 

109 except (ProcessLookupError, PermissionError, OSError): 

110 pass 

111 

112 

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 

123 

124 

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) 

129 

130 

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 ) 

161 

162 

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 

169 

170 

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) 

178 

179 

180REAPER_THREAD_NAME: Final = "litellm-orphan-query-engine-reaper" 

181 

182 

183def start_query_engine_reaper() -> threading.Thread | None: 

184 """Start the reaper daemon thread in the supervisor process. 

185 

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