Files
vesti/sources/sources.py
T
kpa39l 48c25f13d4 feat(sources): форма добавления источника — одно поле Address, авто-название из TG/RSS (address-field-enhancement)
- 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.
2026-09-16 17:03:23 +00:00

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