diff --git a/packages/core/minos-microservice-common/minos/common/builders.py b/packages/core/minos-microservice-common/minos/common/builders.py index 1902e65aa..73df3b18b 100644 --- a/packages/core/minos-microservice-common/minos/common/builders.py +++ b/packages/core/minos-microservice-common/minos/common/builders.py @@ -4,12 +4,14 @@ from abc import ( ABC, - abstractmethod, ) from typing import ( Any, Generic, + Optional, TypeVar, + Union, + get_args, ) from .config import ( @@ -25,25 +27,50 @@ class Builder(SetupMixin, ABC, Generic[Instance]): """Builder class.""" - def __init__(self, *args, **kwargs): + def __init__(self, instance_cls: Optional[type[Instance]] = None, *args, **kwargs): + super().__init__(*args, **kwargs) + if instance_cls is None: + instance_cls = self._get_cls() + self.kwargs = dict() + self.instance_cls = instance_cls + + def _get_cls(self) -> Optional[type]: + # noinspection PyUnresolvedReferences + bases = self.__orig_bases__ + + instance_cls = get_args(next((base for base in bases if len(get_args(base))), None))[0] + + if not isinstance(instance_cls, type): + return None + + return instance_cls def copy(self: type[B]) -> B: """Get a copy of the instance. - :return: A ``BrokerSubscriberBuilder`` instance. + :return: A ``Builder`` instance. """ - return self.new().with_kwargs(self.kwargs) + return self.new().with_cls(self.instance_cls).with_kwargs(self.kwargs) @classmethod def new(cls: type[B]) -> B: """Get a new instance. - :return: A ``BrokerSubscriberBuilder`` instance. + :return: A ``Builder`` instance. """ return cls() + def with_cls(self: B, cls: type) -> B: + """Set class to be built. + + :param cls: The class to be set. + :return: This method return the builder instance. + """ + self.instance_cls = cls + return self + def with_kwargs(self: B, kwargs: dict[str, Any]) -> B: """Set kwargs. @@ -62,12 +89,18 @@ def with_config(self: B, config: Config) -> B: """ return self - @abstractmethod def build(self) -> Instance: """Build the instance. - :return: A ``BrokerSubscriber`` instance. + :return: A ``Instance`` instance. """ + return self.instance_cls(**self.kwargs) + + def __eq__(self, other: Any) -> bool: + return isinstance(other, type(self)) and self.instance_cls == other.instance_cls and self.kwargs == other.kwargs + + def __repr__(self) -> str: + return f"{type(self).__name__}({self.instance_cls.__name__}, {self.kwargs!r})" Ins = TypeVar("Ins", bound="BuildableMixin") @@ -76,28 +109,38 @@ def build(self) -> Instance: class BuildableMixin(SetupMixin): """Buildable Mixin class.""" - _builder_cls: type[Builder[Ins]] + _builder: Union[Builder[Ins], type[Builder[Ins]]] = Builder @classmethod - def _from_config(cls: type[Ins], config: Config, **kwargs) -> Ins: - return cls.get_builder().new().with_config(config).with_kwargs(kwargs).build() + def _from_config(cls, config: Config, **kwargs): + return cls.get_builder().with_config(config).with_kwargs(kwargs).build() @classmethod - def set_builder(cls: type[Ins], builder: type[Builder[Ins]]) -> None: + def set_builder(cls: type[Ins], builder: Union[Builder[Ins], type[Builder[Ins]]]) -> None: """Set a builder class. :param builder: The builder class to be set. :return: This method does not return anything. """ - cls._builder_cls = builder + if not isinstance(builder, Builder) and not (isinstance(builder, type) and issubclass(builder, Builder)): + raise ValueError(f"Given builder value is invalid: {builder!r}") + + cls._builder = builder @classmethod - def get_builder(cls) -> type[Builder[Ins]]: + def get_builder(cls) -> Builder[Ins]: """Get the builder class. - :return: A ``Builder`` subclass. + :return: A ``Builder`` instance. """ - return cls._builder_cls + builder = cls._builder + + if isinstance(builder, Builder): + builder = builder.copy() + else: + builder = builder.new() + + return builder.with_cls(cls) B = TypeVar("B", bound=Builder) 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 246feee6e..a151bced5 100644 --- a/packages/core/minos-microservice-common/minos/common/config/v2.py +++ b/packages/core/minos-microservice-common/minos/common/config/v2.py @@ -11,6 +11,7 @@ from typing import ( TYPE_CHECKING, Any, + Union, ) from ..exceptions import ( @@ -39,7 +40,7 @@ def _version(self) -> int: def _get_name(self) -> str: return self.get_by_key("name") - def _get_injections(self) -> list[type[InjectableMixin]]: + def _get_injections(self) -> list[Union[InjectableMixin, type[InjectableMixin]]]: from ..builders import ( BuildableMixin, ) @@ -82,15 +83,15 @@ def _get_injections(self) -> list[type[InjectableMixin]]: if ( not issubclass(type_, InjectableMixin) and issubclass(type_, BuildableMixin) - and issubclass((builder_type := type_.get_builder()), InjectableMixin) + and isinstance((builder_type := type_.get_builder()), InjectableMixin) ): type_ = builder_type - - if not issubclass(type_, InjectableMixin): + elif not issubclass(type_, InjectableMixin): raise MinosConfigException(f"{type_!r} must be subclass of {InjectableMixin!r}.") ans.append(type_) + # noinspection PyTypeChecker return ans def _get_databases(self) -> dict[str, dict[str, Any]]: diff --git a/packages/core/minos-microservice-common/tests/test_common/test_builders.py b/packages/core/minos-microservice-common/tests/test_common/test_builders.py index a8918f1e4..13db1f03f 100644 --- a/packages/core/minos-microservice-common/tests/test_common/test_builders.py +++ b/packages/core/minos-microservice-common/tests/test_common/test_builders.py @@ -1,61 +1,152 @@ import unittest -from abc import ( - ABC, -) from typing import ( - Any, + Generic, + TypeVar, ) from minos.common import ( + BuildableMixin, Builder, Config, - SetupMixin, ) from tests.utils import ( CONFIG_FILE_PATH, ) -class _Builder(Builder[dict[str, Any]]): - def build(self) -> dict[str, Any]: - """For testing purposes.""" - return self.kwargs - - class TestBuilder(unittest.TestCase): - def test_abstract(self): - self.assertTrue(issubclass(Builder, (ABC, SetupMixin))) - # noinspection PyUnresolvedReferences - self.assertEqual({"build"}, Builder.__abstractmethods__) - def test_new(self): - builder = _Builder.new() - self.assertIsInstance(builder, _Builder) + builder = Builder.new().with_cls(dict) + self.assertIsInstance(builder, Builder) self.assertEqual(dict(), builder.kwargs) def test_copy(self): - base = _Builder.new().with_kwargs({"one": "two"}) + base = Builder.new().with_cls(dict).with_kwargs({"one": "two"}) builder = base.copy() self.assertNotEqual(id(base), id(builder)) - self.assertIsInstance(builder, _Builder) + self.assertIsInstance(builder, Builder) self.assertEqual({"one": "two"}, builder.kwargs) + self.assertEqual(dict, builder.instance_cls) def test_with_kwargs(self): - builder = _Builder().with_kwargs({"foo": "bar"}) - self.assertIsInstance(builder, _Builder) + builder = Builder().with_cls(dict).with_kwargs({"foo": "bar"}) + self.assertIsInstance(builder, Builder) self.assertEqual({"foo": "bar"}, builder.kwargs) def test_with_config(self): config = Config(CONFIG_FILE_PATH) - builder = _Builder().with_config(config) - self.assertIsInstance(builder, _Builder) + builder = Builder().with_cls(dict).with_config(config) + self.assertIsInstance(builder, Builder) self.assertEqual(dict(), builder.kwargs) def test_build(self): - builder = _Builder().with_kwargs({"one": "two"}) - self.assertIsInstance(builder, _Builder) + builder = Builder().with_cls(dict).with_kwargs({"one": "two"}) + self.assertIsInstance(builder, Builder) self.assertEqual({"one": "two"}, builder.build()) + def test_str(self): + builder = Builder().with_cls(dict).with_kwargs({"one": "two"}) + + self.assertEqual("Builder(dict, {'one': 'two'})", repr(builder)) + + def test_cmp(self): + base = Builder().with_cls(dict).with_kwargs({"one": "two"}) + + one = Builder().with_cls(dict).with_kwargs({"one": "two"}) + self.assertEqual(base, one) + + two = Builder().with_cls(int).with_kwargs({"one": "two"}) + self.assertNotEqual(base, two) + + three = Builder().with_cls(dict).with_kwargs({"three": "four"}) + self.assertNotEqual(base, three) + + def test_instance_cls_from_generic(self): + class _Builder(Builder): + """For Testing purposes.""" + + class _Builder2(Builder[int]): + """For Testing purposes.""" + + class _Builder3(_Builder2): + """For Testing purposes.""" + + T = TypeVar("T") + + class _Builder4(_Builder2, Generic[T]): + """For Testing purposes.""" + + class _Builder5(_Builder4[float]): + """For Testing purposes.""" + + self.assertEqual(None, _Builder().instance_cls) + self.assertEqual(int, _Builder2().instance_cls) + self.assertEqual(int, _Builder3().instance_cls) + self.assertEqual(None, _Builder4().instance_cls) + self.assertEqual(float, _Builder5().instance_cls) + + +class TestBuildableMixin(unittest.TestCase): + def test_get_builder_default(self): + class _Foo(BuildableMixin): + """For Testing purposes.""" + + self.assertEqual(Builder().with_cls(_Foo), _Foo.get_builder()) + + def test_get_builder_custom_type(self): + class _Foo(BuildableMixin): + """For Testing purposes.""" + + class _Builder(Builder): + """For Testing purposes.""" + + _Foo.set_builder(_Builder) + + self.assertEqual(_Builder().with_cls(_Foo), _Foo.get_builder()) + + def test_get_builder_custom_instance(self): + class _Foo(BuildableMixin): + """For Testing purposes.""" + + class _Builder(Builder): + """For Testing purposes.""" + + _Foo.set_builder(_Builder().with_kwargs({"foo": "bar"})) + + self.assertEqual(_Builder().with_cls(_Foo).with_kwargs({"foo": "bar"}), _Foo.get_builder()) + + def test_set_builder_raises(self): + class _Foo(BuildableMixin): + """For Testing purposes.""" + + with self.assertRaises(ValueError): + # noinspection PyTypeChecker + _Foo.set_builder(int) + + def test_from_config(self): + class _Foo(BuildableMixin): + """For Testing purposes.""" + + # noinspection PyShadowingNames + def __init__(self, config): + super().__init__() + self.config = config + + class _Builder(Builder): + """For Testing purposes.""" + + def with_config(self, config: Config): + """For Testing purposes.""" + self.kwargs["config"] = config + return super().with_config(config) + + _Foo.set_builder(_Builder) + + config = Config(CONFIG_FILE_PATH) + foo = _Foo.from_config(config) + + self.assertEqual(config, foo.config) + if __name__ == "__main__": unittest.main() 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 5a583bedc..0280cbc34 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 @@ -65,7 +65,7 @@ def test_injections(self): FakeBrokerClientPool, FakeHttpConnector, FakeBrokerPublisher, - FakeBrokerSubscriberBuilder, + FakeBrokerSubscriberBuilder(FakeBrokerSubscriber), FakeEventRepository, FakeSnapshotRepository, FakeTransactionRepository, diff --git a/packages/core/minos-microservice-common/tests/utils.py b/packages/core/minos-microservice-common/tests/utils.py index b28a24733..b2f386dff 100644 --- a/packages/core/minos-microservice-common/tests/utils.py +++ b/packages/core/minos-microservice-common/tests/utils.py @@ -132,9 +132,6 @@ class FakeBrokerPublisher(BuildableMixin): class FakeBrokerPublisherBuilder(Builder[FakeBrokerPublisher]): """For testing purposes.""" - def build(self) -> FakeBrokerPublisher: - return FakeBrokerPublisher() - FakeBrokerPublisher.set_builder(FakeBrokerPublisherBuilder) @@ -147,9 +144,6 @@ class FakeBrokerSubscriber(BuildableMixin): class FakeBrokerSubscriberBuilder(Builder[FakeBrokerSubscriber]): """For testing purposes.""" - def build(self) -> FakeBrokerSubscriber: - return FakeBrokerSubscriber() - FakeBrokerSubscriber.set_builder(FakeBrokerSubscriberBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/__init__.py b/packages/core/minos-microservice-networks/minos/networks/__init__.py index 427ae6ad5..bd0e70be5 100644 --- a/packages/core/minos-microservice-networks/minos/networks/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/__init__.py @@ -16,6 +16,7 @@ BrokerMessageV1Status, BrokerMessageV1Strategy, BrokerPublisher, + BrokerPublisherBuilder, BrokerPublisherQueue, BrokerQueue, BrokerRequest, @@ -38,7 +39,9 @@ PostgreSqlBrokerPublisherQueue, PostgreSqlBrokerPublisherQueueQueryFactory, PostgreSqlBrokerQueue, + PostgreSqlBrokerQueueBuilder, PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, PostgreSqlBrokerSubscriberQueue, PostgreSqlBrokerSubscriberQueueBuilder, 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 6681a48fd..382de106f 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/__init__.py @@ -5,6 +5,7 @@ BrokerQueue, InMemoryBrokerQueue, PostgreSqlBrokerQueue, + PostgreSqlBrokerQueueBuilder, ) from .dispatchers import ( BrokerDispatcher, @@ -30,6 +31,7 @@ ) from .publishers import ( BrokerPublisher, + BrokerPublisherBuilder, BrokerPublisherQueue, InMemoryBrokerPublisher, InMemoryBrokerPublisherQueue, @@ -50,6 +52,7 @@ InMemoryBrokerSubscriberQueue, InMemoryBrokerSubscriberQueueBuilder, PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, PostgreSqlBrokerSubscriberQueue, PostgreSqlBrokerSubscriberQueueBuilder, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/__init__.py index bef343971..924871e4d 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/__init__.py @@ -2,5 +2,6 @@ BrokerQueue, InMemoryBrokerQueue, PostgreSqlBrokerQueue, + PostgreSqlBrokerQueueBuilder, PostgreSqlBrokerQueueQueryFactory, ) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/__init__.py index a889b2ff7..6bba1e258 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/__init__.py @@ -6,5 +6,6 @@ ) from .pg import ( PostgreSqlBrokerQueue, + PostgreSqlBrokerQueueBuilder, PostgreSqlBrokerQueueQueryFactory, ) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/abc.py index a63082438..6d81a058b 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/abc.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/abc.py @@ -12,7 +12,7 @@ ) from minos.common import ( - SetupMixin, + BuildableMixin, ) from ...messages import ( @@ -22,7 +22,7 @@ logger = logging.getLogger(__name__) -class BrokerQueue(ABC, SetupMixin): +class BrokerQueue(ABC, BuildableMixin): """Broker Queue class.""" async def enqueue(self, message: BrokerMessage) -> None: diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/pg.py b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/pg.py index d04a70973..fa3de8baf 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/pg.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/collections/queues/pg.py @@ -40,6 +40,7 @@ ) from ....utils import ( + Builder, consume_queue, ) from ...messages import ( @@ -351,3 +352,20 @@ def __lt__(self, other: Any) -> bool: return isinstance(other, type(self)) and self.data < other.data except Exception: return False + + +class PostgreSqlBrokerQueueBuilder(Builder): + """PostgreSql Broker Queue Builder class.""" + + def with_config(self, config: Config): + """Set config. + + :param config: The config to be set. + :return: This method return the builder instance. + """ + self.kwargs |= config.get_database_by_name("broker") + self.kwargs |= config.get_interface_by_name("broker").get("common", dict()).get("queue", dict()) + return super().with_config(config) + + +PostgreSqlBrokerQueue.set_builder(PostgreSqlBrokerQueueBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/__init__.py b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/__init__.py index 50253e07f..7d3037037 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/__init__.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/__init__.py @@ -1,5 +1,6 @@ from .abc import ( BrokerPublisher, + BrokerPublisherBuilder, ) from .memory import ( InMemoryBrokerPublisher, diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/abc.py index 44c033b7e..89cc8e7d7 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/abc.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/abc.py @@ -1,23 +1,44 @@ +from __future__ import ( + annotations, +) + import logging from abc import ( ABC, abstractmethod, ) +from typing import ( + TYPE_CHECKING, + Any, + Generic, + Optional, + TypeVar, + Union, +) from minos.common import ( + BuildableMixin, + Builder, + Config, Injectable, - SetupMixin, + MinosConfigException, ) from ..messages import ( BrokerMessage, ) +if TYPE_CHECKING: + from .queued import ( + BrokerPublisherQueue, + QueuedBrokerPublisher, + ) + logger = logging.getLogger(__name__) @Injectable("broker_publisher") -class BrokerPublisher(ABC, SetupMixin): +class BrokerPublisher(ABC, BuildableMixin): """Broker Publisher class.""" async def send(self, message: BrokerMessage) -> None: @@ -32,3 +53,101 @@ async def send(self, message: BrokerMessage) -> None: @abstractmethod async def _send(self, message: BrokerMessage) -> None: raise NotImplementedError + + +BrokerPublisherCls = TypeVar("BrokerPublisherCls", bound=BrokerPublisher) + + +class BrokerPublisherBuilder(Builder[BrokerPublisher], Generic[BrokerPublisherCls]): + """Broker Publisher Builder class.""" + + def __init__( + self, + *args, + queue_builder: Optional[Builder] = None, + queued_cls: Optional[type[QueuedBrokerPublisher]] = None, + **kwargs, + ): + super().__init__(*args, **kwargs) + + if queued_cls is None: + from .queued import ( + QueuedBrokerPublisher, + ) + + queued_cls = QueuedBrokerPublisher + + self.queue_builder = queue_builder + + self.queued_cls = queued_cls + + def with_queued_cls(self, queued_cls: type[QueuedBrokerPublisher]): + """Set the queued class. + + :param queued_cls: A subclass of ``QueuedBrokerPublisher``. + :return: This method return the builder instance. + """ + self.queued_cls = queued_cls + + return self + + def with_config(self, config: Config): + """Set config. + + :param config: The config to be set. + :return: This method return the builder instance. + """ + self._with_builders_from_config(config) + + if self.queue_builder is not None: + self.queue_builder.with_config(config) + return super().with_config(config) + + def _with_builders_from_config(self, config): + try: + broker_config = config.get_interface_by_name("broker") + except MinosConfigException: + return + + broker_publisher_config = broker_config["publisher"] + + if "queue" in broker_publisher_config: + self.with_queue(broker_publisher_config["queue"]) + + def with_queue(self, queue: Union[type[BrokerPublisherQueue], Builder[BrokerPublisherQueue]]): + """Set the queue builder. + + :param queue: The queue builder to be set. + :return: This method return the builder instance. + """ + if not isinstance(queue, Builder): + queue = queue.get_builder() + self.queue_builder = queue.copy() + return self + + def with_kwargs(self, kwargs: dict[str, Any]): + """Set kwargs. + + :param kwargs: The kwargs to be set. + :return: This method return the builder instance. + """ + if self.queue_builder is not None: + self.queue_builder.with_kwargs(kwargs) + + return super().with_kwargs(kwargs) + + def build(self) -> BrokerPublisher: + """Build the instance. + + :return: A ``QueuedBrokerSubscriber`` instance. + """ + impl = super().build() + + if self.queue_builder is not None: + queue = self.queue_builder.build() + impl = self.queued_cls(impl=impl, queue=queue, **self.kwargs) + + return impl + + +BrokerPublisher.set_builder(BrokerPublisherBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/queued/impl.py b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/queued/impl.py index f4c1a44b3..2d709c2e1 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/queued/impl.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/publishers/queued/impl.py @@ -11,6 +11,10 @@ NoReturn, ) +from minos.common import ( + Builder, +) + from ...messages import ( BrokerMessage, ) @@ -65,3 +69,6 @@ async def _run(self) -> NoReturn: async def _send(self, message: BrokerMessage) -> None: await self.queue.enqueue(message) + + +QueuedBrokerPublisher.set_builder(Builder) 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 9726686d0..3b4bbdf56 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, + PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, ) from .memory import ( 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 ceff83283..a6e51b16e 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 @@ -12,19 +12,38 @@ Iterable, ) from typing import ( + TYPE_CHECKING, + Any, + Generic, Optional, + TypeVar, + Union, ) from minos.common import ( BuildableMixin, Builder, + Config, Injectable, + MinosConfigException, ) from ..messages import ( BrokerMessage, ) +if TYPE_CHECKING: + from .idempotent import ( + BrokerSubscriberDuplicateDetector, + IdempotentBrokerSubscriber, + ) + from .queued import ( + BrokerSubscriberQueue, + BrokerSubscriberQueueBuilder, + QueuedBrokerSubscriber, + ) + + logger = logging.getLogger(__name__) @@ -65,10 +84,131 @@ async def _receive(self) -> BrokerMessage: raise NotImplementedError +BrokerSubscriberCls = TypeVar("BrokerSubscriberCls", bound=BrokerSubscriber) + + @Injectable("broker_subscriber_builder") -class BrokerSubscriberBuilder(Builder[BrokerSubscriber], ABC): +class BrokerSubscriberBuilder(Builder[BrokerSubscriberCls], Generic[BrokerSubscriberCls]): """Broker Subscriber Builder class.""" + def __init__( + self, + *args, + idempotent_builder: Optional[Builder] = None, + queue_builder: Optional[BrokerSubscriberQueueBuilder] = None, + idempotent_cls: Optional[type[IdempotentBrokerSubscriber]] = None, + queued_cls: Optional[type[QueuedBrokerSubscriber]] = None, + **kwargs, + ): + super().__init__(*args, **kwargs) + + if idempotent_cls is None: + from .idempotent import ( + IdempotentBrokerSubscriber, + ) + + idempotent_cls = IdempotentBrokerSubscriber + + if queued_cls is None: + from .queued import ( + QueuedBrokerSubscriber, + ) + + queued_cls = QueuedBrokerSubscriber + + self.duplicate_detector_builder = idempotent_builder + self.queue_builder = queue_builder + + self.idempotent_cls = idempotent_cls + self.queued_cls = queued_cls + + def with_idempotent_cls(self, idempotent_cls: type[IdempotentBrokerSubscriber]): + """Set the idempotent class. + + :param idempotent_cls: A subclass of ``IdempotentBrokerSubscriber``. + :return: This method return the builder instance. + """ + self.idempotent_cls = idempotent_cls + + return self + + def with_queued_cls(self, queued_cls: type[QueuedBrokerSubscriber]): + """Set the queued class. + + :param queued_cls: A subclass of ``QueuedBrokerSubscriber``. + :return: This method return the builder instance. + """ + self.queued_cls = queued_cls + + return self + + def with_config(self, config: Config): + """Set config. + + :param config: The config to be set. + :return: This method return the builder instance. + """ + self._with_builders_from_config(config) + + if self.duplicate_detector_builder is not None: + self.duplicate_detector_builder.with_config(config) + if self.queue_builder is not None: + self.queue_builder.with_config(config) + return super().with_config(config) + + def _with_builders_from_config(self, config): + try: + broker_config = config.get_interface_by_name("broker") + except MinosConfigException: + return + + broker_subscriber_config = broker_config["subscriber"] + + if "idempotent" in broker_subscriber_config: + self.with_duplicate_detector(broker_subscriber_config["idempotent"]) + + if "queue" in broker_subscriber_config: + self.with_queue(broker_subscriber_config["queue"]) + + def with_duplicate_detector( + self, + duplicate_detector: Union[type[BrokerSubscriberDuplicateDetector], Builder[BrokerSubscriberDuplicateDetector]], + ): + """Set the duplicate detector. + + :param duplicate_detector: 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() + return self + + def with_queue(self, queue: Union[type[BrokerSubscriberQueue], BrokerSubscriberQueueBuilder]): + """Set the queue builder. + + :param queue: The queue to be set. + :return: This method return the builder instance. + """ + if not isinstance(queue, Builder): + queue = queue.get_builder() + self.queue_builder = queue.copy() + return self + + def with_kwargs(self, kwargs: dict[str, Any]): + """Set kwargs. + + :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.queue_builder is not None: + self.queue_builder.with_kwargs(kwargs) + + return super().with_kwargs(kwargs) + def with_group_id(self, group_id: Optional[str]): """Set group_id. @@ -76,6 +216,7 @@ def with_group_id(self, group_id: Optional[str]): :return: This method return the builder instance. """ self.kwargs["group_id"] = group_id + return self def with_remove_topics_on_destroy(self, remove_topics_on_destroy: bool): @@ -85,6 +226,7 @@ def with_remove_topics_on_destroy(self, remove_topics_on_destroy: bool): :return: This method return the builder instance. """ self.kwargs["remove_topics_on_destroy"] = remove_topics_on_destroy + return self def with_topics(self, topics: Iterable[str]): @@ -93,8 +235,30 @@ def with_topics(self, topics: Iterable[str]): :param topics: The topics to be set. :return: This method return the builder instance. """ + topics = set(topics) self.kwargs["topics"] = set(topics) + + if self.queue_builder is not None: + self.queue_builder.with_topics(topics) + return self + def build(self) -> BrokerSubscriber: + """Build the instance. + + :return: A ``QueuedBrokerSubscriber`` instance. + """ + 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.queue_builder is not None: + queue = self.queue_builder.build() + impl = self.queued_cls(impl=impl, queue=queue, **self.kwargs) + + return impl + BrokerSubscriber.set_builder(BrokerSubscriberBuilder) 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 3cdac8c73..10ac806c8 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, + PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory, ) from .impl import ( 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 5bb630ea7..9a9e01451 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,5 +6,6 @@ ) from .pg import ( PostgreSqlBrokerSubscriberDuplicateDetector, + PostgreSqlBrokerSubscriberDuplicateDetectorBuilder, 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 f4d0bd4ff..154d5c3fa 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 ( - SetupMixin, + BuildableMixin, ) from ....messages import ( @@ -19,7 +19,7 @@ ) -class BrokerSubscriberDuplicateDetector(ABC, SetupMixin): +class BrokerSubscriberDuplicateDetector(BuildableMixin, ABC): """Broker Subscriber Duplicate Detector class.""" async def is_valid(self, message: BrokerMessage) -> bool: 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 76f253962..079219814 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 @@ -17,6 +17,7 @@ ) from minos.common import ( + Builder, Config, PostgreSqlMinosDatabase, ) @@ -37,10 +38,6 @@ def __init__( super().__init__(*args, **kwargs) self._query_factory = query_factory - @classmethod - def _from_config(cls, config: Config, **kwargs) -> PostgreSqlBrokerSubscriberDuplicateDetector: - return cls(**config.get_database_by_name("broker"), **kwargs) - async def _setup(self) -> None: await super()._setup() await self._create_table() @@ -71,6 +68,22 @@ async def _is_valid(self, topic: str, uuid: UUID) -> bool: return False +class PostgreSqlBrokerSubscriberDuplicateDetectorBuilder(Builder[PostgreSqlBrokerSubscriberDuplicateDetector]): + """PostgreSql Broker Subscriber Duplicate Detector Builder class.""" + + def with_config(self, config: Config): + """Set config. + + :param config: The config to be set. + :return: This method return the builder instance. + """ + self.kwargs |= config.get_database_by_name("broker") + return super().with_config(config) + + +PostgreSqlBrokerSubscriberDuplicateDetector.set_builder(PostgreSqlBrokerSubscriberDuplicateDetectorBuilder) + + class PostgreSqlBrokerSubscriberDuplicateDetectorQueryFactory: """PostgreSql Broker Subscriber Duplicate Detector Query Factory class.""" 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 63a10cd00..93d3511f9 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 @@ -1,3 +1,7 @@ +from minos.common import ( + Builder, +) + from ...messages import ( BrokerMessage, ) @@ -37,3 +41,6 @@ async def _receive(self) -> BrokerMessage: while message is _sentinel or not (await self.duplicate_detector.is_valid(message)): message = await self.impl.receive() return message + + +IdempotentBrokerSubscriber.set_builder(Builder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/memory.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/memory.py index d0cdc3a85..135da099d 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/memory.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/memory.py @@ -41,7 +41,7 @@ async def _receive(self) -> BrokerMessage: return await self._queue.get() -class InMemoryBrokerSubscriberBuilder(BrokerSubscriberBuilder): +class InMemoryBrokerSubscriberBuilder(BrokerSubscriberBuilder[InMemoryBrokerSubscriber]): """In Memory Broker Subscriber Builder class.""" def with_messages(self, messages: Iterable[BrokerMessage]) -> InMemoryBrokerSubscriberBuilder: @@ -53,9 +53,5 @@ def with_messages(self, messages: Iterable[BrokerMessage]) -> InMemoryBrokerSubs self.kwargs["messages"] = messages return self - def build(self) -> BrokerSubscriber: - """Build the instance. - :return: An ``InMemoryBrokerSubscriber`` instance. - """ - return InMemoryBrokerSubscriber(**self.kwargs) +InMemoryBrokerSubscriber.set_builder(InMemoryBrokerSubscriberBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/impl.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/impl.py index 607365431..c6fbd6737 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/impl.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/impl.py @@ -1,3 +1,8 @@ +from __future__ import ( + annotations, +) + +import warnings from asyncio import ( CancelledError, TimeoutError, @@ -81,17 +86,19 @@ def _receive(self) -> Awaitable[BrokerMessage]: return self.queue.dequeue() -class QueuedBrokerSubscriberBuilder(BrokerSubscriberBuilder): +class QueuedBrokerSubscriberBuilder(BrokerSubscriberBuilder[QueuedBrokerSubscriber]): """Queued Broker Subscriber Publisher class.""" def __init__( self, *args, impl_builder: BrokerSubscriberBuilder, queue_builder: BrokerSubscriberQueueBuilder, **kwargs ): + warnings.warn(f"{type(self)!r} has been deprecated. Use {BrokerSubscriberBuilder} instead.", DeprecationWarning) + super().__init__(*args, **kwargs) self.impl_builder = impl_builder self.queue_builder = queue_builder - def with_config(self, config: Config) -> BrokerSubscriberBuilder: + def with_config(self, config: Config) -> QueuedBrokerSubscriberBuilder: """Set config. :param config: The config to be set. @@ -99,9 +106,9 @@ def with_config(self, config: Config) -> BrokerSubscriberBuilder: """ self.impl_builder.with_config(config) self.queue_builder.with_config(config) - return super().with_config(config) + return self - def with_kwargs(self, kwargs: dict[str, Any]) -> BrokerSubscriberBuilder: + def with_kwargs(self, kwargs: dict[str, Any]) -> QueuedBrokerSubscriberBuilder: """Set kwargs. :param kwargs: The kwargs to be set. @@ -109,9 +116,9 @@ def with_kwargs(self, kwargs: dict[str, Any]) -> BrokerSubscriberBuilder: """ self.impl_builder.with_kwargs(kwargs) self.queue_builder.with_kwargs(kwargs) - return super().with_kwargs(kwargs) + return self - def with_topics(self, topics: Iterable[str]) -> BrokerSubscriberBuilder: + def with_topics(self, topics: Iterable[str]) -> QueuedBrokerSubscriberBuilder: """Set topics. :param topics: The topics to be set. @@ -120,27 +127,27 @@ def with_topics(self, topics: Iterable[str]) -> BrokerSubscriberBuilder: topics = set(topics) self.impl_builder.with_topics(topics) self.queue_builder.with_topics(topics) - return super().with_topics(topics) + return self - def with_group_id(self, group_id: Optional[str]) -> BrokerSubscriberBuilder: + def with_group_id(self, group_id: Optional[str]) -> QueuedBrokerSubscriberBuilder: """Set group_id. :param group_id: The group_id to be set. :return: This method return the builder instance. """ self.impl_builder.with_group_id(group_id) - return super().with_group_id(group_id) + return self - def with_remove_topics_on_destroy(self, remove_topics_on_destroy: bool) -> BrokerSubscriberBuilder: + def with_remove_topics_on_destroy(self, remove_topics_on_destroy: bool) -> QueuedBrokerSubscriberBuilder: """Set remove_topics_on_destroy. :param remove_topics_on_destroy: The remove_topics_on_destroy flag to be set. :return: This method return the builder instance. """ self.impl_builder.with_remove_topics_on_destroy(remove_topics_on_destroy) - return super().with_remove_topics_on_destroy(remove_topics_on_destroy) + return self - def build(self) -> BrokerSubscriber: + def build(self) -> QueuedBrokerSubscriber: """Build the instance. :return: A ``QueuedBrokerSubscriber`` instance. diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/abc.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/abc.py index 55bc70553..2e933ee97 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/abc.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/abc.py @@ -10,6 +10,7 @@ Iterable, ) from typing import ( + Generic, TypeVar, ) @@ -43,7 +44,10 @@ def topics(self) -> set[str]: return self._topics -class BrokerSubscriberQueueBuilder(Builder[BrokerSubscriberQueue], ABC): +BrokerSubscriberQueueCls = TypeVar("BrokerSubscriberQueueCls", bound=BrokerSubscriberQueue) + + +class BrokerSubscriberQueueBuilder(Builder[BrokerSubscriberQueueCls], Generic[BrokerSubscriberQueueCls]): """Broker Subscriber Queue Builder class.""" def with_topics(self: B, topics: Iterable[str]) -> B: @@ -56,4 +60,6 @@ def with_topics(self: B, topics: Iterable[str]) -> B: return self +BrokerSubscriberQueue.set_builder(BrokerSubscriberQueueBuilder) + B = TypeVar("B", bound=BrokerSubscriberQueueBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/memory.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/memory.py index fe0b900a7..1d8e5c474 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/memory.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/memory.py @@ -19,12 +19,8 @@ class InMemoryBrokerSubscriberQueue(InMemoryBrokerQueue, BrokerSubscriberQueue): """In Memory Broker Subscriber Queue class.""" -class InMemoryBrokerSubscriberQueueBuilder(BrokerSubscriberQueueBuilder): +class InMemoryBrokerSubscriberQueueBuilder(BrokerSubscriberQueueBuilder[InMemoryBrokerSubscriberQueue]): """In Memory Broker Subscriber Queue Builder class.""" - def build(self) -> BrokerSubscriberQueue: - """Build the instance. - :return: An ``InMemoryBrokerSubscriberQueue`` instance. - """ - return InMemoryBrokerSubscriberQueue(**self.kwargs) +InMemoryBrokerSubscriberQueue.set_builder(InMemoryBrokerSubscriberQueueBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/pg.py b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/pg.py index 239b6522e..1cd1cdeec 100644 --- a/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/pg.py +++ b/packages/core/minos-microservice-networks/minos/networks/brokers/subscribers/queued/queues/pg.py @@ -16,12 +16,9 @@ Identifier, ) -from minos.common import ( - Config, -) - from ....collections import ( PostgreSqlBrokerQueue, + PostgreSqlBrokerQueueBuilder, PostgreSqlBrokerQueueQueryFactory, ) from ....messages import ( @@ -127,22 +124,10 @@ def build_select_not_processed(self) -> SQL: ) -class PostgreSqlBrokerSubscriberQueueBuilder(BrokerSubscriberQueueBuilder): +class PostgreSqlBrokerSubscriberQueueBuilder( + BrokerSubscriberQueueBuilder[PostgreSqlBrokerSubscriberQueue], PostgreSqlBrokerQueueBuilder +): """PostgreSql Broker Subscriber Queue Builder class.""" - def with_config(self, config: Config): - """Set config. - - :param config: The config to be set. - :return: This method return the builder instance. - """ - self.kwargs |= config.get_database_by_name("broker") - self.kwargs |= config.get_interface_by_name("broker").get("common", dict()).get("queue", dict()) - return super().with_config(config) - - def build(self) -> PostgreSqlBrokerSubscriberQueue: - """Build the instance. - :return: A ``BrokerSubscriberQueue`` instance. - """ - return PostgreSqlBrokerSubscriberQueue(**self.kwargs) +PostgreSqlBrokerSubscriberQueue.set_builder(PostgreSqlBrokerSubscriberQueueBuilder) diff --git a/packages/core/minos-microservice-networks/minos/networks/utils.py b/packages/core/minos-microservice-networks/minos/networks/utils.py index f4a111ee4..d2984217d 100644 --- a/packages/core/minos-microservice-networks/minos/networks/utils.py +++ b/packages/core/minos-microservice-networks/minos/networks/utils.py @@ -5,9 +5,6 @@ import re import socket import warnings -from abc import ( - ABC, -) from asyncio import ( QueueEmpty, ) @@ -59,7 +56,7 @@ async def consume_queue(queue, max_count: int) -> None: break -class Builder(CommonBuilder, ABC): +class Builder(CommonBuilder): """Builder class.""" def __init__(self, *args, **kwargs): diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_publishers/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_publishers/test_abc.py index d6032c2c1..ff6944caf 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_publishers/test_abc.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_publishers/test_abc.py @@ -4,10 +4,14 @@ ) from unittest.mock import ( AsyncMock, + MagicMock, call, ) from minos.common import ( + Builder, + Config, + MinosConfigException, SetupMixin, ) from minos.networks import ( @@ -15,6 +19,13 @@ BrokerMessageV1, BrokerMessageV1Payload, BrokerPublisher, + BrokerPublisherBuilder, + InMemoryBrokerPublisher, + InMemoryBrokerPublisherQueue, + QueuedBrokerPublisher, +) +from tests.utils import ( + CONFIG_FILE_PATH, ) @@ -42,5 +53,96 @@ async def test_send(self): self.assertEqual([call(message)], mock.call_args_list) +class TestBrokerPublisherBuilder(unittest.TestCase): + def test_constructor(self): + builder = BrokerPublisherBuilder() + self.assertEqual(None, builder.queue_builder) + self.assertEqual(QueuedBrokerPublisher, builder.queued_cls) + + def test_with_queued_cls(self): + # noinspection PyTypeChecker + builder = BrokerPublisherBuilder().with_queued_cls(int) + self.assertEqual(int, builder.queued_cls) + + def test_constructor_with_queue_builder(self): + queue_builder = Builder().with_cls(InMemoryBrokerPublisherQueue) + builder = BrokerPublisherBuilder(queue_builder=queue_builder) + self.assertEqual(queue_builder, builder.queue_builder) + self.assertEqual(QueuedBrokerPublisher, builder.queued_cls) + + def test_with_config_none(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock(side_effect=MinosConfigException("")) + config.get_interface_by_name = mock + + builder = BrokerPublisherBuilder().with_config(config) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({}, builder.kwargs) + + def test_with_config_empty(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock(return_value={"publisher": {}}) + config.get_interface_by_name = mock + + builder = BrokerPublisherBuilder().with_config(config) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({}, builder.kwargs) + + def test_with_config(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock(return_value={"publisher": {"queue": InMemoryBrokerPublisherQueue}}) + config.get_interface_by_name = mock + + builder = BrokerPublisherBuilder().with_config(config) + self.assertEqual(Builder().with_cls(InMemoryBrokerPublisherQueue), builder.queue_builder) + self.assertEqual({}, builder.kwargs) + + def test_with_queue_with_config(self): + config = Config(CONFIG_FILE_PATH) + + builder = BrokerPublisherBuilder().with_queue(InMemoryBrokerPublisherQueue).with_config(config) + self.assertEqual({}, builder.kwargs) + self.assertEqual(Builder().with_cls(InMemoryBrokerPublisherQueue), builder.queue_builder) + + def test_with_kwargs(self): + builder = BrokerPublisherBuilder().with_kwargs({"foo": "bar"}) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({"foo": "bar"}, builder.kwargs) + + def test_with_queue_with_kwargs(self): + builder = BrokerPublisherBuilder().with_queue(InMemoryBrokerPublisherQueue).with_kwargs({"foo": "bar"}) + self.assertEqual( + Builder().with_cls(InMemoryBrokerPublisherQueue).with_kwargs({"foo": "bar"}), builder.queue_builder + ) + self.assertEqual({"foo": "bar"}, builder.kwargs) + + def test_with_queue_cls(self): + queue_builder = Builder().with_cls(InMemoryBrokerPublisherQueue) + builder = BrokerPublisherBuilder().with_queue(InMemoryBrokerPublisherQueue) + self.assertEqual(queue_builder, builder.queue_builder) + + def test_with_queue_builder(self): + queue_builder = Builder().with_cls(InMemoryBrokerPublisherQueue) + builder = BrokerPublisherBuilder().with_queue(queue_builder) + self.assertEqual(queue_builder, builder.queue_builder) + + def test_build(self): + publisher = BrokerPublisherBuilder().with_cls(InMemoryBrokerPublisher).build() + + self.assertIsInstance(publisher, InMemoryBrokerPublisher) + + def test_build_with_queue(self): + publisher = ( + BrokerPublisherBuilder().with_cls(InMemoryBrokerPublisher).with_queue(InMemoryBrokerPublisherQueue).build() + ) + + self.assertIsInstance(publisher, QueuedBrokerPublisher) + self.assertIsInstance(publisher.impl, InMemoryBrokerPublisher) + self.assertIsInstance(publisher.queue, InMemoryBrokerPublisherQueue) + + if __name__ == "__main__": unittest.main() 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 a8919f37a..aa1631d61 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 @@ -4,9 +4,13 @@ ) from unittest.mock import ( AsyncMock, + MagicMock, ) from minos.common import ( + Builder, + Config, + MinosConfigException, SetupMixin, ) from minos.networks import ( @@ -15,16 +19,21 @@ BrokerMessageV1Payload, BrokerSubscriber, BrokerSubscriberBuilder, + IdempotentBrokerSubscriber, + InMemoryBrokerSubscriber, + InMemoryBrokerSubscriberDuplicateDetector, + InMemoryBrokerSubscriberQueue, + InMemoryBrokerSubscriberQueueBuilder, + QueuedBrokerSubscriber, +) +from tests.utils import ( + CONFIG_FILE_PATH, ) - - -class _BrokerSubscriberBuilder(BrokerSubscriberBuilder): - def build(self) -> BrokerSubscriber: - """For testing purposes.""" - return _BrokerSubscriber(**self.kwargs) class _BrokerSubscriber(BrokerSubscriber): + """For testing purposes.""" + async def _receive(self) -> BrokerMessage: """For testing purposes.""" @@ -70,32 +79,203 @@ async def test_aiter(self): class TestBrokerSubscriberBuilder(unittest.TestCase): - def test_abstract(self): - self.assertTrue(issubclass(BrokerSubscriberBuilder, (ABC, SetupMixin))) - # noinspection PyUnresolvedReferences - self.assertEqual({"build"}, BrokerSubscriberBuilder.__abstractmethods__) + def test_constructor(self): + builder = BrokerSubscriberBuilder() + self.assertEqual(None, builder.queue_builder) + self.assertEqual(None, builder.duplicate_detector_builder) + self.assertEqual(QueuedBrokerSubscriber, builder.queued_cls) + + def test_with_queued_cls(self): + # noinspection PyTypeChecker + builder = BrokerSubscriberBuilder().with_queued_cls(int) + self.assertEqual(int, builder.queued_cls) + + def test_with_idempotent_cls(self): + # noinspection PyTypeChecker + builder = BrokerSubscriberBuilder().with_idempotent_cls(int) + self.assertEqual(int, builder.idempotent_cls) + + def test_constructor_with_queue_builder(self): + queue_builder = InMemoryBrokerSubscriberQueueBuilder() + builder = BrokerSubscriberBuilder(queue_builder=queue_builder) + 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) + self.assertEqual(QueuedBrokerSubscriber, builder.queued_cls) + + def test_with_config_none(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock(side_effect=MinosConfigException("")) + config.get_interface_by_name = mock + + builder = BrokerSubscriberBuilder().with_config(config) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({}, builder.kwargs) + + def test_with_config_empty(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock(return_value={"subscriber": {}}) + config.get_interface_by_name = mock + + builder = BrokerSubscriberBuilder().with_config(config) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({}, builder.kwargs) + + def test_with_config(self): + config = Config(CONFIG_FILE_PATH) + + mock = MagicMock( + return_value={ + "subscriber": { + "queue": InMemoryBrokerSubscriberQueue, + "idempotent": InMemoryBrokerSubscriberDuplicateDetector, + } + } + ) + config.get_interface_by_name = mock + + builder = BrokerSubscriberBuilder().with_config(config) + self.assertEqual(InMemoryBrokerSubscriberQueueBuilder(), builder.queue_builder) + self.assertEqual( + Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector), builder.duplicate_detector_builder + ) + self.assertEqual({}, builder.kwargs) + + def test_with_queue_with_config(self): + config = Config(CONFIG_FILE_PATH) + + builder = BrokerSubscriberBuilder().with_queue(InMemoryBrokerSubscriberQueue).with_config(config) + self.assertEqual({}, builder.kwargs) + self.assertEqual(InMemoryBrokerSubscriberQueueBuilder(), builder.queue_builder) + + def test_with_duplicate_with_config(self): + config = Config(CONFIG_FILE_PATH) + + builder = ( + BrokerSubscriberBuilder() + .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_config(config) + ) + self.assertEqual({}, builder.kwargs) + self.assertEqual( + Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector), builder.duplicate_detector_builder + ) + + def test_with_kwargs(self): + builder = BrokerSubscriberBuilder().with_kwargs({"foo": "bar"}) + self.assertEqual(None, builder.queue_builder) + self.assertEqual({"foo": "bar"}, builder.kwargs) + + def test_with_queue_with_kwargs(self): + builder = BrokerSubscriberBuilder().with_queue(InMemoryBrokerSubscriberQueue).with_kwargs({"foo": "bar"}) + self.assertEqual(InMemoryBrokerSubscriberQueueBuilder().with_kwargs({"foo": "bar"}), builder.queue_builder) + self.assertEqual({"foo": "bar"}, builder.kwargs) + + def test_with_duplicate_detector_with_kwargs(self): + builder = ( + BrokerSubscriberBuilder() + .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_kwargs({"foo": "bar"}) + ) + self.assertEqual( + Builder().with_cls(InMemoryBrokerSubscriberDuplicateDetector).with_kwargs({"foo": "bar"}), + builder.duplicate_detector_builder, + ) + self.assertEqual({"foo": "bar"}, builder.kwargs) + + def test_with_queue_cls(self): + queue_builder = InMemoryBrokerSubscriberQueueBuilder() + builder = BrokerSubscriberBuilder().with_queue(InMemoryBrokerSubscriberQueue) + self.assertEqual(queue_builder, builder.queue_builder) + + def test_with_queue_builder(self): + queue_builder = InMemoryBrokerSubscriberQueueBuilder() + 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_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_build(self): + subscriber = BrokerSubscriberBuilder().with_topics({"one", "two"}).with_cls(InMemoryBrokerSubscriber).build() + + self.assertIsInstance(subscriber, InMemoryBrokerSubscriber) + + def test_build_with_queue(self): + subscriber = ( + BrokerSubscriberBuilder() + .with_cls(InMemoryBrokerSubscriber) + .with_queue(InMemoryBrokerSubscriberQueue) + .with_topics({"one", "two"}) + .build() + ) + + self.assertIsInstance(subscriber, QueuedBrokerSubscriber) + self.assertIsInstance(subscriber.impl, InMemoryBrokerSubscriber) + self.assertIsInstance(subscriber.queue, InMemoryBrokerSubscriberQueue) + + def test_build_with_duplicate_detector(self): + subscriber = ( + BrokerSubscriberBuilder() + .with_cls(InMemoryBrokerSubscriber) + .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_topics({"one", "two"}) + .build() + ) + + self.assertIsInstance(subscriber, IdempotentBrokerSubscriber) + self.assertIsInstance(subscriber.impl, InMemoryBrokerSubscriber) + self.assertIsInstance(subscriber.duplicate_detector, InMemoryBrokerSubscriberDuplicateDetector) + + def test_build_with_duplicate_detector_with_queue(self): + subscriber = ( + BrokerSubscriberBuilder() + .with_cls(InMemoryBrokerSubscriber) + .with_duplicate_detector(InMemoryBrokerSubscriberDuplicateDetector) + .with_queue(InMemoryBrokerSubscriberQueue) + .with_topics({"one", "two"}) + .build() + ) + self.assertIsInstance(subscriber, QueuedBrokerSubscriber) + self.assertIsInstance(subscriber.queue, InMemoryBrokerSubscriberQueue) + + self.assertIsInstance(subscriber.impl, IdempotentBrokerSubscriber) + self.assertIsInstance(subscriber.impl.impl, InMemoryBrokerSubscriber) + self.assertIsInstance(subscriber.impl.duplicate_detector, InMemoryBrokerSubscriberDuplicateDetector) def test_with_group_id(self): - builder = _BrokerSubscriberBuilder().with_group_id("foobar") - self.assertIsInstance(builder, _BrokerSubscriberBuilder) + builder = BrokerSubscriberBuilder().with_group_id("foobar") + self.assertIsInstance(builder, BrokerSubscriberBuilder) self.assertEqual({"group_id": "foobar"}, builder.kwargs) def test_with_remove_topics_on_destroy(self): - builder = _BrokerSubscriberBuilder().with_remove_topics_on_destroy(False) - self.assertIsInstance(builder, _BrokerSubscriberBuilder) + builder = BrokerSubscriberBuilder().with_remove_topics_on_destroy(False) + self.assertIsInstance(builder, BrokerSubscriberBuilder) self.assertEqual({"remove_topics_on_destroy": False}, builder.kwargs) def test_with_topics(self): - builder = _BrokerSubscriberBuilder().with_topics({"one", "two"}) - self.assertIsInstance(builder, _BrokerSubscriberBuilder) + builder = BrokerSubscriberBuilder().with_topics({"one", "two"}) + self.assertIsInstance(builder, BrokerSubscriberBuilder) self.assertEqual({"topics": {"one", "two"}}, builder.kwargs) - def test_build(self): - builder = _BrokerSubscriberBuilder().with_topics({"one", "two"}) - self.assertIsInstance(builder, _BrokerSubscriberBuilder) - subscriber = builder.build() - self.assertIsInstance(subscriber, _BrokerSubscriber) - self.assertEqual({"one", "two"}, subscriber.topics) + def test_with_topics_with_queue(self): + builder = BrokerSubscriberBuilder().with_queue(InMemoryBrokerSubscriberQueue).with_topics({"one", "two"}) + self.assertIsInstance(builder, BrokerSubscriberBuilder) + self.assertEqual(InMemoryBrokerSubscriberQueueBuilder().with_topics({"one", "two"}), builder.queue_builder) + self.assertEqual({"topics": {"one", "two"}}, builder.kwargs) if __name__ == "__main__": diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_impl.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_impl.py index 8eec8e3fd..2c0f2bd92 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_impl.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_impl.py @@ -1,4 +1,5 @@ import unittest +import warnings from asyncio import ( sleep, ) @@ -115,7 +116,9 @@ def test_with_kwargs(self): self.impl_builder.with_kwargs = impl_mock self.queue_builder.with_kwargs = queue_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_kwargs({"foo": "bar"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_kwargs({"foo": "bar"}) self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) self.assertEqual([call({"foo": "bar"})], impl_mock.call_args_list) @@ -127,7 +130,9 @@ def test_with_config(self): self.impl_builder.with_config = impl_mock self.queue_builder.with_config = queue_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_config(self.config) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_config(self.config) self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) self.assertEqual([call(self.config)], impl_mock.call_args_list) @@ -137,9 +142,11 @@ def test_with_group_id(self): impl_mock = MagicMock(side_effect=self.impl_builder.with_group_id) self.impl_builder.with_group_id = impl_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_group_id("foobar") + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_group_id("foobar") + self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) - self.assertEqual({"group_id": "foobar"}, builder.kwargs) self.assertEqual([call("foobar")], impl_mock.call_args_list) @@ -147,9 +154,10 @@ def test_with_remove_topics_on_destroy(self): impl_mock = MagicMock(side_effect=self.impl_builder.with_remove_topics_on_destroy) self.impl_builder.with_remove_topics_on_destroy = impl_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_remove_topics_on_destroy(False) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_remove_topics_on_destroy(False) self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) - self.assertEqual({"remove_topics_on_destroy": False}, builder.kwargs) self.assertEqual([call(False)], impl_mock.call_args_list) @@ -159,9 +167,10 @@ def test_with_topics(self): self.impl_builder.with_topics = impl_mock self.queue_builder.with_topics = queue_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_topics({"one", "two"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_topics({"one", "two"}) self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) - self.assertEqual({"topics": {"one", "two"}}, builder.kwargs) self.assertEqual([call({"one", "two"})], impl_mock.call_args_list) self.assertEqual([call({"one", "two"})], queue_mock.call_args_list) @@ -172,7 +181,9 @@ def test_build(self): self.impl_builder.build = impl_mock self.queue_builder.build = queue_mock - builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_topics({"one", "two"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = QueuedBrokerSubscriberBuilder(**self._kwargs).with_topics({"one", "two"}) self.assertIsInstance(builder, QueuedBrokerSubscriberBuilder) subscriber = builder.build() self.assertIsInstance(subscriber, QueuedBrokerSubscriber) diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_abc.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_abc.py index 4ac0151fe..5de2334a3 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_abc.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_abc.py @@ -3,9 +3,6 @@ ABC, ) -from minos.common import ( - SetupMixin, -) from minos.networks import ( BrokerMessage, BrokerQueue, @@ -24,14 +21,6 @@ async def _dequeue(self) -> BrokerMessage: """For testing purposes.""" -class _BrokerSubscriberQueueBuilder(BrokerSubscriberQueueBuilder): - """For testing purposes.""" - - def build(self) -> BrokerSubscriberQueue: - """For testing purposes.""" - return _BrokerSubscriberQueue(**self.kwargs) - - class TestBrokerSubscriberQueue(unittest.IsolatedAsyncioTestCase): def setUp(self) -> None: self.topics = {"foo", "bar"} @@ -50,20 +39,15 @@ def test_topics_raises(self): _BrokerSubscriberQueue([]) -class TestBrokerSubscriberBuilder(unittest.TestCase): - def test_abstract(self): - self.assertTrue(issubclass(BrokerSubscriberQueueBuilder, (ABC, SetupMixin))) - # noinspection PyUnresolvedReferences - self.assertEqual({"build"}, BrokerSubscriberQueueBuilder.__abstractmethods__) - +class TestBrokerSubscriberQueueBuilder(unittest.TestCase): def test_with_topics(self): - builder = _BrokerSubscriberQueueBuilder().with_topics({"one", "two"}) - self.assertIsInstance(builder, _BrokerSubscriberQueueBuilder) + builder = BrokerSubscriberQueueBuilder().with_topics({"one", "two"}) + self.assertIsInstance(builder, BrokerSubscriberQueueBuilder) self.assertEqual({"topics": {"one", "two"}}, builder.kwargs) def test_build(self): - builder = _BrokerSubscriberQueueBuilder().with_topics({"one", "two"}) - self.assertIsInstance(builder, _BrokerSubscriberQueueBuilder) + builder = BrokerSubscriberQueueBuilder().with_topics({"one", "two"}).with_cls(_BrokerSubscriberQueue) + self.assertIsInstance(builder, BrokerSubscriberQueueBuilder) subscriber = builder.build() self.assertIsInstance(subscriber, _BrokerSubscriberQueue) self.assertEqual({"one", "two"}, subscriber.topics) diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_pg.py b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_pg.py index bbf82b846..45e7c9786 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_pg.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_brokers/test_subscribers/test_queued/test_queues/test_pg.py @@ -1,4 +1,5 @@ import unittest +import warnings from asyncio import ( sleep, ) @@ -90,7 +91,9 @@ def setUp(self) -> None: self.config = Config(CONFIG_FILE_PATH) def test_build(self): - builder = PostgreSqlBrokerSubscriberQueueBuilder().with_config(self.config).with_topics({"one", "two"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = PostgreSqlBrokerSubscriberQueueBuilder().with_config(self.config).with_topics({"one", "two"}) subscriber = builder.build() self.assertIsInstance(subscriber, PostgreSqlBrokerSubscriberQueue) diff --git a/packages/core/minos-microservice-networks/tests/test_networks/test_utils.py b/packages/core/minos-microservice-networks/tests/test_networks/test_utils.py index 41eb5c0b5..5e4cc21fd 100644 --- a/packages/core/minos-microservice-networks/tests/test_networks/test_utils.py +++ b/packages/core/minos-microservice-networks/tests/test_networks/test_utils.py @@ -37,13 +37,9 @@ def test_is_subclass(self): self.assertTrue(issubclass(Builder, CommonBuilder)) def test_warnings(self): - class _Builder(Builder): - def build(self) -> None: - """For testing purpose""" - with warnings.catch_warnings(): warnings.simplefilter("ignore", DeprecationWarning) - builder = _Builder() + builder = Builder() self.assertIsInstance(builder, CommonBuilder) diff --git a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/__init__.py b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/__init__.py index d4dff278b..b6f11bf23 100644 --- a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/__init__.py +++ b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/__init__.py @@ -2,12 +2,14 @@ __email__ = "hey@minos.run" __version__ = "0.5.1" -from .mixins import ( +from .common import ( + KafkaBrokerBuilderMixin, KafkaCircuitBreakerMixin, ) from .publisher import ( InMemoryQueuedKafkaBrokerPublisher, KafkaBrokerPublisher, + KafkaBrokerPublisherBuilder, PostgreSqlQueuedKafkaBrokerPublisher, ) from .subscriber import ( diff --git a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/common.py b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/common.py new file mode 100644 index 000000000..3b6a19d45 --- /dev/null +++ b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/common.py @@ -0,0 +1,44 @@ +from __future__ import ( + annotations, +) + +from collections.abc import ( + Iterable, +) + +from kafka.errors import ( + KafkaError, +) + +from minos.common import ( + Builder, + CircuitBreakerMixin, + Config, +) + + +class KafkaCircuitBreakerMixin(CircuitBreakerMixin): + """Kafka Circuit Breaker Mixin class.""" + + def __init__(self, *args, circuit_breaker_exceptions: Iterable[type] = tuple(), **kwargs): + super().__init__(*args, circuit_breaker_exceptions=(KafkaError, *circuit_breaker_exceptions), **kwargs) + + +class KafkaBrokerBuilderMixin(Builder): + """Kafka Broker Builder Mixin class.""" + + def with_config(self, config: Config): + """Set config. + + :param config: The config to be set. + :return: This method return the builder instance. + """ + broker_config = config.get_interface_by_name("broker") + common_config = broker_config.get("common", dict()) + + self.kwargs |= { + "group_id": config.get_name(), + "host": common_config.get("host"), + "port": common_config.get("port"), + } + return super().with_config(config) diff --git a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/mixins.py b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/mixins.py deleted file mode 100644 index 1879f5ecb..000000000 --- a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/mixins.py +++ /dev/null @@ -1,18 +0,0 @@ -from collections.abc import ( - Iterable, -) - -from kafka.errors import ( - KafkaError, -) - -from minos.common import ( - CircuitBreakerMixin, -) - - -class KafkaCircuitBreakerMixin(CircuitBreakerMixin): - """Kafka Circuit Breaker Mixin class.""" - - def __init__(self, *args, circuit_breaker_exceptions: Iterable[type] = tuple(), **kwargs): - super().__init__(*args, circuit_breaker_exceptions=(KafkaError, *circuit_breaker_exceptions), **kwargs) diff --git a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/publisher.py b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/publisher.py index 44857db05..e2a1cc760 100644 --- a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/publisher.py +++ b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/publisher.py @@ -3,6 +3,7 @@ ) import logging +import warnings from asyncio import ( TimeoutError, wait_for, @@ -27,12 +28,14 @@ from minos.networks import ( BrokerMessage, BrokerPublisher, + BrokerPublisherBuilder, InMemoryBrokerPublisherQueue, PostgreSqlBrokerPublisherQueue, QueuedBrokerPublisher, ) -from .mixins import ( +from .common import ( + KafkaBrokerBuilderMixin, KafkaCircuitBreakerMixin, ) @@ -42,6 +45,10 @@ class PostgreSqlQueuedKafkaBrokerPublisher(QueuedBrokerPublisher): """PostgreSql Queued Kafka Broker Publisher class.""" + def __init__(self, *args, **kwargs): + warnings.warn(f"{PostgreSqlQueuedKafkaBrokerPublisher!r} has been deprecated.", DeprecationWarning) + super().__init__(*args, **kwargs) + @classmethod def _from_config(cls, config: Config, **kwargs) -> PostgreSqlQueuedKafkaBrokerPublisher: impl = KafkaBrokerPublisher.from_config(config, **kwargs) @@ -52,6 +59,10 @@ def _from_config(cls, config: Config, **kwargs) -> PostgreSqlQueuedKafkaBrokerPu class InMemoryQueuedKafkaBrokerPublisher(QueuedBrokerPublisher): """In Memory Queued Kafka Broker Publisher class.""" + def __init__(self, *args, **kwargs): + warnings.warn(f"{InMemoryQueuedKafkaBrokerPublisher!r} has been deprecated.", DeprecationWarning) + super().__init__(*args, **kwargs) + @classmethod def _from_config(cls, config: Config, **kwargs) -> InMemoryQueuedKafkaBrokerPublisher: impl = KafkaBrokerPublisher.from_config(config, **kwargs) @@ -92,14 +103,6 @@ def port(self) -> int: """ return self._port - @classmethod - def _from_config(cls, config: Config, **kwargs) -> KafkaBrokerPublisher: - broker_config = config.get_interface_by_name("broker") - common_config = broker_config.get("common", dict()) - kwargs["host"] = common_config.get("host") - kwargs["port"] = common_config.get("port") - return cls(**kwargs) - async def _setup(self) -> None: await super()._setup() await self._start_client() @@ -140,3 +143,10 @@ def _build_client(self) -> AIOKafkaProducer: @property def _bootstrap_servers(self): return f"{self.host}:{self.port}" + + +class KafkaBrokerPublisherBuilder(BrokerPublisherBuilder[KafkaBrokerPublisher], KafkaBrokerBuilderMixin): + """Kafka Broker Publisher Builder class.""" + + +KafkaBrokerPublisher.set_builder(KafkaBrokerPublisherBuilder) diff --git a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/subscriber.py b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/subscriber.py index 36937e54c..f63e20ef0 100644 --- a/packages/plugins/minos-broker-kafka/minos/plugins/kafka/subscriber.py +++ b/packages/plugins/minos-broker-kafka/minos/plugins/kafka/subscriber.py @@ -36,9 +36,6 @@ TopicAlreadyExistsError, ) -from minos.common import ( - Config, -) from minos.networks import ( BrokerMessage, BrokerSubscriber, @@ -48,7 +45,8 @@ QueuedBrokerSubscriberBuilder, ) -from .mixins import ( +from .common import ( + KafkaBrokerBuilderMixin, KafkaCircuitBreakerMixin, ) @@ -187,32 +185,9 @@ def client(self) -> AIOKafkaConsumer: ) -class KafkaBrokerSubscriberBuilder(BrokerSubscriberBuilder): +class KafkaBrokerSubscriberBuilder(BrokerSubscriberBuilder[KafkaBrokerSubscriber], KafkaBrokerBuilderMixin): """Kafka Broker Subscriber Builder class.""" - def with_config(self, config: Config) -> BrokerSubscriberBuilder: - """Set config. - - :param config: The config to be set. - :return: This method return the builder instance. - """ - broker_config = config.get_interface_by_name("broker") - common_config = broker_config.get("common", dict()) - - self.kwargs |= { - "group_id": config.get_name(), - "host": common_config.get("host"), - "port": common_config.get("port"), - } - return self - - def build(self) -> BrokerSubscriber: - """Build the instance. - - :return: A ``KafkaBrokerSubscriber`` instance. - """ - return KafkaBrokerSubscriber(**self.kwargs) - KafkaBrokerSubscriber.set_builder(KafkaBrokerSubscriberBuilder) diff --git a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_common.py b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_common.py new file mode 100644 index 000000000..a77050094 --- /dev/null +++ b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_common.py @@ -0,0 +1,43 @@ +import unittest + +from kafka.errors import ( + KafkaError, +) + +from minos.common import ( + Config, +) +from minos.plugins.kafka import ( + KafkaBrokerBuilderMixin, + KafkaCircuitBreakerMixin, +) +from tests.utils import ( + CONFIG_FILE_PATH, +) + + +class TestKafkaCircuitBreakerMixin(unittest.IsolatedAsyncioTestCase): + def test_constructor(self): + mixin = KafkaCircuitBreakerMixin() + self.assertEqual((KafkaError,), mixin.circuit_breaker_exceptions) + + +class TestKafkaBrokerBuilderMixin(unittest.IsolatedAsyncioTestCase): + def test_constructor(self): + mixin = KafkaBrokerBuilderMixin() + + config = Config(CONFIG_FILE_PATH) + mixin.with_config(config) + + common_config = config.get_interface_by_name("broker")["common"] + + expected = { + "group_id": config.get_name(), + "host": common_config["host"], + "port": common_config["port"], + } + self.assertEqual(expected, mixin.kwargs) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_mixins.py b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_mixins.py deleted file mode 100644 index 485ae199a..000000000 --- a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_mixins.py +++ /dev/null @@ -1,19 +0,0 @@ -import unittest - -from kafka.errors import ( - KafkaError, -) - -from minos.plugins.kafka import ( - KafkaCircuitBreakerMixin, -) - - -class TestKafkaCircuitBreakerMixin(unittest.IsolatedAsyncioTestCase): - def test_constructor(self): - mixin = KafkaCircuitBreakerMixin() - self.assertEqual((KafkaError,), mixin.circuit_breaker_exceptions) - - -if __name__ == "__main__": - unittest.main() diff --git a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_publisher.py b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_publisher.py index 4d3d5a8e5..c734086da 100644 --- a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_publisher.py +++ b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_publisher.py @@ -1,4 +1,5 @@ import unittest +import warnings from unittest.mock import ( AsyncMock, ) @@ -24,6 +25,7 @@ from minos.plugins.kafka import ( InMemoryQueuedKafkaBrokerPublisher, KafkaBrokerPublisher, + KafkaBrokerPublisherBuilder, PostgreSqlQueuedKafkaBrokerPublisher, ) from tests.utils import ( @@ -124,9 +126,36 @@ async def test_setup_destroy(self): self.assertEqual(1, stop_mock.call_count) +class TestKafkaBrokerPublisherBuilder(unittest.TestCase): + def setUp(self) -> None: + self.config = Config(CONFIG_FILE_PATH) + + def test_with_config(self): + builder = KafkaBrokerPublisherBuilder().with_config(self.config) + common_config = self.config.get_interface_by_name("broker")["common"] + + expected = { + "group_id": self.config.get_name(), + "host": common_config["host"], + "port": common_config["port"], + } + self.assertEqual(expected, builder.kwargs) + + def test_build(self): + common_config = self.config.get_interface_by_name("broker")["common"] + builder = KafkaBrokerPublisherBuilder().with_config(self.config) + publisher = builder.build() + + self.assertIsInstance(publisher, KafkaBrokerPublisher) + self.assertEqual(common_config["host"], publisher.host) + self.assertEqual(common_config["port"], publisher.port) + + class TestPostgreSqlQueuedKafkaBrokerPublisher(unittest.IsolatedAsyncioTestCase): def test_from_config(self): - publisher = PostgreSqlQueuedKafkaBrokerPublisher.from_config(CONFIG_FILE_PATH) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + publisher = PostgreSqlQueuedKafkaBrokerPublisher.from_config(CONFIG_FILE_PATH) self.assertIsInstance(publisher, PostgreSqlQueuedKafkaBrokerPublisher) self.assertIsInstance(publisher.impl, KafkaBrokerPublisher) self.assertIsInstance(publisher.queue, PostgreSqlBrokerPublisherQueue) @@ -134,7 +163,9 @@ def test_from_config(self): class TestInMemoryQueuedKafkaBrokerPublisher(unittest.IsolatedAsyncioTestCase): def test_from_config(self): - publisher = InMemoryQueuedKafkaBrokerPublisher.from_config(CONFIG_FILE_PATH) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + publisher = InMemoryQueuedKafkaBrokerPublisher.from_config(CONFIG_FILE_PATH) self.assertIsInstance(publisher, InMemoryQueuedKafkaBrokerPublisher) self.assertIsInstance(publisher.impl, KafkaBrokerPublisher) self.assertIsInstance(publisher.queue, InMemoryBrokerPublisherQueue) diff --git a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_subscriber.py b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_subscriber.py index 05d32cf7e..47a4b1900 100644 --- a/packages/plugins/minos-broker-kafka/tests/test_kafka/test_subscriber.py +++ b/packages/plugins/minos-broker-kafka/tests/test_kafka/test_subscriber.py @@ -1,4 +1,5 @@ import unittest +import warnings from collections import ( namedtuple, ) @@ -237,7 +238,12 @@ def setUp(self) -> None: self.config = Config(CONFIG_FILE_PATH) def test_build(self): - builder = PostgreSqlQueuedKafkaBrokerSubscriberBuilder().with_config(self.config).with_topics({"one", "two"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = ( + PostgreSqlQueuedKafkaBrokerSubscriberBuilder().with_config(self.config).with_topics({"one", "two"}) + ) + subscriber = builder.build() self.assertIsInstance(subscriber, QueuedBrokerSubscriber) @@ -250,7 +256,10 @@ def setUp(self) -> None: self.config = Config(CONFIG_FILE_PATH) def test_build(self): - builder = InMemoryQueuedKafkaBrokerSubscriberBuilder().with_config(self.config).with_topics({"one", "two"}) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", DeprecationWarning) + builder = InMemoryQueuedKafkaBrokerSubscriberBuilder().with_config(self.config).with_topics({"one", "two"}) + subscriber = builder.build() self.assertIsInstance(subscriber, QueuedBrokerSubscriber)