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

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

227 statements  

1############################################################################### 

2# Basic context management with LokyContext 

3# 

4# author: Thomas Moreau and Olivier Grisel 

5# 

6# adapted from multiprocessing/context.py 

7# * Create a context ensuring loky uses only objects that are compatible 

8# * Add LokyContext to the list of context of multiprocessing so loky can be 

9# used with multiprocessing.set_start_method 

10# * Implement a CFS-aware amd physical-core aware cpu_count function. 

11# 

12import ctypes 

13import math 

14import multiprocessing as mp 

15import os 

16import subprocess 

17import sys 

18import traceback 

19import warnings 

20 

21from ctypes import wintypes 

22from multiprocessing import get_context as mp_get_context 

23from multiprocessing.context import BaseContext 

24from concurrent.futures.process import _MAX_WINDOWS_WORKERS 

25 

26 

27from .process import LokyProcess, LokyInitMainProcess 

28 

29# Apparently, on older Python versions, loky cannot work 61 workers on Windows 

30# but instead 60: ¯\_(ツ)_/¯ 

31if sys.version_info < (3, 10): 

32 _MAX_WINDOWS_WORKERS = _MAX_WINDOWS_WORKERS - 1 

33 

34START_METHODS = ["loky", "loky_init_main", "spawn"] 

35if sys.platform != "win32": 

36 START_METHODS += ["fork", "forkserver"] 

37 

38_DEFAULT_START_METHOD = None 

39 

40# Cache for the number of physical cores to avoid repeating subprocess calls. 

41# It should not change during the lifetime of the program. 

42physical_cores_cache = None 

43 

44 

45def get_context(method=None): 

46 # Try to overload the default context 

47 method = method or _DEFAULT_START_METHOD or "loky" 

48 if method == "fork": 

49 # If 'fork' is explicitly requested, warn user about potential issues. 

50 warnings.warn( 

51 "`fork` start method should not be used with " 

52 "`loky` as it does not respect POSIX. Try using " 

53 "`spawn` or `loky` instead.", 

54 UserWarning, 

55 ) 

56 try: 

57 return mp_get_context(method) 

58 except ValueError: 

59 raise ValueError( 

60 f"Unknown context '{method}'. Value should be in " 

61 f"{START_METHODS}." 

62 ) 

63 

64 

65def set_start_method(method, force=False): 

66 global _DEFAULT_START_METHOD 

67 if _DEFAULT_START_METHOD is not None and not force: 

68 raise RuntimeError("context has already been set") 

69 assert method is None or method in START_METHODS, ( 

70 f"'{method}' is not a valid start_method. It should be in " 

71 f"{START_METHODS}" 

72 ) 

73 

74 _DEFAULT_START_METHOD = method 

75 

76 

77def get_start_method(): 

78 return _DEFAULT_START_METHOD 

79 

80 

81def cpu_count(only_physical_cores=False): 

82 """Return the number of CPUs the current process can use. 

83 

84 The returned number of CPUs accounts for: 

85 * the number of CPUs in the system, as given by 

86 ``multiprocessing.cpu_count``; 

87 * the CPU affinity settings of the current process 

88 (available on some Unix systems); 

89 * Cgroup CPU bandwidth limit (available on Linux only, typically 

90 set by docker and similar container orchestration systems); 

91 * the value of the LOKY_MAX_CPU_COUNT environment variable if defined. 

92 and is given as the minimum of these constraints. 

93 

94 If ``only_physical_cores`` is True, return the number of physical cores 

95 instead of the number of logical cores (hyperthreading / SMT). Note that 

96 this option is not enforced if the number of usable cores is controlled in 

97 any other way such as: process affinity, Cgroup restricted CPU bandwidth 

98 or the LOKY_MAX_CPU_COUNT environment variable. If the number of physical 

99 cores is not found, return the number of logical cores. 

100 

101 Note that on Windows, the returned number of CPUs cannot exceed 61 (or 60 for 

102 Python < 3.10), see: 

103 https://bugs.python.org/issue26903. 

104 

105 It is also always larger or equal to 1. 

106 """ 

