"""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() (50–70 с — спейсинг авто-вступлений), ПОСЛЕ паузы кандидат перечитывается — 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