"""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)}