/
Lexus_666
/
PIR_LAB_4
Обзор
Документация
Войти
/
Lexus_666
/
PIR_LAB_4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
client/main.py
161 строка
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 uuid import logging from typing import Optional, Dict, Any from common.config import config from common.models import RequestMessage import time from common.utils import json_dumps, json_loads logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class RabbitMQClient: def __init__(self): self.connection = None self.channel = None self.callback_queue = None self.response = None self.correlation_id = None self.setup_connection() def setup_connection(self): credentials = pika.PlainCredentials( config.RABBITMQ_USER, config.RABBITMQ_PASSWORD ) self.connection = pika.BlockingConnection( pika.ConnectionParameters( host=config.RABBITMQ_HOST, port=config.RABBITMQ_PORT, credentials=credentials ) ) self.channel = self.connection.channel() result = self.channel.queue_declare(queue='', exclusive=True) self.callback_queue = result.method.queue self.channel.basic_consume( queue=self.callback_queue, on_message_callback=self.on_response, auto_ack=True ) def on_response(self, ch, method, properties, body): if self.correlation_id == properties.correlation_id: self.response = json_loads(body) logger.info(f"Response received: {self.response}") def send_request(self, action: str, data: Dict[str, Any], auth: str = "test-api-key-123") -> Optional[ Dict[str, Any]]: request = RequestMessage(action, data, auth) message = request.to_dict() self.response = None self.correlation_id = message['id'] self.channel.basic_publish( exchange='', routing_key=config.REQUEST_QUEUE, body=json_dumps(message), properties=pika.BasicProperties( reply_to=self.callback_queue, correlation_id=self.correlation_id, delivery_mode=2 ) ) logger.info(f"Sent request: {action} with id: {self.correlation_id}") timeout = 30 start_time = time.time() while self.response is None: self.connection.process_data_events() if time.time() - start_time > timeout: logger.error(f"Timeout waiting for response to request: {self.correlation_id}") return None time.sleep(0.1) logger.info(f"Received response for request: {self.correlation_id}") return self.response def close(self): if self.connection and not self.connection.is_closed: self.connection.close() def test_create_user(): client = RabbitMQClient() try: response = client.send_request( action="create_user", data={ "username": "testuser", "email": "test@example.com", "full_name": "Test User" } ) if response: print(f"Response: {json.dumps(response, indent=2)}") else: print("No response received") finally: client.close() def test_get_all_users(): client = RabbitMQClient() try: response = client.send_request( action="get_all_users", data={} ) if response: print(f"Response: {json.dumps(response, indent=2)}") else: print("No response received") finally: client.close() def test_health_check(): client = RabbitMQClient() try: response = client.send_request( action="health_check", data={} ) if response: print(f"Response: {json.dumps(response, indent=2)}") else: print("No response received") finally: client.close() if __name__ == "__main__": print("Testing RabbitMQ Client") print("=" * 50) print("\n1. Testing health check:") test_health_check() print("\n2. Testing create user:") test_create_user() print("\n3. Testing get all users:") test_get_all_users()