_eventloop.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428
  1. from __future__ import annotations
  2. import math
  3. import sys
  4. from abc import ABCMeta, abstractmethod
  5. from collections.abc import (
  6. AsyncIterator,
  7. Awaitable,
  8. Callable,
  9. Coroutine,
  10. Iterable,
  11. Mapping,
  12. Sequence,
  13. )
  14. from contextlib import AbstractContextManager
  15. from os import PathLike
  16. from signal import Signals
  17. from socket import AddressFamily, SocketKind, socket
  18. from typing import (
  19. IO,
  20. TYPE_CHECKING,
  21. Any,
  22. TypeAlias,
  23. TypeVar,
  24. overload,
  25. )
  26. if sys.version_info >= (3, 11):
  27. from typing import TypeVarTuple, Unpack
  28. else:
  29. from typing_extensions import TypeVarTuple, Unpack
  30. if TYPE_CHECKING:
  31. from _typeshed import FileDescriptorLike
  32. from .._core._synchronization import CapacityLimiter, Event, Lock, Semaphore
  33. from .._core._tasks import CancelScope
  34. from .._core._testing import TaskInfo
  35. from ._sockets import (
  36. ConnectedUDPSocket,
  37. ConnectedUNIXDatagramSocket,
  38. IPSockAddrType,
  39. SocketListener,
  40. SocketStream,
  41. UDPSocket,
  42. UNIXDatagramSocket,
  43. UNIXSocketStream,
  44. )
  45. from ._subprocesses import Process
  46. from ._tasks import TaskGroup
  47. from ._testing import TestRunner
  48. T_Retval = TypeVar("T_Retval")
  49. T_co = TypeVar("T_co", covariant=True)
  50. PosArgsT = TypeVarTuple("PosArgsT")
  51. StrOrBytesPath: TypeAlias = str | bytes | PathLike[str] | PathLike[bytes]
  52. class AsyncBackend(metaclass=ABCMeta):
  53. @classmethod
  54. @abstractmethod
  55. def run(
  56. cls,
  57. func: Callable[[Unpack[PosArgsT]], Awaitable[T_Retval]],
  58. args: tuple[Unpack[PosArgsT]],
  59. kwargs: dict[str, Any],
  60. options: dict[str, Any],
  61. ) -> T_Retval:
  62. """
  63. Run the given coroutine function in an asynchronous event loop.
  64. The current thread must not be already running an event loop.
  65. :param func: a coroutine function
  66. :param args: positional arguments to ``func``
  67. :param kwargs: positional arguments to ``func``
  68. :param options: keyword arguments to call the backend ``run()`` implementation
  69. with
  70. :return: the return value of the coroutine function
  71. """
  72. @classmethod
  73. @abstractmethod
  74. def current_token(cls) -> object:
  75. """
  76. Return an object that allows other threads to run code inside the event loop.
  77. :return: a token object, specific to the event loop running in the current
  78. thread
  79. """
  80. @classmethod
  81. @abstractmethod
  82. def current_time(cls) -> float:
  83. """
  84. Return the current value of the event loop's internal clock.
  85. :return: the clock value (seconds)
  86. """
  87. @classmethod
  88. @abstractmethod
  89. def cancelled_exception_class(cls) -> type[BaseException]:
  90. """Return the exception class that is raised in a task if it's cancelled."""
  91. @classmethod
  92. @abstractmethod
  93. async def checkpoint(cls) -> None:
  94. """
  95. Check if the task has been cancelled, and allow rescheduling of other tasks.
  96. This is effectively the same as running :meth:`checkpoint_if_cancelled` and then
  97. :meth:`cancel_shielded_checkpoint`.
  98. """
  99. @classmethod
  100. async def checkpoint_if_cancelled(cls) -> None:
  101. """
  102. Check if the current task group has been cancelled.
  103. This will check if the task has been cancelled, but will not allow other tasks
  104. to be scheduled if not.
  105. """
  106. if cls.current_effective_deadline() == -math.inf:
  107. await cls.checkpoint()
  108. @classmethod
  109. async def cancel_shielded_checkpoint(cls) -> None:
  110. """
  111. Allow the rescheduling of other tasks.
  112. This will give other tasks the opportunity to run, but without checking if the
  113. current task group has been cancelled, unlike with :meth:`checkpoint`.
  114. """
  115. with cls.create_cancel_scope(shield=True):
  116. await cls.sleep(0)
  117. @classmethod
  118. @abstractmethod
  119. async def sleep(cls, delay: float) -> None:
  120. """
  121. Pause the current task for the specified duration.
  122. :param delay: the duration, in seconds
  123. """
  124. @classmethod
  125. @abstractmethod
  126. def create_cancel_scope(
  127. cls, *, deadline: float = math.inf, shield: bool = False
  128. ) -> CancelScope:
  129. pass
  130. @classmethod
  131. @abstractmethod
  132. def current_effective_deadline(cls) -> float:
  133. """
  134. Return the nearest deadline among all the cancel scopes effective for the
  135. current task.
  136. :return:
  137. - a clock value from the event loop's internal clock
  138. - ``inf`` if there is no deadline in effect
  139. - ``-inf`` if the current scope has been cancelled
  140. :rtype: float
  141. """
  142. @classmethod
  143. @abstractmethod
  144. def create_task_group(cls) -> TaskGroup:
  145. pass
  146. @classmethod
  147. @abstractmethod
  148. def create_event(cls) -> Event:
  149. pass
  150. @classmethod
  151. @abstractmethod
  152. def create_lock(cls, *, fast_acquire: bool) -> Lock:
  153. pass
  154. @classmethod
  155. @abstractmethod
  156. def create_semaphore(
  157. cls,
  158. initial_value: int,
  159. *,
  160. max_value: int | None = None,
  161. fast_acquire: bool = False,
  162. ) -> Semaphore:
  163. pass
  164. @classmethod
  165. @abstractmethod
  166. def create_capacity_limiter(cls, total_tokens: float) -> CapacityLimiter:
  167. pass
  168. @classmethod
  169. @abstractmethod
  170. async def run_sync_in_worker_thread(
  171. cls,
  172. func: Callable[[Unpack[PosArgsT]], T_Retval],
  173. args: tuple[Unpack[PosArgsT]],
  174. abandon_on_cancel: bool = False,
  175. limiter: CapacityLimiter | None = None,
  176. ) -> T_Retval:
  177. pass
  178. @classmethod
  179. @abstractmethod
  180. def check_cancelled(cls) -> None:
  181. pass
  182. @classmethod
  183. @abstractmethod
  184. def run_async_from_thread(
  185. cls,
  186. func: Callable[[Unpack[PosArgsT]], Coroutine[Any, Any, T_co]],
  187. args: tuple[Unpack[PosArgsT]],
  188. token: object,
  189. ) -> T_co:
  190. pass
  191. @classmethod
  192. @abstractmethod
  193. def run_sync_from_thread(
  194. cls,
  195. func: Callable[[Unpack[PosArgsT]], T_Retval],
  196. args: tuple[Unpack[PosArgsT]],
  197. token: object,
  198. ) -> T_Retval:
  199. pass
  200. @classmethod
  201. @abstractmethod
  202. async def open_process(
  203. cls,
  204. command: StrOrBytesPath | Sequence[StrOrBytesPath],
  205. *,
  206. stdin: int | IO[Any] | None,
  207. stdout: int | IO[Any] | None,
  208. stderr: int | IO[Any] | None,
  209. cwd: StrOrBytesPath | None = None,
  210. env: Mapping[str, str] | None = None,
  211. startupinfo: Any = None,
  212. creationflags: int = 0,
  213. start_new_session: bool = False,
  214. pass_fds: Sequence[int] = (),
  215. user: str | int | None = None,
  216. group: str | int | None = None,
  217. extra_groups: Iterable[str | int] | None = None,
  218. umask: int = -1,
  219. **kwargs: Any,
  220. ) -> Process:
  221. pass
  222. @classmethod
  223. @abstractmethod
  224. def setup_process_pool_exit_at_shutdown(cls, workers: set[Process]) -> None:
  225. pass
  226. @classmethod
  227. @abstractmethod
  228. async def connect_tcp(
  229. cls, host: str, port: int, local_address: IPSockAddrType | None = None
  230. ) -> SocketStream:
  231. pass
  232. @classmethod
  233. @abstractmethod
  234. async def connect_unix(cls, path: str | bytes) -> UNIXSocketStream:
  235. pass
  236. @classmethod
  237. @abstractmethod
  238. def create_tcp_listener(cls, sock: socket) -> SocketListener:
  239. pass
  240. @classmethod
  241. @abstractmethod
  242. def create_unix_listener(cls, sock: socket) -> SocketListener:
  243. pass
  244. @classmethod
  245. @abstractmethod
  246. async def create_udp_socket(
  247. cls,
  248. family: AddressFamily,
  249. local_address: IPSockAddrType | None,
  250. remote_address: IPSockAddrType | None,
  251. reuse_port: bool,
  252. ) -> UDPSocket | ConnectedUDPSocket:
  253. pass
  254. @classmethod
  255. @overload
  256. async def create_unix_datagram_socket(
  257. cls, raw_socket: socket, remote_path: None
  258. ) -> UNIXDatagramSocket: ...
  259. @classmethod
  260. @overload
  261. async def create_unix_datagram_socket(
  262. cls, raw_socket: socket, remote_path: str | bytes
  263. ) -> ConnectedUNIXDatagramSocket: ...
  264. @classmethod
  265. @abstractmethod
  266. async def create_unix_datagram_socket(
  267. cls, raw_socket: socket, remote_path: str | bytes | None
  268. ) -> UNIXDatagramSocket | ConnectedUNIXDatagramSocket:
  269. pass
  270. @classmethod
  271. @abstractmethod
  272. async def getaddrinfo(
  273. cls,
  274. host: bytes | str | None,
  275. port: str | int | None,
  276. *,
  277. family: int | AddressFamily = 0,
  278. type: int | SocketKind = 0,
  279. proto: int = 0,
  280. flags: int = 0,
  281. ) -> Sequence[
  282. tuple[
  283. AddressFamily,
  284. SocketKind,
  285. int,
  286. str,
  287. tuple[str, int] | tuple[str, int, int, int] | tuple[int, bytes],
  288. ]
  289. ]:
  290. pass
  291. @classmethod
  292. @abstractmethod
  293. async def getnameinfo(
  294. cls, sockaddr: IPSockAddrType, flags: int = 0
  295. ) -> tuple[str, str]:
  296. pass
  297. @classmethod
  298. @abstractmethod
  299. async def wait_readable(cls, obj: FileDescriptorLike) -> None:
  300. pass
  301. @classmethod
  302. @abstractmethod
  303. async def wait_writable(cls, obj: FileDescriptorLike) -> None:
  304. pass
  305. @classmethod
  306. @abstractmethod
  307. def notify_closing(cls, obj: FileDescriptorLike) -> None:
  308. pass
  309. @classmethod
  310. @abstractmethod
  311. async def wrap_listener_socket(cls, sock: socket) -> SocketListener:
  312. pass
  313. @classmethod
  314. @abstractmethod
  315. async def wrap_stream_socket(cls, sock: socket) -> SocketStream:
  316. pass
  317. @classmethod
  318. @abstractmethod
  319. async def wrap_unix_stream_socket(cls, sock: socket) -> UNIXSocketStream:
  320. pass
  321. @classmethod
  322. @abstractmethod
  323. async def wrap_udp_socket(cls, sock: socket) -> UDPSocket:
  324. pass
  325. @classmethod
  326. @abstractmethod
  327. async def wrap_connected_udp_socket(cls, sock: socket) -> ConnectedUDPSocket:
  328. pass
  329. @classmethod
  330. @abstractmethod
  331. async def wrap_unix_datagram_socket(cls, sock: socket) -> UNIXDatagramSocket:
  332. pass
  333. @classmethod
  334. @abstractmethod
  335. async def wrap_connected_unix_datagram_socket(
  336. cls, sock: socket
  337. ) -> ConnectedUNIXDatagramSocket:
  338. pass
  339. @classmethod
  340. @abstractmethod
  341. def current_default_thread_limiter(cls) -> CapacityLimiter:
  342. pass
  343. @classmethod
  344. @abstractmethod
  345. def open_signal_receiver(
  346. cls, *signals: Signals
  347. ) -> AbstractContextManager[AsyncIterator[Signals]]:
  348. pass
  349. @classmethod
  350. @abstractmethod
  351. def get_current_task(cls) -> TaskInfo:
  352. pass
  353. @classmethod
  354. @abstractmethod
  355. def get_running_tasks(cls) -> Sequence[TaskInfo]:
  356. pass
  357. @classmethod
  358. @abstractmethod
  359. async def wait_all_tasks_blocked(cls) -> None:
  360. pass
  361. @classmethod
  362. @abstractmethod
  363. def create_test_runner(cls, options: dict[str, Any]) -> TestRunner:
  364. pass