# Spec: crawler-queue ## Purpose Двухфазный пайплайн сбора и обработки новостей: краулеры (фаза 1) собирают кандидатов в очередь (`posts` со `status='new'` и неклассифицированные), воркеры (фаза 3) разбирают её параллельно (ThreadPoolExecutor). Отдельная страница `/crawlers` показывает статус источников и размер очереди. ## ADDED Requirements ### Requirement: Параллельная фаза 3 (воркеры) Система MUST предоставлять `crawler/worker.py` с пулом `ThreadPoolExecutor(max_workers=N)`, где каждый воркер атомарно забирает пост из очереди (`UPDATE posts SET classified=-1 WHERE id=? AND classified=0`), классифицирует его (`classify_text`: keywords → Ollama, trafilatura для summary-only) и пишет результат (direction/relevance/interest/summary, classified=1). При исключении посте MUST возвращаться в очередь (classified=0). #### Scenario: Параллельная обработка очереди - **GIVEN** 100 постов со `classified=0` - **WHEN** `python -m crawler.worker --workers 8 --limit 200` - **THEN** все 100 постов обработаны (classified=1), дублей нет (каждый обработан ровно 1 раз) #### Scenario: Сбой воркера - **GIVEN** пост, у которого `classify_text` бросает исключение (Ollama недоступна) - **WHEN** воркер обрабатывает пост - **THEN** пост возвращается в очередь (classified=0), воркер продолжает работу, запуск не падает ### Requirement: Очередь на основе posts (без новой таблицы) «Размер очереди» MUST вычисляться из существующей таблицы `posts`: `SELECT COUNT(*) FROM posts WHERE classified IS NULL OR classified=0` (pending), `classified=-1` (processing), `classified=1 AND fetched_at >= date('now')` (done today). Новая таблица для очереди НЕ создаётся — она дублировала бы `posts`. #### Scenario: Размер очереди - **GIVEN** в posts 10 новых (classified=0) и 2 в обработке (classified=-1) - **WHEN** страница /crawlers запрашивает размер очереди - **THEN** pending=10, processing=2, done today=0 ### Requirement: Страница /crawlers Веб MUST предоставлять `GET /crawlers` (за аутентификацией): таблица источников (slug, name, crawler, status, last_fetch, last_error за 200 симв., error_count, priority; для rss — etag/modified/last_build_date из rss_state) и карточки очереди (pending/processing/done). Кнопка «Сбросить dead» (POST /crawlers/reset) MUST устанавливать sources.status='alive', error_count=0, last_error=NULL для всех источников со status='dead'. #### Scenario: Просмотр статуса - **GIVEN** источник lwn (rss, alive) и 15 новых постов в очереди - **WHEN** GET /crawlers - **THEN** страница показывает lwn с статусом alive, размер очереди 15 #### Scenario: Сброс dead-источников - **GIVEN** источник со status='dead', error_count=7 - **WHEN** POST /crawlers/reset - **THEN** источник становится alive, error_count=0, last_error=NULL ### Requirement: Фаза 1 (сбор) отдельно от фазы 3 (обработка) Система MUST предоставлять `crawler/crawl_sources.py` — запуск сбора всех включённых источников (telegram_crawler + rss_crawler) без классификации; новые посты попадают в очередь (posts, classified=0). Обработка (фаза 3) запускается отдельно (`crawler/worker.py`). #### Scenario: Сбор без классификации - **GIVEN** включённые источники telegram и rss - **WHEN** `python -m crawler.crawl_sources --all` - **THEN** новые посты добавлены в posts с classified=0, классификация НЕ запущена ## NOT Requirements - НЕ создаём отдельную таблицу очереди (`crawl_queue`) — используется `posts`. - НЕ вводим Redis/RabbitMQ/брокеры — атомарный claim через SQLite достаточен. - НЕ переписываем telegram_crawler/rss_crawler (фаза 1) — они уже работают. - НЕ выносим воркеры в отдельный процесс/сервис — модуль в том же проекте.