diff --git a/packages/core/minos-microservice-networks/minos/networks/__init__.py b/packages/core/minos-microservice-networks/minos/networks/__init__.py index 3a37a162f..897157fc9 100644 --- a/packages/core/minos-microservice-networks/minos/networks/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/__init__.py @@ -23,18 +23,23 @@ BrokerResponseException, BrokerSubscriber, BrokerSubscriberBuilder, + BrokerSubscriberDuplicateDetector, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, + IdempotentBrokerSubscriber, InMemoryBrokerPublisher, InMemoryBrokerPublisherQueue, InMemoryBrokerQueue, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, + InMemoryBrokerSubscriberDuplicateDetector, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, PostgreSqlBrokerPublisherQueue, PostgreSqlBrokerPublisherQueueQueryFactory, PostgreSqlBrokerQueue, + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, PostgreSqlBrokerSubscriberQueue, PostgreSqlBrokerSubscriberQueueBuilder, PostgreSqlBrokerSubscriberQueueQueryFactory, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py index 97931078e..956943326 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py @@ -40,12 +40,17 @@ from .subscribers import ( BrokerSubscriber, BrokerSubscriberBuilder, + BrokerSubscriberDuplicateDetector, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, + IdempotentBrokerSubscriber, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, + InMemoryBrokerSubscriberDuplicateDetector, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, PostgreSqlBrokerSubscriberQueue, PostgreSqlBrokerSubscriberQueueBuilder, PostgreSqlBrokerSubscriberQueueQueryFactory, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/__init__.py index f5282ff3b..9726686d0 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/__init__.py @@ -2,6 +2,13 @@ BrokerSubscriber, BrokerSubscriberBuilder, ) +from .idempotent import ( + BrokerSubscriberDuplicateDetector, + IdempotentBrokerSubscriber, + InMemoryBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, +) from .memory import ( InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py new file mode 100644 index 000000000..3cdac8c73 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py @@ -0,0 +1,9 @@ +from .detectors import ( + BrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, +) +from .impl import ( + IdempotentBrokerSubscriber, +) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py new file mode 100644 index 000000000..5bb630ea7 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py @@ -0,0 +1,10 @@ +from .abc import ( + BrokerSubscriberDuplicateDetector, +) +from .memory import ( + InMemoryBrokerSubscriberDuplicateDetector, +) +from .pg import ( + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, +) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/abc.py new file mode 100644 index 000000000..f4d0bd4ff --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/abc.py @@ -0,0 +1,35 @@ +from __future__ import ( + annotations, +) + +from abc import ( + ABC, + abstractmethod, +) +from uuid import ( + UUID, +) + +from minos.common import ( + SetupMixin, +) + +from ....messages import ( + BrokerMessage, +) + + +class BrokerSubscriberDuplicateDetector(ABC, SetupMixin): + """Broker Subscriber Duplicate Detector class.""" + + async def is_valid(self, message: BrokerMessage) -> bool: + """Check if the given message is valid. + + :param message: The message to be checked. + :return: ``True`` if it is valid or ``False`` otherwise. + """ + return await self._is_valid(message.topic, message.identifier) + + @abstractmethod + async def _is_valid(self, topic: str, uuid: UUID) -> bool: + raise NotImplementedError diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/memory.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/memory.py new file mode 100644 index 000000000..62b605b35 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/memory.py @@ -0,0 +1,34 @@ +from collections.abc import ( + Iterable, +) +from uuid import ( + UUID, +) + +from .abc import ( + BrokerSubscriberDuplicateDetector, +) + + +class InMemoryBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): + """In Memory Broker Subscriber Duplicate Detector class.""" + + def __init__(self, seen: Iterable[tuple[str, UUID]] = None, *args, **kwargs): + super().__init__(*args, **kwargs) + if seen is None: + seen = set() + self._seen = set(seen) + + @property + def seen(self) -> set[tuple[str, UUID]]: + """Get the seen pairs. + + :return: A ``set`` of ``tuple`` instances in which the first value is a ``str`` and the second an ``UUID``. + """ + return self._seen + + async def _is_valid(self, topic: str, uuid: UUID) -> bool: + if (topic, uuid) not in self._seen: + self._seen.add((topic, uuid)) + return True + return False diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py new file mode 100644 index 000000000..da2d4278e --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py @@ -0,0 +1,120 @@ +from __future__ import ( + annotations, +) + +from typing import ( + Optional, +) +from uuid import ( + UUID, +) + +from psycopg2 import ( + IntegrityError, +) +from psycopg2.sql import ( + SQL, +) + +from minos.common import ( + MinosConfig, + PostgreSqlMinosDatabase, +) + +from .abc import ( + BrokerSubscriberDuplicateDetector, +) + + +class PostgreSqlBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector, PostgreSqlMinosDatabase): + """PostgreSql Broker Subscriber Duplicate Detector class.""" + + def __init__( + self, query_factory: Optional[PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory] = None, *args, **kwargs + ): + if query_factory is None: + query_factory = PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory() + super().__init__(*args, **kwargs) + self._query_factory = query_factory + + @classmethod + def _from_config(cls, config: MinosConfig, **kwargs) -> PostgreSqlBrokerSubscriberDuplicateDetector: + # noinspection PyProtectedMember + return cls(**config.broker.queue._asdict(), **kwargs) + + async def _setup(self) -> None: + await super()._setup() + await self._create_table() + + async def _create_table(self) -> None: + await self.submit_query( + self._query_factory.build_activate_uuid_extension(), + lock=self._query_factory.build_uuid_extension_name(), + ) + await self.submit_query( + self._query_factory.build_create_table(), + lock=self._query_factory.build_table_name(), + ) + + @property + def query_factory(self) -> PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory: + """Get the query factory. + + :return: A ``PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory`` instance. + """ + return self._query_factory + + async def _is_valid(self, topic: str, uuid: UUID) -> bool: + try: + await self.submit_query(self._query_factory.build_insert_row(), {"topic": topic, "uuid": uuid}) + return True + except IntegrityError: + return False + + +class PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory: + """PostgreSql Broker Subscriber Duplicate Detector Query Factory class.""" + + @staticmethod + def build_uuid_extension_name() -> str: + """Build the uuid extension name. + + :return: A ``str`` instance. + """ + return "uuid-ossp" + + def build_activate_uuid_extension(self) -> SQL: + """Build activate uuid extension query. + + :return: A ``SQL`` instance. + """ + return SQL(f'CREATE EXTENSION IF NOT EXISTS "{self.build_uuid_extension_name()}";') + + @staticmethod + def build_table_name() -> str: + """Build the table name. + + :return: A ``str`` instance. + """ + return "broker_subscriber_processed_messages" + + def build_create_table(self) -> SQL: + """Build the "create table" query. + + :return: A ``SQL`` instance. + """ + return SQL( + f"CREATE TABLE IF NOT EXISTS {self.build_table_name()} (" + " topic VARCHAR(255) NOT NULL, " + " uuid UUID NOT NULL, " + " created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()," + " PRIMARY KEY (topic, uuid)" + ")" + ) + + def build_insert_row(self) -> SQL: + """Build the "insert row" query. + + :return: A ``SQL`` instance. + """ + return SQL(f"INSERT INTO {self.build_table_name()}(topic, uuid) VALUES(%(topic)s, %(uuid)s)") diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py new file mode 100644 index 000000000..63a10cd00 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py @@ -0,0 +1,39 @@ +from ...messages import ( + BrokerMessage, +) +from ..abc import ( + BrokerSubscriber, +) +from .detectors import ( + BrokerSubscriberDuplicateDetector, +) + +_sentinel = object() + + +class IdempotentBrokerSubscriber(BrokerSubscriber): + """Idempotent Broker Subscriber class.""" + + impl: BrokerSubscriber + duplicate_detector: BrokerSubscriberDuplicateDetector + + def __init__(self, impl: BrokerSubscriber, duplicate_detector: BrokerSubscriberDuplicateDetector, **kwargs): + super().__init__(**(kwargs | {"topics": impl.topics})) + self.impl = impl + self.duplicate_detector = duplicate_detector + + async def _setup(self) -> None: + await super()._setup() + await self.duplicate_detector.setup() + await self.impl.setup() + + async def _destroy(self) -> None: + await self.impl.destroy() + await self.duplicate_detector.destroy() + await super()._destroy() + + async def _receive(self) -> BrokerMessage: + message = _sentinel + while message is _sentinel or not (await self.duplicate_detector.is_valid(message)): + message = await self.impl.receive() + return message diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/__init__.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/__init__.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py new file mode 100644 index 000000000..54389f775 --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py @@ -0,0 +1,52 @@ +import unittest +from abc import ( + ABC, +) +from unittest.mock import ( + AsyncMock, + call, +) +from uuid import ( + UUID, +) + +from minos.common import ( + SetupMixin, +) +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriberDuplicateDetector, +) + + +class _BrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): + """For testing purposes.""" + + async def _is_valid(self, topic: str, uuid: UUID) -> bool: + """For testing purposes.""" + + +class TestBrokerSubscriberDuplicateDetector(unittest.IsolatedAsyncioTestCase): + def test_abstract(self): + self.assertTrue(issubclass(BrokerSubscriberDuplicateDetector, (ABC, SetupMixin))) + # noinspection PyUnresolvedReferences + self.assertEqual({"_is_valid"}, BrokerSubscriberDuplicateDetector.__abstractmethods__) + + async def test_is_valid(self): + message = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + detector = _BrokerSubscriberDuplicateDetector() + + mock = AsyncMock(side_effect=[True, False]) + detector._is_valid = mock + + self.assertTrue(await detector.is_valid(message)) + self.assertFalse(await detector.is_valid(message)) + + self.assertEqual( + [call(message.topic, message.identifier), call(message.topic, message.identifier)], mock.call_args_list + ) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py new file mode 100644 index 000000000..62f355485 --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py @@ -0,0 +1,36 @@ +import unittest +from uuid import ( + uuid4, +) + +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + InMemoryBrokerSubscriberDuplicateDetector, +) + + +class TestInMemoryBrokerSubscriberDuplicateDetector(unittest.IsolatedAsyncioTestCase): + async def test_constructor(self): + detector = InMemoryBrokerSubscriberDuplicateDetector() + self.assertEqual(set(), detector.seen) + + async def test_constructor_extended(self): + one = ("foo", uuid4()) + detector = InMemoryBrokerSubscriberDuplicateDetector([one]) + self.assertEqual({one}, detector.seen) + + async def test_is_valid(self): + one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + two = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + three = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + + detector = InMemoryBrokerSubscriberDuplicateDetector() + self.assertTrue(await detector.is_valid(one)) + self.assertTrue(await detector.is_valid(two)) + self.assertFalse(await detector.is_valid(one)) + self.assertTrue(await detector.is_valid(three)) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py new file mode 100644 index 000000000..b3fd7f04d --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py @@ -0,0 +1,42 @@ +import unittest + +from minos.common.testing import ( + PostgresAsyncTestCase, +) +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, +) +from tests.utils import ( + CONFIG_FILE_PATH, +) + + +class TestPostgreSqlBrokerSubscriberDuplicateDetector(PostgresAsyncTestCase): + CONFIG_FILE_PATH = CONFIG_FILE_PATH + + def test_is_subclass(self): + self.assertTrue(issubclass(PostgreSqlBrokerSubscriberDuplicateDetector, BrokerSubscriberDuplicateDetector)) + + async def test_query_factory(self): + detector = PostgreSqlBrokerSubscriberDuplicateDetector.from_config(self.config) + + self.assertIsInstance(detector.query_factory, PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory) + + async def test_is_valid(self): + one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + two = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + three = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + + async with PostgreSqlBrokerSubscriberDuplicateDetector.from_config(self.config) as detector: + self.assertTrue(await detector.is_valid(one)) + self.assertTrue(await detector.is_valid(two)) + self.assertFalse(await detector.is_valid(one)) + self.assertTrue(await detector.is_valid(three)) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_impl.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_impl.py new file mode 100644 index 000000000..659c09719 --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_impl.py @@ -0,0 +1,74 @@ +import unittest +from unittest.mock import ( + AsyncMock, +) + +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriber, + IdempotentBrokerSubscriber, + InMemoryBrokerSubscriber, + InMemoryBrokerSubscriberDuplicateDetector, +) + + +class TestIdempotentBrokerSubscriber(unittest.IsolatedAsyncioTestCase): + def setUp(self) -> None: + self.topics = {"foo", "bar"} + self.impl = InMemoryBrokerSubscriber(self.topics) + self.duplicate_detector = InMemoryBrokerSubscriberDuplicateDetector() + + def test_is_subclass(self): + self.assertTrue(issubclass(IdempotentBrokerSubscriber, BrokerSubscriber)) + + def test_impl(self): + subscriber = IdempotentBrokerSubscriber(self.impl, self.duplicate_detector) + self.assertEqual(self.impl, subscriber.impl) + + def test_queue(self): + subscriber = IdempotentBrokerSubscriber(self.impl, self.duplicate_detector) + self.assertEqual(self.duplicate_detector, subscriber.duplicate_detector) + + async def test_setup_destroy(self): + impl_setup_mock = AsyncMock() + impl_destroy_mock = AsyncMock() + duplicate_detector_setup_mock = AsyncMock() + duplicate_detector_destroy_mock = AsyncMock() + + self.impl.setup = impl_setup_mock + self.impl.destroy = impl_destroy_mock + self.duplicate_detector.setup = duplicate_detector_setup_mock + self.duplicate_detector.destroy = duplicate_detector_destroy_mock + + async with IdempotentBrokerSubscriber(self.impl, self.duplicate_detector): + self.assertEqual(1, impl_setup_mock.call_count) + self.assertEqual(0, impl_destroy_mock.call_count) + self.assertEqual(1, duplicate_detector_setup_mock.call_count) + self.assertEqual(0, duplicate_detector_destroy_mock.call_count) + + impl_setup_mock.reset_mock() + impl_destroy_mock.reset_mock() + duplicate_detector_setup_mock.reset_mock() + duplicate_detector_destroy_mock.reset_mock() + + self.assertEqual(0, impl_setup_mock.call_count) + self.assertEqual(1, impl_destroy_mock.call_count) + self.assertEqual(0, duplicate_detector_setup_mock.call_count) + self.assertEqual(1, duplicate_detector_destroy_mock.call_count) + + async def test_receive(self): + one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + two = BrokerMessageV1("bar", BrokerMessageV1Payload("foo")) + + self.impl.add_message(one) + self.impl.add_message(one) + self.impl.add_message(two) + + async with IdempotentBrokerSubscriber(self.impl, self.duplicate_detector) as subscriber: + self.assertEqual(one, await subscriber.receive()) + self.assertEqual(two, await subscriber.receive()) + + +if __name__ == "__main__": + unittest.main()