diff --git a/.env.example b/.env.example index f7df931..bbb429f 100644 --- a/.env.example +++ b/.env.example @@ -26,4 +26,8 @@ ADMIN_USER=admin ADMIN_PASSWORD=change_me # Ollama OLLAMA_URL=http://127.0.0.1:11434 -OLLAMA_MODEL=qwen3:8b-nothink \ No newline at end of file +OLLAMA_MODEL=qwen3:8b-nothink +# RSS-краулер VESTI (честный User-Agent, обязателен для публичных лент) +RSS_USER_AGENT=vesti-rss/0.1 (+https://vesti.nixg.ru) +# Таймаут HTTP-запроса к ленте, сек +RSS_DEFAULT_TIMEOUT=20 \ No newline at end of file diff --git a/AGENT.MD b/AGENT.MD index ffc6a49..c25238c 100644 --- a/AGENT.MD +++ b/AGENT.MD @@ -1,6 +1,7 @@ # AGENT.MD — Правила проекта VESTI ## Правило №1: ВСЕ изменения — через OpenSpec +- **Перед каждым планируемым изменением — обязательно grill-with-docs** (интервью по дизайн-дереву: глоссарий в CONTEXT.md + ADR при необходимости). Только после достижения общего понимания с пользователем — создавать change. - **Любое изменение проекта (багфикс, фича, рефакторинг, конфиг, деплой, docs) оформляется как отдельный change** в `/opt/vesti/openspec/changes//`. - **НЕ начинать правки кода/файлов, пока change не создан и не провалидирован (`openspec validate` чисто).** Это обязательное требование (у VESTI нет бэкапа — OpenSpec фиксирует состояние и откат). - **Сразу после правок — обновить STATUS.md/TODO.md и сделать бэкап** (см. «Бэкап» ниже). @@ -16,6 +17,7 @@ ## Стек и архитектура - Новостной агрегатор: TG-краулер → классификатор (Ollama qwen3:8b-nothink) → веб-модерация → публикация в Telegram. +- Веб: :8400 — страницы /candidates (модерация), /crawlers (статус краулеров, запуск воркера), /sources (управление источниками: add/edit/toggle-pause/delete; источник истины — sources.yaml). - `vesti-web` (FastAPI+Jinja2, :8400, systemd-юнит) --HTTP POST--> `publisher-service` (FastAPI, Docker-контейнер, :8410) --Bot API (SOCKS5 127.0.0.1:1080)--> Telegram @dedinit_vesti. - БД: SQLite `/opt/vesti/db/vesti.db` — открывать с WAL + busy_timeout (иначе `database is locked` при фоновом классификаторе). - Фронтенд: **без внешних CDN** (никаких unpkg/jsdelivr). Bootstrap локален: `web/static/bootstrap.min.css`, htmx запрещён. approve/reject — обычные POST-формы. @@ -43,6 +45,11 @@ sudo systemctl restart vesti-web # веб :8400 (0.0.0.0, uvicorn web.a docker compose -f services/publisher/docker-compose.yml up -d --build # publisher :8410 curl -s http://127.0.0.1:8410/healthz # ok, bot=dedinit_controller_bot, proxy, channels .venv/bin/python -m crawler.telegram_crawler --all # краулер (инкрементально) +.venv/bin/python -m crawler.rss_crawler --all # RSS/Atom-краулер (условные GET, etag/modified; rss_state) +.venv/bin/python -m crawler.rss_crawler --source lwn --dry-run # сухой прогон без записи в БД +.venv/bin/python -m crawler.crawl_sources --all # фаза 1: сбор кандидатов в очередь (telegram+rss) +.venv/bin/python -m crawler.worker --workers 8 --limit 200 # фаза 3: параллельная обработка очереди (ThreadPoolExecutor) +MAX_CONCURRENT_LLM=2 .venv/bin/python -m crawler.worker --workers 8 # лимит одновременных LLM-запросов к Ollama CLASSIFY_TIMEOUT=20 .venv/bin/python -m classifier.classify --db db/vesti.db --limit 200 openspec validate # валидация change ``` diff --git a/config.py b/config.py index 78e82c9..a8d3be3 100644 --- a/config.py +++ b/config.py @@ -26,6 +26,10 @@ ADMIN_PASSWORD = os.getenv("ADMIN_PASSWORD", "change_me") OLLAMA_URL = os.getenv("OLLAMA_URL", "http://127.0.0.1:11434") OLLAMA_MODEL = os.getenv("OLLAMA_MODEL", "qwen3:8b-nothink") +# RSS-краулер: честный User-Agent (некоторые издатели режут дефолтные), таймаут +RSS_USER_AGENT = os.getenv("RSS_USER_AGENT", "vesti-rss/0.1 (+https://vesti.nixg.ru)") +RSS_DEFAULT_TIMEOUT = float(os.getenv("RSS_DEFAULT_TIMEOUT", "20")) + # Направления (категории) DIRECTIONS = ["tech", "politics", "games", "electronics", "llm", "linux"] diff --git a/crawler/crawl_sources.py b/crawler/crawl_sources.py new file mode 100644 index 0000000..ac3c3d6 --- /dev/null +++ b/crawler/crawl_sources.py @@ -0,0 +1,43 @@ +# VESTI — фаза 1: сбор кандидатов с источников в очередь (posts). +# Запускает telegram_crawler и rss_crawler по включённым источникам; +# классификация НЕ выполняется — новые посты уходят в очередь (classified=0), +# которую разбирает crawler.worker (фаза 3, параллельно). +# +# Запуск: python -m crawler.crawl_sources --all +# python -m crawler.crawl_sources --crawler rss # только RSS +import argparse +import subprocess +import sys +from pathlib import Path + +BASE = Path(__file__).resolve().parent.parent +PY = str(BASE / ".venv" / "bin" / "python") + + +def run_crawler(mod: str, extra: list[str] | None = None) -> int: + """Запускает модуль-краулер в подпроцессе. Возвращает exit code.""" + cmd = [PY, "-m", mod] + (extra or []) + print(f"--- {mod} {' '.join(extra or [])} ---", flush=True) + return subprocess.run(cmd, cwd=str(BASE)).returncode + + +def main(): + ap = argparse.ArgumentParser(description="VESTI фаза 1: сбор кандидатов с источников в очередь") + ap.add_argument("--all", action="store_true", help="сбор со всех включённых краулеров") + ap.add_argument("--crawler", choices=["telegram", "rss"], help="только один тип краулера") + args = ap.parse_args() + + if not args.all and not args.crawler: + ap.print_help() + sys.exit(2) + + rc = 0 + if args.all or args.crawler == "telegram": + rc |= run_crawler("crawler.telegram_crawler", ["--all"]) + if args.all or args.crawler == "rss": + rc |= run_crawler("crawler.rss_crawler", ["--all"]) + sys.exit(rc) + + +if __name__ == "__main__": + main() \ No newline at end of file diff --git a/crawler/worker.py b/crawler/worker.py new file mode 100644 index 0000000..94551bc --- /dev/null +++ b/crawler/worker.py @@ -0,0 +1,170 @@ +# VESTI — воркер (фаза 3): параллельная обработка очереди кандидатов. +# ThreadPoolExecutor(N): каждый воркер атомарно забирает пост из очереди +# (posts WHERE classified=0), классифицирует (keywords → Ollama), пишет +# direction/relevance/interest/summary, ставит classified=1. При сбое — возвращает +# в очередь (classified=0). Очередь — сама таблица posts, отдельная таблица не нужна. +# +# Запуск: python -m crawler.worker --workers 8 --limit 200 +import argparse +import os +import sqlite3 +import sys +import threading +import time +from concurrent.futures import ThreadPoolExecutor, as_completed +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) + +from config import DB_PATH, ensure_dirs +from classifier.classify import classify_text, call_ollama # noqa: F401 (переиспользуем) + +# Лимит одновременных LLM-запросов: локальная Ollama (qwen3:8b на CPU/GPU) держит +# один слот — параллельные 8 запросов = каждый ждёт >CLASSIFY_TIMEOUT и падает. +# 8 воркеров = 8 потоков, но к LLM пускаем максимум N (обычно 2) одновременных. +LLM_SEM = threading.BoundedSemaphore(int(os.getenv("MAX_CONCURRENT_LLM", "2"))) + + +def claim_post(conn: sqlite3.Connection) -> sqlite3.Row | None: + """Атомарный claim поста из очереди. + + BEGIN IMMEDIATE блокирует других воркеров на время claim; SELECT→UPDATE + гарантирует, что строку заберёт ровно один воркер (classified: 0 → -1). + Возвращает строку поста или None (очередь пуста). + """ + conn.execute("BEGIN IMMEDIATE") + try: + row = conn.execute( + "SELECT p.id, p.text, p.views, p.reactions_total, p.is_own, s.slug AS source_slug " + "FROM posts p LEFT JOIN sources s ON s.id = p.source_id " + "WHERE p.classified IS NULL OR p.classified=0 " + "ORDER BY p.is_own DESC, p.id LIMIT 1" + ).fetchone() + if row: + conn.execute("UPDATE posts SET classified=-1 WHERE id=?", (row["id"],)) + conn.commit() + return row + except Exception: + conn.rollback() + raise + + +def process_post(conn: sqlite3.Connection, post_id: int, text: str, + views: int, reactions: int, is_own: int, source_slug: str | None = None) -> dict: + """Классифицирует пост и пишет результат. Возвращает result classify_text. + + При исключении — пост возвращается в очередь (classified=0), НЕ теряется. + Мультинаправления из словаря (find_directions_all) — в classifications. + source_slug (правило SOURCE_RULES, напр. lwn→tech) — направление принудительное. + """ + res = classify_text({"text": text, "views": views, + "reactions_total": reactions, "is_own": is_own, + "source_slug": source_slug}) + conn.execute( + "UPDATE posts SET direction=?, relevance=?, interest=?, summary=?, classified=? WHERE id=?", + (res["direction"], res["relevance"], res["interest"], res["summary"], + 1 if res["classified"] else 0, post_id), + ) + # мультинаправления: все словарные попадания → classifications (для fan-out) + from classifier.keywords import find_directions_all + dirs_all = find_directions_all(text or "", source_slug) + if res["direction"] and res["direction"] not in dirs_all: + dirs_all.append(res["direction"]) + for d in dirs_all: + exists = conn.execute( + "SELECT 1 FROM classifications WHERE post_id=? AND direction=? LIMIT 1", + (post_id, d), + ).fetchone() + if not exists: + conn.execute( + """INSERT INTO classifications (post_id, direction, relevance, interest, summary, model) + VALUES (?,?,?,?,?,?)""", + (post_id, d, res["relevance"], res["interest"], res["summary"], + res.get("method") or "classify"), + ) + conn.commit() + return res + + +def work_loop(worker_id: int, stats: dict, stop: threading.Event, limit: int): + """Цикл воркера: claim → классификация → запись, пока очередь не пуста/лимит. + + Соединение создаётся ВНУТРИ потока (SQLite: соединение привязано к thread). + """ + conn = _new_conn() + processed = 0 + try: + while not stop.is_set() and processed < limit: + row = claim_post(conn) + if row is None: + break + try: + # LLM-вызовы ограничены семифором (не перегружаем Ollama) + with LLM_SEM: + res = process_post(conn, row["id"], row["text"], row["views"], + row["reactions_total"], row["is_own"], + row["source_slug"]) + with stats["lock"]: + stats["processed"] += 1 + if res["classified"]: + stats["classified"] += 1 + if res.get("method") == "llm": + stats["llm_ok"] += 1 + except Exception as e: + # вернуть в очередь — пост не теряется; зафиксировать счётчик ошибок + try: + conn.execute("UPDATE posts SET classified=0 WHERE id=?", (row["id"],)) + conn.commit() + except Exception: + pass + with stats["lock"]: + stats["errors"] += 1 + print(f"[worker {worker_id}] error post {row['id']}: {e}", flush=True) + processed += 1 + finally: + conn.close() + return processed + + +def _new_conn() -> sqlite3.Connection: + """Отдельное соединение для воркера (WAL — читатели/писатели параллельны).""" + ensure_dirs() + conn = sqlite3.connect(DB_PATH, timeout=30) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + conn.execute("PRAGMA busy_timeout=10000") + return conn + + +def run(workers: int = 8, limit: int = 200): + """Параллельная обработка очереди кандидатов. Возвращает (processed, classified, llm_ok, errors).""" + stats = {"processed": 0, "classified": 0, "llm_ok": 0, "errors": 0, "lock": threading.Lock()} + stop = threading.Event() + + # каждый воркер сам создаёт соединение внутри своего потока + print(f"workers={workers}, limit={limit}") + with ThreadPoolExecutor(max_workers=workers) as ex: + futs = {ex.submit(work_loop, i, stats, stop, limit): i + for i in range(workers)} + try: + for f in as_completed(futs): + f.result() + except KeyboardInterrupt: + stop.set() + print("\nостановлено (Ctrl+C)") + return stats["processed"], stats["classified"], stats["llm_ok"], stats["errors"] + + +def main(): + ap = argparse.ArgumentParser(description="VESTI worker — параллельная обработка очереди кандидатов") + ap.add_argument("--workers", type=int, default=8) + ap.add_argument("--limit", type=int, default=200) + args = ap.parse_args() + + t0 = time.monotonic() + proc, cls, llm, err = run(workers=args.workers, limit=args.limit) + print(f"Обработано: {proc}, классифицировано: {cls}, через LLM: {llm}, ошибок: {err}, время: {time.monotonic()-t0:.1f}s", flush=True) + + +if __name__ == "__main__": + main() \ No newline at end of file diff --git a/web/templates/crawlers.html b/web/templates/crawlers.html new file mode 100644 index 0000000..84d59ea --- /dev/null +++ b/web/templates/crawlers.html @@ -0,0 +1,98 @@ +{% extends "base.html" %} +{% block title %}Краулеры{% endblock %} +{% block content %} +

