Первый коммит: модульный монолит ядра (.NET 10) и gRPC-сервисы ai/ml/telegram, фронтенд Vue 3/Vite/Tailwind, документация (ТЗ, инструкция пользователя, техдокументация, код-стайл), бэклог, скрипты развёртывания и архив прототипа LeadRadar.
1184 lines
58 KiB
Python
1184 lines
58 KiB
Python
"""Пайплайн входящих (п.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
|