Add developer docs for CA/CP/PI flows and wire cp-workers to the host SOCKS proxy so local Telegram auth and parsing work reliably. Co-authored-by: Cursor <cursoragent@cursor.com>
5.0 KiB
5.0 KiB
Поток данных MapMil
Пошаговый путь события от настройки парсера до карты и внешнего API.
См. также: обзор архитектуры, ЦП.
Схема end-to-end
sequenceDiagram
participant UI as Admin UI
participant CA as ca-api
participant DB as PostgreSQL
participant R as Redis
participant CP as cp-workers
participant Src as Source
participant Map as Map UI
participant PI as External client
UI->>CA: POST /admin/jobs
CA->>DB: INSERT ParseJob queued
CA->>R: RPUSH cp:jobs:family
R->>CP: BLPOP
CP->>CA: PATCH /internal/jobs/id running
CP->>Src: fetch / listen
Src-->>CP: raw posts/pages
CP->>CA: POST /internal/ingest
CA->>DB: INSERT Event if new source_url
CA->>DB: sync MapObject if coords
CA->>CP: 200 ingested/skipped
CP->>CA: PATCH job completed or failed
Map->>CA: GET /api/map/objects
PI->>CA: GET /api/v1/events Bearer
1. Создание парсера (ЦА)
- Пользователь в UI
/parsersилиPOST /admin/jobsпередаёт:source_type(telegram|crawl4ai|viina);source_config(валидируетсяcontracts/sources.py);- опционально
interval_seconds,is_active.
- ЦА пишет строку
ParseJobсо статусомqueued. - ЦА делает
RPUSHв Redis-ключcp:jobs:{family}(contracts/queues.py).
Примеры source_config:
{ "channel": "example", "limit": 50, "extract_mode": "heuristic" }
{ "urls": ["https://news.example/"], "extract_mode": "llm" }
2. Batch-обработка (ЦП)
- Воркер нужного семейства делает
BLPOPсвоей очереди. PATCH /internal/jobs/{id}?status=running.- Registry выбирает адаптер по
source_type→adapter.run(...). - Адаптер возвращает список dict в форме
IngestEventItem(+ опциональный warning). - ЦП шлёт
POST /internal/ingestс{ job_id, events }. - При успехе/ошибке —
PATCHстатусаcompleted/failed.
Если job попал не в ту очередь (чужой family), воркер перекладывает его в правильную очередь.
3. Real-time Telegram (параллельно)
Только cp-workers при TELEGRAM_LISTENER_ENABLED=true:
- Listener опрашивает
GET /internal/listener/subscriptions(активные telegram-jobs). - Telethon получает
NewMessage/Album. - Пост парсится (эвристика) → один event →
POST /internal/ingestсlistener: true. - Статус
ParseJobот listener не переводится вcompleted(в отличие от batch).
4. Ingest в ЦА
Сервис centers/analytics/api/app/services/ingest.py:
- Для каждого item проверяет уникальность
source_url. - Уже есть → skip (дедуп).
- Нет → INSERT
Event. - Если есть
latitude/longitude→ создать/обновить связанныйMapObject. - Batch (
listenerне true) обновляет статус job иlast_run_at.
5. Карта
- UI периодически (около 30 с) и по фильтрам вызывает
GET /api/map/objects. - Ответ — точки с данными события (даты, регион, тема, источник).
- Ручные объекты карты живут в
map_objectsбез обязательногоevent_id(CRUD через UI //api/objects).
6. Внешний потребитель (ПИ)
- В UI
/consumersсоздаётся consumer + API-ключ (в БД хранится hash). - Клиент:
GET /api/v1/events+Authorization: Bearer <key>. - ЦА применяет
ConsumerFilter(регионы, темы,date_from) и отдаёт срез событий.
7. Периодический перезапуск
Scheduler в ca-api (~30 с):
- активный job;
- статус
completedилиfailed; - прошло ≥
interval_secondsс последнего запуска
→ снова queued + enqueue в Redis.
Точки отказа (для отладки)
| Симптом | Куда смотреть |
|---|---|
Job висит в queued |
Redis, нужный cp-workers*, WORKER_FAMILIES |
Job failed |
last_error в admin, логи контейнера адаптера |
| Событий 0 при успехе | Парсер/LLM не извлёк поля; дедуп по source_url |
| Нет точек на карте | Нет координат у Event; фильтры даты/региона в UI |
401 на /api/v1/events |
Ключ, is_active consumer |
| Telegram auth error | data/telegram.session, TELEGRAM_API_ID/HASH |