diff --git a/configs/daemons.json b/configs/daemons.json new file mode 100644 index 0000000..f438df0 --- /dev/null +++ b/configs/daemons.json @@ -0,0 +1,7 @@ +{ + "daemons": + [ + "welcome_email", + "reset_email" + ] +} \ No newline at end of file diff --git a/daemon_run.py b/daemon_run.py index 354f527..4d46501 100644 --- a/daemon_run.py +++ b/daemon_run.py @@ -1,11 +1,56 @@ import asyncio +import os +import sys -from src.messaging.consumers.consumers import email_consumer +from src.daemons.registry import DAEMONS +from src.models.configs_read.daemons_json import daemons_config + + +async def run_one(daemon_name: str) -> None: + daemon_cls = DAEMONS.get(daemon_name) + if daemon_cls is None: + print(f"Unknown daemon: {daemon_name}. Available: {list(DAEMONS.keys())}") + sys.exit(1) + daemon = daemon_cls() + await daemon.run() + + +async def run_enabled_from_config() -> None: + + daemons = [] + + for name in daemons_config.daemons: + name = name.strip() + if name not in DAEMONS: + print(f"Warning: unknown daemon '{name}' in config, skipping") + continue + daemons.append(DAEMONS[name]()) + + if not daemons: + print("No enabled daemons found in config") + sys.exit(1) + + await asyncio.gather(*(d.run() for d in daemons)) async def main() -> None: - await email_consumer.start_consuming() + if len(sys.argv) < 2: + print("Usage: python run_daemon.py | --all") + sys.exit(1) + + arg = sys.argv[1] + if arg == "--all": + await run_enabled_from_config() + else: + await run_one(arg) if __name__ == "__main__": - asyncio.run(main()) \ No newline at end of file + try: + asyncio.run(main()) + except KeyboardInterrupt: + print('Interrupted') + try: + sys.exit(0) + except SystemExit: + os._exit(0) \ No newline at end of file diff --git a/src/daemons/base.py b/src/daemons/base.py index e69de29..249ed79 100644 --- a/src/daemons/base.py +++ b/src/daemons/base.py @@ -0,0 +1,9 @@ +from abc import ABC, abstractmethod + + +class BaseDaemon(ABC): + name: str + + @abstractmethod + async def run(self) -> None: + ... \ No newline at end of file diff --git a/src/daemons/email_consumer.py b/src/daemons/email_consumer.py deleted file mode 100644 index e69de29..0000000 diff --git a/src/daemons/email_daemons.py b/src/daemons/email_daemons.py new file mode 100644 index 0000000..2dbeda5 --- /dev/null +++ b/src/daemons/email_daemons.py @@ -0,0 +1,19 @@ +from src.daemons.base import BaseDaemon +from src.messaging.consumers.consumers import ( + reset_email_consumer, + welcome_email_consumer, +) + + +class WelcomeEmailDaemon(BaseDaemon): + name = "welcome_email" + + async def run(self) -> None: + await welcome_email_consumer.start_consuming() + + +class ResetEmailDaemon(BaseDaemon): + name = "reset_email" + + async def run(self) -> None: + await reset_email_consumer.start_consuming() \ No newline at end of file diff --git a/src/daemons/registry.py b/src/daemons/registry.py index e69de29..580c33c 100644 --- a/src/daemons/registry.py +++ b/src/daemons/registry.py @@ -0,0 +1,6 @@ +from src.daemons.email_daemons import ResetEmailDaemon, WelcomeEmailDaemon + +DAEMONS = { + "welcome_email": WelcomeEmailDaemon, + "reset_email": ResetEmailDaemon, +} \ No newline at end of file diff --git a/src/messaging/consumers/consumers.py b/src/messaging/consumers/consumers.py index 6ac0451..068502d 100644 --- a/src/messaging/consumers/consumers.py +++ b/src/messaging/consumers/consumers.py @@ -1,9 +1,11 @@ import json +import aio_pika + from src.messaging.rabbitmq_client import rabbitmq_client -class EmailConsumer: +class WelcomeEmailConsumer: def __init__(self) -> None: self.channel = None @@ -13,23 +15,65 @@ class EmailConsumer: self.channel = await rabbitmq_client.get_channel() await self.channel.set_qos(prefetch_count=10) - self.queue = await self.channel.declare_queue("email_queue", durable=True) + + exchange = await self.channel.declare_exchange("email", aio_pika.ExchangeType.TOPIC) + self.queue = await self.channel.declare_queue("queue_welcome_email", durable=True, arguments={"x-queue-type": "quorum"}) + + await self.queue.bind(exchange, routing_key="email.welcome") async def process_message(self, message) -> None: + async with message.process(): data = json.loads(message.body) - print(f"Отправляю email на {data['email']}") + print(f"Обрабатываю: {data}, метка: {message.routing_key}") async def start_consuming(self)->None: if self.queue is None: await self.setup() - - if self.channel is None or self.queue is None: - raise RuntimeError("Failed to set up RabbitMQ channel/queue") - - async with self.queue.iterator() as queue_iter: + + queue = self.queue + if queue is None: + raise RuntimeError("Failed to set up RabbitMQ queue") + + async with queue.iterator() as queue_iter: async for message in queue_iter: await self.process_message(message) -email_consumer=EmailConsumer() \ No newline at end of file +class ResetEmailConsumer: + + def __init__(self) -> None: + self.channel = None + self.queue = None + + async def setup(self)->None: + + self.channel = await rabbitmq_client.get_channel() + await self.channel.set_qos(prefetch_count=10) + + exchange = await self.channel.declare_exchange("email", aio_pika.ExchangeType.TOPIC) + self.queue = await self.channel.declare_queue("queue_reset_email", durable=True, arguments={"x-queue-type": "quorum"}) + + await self.queue.bind(exchange, routing_key="email.reset") + + async def process_message(self, message) -> None: + + async with message.process(): + data = json.loads(message.body) + print(f"Обрабатываю: {data}, метка: {message.routing_key}") + + async def start_consuming(self)->None: + + if self.queue is None: + await self.setup() + + queue = self.queue + if queue is None: + raise RuntimeError("Failed to set up RabbitMQ queue") + + async with queue.iterator() as queue_iter: + async for message in queue_iter: + await self.process_message(message) + +welcome_email_consumer=WelcomeEmailConsumer() +reset_email_consumer=ResetEmailConsumer() \ No newline at end of file diff --git a/src/messaging/producers/producers.py b/src/messaging/producers/producers.py index 6ed6e1d..1c49fc0 100644 --- a/src/messaging/producers/producers.py +++ b/src/messaging/producers/producers.py @@ -14,20 +14,28 @@ class EmailProducer: async def setup(self)->None: self.channel = await rabbitmq_client.get_channel() - self.queue = await self.channel.declare_queue("email_queue", durable=True) + self.exchange = await self.channel.declare_exchange("email", aio_pika.ExchangeType.TOPIC) async def send_welcome_email(self, email:str)->None: - if self.queue is None: - await self.setup() + await self._publish({"email":email}, routing_key="email.welcome") - if self.channel is None or self.queue is None: - raise RuntimeError("Failed to set up RabbitMQ channel/queue") + async def send_reset_email(self, email:str)->None: + + await self._publish({"email":email}, routing_key="email.reset") + + async def _publish(self,data:dict, routing_key:str)->None: + + if self.exchange is None: + await self.setup() - message=aio_pika.Message( - body=json.dumps({"email": email}).encode(), + if self.exchange is None: + raise RuntimeError("Failed to set up RabbitMQ exchange") + + message = aio_pika.Message( + body=json.dumps(data).encode(), delivery_mode=aio_pika.DeliveryMode.PERSISTENT, ) - await self.channel.default_exchange.publish(message=message, routing_key=self.queue.name) + await self.exchange.publish(message=message, routing_key=routing_key) email_producer=EmailProducer() \ No newline at end of file diff --git a/src/models/configs_read/daemons_json.py b/src/models/configs_read/daemons_json.py new file mode 100644 index 0000000..684c087 --- /dev/null +++ b/src/models/configs_read/daemons_json.py @@ -0,0 +1,12 @@ +from pydantic_settings import SettingsConfigDict + +from src.models.configs_read.env import Base + + +class DaemonsConfig(Base): + + daemons: list[str] + + model_config = SettingsConfigDict(json_file="configs/daemons.json") + +daemons_config = DaemonsConfig() # type: ignore[call-arg] \ No newline at end of file