Решение по TD-STYLE-ANALYZERS: LF — инструменты проекта (Python/Node) пишут LF, CRLF-.sh не работают на Linux CI (sh scripts/ci.sh), большинство файлов уже были LF. Добавлен .gitattributes (* text=auto eol=lf, бинарные исключения), .editorconfig переведён на lf, 1029 файлов конвертированы, git add --renormalize. Из индекса убраны закравшиеся archive/**/__pycache__/*.pyc.
286 lines
12 KiB
Python
286 lines
12 KiB
Python
"""API Discovery (Task 7): задачи поиска каналов, кандидаты, чёрный список, лог.
|
||
|
||
Prefix /api/discovery, авторизация — current_login (как в соседних роутерах).
|
||
Сервис discovery отдаёт наружу camelCase-словари (см. его docstring), поэтому
|
||
Pydantic-модели повторяют имена полей API без алиасов (как PreviewBody в tg_routes).
|
||
Списки наружу — {"items": [...]}, единичные объекты — как есть (конвенция проекта).
|
||
|
||
Обработка ошибок контракта: ValueError -> HTTP 400, KeyError -> HTTP 404.
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import logging
|
||
from typing import Literal
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException
|
||
from pydantic import BaseModel
|
||
|
||
from ..auth import current_login
|
||
from ..db import store
|
||
from ..services import ai as ai_service
|
||
from ..services import discovery
|
||
from ..services.telegram import _spawn, tg
|
||
|
||
log = logging.getLogger("leadradar.discovery_api")
|
||
|
||
router = APIRouter(prefix="/api/discovery", tags=["discovery"])
|
||
|
||
# потолок текста описания, уходящего ИИ-генератору ключей
|
||
_AI_DESCRIPTION_LIMIT = 4000
|
||
# страховочный потолок числа сгенерированных ключей (промпт просит 10–16)
|
||
_KEYWORDS_LIMIT = 30
|
||
# потолок длины одного ключа (короткие фразы для поиска Telegram)
|
||
_KEYWORD_LENGTH_LIMIT = 60
|
||
|
||
# промпт генерации ключевых слов по описанию задачи (RU+EN, глобальный поиск)
|
||
_KEYWORDS_PROMPT = (
|
||
"Ты — эксперт по поиску Telegram-каналов и групп. По описанию ниши/задачи "
|
||
"составь поисковые ключевые слова, по которым в глобальном поиске Telegram "
|
||
"находят подходящие источники. Верни строго JSON вида "
|
||
'{"keywords": ["...", "..."]}. Требования к списку:\n'
|
||
"- 10–16 ключей;\n"
|
||
"- примерно поровну русских и английских (английские — популярные в нише термины);\n"
|
||
"- короткие фразы 1–4 слова;\n"
|
||
"- без #, @, кавычек и лишней пунктуации;\n"
|
||
"- конкретные для ниши, включая сленг заказчиков и подрядчиков;\n"
|
||
"- без дублей и близких по смыслу повторов."
|
||
)
|
||
|
||
|
||
class TaskCreate(BaseModel):
|
||
name: str
|
||
description: str = ""
|
||
keywords: list[str] = []
|
||
minSubscribers: int = 0
|
||
lang: str = "ru"
|
||
threshold: int | None = None
|
||
sampleSize: int | None = None
|
||
planJoins: int = 1
|
||
autoJoin: bool = False
|
||
|
||
|
||
class TaskPatch(BaseModel):
|
||
name: str | None = None
|
||
description: str | None = None
|
||
keywords: list[str] | None = None
|
||
minSubscribers: int | None = None
|
||
lang: str | None = None
|
||
threshold: int | None = None
|
||
sampleSize: int | None = None
|
||
planJoins: int | None = None
|
||
autoJoin: bool | None = None
|
||
|
||
|
||
def _payload(body: BaseModel) -> dict:
|
||
"""Поля модели -> payload сервиса, без явных None.
|
||
|
||
discovery.create_task/patch_task сами подставляют значения по умолчанию
|
||
(в т.ч. threshold/sampleSize из настроек), поэтому None-поля пропускаем.
|
||
"""
|
||
return body.model_dump(exclude_none=True)
|
||
|
||
|
||
def _task_or_404(task_id: str) -> dict:
|
||
task = discovery.get_task(task_id)
|
||
if task is None:
|
||
raise HTTPException(404, "Задача не найдена")
|
||
return task
|
||
|
||
|
||
def _candidate_or_404(dialog_id: str) -> dict:
|
||
row = store.query_one("SELECT * FROM disc_candidates WHERE dialog_id = ?", [dialog_id])
|
||
if row is None:
|
||
raise HTTPException(404, "Кандидат не найден")
|
||
return row
|
||
|
||
|
||
def _ai_unavailable_reason() -> str | None:
|
||
"""Причина недоступности ИИ (None — можно вызывать)."""
|
||
if not store.get_setting("aiEnabled"):
|
||
return "ИИ выключен в настройках (aiEnabled)"
|
||
try:
|
||
status = ai_service.provider_status()
|
||
except Exception as exc: # noqa: BLE001 — статус не читается = ИИ недоступен
|
||
log.debug("generate-keywords: статус ИИ недоступен (%s)", exc)
|
||
return "Не удалось прочитать статус ИИ-провайдера"
|
||
if not (status.get("local") or status.get("keySet")):
|
||
return "Не задан API-ключ ИИ-провайдера"
|
||
return None
|
||
|
||
|
||
def _clean_keywords(raw) -> list[str]:
|
||
"""Ключи из ответа ИИ: строки без пустых/длинных и повторов (casefold)."""
|
||
seen: set[str] = set()
|
||
out: list[str] = []
|
||
for item in raw or []:
|
||
if not isinstance(item, str):
|
||
continue
|
||
keyword = item.strip()
|
||
if not keyword or len(keyword) > _KEYWORD_LENGTH_LIMIT:
|
||
continue
|
||
key = keyword.casefold()
|
||
if key in seen:
|
||
continue
|
||
seen.add(key)
|
||
out.append(keyword)
|
||
if len(out) >= _KEYWORDS_LIMIT:
|
||
break
|
||
return out
|
||
|
||
|
||
async def _backfill_quiet(dialog_id: str) -> None:
|
||
"""Догон последних сообщений вступившего источника (фон, best-effort)."""
|
||
try:
|
||
await tg.backfill_dialog(dialog_id)
|
||
except Exception as exc: # noqa: BLE001 — вступление уже состоялось
|
||
log.warning("join %s: backfill не удался: %s", dialog_id, exc)
|
||
|
||
|
||
# ─── задачи ────────────────────────────────────────────────────────────────
|
||
|
||
@router.get("/tasks")
|
||
def list_tasks(_: str = Depends(current_login)) -> dict:
|
||
return {"items": discovery.list_tasks()}
|
||
|
||
|
||
@router.post("/tasks")
|
||
def create_task(body: TaskCreate, _: str = Depends(current_login)) -> dict:
|
||
try:
|
||
return discovery.create_task(_payload(body))
|
||
except ValueError as exc:
|
||
raise HTTPException(400, str(exc)) from exc
|
||
|
||
|
||
@router.patch("/tasks/{task_id}")
|
||
def patch_task(task_id: str, body: TaskPatch, _: str = Depends(current_login)) -> dict:
|
||
try:
|
||
return discovery.patch_task(task_id, _payload(body))
|
||
except KeyError as exc:
|
||
raise HTTPException(404, "Задача не найдена") from exc
|
||
except ValueError as exc:
|
||
raise HTTPException(400, str(exc)) from exc
|
||
|
||
|
||
@router.delete("/tasks/{task_id}")
|
||
def delete_task(task_id: str, _: str = Depends(current_login)) -> dict:
|
||
_task_or_404(task_id)
|
||
discovery.delete_task(task_id)
|
||
return {"ok": True}
|
||
|
||
|
||
@router.post("/tasks/{task_id}/start")
|
||
def start_task(task_id: str, _: str = Depends(current_login)) -> dict:
|
||
try:
|
||
return discovery.start_task(task_id)
|
||
except KeyError as exc:
|
||
raise HTTPException(404, "Задача не найдена") from exc
|
||
except ValueError as exc:
|
||
raise HTTPException(400, str(exc)) from exc
|
||
|
||
|
||
@router.post("/tasks/{task_id}/pause")
|
||
def pause_task(task_id: str, _: str = Depends(current_login)) -> dict:
|
||
try:
|
||
return discovery.pause_task(task_id)
|
||
except KeyError as exc:
|
||
raise HTTPException(404, "Задача не найдена") from exc
|
||
|
||
|
||
@router.post("/tasks/{task_id}/generate-keywords")
|
||
async def generate_keywords(task_id: str, _: str = Depends(current_login)) -> dict:
|
||
"""ИИ-генерация ключей по описанию задачи: RU+EN, 10–16 строк.
|
||
|
||
ИИ выключен/не настроен/ответил ошибкой — {"keywords": [], "error": "..."}
|
||
с HTTP 200, чтобы UI показал причину, а не падал.
|
||
"""
|
||
task = _task_or_404(task_id)
|
||
reason = _ai_unavailable_reason()
|
||
if reason:
|
||
return {"keywords": [], "error": reason}
|
||
description = str(task.get("description") or "").strip()
|
||
if not description:
|
||
return {"keywords": [], "error": "У задачи нет описания — по нему генерируются ключи"}
|
||
try:
|
||
out = await ai_service.chat_json(
|
||
_KEYWORDS_PROMPT,
|
||
f"Описание ниши/задачи:\n{description[:_AI_DESCRIPTION_LIMIT]}",
|
||
)
|
||
except Exception as exc: # noqa: BLE001 — сбой провайдера не роняет API
|
||
log.warning("generate-keywords задача %s: ИИ не ответил: %s", task_id, exc)
|
||
return {"keywords": [], "error": str(exc)}
|
||
return {"keywords": _clean_keywords(out.get("keywords", []) if isinstance(out, dict) else [])}
|
||
|
||
|
||
# ─── кандидаты ─────────────────────────────────────────────────────────────
|
||
|
||
@router.get("/tasks/{task_id}/candidates")
|
||
def list_candidates(
|
||
task_id: str,
|
||
status: Literal["new", "review", "joined", "rejected"] | None = None,
|
||
_: str = Depends(current_login),
|
||
) -> dict:
|
||
_task_or_404(task_id)
|
||
return {"items": discovery.list_candidates(task_id, status)}
|
||
|
||
|
||
@router.post("/candidates/{dialog_id}/join")
|
||
async def join_candidate(dialog_id: str, _: str = Depends(current_login)) -> dict:
|
||
"""Ручное вступление (вне квот и пауз воркера).
|
||
|
||
tg.discovery_join -> add_dialog_monitored -> backfill_dialog (последние
|
||
сообщения, best-effort) -> mark_joined(auto=False); источник снимается с
|
||
чёрного списка. Ошибка Telegram -> 400 с текстом причины.
|
||
"""
|
||
row = _candidate_or_404(dialog_id)
|
||
if row["status"] == "joined":
|
||
raise HTTPException(400, "Уже вступили в этот источник")
|
||
username = str(row.get("username") or "")
|
||
try:
|
||
await tg.discovery_join(username)
|
||
except Exception as exc: # текст ошибки уходит наружу
|
||
raise HTTPException(400, f"Не удалось вступить в @{username}: {exc}") from exc
|
||
tg.add_dialog_monitored(dialog_id, row.get("name"), username, row.get("kind"), row.get("hue"))
|
||
# догон последних сообщений — в фоне: join из UI не должен висеть на
|
||
# паузах backfill (10 сообщений × 1.5–3 с); источник уже в мониторинге
|
||
_spawn(_backfill_quiet(dialog_id))
|
||
discovery.remove_blacklist(dialog_id)
|
||
try:
|
||
return discovery.mark_joined(dialog_id, auto=False)
|
||
except KeyError as exc:
|
||
raise HTTPException(404, "Кандидат не найден") from exc
|
||
|
||
|
||
@router.post("/candidates/{dialog_id}/reject")
|
||
def reject_candidate(dialog_id: str, _: str = Depends(current_login)) -> dict:
|
||
"""Отклонить кандидата (в чёрный список). Уже вступившего — нельзя."""
|
||
row = _candidate_or_404(dialog_id)
|
||
if row["status"] == "joined":
|
||
raise HTTPException(400, "Уже вступили — удалите источник из каналов")
|
||
try:
|
||
return discovery.mark_rejected(dialog_id, reason="отклонено вручную")
|
||
except KeyError as exc:
|
||
raise HTTPException(404, "Кандидат не найден") from exc
|
||
except ValueError as exc:
|
||
raise HTTPException(400, str(exc)) from exc
|
||
|
||
|
||
# ─── чёрный список ─────────────────────────────────────────────────────────
|
||
|
||
@router.get("/blacklist")
|
||
def list_blacklist(_: str = Depends(current_login)) -> dict:
|
||
return {"items": discovery.list_blacklist()}
|
||
|
||
|
||
@router.delete("/blacklist/{dialog_id}")
|
||
def remove_blacklist(dialog_id: str, _: str = Depends(current_login)) -> dict:
|
||
discovery.remove_blacklist(dialog_id)
|
||
return {"ok": True}
|
||
|
||
|
||
# ─── лог задачи ────────────────────────────────────────────────────────────
|
||
|
||
@router.get("/tasks/{task_id}/log")
|
||
def task_log(task_id: str, _: str = Depends(current_login)) -> dict:
|
||
_task_or_404(task_id)
|
||
return {"items": discovery.task_log(task_id)}
|