Files
Deal/archive/leadradar-legacy/backend/app/services/discovery_worker.py
T
Rustam Khalimov 713d554dc2 Нормализовать переводы строк в LF
Решение по 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.
2026-09-11 19:01:42 +03:00

485 lines
26 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Discovery worker: поиск → оценка → авто-вступление (Task 6).
Фоновый цикл main._discovery_loop вызывает `tick()` каждые ~5 секунд; каждый
вызов выполняет ОДНО действие для самой старой running-задачи и возвращает
{"action": "search"|"review"|"skip"|"join"|"reject"|"flood"|"error"|"done"|"none",
"taskId": ...}. Если работы нет — {"action": "none"}.
Приоритеты внутри tick:
0. ban_guard.global_paused() — ручной стоп-кран: возвращаем none;
1. задача достигла плана вступлений (joined >= planJoins) → status=done
(+ лог done) — занимаемый ею бюджет планов освобождается сразу;
2. шаг поиска: search_done=False → следующий ключ keywords[search_idx],
tg.discovery_search, каждый результат — discovery.add_candidate,
discovery.advance_search; при переходе search_done=True — лог search
«поиск завершён: N кандидатов» (N = счётчик found задачи). Личные чаты/
боты (kind «чат») пропускаются (лог skip). FloodWaitError → note_flood +
лог flood; прочие ошибки поиска → лог error; индекс ключей не двигается
(тик повторит ключ позже), но после 3 ошибок подряд ключ пропускается
(advance_search + лог error «ключ пропущен») — битый ключ не должен
зацикливать поиск навсегда;
2a. флуд-стоп: пока действует flood-блокировка дня (ban_guard.flood_today())
discovery-воркер полностью стоит (никаких сетевых действий поиска/оценки/
вступлений) — это и есть «стоп до конца суток» из лога flood;
3. шаг оценки: первый кандидат status='new' → tg.discovery_info (участники,
kind; форум — если is_forum), фильтр minSubscribers (меньше — delete +
лог skip; не получено — метка), tg.discovery_read: история недоступна →
review с меткой («канал…»/«закрытая группа…», контент не оцениваем),
язык (для ru-задач), evaluate_sample (форумы — по темам через
group_by_topic) → passed → review с fitRatio/topics, иначе
delete_candidate + лог skip «мало подходящих (X из N)»;
4. авто-вступление (отдельный проход, приоритет ниже оценки): у running-
задачи с autoJoin и кандидатом status='review', если
ban_guard.can_auto_join() → повторная проверка «мы не состоим»
(dialogs/disc_blacklist — могли вступить между оценкой и join) → если
уже состоим/в чёрном списке — discovery.mark_rejected + лог;
ban_guard.wait_join_delay() (5070 с — спейсинг авто-вступлений), ПОСЛЕ
паузы кандидат перечитывается — join выполняется, только если запись
ещё есть и в review, задача ещё running с autoJoin и мы не состоим
(иначе — тик выходит без join);
tg.discovery_join(username) → discovery.mark_joined(auto=True) →
tg.add_dialog_monitored(...) → tg.backfill_dialog(dialog_id) →
discovery.remove_blacklist. FloodWaitError → ban_guard.note_flood()
+ лог flood (кандидат остаётся review); прочие ошибки join →
join_failures += 1, лог error, кандидат остаётся review для повтора
(ретраи ограничены: после 3-й неудачи кандидат удаляется, лог skip).
Метки кандидата (marks) — список строк-чипов:
«участники не подтверждены», «язык не подтверждён»,
«канал: история недоступна», «закрытая группа (история скрыта) — вступите сами»,
«мало сообщений».
Темы форума (topics) — список dict (заполняется только для kind='forum'):
{"topicId": str|int, "title": str, "fitCount": int, "total": int,
"fitRatio": float, "passed": bool}
fit каждой темы считается по своей выборке (evaluate_sample + passed);
общий вердикт форума — есть хотя бы одна проходная тема. fitRatio кандидата
— агрегат по всей выборке (fit из N). Для не-форумов topics не заполняется.
"""
from __future__ import annotations
import logging
import time
from contextlib import suppress
from telethon.errors.rpcerrorlist import FloodWaitError
from ..db import store
from . import ban_guard, discovery
from . import discovery_eval as eval_svc
from .telegram import tg
log = logging.getLogger("leadradar.discovery_worker")
_ACTION_NONE = {"action": "none"}
# минимальный объём содержательной выборки для вердикта оценки
_MIN_CONTENT = 3
# ошибки поиска одного ключа подряд, после которых ключ пропускается
_SEARCH_ERRORS_TO_SKIP = 3
# счётчик ошибок поиска по задачам (память процесса; при рестарте сбрасывается)
_search_errors: dict[str, int] = {}
# метки кандидата (marks) — чипы в UI
_MARK_PARTICIPANTS = "участники не подтверждены"
_MARK_LANG = "язык не подтверждён"
_MARK_CHANNEL_NO_HISTORY = "канал: история недоступна"
_MARK_CLOSED_GROUP = "закрытая группа (история скрыта) — вступите сами"
_MARK_FEW_MESSAGES = "мало сообщений"
# kind из Telegram (_kind_of: канал/группа/чат) → код кандидата (channel/group/forum)
_KIND_RU_TO_CODE = {"канал": "channel", "группа": "group", "чат": "group"}
def _now_ms() -> int:
return time.time_ns() // 1_000_000
def _kind_code(kind: str, is_forum: bool = False) -> str:
"""Нормализация kind источника: 'channel'|'group'|'forum' (см. _KIND_RU_TO_CODE)."""
if is_forum:
return "forum"
code = str(kind or "").strip().casefold()
if code in ("channel", "group", "forum"):
return code
return _KIND_RU_TO_CODE.get(code, "group")
def _running_tasks() -> list[dict]:
"""Running-задачи, старые первыми (list_tasks сортирует по created_at)."""
return [t for t in discovery.list_tasks() if t["status"] == "running"]
def _we_are_in(dialog_id: str) -> bool:
"""Уже состоим/отклонили: источник в dialogs или в чёрном списке.
Повторная проверка «мы не состоим» перед авто-вступлением (диалог мог
появиться между оценкой кандидата и join). dialog_id — подписанный peer id,
как в dialogs.id (конвенция discovery_search, Task 4).
"""
in_dialogs = store.scalar("SELECT 1 FROM dialogs WHERE id = ? LIMIT 1", [dialog_id])
in_blacklist = store.scalar("SELECT 1 FROM disc_blacklist WHERE dialog_id = ? LIMIT 1", [dialog_id])
return bool(in_dialogs or in_blacklist)
def _finish_done(task: dict) -> None:
"""Задача выполнила план вступлений: status=done + лог done."""
store.execute(
"UPDATE disc_tasks SET status = 'done', updated_at = ? WHERE id = ?",
[_now_ms(), task["id"]],
)
discovery.add_log(
task["id"],
"done",
f"план выполнен: вступили {task['joined']} из {task['planJoins']}",
)
def _close_search(task_id: str) -> None:
"""Закрыть проход по ключам (пустой список ключей / индекс за границей)."""
discovery.advance_search(task_id)
task = discovery.get_task(task_id)
if task and task["searchDone"]:
discovery.add_log(task_id, "search", f"поиск завершён: {task['found']} кандидатов")
def _log_search_done(task_id: str) -> None:
"""Лог завершения поиска, если advance_search перевёл задачу в search_done."""
task = discovery.get_task(task_id)
if task and task["searchDone"]:
discovery.add_log(task_id, "search", f"поиск завершён: {task['found']} кандидатов")
def _finish_review(
task_id: str,
dialog_id: str,
*,
marks: list[str],
lang_ru: bool | None = None,
fit_ratio: float | None = None,
topics: list[dict] | None = None,
) -> None:
"""Перевести кандидата в review с метками/оценкой (статус пишет лог review)."""
patch: dict = {"marks": [m for m in marks if m]}
if lang_ru is not None:
patch["langRu"] = lang_ru
if fit_ratio is not None:
patch["fitRatio"] = fit_ratio
if topics is not None:
patch["topics"] = topics
discovery.set_candidate(task_id, dialog_id, patch)
discovery.set_candidate_status(dialog_id, "review")
# ─── шаги tick ─────────────────────────────────────────────────────────────
async def _search_step(task: dict) -> dict:
"""Шаг поиска: один ключ keywords[searchIdx] → кандидаты + advance_search."""
task_id = task["id"]
keywords = list(task.get("keywords") or [])
idx = int(task.get("searchIdx") or 0)
if not keywords or idx >= len(keywords):
# ключи закончились/пустой список: закрываем проход без сетевого вызова
_close_search(task_id)
return {"action": "search", "taskId": task_id}
keyword = keywords[idx]
try:
results = await tg.discovery_search(keyword)
except FloodWaitError:
# флуд: стоп авто-вступлений до конца суток; ключ не двигаем — повторим позже
ban_guard.note_flood()
discovery.add_log(task_id, "flood", f"поиск «{keyword}»: flood — стоп до конца суток")
return {"action": "flood", "taskId": task_id}
except Exception as exc: # noqa: BLE001 — сбой поиска не двигает индекс ключей
errors = _search_errors.get(task_id, 0) + 1
if errors >= _SEARCH_ERRORS_TO_SKIP:
# 3 ошибки подряд одного ключа: пропускаем (битый ключ не должен
# зацикливать поиск и блокировать оценку/вступления других задач)
_search_errors.pop(task_id, None)
discovery.add_log(task_id, "error", f"поиск «{keyword}»: {exc} — ключ пропущен ({errors} ошибки подряд)")
discovery.advance_search(task_id)
_log_search_done(task_id)
return {"action": "error", "taskId": task_id}
_search_errors[task_id] = errors
discovery.add_log(task_id, "error", f"поиск «{keyword}»: {exc}")
return {"action": "error", "taskId": task_id}
_search_errors.pop(task_id, None) # успешный поиск — сброс счётчика ошибок ключа
for item in results:
kind_raw = str(item.get("kind") or "").strip().casefold()
name = item.get("name") or ""
if kind_raw == "чат":
# люди/личные чаты и боты глобальным поиском не предлагаются
discovery.add_log(task_id, "skip", f"{name}: личный чат/бот")
continue
discovery.add_candidate(
task_id,
str(item.get("id") or ""),
name,
item.get("username") or "",
_kind_code(kind_raw),
item.get("hue") or "#666",
)
discovery.advance_search(task_id)
_log_search_done(task_id)
return {"action": "search", "taskId": task_id}
async def _eval_step(task: dict, cand: dict) -> dict:
"""Шаг оценки первого кандидата status='new' (все ветки — одно действие)."""
task_id = task["id"]
dialog_id = cand["dialogId"]
marks: list[str] = []
# ── инфо об источнике: kind/forum, участники, имя/username ────────────
info = await tg.discovery_info(dialog_id)
is_forum = bool(info.get("is_forum"))
resolved = info.get("kind") or is_forum
kind = _kind_code(info.get("kind", ""), is_forum) if resolved else (cand["kind"] or "channel")
cand_patch: dict = {
"kind": kind,
"hue": info.get("hue") or "",
"participants": info.get("participants"),
}
if resolved:
# имя/username обновляем только при успешном резолве: при fallback
# discovery_info возвращает name=dialog_id и не должен затирать имя
cand_patch["name"] = info.get("name") or ""
cand_patch["username"] = info.get("username") or ""
discovery.set_candidate(task_id, dialog_id, cand_patch)
participants = info.get("participants")
# ── фильтр minSubscribers ──────────────────────────────────────────────
min_sub = int(task.get("minSubscribers") or 0)
if min_sub > 0:
if participants is None:
marks.append(_MARK_PARTICIPANTS)
elif int(participants) < min_sub:
discovery.delete_candidate(dialog_id)
discovery.add_log(
task_id,
"skip",
f"{dialog_id}: мало участников ({participants} < {min_sub})",
)
discovery.bump_counter(task_id, "evaluated")
return {"action": "skip", "taskId": task_id}
# ── чтение истории для оценки ──────────────────────────────────────────
read = await tg.discovery_read(dialog_id, int(task.get("sampleSize") or 10))
if not read.get("ok"):
# история недоступна без членства: контент не оцениваем, фильтры помечаем
marks.append(_MARK_CHANNEL_NO_HISTORY if kind == "channel" else _MARK_CLOSED_GROUP)
if task.get("lang") == "ru":
marks.append(_MARK_LANG)
_finish_review(task_id, dialog_id, marks=marks)
discovery.bump_counter(task_id, "evaluated")
return {"action": "review", "taskId": task_id}
messages = read.get("messages") or []
# ── язык (только для ru-задач) ─────────────────────────────────────────
lang_ru: bool | None = None
if task.get("lang") == "ru":
lang_ru = eval_svc.detect_lang_ru([str(m.get("text") or "") for m in messages])
if lang_ru is False:
discovery.delete_candidate(dialog_id)
discovery.add_log(task_id, "skip", f"{dialog_id}: язык не русский")
discovery.bump_counter(task_id, "evaluated")
return {"action": "skip", "taskId": task_id}
if lang_ru is None:
marks.append(_MARK_LANG)
# ── объём выборки: меньше 3 содержательных — решает человек ────────────
if len(messages) < _MIN_CONTENT:
marks.append(_MARK_FEW_MESSAGES)
_finish_review(task_id, dialog_id, marks=marks, lang_ru=lang_ru)
discovery.bump_counter(task_id, "evaluated")
return {"action": "review", "taskId": task_id}
# ── оценка содержания (форумы — по темам) ──────────────────────────────
fit_count, total, fit_ratio, topics, ok = await _evaluate_content(task, kind, messages)
if ok:
_finish_review(
task_id,
dialog_id,
marks=marks,
lang_ru=lang_ru,
fit_ratio=fit_ratio,
topics=topics if kind == "forum" else None,
)
discovery.bump_counter(task_id, "evaluated")
return {"action": "review", "taskId": task_id}
discovery.delete_candidate(dialog_id)
discovery.add_log(task_id, "skip", f"{dialog_id}: мало подходящих ({fit_count} из {total})")
discovery.bump_counter(task_id, "evaluated")
return {"action": "skip", "taskId": task_id}
async def _evaluate_content(task: dict, kind: str, messages: list[dict]) -> tuple[int, int, float, list[dict], bool]:
"""evaluate_sample по выборке кандидата.
Не-форум: один прогон по всем сообщениям, topics пуст, вердикт — passed().
Форум: прогон по каждой теме (group_by_topic), topics заполняется
({topicId, title, fitCount, total, fitRatio, passed}), вердикт — есть хотя
бы одна проходная тема; fitRatio — агрегат fit из N по всей выборке.
"""
if kind == "forum":
topics: list[dict] = []
fit_count = 0
total = 0
any_passed = False
for group in eval_svc.group_by_topic(messages):
ev = await eval_svc.evaluate_sample(task, group["messages"])
t_ok = eval_svc.passed(ev, task)
fit_count += int(ev.get("fit_count") or 0)
total += int(ev.get("total") or 0)
topics.append(
{
"topicId": group["topic_id"],
"title": group["title"] or "",
"fitCount": int(ev.get("fit_count") or 0),
"total": int(ev.get("total") or 0),
"fitRatio": float(ev.get("fit_ratio") or 0.0),
"passed": bool(t_ok),
}
)
any_passed = any_passed or bool(t_ok)
return fit_count, total, (fit_count / total) if total else 0.0, topics, any_passed
ev = await eval_svc.evaluate_sample(task, messages)
return (
int(ev.get("fit_count") or 0),
int(ev.get("total") or 0),
float(ev.get("fit_ratio") or 0.0),
[],
eval_svc.passed(ev, task),
)
async def _join_step(task: dict, cand: dict) -> dict:
"""Шаг авто-вступления одного кандидата status='review'."""
task_id = task["id"]
dialog_id = cand["dialogId"]
# между оценкой и вступлением могли вступить/отклонить источник
if _we_are_in(dialog_id):
already = bool(store.scalar("SELECT 1 FROM dialogs WHERE id = ? LIMIT 1", [dialog_id]))
reason = (
"уже вступили между оценкой и авто-вступлением"
if already
else "источник в чёрном списке (повторная проверка перед авто-вступлением)"
)
with suppress(KeyError, ValueError):
discovery.mark_rejected(dialog_id, reason=reason)
return {"action": "reject", "taskId": task_id}
# спейсинг авто-вступлений (сек из настроек discJoinDelayMin/Max)
await ban_guard.wait_join_delay()
# за время паузы задача/кандидат/состояние BanGuard могли измениться:
# вступаем только если кандидат всё ещё есть и в review, задача ещё running
# с autoJoin, мы не состоим и авто-вступления по-прежнему разрешены (стоп-
# кран/flood/лимит могли включиться во время паузы) — иначе выходим без join
fresh = store.query_one("SELECT * FROM disc_candidates WHERE dialog_id = ?", [dialog_id])
task_now = discovery.get_task(task_id)
if (
fresh is None
or fresh["status"] != "review"
or task_now is None
or task_now["status"] != "running"
or not task_now["autoJoin"]
or _we_are_in(dialog_id)
or not ban_guard.can_auto_join()
):
return _ACTION_NONE
username = str(fresh.get("username") or "").strip().lstrip("@")
try:
await tg.discovery_join(username)
except FloodWaitError:
ban_guard.note_flood() # идемпотентно: discovery_join тоже фиксирует флуд
discovery.add_log(task_id, "flood", f"авто-вступление {dialog_id}: flood — стоп до конца суток")
return {"action": "flood", "taskId": task_id}
except Exception as exc: # noqa: BLE001 — ретраи ограничены счётчиком join_failures
# между паузой и неудачным join кандидата могли отклонить/удалить:
# счётчик и удаление трогаем только у живой записи в статусе review
row_now = store.query_one(
"SELECT status FROM disc_candidates WHERE dialog_id = ?", [dialog_id]
)
if not row_now or row_now["status"] != "review":
return _ACTION_NONE
failures = int(fresh.get("join_failures") or 0) + 1
store.execute(
"UPDATE disc_candidates SET join_failures = ?, updated_at = ? "
"WHERE dialog_id = ? AND status = 'review'",
[failures, _now_ms(), dialog_id],
)
if failures >= 3:
discovery.delete_candidate(dialog_id)
discovery.add_log(task_id, "skip", f"{dialog_id}: не удалось вступить (3 попытки): {exc}")
return {"action": "skip", "taskId": task_id}
discovery.add_log(task_id, "error", f"авто-вступление {dialog_id}: {exc}")
return {"action": "error", "taskId": task_id}
discovery.mark_joined(dialog_id, auto=True)
tg.add_dialog_monitored(
dialog_id,
fresh.get("name"),
username,
fresh.get("kind"),
fresh.get("hue"),
)
# разбор последних сообщений источника (спейсинг/read-ack внутри метода);
# вступление уже состоялось — сбой backfill не роняет шаг
try:
await tg.backfill_dialog(dialog_id)
except Exception as exc: # noqa: BLE001
log.warning("join %s: backfill не удался: %s", dialog_id, exc)
discovery.remove_blacklist(dialog_id)
return {"action": "join", "taskId": task_id}
# ─── tick ──────────────────────────────────────────────────────────────────
async def tick() -> dict:
"""Одно действие discovery-воркера (см. docstring модуля)."""
if ban_guard.global_paused():
return _ACTION_NONE
if ban_guard.flood_today():
# флуд-блокировка дня: никаких сетевых действий (поиск/оценка/join),
# пока действует discFloodDay — воркер просто стоит
return _ACTION_NONE
running = _running_tasks()
if not running:
return _ACTION_NONE
# 1. план достигнут — закрываем задачу (важно до поиска/оценки/join:
# задачу с выполненным планом нельзя продолжать обрабатывать)
for task in running:
if int(task["joined"]) >= int(task["planJoins"]):
_finish_done(task)
return {"action": "done", "taskId": task["id"]}
# 2. поиск: следующая running-задача с незавершённым проходом по ключам
for task in running:
if not task["searchDone"]:
return await _search_step(task)
# 3. оценка: первый кандидат status='new' (самая старая задача — первой)
for task in running:
new_cands = discovery.list_candidates(task["id"], status="new")
if new_cands:
return await _eval_step(task, new_cands[0])
# 4. авто-вступление: отдельный проход, приоритет ниже оценки
for task in running:
if not task["autoJoin"]:
continue
review_cands = discovery.list_candidates(task["id"], status="review")
if review_cands:
if not ban_guard.can_auto_join():
return _ACTION_NONE # суточный лимит/флуд/пауза — join никому нельзя
return await _join_step(task, review_cands[0])
return _ACTION_NONE