/
Frozin
/
DIS-4
Обзор
Документация
Войти
/
Frozin
/
DIS-4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
server/main.py
145 строк
5 KB
Frozin-p
Fixed database operations. Added files for testing
16 дек 2025, 16:45
16 дек 2025, 16:45
109ec6f
Код
Авторство
О чём код?
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 ) # Объявление Dead Letter Queue 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()