/
vladnikan
/
lab4
Обзор
Документация
Войти
/
vladnikan
/
lab4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
message_processor.py
384 строки
14 KB
Vlad Nikanorov
first_commit
30 дек 2025, 00:57
30 дек 2025, 00:57
562988a
Код
Авторство
О чём код?
import json import logging import time import uuid from typing import Dict, Any, Optional, Callable from datetime import datetime, timedelta import jwt from pydantic import BaseModel, Field, validator from passlib.context import CryptContext from config import Config # Настройка логирования logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) # Хранилище в памяти (для простоты) users_db = [] projects_db = [] tasks_db = [] # Криптография pwd_context = CryptContext(schemes=["bcrypt"], deprecated="auto") class RequestMessage(BaseModel): """Модель входящего сообщения""" id: str = Field(default_factory=lambda: str(uuid.uuid4())) version: str = Field("v1", regex="^(v1|v2)$") action: str data: Dict[str, Any] = {} auth: Optional[str] = None timestamp: datetime = Field(default_factory=datetime.utcnow) class ResponseMessage(BaseModel): """Модель исходящего сообщения""" correlation_id: str status: str = Field("ok", regex="^(ok|error)$") data: Optional[Dict[str, Any]] = None error: Optional[str] = None timestamp: datetime = Field(default_factory=datetime.utcnow) class MessageProcessor: """Обработчик сообщений""" def __init__(self): self.processors = { "v1": { "create_user": self._create_user_v1, "get_users": self._get_users_v1, "get_user": self._get_user_v1, "update_user": self._update_user_v1, "delete_user": self._delete_user_v1, "create_project": self._create_project_v1, "get_projects": self._get_projects_v1, "get_project": self._get_project_v1, "update_project": self._update_project_v1, "delete_project": self._delete_project_v1, "create_task": self._create_task_v1, "get_tasks": self._get_tasks_v1, "get_task": self._get_task_v1, "update_task": self._update_task_v1, "delete_task": self._delete_task_v1, "login": self._login_v1, }, "v2": { "get_tasks": self._get_tasks_v2, } } def _hash_password(self, password: str) -> str: return pwd_context.hash(password) def _verify_password(self, plain: str, hashed: str) -> bool: return pwd_context.verify(plain, hashed) def _create_token(self, user_id: int) -> str: payload = { "user_id": user_id, "exp": datetime.utcnow() + timedelta(hours=1) } return jwt.encode(payload, Config.SECRET_KEY, algorithm="HS256") def _verify_auth(self, auth_token: Optional[str], required_user_id: Optional[int] = None) -> int: """Проверка авторизации""" if not auth_token: raise ValueError("Authentication required") if auth_token == Config.API_KEY: # API ключ для внутренних запросов return 0 try: payload = jwt.decode(auth_token, Config.SECRET_KEY, algorithms=["HS256"]) user_id = payload.get("user_id") if required_user_id and user_id != required_user_id: raise ValueError("Not authorized") return user_id except jwt.ExpiredSignatureError: raise ValueError("Token expired") except jwt.InvalidTokenError: raise ValueError("Invalid token") # V1 Handlers def _create_user_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: if any(u["username"] == data["username"] for u in users_db): raise ValueError("Username exists") new_user = { "id": len(users_db) + 1, "username": data["username"], "email": data["email"], "password_hash": self._hash_password(data["password"]) } users_db.append(new_user) return { "id": new_user["id"], "username": new_user["username"], "email": new_user["email"] } def _get_users_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) return { "users": [ {"id": u["id"], "username": u["username"], "email": u["email"]} for u in users_db ] } def _get_user_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) user_id = data["user_id"] user = next((u for u in users_db if u["id"] == user_id), None) if not user: raise ValueError("User not found") return { "id": user["id"], "username": user["username"], "email": user["email"] } def _update_user_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: user_id = data["user_id"] current_user = self._verify_auth(auth, user_id) user = next((u for u in users_db if u["id"] == user_id), None) if not user: raise ValueError("User not found") user.update({ "username": data["username"], "email": data["email"] }) return { "id": user["id"], "username": user["username"], "email": user["email"] } def _delete_user_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: user_id = data["user_id"] current_user = self._verify_auth(auth, user_id) global users_db users_db = [u for u in users_db if u["id"] != user_id] return {"success": True} def _create_project_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) if data["owner_id"] != current_user: raise ValueError("Not authorized") new_project = { "id": len(projects_db) + 1, "name": data["name"], "description": data["description"], "owner_id": data["owner_id"] } projects_db.append(new_project) return new_project def _get_projects_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) return { "projects": [ p for p in projects_db if p["owner_id"] == current_user ] } def _get_project_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) project_id = data["project_id"] project = next((p for p in projects_db if p["id"] == project_id and p["owner_id"] == current_user), None) if not project: raise ValueError("Project not found") return project def _update_project_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) project_id = data["project_id"] project = next((p for p in projects_db if p["id"] == project_id and p["owner_id"] == current_user), None) if not project: raise ValueError("Project not found") project.update({ "name": data["name"], "description": data["description"] }) return project def _delete_project_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) project_id = data["project_id"] global projects_db projects_db = [p for p in projects_db if p["id"] != project_id or p["owner_id"] != current_user] return {"success": True} def _create_task_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) # Проверка проекта project = next((p for p in projects_db if p["id"] == data["project_id"] and p["owner_id"] == current_user), None) if not project: raise ValueError("Not authorized for this project") new_task = { "id": len(tasks_db) + 1, "title": data["title"], "description": data["description"], "status": data["status"], "project_id": data["project_id"], "assignee_id": data.get("assignee_id") } tasks_db.append(new_task) return new_task def _get_tasks_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) user_tasks = [ t for t in tasks_db if t.get("assignee_id") == current_user or any(p["id"] == t["project_id"] and p["owner_id"] == current_user for p in projects_db) ] return {"tasks": user_tasks} def _get_task_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) task_id = data["task_id"] task = next( (t for t in tasks_db if t["id"] == task_id and (t.get("assignee_id") == current_user or any(p["id"] == t["project_id"] and p["owner_id"] == current_user for p in projects_db))), None ) if not task: raise ValueError("Task not found") return task def _update_task_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) task_id = data["task_id"] task = next( (t for t in tasks_db if t["id"] == task_id and (t.get("assignee_id") == current_user or any(p["id"] == t["project_id"] and p["owner_id"] == current_user for p in projects_db))), None ) if not task: raise ValueError("Task not found") task.update({ "title": data["title"], "description": data["description"], "status": data["status"], "assignee_id": data.get("assignee_id") }) return task def _delete_task_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) task_id = data["task_id"] global tasks_db tasks_db = [ t for t in tasks_db if t["id"] != task_id or not (t.get("assignee_id") == current_user or any(p["id"] == t["project_id"] and p["owner_id"] == current_user for p in projects_db)) ] return {"success": True} def _login_v1(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: user = next((u for u in users_db if u["username"] == data["username"]), None) if not user or not self._verify_password(data["password"], user["password_hash"]): raise ValueError("Invalid credentials") token = self._create_token(user["id"]) return {"token": token} # V2 Handlers def _get_tasks_v2(self, data: Dict[str, Any], auth: Optional[str]) -> Dict[str, Any]: current_user = self._verify_auth(auth) user_tasks = [ t for t in tasks_db if t.get("assignee_id") == current_user or any(p["id"] == t["project_id"] and p["owner_id"] == current_user for p in projects_db) ] # Пагинация page = data.get("page", 1) size = data.get("size", 10) start = (page - 1) * size end = start + size paginated = user_tasks[start:end] # Опциональные поля include = data.get("include") if include: fields = include.split(",") paginated = [{k: v for k, v in t.items() if k in fields} for t in paginated] return {"tasks": paginated, "page": page, "size": size, "total": len(user_tasks)} def process_message(self, request_msg: RequestMessage) -> ResponseMessage: """Обработка входящего сообщения""" try: logger.info(f"Processing message {request_msg.id}, action: {request_msg.action}") # Проверка версии и действия if request_msg.version not in self.processors: raise ValueError(f"Unsupported version: {request_msg.version}") if request_msg.action not in self.processors[request_msg.version]: raise ValueError(f"Unknown action: {request_msg.action}") # Выполнение обработчика handler = self.processors[request_msg.version][request_msg.action] result = handler(request_msg.data, request_msg.auth) return ResponseMessage( correlation_id=request_msg.id, status="ok", data=result, error=None ) except Exception as e: logger.error(f"Error processing message {request_msg.id}: {str(e)}") return ResponseMessage( correlation_id=request_msg.id, status="error", data=None, error=str(e) )