codekingpro/portable-devtools
114k
1from __future__ import annotations2 3__all__ = (4 "current_default_process_limiter",5 "process_worker",6 "run_sync",7)8 9import os10import pickle11import subprocess12import sys13from collections import deque14from collections.abc import Callable15from importlib.util import module_from_spec, spec_from_file_location16from typing import TypeVar, cast17 18from ._core._eventloop import current_time, get_async_backend, get_cancelled_exc_class19from ._core._exceptions import BrokenWorkerProcess20from ._core._subprocesses import open_process21from ._core._synchronization import CapacityLimiter22from ._core._tasks import CancelScope, fail_after23from .abc import ByteReceiveStream, ByteSendStream, Process24from .lowlevel import RunVar, checkpoint_if_cancelled25from .streams.buffered import BufferedByteReceiveStream26 27if sys.version_info >= (3, 11):28 from typing import TypeVarTuple, Unpack29else:30 from typing_extensions import TypeVarTuple, Unpack31 32WORKER_MAX_IDLE_TIME = 300 # 5 minutes33 34T_Retval = TypeVar("T_Retval")35PosArgsT = TypeVarTuple("PosArgsT")36 37_process_pool_workers: RunVar[set[Process]] = RunVar("_process_pool_workers")38_process_pool_idle_workers: RunVar[deque[tuple[Process, float]]] = RunVar(39 "_process_pool_idle_workers"40)41_default_process_limiter: RunVar[CapacityLimiter] = RunVar("_default_process_limiter")42 43 44async def run_sync( # type: ignore[return]45 func: Callable[[Unpack[PosArgsT]], T_Retval],46 *args: Unpack[PosArgsT],47 cancellable: bool = False,48 limiter: CapacityLimiter | None = None,49) -> T_Retval:50 """51 Call the given function with the given arguments in a worker process.52 53 If the ``cancellable`` option is enabled and the task waiting for its completion is54 cancelled, the worker process running it will be abruptly terminated using SIGKILL55 (or ``terminateProcess()`` on Windows).56 57 :param func: a callable58 :param args: positional arguments for the callable59 :param cancellable: ``True`` to allow cancellation of the operation while it's60 running61 :param limiter: capacity limiter to use to limit the total amount of processes62 running (if omitted, the default limiter is used)63 :raises NoEventLoopError: if no supported asynchronous event loop is running in the64 current thread65 :return: an awaitable that yields the return value of the function.66 67 """68 69 async def send_raw_command(pickled_cmd: bytes) -> object:70 try:71 await stdin.send(pickled_cmd)72 response = await buffered.receive_until(b"\n", 50)73 status, length = response.split(b" ")74 if status not in (b"RETURN", b"EXCEPTION"):75 raise RuntimeError(76 f"Worker process returned unexpected response: {response!r}"77 )78 79 pickled_response = await buffered.receive_exactly(int(length))80 except BaseException as exc:81 workers.discard(process)82 try:83 process.kill()84 with CancelScope(shield=True):85 await process.aclose()86 except ProcessLookupError:87 pass88 89 if isinstance(exc, get_cancelled_exc_class()):90 raise91 else:92 raise BrokenWorkerProcess from exc93 94 retval = pickle.loads(pickled_response)95 if status == b"EXCEPTION":96 assert isinstance(retval, BaseException)97 raise retval98 else:99 return retval100 101 # First pickle the request before trying to reserve a worker process102 await checkpoint_if_cancelled()103 request = pickle.dumps(("run", func, args), protocol=pickle.HIGHEST_PROTOCOL)104 105 # If this is the first run in this event loop thread, set up the necessary variables106 try:107 workers = _process_pool_workers.get()108 idle_workers = _process_pool_idle_workers.get()109 except LookupError:110 workers = set()111 idle_workers = deque()112 _process_pool_workers.set(workers)113 _process_pool_idle_workers.set(idle_workers)114 get_async_backend().setup_process_pool_exit_at_shutdown(workers)115 116 async with limiter or current_default_process_limiter():117 # Pop processes from the pool (starting from the most recently used) until we118 # find one that hasn't exited yet119 process: Process120 while idle_workers:121 process, idle_since = idle_workers.pop()122 if process.returncode is None:123 stdin = cast(ByteSendStream, process.stdin)124 buffered = BufferedByteReceiveStream(125 cast(ByteReceiveStream, process.stdout)126 )127 128 # Prune any other workers that have been idle for WORKER_MAX_IDLE_TIME129 # seconds or longer130 now = current_time()131 killed_processes: list[Process] = []132 while idle_workers:133 if now - idle_workers[0][1] < WORKER_MAX_IDLE_TIME:134 break135 136 process_to_kill, idle_since = idle_workers.popleft()137 process_to_kill.kill()138 workers.remove(process_to_kill)139 killed_processes.append(process_to_kill)140 141 with CancelScope(shield=True):142 for killed_process in killed_processes:143 await killed_process.aclose()144 145 break146 147 workers.remove(process)148 else:149 command = [sys.executable, "-u", "-m", __name__]150 process = await open_process(151 command, stdin=subprocess.PIPE, stdout=subprocess.PIPE152 )153 try:154 stdin = cast(ByteSendStream, process.stdin)155 buffered = BufferedByteReceiveStream(156 cast(ByteReceiveStream, process.stdout)157 )158 with fail_after(20):159 message = await buffered.receive(6)160 161 if message != b"READY\n":162 raise BrokenWorkerProcess(163 f"Worker process returned unexpected response: {message!r}"164 )165 166 main_module_path = getattr(sys.modules["__main__"], "__file__", None)167 pickled = pickle.dumps(168 ("init", sys.path, main_module_path),169 protocol=pickle.HIGHEST_PROTOCOL,170 )171 await send_raw_command(pickled)172 except (BrokenWorkerProcess, get_cancelled_exc_class()):173 raise174 except BaseException as exc:175 process.kill()176 raise BrokenWorkerProcess(177 "Error during worker process initialization"178 ) from exc179 180 workers.add(process)181 182 with CancelScope(shield=not cancellable):183 try:184 return cast(T_Retval, await send_raw_command(request))185 finally:186 if process in workers:187 idle_workers.append((process, current_time()))188 189 190def current_default_process_limiter() -> CapacityLimiter:191 """192 Return the capacity limiter that is used by default to limit the number of worker193 processes.194 195 :return: a capacity limiter object196 197 """198 try:199 return _default_process_limiter.get()200 except LookupError:201 limiter = CapacityLimiter(os.cpu_count() or 2)202 _default_process_limiter.set(limiter)203 return limiter204 205 206def process_worker() -> None:207 # Redirect standard streams to os.devnull so that user code won't interfere with the208 # parent-worker communication209 stdin = sys.stdin210 stdout = sys.stdout211 sys.stdin = open(os.devnull)212 sys.stdout = open(os.devnull, "w")213 214 stdout.buffer.write(b"READY\n")215 while True:216 retval = exception = None217 try:218 command, *args = pickle.load(stdin.buffer)219 except EOFError:220 return221 except BaseException as exc:222 exception = exc223 else:224 if command == "run":225 func, args = args226 try:227 retval = func(*args)228 except BaseException as exc:229 exception = exc230 elif command == "init":231 main_module_path: str | None232 sys.path, main_module_path = args233 del sys.modules["__main__"]234 if main_module_path and os.path.isfile(main_module_path):235 # Load the parent's main module but as __mp_main__ instead of236 # __main__ (like multiprocessing does) to avoid infinite recursion237 try:238 spec = spec_from_file_location("__mp_main__", main_module_path)239 if spec and spec.loader:240 main = module_from_spec(spec)241 spec.loader.exec_module(main)242 sys.modules["__main__"] = main243 except BaseException as exc:244 exception = exc245 try:246 if exception is not None:247 status = b"EXCEPTION"248 pickled = pickle.dumps(exception, pickle.HIGHEST_PROTOCOL)249 else:250 status = b"RETURN"251 pickled = pickle.dumps(retval, pickle.HIGHEST_PROTOCOL)252 except BaseException as exc:253 exception = exc254 status = b"EXCEPTION"255 pickled = pickle.dumps(exc, pickle.HIGHEST_PROTOCOL)256 257 stdout.buffer.write(b"%s %d\n" % (status, len(pickled)))258 stdout.buffer.write(pickled)259 260 # Respect SIGTERM261 if isinstance(exception, SystemExit):262 raise exception263 264 265if __name__ == "__main__":266 process_worker()267 