🤖 Краулеры

+ +
+
+
+
+
Очередь (pending)
+

{{ pending }}

+
+
+
+
+
+
+
В обработке (processing)
+

{{ processing }}

+
+
+
+
+
+
+
Обработано сегодня
+

{{ done_today }}

+
+
+
+
+
+
+
Источники
+

{{ alive }}/{{ total_sources}} + {% if dead %} {{ dead }} dead{% endif %}

+
+
+
+
+ +
+
+ + + +
+
+ +
+
+ +{% if request.query_params.get('run') == 'started' %} +
Воркер запущен ({{ request.query_params.get('workers') }} потоков, лимит {{ request.query_params.get('limit') }}). Следите за processing → done.
+{% elif request.query_params.get('run') == 'error' %} +
Ошибка запуска: {{ request.query_params.get('msg') }}
+{% elif request.query_params.get('reset') %} +
Сброшено {{ request.query_params.get('reset') }} source(s) в alive.
+{% endif %} + +
Источники
+ + + + {% for s in sources %} + + + + + + + + + {% endfor %} + +
ИсточникТипСтатусПоследний сборОшибкиПриоритет
{{ s.name }}
{{ s.slug }}
{{ s.crawler }} + {% if s.status == 'dead' %}dead + {% elif s.enabled %}alive + {% else %}off{% endif %} + {{ s.last_fetch | dt if s.last_fetch else '—' }}{% if s.error_count %}{{ s.error_count }}{% else %}—{% endif %} + {% if s.last_error %}
{{ s.last_error }}{% endif %}
{{ s.priority }}
+ +
Последние запуски
+ + + + {% for r in runs %} + + + + {% endfor %} + +
ВремяИсточникСтатусПостовДлительность
{{ r.started_at | dt }}{{ r.source_name or r.source_id }}{{ r.status }}{{ r.posts_new }}/{{ r.posts_fetched }}{% if r.duration_ms %}{{ "%.1f"|format(r.duration_ms/1000) }}s{% else %}—{% endif %}
+{% endblock %} \ No newline at end of file