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()