Нормализовать переводы строк в LF
Решение по TD-STYLE-ANALYZERS: LF — инструменты проекта (Python/Node) пишут LF, CRLF-.sh не работают на Linux CI (sh scripts/ci.sh), большинство файлов уже были LF. Добавлен .gitattributes (* text=auto eol=lf, бинарные исключения), .editorconfig переведён на lf, 1029 файлов конвертированы, git add --renormalize. Из индекса убраны закравшиеся archive/**/__pycache__/*.pyc.
This commit is contained in:
@@ -1,320 +1,320 @@
|
||||
"""Мониторинг пайплайна — вкладка «Обработка» (очередь и отсев).
|
||||
|
||||
Очередь — сырые сообщения из каналов, ждущие разбора (pipeline_msg).
|
||||
Отсев — сообщения, отброшенные на любом этапе: стоп-фразы/резюме/тип заявки/
|
||||
без суммы (source='stop'), устарело ('stale'), ML ('ml'), ИИ ('ai'), повтор
|
||||
('dup'). Для каждой записи храним этап, причину и конкретное слово/фразу
|
||||
(kw), если отсев по стоп-списку.
|
||||
|
||||
Автоочистка отсева — раз в 3 суток (вызывается из leads.tick_storage),
|
||||
плюс ручная очистка и удаление отдельных записей из UI.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
|
||||
from ..db import store
|
||||
from . import fts as fts_svc
|
||||
|
||||
log = logging.getLogger("leadradar.processing")
|
||||
|
||||
# отсев живёт 3 суток, дальше удаляется автоматически
|
||||
RETENTION_DAYS = 3
|
||||
|
||||
# человекочитаемые подписи этапов (для UI; source хранится отдельно)
|
||||
_STAGE_LABELS = {
|
||||
"length": "короткое сообщение",
|
||||
"stop": "стоп-фраза",
|
||||
"resume": "резюме соискателя",
|
||||
"type": "тип заявки",
|
||||
"budget": "нет суммы",
|
||||
"stale": "устарело",
|
||||
"spam_ml": "спам (ML)",
|
||||
"spam_ai": "спам (ИИ)",
|
||||
"filter_ai": "ИИ-фильтр",
|
||||
"dup": "повтор",
|
||||
}
|
||||
|
||||
# «чьё» решение: используется в UI как источник метки
|
||||
_SOURCE_LABELS = {
|
||||
"stop": "правила",
|
||||
"ml": "ML",
|
||||
"ai": "ИИ",
|
||||
"stale": "система",
|
||||
"dup": "система",
|
||||
}
|
||||
|
||||
DEFAULT_LIMIT = 100
|
||||
MAX_LIMIT = 500
|
||||
|
||||
|
||||
def _now() -> int:
|
||||
return time.time_ns() // 1_000_000
|
||||
|
||||
|
||||
def stage_label(stage: str) -> str:
|
||||
return _STAGE_LABELS.get(stage, stage or "отсев")
|
||||
|
||||
|
||||
def source_label(source: str) -> str:
|
||||
return _SOURCE_LABELS.get(source, source or "система")
|
||||
|
||||
|
||||
# ─── Запись отсева ────────────────────────────────────────────────────────
|
||||
|
||||
def record(row: dict, source: str, stage: str, reason: str, kw: str = "") -> None:
|
||||
"""Сохраняет отброшенное сообщение в таблицу отсева.
|
||||
|
||||
Ид записи детерминирован по (dialog_id, msg_id): повторное отбрасывание
|
||||
того же сообщения (перечитывание каналов) обновляет запись, а не копит
|
||||
дубликаты в списке отсева.
|
||||
"""
|
||||
if not row or not (row.get("text") or "").strip():
|
||||
return
|
||||
msg_id = row.get("msg_id")
|
||||
dialog_id = row.get("dialog_id") or ""
|
||||
rid = f"r_{dialog_id}_{msg_id}" if (msg_id is not None and dialog_id) else store.uid("r_")
|
||||
now = _now()
|
||||
store.execute(
|
||||
"INSERT INTO rejected_msgs(id, dialog_id, msg_id, text, ch_name, ch_handle, ch_hue, stage, reason, kw, source, msg_at, rejected_at) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) "
|
||||
"ON CONFLICT(id) DO UPDATE SET "
|
||||
"text = excluded.text, ch_name = excluded.ch_name, ch_handle = excluded.ch_handle, "
|
||||
"ch_hue = excluded.ch_hue, stage = excluded.stage, reason = excluded.reason, kw = excluded.kw, "
|
||||
"source = excluded.source, msg_at = excluded.msg_at, rejected_at = excluded.rejected_at",
|
||||
[
|
||||
rid,
|
||||
dialog_id,
|
||||
msg_id,
|
||||
str(row["text"])[:6000],
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
stage,
|
||||
str(reason or "")[:500],
|
||||
str(kw or "")[:200],
|
||||
source,
|
||||
row.get("msg_at"),
|
||||
now,
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
def purge_expired(days: int = RETENTION_DAYS) -> int:
|
||||
"""Автоочистка отсева: записи старше N суток удаляются безвозвратно."""
|
||||
cut = _now() - days * 24 * 3600 * 1000
|
||||
rows = store.query(
|
||||
"SELECT id FROM rejected_msgs WHERE rejected_at < ?", [cut]
|
||||
)
|
||||
if not rows:
|
||||
return 0
|
||||
ids = [r["id"] for r in rows]
|
||||
store.execute(
|
||||
"DELETE FROM rejected_msgs WHERE id IN (" + ",".join(["?"] * len(ids)) + ")",
|
||||
ids,
|
||||
)
|
||||
return len(ids)
|
||||
|
||||
|
||||
def clear_all() -> int:
|
||||
rows = store.query("SELECT count(*) AS c FROM rejected_msgs")
|
||||
total = int(rows[0]["c"]) if rows else 0
|
||||
if total:
|
||||
store.execute("DELETE FROM rejected_msgs")
|
||||
return total
|
||||
|
||||
|
||||
def return_to_queue(rej_id: str, reason: str = "") -> dict:
|
||||
"""Вернуть отсеянное сообщение в обработку (кнопка в «Обработке»).
|
||||
|
||||
Строка очереди помечается force: этап 1, устарело, ML-решения и ИИ-отсев
|
||||
для неё игнорируются — сообщение уходит на классификацию и создаёт карточку.
|
||||
Запись в отсеве не удаляется, а помечается «возвращено» с причиной (аудит).
|
||||
Если отсев был по решению «спам» (ML/ИИ) — снимаем у ML вес спама для текста.
|
||||
"""
|
||||
row = store.query_one("SELECT * FROM rejected_msgs WHERE id = ?", [rej_id])
|
||||
if not row:
|
||||
raise KeyError(rej_id)
|
||||
if bool(row.get("returned")):
|
||||
raise ValueError("Сообщение уже возвращено в обработку")
|
||||
if str(row.get("source") or "") == "dup":
|
||||
raise ValueError("Повтор: карточка с таким текстом уже есть в системе — возвращать нечего")
|
||||
dialog_id = str(row.get("dialog_id") or "")
|
||||
msg_id = row.get("msg_id")
|
||||
text = str(row.get("text") or "").strip()
|
||||
if not text:
|
||||
raise ValueError("В записи нет текста сообщения")
|
||||
|
||||
if str(row.get("stage") or "") in ("spam_ml", "spam_ai", "filter_ai"):
|
||||
from . import ml_client
|
||||
|
||||
# реальное действие пользователя: этот текст НЕ спам
|
||||
ml_client.push(text, "spam", delta=-1.0)
|
||||
|
||||
now = _now()
|
||||
store.execute(
|
||||
"UPDATE rejected_msgs SET returned = TRUE, returned_at = ?, return_reason = ? WHERE id = ?",
|
||||
[now, str(reason or "").strip()[:500], rej_id],
|
||||
)
|
||||
|
||||
from .pipeline import enqueue # локальный импорт: pipeline импортирует processing
|
||||
|
||||
if dialog_id and msg_id is not None:
|
||||
enqueue(
|
||||
dialog_id,
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
msg_id,
|
||||
text,
|
||||
row.get("msg_at") or now,
|
||||
force=True,
|
||||
)
|
||||
else:
|
||||
# старые записи (до сохранения dialog_id/msg_id): текст сохранился,
|
||||
# возвращаем без ссылки на исходное сообщение (force=True)
|
||||
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 (?, ?, ?, ?, ?, ?, ?, ?, 'new', TRUE, ?, ?)",
|
||||
[
|
||||
store.uid("p_"),
|
||||
dialog_id,
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
text[:6000],
|
||||
msg_id,
|
||||
row.get("msg_at") or now,
|
||||
now,
|
||||
now,
|
||||
],
|
||||
)
|
||||
return {"id": rej_id, "returned": True, "returnedAt": now}
|
||||
|
||||
|
||||
def delete_one(rej_id: str) -> bool:
|
||||
store.execute("DELETE FROM rejected_msgs WHERE id = ?", [rej_id])
|
||||
return True
|
||||
|
||||
|
||||
def rejected_count() -> int:
|
||||
return int(store.scalar("SELECT count(*) FROM rejected_msgs") or 0)
|
||||
|
||||
|
||||
# ─── Очередь (pipeline_msg) ───────────────────────────────────────────────
|
||||
|
||||
def queue_counts() -> dict:
|
||||
out = {"new": 0, "ai": 0}
|
||||
for r in store.query("SELECT status, count(*) AS c FROM pipeline_msg GROUP BY status"):
|
||||
if r["status"] == "new":
|
||||
out["new"] = int(r["c"])
|
||||
elif r["status"] == "filtered":
|
||||
out["ai"] = int(r["c"])
|
||||
out["total"] = out["new"] + out["ai"]
|
||||
return out
|
||||
|
||||
|
||||
def list_queue(limit: int = DEFAULT_LIMIT) -> list[dict]:
|
||||
limit = min(max(1, limit), MAX_LIMIT)
|
||||
rows = store.query(
|
||||
"SELECT * FROM pipeline_msg ORDER BY created_at LIMIT ?", [limit]
|
||||
)
|
||||
items = []
|
||||
for r in rows:
|
||||
items.append(
|
||||
{
|
||||
"id": r["id"],
|
||||
"dialogId": r.get("dialog_id") or "",
|
||||
"msgId": r.get("msg_id"),
|
||||
"text": r["text"],
|
||||
"status": r["status"], # new | filtered
|
||||
"ch": {
|
||||
"name": r["ch_name"],
|
||||
"handle": r["ch_handle"],
|
||||
"hue": r["ch_hue"],
|
||||
},
|
||||
"msgAt": r["msg_at"],
|
||||
"queuedAt": r["created_at"],
|
||||
}
|
||||
)
|
||||
return items
|
||||
|
||||
|
||||
# ─── Отсев (rejected_msgs) ────────────────────────────────────────────────
|
||||
|
||||
def list_rejected(q: str = "", offset: int = 0, limit: int = DEFAULT_LIMIT) -> dict:
|
||||
offset = max(0, offset)
|
||||
limit = min(max(1, limit), MAX_LIMIT)
|
||||
qq = (q or "").strip().lower()
|
||||
ids: list[str] = []
|
||||
|
||||
if qq:
|
||||
# FTS-кандидаты + LIKE-дополнение (свежие записи после последнего rebuild)
|
||||
if fts_svc.is_ready():
|
||||
try:
|
||||
ids = fts_svc.search(qq, limit=limit)["rejected"]
|
||||
except Exception: # noqa: BLE001
|
||||
ids = []
|
||||
pattern = f"%{qq}%"
|
||||
like = store.query(
|
||||
"SELECT id FROM rejected_msgs WHERE "
|
||||
"lower(text) LIKE ? OR lower(reason) LIKE ? OR lower(kw) LIKE ? OR lower(ch_name) LIKE ? "
|
||||
"ORDER BY rejected_at DESC LIMIT ?",
|
||||
[pattern, pattern, pattern, pattern, limit * 2],
|
||||
)
|
||||
for r in like:
|
||||
if r["id"] not in ids:
|
||||
ids.append(r["id"])
|
||||
total = len(ids) # итог по условию поиска (все кандидаты)
|
||||
page = ids[offset : offset + limit]
|
||||
else:
|
||||
total = rejected_count()
|
||||
page_rows = store.query(
|
||||
"SELECT id FROM rejected_msgs ORDER BY rejected_at DESC LIMIT ? OFFSET ?",
|
||||
[limit, offset],
|
||||
)
|
||||
page = [r["id"] for r in page_rows]
|
||||
|
||||
by_id: dict[str, dict] = {}
|
||||
if page:
|
||||
ph = ",".join(["?"] * len(page))
|
||||
rows = store.query(
|
||||
f"SELECT * FROM rejected_msgs WHERE id IN ({ph})", page
|
||||
)
|
||||
for r in rows:
|
||||
by_id[r["id"]] = r
|
||||
items = []
|
||||
for rid in page:
|
||||
r = by_id.get(rid)
|
||||
if not r:
|
||||
continue
|
||||
items.append(
|
||||
{
|
||||
"id": r["id"],
|
||||
"dialogId": r.get("dialog_id") or "",
|
||||
"msgId": r.get("msg_id"),
|
||||
"text": r["text"],
|
||||
"stage": r["stage"],
|
||||
"stageLabel": stage_label(r["stage"]),
|
||||
"reason": r["reason"],
|
||||
"kw": r["kw"],
|
||||
"source": r["source"],
|
||||
"sourceLabel": source_label(r["source"]),
|
||||
"ch": {"name": r["ch_name"], "handle": r["ch_handle"], "hue": r.get("ch_hue") or "#666"},
|
||||
"msgAt": r["msg_at"],
|
||||
"rejectedAt": r["rejected_at"],
|
||||
"returned": bool(r.get("returned")),
|
||||
"returnedAt": r.get("returned_at"),
|
||||
"returnReason": r.get("return_reason") or "",
|
||||
}
|
||||
)
|
||||
return {"items": items, "total": total, "offset": offset, "limit": limit}
|
||||
|
||||
|
||||
def stats() -> dict:
|
||||
q = queue_counts()
|
||||
return {
|
||||
"queue": q,
|
||||
"rejected": rejected_count(),
|
||||
}
|
||||
"""Мониторинг пайплайна — вкладка «Обработка» (очередь и отсев).
|
||||
|
||||
Очередь — сырые сообщения из каналов, ждущие разбора (pipeline_msg).
|
||||
Отсев — сообщения, отброшенные на любом этапе: стоп-фразы/резюме/тип заявки/
|
||||
без суммы (source='stop'), устарело ('stale'), ML ('ml'), ИИ ('ai'), повтор
|
||||
('dup'). Для каждой записи храним этап, причину и конкретное слово/фразу
|
||||
(kw), если отсев по стоп-списку.
|
||||
|
||||
Автоочистка отсева — раз в 3 суток (вызывается из leads.tick_storage),
|
||||
плюс ручная очистка и удаление отдельных записей из UI.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
|
||||
from ..db import store
|
||||
from . import fts as fts_svc
|
||||
|
||||
log = logging.getLogger("leadradar.processing")
|
||||
|
||||
# отсев живёт 3 суток, дальше удаляется автоматически
|
||||
RETENTION_DAYS = 3
|
||||
|
||||
# человекочитаемые подписи этапов (для UI; source хранится отдельно)
|
||||
_STAGE_LABELS = {
|
||||
"length": "короткое сообщение",
|
||||
"stop": "стоп-фраза",
|
||||
"resume": "резюме соискателя",
|
||||
"type": "тип заявки",
|
||||
"budget": "нет суммы",
|
||||
"stale": "устарело",
|
||||
"spam_ml": "спам (ML)",
|
||||
"spam_ai": "спам (ИИ)",
|
||||
"filter_ai": "ИИ-фильтр",
|
||||
"dup": "повтор",
|
||||
}
|
||||
|
||||
# «чьё» решение: используется в UI как источник метки
|
||||
_SOURCE_LABELS = {
|
||||
"stop": "правила",
|
||||
"ml": "ML",
|
||||
"ai": "ИИ",
|
||||
"stale": "система",
|
||||
"dup": "система",
|
||||
}
|
||||
|
||||
DEFAULT_LIMIT = 100
|
||||
MAX_LIMIT = 500
|
||||
|
||||
|
||||
def _now() -> int:
|
||||
return time.time_ns() // 1_000_000
|
||||
|
||||
|
||||
def stage_label(stage: str) -> str:
|
||||
return _STAGE_LABELS.get(stage, stage or "отсев")
|
||||
|
||||
|
||||
def source_label(source: str) -> str:
|
||||
return _SOURCE_LABELS.get(source, source or "система")
|
||||
|
||||
|
||||
# ─── Запись отсева ────────────────────────────────────────────────────────
|
||||
|
||||
def record(row: dict, source: str, stage: str, reason: str, kw: str = "") -> None:
|
||||
"""Сохраняет отброшенное сообщение в таблицу отсева.
|
||||
|
||||
Ид записи детерминирован по (dialog_id, msg_id): повторное отбрасывание
|
||||
того же сообщения (перечитывание каналов) обновляет запись, а не копит
|
||||
дубликаты в списке отсева.
|
||||
"""
|
||||
if not row or not (row.get("text") or "").strip():
|
||||
return
|
||||
msg_id = row.get("msg_id")
|
||||
dialog_id = row.get("dialog_id") or ""
|
||||
rid = f"r_{dialog_id}_{msg_id}" if (msg_id is not None and dialog_id) else store.uid("r_")
|
||||
now = _now()
|
||||
store.execute(
|
||||
"INSERT INTO rejected_msgs(id, dialog_id, msg_id, text, ch_name, ch_handle, ch_hue, stage, reason, kw, source, msg_at, rejected_at) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) "
|
||||
"ON CONFLICT(id) DO UPDATE SET "
|
||||
"text = excluded.text, ch_name = excluded.ch_name, ch_handle = excluded.ch_handle, "
|
||||
"ch_hue = excluded.ch_hue, stage = excluded.stage, reason = excluded.reason, kw = excluded.kw, "
|
||||
"source = excluded.source, msg_at = excluded.msg_at, rejected_at = excluded.rejected_at",
|
||||
[
|
||||
rid,
|
||||
dialog_id,
|
||||
msg_id,
|
||||
str(row["text"])[:6000],
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
stage,
|
||||
str(reason or "")[:500],
|
||||
str(kw or "")[:200],
|
||||
source,
|
||||
row.get("msg_at"),
|
||||
now,
|
||||
],
|
||||
)
|
||||
|
||||
|
||||
def purge_expired(days: int = RETENTION_DAYS) -> int:
|
||||
"""Автоочистка отсева: записи старше N суток удаляются безвозвратно."""
|
||||
cut = _now() - days * 24 * 3600 * 1000
|
||||
rows = store.query(
|
||||
"SELECT id FROM rejected_msgs WHERE rejected_at < ?", [cut]
|
||||
)
|
||||
if not rows:
|
||||
return 0
|
||||
ids = [r["id"] for r in rows]
|
||||
store.execute(
|
||||
"DELETE FROM rejected_msgs WHERE id IN (" + ",".join(["?"] * len(ids)) + ")",
|
||||
ids,
|
||||
)
|
||||
return len(ids)
|
||||
|
||||
|
||||
def clear_all() -> int:
|
||||
rows = store.query("SELECT count(*) AS c FROM rejected_msgs")
|
||||
total = int(rows[0]["c"]) if rows else 0
|
||||
if total:
|
||||
store.execute("DELETE FROM rejected_msgs")
|
||||
return total
|
||||
|
||||
|
||||
def return_to_queue(rej_id: str, reason: str = "") -> dict:
|
||||
"""Вернуть отсеянное сообщение в обработку (кнопка в «Обработке»).
|
||||
|
||||
Строка очереди помечается force: этап 1, устарело, ML-решения и ИИ-отсев
|
||||
для неё игнорируются — сообщение уходит на классификацию и создаёт карточку.
|
||||
Запись в отсеве не удаляется, а помечается «возвращено» с причиной (аудит).
|
||||
Если отсев был по решению «спам» (ML/ИИ) — снимаем у ML вес спама для текста.
|
||||
"""
|
||||
row = store.query_one("SELECT * FROM rejected_msgs WHERE id = ?", [rej_id])
|
||||
if not row:
|
||||
raise KeyError(rej_id)
|
||||
if bool(row.get("returned")):
|
||||
raise ValueError("Сообщение уже возвращено в обработку")
|
||||
if str(row.get("source") or "") == "dup":
|
||||
raise ValueError("Повтор: карточка с таким текстом уже есть в системе — возвращать нечего")
|
||||
dialog_id = str(row.get("dialog_id") or "")
|
||||
msg_id = row.get("msg_id")
|
||||
text = str(row.get("text") or "").strip()
|
||||
if not text:
|
||||
raise ValueError("В записи нет текста сообщения")
|
||||
|
||||
if str(row.get("stage") or "") in ("spam_ml", "spam_ai", "filter_ai"):
|
||||
from . import ml_client
|
||||
|
||||
# реальное действие пользователя: этот текст НЕ спам
|
||||
ml_client.push(text, "spam", delta=-1.0)
|
||||
|
||||
now = _now()
|
||||
store.execute(
|
||||
"UPDATE rejected_msgs SET returned = TRUE, returned_at = ?, return_reason = ? WHERE id = ?",
|
||||
[now, str(reason or "").strip()[:500], rej_id],
|
||||
)
|
||||
|
||||
from .pipeline import enqueue # локальный импорт: pipeline импортирует processing
|
||||
|
||||
if dialog_id and msg_id is not None:
|
||||
enqueue(
|
||||
dialog_id,
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
msg_id,
|
||||
text,
|
||||
row.get("msg_at") or now,
|
||||
force=True,
|
||||
)
|
||||
else:
|
||||
# старые записи (до сохранения dialog_id/msg_id): текст сохранился,
|
||||
# возвращаем без ссылки на исходное сообщение (force=True)
|
||||
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 (?, ?, ?, ?, ?, ?, ?, ?, 'new', TRUE, ?, ?)",
|
||||
[
|
||||
store.uid("p_"),
|
||||
dialog_id,
|
||||
str(row.get("ch_name") or ""),
|
||||
str(row.get("ch_handle") or ""),
|
||||
str(row.get("ch_hue") or "#666"),
|
||||
text[:6000],
|
||||
msg_id,
|
||||
row.get("msg_at") or now,
|
||||
now,
|
||||
now,
|
||||
],
|
||||
)
|
||||
return {"id": rej_id, "returned": True, "returnedAt": now}
|
||||
|
||||
|
||||
def delete_one(rej_id: str) -> bool:
|
||||
store.execute("DELETE FROM rejected_msgs WHERE id = ?", [rej_id])
|
||||
return True
|
||||
|
||||
|
||||
def rejected_count() -> int:
|
||||
return int(store.scalar("SELECT count(*) FROM rejected_msgs") or 0)
|
||||
|
||||
|
||||
# ─── Очередь (pipeline_msg) ───────────────────────────────────────────────
|
||||
|
||||
def queue_counts() -> dict:
|
||||
out = {"new": 0, "ai": 0}
|
||||
for r in store.query("SELECT status, count(*) AS c FROM pipeline_msg GROUP BY status"):
|
||||
if r["status"] == "new":
|
||||
out["new"] = int(r["c"])
|
||||
elif r["status"] == "filtered":
|
||||
out["ai"] = int(r["c"])
|
||||
out["total"] = out["new"] + out["ai"]
|
||||
return out
|
||||
|
||||
|
||||
def list_queue(limit: int = DEFAULT_LIMIT) -> list[dict]:
|
||||
limit = min(max(1, limit), MAX_LIMIT)
|
||||
rows = store.query(
|
||||
"SELECT * FROM pipeline_msg ORDER BY created_at LIMIT ?", [limit]
|
||||
)
|
||||
items = []
|
||||
for r in rows:
|
||||
items.append(
|
||||
{
|
||||
"id": r["id"],
|
||||
"dialogId": r.get("dialog_id") or "",
|
||||
"msgId": r.get("msg_id"),
|
||||
"text": r["text"],
|
||||
"status": r["status"], # new | filtered
|
||||
"ch": {
|
||||
"name": r["ch_name"],
|
||||
"handle": r["ch_handle"],
|
||||
"hue": r["ch_hue"],
|
||||
},
|
||||
"msgAt": r["msg_at"],
|
||||
"queuedAt": r["created_at"],
|
||||
}
|
||||
)
|
||||
return items
|
||||
|
||||
|
||||
# ─── Отсев (rejected_msgs) ────────────────────────────────────────────────
|
||||
|
||||
def list_rejected(q: str = "", offset: int = 0, limit: int = DEFAULT_LIMIT) -> dict:
|
||||
offset = max(0, offset)
|
||||
limit = min(max(1, limit), MAX_LIMIT)
|
||||
qq = (q or "").strip().lower()
|
||||
ids: list[str] = []
|
||||
|
||||
if qq:
|
||||
# FTS-кандидаты + LIKE-дополнение (свежие записи после последнего rebuild)
|
||||
if fts_svc.is_ready():
|
||||
try:
|
||||
ids = fts_svc.search(qq, limit=limit)["rejected"]
|
||||
except Exception: # noqa: BLE001
|
||||
ids = []
|
||||
pattern = f"%{qq}%"
|
||||
like = store.query(
|
||||
"SELECT id FROM rejected_msgs WHERE "
|
||||
"lower(text) LIKE ? OR lower(reason) LIKE ? OR lower(kw) LIKE ? OR lower(ch_name) LIKE ? "
|
||||
"ORDER BY rejected_at DESC LIMIT ?",
|
||||
[pattern, pattern, pattern, pattern, limit * 2],
|
||||
)
|
||||
for r in like:
|
||||
if r["id"] not in ids:
|
||||
ids.append(r["id"])
|
||||
total = len(ids) # итог по условию поиска (все кандидаты)
|
||||
page = ids[offset : offset + limit]
|
||||
else:
|
||||
total = rejected_count()
|
||||
page_rows = store.query(
|
||||
"SELECT id FROM rejected_msgs ORDER BY rejected_at DESC LIMIT ? OFFSET ?",
|
||||
[limit, offset],
|
||||
)
|
||||
page = [r["id"] for r in page_rows]
|
||||
|
||||
by_id: dict[str, dict] = {}
|
||||
if page:
|
||||
ph = ",".join(["?"] * len(page))
|
||||
rows = store.query(
|
||||
f"SELECT * FROM rejected_msgs WHERE id IN ({ph})", page
|
||||
)
|
||||
for r in rows:
|
||||
by_id[r["id"]] = r
|
||||
items = []
|
||||
for rid in page:
|
||||
r = by_id.get(rid)
|
||||
if not r:
|
||||
continue
|
||||
items.append(
|
||||
{
|
||||
"id": r["id"],
|
||||
"dialogId": r.get("dialog_id") or "",
|
||||
"msgId": r.get("msg_id"),
|
||||
"text": r["text"],
|
||||
"stage": r["stage"],
|
||||
"stageLabel": stage_label(r["stage"]),
|
||||
"reason": r["reason"],
|
||||
"kw": r["kw"],
|
||||
"source": r["source"],
|
||||
"sourceLabel": source_label(r["source"]),
|
||||
"ch": {"name": r["ch_name"], "handle": r["ch_handle"], "hue": r.get("ch_hue") or "#666"},
|
||||
"msgAt": r["msg_at"],
|
||||
"rejectedAt": r["rejected_at"],
|
||||
"returned": bool(r.get("returned")),
|
||||
"returnedAt": r.get("returned_at"),
|
||||
"returnReason": r.get("return_reason") or "",
|
||||
}
|
||||
)
|
||||
return {"items": items, "total": total, "offset": offset, "limit": limit}
|
||||
|
||||
|
||||
def stats() -> dict:
|
||||
q = queue_counts()
|
||||
return {
|
||||
"queue": q,
|
||||
"rejected": rejected_count(),
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user