diff --git a/AGENTS.md b/AGENTS.md index 6ea925f..e2bc1a4 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -86,7 +86,7 @@ These come from the LoRa channel's physical constraints and must be preserved by - Built on `anyio` (asyncio backend). `Bridge.run` creates a single task group: one consumer task per transport, one egress worker per node, one notifier flush loop. No raw `asyncio.create_task` — use the task group so cancellation propagates correctly. - `CommitQueue` uses `anyio.create_memory_object_stream` for the queue and exposes `offer()` (non-blocking; returns `False` on full/rate-limited) and async iteration on the receive side. Don't await inside `offer`. -- Mirror-to-messenger errors are swallowed in `Bridge.mirror_to_messenger` on purpose (`# noqa: BLE001`) — a flaky messenger must not stall the LoRa pipeline. +- Mirror-to-messenger errors are swallowed in `Bridge.mirror_to_messenger` on purpose — a flaky messenger must not stall the LoRa pipeline. ### Configuration model diff --git a/docs/gen_pages.py b/docs/gen_pages.py index a97cef8..d51d7ec 100644 --- a/docs/gen_pages.py +++ b/docs/gen_pages.py @@ -18,14 +18,14 @@ import inspect import typing -from typing import Any, Union, get_args, get_origin +from typing import Any, get_args, get_origin import mkdocs_gen_files from pydantic import BaseModel from pydantic.fields import FieldInfo from lora_bridge.config import schema - +from lora_bridge.config.introspect import is_union_origin, strip_annotated # --------------------------------------------------------------------------- # Описания секций (то, что не выводится из типов) @@ -185,13 +185,13 @@ def _collect_models(root: type[BaseModel]) -> list[type[BaseModel]]: def _models_in(t: Any) -> list[type[BaseModel]]: """Все BaseModel-классы, до которых можно дотянуться, развернув ``t``.""" - t = _unwrap_annotated(t) + t = strip_annotated(t) if isinstance(t, type) and issubclass(t, BaseModel): return [t] origin = get_origin(t) if origin in (list, dict, tuple, set, frozenset): return [m for a in get_args(t) for m in _models_in(a)] - if origin in (Union, typing.Union): + if is_union_origin(origin): return [m for a in get_args(t) for m in _models_in(a)] return [] @@ -235,7 +235,7 @@ def _render_type(fi: FieldInfo) -> str: def _render_discriminated(fi: FieldInfo) -> str: """Discriminated union → «один из (тегов)» с ссылками на варианты.""" - variants = [_unwrap_annotated(a) for a in get_args(_unwrap_annotated(fi.annotation))] + variants = [strip_annotated(a) for a in get_args(strip_annotated(fi.annotation))] discr_field = fi.discriminator if isinstance(fi.discriminator, str) else "type" parts: list[str] = [] for v in variants: @@ -254,7 +254,7 @@ def _discriminator_value(model: type[BaseModel], field: str) -> str | None: fi = model.model_fields.get(field) if fi is None: return None - ann = _unwrap_annotated(fi.annotation) + ann = strip_annotated(fi.annotation) if get_origin(ann) is typing.Literal: args = get_args(ann) return repr(args[0]) if args else None @@ -263,7 +263,7 @@ def _discriminator_value(model: type[BaseModel], field: str) -> str | None: def _pretty_type(t: Any) -> str: """Markdown-friendly рендер аннотации типа для ячейки таблицы.""" - t = _unwrap_annotated(t) + t = strip_annotated(t) sup = getattr(t, "__supertype__", None) if sup is not None: # NewType — рендерим имя без ссылки: смысл id виден из описания соседних @@ -288,7 +288,7 @@ def _pretty_type(t: Any) -> str: return "кортеж (" + ", ".join(_pretty_type(a) for a in args) + ")" if origin is typing.Literal: return " \\| ".join(f"`{a!r}`" for a in args) - if origin in (Union, typing.Union): + if is_union_origin(origin): non_none = [a for a in args if a is not type(None)] rendered = " \\| ".join(_pretty_type(a) for a in non_none) if len(non_none) < len(args): @@ -303,7 +303,7 @@ def _render_default(fi: FieldInfo) -> str: if fi.default_factory is not None: try: v = fi.default_factory() # type: ignore[call-arg] - except Exception: + except Exception: # noqa: BLE001 — default_factory может кинуть; показываем прочерк return "—" return f"`{_repr_default(v)}`" if fi.default is None: @@ -324,12 +324,6 @@ def _escape_cell(text: str) -> str: return text.replace("|", r"\|").replace("\n", " ").strip() -def _unwrap_annotated(t: Any) -> Any: - while hasattr(t, "__metadata__"): - t = t.__origin__ - return t - - def emit_commands_page(*, path: str, title: str) -> None: from lora_bridge.transports.telegram.commands import ALL_COMMAND_METAS diff --git a/lora_bridge/config/errors.py b/lora_bridge/config/errors.py index 4a67c63..66c0674 100644 --- a/lora_bridge/config/errors.py +++ b/lora_bridge/config/errors.py @@ -15,10 +15,11 @@ from __future__ import annotations import typing -from typing import Any, Union, get_args, get_origin +from typing import Any, get_args, get_origin from pydantic import BaseModel, ValidationError +from .introspect import is_union_origin, strip_annotated from .schema import AppConfig __all__ = ["format_validation_error"] @@ -278,7 +279,7 @@ def _humanize_model(model: type[BaseModel]) -> str: def _pretty_type(t: Any) -> str: - t = _strip_annotated(t) + t = strip_annotated(t) if t is type(None): return "null" # NewType — показываем «NodeId (строка)»: семантика + базовый тип @@ -293,7 +294,7 @@ def _pretty_type(t: Any) -> str: if origin is dict: k, v = get_args(t) return f"словарь {_pretty_type(k)} → {_pretty_type(v)}" - if origin in (Union, typing.Union): + if is_union_origin(origin): inner = [_pretty_type(a) for a in get_args(t) if a is not type(None)] return " | ".join(inner) if origin is typing.Literal: @@ -304,12 +305,6 @@ def _pretty_type(t: Any) -> str: # --- резолвер: loc → набор кандидатов BaseModel ---------------------------- -def _strip_annotated(t: Any) -> Any: - while hasattr(t, "__metadata__"): - t = t.__origin__ - return t - - def _resolve_models(loc: tuple[Any, ...]) -> list[type[BaseModel]]: """Идём по ``loc`` в типовом дереве ``AppConfig``; возвращаем модели на конце пути. @@ -341,28 +336,28 @@ def _model_with_field(models: list[type[BaseModel]], field: Any) -> type[BaseMod def _step(node: Any, step: Any) -> list[Any]: - node = _strip_annotated(node) + node = strip_annotated(node) origin = get_origin(node) if isinstance(node, type) and issubclass(node, BaseModel): if isinstance(step, str) and step in node.model_fields: - return [_strip_annotated(node.model_fields[step].annotation)] + return [strip_annotated(node.model_fields[step].annotation)] # smart-union: loc содержит имя класса варианта прямо здесь if isinstance(step, str) and step == node.__name__: return [node] return [] if origin is list and isinstance(step, int): - return [_strip_annotated(get_args(node)[0])] + return [strip_annotated(get_args(node)[0])] if origin is dict and isinstance(step, str): - return [_strip_annotated(get_args(node)[1])] + return [strip_annotated(get_args(node)[1])] - if origin in (Union, typing.Union): + if is_union_origin(origin): # discriminator-тег: сузим Union до варианта, у которого Literal[step] if isinstance(step, str): for arg in get_args(node): - arg_t = _strip_annotated(arg) + arg_t = strip_annotated(arg) if not (isinstance(arg_t, type) and issubclass(arg_t, BaseModel)): continue if arg_t.__name__ == step: @@ -381,8 +376,8 @@ def _step(node: Any, step: Any) -> list[Any]: def _expand_union(t: Any) -> list[Any]: - t = _strip_annotated(t) - if get_origin(t) in (Union, typing.Union): + t = strip_annotated(t) + if is_union_origin(get_origin(t)): out: list[Any] = [] for arg in get_args(t): if arg is type(None): @@ -394,7 +389,7 @@ def _expand_union(t: Any) -> list[Any]: def _variant_matches_tag(variant: type[BaseModel], tag: str) -> bool: for fi in variant.model_fields.values(): - ann = _strip_annotated(fi.annotation) + ann = strip_annotated(fi.annotation) if get_origin(ann) is typing.Literal and tag in get_args(ann): return True return False @@ -406,7 +401,7 @@ def _collect_discriminator_tags(root: type[BaseModel] | None = None) -> set[str] seen: set[Any] = set() def visit(t: Any) -> None: - t = _strip_annotated(t) + t = strip_annotated(t) if t in seen: return seen.add(t) @@ -415,12 +410,12 @@ def visit(t: Any) -> None: visit(fi.annotation) return origin = get_origin(t) - if origin in (Union, typing.Union): + if is_union_origin(origin): for arg in get_args(t): - arg_t = _strip_annotated(arg) + arg_t = strip_annotated(arg) if isinstance(arg_t, type) and issubclass(arg_t, BaseModel): for fi in arg_t.model_fields.values(): - ann = _strip_annotated(fi.annotation) + ann = strip_annotated(fi.annotation) if get_origin(ann) is typing.Literal: for v in get_args(ann): if isinstance(v, str): diff --git a/lora_bridge/config/introspect.py b/lora_bridge/config/introspect.py new file mode 100644 index 0000000..057d999 --- /dev/null +++ b/lora_bridge/config/introspect.py @@ -0,0 +1,23 @@ +"""Примитивы интроспекции типового дерева конфиг-схемы. + +Общие для рендера конфиг-ошибок (``config/errors.py``) и генератора +справочника конфига (``docs/gen_pages.py``). +""" + +from __future__ import annotations + +import types +import typing +from typing import Any + + +def is_union_origin(origin: object) -> bool: + """Union в обеих формах: ``typing.Union[X, Y]`` и PEP 604 ``X | Y`` (types.UnionType).""" + return origin is typing.Union or origin is types.UnionType + + +def strip_annotated(t: Any) -> Any: + """Снять слои ``Annotated[...]``, добравшись до базового типа.""" + while hasattr(t, "__metadata__"): + t = t.__origin__ + return t diff --git a/lora_bridge/config/schema/app_config.py b/lora_bridge/config/schema/app_config.py index 7639c41..8f8ef1f 100644 --- a/lora_bridge/config/schema/app_config.py +++ b/lora_bridge/config/schema/app_config.py @@ -6,11 +6,11 @@ from pydantic import BaseModel, Field, model_validator +from ...domain.models import messenger_channel from .ids import EndpointName, NodeId from .messengers import MessengerConfig from .nodes import LoraNode from .rooms import LoraRef, LoraSubscriber, MessengerSubscriber, RoomConfig -from ...domain.models import messenger_channel def validate_lora_ref( diff --git a/lora_bridge/config/schema/connections.py b/lora_bridge/config/schema/connections.py index 0323547..7bd0d4d 100644 --- a/lora_bridge/config/schema/connections.py +++ b/lora_bridge/config/schema/connections.py @@ -8,7 +8,7 @@ from __future__ import annotations -from typing import Annotated, Literal, Union +from typing import Annotated, Literal from pydantic import BaseModel, Field @@ -64,7 +64,7 @@ class BleConnection(ConnectionBase): Connection = Annotated[ - Union[UsbConnection, SerialConnection, TcpConnection, BleConnection], + UsbConnection | SerialConnection | TcpConnection | BleConnection, Field(discriminator="type"), ] """Способ физического подключения к LoRa-узлу. diff --git a/lora_bridge/config/schema/endpoints.py b/lora_bridge/config/schema/endpoints.py index 7fa6e90..f43af93 100644 --- a/lora_bridge/config/schema/endpoints.py +++ b/lora_bridge/config/schema/endpoints.py @@ -7,7 +7,7 @@ from __future__ import annotations -from typing import Annotated, Literal, Optional, Union +from typing import Annotated, Literal from pydantic import BaseModel, Field, field_validator @@ -103,7 +103,7 @@ def normalize_pubkey(cls, value: str) -> str: Нормализуем здесь, чтобы ключи join'ились независимо от регистра. """ return value.lower() - password: Optional[str] = Field( + password: str | None = Field( default=None, description=( "Гостевой пароль. Если опущен — доступ read-only (постинг недоступен)." @@ -112,7 +112,7 @@ def normalize_pubkey(cls, value: str) -> str: Endpoint = Annotated[ - Union[PublicEndpoint, PrivateEndpoint, RoomServerEndpoint], + PublicEndpoint | PrivateEndpoint | RoomServerEndpoint, Field(discriminator="type"), ] """Тип LoRa-эндпоинта в MeshCore-ноде. diff --git a/lora_bridge/config/schema/messengers.py b/lora_bridge/config/schema/messengers.py index 8c3c11d..aaddf68 100644 --- a/lora_bridge/config/schema/messengers.py +++ b/lora_bridge/config/schema/messengers.py @@ -5,7 +5,7 @@ from __future__ import annotations -from typing import Annotated, Literal, Optional, Union +from typing import Annotated, Literal from pydantic import BaseModel, Field @@ -24,7 +24,7 @@ class BaseMessengerConfig(BaseModel): kind: str = Field( description="Тип мессенджера. Перекрывается ``Literal`` в подклассах.", ) - tag: Optional[str] = Field( + tag: str | None = Field( default=None, description=( "Переопределение тега источника в префиксе ``[тип:ник]`` при выгрузке " @@ -64,14 +64,14 @@ class TelegramMessengerConfig(BaseMessengerConfig): description="Тег дискриминатора — должно быть ``telegram``." ) token: str = Field(description="Telegram Bot API token, выданный BotFather.") - commands: Optional[TelegramCommandsConfig] = Field( + commands: TelegramCommandsConfig | None = Field( default=None, description="Блок команд; отсутствие или null отключает командный роутер.", ) MessengerConfig = Annotated[ - Union[TelegramMessengerConfig], # расширять Union при добавлении мессенджеров + TelegramMessengerConfig, # расширять Union при добавлении мессенджеров Field(discriminator="kind"), ] """Конфиг одного мессенджер-транспорта. diff --git a/lora_bridge/config/schema/rooms.py b/lora_bridge/config/schema/rooms.py index 14e0600..9ba7c3b 100644 --- a/lora_bridge/config/schema/rooms.py +++ b/lora_bridge/config/schema/rooms.py @@ -6,8 +6,6 @@ from __future__ import annotations -from typing import Optional, Union - from pydantic import BaseModel, ConfigDict, Field, model_validator from .ids import EndpointName, MessengerId, NodeId @@ -36,7 +34,7 @@ class MessengerSubscriber(BaseModel): "(узнаётся через @userinfobot или getUpdates)." ) ) - topic: Optional[str] = Field( + topic: str | None = Field( default=None, description=( "Тема (thread) внутри чата. Если опущена — работаем только с General. " @@ -53,7 +51,7 @@ class LoraSubscriber(BaseModel): lora: LoraRef = Field(description="LoRa-эндпоинт-получатель.") -Subscriber = Union[MessengerSubscriber, LoraSubscriber] +Subscriber = MessengerSubscriber | LoraSubscriber """Подписчик комнаты — либо чат мессенджера, либо peer LoRa-эндпоинт. Smart union без явного дискриминатора: pydantic выбирает форму по набору полей diff --git a/lora_bridge/core/bridge.py b/lora_bridge/core/bridge.py index a7d89f2..fff4958 100644 --- a/lora_bridge/core/bridge.py +++ b/lora_bridge/core/bridge.py @@ -12,10 +12,18 @@ import functools import logging from dataclasses import dataclass, replace -from typing import Optional, assert_never +from typing import assert_never import anyio +from ..domain.models import ( + ChannelRef, + DeliveryStatus, + LabelFormat, + Message, + RejectReason, +) +from ..domain.ports import AdmissionPolicy, Transport from .dedup import TtlDedup from .egress import RadioWorker from .journal import JournalEntry, OutboundJournal @@ -26,14 +34,6 @@ from .status import StatusDispatcher from .supervisor import Supervisor from .transform import build_lora_text, oversize_bytes, relay_lora_text -from ..domain.models import ( - ChannelRef, - DeliveryStatus, - LabelFormat, - Message, - RejectReason, -) -from ..domain.ports import AdmissionPolicy, Transport log = logging.getLogger(__name__) @@ -70,7 +70,7 @@ def __init__( rooms: RoomRegistry, status: StatusDispatcher, journal: OutboundJournal, - admission_policy: Optional[AdmissionPolicy] = None, + admission_policy: AdmissionPolicy | None = None, readonly_endpoints: frozenset[ChannelRef] = frozenset(), ) -> None: self._nodes = nodes @@ -111,7 +111,7 @@ async def run(self) -> None: for transport in reversed(all_transports): try: await transport.stop() - except Exception: # noqa: BLE001 — на shutdown не валим остальные stop() + except Exception: # на shutdown не валим остальные stop() log.exception("остановка транспорта '%s' упала", transport.id) def build_worker(self, node: NodeRuntime) -> RadioWorker: @@ -264,7 +264,7 @@ async def mirror_to_messenger(self, member: MessengerMember, msg: Message) -> No return try: await binding.transport.send(member.ref, msg) # best-effort, без статуса (A2) - except Exception: # noqa: BLE001 — миррор не должен валить поток + except Exception: # миррор не должен валить поток log.exception("миррор в %s не удался", member.ref) async def reject( diff --git a/lora_bridge/core/dedup.py b/lora_bridge/core/dedup.py index 15a46a3..7bec5de 100644 --- a/lora_bridge/core/dedup.py +++ b/lora_bridge/core/dedup.py @@ -7,10 +7,10 @@ from __future__ import annotations import time -from typing import Callable +from collections.abc import Callable -from .ttl_window import TtlWindow from ..domain.models import Message +from .ttl_window import TtlWindow class TtlDedup: diff --git a/lora_bridge/core/egress.py b/lora_bridge/core/egress.py index ebdb17e..1f17267 100644 --- a/lora_bridge/core/egress.py +++ b/lora_bridge/core/egress.py @@ -6,16 +6,16 @@ from __future__ import annotations -from typing import Awaitable, Callable +from collections.abc import Awaitable, Callable import anyio +from ..domain.models import DeliveryStatus, RejectReason, SendResult +from ..domain.ports import Transport from .journal import OutboundJournal from .loopguard import LoopGuard from .queue import CommitQueue, QueueItem from .status import StatusDispatcher -from ..domain.models import DeliveryStatus, RejectReason, SendResult -from ..domain.ports import Transport OnCommitted = Callable[[QueueItem], Awaitable[None]] OnReject = Callable[[QueueItem, RejectReason], Awaitable[None]] diff --git a/lora_bridge/core/journal.py b/lora_bridge/core/journal.py index a5a2628..d6c7776 100644 --- a/lora_bridge/core/journal.py +++ b/lora_bridge/core/journal.py @@ -8,8 +8,9 @@ from __future__ import annotations import time +from collections.abc import Callable from dataclasses import dataclass -from typing import Callable, Optional, Protocol +from typing import Protocol import aiosqlite @@ -41,7 +42,7 @@ class JournalEntry: target_endpoint: str status: DeliveryStatus enqueued_at: float - tx_started_at: Optional[float] + tx_started_at: float | None payload: str @@ -59,7 +60,7 @@ class SqliteJournal: def __init__(self, db_path: str, *, _clock: Callable[[], float] = time.monotonic) -> None: self._db_path = db_path self._clock = _clock - self._db: Optional[aiosqlite.Connection] = None + self._db: aiosqlite.Connection | None = None async def start(self) -> None: self._db = await aiosqlite.connect(self._db_path) diff --git a/lora_bridge/core/loopguard.py b/lora_bridge/core/loopguard.py index feec21f..b545fb1 100644 --- a/lora_bridge/core/loopguard.py +++ b/lora_bridge/core/loopguard.py @@ -7,10 +7,10 @@ from __future__ import annotations import time -from typing import Callable +from collections.abc import Callable -from .ttl_window import TtlWindow from ..domain.models import Message +from .ttl_window import TtlWindow class LoopGuard: diff --git a/lora_bridge/core/notifier.py b/lora_bridge/core/notifier.py index 8484436..b583c02 100644 --- a/lora_bridge/core/notifier.py +++ b/lora_bridge/core/notifier.py @@ -8,7 +8,8 @@ import time from collections import defaultdict -from typing import Awaitable, Callable, NamedTuple +from collections.abc import Awaitable, Callable +from typing import NamedTuple import anyio diff --git a/lora_bridge/core/queue.py b/lora_bridge/core/queue.py index 1386e8a..c79f042 100644 --- a/lora_bridge/core/queue.py +++ b/lora_bridge/core/queue.py @@ -8,14 +8,14 @@ from __future__ import annotations import time +from collections.abc import AsyncIterator, Callable from dataclasses import dataclass, field -from typing import AsyncIterator, Callable, Optional import anyio from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream -from .ratelimit import TokenBucket from ..domain.models import ChannelRef, Message, RateSpec +from .ratelimit import TokenBucket @dataclass @@ -39,7 +39,7 @@ class CommitQueue: def __init__( self, capacity: int, - rate: Optional[RateSpec], + rate: RateSpec | None, ttl_seconds: float, *, _clock: Callable[[], float] = time.monotonic, diff --git a/lora_bridge/core/ratelimit.py b/lora_bridge/core/ratelimit.py index 2cb94de..cc5e6bd 100644 --- a/lora_bridge/core/ratelimit.py +++ b/lora_bridge/core/ratelimit.py @@ -3,7 +3,7 @@ from __future__ import annotations import time -from typing import Callable +from collections.abc import Callable from ..domain.models import RateSpec diff --git a/lora_bridge/core/routing.py b/lora_bridge/core/routing.py index 1529b48..f447a4b 100644 --- a/lora_bridge/core/routing.py +++ b/lora_bridge/core/routing.py @@ -8,7 +8,6 @@ from __future__ import annotations from dataclasses import dataclass -from typing import Optional, Union from ..domain.models import ChannelRef, messenger_channel @@ -27,14 +26,14 @@ def ref(self) -> ChannelRef: class MessengerMember: transport_id: str chat: str - topic: Optional[str] = None + topic: str | None = None @property def ref(self) -> ChannelRef: return ChannelRef(self.transport_id, messenger_channel(self.chat, self.topic)) -Member = Union[LoraMember, MessengerMember] +Member = LoraMember | MessengerMember @dataclass(frozen=True) @@ -64,5 +63,5 @@ def __init__(self, routes: list[RoomRoute]) -> None: raise ValueError(f"эндпоинт {m.ref} состоит более чем в одной комнате") self._by_ref[m.ref] = route - def for_source(self, source: ChannelRef) -> Optional[RoomRoute]: + def for_source(self, source: ChannelRef) -> RoomRoute | None: return self._by_ref.get(source) diff --git a/lora_bridge/core/supervisor.py b/lora_bridge/core/supervisor.py index 5bc3ae8..c9070a9 100644 --- a/lora_bridge/core/supervisor.py +++ b/lora_bridge/core/supervisor.py @@ -9,7 +9,8 @@ from __future__ import annotations import logging -from typing import Callable, Coroutine, Any +from collections.abc import Callable, Coroutine +from typing import Any import anyio diff --git a/lora_bridge/core/ttl_window.py b/lora_bridge/core/ttl_window.py index 737f95e..6bbc44c 100644 --- a/lora_bridge/core/ttl_window.py +++ b/lora_bridge/core/ttl_window.py @@ -2,7 +2,7 @@ import time from collections import OrderedDict -from typing import Callable +from collections.abc import Callable class TtlWindow: diff --git a/lora_bridge/domain/models.py b/lora_bridge/domain/models.py index 593ee5c..a66bf09 100644 --- a/lora_bridge/domain/models.py +++ b/lora_bridge/domain/models.py @@ -10,7 +10,6 @@ import datetime as dt from dataclasses import dataclass from enum import Enum -from typing import Optional BRIDGE_TRANSPORT_UID = "__bridge__" LORA_SENDER_UID = "__lora__" # transport_uid для сообщений из эфира (нет отправителя на уровне протокола) @@ -22,7 +21,7 @@ class ChannelRef: channel: str # opaque id эндпоинта; топик — забота адаптера -def messenger_channel(chat: str, topic: Optional[str]) -> str: +def messenger_channel(chat: str, topic: str | None) -> str: """Канонический opaque ``ChannelRef.channel`` для мессенджер-эндпоинта. Единый контракт для RoomRegistry (ядро) и мессенджер-адаптера (транспорт) — @@ -44,8 +43,8 @@ class Message: sender: Identity text: str # Время источника, если извлекается. Для LoRa часто None — не выдумываем. - timestamp: Optional[dt.datetime] = None - origin_tag: Optional[str] = None # loop-guard (только LoRa-путь) + timestamp: dt.datetime | None = None + origin_tag: str | None = None # loop-guard (только LoRa-путь) class DeliveryStatus(Enum): @@ -74,7 +73,7 @@ class RateSpec: @dataclass(frozen=True) class Capabilities: max_text_bytes: int - egress_rate: Optional[RateSpec] = None + egress_rate: RateSpec | None = None supports_status_feedback: bool = False # умеет показать статус (реакция) emits_tx_done: bool = False # узел отдаёт TX-done (commit); у MeshCore False (§5.1) diff --git a/lora_bridge/domain/ports.py b/lora_bridge/domain/ports.py index 50dea0e..1f889bf 100644 --- a/lora_bridge/domain/ports.py +++ b/lora_bridge/domain/ports.py @@ -7,7 +7,8 @@ from __future__ import annotations from abc import ABC, abstractmethod -from typing import AsyncIterator, Optional, Protocol +from collections.abc import AsyncIterator +from typing import Protocol from .models import ( Capabilities, @@ -25,7 +26,7 @@ class AdmissionPolicy(Protocol): Возвращает ``None`` — допустить; ``RejectReason`` — отклонить с этой причиной. """ - async def check(self, msg: Message) -> Optional[RejectReason]: ... + async def check(self, msg: Message) -> RejectReason | None: ... class Transport(ABC): @@ -59,5 +60,5 @@ async def report_status( origin: ChannelRef, message_id: str, status: DeliveryStatus, - reason: Optional[RejectReason] = None, + reason: RejectReason | None = None, ) -> None: ... diff --git a/lora_bridge/transports/hub.py b/lora_bridge/transports/hub.py index 75f0511..4d12e97 100644 --- a/lora_bridge/transports/hub.py +++ b/lora_bridge/transports/hub.py @@ -7,8 +7,8 @@ from __future__ import annotations +from collections.abc import AsyncIterator, Iterator from contextlib import contextmanager -from typing import AsyncIterator, Iterator import anyio from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream diff --git a/lora_bridge/transports/meshcore/connection.py b/lora_bridge/transports/meshcore/connection.py index 1e2e9bc..2e50444 100644 --- a/lora_bridge/transports/meshcore/connection.py +++ b/lora_bridge/transports/meshcore/connection.py @@ -9,7 +9,8 @@ import logging import time -from typing import Awaitable, assert_never +from collections.abc import Awaitable +from typing import assert_never from meshcore import MeshCore from serial.tools import list_ports diff --git a/lora_bridge/transports/meshcore/mappers/__init__.py b/lora_bridge/transports/meshcore/mappers/__init__.py index 5448124..e912fa8 100644 --- a/lora_bridge/transports/meshcore/mappers/__init__.py +++ b/lora_bridge/transports/meshcore/mappers/__init__.py @@ -8,7 +8,8 @@ from __future__ import annotations -from typing import Iterable, assert_never +from collections.abc import Iterable +from typing import assert_never from ....config.schema import ( Endpoint, @@ -52,11 +53,11 @@ def collect_channel_names(handlers: Iterable[EndpointHandler]) -> frozenset[str] __all__ = [ "EndpointHandler", - "ResolveContext", - "PublicChannelHandler", "PrivateChannelHandler", + "PublicChannelHandler", + "ResolveContext", "RoomServerHandler", - "init_endpoint_handler", "collect_channel_names", + "init_endpoint_handler", "route_rx", ] diff --git a/lora_bridge/transports/meshcore/mappers/channel_util.py b/lora_bridge/transports/meshcore/mappers/channel_util.py index 9e26074..5fa86db 100644 --- a/lora_bridge/transports/meshcore/mappers/channel_util.py +++ b/lora_bridge/transports/meshcore/mappers/channel_util.py @@ -20,9 +20,9 @@ from meshcore import MeshCore from ....domain.models import ( + LORA_SENDER_UID, ChannelRef, Identity, - LORA_SENDER_UID, Message, ) diff --git a/lora_bridge/transports/meshcore/mappers/handler.py b/lora_bridge/transports/meshcore/mappers/handler.py index c2e0f0a..fb2d3a9 100644 --- a/lora_bridge/transports/meshcore/mappers/handler.py +++ b/lora_bridge/transports/meshcore/mappers/handler.py @@ -11,10 +11,12 @@ import logging from abc import ABC, abstractmethod +from collections.abc import Callable, Iterable from dataclasses import dataclass -from typing import Any, Callable, ClassVar, Iterable +from typing import Any, ClassVar -from meshcore import EventType as McEventType, MeshCore +from meshcore import EventType as McEventType +from meshcore import MeshCore from ....domain.models import Message diff --git a/lora_bridge/transports/meshcore/mappers/private.py b/lora_bridge/transports/meshcore/mappers/private.py index 21dbf32..06e3886 100644 --- a/lora_bridge/transports/meshcore/mappers/private.py +++ b/lora_bridge/transports/meshcore/mappers/private.py @@ -11,9 +11,9 @@ from meshcore import MeshCore -from . import channel_util -from .handler import AuthorResolver, EV_CHANNEL_MSG, EndpointHandler, ResolveContext from ....domain.models import Message +from . import channel_util +from .handler import EV_CHANNEL_MSG, AuthorResolver, EndpointHandler, ResolveContext @dataclass diff --git a/lora_bridge/transports/meshcore/mappers/public.py b/lora_bridge/transports/meshcore/mappers/public.py index 4cffbde..1757b12 100644 --- a/lora_bridge/transports/meshcore/mappers/public.py +++ b/lora_bridge/transports/meshcore/mappers/public.py @@ -11,9 +11,9 @@ from meshcore import MeshCore -from . import channel_util -from .handler import AuthorResolver, EV_CHANNEL_MSG, EndpointHandler, ResolveContext from ....domain.models import Message +from . import channel_util +from .handler import EV_CHANNEL_MSG, AuthorResolver, EndpointHandler, ResolveContext @dataclass diff --git a/lora_bridge/transports/meshcore/mappers/room_server.py b/lora_bridge/transports/meshcore/mappers/room_server.py index c8206f3..6f845f1 100644 --- a/lora_bridge/transports/meshcore/mappers/room_server.py +++ b/lora_bridge/transports/meshcore/mappers/room_server.py @@ -10,18 +10,20 @@ from __future__ import annotations import logging +from collections.abc import Awaitable, Callable from dataclasses import dataclass -from typing import Any, Awaitable, Callable, ClassVar +from typing import Any, ClassVar -from meshcore import EventType as McEventType, MeshCore +from meshcore import EventType as McEventType +from meshcore import MeshCore -from .handler import AuthorResolver, EV_CONTACT_MSG, EndpointHandler, ResolveContext from ....domain.models import ( + LORA_SENDER_UID, ChannelRef, Identity, - LORA_SENDER_UID, Message, ) +from .handler import EV_CONTACT_MSG, AuthorResolver, EndpointHandler, ResolveContext log = logging.getLogger(__name__) diff --git a/lora_bridge/transports/meshcore/transport.py b/lora_bridge/transports/meshcore/transport.py index b6525c3..6376b11 100644 --- a/lora_bridge/transports/meshcore/transport.py +++ b/lora_bridge/transports/meshcore/transport.py @@ -19,22 +19,13 @@ from __future__ import annotations import logging -from typing import Any, AsyncIterator, Optional, TYPE_CHECKING +from collections.abc import AsyncIterator +from typing import TYPE_CHECKING, Any import anyio -from meshcore import EventType as McEventType, MeshCore +from meshcore import EventType as McEventType +from meshcore import MeshCore -from . import connection -from .mappers import ( - EndpointHandler, - ResolveContext, - collect_channel_names, - init_endpoint_handler, - route_rx, -) -from .result import classify -from ..hub import Hub -from ...domain.ports import Transport from ...domain.models import ( Capabilities, ChannelRef, @@ -44,6 +35,17 @@ RejectReason, SendResult, ) +from ...domain.ports import Transport +from ..hub import Hub +from . import connection +from .mappers import ( + EndpointHandler, + ResolveContext, + collect_channel_names, + init_endpoint_handler, + route_rx, +) +from .result import classify if TYPE_CHECKING: from ...config.schema import MeshCoreNode @@ -133,7 +135,7 @@ async def run(self) -> None: await self.start() except anyio.get_cancelled_exc_class(): raise - except Exception as exc: + except Exception as exc: # noqa: BLE001 — реконнект: любой сбой драйвера = retry с backoff log.warning("нода '%s': реконнект не удался (%s), следующая попытка через %.0f с", self.id, exc, delay) delay = min(delay * 2, 60.0) # start() уже выставил свежий _disconnect_ev — сбрасываем, @@ -154,7 +156,7 @@ async def _teardown(self) -> None: if mc is not None: try: await mc.disconnect() # type: ignore[no-untyped-call] # verify; у метода либы нет аннотаций - except Exception: + except Exception: # noqa: BLE001, S110 — коннект уже мёртв, ошибки disconnect неинтересны pass def _signal_disconnect(self) -> None: @@ -186,7 +188,7 @@ async def send(self, target: ChannelRef, msg: Message) -> SendResult: self.id, target.channel, result.detail, res.payload, ) return result - except Exception as exc: # noqa: BLE001 + except Exception as exc: log.exception("MeshCore send в %s упал", target.channel) return SendResult.failure(str(exc)) @@ -222,6 +224,6 @@ async def report_status( origin: ChannelRef, message_id: str, status: DeliveryStatus, - reason: Optional[RejectReason] = None, + reason: RejectReason | None = None, ) -> None: return None # LoRa не показывает статус diff --git a/lora_bridge/transports/telegram/commands/framework.py b/lora_bridge/transports/telegram/commands/framework.py index 3fc2530..b86fde3 100644 --- a/lora_bridge/transports/telegram/commands/framework.py +++ b/lora_bridge/transports/telegram/commands/framework.py @@ -29,8 +29,7 @@ from aiogram import Router from aiogram.filters import Command -from aiogram.types import BotCommand -from aiogram.types import CallbackQuery +from aiogram.types import BotCommand, CallbackQuery from aiogram.types import Message as TgMessage from ..ephemeral import delete_after @@ -60,7 +59,7 @@ class CommandMeta: name: str description: str - min_role: "Role" + min_role: Role @dataclass(frozen=True) @@ -76,10 +75,10 @@ class CallbackSpec: prefix: str handler: CallbackHandler - min_role: "Role" + min_role: Role -def render_help(commands: list[CommandMeta], role: "Role | None" = None) -> str: +def render_help(commands: list[CommandMeta], role: Role | None = None) -> str: """Текст ``/help`` из переданного (уже отфильтрованного) реестра.""" lines = [f"/{spec.name} — {spec.description}" for spec in commands] header = "Доступные команды:" @@ -88,7 +87,7 @@ def render_help(commands: list[CommandMeta], role: "Role | None" = None) -> str: return header + "\n" + "\n".join(lines) -def command_menu(commands: list[CommandMeta], role: "Role") -> list[BotCommand]: +def command_menu(commands: list[CommandMeta], role: Role) -> list[BotCommand]: """Меню для ``Bot.set_my_commands`` — фильтрует по роли вызывающего.""" visible = [c for c in commands if c.min_role <= role] return [BotCommand(command=spec.name, description=spec.description) for spec in visible] @@ -97,7 +96,7 @@ def command_menu(commands: list[CommandMeta], role: "Role") -> list[BotCommand]: def build_command_router( transport_id: str, commands: list[CommandSpec], - store: "ModerationStore | None" = None, + store: ModerationStore | None = None, owner_id: int = 0, callbacks: list[CallbackSpec] | None = None, *, diff --git a/lora_bridge/transports/telegram/commands/handlers.py b/lora_bridge/transports/telegram/commands/handlers.py index 765a89a..7221ff5 100644 --- a/lora_bridge/transports/telegram/commands/handlers.py +++ b/lora_bridge/transports/telegram/commands/handlers.py @@ -11,8 +11,8 @@ from aiogram.types import Message as TgMessage -from .framework import CommandMeta, CommandSpec, render_help from ..moderation.roles import Role +from .framework import CommandMeta, CommandSpec, render_help if TYPE_CHECKING: from ..moderation.store import ModerationStore @@ -24,7 +24,7 @@ def make_basic_commands( - store: "ModerationStore", + store: ModerationStore, owner_id: int, all_metas: list[CommandMeta], ) -> list[CommandSpec]: diff --git a/lora_bridge/transports/telegram/commands/moderation.py b/lora_bridge/transports/telegram/commands/moderation.py index 08fb669..ab420e1 100644 --- a/lora_bridge/transports/telegram/commands/moderation.py +++ b/lora_bridge/transports/telegram/commands/moderation.py @@ -5,17 +5,19 @@ import math import time from collections.abc import Awaitable, Callable -from typing import Optional, TYPE_CHECKING +from typing import TYPE_CHECKING from aiogram.types import ( CallbackQuery, InlineKeyboardButton, InlineKeyboardMarkup, +) +from aiogram.types import ( Message as TgMessage, ) -from .framework import CallbackSpec, CommandMeta, CommandSpec from ..moderation.roles import Role +from .framework import CallbackSpec, CommandMeta, CommandSpec if TYPE_CHECKING: from ..moderation.store import ModerationStore @@ -34,7 +36,7 @@ _PAGE_SIZE = 10 -async def resolve_target(message: TgMessage) -> Optional[tuple[int, Optional[str]]]: +async def resolve_target(message: TgMessage) -> tuple[int, str | None] | None: """Reply или числовой аргумент → (tg_id, display_name|None). None = неопределимо.""" if message.reply_to_message and message.reply_to_message.from_user: u = message.reply_to_message.from_user @@ -52,13 +54,13 @@ async def resolve_target(message: TgMessage) -> Optional[tuple[int, Optional[str return None -def _mention(tg_id: int, name: Optional[str]) -> str: +def _mention(tg_id: int, name: str | None) -> str: label = html.escape(name) if name else str(tg_id) return f'{label}' async def _audit_text_and_kb( - page: int, store: "ModerationStore" + page: int, store: ModerationStore ) -> tuple[str, InlineKeyboardMarkup]: import datetime as _dt total = await store.count_audit_entries() @@ -68,7 +70,7 @@ async def _audit_text_and_kb( lines = [] for e in entries: - ts_str = _dt.datetime.utcfromtimestamp(e.ts).strftime("%Y-%m-%d %H:%M") + ts_str = _dt.datetime.fromtimestamp(e.ts, tz=_dt.UTC).strftime("%Y-%m-%d %H:%M") actor = _mention(e.actor_id, e.actor_name) target = _mention(e.target_id, e.target_name) if e.target_id else "" detail = f" [{html.escape(e.detail)}]" if e.detail else "" @@ -90,7 +92,7 @@ async def _audit_text_and_kb( return text, kb -async def _send_audit_page_edit(query: CallbackQuery, page: int, store: "ModerationStore") -> None: +async def _send_audit_page_edit(query: CallbackQuery, page: int, store: ModerationStore) -> None: text, kb = await _audit_text_and_kb(page, store) msg = query.message if isinstance(msg, TgMessage): @@ -98,7 +100,7 @@ async def _send_audit_page_edit(query: CallbackQuery, page: int, store: "Moderat def make_moderation_commands( - store: "ModerationStore", + store: ModerationStore, cfg: object, *, on_role_changed: Callable[[int, int], Awaitable[None]] | None = None, @@ -167,8 +169,8 @@ async def set_alias(message: TgMessage) -> None: actor_id = actor.id if actor else 0 target_id: int = actor_id - target_name: Optional[str] = None - new_alias: Optional[str] = None + target_name: str | None = None + new_alias: str | None = None if args: if message.reply_to_message and message.reply_to_message.from_user: @@ -304,7 +306,7 @@ async def audit(message: TgMessage) -> None: ] -def make_audit_callbacks(store: "ModerationStore") -> list[CallbackSpec]: +def make_audit_callbacks(store: ModerationStore) -> list[CallbackSpec]: """Фабрика CallbackSpec для пагинации /audit.""" async def audit_page(query: CallbackQuery) -> None: data = query.data or "" diff --git a/lora_bridge/transports/telegram/ephemeral.py b/lora_bridge/transports/telegram/ephemeral.py index f240d06..2856714 100644 --- a/lora_bridge/transports/telegram/ephemeral.py +++ b/lora_bridge/transports/telegram/ephemeral.py @@ -8,7 +8,7 @@ from aiogram.types import Message as TgMessage -async def delete_after(delay: float, *messages: "TgMessage") -> None: +async def delete_after(delay: float, *messages: TgMessage) -> None: await asyncio.sleep(delay) for msg in messages: with suppress(Exception): diff --git a/lora_bridge/transports/telegram/moderation/roles.py b/lora_bridge/transports/telegram/moderation/roles.py index 62c3656..03320a1 100644 --- a/lora_bridge/transports/telegram/moderation/roles.py +++ b/lora_bridge/transports/telegram/moderation/roles.py @@ -1,4 +1,5 @@ from __future__ import annotations + from enum import IntEnum diff --git a/lora_bridge/transports/telegram/moderation/store.py b/lora_bridge/transports/telegram/moderation/store.py index c82b9d1..71aabf2 100644 --- a/lora_bridge/transports/telegram/moderation/store.py +++ b/lora_bridge/transports/telegram/moderation/store.py @@ -1,7 +1,6 @@ from __future__ import annotations from dataclasses import dataclass -from typing import Optional import aiosqlite @@ -35,10 +34,10 @@ @dataclass(frozen=True) class UserSettings: - alias: Optional[str] = None + alias: str | None = None transliter: bool = False disabled: bool = False - banned_name: Optional[str] = None + banned_name: str | None = None @dataclass(frozen=True) @@ -46,17 +45,17 @@ class AuditEntry: id: int ts: int actor_id: int - actor_name: Optional[str] + actor_name: str | None action: str - target_id: Optional[int] - target_name: Optional[str] - detail: Optional[str] + target_id: int | None + target_name: str | None + detail: str | None class ModerationStore: def __init__(self, db_path: str) -> None: self._db_path = db_path - self._db: Optional[aiosqlite.Connection] = None + self._db: aiosqlite.Connection | None = None async def start(self) -> None: self._db = await aiosqlite.connect(self._db_path) @@ -85,7 +84,7 @@ async def get_role(self, owner_id: int, tg_id: int) -> Role: return Role.USER return Role.ADMIN if row[0] == "admin" else Role.MODERATOR - async def set_role(self, tg_id: int, role: str, chat_id: Optional[int] = None) -> None: + async def set_role(self, tg_id: int, role: str, chat_id: int | None = None) -> None: await self._conn.execute( "INSERT OR REPLACE INTO roles (tg_id, role, last_chat_id) VALUES (?,?,?)", (tg_id, role, chat_id), @@ -96,7 +95,7 @@ async def remove_role(self, tg_id: int) -> None: await self._conn.execute("DELETE FROM roles WHERE tg_id=?", (tg_id,)) await self._conn.commit() - async def get_all_privileged(self) -> list[tuple[int, str, Optional[int]]]: + async def get_all_privileged(self) -> list[tuple[int, str, int | None]]: cur = await self._conn.execute("SELECT tg_id, role, last_chat_id FROM roles") rows = await cur.fetchall() return [(r[0], r[1], r[2]) for r in rows] @@ -125,7 +124,7 @@ async def get_user_settings(self, tg_id: int) -> UserSettings: banned_name=row[3], ) - async def ban_user(self, tg_id: int, banned_name: Optional[str]) -> None: + async def ban_user(self, tg_id: int, banned_name: str | None) -> None: await self._conn.execute( "INSERT INTO user_settings (tg_id, disabled, banned_name) VALUES (?,1,?) " "ON CONFLICT(tg_id) DO UPDATE SET disabled=1, banned_name=excluded.banned_name", @@ -139,14 +138,14 @@ async def unban_user(self, tg_id: int) -> None: ) await self._conn.commit() - async def get_banned_users(self) -> list[tuple[int, Optional[str], Optional[str]]]: + async def get_banned_users(self) -> list[tuple[int, str | None, str | None]]: cur = await self._conn.execute( "SELECT tg_id, banned_name, alias FROM user_settings WHERE disabled=1" ) rows = await cur.fetchall() return [(r[0], r[1], r[2]) for r in rows] - async def set_alias(self, tg_id: int, alias: Optional[str]) -> None: + async def set_alias(self, tg_id: int, alias: str | None) -> None: await self._conn.execute( "INSERT INTO user_settings (tg_id, alias) VALUES (?,?) " "ON CONFLICT(tg_id) DO UPDATE SET alias=excluded.alias", @@ -174,11 +173,11 @@ async def log_action( self, ts: int, actor_id: int, - actor_name: Optional[str], + actor_name: str | None, action: str, - target_id: Optional[int] = None, - target_name: Optional[str] = None, - detail: Optional[str] = None, + target_id: int | None = None, + target_name: str | None = None, + detail: str | None = None, ) -> None: await self._conn.execute( "INSERT INTO audit_log " diff --git a/lora_bridge/transports/telegram/reactions.py b/lora_bridge/transports/telegram/reactions.py index de465a2..37d4aa2 100644 --- a/lora_bridge/transports/telegram/reactions.py +++ b/lora_bridge/transports/telegram/reactions.py @@ -13,13 +13,13 @@ import asyncio import logging -from typing import Optional from aiogram import Bot -from aiogram.types import Message as TgMessage, ReactionTypeEmoji, ReactionTypeUnion +from aiogram.types import Message as TgMessage +from aiogram.types import ReactionTypeEmoji, ReactionTypeUnion -from .ephemeral import delete_after from ...domain.models import DeliveryStatus, RejectReason +from .ephemeral import delete_after log = logging.getLogger(__name__) @@ -97,7 +97,7 @@ async def clear_now(self, key: tuple[int, str], bot: Bot) -> None: chat_id, message_id = key try: await bot.set_message_reaction(chat_id, int(message_id), reaction=[]) - except Exception: # noqa: BLE001 + except Exception: log.debug("clear реакции не удался для %s/%s", chat_id, message_id, exc_info=True) async def _delayed_apply( @@ -124,7 +124,7 @@ async def _delayed_apply( try: await bot.set_message_reaction(chat_id, int(message_id), reaction=reaction) self._sent.add(key) # реакция выставлена — теперь clear_now знает что чистить - except Exception: # noqa: BLE001 + except Exception: log.debug("set_message_reaction не удался для %s/%s", chat_id, message_id, exc_info=True) @@ -143,7 +143,7 @@ async def report( chat_id: int, message_id: str, status: DeliveryStatus, - reason: Optional[RejectReason] = None, + reason: RejectReason | None = None, ) -> None: if status == DeliveryStatus.TRANSMITTING: return # промежуточный — не трогаем реакцию @@ -160,7 +160,7 @@ async def report( if reaction: self._debouncer.schedule(key, reaction, self._bot) - async def report_disabled(self, message: "TgMessage") -> None: + async def report_disabled(self, message: TgMessage) -> None: """Реакция 🚫 на сообщение забаненного пользователя (best-effort).""" try: await self._bot.set_message_reaction( # verify @@ -168,10 +168,10 @@ async def report_disabled(self, message: "TgMessage") -> None: message.message_id, reaction=[ReactionTypeEmoji(emoji="🚫")], ) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001, S110 — реакция best-effort pass - async def report_alias_required(self, message: "TgMessage") -> None: + async def report_alias_required(self, message: TgMessage) -> None: """Реакция 🪪 на сообщение без alias, когда он обязателен (best-effort).""" try: await self._bot.set_message_reaction( @@ -179,11 +179,11 @@ async def report_alias_required(self, message: "TgMessage") -> None: message.message_id, reaction=[ReactionTypeEmoji(emoji="🪪")], ) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001, S110 — реакция best-effort pass async def send_expiring_reply( - self, message: "TgMessage", text: str, delay: float = ALIAS_REPLY_TTL_S + self, message: TgMessage, text: str, delay: float = ALIAS_REPLY_TTL_S ) -> None: """Reply, который сам удаляется через ``delay`` секунд. Исходное сообщение не трогаем.""" try: @@ -194,7 +194,7 @@ async def send_expiring_reply( @staticmethod def _reaction_for( - status: DeliveryStatus, reason: Optional[RejectReason] + status: DeliveryStatus, reason: RejectReason | None ) -> list[ReactionTypeUnion]: if status == DeliveryStatus.PENDING: return [ReactionTypeEmoji(emoji=_PENDING_EMOJI)] diff --git a/lora_bridge/transports/telegram/transport.py b/lora_bridge/transports/telegram/transport.py index 4415a59..73f25ce 100644 --- a/lora_bridge/transports/telegram/transport.py +++ b/lora_bridge/transports/telegram/transport.py @@ -12,25 +12,12 @@ import asyncio import html import logging -from typing import AsyncIterator, Optional, TYPE_CHECKING +from collections.abc import AsyncIterator +from typing import TYPE_CHECKING from aiogram import Bot, Dispatcher, F, Router from aiogram.types import Message as TgMessage -from .commands import ( - ALL_COMMAND_METAS, - build_command_router, - command_menu, - make_audit_callbacks, - make_basic_commands, - make_moderation_commands, -) -from .moderation.roles import Role -from .moderation.store import ModerationStore, UserSettings -from .moderation.transliterate import transliterate -from .reactions import ReactionFeedback -from ..hub import Hub -from ...domain.ports import Transport from ...domain.models import ( BRIDGE_TRANSPORT_UID, Capabilities, @@ -43,6 +30,20 @@ SendResult, messenger_channel, ) +from ...domain.ports import Transport +from ..hub import Hub +from .commands import ( + ALL_COMMAND_METAS, + build_command_router, + command_menu, + make_audit_callbacks, + make_basic_commands, + make_moderation_commands, +) +from .moderation.roles import Role +from .moderation.store import ModerationStore, UserSettings +from .moderation.transliterate import transliterate +from .reactions import ReactionFeedback if TYPE_CHECKING: from ...config.schema import TelegramMessengerConfig @@ -55,7 +56,7 @@ ) -def split_channel(channel: str) -> tuple[int, Optional[int]]: +def split_channel(channel: str) -> tuple[int, int | None]: """``"chat"`` / ``"chat#topic"`` → ``(chat_id, thread_id|None)``. Инверсия ``messenger_channel``. Декод нужен только на send-стороне адаптера (chat_id/thread_id для ``bot.send_message``), @@ -81,15 +82,15 @@ class TelegramTransport(Transport): def __init__( self, transport_id: str, - config: "TelegramMessengerConfig", + config: TelegramMessengerConfig, *, - _store: Optional[ModerationStore] = None, + _store: ModerationStore | None = None, ) -> None: self.id = transport_id self._hub = Hub() self._bot = Bot(config.token) self._dp = Dispatcher() - self._store: Optional[ModerationStore] = None + self._store: ModerationStore | None = None self._owner_id: int = 0 self._require_alias: bool = False # (tg_id, chat_id) — уже обновлённые scope; избегаем лишних API-вызовов @@ -147,7 +148,7 @@ async def _set_user_commands(self, tg_id: int, role: Role) -> None: try: await self._bot.set_my_commands(menu, scope=BotCommandScopeChat(chat_id=tg_id)) self._cmd_scope_done.add((tg_id, tg_id)) - except Exception: + except Exception: # noqa: BLE001 — меню команд best-effort, сбой не критичен log.debug("Не удалось обновить меню команд user=%d", tg_id) async def _clear_user_group_commands(self, tg_id: int, group_chat_id: int) -> None: @@ -157,7 +158,7 @@ async def _clear_user_group_commands(self, tg_id: int, group_chat_id: int) -> No await self._bot.set_my_commands( [], scope=BotCommandScopeChatMember(chat_id=group_chat_id, user_id=tg_id) ) - except Exception: + except Exception: # noqa: BLE001 — меню команд best-effort, сбой не критичен log.debug("Не удалось скрыть меню в группе user=%d chat=%d", tg_id, group_chat_id) async def start(self) -> None: @@ -208,7 +209,7 @@ async def on_message(self, message: TgMessage) -> None: await self._hub.publish(self.normalize(message, settings)) def normalize( - self, message: TgMessage, settings: Optional[UserSettings] = None + self, message: TgMessage, settings: UserSettings | None = None ) -> Message: thread = message.message_thread_id chat_id = str(message.chat.id) @@ -252,7 +253,7 @@ async def send(self, target: ChannelRef, msg: Message) -> SendResult: chat_id, text, message_thread_id=thread_id, parse_mode=parse_mode ) return SendResult.success() - except Exception as exc: # noqa: BLE001 + except Exception as exc: log.exception("Telegram send в %s упал", target.channel) return SendResult.failure(str(exc)) @@ -264,7 +265,7 @@ async def report_status( origin: ChannelRef, message_id: str, status: DeliveryStatus, - reason: Optional[RejectReason] = None, + reason: RejectReason | None = None, ) -> None: chat_id, _ = split_channel(origin.channel) await self._reactions.report(chat_id, message_id, status, reason) diff --git a/lora_bridge/wiring.py b/lora_bridge/wiring.py index 883ddd6..22d99d4 100644 --- a/lora_bridge/wiring.py +++ b/lora_bridge/wiring.py @@ -13,14 +13,14 @@ from .config.schema import ( AppConfig, LoraSubscriber, - MessengerSubscriber, MeshCoreNode, MessengerConfig, + MessengerSubscriber, ) from .core.bridge import MessengerBinding, NodeRuntime -from .core.notifier import DropNotifier, NotifySink from .core.dedup import TtlDedup from .core.loopguard import LoopGuard +from .core.notifier import DropNotifier, NotifySink from .core.queue import CommitQueue from .core.routing import LoraMember, Member, MessengerMember, RoomRegistry, RoomRoute from .domain.models import ChannelRef, Identity, LabelFormat, Message, RateSpec diff --git a/pyproject.toml b/pyproject.toml index cf6401d..77516a8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -22,7 +22,7 @@ dependencies = [ dev = [ "pytest>=8.0", "pytest-anyio>=0.0", - "ruff>=0.4,<0.16", # 0.16.0 расширил дефолтные правила (~170 ошибок); снятие пина — LoRa-Bridge-9qw + "ruff>=0.16,<0.17", # минорные релизы ruff меняют дефолтные правила — бамп границы осознанным PR "mypy>=1.14", # per-module follow_untyped_imports (см. [tool.mypy] ниже) — добавлен в 1.14.0 "types-pyserial", ] diff --git a/tests/helpers/fakes.py b/tests/helpers/fakes.py index 47fe827..b55e13d 100644 --- a/tests/helpers/fakes.py +++ b/tests/helpers/fakes.py @@ -2,7 +2,7 @@ from __future__ import annotations -from typing import AsyncIterator, Optional +from collections.abc import AsyncIterator import anyio @@ -55,7 +55,7 @@ def __init__( self.capabilities = capabilities self._hub = Hub() self.sent: list[tuple[ChannelRef, Message]] = [] - self.statuses: list[tuple[str, DeliveryStatus, Optional[RejectReason]]] = [] + self.statuses: list[tuple[str, DeliveryStatus, RejectReason | None]] = [] self.started = False self._fail = fail self._busy_left = busy_times @@ -89,7 +89,7 @@ async def report_status( origin: ChannelRef, message_id: str, status: DeliveryStatus, - reason: Optional[RejectReason] = None, + reason: RejectReason | None = None, ) -> None: self.statuses.append((message_id, status, reason)) diff --git a/tests/test_config_descriptions.py b/tests/test_config_descriptions.py index 4eb5aa0..70df2f9 100644 --- a/tests/test_config_descriptions.py +++ b/tests/test_config_descriptions.py @@ -39,7 +39,6 @@ UsbConnection, ) - # --------------------------------------------------------------------------- # Покрытие описаниями # --------------------------------------------------------------------------- diff --git a/tests/test_config_errors_unit.py b/tests/test_config_errors_unit.py index 6b1a3c1..d96cad1 100644 --- a/tests/test_config_errors_unit.py +++ b/tests/test_config_errors_unit.py @@ -25,9 +25,9 @@ _pretty_type, _resolve_models, _smart_union_variant_index, - _strip_annotated, format_validation_error, ) +from lora_bridge.config.introspect import strip_annotated from lora_bridge.config.schema import ( AppConfig, MessengerSubscriber, @@ -39,7 +39,6 @@ UsbConnection, ) - # =========================================================================== # 1. Чистые хелперы # =========================================================================== @@ -122,7 +121,7 @@ def test_pretty_type_renders_russian_label(t, expected): ], ) def test_strip_annotated_unwraps_to_raw_type(t, expected): - assert _strip_annotated(t) is expected + assert strip_annotated(t) is expected @pytest.mark.parametrize( @@ -342,7 +341,7 @@ class _M(BaseModel): try: _M.model_validate({"x": "not an int"}) except ValidationError as exc: - err = {**list(exc.errors())[0], "type": "totally_unknown_pydantic_kind"} + err = {**next(iter(exc.errors())), "type": "totally_unknown_pydantic_kind"} rendered = "\n".join(_format_one(err, 1)) assert rendered.startswith("1. ") # сырое сообщение pydantic пробрасывается как есть diff --git a/tests/test_config_schema.py b/tests/test_config_schema.py index ac53217..a29a979 100644 --- a/tests/test_config_schema.py +++ b/tests/test_config_schema.py @@ -14,8 +14,8 @@ ConnectionBase, Endpoint, EndpointBase, - MessengerConfig, MeshCoreNode, + MessengerConfig, TelegramCommandsConfig, TelegramMessengerConfig, ) diff --git a/tests/test_pipeline.py b/tests/test_pipeline.py index de6e7d1..c545c21 100644 --- a/tests/test_pipeline.py +++ b/tests/test_pipeline.py @@ -15,16 +15,16 @@ from lora_bridge.core.routing import LoraMember, MessengerMember, RoomRegistry, RoomRoute from lora_bridge.core.status import StatusDispatcher from lora_bridge.domain.models import ( + LORA_SENDER_UID, ChannelRef, DeliveryStatus, Identity, - LORA_SENDER_UID, LabelFormat, Message, RateSpec, RejectReason, ) -from tests.helpers.fakes import FakeTransport, LORA_CAPS, MSG_CAPS +from tests.helpers.fakes import LORA_CAPS, MSG_CAPS, FakeTransport pytestmark = pytest.mark.anyio @@ -54,9 +54,12 @@ def _msg(transport_id: str, channel: str, text: str, mid: str = "m1") -> Message ) +_DEFAULT_RATE = RateSpec(100, 60) + + async def _build( routes, nodes_transports, messengers, *, - capacity=16, rate=RateSpec(100, 60), readonly_endpoints=frozenset(), + capacity=16, rate=_DEFAULT_RATE, readonly_endpoints=frozenset(), ): journal = SqliteJournal(":memory:") await journal.start() @@ -213,7 +216,7 @@ async def test_rate_limit_rejected(): m1 = FakeTransport("tg", MSG_CAPS) room = RoomRoute(members=(LoraMember("n1", "emergency"), MessengerMember("tg", "-100", None))) # ёмкость 1, бёрст 1 → второе сообщение отвергается - bridge, nodes, notices = await _build( + bridge, _, _ = await _build( [room], {"n1": lora}, {"tg": m1}, capacity=1, rate=RateSpec(1, 60, burst=1) ) diff --git a/tests/test_reconnect.py b/tests/test_reconnect.py index d9544ca..777bdb6 100644 --- a/tests/test_reconnect.py +++ b/tests/test_reconnect.py @@ -10,8 +10,7 @@ import anyio -from lora_bridge.transports.meshcore.transport import MeshCoreTransport, EV_DISCONNECTED - +from lora_bridge.transports.meshcore.transport import EV_DISCONNECTED, MeshCoreTransport # --------------------------------------------------------------------------- # Helpers @@ -62,7 +61,7 @@ async def test_signal_disconnect_idempotent(): # --------------------------------------------------------------------------- async def test_send_returns_overloaded_when_mc_is_none(): - from lora_bridge.domain.models import ChannelRef, Message, Identity + from lora_bridge.domain.models import ChannelRef, Identity, Message t = _make_transport() assert t._mc is None # начальное состояние target = ChannelRef("test", "ep") diff --git a/tests/test_shutdown.py b/tests/test_shutdown.py index 37cd3bb..30bb590 100644 --- a/tests/test_shutdown.py +++ b/tests/test_shutdown.py @@ -19,7 +19,7 @@ from lora_bridge.core.routing import RoomRegistry from lora_bridge.core.status import StatusDispatcher from lora_bridge.domain.models import LabelFormat, RateSpec -from tests.helpers.fakes import FakeTransport, LORA_CAPS, MSG_CAPS +from tests.helpers.fakes import LORA_CAPS, MSG_CAPS, FakeTransport pytestmark = pytest.mark.anyio diff --git a/tests/test_telegram_commands.py b/tests/test_telegram_commands.py index 91f4f73..236fa0c 100644 --- a/tests/test_telegram_commands.py +++ b/tests/test_telegram_commands.py @@ -50,7 +50,7 @@ def _update(text: str, user_id: int = 2, chat_type: str = "private", chat_id: in update_id=1, message=Message( message_id=10, - date=dt.datetime(2024, 1, 1), + date=dt.datetime(2024, 1, 1, tzinfo=dt.UTC), chat=Chat(id=chat_id, type=chat_type), from_user=User(id=user_id, is_bot=False, first_name="tester"), text=text, diff --git a/tests/test_telegram_moderation_commands.py b/tests/test_telegram_moderation_commands.py index 82493c7..17da4ed 100644 --- a/tests/test_telegram_moderation_commands.py +++ b/tests/test_telegram_moderation_commands.py @@ -2,11 +2,10 @@ from __future__ import annotations import datetime as dt +from collections.abc import AsyncGenerator from types import SimpleNamespace from unittest.mock import AsyncMock, patch -from collections.abc import AsyncGenerator - import pytest from aiogram.types import Chat, Message, User @@ -32,14 +31,14 @@ def _msg(text: str, user_id: int = 10, reply_user_id: int | None = None) -> Mess if reply_user_id is not None: reply = Message( message_id=5, - date=dt.datetime(2024, 1, 1), + date=dt.datetime(2024, 1, 1, tzinfo=dt.UTC), chat=Chat(id=1, type="group"), from_user=User(id=reply_user_id, is_bot=False, first_name="Target"), text="some text", ) return Message( message_id=10, - date=dt.datetime(2024, 1, 1), + date=dt.datetime(2024, 1, 1, tzinfo=dt.UTC), chat=Chat(id=1, type="group"), from_user=User(id=user_id, is_bot=False, first_name="Actor"), text=text, diff --git a/tests/test_telegram_moderation_store.py b/tests/test_telegram_moderation_store.py index ed8f25a..4f59d1c 100644 --- a/tests/test_telegram_moderation_store.py +++ b/tests/test_telegram_moderation_store.py @@ -2,9 +2,11 @@ from collections.abc import AsyncGenerator import pytest + from lora_bridge.transports.telegram.moderation.roles import Role from lora_bridge.transports.telegram.moderation.store import ModerationStore + @pytest.fixture async def store() -> AsyncGenerator[ModerationStore, None]: s = ModerationStore(":memory:") diff --git a/tests/test_telegram_reactions.py b/tests/test_telegram_reactions.py index 9761a03..59b92f2 100644 --- a/tests/test_telegram_reactions.py +++ b/tests/test_telegram_reactions.py @@ -18,8 +18,8 @@ from lora_bridge.domain.models import DeliveryStatus, RejectReason from lora_bridge.transports.telegram.reactions import ( - REJECT_EMOJI, _PENDING_EMOJI, + REJECT_EMOJI, ReactionDebouncer, ReactionFeedback, ) diff --git a/tests/test_telegram_require_alias.py b/tests/test_telegram_require_alias.py index 3c14974..8c52e6f 100644 --- a/tests/test_telegram_require_alias.py +++ b/tests/test_telegram_require_alias.py @@ -38,7 +38,7 @@ def _group_update(text: str, user_id: int = _USER_ID) -> Update: update_id=1, message=Message( message_id=10, - date=dt.datetime(2024, 1, 1), + date=dt.datetime(2024, 1, 1, tzinfo=dt.UTC), chat=Chat(id=_GROUP_CHAT_ID, type="supergroup"), from_user=User(id=user_id, is_bot=False, first_name="tester"), text=text, diff --git a/tests/test_telegram_send_format.py b/tests/test_telegram_send_format.py index 25093c0..9facaef 100644 --- a/tests/test_telegram_send_format.py +++ b/tests/test_telegram_send_format.py @@ -12,9 +12,9 @@ from lora_bridge.domain.models import ( BRIDGE_TRANSPORT_UID, + LORA_SENDER_UID, ChannelRef, Identity, - LORA_SENDER_UID, Message, messenger_channel, ) diff --git a/tests/test_ttl_scenarios.py b/tests/test_ttl_scenarios.py index 241d5d4..2e70e25 100644 --- a/tests/test_ttl_scenarios.py +++ b/tests/test_ttl_scenarios.py @@ -13,6 +13,7 @@ import time import anyio +import anyio.lowlevel import pytest import lora_bridge.core.egress as egress_mod @@ -33,7 +34,7 @@ RateSpec, RejectReason, ) -from tests.helpers.fakes import FakeClock, FakeTransport, LORA_CAPS, MSG_CAPS +from tests.helpers.fakes import LORA_CAPS, MSG_CAPS, FakeClock, FakeTransport pytestmark = pytest.mark.anyio @@ -77,13 +78,16 @@ def lora_msg(text: str, mid: str = "l1", *, origin_tag: str | None = None) -> Me ) +_DEFAULT_RATE = RateSpec(100, 60) + + async def build_bridge( lora: FakeTransport, messenger: FakeTransport, *, queue_ttl: float = 45.0, commit_timeout: float = 5.0, - rate: RateSpec = RateSpec(100, 60), + rate: RateSpec = _DEFAULT_RATE, notify_window: float = 60.0, notify_clock=None, ) -> tuple[Bridge, NodeRuntime, list]: @@ -218,16 +222,16 @@ async def test_dedup_drops_duplicate_lora_message(): """Два одинаковых LoRa-сообщения подряд: второе silently отбрасывается dedup.""" lora = FakeTransport("n1", LORA_CAPS) tg = FakeTransport("tg", MSG_CAPS) - bridge, node, _ = await build_bridge(lora, tg) + bridge, _, _ = await build_bridge(lora, tg) m = lora_msg("mesh broadcast", mid="same-id") async with anyio.create_task_group() as tg_scope: tg_scope.start_soon(bridge.consume, lora) - await anyio.sleep(0) # даём consume() запуститься и подписаться на hub + await anyio.lowlevel.checkpoint() # даём consume() запуститься и подписаться на hub await lora.inject(m) await lora.inject(m) # дубль — должен быть проглочен dedup - await anyio.sleep(0) # даём consume() обработать оба сообщения + await anyio.lowlevel.checkpoint() # даём consume() обработать оба сообщения tg_scope.cancel_scope.cancel() assert len(tg.sent) == 1, "дубль не должен зеркалироваться в мессенджер" @@ -249,9 +253,9 @@ async def test_loopguard_suppresses_own_echo(): async with anyio.create_task_group() as tg_scope: tg_scope.start_soon(bridge.consume, lora) - await anyio.sleep(0) # даём consume() запуститься и подписаться на hub + await anyio.lowlevel.checkpoint() # даём consume() запуститься и подписаться на hub await lora.inject(lora_msg(sent_text)) # то же самое вернулось назад из эфира - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg_scope.cancel_scope.cancel() assert tg.sent == [], "эхо собственной передачи не должно уходить в мессенджер" @@ -267,9 +271,9 @@ async def test_loopguard_passes_other_messages(): async with anyio.create_task_group() as tg_scope: tg_scope.start_soon(bridge.consume, lora) - await anyio.sleep(0) # даём consume() запуститься и подписаться на hub + await anyio.lowlevel.checkpoint() # даём consume() запуститься и подписаться на hub await lora.inject(lora_msg("совершенно другое сообщение")) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg_scope.cancel_scope.cancel() assert len(tg.sent) == 1 diff --git a/tests/test_wiring_integration.py b/tests/test_wiring_integration.py index 283d18e..24aa736 100644 --- a/tests/test_wiring_integration.py +++ b/tests/test_wiring_integration.py @@ -12,9 +12,10 @@ from __future__ import annotations -import yaml -import pytest import anyio +import anyio.lowlevel +import pytest +import yaml import lora_bridge.wiring as wiring_mod from lora_bridge.config.schema import AppConfig @@ -22,7 +23,12 @@ from lora_bridge.core.journal import SqliteJournal from lora_bridge.core.status import StatusDispatcher from lora_bridge.domain.models import ( - BRIDGE_TRANSPORT_UID, ChannelRef, DeliveryStatus, Identity, Message, RejectReason, + BRIDGE_TRANSPORT_UID, + ChannelRef, + DeliveryStatus, + Identity, + Message, + RejectReason, ) from lora_bridge.wiring import ( build_lora_nodes, @@ -31,7 +37,7 @@ build_readonly_endpoints, build_rooms, ) -from tests.helpers.fakes import FakeTransport, LORA_CAPS, MSG_CAPS +from tests.helpers.fakes import LORA_CAPS, MSG_CAPS, FakeTransport _NOTICE_SENDER = Identity(display_name="bridge", transport_uid=BRIDGE_TRANSPORT_UID) @@ -437,9 +443,9 @@ async def test_lora_to_tg_full_path(wire_fakes): async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_transport) - await anyio.sleep(0) # даём consume() подписаться на hub + await anyio.lowlevel.checkpoint() # даём consume() подписаться на hub await lora_transport.inject(incoming) - await anyio.sleep(0) # даём consume() обработать + await anyio.lowlevel.checkpoint() # даём consume() обработать tg.cancel_scope.cancel() await journal.stop() @@ -468,9 +474,9 @@ async def test_lora_to_lora_full_path(wire_fakes): async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() await lora_fakes["mc-1"].inject(incoming) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg.cancel_scope.cancel() node2 = runtimes["mc-2"] @@ -490,7 +496,7 @@ async def test_lora_to_lora_full_path(wire_fakes): async def test_rate_limit_from_config(wire_fakes): """egress_rate: 1 msg/60s в конфиге → второе сообщение отклоняется с RATE_LIMIT.""" _, tg_fakes = wire_fakes - bridge, runtimes, journal = await assemble(_CFG_TIGHT_RATE) + bridge, _, journal = await assemble(_CFG_TIGHT_RATE) def _msg(mid): return Message( @@ -512,7 +518,7 @@ def _msg(mid): async def test_too_long_message_rejected(wire_fakes): """Сообщение длиннее 150 байт отклоняется с TOO_LONG до попадания в очередь.""" lora_fakes, tg_fakes = wire_fakes - bridge, runtimes, journal = await assemble(_CFG_ONE_ROOM) + bridge, _, journal = await assemble(_CFG_ONE_ROOM) await bridge.admit( Message( @@ -532,7 +538,7 @@ async def test_too_long_message_rejected(wire_fakes): async def test_read_only_endpoint_from_config_rejects_post(wire_fakes): """endpoints..read_only: true в YAML → постинг из мессенджера отклоняется с READONLY.""" lora_fakes, tg_fakes = wire_fakes - bridge, runtimes, journal = await assemble(_CFG_READONLY) + bridge, _, journal = await assemble(_CFG_READONLY) await bridge.admit( Message( @@ -577,7 +583,7 @@ async def test_short_ttl_from_config_expires_message(wire_fakes): async def test_message_from_unconfigured_chat_is_dropped(wire_fakes): """Сообщение из чата вне rooms тихо игнорируется — без статуса, без очереди.""" lora_fakes, tg_fakes = wire_fakes - bridge, runtimes, journal = await assemble(_CFG_ONE_ROOM) + bridge, _, journal = await assemble(_CFG_ONE_ROOM) await bridge.admit( Message( @@ -607,9 +613,9 @@ async def test_lora_message_to_unconfigured_endpoint_is_dropped(wire_fakes): async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() await lora_fakes["mc-1"].inject(ghost_msg) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg.cancel_scope.cancel() await journal.stop() @@ -739,9 +745,9 @@ async def test_lora_rx_mirrors_to_both_tg_chats(wire_fakes): async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) # даём consume() подписаться на hub + await anyio.lowlevel.checkpoint() # даём consume() подписаться на hub await lora_fakes["mc-1"].inject(incoming) - await anyio.sleep(0) # даём consume() обработать + await anyio.lowlevel.checkpoint() # даём consume() обработать tg.cancel_scope.cancel() await journal.stop() @@ -779,9 +785,9 @@ async def send_failing_first_chat(target, msg): try: async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() await lora_fakes["mc-1"].inject(incoming) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg.cancel_scope.cancel() finally: await journal.stop() @@ -829,10 +835,10 @@ async def test_dedup_same_lora_message_twice(wire_fakes): async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() await lora_fakes["mc-1"].inject(msg) await lora_fakes["mc-1"].inject(msg) # тот же текст → dedup - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg.cancel_scope.cancel() await journal.stop() @@ -870,9 +876,9 @@ async def test_loopguard_suppresses_echo_through_wiring(wire_fakes): ) async with anyio.create_task_group() as tg: tg.start_soon(bridge.consume, lora_fakes["mc-1"]) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() await lora_fakes["mc-1"].inject(echo) - await anyio.sleep(0) + await anyio.lowlevel.checkpoint() tg.cancel_scope.cancel() await journal.stop()