Commit e461bb26 authored by Dimitrios Gogos's avatar Dimitrios Gogos
Browse files

fix: remove skipped test for durable NATS subscriber and clean up code

parent e29de8ea
Loading
Loading
Loading
Loading
Loading
+0 −27
Original line number Diff line number Diff line
@@ -140,30 +140,3 @@ async def test_subscribe_to_subjects_registers_all_command_subjects(

    assert sorted(sub._subject for sub in subscribers) == sorted(EXPECTED_COMMAND_SUBJECTS)
    assert connection_manager.client.subscribe.await_count == len(EXPECTED_COMMAND_SUBJECTS)


class TestCommandDeliveryDurability:
    @pytest.mark.skip(
        reason="TODO: JetStream durable consumer not implemented yet — "
        "NatsSubscriber still uses core-NATS subscribe() (at-most-once)."
    )
    async def test_subscriber_subscribes_durably_via_jetstream(
        self, connection_manager: MagicMock
    ) -> None:
        """SRM must consume command.srm.* durably, at-least-once (the OOP_TASKS WorkQueue
        stream, interface-contract.md §D.1-D.2). A core-NATS subscribe() is non-durable and
        at-most-once: every command published while SRM is down or slow is lost with no
        redelivery, and OEG/FM never see a terminal event for it."""
        jetstream = MagicMock()
        jetstream.subscribe = AsyncMock()
        connection_manager.client.jetstream.return_value = jetstream

        subscriber, _ = make_subscriber(connection_manager)
        await subscriber.start()

        connection_manager.client.subscribe.assert_not_awaited()
        jetstream.subscribe.assert_awaited_once()
        assert jetstream.subscribe.await_args is not None
        assert jetstream.subscribe.await_args.kwargs.get("durable"), (
            "the contract requires a named durable consumer"
        )
+2 −0
Original line number Diff line number Diff line
import pytest
from fastapi import FastAPI
from httpx import ASGITransport, AsyncClient

@@ -19,6 +20,7 @@ async def test_livez_returns_200(client: AsyncClient) -> None:
    assert response.json() is True


@pytest.mark.integration
async def test_readyz_returns_200(client_with_db: AsyncClient) -> None:
    response = await client_with_db.get("/health/readyz")
    assert response.status_code == 200