/
Pepegator
/
task-api
Обзор
Документация
Войти
/
Pepegator
/
task-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
lab4
pythonProject2/rabbitmq_service.py
372 строки
12 KB
Pol136
add rabbitmq
30 дек 2025, 17:44
30 дек 2025, 17:44
7c8d717
Код
Авторство
О чём код?
import asyncio import json import time import logging from datetime import datetime, timedelta from typing import Optional, Dict, Any import aio_pika from aio_pika import Channel, Connection, Queue, IncomingMessage from sqlalchemy.orm import Session from database import SessionLocal, ProcessedRequest, MessageError, get_db logger = logging.getLogger(__name__) class RabbitMQConfig: RABBITMQ_URL = "amqp://guest:guest@localhost:5672/" REQUEST_QUEUE = "api.requests" RESPONSE_QUEUE = "api.responses" DLQ_QUEUE = "api.requests.dlq" MAX_RETRIES = 3 INITIAL_RETRY_DELAY = 1000 RETRY_BACKOFF = 2 class RabbitMQPublisher: def __init__(self, channel: Channel): self.channel = channel async def initialize(self): await self.channel.declare_queue( name=RabbitMQConfig.RESPONSE_QUEUE, durable=True, auto_delete=False ) await self.channel.declare_queue( name=RabbitMQConfig.DLQ_QUEUE, durable=True, auto_delete=False ) async def publish_response(self, response_data: Dict) -> None: try: queue = await self.channel.get_queue(RabbitMQConfig.RESPONSE_QUEUE) message = aio_pika.Message( body=json.dumps(response_data).encode(), delivery_mode=aio_pika.DeliveryMode.PERSISTENT, content_type='application/json' ) await queue.put(message) logger.info(f"Response published: {response_data.get('correlation_id')}") except Exception as e: logger.error(f"Failed to publish response: {e}") async def publish_to_dlq(self, message_data: Dict, error: str) -> None: try: dlq = { "request_id": message_data.get('id', 'unknown'), "original_message": message_data, "error_message": error, "timestamp": datetime.utcnow().isoformat() } queue = await self.channel.get_queue(RabbitMQConfig.DLQ_QUEUE) message = aio_pika.Message( body=json.dumps(dlq).encode(), delivery_mode=aio_pika.DeliveryMode.PERSISTENT, content_type='application/json' ) await queue.put(message) logger.warning(f"Message sent to DLQ: {dlq['request_id']}") except Exception as e: logger.error(f"Failed to publish to DLQ: {e}") class IdempotencyService: @staticmethod def save_processed_request(request_id: str, response_data: Dict, status_code: int, db: Session, ttl_hours: int = 24) -> None: try: existing = db.query(ProcessedRequest).filter( ProcessedRequest.id == request_id ).first() if existing: db.delete(existing) expires_at = datetime.utcnow() + timedelta(hours=ttl_hours) processed = ProcessedRequest( id=request_id, response_data=response_data, status_code=status_code, created_at=datetime.utcnow(), expires_at=expires_at ) db.add(processed) db.commit() logger.info(f"Processed request saved: {request_id}") except Exception as e: db.rollback() logger.error(f"Failed to save processed request: {e}") @staticmethod def get_cached_response(request_id: str, db: Session) -> Optional[Dict]: try: processed = db.query(ProcessedRequest).filter( ProcessedRequest.id == request_id, ProcessedRequest.expires_at > datetime.utcnow() ).first() if processed: logger.info(f"Cache hit for request: {request_id}") return { "data": processed.response_data, "status_code": processed.status_code } return None except Exception as e: logger.error(f"Failed to get cached response: {e}") return None class AuthService: VALID_API_KEYS = { "sk_test_abc123xyz": {"name": "Test Client", "created_at": "2024-01-01"}, "sk_live_abc123xyz": {"name": "Production Client", "created_at": "2024-01-01"} } @staticmethod def validate_api_key(api_key: str) -> bool: is_valid = api_key in AuthService.VALID_API_KEYS if is_valid: logger.info(f"API key validated: {api_key[:10]}...") else: logger.warning(f"Invalid API key: {api_key[:10]}...") return is_valid class MessageHandler: def __init__(self, db: Session): self.db = db async def handle_request(self, request: Dict) -> Dict: start_time = time.time() request_id = request.get('id') try: cached = IdempotencyService.get_cached_response(request_id, self.db) if cached: return { "correlation_id": request_id, "status": "success", "code": cached["status_code"], "data": cached["data"], "error": None, "processing_time_ms": 0 } api_key = request.get('auth', {}).get('api_key') if not AuthService.validate_api_key(api_key): return { "correlation_id": request_id, "status": "error", "code": 401, "data": None, "error": { "type": "AUTH_ERROR", "message": "Invalid API key" }, "processing_time_ms": int((time.time() - start_time) * 1000) } action = request.get('action') response_data, status_code = await self._route_action(action, request.get('data', {})) IdempotencyService.save_processed_request( request_id, response_data, status_code, self.db ) processing_time = int((time.time() - start_time) * 1000) return { "correlation_id": request_id, "status": "success", "code": status_code, "data": response_data, "error": None, "processing_time_ms": processing_time } except Exception as e: logger.error(f"Error handling request {request_id}: {e}") processing_time = int((time.time() - start_time) * 1000) return { "correlation_id": request_id, "status": "error", "code": 500, "data": None, "error": { "type": "INTERNAL_ERROR", "message": str(e) }, "processing_time_ms": processing_time } async def _route_action(self, action: str, data: Dict) -> tuple: logger.info(f"Routing action: {action}") if action == "create_task": return { "id": 1, "title": data.get("title"), "description": data.get("description"), "priority": data.get("priority", "medium"), "status": "todo", "created_at": datetime.utcnow().isoformat() }, 201 elif action == "list_tasks": return { "items": [], "total": 0, "skip": data.get("skip", 0), "limit": data.get("limit", 10) }, 200 elif action == "get_task": return { "id": data.get("task_id"), "title": "Task" }, 200 elif action == "create_user": return { "id": 1, "username": data.get("username"), "email": data.get("email"), "created_at": datetime.utcnow().isoformat() }, 201 else: raise ValueError(f"Unknown action: {action}") class RabbitMQConsumer: def __init__(self): self.connection: Optional[Connection] = None self.channel: Optional[Channel] = None self.publisher: Optional[RabbitMQPublisher] = None self.running = False async def connect(self) -> None: try: self.connection = await aio_pika.connect_robust( RabbitMQConfig.RABBITMQ_URL ) self.channel = await self.connection.channel() self.publisher = RabbitMQPublisher(self.channel) await self.publisher.initialize() logger.info("Connected to RabbitMQ") except Exception as e: logger.error(f"Failed to connect to RabbitMQ: {e}") raise async def declare_queues(self) -> None: try: await self.channel.declare_queue( name=RabbitMQConfig.REQUEST_QUEUE, durable=True, auto_delete=False ) await self.channel.declare_queue( name=RabbitMQConfig.DLQ_QUEUE, durable=True, auto_delete=False ) logger.info("Queues declared") except Exception as e: logger.error(f"Failed to declare queues: {e}") raise async def start_consuming(self) -> None: try: queue = await self.channel.get_queue(RabbitMQConfig.REQUEST_QUEUE) self.running = True logger.info(f"Started consuming from {RabbitMQConfig.REQUEST_QUEUE}") async with queue.iterator() as queue_iter: async for message in queue_iter: if not self.running: break await self._process_message(message) except asyncio.CancelledError: logger.info("Consumer cancelled") except Exception as e: logger.error(f"Error in consumer: {e}") async def _process_message(self, message: IncomingMessage) -> None: try: request_data = json.loads(message.body.decode()) request_id = request_data.get('id') logger.info(f"Processing message: {request_id}") db = SessionLocal() try: handler = MessageHandler(db) response = await handler.handle_request(request_data) finally: db.close() await self.publisher.publish_response(response) await message.ack() except json.JSONDecodeError as e: logger.error(f"Invalid JSON in message: {e}") await message.nack(requeue=False) await self.publisher.publish_to_dlq( {"body": message.body.decode()}, f"JSON_DECODE_ERROR: {str(e)}" ) except Exception as e: logger.error(f"Error processing message: {e}") await message.nack(requeue=False) try: request_data = json.loads(message.body.decode()) await self.publisher.publish_to_dlq( request_data, f"{type(e).__name__}: {str(e)}" ) except: pass async def stop(self) -> None: self.running = False if self.connection: await self.connection.close() logger.info("Disconnected from RabbitMQ") _consumer: Optional[RabbitMQConsumer] = None async def start_consumer() -> None: global _consumer _consumer = RabbitMQConsumer() try: await _consumer.connect() await _consumer.declare_queues() await _consumer.start_consuming() except Exception as e: logger.error(f"Consumer failed: {e}") raise async def stop_consumer() -> None: global _consumer if _consumer: await _consumer.stop()