diff --git a/services/worker-indexer/app/db.py b/services/worker-indexer/app/db.py index 94cf6c9..314f64b 100644 --- a/services/worker-indexer/app/db.py +++ b/services/worker-indexer/app/db.py @@ -1,9 +1,11 @@ -"""Синхронное подключение к PostgreSQL и MinIO для индексер-воркера.""" +"""Синхронное подключение к 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 @@ -65,3 +67,47 @@ def update_task_status(task_id: str, status: str, error: str | None = None) -> N 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}") diff --git a/services/worker-indexer/app/tasks/index.py b/services/worker-indexer/app/tasks/index.py index f710112..c90e73d 100644 --- a/services/worker-indexer/app/tasks/index.py +++ b/services/worker-indexer/app/tasks/index.py @@ -13,7 +13,7 @@ from app.algorithms.minhash import add_to_lsh, find_similar from app.algorithms.winnowing import winnow from app.celery_app import celery_app from app.config import settings -from app.db import db_session, get_minio, update_task_status +from app.db import db_session, get_minio, refund_plagiarism_quota, update_task_status from app.extractors.docx import extract_text_from_docx, extract_text_from_txt from app.extractors.pdf import extract_text_from_pdf from app.fragments import split_into_fragments @@ -202,6 +202,16 @@ def extract_and_check( "level2": len(level2_matches), } + except ValueError as exc: + # Детерминированная ошибка (битый файл, не тот формат, пустой текст) — + # повтор не поможет: результат будет тем же на 2-й и 3-й попытке. Не + # ретраим (быстрее фидбек юзеру) и возвращаем месячную квоту — реальная + # проверка так и не началась, списывать не за что. + logger.warning(f"Задача {task_id!r}: не удалось обработать файл: {exc}") + update_task_status(task_id, "failed", str(exc)) + refund_plagiarism_quota(task_id) + return {"task_id": task_id, "status": "failed", "error": str(exc)} + except Exception as exc: logger.error(f"Ошибка при обработке задачи {task_id!r}: {exc}", exc_info=True) update_task_status(task_id, "failed", str(exc))