Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
zmq_loop.py278 linesDownload Raw Back to event_loop
1# Urwid main loop code using ZeroMQ queues2#    Copyright (C) 2019 Dave Jones3#4#    This library is free software; you can redistribute it and/or5#    modify it under the terms of the GNU Lesser General Public6#    License as published by the Free Software Foundation; either7#    version 2.1 of the License, or (at your option) any later version.8#9#    This library is distributed in the hope that it will be useful,10#    but WITHOUT ANY WARRANTY; without even the implied warranty of11#    MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the GNU12#    Lesser General Public License for more details.13#14#    You should have received a copy of the GNU Lesser General Public15#    License along with this library; if not, write to the Free Software16#    Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA  02111-1307  USA17#18# Urwid web site: https://urwid.org/19 20"""ZeroMQ based urwid EventLoop implementation.21 22`ZeroMQ <https://zeromq.org>`_ library is required.23"""24 25from __future__ import annotations26 27import contextlib28import errno29import heapq30import logging31import os32import time33import typing34from itertools import count35 36import zmq37 38from .abstract_loop import EventLoop, ExitMainLoop39 40if typing.TYPE_CHECKING:41    import io42    from collections.abc import Callable43    from concurrent.futures import Executor, Future44 45    from typing_extensions import ParamSpec46 47    ZMQAlarmHandle = typing.TypeVar("ZMQAlarmHandle")48    _T = typing.TypeVar("_T")49    _Spec = ParamSpec("_Spec")50 51 52class ZMQEventLoop(EventLoop):53    """54    This class is an urwid event loop for `ZeroMQ`_ applications. It is very55    similar to :class:`SelectEventLoop`, supporting the usual :meth:`alarm`56    events and file watching (:meth:`watch_file`) capabilities, but also57    incorporates the ability to watch zmq queues for events58    (:meth:`watch_queue`).59 60    .. _ZeroMQ: https://zeromq.org/61    """62 63    _alarm_break = count()64 65    def __init__(self) -> None:66        super().__init__()67        self.logger = logging.getLogger(__name__).getChild(self.__class__.__name__)68        self._did_something = True69        self._alarms: list[tuple[float, int, Callable[[], typing.Any]]] = []70        self._poller = zmq.Poller()71        self._queue_callbacks: dict[int, Callable[[], typing.Any]] = {}72        self._idle_handle = 073        self._idle_callbacks: dict[int, Callable[[], typing.Any]] = {}74 75    def run_in_executor(76        self,77        executor: Executor,78        func: Callable[_Spec, _T],79        *args: _Spec.args,80        **kwargs: _Spec.kwargs,81    ) -> Future[_T]:82        """Run callable in executor.83 84        :param executor: Executor to use for running the function85        :type executor: concurrent.futures.Executor86        :param func: function to call87        :type func: Callable88        :param args: positional arguments to function89        :type args: object90        :param kwargs: keyword arguments to function91        :type kwargs: object92        :return: future object for the function call outcome.93        :rtype: concurrent.futures.Future94        """95        return executor.submit(func, *args, **kwargs)96 97    def alarm(self, seconds: float, callback: Callable[[], typing.Any]) -> ZMQAlarmHandle:98        """99        Call *callback* a given time from now. No parameters are passed to100        callback. Returns a handle that may be passed to :meth:`remove_alarm`.101 102        :param float seconds:103            floating point time to wait before calling callback.104 105        :param callback:106            function to call from event loop.107        """108        handle = (time.time() + seconds, next(self._alarm_break), callback)109        heapq.heappush(self._alarms, handle)110        return handle111 112    def remove_alarm(self, handle: ZMQAlarmHandle) -> bool:113        """114        Remove an alarm. Returns ``True`` if the alarm exists, ``False``115        otherwise.116        """117        try:118            self._alarms.remove(handle)119            heapq.heapify(self._alarms)120 121        except ValueError:122            return False123 124        return True125 126    def watch_queue(127        self,128        queue: zmq.Socket,129        callback: Callable[[], typing.Any],130        flags: int = zmq.POLLIN,131    ) -> zmq.Socket:132        """133        Call *callback* when zmq *queue* has something to read (when *flags* is134        set to ``POLLIN``, the default) or is available to write (when *flags*135        is set to ``POLLOUT``). No parameters are passed to the callback.136        Returns a handle that may be passed to :meth:`remove_watch_queue`.137 138        :param queue:139            The zmq queue to poll.140 141        :param callback:142            The function to call when the poll is successful.143 144        :param int flags:145            The condition to monitor on the queue (defaults to ``POLLIN``).146        """147        if queue in self._queue_callbacks:148            raise ValueError(f"already watching {queue!r}")149        self._poller.register(queue, flags)150        self._queue_callbacks[queue] = callback151        return queue152 153    def watch_file(154        self,155        fd: int | io.TextIOWrapper,156        callback: Callable[[], typing.Any],157        flags: int = zmq.POLLIN,158    ) -> io.TextIOWrapper:159        """160        Call *callback* when *fd* has some data to read. No parameters are161        passed to the callback. The *flags* are as for :meth:`watch_queue`.162        Returns a handle that may be passed to :meth:`remove_watch_file`.163 164        :param fd:165            The file-like object, or fileno to monitor.166 167        :param callback:168            The function to call when the file has data available.169 170        :param int flags:171            The condition to monitor on the file (defaults to ``POLLIN``).172        """173        if isinstance(fd, int):174            fd = os.fdopen(fd)175        self._poller.register(fd, flags)176        self._queue_callbacks[fd.fileno()] = callback177        return fd178 179    def remove_watch_queue(self, handle: zmq.Socket) -> bool:180        """181        Remove a queue from background polling. Returns ``True`` if the queue182        was being monitored, ``False`` otherwise.183        """184        try:185            try:186                self._poller.unregister(handle)187            finally:188                self._queue_callbacks.pop(handle, None)189 190        except KeyError:191            return False192 193        return True194 195    def remove_watch_file(self, handle: io.TextIOWrapper) -> bool:196        """197        Remove a file from background polling. Returns ``True`` if the file was198        being monitored, ``False`` otherwise.199        """200        try:201            try:202                self._poller.unregister(handle)203            finally:204                self._queue_callbacks.pop(handle.fileno(), None)205 206        except KeyError:207            return False208 209        return True210 211    def enter_idle(self, callback: Callable[[], typing.Any]) -> int:212        """213        Add a *callback* to be executed when the event loop detects it is idle.214        Returns a handle that may be passed to :meth:`remove_enter_idle`.215        """216        self._idle_handle += 1217        self._idle_callbacks[self._idle_handle] = callback218        return self._idle_handle219 220    def remove_enter_idle(self, handle: int) -> bool:221        """222        Remove an idle callback. Returns ``True`` if *handle* was removed,223        ``False`` otherwise.224        """225        try:226            del self._idle_callbacks[handle]227        except KeyError:228            return False229 230        return True231 232    def _entering_idle(self) -> None:233        for callback in list(self._idle_callbacks.values()):234            callback()235 236    def run(self) -> None:237        """238        Start the event loop. Exit the loop when any callback raises an239        exception. If :exc:`ExitMainLoop` is raised, exit cleanly.240        """241        with contextlib.suppress(ExitMainLoop):242            while True:243                try:244                    self._loop()245                except zmq.error.ZMQError as exc:  # noqa: PERF203246                    if exc.errno != errno.EINTR:247                        raise248 249    def _loop(self) -> None:250        """251        A single iteration of the event loop.252        """253        state = "wait"  # default state not expecting any action254        if self._alarms or self._did_something:255            timeout = 0256            if self._alarms:257                state = "alarm"258                timeout = max(0.0, self._alarms[0][0] - time.time())259            if self._did_something and (not self._alarms or (self._alarms and timeout > 0)):260                state = "idle"261                timeout = 0262            ready = dict(self._poller.poll(timeout * 1000))263        else:264            ready = dict(self._poller.poll())265 266        if not ready:267            if state == "idle":268                self._entering_idle()269                self._did_something = False270            elif state == "alarm":271                _due, _tie_break, callback = heapq.heappop(self._alarms)272                callback()273                self._did_something = True274 275        for queue in ready:276            self._queue_callbacks[queue]()277            self._did_something = True278 
codekingpro/portable-devtools · Team Ai