From 8fbabd3c11245c68a16dfd9a3a0aa457324bea86 Mon Sep 17 00:00:00 2001 From: gitrusprus Date: Fri, 14 Aug 2026 11:34:28 +0300 Subject: [PATCH] Add CP source adapter registry with multi-worker queues and LLM extract. Replace legacy root backend/frontend with Telegram, Crawl4AI, and VIINA adapters routed by Redis job families. Co-authored-by: Cursor --- .env.example | 9 + .gitignore | 1 + README.md | 23 +- backend/Dockerfile | 14 - backend/app/database.py | 23 -- backend/app/main.py | 195 ---------- backend/app/models.py | 47 --- backend/app/schemas.py | 48 --- backend/app/seed.py | 54 --- backend/app/storage.py | 39 -- backend/requirements.txt | 5 - centers/analytics/api/Dockerfile | 7 +- centers/analytics/api/app/main.py | 9 + centers/analytics/api/app/routers/admin.py | 16 +- centers/analytics/api/app/services/jobs.py | 26 +- .../frontend/src/views/ParsersView.vue | 318 ++++++++++++++-- centers/parsing/ARCHITECTURE.md | 146 ++++++++ centers/parsing/workers/Dockerfile | 15 +- centers/parsing/workers/requirements-nlp.txt | 5 + centers/parsing/workers/requirements-web.txt | 7 + centers/parsing/workers/requirements.txt | 1 + centers/parsing/workers/worker.py | 127 ++++--- .../workers/workers/adapters/__init__.py | 5 + .../parsing/workers/workers/adapters/base.py | 28 ++ .../workers/adapters/crawl4ai_adapter.py | 286 ++++++++++++++ .../workers/workers/adapters/registry.py | 89 +++++ .../workers/workers/adapters/telegram.py | 103 ++++++ .../parsing/workers/workers/adapters/viina.py | 155 ++++++++ .../parsing/workers/workers/llm_extract.py | 179 +++++++++ contracts/__init__.py | 1 + contracts/jobs.py | 5 + contracts/queues.py | 32 ++ contracts/sources.py | 86 +++++ docker-compose.yml | 53 ++- frontend/Dockerfile | 18 - frontend/index.html | 18 - frontend/nginx.conf | 19 - frontend/package.json | 22 -- frontend/src/App.vue | 349 ------------------ frontend/src/api/objects.ts | 82 ---- frontend/src/components/ContextMenu.vue | 81 ---- frontend/src/components/CreateObjectModal.vue | 253 ------------- frontend/src/components/EditObjectModal.vue | 239 ------------ frontend/src/components/MapView.vue | 232 ------------ frontend/src/components/ObjectPanel.vue | 343 ----------------- frontend/src/components/TimelineBar.vue | 199 ---------- frontend/src/main.ts | 5 - frontend/src/style.css | 22 -- frontend/src/types/object.ts | 50 --- frontend/src/vite-env.d.ts | 7 - frontend/tsconfig.app.json | 8 - frontend/tsconfig.json | 21 -- frontend/tsconfig.node.json | 10 - frontend/vite.config.ts | 11 - 54 files changed, 1625 insertions(+), 2521 deletions(-) delete mode 100644 backend/Dockerfile delete mode 100644 backend/app/database.py delete mode 100644 backend/app/main.py delete mode 100644 backend/app/models.py delete mode 100644 backend/app/schemas.py delete mode 100644 backend/app/seed.py delete mode 100644 backend/app/storage.py delete mode 100644 backend/requirements.txt create mode 100644 centers/parsing/ARCHITECTURE.md create mode 100644 centers/parsing/workers/requirements-nlp.txt create mode 100644 centers/parsing/workers/requirements-web.txt create mode 100644 centers/parsing/workers/workers/adapters/__init__.py create mode 100644 centers/parsing/workers/workers/adapters/base.py create mode 100644 centers/parsing/workers/workers/adapters/crawl4ai_adapter.py create mode 100644 centers/parsing/workers/workers/adapters/registry.py create mode 100644 centers/parsing/workers/workers/adapters/telegram.py create mode 100644 centers/parsing/workers/workers/adapters/viina.py create mode 100644 centers/parsing/workers/workers/llm_extract.py create mode 100644 contracts/queues.py create mode 100644 contracts/sources.py delete mode 100644 frontend/Dockerfile delete mode 100644 frontend/index.html delete mode 100644 frontend/nginx.conf delete mode 100644 frontend/package.json delete mode 100644 frontend/src/App.vue delete mode 100644 frontend/src/api/objects.ts delete mode 100644 frontend/src/components/ContextMenu.vue delete mode 100644 frontend/src/components/CreateObjectModal.vue delete mode 100644 frontend/src/components/EditObjectModal.vue delete mode 100644 frontend/src/components/MapView.vue delete mode 100644 frontend/src/components/ObjectPanel.vue delete mode 100644 frontend/src/components/TimelineBar.vue delete mode 100644 frontend/src/main.ts delete mode 100644 frontend/src/style.css delete mode 100644 frontend/src/types/object.ts delete mode 100644 frontend/src/vite-env.d.ts delete mode 100644 frontend/tsconfig.app.json delete mode 100644 frontend/tsconfig.json delete mode 100644 frontend/tsconfig.node.json delete mode 100644 frontend/vite.config.ts diff --git a/.env.example b/.env.example index c47a2a7..e49efc1 100644 --- a/.env.example +++ b/.env.example @@ -11,3 +11,12 @@ TELEGRAM_SESSION_PATH=/data/telegram.session # Platform internals INTERNAL_TOKEN=dev-internal-token TEST_PI_API_KEY=test-pi-api-key-change-me + +# DeepSeek LLM for extract_mode=llm (Telegram unstructured + Crawl4AI LLM) +# DEEPSEEK_API_KEY=sk-... +# DEEPSEEK_BASE_URL=https://api.deepseek.com +# DEEPSEEK_MODEL=deepseek-chat + +# CP adapter workers (set in docker-compose; override locally if needed) +# ENABLED_ADAPTERS=telegram +# WORKER_FAMILIES=telegram diff --git a/.gitignore b/.gitignore index b124d43..8a1d2bb 100644 --- a/.gitignore +++ b/.gitignore @@ -1,5 +1,6 @@ .env data/ +mapmil-secrets.tar.gz __pycache__/ *.pyc node_modules/ diff --git a/README.md b/README.md index b8099de..9bdc62d 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # MapMil Platform (ЦП → ЦА → ПИ) -Единая платформа: ЦА (аналитика и карта) + ЦП (парсинг Telegram). +Единая платформа: ЦА (аналитика и карта) + ЦП (адаптеры парсинга: Telegram, Crawl4AI, VIINA). ## Архитектура @@ -10,7 +10,7 @@ flowchart LR CP[ЦП Parsing Center] PI[ПИ External Consumers] - CA -->|jobs via Redis| CP + CA -->|jobs via Redis families| CP CP -->|POST /internal/ingest| CA CA -->|GET /api/v1/events| PI ``` @@ -18,8 +18,10 @@ flowchart LR | Центр | Контейнеры | Назначение | |-------|------------|------------| | **ЦА** | `ca-db`, `ca-api`, `ca-frontend` | PostgreSQL, ingest API, карта, distribution API | -| **ЦП** | `cp-workers` | Парсинг Telegram: real-time listener (Telethon) + batch-задания из Redis | -| **Общее** | `redis` | Очередь заданий ЦА → ЦП | +| **ЦП** | `cp-workers`, `cp-workers-web`, `cp-workers-nlp` | Адаптеры: Telegram / Crawl4AI / VIINA | +| **Общее** | `redis` | Очереди `cp:jobs:{telegram\|web\|nlp}` | + +Подробности ЦП: [`centers/parsing/ARCHITECTURE.md`](centers/parsing/ARCHITECTURE.md). ## Структура monorepo @@ -30,11 +32,12 @@ MapMil/ │ │ ├── api/ # CA backend (FastAPI + PostgreSQL) │ │ └── frontend/ # CA admin UI (Vue + Leaflet) │ └── parsing/ -│ └── workers/ # CP workers (Telethon) -├── contracts/ # Shared schemas (ingest, jobs) +│ ├── ARCHITECTURE.md +│ └── workers/ # CP workers + adapters +├── contracts/ # Shared schemas (ingest, jobs, sources, queues) ├── data/ # telegram.session (локально, не в git) ├── docker-compose.yml -└── .env # TELEGRAM_API_ID, TELEGRAM_API_HASH, … +└── .env ``` ## Быстрый старт @@ -67,7 +70,7 @@ docker compose up --build | Раздел | Путь | Описание | |--------|------|----------| | **Карта** | `/` | Интерактивная карта событий: навигация по датам (flatpickr), пресеты периода, фильтры региона/темы/источника, подложки Яндекс/OSM/Topo/ESRI, линейка, полноэкранный режим, центрирование по координатам и городам, поиск населённых пунктов (Nominatim). CRUD для ручных объектов (ПКМ). Поддерживает `?eventId=` | -| **Парсеры** | `/parsers` | Telegram-парсеры с периодическим запуском, настройка интервала, редактирование и удаление; дубликаты по `source_url` не записываются | +| **Парсеры** | `/parsers` | Адаптеры `telegram` / `crawl4ai` / `viina`, интервал, CRUD; дедуп по `source_url` | | **События** | `/events` | Фильтрация, пагинация, просмотр деталей, ссылка «На карте» для событий с координатами | | **Аналитика** | `/analytics` | KPI-карточки, график динамики ingest за 30 дней, топ населённых пунктов и регионов | | **ПИ** | `/consumers` | CRUD подписчиков distribution API, ротация ключей, тест среза через `/api/v1/events` | @@ -198,7 +201,3 @@ docker compose down ``` Данные PostgreSQL сохраняются в volume `pgdata`. - -## Legacy - -Старые каталоги `backend/` и `frontend/` в корне оставлены для справки; активная разработка — в `centers/`. diff --git a/backend/Dockerfile b/backend/Dockerfile deleted file mode 100644 index fcdf862..0000000 --- a/backend/Dockerfile +++ /dev/null @@ -1,14 +0,0 @@ -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 deleted file mode 100644 index 5dbc49a..0000000 --- a/backend/app/database.py +++ /dev/null @@ -1,23 +0,0 @@ -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 deleted file mode 100644 index 5ffb9e8..0000000 --- a/backend/app/main.py +++ /dev/null @@ -1,195 +0,0 @@ -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 deleted file mode 100644 index 6b9efac..0000000 --- a/backend/app/models.py +++ /dev/null @@ -1,47 +0,0 @@ -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 deleted file mode 100644 index dc97ac2..0000000 --- a/backend/app/schemas.py +++ /dev/null @@ -1,48 +0,0 @@ -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 deleted file mode 100644 index 96cb58c..0000000 --- a/backend/app/seed.py +++ /dev/null @@ -1,54 +0,0 @@ -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 deleted file mode 100644 index 99271e8..0000000 --- a/backend/app/storage.py +++ /dev/null @@ -1,39 +0,0 @@ -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 deleted file mode 100644 index 677e9c2..0000000 --- a/backend/requirements.txt +++ /dev/null @@ -1,5 +0,0 @@ -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 index fcdf862..32533cb 100644 --- a/centers/analytics/api/Dockerfile +++ b/centers/analytics/api/Dockerfile @@ -2,13 +2,16 @@ FROM python:3.12-slim WORKDIR /app -COPY requirements.txt . +COPY centers/analytics/api/requirements.txt . RUN pip install --no-cache-dir -r requirements.txt -COPY app/ ./app/ +COPY contracts/ ./contracts/ +COPY centers/analytics/api/app/ ./app/ RUN mkdir -p /data +ENV PYTHONPATH=/app + EXPOSE 8000 CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"] diff --git a/centers/analytics/api/app/main.py b/centers/analytics/api/app/main.py index 606d912..68c80f3 100644 --- a/centers/analytics/api/app/main.py +++ b/centers/analytics/api/app/main.py @@ -1,8 +1,17 @@ from contextlib import asynccontextmanager +import sys +from pathlib import Path from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware +# Monorepo / Docker: contracts live next to app or at repo root +_HERE = Path(__file__).resolve() +for _candidate in (_HERE.parent, *_HERE.parents): + if (_candidate / "contracts").is_dir() and str(_candidate) not in sys.path: + sys.path.insert(0, str(_candidate)) + break + from .database import Base, engine, get_db from .routers import admin, internal, map, objects, v1 from .seed import seed_objects, seed_test_consumer diff --git a/centers/analytics/api/app/routers/admin.py b/centers/analytics/api/app/routers/admin.py index e38bb70..9f4abde 100644 --- a/centers/analytics/api/app/routers/admin.py +++ b/centers/analytics/api/app/routers/admin.py @@ -37,11 +37,23 @@ from ..services.jobs import enqueue_job router = APIRouter(prefix="/admin", tags=["admin"]) +def _validated_source_config(source_type: str, source_config: dict) -> dict: + try: + from contracts.queues import family_for_source + from contracts.sources import parse_source_config + + family_for_source(source_type) + return parse_source_config(source_type, source_config or {}).model_dump() + except Exception as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + @router.post("/jobs", response_model=ParseJobRead, status_code=201) def create_parse_job(payload: ParseJobCreate, db: Session = Depends(get_db)): + source_config = _validated_source_config(payload.source_type, payload.source_config) job = ParseJob( source_type=payload.source_type, - source_config=payload.source_config, + source_config=source_config, schedule=payload.schedule, interval_seconds=payload.interval_seconds, is_active=payload.is_active, @@ -96,7 +108,7 @@ def update_parse_job( raise HTTPException(status_code=404, detail="Job not found") if payload.source_config is not None: - job.source_config = payload.source_config + job.source_config = _validated_source_config(job.source_type, payload.source_config) if payload.interval_seconds is not None: job.interval_seconds = payload.interval_seconds if payload.is_active is not None: diff --git a/centers/analytics/api/app/services/jobs.py b/centers/analytics/api/app/services/jobs.py index 9945a33..69187f2 100644 --- a/centers/analytics/api/app/services/jobs.py +++ b/centers/analytics/api/app/services/jobs.py @@ -1,10 +1,22 @@ import json import os +import sys +from pathlib import Path import redis +# Allow importing shared contracts from monorepo root in local runs +_HERE = Path(__file__).resolve() +for _candidate in (_HERE.parent, *_HERE.parents): + if (_candidate / "contracts").is_dir(): + if str(_candidate) not in sys.path: + sys.path.insert(0, str(_candidate)) + break + +from contracts.queues import queue_key_for_source # noqa: E402 +from contracts.sources import parse_source_config # noqa: E402 + REDIS_URL = os.getenv("REDIS_URL", "redis://redis:6379/0") -JOB_QUEUE_KEY = "cp:jobs" def get_redis() -> redis.Redis: @@ -12,9 +24,17 @@ def get_redis() -> redis.Redis: def enqueue_job(job_id: int, source_type: str, source_config: dict) -> None: + # Validate known configs early; unknown types still raise from queue_key_for_source + try: + parse_source_config(source_type, source_config or {}) + except Exception: + # Keep enqueue resilient for legacy rows; worker validates again + pass + payload = { "job_id": job_id, "source_type": source_type, - "source_config": source_config, + "source_config": source_config or {}, } - get_redis().rpush(JOB_QUEUE_KEY, json.dumps(payload)) + key = queue_key_for_source(source_type) + get_redis().rpush(key, json.dumps(payload)) diff --git a/centers/analytics/frontend/src/views/ParsersView.vue b/centers/analytics/frontend/src/views/ParsersView.vue index 3c1d392..0d2a547 100644 --- a/centers/analytics/frontend/src/views/ParsersView.vue +++ b/centers/analytics/frontend/src/views/ParsersView.vue @@ -1,5 +1,5 @@ - - diff --git a/frontend/nginx.conf b/frontend/nginx.conf deleted file mode 100644 index e82ec61..0000000 --- a/frontend/nginx.conf +++ /dev/null @@ -1,19 +0,0 @@ -server { - listen 80; - server_name localhost; - root /usr/share/nginx/html; - index index.html; - - location /api/ { - client_max_body_size 50M; - proxy_pass http://backend: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 / { - try_files $uri $uri/ /index.html; - } -} diff --git a/frontend/package.json b/frontend/package.json deleted file mode 100644 index a0c2ffd..0000000 --- a/frontend/package.json +++ /dev/null @@ -1,22 +0,0 @@ -{ - "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/frontend/src/App.vue b/frontend/src/App.vue deleted file mode 100644 index 42b19e2..0000000 --- a/frontend/src/App.vue +++ /dev/null @@ -1,349 +0,0 @@ - - - - - diff --git a/frontend/src/api/objects.ts b/frontend/src/api/objects.ts deleted file mode 100644 index e03ff13..0000000 --- a/frontend/src/api/objects.ts +++ /dev/null @@ -1,82 +0,0 @@ -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/frontend/src/components/ContextMenu.vue b/frontend/src/components/ContextMenu.vue deleted file mode 100644 index 424fb31..0000000 --- a/frontend/src/components/ContextMenu.vue +++ /dev/null @@ -1,81 +0,0 @@ - - - - - diff --git a/frontend/src/components/CreateObjectModal.vue b/frontend/src/components/CreateObjectModal.vue deleted file mode 100644 index 7827832..0000000 --- a/frontend/src/components/CreateObjectModal.vue +++ /dev/null @@ -1,253 +0,0 @@ - - -