fix(indexer): возвращать месячную квоту при провале извлечения файла
Живой репорт: юзер получил "Превышен месячный лимит проверок (тариф 'free'):
1/1" сразу после ЕДИНСТВЕННОЙ попытки — а та попытка провалилась ещё на
"файл не является PDF" (см. предыдущий коммит df2dfcc). Причина: API списывает
месячную квоту plagiarism синхронно при ЗАГРУЗКЕ файла (check_and_increment_
limit в documents.py), ДО того как воркер вообще попытается его распарсить —
реальной проверки не было, а квота уже списана навсегда (сброс только в
следующем месяце). У free-тарифа лимит 1/мес — то есть один неверный формат
файла сжигал единственную попытку целиком.
Плюс сопутствующая неэффективность: extract_and_check ретраил (3 попытки,
60с задержка) даже детерминированные ошибки формата — на 2-й и 3-й попытке
результат будет тем же, ретрай только откладывает финальный фидбек юзеру
на пару минут без всякого смысла.
Фикс:
- ValueError (битый файл/пустой текст) теперь ловится ДО общего Exception:
без ретрая (не поможет), с возвратом квоты (реальной проверки не было).
- db.refund_plagiarism_quota(): декремент того же Redis-ключа
rl:{user_id}:plagiarism:{YYYY-MM}, что инкрементит api/rate_limiter.py —
тот же формат ключа, декремент виден мгновенно и там, и там (общий Redis).
- Транзиентные ошибки (сеть/MinIO/БД) — поведение прежнее (ретрай, без
возврата квоты, т.к. задача может ещё успешно завершиться).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -1,9 +1,11 @@
|
|||||||
"""Синхронное подключение к PostgreSQL и MinIO для индексер-воркера."""
|
"""Синхронное подключение к PostgreSQL, MinIO и Redis для индексер-воркера."""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
from collections.abc import Generator
|
from collections.abc import Generator
|
||||||
from contextlib import contextmanager
|
from contextlib import contextmanager
|
||||||
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
|
import redis
|
||||||
from minio import Minio
|
from minio import Minio
|
||||||
from sqlalchemy import create_engine
|
from sqlalchemy import create_engine
|
||||||
from sqlalchemy.orm import Session, sessionmaker
|
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:
|
if error:
|
||||||
task.error = error
|
task.error = error
|
||||||
session.commit()
|
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}")
|
||||||
|
|||||||
@@ -13,7 +13,7 @@ from app.algorithms.minhash import add_to_lsh, find_similar
|
|||||||
from app.algorithms.winnowing import winnow
|
from app.algorithms.winnowing import winnow
|
||||||
from app.celery_app import celery_app
|
from app.celery_app import celery_app
|
||||||
from app.config import settings
|
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.docx import extract_text_from_docx, extract_text_from_txt
|
||||||
from app.extractors.pdf import extract_text_from_pdf
|
from app.extractors.pdf import extract_text_from_pdf
|
||||||
from app.fragments import split_into_fragments
|
from app.fragments import split_into_fragments
|
||||||
@@ -202,6 +202,16 @@ def extract_and_check(
|
|||||||
"level2": len(level2_matches),
|
"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:
|
except Exception as exc:
|
||||||
logger.error(f"Ошибка при обработке задачи {task_id!r}: {exc}", exc_info=True)
|
logger.error(f"Ошибка при обработке задачи {task_id!r}: {exc}", exc_info=True)
|
||||||
update_task_status(task_id, "failed", str(exc))
|
update_task_status(task_id, "failed", str(exc))
|
||||||
|
|||||||
Reference in New Issue
Block a user