Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
to_process.py267 linesDownload Raw Back to anyio
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 
codekingpro/portable-devtools · Team Ai