# VESTI — работа с реестром источников (sources.yaml) и таблицей sources/runs. import datetime 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, 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, 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("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")