From 6dff3c1c3d0bc0bc9accde8b13be386ccc2e5338 Mon Sep 17 00:00:00 2001 From: gitrusprus Date: Sat, 15 Aug 2026 22:57:07 +0300 Subject: [PATCH] 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 --- AGENTS.md | 2 +- README.md | 28 +-- centers/analytics/ARCHITECTURE.md | 154 ++++++++++++++++ centers/analytics/DISTRIBUTION.md | 91 ++++++++++ docker-compose.yml | 5 + docs/README.md | 35 ++++ docs/architecture-overview.md | 288 ++++++++++++++++++++++++++++++ docs/contracts.md | 88 +++++++++ docs/data-flow.md | 134 ++++++++++++++ docs/local-dev.md | 118 ++++++++++++ scripts/telegram_auth.py | 115 ++++++++++++ 11 files changed, 1046 insertions(+), 12 deletions(-) create mode 100644 centers/analytics/ARCHITECTURE.md create mode 100644 centers/analytics/DISTRIBUTION.md create mode 100644 docs/README.md create mode 100644 docs/architecture-overview.md create mode 100644 docs/contracts.md create mode 100644 docs/data-flow.md create mode 100644 docs/local-dev.md create mode 100644 scripts/telegram_auth.py diff --git a/AGENTS.md b/AGENTS.md index 654d23f..2cfbbb9 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -86,7 +86,7 @@ Adapters return dicts matching `IngestEventItem` (`contracts/ingest.py`). Requir - Rule: `.cursor/rules/platform.mdc` (always on) - Skill: `.cursor/skills/add-parser-adapter/` for end-to-end new sources -- Deep docs: `README.md`, `centers/parsing/ARCHITECTURE.md` +- Deep docs: `docs/` (overview + data-flow), `centers/analytics/ARCHITECTURE.md`, `centers/parsing/ARCHITECTURE.md` ## Commit style diff --git a/README.md b/README.md index 9bdc62d..e2c3cf1 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,8 @@ Единая платформа: ЦА (аналитика и карта) + ЦП (адаптеры парсинга: Telegram, Crawl4AI, VIINA). +**Документация для разработчиков:** [`docs/`](docs/README.md) — [обзор архитектуры](docs/architecture-overview.md), [поток данных](docs/data-flow.md), [локальный запуск](docs/local-dev.md). + ## Архитектура ```mermaid @@ -15,13 +17,14 @@ flowchart LR CA -->|GET /api/v1/events| PI ``` -| Центр | Контейнеры | Назначение | -|-------|------------|------------| -| **ЦА** | `ca-db`, `ca-api`, `ca-frontend` | PostgreSQL, ingest API, карта, distribution API | -| **ЦП** | `cp-workers`, `cp-workers-web`, `cp-workers-nlp` | Адаптеры: Telegram / Crawl4AI / VIINA | -| **Общее** | `redis` | Очереди `cp:jobs:{telegram\|web\|nlp}` | +| Центр | Контейнеры | Назначение | Документ | +|-------|------------|------------|----------| +| **ЦА** | `ca-db`, `ca-api`, `ca-frontend` | PostgreSQL, ingest API, карта, distribution API | [`centers/analytics/ARCHITECTURE.md`](centers/analytics/ARCHITECTURE.md) | +| **ЦП** | `cp-workers`, `cp-workers-web`, `cp-workers-nlp` | Адаптеры: Telegram / Crawl4AI / VIINA | [`centers/parsing/ARCHITECTURE.md`](centers/parsing/ARCHITECTURE.md) | +| **ПИ** | (маршрут в ЦА) | `GET /api/v1/events` | [`centers/analytics/DISTRIBUTION.md`](centers/analytics/DISTRIBUTION.md) | +| **Общее** | `redis` | Очереди `cp:jobs:{telegram\|web\|nlp}` | [`docs/contracts.md`](docs/contracts.md) | -Подробности ЦП: [`centers/parsing/ARCHITECTURE.md`](centers/parsing/ARCHITECTURE.md). +Подробный обзор платформы: [`docs/architecture-overview.md`](docs/architecture-overview.md). ## Структура monorepo @@ -29,13 +32,16 @@ flowchart LR MapMil/ ├── centers/ │ ├── analytics/ -│ │ ├── api/ # CA backend (FastAPI + PostgreSQL) -│ │ └── frontend/ # CA admin UI (Vue + Leaflet) +│ │ ├── api/ # CA backend (FastAPI + PostgreSQL) +│ │ ├── frontend/ # CA admin UI (Vue + Leaflet) +│ │ ├── ARCHITECTURE.md +│ │ └── DISTRIBUTION.md # ПИ │ └── parsing/ │ ├── ARCHITECTURE.md -│ └── workers/ # CP workers + adapters -├── contracts/ # Shared schemas (ingest, jobs, sources, queues) -├── data/ # telegram.session (локально, не в git) +│ └── workers/ # CP workers + adapters +├── contracts/ # Shared schemas (ingest, jobs, sources, queues) +├── docs/ # Документация для разработчиков +├── data/ # telegram.session (локально, не в git) ├── docker-compose.yml └── .env ``` diff --git a/centers/analytics/ARCHITECTURE.md b/centers/analytics/ARCHITECTURE.md new file mode 100644 index 0000000..84b0391 --- /dev/null +++ b/centers/analytics/ARCHITECTURE.md @@ -0,0 +1,154 @@ +# Архитектура ЦА (Analytics Center) + +Центр аналитики — ядро MapMil: PostgreSQL, FastAPI, Vue admin/карта, distribution API для ПИ. + +Общий контекст: [docs/architecture-overview.md](../../docs/architecture-overview.md). +ПИ отдельно: [DISTRIBUTION.md](DISTRIBUTION.md). + +--- + +## Общими чертами + +ЦА: + +- хранит события, jobs, объекты карты, consumers; +- отдаёт единственный UI (`ca-frontend`); +- ставит задания в Redis для ЦП; +- принимает ingest от ЦП; +- отдаёт срезы внешним системам через `/api/v1/events`. + +```mermaid +flowchart TB + UI[ca-frontend Vue] + API[ca-api FastAPI] + DB[(PostgreSQL)] + Redis[(Redis)] + CP[CP workers] + + UI --> API + API --> DB + API -->|enqueue| Redis + Redis --> CP + CP -->|/internal/*| API +``` + +Контейнеры: `ca-db`, `ca-api`, `ca-frontend` (+ общий `redis`). + +--- + +## Подробнее + +### Структура + +```text +centers/analytics/ +├── ARCHITECTURE.md +├── DISTRIBUTION.md +├── api/ +│ ├── Dockerfile +│ ├── requirements.txt +│ └── app/ +│ ├── main.py # lifespan: migrations, scheduler, seed +│ ├── database.py +│ ├── models.py +│ ├── schemas.py +│ ├── deps.py # X-Internal-Token +│ ├── seed.py +│ ├── storage.py # uploads +│ ├── routers/ +│ │ ├── objects.py # /api/health, /api/objects, media +│ │ ├── map.py # /api/map/* +│ │ ├── admin.py # /admin/* +│ │ ├── internal.py # /internal/* (только ЦП) +│ │ └── v1.py # /api/v1/* (ПИ) +│ └── services/ +│ ├── jobs.py # Redis RPUSH +│ ├── scheduler.py # периодический re-queue +│ ├── ingest.py # дедуп + map sync +│ ├── filtering.py +│ ├── map_query.py +│ ├── events_query.py +│ ├── analytics.py +│ └── migrations.py +└── frontend/ + ├── Dockerfile # Vite build + nginx + ├── nginx.conf # proxy /api /admin /internal → ca-api + └── src/ + ├── views/ # Map, Parsers, Events, Analytics, Consumers + ├── components/ # карта, CRUD объектов + ├── api/ # HTTP-клиенты + └── router/index.ts +``` + +### Модели данных + +| Модель | Таблица | Назначение | +|--------|--------|------------| +| `Event` | `events` | Нормализованное событие; UK `source_url` | +| `ParseJob` | `parse_jobs` | Конфиг парсера, интервал, статус | +| `MapObject` | `map_objects` | Точка на карте (event или ручная) | +| `ObjectMedia` | `object_media` | Файлы к объектам | +| `Consumer` | `consumers` | Подписчик ПИ (hash ключа) | +| `ConsumerFilter` | `consumer_filters` | Фильтры среза ПИ | + +### HTTP-поверхности + +| Prefix | Кто вызывает | Содержание | +|--------|--------------|------------| +| `/api/*` | UI, публичный health | Карта, объекты, медиа | +| `/admin/*` | UI admin | Jobs, events, analytics, consumers | +| `/internal/*` | Только ЦП | ingest, job status, listener subscriptions | +| `/api/v1/*` | Внешние клиенты | Events с Bearer-ключом | + +Internal защищён заголовком `X-Internal-Token` (`INTERNAL_TOKEN`). + +### Jobs и scheduler + +1. `POST /admin/jobs` / retry → запись `ParseJob` + `enqueue_job` (`services/jobs.py`). +2. Очередь: `cp:jobs:{family}` из `contracts/queues.py`. +3. `services/scheduler.py` — тик ~30 с, повторная постановка активных jobs по `interval_seconds`. + +### Ingest + +`POST /internal/ingest` → `services/ingest.py`: + +- дедуп по `source_url`; +- создание `Event`; +- при координатах — sync `MapObject`; +- batch обновляет статус job; `listener: true` — нет. + +### Frontend (маршруты) + +| Path | View | +|------|------| +| `/` | `MapViewPage.vue` | +| `/parsers` | `ParsersView.vue` | +| `/events` | `EventsView.vue` | +| `/analytics` | `AnalyticsView.vue` | +| `/consumers` | `ConsumersView.vue` | + +Карта: Leaflet, фильтры дат/региона/темы/источника, CRUD объектов (ПКМ), медиа, таймлайн появления (если включён в UI). + +### Nginx + +`ca-frontend` слушает `:80` (с хоста `:8080`), проксирует backend-пути на `http://ca-api:8000`. Лимит тела для медиа задаётся в `nginx.conf`. + +### Env (ЦА) + +| Переменная | Назначение | +|------------|------------| +| `DATABASE_URL` | PostgreSQL | +| `REDIS_URL` | Очереди | +| `INTERNAL_TOKEN` | Auth ЦП ↔ ЦА | +| `TEST_PI_API_KEY` | Seed consumer `test-pi` | +| `UPLOAD_DIR` | Медиа (по умолчанию `/data/uploads`) | + +### Типовые точки входа в код + +| Задача | Файл | +|--------|------| +| Новый admin endpoint | `routers/admin.py` | +| Логика ingest | `services/ingest.py` | +| Фильтры карты | `services/map_query.py` + `routers/map.py` | +| Новый экран UI | `frontend/src/views/` + `router/index.ts` | +| Поля парсера в форме | `ParsersView.vue` (+ contracts) | diff --git a/centers/analytics/DISTRIBUTION.md b/centers/analytics/DISTRIBUTION.md new file mode 100644 index 0000000..937f124 --- /dev/null +++ b/centers/analytics/DISTRIBUTION.md @@ -0,0 +1,91 @@ +# Distribution API (ПИ) + +**ПИ** — не отдельный Docker-сервис, а внешний HTTP-срез поверх данных ЦА. + +Общий контекст: [docs/architecture-overview.md](../../docs/architecture-overview.md). + +--- + +## Общими чертами + +1. В UI «ПИ» (`/consumers`) создаётся подписчик. +2. Ему выдаётся API-ключ (в БД хранится только SHA-256 hash). +3. Клиент читает события: + +```http +GET /api/v1/events +Authorization: Bearer +``` + +4. ЦА отдаёт события с учётом фильтров consumer (регионы, темы, дата). + +```mermaid +flowchart LR + Client[External system] + V1[GET /api/v1/events] + Cons[Consumer + Filter] + Events[(events)] + + Client -->|Bearer key| V1 + V1 --> Cons + Cons --> Events +``` + +--- + +## Подробнее + +### Код + +| Часть | Путь | +|-------|------| +| Роут | `centers/analytics/api/app/routers/v1.py` | +| Модели | `Consumer`, `ConsumerFilter` в `models.py` | +| Фильтрация | `services/filtering.py` | +| Admin CRUD | `routers/admin.py` — `/admin/consumers`, `…/rotate-key` | +| UI | `frontend/src/views/ConsumersView.vue` | +| Seed | `seed.py` — consumer `test-pi` из `TEST_PI_API_KEY` | + +### Аутентификация + +- Заголовок: `Authorization: Bearer `. +- Сравнение с `consumers.api_key_hash` (SHA-256). +- Неактивный consumer → отказ. + +### Фильтры подписчика + +Типичные ограничения (через `ConsumerFilter`): + +- список регионов; +- список тем; +- `date_from` (нижняя граница даты события). + +Точный набор полей — в модели/схемах admin API; при изменении фильтров обновляйте UI consumers и `filtering.py`. + +### Admin операции + +| Метод | Путь | Действие | +|-------|------|----------| +| GET | `/admin/consumers` | Список | +| POST | `/admin/consumers` | Создать (+ ключ в ответе один раз) | +| PATCH | `/admin/consumers/{id}` | Обновить | +| POST | `/admin/consumers/{id}/rotate-key` | Новый ключ | + +### Локальный тест + +После `docker compose up` (если не меняли seed): + +```bash +curl -s -H "Authorization: Bearer test-pi-api-key-change-me" \ + "http://localhost:8080/api/v1/events" | head +``` + +В проде обязательно смените `TEST_PI_API_KEY` и ротируйте ключи. + +### Что ПИ не делает + +- не пишет в БД; +- не ставит parse jobs; +- не ходит в Redis / ЦП. + +Только чтение уже ingest'нутых событий через контракт `/api/v1`. diff --git a/docker-compose.yml b/docker-compose.yml index 7587517..ac208d1 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -57,11 +57,16 @@ services: TELEGRAM_SESSION_PATH: /data/telegram.session TELEGRAM_LISTENER_ENABLED: "true" TELEGRAM_LISTENER_REFRESH_SECONDS: "60" + # Host xray is 127.0.0.1:10808; socat on host forwards 0.0.0.0:11080 → that SOCKS + TELEGRAM_PROXY_HOST: host.docker.internal + TELEGRAM_PROXY_PORT: "11080" ENABLED_ADAPTERS: telegram WORKER_FAMILIES: telegram CA_API_URL: http://ca-api:8000 REDIS_URL: redis://redis:6379/0 INTERNAL_TOKEN: dev-internal-token + extra_hosts: + - "host.docker.internal:host-gateway" volumes: - ./data:/data depends_on: diff --git a/docs/README.md b/docs/README.md new file mode 100644 index 0000000..022625b --- /dev/null +++ b/docs/README.md @@ -0,0 +1,35 @@ +# Документация MapMil для разработчиков + +Карта документов. Начните с обзора, затем углубитесь в нужный центр. + +## С чего начать + +1. [Обзор архитектуры](architecture-overview.md) — общая картина, затем детали +2. [Поток данных](data-flow.md) — от парсера до карты и ПИ +3. [Локальный запуск](local-dev.md) — `.env`, сессия, Docker + +## По зонам + +| Документ | Содержание | +|----------|------------| +| [architecture-overview.md](architecture-overview.md) | Платформа целиком: ЦА, ЦП, ПИ, границы, сервисы | +| [data-flow.md](data-flow.md) | E2E: admin → Redis → адаптер → ingest → карта → ПИ | +| [contracts.md](contracts.md) | Shared-схемы `contracts/` | +| [local-dev.md](local-dev.md) | Разработка и отладка локально | +| [ЦА: ARCHITECTURE](../centers/analytics/ARCHITECTURE.md) | API, модели, UI, scheduler | +| [ЦА: DISTRIBUTION (ПИ)](../centers/analytics/DISTRIBUTION.md) | External API для потребителей | +| [ЦП: ARCHITECTURE](../centers/parsing/ARCHITECTURE.md) | Адаптеры, очереди, Telegram listener | + +## Корень репозитория + +| Файл | Для кого | +|------|----------| +| [`README.md`](../README.md) | Быстрый старт и краткий API | +| [`AGENTS.md`](../AGENTS.md) | Ориентация для AI-агентов | +| [`.cursor/skills/add-parser-adapter/`](../.cursor/skills/add-parser-adapter/) | Чеклист нового `source_type` | + +## Принцип + +- **Contracts first** — изменение формы данных или маршрутизации начинается в `contracts/`. +- **Одна зона за раз** — не смешивать контракты + UI + инфра в одном несвязанном изменении. +- **ЦП не трогает БД ЦА** — только Redis и `POST/PATCH /internal/*`. diff --git a/docs/architecture-overview.md b/docs/architecture-overview.md new file mode 100644 index 0000000..54b2d38 --- /dev/null +++ b/docs/architecture-overview.md @@ -0,0 +1,288 @@ +# Обзор архитектуры MapMil + +Документ для разработчиков: сначала общая картина, затем устройство платформы подробнее. + +Связанные материалы: [поток данных](data-flow.md), [контракты](contracts.md), [ЦА](../centers/analytics/ARCHITECTURE.md), [ЦП](../centers/parsing/ARCHITECTURE.md), [ПИ](../centers/analytics/DISTRIBUTION.md). + +--- + +## 1. Общими чертами + +**MapMil** — monorepo платформы сбора и отображения событий: + +```text +Источники → ЦП (парсинг) → ЦА (хранение + карта + admin) → ПИ (внешние клиенты) +``` + +| Роль | Что это | Где код | +|------|---------|---------| +| **ЦА** (аналитика) | Единственный UI, PostgreSQL, ingest, admin, карта, distribution API | `centers/analytics/` | +| **ЦП** (парсинг) | Адаптеры источников без своей БД и без UI | `centers/parsing/workers/` | +| **ПИ** (потребители) | Не отдельный сервис: `GET /api/v1/events` в ЦА + ключи consumers | `centers/analytics/api/.../v1.py` | +| **Contracts** | Общие Pydantic-схемы очередей и payload | `contracts/` | + +### Главный поток + +```mermaid +flowchart LR + Admin[Admin UI / API] + Redis[(Redis queues)] + CP[CP adapters] + Src[Telegram / Web / NLP] + CA[(PostgreSQL)] + Map[Map UI] + PI[External API] + + Admin -->|enqueue job| Redis + Redis -->|BLPOP| CP + CP --> Src + CP -->|POST /internal/ingest| CA + CA --> Map + CA --> PI +``` + +### Сценарий: создание парсера и взаимодействие + +Объекты: **ЦА** (UI + API), **ЦП**, **ПИ**, **PostgreSQL**, **Redis**. Кто кому что передаёт при создании парсера и дальнейшем чтении данных. + +```mermaid +sequenceDiagram + participant UI as ЦА UI + participant CA as ЦА API + participant PG as PostgreSQL + participant R as Redis + participant CP as ЦП + participant PI as ПИ клиент + + Note over UI,CP: Создание парсера и batch-прогон + UI->>CA: POST /admin/jobs
source_type + source_config + CA->>PG: INSERT ParseJob status=queued + CA->>R: RPUSH cp:jobs:family
JobPayload job_id type config + R->>CP: BLPOP задание + CP->>CA: PATCH /internal/jobs/id
status=running + CA->>PG: UPDATE ParseJob status + CP->>CP: fetch источника
адаптер → события + CP->>CA: POST /internal/ingest
events IngestEventItem + CA->>PG: INSERT Event
дедуп по source_url + CA->>PG: sync MapObject
если есть coords + CP->>CA: PATCH job completed/failed + CA->>PG: UPDATE ParseJob + + Note over UI,PI: Чтение результата + UI->>CA: GET /api/map/objects + CA->>PG: SELECT map/events + CA-->>UI: точки на карте + + PI->>CA: GET /api/v1/events
Bearer API key + CA->>PG: SELECT events
фильтры Consumer + CA-->>PI: срез событий +``` + +| От | Кому | Что передаётся | +|----|------|----------------| +| ЦА UI → ЦА API | Конфиг парсера (`source_type`, `source_config`) | +| ЦА API → PostgreSQL | `ParseJob`, затем `Event` / `MapObject` | +| ЦА API → Redis | `JobPayload` (`job_id`, type, config) | +| Redis → ЦП | То же задание (`BLPOP`) | +| ЦП → ЦА API | Статус job + пакет событий (ingest) | +| ЦА API → ЦА UI | Объекты карты | +| ЦА API → ПИ | Отфильтрованный срез событий | + +### Жёсткие границы + +1. ЦП **никогда** не подключается к PostgreSQL ЦА. +2. Связь ЦА ↔ ЦП: Redis (`cp:jobs:{family}`) + HTTP `/internal/*` с `X-Internal-Token`. +3. Весь фронтенд — только ЦА (`ca-frontend` на порту **8080**). +4. Дедуп событий в ЦА по стабильному `source_url`. +5. Секреты (`.env`, `data/telegram.session`) не коммитятся. + +### Docker-сервисы одной командой + +| Сервис | Назначение | +|--------|------------| +| `ca-db` | PostgreSQL | +| `redis` | Очереди заданий | +| `ca-api` | FastAPI ЦА | +| `ca-frontend` | Vue admin + карта | +| `cp-workers` | Telegram (batch + listener) | +| `cp-workers-web` | Crawl4AI | +| `cp-workers-nlp` | VIINA | + +Запуск: `docker compose up --build` → http://localhost:8080 + +--- + +## 2. Подробнее + +### 2.1. Структура monorepo + +```text +MapMil/ +├── centers/ +│ ├── analytics/ # ЦА +│ │ ├── api/ # FastAPI +│ │ ├── frontend/ # Vue 3 + Leaflet +│ │ ├── ARCHITECTURE.md +│ │ └── DISTRIBUTION.md +│ └── parsing/ # ЦП +│ ├── ARCHITECTURE.md +│ └── workers/ # adapters + Telethon +├── contracts/ # ingest, jobs, sources, queues +├── docs/ # эта документация +├── data/ # telegram.session (локально) +├── docker-compose.yml +├── .env.example +├── README.md +└── AGENTS.md +``` + +### 2.2. Центр аналитики (ЦА) + +**Ответственность:** владение данными, UI, постановка jobs, приём ingest, отдача срезов ПИ. + +**Backend** (`centers/analytics/api/`): + +| Слой | Содержание | +|------|------------| +| Routers | `/api/*` (карта, объекты), `/admin/*`, `/internal/*`, `/api/v1/*` | +| Models | `Event`, `ParseJob`, `MapObject`, `ObjectMedia`, `Consumer`, `ConsumerFilter` | +| Services | enqueue → Redis, scheduler, ingest + map sync, analytics queries | +| Auth internal | `X-Internal-Token` | +| Auth ПИ | Bearer API key (SHA-256 hash в БД) | + +**Frontend** (`centers/analytics/frontend/`): + +| Маршрут | Экран | +|---------|-------| +| `/` | Карта событий и ручных объектов | +| `/parsers` | CRUD парсеров (`telegram` / `crawl4ai` / `viina`) | +| `/events` | Список событий | +| `/analytics` | KPI и таймлайны | +| `/consumers` | Подписчики ПИ | + +Nginx проксирует `/api/`, `/admin/`, `/internal/` на `ca-api:8000`. + +Подробности: [centers/analytics/ARCHITECTURE.md](../centers/analytics/ARCHITECTURE.md). + +### 2.3. Центр парсинга (ЦП) + +**Ответственность:** забрать данные из внешнего источника и вернуть список `IngestEventItem`. + +Адаптерный контракт: + +```python +async def run(job_id, source_config, *, ctx) -> tuple[list[dict], str | None] +``` + +| source_type | Очередь | Контейнер | +|-------------|---------|-----------| +| `telegram` | `cp:jobs:telegram` | `cp-workers` | +| `crawl4ai` | `cp:jobs:web` | `cp-workers-web` | +| `viina` | `cp:jobs:nlp` | `cp-workers-nlp` | + +Дополнительно для Telegram: **real-time listener** (Telethon `NewMessage` / `Album`) параллельно с batch, если `TELEGRAM_LISTENER_ENABLED=true`. + +Опционально: `extract_mode: llm` (DeepSeek) для telegram/crawl4ai в batch. + +Подробности: [centers/parsing/ARCHITECTURE.md](../centers/parsing/ARCHITECTURE.md). + +### 2.4. Distribution (ПИ) + +Отдельного центра нет. Внешний клиент: + +1. Получает API-ключ в UI «ПИ» / admin API. +2. Вызывает `GET /api/v1/events` с `Authorization: Bearer `. +3. Получает события с учётом фильтров consumer (регионы, темы, `date_from`). + +Подробности: [centers/analytics/DISTRIBUTION.md](../centers/analytics/DISTRIBUTION.md). + +### 2.5. Contracts + +Единый источник правды для формы обмена: + +| Модуль | Назначение | +|--------|------------| +| `contracts/jobs.py` | `JobPayload` в Redis | +| `contracts/ingest.py` | `IngestEventItem` / payload ingest | +| `contracts/sources.py` | Валидация `source_config` по `source_type` | +| `contracts/queues.py` | `SOURCE_FAMILY` → ключ очереди | + +Правило: меняете форму события / конфиг источника / очередь — сначала `contracts/`, потом CA/CP/UI. + +Подробности: [contracts.md](contracts.md). + +### 2.6. Данные в PostgreSQL (ядро) + +```mermaid +erDiagram + ParseJob ||--o{ Event : produces + Event ||--o| MapObject : may_create + MapObject ||--o{ ObjectMedia : has + Consumer ||--o| ConsumerFilter : has + + ParseJob { + int id + string source_type + json source_config + int interval_seconds + bool is_active + string status + } + Event { + int id + string source_url UK + string source_type + float latitude + float longitude + string region + string topic + } + MapObject { + int id + int event_id FK + float latitude + float longitude + } + Consumer { + int id + string name + string api_key_hash + } +``` + +- **ParseJob** — конфигурация парсера и статус batch-прогона. +- **Event** — нормализованное событие; уникальность по `source_url`. +- **MapObject** — точка на карте (из события с координатами или ручная). +- **Consumer** — внешний подписчик ПИ. + +### 2.7. Планировщик + +В `ca-api` фоновый scheduler (~каждые 30 с): + +- берёт активные `ParseJob` со статусом `completed` / `failed`; +- если прошло ≥ `interval_seconds` с `last_run_at` — снова ставит в Redis (`queued`). + +Так работают периодические парсеры без cron снаружи. + +### 2.8. Секреты и тома + +| Артефакт | Где | Заметка | +|----------|-----|---------| +| `.env` | корень | `TELEGRAM_*`, опционально `DEEPSEEK_*` | +| `data/telegram.session` | volume `./data` → `/data` в `cp-workers` | Telethon SQLite | +| `pgdata` | Docker volume | PostgreSQL | +| `ca_uploads` | Docker volume | медиа объектов карты | + +--- + +## 3. Куда смотреть при типовых задачах + +| Задача | Документ / код | +|--------|----------------| +| Новый источник парсинга | skill `add-parser-adapter` + `contracts/sources.py` + `queues.py` | +| Изменить поля события | `contracts/ingest.py` → ingest CA → map/ПИ | +| Карта / UI | `centers/analytics/frontend/src/` | +| Admin API | `centers/analytics/api/app/routers/admin.py` | +| Ingest / дедуп | `centers/analytics/api/app/services/ingest.py` | +| Telegram listener | `centers/parsing/.../telegram_listener.py` | +| Ключи для внешних систем | [DISTRIBUTION.md](../centers/analytics/DISTRIBUTION.md) | diff --git a/docs/contracts.md b/docs/contracts.md new file mode 100644 index 0000000..e4fda69 --- /dev/null +++ b/docs/contracts.md @@ -0,0 +1,88 @@ +# Contracts — общие схемы + +Каталог [`contracts/`](../contracts/) — shared truth для формы данных и маршрутизации между ЦА и ЦП. + +Правило платформы: **сначала contracts**, потом реализация в CA/CP/UI. + +--- + +## Модули + +| Файл | Назначение | +|------|------------| +| [`jobs.py`](../contracts/jobs.py) | Контракт задания. Payload задания в Redis | +| [`ingest.py`](../contracts/ingest.py) |Контракт результата. Событие и пакет ingest | +| [`sources.py`](../contracts/sources.py) | Контракт настроек парсера. Схемы `source_config` по `source_type` | +| [`queues.py`](../contracts/queues.py) | Контракт доставки задания нужному воркеру. `source_type` → family → Redis key | + +ЦА при admin CRUD валидирует конфиг через `parse_source_config`. +ЦП адаптеры должны отдавать dict, совместимые с `IngestEventItem`. + +> На практике CA дублирует часть DTO в `app/schemas.py` для FastAPI; при изменении формы сверяйте оба места и UI. + +--- + +## JobPayload (`jobs.py`) + +```python +class JobPayload(BaseModel): + job_id: int + source_type: str + source_config: dict = {} +``` + +Уходит в Redis как JSON. Семейство очереди выбирается по `source_type`, не по содержимому config. + +--- + +## Ingest (`ingest.py`) + +Ключевые поля `IngestEventItem`: + +| Поле | Обязательность | Заметка | +|------|----------------|---------| +| `source_url` | да | Стабильный URL; дедуп в ЦА | +| `source_type` | да (часто default) | Должен соответствовать адаптеру | +| `raw_text` / `title` / `description` | нет | Текст события | +| `latitude` / `longitude` | нет | Без них точка на карту не создаётся | +| `event_date`, `locality`, `region`, `topic` | нет | Фильтры карты и ПИ | +| `tags`, `metadata` | нет | Расширения | + +Пакет: `{ job_id?, events: [...] }`. +Listener добавляет флаг `listener: true` на стороне CA API (см. internal schemas). + +--- + +## Source config (`sources.py`) + +| source_type | Модель | Главные поля | +|-------------|--------|--------------| +| `telegram` | `TelegramSourceConfig` | `channel`, `limit`, `extract_mode` (`heuristic`\|`llm`) | +| `crawl4ai` | `Crawl4AISourceConfig` | `urls`, `extract_mode`, `extract_schema`, `domain_profile` | +| `viina` | `ViinaSourceConfig` | `urls` / `texts`, `input_mode` | + +Реестр: `CONFIG_MODELS` + `parse_source_config(source_type, raw)`. + +--- + +## Очереди (`queues.py`) + +```text +SOURCE_FAMILY = { + "telegram": "telegram", # cp:jobs:telegram + "crawl4ai": "web", # cp:jobs:web + "viina": "nlp", # cp:jobs:nlp +} +``` + +- `queue_key_for_source(source_type)` — куда enqueue из ЦА. +- Legacy-ключ `cp:jobs` ещё может дренироваться telegram-воркерами (совместимость). + +При добавлении нового `source_type` обязательно: + +1. схема в `sources.py`; +2. запись в `SOURCE_FAMILY`; +3. адаптер + registry + (при необходимости) новый Docker-сервис; +4. поля формы в `ParsersView.vue`. + +Чеклист: [`.cursor/skills/add-parser-adapter/`](../.cursor/skills/add-parser-adapter/). diff --git a/docs/data-flow.md b/docs/data-flow.md new file mode 100644 index 0000000..7955b2c --- /dev/null +++ b/docs/data-flow.md @@ -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 `. +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` | diff --git a/docs/local-dev.md b/docs/local-dev.md new file mode 100644 index 0000000..c389262 --- /dev/null +++ b/docs/local-dev.md @@ -0,0 +1,118 @@ +# Локальная разработка + +## Требования + +- Docker + Docker Compose +- Файл `.env` (из `.env.example`) +- Для Telegram: `data/telegram.session` + `TELEGRAM_API_ID` / `TELEGRAM_API_HASH` +- Опционально: `DEEPSEEK_API_KEY` для `extract_mode: llm` + +## Быстрый старт + +```bash +cp .env.example .env +# заполнить TELEGRAM_* (и при необходимости DEEPSEEK_*) + +mkdir -p data +# положить telegram.session в data/ + +docker compose up --build +``` + +UI: http://localhost:8080 +Health: `curl http://localhost:8080/api/health` + +## Сервисы + +| Сервис | Порт снаружи | Когда нужен | +|--------|--------------|-------------| +| `ca-frontend` | 8080 | всегда (UI) | +| `ca-api` | внутренний 8000 | всегда | +| `ca-db` / `redis` | — | всегда | +| `cp-workers` | — | Telegram | +| `cp-workers-web` | — | Crawl4AI | +| `cp-workers-nlp` | — | VIINA | + +Пересборка одного воркера: + +```bash +docker compose up --build cp-workers-web +``` + +Проверка compose-файла: + +```bash +docker compose config +``` + +## Telegram-сессия + +- Путь в контейнере: `/data/telegram.session` (`TELEGRAM_SESSION_PATH`) +- Локально: `./data` монтируется в `cp-workers` +- Файл **не** коммитить (`.gitignore`) +- API id/hash в `.env` должны совпадать с теми, под которыми создавалась сессия + +### Создать / обновить сессию (локальный прокси) + +1. Остановите Telegram-воркер, чтобы не делить session-файл: + `docker compose stop cp-workers` +2. В `.env`: `TELEGRAM_PROXY_TYPE=socks5`, `TELEGRAM_PROXY_HOST=127.0.0.1`, `TELEGRAM_PROXY_PORT=10808` (Happ/xray). +3. Проброс SOCKS в Docker (xray слушает только `127.0.0.1:10808`): + ```bash + socat TCP-LISTEN:11080,bind=0.0.0.0,fork,reuseaddr TCP:127.0.0.1:10808 & + ``` + В compose у `cp-workers`: `TELEGRAM_PROXY_HOST=host.docker.internal`, порт `11080`. +4. Авторизация (интерактивно, код из Telegram): + ```bash + docker run --rm -it --network host \ + --env-file .env \ + -e TELEGRAM_SESSION_PATH=/data/telegram.session \ + -e TELEGRAM_PROXY_HOST=127.0.0.1 \ + -v "$PWD/data:/data" \ + -v "$PWD/scripts:/scripts:ro" \ + mapmil-cp-workers \ + python /scripts/telegram_auth.py + ``` +5. `docker compose up -d cp-workers` + +Без сессии batch/listener Telegram не авторизуются; остальные адаптеры могут работать. + +## Полезные curl + +Создать telegram-парсер: + +```bash +curl -X POST http://localhost:8080/admin/jobs \ + -H 'Content-Type: application/json' \ + -d '{"source_type":"telegram","source_config":{"channel":"example","limit":50}}' +``` + +Срез ПИ (подставьте ключ из seed / UI): + +```bash +curl -H "Authorization: Bearer test-pi-api-key-change-me" \ + "http://localhost:8080/api/v1/events" +``` + +## Логи + +```bash +docker compose logs -f ca-api +docker compose logs -f cp-workers +docker compose logs -f cp-workers-web +``` + +## Типичные проблемы + +| Проблема | Действие | +|----------|----------| +| Frontend 000 / нет контейнеров | `docker compose up -d` (без `--build`, если registry недоступен, но образы уже есть) | +| Job не берётся | Смотреть `ENABLED_ADAPTERS` / `WORKER_FAMILIES` нужного сервиса | +| LLM не работает | `DEEPSEEK_API_KEY` в `.env`, перезапуск `cp-workers` / `cp-workers-web` | +| Изменения UI не видны | Пересобрать `ca-frontend` | + +## Границы при разработке + +- Не добавлять прямой доступ к БД из CP workers. +- Не коммитить `.env`, `data/`, API-ключи. +- Изменение формы данных — через `contracts/` (см. [contracts.md](contracts.md)). diff --git a/scripts/telegram_auth.py b/scripts/telegram_auth.py new file mode 100644 index 0000000..4c8d35b --- /dev/null +++ b/scripts/telegram_auth.py @@ -0,0 +1,115 @@ +#!/usr/bin/env python3 +"""Интерактивная авторизация Telethon → data/telegram.session. + +Запуск через образ cp-workers и host-сеть (локальный xray/Happ на 127.0.0.1:10808): + + docker compose stop cp-workers + docker run --rm -it --network host \\ + --env-file .env \\ + -e TELEGRAM_SESSION_PATH=/data/telegram.session \\ + -e TELEGRAM_PROXY_HOST=127.0.0.1 \\ + -v \"$PWD/data:/data\" \\ + -v \"$PWD/scripts:/scripts:ro\" \\ + mapmil-cp-workers \\ + python /scripts/telegram_auth.py + docker compose up -d cp-workers +""" + +from __future__ import annotations + +import asyncio +import os +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +WORKERS = ROOT / "centers" / "parsing" / "workers" +if WORKERS.exists() and str(WORKERS) not in sys.path: + sys.path.insert(0, str(WORKERS)) + + +def _load_dotenv(path: Path) -> None: + if not path.is_file(): + return + for line in path.read_text(encoding="utf-8").splitlines(): + line = line.strip() + if not line or line.startswith("#") or "=" not in line: + continue + key, _, value = line.partition("=") + key = key.strip() + value = value.strip().strip("'").strip('"') + os.environ.setdefault(key, value) + + +async def main() -> int: + _load_dotenv(ROOT / ".env") + + # Prefer repo-local session when running outside Docker /data mount. + if not os.environ.get("TELEGRAM_SESSION_PATH"): + local = ROOT / "data" / "telegram.session" + os.environ["TELEGRAM_SESSION_PATH"] = str(local) + + from workers.sources.telegram_settings import ( + create_client, + describe_connection, + get_api_credentials, + get_session_path, + ) + + api_id, api_hash = get_api_credentials() + session_path = get_session_path() + Path(session_path).parent.mkdir(parents=True, exist_ok=True) + + print(f"Сессия: {session_path}") + print(f"Подключение: {describe_connection()}") + + client = create_client(session_path, api_id, api_hash) + await client.connect() + + if await client.is_user_authorized(): + me = await client.get_me() + print( + f"Уже авторизовано: id={me.id} " + f"username={getattr(me, 'username', None) or '—'} " + f"phone={getattr(me, 'phone', None) or '—'}" + ) + await client.disconnect() + return 0 + + print("Сессия не авторизована — вход в Telegram.") + phone = input("Номер телефона (+7...): ").strip() + if not phone: + print("Номер не указан.", file=sys.stderr) + await client.disconnect() + return 1 + + await client.send_code_request(phone) + code = input("Код из Telegram/SMS: ").strip() + try: + await client.sign_in(phone=phone, code=code) + except Exception as exc: + # 2FA + from telethon.errors import SessionPasswordNeededError + + if not isinstance(exc, SessionPasswordNeededError): + raise + password = input("Пароль 2FA: ").strip() + await client.sign_in(password=password) + + me = await client.get_me() + print( + f"Готово. Авторизован: id={me.id} " + f"username={getattr(me, 'username', None) or '—'} " + f"phone={getattr(me, 'phone', None) or '—'}" + ) + print(f"Файл сессии: {session_path}") + await client.disconnect() + return 0 + + +if __name__ == "__main__": + try: + raise SystemExit(asyncio.run(main())) + except KeyboardInterrupt: + print("\nОтменено.", file=sys.stderr) + raise SystemExit(130)