codekingpro/portable-devtools
114k
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 