commit 9dfe7136687ad8fdf23da99b342fe40d3f5d33b9 Author: gitrusprus Date: Tue Jun 30 11:43:27 2026 +0300 Add CP→CA→PI platform foundation on MapMil monorepo. Unify parsing workers, analytics API with PostgreSQL, map UI, and PI distribution into centers/ with Docker Compose. Co-authored-by: Cursor diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..902ed8f --- /dev/null +++ b/.env.example @@ -0,0 +1,13 @@ +# Telegram (CP workers) — copy from SocialParser/.env +TELEGRAM_API_ID=12345678 +TELEGRAM_API_HASH=your_api_hash_here +TELEGRAM_SESSION_PATH=/data/telegram.session + +# Optional Telegram proxy +# TELEGRAM_PROXY_TYPE=socks5 +# TELEGRAM_PROXY_HOST=127.0.0.1 +# TELEGRAM_PROXY_PORT=1080 + +# Platform internals +INTERNAL_TOKEN=dev-internal-token +TEST_PI_API_KEY=test-pi-api-key-change-me diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..037b4ae --- /dev/null +++ b/.gitignore @@ -0,0 +1,8 @@ +.env +data/telegram.session +data/*.session +__pycache__/ +*.pyc +node_modules/ +dist/ +.venv/ diff --git a/README.md b/README.md new file mode 100644 index 0000000..468e9a6 --- /dev/null +++ b/README.md @@ -0,0 +1,153 @@ +# MapMil Platform (ЦП → ЦА → ПИ) + +Единая платформа на базе MapMil (ЦА — аналитика) и SocialParser (ЦП — парсинг). + +## Архитектура + +```mermaid +flowchart LR + CA[ЦА Analytics Center] + CP[ЦП Parsing Center] + PI[ПИ External Consumers] + + CA -->|jobs via Redis| CP + CP -->|POST /internal/ingest| CA + CA -->|GET /api/v1/events| PI +``` + +| Центр | Контейнеры | Назначение | +|-------|------------|------------| +| **ЦА** | `ca-db`, `ca-api`, `ca-frontend` | PostgreSQL, ingest API, карта, distribution API | +| **ЦП** | `cp-workers` | Парсинг Telegram (код из SocialParser) | +| **Общее** | `redis` | Очередь заданий ЦА → ЦП | + +## Структура monorepo + +``` +MapMil/ +├── centers/ +│ ├── analytics/ +│ │ ├── api/ # CA backend (FastAPI + PostgreSQL) +│ │ └── frontend/ # CA admin UI (Vue + Leaflet) +│ └── parsing/ +│ └── workers/ # CP workers (Telethon) +├── contracts/ # Shared schemas (ingest, jobs) +├── data/ # telegram.session (symlink → SocialParser) +├── docker-compose.yml +└── .env # TELEGRAM_* из SocialParser +``` + +## Быстрый старт + +1. Скопируйте `.env` из SocialParser (или создайте из `.env.example`): + +```bash +cp ../SocialParser/.env .env +``` + +2. Убедитесь, что сессия Telegram доступна: + +```bash +ls -la data/telegram.session +# symlink → ../SocialParser/data/telegram.session +``` + +3. Запуск: + +```bash +docker compose up --build +``` + +4. Откройте карту: [http://localhost:8080](http://localhost:8080) + +## Сохранение Telegram-сессии + +**Важно:** существующий файл сессии **не удаляется и не пересоздаётся**. + +- Оригинал: `SocialParser/data/telegram.session` +- В репозитории MapMil: `data/telegram.session` — симлинк для локальной разработки +- В Docker `cp-workers`: файл монтируется напрямую как `../SocialParser/data/telegram.session:/data/telegram.session:ro` +- Путь в контейнере: `/data/telegram.session` +- Переменные из `.env`: `TELEGRAM_API_ID`, `TELEGRAM_API_HASH`, `TELEGRAM_SESSION_PATH=/data/telegram.session` + +> Симлинк не работает внутри Docker — compose монтирует исходный файл из SocialParser. + +Файл сессии и `.env` добавлены в `.gitignore` и не коммитятся. + +## API + +### Карта (совместимость MapMil) + +| Метод | Путь | Описание | +|-------|------|----------| +| GET | `/api/health` | Health check | +| GET | `/api/objects` | Объекты на карте (включая события с координатами) | +| POST | `/api/objects` | Создать объект вручную | + +### ЦА Admin + +| Метод | Путь | Описание | +|-------|------|----------| +| POST | `/admin/jobs` | Создать задание парсинга (ставится в Redis) | +| GET | `/admin/jobs` | Список заданий | +| GET | `/admin/events` | Все события | +| POST | `/admin/consumers` | Создать подписчика ПИ | + +Пример задания Telegram: + +```bash +curl -X POST http://localhost:8080/admin/jobs \ + -H 'Content-Type: application/json' \ + -d '{"source_type":"telegram","source_config":{"channel":"creamy_caprice","limit":50}}' +``` + +### ЦП → ЦА (internal) + +| Метод | Путь | Описание | +|-------|------|----------| +| POST | `/internal/ingest` | Приём batch событий (заголовок `X-Internal-Token`) | + +### ПИ Distribution API + +| Метод | Путь | Описание | +|-------|------|----------| +| GET | `/api/v1/events` | События с фильтром по API-ключу | + +Тестовый consumer `test-pi` создаётся при старте с ключом из `TEST_PI_API_KEY` (по умолчанию `test-pi-api-key-change-me`): + +```bash +curl http://localhost:8080/api/v1/events \ + -H 'Authorization: Bearer test-pi-api-key-change-me' +``` + +## Поток данных + +1. Аналитик создаёт задание: `POST /admin/jobs` +2. `ca-api` ставит задание в Redis (`cp:jobs`) +3. `cp-workers` забирает задание, парсит Telegram через существующую сессию +4. Результаты отправляются в `POST /internal/ingest` +5. События с координатами автоматически появляются на карте как `MapObject` +6. Внешние ПИ получают отфильтрованный срез через `/api/v1/events` + +## Миграция EventRecord → Event + +| SocialParser | CA Event | +|--------------|----------| +| `event` | `description` / `title` | +| `date` (dd.mm.yy) | `event_date` | +| `geolocation` | `latitude`, `longitude` | +| `locality` | `locality`, `region` | +| `source_url` | `source_url` | +| — | `source_type = "telegram"` | + +## Остановка + +```bash +docker compose down +``` + +Данные PostgreSQL сохраняются в volume `pgdata`. + +## Legacy + +Старые каталоги `backend/` и `frontend/` в корне оставлены для справки; активная разработка — в `centers/`. diff --git a/backend/Dockerfile b/backend/Dockerfile new file mode 100644 index 0000000..fcdf862 --- /dev/null +++ b/backend/Dockerfile @@ -0,0 +1,14 @@ +FROM python:3.12-slim + +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY app/ ./app/ + +RUN mkdir -p /data + +EXPOSE 8000 + +CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/backend/app/database.py b/backend/app/database.py new file mode 100644 index 0000000..5dbc49a --- /dev/null +++ b/backend/app/database.py @@ -0,0 +1,23 @@ +import os +from sqlalchemy import create_engine +from sqlalchemy.orm import DeclarativeBase, sessionmaker + +DATABASE_URL = os.getenv("DATABASE_URL", "sqlite:////data/mapmil.db") + +engine = create_engine( + DATABASE_URL, + connect_args={"check_same_thread": False}, +) +SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine) + + +class Base(DeclarativeBase): + pass + + +def get_db(): + db = SessionLocal() + try: + yield db + finally: + db.close() diff --git a/backend/app/main.py b/backend/app/main.py new file mode 100644 index 0000000..5ffb9e8 --- /dev/null +++ b/backend/app/main.py @@ -0,0 +1,195 @@ +from contextlib import asynccontextmanager + +from fastapi import Depends, FastAPI, File, HTTPException, UploadFile +from fastapi.middleware.cors import CORSMiddleware +from fastapi.responses import FileResponse +from sqlalchemy.orm import Session + +from .database import Base, engine, get_db +from .models import MapObject, ObjectMedia +from .schemas import MapObjectCreate, MapObjectRead, MapObjectUpdate, ObjectMediaRead +from .seed import seed_objects +from .storage import ( + ALLOWED_CONTENT_TYPES, + MAX_FILE_SIZE, + build_stored_name, + ensure_upload_dir, + is_allowed_content_type, + media_file_path, + remove_media_file, +) + + +def media_to_read(media: ObjectMedia) -> ObjectMediaRead: + return ObjectMediaRead( + id=media.id, + object_id=media.object_id, + original_name=media.original_name, + content_type=media.content_type, + size=media.size, + created_at=media.created_at, + url=f"/api/media/{media.id}/file", + ) + + +@asynccontextmanager +async def lifespan(_: FastAPI): + ensure_upload_dir() + Base.metadata.create_all(bind=engine) + db = next(get_db()) + try: + seed_objects(db) + finally: + db.close() + yield + + +app = FastAPI(title="MapMil API", lifespan=lifespan) + +app.add_middleware( + CORSMiddleware, + allow_origins=["*"], + allow_credentials=True, + allow_methods=["*"], + allow_headers=["*"], +) + + +@app.get("/api/health") +def health(): + return {"status": "ok"} + + +@app.get("/api/objects", response_model=list[MapObjectRead]) +def list_objects(db: Session = Depends(get_db)): + return db.query(MapObject).order_by(MapObject.id).all() + + +@app.get("/api/objects/{object_id}", response_model=MapObjectRead) +def get_object(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + return obj + + +@app.post("/api/objects", response_model=MapObjectRead, status_code=201) +def create_object(payload: MapObjectCreate, db: Session = Depends(get_db)): + data = payload.model_dump() + created_at = data.pop("created_at", None) + obj = MapObject(**data) + if created_at is not None: + obj.created_at = created_at + db.add(obj) + db.commit() + db.refresh(obj) + return obj + + +@app.patch("/api/objects/{object_id}", response_model=MapObjectRead) +def update_object( + object_id: int, + payload: MapObjectUpdate, + db: Session = Depends(get_db), +): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + for field, value in payload.model_dump(exclude_unset=True).items(): + setattr(obj, field, value) + + db.commit() + db.refresh(obj) + return obj + + +@app.delete("/api/objects/{object_id}", status_code=204) +def delete_object(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + for media in obj.media: + remove_media_file(media.stored_name) + + db.delete(obj) + db.commit() + + +@app.get("/api/objects/{object_id}/media", response_model=list[ObjectMediaRead]) +def list_object_media(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + media_items = ( + db.query(ObjectMedia) + .filter(ObjectMedia.object_id == object_id) + .order_by(ObjectMedia.id) + .all() + ) + return [media_to_read(item) for item in media_items] + + +@app.post("/api/objects/{object_id}/media", response_model=ObjectMediaRead, status_code=201) +async def upload_object_media( + object_id: int, + file: UploadFile = File(...), + db: Session = Depends(get_db), +): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + content_type = file.content_type or "application/octet-stream" + if not is_allowed_content_type(content_type): + raise HTTPException( + status_code=400, + detail=f"Неподдерживаемый тип файла. Разрешены: {', '.join(sorted(ALLOWED_CONTENT_TYPES))}", + ) + + content = await file.read() + if len(content) > MAX_FILE_SIZE: + raise HTTPException(status_code=400, detail="Файл слишком большой (макс. 20 МБ)") + + original_name = file.filename or "file" + stored_name = build_stored_name(original_name) + path = media_file_path(stored_name) + path.write_bytes(content) + + media = ObjectMedia( + object_id=object_id, + stored_name=stored_name, + original_name=original_name, + content_type=content_type, + size=len(content), + ) + db.add(media) + db.commit() + db.refresh(media) + return media_to_read(media) + + +@app.get("/api/media/{media_id}/file") +def get_media_file(media_id: int, db: Session = Depends(get_db)): + media = db.query(ObjectMedia).filter(ObjectMedia.id == media_id).first() + if not media: + raise HTTPException(status_code=404, detail="Медиафайл не найден") + + path = media_file_path(media.stored_name) + if not path.exists(): + raise HTTPException(status_code=404, detail="Файл не найден на диске") + + return FileResponse(path, media_type=media.content_type, filename=media.original_name) + + +@app.delete("/api/media/{media_id}", status_code=204) +def delete_media(media_id: int, db: Session = Depends(get_db)): + media = db.query(ObjectMedia).filter(ObjectMedia.id == media_id).first() + if not media: + raise HTTPException(status_code=404, detail="Медиафайл не найден") + + remove_media_file(media.stored_name) + db.delete(media) + db.commit() diff --git a/backend/app/models.py b/backend/app/models.py new file mode 100644 index 0000000..6b9efac --- /dev/null +++ b/backend/app/models.py @@ -0,0 +1,47 @@ +from datetime import datetime, timezone + +from sqlalchemy import DateTime, Float, ForeignKey, Integer, String, Text +from sqlalchemy.orm import Mapped, mapped_column, relationship + +from .database import Base + + +class MapObject(Base): + __tablename__ = "map_objects" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + name: Mapped[str] = mapped_column(String(255), nullable=False) + description: Mapped[str] = mapped_column(Text, default="") + type: Mapped[str] = mapped_column(String(50), nullable=False) + latitude: Mapped[float] = mapped_column(Float, nullable=False) + longitude: Mapped[float] = mapped_column(Float, nullable=False) + created_at: Mapped[datetime] = mapped_column( + DateTime, + default=lambda: datetime.now(timezone.utc), + ) + + media: Mapped[list["ObjectMedia"]] = relationship( + back_populates="object", + cascade="all, delete-orphan", + ) + + +class ObjectMedia(Base): + __tablename__ = "object_media" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + object_id: Mapped[int] = mapped_column( + ForeignKey("map_objects.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + stored_name: Mapped[str] = mapped_column(String(255), nullable=False) + original_name: Mapped[str] = mapped_column(String(255), nullable=False) + content_type: Mapped[str] = mapped_column(String(100), nullable=False) + size: Mapped[int] = mapped_column(Integer, nullable=False) + created_at: Mapped[datetime] = mapped_column( + DateTime, + default=lambda: datetime.now(timezone.utc), + ) + + object: Mapped["MapObject"] = relationship(back_populates="media") diff --git a/backend/app/schemas.py b/backend/app/schemas.py new file mode 100644 index 0000000..dc97ac2 --- /dev/null +++ b/backend/app/schemas.py @@ -0,0 +1,48 @@ +from datetime import datetime +from typing import Literal + +from pydantic import BaseModel, ConfigDict, Field + +ObjectType = Literal["point", "marker", "zone", "other"] + + +class MapObjectCreate(BaseModel): + name: str = Field(min_length=1, max_length=255) + description: str = "" + type: ObjectType + latitude: float = Field(ge=-90, le=90) + longitude: float = Field(ge=-180, le=180) + created_at: datetime | None = None + + +class MapObjectRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + name: str + description: str + type: ObjectType + latitude: float + longitude: float + created_at: datetime + + +class MapObjectUpdate(BaseModel): + name: str | None = Field(default=None, min_length=1, max_length=255) + description: str | None = None + type: ObjectType | None = None + latitude: float | None = Field(default=None, ge=-90, le=90) + longitude: float | None = Field(default=None, ge=-180, le=180) + created_at: datetime | None = None + + +class ObjectMediaRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + object_id: int + original_name: str + content_type: str + size: int + created_at: datetime + url: str diff --git a/backend/app/seed.py b/backend/app/seed.py new file mode 100644 index 0000000..96cb58c --- /dev/null +++ b/backend/app/seed.py @@ -0,0 +1,54 @@ +from sqlalchemy.orm import Session + +from .models import MapObject + +from datetime import datetime, timezone + +from sqlalchemy.orm import Session + +from .models import MapObject + +SEED_OBJECTS = [ + { + "name": "Красная площадь", + "description": "Главная площадь Москвы, исторический центр города.", + "type": "marker", + "latitude": 55.7539, + "longitude": 37.6208, + "created_at": datetime(2018, 5, 9, 12, 0, tzinfo=timezone.utc), + }, + { + "name": "ВДНХ", + "description": "Выставка достижений народного хозяйства — крупный выставочный комплекс.", + "type": "zone", + "latitude": 55.8298, + "longitude": 37.6361, + "created_at": datetime(2020, 8, 15, 10, 30, tzinfo=timezone.utc), + }, + { + "name": "МГУ", + "description": "Московский государственный университет имени М.В. Ломоносова.", + "type": "point", + "latitude": 55.7033, + "longitude": 37.5307, + "created_at": datetime(2022, 2, 12, 14, 0, tzinfo=timezone.utc), + }, + { + "name": "Парк Горького", + "description": "Центральный парк культуры и отдыха имени М. Горького.", + "type": "other", + "latitude": 55.7312, + "longitude": 37.6013, + "created_at": datetime(2024, 6, 1, 9, 0, tzinfo=timezone.utc), + }, +] + + +def seed_objects(db: Session) -> None: + if db.query(MapObject).count() > 0: + return + + for item in SEED_OBJECTS: + db.add(MapObject(**item)) + + db.commit() diff --git a/backend/app/storage.py b/backend/app/storage.py new file mode 100644 index 0000000..99271e8 --- /dev/null +++ b/backend/app/storage.py @@ -0,0 +1,39 @@ +import os +import uuid +from pathlib import Path + +UPLOAD_DIR = Path(os.getenv("UPLOAD_DIR", "/data/uploads")) + +ALLOWED_CONTENT_TYPES = { + "image/jpeg", + "image/png", + "image/gif", + "image/webp", + "video/mp4", + "video/webm", +} + +MAX_FILE_SIZE = 20 * 1024 * 1024 + + +def ensure_upload_dir() -> None: + UPLOAD_DIR.mkdir(parents=True, exist_ok=True) + + +def is_allowed_content_type(content_type: str) -> bool: + return content_type in ALLOWED_CONTENT_TYPES + + +def build_stored_name(original_name: str) -> str: + suffix = Path(original_name).suffix.lower() + return f"{uuid.uuid4().hex}{suffix}" + + +def media_file_path(stored_name: str) -> Path: + return UPLOAD_DIR / stored_name + + +def remove_media_file(stored_name: str) -> None: + path = media_file_path(stored_name) + if path.exists(): + path.unlink() diff --git a/backend/requirements.txt b/backend/requirements.txt new file mode 100644 index 0000000..677e9c2 --- /dev/null +++ b/backend/requirements.txt @@ -0,0 +1,5 @@ +fastapi==0.115.6 +uvicorn[standard]==0.34.0 +sqlalchemy==2.0.36 +pydantic==2.10.3 +python-multipart==0.0.20 diff --git a/centers/analytics/api/Dockerfile b/centers/analytics/api/Dockerfile new file mode 100644 index 0000000..fcdf862 --- /dev/null +++ b/centers/analytics/api/Dockerfile @@ -0,0 +1,14 @@ +FROM python:3.12-slim + +WORKDIR /app + +COPY requirements.txt . +RUN pip install --no-cache-dir -r requirements.txt + +COPY app/ ./app/ + +RUN mkdir -p /data + +EXPOSE 8000 + +CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/centers/analytics/api/app/database.py b/centers/analytics/api/app/database.py new file mode 100644 index 0000000..94cc6dc --- /dev/null +++ b/centers/analytics/api/app/database.py @@ -0,0 +1,28 @@ +import os + +from sqlalchemy import create_engine +from sqlalchemy.orm import DeclarativeBase, sessionmaker + +DATABASE_URL = os.getenv( + "DATABASE_URL", + "postgresql://mapmil:mapmil@ca-db:5432/mapmil", +) + +connect_args: dict = {} +if DATABASE_URL.startswith("sqlite"): + connect_args = {"check_same_thread": False} + +engine = create_engine(DATABASE_URL, connect_args=connect_args) +SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine) + + +class Base(DeclarativeBase): + pass + + +def get_db(): + db = SessionLocal() + try: + yield db + finally: + db.close() diff --git a/centers/analytics/api/app/deps.py b/centers/analytics/api/app/deps.py new file mode 100644 index 0000000..1a68038 --- /dev/null +++ b/centers/analytics/api/app/deps.py @@ -0,0 +1,11 @@ +import os + +from fastapi import Header, HTTPException + + +def verify_internal_token( + x_internal_token: str | None = Header(default=None, alias="X-Internal-Token"), +) -> None: + expected = os.getenv("INTERNAL_TOKEN", "dev-internal-token") + if not x_internal_token or x_internal_token != expected: + raise HTTPException(status_code=401, detail="Invalid internal token") diff --git a/centers/analytics/api/app/main.py b/centers/analytics/api/app/main.py new file mode 100644 index 0000000..f031e21 --- /dev/null +++ b/centers/analytics/api/app/main.py @@ -0,0 +1,38 @@ +from contextlib import asynccontextmanager + +from fastapi import FastAPI +from fastapi.middleware.cors import CORSMiddleware + +from .database import Base, engine, get_db +from .routers import admin, internal, objects, v1 +from .seed import seed_objects, seed_test_consumer +from .storage import ensure_upload_dir + + +@asynccontextmanager +async def lifespan(_: FastAPI): + ensure_upload_dir() + Base.metadata.create_all(bind=engine) + db = next(get_db()) + try: + seed_objects(db) + seed_test_consumer(db) + finally: + db.close() + yield + + +app = FastAPI(title="CA API (Analytics Center)", lifespan=lifespan) + +app.add_middleware( + CORSMiddleware, + allow_origins=["*"], + allow_credentials=True, + allow_methods=["*"], + allow_headers=["*"], +) + +app.include_router(objects.router) +app.include_router(internal.router) +app.include_router(admin.router) +app.include_router(v1.router) diff --git a/centers/analytics/api/app/models.py b/centers/analytics/api/app/models.py new file mode 100644 index 0000000..a9466cf --- /dev/null +++ b/centers/analytics/api/app/models.py @@ -0,0 +1,145 @@ +from datetime import datetime, timezone + +from sqlalchemy import ( + JSON, + Boolean, + DateTime, + Float, + ForeignKey, + Integer, + String, + Text, + UniqueConstraint, +) +from sqlalchemy.orm import Mapped, mapped_column, relationship + +from .database import Base + + +class MapObject(Base): + __tablename__ = "map_objects" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + name: Mapped[str] = mapped_column(String(255), nullable=False) + description: Mapped[str] = mapped_column(Text, default="") + type: Mapped[str] = mapped_column(String(50), nullable=False) + latitude: Mapped[float] = mapped_column(Float, nullable=False) + longitude: Mapped[float] = mapped_column(Float, nullable=False) + event_id: Mapped[int | None] = mapped_column( + ForeignKey("events.id", ondelete="SET NULL"), + nullable=True, + index=True, + ) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + ) + + media: Mapped[list["ObjectMedia"]] = relationship( + back_populates="object", + cascade="all, delete-orphan", + ) + event: Mapped["Event | None"] = relationship(back_populates="map_object") + + +class ObjectMedia(Base): + __tablename__ = "object_media" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + object_id: Mapped[int] = mapped_column( + ForeignKey("map_objects.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + stored_name: Mapped[str] = mapped_column(String(255), nullable=False) + original_name: Mapped[str] = mapped_column(String(255), nullable=False) + content_type: Mapped[str] = mapped_column(String(100), nullable=False) + size: Mapped[int] = mapped_column(Integer, nullable=False) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + ) + + object: Mapped["MapObject"] = relationship(back_populates="media") + + +class Event(Base): + __tablename__ = "events" + __table_args__ = (UniqueConstraint("source_url", name="uq_events_source_url"),) + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + source_type: Mapped[str] = mapped_column(String(50), nullable=False, index=True) + source_url: Mapped[str] = mapped_column(String(512), nullable=False) + raw_text: Mapped[str] = mapped_column(Text, default="") + title: Mapped[str] = mapped_column(String(512), default="") + description: Mapped[str] = mapped_column(Text, default="") + locality: Mapped[str] = mapped_column(String(255), default="") + latitude: Mapped[float | None] = mapped_column(Float, nullable=True) + longitude: Mapped[float | None] = mapped_column(Float, nullable=True) + event_date: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + ingested_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + index=True, + ) + region: Mapped[str | None] = mapped_column(String(100), nullable=True, index=True) + topic: Mapped[str | None] = mapped_column(String(100), nullable=True, index=True) + tags: Mapped[list | None] = mapped_column(JSON, nullable=True) + metadata_: Mapped[dict | None] = mapped_column("metadata", JSON, nullable=True) + + map_object: Mapped["MapObject | None"] = relationship( + back_populates="event", + uselist=False, + ) + + +class ParseJob(Base): + __tablename__ = "parse_jobs" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + source_type: Mapped[str] = mapped_column(String(50), nullable=False) + source_config: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict) + schedule: Mapped[str | None] = mapped_column(String(100), nullable=True) + status: Mapped[str] = mapped_column(String(50), default="pending", index=True) + last_run_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + last_error: Mapped[str | None] = mapped_column(Text, nullable=True) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + ) + + +class Consumer(Base): + __tablename__ = "consumers" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + name: Mapped[str] = mapped_column(String(255), nullable=False, unique=True) + api_key_hash: Mapped[str] = mapped_column(String(64), nullable=False, unique=True) + is_active: Mapped[bool] = mapped_column(Boolean, default=True) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=lambda: datetime.now(timezone.utc), + ) + + filter: Mapped["ConsumerFilter | None"] = relationship( + back_populates="consumer", + uselist=False, + cascade="all, delete-orphan", + ) + + +class ConsumerFilter(Base): + __tablename__ = "consumer_filters" + + id: Mapped[int] = mapped_column(Integer, primary_key=True, index=True) + consumer_id: Mapped[int] = mapped_column( + ForeignKey("consumers.id", ondelete="CASCADE"), + nullable=False, + unique=True, + ) + regions: Mapped[list | None] = mapped_column(JSON, nullable=True) + topics: Mapped[list | None] = mapped_column(JSON, nullable=True) + min_access_level: Mapped[int] = mapped_column(Integer, default=0) + date_from: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + + consumer: Mapped["Consumer"] = relationship(back_populates="filter") diff --git a/centers/analytics/api/app/routers/__init__.py b/centers/analytics/api/app/routers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/centers/analytics/api/app/routers/admin.py b/centers/analytics/api/app/routers/admin.py new file mode 100644 index 0000000..acb5217 --- /dev/null +++ b/centers/analytics/api/app/routers/admin.py @@ -0,0 +1,78 @@ +from fastapi import APIRouter, Depends, HTTPException, Query +from sqlalchemy.orm import Session + +from ..database import get_db +from ..models import Event, ParseJob +from ..schemas import ( + ConsumerCreate, + ConsumerRead, + EventRead, + ParseJobCreate, + ParseJobRead, +) +from ..services.filtering import create_consumer +from ..services.jobs import enqueue_job + +router = APIRouter(prefix="/admin", tags=["admin"]) + + +@router.post("/jobs", response_model=ParseJobRead, status_code=201) +def create_parse_job(payload: ParseJobCreate, db: Session = Depends(get_db)): + job = ParseJob( + source_type=payload.source_type, + source_config=payload.source_config, + schedule=payload.schedule, + status="queued", + ) + db.add(job) + db.commit() + db.refresh(job) + + enqueue_job(job.id, job.source_type, job.source_config) + return job + + +@router.get("/jobs", response_model=list[ParseJobRead]) +def list_parse_jobs(db: Session = Depends(get_db)): + return db.query(ParseJob).order_by(ParseJob.id.desc()).all() + + +@router.get("/jobs/{job_id}", response_model=ParseJobRead) +def get_parse_job(job_id: int, db: Session = Depends(get_db)): + job = db.query(ParseJob).filter(ParseJob.id == job_id).first() + if not job: + raise HTTPException(status_code=404, detail="Job not found") + return job + + +@router.get("/events", response_model=list[EventRead]) +def list_events( + limit: int = Query(default=100, ge=1, le=1000), + offset: int = Query(default=0, ge=0), + db: Session = Depends(get_db), +): + return ( + db.query(Event) + .order_by(Event.ingested_at.desc()) + .offset(offset) + .limit(limit) + .all() + ) + + +@router.post("/consumers", response_model=ConsumerRead, status_code=201) +def create_pi_consumer(payload: ConsumerCreate, db: Session = Depends(get_db)): + consumer, api_key = create_consumer( + db, + name=payload.name, + regions=payload.regions, + topics=payload.topics, + date_from=payload.date_from, + ) + return ConsumerRead( + id=consumer.id, + name=consumer.name, + is_active=consumer.is_active, + created_at=consumer.created_at, + api_key=api_key, + ) diff --git a/centers/analytics/api/app/routers/internal.py b/centers/analytics/api/app/routers/internal.py new file mode 100644 index 0000000..5e0c116 --- /dev/null +++ b/centers/analytics/api/app/routers/internal.py @@ -0,0 +1,45 @@ +from datetime import datetime, timezone + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session + +from ..database import get_db +from ..deps import verify_internal_token +from ..models import ParseJob +from ..schemas import IngestRequest, IngestResponse +from ..services.ingest import ingest_events + +router = APIRouter(prefix="/internal", tags=["internal"]) + + +@router.post("/ingest", response_model=IngestResponse) +def internal_ingest( + payload: IngestRequest, + _: None = Depends(verify_internal_token), + db: Session = Depends(get_db), +): + ingested, updated, map_synced = ingest_events(db, payload.events, payload.job_id) + return IngestResponse( + ingested=ingested, + updated=updated, + map_objects_synced=map_synced, + ) + + +@router.patch("/jobs/{job_id}") +def update_job_status( + job_id: int, + status: str, + error: str | None = None, + _: None = Depends(verify_internal_token), + db: Session = Depends(get_db), +): + job = db.query(ParseJob).filter(ParseJob.id == job_id).first() + if not job: + raise HTTPException(status_code=404, detail="Job not found") + + job.status = status + job.last_run_at = datetime.now(timezone.utc) + job.last_error = error + db.commit() + return {"id": job.id, "status": job.status} diff --git a/centers/analytics/api/app/routers/objects.py b/centers/analytics/api/app/routers/objects.py new file mode 100644 index 0000000..c201c74 --- /dev/null +++ b/centers/analytics/api/app/routers/objects.py @@ -0,0 +1,169 @@ +from fastapi import APIRouter, Depends, File, HTTPException, UploadFile +from fastapi.responses import FileResponse +from sqlalchemy.orm import Session + +from ..database import get_db +from ..models import MapObject, ObjectMedia +from ..schemas import MapObjectCreate, MapObjectRead, MapObjectUpdate, ObjectMediaRead +from ..storage import ( + ALLOWED_CONTENT_TYPES, + MAX_FILE_SIZE, + build_stored_name, + is_allowed_content_type, + media_file_path, + remove_media_file, +) + +router = APIRouter(tags=["objects"]) + + +def media_to_read(media: ObjectMedia) -> ObjectMediaRead: + return ObjectMediaRead( + id=media.id, + object_id=media.object_id, + original_name=media.original_name, + content_type=media.content_type, + size=media.size, + created_at=media.created_at, + url=f"/api/media/{media.id}/file", + ) + + +@router.get("/api/health") +def health(): + return {"status": "ok"} + + +@router.get("/api/objects", response_model=list[MapObjectRead]) +def list_objects(db: Session = Depends(get_db)): + return db.query(MapObject).order_by(MapObject.id).all() + + +@router.get("/api/objects/{object_id}", response_model=MapObjectRead) +def get_object(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + return obj + + +@router.post("/api/objects", response_model=MapObjectRead, status_code=201) +def create_object(payload: MapObjectCreate, db: Session = Depends(get_db)): + data = payload.model_dump() + created_at = data.pop("created_at", None) + obj = MapObject(**data) + if created_at is not None: + obj.created_at = created_at + db.add(obj) + db.commit() + db.refresh(obj) + return obj + + +@router.patch("/api/objects/{object_id}", response_model=MapObjectRead) +def update_object( + object_id: int, + payload: MapObjectUpdate, + db: Session = Depends(get_db), +): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + for field, value in payload.model_dump(exclude_unset=True).items(): + setattr(obj, field, value) + + db.commit() + db.refresh(obj) + return obj + + +@router.delete("/api/objects/{object_id}", status_code=204) +def delete_object(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + for media in obj.media: + remove_media_file(media.stored_name) + + db.delete(obj) + db.commit() + + +@router.get("/api/objects/{object_id}/media", response_model=list[ObjectMediaRead]) +def list_object_media(object_id: int, db: Session = Depends(get_db)): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + media_items = ( + db.query(ObjectMedia) + .filter(ObjectMedia.object_id == object_id) + .order_by(ObjectMedia.id) + .all() + ) + return [media_to_read(item) for item in media_items] + + +@router.post("/api/objects/{object_id}/media", response_model=ObjectMediaRead, status_code=201) +async def upload_object_media( + object_id: int, + file: UploadFile = File(...), + db: Session = Depends(get_db), +): + obj = db.query(MapObject).filter(MapObject.id == object_id).first() + if not obj: + raise HTTPException(status_code=404, detail="Объект не найден") + + content_type = file.content_type or "application/octet-stream" + if not is_allowed_content_type(content_type): + raise HTTPException( + status_code=400, + detail=f"Неподдерживаемый тип файла. Разрешены: {', '.join(sorted(ALLOWED_CONTENT_TYPES))}", + ) + + content = await file.read() + if len(content) > MAX_FILE_SIZE: + raise HTTPException(status_code=400, detail="Файл слишком большой (макс. 20 МБ)") + + original_name = file.filename or "file" + stored_name = build_stored_name(original_name) + path = media_file_path(stored_name) + path.write_bytes(content) + + media = ObjectMedia( + object_id=object_id, + stored_name=stored_name, + original_name=original_name, + content_type=content_type, + size=len(content), + ) + db.add(media) + db.commit() + db.refresh(media) + return media_to_read(media) + + +@router.get("/api/media/{media_id}/file") +def get_media_file(media_id: int, db: Session = Depends(get_db)): + media = db.query(ObjectMedia).filter(ObjectMedia.id == media_id).first() + if not media: + raise HTTPException(status_code=404, detail="Медиафайл не найден") + + path = media_file_path(media.stored_name) + if not path.exists(): + raise HTTPException(status_code=404, detail="Файл не найден на диске") + + return FileResponse(path, media_type=media.content_type, filename=media.original_name) + + +@router.delete("/api/media/{media_id}", status_code=204) +def delete_media(media_id: int, db: Session = Depends(get_db)): + media = db.query(ObjectMedia).filter(ObjectMedia.id == media_id).first() + if not media: + raise HTTPException(status_code=404, detail="Медиафайл не найден") + + remove_media_file(media.stored_name) + db.delete(media) + db.commit() diff --git a/centers/analytics/api/app/routers/v1.py b/centers/analytics/api/app/routers/v1.py new file mode 100644 index 0000000..e99c9dc --- /dev/null +++ b/centers/analytics/api/app/routers/v1.py @@ -0,0 +1,33 @@ +from fastapi import APIRouter, Depends, Header, HTTPException, Query +from sqlalchemy.orm import Session + +from ..database import get_db +from ..schemas import EventRead +from ..services.filtering import apply_consumer_filter, find_consumer_by_api_key + +router = APIRouter(prefix="/api/v1", tags=["distribution"]) + + +def get_consumer_from_api_key( + authorization: str | None = Header(default=None), + db: Session = Depends(get_db), +): + if not authorization or not authorization.startswith("Bearer "): + raise HTTPException(status_code=401, detail="Missing or invalid Authorization header") + + api_key = authorization.removeprefix("Bearer ").strip() + consumer = find_consumer_by_api_key(db, api_key) + if not consumer: + raise HTTPException(status_code=401, detail="Invalid API key") + return consumer + + +@router.get("/events", response_model=list[EventRead]) +def list_filtered_events( + limit: int = Query(default=100, ge=1, le=1000), + offset: int = Query(default=0, ge=0), + consumer=Depends(get_consumer_from_api_key), + db: Session = Depends(get_db), +): + events = apply_consumer_filter(db, consumer) + return events[offset : offset + limit] diff --git a/centers/analytics/api/app/schemas.py b/centers/analytics/api/app/schemas.py new file mode 100644 index 0000000..7aeaca6 --- /dev/null +++ b/centers/analytics/api/app/schemas.py @@ -0,0 +1,132 @@ +from datetime import datetime +from typing import Any, Literal + +from pydantic import BaseModel, ConfigDict, Field + +ObjectType = Literal["point", "marker", "zone", "other"] + + +class MapObjectCreate(BaseModel): + name: str = Field(min_length=1, max_length=255) + description: str = "" + type: ObjectType + latitude: float = Field(ge=-90, le=90) + longitude: float = Field(ge=-180, le=180) + created_at: datetime | None = None + + +class MapObjectRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + name: str + description: str + type: ObjectType + latitude: float + longitude: float + event_id: int | None = None + created_at: datetime + + +class MapObjectUpdate(BaseModel): + name: str | None = Field(default=None, min_length=1, max_length=255) + description: str | None = None + type: ObjectType | None = None + latitude: float | None = Field(default=None, ge=-90, le=90) + longitude: float | None = Field(default=None, ge=-180, le=180) + created_at: datetime | None = None + + +class ObjectMediaRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + object_id: int + original_name: str + content_type: str + size: int + created_at: datetime + url: str + + +class EventRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + source_type: str + source_url: str + raw_text: str + title: str + description: str + locality: str + latitude: float | None + longitude: float | None + event_date: datetime | None + ingested_at: datetime + region: str | None + topic: str | None + tags: list[str] | None + metadata: dict[str, Any] | None = Field(validation_alias="metadata_") + + +class IngestEventItem(BaseModel): + source_type: str = "telegram" + source_url: str + raw_text: str = "" + title: str = "" + description: str = "" + locality: str = "" + latitude: float | None = None + longitude: float | None = None + event_date: datetime | None = None + region: str | None = None + topic: str | None = None + tags: list[str] | None = None + metadata: dict[str, Any] | None = None + + +class IngestRequest(BaseModel): + job_id: int | None = None + events: list[IngestEventItem] = Field(default_factory=list) + + +class IngestResponse(BaseModel): + ingested: int + updated: int + map_objects_synced: int + + +class ParseJobCreate(BaseModel): + source_type: str = "telegram" + source_config: dict[str, Any] = Field(default_factory=dict) + schedule: str | None = None + + +class ParseJobRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + source_type: str + source_config: dict[str, Any] + schedule: str | None + status: str + last_run_at: datetime | None + last_error: str | None + created_at: datetime + + +class ConsumerCreate(BaseModel): + name: str + regions: list[str] | None = None + topics: list[str] | None = None + date_from: datetime | None = None + + +class ConsumerRead(BaseModel): + model_config = ConfigDict(from_attributes=True) + + id: int + name: str + is_active: bool + created_at: datetime + api_key: str | None = None diff --git a/centers/analytics/api/app/seed.py b/centers/analytics/api/app/seed.py new file mode 100644 index 0000000..75ea34d --- /dev/null +++ b/centers/analytics/api/app/seed.py @@ -0,0 +1,68 @@ +import os +from datetime import datetime, timezone + +from sqlalchemy.orm import Session + +from .models import Consumer, ConsumerFilter, MapObject +from .services.filtering import hash_api_key + +SEED_OBJECTS = [ + { + "name": "Красная площадь", + "description": "Главная площадь Москвы, исторический центр города.", + "type": "marker", + "latitude": 55.7539, + "longitude": 37.6208, + "created_at": datetime(2018, 5, 9, 12, 0, tzinfo=timezone.utc), + }, + { + "name": "ВДНХ", + "description": "Выставка достижений народного хозяйства — крупный выставочный комплекс.", + "type": "zone", + "latitude": 55.8298, + "longitude": 37.6361, + "created_at": datetime(2020, 8, 15, 10, 30, tzinfo=timezone.utc), + }, + { + "name": "МГУ", + "description": "Московский государственный университет имени М.В. Ломоносова.", + "type": "point", + "latitude": 55.7033, + "longitude": 37.5307, + "created_at": datetime(2022, 2, 12, 14, 0, tzinfo=timezone.utc), + }, + { + "name": "Парк Горького", + "description": "Центральный парк культуры и отдыха имени М. Горького.", + "type": "other", + "latitude": 55.7312, + "longitude": 37.6013, + "created_at": datetime(2024, 6, 1, 9, 0, tzinfo=timezone.utc), + }, +] + + +def seed_objects(db: Session) -> None: + if db.query(MapObject).count() > 0: + return + + for item in SEED_OBJECTS: + db.add(MapObject(**item)) + + db.commit() + + +def seed_test_consumer(db: Session) -> None: + if db.query(Consumer).filter(Consumer.name == "test-pi").first(): + return + + api_key = os.getenv("TEST_PI_API_KEY", "test-pi-api-key-change-me") + consumer = Consumer( + name="test-pi", + api_key_hash=hash_api_key(api_key), + is_active=True, + ) + db.add(consumer) + db.flush() + db.add(ConsumerFilter(consumer_id=consumer.id, regions=None, topics=None)) + db.commit() diff --git a/centers/analytics/api/app/services/__init__.py b/centers/analytics/api/app/services/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/centers/analytics/api/app/services/filtering.py b/centers/analytics/api/app/services/filtering.py new file mode 100644 index 0000000..be3e236 --- /dev/null +++ b/centers/analytics/api/app/services/filtering.py @@ -0,0 +1,109 @@ +import hashlib +import secrets +from datetime import datetime + +from sqlalchemy.orm import Session + +from ..models import Consumer, ConsumerFilter, Event, MapObject + + +def hash_api_key(api_key: str) -> str: + return hashlib.sha256(api_key.encode()).hexdigest() + + +def generate_api_key() -> str: + return secrets.token_urlsafe(32) + + +def create_consumer( + db: Session, + name: str, + regions: list[str] | None = None, + topics: list[str] | None = None, + date_from: datetime | None = None, +) -> tuple[Consumer, str]: + api_key = generate_api_key() + consumer = Consumer(name=name, api_key_hash=hash_api_key(api_key), is_active=True) + db.add(consumer) + db.flush() + + consumer_filter = ConsumerFilter( + consumer_id=consumer.id, + regions=regions, + topics=topics, + date_from=date_from, + ) + db.add(consumer_filter) + db.commit() + db.refresh(consumer) + return consumer, api_key + + +def find_consumer_by_api_key(db: Session, api_key: str) -> Consumer | None: + key_hash = hash_api_key(api_key) + return ( + db.query(Consumer) + .filter(Consumer.api_key_hash == key_hash, Consumer.is_active.is_(True)) + .first() + ) + + +def apply_consumer_filter(db: Session, consumer: Consumer) -> list[Event]: + query = db.query(Event).order_by(Event.ingested_at.desc()) + consumer_filter = consumer.filter + + if consumer_filter: + if consumer_filter.regions: + query = query.filter(Event.region.in_(consumer_filter.regions)) + if consumer_filter.topics: + query = query.filter(Event.topic.in_(consumer_filter.topics)) + if consumer_filter.date_from: + query = query.filter(Event.event_date >= consumer_filter.date_from) + + return query.all() + + +def sync_event_to_map_object(db: Session, event: Event) -> MapObject | None: + if event.latitude is None or event.longitude is None: + return None + + name = event.title or event.locality or f"Событие #{event.id}" + description = event.description or event.raw_text or "" + + existing = ( + db.query(MapObject) + .filter(MapObject.event_id == event.id) + .first() + ) + if existing: + existing.name = name + existing.description = description + existing.latitude = event.latitude + existing.longitude = event.longitude + return existing + + by_url = ( + db.query(MapObject) + .join(Event, MapObject.event_id == Event.id) + .filter(Event.source_url == event.source_url) + .first() + ) + if by_url: + by_url.event_id = event.id + by_url.name = name + by_url.description = description + by_url.latitude = event.latitude + by_url.longitude = event.longitude + return by_url + + obj = MapObject( + name=name, + description=description, + type="marker", + latitude=event.latitude, + longitude=event.longitude, + event_id=event.id, + created_at=event.event_date or event.ingested_at, + ) + db.add(obj) + return obj diff --git a/centers/analytics/api/app/services/ingest.py b/centers/analytics/api/app/services/ingest.py new file mode 100644 index 0000000..63e1576 --- /dev/null +++ b/centers/analytics/api/app/services/ingest.py @@ -0,0 +1,72 @@ +from sqlalchemy.orm import Session + +from ..models import Event, ParseJob +from ..schemas import IngestEventItem +from .filtering import sync_event_to_map_object + + +def ingest_events( + db: Session, + items: list[IngestEventItem], + job_id: int | None = None, +) -> tuple[int, int, int]: + ingested = 0 + updated = 0 + map_synced = 0 + + for item in items: + existing = ( + db.query(Event) + .filter(Event.source_url == item.source_url) + .first() + ) + + if existing: + existing.source_type = item.source_type + existing.raw_text = item.raw_text + existing.title = item.title + existing.description = item.description + existing.locality = item.locality + existing.latitude = item.latitude + existing.longitude = item.longitude + existing.event_date = item.event_date + existing.region = item.region + existing.topic = item.topic + existing.tags = item.tags + existing.metadata_ = item.metadata + event = existing + updated += 1 + else: + event = Event( + source_type=item.source_type, + source_url=item.source_url, + raw_text=item.raw_text, + title=item.title, + description=item.description, + locality=item.locality, + latitude=item.latitude, + longitude=item.longitude, + event_date=item.event_date, + region=item.region, + topic=item.topic, + tags=item.tags, + metadata_=item.metadata, + ) + db.add(event) + ingested += 1 + + db.flush() + if sync_event_to_map_object(db, event): + map_synced += 1 + + if job_id is not None: + job = db.query(ParseJob).filter(ParseJob.id == job_id).first() + if job: + from datetime import datetime, timezone + + job.status = "completed" + job.last_run_at = datetime.now(timezone.utc) + job.last_error = None + + db.commit() + return ingested, updated, map_synced diff --git a/centers/analytics/api/app/services/jobs.py b/centers/analytics/api/app/services/jobs.py new file mode 100644 index 0000000..9945a33 --- /dev/null +++ b/centers/analytics/api/app/services/jobs.py @@ -0,0 +1,20 @@ +import json +import os + +import redis + +REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0") +JOB_QUEUE_KEY = "cp:jobs" + + +def get_redis() -> redis.Redis: + return redis.from_url(REDIS_URL, decode_responses=True) + + +def enqueue_job(job_id: int, source_type: str, source_config: dict) -> None: + payload = { + "job_id": job_id, + "source_type": source_type, + "source_config": source_config, + } + get_redis().rpush(JOB_QUEUE_KEY, json.dumps(payload)) diff --git a/centers/analytics/api/app/storage.py b/centers/analytics/api/app/storage.py new file mode 100644 index 0000000..99271e8 --- /dev/null +++ b/centers/analytics/api/app/storage.py @@ -0,0 +1,39 @@ +import os +import uuid +from pathlib import Path + +UPLOAD_DIR = Path(os.getenv("UPLOAD_DIR", "/data/uploads")) + +ALLOWED_CONTENT_TYPES = { + "image/jpeg", + "image/png", + "image/gif", + "image/webp", + "video/mp4", + "video/webm", +} + +MAX_FILE_SIZE = 20 * 1024 * 1024 + + +def ensure_upload_dir() -> None: + UPLOAD_DIR.mkdir(parents=True, exist_ok=True) + + +def is_allowed_content_type(content_type: str) -> bool: + return content_type in ALLOWED_CONTENT_TYPES + + +def build_stored_name(original_name: str) -> str: + suffix = Path(original_name).suffix.lower() + return f"{uuid.uuid4().hex}{suffix}" + + +def media_file_path(stored_name: str) -> Path: + return UPLOAD_DIR / stored_name + + +def remove_media_file(stored_name: str) -> None: + path = media_file_path(stored_name) + if path.exists(): + path.unlink() diff --git a/centers/analytics/api/requirements.txt b/centers/analytics/api/requirements.txt new file mode 100644 index 0000000..d44b3a0 --- /dev/null +++ b/centers/analytics/api/requirements.txt @@ -0,0 +1,7 @@ +fastapi==0.115.6 +uvicorn[standard]==0.34.0 +sqlalchemy==2.0.36 +pydantic==2.10.3 +python-multipart==0.0.20 +psycopg2-binary==2.9.10 +redis==5.2.1 diff --git a/centers/analytics/frontend/Dockerfile b/centers/analytics/frontend/Dockerfile new file mode 100644 index 0000000..54c360c --- /dev/null +++ b/centers/analytics/frontend/Dockerfile @@ -0,0 +1,18 @@ +FROM node:22-alpine AS build + +WORKDIR /app + +COPY package.json ./ +RUN npm install + +COPY . . +RUN npm run build + +FROM nginx:alpine + +COPY nginx.conf /etc/nginx/conf.d/default.conf +COPY --from=build /app/dist /usr/share/nginx/html + +EXPOSE 80 + +CMD ["nginx", "-g", "daemon off;"] diff --git a/centers/analytics/frontend/index.html b/centers/analytics/frontend/index.html new file mode 100644 index 0000000..42b6d3c --- /dev/null +++ b/centers/analytics/frontend/index.html @@ -0,0 +1,18 @@ + + + + + + MapMil + + + +
+ + + diff --git a/centers/analytics/frontend/nginx.conf b/centers/analytics/frontend/nginx.conf new file mode 100644 index 0000000..d293c77 --- /dev/null +++ b/centers/analytics/frontend/nginx.conf @@ -0,0 +1,35 @@ +server { + listen 80; + server_name localhost; + root /usr/share/nginx/html; + index index.html; + + location /api/ { + client_max_body_size 50M; + proxy_pass http://ca-api:8000/api/; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + + location /admin/ { + proxy_pass http://ca-api:8000/admin/; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + + location /internal/ { + proxy_pass http://ca-api:8000/internal/; + proxy_set_header Host $host; + proxy_set_header X-Real-IP $remote_addr; + proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; + proxy_set_header X-Forwarded-Proto $scheme; + } + + location / { + try_files $uri $uri/ /index.html; + } +} diff --git a/centers/analytics/frontend/package.json b/centers/analytics/frontend/package.json new file mode 100644 index 0000000..a0c2ffd --- /dev/null +++ b/centers/analytics/frontend/package.json @@ -0,0 +1,22 @@ +{ + "name": "mapmil-frontend", + "private": true, + "version": "1.0.0", + "type": "module", + "scripts": { + "dev": "vite", + "build": "vite build", + "preview": "vite preview" + }, + "dependencies": { + "leaflet": "^1.9.4", + "vue": "^3.5.13" + }, + "devDependencies": { + "@types/leaflet": "^1.9.15", + "@vitejs/plugin-vue": "^5.2.1", + "typescript": "~5.6.3", + "vite": "^6.0.3", + "vue-tsc": "^2.1.10" + } +} diff --git a/centers/analytics/frontend/src/App.vue b/centers/analytics/frontend/src/App.vue new file mode 100644 index 0000000..42b19e2 --- /dev/null +++ b/centers/analytics/frontend/src/App.vue @@ -0,0 +1,349 @@ + + + + + diff --git a/centers/analytics/frontend/src/api/objects.ts b/centers/analytics/frontend/src/api/objects.ts new file mode 100644 index 0000000..e03ff13 --- /dev/null +++ b/centers/analytics/frontend/src/api/objects.ts @@ -0,0 +1,82 @@ +import type { MapObject, MapObjectCreate, MapObjectUpdate, ObjectMedia } from "../types/object"; + +const API_BASE = "/api"; + +async function request(url: string, options?: RequestInit): Promise { + const headers = new Headers(options?.headers); + const isFormData = options?.body instanceof FormData; + + if (!isFormData && !headers.has("Content-Type")) { + headers.set("Content-Type", "application/json"); + } + + const response = await fetch(`${API_BASE}${url}`, { + ...options, + headers, + }); + + if (!response.ok) { + const message = await response.text(); + throw new Error(message || `Ошибка запроса: ${response.status}`); + } + + if (response.status === 204) { + return undefined as T; + } + + return response.json() as Promise; +} + +export function fetchObjects(): Promise { + return request("/objects"); +} + +export function fetchObject(id: number): Promise { + return request(`/objects/${id}`); +} + +export function createObject(payload: MapObjectCreate): Promise { + return request("/objects", { + method: "POST", + body: JSON.stringify(payload), + }); +} + +export function updateObject(id: number, payload: MapObjectUpdate): Promise { + return request(`/objects/${id}`, { + method: "PATCH", + body: JSON.stringify(payload), + }); +} + +export function deleteObject(id: number): Promise { + return request(`/objects/${id}`, { + method: "DELETE", + }); +} + +export function fetchObjectMedia(objectId: number): Promise { + return request(`/objects/${objectId}/media`); +} + +export function uploadObjectMedia(objectId: number, file: File): Promise { + const form = new FormData(); + form.append("file", file); + + return request(`/objects/${objectId}/media`, { + method: "POST", + body: form, + }); +} + +export function deleteObjectMedia(mediaId: number): Promise { + return request(`/media/${mediaId}`, { + method: "DELETE", + }); +} + +export function formatFileSize(bytes: number): string { + if (bytes < 1024) return `${bytes} Б`; + if (bytes < 1024 * 1024) return `${(bytes / 1024).toFixed(1)} КБ`; + return `${(bytes / (1024 * 1024)).toFixed(1)} МБ`; +} diff --git a/centers/analytics/frontend/src/components/ContextMenu.vue b/centers/analytics/frontend/src/components/ContextMenu.vue new file mode 100644 index 0000000..424fb31 --- /dev/null +++ b/centers/analytics/frontend/src/components/ContextMenu.vue @@ -0,0 +1,81 @@ + + + + + diff --git a/centers/analytics/frontend/src/components/CreateObjectModal.vue b/centers/analytics/frontend/src/components/CreateObjectModal.vue new file mode 100644 index 0000000..7827832 --- /dev/null +++ b/centers/analytics/frontend/src/components/CreateObjectModal.vue @@ -0,0 +1,253 @@ + + +