1###############################################################################
2# Server process to keep track of unlinked resources, like folders and
3# semaphores and clean them.
4#
5# author: Thomas Moreau
6#
7# Adapted from multiprocessing/resource_tracker.py
8# * add some VERBOSE logging,
9# * add support to track folders,
10# * add Windows support,
11# * refcounting scheme to avoid unlinking resources still in use.
12#
13# On Unix we run a server process which keeps track of unlinked
14# resources. The server ignores SIGINT and SIGTERM and reads from a
15# pipe. The resource_tracker implements a reference counting scheme: each time
16# a Python process anticipates the shared usage of a resource by another
17# process, it signals the resource_tracker of this shared usage, and in return,
18# the resource_tracker increments the resource's reference count by 1.
19# Similarly, when access to a resource is closed by a Python process, the
20# process notifies the resource_tracker by asking it to decrement the
21# resource's reference count by 1. When the reference count drops to 0, the
22# resource_tracker attempts to clean up the underlying resource.
23
24# Finally, every other process connected to the resource tracker has a copy of
25# the writable end of the pipe used to communicate with it, so the resource
26# tracker gets EOF when all other processes have exited. Then the
27# resource_tracker process unlinks any remaining leaked resources (with
28# reference count above 0)
29
30# For semaphores, this is important because the system only supports a limited
31# number of named semaphores, and they will not be automatically removed till
32# the next reboot. Without this resource tracker process, "killall python"
33# would probably leave unlinked semaphores.
34
35# Note that this behavior differs from CPython's resource_tracker, which only
36# implements list of shared resources, and not a proper refcounting scheme.
37# Also, CPython's resource tracker will only attempt to cleanup those shared
38# resources once all processes connected to the resource tracker have exited.
39#
40# Source code organization and maintenance:
41#
42# This file is derived from the matching file in the standard library with
43# modifications to add the new features described above.
44#
45# It defines a subclass of a vendored copy of the ResourceTracker from the
46# standard library to make loky less likely to break when internals of
47# change in a new version of CPython.
48#
49# The new features are implemented in the overriden method of the ResourceTracker
50# class as well as in the custom `main` function.
51#
52# The vendored copy is minimally patched and should not be edited by hand.
53# Instead, it should be revendored from time to time using the script
54# in the `tools/` folder at the root of this repo.
55#
56# When vendoring a new copy of the resource_tracker module of the standard
57# library, it is important to check of the `main` function has evolved
58# upstream. If so, doing a diff of that particular function in the upstream
59# file and in the current file might be helpful.
60#
61# Make sure to update the inline comments to help future maintainers understand
62# what loky-specific changes were made.
63
64import os
65import shutil
66import sys
67import signal
68import warnings
69from multiprocessing import util
70import base64
71import json
72import threading
73
74from .stdlib_py314_resource_tracker import (
75 ResourceTracker as StdLibResourceTracker,
76 _decode_message,
77)
78
79from . import spawn
80
81if sys.platform == "win32":
82 import _winapi
83 import msvcrt
84 from multiprocessing.reduction import duplicate
85
86# To minimize diff vs stdlib
87# fmt:off
88__all__ = ['ensure_running', 'register', 'unregister']
89
90_HAVE_SIGMASK = hasattr(signal, 'pthread_sigmask')
91_IGNORED_SIGNALS = (signal.SIGINT, signal.SIGTERM)
92
93def cleanup_noop(name):
94 raise RuntimeError('noop should never be registered or cleaned up')
95
96
97_CLEANUP_FUNCS = {
98 'noop': cleanup_noop,
99 'dummy': lambda name: None, # Dummy resource used in tests
100 # loky: add 'folder' and 'file' resources
101 'folder': shutil.rmtree,
102 'file': os.unlink,
103}
104
105if os.name == 'posix':
106 import _multiprocessing
107
108 # Use sem_unlink() to clean up named semaphores.
109 #
110 # sem_unlink() may be missing if the Python build process detected the
111 # absence of POSIX named semaphores. In that case, no named semaphores were
112 # ever opened, so no cleanup would be necessary.
113 if hasattr(_multiprocessing, 'sem_unlink'):
114 _CLEANUP_FUNCS.update(
115 {
116 'semlock': _multiprocessing.sem_unlink,
117 }
118 )
119# fmt: on
120
121# loky: logging
122VERBOSE = False
123
124
125# loky: compatibility for CPython versions that don't have _RLock._recursion_count
126# This was done in CPython 3.13 in
127# 'Fix reentrancy issue in multiprocessing resource_tracker'
128# https://github.com/python/cpython/pull/109629
129# This was back-ported in 3.11.6 and 3.12.1
130# TODO Remove work-around when Python 3.13 is our minimum supported version
131class LokyRLock(type(threading.RLock())):
132 def _recursion_count(self):
133 return 1
134
135
136class ResourceTracker(StdLibResourceTracker):
137 """Resource tracker with refcounting scheme.
138
139 This class is an extension of the multiprocessing ResourceTracker class
140 which implements a reference counting scheme to avoid unlinking shared
141 resources still in use in other processes.
142
143 This feature is notably used by `joblib.Parallel` to share temporary
144 folders and memory mapped files between the main process and the worker
145 processes.
146
147 The actual implementation of the refcounting scheme is in the main
148 function, which is run in a dedicated process.
149 """
150
151 def __init__(self):
152 super().__init__()
153 # Windows process handle returned by CreateProcess. A pid cannot be
154 # passed to os.waitpid on Windows because the underlying _cwait expects
155 # a process handle.
156 self._proc_handle = None
157 # TODO Remove block when Python 3.13 is our minimum supported version
158 # see above comment about _recursion_count
159 if not hasattr(self._lock, "_recursion_count"):
160 self._lock = LokyRLock()
161
162 def maybe_unlink(self, name, rtype):
163 """Decrement the refcount of a resource, and delete it if it hits 0"""
164 self._send("MAYBE_UNLINK", name, rtype)
165
166 def _teardown_dead_process(self):
167 if os.name == "posix":
168 super()._teardown_dead_process()
169 elif sys.platform == "win32":
170 os.close(self._fd)
171 if (proc_handle := self._proc_handle) is not None:
172 _winapi.CloseHandle(proc_handle)
173 # All 3 lines copied from stdlib _teardown_dead_processes
174 self._fd = None
175 self._pid = None
176 self._exitcode = None
177 self._proc_handle = None
178
179 warnings.warn(
180 "resource_tracker: process died unexpectedly, "
181 "relaunching. Some resources might leak."
182 )
183
184 # To minimize the diff with stdlib ResourceTracker._launch
185 # fmt: off
186 def _launch(self):
187 # This is copied from Python 3.14.7 with loky additions/modifications
188 # mostly for Windows support and logging.
189 # Added or changed lines have a comment that starts with "# loky:"
190 fds_to_pass = []
191 try:
192 fds_to_pass.append(sys.stderr.fileno())
193 except Exception:
194 pass
195 r, w = os.pipe()
196 # loky: Windows support
197 if sys.platform == "win32":
198 _r = duplicate(msvcrt.get_osfhandle(r), inheritable=True)
199 os.close(r)
200 r = _r
201
202 try:
203 fds_to_pass.append(r)
204 # process will out live us, so no need to wait on pid
205 exe = spawn.get_executable()
206 args = [
207 exe,
208 *util._args_from_interpreter_flags(),
209 '-c',
210 # loky: use loky main function rather than stdlib one
211 f'from {main.__module__} import main; main({r}, {VERBOSE})'
212 ]
213 # loky: logging
214 util.debug(f"launching resource tracker: {args}")
215 # bpo-33613: Register a signal mask that will block the signals.
216 # This signal mask will be inherited by the child that is going
217 # to be spawned and will protect the child from a race condition
218 # that can make the child die before it registers signal handlers
219 # for SIGINT and SIGTERM. The mask is unregistered after spawning
220 # the child.
221 prev_sigmask = None
222 try:
223 if _HAVE_SIGMASK:
224 prev_sigmask = signal.pthread_sigmask(signal.SIG_BLOCK, _IGNORED_SIGNALS)
225 # loky: call loky spawnv_passfds which supports Windows
226 pid, proc_handle = spawnv_passfds(exe, args, fds_to_pass)
227 finally:
228 if prev_sigmask is not None:
229 signal.pthread_sigmask(signal.SIG_SETMASK, prev_sigmask)
230 except:
231 os.close(w)
232 raise
233 else:
234 self._fd = w
235 self._pid = pid
236 self._proc_handle = proc_handle
237 finally:
238 # loky: Windows support
239 if sys.platform == "win32":
240 _winapi.CloseHandle(r)
241 else:
242 os.close(r)
243 # fmt: on
244
245 if sys.platform == "win32":
246 # stdlib ResourceTracker._stop_locked is POSIX-specific, so we override
247 # to be able to use Windows-specific primitives.
248 # This is loosely inspired from stdlib ResourceTracker._stop_locked but
249 # there are a number of changes. TODO: could the structure be closer?
250 def _stop_locked(
251 self,
252 close=os.close,
253 wait_timeout=None,
254 wait_for_single_object=_winapi.WaitForSingleObject,
255 get_exit_code_process=_winapi.GetExitCodeProcess,
256 close_handle=_winapi.CloseHandle,
257 winapi_infinite=_winapi.INFINITE,
258 wait_timeout_code=_winapi.WAIT_TIMEOUT,
259 ):
260 # This shouldn't happen (it might when called by a finalizer)
261 # so we check for it anyway.
262 if self._lock._recursion_count() > 1:
263 raise self._reentrant_call_error()
264 if self._fd is None:
265 # not running
266 return
267 if self._pid is None:
268 return
269
270 # Closing the "alive" file descriptor asks the tracker to stop.
271 close(self._fd)
272 self._fd = None
273
274 proc_handle = self._proc_handle
275 try:
276 if proc_handle is not None:
277 timeout_ms = (
278 winapi_infinite
279 if wait_timeout is None
280 else round(wait_timeout * 1000)
281 )
282 wait_result = wait_for_single_object(
283 proc_handle, timeout_ms
284 )
285 if wait_result == wait_timeout_code:
286 self._pid = None
287 self._exitcode = None
288 self._waitpid_timed_out = True
289 return
290
291 self._pid = None
292 self._exitcode = get_exit_code_process(proc_handle)
293 else:
294 self._pid = None
295 self._exitcode = None
296 finally:
297 if proc_handle is not None:
298 close_handle(proc_handle)
299 self._proc_handle = None
300
301
302# Copied from Python 3.14.7
303# fmt: off
304_resource_tracker = ResourceTracker()
305ensure_running = _resource_tracker.ensure_running
306register = _resource_tracker.register
307maybe_unlink = _resource_tracker.maybe_unlink
308unregister = _resource_tracker.unregister
309getfd = _resource_tracker.getfd
310
311# gh-146313: See _after_fork_in_child docstring.
312if hasattr(os, 'register_at_fork'):
313 os.register_at_fork(after_in_child=_resource_tracker._after_fork_in_child)
314# fmt: on
315
316# fmt: off
317# The main function has been copied from Python 3.14.7 and modified, mostly for
318# Windows support, logging and refcount functionality.
319# Added or changed lines have a comment that starts with "# loky:"
320# loky: add verbose argument for logging
321def main(fd, verbose=0):
322 '''Run resource tracker.'''
323
324 # loky: logging
325 if verbose:
326 util.log_to_stderr(level=util.DEBUG)
327
328 # protect the process from ^C and "killall python" etc
329 signal.signal(signal.SIGINT, signal.SIG_IGN)
330 signal.signal(signal.SIGTERM, signal.SIG_IGN)
331
332 if _HAVE_SIGMASK:
333 signal.pthread_sigmask(signal.SIG_UNBLOCK, _IGNORED_SIGNALS)
334
335 for f in (sys.stdin, sys.stdout):
336 try:
337 f.close()
338 except Exception:
339 pass
340
341 # loky: logging
342 if verbose:
343 util.debug("Main resource tracker is running")
344
345 # loky: change for refcount functionality we want a dict[str, dict] rather
346 # than a dict[str, set] so that cache[folder]['resource'] is the refcount
347 # associated with it
348 cache = {rtype: dict() for rtype in _CLEANUP_FUNCS.keys()}
349 exit_code = 0
350
351 try:
352 # loky: Windows support
353 if sys.platform == "win32":
354 fd = msvcrt.open_osfhandle(fd, os.O_RDONLY)
355 # keep track of registered/unregistered resources
356 with open(fd, 'rb') as f:
357 for line in f:
358 try:
359 cmd, rtype, name = _decode_message(line)
360 cleanup_func = _CLEANUP_FUNCS.get(rtype, None)
361 if cleanup_func is None:
362 raise ValueError(
363 f'Cannot register {name} for automatic cleanup: '
364 f'unknown resource type ({rtype})'
365 # loky: additional info for possible keys
366 '. Resource type should be one of the following: '
367 f'{list(_CLEANUP_FUNCS.keys())}'
368 )
369
370 if cmd == 'REGISTER':
371 # loky: refcount functionality
372 if name not in cache[rtype]:
373 cache[rtype][name] = 1
374 else:
375 cache[rtype][name] += 1
376
377 # loky: logging
378 if verbose:
379 util.debug(
380 '[ResourceTracker] incremented refcount of '
381 f'{rtype} {name} '
382 f'(current {cache[rtype][name]})'
383 )
384 elif cmd == 'UNREGISTER':
385 # loky: refcount functionality
386 del cache[rtype][name]
387 # loky: logging
388 if verbose:
389 util.debug(
390 f'[ResourceTracker] unregister {name} {rtype}: '
391 f'cache({len(cache)})'
392 )
393 elif cmd == 'PROBE':
394 pass
395 # loky: refcount functionality with logging
396 elif cmd == 'MAYBE_UNLINK':
397 cache[rtype][name] -= 1
398 if verbose:
399 util.debug(
400 '[ResourceTracker] decremented refcount of '
401 f'{rtype} {name} '
402 f'(current {cache[rtype][name]})'
403 )
404
405 if cache[rtype][name] == 0:
406 del cache[rtype][name]
407 try:
408 if verbose:
409 util.debug(
410 f'[ResourceTracker] unlink {name}'
411 )
412 _CLEANUP_FUNCS[rtype](name)
413 except Exception as e:
414 warnings.warn(
415 f"resource_tracker: {name}: {e!r}"
416 )
417
418 else:
419 raise RuntimeError('unrecognized command %r' % cmd)
420 except Exception:
421 exit_code = 3
422 try:
423 sys.excepthook(*sys.exc_info())
424 except:
425 pass
426 finally:
427 # all processes have terminated; cleanup any remaining resources
428
429 # loky: loky wants to clean ressources first and folder last because
430 # there can be tracked resources inside tracked folders.
431 # _unlink_resources is the stdlib code with some additional logging, it
432 # is called for all resources except folders and then at the end for
433 # all folders
434 def _unlink_resources(rtype_cache, rtype):
435 nonlocal exit_code
436 if rtype_cache:
437 try:
438 exit_code = 1
439 if rtype == 'dummy':
440 # The test 'dummy' resource is expected to leak.
441 # We skip the warning (and *only* the warning) for it.
442 pass
443 else:
444 warnings.warn(
445 f'resource_tracker: There appear to be '
446 f'{len(rtype_cache)} leaked {rtype} objects to '
447 f'clean up at shutdown: {rtype_cache}'
448 )
449 except Exception:
450 pass
451 for name in rtype_cache:
452 # For some reason the process which created and registered this
453 # resource has failed to unregister it. Presumably it has
454 # died. We therefore unlink it.
455 try:
456 try:
457 _CLEANUP_FUNCS[rtype](name)
458 # loky: logging
459 if verbose:
460 util.debug(f'[ResourceTracker] unlink {name}')
461 except Exception as e:
462 exit_code = 2
463 # loky: tweaked formatting (%r instead of %s for
464 # exception and %s instead of %s for name) I guess you
465 # get the exact exception type with %r. %s instead of
466 # %r for name, may be because on Windows %r doubles the
467 # backslashes for path-like ressources
468 warnings.warn('resource_tracker: %s: %r' % (name, e))
469 finally:
470 pass
471
472 # loky: 2 stage cleaning process
473 # The default cleanup routine for folders deletes everything inside
474 # those folders recursively, which can include other resources tracked
475 # by the resource tracker). To limit the risk of the resource tracker
476 # attempting to delete twice a resource (once as part of a tracked
477 # folder, and once as a resource), we delete the folders after all
478 # other resource types.
479 for rtype, rtype_cache in cache.items():
480 if rtype == 'folder':
481 continue
482 _unlink_resources(rtype_cache, rtype)
483
484 if 'folder' in cache:
485 _unlink_resources(cache['folder'], 'folder')
486
487 # loky: logging
488 if verbose:
489 util.debug("resource tracker shut down")
490
491 # TODO not sure about exit_code management with 2-stage clean-up
492 # (first non-folder resources then folder resources) but seems good enough
493 # for now. We did not have any exit code management until
494 # https://github.com/joblib/loky/pull/472
495 sys.exit(exit_code)
496# fmt: on
497
498
499def spawnv_passfds(path, args, passfds):
500 """Loky version multiprocessing.util.spawnv_passfds with added Windows support.
501
502 Returns (pid, process_handle) because os.waitpid needs handle on Windows.
503 On Linux process_handle is None.
504 """
505 if sys.platform != "win32":
506 # loky: additional encoding needed here
507 # TODO We should fix loky.backend.spawn.get_executable to return bytes
508 # on POSIX so that this can be removed
509 path = path.encode("utf-8")
510 return util.spawnv_passfds(path, args, passfds), None
511 else:
512 # loky: Windows support
513 passfds = sorted(passfds)
514 cmd = " ".join(f'"{x}"' for x in args)
515 hp, ht, pid, _ = _winapi.CreateProcess(
516 path, cmd, None, None, True, 0, None, None, None
517 )
518 _winapi.CloseHandle(ht)
519 # Keep the process handle for a safe wait during teardown. Pids and
520 # handles are different namespaces on Windows.
521 return pid, hp