"""Синхронное подключение к PostgreSQL, MinIO и Redis для индексер-воркера.""" import logging from collections.abc import Generator from contextlib import contextmanager from datetime import UTC, datetime import redis from minio import Minio from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker from app.config import settings logger = logging.getLogger(__name__) # Синхронный движок SQLAlchemy engine = create_engine( settings.database_url_sync, pool_size=5, max_overflow=10, pool_pre_ping=True, pool_recycle=3600, ) SessionLocal = sessionmaker(bind=engine, autocommit=False, autoflush=False) @contextmanager def db_session() -> Generator[Session, None, None]: """Контекстный менеджер для сессии БД.""" session = SessionLocal() try: yield session session.commit() except Exception: session.rollback() raise finally: session.close() _minio_client: Minio | None = None def get_minio() -> Minio: """Получить или создать MinIO клиент.""" global _minio_client if _minio_client is None: _minio_client = Minio( settings.MINIO_ENDPOINT, access_key=settings.MINIO_ACCESS_KEY, secret_key=settings.MINIO_SECRET_KEY, secure=False, ) return _minio_client def update_task_status(task_id: str, status: str, error: str | None = None) -> None: """Обновить статус задачи в БД.""" from app.models import Task with db_session() as session: task = session.get(Task, task_id) if task: task.status = status if error: task.error = error session.commit() _redis_client: redis.Redis | None = None def get_redis() -> redis.Redis: """Получить или создать sync Redis клиент (тот же Redis, что и у API — там считаются rate-limit'ы, ключи вида rl:{user_id}:{action}:{period}).""" global _redis_client if _redis_client is None: _redis_client = redis.Redis.from_url(settings.REDIS_URL, decode_responses=True) return _redis_client def refund_plagiarism_quota(task_id: str) -> None: """Вернуть месячную квоту проверок плагиата владельцу задачи. Вызывается, когда задача провалилась ДО начала реальной проверки (битый файл, не тот формат) — API списывает квоту синхронно при загрузке файла, заранее, не дожидаясь, распознается ли он вообще. Без возврата пользователь терял бы месячный лимит (у free — 1 в месяц) за одну неудачную попытку с неправильным файлом. Ключ и формат периода — как в api/app/core/rate_limiter.py (rl:{user_id}:plagiarism:{YYYY-MM}), чтобы декремент попадал в тот же счётчик, что инкрементил API. """ from app.models import Task try: with db_session() as session: task = session.get(Task, task_id) if task is None: return user_id = task.user_id period = datetime.now(UTC).strftime("%Y-%m") key = f"rl:{user_id}:plagiarism:{period}" r = get_redis() current = r.get(key) if current and int(current) > 0: r.decr(key) logger.info(f"Квота plagiarism возвращена пользователю {user_id} (задача {task_id!r})") except Exception as e: # Невозврат квоты — не повод валить обработку ошибки задачи logger.warning(f"Не удалось вернуть квоту для задачи {task_id!r}: {e}")