Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/joblib/externals/loky/backend/resource_tracker.py: 20%

Shortcuts on this page

r m x   toggle line displays

j k   next/prev highlighted chunk

0   (zero) top of page

1   (one) first highlighted chunk

211 statements  

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