/
Pepegator
/
task-api
Обзор
Документация
Войти
/
Pepegator
/
task-api
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
lab4
pythonProject2/test_rabbit.py
301 строка
10 KB
Pol136
add rabbitmq
30 дек 2025, 17:44
30 дек 2025, 17:44
7c8d717
Код
Авторство
О чём код?
import asyncio import aio_pika import json import uuid from datetime import datetime RABBITMQ_URL = "amqp://guest:guest@localhost/" TASK_QUEUE = "task_queue" RESPONSE_QUEUE = "response_queue" VALID_API_KEY = "sk_test_abc123xyz" INVALID_API_KEY = "sk_invalid_key" def create_request_message(action, data, api_key=VALID_API_KEY): return { "id": str(uuid.uuid4()), "version": "v1", "action": action, "timestamp": datetime.utcnow().isoformat(), "data": data, "auth": { "api_key": api_key, "signature": None }, "reply_to": RESPONSE_QUEUE, "timeout_ms": 5000 } async def connect_to_rabbitmq(): try: connection = await aio_pika.connect_robust(RABBITMQ_URL) channel = await connection.channel() return connection, channel except Exception as e: print(f"[ERROR] Не удалось подключиться к RabbitMQ: {e}") print("[ERROR] Убедитесь что RabbitMQ запущен: docker-compose up -d") raise async def send_message(channel, message, queue_name=TASK_QUEUE): queue = await channel.get_queue(queue_name) body = json.dumps(message).encode() amessage = aio_pika.Message( body=body, content_type="application/json", correlation_id=message["id"] ) await channel.default_exchange.publish(amessage, routing_key=queue_name) async def test_create_task(): print("\n" + "="*70) print("TEST 1: Создание задачи (create_task)") print("="*70) connection, channel = await connect_to_rabbitmq() try: task_queue = await channel.declare_queue(TASK_QUEUE, durable=True) response_queue = await channel.declare_queue(RESPONSE_QUEUE, durable=True) request = create_request_message( action="create_task", data={ "title": "Купить молоко", "description": "Купить молоко в магазине", "priority": "high" } ) await send_message(channel, request, TASK_QUEUE) print(f"[OK] Сообщение отправлено в очередь '{TASK_QUEUE}'") print(f"[OK] Request ID: {request['id']}") print(f"[OK] Action: {request['action']}") print(f"[OK] API Key: {request['auth']['api_key']}") message_count = task_queue.declaration_result.method.message_count print(f"[OK] Сообщений в очереди: {message_count}") await task_queue.purge() print("[PASS] Тест пройден успешно") except Exception as e: print(f"[FAIL] Ошибка: {e}") finally: await connection.close() async def test_idempotency(): print("\n" + "="*70) print("TEST 2: Идемпотентность (дубликаты)") print("="*70) connection, channel = await connect_to_rabbitmq() try: task_queue = await channel.declare_queue(TASK_QUEUE, durable=True) request_id = str(uuid.uuid4()) for attempt in range(1, 3): request = create_request_message( action="update_task", data={ "task_id": 1, "status": "completed" } ) request["id"] = request_id await send_message(channel, request, TASK_QUEUE) print(f"[OK] Попытка {attempt}: Отправлено сообщение с ID {request_id}") message_count = task_queue.declaration_result.method.message_count print(f"[OK] Всего сообщений в очереди: {message_count}") print(f"[OK] Request ID для отслеживания: {request_id}") await task_queue.purge() print("[PASS] Тест пройден успешно") except Exception as e: print(f"[FAIL] Ошибка: {e}") finally: await connection.close() async def test_invalid_api_key(): print("\n" + "="*70) print("TEST 3: Аутентификация (невалидный API ключ)") print("="*70) connection, channel = await connect_to_rabbitmq() try: task_queue = await channel.declare_queue(TASK_QUEUE, durable=True) request = create_request_message( action="create_task", data={ "title": "Тестовая задача", "priority": "medium" }, api_key=INVALID_API_KEY ) await send_message(channel, request, TASK_QUEUE) print(f"[OK] Сообщение с невалидным ключом отправлено") print(f"[OK] API Key: {INVALID_API_KEY}") print(f"[INFO] При обработке это сообщение будет отклонено") print(f"[INFO] Оно попадёт в Dead Letter Queue для retry") message_count = task_queue.declaration_result.method.message_count print(f"[OK] Сообщений в очереди: {message_count}") await task_queue.purge() print("[PASS] Тест пройден успешно") except Exception as e: print(f"[FAIL] Ошибка: {e}") finally: await connection.close() async def test_list_tasks(): print("\n" + "="*70) print("TEST 4: Получение списка задач (list_tasks)") print("="*70) connection, channel = await connect_to_rabbitmq() try: task_queue = await channel.declare_queue(TASK_QUEUE, durable=True) request = create_request_message( action="list_tasks", data={ "skip": 0, "limit": 10, "status": "active" } ) await send_message(channel, request, TASK_QUEUE) print(f"[OK] Запрос на получение списка отправлен") print(f"[OK] Action: {request['action']}") print(f"[OK] Skip: {request['data']['skip']}") print(f"[OK] Limit: {request['data']['limit']}") print(f"[OK] Status filter: {request['data']['status']}") message_count = task_queue.declaration_result.method.message_count print(f"[OK] Сообщений в очереди: {message_count}") await task_queue.purge() print("[PASS] Тест пройден успешно") except Exception as e: print(f"[FAIL] Ошибка: {e}") finally: await connection.close() async def test_message_format(): print("\n" + "="*70) print("TEST 5: Формат сообщения (валидация структуры)") print("="*70) try: request = create_request_message( action="create_task", data={ "title": "Новая задача", "priority": "low" } ) required_fields = ["id", "version", "action", "timestamp", "data", "auth"] for field in required_fields: if field not in request: raise ValueError(f"Отсутствует обязательное поле: {field}") print(f"[OK] Поле '{field}' присутствует") auth_fields = ["api_key"] for field in auth_fields: if field not in request["auth"]: raise ValueError(f"Отсутствует обязательное поле в auth: {field}") print(f"[OK] Поле 'auth.{field}' присутствует") json_str = json.dumps(request) print(f"[OK] Сообщение успешно сериализовано в JSON") print(f"[OK] Размер: {len(json_str)} байт") decoded = json.loads(json_str) print(f"[OK] Сообщение успешно десериализовано из JSON") assert isinstance(request["id"], str), "id должен быть строкой" assert isinstance(request["version"], str), "version должна быть строкой" assert isinstance(request["action"], str), "action должна быть строкой" assert isinstance(request["data"], dict), "data должна быть dict" assert isinstance(request["auth"], dict), "auth должна быть dict" print(f"[OK] Все типы данных корректны") print("[PASS] Тест пройден успешно") except Exception as e: print(f"[FAIL] Ошибка: {e}") async def main(): print("\n") print("="*70) print("RabbitMQ INTEGRATION TESTS") print("="*70) print("Запуск набора тестов для проверки RabbitMQ интеграции") print("="*70) tests = [ ("Создание задачи", test_create_task), ("Идемпотентность", test_idempotency), ("Аутентификация", test_invalid_api_key), ("Получение списка", test_list_tasks), ("Формат сообщения", test_message_format), ] passed = 0 failed = 0 for name, test_func in tests: try: await test_func() passed += 1 except Exception as e: print(f"\n[FAIL] Тест '{name}' провалился: {e}") failed += 1 print("\n" + "="*70) print("РЕЗУЛЬТАТЫ") print("="*70) print(f"Всего тестов: {len(tests)}") print(f"Пройдено: {passed}") print(f"Провалено: {failed}") if failed == 0: print("\nВсе тесты пройдены успешно!") else: print(f"\n{failed} тест(ов) провалились") print("="*70 + "\n") if __name__ == "__main__": asyncio.run(main())