Support kind=llm profiles (instruction/schema), optional multi-event posts via #eN URLs, and recover stale running/queued parse jobs after worker crashes. Co-authored-by: Cursor <cursoragent@cursor.com>
178 lines
8.9 KiB
Markdown
178 lines
8.9 KiB
Markdown
# Архитектура парсинга (ЦП)
|
||
|
||
Центр парсинга (**ЦП**) — сервисы `cp-workers*` с **реестром адаптеров** источников. Собирают события и отправляют их в ЦА через internal API.
|
||
|
||
## Роль в платформе
|
||
|
||
```mermaid
|
||
flowchart LR
|
||
CA[ЦА ca-api]
|
||
Redis[(Redis cp:jobs:family)]
|
||
CP[ЦП adapters]
|
||
Src[Sources]
|
||
|
||
CA -->|"RPUSH JobPayload"| Redis
|
||
Redis -->|"BLPOP"| CP
|
||
CP --> Src
|
||
CP -->|"POST /internal/ingest"| CA
|
||
CP -->|"PATCH /internal/jobs/{id}"| CA
|
||
CP -->|"GET /internal/listener/subscriptions"| CA
|
||
```
|
||
|
||
| Направление | Механизм | Назначение |
|
||
|-------------|----------|------------|
|
||
| ЦА → ЦП | Redis `cp:jobs:{family}` | Batch-задания по семейству адаптеров |
|
||
| ЦП → источники | Адаптер (`telegram` / `crawl4ai` / `viina`) | Fetch + extract |
|
||
| ЦП → ЦА | `POST /internal/ingest` | Запись событий (`IngestEventItem`) |
|
||
| ЦП → ЦА | `PATCH /internal/jobs/{id}` | Статус batch-задания |
|
||
| ЦП ← ЦА | `GET /internal/listener/subscriptions` | Каналы real-time (только Telegram) |
|
||
|
||
Общие контракты: [`contracts/jobs.py`](../../contracts/jobs.py), [`contracts/ingest.py`](../../contracts/ingest.py), [`contracts/sources.py`](../../contracts/sources.py), [`contracts/queues.py`](../../contracts/queues.py).
|
||
|
||
## Адаптеры источников
|
||
|
||
Каждый `source_type` реализует `SourceAdapter`:
|
||
|
||
```python
|
||
async def run(job_id, source_config, *, ctx) -> tuple[list[dict], str | None]
|
||
```
|
||
|
||
Выход — список dict в форме `IngestEventItem`. ЦА **не** знает про Crawl4AI/VIINA.
|
||
|
||
| source_type | Семейство / очередь | Сервис | Зависимости |
|
||
|-------------|---------------------|--------|-------------|
|
||
| `telegram` | `telegram` → `cp:jobs:telegram` | `cp-workers` | Telethon |
|
||
| `crawl4ai` | `web` → `cp:jobs:web` | `cp-workers-web` | Crawl4AI + Playwright |
|
||
| `viina` | `nlp` → `cp:jobs:nlp` | `cp-workers-nlp` | httpx + BeautifulSoup |
|
||
|
||
Маршрутизация при enqueue в ЦА: [`queue_key_for_source`](../../contracts/queues.py).
|
||
Воркер слушает `WORKER_FAMILIES` / очереди своих `ENABLED_ADAPTERS`. Чужой job → requeue в нужную очередь.
|
||
|
||
Правила для всех адаптеров:
|
||
|
||
- стабильный `source_url` (дедуп в ЦА);
|
||
- `source_type` события = тип адаптера;
|
||
- `source_config` валидируется схемами из `contracts/sources.py`.
|
||
|
||
### Режимы извлечения (Telegram)
|
||
|
||
Цель всегда одна: фиксированные поля `IngestEventItem` / Event в ЦА (`title`, `description`, `locality`, `coords`→lat/lng, `event_date`, `topic`, …). Кастомные пользовательские таблицы — roadmap.
|
||
|
||
| `extract_mode` | Что делает | Откуда в UI |
|
||
|----------------|------------|-------------|
|
||
| `heuristic` | Legacy regex-парсер постов | raw job API (не пара) |
|
||
| `llm` | DeepSeek **на каждый** пост (batch) | профиль `kind=llm` → pair flatten |
|
||
| `profile` | Статичные правила `HeuristicProfile` | профиль `kind=heuristic` → pair flatten |
|
||
|
||
### LLM-режим (`extract_mode: llm`)
|
||
|
||
Для неструктурированных Telegram-постов и Crawl4AI:
|
||
|
||
- ключ `DEEPSEEK_API_KEY` в `.env` (воркеры `cp-workers` / `cp-workers-web`; preview в `ca-api`);
|
||
- Telegram batch: текст → DeepSeek JSON → один или несколько `IngestEventItem` (`workers/llm_extract.py`);
|
||
- опционально `extract_schema`, `instruction`, `required_fields` (post-extract gate поверх `is_event`);
|
||
- **`multi_event`** (opt-in в `llm_profile` / `TelegramSourceConfig`): модель возвращает массив `events`; каждое событие — отдельный ingest. URL: один event → `post.url`; несколько → `post.url#e1`, `#e2`, … (дедуп ЦА по `source_url`);
|
||
- heuristic / profile по-прежнему **1 пост → 1 событие**;
|
||
- Crawl4AI: страница → `LLMExtractionStrategy` (DeepSeek) с fallback на тот же DeepSeek по markdown;
|
||
- **listener:** для `llm` пока fallback на heuristic (LLM — batch-only by design).
|
||
|
||
Reusable LLM-профиль в ЦА: `ParserProfile.kind=llm` + JSON `llm_profile` (`contracts/llm_profile.py`). При enqueue flatten → `extract_mode=llm` + schema/instruction/required_fields/`multi_event`.
|
||
|
||
### Profile-режим (`extract_mode: profile`)
|
||
|
||
Статичный парсер из конструктора ЦА:
|
||
|
||
- `heuristic_profile` в `source_config` (схема `contracts/heuristic_profile.py`);
|
||
- интерпретатор: `workers/heuristic_profile.py` (те же правила, что preview в ЦА);
|
||
- `required_fields`: пост без заполненных обязательных полей не ингестится (batch + listener);
|
||
- генерация правил — один раз в админке (`/admin/parser-profiles/generate`); runtime **без** LLM.
|
||
|
||
Reusable heuristic-профиль: `ParserProfile.kind=heuristic` + `heuristic_profile`. Pair на «Парсеры» → flatten `extract_mode=profile`.
|
||
|
||
UI «Профили»: выбор `kind` (heuristic | llm). Связка канал+профиль на «Парсеры».
|
||
|
||
```mermaid
|
||
flowchart LR
|
||
Job[JobPayload]
|
||
Reg[AdapterRegistry]
|
||
TG[TelegramAdapter]
|
||
C4[Crawl4AIAdapter]
|
||
VI[ViinaAdapter]
|
||
Ingest[IngestEventItem]
|
||
|
||
Job --> Reg
|
||
Reg --> TG --> Ingest
|
||
Reg --> C4 --> Ingest
|
||
Reg --> VI --> Ingest
|
||
```
|
||
|
||
## Структура каталога
|
||
|
||
```
|
||
centers/parsing/
|
||
├── ARCHITECTURE.md
|
||
└── workers/
|
||
├── Dockerfile # context = repo root; ARG REQUIREMENTS_FILE
|
||
├── requirements.txt # telegram
|
||
├── requirements-web.txt # crawl4ai
|
||
├── requirements-nlp.txt # viina
|
||
├── worker.py
|
||
└── workers/
|
||
├── adapters/
|
||
│ ├── base.py # SourceAdapter, WorkerContext
|
||
│ ├── registry.py # ENABLED_ADAPTERS
|
||
│ ├── telegram.py
|
||
│ ├── crawl4ai_adapter.py
|
||
│ └── viina.py
|
||
├── converter.py
|
||
├── parsers/telegram_events.py
|
||
└── sources/ # Telethon session / listener / client
|
||
```
|
||
|
||
## Режимы работы
|
||
|
||
| Режим | Условие | Поведение |
|
||
|-------|---------|-----------|
|
||
| Listener + batch | `ENABLED_ADAPTERS` включает `telegram` и `TELEGRAM_LISTENER_ENABLED=true` | Shared Telethon + listener + `worker_loop` |
|
||
| Только batch | listener выключен или нет telegram | Только `BLPOP` по очередям семейства |
|
||
|
||
## Поток batch-заданий
|
||
|
||
1. ЦА `enqueue_job` → Redis `cp:jobs:{family}` с `{ job_id, source_type, source_config }`
|
||
2. Воркер семейства: `BLPOP` → `handle_job` → `registry.get(source_type).run(...)`
|
||
3. `POST /internal/ingest` + статус job
|
||
|
||
Legacy-ключ `cp:jobs` по-прежнему дренируется telegram-воркером (совместимость).
|
||
|
||
## Telegram real-time
|
||
|
||
`TelegramListener`: подписки из ЦА, ingest с `listener: true` (статус `ParseJob` не трогается).
|
||
|
||
- `extract_mode=profile` — те же heuristic-правила, что batch;
|
||
- `extract_mode=llm` — **не** вызывает DeepSeek; fallback на legacy heuristic (LLM только в batch).
|
||
|
||
## Переменные окружения
|
||
|
||
| Переменная | Назначение |
|
||
|------------|------------|
|
||
| `ENABLED_ADAPTERS` | Список адаптеров через запятую (`telegram`, `crawl4ai`, `viina`) |
|
||
| `WORKER_FAMILIES` | Какие семейства очередей слушать (`telegram`, `web`, `nlp`) |
|
||
| `REDIS_URL` / `CA_API_URL` / `INTERNAL_TOKEN` | Как раньше |
|
||
| `TELEGRAM_*` | Только для `cp-workers` |
|
||
|
||
## Как добавить новый источник
|
||
|
||
1. Схема `source_config` в [`contracts/sources.py`](../../contracts/sources.py) + запись в `SOURCE_FAMILY` ([`queues.py`](../../contracts/queues.py))
|
||
2. Класс адаптера в `workers/adapters/` + factory в `registry.py`
|
||
3. При тяжёлых deps — `requirements-*.txt` и сервис в `docker-compose.yml`
|
||
4. Поля формы в UI «Парсеры»
|
||
|
||
## Связанные части ЦА
|
||
|
||
| Файл ЦА | Роль |
|
||
|---------|------|
|
||
| `services/jobs.py` | `enqueue_job` → `cp:jobs:{family}` |
|
||
| `routers/admin.py` | CRUD + валидация `source_config` |
|
||
| `services/ingest.py` | Сохранение Event + карта |
|
||
| UI `/parsers` | Выбор `telegram` / `crawl4ai` / `viina` |
|