From 0407a05b9d2d1396f066b143c6e3eaf8af3824ed Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Sergio=20Garc=C3=ADa=20Prado?= Date: Thu, 24 Mar 2022 15:58:45 +0100 Subject: [PATCH] ISSUE #87 * Rename `BrokerSubscriberDuplicateDetector` as `BrokerSubscriberDuplicateValidator`. * Rename `IdempotentBrokerSubscriber` as `FilteredBrokerSubscriber` --- .../minos/common/config/v2.py | 4 +- .../tests/config/v2.yml | 2 +- .../test_config/test_v2/test_base.py | 2 +- .../minos/networks/__init__.py | 13 ++-- .../minos/networks/brokers/__init__.py | 13 ++-- .../networks/brokers/subscribers/__init__.py | 15 ++-- .../minos/networks/brokers/subscribers/abc.py | 60 +++++++-------- .../brokers/subscribers/filtered/__init__.py | 11 +++ .../{idempotent => filtered}/impl.py | 22 +++--- .../filtered/validators/__init__.py | 10 +++ .../detectors => filtered/validators}/abc.py | 9 +-- .../validators/duplicates/__init__.py | 11 +++ .../filtered/validators/duplicates/abc.py | 34 +++++++++ .../validators/duplicates}/memory.py | 6 +- .../validators/duplicates}/pg.py | 20 ++--- .../subscribers/idempotent/__init__.py | 10 --- .../idempotent/detectors/__init__.py | 11 --- .../test_brokers/test_subscribers/test_abc.py | 74 +++++++++---------- .../__init__.py | 0 .../test_impl.py | 40 +++++----- .../test_validators}/__init__.py | 0 .../test_filtered/test_validators/test_abc.py | 48 ++++++++++++ .../test_duplicates/__init__.py | 0 .../test_duplicates/test_abc.py | 50 +++++++++++++ .../test_duplicates/test_memory.py | 36 +++++++++ .../test_duplicates/test_pg.py | 42 +++++++++++ .../test_detectors/test_abc.py | 52 ------------- .../test_detectors/test_memory.py | 36 --------- .../test_idempotent/test_detectors/test_pg.py | 42 ----------- 29 files changed, 379 insertions(+), 294 deletions(-) create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/__init__.py rename packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/{idempotent => filtered}/impl.py (50%) create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/__init__.py rename packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/{idempotent/detectors => filtered/validators}/abc.py (70%) create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/__init__.py create mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/abc.py rename packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/{idempotent/detectors => filtered/validators/duplicates}/memory.py (79%) rename packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/{idempotent/detectors => filtered/validators/duplicates}/pg.py (80%) delete mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py delete mode 100644 packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py rename packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/{test_idempotent => test_filtered}/__init__.py (100%) rename packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/{test_idempotent => test_filtered}/test_impl.py (50%) rename packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/{test_idempotent/test_detectors => test_filtered/test_validators}/__init__.py (100%) create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_abc.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/__init__.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_abc.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_memory.py create mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_pg.py delete mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py delete mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py delete mode 100644 packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py diff --git a/packages/core/minos-microservice-common/minos/common/config/v2.py b/packages/core/minos-microservice-common/minos/common/config/v2.py index a151bced5..3eae762f8 100644 --- a/packages/core/minos-microservice-common/minos/common/config/v2.py +++ b/packages/core/minos-microservice-common/minos/common/config/v2.py @@ -130,8 +130,8 @@ def _parse_broker_interface(data: dict[str, Any]) -> dict[str, Any]: data["subscriber"]["client"] = import_module(data["subscriber"]["client"]) if "queue" in data["subscriber"]: data["subscriber"]["queue"] = import_module(data["subscriber"]["queue"]) - if "idempotent" in data["subscriber"]: - data["subscriber"]["idempotent"] = import_module(data["subscriber"]["idempotent"]) + if "validator" in data["subscriber"]: + data["subscriber"]["validator"] = import_module(data["subscriber"]["validator"]) return data diff --git a/packages/core/minos-microservice-common/tests/config/v2.yml b/packages/core/minos-microservice-common/tests/config/v2.yml index 9073f058e..8c7771bd0 100644 --- a/packages/core/minos-microservice-common/tests/config/v2.yml +++ b/packages/core/minos-microservice-common/tests/config/v2.yml @@ -40,7 +40,7 @@ interfaces: subscriber: client: tests.utils.FakeBrokerSubscriber queue: builtins.int - idempotent: builtins.float + validator: builtins.float periodic: port: tests.utils.FakePeriodicPort pools: diff --git a/packages/core/minos-microservice-common/tests/test_common/test_config/test_v2/test_base.py b/packages/core/minos-microservice-common/tests/test_common/test_config/test_v2/test_base.py index 0280cbc34..2a565ab84 100644 --- a/packages/core/minos-microservice-common/tests/test_common/test_config/test_v2/test_base.py +++ b/packages/core/minos-microservice-common/tests/test_common/test_config/test_v2/test_base.py @@ -115,7 +115,7 @@ def test_interface_broker(self): "queue": {"records": 10, "retry": 2}, }, "publisher": {"client": FakeBrokerPublisher, "queue": int}, - "subscriber": {"client": FakeBrokerSubscriber, "queue": int, "idempotent": float}, + "subscriber": {"client": FakeBrokerSubscriber, "queue": int, "validator": float}, } self.assertEqual(expected, broker) diff --git a/packages/core/minos-microservice-networks/minos/networks/__init__.py b/packages/core/minos-microservice-networks/minos/networks/__init__.py index bd0e70be5..ac8e5d0af 100644 --- a/packages/core/minos-microservice-networks/minos/networks/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/__init__.py @@ -24,25 +24,26 @@ BrokerResponseException, BrokerSubscriber, BrokerSubscriberBuilder, - BrokerSubscriberDuplicateDetector, + BrokerSubscriberDuplicateValidator, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, - IdempotentBrokerSubscriber, + BrokerSubscriberValidator, + FilteredBrokerSubscriber, InMemoryBrokerPublisher, InMemoryBrokerPublisherQueue, InMemoryBrokerQueue, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, - InMemoryBrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateValidator, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, PostgreSqlBrokerPublisherQueue, PostgreSqlBrokerPublisherQueueQueryFactory, PostgreSqlBrokerQueue, PostgreSqlBrokerQueueBuilder, - PostgreSqlBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, - PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, 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 382de106f..dc2e473ec 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py @@ -42,18 +42,19 @@ from .subscribers import ( BrokerSubscriber, BrokerSubscriberBuilder, - BrokerSubscriberDuplicateDetector, + BrokerSubscriberDuplicateValidator, BrokerSubscriberQueue, BrokerSubscriberQueueBuilder, - IdempotentBrokerSubscriber, + BrokerSubscriberValidator, + FilteredBrokerSubscriber, InMemoryBrokerSubscriber, InMemoryBrokerSubscriberBuilder, - InMemoryBrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateValidator, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, - PostgreSqlBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, - PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, 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 3b4bbdf56..0a94221e5 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,13 +2,14 @@ BrokerSubscriber, BrokerSubscriberBuilder, ) -from .idempotent import ( - BrokerSubscriberDuplicateDetector, - IdempotentBrokerSubscriber, - InMemoryBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, - PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, +from .filtered import ( + BrokerSubscriberDuplicateValidator, + BrokerSubscriberValidator, + FilteredBrokerSubscriber, + InMemoryBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, ) from .memory import ( InMemoryBrokerSubscriber, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/abc.py index a6e51b16e..3fb2b5d1b 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/abc.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/abc.py @@ -33,9 +33,9 @@ ) if TYPE_CHECKING: - from .idempotent import ( - BrokerSubscriberDuplicateDetector, - IdempotentBrokerSubscriber, + from .filtered import ( + BrokerSubscriberValidator, + FilteredBrokerSubscriber, ) from .queued import ( BrokerSubscriberQueue, @@ -94,20 +94,20 @@ class BrokerSubscriberBuilder(Builder[BrokerSubscriberCls], Generic[BrokerSubscr def __init__( self, *args, - idempotent_builder: Optional[Builder] = None, + validator_builder: Optional[Builder] = None, queue_builder: Optional[BrokerSubscriberQueueBuilder] = None, - idempotent_cls: Optional[type[IdempotentBrokerSubscriber]] = None, + filtered_cls: Optional[type[FilteredBrokerSubscriber]] = None, queued_cls: Optional[type[QueuedBrokerSubscriber]] = None, **kwargs, ): super().__init__(*args, **kwargs) - if idempotent_cls is None: - from .idempotent import ( - IdempotentBrokerSubscriber, + if filtered_cls is None: + from .filtered import ( + FilteredBrokerSubscriber, ) - idempotent_cls = IdempotentBrokerSubscriber + filtered_cls = FilteredBrokerSubscriber if queued_cls is None: from .queued import ( @@ -116,19 +116,19 @@ def __init__( queued_cls = QueuedBrokerSubscriber - self.duplicate_detector_builder = idempotent_builder + self.validator_builder = validator_builder self.queue_builder = queue_builder - self.idempotent_cls = idempotent_cls + self.filtered_cls = filtered_cls self.queued_cls = queued_cls - def with_idempotent_cls(self, idempotent_cls: type[IdempotentBrokerSubscriber]): - """Set the idempotent class. + def with_filtered_cls(self, filtered_cls: type[FilteredBrokerSubscriber]): + """Set the filtered class. - :param idempotent_cls: A subclass of ``IdempotentBrokerSubscriber``. + :param filtered_cls: A subclass of ``FilteredBrokerSubscriber``. :return: This method return the builder instance. """ - self.idempotent_cls = idempotent_cls + self.filtered_cls = filtered_cls return self @@ -150,8 +150,8 @@ def with_config(self, config: Config): """ self._with_builders_from_config(config) - if self.duplicate_detector_builder is not None: - self.duplicate_detector_builder.with_config(config) + if self.validator_builder is not None: + self.validator_builder.with_config(config) if self.queue_builder is not None: self.queue_builder.with_config(config) return super().with_config(config) @@ -164,24 +164,24 @@ def _with_builders_from_config(self, config): broker_subscriber_config = broker_config["subscriber"] - if "idempotent" in broker_subscriber_config: - self.with_duplicate_detector(broker_subscriber_config["idempotent"]) + if "validator" in broker_subscriber_config: + self.with_validator(broker_subscriber_config["validator"]) if "queue" in broker_subscriber_config: self.with_queue(broker_subscriber_config["queue"]) - def with_duplicate_detector( + def with_validator( self, - duplicate_detector: Union[type[BrokerSubscriberDuplicateDetector], Builder[BrokerSubscriberDuplicateDetector]], + validator: Union[type[BrokerSubscriberValidator], Builder[BrokerSubscriberValidator]], ): """Set the duplicate detector. - :param duplicate_detector: The duplicate detector to be set. + :param validator: The duplicate detector to be set. :return: This method return the builder instance. """ - if not isinstance(duplicate_detector, Builder): - duplicate_detector = duplicate_detector.get_builder() - self.duplicate_detector_builder = duplicate_detector.copy() + if not isinstance(validator, Builder): + validator = validator.get_builder() + self.validator_builder = validator.copy() return self def with_queue(self, queue: Union[type[BrokerSubscriberQueue], BrokerSubscriberQueueBuilder]): @@ -201,8 +201,8 @@ def with_kwargs(self, kwargs: dict[str, Any]): :param kwargs: The kwargs to be set. :return: This method return the builder instance. """ - if self.duplicate_detector_builder is not None: - self.duplicate_detector_builder.with_kwargs(kwargs) + if self.validator_builder is not None: + self.validator_builder.with_kwargs(kwargs) if self.queue_builder is not None: self.queue_builder.with_kwargs(kwargs) @@ -250,9 +250,9 @@ def build(self) -> BrokerSubscriber: """ impl = super().build() - if self.duplicate_detector_builder is not None: - duplicate_detector = self.duplicate_detector_builder.build() - impl = self.idempotent_cls(impl=impl, duplicate_detector=duplicate_detector, **self.kwargs) + if self.validator_builder is not None: + validator = self.validator_builder.build() + impl = self.filtered_cls(impl=impl, validator=validator, **self.kwargs) if self.queue_builder is not None: queue = self.queue_builder.build() diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/__init__.py new file mode 100644 index 000000000..12e98c065 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/__init__.py @@ -0,0 +1,11 @@ +from .impl import ( + FilteredBrokerSubscriber, +) +from .validators import ( + BrokerSubscriberDuplicateValidator, + BrokerSubscriberValidator, + InMemoryBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, +) 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/filtered/impl.py similarity index 50% rename from packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py rename to packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/impl.py index 93d3511f9..a40c652dc 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/impl.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/impl.py @@ -8,39 +8,39 @@ from ..abc import ( BrokerSubscriber, ) -from .detectors import ( - BrokerSubscriberDuplicateDetector, +from .validators import ( + BrokerSubscriberValidator, ) _sentinel = object() -class IdempotentBrokerSubscriber(BrokerSubscriber): - """Idempotent Broker Subscriber class.""" +class FilteredBrokerSubscriber(BrokerSubscriber): + """Filtered Broker Subscriber class.""" impl: BrokerSubscriber - duplicate_detector: BrokerSubscriberDuplicateDetector + validator: BrokerSubscriberValidator - def __init__(self, impl: BrokerSubscriber, duplicate_detector: BrokerSubscriberDuplicateDetector, **kwargs): + def __init__(self, impl: BrokerSubscriber, validator: BrokerSubscriberValidator, **kwargs): super().__init__(**(kwargs | {"topics": impl.topics})) self.impl = impl - self.duplicate_detector = duplicate_detector + self.validator = validator async def _setup(self) -> None: await super()._setup() - await self.duplicate_detector.setup() + await self.validator.setup() await self.impl.setup() async def _destroy(self) -> None: await self.impl.destroy() - await self.duplicate_detector.destroy() + await self.validator.destroy() await super()._destroy() async def _receive(self) -> BrokerMessage: message = _sentinel - while message is _sentinel or not (await self.duplicate_detector.is_valid(message)): + while message is _sentinel or not (await self.validator.is_valid(message)): message = await self.impl.receive() return message -IdempotentBrokerSubscriber.set_builder(Builder) +FilteredBrokerSubscriber.set_builder(Builder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/__init__.py new file mode 100644 index 000000000..1ac774fae --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/__init__.py @@ -0,0 +1,10 @@ +from .abc import ( + BrokerSubscriberValidator, +) +from .duplicates import ( + BrokerSubscriberDuplicateValidator, + InMemoryBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, +) 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/filtered/validators/abc.py similarity index 70% rename from packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/abc.py rename to packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/abc.py index 154d5c3fa..d0e5f6f7f 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/filtered/validators/abc.py @@ -6,9 +6,6 @@ ABC, abstractmethod, ) -from uuid import ( - UUID, -) from minos.common import ( BuildableMixin, @@ -19,7 +16,7 @@ ) -class BrokerSubscriberDuplicateDetector(BuildableMixin, ABC): +class BrokerSubscriberValidator(BuildableMixin, ABC): """Broker Subscriber Duplicate Detector class.""" async def is_valid(self, message: BrokerMessage) -> bool: @@ -28,8 +25,8 @@ async def is_valid(self, message: BrokerMessage) -> bool: :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) + return await self._is_valid(message) @abstractmethod - async def _is_valid(self, topic: str, uuid: UUID) -> bool: + async def _is_valid(self, message: BrokerMessage) -> bool: raise NotImplementedError diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/__init__.py new file mode 100644 index 000000000..114d74268 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/__init__.py @@ -0,0 +1,11 @@ +from .abc import ( + BrokerSubscriberDuplicateValidator, +) +from .memory import ( + InMemoryBrokerSubscriberDuplicateValidator, +) +from .pg import ( + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorBuilder, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, +) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/abc.py new file mode 100644 index 000000000..5dcdf5199 --- /dev/null +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/abc.py @@ -0,0 +1,34 @@ +from __future__ import ( + annotations, +) + +from abc import ( + ABC, + abstractmethod, +) +from uuid import ( + UUID, +) + +from .....messages import ( + BrokerMessage, +) +from ..abc import ( + BrokerSubscriberValidator, +) + + +class BrokerSubscriberDuplicateValidator(BrokerSubscriberValidator, ABC): + """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_unique(message.topic, message.identifier) + + @abstractmethod + async def _is_unique(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/filtered/validators/duplicates/memory.py similarity index 79% rename from packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/memory.py rename to packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/memory.py index 62b605b35..87f331d19 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/filtered/validators/duplicates/memory.py @@ -6,11 +6,11 @@ ) from .abc import ( - BrokerSubscriberDuplicateDetector, + BrokerSubscriberDuplicateValidator, ) -class InMemoryBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector): +class InMemoryBrokerSubscriberDuplicateValidator(BrokerSubscriberDuplicateValidator): """In Memory Broker Subscriber Duplicate Detector class.""" def __init__(self, seen: Iterable[tuple[str, UUID]] = None, *args, **kwargs): @@ -27,7 +27,7 @@ def seen(self) -> set[tuple[str, UUID]]: """ return self._seen - async def _is_valid(self, topic: str, uuid: UUID) -> bool: + async def _is_unique(self, topic: str, uuid: UUID) -> bool: if (topic, uuid) not in self._seen: self._seen.add((topic, uuid)) return True 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/filtered/validators/duplicates/pg.py similarity index 80% rename from packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/pg.py rename to packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/filtered/validators/duplicates/pg.py index 079219814..fb2ee854d 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/filtered/validators/duplicates/pg.py @@ -23,18 +23,18 @@ ) from .abc import ( - BrokerSubscriberDuplicateDetector, + BrokerSubscriberDuplicateValidator, ) -class PostgreSqlBrokerSubscriberDuplicateDetector(BrokerSubscriberDuplicateDetector, PostgreSqlMinosDatabase): +class PostgreSqlBrokerSubscriberDuplicateValidator(BrokerSubscriberDuplicateValidator, PostgreSqlMinosDatabase): """PostgreSql Broker Subscriber Duplicate Detector class.""" def __init__( - self, query_factory: Optional[PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory] = None, *args, **kwargs + self, query_factory: Optional[PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory] = None, *args, **kwargs ): if query_factory is None: - query_factory = PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory() + query_factory = PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory() super().__init__(*args, **kwargs) self._query_factory = query_factory @@ -53,14 +53,14 @@ async def _create_table(self) -> None: ) @property - def query_factory(self) -> PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory: + def query_factory(self) -> PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory: """Get the query factory. - :return: A ``PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory`` instance. + :return: A ``PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory`` instance. """ return self._query_factory - async def _is_valid(self, topic: str, uuid: UUID) -> bool: + async def _is_unique(self, topic: str, uuid: UUID) -> bool: try: await self.submit_query(self._query_factory.build_insert_row(), {"topic": topic, "uuid": uuid}) return True @@ -68,7 +68,7 @@ async def _is_valid(self, topic: str, uuid: UUID) -> bool: return False -class PostgreSqlBrokerSubscriberDuplicateDetectorBuilder(Builder[PostgreSqlBrokerSubscriberDuplicateDetector]): +class PostgreSqlBrokerSubscriberDuplicateValidatorBuilder(Builder[PostgreSqlBrokerSubscriberDuplicateValidator]): """PostgreSql Broker Subscriber Duplicate Detector Builder class.""" def with_config(self, config: Config): @@ -81,10 +81,10 @@ def with_config(self, config: Config): return super().with_config(config) -PostgreSqlBrokerSubscriberDuplicateDetector.set_builder(PostgreSqlBrokerSubscriberDuplicateDetectorBuilder) +PostgreSqlBrokerSubscriberDuplicateValidator.set_builder(PostgreSqlBrokerSubscriberDuplicateValidatorBuilder) -class PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory: +class PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory: """PostgreSql Broker Subscriber Duplicate Detector Query Factory class.""" @staticmethod 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 deleted file mode 100644 index 10ac806c8..000000000 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/__init__.py +++ /dev/null @@ -1,10 +0,0 @@ -from .detectors import ( - BrokerSubscriberDuplicateDetector, - InMemoryBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, - 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 deleted file mode 100644 index 9a9e01451..000000000 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/idempotent/detectors/__init__.py +++ /dev/null @@ -1,11 +0,0 @@ -from .abc import ( - BrokerSubscriberDuplicateDetector, -) -from .memory import ( - InMemoryBrokerSubscriberDuplicateDetector, -) -from .pg import ( - PostgreSqlBrokerSubscriberDuplicateDetector, - PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, - PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, -) diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_abc.py index aa1631d61..bb48730af 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_abc.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_abc.py @@ -19,9 +19,9 @@ BrokerMessageV1Payload, BrokerSubscriber, BrokerSubscriberBuilder, - IdempotentBrokerSubscriber, + FilteredBrokerSubscriber, InMemoryBrokerSubscriber, - InMemoryBrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateValidator, InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, QueuedBrokerSubscriber, @@ -82,7 +82,7 @@ class TestBrokerSubscriberBuilder(unittest.TestCase): def test_constructor(self): builder = BrokerSubscriberBuilder() self.assertEqual(None, builder.queue_builder) - self.assertEqual(None, builder.duplicate_detector_builder) + self.assertEqual(None, builder.validator_builder) self.assertEqual(QueuedBrokerSubscriber, builder.queued_cls) def test_with_queued_cls(self): @@ -90,10 +90,10 @@ def test_with_queued_cls(self): builder = BrokerSubscriberBuilder().with_queued_cls(int) self.assertEqual(int, builder.queued_cls) - def test_with_idempotent_cls(self): + def test_with_filtered_cls(self): # noinspection PyTypeChecker - builder = BrokerSubscriberBuilder().with_idempotent_cls(int) - self.assertEqual(int, builder.idempotent_cls) + builder = BrokerSubscriberBuilder().with_filtered_cls(int) + self.assertEqual(int, builder.filtered_cls) def test_constructor_with_queue_builder(self): queue_builder = InMemoryBrokerSubscriberQueueBuilder() @@ -101,10 +101,10 @@ def test_constructor_with_queue_builder(self): self.assertEqual(queue_builder, builder.queue_builder) self.assertEqual(QueuedBrokerSubscriber, builder.queued_cls) - def test_constructor_with_duplicate_detector(self): - idempotent_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector) - builder = BrokerSubscriberBuilder(idempotent_builder=idempotent_builder) - self.assertEqual(idempotent_builder, builder.duplicate_detector_builder) + def test_constructor_with_validator(self): + validator_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator) + builder = BrokerSubscriberBuilder(validator_builder=validator_builder) + self.assertEqual(validator_builder, builder.validator_builder) self.assertEqual(QueuedBrokerSubscriber, builder.queued_cls) def test_with_config_none(self): @@ -134,7 +134,7 @@ def test_with_config(self): return_value={ "subscriber": { "queue": InMemoryBrokerSubscriberQueue, - "idempotent": InMemoryBrokerSubscriberDuplicateDetector, + "validator": InMemoryBrokerSubscriberDuplicateValidator, } } ) @@ -142,9 +142,7 @@ def test_with_config(self): builder = BrokerSubscriberBuilder().with_config(config) self.assertEqual(InMemoryBrokerSubscriberQueueBuilder(), builder.queue_builder) - self.assertEqual( - Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector), builder.duplicate_detector_builder - ) + self.assertEqual(Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator), builder.validator_builder) self.assertEqual({}, builder.kwargs) def test_with_queue_with_config(self): @@ -158,14 +156,10 @@ def test_with_duplicate_with_config(self): config = Config(CONFIG_FILE_PATH) builder = ( - BrokerSubscriberBuilder() - .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) - .with_config(config) + BrokerSubscriberBuilder().with_validator(InMemoryBrokerSubscriberDuplicateValidator).with_config(config) ) self.assertEqual({}, builder.kwargs) - self.assertEqual( - Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector), builder.duplicate_detector_builder - ) + self.assertEqual(Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator), builder.validator_builder) def test_with_kwargs(self): builder = BrokerSubscriberBuilder().with_kwargs({"foo": "bar"}) @@ -177,15 +171,15 @@ def test_with_queue_with_kwargs(self): self.assertEqual(InMemoryBrokerSubscriberQueueBuilder().with_kwargs({"foo": "bar"}), builder.queue_builder) self.assertEqual({"foo": "bar"}, builder.kwargs) - def test_with_duplicate_detector_with_kwargs(self): + def test_with_validator_with_kwargs(self): builder = ( BrokerSubscriberBuilder() - .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_validator(InMemoryBrokerSubscriberDuplicateValidator) .with_kwargs({"foo": "bar"}) ) self.assertEqual( - Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector).with_kwargs({"foo": "bar"}), - builder.duplicate_detector_builder, + Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator).with_kwargs({"foo": "bar"}), + builder.validator_builder, ) self.assertEqual({"foo": "bar"}, builder.kwargs) @@ -199,15 +193,15 @@ def test_with_queue_builder(self): builder = BrokerSubscriberBuilder().with_queue(queue_builder) self.assertEqual(queue_builder, builder.queue_builder) - def test_with_duplicate_detector_cls(self): - duplicate_detector_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector) - builder = BrokerSubscriberBuilder().with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) - self.assertEqual(duplicate_detector_builder, builder.duplicate_detector_builder) + def test_with_validator_cls(self): + validator_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator) + builder = BrokerSubscriberBuilder().with_validator(InMemoryBrokerSubscriberDuplicateValidator) + self.assertEqual(validator_builder, builder.validator_builder) - def test_with_duplicate_detector_builder(self): - duplicate_detector_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector) - builder = BrokerSubscriberBuilder().with_duplicate_detector(duplicate_detector_builder) - self.assertEqual(duplicate_detector_builder, builder.duplicate_detector_builder) + def test_with_validator_builder(self): + validator_builder = Builder().with_cls(InMemoryBrokerSubscriberDuplicateValidator) + builder = BrokerSubscriberBuilder().with_validator(validator_builder) + self.assertEqual(validator_builder, builder.validator_builder) def test_build(self): subscriber = BrokerSubscriberBuilder().with_topics({"one", "two"}).with_cls(InMemoryBrokerSubscriber).build() @@ -227,24 +221,24 @@ def test_build_with_queue(self): self.assertIsInstance(subscriber.impl, InMemoryBrokerSubscriber) self.assertIsInstance(subscriber.queue, InMemoryBrokerSubscriberQueue) - def test_build_with_duplicate_detector(self): + def test_build_with_validator(self): subscriber = ( BrokerSubscriberBuilder() .with_cls(InMemoryBrokerSubscriber) - .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_validator(InMemoryBrokerSubscriberDuplicateValidator) .with_topics({"one", "two"}) .build() ) - self.assertIsInstance(subscriber, IdempotentBrokerSubscriber) + self.assertIsInstance(subscriber, FilteredBrokerSubscriber) self.assertIsInstance(subscriber.impl, InMemoryBrokerSubscriber) - self.assertIsInstance(subscriber.duplicate_detector, InMemoryBrokerSubscriberDuplicateDetector) + self.assertIsInstance(subscriber.validator, InMemoryBrokerSubscriberDuplicateValidator) - def test_build_with_duplicate_detector_with_queue(self): + def test_build_with_validator_with_queue(self): subscriber = ( BrokerSubscriberBuilder() .with_cls(InMemoryBrokerSubscriber) - .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_validator(InMemoryBrokerSubscriberDuplicateValidator) .with_queue(InMemoryBrokerSubscriberQueue) .with_topics({"one", "two"}) .build() @@ -252,9 +246,9 @@ def test_build_with_duplicate_detector_with_queue(self): self.assertIsInstance(subscriber, QueuedBrokerSubscriber) self.assertIsInstance(subscriber.queue, InMemoryBrokerSubscriberQueue) - self.assertIsInstance(subscriber.impl, IdempotentBrokerSubscriber) + self.assertIsInstance(subscriber.impl, FilteredBrokerSubscriber) self.assertIsInstance(subscriber.impl.impl, InMemoryBrokerSubscriber) - self.assertIsInstance(subscriber.impl.duplicate_detector, InMemoryBrokerSubscriberDuplicateDetector) + self.assertIsInstance(subscriber.impl.validator, InMemoryBrokerSubscriberDuplicateValidator) def test_with_group_id(self): builder = BrokerSubscriberBuilder().with_group_id("foobar") 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_filtered/__init__.py similarity index 100% rename from packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/__init__.py rename to packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/__init__.py 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_filtered/test_impl.py similarity index 50% rename from packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_impl.py rename to packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_impl.py index 659c09719..4949cb70d 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_filtered/test_impl.py @@ -7,55 +7,55 @@ BrokerMessageV1, BrokerMessageV1Payload, BrokerSubscriber, - IdempotentBrokerSubscriber, + FilteredBrokerSubscriber, InMemoryBrokerSubscriber, - InMemoryBrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberDuplicateValidator, ) -class TestIdempotentBrokerSubscriber(unittest.IsolatedAsyncioTestCase): +class TestFilteredBrokerSubscriber(unittest.IsolatedAsyncioTestCase): def setUp(self) -> None: self.topics = {"foo", "bar"} self.impl = InMemoryBrokerSubscriber(self.topics) - self.duplicate_detector = InMemoryBrokerSubscriberDuplicateDetector() + self.validator = InMemoryBrokerSubscriberDuplicateValidator() def test_is_subclass(self): - self.assertTrue(issubclass(IdempotentBrokerSubscriber, BrokerSubscriber)) + self.assertTrue(issubclass(FilteredBrokerSubscriber, BrokerSubscriber)) def test_impl(self): - subscriber = IdempotentBrokerSubscriber(self.impl, self.duplicate_detector) + subscriber = FilteredBrokerSubscriber(self.impl, self.validator) 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) + subscriber = FilteredBrokerSubscriber(self.impl, self.validator) + self.assertEqual(self.validator, subscriber.validator) async def test_setup_destroy(self): impl_setup_mock = AsyncMock() impl_destroy_mock = AsyncMock() - duplicate_detector_setup_mock = AsyncMock() - duplicate_detector_destroy_mock = AsyncMock() + validator_setup_mock = AsyncMock() + validator_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 + self.validator.setup = validator_setup_mock + self.validator.destroy = validator_destroy_mock - async with IdempotentBrokerSubscriber(self.impl, self.duplicate_detector): + async with FilteredBrokerSubscriber(self.impl, self.validator): 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) + self.assertEqual(1, validator_setup_mock.call_count) + self.assertEqual(0, validator_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() + validator_setup_mock.reset_mock() + validator_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) + self.assertEqual(0, validator_setup_mock.call_count) + self.assertEqual(1, validator_destroy_mock.call_count) async def test_receive(self): one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) @@ -65,7 +65,7 @@ async def test_receive(self): self.impl.add_message(one) self.impl.add_message(two) - async with IdempotentBrokerSubscriber(self.impl, self.duplicate_detector) as subscriber: + async with FilteredBrokerSubscriber(self.impl, self.validator) as subscriber: self.assertEqual(one, await subscriber.receive()) self.assertEqual(two, await subscriber.receive()) 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_filtered/test_validators/__init__.py similarity index 100% rename from packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/__init__.py rename to packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/__init__.py diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_abc.py new file mode 100644 index 000000000..1289df79b --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_abc.py @@ -0,0 +1,48 @@ +import unittest +from abc import ( + ABC, +) +from unittest.mock import ( + AsyncMock, + call, +) + +from minos.common import ( + SetupMixin, +) +from minos.networks import ( + BrokerMessage, + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriberValidator, +) + + +class _BrokerSubscriberValidator(BrokerSubscriberValidator): + """For testing purposes.""" + + async def _is_valid(self, message: BrokerMessage) -> bool: + """For testing purposes.""" + + +class TestBrokerSubscriberValidator(unittest.IsolatedAsyncioTestCase): + def test_abstract(self): + self.assertTrue(issubclass(BrokerSubscriberValidator, (ABC, SetupMixin))) + # noinspection PyUnresolvedReferences + self.assertEqual({"_is_valid"}, BrokerSubscriberValidator.__abstractmethods__) + + async def test_is_valid(self): + message = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + validator = _BrokerSubscriberValidator() + + mock = AsyncMock(side_effect=[True, False]) + validator._is_valid = mock + + self.assertTrue(await validator.is_valid(message)) + self.assertFalse(await validator.is_valid(message)) + + self.assertEqual([call(message), call(message)], 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_filtered/test_validators/test_duplicates/__init__.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/__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_filtered/test_validators/test_duplicates/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_abc.py new file mode 100644 index 000000000..c3da76e57 --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_abc.py @@ -0,0 +1,50 @@ +import unittest +from abc import ( + ABC, +) +from unittest.mock import ( + AsyncMock, + call, +) +from uuid import ( + UUID, +) + +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriberDuplicateValidator, + BrokerSubscriberValidator, +) + + +class _BrokerSubscriberDuplicateValidator(BrokerSubscriberDuplicateValidator): + """For testing purposes.""" + + async def _is_unique(self, topic: str, uuid: UUID) -> bool: + """For testing purposes.""" + + +class TestBrokerSubscriberDuplicateValidator(unittest.IsolatedAsyncioTestCase): + def test_abstract(self): + self.assertTrue(issubclass(BrokerSubscriberDuplicateValidator, (ABC, BrokerSubscriberValidator))) + # noinspection PyUnresolvedReferences + self.assertEqual({"_is_unique"}, BrokerSubscriberDuplicateValidator.__abstractmethods__) + + async def test_is_unique(self): + message = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + validator = _BrokerSubscriberDuplicateValidator() + + mock = AsyncMock(side_effect=[True, False]) + validator._is_unique = mock + + self.assertTrue(await validator.is_valid(message)) + self.assertFalse(await validator.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_filtered/test_validators/test_duplicates/test_memory.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_memory.py new file mode 100644 index 000000000..6090b0a4c --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_memory.py @@ -0,0 +1,36 @@ +import unittest +from uuid import ( + uuid4, +) + +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + InMemoryBrokerSubscriberDuplicateValidator, +) + + +class TestInMemoryBrokerSubscriberDuplicateValidator(unittest.IsolatedAsyncioTestCase): + async def test_constructor(self): + validator = InMemoryBrokerSubscriberDuplicateValidator() + self.assertEqual(set(), validator.seen) + + async def test_constructor_extended(self): + one = ("foo", uuid4()) + validator = InMemoryBrokerSubscriberDuplicateValidator([one]) + self.assertEqual({one}, validator.seen) + + async def test_is_valid(self): + one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + two = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + three = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + + validator = InMemoryBrokerSubscriberDuplicateValidator() + self.assertTrue(await validator.is_valid(one)) + self.assertTrue(await validator.is_valid(two)) + self.assertFalse(await validator.is_valid(one)) + self.assertTrue(await validator.is_valid(three)) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_pg.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_pg.py new file mode 100644 index 000000000..7996b8e74 --- /dev/null +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_filtered/test_validators/test_duplicates/test_pg.py @@ -0,0 +1,42 @@ +import unittest + +from minos.common.testing import ( + PostgresAsyncTestCase, +) +from minos.networks import ( + BrokerMessageV1, + BrokerMessageV1Payload, + BrokerSubscriberValidator, + PostgreSqlBrokerSubscriberDuplicateValidator, + PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory, +) +from tests.utils import ( + CONFIG_FILE_PATH, +) + + +class TestPostgreSqlBrokerSubscriberDuplicateValidator(PostgresAsyncTestCase): + CONFIG_FILE_PATH = CONFIG_FILE_PATH + + def test_is_subclass(self): + self.assertTrue(issubclass(PostgreSqlBrokerSubscriberDuplicateValidator, BrokerSubscriberValidator)) + + async def test_query_factory(self): + validator = PostgreSqlBrokerSubscriberDuplicateValidator.from_config(self.config) + + self.assertIsInstance(validator.query_factory, PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory) + + async def test_is_valid(self): + one = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + two = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + three = BrokerMessageV1("foo", BrokerMessageV1Payload("bar")) + + async with PostgreSqlBrokerSubscriberDuplicateValidator.from_config(self.config) as validator: + self.assertTrue(await validator.is_valid(one)) + self.assertTrue(await validator.is_valid(two)) + self.assertFalse(await validator.is_valid(one)) + self.assertTrue(await validator.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_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py deleted file mode 100644 index 54389f775..000000000 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_abc.py +++ /dev/null @@ -1,52 +0,0 @@ -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 deleted file mode 100644 index 62f355485..000000000 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_memory.py +++ /dev/null @@ -1,36 +0,0 @@ -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 deleted file mode 100644 index b3fd7f04d..000000000 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_idempotent/test_detectors/test_pg.py +++ /dev/null @@ -1,42 +0,0 @@ -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()