Document platform architecture and Telegram proxy local setup.
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>
This commit is contained in:
@@ -0,0 +1,134 @@
|
||||
# Поток данных MapMil
|
||||
|
||||
Пошаговый путь события от настройки парсера до карты и внешнего API.
|
||||
|
||||
См. также: [обзор архитектуры](architecture-overview.md), [ЦП](../centers/parsing/ARCHITECTURE.md).
|
||||
|
||||
---
|
||||
|
||||
## Схема end-to-end
|
||||
|
||||
```mermaid
|
||||
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. Создание парсера (ЦА)
|
||||
|
||||
1. Пользователь в UI `/parsers` или `POST /admin/jobs` передаёт:
|
||||
- `source_type` (`telegram` | `crawl4ai` | `viina`);
|
||||
- `source_config` (валидируется `contracts/sources.py`);
|
||||
- опционально `interval_seconds`, `is_active`.
|
||||
2. ЦА пишет строку `ParseJob` со статусом `queued`.
|
||||
3. ЦА делает `RPUSH` в Redis-ключ `cp:jobs:{family}` (`contracts/queues.py`).
|
||||
|
||||
Примеры `source_config`:
|
||||
|
||||
```json
|
||||
{ "channel": "example", "limit": 50, "extract_mode": "heuristic" }
|
||||
```
|
||||
|
||||
```json
|
||||
{ "urls": ["https://news.example/"], "extract_mode": "llm" }
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 2. Batch-обработка (ЦП)
|
||||
|
||||
1. Воркер нужного семейства делает `BLPOP` своей очереди.
|
||||
2. `PATCH /internal/jobs/{id}?status=running`.
|
||||
3. Registry выбирает адаптер по `source_type` → `adapter.run(...)`.
|
||||
4. Адаптер возвращает список dict в форме `IngestEventItem` (+ опциональный warning).
|
||||
5. ЦП шлёт `POST /internal/ingest` с `{ job_id, events }`.
|
||||
6. При успехе/ошибке — `PATCH` статуса `completed` / `failed`.
|
||||
|
||||
Если job попал не в ту очередь (чужой family), воркер **перекладывает** его в правильную очередь.
|
||||
|
||||
---
|
||||
|
||||
## 3. Real-time Telegram (параллельно)
|
||||
|
||||
Только `cp-workers` при `TELEGRAM_LISTENER_ENABLED=true`:
|
||||
|
||||
1. Listener опрашивает `GET /internal/listener/subscriptions` (активные telegram-jobs).
|
||||
2. Telethon получает `NewMessage` / `Album`.
|
||||
3. Пост парсится (эвристика) → один event → `POST /internal/ingest` с `listener: true`.
|
||||
4. Статус `ParseJob` от listener **не** переводится в `completed` (в отличие от batch).
|
||||
|
||||
---
|
||||
|
||||
## 4. Ingest в ЦА
|
||||
|
||||
Сервис `centers/analytics/api/app/services/ingest.py`:
|
||||
|
||||
1. Для каждого item проверяет уникальность `source_url`.
|
||||
2. Уже есть → skip (дедуп).
|
||||
3. Нет → INSERT `Event`.
|
||||
4. Если есть `latitude` / `longitude` → создать/обновить связанный `MapObject`.
|
||||
5. Batch (`listener` не true) обновляет статус job и `last_run_at`.
|
||||
|
||||
---
|
||||
|
||||
## 5. Карта
|
||||
|
||||
1. UI периодически (около 30 с) и по фильтрам вызывает `GET /api/map/objects`.
|
||||
2. Ответ — точки с данными события (даты, регион, тема, источник).
|
||||
3. Ручные объекты карты живут в `map_objects` без обязательного `event_id` (CRUD через UI / `/api/objects`).
|
||||
|
||||
---
|
||||
|
||||
## 6. Внешний потребитель (ПИ)
|
||||
|
||||
1. В UI `/consumers` создаётся consumer + API-ключ (в БД хранится hash).
|
||||
2. Клиент: `GET /api/v1/events` + `Authorization: Bearer <key>`.
|
||||
3. ЦА применяет `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` |
|
||||
Reference in New Issue
Block a user