/
Lexus_666
/
PIR_LAB_4
Обзор
Документация
Войти
/
Lexus_666
/
PIR_LAB_4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
server/main.py
129 строк
4 KB
Lexus-666
Lab_4_done added client + server + test run
25 дек 2025, 02:27
25 дек 2025, 02:27
1e2f432
Код
Авторство
О чём код?
import pika import json import logging from common.config import config from server.message_handler import MessageHandler import time from common.utils import json_dumps, json_loads logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' ) logger = logging.getLogger(__name__) class RabbitMQServer: def __init__(self): self.connection = None self.channel = None self.message_handler = MessageHandler() self.setup_connection() def setup_connection(self): credentials = pika.PlainCredentials( config.RABBITMQ_USER, config.RABBITMQ_PASSWORD ) while True: try: self.connection = pika.BlockingConnection( pika.ConnectionParameters( host=config.RABBITMQ_HOST, port=config.RABBITMQ_PORT, credentials=credentials ) ) self.channel = self.connection.channel() self.channel.queue_declare( queue=config.REQUEST_QUEUE, durable=True, arguments={ 'x-dead-letter-exchange': '', 'x-dead-letter-routing-key': config.DEAD_LETTER_QUEUE } ) self.channel.queue_declare( queue=config.RESPONSE_QUEUE, durable=True ) self.channel.queue_declare( queue=config.DEAD_LETTER_QUEUE, durable=True ) logger.info("Connected to RabbitMQ") break except Exception as e: logger.error(f"Failed to connect to RabbitMQ: {e}") logger.info("Retrying in 5 seconds...") time.sleep(5) def callback(self, ch, method, properties, body): try: logger.info(f"Received message: {body}") message = json_loads(body) response = self.message_handler.process_message(message) self.send_response(response, properties.reply_to, properties.correlation_id) ch.basic_ack(delivery_tag=method.delivery_tag) logger.info(f"Processed request: {message.get('id')}") except json.JSONDecodeError as e: logger.error(f"Invalid JSON: {e}") ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False) except Exception as e: logger.error(f"Error processing message: {e}") ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True) def send_response(self, response: dict, reply_to: str, correlation_id: str): if not reply_to: logger.error("No reply_to queue specified") return self.channel.basic_publish( exchange='', routing_key=reply_to, body=json_dumps(response), properties=pika.BasicProperties( correlation_id=correlation_id, delivery_mode=2 ) ) logger.info(f"Sent response to {reply_to}, correlation_id: {correlation_id}") def start(self): logger.info(f"Starting server. Listening on queue: {config.REQUEST_QUEUE}") self.channel.basic_qos(prefetch_count=1) self.channel.basic_consume( queue=config.REQUEST_QUEUE, on_message_callback=self.callback ) try: self.channel.start_consuming() except KeyboardInterrupt: logger.info("Server stopped by user") self.stop() except Exception as e: logger.error(f"Server error: {e}") self.stop() def stop(self): if self.connection and not self.connection.is_closed: self.connection.close() logger.info("Server stopped") if __name__ == "__main__": server = RabbitMQServer() server.start()