# 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 ```bash # Новое 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) ```python # 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 ```bash .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