Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions benchmarks/probes.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
6 changes: 3 additions & 3 deletions faststream_outbox/subscriber/usecase.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
6 changes: 3 additions & 3 deletions faststream_outbox/testing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
6 changes: 3 additions & 3 deletions tests/test_fastapi.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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
Expand Down
9 changes: 5 additions & 4 deletions tests/test_unit.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 {}
Expand Down Expand Up @@ -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()
Expand Down
Loading