107 # Note: os.cpu_count() is allowed to return None in its docstring 

108 os_cpu_count = os.cpu_count() or 1 

109 if sys.platform == "win32": 

110 # On Windows, attempting to use more than 61 CPUs would result in a 

111 # OS-level error. See https://bugs.python.org/issue26903. According to 

112 # https://learn.microsoft.com/en-us/windows/win32/procthread/processor-groups 

113 # it might be possible to go beyond with a lot of extra work but this 

114 # does not look easy. 

115 os_cpu_count = min(os_cpu_count, _MAX_WINDOWS_WORKERS) 

116 

117 cpu_count_user = _cpu_count_user(os_cpu_count) 

118 aggregate_cpu_count = max(min(os_cpu_count, cpu_count_user), 1) 

119 

120 if not only_physical_cores: 

121 return aggregate_cpu_count 

122 

123 if cpu_count_user < os_cpu_count: 

124 # Respect user setting 

125 return max(cpu_count_user, 1) 

126 

127 cpu_count_physical, exception = _count_physical_cores() 

128 if cpu_count_physical != "not found": 

129 return cpu_count_physical 

130 

131 # Fallback to default behavior 

132 if exception is not None: 

133 # warns only the first time 

134 warnings.warn( 

135 "Could not find the number of physical cores for the " 

136 f"following reason:\n{exception}\n" 

137 "Returning the number of logical cores instead. You can " 

138 "silence this warning by setting LOKY_MAX_CPU_COUNT to " 

139 "the number of cores you want to use." 

140 ) 

141 traceback.print_tb(exception.__traceback__) 

142 

143 return aggregate_cpu_count 

144 

145 

146def _cpu_count_cgroup(os_cpu_count): 

147 # Cgroup CPU bandwidth limit available in Linux since 2.6 kernel 

148 cpu_max_fname = "/sys/fs/cgroup/cpu.max" 

149 cfs_quota_fname = "/sys/fs/cgroup/cpu/cpu.cfs_quota_us" 

150 cfs_period_fname = "/sys/fs/cgroup/cpu/cpu.cfs_period_us" 

151 

152 cpu_quota_us = None 

153 cpu_period_us = None 

154 

155 if os.path.exists(cpu_max_fname): 

156 # cgroup v2 

157 # https://www.kernel.org/doc/html/latest/admin-guide/cgroup-v2.html 

158 with open(cpu_max_fname) as fh: 

159 # Parse the quota and period values 

160 parts = fh.read().strip().split() 

161 if len(parts) == 2: 

162 cpu_quota_us, cpu_period_us = parts 

163 # If len(parts) != 2, leave as None and fall back to v1 

164 

165 # If we didn't get values from cgroup v2, try cgroup v1 

166 if cpu_quota_us is None or cpu_period_us is None: 

167 if os.path.exists(cfs_quota_fname) and os.path.exists( 

168 cfs_period_fname 

169 ): 

170 # cgroup v1 

171 # https://www.kernel.org/doc/html/latest/scheduler/sched-bwc.html#management 

172 with open(cfs_quota_fname) as fh: 

173 cpu_quota_us = fh.read().strip() 

174 with open(cfs_period_fname) as fh: 

175 cpu_period_us = fh.read().strip() 

176 else: 

177 # No Cgroup CPU bandwidth limit (e.g. non-Linux platform) 

178 cpu_quota_us = "max" 

179 

180 if cpu_quota_us == "max": 

181 # No active Cgroup quota on a Cgroup-capable platform 

182 return os_cpu_count 

183 else: 

184 cpu_quota_us = int(cpu_quota_us) 

185 cpu_period_us = int(cpu_period_us) 

186 if cpu_quota_us > 0 and cpu_period_us > 0: 

187 return math.ceil(cpu_quota_us / cpu_period_us) 

188 else: # pragma: no cover 

189 # Setting a negative cpu_quota_us value is a valid way to disable 

190 # cgroup CPU bandwidth limits 

191 return os_cpu_count 

192 

193 

194def _cpu_count_affinity(os_cpu_count): 

195 # Number of available CPUs given affinity settings 

