Files
anti-plagiarism/services/api/app/core/websocket_manager.py
jze9 7758315632 feat: initial microservices project structure
Services:
- api: FastAPI gateway with JWT auth, async endpoints, WebSocket
- worker-gpu: CUDA sentence-transformers, FAISS IVFFlat, Ollama LLM
- worker-indexer: Winnowing+MinHash plagiarism detection, PDF/DOCX extraction
- worker-notifier: SMTP email notifications
- worker-gost: GOST 7.1-2003 and GOST R 7.0.5-2008 formatting

Infrastructure:
- docker-compose.yml (production) + docker-compose.dev.yml (hot reload)
- Nginx reverse proxy + WebSocket support
- PostgreSQL 16 with Alembic migrations
- Elasticsearch 8 with Russian/English analyzers
- MinIO, RabbitMQ, Redis, Ollama

Frontend:
- React 18 + Vite + TypeScript + TailwindCSS + Zustand + React Query v5
- 9 pages: Home, Search, Cabinet, Task, Check, Bibliography, Pricing, Login, Register

Scripts:
- Parser stubs: OpenAlex, КиберЛенинка, arXiv (Phase 0 - to be filled)

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-24 19:42:39 +05:00

87 lines
3.3 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""WebSocket менеджер для real-time обновлений статуса задач."""
import json
import logging
from typing import Any
from fastapi import WebSocket
logger = logging.getLogger(__name__)
class ConnectionManager:
"""Менеджер WebSocket соединений с поддержкой подписки на задачи."""
def __init__(self) -> None:
# task_id -> список активных соединений
self._connections: dict[str, list[WebSocket]] = {}
async def connect(self, task_id: str, websocket: WebSocket) -> None:
"""Принять WebSocket соединение и зарегистрировать его для задачи."""
await websocket.accept()
if task_id not in self._connections:
self._connections[task_id] = []
self._connections[task_id].append(websocket)
logger.info(f"WebSocket подключён к задаче {task_id!r}")
def disconnect(self, task_id: str, websocket: WebSocket) -> None:
"""Удалить соединение из реестра."""
if task_id in self._connections:
try:
self._connections[task_id].remove(websocket)
except ValueError:
pass
if not self._connections[task_id]:
del self._connections[task_id]
logger.info(f"WebSocket отключён от задачи {task_id!r}")
async def send_task_update(self, task_id: str, data: dict[str, Any]) -> None:
"""
Отправить обновление статуса задачи всем подключённым клиентам.
Args:
task_id: ID задачи
data: Словарь с обновлением (status, queue_position, result и т.д.)
"""
connections = self._connections.get(task_id, [])
if not connections:
return
message = json.dumps(data, ensure_ascii=False, default=str)
dead_connections = []
for websocket in connections:
try:
await websocket.send_text(message)
except Exception as e:
logger.warning(f"Ошибка отправки WebSocket сообщения: {e}")
dead_connections.append(websocket)
# Очистить мёртвые соединения
for ws in dead_connections:
self.disconnect(task_id, ws)
async def broadcast(self, data: dict[str, Any]) -> None:
"""Отправить сообщение всем подключённым клиентам."""
message = json.dumps(data, ensure_ascii=False, default=str)
all_dead = []
for task_id, connections in self._connections.items():
for websocket in connections:
try:
await websocket.send_text(message)
except Exception:
all_dead.append((task_id, websocket))
for task_id, ws in all_dead:
self.disconnect(task_id, ws)
@property
def active_connections_count(self) -> int:
"""Количество активных WebSocket соединений."""
return sum(len(conns) for conns in self._connections.values())
# Глобальный экземпляр менеджера
ws_manager = ConnectionManager()