Skip to content
Open
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
62 changes: 38 additions & 24 deletions faststream/nats/publisher/producer.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import asyncio
from abc import abstractmethod
from contextlib import suppress
from typing import TYPE_CHECKING, Any, Optional

import anyio
Expand Down Expand Up @@ -183,31 +184,44 @@ async def request(self, cmd: "NatsPublishCommand") -> "Msg":
future=future,
max_msgs=1,
)
await sub.unsubscribe(limit=1)

headers_to_send = {
"content-type": content_type or "",
"reply_to": reply_to,
**cmd.headers_to_publish(js=False),
}

with anyio.fail_after(cmd.timeout):
await self.__state.connection.publish(
subject=cmd.destination,
payload=payload,
headers=headers_to_send,
stream=cmd.stream,
timeout=cmd.timeout,
)

msg = await future

if ( # pragma: no cover
msg.headers and (msg.headers.get(Header.STATUS) == NO_RESPONDERS_STATUS)
try:
await sub.unsubscribe(limit=1)

headers_to_send = {
"content-type": content_type or "",
"reply_to": reply_to,
**cmd.headers_to_publish(js=False),
}

with anyio.fail_after(cmd.timeout):
await self.__state.connection.publish(
subject=cmd.destination,
payload=payload,
headers=headers_to_send,
stream=cmd.stream,
timeout=cmd.timeout,
)

msg = await future

if ( # pragma: no cover
msg.headers
and (msg.headers.get(Header.STATUS) == NO_RESPONDERS_STATUS)
):
raise nats.errors.NoRespondersError

return msg
finally:
future.cancel()
# The one-message limit cannot remove an inbox that never receives a reply.
with (
anyio.CancelScope(shield=True),
suppress(
nats.errors.ConnectionClosedError,
nats.errors.ConnectionDrainingError,
),
):
raise nats.errors.NoRespondersError

return msg
await sub.unsubscribe()


class FakeNatsFastProducer(NatsFastProducer):
Expand Down
121 changes: 121 additions & 0 deletions tests/brokers/nats/test_request_cleanup.py

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Add type hints

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don’t like your parameterized tests because of the large number of branches.

Let’s write a separate test for each parameter.I don’t like your parameterized tests because of the large number of branches.

Let’s write a separate test for each parameter.

@Kuang-xianxin Kuang-xianxin Sep 8, 2026 •

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Updated in c7a7fd4: timeout, publish failure, reply-limit failure, and cancellation now have separate tests with fixed expectations. Shared real-server setup is in a typed fixture; test arguments, fixture yields, and callbacks have type hints. The connected tests still repeat each failure three times to check for accumulating subscriptions.

Validation: 27 cleanup/request tests passed with real NATS 2.14.6, and the changed test module passes strict mypy plus Ruff lint/format. The production fix is unchanged.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reworked in a99dd63 following the repository testing-patterns skill: reduced nine cleanup cases to four distinct boundaries, removed duplicate success/failure coverage and the single-use fixture, used the existing settings/queue fixtures, and consolidated the subscription assertions. Existing request tests cover successful replies.

The named cleanup/request files pass all 22 tests with real NATS 2.14.6; reverting the production fix gives 3 failures. Disabling shielding and removing connection-error suppression each fail their specific retained test. Ruff/format, codespell, strict test typing, and configured mypy over 516 source files pass. Production code is unchanged.

Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
from collections.abc import AsyncIterator
from typing import TypeAlias
from unittest.mock import AsyncMock

import anyio
import nats
import pytest
import pytest_asyncio
from nats.aio.client import Client
from nats.js import JetStreamContext
from nats.js.api import PubAck

from faststream.nats import NatsBroker

from .settings import Settings

pytestmark = [pytest.mark.asyncio(), pytest.mark.nats()]

RequestTransport: TypeAlias = tuple[NatsBroker, Client, list[bytes]]


@pytest_asyncio.fixture()
async def request_transport(
monkeypatch: pytest.MonkeyPatch,
) -> AsyncIterator[RequestTransport]:
connection = Client()
commands: list[bytes] = []

async def send_command(command: bytes) -> None:
await anyio.lowlevel.checkpoint()
commands.append(command)

monkeypatch.setattr(connection, "_send_command", send_command)
monkeypatch.setattr(connection, "_flush_pending", AsyncMock())
monkeypatch.setattr(nats, "connect", AsyncMock(return_value=connection))
broker = NatsBroker()
await broker.connect()
try:
yield broker, connection, commands
finally:
await connection.close()


async def test_failed_reply_limit_setup_removes_reply_subscription(
request_transport: RequestTransport,
monkeypatch: pytest.MonkeyPatch,
) -> None:
broker, connection, commands = request_transport
send_unsubscribe = connection._send_unsubscribe

async def fail_initial_unsubscribe(sid: int, limit: int = 0) -> None:
if limit:
raise RuntimeError
await send_unsubscribe(sid, limit)

monkeypatch.setattr(connection, "_send_unsubscribe", fail_initial_unsubscribe)

with pytest.raises(RuntimeError):
await broker.request("request", "jobs", stream="jobs", timeout=0.01)

assert (connection._subs, commands[-1].split()) == ({}, [b"UNSUB", b"1"])


@pytest.mark.connected()
async def test_real_timed_out_requests_remove_reply_subscriptions(
settings: Settings,
queue: str,
) -> None:
async with NatsBroker(settings.url) as broker:
connection = broker._connection
assert connection is not None
jetstream = connection.jetstream()
await jetstream.add_stream(name=queue, subjects=[queue])
try:
# Creating the stream first includes the shared JetStream reply mux.
subscriptions = set(connection._subs)
with pytest.raises(TimeoutError):
await broker.request("request", queue, stream=queue, timeout=0.05)

await connection.flush()
assert set(connection._subs) == subscriptions
finally:
await jetstream.delete_stream(queue)


async def test_cancelled_jetstream_request_unsubscribes_under_cancel_scope(
request_transport: RequestTransport,
monkeypatch: pytest.MonkeyPatch,
) -> None:
broker, connection, commands = request_transport

async def publish(*args: object, **kwargs: object) -> None:
scope.cancel()
await anyio.sleep_forever()

monkeypatch.setattr(JetStreamContext, "publish", publish)
with anyio.CancelScope() as scope:
await broker.request("request", "jobs", stream="jobs")

assert (scope.cancelled_caught, connection._subs, commands[-1].split()) == (
True,
{},
[b"UNSUB", b"1"],
)


async def test_closed_connection_preserves_publish_error(
request_transport: RequestTransport,
monkeypatch: pytest.MonkeyPatch,
) -> None:
broker, connection, _ = request_transport

async def publish(*args: object, **kwargs: object) -> PubAck:
await connection.close()
raise nats.errors.TimeoutError

monkeypatch.setattr(JetStreamContext, "publish", publish)
with pytest.raises(nats.errors.TimeoutError):
await broker.request("request", "jobs", stream="jobs")

assert not connection._subs
Loading