from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException # BackgroundTasks используется в history endpoints from pydantic import BaseModel from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from bd.database import get_db from bd.models import VkSource from route.deps import require_admin from vk_parser import run_import, run_history_import, VK_TOKEN router = APIRouter() class SourceIn(BaseModel): group_id: str group_name: str enabled: bool = True def _source_dict(s: VkSource) -> dict: return { "id": str(s.id), "group_id": s.group_id, "group_name": s.group_name, "enabled": s.enabled, "last_run": s.last_run.isoformat() if s.last_run else None, "created_at": s.created_at.isoformat(), } @router.get("/status") async def vk_status(_: str = Depends(require_admin)): return {"configured": bool(VK_TOKEN)} # ── CRUD источников ─────────────────────────────────────────────────────────── @router.get("/sources") async def list_sources(db: AsyncSession = Depends(get_db), _: str = Depends(require_admin)): rows = (await db.execute(select(VkSource).order_by(VkSource.created_at))).scalars().all() return [_source_dict(r) for r in rows] def _clean_group_id(raw: str) -> str: """Из строки вида '-95503797', '95503797_6512', 'wall-95503797_12' вытащить только числовой ID группы.""" import re raw = raw.strip().lstrip("-") # Берём только первое число (до '_' если есть) m = re.search(r"(\d+)", raw) return m.group(1) if m else raw @router.post("/sources", status_code=201) async def add_source(data: SourceIn, db: AsyncSession = Depends(get_db), _: str = Depends(require_admin)): group_id = _clean_group_id(data.group_id) exists = (await db.execute(select(VkSource).where(VkSource.group_id == group_id))).scalar_one_or_none() if exists: raise HTTPException(status_code=409, detail="Источник с таким group_id уже есть") source = VkSource(group_id=group_id, group_name=data.group_name, enabled=data.enabled) db.add(source) await db.commit() await db.refresh(source) return _source_dict(source) @router.patch("/sources/{source_id}") async def update_source(source_id: str, data: SourceIn, db: AsyncSession = Depends(get_db), _: str = Depends(require_admin)): source = (await db.execute(select(VkSource).where(VkSource.id == source_id))).scalar_one_or_none() if not source: raise HTTPException(status_code=404, detail="Источник не найден") source.group_id = _clean_group_id(data.group_id) source.group_name = data.group_name source.enabled = data.enabled await db.commit() await db.refresh(source) return _source_dict(source) @router.delete("/sources/{source_id}", status_code=204) async def delete_source(source_id: str, db: AsyncSession = Depends(get_db), _: str = Depends(require_admin)): source = (await db.execute(select(VkSource).where(VkSource.id == source_id))).scalar_one_or_none() if not source: raise HTTPException(status_code=404, detail="Источник не найден") await db.delete(source) await db.commit() # ── Ручной запуск импорта ───────────────────────────────────────────────────── @router.post("/import") async def vk_manual_import(_: str = Depends(require_admin)): result = await run_import() return result @router.post("/import/history") async def vk_history_import_all(background_tasks: BackgroundTasks, _: str = Depends(require_admin)): background_tasks.add_task(run_history_import) return {"ok": True, "message": "Исторический импорт всех групп запущен в фоне (до 5 мин)"} @router.post("/import/history/{source_id}") async def vk_history_import_one( source_id: str, background_tasks: BackgroundTasks, db: AsyncSession = Depends(get_db), _: str = Depends(require_admin), ): source = (await db.execute(select(VkSource).where(VkSource.id == source_id))).scalar_one_or_none() if not source: raise HTTPException(status_code=404, detail="Источник не найден") background_tasks.add_task(run_history_import, group_id=source.group_id) return {"ok": True, "message": f"Исторический импорт группы «{source.group_name}» запущен в фоне"}