114 lines
4.7 KiB
Python
114 lines
4.7 KiB
Python
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}» запущен в фоне"}
|