Files
MapMil/docs/architecture-overview.md
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

12 KiB
Raw Permalink Blame History

Обзор архитектуры MapMil

Документ для разработчиков: сначала общая картина, затем устройство платформы подробнее.

Связанные материалы: поток данных, контракты, ЦА, ЦП, ПИ.


1. Общими чертами

MapMil — monorepo платформы сбора и отображения событий:

Источники  →  ЦП (парсинг)  →  ЦА (хранение + карта + 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/

Главный поток

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. Кто кому что передаёт при создании парсера и дальнейшем чтении данных.

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<br/>source_type + source_config
  CA->>PG: INSERT ParseJob status=queued
  CA->>R: RPUSH cp:jobs:family<br/>JobPayload job_id type config
  R->>CP: BLPOP задание
  CP->>CA: PATCH /internal/jobs/id<br/>status=running
  CA->>PG: UPDATE ParseJob status
  CP->>CP: fetch источника<br/>адаптер → события
  CP->>CA: POST /internal/ingest<br/>events IngestEventItem
  CA->>PG: INSERT Event<br/>дедуп по source_url
  CA->>PG: sync MapObject<br/>если есть 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<br/>Bearer API key
  CA->>PG: SELECT events<br/>фильтры 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

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.

2.3. Центр парсинга (ЦП)

Ответственность: забрать данные из внешнего источника и вернуть список IngestEventItem.

Адаптерный контракт:

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.

Профиль + канал: UI /parser-profiles (CRUD, Generate/Preview для heuristic; instruction/schema/Preview для llm) и /channels; связка на /parsers → ParseJob с FK. При enqueue ЦА flatten по ParserProfile.kind:

  • heuristic → extract_mode=profile + heuristic_profile (runtime без LLM; listener поддерживает);
  • llm → extract_mode=llm + extract_schema / instruction / required_fields (DeepSeek на каждый пост в batch; listener пока fallback на heuristic).

Целевые поля — фиксированная схема Event/IngestEventItem. Кастомные пользовательские таблицы — roadmap.

Подробности: centers/parsing/ARCHITECTURE.md, centers/analytics/ARCHITECTURE.md.

2.4. Distribution (ПИ)

Отдельного центра нет. Внешний клиент:

  1. Получает API-ключ в UI «ПИ» / admin API.
  2. Вызывает GET /api/v1/events с Authorization: Bearer <key>.
  3. Получает события с учётом фильтров consumer (регионы, темы, date_from).

Подробности: 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/heuristic_profile.py Статичный профиль конструктора + apply_profile
contracts/llm_profile.py Reusable LLM instruction/schema + match_llm_required

Правило: меняете форму события / конфиг источника / очередь — сначала contracts/, потом CA/CP/UI.

Подробности: contracts.md.

2.6. Данные в PostgreSQL (ядро)

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