Files
news_all_spo/api/route/admin_vk.py
jze9 615c819e97 fix: корректные source_url для VK постов (screen_name формат)
- VkSource: добавлено поле screen_name, заполняется из VK API при добавлении источника
- vk_parser: source_url теперь в формате vk.com/{screen_name}?w=wall{owner}_{post}
  вместо устаревшего vk.com/wall{owner}_{post}
- _vk_post_exists: проверка через LIKE чтобы находить оба формата URL
- admin_vk: автоматически получает screen_name при создании/обновлении источника
- admin_system: исправлен regex в repair-media-urls — только /news-media/ URL
- Alembic миграция: b2c3d4e5f6a7 добавляет колонку screen_name в vk_sources

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-05-20 19:32:28 +05:00

145 lines
5.8 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import httpx
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
strict_filter: bool = True
async def _fetch_screen_name(group_id: str) -> str | None:
if not VK_TOKEN:
return None
try:
async with httpx.AsyncClient(timeout=5) as client:
r = await client.get(
"https://api.vk.com/method/groups.getById",
params={"group_ids": group_id, "access_token": VK_TOKEN, "v": "5.131"},
)
data = r.json()
groups = data.get("response", [])
if groups:
return groups[0].get("screen_name")
except Exception:
pass
return None
def _source_dict(s: VkSource) -> dict:
return {
"id": str(s.id),
"group_id": s.group_id,
"group_name": s.group_name,
"screen_name": s.screen_name,
"enabled": s.enabled,
"strict_filter": s.strict_filter,
"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 уже есть")
screen_name = await _fetch_screen_name(group_id)
source = VkSource(group_id=group_id, group_name=data.group_name, screen_name=screen_name, enabled=data.enabled, strict_filter=data.strict_filter)
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
source.strict_filter = data.strict_filter
if not source.screen_name:
source.screen_name = await _fetch_screen_name(source.group_id)
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,
since_days: int = 365,
_: str = Depends(require_admin),
):
background_tasks.add_task(run_history_import, since_days=since_days)
return {"ok": True, "message": f"Исторический импорт всех групп запущен в фоне (глубина {since_days} дн.)"}
@router.post("/import/history/{source_id}")
async def vk_history_import_one(
source_id: str,
background_tasks: BackgroundTasks,
since_days: int = 365,
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, since_days=since_days)
return {"ok": True, "message": f"Исторический импорт группы «{source.group_name}» запущен в фоне"}