/
Evdik
/
Integration_System_labs4
Обзор
Документация
Войти
/
Evdik
/
Integration_System_labs4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
server.py
616 строк
23 KB
Evdikkom
first
22 ноя 2025, 12:27
22 ноя 2025, 12:27
8ded1dc
Код
Авторство
О чём код?
import pika import json import uuid import logging from datetime import datetime, timedelta from sqlalchemy.orm import sessionmaker import jwt import redis from models import init_database, User, Game, Order, OrderItem, GameKey, Review, ProcessedMessage, hash_password, verify_password from config import Config from sqlalchemy import text, func logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class MessageProcessor: def __init__(self): self.engine = init_database() self.Session = sessionmaker(bind=self.engine) self.redis_client = redis.Redis( host=Config.REDIS_HOST, port=Config.REDIS_PORT, db=Config.REDIS_DB, decode_responses=True ) def authenticate(self, api_key): """Аутентификация по API-ключу""" return Config.API_KEYS.get(api_key) def is_duplicate(self, message_id): """Проверка идемпотентности через Redis""" key = f"processed:{message_id}" if self.redis_client.exists(key): return True # Также проверяем в БД для надежности session = self.Session() try: existing = session.query(ProcessedMessage).filter_by(message_id=message_id).first() return existing is not None finally: session.close() def mark_processed(self, message_id, response_data): """Отметка сообщения как обработанного""" key = f"processed:{message_id}" self.redis_client.setex(key, 3600, "processed") # TTL 1 час session = self.Session() try: processed_msg = ProcessedMessage( message_id=message_id, response_data=json.dumps(response_data) ) session.add(processed_msg) session.commit() except Exception as e: session.rollback() logger.error(f"Error saving processed message: {e}") finally: session.close() def get_cached_response(self, message_id): """Получение кэшированного ответа""" session = self.Session() try: processed = session.query(ProcessedMessage).filter_by(message_id=message_id).first() if processed: return json.loads(processed.response_data) finally: session.close() return None def process_message(self, message): """Основной обработчик сообщений""" message_id = message.get('id') action = message.get('action') data = message.get('data', {}) auth = message.get('auth') # Аутентификация auth_info = self.authenticate(auth) if not auth_info: return { 'correlation_id': message_id, 'status': 'error', 'data': None, 'error': 'Authentication failed' } # Проверка идемпотентности if self.is_duplicate(message_id): logger.info(f"Duplicate message {message_id}, returning cached response") cached_response = self.get_cached_response(message_id) if cached_response: return cached_response # Обработка действий try: if action == 'create_user': result = self.handle_create_user(data) elif action == 'login': result = self.handle_login(data) elif action == 'get_profile': result = self.handle_get_profile(data, auth_info) elif action == 'deposit_balance': result = self.handle_deposit_balance(data, auth_info) elif action == 'create_order': result = self.handle_create_order(data, auth_info) elif action == 'get_games': result = self.handle_get_games(data) elif action == 'get_game': result = self.handle_get_game(data) elif action == 'add_game_keys': result = self.handle_add_game_keys(data, auth_info) elif action == 'create_review': result = self.handle_create_review(data, auth_info) elif action == 'get_orders': result = self.handle_get_orders(data, auth_info) elif action == 'health_check': result = self.handle_health_check() elif action == 'sales_analytics': result = self.handle_sales_analytics(data) else: result = {'error': f'Unknown action: {action}'} response = { 'correlation_id': message_id, 'status': 'error' if 'error' in result else 'ok', 'data': result if 'error' not in result else None, 'error': result.get('error') } # Сохраняем результат для идемпотентности if response['status'] == 'ok': self.mark_processed(message_id, response) return response except Exception as e: logger.error(f"Error processing message: {e}") return { 'correlation_id': message_id, 'status': 'error', 'data': None, 'error': str(e) } def handle_create_user(self, data): session = self.Session() try: # Валидация обязательных полей if not data.get('username') or not data.get('email') or not data.get('password'): return {'error': 'Username, email and password are required'} # Проверка уникальности if session.query(User).filter_by(username=data['username']).first(): return {'error': 'Username already exists'} if session.query(User).filter_by(email=data['email']).first(): return {'error': 'Email already exists'} user = User( username=data['username'], email=data['email'], password_hash=hash_password(data['password']) ) session.add(user) session.commit() return {'message': 'User created successfully', 'user_id': user.id} except Exception as e: session.rollback() logger.error(f"Error creating user: {e}") return {'error': f'Failed to create user: {str(e)}'} finally: session.close() def handle_login(self, data): session = self.Session() try: user = session.query(User).filter_by(username=data['username']).first() if user and verify_password(data['password'], user.password_hash): user.last_login = datetime.utcnow() session.commit() token = jwt.encode({ 'user_id': user.id, 'exp': datetime.utcnow() + timedelta(hours=24) }, Config.JWT_SECRET_KEY, algorithm='HS256') return { 'access_token': token, 'user': { 'id': user.id, 'username': user.username, 'email': user.email, 'balance': user.balance, 'role': user.role } } return {'error': 'Invalid credentials'} finally: session.close() def handle_get_profile(self, data, auth_info): session = self.Session() try: user_id = data.get('user_id', auth_info['user_id']) # Используем session.get вместо query.get user = session.get(User, user_id) if not user: return {'error': 'User not found'} return { 'id': user.id, 'username': user.username, 'email': user.email, 'balance': user.balance, 'role': user.role } finally: session.close() def handle_deposit_balance(self, data, auth_info): logger.info(f"Deposit request: user_id={data.get('user_id', auth_info['user_id'])}, amount={data.get('amount')}") session = self.Session() try: user_id = data.get('user_id', auth_info['user_id']) # Используем session.get вместо query.get user = session.get(User, user_id) if not user: return {'error': 'User not found'} amount = float(data['amount']) user.balance += amount # Явно коммитим изменения session.commit() # Обновляем объект из базы для получения актуальных данных session.refresh(user) return {'new_balance': user.balance, 'message': 'Balance updated successfully'} except Exception as e: session.rollback() logger.error(f"Error in deposit_balance: {e}") return {'error': f'Failed to deposit balance: {str(e)}'} finally: session.close() def handle_create_order(self, data, auth_info): logger.info(f"Create order request: user_id={data.get('user_id', auth_info['user_id'])}, items={data.get('items')}") session = self.Session() try: user_id = data.get('user_id', auth_info['user_id']) # Используем session.get вместо query.get user = session.get(User, user_id) if not user: return {'error': 'User not found'} total_amount = 0 items = data.get('items', []) if not items: return {'error': 'No items in order'} # Проверяем доступность игр и рассчитываем сумму for item in items: # Используем session.get вместо query.get game = session.get(Game, item['game_id']) if not game: return {'error': f'Game {item["game_id"]} not found'} if not game.is_available: return {'error': f'Game {game.title} is not available'} quantity = item.get('quantity', 1) if quantity <= 0: return {'error': f'Invalid quantity for game {game.title}'} total_amount += game.price * quantity # Обновляем пользователя из базы для получения актуального баланса session.refresh(user) # Проверяем баланс if user.balance < total_amount: return {'error': f'Insufficient balance. Available: ${user.balance}, Required: ${total_amount}'} # Создаем заказ order = Order(user_id=user.id, total_amount=total_amount, status='completed') session.add(order) session.flush() # Получаем ID заказа # Создаем элементы заказа for item in items: game = session.get(Game, item['game_id']) order_item = OrderItem( order_id=order.id, game_id=game.id, quantity=item.get('quantity', 1), price_at_purchase=game.price ) session.add(order_item) # Списание средств user.balance -= total_amount # Явно коммитим все изменения session.commit() # Обновляем объекты для получения актуальных данных session.refresh(user) session.refresh(order) return { 'order_id': order.id, 'total_amount': total_amount, 'status': order.status } except Exception as e: session.rollback() logger.error(f"Error creating order: {e}") return {'error': f'Failed to create order: {str(e)}'} finally: session.close() def handle_get_games(self, data): session = self.Session() try: query = session.query(Game) if data.get('genre'): query = query.filter(Game.genre == data['genre']) if data.get('platform'): query = query.filter(Game.platform == data['platform']) if data.get('available') is not None: query = query.filter(Game.is_available == data['available']) page = data.get('page', 1) per_page = min(data.get('per_page', 10), 100) offset = (page - 1) * per_page games = query.offset(offset).limit(per_page).all() total = query.count() games_data = [] for game in games: games_data.append({ 'id': game.id, 'title': game.title, 'price': game.price, 'developer': game.developer, 'genre': game.genre, 'platform': game.platform, 'is_available': game.is_available }) return { 'games': games_data, 'total': total, 'page': page, 'per_page': per_page, 'pages': (total + per_page - 1) // per_page } finally: session.close() def handle_get_game(self, data): session = self.Session() try: # Используем session.get вместо query.get game = session.get(Game, data['game_id']) if not game: return {'error': 'Game not found'} return { 'id': game.id, 'title': game.title, 'description': game.description, 'price': game.price, 'developer': game.developer, 'publisher': game.publisher, 'genre': game.genre, 'platform': game.platform, 'release_date': game.release_date.isoformat() if game.release_date else None, 'is_available': game.is_available } finally: session.close() def handle_add_game_keys(self, data, auth_info): if auth_info.get('role') != 'admin': return {'error': 'Admin access required'} session = self.Session() try: # Используем session.get вместо query.get game = session.get(Game, data['game_id']) if not game: return {'error': 'Game not found'} added_keys = [] for key_value in data['keys']: if not session.query(GameKey).filter_by(key=key_value).first(): game_key = GameKey(game_id=game.id, key=key_value) session.add(game_key) added_keys.append(key_value) session.commit() return {'added_count': len(added_keys), 'added_keys': added_keys} finally: session.close() def handle_create_review(self, data, auth_info): session = self.Session() try: # Проверяем, что пользователь покупал игру has_purchased = session.query(OrderItem).join(Order).filter( Order.user_id == auth_info['user_id'], OrderItem.game_id == data['game_id'] ).first() if not has_purchased: return {'error': 'You must purchase the game before reviewing'} existing_review = session.query(Review).filter_by( user_id=auth_info['user_id'], game_id=data['game_id'] ).first() if existing_review: return {'error': 'You have already reviewed this game'} review = Review( user_id=auth_info['user_id'], game_id=data['game_id'], rating=data['rating'], comment=data.get('comment', ''), is_verified=True ) session.add(review) session.commit() return {'review_id': review.id, 'message': 'Review added successfully'} finally: session.close() def handle_get_orders(self, data, auth_info): session = self.Session() try: query = session.query(Order).filter_by(user_id=auth_info['user_id']) page = data.get('page', 1) per_page = min(data.get('per_page', 10), 100) offset = (page - 1) * per_page orders = query.offset(offset).limit(per_page).all() total = query.count() orders_data = [] for order in orders: orders_data.append({ 'id': order.id, 'total_amount': order.total_amount, 'status': order.status, 'created_at': order.created_at.isoformat(), 'items_count': len(order.order_items) }) return { 'orders': orders_data, 'total': total, 'page': page, 'per_page': per_page } finally: session.close() def handle_health_check(self): session = self.Session() try: # Используем text() для SQL выражений session.execute(text('SELECT 1')) db_status = 'healthy' # Проверяем Redis try: self.redis_client.ping() redis_status = 'healthy' except: redis_status = 'unhealthy' return { 'status': 'healthy' if db_status == 'healthy' and redis_status == 'healthy' else 'degraded', 'database': db_status, 'redis': redis_status, 'timestamp': datetime.utcnow().isoformat() } except Exception as e: return {'status': 'unhealthy', 'error': str(e)} finally: session.close() def handle_sales_analytics(self, data): session = self.Session() try: period = data.get('period', 'week') if period == 'day': start_date = datetime.utcnow() - timedelta(days=1) elif period == 'week': start_date = datetime.utcnow() - timedelta(weeks=1) elif period == 'month': start_date = datetime.utcnow() - timedelta(days=30) else: start_date = datetime.utcnow() - timedelta(days=365) total_orders = session.query(Order).filter(Order.created_at >= start_date).count() # Используем правильный способ получения суммы total_revenue_result = session.query(func.sum(Order.total_amount)).filter( Order.created_at >= start_date ).first() total_revenue = total_revenue_result[0] if total_revenue_result[0] is not None else 0 return { 'period': period, 'total_orders': total_orders, 'total_revenue': float(total_revenue), 'average_order_value': float(total_revenue / total_orders) if total_orders > 0 else 0 } except Exception as e: logger.error(f"Error in sales analytics: {e}") return {'error': str(e)} finally: session.close() class RabbitMQServer: def __init__(self): self.processor = MessageProcessor() self.connection = None self.channel = None def connect(self): credentials = pika.PlainCredentials(Config.RABBITMQ_USER, Config.RABBITMQ_PASS) parameters = pika.ConnectionParameters( host=Config.RABBITMQ_HOST, port=Config.RABBITMQ_PORT, credentials=credentials, heartbeat=600 ) self.connection = pika.BlockingConnection(parameters) self.channel = self.connection.channel() # Объявляем очереди self.channel.queue_declare(queue=Config.REQUEST_QUEUE, durable=True) self.channel.queue_declare(queue=Config.RESPONSE_QUEUE, durable=True) self.channel.queue_declare(queue=Config.DEAD_LETTER_QUEUE, durable=True) # Настраиваем DLQ self.channel.queue_bind( exchange='amq.direct', queue=Config.DEAD_LETTER_QUEUE ) def on_request(self, ch, method, props, body): try: message = json.loads(body) logger.info(f"Received message: {message['id']} - {message['action']}") response = self.processor.process_message(message) ch.basic_publish( exchange='', routing_key=props.reply_to, properties=pika.BasicProperties( correlation_id=props.correlation_id, ), body=json.dumps(response) ) ch.basic_ack(delivery_tag=method.delivery_tag) except json.JSONDecodeError as e: logger.error(f"JSON decode error: {e}") # Отправляем в DLQ ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False) except Exception as e: logger.error(f"Error processing request: {e}") # Повторная обработка ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True) def start(self): self.connect() self.channel.basic_qos(prefetch_count=1) self.channel.basic_consume( queue=Config.REQUEST_QUEUE, on_message_callback=self.on_request ) logger.info("Server started. Waiting for messages...") self.channel.start_consuming() def stop(self): if self.connection: self.connection.close() if __name__ == '__main__': server = RabbitMQServer() try: server.start() except KeyboardInterrupt: server.stop() logger.info("Server stopped")