Files
gitrusprusandCursor 3f9dc6643b Add reusable LLM parser profiles with multi-event extract.
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>
2026-09-13 20:22:38 +03:00

8.9 KiB
Raw Permalink Blame History

Архитектура парсинга (ЦП)

Центр парсинга (ЦП) — сервисы cp-workers* с реестром адаптеров источников. Собирают события и отправляют их в ЦА через internal API.

Роль в платформе

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/ingest.py, contracts/sources.py, contracts/queues.py.

Адаптеры источников

Каждый source_type реализует SourceAdapter:

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.
Воркер слушает 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). Связка канал+профиль на «Парсеры».

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 + запись в SOURCE_FAMILY (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