# 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