mirror of
https://gitverse.ru/kpa39l/vesti.git
synced 2026-09-29 09:55:03 +00:00
48c25f13d4
- parse_address(address, crawler) в sources/sources.py: telegram (t.me/name, @name, https://t.me/s/name) → slug/channel/url; rss → feed_url/url/slug из домена. - GET /sources/preview: telegram — название через Telethon get_entity (title), rss — через feedparser (feed.title); slug и name возвращаются в JSON. - _form_source: единое поле address вместо url/channel/feed_url; slug из адреса. - /sources: direction datalist из DIRECTIONS + уникальные значения из БД, свободный ввод остаётся. - Форма: поле Address + кнопка «Найти» (JS fetch → заполнение slug/name), direction datalist динамический. Также в sources.py входит незакоммиченная ранее веб-CRUD-логика (upsert_source_yaml, delete_source_yaml, sync_yaml_to_db_and_back и др.) — работала на проде, теперь в git.
308 lines
14 KiB
Python
308 lines
14 KiB
Python
# VESTI — работа с реестром источников (sources.yaml) и таблицей sources/runs.
|
|
import datetime
|
|
import re
|
|
import sqlite3
|
|
import sys
|
|
from pathlib import Path
|
|
|
|
import yaml
|
|
|
|
sys.path.insert(0, str(Path(__file__).resolve().parent))
|
|
from config import SOURCES_PATH
|
|
from db.db import db
|
|
|
|
|
|
def load_sources_yaml(path=None):
|
|
"""Читает sources.yaml -> list[dict]. Ошибки структуры -> исключение."""
|
|
path = Path(path or SOURCES_PATH)
|
|
with open(path, "r", encoding="utf-8") as f:
|
|
data = yaml.safe_load(f)
|
|
if not isinstance(data, list):
|
|
raise ValueError(f"{path}: ожидается список источников (YAML list)")
|
|
return data
|
|
|
|
|
|
def sync_sources_to_db(data=None):
|
|
"""Загружает sources.yaml в БД (UPSERT по slug). Возвращает {slug: source_id}."""
|
|
data = data if data is not None else load_sources_yaml()
|
|
mapping = {}
|
|
with db() as conn:
|
|
for s in data:
|
|
slug = s.get("slug")
|
|
if not slug:
|
|
continue
|
|
cur = conn.execute(
|
|
"""INSERT INTO sources (slug, name, url, channel, crawler, feed_url, direction, lang, priority, enabled, status, own)
|
|
VALUES (?,?,?,?,?,?,?,?,?,?, 'new', ?)
|
|
ON CONFLICT(slug) DO UPDATE SET
|
|
name=excluded.name, url=excluded.url, channel=excluded.channel,
|
|
crawler=excluded.crawler, feed_url=excluded.feed_url,
|
|
direction=excluded.direction,
|
|
lang=excluded.lang, priority=excluded.priority, enabled=excluded.enabled,
|
|
own=excluded.own""",
|
|
(slug, s.get("name"), s.get("url"), s.get("channel"),
|
|
s.get("crawler", "telegram"), s.get("feed_url"),
|
|
s.get("direction"), s.get("lang", "ru"),
|
|
s.get("priority", "P1"), int(s.get("enabled", True)),
|
|
int(s.get("own", False) or 0)),
|
|
)
|
|
mapping[slug] = cur.lastrowid
|
|
return mapping
|
|
|
|
|
|
def get_enabled_sources(crawler=None, direction=None):
|
|
"""Список включённых источников (для cron). Фильтры: crawler, direction."""
|
|
with db() as conn:
|
|
q = "SELECT * FROM sources WHERE enabled=1"
|
|
args = []
|
|
if crawler:
|
|
q += " AND crawler=?"
|
|
args.append(crawler)
|
|
if direction:
|
|
q += " AND direction=?"
|
|
args.append(direction)
|
|
q += " ORDER BY priority, slug"
|
|
return [dict(r) for r in conn.execute(q, args).fetchall()]
|
|
|
|
|
|
def get_source(slug):
|
|
with db() as conn:
|
|
r = conn.execute("SELECT * FROM sources WHERE slug=?", (slug,)).fetchone()
|
|
return dict(r) if r else None
|
|
|
|
|
|
def start_run(source_id, task, trigger="cron"):
|
|
"""Начало запуска задачи для конкретного источника. Возвращает run_id."""
|
|
with db() as conn:
|
|
cur = conn.execute(
|
|
"INSERT INTO runs (source_id, task, trigger, status) VALUES (?,?,?, 'running')",
|
|
(source_id, task, trigger),
|
|
)
|
|
return cur.lastrowid
|
|
|
|
|
|
def finish_run(run_id, status="ok", posts_fetched=0, posts_new=0, error=None, duration_ms=None):
|
|
"""Завершение запуска."""
|
|
with db() as conn:
|
|
if duration_ms is None:
|
|
conn.execute(
|
|
"UPDATE runs SET status=?, posts_fetched=?, posts_new=?, error=?, finished_at=datetime('now'),"
|
|
" duration_ms=(julianday('now')-julianday(started_at))*86400000 WHERE id=?",
|
|
(status, posts_fetched, posts_new, error, run_id),
|
|
)
|
|
else:
|
|
conn.execute(
|
|
"UPDATE runs SET status=?, posts_fetched=?, posts_new=?, error=?, finished_at=datetime('now'), duration_ms=? WHERE id=?",
|
|
(status, posts_fetched, posts_new, error, duration_ms, run_id),
|
|
)
|
|
|
|
|
|
def recent_runs(limit=50, source_id=None):
|
|
"""Последние запуски (для веб-интерфейса)."""
|
|
with db() as conn:
|
|
if source_id:
|
|
rows = conn.execute(
|
|
"SELECT r.*, s.slug, s.name as source_name FROM runs r LEFT JOIN sources s ON s.id=r.source_id"
|
|
" WHERE r.source_id=? ORDER BY r.id DESC LIMIT ?", (source_id, limit),
|
|
).fetchall()
|
|
else:
|
|
rows = conn.execute(
|
|
"SELECT r.*, s.slug, s.name as source_name FROM runs r LEFT JOIN sources s ON s.id=r.source_id"
|
|
" ORDER BY r.id DESC LIMIT ?", (limit,),
|
|
).fetchall()
|
|
return [dict(r) for r in rows]
|
|
|
|
|
|
def now_iso():
|
|
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# CRUD источников (веб-страница /sources): синхронизация БД + sources.yaml.
|
|
# sources.yaml — источник истины (краулеры при запуске делают sync_sources_to_db
|
|
# из yaml), поэтому КАЖДОЕ изменение пишется и в yaml, и в БД.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
def _write_sources_yaml(data: list, path=None):
|
|
"""Атомарная запись списка источников в YAML (temp + os.replace)."""
|
|
path = Path(path or SOURCES_PATH)
|
|
# сортировка: сначала включённые, потом по (priority, slug) — стабильный, читаемый порядок
|
|
p_order = {"P0": 0, "P1": 1, "P2": 2, "P3": 3}
|
|
data = sorted(data, key=lambda s: (0 if s.get("enabled", True) else 1,
|
|
p_order.get(str(s.get("priority", "P1")), 9),
|
|
str(s.get("slug", ""))))
|
|
tmp = path.with_suffix(".yaml.tmp")
|
|
with open(tmp, "w", encoding="utf-8") as f:
|
|
yaml.safe_dump(data, f, allow_unicode=True, sort_keys=False)
|
|
f.write("\n# Пример (включить, раскомментировав и поправив direction/приоритет):\n"
|
|
"# - slug: opennet\n"
|
|
"# name: \"OpenNET\"\n"
|
|
"# url: https://www.opennet.ru/\n"
|
|
"# feed_url: https://www.opennet.ru/opennews/opennews_all.rss\n"
|
|
"# channel: null\n"
|
|
"# crawler: rss\n"
|
|
"# direction: linux\n"
|
|
"# lang: ru\n"
|
|
"# priority: P1\n"
|
|
"# enabled: true\n")
|
|
tmp.replace(path)
|
|
|
|
|
|
def upsert_source_yaml(source: dict, path=None):
|
|
"""Добавляет/обновляет источник в sources.yaml (по slug)."""
|
|
path = Path(path or SOURCES_PATH)
|
|
data = load_sources_yaml(path) if path.exists() else []
|
|
idx = next((i for i, s in enumerate(data) if s.get("slug") == source.get("slug")), None)
|
|
if idx is not None:
|
|
merged = dict(data[idx]); merged.update({k: v for k, v in source.items() if v is not None or k in ("enabled",)})
|
|
data[idx] = merged
|
|
else:
|
|
data.append(source)
|
|
_write_sources_yaml(data, path)
|
|
|
|
|
|
def delete_source_yaml(slug: str, path=None):
|
|
"""Удаляет источник из sources.yaml (по slug). Возвращает True, если был."""
|
|
path = Path(path or SOURCES_PATH)
|
|
data = load_sources_yaml(path) if path.exists() else []
|
|
out = [s for s in data if s.get("slug") != slug]
|
|
if len(out) == len(data):
|
|
return False
|
|
_write_sources_yaml(out, path)
|
|
return True
|
|
|
|
|
|
def db_source_to_yaml(source: dict) -> dict:
|
|
"""Преобразует строку БД sources в запись sources.yaml (только нужные поля)."""
|
|
return {
|
|
"slug": source["slug"],
|
|
"name": source.get("name"),
|
|
"url": source.get("url"),
|
|
"channel": source.get("channel"),
|
|
"crawler": source.get("crawler", "telegram"),
|
|
"feed_url": source.get("feed_url"),
|
|
"direction": source.get("direction"),
|
|
"lang": source.get("lang", "ru"),
|
|
"priority": source.get("priority", "P1"),
|
|
"enabled": bool(source.get("enabled", 1)),
|
|
"own": bool(source.get("own", 0) or 0),
|
|
}
|
|
|
|
def sync_yaml_to_db_and_back(path=None):
|
|
"""Синхронизирует yaml → БД (как краулеры) — используется после правок yaml извне.
|
|
|
|
Возвращает (mapping, count). Просто вызывает sync_sources_to_db(load_sources_yaml()).
|
|
"""
|
|
data = load_sources_yaml(path)
|
|
return sync_sources_to_db(data), len(data)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Веб-CRUD: изменения пишутся в yaml → синкаются в БД (yaml — источник истины).
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_SLUG_RE = re.compile(r"^[a-z0-9_-]+$")
|
|
|
|
_TG_HANDLE_RE = re.compile(r"^(?:https?://)?(?:www\.)?(?:t\.me|telegram\.me)/(?:s/)?([a-zA-Z][a-zA-Z0-9_]{3,31})$")
|
|
_RSS_URL_RE = re.compile(r"^https?://", re.I)
|
|
|
|
|
|
def parse_address(address: str, crawler: str) -> dict:
|
|
"""Разбирает одно поле Address в словарь полей источника.
|
|
|
|
telegram: address = @handle | t.me/handle | https://t.me/s/handle | handle
|
|
→ {slug: handle, channel: handle, url: https://t.me/handle}
|
|
rss: address = https://...feed.xml | http://...
|
|
→ {feed_url: address, url: address, slug: <домен первого сегмента>}
|
|
Возвращает dict с заполненными полями (пустые отсутствуют). Пустой dict —
|
|
адрес не разобран (невалидный для выбранного crawler).
|
|
"""
|
|
addr = (address or "").strip()
|
|
if not addr:
|
|
return {}
|
|
if crawler == "telegram":
|
|
m = _TG_HANDLE_RE.match(addr)
|
|
if not m:
|
|
# пробуем без схемы/домена: "handle", "@handle"
|
|
m = _TG_HANDLE_RE.match("https://t.me/" + addr.lstrip("@/"))
|
|
if not m:
|
|
return {}
|
|
handle = m.group(1)
|
|
return {"slug": handle, "channel": handle, "url": f"https://t.me/{handle}"}
|
|
# rss
|
|
if not _RSS_URL_RE.match(addr):
|
|
return {}
|
|
from urllib.parse import urlparse
|
|
host = urlparse(addr).netloc.replace("www.", "").split(".")[0] or "feed"
|
|
return {"feed_url": addr, "url": addr, "slug": host}
|
|
|
|
|
|
def _validate_source(data: dict) -> str | None:
|
|
"""Валидация полей источника. Возвращает текст ошибки или None."""
|
|
slug = str(data.get("slug") or "").strip()
|
|
if not slug or not _SLUG_RE.match(slug):
|
|
return "slug: только [a-z0-9_-]"
|
|
if data.get("crawler", "telegram") not in ("telegram", "rss"):
|
|
return 'crawler: только "telegram" или "rss"'
|
|
if str(data.get("priority", "P1")) not in ("P0", "P1", "P2", "P3"):
|
|
return "priority: только P0..P3"
|
|
if data.get("crawler") == "rss" and not str(data.get("feed_url") or "").strip():
|
|
return "rss-источнику нужен feed_url"
|
|
if data.get("crawler") == "telegram" and not str(data.get("channel") or "").strip():
|
|
return "telegram-источнику нужен channel"
|
|
return None
|
|
|
|
|
|
def add_or_update_source(data: dict) -> tuple[bool, str]:
|
|
"""Добавить/обновить источник: валидация → yaml → БД.
|
|
|
|
Возвращает (ok, ошибка_или_slug).
|
|
"""
|
|
err = _validate_source(data)
|
|
if err:
|
|
return False, err
|
|
# нормализация
|
|
data = {
|
|
"slug": str(data["slug"]).strip(),
|
|
"name": (data.get("name") or "").strip() or None,
|
|
"url": (data.get("url") or "").strip() or None,
|
|
"channel": (data.get("channel") or "").strip() or None,
|
|
"crawler": data.get("crawler", "telegram"),
|
|
"feed_url": (data.get("feed_url") or "").strip() or None,
|
|
"direction": (data.get("direction") or "").strip() or None,
|
|
"lang": (data.get("lang") or "ru").strip() or "ru",
|
|
"priority": str(data.get("priority", "P1")).upper(),
|
|
"enabled": bool(data.get("enabled", True)),
|
|
"own": bool(data.get("own", False)),
|
|
}
|
|
upsert_source_yaml(data)
|
|
sync_yaml_to_db_and_back()
|
|
return True, data["slug"]
|
|
|
|
|
|
def delete_source(slug: str) -> bool:
|
|
"""Удаляет источник: yaml → БД (каскад: посты/раны остаются с source_id=NULL)."""
|
|
ok = delete_source_yaml(slug)
|
|
if not ok:
|
|
return False
|
|
with db() as conn:
|
|
conn.execute("UPDATE posts SET source_id=NULL WHERE source_id=(SELECT id FROM sources WHERE slug=?)", (slug,))
|
|
conn.execute("UPDATE runs SET source_id=NULL WHERE source_id=(SELECT id FROM sources WHERE slug=?)", (slug,))
|
|
conn.execute("DELETE FROM rss_state WHERE slug=?", (slug,))
|
|
conn.execute("DELETE FROM sources WHERE slug=?", (slug,))
|
|
return True
|
|
|
|
|
|
def set_source_enabled(slug: str, enabled: bool):
|
|
"""Пауза/снятие: enabled в yaml + БД."""
|
|
with db() as conn:
|
|
conn.row_factory = sqlite3.Row
|
|
row = conn.execute("SELECT * FROM sources WHERE slug=?", (slug,)).fetchone()
|
|
if not row:
|
|
return False
|
|
# обновить yaml (сохранить остальные поля из БД)
|
|
y = db_source_to_yaml(dict(row))
|
|
y["enabled"] = bool(enabled)
|
|
upsert_source_yaml(y)
|
|
sync_yaml_to_db_and_back()
|
|
return True |