1###############################################################################
2# Server process to keep track of unlinked resources (like shared memory
3# segments, semaphores etc.) and clean them.
4#
5# On Unix we run a server process which keeps track of unlinked
6# resources. The server ignores SIGINT and SIGTERM and reads from a
7# pipe. Every other process of the program has a copy of the writable
8# end of the pipe, so we get EOF when all other processes have exited.
9# Then the server process unlinks any remaining resource names.
10#
11# This is important because there may be system limits for such resources: for
12# instance, the system only supports a limited number of named semaphores, and
13# shared-memory segments live in the RAM. If a python process leaks such a
14# resource, this resource will not be removed till the next reboot. Without
15# this resource tracker process, "killall python" would probably leave unlinked
16# resources.
17
18import base64
19import os
20import signal
21import sys
22import threading
23import time
24import warnings
25from collections import deque
26from types import SimpleNamespace
27
28import json
29
30# loky: fix imports
31from multiprocessing import spawn
32from multiprocessing import util
33
34__all__ = ['ensure_running', 'register', 'unregister']
35
36_HAVE_SIGMASK = hasattr(signal, 'pthread_sigmask')
37_IGNORED_SIGNALS = (signal.SIGINT, signal.SIGTERM)
38
39def cleanup_noop(name):
40 raise RuntimeError('noop should never be registered or cleaned up')
41
42_CLEANUP_FUNCS = {
43 'noop': cleanup_noop,
44 'dummy': lambda name: None, # Dummy resource used in tests
45}
46
47if os.name == 'posix':
48 import _multiprocessing
49 import _posixshmem
50
51 # Use sem_unlink() to clean up named semaphores.
52 #
53 # sem_unlink() may be missing if the Python build process detected the
54 # absence of POSIX named semaphores. In that case, no named semaphores were
55 # ever opened, so no cleanup would be necessary.
56 if hasattr(_multiprocessing, 'sem_unlink'):
57 _CLEANUP_FUNCS['semaphore'] = _multiprocessing.sem_unlink
58 _CLEANUP_FUNCS['shared_memory'] = _posixshmem.shm_unlink
59
60
61class ReentrantCallError(RuntimeError):
62 pass
63
64
65class ResourceTracker(object):
66
67 def __init__(self):
68 self._lock = threading.RLock()
69 self._fd = None
70 self._pid = None
71 self._exitcode = None
72 self._reentrant_messages = deque()
73
74 # True to use colon-separated lines, rather than JSON lines,
75 # for internal communication. (Mainly for testing).
76 # Filenames not supported by the simple format will always be sent
77 # using JSON.
78 # The reader should understand all formats.
79 self._use_simple_format = True
80
81 # Set to True by _stop_locked() if the waitpid polling loop ran to
82 # its timeout without reaping the tracker. Exposed for tests.
83 self._waitpid_timed_out = False
84
85 def _reentrant_call_error(self):
86 # gh-109629: this happens if an explicit call to the ResourceTracker
87 # gets interrupted by a garbage collection, invoking a finalizer (*)
88 # that itself calls back into ResourceTracker.
89 # (*) for example the SemLock finalizer
90 raise ReentrantCallError(
91 "Reentrant call into the multiprocessing resource tracker")
92
93 def __del__(self):
94 # making sure child processess are cleaned before ResourceTracker
95 # gets destructed.
96 # see https://github.com/python/cpython/issues/88887
97 # gh-146313: use a timeout to avoid deadlocking if a forked child
98 # still holds the pipe's write end open.
99 self._stop(use_blocking_lock=False, wait_timeout=1.0)
100
101 def _after_fork_in_child(self):
102 # gh-146313: Called in the child right after os.fork().
103 #
104 # The tracker process is a child of the *parent*, not of us, so we
105 # could never waitpid() it anyway. Clearing _pid means our __del__
106 # becomes a no-op (the early return for _pid is None).
107 #
108 # Whether we keep the inherited _fd depends on who forked us:
109 #
110 # - multiprocessing.Process with the 'fork' start method sets
111 # _fork_intent.preserve_fd before forking. The child keeps the
112 # fd and reuses the parent's tracker (gh-80849). This is safe
113 # because multiprocessing's atexit handler joins all children
114 # before the parent's __del__ runs, so by then the fd copies
115 # are gone and the parent can reap the tracker promptly.
116 #
117 # - A raw os.fork() leaves the flag unset. We close the fd in the child after forking so
118 # the parent's __del__ can reap the tracker without waiting
119 # for the child to exit. If we later need a tracker, ensure_running()
120 # will launch a fresh one.
121 self._lock._at_fork_reinit()
122 self._reentrant_messages.clear()
123 self._pid = None
124 self._exitcode = None
125 if (self._fd is not None and
126 not getattr(_fork_intent, 'preserve_fd', False)):
127 fd = self._fd
128 self._fd = None
129 try:
130 os.close(fd)
131 except OSError:
132 pass
133
134 def _stop(self, use_blocking_lock=True, wait_timeout=None):
135 if use_blocking_lock:
136 with self._lock:
137 self._stop_locked(wait_timeout=wait_timeout)
138 else:
139 acquired = self._lock.acquire(blocking=False)
140 try:
141 self._stop_locked(wait_timeout=wait_timeout)
142 finally:
143 if acquired:
144 self._lock.release()
145
146 def _stop_locked(
147 self,
148 close=os.close,
149 waitpid=os.waitpid,
150 waitstatus_to_exitcode=os.waitstatus_to_exitcode,
151 monotonic=time.monotonic,
152 sleep=time.sleep,
153 WNOHANG=getattr(os, 'WNOHANG', None),
154 wait_timeout=None,
155 ):
156 # This shouldn't happen (it might when called by a finalizer)
157 # so we check for it anyway.
158 if self._lock._recursion_count() > 1:
159 raise self._reentrant_call_error()
160 if self._fd is None:
161 # not running
162 return
163 if self._pid is None:
164 return
165
166 # closing the "alive" file descriptor stops main()
167 close(self._fd)
168 self._fd = None
169
170 try:
171 if wait_timeout is None:
172 _, status = waitpid(self._pid, 0)
173 else:
174 # gh-146313: A forked child may still hold the pipe's write
175 # end open, preventing the tracker from seeing EOF and
176 # exiting. Poll with WNOHANG to avoid blocking forever.
177 deadline = monotonic() + wait_timeout
178 delay = 0.001
179 while True:
180 result_pid, status = waitpid(self._pid, WNOHANG)
181 if result_pid != 0:
182 break
183 remaining = deadline - monotonic()
184 if remaining <= 0:
185 # The tracker is still running; it will be
186 # reparented to PID 1 (or the nearest subreaper)
187 # when we exit, and reaped there once all pipe
188 # holders release their fd.
189 self._pid = None
190 self._exitcode = None
191 self._waitpid_timed_out = True
192 return
193 delay = min(delay * 2, remaining, 0.1)
194 sleep(delay)
195 except ChildProcessError:
196 self._pid = None
197 self._exitcode = None
198 return
199
200 self._pid = None
201
202 try:
203 self._exitcode = waitstatus_to_exitcode(status)
204 except ValueError:
205 # os.waitstatus_to_exitcode may raise an exception for invalid values
206 self._exitcode = None
207
208 def getfd(self):
209 self.ensure_running()
210 return self._fd
211
212 def ensure_running(self):
213 '''Make sure that resource tracker process is running.
214
215 This can be run from any process. Usually a child process will use
216 the resource created by its parent.'''
217 return self._ensure_running_and_write()
218
219 def _teardown_dead_process(self):
220 os.close(self._fd)
221
222 # Clean-up to avoid dangling processes.
223 try:
224 # _pid can be None if this process is a child from another
225 # python process, which has started the resource_tracker.
226 if self._pid is not None:
227 os.waitpid(self._pid, 0)
228 except ChildProcessError:
229 # The resource_tracker has already been terminated.
230 pass
231 self._fd = None
232 self._pid = None
233 self._exitcode = None
234
235 warnings.warn('resource_tracker: process died unexpectedly, '
236 'relaunching. Some resources might leak.')
237
238 def _launch(self):
239 fds_to_pass = []
240 try:
241 fds_to_pass.append(sys.stderr.fileno())
242 except Exception:
243 pass
244 r, w = os.pipe()
245 try:
246 fds_to_pass.append(r)
247 # process will out live us, so no need to wait on pid
248 exe = spawn.get_executable()
249 args = [
250 exe,
251 *util._args_from_interpreter_flags(),
252 '-c',
253 f'from multiprocessing.resource_tracker import main;main({r})',
254 ]
255 # bpo-33613: Register a signal mask that will block the signals.
256 # This signal mask will be inherited by the child that is going
257 # to be spawned and will protect the child from a race condition
258 # that can make the child die before it registers signal handlers
259 # for SIGINT and SIGTERM. The mask is unregistered after spawning
260 # the child.
261 prev_sigmask = None
262 try:
263 if _HAVE_SIGMASK:
264 prev_sigmask = signal.pthread_sigmask(signal.SIG_BLOCK, _IGNORED_SIGNALS)
265 pid = util.spawnv_passfds(exe, args, fds_to_pass)
266 finally:
267 if prev_sigmask is not None:
268 signal.pthread_sigmask(signal.SIG_SETMASK, prev_sigmask)
269 except:
270 os.close(w)
271 raise
272 else:
273 self._fd = w
274 self._pid = pid
275 finally:
276 os.close(r)
277
278 def _make_probe_message(self):
279 """Return a probe message."""
280 if self._use_simple_format:
281 return b'PROBE:0:noop\n'
282 return (
283 json.dumps(
284 {"cmd": "PROBE", "rtype": "noop"},
285 ensure_ascii=True,
286 separators=(",", ":"),
287 )
288 + "\n"
289 ).encode("ascii")
290
291 def _ensure_running_and_write(self, msg=None):
292 with self._lock:
293 if self._lock._recursion_count() > 1:
294 # The code below is certainly not reentrant-safe, so bail out
295 if msg is None:
296 raise self._reentrant_call_error()
297 return self._reentrant_messages.append(msg)
298
299 if self._fd is not None:
300 # resource tracker was launched before, is it still running?
301 if msg is None:
302 to_send = self._make_probe_message()
303 else:
304 to_send = msg
305 try:
306 self._write(to_send)
307 except OSError:
308 self._teardown_dead_process()
309 self._launch()
310
311 msg = None # message was sent in probe
312 else:
313 self._launch()
314
315 while True:
316 try:
317 reentrant_msg = self._reentrant_messages.popleft()
318 except IndexError:
319 break
320 self._write(reentrant_msg)
321 if msg is not None:
322 self._write(msg)
323
324 def _check_alive(self):
325 '''Check that the pipe has not been closed by sending a probe.'''
326 try:
327 # We cannot use send here as it calls ensure_running, creating
328 # a cycle.
329 os.write(self._fd, self._make_probe_message())
330 except OSError:
331 return False
332 else:
333 return True
334
335 def register(self, name, rtype):
336 '''Register name of resource with resource tracker.'''
337 self._send('REGISTER', name, rtype)
338
339 def unregister(self, name, rtype):
340 '''Unregister name of resource with resource tracker.'''
341 self._send('UNREGISTER', name, rtype)
342
343 def _write(self, msg):
344 nbytes = os.write(self._fd, msg)
345 assert nbytes == len(msg), f"{nbytes=} != {len(msg)=}"
346
347 def _send(self, cmd, name, rtype):
348 if self._use_simple_format and '\n' not in name:
349 msg = f"{cmd}:{name}:{rtype}\n".encode("ascii")
350 if len(msg) > 512:
351 # posix guarantees that writes to a pipe of less than PIPE_BUF
352 # bytes are atomic, and that PIPE_BUF >= 512
353 raise ValueError('msg too long')
354 self._ensure_running_and_write(msg)
355 return
356
357 # POSIX guarantees that writes to a pipe of less than PIPE_BUF (512 on Linux)
358 # bytes are atomic. Therefore, we want the message to be shorter than 512 bytes.
359 # POSIX shm_open() and sem_open() require the name, including its leading slash,
360 # to be at most NAME_MAX bytes (255 on Linux)
361 # With json.dump(..., ensure_ascii=True) every non-ASCII byte becomes a 6-char
362 # escape like \uDC80.
363 # As we want the overall message to be kept atomic and therefore smaller than 512,
364 # we encode encode the raw name bytes with URL-safe Base64 - so a 255 long name
365 # will not exceed 340 bytes.
366 b = name.encode('utf-8', 'surrogateescape')
367 if len(b) > 255:
368 raise ValueError('shared memory name too long (max 255 bytes)')
369 b64 = base64.urlsafe_b64encode(b).decode('ascii')
370
371 payload = {"cmd": cmd, "rtype": rtype, "base64_name": b64}
372 msg = (json.dumps(payload, ensure_ascii=True, separators=(",", ":")) + "\n").encode("ascii")
373
374 # The entire JSON message is guaranteed < PIPE_BUF (512 bytes) by construction.
375 assert len(msg) <= 512, f"internal error: message too long ({len(msg)} bytes)"
376 assert msg.startswith(b'{')
377
378 self._ensure_running_and_write(msg)
379
380
381# loky: We want to use multiprocessing.resource_tracker._fork_intent rather
382# than create a separate variable.
383# multiprocessing.resource_tracker._fork_intent is used by
384# multiprocessing.popen_fork to distinguish raw forks from multiprocessing
385# forks.
386try:
387 from multiprocessing.resource_tracker import _fork_intent
388except ImportError:
389 # Older Python versions cannot distinguish multiprocessing forks from raw
390 # os.fork() calls. Preserve the inherited tracker fd in both cases.
391 # TODO remove when Python 3.15 is the minimum supported version. For more details about preserve_fd see
392 # https://github.com/python/cpython/pull/146316
393 # which has been merged in Python 3.15 and backported in 3.14.5 and 3.13.14
394 _fork_intent = SimpleNamespace(preserve_fd=True)
395
396# loky: begin commented out section
397# Remove module-level globals and logic which
398# are unused and can cause confusing errors at interpreter exit because of
399# discrepancies in the code from vendored Python version and used Python
400# version
401# _resource_tracker = ResourceTracker()
402# ensure_running = _resource_tracker.ensure_running
403# register = _resource_tracker.register
404# unregister = _resource_tracker.unregister
405# getfd = _resource_tracker.getfd
406
407# # gh-146313: See _after_fork_in_child docstring.
408# if hasattr(os, 'register_at_fork'):
409# os.register_at_fork(after_in_child=_resource_tracker._after_fork_in_child)
410# loky: end commented out section
411
412
413def _decode_message(line):
414 if line.startswith(b'{'):
415 try:
416 obj = json.loads(line.decode('ascii'))
417 except Exception as e:
418 raise ValueError("malformed resource_tracker message: %r" % (line,)) from e
419
420 cmd = obj["cmd"]
421 rtype = obj["rtype"]
422 b64 = obj.get("base64_name", "")
423
424 if not isinstance(cmd, str) or not isinstance(rtype, str) or not isinstance(b64, str):
425 raise ValueError("malformed resource_tracker fields: %r" % (obj,))
426
427 try:
428 name = base64.urlsafe_b64decode(b64).decode('utf-8', 'surrogateescape')
429 except ValueError as e:
430 raise ValueError("malformed resource_tracker base64_name: %r" % (b64,)) from e
431 else:
432 cmd, rest = line.strip().decode('ascii').split(':', maxsplit=1)
433 name, rtype = rest.rsplit(':', maxsplit=1)
434 return cmd, rtype, name
435
436
437def main(fd):
438 '''Run resource tracker.'''
439 # protect the process from ^C and "killall python" etc
440 signal.signal(signal.SIGINT, signal.SIG_IGN)
441 signal.signal(signal.SIGTERM, signal.SIG_IGN)
442 if _HAVE_SIGMASK:
443 signal.pthread_sigmask(signal.SIG_UNBLOCK, _IGNORED_SIGNALS)
444
445 for f in (sys.stdin, sys.stdout):
446 try:
447 f.close()
448 except Exception:
449 pass
450
451 cache = {rtype: set() for rtype in _CLEANUP_FUNCS.keys()}
452 exit_code = 0
453
454 try:
455 # keep track of registered/unregistered resources
456 with open(fd, 'rb') as f:
457 for line in f:
458 try:
459 cmd, rtype, name = _decode_message(line)
460 cleanup_func = _CLEANUP_FUNCS.get(rtype, None)
461 if cleanup_func is None:
462 raise ValueError(
463 f'Cannot register {name} for automatic cleanup: '
464 f'unknown resource type {rtype}')
465
466 if cmd == 'REGISTER':
467 cache[rtype].add(name)
468 elif cmd == 'UNREGISTER':
469 cache[rtype].remove(name)
470 elif cmd == 'PROBE':
471 pass
472 else:
473 raise RuntimeError('unrecognized command %r' % cmd)
474 except Exception:
475 exit_code = 3
476 try:
477 sys.excepthook(*sys.exc_info())
478 except:
479 pass
480 finally:
481 # all processes have terminated; cleanup any remaining resources
482 for rtype, rtype_cache in cache.items():
483 if rtype_cache:
484 try:
485 exit_code = 1
486 if rtype == 'dummy':
487 # The test 'dummy' resource is expected to leak.
488 # We skip the warning (and *only* the warning) for it.
489 pass
490 else:
491 warnings.warn(
492 f'resource_tracker: There appear to be '
493 f'{len(rtype_cache)} leaked {rtype} objects to '
494 f'clean up at shutdown: {rtype_cache}'
495 )
496 except Exception:
497 pass
498 for name in rtype_cache:
499 # For some reason the process which created and registered this
500 # resource has failed to unregister it. Presumably it has
501 # died. We therefore unlink it.
502 try:
503 try:
504 _CLEANUP_FUNCS[rtype](name)
505 except Exception as e:
506 exit_code = 2
507 warnings.warn('resource_tracker: %r: %s' % (name, e))
508 finally:
509 pass
510
511 sys.exit(exit_code)