"""Пайплайн входящих (п.4, п.5 ТЗ) — очередь + фоновый воркер. Все сообщения из групп/каналов попадают в очередь pipeline_msg и разбираются фоновым воркером (main._pipeline_loop и ручной /api/admin/tick): queue(status='new') -> этап 1: стоп-фразы и длина (без ИИ); не прошло -> в отсев (rejected) -> дедупликация по нормализованному тексту (повтор -> в отсев) -> ML-слой: если уверен, решает сам (спам -> отсев, доска -> карточка сразу); иначе status='filtered' queue(status='filtered') -> этап 2: ИИ-фильтр (если включён; иначе пропуск) -> ИИ-классификация (если ML не решил) -> карточка в БД + SSE; строка очереди удаляется В БД оседают только данные, прошедшие фильтры (карточки + их исходные сообщения). Каждый отброс с причиной пишется в rejected_msgs (вкладка «Обработка»: очередь и отсев) и автоочищается раз в 3 суток. """ from __future__ import annotations import asyncio import json import logging import re import time from .. import constants as C from ..db import store from ..sse import broker from . import ai as ai_service from . import ml_client from . import processing as processing_svc from . import rules as rules_svc log = logging.getLogger("leadradar.pipeline") # воркер один: фоновый цикл и ручной /admin/tick не должны разбирать # одну и ту же строку одновременно (иначе возможны дубли карточек) _pump_lock = asyncio.Lock() # статусы очереди ST_NEW = "new" ST_AI = "filtered" def _now() -> int: return time.time_ns() // 1_000_000 # ─── Очередь ───────────────────────────────────────────────────────────── def enqueue( dialog_id: str, ch_name: str, ch_handle: str, ch_hue: str, msg_id: int | None, text: str, msg_at: int, force: bool = False, ) -> None: """Все входящие сообщения из мониторящихся диалогов. force=TRUE — сообщение возвращено пользователем из отсева: этап 1, устарело и ML-решения для него игнорируются (см. _pump_unlocked). """ text = (text or "").strip() if not text or not dialog_id: return now = _now() # защита от повторов (Telethon иногда отдаёт событие дважды): # строка живёт только пока сообщение в очереди/обработке dup = None if msg_id: dup = store.scalar( "SELECT 1 FROM pipeline_msg WHERE dialog_id = ? AND msg_id = ? LIMIT 1", [dialog_id, msg_id], ) if not dup: store.execute( "INSERT INTO pipeline_msg(id, dialog_id, ch_name, ch_handle, ch_hue, text, msg_id, msg_at, status, force, created_at, updated_at) " "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", [store.uid("p_"), dialog_id, ch_name, ch_handle, ch_hue, text[:6000], msg_id, msg_at, ST_NEW, bool(force), now, now], ) def queue_len() -> int: return int(store.scalar("SELECT count(*) FROM pipeline_msg") or 0) # ─── Этап 1: без ИИ ─────────────────────────────────────────────────────── def stage1_plain(text: str) -> dict: """Этап 1 без ИИ: минимальная длина, стоп-фразы, отсев резюме, тип заявки. Возвращает {'pass', 'reason', 'stage', 'kind', 'kw'}: kind — какое именно правило сработало (length|stop|resume|type), kw — конкретное слово/фраза стоп-списка (для мониторинга отсева «почему сюда попало»). """ t = (text or "").strip() min_len = int(store.get_setting("minLen") or C.DEFAULT_MIN_LEN) if len(t) < min_len: return {"pass": False, "reason": f"короче {min_len} символов", "stage": 1, "kind": "length", "kw": ""} lower = t.casefold() for phrase in store.get_setting("stopPhrases") or []: if phrase and phrase.casefold() in lower: return {"pass": False, "reason": f"стоп-фраза «{phrase}»", "stage": 1, "kind": "stop", "kw": phrase} if bool(store.get_setting("blockResumes")): # Система ищет вакансии/заказы (blockResumes включён) — резюме # соискателей отсекаем сразу, до правил колонок и ИИ. marker = _resume_reason(lower) if marker: return {"pass": False, "reason": f"резюме соискателя («{marker}»)", "stage": 1, "kind": "resume", "kw": marker} # Тип заявок: «только вакансии» / «только фриланс и заказы». Считается на # этапе 1 (до ML/ИИ, не тратит токены) по маркерам найма из «Сферы и ключей». wanted = str(store.get_setting("wantedType") or "both").strip().lower() if wanted in ("vacancy", "freelance"): looks_vacancy = any(mk in lower for mk in _hire_markers()) if wanted == "freelance" and looks_vacancy: return {"pass": False, "reason": "ищете только разовые заявки — сообщение похоже на вакансию/занятость", "stage": 1, "kind": "type", "kw": ""} if wanted == "vacancy" and not looks_vacancy: return {"pass": False, "reason": "ищете только занятость — сообщение похоже на разовый заказ/услугу", "stage": 1, "kind": "type", "kw": ""} return {"pass": True, "reason": None, "stage": 1, "kind": "", "kw": ""} # ─── Карточка ───────────────────────────────────────────────────────────── # ── Нормализация ответов ИИ/локального разбора ──────────────────────────── _MD_LINK_RE = re.compile(r"\[([^\]]*)\]\([^)\s]+\)") # [текст](url) -> текст _MD_BOLD2_RE = re.compile(r"__([^_\n]+?)__") _MD_CODE_RE = re.compile(r"`([^`\n]+?)`") _MD_STRIKE_RE = re.compile(r"~~([^~\n]+?)~~") _BARE_URL_RE = re.compile(r"https?://[^\s<>\"']+") # декоративные эмодзи/символы-маркеры (🔥 📍 💎 ‼ ℹ и т.п.) — мусор в структурированном тексте _EMOJI_RE = re.compile( "[" "\\U0001F000-\\U0001FAFF" # доп. пиктограммы/эмодзи "\\U0001F1E6-\\U0001F1FF" # региональные флаги "\\U00002600-\\U000027BF" # разные символы (☀⭐⚠…) "\\U00002B00-\\U00002BFF" # стрелки-символы "\\U0000FE0F" # variation selector (цветные эмодзи) "]+" ) def clean_short(text, limit: int | None = None) -> str: """Чистит текстовые поля (заголовок, суть) от markdown-разметки, markdown-ссылок ([текст](url) -> текст), голых URL и служебных символов. Используется и при сохранении, и при отдаче карточки, чтобы в UI не показывались сырые куски исходника с разметкой. Схлопывает в одну строку (для заголовка); для сути с сохранением переносов см. clean_block. """ return re.sub(r"\n+", " ", clean_block(text, limit)) def clean_block(text, limit: int | None = None) -> str: """Как clean_short, но сохраняет переносы строк: структурированная суть от ИИ (список параметров/пунктов) остаётся читаемой на карточке. """ s = str(text or "") s = _MD_LINK_RE.sub(lambda m: (m.group(1) or "").strip(), s) s = _MD_BOLD.sub(r"\1", s) # **жирный** s = _MD_BOLD2_RE.sub(r"\1", s) # __жирный__ s = _MD_CODE_RE.sub(r"\1", s) # `код` s = _MD_STRIKE_RE.sub(r"\1", s) # ~~зачёркнуто~~ s = s.replace("||", "") # ||спойлер|| s = _BARE_URL_RE.sub(" ", s) # голые ссылки # решётка в составе названия (C#, F#, .NET# нет, но бывают хэштеги) — # защищаем её от вырезания ниже, чтобы «C#» не превращалось в «C» s = re.sub(r"(?i)\b([a-zа-яё])\#", "\\1\u2063", s) s = s.replace("\u200b", "") # zero-width (фото-превью Telegram) s = s.replace("\u00a0", " ") s = re.sub(r"[*`#>~]+", " ", s) # служебные символы markdown s = s.replace("\u2063", "#") s = _EMOJI_RE.sub("", s) # декоративные эмодзи-маркеры # маркеры списков и «декоративные» буллеты в начале строк убираем s = re.sub(r"(?m)^[\s>#*\-–—•▪▫●○‣]+\s*", "", s) s = re.sub(r"[ \t]+", " ", s) s = re.sub(r"\n[ \t]+", "\n", s) s = re.sub(r"\n{2,}", "\n", s) s = s.strip(" \t\n\r-–—·•|:;,") if limit and len(s) > limit: # режем по границе последнего переноса/пробела до лимита cut = s[:limit] br = cut.rfind("\n") sp = cut.rfind(" ") at = br if br > limit // 2 else (sp if sp > limit // 2 else -1) if at != -1: cut = cut[:at] s = cut.rstrip() + "…" return s def _skip_no_budget(raw: dict, text: str) -> bool: """Глобальный фильтр «без указания суммы карточку не создаём». Включается отдельно для найма (вакансий) и для заказов (фриланс/услуги/ товары) — настройки budgetRequiredHire / budgetRequiredOrder. Суммой считаем распознанный бюджет или сумму в тексте (в любой валюте). Если требуется и суммы нет — сообщение просто пропускается (карточка не создаётся, dedup снимается, чтобы после выключения фильтра сообщение можно было обработать заново). """ try: req_hire = bool(store.get_setting("budgetRequiredHire")) req_order = bool(store.get_setting("budgetRequiredOrder")) except (TypeError, ValueError): return False if not (req_hire or req_order): return False is_hire = bool(raw.get("is_vacancy")) if (is_hire and not req_hire) or (not is_hire and not req_order): return False if ai_service.clean_budget(raw.get("budget")): return False return not bool(rules_svc.extract_amounts(text or "")) # Канонический текст блока «О заявке». Карточка всегда собирается из одних и # тех же блоков (единая структура, разная длина): Компания → Формат → О задаче # → Требования → Будет плюсом → Условия. Недостающие блоки пропускаются. def compose_summary(raw: dict, text: str = "") -> str: """Собирает «О заявке» карточки из структурированных полей разбора. raw — результат ИИ-классификатора или локального разбора (_local_fields). Если структурированных полей нет (старые/чужие ответы) — сохраняет исходную summary как есть, чтобы не потерять данные. """ def _line(value: str) -> str: return clean_short(value or "") def _items(key: str) -> list[str]: out: list[str] = [] for x in normalize_list(raw.get(key)): s = _line(x) if s: out.append(s) return out blocks: list[str] = [] company = _line(raw.get("company")) if company: blocks.append(f"Компания: {company}") fmt = _line(raw.get("format")) if fmt: blocks.append(f"Формат: {fmt}") task = _line(raw.get("task")) if task: blocks.append(f"О задаче: {task}") req = _items("requirements") if req: blocks.append("Требования: " + ", ".join(req[:14])) plus = _items("plus") if plus: blocks.append("Будет плюсом: " + ", ".join(plus[:10])) cond = _line(raw.get("conditions")) if cond: blocks.append(f"Условия: {cond}") if blocks: return "\n".join(blocks) # структурированных полей нет. «summary» от модели часто является копией # исходника/шумом — если в нём есть футеры/хэштеги-мусор, не используем его. legacy = _line(raw.get("summary")) if legacy and not any(h in legacy.casefold() for h in _FOOTER_HINTS): return legacy # локальный разбор без ИИ: «О задаче» из содержательных строк (без # хэштег-строк и служебных футеров агрегаторов) raw_lines = [ln.strip() for ln in (text or "").splitlines() if ln.strip()] keep: list[str] = [] for ln in raw_lines: low = ln.casefold() if ln.startswith("#") or low.startswith("**#") or any(h in low for h in _FOOTER_HINTS): continue cl = _clean_line(ln) if cl: keep.append(cl) if keep: short = _local_summary("\n".join(keep), keep) if short: return "О задаче: " + short return clean_short(text or "") # Футеры/служебные строки, по которым «суть» от модели — не структура, а шум _FOOTER_HINTS = ( "откликнуться через", "runello", "больше вакансий", "teletype", "при отклике укажите", "больше заявок", "узнать подробнее", "написать в лс", "пишите в лс", ) def _local_summary(body: str, lines: list[str], skip_first: int = 1) -> str: """Суть карточки при локальном разборе (без ИИ): не весь исходник, а очищенные содержательные строки после заголовка, без меток-полей и мусора. Схлопываем в короткий абзац — карточка не выглядит как сырое сообщение. """ parts: list[str] = [] for ln in lines[skip_first:]: hit = _field_of(ln) if hit: continue # «Стек: …», «Бюджет: …» и т.п. уже разобраны в поля cl = clean_short(ln) if len(cl) < 2 or cl.casefold() in ("вакансия", "вакансию", "фриланс"): continue parts.append(cl) if len(parts) >= 4: break out = " ".join(parts) if len(out) < 40: # мало содержательных строк — берём очищенное начало всего текста out = clean_short(body, 360) return out[:600] def normalize_list(value) -> list[str]: """ИИ иногда возвращает список, иногда строку («Java, Kotlin»). Строку разбиваем на элементы и чистим.""" if value is None: return [] parts = re.split(r"[;|\n]+", value) if isinstance(value, str) else value out: list[str] = [] for p in parts: s = str(p).strip().strip('*`#').strip() s = re.sub(r"\s+([.,])\s*$", r"\1", s).strip().strip(',').strip() if s and len(s) > 1 and s not in out: out.append(s) return out def normalize_stack(raw) -> list[str]: """Стек: чинит «побуквенный» разбор (когда ИИ вернул строку, а не список).""" out: list[str] = [] for s in normalize_list(raw): if len(s) < 2: continue # одиночные буквы/мусор — не технология out.append(s) if len(out) >= 12: break return out # Контакты: квалификация по типу. Храним список {type, value}. _TG_BOT_HINTS = ("bot",) _TG_SERVICE_NAMES = {"joinchat", "share", "s", "c", "addstickers", "addtheme", "proxy", "bg", "login"} _SKIP_SITE_HOSTS = {"teletype.in", "forms.gle", "docs.google.com", "youtube.com", "youtu.be", "clck.ru"} def qualify_contact(raw) -> dict | None: """Классифицирует один сырой контакт → {type, value} или None. type: tg | phone | email | linkedin | whatsapp | site. Отбрасываем ботов (@…bot), сервисные t.me-ссылки (joinchat/+/s/c…), «постовые» сайты (teletype, google-формы и т.п.). """ s = str(raw or "").strip() if not s or len(s) > 300: return None if s.startswith("@"): name = s[1:].strip() if re.fullmatch(r"[A-Za-z0-9_]{4,32}", name) and not name.lower().endswith("bot"): return {"type": "tg", "value": "@" + name} return None low = s.lower() m = re.match(r"https?://(?:www\.)?t\.me/([A-Za-z0-9_]{4,32})/?$", s) if m: name = m.group(1) if name.lower() not in _TG_SERVICE_NAMES and not name.lower().endswith("bot"): return {"type": "tg", "value": "@" + name} return None if re.fullmatch(r"[A-Za-z0-9._%+\-]+@[A-Za-z0-9.\-]+\.[A-Za-z]{2,}", s): return {"type": "email", "value": s.lower()} digits = re.sub(r"\D", "", s) if re.fullmatch(r"\+?[\d\s\-()]{6,20}", s) and 10 <= len(digits) <= 15: return {"type": "phone", "value": ("+" if s.startswith("+") else "") + digits} if re.search(r"linkedin\.com/in/", low): return {"type": "linkedin", "value": s} if re.search(r"wa\.me|api\.whatsapp\.com", low): return {"type": "whatsapp", "value": s} if low.startswith("http"): host = re.sub(r"https?://(?:www\.)?", "", low).split("/")[0].split("?")[0].split(":")[0] if host in _SKIP_SITE_HOSTS or host.endswith(".teletype.in"): return None return {"type": "site", "value": s} return None def build_contacts(raw, text: str = "") -> list[dict]: """Собирает квалифицированные контакты из ответа ИИ/локального разбора. Если контакты не нашлись, но есть текст — вытаскивает кандидатов из него (@username, e-mail, телефон). Возвращает до 6 записей {type, value}. """ cands: list[str] = [] if isinstance(raw, str): cands = re.split(r"[;|\n]+", raw) elif isinstance(raw, list): for x in raw: if isinstance(x, str): cands.extend(re.split(r"[;|\n]+", x)) elif isinstance(x, dict) and x.get("value"): cands.append(str(x["value"])) elif raw: cands = [str(raw)] if not cands and text: cands = _contacts_from(text) out: list[dict] = [] seen: set[str] = set() for c in cands: q = qualify_contact(c) if not q: continue key = q["value"].casefold() if key in seen: continue seen.add(key) out.append(q) if len(out) >= 6: break return out def primary_contact(contacts: list[dict]) -> str: """Основной контакт для быстрого действия: tg → телефон → почта → …""" order = {"tg": 0, "phone": 1, "whatsapp": 2, "email": 3, "linkedin": 4, "site": 5} if not contacts: return "" best = min(contacts, key=lambda c: order.get(str(c.get("type")), 9)) return str(best.get("value") or "") def _store_lead( digest: str, dialog_id: str, ch_name: str, ch_handle: str, ch_hue: str, text: str, raw: dict, msg_at: int, msg_id: int | None = None, ) -> dict | None: now = _now() lead_id = store.uid("l_") board_id = str(raw.get("board") or "").strip() or None # страховка: колонку с активными правилами может назначить только текст, # прошедший эти правила (ИИ/ML не должны класть в неё нерелевантное) if board_id and not rules_svc.board_accepts(board_id, text): board_id = None stack = normalize_stack(raw.get("stack")) budget = ai_service.clean_budget(raw.get("budget")) title = clean_short(raw.get("title") or "", 140) or clean_short(text, 140) # «О заявке» всегда собирается из одинаковых блоков (см. compose_summary): # Компания → Формат → О задаче → Требования → Будет плюсом → Условия. summary = clean_block(compose_summary(raw, text), 2000) or clean_short(text, 2000) if not budget: # ИИ не выделил бюджет отдельным полем, но сумма с валютой есть в исходнике # или в структурированной «О заявке» (часто уходит в «Условия») — показываем # её на карточке. Тот же источник, что и фильтр «не создавать без суммы». for src in (text, summary): amts = rules_svc.extract_amounts(src or "") if amts: a = amts[0] budget = {"from": a["from"], "to": a["to"], "currency": a["cur"]} break conv = ai_service.budget_to_target(budget) contacts = build_contacts(raw.get("contacts"), text) contact = primary_contact(contacts)[:200] # по каким критериям фильтра колонки карточка сюда попала (пусто — колонка без правил) match_hits = rules_svc.hits_for_board(board_id, text) if board_id else [] _cols = ( "id", "col", "is_new", "is_vacancy", "is_vacancy_known", "title", "summary", "stack", "budget_from", "budget_to", "budget_cur", "conv_from", "conv_to", "conv_cur", "contact", "contacts", "ch_name", "ch_handle", "ch_hue", "time_label", "received_at", "source_msg", "source_dialog_id", "source_msg_id", "prev_col", "comments", "match_hits", "created_at", ) known_type = bool(raw.get("is_vacancy_known")) store.execute( "INSERT INTO leads(" + ",".join(_cols) + ") VALUES (" + ",".join(["?"] * len(_cols)) + ")", [ lead_id, board_id or "inbox", True, bool(raw.get("is_vacancy")), known_type, title, summary, json.dumps(stack, ensure_ascii=False), budget.get("from") if budget else None, budget.get("to") if budget else None, budget.get("currency", "") if budget else "", conv["convFrom"], conv["convTo"], conv["convCur"], contact, json.dumps(contacts, ensure_ascii=False), ch_name, ch_handle, ch_hue, human_age(now, msg_at or now), msg_at or now, text[:4000], dialog_id or "", msg_id, "inbox", "[]", json.dumps(match_hits, ensure_ascii=False), now, ], ) store.execute("UPDATE dedup SET lead_id = ? WHERE hash = ?", [lead_id, digest]) _link_message(lead_id, dialog_id, msg_id, text, msg_at or now) return lead_to_dict(lead_id) def _link_message(lead_id: str, dialog_id: str, msg_id: int | None, text: str, msg_at: int) -> None: """В messages оседает только исходник карточки (для предпросмотра/флагов).""" if not dialog_id or msg_id is None: return store.execute( "INSERT OR IGNORE INTO messages(id, dialog_id, text, msg_at) VALUES (?, ?, ?, ?)", [f"m_{dialog_id}_{msg_id}", dialog_id, text[:4000], msg_at], ) store.execute("UPDATE messages SET lead_id = ? WHERE id = ?", [lead_id, f"m_{dialog_id}_{msg_id}"]) def human_age(now_ms: int, ts: int) -> str: delta = max(0, now_ms - ts) minutes = delta // C.MIN_MS if minutes < 60: return "только что" if minutes < 1 else f"{minutes} мин" hours = minutes // 60 if hours < 24: return f"{hours} ч" days = hours // 24 return f"{days} дн" def lead_to_dict(lead_id: str) -> dict: row = store.query_one("SELECT * FROM leads WHERE id = ?", [lead_id]) if not row: return {} comments = json.loads(row["comments"] or "[]") try: match_hits = json.loads(row.get("match_hits") or "[]") except Exception: # noqa: BLE001 match_hits = [] try: contacts = json.loads(row.get("contacts") or "[]") except Exception: # noqa: BLE001 contacts = [] if not contacts and (row.get("contact") or ""): q = qualify_contact(row["contact"]) contacts = [q] if q else [{"type": "other", "value": row["contact"]}] return { "id": row["id"], "col": row["col"], "isNew": bool(row["is_new"]), "isVacancy": bool(row["is_vacancy"]), "isVacancyKnown": bool(row.get("is_vacancy_known")), "title": clean_short(row["title"], 140), "summary": clean_block(row["summary"], 2000), "stack": json.loads(row["stack"] or "[]"), "budget": ( {"from": row["budget_from"], "to": row["budget_to"], "cur": row["budget_cur"]} if row["budget_cur"] else None ), "converted": ( {"from": row["conv_from"], "to": row["conv_to"], "cur": row["conv_cur"]} if row["conv_cur"] else None ), "contact": row["contact"], "contacts": contacts, "ch": {"name": row["ch_name"], "handle": row["ch_handle"], "hue": row["ch_hue"]}, "time": row["time_label"], "receivedAt": row["received_at"], "sourceMsg": row["source_msg"], "sourceDialogId": row["source_dialog_id"] or "", "sourceMsgId": row["source_msg_id"], "prevCol": row["prev_col"], "matchHits": match_hits, "comments": comments, } # Метки-поля, встречающиеся в объявлениях/вакансиях: «Стек: Java, Kotlin», # «Грейд: Middle», «Контакты: @user», «Бюджет: 1 200–1 500$» и т.п. _FIELD_LABELS = { "stack": {"стек", "технологии", "технология", "скиллы", "скилы", "языки", "язык", "инструменты", "tools", "tech stack", "stack"}, "grade": {"грейд", "уровень", "грейд/уровень", "seniority", "level"}, "contacts": {"контакт", "контакты", "связь", "телеграм", "почта", "email", "контакты для связи"}, "budget": {"бюджет", "оплата", "зп", "зарплата", "вилка", "оклад", "ставка", "цена", "цену", "гонорар", "pay", "salary"}, } _MD_EDGES = re.compile(r"^[\s*>#_~]+|[\s*>#_~]+$") _MD_BOLD = re.compile(r"\*\*(.+?)\*\*") _LABEL_RE = re.compile(r"^[\s*>#_~]*([А-Яа-яЁёA-Za-z][А-Яа-яЁёA-Za-z0-9 /+\-]{1,36}?)\s*[:|]\s*(.+)$") _TOKEN_RE = re.compile(r"(?:[A-Za-zА-Яа-яЁё0-9][A-Za-zА-Яа-яЁё0-9#.+\-]*|\.[A-Za-zА-Яа-яЁё][A-Za-zА-Яа-яЁё0-9#.+\-]*)") _CONTACT_RE = re.compile(r"@[A-Za-z0-9_]{3,}") _EMAIL_RE = re.compile(r"[A-Za-z0-9._%+\-]+@[A-Za-z0-9\-]+(?:\.[A-Za-z0-9\-]+)+") _PHONE_RE = re.compile(r"(?:\+7|8|7)[\s\-()]*\d{3}[\s\-()]*\d{3}[\s\-]*\d{2}[\s\-]*\d{2}") _STOP_STACK = { "и", "или", "на", "по", "с", "не", "а", "в", "о", "об", "от", "до", "для", "опыт", "знание", "знания", "уметь", "умение", "умения", "работать", "работы", "работа", "работе", "требуется", "приветствуется", "будет", "плюсом", "разработка", "разработке", "разработчик", "разработчика", "вакансия", "вакансию", "вакансии", "команда", "команду", "команды", "проект", "проекта", "проекты", "приветствуются", "желательно", "уверенное", "хорошее", "понимание", "навыки", "навык", "навыков", } # Маркеры «найма» и «уровней» не зашиты в код: они редактируются в UI # «Сфера и ключи» (hireMarkers / levelTerms). Здесь только чтение настроек # с фолбэком на дефолты из constants.py. def _hire_markers() -> set[str]: raw = store.get_setting("hireMarkers") if raw is None: raw = C.DEFAULT_HIRE_MARKERS if isinstance(raw, str): raw = [raw] return {str(x).strip().casefold() for x in raw if str(x).strip()} def _level_terms() -> set[str]: raw = store.get_setting("levelTerms") if raw is None: raw = C.DEFAULT_LEVEL_TERMS if isinstance(raw, str): raw = [raw] return {str(x).strip().casefold() for x in raw if str(x).strip()} def _resume_markers() -> set[str]: raw = store.get_setting("resumeMarkers") if raw is None: raw = C.DEFAULT_RESUME_MARKERS if isinstance(raw, str): raw = [raw] return {str(x).strip().casefold() for x in raw if str(x).strip()} def _resume_reason(lower: str) -> str | None: """Вернуть маркер, по которому текст опознан как резюме соискателя. Для слова «резюме» есть контекстный guard: объявление работодателя вида «…вакансия…, присылайте резюме» содержит маркер найма (hireMarkers) ДО слова «резюме» — это не резюме, а заявка, её не отсекаем. """ for marker in _resume_markers(): pos = lower.find(marker) if pos < 0: continue if marker == "резюме" and any(h in lower[:pos] for h in _hire_markers()): continue return marker return None def _norm_phone(p: str) -> str: p = p.replace(" ", "").replace("\u00a0", "").replace("-", "").replace("(", "").replace(")", "") return p[:18] def _contacts_from(body: str) -> list[str]: """Универсальные контакты из текста: @username, email, телефоны.""" out: list[str] = [] for m in _CONTACT_RE.findall(body): if m not in out: out.append(m) for m in _EMAIL_RE.findall(body): if m not in out: out.append(m) for m in _PHONE_RE.findall(body): p = _norm_phone(m) if p not in out: out.append(p) return out[:4] def _clean_line(ln: str) -> str: return _MD_EDGES.sub("", ln).strip() def _field_of(line: str) -> tuple[str, str] | None: """Если строка вида «Метка: значение» — вернуть (категория, значение).""" m = _LABEL_RE.match(line) if not m: return None label = m.group(1).strip().casefold() value = m.group(2).strip() if not value: return None for cat, syns in _FIELD_LABELS.items(): if label in syns: return cat, value return None def _pick_stack(value: str) -> list[str]: """Слова из значения метки («Стек:», «Услуги:», «Материалы:» …). Универсально: не только технологии с латиницей, а любые значимые слова — тип работ, услуги, товары, материалы (для не-IT сфер). """ out: list[str] = [] for tok in _TOKEN_RE.findall(value): low = tok.casefold() if low in _STOP_STACK or len(tok) < 2: continue out.append(tok) if len(out) >= 10: break return out def _local_fields(text: str) -> dict: """Локальный структуратор (без ИИ): заголовок, суть, стек, грейд, бюджет, контакты, признак вакансии. Используется для быстрых путей (правила/ML), чтобы карточка не выглядела как сырое сообщение. """ text = text or "" hire = _hire_markers() levels = _level_terms() lines = [_clean_line(ln) for ln in text.splitlines()] lines = [ln for ln in lines if ln] body = "\n".join(lines) lower = body.casefold() fields: dict[str, list[str]] = {"stack": [], "grade": [], "contacts": []} budget_raw = "" for ln in lines[1:]: hit = _field_of(ln) if not hit: continue cat, value = hit if cat == "stack": fields["stack"].extend(_pick_stack(value)) elif cat == "grade": for tok in _TOKEN_RE.findall(value): low = tok.casefold() if low in levels and low not in fields["grade"]: fields["grade"].append(low) elif cat == "contacts": fields["contacts"].extend(_contacts_from(value)) elif cat == "budget": budget_raw = value # fallback-извлечения по всему тексту (если объявление без меток) if not fields["contacts"]: fields["contacts"] = _contacts_from(body) if not fields["stack"]: # метка «Стек: …» может стоять не в начале строки (однострочные объявления) for m in re.finditer(r"(?im)\b(?:стек|технологии|скиллы|скилы|языки|язык|инструменты)\s*[:|]\s*([^\n]{2,120})", body): fields["stack"].extend(_pick_stack(m.group(1))) if fields["stack"]: break if not fields["grade"]: for w in lower.split(): w = w.strip(",;.:«»\"'()").casefold() if w in levels: fields["grade"].append(w) break if not budget_raw: budget_raw = body budget = None for amt in rules_svc.extract_amounts(budget_raw): budget = {"from": amt["from"], "to": amt["to"], "currency": amt["cur"]} break if budget is None and lines: for amt in rules_svc.extract_amounts(body): budget = {"from": amt["from"], "to": amt["to"], "currency": amt["cur"]} break title = clean_short(lines[0] if lines else body, 140) if not title: title = clean_short(body, 140) summary = _local_summary(body, lines) is_vacancy = any(mk in lower for mk in hire) contact = "; ".join(fields["contacts"])[:200] stack = [] for s in fields["stack"]: if s.casefold() not in (x.casefold() for x in stack): stack.append(s) return { "title": title, "summary": summary, "stack": stack[:12], "grade": fields["grade"][:4], "budget": budget, "contacts": contact, "is_vacancy": is_vacancy, "is_vacancy_known": False, # маркерная оценка — не контекст; тип подтверждает ИИ "board": None, } # ─── Фоновый воркер ─────────────────────────────────────────────────────── def _raw_spam(raw: dict) -> bool: """Признак спама в ответе ИИ-классификатора (используется, чтобы не учить ML на карточке, которую ИИ пометил мусором, но она всё же сохранена).""" return bool((raw or {}).get("is_spam")) def _drop_row(row: dict, with_dedup: bool = True) -> None: if with_dedup: digest = ai_service.normalize_dedup(row["text"]) store.execute("DELETE FROM dedup WHERE hash = ? AND lead_id IS NULL", [digest]) store.execute("DELETE FROM pipeline_msg WHERE id = ?", [row["id"]]) def _stale_max_age() -> int: """Срок актуальности входящего сообщения = срок до архива (в мс). Сообщение, опубликованное раньше этого срока, сразу ушло бы в архив — поэтому оно не заводится в систему вообще (ни карточкой, ни в архив/ корзину). Работает только при включённом автоархиве. """ if not store.get_setting("autoArchive"): return 0 days = int(store.get_setting("archiveAfterDays") or 14) return days * C.DAY_MS def _is_stale(msg_at: int | None) -> bool: max_age = _stale_max_age() if max_age <= 0 or not msg_at: return False return _now() - msg_at > max_age def _drop_stale(row: dict) -> None: """Устаревшее сообщение удаляем сразу, без создания карточки и без ML/ИИ.""" days = int(store.get_setting("archiveAfterDays") or 14) log.info("stage1 blocked (устарело, старше срока до архива): %s", row["text"][:80]) processing_svc.record( row, "stale", "stale", f"сообщение старше {days} дн. (срок до автоархива) — не заводим в систему", ) _drop_row(row) def _reject_stage1(row: dict, r1: dict) -> None: """Отсев на этапе 1: пишем причину в мониторинг (вкладка «Обработка»).""" kind = str(r1.get("kind") or "stop") processing_svc.record(row, "stop", kind, r1["reason"], str(r1.get("kw") or "")) def _reject_no_budget(row: dict, _is_hire: bool) -> None: """Глобальный фильтр «без суммы»: сообщение пропущено, карточка не создана.""" processing_svc.record( row, "stop", "budget", "включён фильтр «не создавать карточку без суммы» — в тексте не указан бюджет", ) # сигнатура последнего опубликованного состояния очереди/отсева (SSE-событие # шлём только при изменении — иначе воркер спамил бы эфир каждые 2 секунды) _last_pub_sig: tuple | None = None async def _maybe_publish_stats() -> None: """SSE pipeline_stats для вкладки «Обработка» при изменении очереди/отсева.""" global _last_pub_sig q = processing_svc.queue_counts() sig = (q["new"], q["ai"], processing_svc.rejected_count()) if sig == _last_pub_sig: return _last_pub_sig = sig await broker.publish( "pipeline_stats", {"new": sig[0], "ai": sig[1], "total": sig[0] + sig[1], "rejected": sig[2]}, ) def _pump_gate() -> tuple[int, int]: """Шлагбаум «остановиться после N карточек»: (limit, done). 0 = выключен.""" try: limit = int(store.get_setting("pumpGate") or 0) done = int(store.get_setting("pumpGateDone") or 0) except (TypeError, ValueError): limit, done = 0, 0 return max(0, limit), max(0, done) async def pump_once(new_limit: int = 12, ai_limit: int = 4) -> dict: """Один проход воркера по очереди. Вызывается из фонового цикла и /admin/tick. Если включён шлагбаум (pumpGate > 0), после N созданных карточек воркер останавливается: очередь копится, новые карточки не создаются до снятия шлагбаума (admin/pump-gate, limit=0 или новый лимит). Нужно для «прогона с остановкой после первых N сообщений» (проверка результата пользователем). """ limit, done = _pump_gate() if limit and done >= limit: return {"paused": True, "done": done, "limit": limit} if _pump_lock.locked(): return {} # другой воркер уже разбирает очередь async with _pump_lock: res = await _pump_unlocked(new_limit, ai_limit) if limit: added = int(res.get("rulesStored") or 0) + int(res.get("mlStored") or 0) + int(res.get("aiStored") or 0) done = done + added store.set_setting("pumpGateDone", done) res["gateDone"] = done res["gateLimit"] = limit if done >= limit: res["paused"] = True await broker.publish_toast( f"Пауза: создано {done} карточек (лимит {limit}) — проверьте результат", "clock" ) await _maybe_publish_stats() return res async def _pump_unlocked(new_limit: int, ai_limit: int) -> dict: res = {"staged": 0, "rulesStored": 0, "mlStored": 0, "mlDrop": 0, "typeDrop": 0, "aiStored": 0, "aiDrop": 0, "aiFail": 0, "noBudget": 0} # 1) new -> стоп-фразы/дедуп -> ML (или в очередь на ИИ) rows = store.query( "SELECT * FROM pipeline_msg WHERE status = ? ORDER BY created_at LIMIT ?", [ST_NEW, new_limit], ) for row in rows: force = bool(row.get("force")) if not force and _is_stale(row["msg_at"]): _drop_stale(row) continue if not force: r1 = stage1_plain(row["text"]) if not r1["pass"]: log.info("stage1 blocked (%s): %s", r1["reason"], row["text"][:80]) _reject_stage1(row, r1) _drop_row(row) continue digest = ai_service.normalize_dedup(row["text"]) if store.scalar("SELECT 1 FROM dedup WHERE hash = ?", [digest]): processing_svc.record( row, "dup", "dup", "сообщение уже в системе: карточка создана ранее или этот текст уже обрабатывается", ) _drop_row(row) continue store.execute( "INSERT OR IGNORE INTO dedup(hash, lead_id, created_at) VALUES (?, NULL, ?)", [digest, _now()], ) res["staged"] += 1 # Смысловые колонки НЕ назначаем словарным фильтром до ИИ: словарный # матч не понимает смысл (дайджест из 5 ролей, «desktop» в URL/футере # и т.п.). Сначала смысл оценивает ML (см. ниже), не уверен — ИИ. # Правила колонок остаются критериями для ИИ/проверкой при назначении # и «почему карточка попала» (см. _store_lead: board_accepts + hits). # ML-слой: если включён настройкой и уверен — решает без ИИ. # ML не назначает колонки с активными правилами и ИИ-предложения: # они наполняются ИИ (с проверкой правил) или действиями пользователя. # force (возврат из отсева) идёт мимо ML к ИИ-классификации: пользователь # явно подтвердил, что сообщение релевантно, — ML не должен снова его # отсеять как спам/не тот тип. if ml_client.is_enabled() and not force: dec = await ml_client.predict(row["text"]) label = dec.get("label") if dec.get("take") and label: if label == "spam": score = dec.get("scores", {}).get("spam", 0) log.info("ml drop (spam, score %.2f): %s", score, row["text"][:120].replace("\n", " ")) processing_svc.record( row, "ml", "spam_ml", f"ML уверен, что это спам/не заявка (score {float(score):.2f})", ) _drop_row(row) res["mlDrop"] += 1 continue br = store.query_one("SELECT suggested, rules FROM boards WHERE id = ?", [label]) allowed = False if br and not bool(br["suggested"]): try: br_rules = json.loads(br["rules"] or "{}") if br["rules"] else {} except Exception: # noqa: BLE001 br_rules = {} allowed = not rules_svc.has_active_rules(br_rules) if allowed: raw = _local_fields(row["text"]) raw["board"] = label # ML назначил колонку и уверен в типе — тип тоже его решение typ = dec.get("type") or {} if typ.get("take"): raw["is_vacancy"] = typ.get("label") == "hire" raw["is_vacancy_known"] = True # ML «узнала» в тексте слова, характерные для колонки: # докладываем их в стек (если локальный разбор их пропустил) added = 0 for term in (dec.get("terms") or []): t = str(term).strip().strip("@+#.") low = t.casefold() if not t or len(t) < 2 or t.startswith("~"): continue if low in _STOP_STACK: continue if any(x.casefold() == low for x in raw.get("stack") or []): continue raw.setdefault("stack", []).append(t) added += 1 if added >= 4: break if _skip_no_budget(raw, row["text"]): res["noBudget"] += 1 _reject_no_budget(row, bool(raw.get("is_vacancy"))) _drop_row(row) # без dedup: после выключения фильтра сообщение можно взять снова continue lead = _store_lead( digest, row["dialog_id"], row["ch_name"], row["ch_handle"], row["ch_hue"], row["text"], raw, row["msg_at"], row["msg_id"], ) _drop_row(row, with_dedup=False) # dedup уже привязан к lead if lead: res["mlStored"] += 1 await broker.publish("new_lead", lead) continue # ML уверен в типе (t:hire/t:order), даже если колонку не назначил typ = dec.get("type") or {} if typ.get("take"): is_hire = typ.get("label") == "hire" want = str(store.get_setting("wantedType") or "both").strip().lower() if want in ("vacancy", "freelance"): bad = (want == "freelance" and is_hire) or (want == "vacancy" and not is_hire) if bad: # тип не под режим «что собираем» — не тратим ИИ log.info("ml typeDrop (%s, хотят %s): %s", "найм" if is_hire else "разовое", want, row["text"][:120].replace("\n", " ")) processing_svc.record( row, "ml", "type", f"ML: тип «{'найм/занятость' if is_hire else 'разовые заказы'}», а вы ищете только «{want}»", ) _drop_row(row) res["typeDrop"] += 1 continue if not store.get_setting("aiEnabled"): # ИИ выключен, но тип ML знает: создаём карточку сами (inbox) raw = _local_fields(row["text"]) raw["is_vacancy"] = is_hire raw["is_vacancy_known"] = True if _skip_no_budget(raw, row["text"]): res["noBudget"] += 1 _reject_no_budget(row, is_hire) _drop_row(row) continue lead = _store_lead( digest, row["dialog_id"], row["ch_name"], row["ch_handle"], row["ch_hue"], row["text"], raw, row["msg_at"], row["msg_id"], ) _drop_row(row, with_dedup=False) if lead: res["mlStored"] += 1 await broker.publish("new_lead", lead) continue store.execute( "UPDATE pipeline_msg SET status = ?, updated_at = ? WHERE id = ?", [ST_AI, _now(), row["id"]], ) # 2) filtered -> ИИ-фильтр (если включён) + классификация rows = store.query( "SELECT * FROM pipeline_msg WHERE status = ? ORDER BY created_at LIMIT ?", [ST_AI, ai_limit], ) for row in rows: force = bool(row.get("force")) if not force and _is_stale(row["msg_at"]): _drop_stale(row) continue text = row["text"] digest = ai_service.normalize_dedup(text) # ИИ выключен (aiEnabled=false): классификацию/ИИ-фильтр не вызываем — # карточку собирает локальный разбор. В колонки кладёт только ML. if not store.get_setting("aiEnabled"): raw = _local_fields(text) if not force and _skip_no_budget(raw, text): res["noBudget"] += 1 _reject_no_budget(row, bool(raw.get("is_vacancy"))) _drop_row(row) continue lead = _store_lead( digest, row["dialog_id"], row["ch_name"], row["ch_handle"], row["ch_hue"], text, raw, row["msg_at"], row["msg_id"], ) _drop_row(row, with_dedup=False) if lead: res["aiStored"] += 1 await broker.publish("new_lead", lead) continue if force: # возврат из отсева: ИИ-фильтр не пересматриваем, пользователь уже # подтвердил релевантность — сразу классификация r2 = {"pass": True, "reason": None, "skipped": True} else: try: r2 = await ai_service.filter_incoming(text) except Exception as exc: # noqa: BLE001 log.warning("AI filter error, skip stage2: %s", exc) r2 = {"pass": True, "reason": None, "skipped": True} try: raw = await ai_service.classify(text) if r2["pass"] else {} if raw: # тип определён ИИ по контексту (не маркерной эвристикой) raw["is_vacancy_known"] = True except Exception as exc: # noqa: BLE001 log.warning("classification failed: %s", exc) raw = {} is_spam = bool(raw.get("is_spam") or (not r2["pass"])) if force and raw.get("is_spam"): # пользователь вернул сообщение из отсева: вердикт «спам» отменяется, # поля ИИ уже разложил — карточку создаём (без обучения «спаму») raw["is_spam"] = False is_spam = False if is_spam: # ИИ-решение «спам/мусор» — обучаем ML отличать спам (с весом гипотезы) log.info("ai drop (не заявка: %s): %s", r2.get("reason") or raw.get("is_spam"), text[:120].replace("\n", " ")) if not r2["pass"]: processing_svc.record( row, "ai", "filter_ai", "ИИ-фильтр: " + (str(r2.get("reason") or "сообщение не относится к вашим интересам")), ) else: processing_svc.record( row, "ai", "spam_ai", "ИИ: не заявка — спам, реклама, скам или служебное сообщение", ) ml_client.push(text, "spam", delta=ml_client.AI_WEIGHT) _drop_row(row) res["aiDrop"] += 1 continue if not raw: # ИИ не дал разбора (недоступен/сбой) — но не спам: локальный разбор raw = _local_fields(text) res["aiFail"] += 1 if not force and _skip_no_budget(raw, text): res["noBudget"] += 1 _reject_no_budget(row, bool(raw.get("is_vacancy"))) _drop_row(row) continue lead = _store_lead( digest, row["dialog_id"], row["ch_name"], row["ch_handle"], row["ch_hue"], text, raw, row["msg_at"], row["msg_id"], ) _drop_row(row, with_dedup=False) if lead: res["aiStored"] += 1 await broker.publish("new_lead", lead) # ИИ назначил колонку — это обучающий сигнал для ML (кроме inbox: # «не знаю» не учим). Карточка не в inbox -> учим текст -> колонка. col = str(lead.get("col") or "") if col and col not in ("inbox", "trash", "archive") and not _raw_spam(raw): # ML в своём пути назначает только колонки без активных правил # и не-suggested — такие же примеры и собираем, иначе модель # будет «знать» колонку, но не сможет её применить. br = store.query_one("SELECT suggested, rules FROM boards WHERE id = ?", [col]) free = False if br and not bool(br["suggested"]): try: br_rules = json.loads(br["rules"] or "{}") if br["rules"] else {} except Exception: # noqa: BLE001 br_rules = {} free = not rules_svc.has_active_rules(br_rules) if free: ml_client.push(text, col, delta=ml_client.AI_WEIGHT) # тип известен от ИИ по контексту — учим ML определять его сам # (t:hire = занятость/найм, t:order = разовая сделка) if raw.get("is_vacancy_known"): ml_client.push( text, "t:hire" if bool(raw.get("is_vacancy")) else "t:order", delta=ml_client.AI_WEIGHT, ) ml_client.track_decisions(ml=res["mlStored"] + res["mlDrop"], ai=res["aiStored"] + res["aiDrop"]) return res