openspec: архив 14 завершённых change-ов (веб-фиксы, crawler-queue, own-content-hub, publisher-service); спеки влиты в openspec/specs

This commit is contained in:
kpa39l
2026-09-16 17:01:00 +00:00
parent 584582a48c
commit 771f6a8276
88 changed files with 1632 additions and 3 deletions
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-09-15
@@ -0,0 +1,107 @@
# 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
@@ -0,0 +1,46 @@
# Proposal: crawler-queue
## Why
Сейчас фазы сбора и обработки не разделены: каждый краулер (telegram/rss) сам
собирает посты и сам же (через `classifier.classify`) разбирает их на
классификацию. Это последовательно: весь прогон — один поток, один источник за
раз. При этом самая дорогая часть — фаза 3 (классификация через Ollama,
дотягивание текста трафилатурой) — выполняется последовательно и без видимости
процесса: непонятно, сколько кандидатов в очереди, какие источники живы, что
упало.
Пользователь хочет:
1. Отдельная страница `/crawlers` — статус работы краулеров (alive/dead,
last_fetch, ошибки) и размер очереди.
2. Двухфазный сбор: (1) краулеры проходят по источникам → ищут новых кандидатов;
(2) найденное кладётся в очередь; (3) воркеры разбирают очередь в НЕСКОЛЬКО
ПОТОКОВ.
Анализ текущего кода:
- Очередь фазы 2 уже существует — это `posts` со `status='new'` и
`classified IS NULL OR classified=0` (классификатор выбирает именно их,
`SELECT ... WHERE classified IS NULL OR classified=0 LIMIT ?`).
- Отдельная таблица `crawl_queue` НЕ нужна — она дублировала бы `posts`.
«Размер очереди» = `COUNT(*) FROM posts WHERE classified IS NULL OR classified=0`.
- Чего нет: (а) параллельной фазы 3 (ThreadPoolExecutor), (б) страницы `/crawlers`,
(в) разделения «сбор» и «обработка» как независимых запусков.
## Goal
- Параллельная фаза 3: воркеры (N потоков) разбирают очередь
`posts(classified=0)` — классификация (keywords → Ollama) + дотягивание текста
(trafilatura) по пайплайну `classifier.classify`.
- Страница `/crawlers` (веб, за auth): таблица источников (slug, name, crawler,
status, last_fetch, last_error, error_count, priority), размер очереди
(pending/processing/done), кнопки «Запустить сейчас» и «Сбросить dead».
- Разделение: `crawler/worker.py` (фаза 3, N потоков) и `crawler/crawl_sources.py`
(фаза 1: собрать новых кандидатов в очередь) — независимые запуски.
- Терминология (для доков): true-конкурентность на IO-bound задачах через потоки;
GIL не мешает, т.к. фаза 3 — ожидание сети.
## Non-goals
- Не вводим Redis/RabbitMQ/брокеры — SQLite-очередь (claim по `posts`) достаточна.
- Не выносим в микросервисы — всё в рамках существующего веба/краулера.
- Не переписываем telegram_crawler; RSS-краулер уже работает.
@@ -0,0 +1,74 @@
# 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) — они уже работают.
- НЕ выносим воркеры в отдельный процесс/сервис — модуль в том же проекте.
@@ -0,0 +1,40 @@
# Tasks: crawler-queue
## 1. Воркер (фаза 3)
- [x] 1.1 crawler/worker.py: claim_post (атомарный claim: classified=0 → -1, BEGIN IMMEDIATE)
Проверка: in-memory тест — два вызова claim_post дают разные id (атомарно)
- [x] 1.2 process_post: classify_text + трафилатура для summary-only; UPDATE posts (direction/relevance/interest/summary/classified=1)
Проверка: in-memory — пост «linux» → tech, classified=1; «игры» → games
- [x] 1.3 Исключение → classified=0 (вернуть в очередь)
Проверка: try/except в work_loop, UPDATE classified=0; реализовано
- [x] 1.4 run(workers=8, limit=200): ThreadPoolExecutor, as_completed, счётчики (processed/classified/llm_ok)
Проверка: `python -m crawler.worker --workers 4 --limit 2` → 8 обработано; без ошибок
## 2. Сбор (фаза 1)
- [x] 2.1 crawler/crawl_sources.py: запуск telegram_crawler + rss_crawler (сбор новых кандидатов в очередь)
Проверка: `python -m crawler.crawl_sources --all` → запускает оба; `--crawler rss` — только RSS
## 3. Страница /crawlers
- [x] 3.1 web/app.py: GET /crawlers (за auth) — таблица источников (slug, name, crawler, status, last_fetch, last_error, error_count, priority) + карточки очереди (pending/processing/done)
Проверка: GET /crawlers с auth → 200, таблица с lwn (проверено httpx)
- [x] 3.2 web/templates/crawlers.html — шаблон (таблица + карточки + кнопки)
- [x] 3.3 POST /crawlers/run → запуск worker.run (фаза 3); POST /crawlers/reset → sources.status='alive', error_count=0
Проверка: POST /crawlers/run → воркер реально стартует (logs/worker.log), 200/302; reset — UPDATE
## 4. Интеграция и доки
- [x] 4.1 AGENT.MD — команды worker/crawl_sources (добавлены)
- [x] 4.2 STATUS.md / TODO.md — страница /crawlers, параллельная фаза 3 (обновлены)
- [x] 4.3 `openspec validate crawler-queue` → 0 ошибок
- [x] 4.4 TODO.md: задача «crawler-очередь» закрыта ✅
## Примечания (найденные при реализации)
- SQLite: соединение привязано к потоку — создаётся ВНУТРИ work_loop (нельзя делить между потоками).
- Локальная Ollama (qwen3:8b) держит 1 слот — 8 параллельных LLM-запросов → таймауты (посты возвращаются в очередь, не теряются). Добавлен `LLM_SEM` (threading.BoundedSemaphore, MAX_CONCURRENT_LLM=2).
- Очередь = posts (classified IS NULL OR 0) — отдельная таблица НЕ нужна (дублирование).
- error_count/etag для RSS — из rss_state (в sources их нет); запрос /crawlers использует COALESCE.
- **SOURCE_RULES (classifier/keywords.py)**: правило источника с приоритетом над словарём и LLM. `lwn → tech` — все посты LWN (технологическое СМИ) получают направление tech детерминированно, без LLM (method='source-rule'). Работает в worker.py и старом CLI classifier.classify (оба JOIN sources → source_slug).