196 if hasattr(os, "sched_getaffinity"): 

197 try: 

198 return len(os.sched_getaffinity(0)) 

199 except NotImplementedError: 

200 pass 

201 

202 # On some platforms, os.sched_getaffinity does not exist or raises 

203 # NotImplementedError, let's try with the psutil if installed. 

204 try: 

205 import psutil 

206 

207 p = psutil.Process() 

208 if hasattr(p, "cpu_affinity"): 

209 return len(p.cpu_affinity()) 

210 

211 except ImportError: # pragma: no cover 

212 if ( 

213 sys.platform == "linux" 

214 and os.environ.get("LOKY_MAX_CPU_COUNT") is None 

215 ): 

216 # Some platforms don't implement os.sched_getaffinity on Linux which 

217 # can cause severe oversubscription problems. Better warn the 

218 # user in this particularly pathological case which can wreck 

219 # havoc, typically on CI workers. 

220 warnings.warn( 

221 "Failed to inspect CPU affinity constraints on this system. " 

222 "Please install psutil or explictly set LOKY_MAX_CPU_COUNT." 

223 ) 

224 

225 # This can happen for platforms that do not implement any kind of CPU 

226 # infinity such as macOS-based platforms. 

227 return os_cpu_count 

228 

229 

230def _cpu_count_user(os_cpu_count): 

231 """Number of user defined available CPUs""" 

232 cpu_count_affinity = _cpu_count_affinity(os_cpu_count) 

233 

234 cpu_count_cgroup = _cpu_count_cgroup(os_cpu_count) 

235 

236 # User defined soft-limit passed as a loky specific environment variable. 

237 cpu_count_loky = int(os.environ.get("LOKY_MAX_CPU_COUNT", os_cpu_count)) 

238 

239 return min(cpu_count_affinity, cpu_count_cgroup, cpu_count_loky) 

240 

241 

242def _count_physical_cores(): 

243 """Return a tuple (number of physical cores, exception) 

244 

245 If the number of physical cores is found, exception is set to None. 

246 If it has not been found, return ("not found", exception). 

247 

248 The number of physical cores is cached to avoid repeating subprocess calls. 

249 """ 

250 exception = None 

251 

252 # First check if the value is cached 

253 global physical_cores_cache 

254 if physical_cores_cache is not None: 

255 return physical_cores_cache, exception 

256 

257 # Not cached yet, find it 

258 try: 

259 if sys.platform == "linux": 

260 cpu_count_physical = _count_physical_cores_linux() 

261 elif sys.platform == "win32": 

262 cpu_count_physical = _count_physical_cores_win32() 

263 elif sys.platform == "darwin": 

264 cpu_count_physical = _count_physical_cores_darwin() 

265 elif sys.platform.startswith("freebsd"): 

266 cpu_count_physical = _count_physical_cores_freebsd() 

267 else: 

268 raise NotImplementedError(f"unsupported platform: {sys.platform}") 

269 

270 # if cpu_count_physical < 1, we did not find a valid value 

271 if cpu_count_physical < 1: 

272 raise ValueError(f"found {cpu_count_physical} physical cores < 1") 

273 

274 except Exception as e: 

275 exception = e 

276 cpu_count_physical = "not found" 

277 

278 # Put the result in cache 

279 physical_cores_cache = cpu_count_physical 

280 

281 return cpu_count_physical, exception 

282 

283 

284def _count_physical_cores_linux(): 

285 try: 

286 cpu_info = subprocess.run( 

287 "lscpu --parse=core".split(), capture_output=True, text=True 

288 ) 

289 cpu_info = cpu_info.stdout.splitlines() 

290 cpu_info = {line for line in cpu_info if not line.startswith("#")} 

291 return len(cpu_info) 

292 except Exception: 

293 pass # fallback to /proc/cpuinfo 

294 

295 cpu_info = subprocess.run( 

296 "cat /proc/cpuinfo".split(), capture_output=True, text=True 

297 ) 

298 cpu_info = cpu_info.stdout.splitlines() 

299 cpu_info = {line for line in cpu_info if line.startswith("core id")} 

300 return len(cpu_info) 

