From 14593dc255ae8cfadfbd41c0c27c596425826322 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Wed, 7 Oct 2026 19:15:46 +0300 Subject: [PATCH] fix: annotate context managers with Generator for ty 0.0.85 --- benchmarks/probes.py | 6 +++--- faststream_outbox/subscriber/usecase.py | 6 +++--- faststream_outbox/testing.py | 6 +++--- tests/test_fastapi.py | 6 +++--- tests/test_unit.py | 9 +++++---- 5 files changed, 17 insertions(+), 16 deletions(-) diff --git a/benchmarks/probes.py b/benchmarks/probes.py index 9b4cb9d..61069f1 100644 --- a/benchmarks/probes.py +++ b/benchmarks/probes.py @@ -24,7 +24,7 @@ import contextlib import dataclasses -from collections.abc import AsyncIterator +from collections.abc import AsyncGenerator from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncConnection, AsyncEngine @@ -155,7 +155,7 @@ class ProbeResult: @contextlib.asynccontextmanager -async def _autocommit(engine: AsyncEngine) -> AsyncIterator[AsyncConnection]: +async def _autocommit(engine: AsyncEngine) -> AsyncGenerator[AsyncConnection]: """Open a connection that emits no implicit BEGIN/COMMIT/ROLLBACK. pg_stat_statements tracks utility statements, and none of BEGIN/COMMIT/ROLLBACK @@ -200,7 +200,7 @@ async def assert_owns_database(engine: AsyncEngine) -> None: @contextlib.asynccontextmanager -async def probe(engine: AsyncEngine, table_name: str, schema: str | None = None) -> AsyncIterator[list[ProbeResult]]: +async def probe(engine: AsyncEngine, table_name: str, schema: str | None = None) -> AsyncGenerator[list[ProbeResult]]: """Snapshot the catalogs around the wrapped workload. Yields a list that holds exactly one :class:`ProbeResult` once the block exits. diff --git a/faststream_outbox/subscriber/usecase.py b/faststream_outbox/subscriber/usecase.py index 6dea8f6..a6c6e22 100644 --- a/faststream_outbox/subscriber/usecase.py +++ b/faststream_outbox/subscriber/usecase.py @@ -26,7 +26,7 @@ import random import time import typing -from collections.abc import AsyncIterator, Awaitable, Callable, Mapping, Sequence +from collections.abc import AsyncGenerator, AsyncIterator, Awaitable, Callable, Mapping, Sequence from contextlib import AbstractAsyncContextManager, AsyncExitStack, asynccontextmanager, suppress from itertools import chain @@ -348,7 +348,7 @@ async def _fetch_loop(self) -> None: async def _open_fetch_resources( self, engine: "AsyncEngine | None", - ) -> AsyncIterator[Mapping[str, object]]: + ) -> AsyncGenerator[Mapping[str, object]]: """Yield the kwargs ``_fetch_inner`` needs, owning fetch_conn + listen_conn lifetimes. Production path opens a long-lived ``AsyncConnection`` for the fetch CTE and a @@ -513,7 +513,7 @@ async def _worker_loop(self) -> None: async def _open_worker_resources( self, engine: "AsyncEngine | None", - ) -> AsyncIterator[Mapping[str, object]]: + ) -> AsyncGenerator[Mapping[str, object]]: """Yield ``writer_conn`` for ``_worker_inner``, owning its lifetime across all flushes. One long-lived ``AsyncConnection`` per outer reconnect cycle — every terminal/retry diff --git a/faststream_outbox/testing.py b/faststream_outbox/testing.py index 1210923..dd40c58 100644 --- a/faststream_outbox/testing.py +++ b/faststream_outbox/testing.py @@ -32,7 +32,7 @@ if typing.TYPE_CHECKING: - from collections.abc import Iterator, Sequence + from collections.abc import Generator, Sequence from sqlalchemy.ext.asyncio import AsyncSession @@ -620,7 +620,7 @@ def feed( return row_id @contextmanager - def _patch_producer(self, broker: OutboxBroker) -> "Iterator[None]": + def _patch_producer(self, broker: OutboxBroker) -> "Generator[None]": # Swap the broker's producer slot for one that routes inserts through # the in-memory fake client. ``OutboxPublisher.publish`` flows through # ``_basic_publish(cmd, producer=...)``, so replacing the producer is @@ -644,7 +644,7 @@ def _patch_producer(self, broker: OutboxBroker) -> "Iterator[None]": broker.config.broker_config.producer = original_producer @contextmanager - def _patch_broker(self, broker: OutboxBroker) -> "Iterator[None]": + def _patch_broker(self, broker: OutboxBroker) -> "Generator[None]": original_client = broker.config.broker_config.client broker.config.broker_config.client = self.fake_client # Mirror real publish's serializer wiring so pydantic / dataclass bodies diff --git a/tests/test_fastapi.py b/tests/test_fastapi.py index ce8abb1..e97258d 100644 --- a/tests/test_fastapi.py +++ b/tests/test_fastapi.py @@ -8,7 +8,7 @@ yielding a session-shaped mock from a normal FastAPI dependency. """ -from collections.abc import AsyncIterator, Mapping +from collections.abc import AsyncGenerator, AsyncIterator, Mapping from contextlib import asynccontextmanager from typing import Any from unittest.mock import AsyncMock @@ -48,7 +48,7 @@ def _make_app_with_router(router: OutboxRouter) -> FastAPI: """Build a FastAPI app mounted with the router; wrap the broker in TestOutboxBroker.""" @asynccontextmanager - async def lifespan(app: FastAPI) -> AsyncIterator[None]: + async def lifespan(app: FastAPI) -> AsyncGenerator[None]: del app # Swap in the in-memory fake client for the broker the router owns, # then let the router's own lifespan run (it starts subscribers). @@ -176,7 +176,7 @@ async def handle( test_broker = TestOutboxBroker(router.broker) @asynccontextmanager - async def lifespan(app: FastAPI) -> AsyncIterator[None]: + async def lifespan(app: FastAPI) -> AsyncGenerator[None]: del app async with test_broker: yield diff --git a/tests/test_unit.py b/tests/test_unit.py index c9fe7a9..e5a9097 100644 --- a/tests/test_unit.py +++ b/tests/test_unit.py @@ -8,6 +8,7 @@ import typing import uuid import warnings +from collections.abc import AsyncGenerator from contextlib import asynccontextmanager from unittest.mock import AsyncMock, MagicMock, patch @@ -2896,7 +2897,7 @@ async def test_fetch_reconnect_loop_exits_on_drain_without_churning() -> None: opens = {"n": 0} @asynccontextmanager - async def _open(_engine: object) -> typing.AsyncIterator[dict[str, object]]: + async def _open(_engine: object) -> AsyncGenerator[dict[str, object]]: opens["n"] += 1 # pragma: no cover - the drain guard exits before resources open yield {} # pragma: no cover - the drain guard exits before resources open @@ -2940,7 +2941,7 @@ def _capture_backoff(attempt: int, ceiling: float, *, base: float = 1.0) -> floa monkeypatch.setattr("faststream_outbox.subscriber.usecase._BACKOFF_RESET_THRESHOLD_SECONDS", 0.0) @asynccontextmanager - async def _open(_engine: object) -> typing.AsyncIterator[dict[str, object]]: + async def _open(_engine: object) -> AsyncGenerator[dict[str, object]]: yield {} async with test_broker: @@ -2988,7 +2989,7 @@ def _capture_backoff(attempt: int, ceiling: float, *, base: float = 1.0) -> floa sub.running = True @asynccontextmanager - async def _failing_open(_engine: object) -> typing.AsyncIterator[dict[str, object]]: + async def _failing_open(_engine: object) -> AsyncGenerator[dict[str, object]]: if len(attempts) >= 3: # exit cleanly after observing 3 escalating attempts sub.running = False yield {} @@ -3732,7 +3733,7 @@ async def test_fetch_cte_carries_partial_index_predicates_as_conjuncts() -> None class _CapturingConn: def begin(self) -> object: @asynccontextmanager - async def _cm() -> typing.AsyncIterator[None]: + async def _cm() -> AsyncGenerator[None]: yield return _cm()