Files

5.3 KiB

Design: crawler-queue

Approach

Двухфазный пайплайн поверх существующей модели данных (БЕЗ новой таблицы-очереди):

Фаза 1 (сбор)          Фаза 3 (обработка, ПАРАЛЛЕЛЬНО)
sources.yaml ──►        ┌──────────────┐
  telegram_crawler ──►  │ ThreadPool   │──► posts(classified=0)
  rss_crawler ──►       │  (8 воркеров)│──► keywords → Ollama → direction
        │               └──────────────┘    + trafilatura (полный текст)
        ▼
posts (status='new', classified=0) = ОЧЕРЕДЬ
  • Очередь = таблица posts уже существует (status='new' + classified=0). Размер очереди: SELECT COUNT(*) FROM posts WHERE classified IS NULL OR classified=0.
  • Фаза 3 — новый crawler/worker.py: ThreadPoolExecutor(8), каждый воркер claim'ит посты (UPDATE posts SET classified=-1 WHERE id=? AND classified=0 → атомарный claim, без дублей), классифицирует (keywords → Ollama), пишет direction/relevance/interest/summary, ставит classified=1. Упавшие → classified=0 обратно (ретрай).
  • Фаза 1 — crawler/crawl_sources.py: обёртка, запускающая telegram_crawler и rss_crawler (сбор источников → посты в очередь). Раздельные запуски.
  • Страница /crawlers — FastAPI route в web/app.py + шаблон web/templates/crawlers.html: таблица источников (из sources + rss_state), размер очереди (pending/processing/done), кнопки «Запустить сейчас» (POST /crawlers/run) и «Сбросить dead» (POST /crawlers/reset).

Files

# Новое
crawler/worker.py            # фаза 3: ThreadPoolExecutor, claim, классификация
crawler/crawl_sources.py     # фаза 1: запуск telegram+rss краулеров (сбор в очередь)
web/templates/crawlers.html  # страница /crawlers

# Изменения
web/app.py                   # routes: GET /crawlers, POST /crawlers/run, POST /crawlers/reset
sources/sources.py           # + get_rss_state(), + reset_source_status()
AGENT.MD / STATUS.md / TODO.md / WALKTHROUGH.md   # доки
openspec/changes/crawler-queue/                    # этот change

Worker (фаза 3)

# crawler/worker.py
from concurrent.futures import ThreadPoolExecutor, as_completed
import sqlite3

def claim_post(conn, worker_id) -> row | None:
    # атомарный claim: одна строка — один воркер
    conn.execute("BEGIN IMMEDIATE")
    row = conn.execute(
        "SELECT id, text, views, reactions_total, is_own FROM posts "
        "WHERE classified IS NULL OR classified=0 "
        "ORDER BY is_own DESC, id LIMIT 1").fetchone()
    if row:
        conn.execute("UPDATE posts SET classified=-1 WHERE id=?", (row["id"],))
    conn.commit()
    return row

def process_one(row) -> dict:  # classify_text из classifier.classify
    ...

def run(workers: int = 8, limit: int = 200):
    with ThreadPoolExecutor(max_workers=workers) as ex:
        futs = [ex.submit(work_loop, worker_id=i) for i in range(workers)]
        for f in as_completed(futs):
            ...

Каждый воркер:

  1. claim_post → строка post (classified=0 → -1).
  2. classify_text (keywords → Ollama; трафилатура для summary-only).
  3. UPDATE posts SET direction=?, relevance=?, interest=?, summary=?, classified=1 WHERE id=?
  4. При исключении → UPDATE posts SET classified=0 WHERE id=? (вернуть в очередь).

Страница /crawlers

| Источник | slug | name | crawler | status | last_fetch | last_error | Приоритет | Действия | Снизу — карточки очереди:

  • Pending: COUNT(classified IS NULL OR 0)
  • Processing: COUNT(classified=-1)
  • Done (сегодня): COUNT(classified=1 AND fetched_at >= today)

Кнопки:

  • «Запустить сейчас» → POST /crawlers/run → запускает worker.run(8) синхронно (или через subprocess в фоне), редирект на /crawlers.
  • «Сбросить dead» → POST /crawlers/reset → UPDATE sources SET status='alive', error_count=0 WHERE status='dead'.

CLI

.venv/bin/python -m crawler.worker --workers 8 --limit 200     # фаза 3, параллельно
.venv/bin/python -m crawler.crawl_sources --all                # фаза 1, сбор в очередь

Verification

  • openspec validate crawler-queue → 0 ошибок
  • worker.run(8) обрабатывает ≥5 постов из очереди; ПОВТОРНЫЙ запуск — 0 новых (все classified=1)
  • Дубли не возникают: 100 постов, 8 воркеров → 100 строк обновлены, 0 пропущено
  • Страница /crawlers: показывает таблицу источников, размер очереди; кнопка «Сбросить dead» обнуляет error_count