301 

302 

303def _count_physical_cores_win32(): 

304 try: 

305 return _count_physical_cores_win32_ctypes() 

306 except Exception: 

307 pass # fallback to powershell 

308 try: 

309 return _count_physical_cores_win32_powershell() 

310 except Exception: 

311 pass # fallback to wmic (older Windows versions; deprecated now) 

312 

313 cpu_info = subprocess.run( 

314 "wmic CPU Get NumberOfCores /Format:csv".split(), 

315 capture_output=True, 

316 text=True, 

317 creationflags=subprocess.CREATE_NO_WINDOW, 

318 ) 

319 cpu_info = cpu_info.stdout.splitlines() 

320 cpu_info = [ 

321 l.split(",")[1] for l in cpu_info if (l and l != "Node,NumberOfCores") 

322 ] 

323 return sum(map(int, cpu_info)) 

324 

325 

326def _count_physical_cores_win32_powershell(): 

327 cmd = "-NoProfile -Command (Get-CimInstance -ClassName Win32_Processor).NumberOfCores" 

328 cpu_info = subprocess.run( 

329 f"powershell.exe {cmd}".split(), 

330 capture_output=True, 

331 text=True, 

332 creationflags=subprocess.CREATE_NO_WINDOW, 

333 ) 

334 cpu_info = cpu_info.stdout.splitlines() 

335 return sum(map(int, cpu_info)) 

336 

337 

338def _count_physical_cores_win32_ctypes(): 

339 ERROR_INSUFFICIENT_BUFFER = 122 

340 RelationProcessorCore = 0 

341 

342 kernel32 = ctypes.WinDLL("kernel32", use_last_error=True) 

343 get_logical_processor_information = ( 

344 kernel32.GetLogicalProcessorInformationEx 

345 ) 

346 get_logical_processor_information.argtypes = [ 

347 wintypes.DWORD, 

348 ctypes.c_void_p, 

349 ctypes.POINTER(wintypes.DWORD), 

350 ] 

351 get_logical_processor_information.restype = wintypes.BOOL 

352 

353 # Mirror the header of SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX. The full 

354 # structure is variable-sized, and only Relationship and Size are needed. 

355 # https://learn.microsoft.com/en-us/windows/win32/api/winnt/ns-winnt-system_logical_processor_information_ex 

356 class SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX(ctypes.Structure): 

357 _fields_ = [ 

358 ("Relationship", wintypes.DWORD), 

359 ("Size", wintypes.DWORD), 

360 ] 

361 

362 # First obtain the required buffer size. This call is expected to fail with 

363 # ERROR_INSUFFICIENT_BUFFER and set returned_length. 

364 # https://learn.microsoft.com/en-us/windows/win32/api/sysinfoapi/nf-sysinfoapi-getlogicalprocessorinformationex 

365 returned_length = wintypes.DWORD() 

366 returned_length_ref = ctypes.byref(returned_length) 

367 if get_logical_processor_information( 

368 RelationProcessorCore, None, returned_length_ref 

369 ): 

370 raise RuntimeError("unexpected successful buffer sizing call") 

371 

372 error = ctypes.get_last_error() 

373 if error != ERROR_INSUFFICIENT_BUFFER: 

374 raise ctypes.WinError(error) 

375 

376 buf = ctypes.create_string_buffer(returned_length.value) 

377 if not get_logical_processor_information( 

378 RelationProcessorCore, buf, returned_length_ref 

379 ): 

380 raise ctypes.WinError(ctypes.get_last_error()) 

381 

382 offset = 0 

383 physical_core_count = 0 

384 header_size = ctypes.sizeof(SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX) 

385 

386 while offset < returned_length.value: 

387 remaining = returned_length.value - offset 

388 if remaining < header_size: 

389 raise RuntimeError("truncated processor information record") 

390 

391 processor_core_info = ( 

392 SYSTEM_LOGICAL_PROCESSOR_INFORMATION_EX.from_buffer(buf, offset) 

393 ) 

394 record_size = processor_core_info.Size 

395 if record_size < header_size or record_size > remaining: 

