Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,9 @@
)

if TYPE_CHECKING:
from .idempotent import (
BrokerSubscriberDuplicateDetector,
IdempotentBrokerSubscriber,
from .filtered import (
BrokerSubscriberValidator,
FilteredBrokerSubscriber,
)
from .queued import (
BrokerSubscriberQueue,
Expand Down Expand Up @@ -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 (
Expand All @@ -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

Expand All @@ -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)
Expand All @@ -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]):
Expand All @@ -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)
Expand Down Expand Up @@ -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()
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
from .impl import (
FilteredBrokerSubscriber,
)
from .validators import (
BrokerSubscriberDuplicateValidator,
BrokerSubscriberValidator,
InMemoryBrokerSubscriberDuplicateValidator,
PostgreSqlBrokerSubscriberDuplicateValidator,
PostgreSqlBrokerSubscriberDuplicateValidatorBuilder,
PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory,
)
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
from .abc import (
BrokerSubscriberValidator,
)
from .duplicates import (
BrokerSubscriberDuplicateValidator,
InMemoryBrokerSubscriberDuplicateValidator,
PostgreSqlBrokerSubscriberDuplicateValidator,
PostgreSqlBrokerSubscriberDuplicateValidatorBuilder,
PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory,
)
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,6 @@
ABC,
abstractmethod,
)
from uuid import (
UUID,
)

from minos.common import (
BuildableMixin,
Expand All @@ -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:
Expand All @@ -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
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
from .abc import (
BrokerSubscriberDuplicateValidator,
)
from .memory import (
InMemoryBrokerSubscriberDuplicateValidator,
)
from .pg import (
PostgreSqlBrokerSubscriberDuplicateValidator,
PostgreSqlBrokerSubscriberDuplicateValidatorBuilder,
PostgreSqlBrokerSubscriberDuplicateValidatorQueryFactory,
)
Original file line number Diff line number Diff line change
@@ -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
Loading