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

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

267 statements  

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)