Загрузить файлы в «venv/Lib/site-packages/anyio/_core»
This commit is contained in:
196
venv/Lib/site-packages/anyio/_core/_subprocesses.py
Normal file
196
venv/Lib/site-packages/anyio/_core/_subprocesses.py
Normal file
@@ -0,0 +1,196 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import AsyncIterable, Iterable, Mapping, Sequence
|
||||||
|
from io import BytesIO
|
||||||
|
from os import PathLike
|
||||||
|
from subprocess import PIPE, CalledProcessError, CompletedProcess
|
||||||
|
from typing import IO, Any, TypeAlias, cast
|
||||||
|
|
||||||
|
from ..abc import Process
|
||||||
|
from ._eventloop import get_async_backend
|
||||||
|
from ._tasks import create_task_group
|
||||||
|
|
||||||
|
StrOrBytesPath: TypeAlias = str | bytes | PathLike[str] | PathLike[bytes]
|
||||||
|
|
||||||
|
|
||||||
|
async def run_process(
|
||||||
|
command: StrOrBytesPath | Sequence[StrOrBytesPath],
|
||||||
|
*,
|
||||||
|
input: bytes | None = None,
|
||||||
|
stdin: int | IO[Any] | None = None,
|
||||||
|
stdout: int | IO[Any] | None = PIPE,
|
||||||
|
stderr: int | IO[Any] | None = PIPE,
|
||||||
|
check: bool = True,
|
||||||
|
cwd: StrOrBytesPath | None = None,
|
||||||
|
env: Mapping[str, str] | None = None,
|
||||||
|
startupinfo: Any = None,
|
||||||
|
creationflags: int = 0,
|
||||||
|
start_new_session: bool = False,
|
||||||
|
pass_fds: Sequence[int] = (),
|
||||||
|
user: str | int | None = None,
|
||||||
|
group: str | int | None = None,
|
||||||
|
extra_groups: Iterable[str | int] | None = None,
|
||||||
|
umask: int = -1,
|
||||||
|
) -> CompletedProcess[bytes]:
|
||||||
|
"""
|
||||||
|
Run an external command in a subprocess and wait until it completes.
|
||||||
|
|
||||||
|
.. seealso:: :func:`subprocess.run`
|
||||||
|
|
||||||
|
:param command: either a string to pass to the shell, or an iterable of strings
|
||||||
|
containing the executable name or path and its arguments
|
||||||
|
:param input: bytes passed to the standard input of the subprocess
|
||||||
|
:param stdin: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`,
|
||||||
|
a file-like object, or `None`; ``input`` overrides this
|
||||||
|
:param stdout: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`,
|
||||||
|
a file-like object, or `None`
|
||||||
|
:param stderr: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`,
|
||||||
|
:data:`subprocess.STDOUT`, a file-like object, or `None`
|
||||||
|
:param check: if ``True``, raise :exc:`~subprocess.CalledProcessError` if the
|
||||||
|
process terminates with a return code other than 0
|
||||||
|
:param cwd: If not ``None``, change the working directory to this before running the
|
||||||
|
command
|
||||||
|
:param env: if not ``None``, this mapping replaces the inherited environment
|
||||||
|
variables from the parent process
|
||||||
|
:param startupinfo: an instance of :class:`subprocess.STARTUPINFO` that can be used
|
||||||
|
to specify process startup parameters (Windows only)
|
||||||
|
:param creationflags: flags that can be used to control the creation of the
|
||||||
|
subprocess (see :class:`subprocess.Popen` for the specifics)
|
||||||
|
:param start_new_session: if ``true`` the setsid() system call will be made in the
|
||||||
|
child process prior to the execution of the subprocess. (POSIX only)
|
||||||
|
:param pass_fds: sequence of file descriptors to keep open between the parent and
|
||||||
|
child processes. (POSIX only)
|
||||||
|
:param user: effective user to run the process as (Python >= 3.9, POSIX only)
|
||||||
|
:param group: effective group to run the process as (Python >= 3.9, POSIX only)
|
||||||
|
:param extra_groups: supplementary groups to set in the subprocess (Python >= 3.9,
|
||||||
|
POSIX only)
|
||||||
|
:param umask: if not negative, this umask is applied in the child process before
|
||||||
|
running the given command (Python >= 3.9, POSIX only)
|
||||||
|
:return: an object representing the completed process
|
||||||
|
:raises ~subprocess.CalledProcessError: if ``check`` is ``True`` and the process
|
||||||
|
exits with a nonzero return code
|
||||||
|
|
||||||
|
"""
|
||||||
|
|
||||||
|
async def drain_stream(stream: AsyncIterable[bytes], index: int) -> None:
|
||||||
|
buffer = BytesIO()
|
||||||
|
async for chunk in stream:
|
||||||
|
buffer.write(chunk)
|
||||||
|
|
||||||
|
stream_contents[index] = buffer.getvalue()
|
||||||
|
|
||||||
|
if stdin is not None and input is not None:
|
||||||
|
raise ValueError("only one of stdin and input is allowed")
|
||||||
|
|
||||||
|
async with await open_process(
|
||||||
|
command,
|
||||||
|
stdin=PIPE if input else stdin,
|
||||||
|
stdout=stdout,
|
||||||
|
stderr=stderr,
|
||||||
|
cwd=cwd,
|
||||||
|
env=env,
|
||||||
|
startupinfo=startupinfo,
|
||||||
|
creationflags=creationflags,
|
||||||
|
start_new_session=start_new_session,
|
||||||
|
pass_fds=pass_fds,
|
||||||
|
user=user,
|
||||||
|
group=group,
|
||||||
|
extra_groups=extra_groups,
|
||||||
|
umask=umask,
|
||||||
|
) as process:
|
||||||
|
stream_contents: list[bytes | None] = [None, None]
|
||||||
|
async with create_task_group() as tg:
|
||||||
|
if process.stdout:
|
||||||
|
tg.start_soon(drain_stream, process.stdout, 0)
|
||||||
|
|
||||||
|
if process.stderr:
|
||||||
|
tg.start_soon(drain_stream, process.stderr, 1)
|
||||||
|
|
||||||
|
if process.stdin and input:
|
||||||
|
await process.stdin.send(input)
|
||||||
|
await process.stdin.aclose()
|
||||||
|
|
||||||
|
await process.wait()
|
||||||
|
|
||||||
|
output, errors = stream_contents
|
||||||
|
if check and process.returncode != 0:
|
||||||
|
raise CalledProcessError(cast(int, process.returncode), command, output, errors)
|
||||||
|
|
||||||
|
return CompletedProcess(command, cast(int, process.returncode), output, errors)
|
||||||
|
|
||||||
|
|
||||||
|
async def open_process(
|
||||||
|
command: StrOrBytesPath | Sequence[StrOrBytesPath],
|
||||||
|
*,
|
||||||
|
stdin: int | IO[Any] | None = PIPE,
|
||||||
|
stdout: int | IO[Any] | None = PIPE,
|
||||||
|
stderr: int | IO[Any] | None = PIPE,
|
||||||
|
cwd: StrOrBytesPath | None = None,
|
||||||
|
env: Mapping[str, str] | None = None,
|
||||||
|
startupinfo: Any = None,
|
||||||
|
creationflags: int = 0,
|
||||||
|
start_new_session: bool = False,
|
||||||
|
pass_fds: Sequence[int] = (),
|
||||||
|
user: str | int | None = None,
|
||||||
|
group: str | int | None = None,
|
||||||
|
extra_groups: Iterable[str | int] | None = None,
|
||||||
|
umask: int = -1,
|
||||||
|
) -> Process:
|
||||||
|
"""
|
||||||
|
Start an external command in a subprocess.
|
||||||
|
|
||||||
|
.. seealso:: :class:`subprocess.Popen`
|
||||||
|
|
||||||
|
:param command: either a string to pass to the shell, or an iterable of strings
|
||||||
|
containing the executable name or path and its arguments
|
||||||
|
:param stdin: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`, a
|
||||||
|
file-like object, or ``None``
|
||||||
|
:param stdout: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`,
|
||||||
|
a file-like object, or ``None``
|
||||||
|
:param stderr: one of :data:`subprocess.PIPE`, :data:`subprocess.DEVNULL`,
|
||||||
|
:data:`subprocess.STDOUT`, a file-like object, or ``None``
|
||||||
|
:param cwd: If not ``None``, the working directory is changed before executing
|
||||||
|
:param env: If env is not ``None``, it must be a mapping that defines the
|
||||||
|
environment variables for the new process
|
||||||
|
:param creationflags: flags that can be used to control the creation of the
|
||||||
|
subprocess (see :class:`subprocess.Popen` for the specifics)
|
||||||
|
:param startupinfo: an instance of :class:`subprocess.STARTUPINFO` that can be used
|
||||||
|
to specify process startup parameters (Windows only)
|
||||||
|
:param start_new_session: if ``true`` the setsid() system call will be made in the
|
||||||
|
child process prior to the execution of the subprocess. (POSIX only)
|
||||||
|
:param pass_fds: sequence of file descriptors to keep open between the parent and
|
||||||
|
child processes. (POSIX only)
|
||||||
|
:param user: effective user to run the process as (POSIX only)
|
||||||
|
:param group: effective group to run the process as (POSIX only)
|
||||||
|
:param extra_groups: supplementary groups to set in the subprocess (POSIX only)
|
||||||
|
:param umask: if not negative, this umask is applied in the child process before
|
||||||
|
running the given command (POSIX only)
|
||||||
|
:return: an asynchronous process object
|
||||||
|
|
||||||
|
"""
|
||||||
|
kwargs: dict[str, Any] = {}
|
||||||
|
if user is not None:
|
||||||
|
kwargs["user"] = user
|
||||||
|
|
||||||
|
if group is not None:
|
||||||
|
kwargs["group"] = group
|
||||||
|
|
||||||
|
if extra_groups is not None:
|
||||||
|
kwargs["extra_groups"] = group
|
||||||
|
|
||||||
|
if umask >= 0:
|
||||||
|
kwargs["umask"] = umask
|
||||||
|
|
||||||
|
return await get_async_backend().open_process(
|
||||||
|
command,
|
||||||
|
stdin=stdin,
|
||||||
|
stdout=stdout,
|
||||||
|
stderr=stderr,
|
||||||
|
cwd=cwd,
|
||||||
|
env=env,
|
||||||
|
startupinfo=startupinfo,
|
||||||
|
creationflags=creationflags,
|
||||||
|
start_new_session=start_new_session,
|
||||||
|
pass_fds=pass_fds,
|
||||||
|
**kwargs,
|
||||||
|
)
|
||||||
772
venv/Lib/site-packages/anyio/_core/_synchronization.py
Normal file
772
venv/Lib/site-packages/anyio/_core/_synchronization.py
Normal file
@@ -0,0 +1,772 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import math
|
||||||
|
from collections import deque
|
||||||
|
from collections.abc import Callable
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from types import TracebackType
|
||||||
|
from typing import TypeVar
|
||||||
|
|
||||||
|
from ..lowlevel import checkpoint_if_cancelled
|
||||||
|
from ._eventloop import get_async_backend
|
||||||
|
from ._exceptions import BusyResourceError, NoEventLoopError
|
||||||
|
from ._tasks import CancelScope
|
||||||
|
from ._testing import TaskInfo, get_current_task
|
||||||
|
|
||||||
|
T = TypeVar("T")
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class EventStatistics:
|
||||||
|
"""
|
||||||
|
:ivar int tasks_waiting: number of tasks waiting on :meth:`~.Event.wait`
|
||||||
|
"""
|
||||||
|
|
||||||
|
tasks_waiting: int
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class CapacityLimiterStatistics:
|
||||||
|
"""
|
||||||
|
:ivar int borrowed_tokens: number of tokens currently borrowed by tasks
|
||||||
|
:ivar float total_tokens: total number of available tokens
|
||||||
|
:ivar tuple borrowers: tasks or other objects currently holding tokens borrowed from
|
||||||
|
this limiter
|
||||||
|
:ivar int tasks_waiting: number of tasks waiting on
|
||||||
|
:meth:`~.CapacityLimiter.acquire` or
|
||||||
|
:meth:`~.CapacityLimiter.acquire_on_behalf_of`
|
||||||
|
"""
|
||||||
|
|
||||||
|
borrowed_tokens: int
|
||||||
|
total_tokens: float
|
||||||
|
borrowers: tuple[object, ...]
|
||||||
|
tasks_waiting: int
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class LockStatistics:
|
||||||
|
"""
|
||||||
|
:ivar bool locked: flag indicating if this lock is locked or not
|
||||||
|
:ivar ~anyio.TaskInfo owner: task currently holding the lock (or ``None`` if the
|
||||||
|
lock is not held by any task)
|
||||||
|
:ivar int tasks_waiting: number of tasks waiting on :meth:`~.Lock.acquire`
|
||||||
|
"""
|
||||||
|
|
||||||
|
locked: bool
|
||||||
|
owner: TaskInfo | None
|
||||||
|
tasks_waiting: int
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class ConditionStatistics:
|
||||||
|
"""
|
||||||
|
:ivar int tasks_waiting: number of tasks blocked on :meth:`~.Condition.wait`
|
||||||
|
:ivar ~anyio.LockStatistics lock_statistics: statistics of the underlying
|
||||||
|
:class:`~.Lock`
|
||||||
|
"""
|
||||||
|
|
||||||
|
tasks_waiting: int
|
||||||
|
lock_statistics: LockStatistics
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class SemaphoreStatistics:
|
||||||
|
"""
|
||||||
|
:ivar int tasks_waiting: number of tasks waiting on :meth:`~.Semaphore.acquire`
|
||||||
|
|
||||||
|
"""
|
||||||
|
|
||||||
|
tasks_waiting: int
|
||||||
|
|
||||||
|
|
||||||
|
class Event:
|
||||||
|
__slots__ = ("__weakref__",)
|
||||||
|
|
||||||
|
def __new__(cls) -> Event:
|
||||||
|
try:
|
||||||
|
return get_async_backend().create_event()
|
||||||
|
except NoEventLoopError:
|
||||||
|
return EventAdapter()
|
||||||
|
|
||||||
|
def set(self) -> None:
|
||||||
|
"""Set the flag, notifying all listeners."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def is_set(self) -> bool:
|
||||||
|
"""Return ``True`` if the flag is set, ``False`` if not."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
async def wait(self) -> None:
|
||||||
|
"""
|
||||||
|
Wait until the flag has been set.
|
||||||
|
|
||||||
|
If the flag has already been set when this method is called, it returns
|
||||||
|
immediately.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def statistics(self) -> EventStatistics:
|
||||||
|
"""Return statistics about the current state of this event."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
|
||||||
|
class EventAdapter(Event):
|
||||||
|
__slots__ = "_internal_event", "_is_set"
|
||||||
|
|
||||||
|
def __new__(cls) -> EventAdapter:
|
||||||
|
return object.__new__(cls)
|
||||||
|
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self._internal_event: Event | None = None
|
||||||
|
self._is_set = False
|
||||||
|
|
||||||
|
@property
|
||||||
|
def _event(self) -> Event:
|
||||||
|
if self._internal_event is None:
|
||||||
|
self._internal_event = get_async_backend().create_event()
|
||||||
|
if self._is_set:
|
||||||
|
self._internal_event.set()
|
||||||
|
|
||||||
|
return self._internal_event
|
||||||
|
|
||||||
|
def set(self) -> None:
|
||||||
|
if self._internal_event is None:
|
||||||
|
self._is_set = True
|
||||||
|
else:
|
||||||
|
self._event.set()
|
||||||
|
|
||||||
|
def is_set(self) -> bool:
|
||||||
|
if self._internal_event is None:
|
||||||
|
return self._is_set
|
||||||
|
|
||||||
|
return self._internal_event.is_set()
|
||||||
|
|
||||||
|
async def wait(self) -> None:
|
||||||
|
await self._event.wait()
|
||||||
|
|
||||||
|
def statistics(self) -> EventStatistics:
|
||||||
|
if self._internal_event is None:
|
||||||
|
return EventStatistics(tasks_waiting=0)
|
||||||
|
|
||||||
|
return self._internal_event.statistics()
|
||||||
|
|
||||||
|
|
||||||
|
class Lock:
|
||||||
|
__slots__ = ("__weakref__",)
|
||||||
|
|
||||||
|
def __new__(cls, *, fast_acquire: bool = False) -> Lock:
|
||||||
|
try:
|
||||||
|
return get_async_backend().create_lock(fast_acquire=fast_acquire)
|
||||||
|
except NoEventLoopError:
|
||||||
|
return LockAdapter(fast_acquire=fast_acquire)
|
||||||
|
|
||||||
|
async def __aenter__(self) -> None:
|
||||||
|
await self.acquire()
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
self.release()
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
"""Acquire the lock."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire the lock, without blocking.
|
||||||
|
|
||||||
|
:raises ~anyio.WouldBlock: if the operation would block
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
"""Release the lock."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def locked(self) -> bool:
|
||||||
|
"""Return True if the lock is currently held."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def statistics(self) -> LockStatistics:
|
||||||
|
"""
|
||||||
|
Return statistics about the current state of this lock.
|
||||||
|
|
||||||
|
.. versionadded:: 3.0
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
|
||||||
|
class LockAdapter(Lock):
|
||||||
|
__slots__ = "_internal_lock", "_fast_acquire"
|
||||||
|
|
||||||
|
def __new__(cls, *, fast_acquire: bool = False) -> LockAdapter:
|
||||||
|
return object.__new__(cls)
|
||||||
|
|
||||||
|
def __init__(self, *, fast_acquire: bool = False):
|
||||||
|
self._internal_lock: Lock | None = None
|
||||||
|
self._fast_acquire = fast_acquire
|
||||||
|
|
||||||
|
@property
|
||||||
|
def _lock(self) -> Lock:
|
||||||
|
if self._internal_lock is None:
|
||||||
|
self._internal_lock = get_async_backend().create_lock(
|
||||||
|
fast_acquire=self._fast_acquire
|
||||||
|
)
|
||||||
|
|
||||||
|
return self._internal_lock
|
||||||
|
|
||||||
|
async def __aenter__(self) -> None:
|
||||||
|
await self._lock.acquire()
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
if self._internal_lock is not None:
|
||||||
|
self._internal_lock.release()
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
"""Acquire the lock."""
|
||||||
|
await self._lock.acquire()
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire the lock, without blocking.
|
||||||
|
|
||||||
|
:raises ~anyio.WouldBlock: if the operation would block
|
||||||
|
|
||||||
|
"""
|
||||||
|
self._lock.acquire_nowait()
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
"""Release the lock."""
|
||||||
|
self._lock.release()
|
||||||
|
|
||||||
|
def locked(self) -> bool:
|
||||||
|
"""Return True if the lock is currently held."""
|
||||||
|
return self._lock.locked()
|
||||||
|
|
||||||
|
def statistics(self) -> LockStatistics:
|
||||||
|
"""
|
||||||
|
Return statistics about the current state of this lock.
|
||||||
|
|
||||||
|
.. versionadded:: 3.0
|
||||||
|
|
||||||
|
"""
|
||||||
|
if self._internal_lock is None:
|
||||||
|
return LockStatistics(False, None, 0)
|
||||||
|
|
||||||
|
return self._internal_lock.statistics()
|
||||||
|
|
||||||
|
|
||||||
|
class Condition:
|
||||||
|
__slots__ = "__weakref__", "_owner_task", "_lock", "_waiters"
|
||||||
|
|
||||||
|
def __init__(self, lock: Lock | None = None):
|
||||||
|
self._owner_task: TaskInfo | None = None
|
||||||
|
self._lock = lock or Lock()
|
||||||
|
self._waiters: deque[Event] = deque()
|
||||||
|
|
||||||
|
async def __aenter__(self) -> None:
|
||||||
|
await self.acquire()
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
self.release()
|
||||||
|
|
||||||
|
def _check_acquired(self) -> None:
|
||||||
|
if self._owner_task != get_current_task():
|
||||||
|
raise RuntimeError("The current task is not holding the underlying lock")
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
"""Acquire the underlying lock."""
|
||||||
|
await self._lock.acquire()
|
||||||
|
self._owner_task = get_current_task()
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire the underlying lock, without blocking.
|
||||||
|
|
||||||
|
:raises ~anyio.WouldBlock: if the operation would block
|
||||||
|
|
||||||
|
"""
|
||||||
|
self._lock.acquire_nowait()
|
||||||
|
self._owner_task = get_current_task()
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
"""Release the underlying lock."""
|
||||||
|
self._lock.release()
|
||||||
|
|
||||||
|
def locked(self) -> bool:
|
||||||
|
"""Return True if the lock is set."""
|
||||||
|
return self._lock.locked()
|
||||||
|
|
||||||
|
def notify(self, n: int = 1) -> None:
|
||||||
|
"""Notify exactly n listeners."""
|
||||||
|
self._check_acquired()
|
||||||
|
for _ in range(n):
|
||||||
|
try:
|
||||||
|
event = self._waiters.popleft()
|
||||||
|
except IndexError:
|
||||||
|
break
|
||||||
|
|
||||||
|
event.set()
|
||||||
|
|
||||||
|
def notify_all(self) -> None:
|
||||||
|
"""Notify all the listeners."""
|
||||||
|
self._check_acquired()
|
||||||
|
for event in self._waiters:
|
||||||
|
event.set()
|
||||||
|
|
||||||
|
self._waiters.clear()
|
||||||
|
|
||||||
|
async def wait(self) -> None:
|
||||||
|
"""Wait for a notification."""
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
self._check_acquired()
|
||||||
|
event = Event()
|
||||||
|
self._waiters.append(event)
|
||||||
|
self.release()
|
||||||
|
try:
|
||||||
|
await event.wait()
|
||||||
|
except BaseException:
|
||||||
|
if not event.is_set():
|
||||||
|
self._waiters.remove(event)
|
||||||
|
elif self._waiters:
|
||||||
|
# This task was notified by could not act on it, so pass
|
||||||
|
# it on to the next task
|
||||||
|
self._waiters.popleft().set()
|
||||||
|
|
||||||
|
raise
|
||||||
|
finally:
|
||||||
|
with CancelScope(shield=True):
|
||||||
|
await self.acquire()
|
||||||
|
|
||||||
|
async def wait_for(self, predicate: Callable[[], T]) -> T:
|
||||||
|
"""
|
||||||
|
Wait until a predicate becomes true.
|
||||||
|
|
||||||
|
:param predicate: a callable that returns a truthy value when the condition is
|
||||||
|
met
|
||||||
|
:return: the result of the predicate
|
||||||
|
|
||||||
|
.. versionadded:: 4.11.0
|
||||||
|
|
||||||
|
"""
|
||||||
|
while not (result := predicate()):
|
||||||
|
await self.wait()
|
||||||
|
|
||||||
|
return result
|
||||||
|
|
||||||
|
def statistics(self) -> ConditionStatistics:
|
||||||
|
"""
|
||||||
|
Return statistics about the current state of this condition.
|
||||||
|
|
||||||
|
.. versionadded:: 3.0
|
||||||
|
"""
|
||||||
|
return ConditionStatistics(len(self._waiters), self._lock.statistics())
|
||||||
|
|
||||||
|
|
||||||
|
class Semaphore:
|
||||||
|
__slots__ = "__weakref__", "_fast_acquire"
|
||||||
|
|
||||||
|
def __new__(
|
||||||
|
cls,
|
||||||
|
initial_value: int,
|
||||||
|
*,
|
||||||
|
max_value: int | None = None,
|
||||||
|
fast_acquire: bool = False,
|
||||||
|
) -> Semaphore:
|
||||||
|
try:
|
||||||
|
return get_async_backend().create_semaphore(
|
||||||
|
initial_value, max_value=max_value, fast_acquire=fast_acquire
|
||||||
|
)
|
||||||
|
except NoEventLoopError:
|
||||||
|
return SemaphoreAdapter(initial_value, max_value=max_value)
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
initial_value: int,
|
||||||
|
*,
|
||||||
|
max_value: int | None = None,
|
||||||
|
fast_acquire: bool = False,
|
||||||
|
):
|
||||||
|
if not isinstance(initial_value, int):
|
||||||
|
raise TypeError("initial_value must be an integer")
|
||||||
|
if initial_value < 0:
|
||||||
|
raise ValueError("initial_value must be >= 0")
|
||||||
|
if max_value is not None:
|
||||||
|
if not isinstance(max_value, int):
|
||||||
|
raise TypeError("max_value must be an integer or None")
|
||||||
|
if max_value < initial_value:
|
||||||
|
raise ValueError(
|
||||||
|
"max_value must be equal to or higher than initial_value"
|
||||||
|
)
|
||||||
|
|
||||||
|
self._fast_acquire = fast_acquire
|
||||||
|
|
||||||
|
async def __aenter__(self) -> Semaphore:
|
||||||
|
await self.acquire()
|
||||||
|
return self
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
self.release()
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
"""Decrement the semaphore value, blocking if necessary."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire the underlying lock, without blocking.
|
||||||
|
|
||||||
|
:raises ~anyio.WouldBlock: if the operation would block
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
"""Increment the semaphore value."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def value(self) -> int:
|
||||||
|
"""The current value of the semaphore."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def max_value(self) -> int | None:
|
||||||
|
"""The maximum value of the semaphore."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def statistics(self) -> SemaphoreStatistics:
|
||||||
|
"""
|
||||||
|
Return statistics about the current state of this semaphore.
|
||||||
|
|
||||||
|
.. versionadded:: 3.0
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
|
||||||
|
class SemaphoreAdapter(Semaphore):
|
||||||
|
__slots__ = "_internal_semaphore", "_initial_value", "_max_value"
|
||||||
|
|
||||||
|
def __new__(
|
||||||
|
cls,
|
||||||
|
initial_value: int,
|
||||||
|
*,
|
||||||
|
max_value: int | None = None,
|
||||||
|
fast_acquire: bool = False,
|
||||||
|
) -> SemaphoreAdapter:
|
||||||
|
return object.__new__(cls)
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
initial_value: int,
|
||||||
|
*,
|
||||||
|
max_value: int | None = None,
|
||||||
|
fast_acquire: bool = False,
|
||||||
|
) -> None:
|
||||||
|
super().__init__(initial_value, max_value=max_value, fast_acquire=fast_acquire)
|
||||||
|
self._internal_semaphore: Semaphore | None = None
|
||||||
|
self._initial_value = initial_value
|
||||||
|
self._max_value = max_value
|
||||||
|
|
||||||
|
@property
|
||||||
|
def _semaphore(self) -> Semaphore:
|
||||||
|
if self._internal_semaphore is None:
|
||||||
|
self._internal_semaphore = get_async_backend().create_semaphore(
|
||||||
|
self._initial_value, max_value=self._max_value
|
||||||
|
)
|
||||||
|
|
||||||
|
return self._internal_semaphore
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
await self._semaphore.acquire()
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
self._semaphore.acquire_nowait()
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
self._semaphore.release()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def value(self) -> int:
|
||||||
|
if self._internal_semaphore is None:
|
||||||
|
return self._initial_value
|
||||||
|
|
||||||
|
return self._semaphore.value
|
||||||
|
|
||||||
|
@property
|
||||||
|
def max_value(self) -> int | None:
|
||||||
|
return self._max_value
|
||||||
|
|
||||||
|
def statistics(self) -> SemaphoreStatistics:
|
||||||
|
if self._internal_semaphore is None:
|
||||||
|
return SemaphoreStatistics(tasks_waiting=0)
|
||||||
|
|
||||||
|
return self._semaphore.statistics()
|
||||||
|
|
||||||
|
|
||||||
|
class CapacityLimiter:
|
||||||
|
__slots__ = ("__weakref__",)
|
||||||
|
|
||||||
|
def __new__(cls, total_tokens: float) -> CapacityLimiter:
|
||||||
|
try:
|
||||||
|
return get_async_backend().create_capacity_limiter(total_tokens)
|
||||||
|
except NoEventLoopError:
|
||||||
|
return CapacityLimiterAdapter(total_tokens)
|
||||||
|
|
||||||
|
async def __aenter__(self) -> None:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def total_tokens(self) -> float:
|
||||||
|
"""
|
||||||
|
The total number of tokens available for borrowing.
|
||||||
|
|
||||||
|
This is a read-write property. If the total number of tokens is increased, the
|
||||||
|
proportionate number of tasks waiting on this limiter will be granted their
|
||||||
|
tokens.
|
||||||
|
|
||||||
|
.. versionchanged:: 3.0
|
||||||
|
The property is now writable.
|
||||||
|
.. versionchanged:: 4.12
|
||||||
|
The value can now be set to 0.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@total_tokens.setter
|
||||||
|
def total_tokens(self, value: float) -> None:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def borrowed_tokens(self) -> int:
|
||||||
|
"""The number of tokens that have currently been borrowed."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def available_tokens(self) -> float:
|
||||||
|
"""The number of tokens currently available to be borrowed"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire a token for the current task without waiting for one to become
|
||||||
|
available.
|
||||||
|
|
||||||
|
:raises ~anyio.WouldBlock: if there are no tokens available for borrowing
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def acquire_on_behalf_of_nowait(self, borrower: object) -> None:
|
||||||
|
"""
|
||||||
|
Acquire a token without waiting for one to become available.
|
||||||
|
|
||||||
|
:param borrower: the entity borrowing a token
|
||||||
|
:raises ~anyio.WouldBlock: if there are no tokens available for borrowing
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
"""
|
||||||
|
Acquire a token for the current task, waiting if necessary for one to become
|
||||||
|
available.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
async def acquire_on_behalf_of(self, borrower: object) -> None:
|
||||||
|
"""
|
||||||
|
Acquire a token, waiting if necessary for one to become available.
|
||||||
|
|
||||||
|
:param borrower: the entity borrowing a token
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
"""
|
||||||
|
Release the token held by the current task.
|
||||||
|
|
||||||
|
:raises RuntimeError: if the current task has not borrowed a token from this
|
||||||
|
limiter.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def release_on_behalf_of(self, borrower: object) -> None:
|
||||||
|
"""
|
||||||
|
Release the token held by the given borrower.
|
||||||
|
|
||||||
|
:raises RuntimeError: if the borrower has not borrowed a token from this
|
||||||
|
limiter.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def statistics(self) -> CapacityLimiterStatistics:
|
||||||
|
"""
|
||||||
|
Return statistics about the current state of this limiter.
|
||||||
|
|
||||||
|
.. versionadded:: 3.0
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
|
||||||
|
class CapacityLimiterAdapter(CapacityLimiter):
|
||||||
|
__slots__ = "_internal_limiter", "_total_tokens"
|
||||||
|
|
||||||
|
def __new__(cls, total_tokens: float) -> CapacityLimiterAdapter:
|
||||||
|
return object.__new__(cls)
|
||||||
|
|
||||||
|
def __init__(self, total_tokens: float) -> None:
|
||||||
|
self._internal_limiter: CapacityLimiter | None = None
|
||||||
|
self.total_tokens = total_tokens
|
||||||
|
|
||||||
|
@property
|
||||||
|
def _limiter(self) -> CapacityLimiter:
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
self._internal_limiter = get_async_backend().create_capacity_limiter(
|
||||||
|
self._total_tokens
|
||||||
|
)
|
||||||
|
|
||||||
|
return self._internal_limiter
|
||||||
|
|
||||||
|
async def __aenter__(self) -> None:
|
||||||
|
await self._limiter.__aenter__()
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
return await self._limiter.__aexit__(exc_type, exc_val, exc_tb)
|
||||||
|
|
||||||
|
@property
|
||||||
|
def total_tokens(self) -> float:
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
return self._total_tokens
|
||||||
|
|
||||||
|
return self._internal_limiter.total_tokens
|
||||||
|
|
||||||
|
@total_tokens.setter
|
||||||
|
def total_tokens(self, value: float) -> None:
|
||||||
|
if not isinstance(value, int) and value is not math.inf:
|
||||||
|
raise TypeError("total_tokens must be an int or math.inf")
|
||||||
|
elif value < 0:
|
||||||
|
raise ValueError("total_tokens must be >= 0")
|
||||||
|
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
self._total_tokens = value
|
||||||
|
return
|
||||||
|
|
||||||
|
self._limiter.total_tokens = value
|
||||||
|
|
||||||
|
@property
|
||||||
|
def borrowed_tokens(self) -> int:
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
return 0
|
||||||
|
|
||||||
|
return self._internal_limiter.borrowed_tokens
|
||||||
|
|
||||||
|
@property
|
||||||
|
def available_tokens(self) -> float:
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
return self._total_tokens
|
||||||
|
|
||||||
|
return self._internal_limiter.available_tokens
|
||||||
|
|
||||||
|
def acquire_nowait(self) -> None:
|
||||||
|
self._limiter.acquire_nowait()
|
||||||
|
|
||||||
|
def acquire_on_behalf_of_nowait(self, borrower: object) -> None:
|
||||||
|
self._limiter.acquire_on_behalf_of_nowait(borrower)
|
||||||
|
|
||||||
|
async def acquire(self) -> None:
|
||||||
|
await self._limiter.acquire()
|
||||||
|
|
||||||
|
async def acquire_on_behalf_of(self, borrower: object) -> None:
|
||||||
|
await self._limiter.acquire_on_behalf_of(borrower)
|
||||||
|
|
||||||
|
def release(self) -> None:
|
||||||
|
self._limiter.release()
|
||||||
|
|
||||||
|
def release_on_behalf_of(self, borrower: object) -> None:
|
||||||
|
self._limiter.release_on_behalf_of(borrower)
|
||||||
|
|
||||||
|
def statistics(self) -> CapacityLimiterStatistics:
|
||||||
|
if self._internal_limiter is None:
|
||||||
|
return CapacityLimiterStatistics(
|
||||||
|
borrowed_tokens=0,
|
||||||
|
total_tokens=self.total_tokens,
|
||||||
|
borrowers=(),
|
||||||
|
tasks_waiting=0,
|
||||||
|
)
|
||||||
|
|
||||||
|
return self._internal_limiter.statistics()
|
||||||
|
|
||||||
|
|
||||||
|
class ResourceGuard:
|
||||||
|
"""
|
||||||
|
A context manager for ensuring that a resource is only used by a single task at a
|
||||||
|
time.
|
||||||
|
|
||||||
|
Entering this context manager while the previous has not exited it yet will trigger
|
||||||
|
:exc:`BusyResourceError`.
|
||||||
|
|
||||||
|
:param action: the action to guard against (visible in the :exc:`BusyResourceError`
|
||||||
|
when triggered, e.g. "Another task is already {action} this resource")
|
||||||
|
|
||||||
|
.. versionadded:: 4.1
|
||||||
|
"""
|
||||||
|
|
||||||
|
__slots__ = "__weakref__", "action", "_guarded"
|
||||||
|
|
||||||
|
def __init__(self, action: str = "using"):
|
||||||
|
self.action: str = action
|
||||||
|
self._guarded = False
|
||||||
|
|
||||||
|
def __enter__(self) -> None:
|
||||||
|
if self._guarded:
|
||||||
|
raise BusyResourceError(self.action)
|
||||||
|
|
||||||
|
self._guarded = True
|
||||||
|
|
||||||
|
def __exit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
self._guarded = False
|
||||||
415
venv/Lib/site-packages/anyio/_core/_tasks.py
Normal file
415
venv/Lib/site-packages/anyio/_core/_tasks.py
Normal file
@@ -0,0 +1,415 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import math
|
||||||
|
import sys
|
||||||
|
from collections.abc import (
|
||||||
|
Coroutine,
|
||||||
|
Generator,
|
||||||
|
)
|
||||||
|
from contextlib import (
|
||||||
|
contextmanager,
|
||||||
|
)
|
||||||
|
from contextvars import ContextVar
|
||||||
|
from enum import Enum, auto
|
||||||
|
from inspect import iscoroutine
|
||||||
|
from types import TracebackType
|
||||||
|
from typing import Any, Generic, final
|
||||||
|
|
||||||
|
from ..abc import TaskGroup, TaskStatus
|
||||||
|
from ._eventloop import get_async_backend, get_cancelled_exc_class
|
||||||
|
from ._exceptions import TaskCancelled, TaskFailed, TaskNotFinished
|
||||||
|
|
||||||
|
if sys.version_info >= (3, 13):
|
||||||
|
from typing import TypeVar
|
||||||
|
else:
|
||||||
|
from typing_extensions import TypeVar
|
||||||
|
|
||||||
|
if sys.version_info >= (3, 11):
|
||||||
|
from typing import Never, TypeVarTuple
|
||||||
|
else:
|
||||||
|
from typing_extensions import Never, TypeVarTuple
|
||||||
|
|
||||||
|
T = TypeVar("T")
|
||||||
|
T_co = TypeVar("T_co", covariant=True)
|
||||||
|
T_startval = TypeVar("T_startval", covariant=True, default=Never)
|
||||||
|
PosArgsT = TypeVarTuple("PosArgsT")
|
||||||
|
|
||||||
|
_current_task_handle: ContextVar[TaskHandle] = ContextVar("_current_task_handle")
|
||||||
|
|
||||||
|
|
||||||
|
class _IgnoredTaskStatus(TaskStatus[object]):
|
||||||
|
def started(self, value: object = None) -> None:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
TASK_STATUS_IGNORED = _IgnoredTaskStatus()
|
||||||
|
|
||||||
|
|
||||||
|
class CancelScope:
|
||||||
|
"""
|
||||||
|
Wraps a unit of work that can be made separately cancellable.
|
||||||
|
|
||||||
|
:param deadline: The time (clock value) when this scope is cancelled automatically
|
||||||
|
:param shield: ``True`` to shield the cancel scope from external cancellation
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
"""
|
||||||
|
|
||||||
|
__slots__ = ("__weakref__",)
|
||||||
|
|
||||||
|
def __new__(
|
||||||
|
cls, *, deadline: float = math.inf, shield: bool = False
|
||||||
|
) -> CancelScope:
|
||||||
|
return get_async_backend().create_cancel_scope(shield=shield, deadline=deadline)
|
||||||
|
|
||||||
|
def cancel(self, reason: str | None = None) -> None:
|
||||||
|
"""
|
||||||
|
Cancel this scope immediately.
|
||||||
|
|
||||||
|
:param reason: a message describing the reason for the cancellation
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def deadline(self) -> float:
|
||||||
|
"""
|
||||||
|
The time (clock value) when this scope is cancelled automatically.
|
||||||
|
|
||||||
|
Will be ``float('inf')`` if no timeout has been set.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@deadline.setter
|
||||||
|
def deadline(self, value: float) -> None:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def cancel_called(self) -> bool:
|
||||||
|
"""``True`` if :meth:`cancel` has been called."""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def cancelled_caught(self) -> bool:
|
||||||
|
"""
|
||||||
|
``True`` if this scope suppressed a cancellation exception it itself raised.
|
||||||
|
|
||||||
|
This is typically used to check if any work was interrupted, or to see if the
|
||||||
|
scope was cancelled due to its deadline being reached. The value will, however,
|
||||||
|
only be ``True`` if the cancellation was triggered by the scope itself (and not
|
||||||
|
an outer scope).
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@property
|
||||||
|
def shield(self) -> bool:
|
||||||
|
"""
|
||||||
|
``True`` if this scope is shielded from external cancellation.
|
||||||
|
|
||||||
|
While a scope is shielded, it will not receive cancellations from outside.
|
||||||
|
|
||||||
|
"""
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
@shield.setter
|
||||||
|
def shield(self, value: bool) -> None:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def __enter__(self) -> CancelScope:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
def __exit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_val: BaseException | None,
|
||||||
|
exc_tb: TracebackType | None,
|
||||||
|
) -> bool:
|
||||||
|
raise NotImplementedError
|
||||||
|
|
||||||
|
|
||||||
|
@contextmanager
|
||||||
|
def fail_after(
|
||||||
|
delay: float | None, shield: bool = False
|
||||||
|
) -> Generator[CancelScope, None, None]:
|
||||||
|
"""
|
||||||
|
Create a context manager which raises a :class:`TimeoutError` if does not finish in
|
||||||
|
time.
|
||||||
|
|
||||||
|
:param delay: maximum allowed time (in seconds) before raising the exception, or
|
||||||
|
``None`` to disable the timeout
|
||||||
|
:param shield: ``True`` to shield the cancel scope from external cancellation
|
||||||
|
:return: a context manager that yields a cancel scope
|
||||||
|
:rtype: :class:`~typing.ContextManager`\\[:class:`~anyio.CancelScope`\\]
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
current_time = get_async_backend().current_time
|
||||||
|
deadline = (current_time() + delay) if delay is not None else math.inf
|
||||||
|
with get_async_backend().create_cancel_scope(
|
||||||
|
deadline=deadline, shield=shield
|
||||||
|
) as cancel_scope:
|
||||||
|
yield cancel_scope
|
||||||
|
|
||||||
|
if cancel_scope.cancelled_caught and current_time() >= cancel_scope.deadline:
|
||||||
|
raise TimeoutError
|
||||||
|
|
||||||
|
|
||||||
|
def move_on_after(delay: float | None, shield: bool = False) -> CancelScope:
|
||||||
|
"""
|
||||||
|
Create a cancel scope with a deadline that expires after the given delay.
|
||||||
|
|
||||||
|
:param delay: maximum allowed time (in seconds) before exiting the context block, or
|
||||||
|
``None`` to disable the timeout
|
||||||
|
:param shield: ``True`` to shield the cancel scope from external cancellation
|
||||||
|
:return: a cancel scope
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
deadline = (
|
||||||
|
(get_async_backend().current_time() + delay) if delay is not None else math.inf
|
||||||
|
)
|
||||||
|
return get_async_backend().create_cancel_scope(deadline=deadline, shield=shield)
|
||||||
|
|
||||||
|
|
||||||
|
def current_effective_deadline() -> float:
|
||||||
|
"""
|
||||||
|
Return the nearest deadline among all the cancel scopes effective for the current
|
||||||
|
task.
|
||||||
|
|
||||||
|
:return: a clock value from the event loop's internal clock (or ``float('inf')`` if
|
||||||
|
there is no deadline in effect, or ``float('-inf')`` if the current scope has
|
||||||
|
been cancelled)
|
||||||
|
:rtype: float
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
return get_async_backend().current_effective_deadline()
|
||||||
|
|
||||||
|
|
||||||
|
def create_task_group() -> TaskGroup:
|
||||||
|
"""
|
||||||
|
Create a task group.
|
||||||
|
|
||||||
|
:return: a task group
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
return get_async_backend().create_task_group()
|
||||||
|
|
||||||
|
|
||||||
|
@final
|
||||||
|
class TaskHandle(Generic[T_co, T_startval]):
|
||||||
|
"""
|
||||||
|
Returned from the task-spawning methods of :class:`TaskGroup`. Can be awaited on to
|
||||||
|
get the return value of the task (or the raised exception). If the task was
|
||||||
|
terminated by a :exc:`BaseException`, :exc:`TaskFailed` will be raised (or its
|
||||||
|
subclass :exc:`TaskCancelled` if the task was cancelled).
|
||||||
|
|
||||||
|
.. versionadded:: 4.14.0
|
||||||
|
"""
|
||||||
|
|
||||||
|
class Status(Enum):
|
||||||
|
"""
|
||||||
|
The status of a task handle.
|
||||||
|
|
||||||
|
.. attribute:: PENDING
|
||||||
|
|
||||||
|
The task has not finished yet.
|
||||||
|
.. attribute:: FINISHED
|
||||||
|
|
||||||
|
The task has finished with a return value.
|
||||||
|
.. attribute:: CANCELLING
|
||||||
|
|
||||||
|
The task has been cancelled but has not finished yet.
|
||||||
|
.. attribute:: CANCELLED
|
||||||
|
|
||||||
|
The task was cancelled and has finished since.
|
||||||
|
.. attribute:: FAILED
|
||||||
|
|
||||||
|
The task raised an exception.
|
||||||
|
"""
|
||||||
|
|
||||||
|
PENDING = auto()
|
||||||
|
FINISHED = auto()
|
||||||
|
CANCELLING = auto()
|
||||||
|
CANCELLED = auto()
|
||||||
|
FAILED = auto()
|
||||||
|
|
||||||
|
__slots__ = (
|
||||||
|
"__weakref__",
|
||||||
|
"_coro",
|
||||||
|
"_name",
|
||||||
|
"_cancel_scope",
|
||||||
|
"_finished_event",
|
||||||
|
"_return_value",
|
||||||
|
"_start_value",
|
||||||
|
"_exception",
|
||||||
|
)
|
||||||
|
|
||||||
|
_return_value: T_co
|
||||||
|
_start_value: T_startval
|
||||||
|
|
||||||
|
def __init__(self, coro: Coroutine[Any, Any, T_co], name: object) -> None:
|
||||||
|
from ._synchronization import Event
|
||||||
|
|
||||||
|
self._coro = coro
|
||||||
|
self._cancel_scope = CancelScope()
|
||||||
|
self._finished_event = Event()
|
||||||
|
self._exception: BaseException | None = None
|
||||||
|
|
||||||
|
if name is not None:
|
||||||
|
self._name = str(name)
|
||||||
|
elif iscoroutine(coro):
|
||||||
|
self._name = coro.__qualname__
|
||||||
|
else:
|
||||||
|
self._name = str(coro) # coroutine-like object (e.g. asend() objects)
|
||||||
|
|
||||||
|
async def _run_coro(self) -> None:
|
||||||
|
__tracebackhide__ = True
|
||||||
|
|
||||||
|
with self._cancel_scope:
|
||||||
|
try:
|
||||||
|
retval = await self._coro
|
||||||
|
except BaseException as exc:
|
||||||
|
self._exception = exc
|
||||||
|
raise
|
||||||
|
else:
|
||||||
|
self._return_value = retval
|
||||||
|
finally:
|
||||||
|
self._finished_event.set()
|
||||||
|
del self # Break the reference cycle
|
||||||
|
|
||||||
|
def cancel(self) -> None:
|
||||||
|
"""
|
||||||
|
Set the task to a cancelled state.
|
||||||
|
|
||||||
|
This will interrupt any interruptible asynchronous operation, and will cause
|
||||||
|
any further awaits on this task to get immediately cancelled, unless done in
|
||||||
|
a shielded cancel scope.
|
||||||
|
|
||||||
|
If the task has already finished, this method has no effect.
|
||||||
|
"""
|
||||||
|
if not self._finished_event.is_set():
|
||||||
|
self._cancel_scope.cancel()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def coro(self) -> Coroutine[Any, Any, T_co]:
|
||||||
|
"""
|
||||||
|
The coroutine object that was passed to one of the task-spawning methods in
|
||||||
|
:class:`TaskGroup`.
|
||||||
|
"""
|
||||||
|
return self._coro
|
||||||
|
|
||||||
|
@property
|
||||||
|
def status(self) -> TaskHandle.Status:
|
||||||
|
"""
|
||||||
|
The current status of the task.
|
||||||
|
|
||||||
|
Every task starts in the :attr:`~TaskHandle.Status.PENDING` state.
|
||||||
|
If a task is cancelled while in this state, it will transition to the
|
||||||
|
:attr:`~TaskHandle.Status.CANCELLING` state. When the task finishes, it will
|
||||||
|
transition to one of the three final states (
|
||||||
|
:attr:`~TaskHandle.Status.FINISHED`, :attr:`~TaskHandle.Status.FAILED`, or
|
||||||
|
:attr:`~TaskHandle.Status.CANCELLING`) depending on the exception the task
|
||||||
|
raised, if any. No other status transitions will happen.
|
||||||
|
"""
|
||||||
|
if not self._finished_event.is_set():
|
||||||
|
if self._cancel_scope.cancel_called:
|
||||||
|
return TaskHandle.Status.CANCELLING
|
||||||
|
else:
|
||||||
|
return TaskHandle.Status.PENDING
|
||||||
|
elif self._exception is not None:
|
||||||
|
if isinstance(self._exception, get_cancelled_exc_class()):
|
||||||
|
return TaskHandle.Status.CANCELLED
|
||||||
|
else:
|
||||||
|
return TaskHandle.Status.FAILED
|
||||||
|
else:
|
||||||
|
return TaskHandle.Status.FINISHED
|
||||||
|
|
||||||
|
@property
|
||||||
|
def name(self) -> str:
|
||||||
|
"""The name of the task."""
|
||||||
|
return self._name
|
||||||
|
|
||||||
|
@property
|
||||||
|
def exception(self) -> BaseException | None:
|
||||||
|
"""
|
||||||
|
The exception raised by the task, or ``None`` if it finished without raising.
|
||||||
|
|
||||||
|
:raises TaskNotFinished: if the task has not finished yet
|
||||||
|
:raises TaskCancelled: if the task was cancelled
|
||||||
|
|
||||||
|
"""
|
||||||
|
match self.status:
|
||||||
|
case TaskHandle.Status.PENDING:
|
||||||
|
raise TaskNotFinished("the task has not finished yet")
|
||||||
|
case TaskHandle.Status.FINISHED:
|
||||||
|
return None
|
||||||
|
case TaskHandle.Status.CANCELLING:
|
||||||
|
raise TaskCancelled("the task was cancelled")
|
||||||
|
case TaskHandle.Status.CANCELLED:
|
||||||
|
raise TaskCancelled("the task was cancelled") from self._exception
|
||||||
|
case TaskHandle.Status.FAILED:
|
||||||
|
return self._exception
|
||||||
|
|
||||||
|
@property
|
||||||
|
def return_value(self) -> T_co:
|
||||||
|
"""
|
||||||
|
The return value of the task.
|
||||||
|
|
||||||
|
:raises TaskNotFinished: if the task has not finished yet
|
||||||
|
:raises TaskCancelled: if the task was cancelled
|
||||||
|
:raises TaskFailed: if the task raised an exception
|
||||||
|
|
||||||
|
"""
|
||||||
|
match self.status:
|
||||||
|
case TaskHandle.Status.PENDING:
|
||||||
|
raise TaskNotFinished("the task has not finished yet")
|
||||||
|
case TaskHandle.Status.FINISHED:
|
||||||
|
return self._return_value
|
||||||
|
case TaskHandle.Status.CANCELLING:
|
||||||
|
raise TaskCancelled("the task was cancelled")
|
||||||
|
case TaskHandle.Status.CANCELLED:
|
||||||
|
raise TaskCancelled("the task was cancelled") from self._exception
|
||||||
|
case TaskHandle.Status.FAILED:
|
||||||
|
raise TaskFailed("the task raised an exception") from self._exception
|
||||||
|
|
||||||
|
@property
|
||||||
|
def start_value(self) -> T_startval:
|
||||||
|
"""
|
||||||
|
The value passed to :meth:`task_status.started() <.abc.TaskStatus.started>`,
|
||||||
|
|
||||||
|
:raises RuntimeError: if the task was not started with :meth:`TaskGroup.start()
|
||||||
|
<.abc.TaskGroup.start>`
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
return self._start_value
|
||||||
|
except AttributeError:
|
||||||
|
raise RuntimeError(
|
||||||
|
"the task was not started with TaskGroup.start()"
|
||||||
|
) from None
|
||||||
|
|
||||||
|
async def wait(self) -> None:
|
||||||
|
"""
|
||||||
|
Wait for the task to finish.
|
||||||
|
|
||||||
|
This method will return as soon as the task has finished, no matter how it
|
||||||
|
happened.
|
||||||
|
"""
|
||||||
|
await self._finished_event.wait()
|
||||||
|
|
||||||
|
def __await__(self) -> Generator[Any, Any, T_co]:
|
||||||
|
yield from self._finished_event.wait().__await__()
|
||||||
|
return self.return_value
|
||||||
|
|
||||||
|
def __repr__(self) -> str:
|
||||||
|
return (
|
||||||
|
f"<{self.__class__.__name__} {self.status.name.lower()} "
|
||||||
|
f"name={self._name!r} coro={self._coro!r}>"
|
||||||
|
)
|
||||||
613
venv/Lib/site-packages/anyio/_core/_tempfile.py
Normal file
613
venv/Lib/site-packages/anyio/_core/_tempfile.py
Normal file
@@ -0,0 +1,613 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
from collections.abc import Iterable
|
||||||
|
from io import BytesIO, TextIOWrapper
|
||||||
|
from types import TracebackType
|
||||||
|
from typing import (
|
||||||
|
TYPE_CHECKING,
|
||||||
|
Any,
|
||||||
|
AnyStr,
|
||||||
|
Generic,
|
||||||
|
overload,
|
||||||
|
)
|
||||||
|
|
||||||
|
from .. import to_thread
|
||||||
|
from .._core._fileio import AsyncFile
|
||||||
|
from ..lowlevel import checkpoint_if_cancelled
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
from _typeshed import OpenBinaryMode, OpenTextMode, ReadableBuffer, WriteableBuffer
|
||||||
|
|
||||||
|
|
||||||
|
class TemporaryFile(Generic[AnyStr]):
|
||||||
|
"""
|
||||||
|
An asynchronous temporary file that is automatically created and cleaned up.
|
||||||
|
|
||||||
|
This class provides an asynchronous context manager interface to a temporary file.
|
||||||
|
The file is created using Python's standard `tempfile.TemporaryFile` function in a
|
||||||
|
background thread, and is wrapped as an asynchronous file using `AsyncFile`.
|
||||||
|
|
||||||
|
:param mode: The mode in which the file is opened. Defaults to "w+b".
|
||||||
|
:param buffering: The buffering policy (-1 means the default buffering).
|
||||||
|
:param encoding: The encoding used to decode or encode the file. Only applicable in
|
||||||
|
text mode.
|
||||||
|
:param newline: Controls how universal newlines mode works (only applicable in text
|
||||||
|
mode).
|
||||||
|
:param suffix: The suffix for the temporary file name.
|
||||||
|
:param prefix: The prefix for the temporary file name.
|
||||||
|
:param dir: The directory in which the temporary file is created.
|
||||||
|
:param errors: The error handling scheme used for encoding/decoding errors.
|
||||||
|
"""
|
||||||
|
|
||||||
|
_async_file: AsyncFile[AnyStr]
|
||||||
|
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: TemporaryFile[bytes],
|
||||||
|
mode: OpenBinaryMode = ...,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
): ...
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: TemporaryFile[str],
|
||||||
|
mode: OpenTextMode,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
): ...
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
mode: OpenTextMode | OpenBinaryMode = "w+b",
|
||||||
|
buffering: int = -1,
|
||||||
|
encoding: str | None = None,
|
||||||
|
newline: str | None = None,
|
||||||
|
suffix: str | None = None,
|
||||||
|
prefix: str | None = None,
|
||||||
|
dir: str | None = None,
|
||||||
|
*,
|
||||||
|
errors: str | None = None,
|
||||||
|
) -> None:
|
||||||
|
self.mode = mode
|
||||||
|
self.buffering = buffering
|
||||||
|
self.encoding = encoding
|
||||||
|
self.newline = newline
|
||||||
|
self.suffix: str | None = suffix
|
||||||
|
self.prefix: str | None = prefix
|
||||||
|
self.dir: str | None = dir
|
||||||
|
self.errors = errors
|
||||||
|
|
||||||
|
async def __aenter__(self) -> AsyncFile[AnyStr]:
|
||||||
|
fp = await to_thread.run_sync(
|
||||||
|
lambda: tempfile.TemporaryFile(
|
||||||
|
self.mode,
|
||||||
|
self.buffering,
|
||||||
|
self.encoding,
|
||||||
|
self.newline,
|
||||||
|
self.suffix,
|
||||||
|
self.prefix,
|
||||||
|
self.dir,
|
||||||
|
errors=self.errors,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
self._async_file = AsyncFile(fp)
|
||||||
|
return self._async_file
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_value: BaseException | None,
|
||||||
|
traceback: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
await self._async_file.aclose()
|
||||||
|
|
||||||
|
|
||||||
|
class NamedTemporaryFile(Generic[AnyStr]):
|
||||||
|
"""
|
||||||
|
An asynchronous named temporary file that is automatically created and cleaned up.
|
||||||
|
|
||||||
|
This class provides an asynchronous context manager for a temporary file with a
|
||||||
|
visible name in the file system. It uses Python's standard
|
||||||
|
:func:`~tempfile.NamedTemporaryFile` function and wraps the file object with
|
||||||
|
:class:`AsyncFile` for asynchronous operations.
|
||||||
|
|
||||||
|
:param mode: The mode in which the file is opened. Defaults to "w+b".
|
||||||
|
:param buffering: The buffering policy (-1 means the default buffering).
|
||||||
|
:param encoding: The encoding used to decode or encode the file. Only applicable in
|
||||||
|
text mode.
|
||||||
|
:param newline: Controls how universal newlines mode works (only applicable in text
|
||||||
|
mode).
|
||||||
|
:param suffix: The suffix for the temporary file name.
|
||||||
|
:param prefix: The prefix for the temporary file name.
|
||||||
|
:param dir: The directory in which the temporary file is created.
|
||||||
|
:param delete: Whether to delete the file when it is closed.
|
||||||
|
:param errors: The error handling scheme used for encoding/decoding errors.
|
||||||
|
:param delete_on_close: (Python 3.12+) Whether to delete the file on close.
|
||||||
|
"""
|
||||||
|
|
||||||
|
_async_file: AsyncFile[AnyStr]
|
||||||
|
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: NamedTemporaryFile[bytes],
|
||||||
|
mode: OpenBinaryMode = ...,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
delete: bool = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
delete_on_close: bool = ...,
|
||||||
|
): ...
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: NamedTemporaryFile[str],
|
||||||
|
mode: OpenTextMode,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
delete: bool = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
delete_on_close: bool = ...,
|
||||||
|
): ...
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
mode: OpenBinaryMode | OpenTextMode = "w+b",
|
||||||
|
buffering: int = -1,
|
||||||
|
encoding: str | None = None,
|
||||||
|
newline: str | None = None,
|
||||||
|
suffix: str | None = None,
|
||||||
|
prefix: str | None = None,
|
||||||
|
dir: str | None = None,
|
||||||
|
delete: bool = True,
|
||||||
|
*,
|
||||||
|
errors: str | None = None,
|
||||||
|
delete_on_close: bool = True,
|
||||||
|
) -> None:
|
||||||
|
self._params: dict[str, Any] = {
|
||||||
|
"mode": mode,
|
||||||
|
"buffering": buffering,
|
||||||
|
"encoding": encoding,
|
||||||
|
"newline": newline,
|
||||||
|
"suffix": suffix,
|
||||||
|
"prefix": prefix,
|
||||||
|
"dir": dir,
|
||||||
|
"delete": delete,
|
||||||
|
"errors": errors,
|
||||||
|
}
|
||||||
|
if sys.version_info >= (3, 12):
|
||||||
|
self._params["delete_on_close"] = delete_on_close
|
||||||
|
|
||||||
|
async def __aenter__(self) -> AsyncFile[AnyStr]:
|
||||||
|
fp = await to_thread.run_sync(
|
||||||
|
lambda: tempfile.NamedTemporaryFile(**self._params)
|
||||||
|
)
|
||||||
|
self._async_file = AsyncFile(fp)
|
||||||
|
return self._async_file
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_value: BaseException | None,
|
||||||
|
traceback: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
await self._async_file.aclose()
|
||||||
|
|
||||||
|
|
||||||
|
class SpooledTemporaryFile(AsyncFile[AnyStr]):
|
||||||
|
"""
|
||||||
|
An asynchronous spooled temporary file that starts in memory and is spooled to disk.
|
||||||
|
|
||||||
|
This class provides an asynchronous interface to a spooled temporary file, much like
|
||||||
|
Python's standard :class:`~tempfile.SpooledTemporaryFile`. It supports asynchronous
|
||||||
|
write operations and provides a method to force a rollover to disk.
|
||||||
|
|
||||||
|
:param max_size: Maximum size in bytes before the file is rolled over to disk.
|
||||||
|
:param mode: The mode in which the file is opened. Defaults to "w+b".
|
||||||
|
:param buffering: The buffering policy (-1 means the default buffering).
|
||||||
|
:param encoding: The encoding used to decode or encode the file (text mode only).
|
||||||
|
:param newline: Controls how universal newlines mode works (text mode only).
|
||||||
|
:param suffix: The suffix for the temporary file name.
|
||||||
|
:param prefix: The prefix for the temporary file name.
|
||||||
|
:param dir: The directory in which the temporary file is created.
|
||||||
|
:param errors: The error handling scheme used for encoding/decoding errors.
|
||||||
|
"""
|
||||||
|
|
||||||
|
_rolled: bool = False
|
||||||
|
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: SpooledTemporaryFile[bytes],
|
||||||
|
max_size: int = ...,
|
||||||
|
mode: OpenBinaryMode = ...,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
): ...
|
||||||
|
@overload
|
||||||
|
def __init__(
|
||||||
|
self: SpooledTemporaryFile[str],
|
||||||
|
max_size: int = ...,
|
||||||
|
mode: OpenTextMode = ...,
|
||||||
|
buffering: int = ...,
|
||||||
|
encoding: str | None = ...,
|
||||||
|
newline: str | None = ...,
|
||||||
|
suffix: str | None = ...,
|
||||||
|
prefix: str | None = ...,
|
||||||
|
dir: str | None = ...,
|
||||||
|
*,
|
||||||
|
errors: str | None = ...,
|
||||||
|
): ...
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
max_size: int = 0,
|
||||||
|
mode: OpenBinaryMode | OpenTextMode = "w+b",
|
||||||
|
buffering: int = -1,
|
||||||
|
encoding: str | None = None,
|
||||||
|
newline: str | None = None,
|
||||||
|
suffix: str | None = None,
|
||||||
|
prefix: str | None = None,
|
||||||
|
dir: str | None = None,
|
||||||
|
*,
|
||||||
|
errors: str | None = None,
|
||||||
|
) -> None:
|
||||||
|
self._tempfile_params: dict[str, Any] = {
|
||||||
|
"mode": mode,
|
||||||
|
"buffering": buffering,
|
||||||
|
"encoding": encoding,
|
||||||
|
"newline": newline,
|
||||||
|
"suffix": suffix,
|
||||||
|
"prefix": prefix,
|
||||||
|
"dir": dir,
|
||||||
|
"errors": errors,
|
||||||
|
}
|
||||||
|
self._max_size = max_size
|
||||||
|
if "b" in mode:
|
||||||
|
super().__init__(BytesIO()) # type: ignore[arg-type]
|
||||||
|
else:
|
||||||
|
super().__init__(
|
||||||
|
TextIOWrapper( # type: ignore[arg-type]
|
||||||
|
BytesIO(),
|
||||||
|
encoding=encoding,
|
||||||
|
errors=errors,
|
||||||
|
newline=newline,
|
||||||
|
write_through=True,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
async def aclose(self) -> None:
|
||||||
|
if not self._rolled:
|
||||||
|
self._fp.close()
|
||||||
|
return
|
||||||
|
|
||||||
|
await super().aclose()
|
||||||
|
|
||||||
|
async def _check(self) -> None:
|
||||||
|
if self._rolled or self._fp.tell() <= self._max_size:
|
||||||
|
return
|
||||||
|
|
||||||
|
await self.rollover()
|
||||||
|
|
||||||
|
async def rollover(self) -> None:
|
||||||
|
if self._rolled:
|
||||||
|
return
|
||||||
|
|
||||||
|
self._rolled = True
|
||||||
|
buffer = self._fp
|
||||||
|
buffer.seek(0)
|
||||||
|
self._fp = await to_thread.run_sync(
|
||||||
|
lambda: tempfile.TemporaryFile(**self._tempfile_params)
|
||||||
|
)
|
||||||
|
await self.write(buffer.read())
|
||||||
|
buffer.close()
|
||||||
|
|
||||||
|
@property
|
||||||
|
def closed(self) -> bool:
|
||||||
|
return self._fp.closed
|
||||||
|
|
||||||
|
async def read(self, size: int = -1) -> AnyStr:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.read(size)
|
||||||
|
|
||||||
|
return await super().read(size) # type: ignore[return-value]
|
||||||
|
|
||||||
|
async def read1(self: SpooledTemporaryFile[bytes], size: int = -1) -> bytes:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.read1(size)
|
||||||
|
|
||||||
|
return await super().read1(size)
|
||||||
|
|
||||||
|
async def readline(self) -> AnyStr:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.readline()
|
||||||
|
|
||||||
|
return await super().readline() # type: ignore[return-value]
|
||||||
|
|
||||||
|
async def readlines(self) -> list[AnyStr]:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.readlines()
|
||||||
|
|
||||||
|
return await super().readlines() # type: ignore[return-value]
|
||||||
|
|
||||||
|
async def readinto(self: SpooledTemporaryFile[bytes], b: WriteableBuffer) -> int:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
self._fp.readinto(b)
|
||||||
|
|
||||||
|
return await super().readinto(b)
|
||||||
|
|
||||||
|
async def readinto1(self: SpooledTemporaryFile[bytes], b: WriteableBuffer) -> int:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
self._fp.readinto(b)
|
||||||
|
|
||||||
|
return await super().readinto1(b)
|
||||||
|
|
||||||
|
async def seek(self, offset: int, whence: int | None = os.SEEK_SET) -> int:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.seek(offset, whence)
|
||||||
|
|
||||||
|
return await super().seek(offset, whence)
|
||||||
|
|
||||||
|
async def tell(self) -> int:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.tell()
|
||||||
|
|
||||||
|
return await super().tell()
|
||||||
|
|
||||||
|
async def truncate(self, size: int | None = None) -> int:
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
return self._fp.truncate(size)
|
||||||
|
|
||||||
|
return await super().truncate(size)
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def write(self: SpooledTemporaryFile[bytes], b: ReadableBuffer) -> int: ...
|
||||||
|
@overload
|
||||||
|
async def write(self: SpooledTemporaryFile[str], b: str) -> int: ...
|
||||||
|
|
||||||
|
async def write(self, b: ReadableBuffer | str) -> int:
|
||||||
|
"""
|
||||||
|
Asynchronously write data to the spooled temporary file.
|
||||||
|
|
||||||
|
If the file has not yet been rolled over, the data is written synchronously,
|
||||||
|
and a rollover is triggered if the size exceeds the maximum size.
|
||||||
|
|
||||||
|
:param s: The data to write.
|
||||||
|
:return: The number of bytes written.
|
||||||
|
:raises RuntimeError: If the underlying file is not initialized.
|
||||||
|
|
||||||
|
"""
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
result = self._fp.write(b)
|
||||||
|
await self._check()
|
||||||
|
return result
|
||||||
|
|
||||||
|
return await super().write(b) # type: ignore[misc]
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def writelines(
|
||||||
|
self: SpooledTemporaryFile[bytes], lines: Iterable[ReadableBuffer]
|
||||||
|
) -> None: ...
|
||||||
|
@overload
|
||||||
|
async def writelines(
|
||||||
|
self: SpooledTemporaryFile[str], lines: Iterable[str]
|
||||||
|
) -> None: ...
|
||||||
|
|
||||||
|
async def writelines(self, lines: Iterable[str] | Iterable[ReadableBuffer]) -> None:
|
||||||
|
"""
|
||||||
|
Asynchronously write a list of lines to the spooled temporary file.
|
||||||
|
|
||||||
|
If the file has not yet been rolled over, the lines are written synchronously,
|
||||||
|
and a rollover is triggered if the size exceeds the maximum size.
|
||||||
|
|
||||||
|
:param lines: An iterable of lines to write.
|
||||||
|
:raises RuntimeError: If the underlying file is not initialized.
|
||||||
|
|
||||||
|
"""
|
||||||
|
if not self._rolled:
|
||||||
|
await checkpoint_if_cancelled()
|
||||||
|
result = self._fp.writelines(lines)
|
||||||
|
await self._check()
|
||||||
|
return result
|
||||||
|
|
||||||
|
return await super().writelines(lines) # type: ignore[misc]
|
||||||
|
|
||||||
|
|
||||||
|
class TemporaryDirectory(Generic[AnyStr]):
|
||||||
|
"""
|
||||||
|
An asynchronous temporary directory that is created and cleaned up automatically.
|
||||||
|
|
||||||
|
This class provides an asynchronous context manager for creating a temporary
|
||||||
|
directory. It wraps Python's standard :class:`~tempfile.TemporaryDirectory` to
|
||||||
|
perform directory creation and cleanup operations in a background thread.
|
||||||
|
|
||||||
|
:param suffix: Suffix to be added to the temporary directory name.
|
||||||
|
:param prefix: Prefix to be added to the temporary directory name.
|
||||||
|
:param dir: The parent directory where the temporary directory is created.
|
||||||
|
:param ignore_cleanup_errors: Whether to ignore errors during cleanup
|
||||||
|
:param delete: Whether to delete the directory upon closing (Python 3.12+).
|
||||||
|
"""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
suffix: AnyStr | None = None,
|
||||||
|
prefix: AnyStr | None = None,
|
||||||
|
dir: AnyStr | None = None,
|
||||||
|
*,
|
||||||
|
ignore_cleanup_errors: bool = False,
|
||||||
|
delete: bool = True,
|
||||||
|
) -> None:
|
||||||
|
self.suffix: AnyStr | None = suffix
|
||||||
|
self.prefix: AnyStr | None = prefix
|
||||||
|
self.dir: AnyStr | None = dir
|
||||||
|
self.ignore_cleanup_errors = ignore_cleanup_errors
|
||||||
|
self.delete = delete
|
||||||
|
|
||||||
|
self._tempdir: tempfile.TemporaryDirectory | None = None
|
||||||
|
|
||||||
|
async def __aenter__(self) -> str:
|
||||||
|
params: dict[str, Any] = {
|
||||||
|
"suffix": self.suffix,
|
||||||
|
"prefix": self.prefix,
|
||||||
|
"dir": self.dir,
|
||||||
|
"ignore_cleanup_errors": self.ignore_cleanup_errors,
|
||||||
|
}
|
||||||
|
if sys.version_info >= (3, 12):
|
||||||
|
params["delete"] = self.delete
|
||||||
|
|
||||||
|
self._tempdir = await to_thread.run_sync(
|
||||||
|
lambda: tempfile.TemporaryDirectory(**params)
|
||||||
|
)
|
||||||
|
return await to_thread.run_sync(self._tempdir.__enter__)
|
||||||
|
|
||||||
|
async def __aexit__(
|
||||||
|
self,
|
||||||
|
exc_type: type[BaseException] | None,
|
||||||
|
exc_value: BaseException | None,
|
||||||
|
traceback: TracebackType | None,
|
||||||
|
) -> None:
|
||||||
|
if self._tempdir is not None:
|
||||||
|
await to_thread.run_sync(
|
||||||
|
self._tempdir.__exit__, exc_type, exc_value, traceback
|
||||||
|
)
|
||||||
|
|
||||||
|
async def cleanup(self) -> None:
|
||||||
|
if self._tempdir is not None:
|
||||||
|
await to_thread.run_sync(self._tempdir.cleanup)
|
||||||
|
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def mkstemp(
|
||||||
|
suffix: str | None = None,
|
||||||
|
prefix: str | None = None,
|
||||||
|
dir: str | None = None,
|
||||||
|
text: bool = False,
|
||||||
|
) -> tuple[int, str]: ...
|
||||||
|
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def mkstemp(
|
||||||
|
suffix: bytes | None = None,
|
||||||
|
prefix: bytes | None = None,
|
||||||
|
dir: bytes | None = None,
|
||||||
|
text: bool = False,
|
||||||
|
) -> tuple[int, bytes]: ...
|
||||||
|
|
||||||
|
|
||||||
|
async def mkstemp(
|
||||||
|
suffix: AnyStr | None = None,
|
||||||
|
prefix: AnyStr | None = None,
|
||||||
|
dir: AnyStr | None = None,
|
||||||
|
text: bool = False,
|
||||||
|
) -> tuple[int, str | bytes]:
|
||||||
|
"""
|
||||||
|
Asynchronously create a temporary file and return an OS-level handle and the file
|
||||||
|
name.
|
||||||
|
|
||||||
|
This function wraps `tempfile.mkstemp` and executes it in a background thread.
|
||||||
|
|
||||||
|
:param suffix: Suffix to be added to the file name.
|
||||||
|
:param prefix: Prefix to be added to the file name.
|
||||||
|
:param dir: Directory in which the temporary file is created.
|
||||||
|
:param text: Whether the file is opened in text mode.
|
||||||
|
:return: A tuple containing the file descriptor and the file name.
|
||||||
|
|
||||||
|
"""
|
||||||
|
return await to_thread.run_sync(tempfile.mkstemp, suffix, prefix, dir, text)
|
||||||
|
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def mkdtemp(
|
||||||
|
suffix: str | None = None,
|
||||||
|
prefix: str | None = None,
|
||||||
|
dir: str | None = None,
|
||||||
|
) -> str: ...
|
||||||
|
|
||||||
|
|
||||||
|
@overload
|
||||||
|
async def mkdtemp(
|
||||||
|
suffix: bytes | None = None,
|
||||||
|
prefix: bytes | None = None,
|
||||||
|
dir: bytes | None = None,
|
||||||
|
) -> bytes: ...
|
||||||
|
|
||||||
|
|
||||||
|
async def mkdtemp(
|
||||||
|
suffix: AnyStr | None = None,
|
||||||
|
prefix: AnyStr | None = None,
|
||||||
|
dir: AnyStr | None = None,
|
||||||
|
) -> str | bytes:
|
||||||
|
"""
|
||||||
|
Asynchronously create a temporary directory and return its path.
|
||||||
|
|
||||||
|
This function wraps `tempfile.mkdtemp` and executes it in a background thread.
|
||||||
|
|
||||||
|
:param suffix: Suffix to be added to the directory name.
|
||||||
|
:param prefix: Prefix to be added to the directory name.
|
||||||
|
:param dir: Parent directory where the temporary directory is created.
|
||||||
|
:return: The path of the created temporary directory.
|
||||||
|
|
||||||
|
"""
|
||||||
|
return await to_thread.run_sync(tempfile.mkdtemp, suffix, prefix, dir)
|
||||||
|
|
||||||
|
|
||||||
|
async def gettempdir() -> str:
|
||||||
|
"""
|
||||||
|
Asynchronously return the name of the directory used for temporary files.
|
||||||
|
|
||||||
|
This function wraps `tempfile.gettempdir` and executes it in a background thread.
|
||||||
|
|
||||||
|
:return: The path of the temporary directory as a string.
|
||||||
|
|
||||||
|
"""
|
||||||
|
return await to_thread.run_sync(tempfile.gettempdir)
|
||||||
|
|
||||||
|
|
||||||
|
async def gettempdirb() -> bytes:
|
||||||
|
"""
|
||||||
|
Asynchronously return the name of the directory used for temporary files in bytes.
|
||||||
|
|
||||||
|
This function wraps `tempfile.gettempdirb` and executes it in a background thread.
|
||||||
|
|
||||||
|
:return: The path of the temporary directory as bytes.
|
||||||
|
|
||||||
|
"""
|
||||||
|
return await to_thread.run_sync(tempfile.gettempdirb)
|
||||||
82
venv/Lib/site-packages/anyio/_core/_testing.py
Normal file
82
venv/Lib/site-packages/anyio/_core/_testing.py
Normal file
@@ -0,0 +1,82 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import Awaitable, Generator
|
||||||
|
from typing import Any, cast
|
||||||
|
|
||||||
|
from ._eventloop import get_async_backend
|
||||||
|
|
||||||
|
|
||||||
|
class TaskInfo:
|
||||||
|
"""
|
||||||
|
Represents an asynchronous task.
|
||||||
|
|
||||||
|
:ivar int id: the unique identifier of the task
|
||||||
|
:ivar parent_id: the identifier of the parent task, if any
|
||||||
|
:vartype parent_id: Optional[int]
|
||||||
|
:ivar str name: the description of the task (if any)
|
||||||
|
:ivar ~collections.abc.Coroutine coro: the coroutine object of the task
|
||||||
|
"""
|
||||||
|
|
||||||
|
__slots__ = "_name", "id", "parent_id", "name", "coro"
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
id: int,
|
||||||
|
parent_id: int | None,
|
||||||
|
name: str | None,
|
||||||
|
coro: Generator[Any, Any, Any] | Awaitable[Any],
|
||||||
|
):
|
||||||
|
func = get_current_task
|
||||||
|
self._name = f"{func.__module__}.{func.__qualname__}"
|
||||||
|
self.id: int = id
|
||||||
|
self.parent_id: int | None = parent_id
|
||||||
|
self.name: str | None = name
|
||||||
|
self.coro: Generator[Any, Any, Any] | Awaitable[Any] = coro
|
||||||
|
|
||||||
|
def __eq__(self, other: object) -> bool:
|
||||||
|
if isinstance(other, TaskInfo):
|
||||||
|
return self.id == other.id
|
||||||
|
|
||||||
|
return NotImplemented
|
||||||
|
|
||||||
|
def __hash__(self) -> int:
|
||||||
|
return hash(self.id)
|
||||||
|
|
||||||
|
def __repr__(self) -> str:
|
||||||
|
return f"{self.__class__.__name__}(id={self.id!r}, name={self.name!r})"
|
||||||
|
|
||||||
|
def has_pending_cancellation(self) -> bool:
|
||||||
|
"""
|
||||||
|
Return ``True`` if the task has a cancellation pending, ``False`` otherwise.
|
||||||
|
|
||||||
|
"""
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def get_current_task() -> TaskInfo:
|
||||||
|
"""
|
||||||
|
Return the current task.
|
||||||
|
|
||||||
|
:return: a representation of the current task
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
return get_async_backend().get_current_task()
|
||||||
|
|
||||||
|
|
||||||
|
def get_running_tasks() -> list[TaskInfo]:
|
||||||
|
"""
|
||||||
|
Return a list of running tasks in the current event loop.
|
||||||
|
|
||||||
|
:return: a list of task info objects
|
||||||
|
:raises NoEventLoopError: if no supported asynchronous event loop is running in the
|
||||||
|
current thread
|
||||||
|
|
||||||
|
"""
|
||||||
|
return cast("list[TaskInfo]", get_async_backend().get_running_tasks())
|
||||||
|
|
||||||
|
|
||||||
|
async def wait_all_tasks_blocked() -> None:
|
||||||
|
"""Wait until all other tasks are waiting for something."""
|
||||||
|
await get_async_backend().wait_all_tasks_blocked()
|
||||||
Reference in New Issue
Block a user