From 5eb33501c813a4f10396171548a72ad528025f40 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 10:00:36 +0100 Subject: [PATCH 1/6] ISSUE #87 * Add empty files. --- .../minos/networks/brokers/subscribers/idempotent/__init__.py | 0 .../networks/brokers/subscribers/idempotent/detectors/__init__.py | 0 .../networks/brokers/subscribers/idempotent/detectors/abc.py | 0 .../networks/brokers/subscribers/idempotent/detectors/memory.py | 0 .../minos/networks/brokers/subscribers/idempotent/detectors/pg.py | 0 .../minos/networks/brokers/subscribers/idempotent/impl.py | 0 6 files changed, 0 insertions(+), 0 deletions(-) create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/abc.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/memory.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py 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..e69de29bb 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..e69de29bb 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..e69de29bb 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..e69de29bb 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..e69de29bb 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..e69de29bb From 92c6a38659e0254dbe0f8892b0e884babd10d6d8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 11:28:37 +0100 Subject: [PATCH 2/6] ISSUE #87 * Add class skeletons. --- .../minos/networks/__init__.py | 4 +++ .../minos/networks/brokers/__init__.py | 4 +++ .../networks/brokers/subscribers/__init__.py | 6 ++++ .../subscribers/idempotent/__init__.py | 8 +++++ .../idempotent/detectors/__init__.py | 9 +++++ .../subscribers/idempotent/detectors/abc.py | 35 +++++++++++++++++++ .../idempotent/detectors/memory.py | 14 ++++++++ .../subscribers/idempotent/detectors/pg.py | 14 ++++++++ .../brokers/subscribers/idempotent/impl.py | 32 +++++++++++++++++ 9 files changed, 126 insertions(+) diff --git a/packages/core/minos-microservice-networks/minos/networks/__init__.py b/packages/core/minos-microservice-networks/minos/networks/__init__.py index 3a37a162f..a1f0d0920 100644 --- a/packages/core/minos-microservice-networks/minos/networks/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/__init__.py @@ -23,18 +23,22 @@ BrokerResponseException, BrokerSubscriber, BrokerSubscriberBuilder, + BrokerSubscriberDuplicateDetector, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, + IdempotentBrokerSubscriber, InMemoryBrokerPublisher, InMemoryBrokerPublisherQueue, InMemoryBrokerQueue, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, + InMemoryBrokerSubscriberDuplicateDetector, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, PostgreSqlBrokerPublisherQueue, PostgreSqlBrokerPublisherQueueQueryFactory, PostgreSqlBrokerQueue, + PostgreSqlBrokerSubscriberDuplicateDetector, 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..e4d767f59 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,16 @@ from .subscribers import ( BrokerSubscriber, BrokerSubscriberBuilder, + BrokerSubscriberDuplicateDetector, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, + IdempotentBrokerSubscriber, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, + InMemoryBrokerSubscriberDuplicateDetector, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, + PostgreSqlBrokerSubscriberDuplicateDetector, 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..2a2208fa7 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,12 @@ BrokerSubscriber, BrokerSubscriberBuilder, ) +from .idempotent import ( + BrokerSubscriberDuplicateDetector, + IdempotentBrokerSubscriber, + InMemoryBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetector, +) 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 index e69de29bb..4e02df13f 100644 --- 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 @@ -0,0 +1,8 @@ +from .detectors import ( + BrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetector, +) +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 index e69de29bb..76f838de4 100644 --- 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 @@ -0,0 +1,9 @@ +from .abc import ( + BrokerSubscriberDuplicateDetector, +) +from .memory import ( + InMemoryBrokerSubscriberDuplicateDetector, +) +from .pg import ( + PostgreSqlBrokerSubscriberDuplicateDetector, +) 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 index e69de29bb..229453b4e 100644 --- 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 @@ -0,0 +1,35 @@ +from __future__ import ( + annotations, +) + +from abc import ( + ABC, + abstractmethod, +) +from uuid import ( + UUID, +) + +from minos.common import ( + MinosSetup, +) + +from ....messages import ( + BrokerMessage, +) + + +class BrokerSubscriberDuplicateDetector(ABC, MinosSetup): + """TODO""" + + async def is_valid(self, message: BrokerMessage) -> bool: + """TODO + + :param message: TODO + :return: TODO + """ + 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 index e69de29bb..65a3fe6c3 100644 --- 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 @@ -0,0 +1,14 @@ +from uuid import ( + UUID, +) + +from .abc import ( + BrokerSubscriberDuplicateDetector, +) + + +class InMemoryBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): + """TODO""" + + 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/pg.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py index e69de29bb..b83a7aa5e 100644 --- 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 @@ -0,0 +1,14 @@ +from uuid import ( + UUID, +) + +from .abc import ( + BrokerSubscriberDuplicateDetector, +) + + +class PostgreSqlBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): + """TODO""" + + 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/impl.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py index e69de29bb..dbc47df88 100644 --- 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 @@ -0,0 +1,32 @@ +from ...messages import ( + BrokerMessage, +) +from ..abc import ( + BrokerSubscriber, +) +from .detectors import ( + BrokerSubscriberDuplicateDetector, +) + + +class IdempotentBrokerSubscriber(BrokerSubscriber): + """TODO""" + + impl: BrokerSubscriber + duplicates_detector: BrokerSubscriberDuplicateDetector + + async def _setup(self) -> None: + await super()._setup() + await self.duplicates_detector.setup() + await self.impl.setup() + + async def _destroy(self) -> None: + await self.impl.destroy() + await self.duplicates_detector.destroy() + await super()._destroy() + + async def _receive(self) -> BrokerMessage: + message = None + while message is None or not self.duplicates_detector.is_valid(message): + message = await self.impl.receive() + return message From 8662c9f43aca430a56c8e6cc8dc39bbcf4eab846 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 13:49:39 +0100 Subject: [PATCH 3/6] ISSUE #87 * Use `SetupMixin` instead of `MinosSetup`. --- .../networks/brokers/subscribers/idempotent/detectors/abc.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 index 229453b4e..5b1f14188 100644 --- 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 @@ -11,7 +11,7 @@ ) from minos.common import ( - MinosSetup, + SetupMixin, ) from ....messages import ( @@ -19,7 +19,7 @@ ) -class BrokerSubscriberDuplicateDetector(ABC, MinosSetup): +class BrokerSubscriberDuplicateDetector(ABC, SetupMixin): """TODO""" async def is_valid(self, message: BrokerMessage) -> bool: From c3ede0d4bfff1a87d228f24bc9ffc5441fcdf361 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 15:25:42 +0100 Subject: [PATCH 4/6] ISSUE #87 * Minor change. --- .../minos/networks/brokers/subscribers/idempotent/impl.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) 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 index dbc47df88..bad9f10aa 100644 --- 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 @@ -8,6 +8,8 @@ BrokerSubscriberDuplicateDetector, ) +_sentinel = object() + class IdempotentBrokerSubscriber(BrokerSubscriber): """TODO""" @@ -26,7 +28,7 @@ async def _destroy(self) -> None: await super()._destroy() async def _receive(self) -> BrokerMessage: - message = None - while message is None or not self.duplicates_detector.is_valid(message): + message = _sentinel + while message is _sentinel or not self.duplicates_detector.is_valid(message): message = await self.impl.receive() return message From a85be4f5c59314e6bb8db41538840a7281f604e0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 17:58:04 +0100 Subject: [PATCH 5/6] ISSUE #87 * Add initial implementation of "minos.networks.brokers.subscribers.idempotent" module. --- .../minos/networks/__init__.py | 1 + .../minos/networks/brokers/__init__.py | 1 + .../networks/brokers/subscribers/__init__.py | 1 + .../subscribers/idempotent/__init__.py | 1 + .../idempotent/detectors/__init__.py | 1 + .../subscribers/idempotent/detectors/abc.py | 8 +- .../idempotent/detectors/memory.py | 24 +++- .../subscribers/idempotent/detectors/pg.py | 112 +++++++++++++++++- .../brokers/subscribers/idempotent/impl.py | 15 ++- .../test_idempotent/__init__.py | 0 .../test_detectors/__init__.py | 0 .../test_detectors/test_abc.py | 52 ++++++++ .../test_detectors/test_memory.py | 36 ++++++ .../test_idempotent/test_detectors/test_pg.py | 42 +++++++ .../test_idempotent/test_impl.py | 74 ++++++++++++ 15 files changed, 354 insertions(+), 14 deletions(-) create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/__init__.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/__init__.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_impl.py diff --git a/packages/core/minos-microservice-networks/minos/networks/__init__.py b/packages/core/minos-microservice-networks/minos/networks/__init__.py index a1f0d0920..897157fc9 100644 --- a/packages/core/minos-microservice-networks/minos/networks/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/__init__.py @@ -39,6 +39,7 @@ 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 e4d767f59..956943326 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py @@ -50,6 +50,7 @@ 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 2a2208fa7..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 @@ -7,6 +7,7 @@ IdempotentBrokerSubscriber, InMemoryBrokerSubscriberDuplicateDetector, PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, ) from .memory import ( InMemoryBrokerSubscriber, 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 index 4e02df13f..3cdac8c73 100644 --- 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 @@ -2,6 +2,7 @@ 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 index 76f838de4..5bb630ea7 100644 --- 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 @@ -6,4 +6,5 @@ ) 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 index 5b1f14188..f4d0bd4ff 100644 --- 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 @@ -20,13 +20,13 @@ class BrokerSubscriberDuplicateDetector(ABC, SetupMixin): - """TODO""" + """Broker Subscriber Duplicate Detector class.""" async def is_valid(self, message: BrokerMessage) -> bool: - """TODO + """Check if the given message is valid. - :param message: TODO - :return: TODO + :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) 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 index 65a3fe6c3..62b605b35 100644 --- 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 @@ -1,3 +1,6 @@ +from collections.abc import ( + Iterable, +) from uuid import ( UUID, ) @@ -8,7 +11,24 @@ class InMemoryBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): - """TODO""" + """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: - raise NotImplementedError + 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 index b83a7aa5e..da2d4278e 100644 --- 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 @@ -1,14 +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): - """TODO""" +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: - raise NotImplementedError + 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 index bad9f10aa..63a10cd00 100644 --- 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 @@ -12,23 +12,28 @@ class IdempotentBrokerSubscriber(BrokerSubscriber): - """TODO""" + """Idempotent Broker Subscriber class.""" impl: BrokerSubscriber - duplicates_detector: BrokerSubscriberDuplicateDetector + 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.duplicates_detector.setup() + await self.duplicate_detector.setup() await self.impl.setup() async def _destroy(self) -> None: await self.impl.destroy() - await self.duplicates_detector.destroy() + await self.duplicate_detector.destroy() await super()._destroy() async def _receive(self) -> BrokerMessage: message = _sentinel - while message is _sentinel or not self.duplicates_detector.is_valid(message): + 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..0a95fed2e --- /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() + queue_setup_mock = AsyncMock() + queue_destroy_mock = AsyncMock() + + self.impl.setup = impl_setup_mock + self.impl.destroy = impl_destroy_mock + self.duplicate_detector.setup = queue_setup_mock + self.duplicate_detector.destroy = queue_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, queue_setup_mock.call_count) + self.assertEqual(0, queue_destroy_mock.call_count) + + impl_setup_mock.reset_mock() + impl_destroy_mock.reset_mock() + queue_setup_mock.reset_mock() + queue_destroy_mock.reset_mock() + + self.assertEqual(0, impl_setup_mock.call_count) + self.assertEqual(1, impl_destroy_mock.call_count) + self.assertEqual(0, queue_setup_mock.call_count) + self.assertEqual(1, queue_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() From 09656329c1c09c1a83defdaada0cdeae56811b28 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Mon, 14 Mar 2022 18:02:42 +0100 Subject: [PATCH 6/6] ISSUE #87 * Minor change. --- .../test_idempotent/test_impl.py | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) 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 index 0a95fed2e..659c09719 100644 --- 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 @@ -33,29 +33,29 @@ def test_queue(self): async def test_setup_destroy(self): impl_setup_mock = AsyncMock() impl_destroy_mock = AsyncMock() - queue_setup_mock = AsyncMock() - queue_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 = queue_setup_mock - self.duplicate_detector.destroy = queue_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, queue_setup_mock.call_count) - self.assertEqual(0, queue_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() - queue_setup_mock.reset_mock() - queue_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, queue_setup_mock.call_count) - self.assertEqual(1, queue_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"))