396 raise RuntimeError("invalid processor information record size") 

397 if processor_core_info.Relationship != RelationProcessorCore: 

398 raise RuntimeError("unexpected logical processor relationship") 

399 

400 physical_core_count += 1 

401 offset += record_size 

402 

403 if physical_core_count == 0: 

404 raise RuntimeError("Windows reported no active physical cores") 

405 

406 if physical_core_count < 1: 

407 raise RuntimeError( 

408 "GetLogicalProcessorInformationEx returned no physical cores" 

409 ) 

410 return physical_core_count 

411 

412 

413def _count_physical_cores_darwin(): 

414 cpu_info = subprocess.run( 

415 "sysctl -n hw.physicalcpu".split(), 

416 capture_output=True, 

417 text=True, 

418 ) 

419 cpu_info = cpu_info.stdout 

420 return int(cpu_info) 

421 

422 

423def _count_physical_cores_freebsd(): 

424 cpu_info = subprocess.run( 

425 "sysctl -n kern.smp.cores".split(), 

426 capture_output=True, 

427 text=True, 

428 ) 

429 cpu_info = cpu_info.stdout 

430 return int(cpu_info) 

431 

432 

433class LokyContext(BaseContext): 

434 """Context relying on the LokyProcess.""" 

435 

436 _name = "loky" 

437 Process = LokyProcess 

438 cpu_count = staticmethod(cpu_count) 

439 

440 def Queue(self, maxsize=0, reducers=None): 

441 """Returns a queue object""" 

442 from .queues import Queue 

443 

444 return Queue(maxsize, reducers=reducers, ctx=self.get_context()) 

445 

446 def SimpleQueue(self, reducers=None): 

447 """Returns a queue object""" 

448 from .queues import SimpleQueue 

449 

450 return SimpleQueue(reducers=reducers, ctx=self.get_context()) 

451 

452 if sys.platform != "win32": 

453 """For Unix platform, use our custom implementation of synchronize 

454 ensuring that we use the loky.backend.resource_tracker to clean-up 

455 the semaphores in case of a worker crash. 

456 """ 

457 

458 def Semaphore(self, value=1): 

459 """Returns a semaphore object""" 

460 from .synchronize import Semaphore 

461 

462 return Semaphore(value=value) 

463 

464 def BoundedSemaphore(self, value): 

465 """Returns a bounded semaphore object""" 

466 from .synchronize import BoundedSemaphore 

467 

468 return BoundedSemaphore(value) 

469 

470 def Lock(self): 

471 """Returns a lock object""" 

472 from .synchronize import Lock 

473 

474 return Lock() 

475 

476 def RLock(self): 

477 """Returns a recurrent lock object""" 

478 from .synchronize import RLock 

479 

480 return RLock() 

481 

482 def Condition(self, lock=None): 

483 """Returns a condition object""" 

484 from .synchronize import Condition 

485 

486 return Condition(lock) 

487 

488 def Event(self): 

489 """Returns an event object""" 

490 from .synchronize import Event 

491 

492 return Event() 

493 

494 

495class LokyInitMainContext(LokyContext): 

496 """Extra context with LokyProcess, which does load the main module 

497 

498 This context is used for compatibility in the case ``cloudpickle`` is not 

499 present on the running system. This permits to load functions defined in 

500 the ``main`` module, using proper safeguards. The declaration of the 

501 ``executor`` should be protected by ``if __name__ == "__main__":`` and the 

502 functions and variable used from main should be out of this block. 

503 

504 This mimics the default behavior of multiprocessing under Windows and the 

505 behavior of the ``spawn`` start method on a posix system. 

506 For more details, see the end of the following section of python doc 

507 https://docs.python.org/3/library/multiprocessing.html#multiprocessing-programming 

508 """ 

509 

510 _name = "loky_init_main" 

511 Process = LokyInitMainProcess 

512 

513 

514# Register loky context so it works with multiprocessing.get_context 

515ctx_loky = LokyContext() 

516mp.context._concrete_contexts["loky"] = ctx_loky 

517mp.context._concrete_contexts["loky_init_main"] = LokyInitMainContext()