infinity serch
This commit is contained in:
@@ -83,5 +83,6 @@ class VkSource(Base):
|
||||
group_id: Mapped[str] = mapped_column(String(50), unique=True) # числовой ID без минуса
|
||||
group_name: Mapped[str] = mapped_column(String(200)) # станет названием категории
|
||||
enabled: Mapped[bool] = mapped_column(Boolean, default=True)
|
||||
strict_filter: Mapped[bool] = mapped_column(Boolean, default=True) # True = только посты с тегами
|
||||
last_run: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
|
||||
|
||||
@@ -17,6 +17,7 @@ from route.admin_articles import router as articles_router
|
||||
from route.admin_categories import router as categories_router
|
||||
from route.admin_media import router as media_router
|
||||
from route.admin_vk import router as vk_router
|
||||
from route.admin_system import router as system_router
|
||||
import scheduler
|
||||
|
||||
DOCS_USER = os.getenv("DOCS_USERNAME", "admin")
|
||||
@@ -102,6 +103,7 @@ app.include_router(articles_router, prefix="/admin/articles", tags=["admin-art
|
||||
app.include_router(categories_router,prefix="/admin/categories",tags=["admin-categories"])
|
||||
app.include_router(media_router, prefix="/admin/media", tags=["admin-media"])
|
||||
app.include_router(vk_router, prefix="/admin/vk", tags=["admin-vk"])
|
||||
app.include_router(system_router, prefix="/admin/system", tags=["admin-system"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -0,0 +1,26 @@
|
||||
"""add strict_filter to vk_sources
|
||||
|
||||
Revision ID: a1b2c3d4e5f6
|
||||
Revises: 25d6ed4a524a
|
||||
Create Date: 2026-05-19 10:00:00.000000
|
||||
|
||||
"""
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
revision: str = 'a1b2c3d4e5f6'
|
||||
down_revision: Union[str, None] = '25d6ed4a524a'
|
||||
branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.add_column('vk_sources',
|
||||
sa.Column('strict_filter', sa.Boolean(), nullable=False, server_default=sa.true())
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_column('vk_sources', 'strict_filter')
|
||||
115
api/route/admin_system.py
Normal file
115
api/route/admin_system.py
Normal file
@@ -0,0 +1,115 @@
|
||||
import os
|
||||
import time
|
||||
import asyncio
|
||||
from urllib.parse import urlparse
|
||||
from fastapi import APIRouter, Depends
|
||||
from sqlalchemy import text
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
|
||||
from bd.database import get_db, get_redis, DATABASE_URL, REDIS_URL
|
||||
from route.deps import require_admin
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
MINIO_ENDPOINT = os.getenv("MINIO_ENDPOINT", "minio:9000")
|
||||
MINIO_ACCESS_KEY = os.getenv("MINIO_ACCESS_KEY", "minioadmin")
|
||||
MINIO_SECRET_KEY = os.getenv("MINIO_SECRET_KEY", "minioadmin")
|
||||
MINIO_BUCKET = os.getenv("MINIO_BUCKET", "news-media")
|
||||
MINIO_PUBLIC_URL = os.getenv("MINIO_PUBLIC_URL", "http://localhost:9000")
|
||||
|
||||
|
||||
def _parse_db_url(url: str) -> dict:
|
||||
"""Вытащить host/port/db/user из DATABASE_URL без пароля."""
|
||||
try:
|
||||
# postgresql+asyncpg://user:pass@host:port/db
|
||||
p = urlparse(url.replace("+asyncpg", ""))
|
||||
return {
|
||||
"host": p.hostname or "?",
|
||||
"port": p.port or 5432,
|
||||
"database": (p.path or "/").lstrip("/") or "?",
|
||||
"user": p.username or "?",
|
||||
}
|
||||
except Exception:
|
||||
return {"host": "?", "port": 5432, "database": "?", "user": "?"}
|
||||
|
||||
|
||||
def _parse_redis_url(url: str) -> dict:
|
||||
"""Вытащить host/port из REDIS_URL без пароля."""
|
||||
try:
|
||||
p = urlparse(url)
|
||||
return {
|
||||
"host": p.hostname or "?",
|
||||
"port": p.port or 6379,
|
||||
"db": (p.path or "/0").lstrip("/") or "0",
|
||||
}
|
||||
except Exception:
|
||||
return {"host": "?", "port": 6379, "db": "0"}
|
||||
|
||||
|
||||
async def _check_postgres(db: AsyncSession) -> dict:
|
||||
info = _parse_db_url(DATABASE_URL)
|
||||
t0 = time.monotonic()
|
||||
try:
|
||||
result = await db.execute(text("SELECT version()"))
|
||||
version = result.scalar() or ""
|
||||
latency_ms = round((time.monotonic() - t0) * 1000, 1)
|
||||
# Извлекаем короткую версию: "PostgreSQL 15.3 ..."
|
||||
short_ver = version.split(",")[0] if version else "?"
|
||||
return {"ok": True, "latency_ms": latency_ms, "version": short_ver, **info}
|
||||
except Exception as e:
|
||||
return {"ok": False, "error": str(e), "latency_ms": None, **info}
|
||||
|
||||
|
||||
async def _check_redis() -> dict:
|
||||
info = _parse_redis_url(REDIS_URL)
|
||||
t0 = time.monotonic()
|
||||
try:
|
||||
redis = await get_redis()
|
||||
pong = await redis.ping()
|
||||
latency_ms = round((time.monotonic() - t0) * 1000, 1)
|
||||
server_info = await redis.info("server")
|
||||
version = server_info.get("redis_version", "?")
|
||||
return {"ok": bool(pong), "latency_ms": latency_ms, "version": f"Redis {version}", **info}
|
||||
except Exception as e:
|
||||
return {"ok": False, "error": str(e), "latency_ms": None, **info}
|
||||
|
||||
|
||||
async def _check_minio() -> dict:
|
||||
t0 = time.monotonic()
|
||||
try:
|
||||
from minio import Minio
|
||||
loop = asyncio.get_event_loop()
|
||||
client = Minio(MINIO_ENDPOINT, access_key=MINIO_ACCESS_KEY, secret_key=MINIO_SECRET_KEY, secure=False)
|
||||
exists = await loop.run_in_executor(None, client.bucket_exists, MINIO_BUCKET)
|
||||
latency_ms = round((time.monotonic() - t0) * 1000, 1)
|
||||
return {
|
||||
"ok": True,
|
||||
"latency_ms": latency_ms,
|
||||
"endpoint": MINIO_ENDPOINT,
|
||||
"bucket": MINIO_BUCKET,
|
||||
"public_url": MINIO_PUBLIC_URL,
|
||||
"bucket_exists": exists,
|
||||
}
|
||||
except Exception as e:
|
||||
return {
|
||||
"ok": False,
|
||||
"error": str(e),
|
||||
"latency_ms": None,
|
||||
"endpoint": MINIO_ENDPOINT,
|
||||
"bucket": MINIO_BUCKET,
|
||||
"public_url": MINIO_PUBLIC_URL,
|
||||
}
|
||||
|
||||
|
||||
@router.get("/status")
|
||||
async def system_status(db: AsyncSession = Depends(get_db), _: str = Depends(require_admin)):
|
||||
pg, redis, minio = await asyncio.gather(
|
||||
_check_postgres(db),
|
||||
_check_redis(),
|
||||
_check_minio(),
|
||||
)
|
||||
return {
|
||||
"postgres": pg,
|
||||
"redis": redis,
|
||||
"minio": minio,
|
||||
}
|
||||
@@ -15,16 +15,18 @@ class SourceIn(BaseModel):
|
||||
group_id: str
|
||||
group_name: str
|
||||
enabled: bool = True
|
||||
strict_filter: 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(),
|
||||
"id": str(s.id),
|
||||
"group_id": s.group_id,
|
||||
"group_name": s.group_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(),
|
||||
}
|
||||
|
||||
|
||||
@@ -56,7 +58,7 @@ async def add_source(data: SourceIn, db: AsyncSession = Depends(get_db), _: str
|
||||
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)
|
||||
source = VkSource(group_id=group_id, group_name=data.group_name, enabled=data.enabled, strict_filter=data.strict_filter)
|
||||
db.add(source)
|
||||
await db.commit()
|
||||
await db.refresh(source)
|
||||
@@ -68,9 +70,10 @@ async def update_source(source_id: str, data: SourceIn, db: AsyncSession = Depen
|
||||
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.group_id = _clean_group_id(data.group_id)
|
||||
source.group_name = data.group_name
|
||||
source.enabled = data.enabled
|
||||
source.strict_filter = data.strict_filter
|
||||
await db.commit()
|
||||
await db.refresh(source)
|
||||
return _source_dict(source)
|
||||
@@ -94,20 +97,25 @@ async def vk_manual_import(_: str = Depends(require_admin)):
|
||||
|
||||
|
||||
@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 мин)"}
|
||||
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)
|
||||
background_tasks.add_task(run_history_import, group_id=source.group_id, since_days=since_days)
|
||||
return {"ok": True, "message": f"Исторический импорт группы «{source.group_name}» запущен в фоне"}
|
||||
|
||||
@@ -359,6 +359,12 @@ async def _slug_exists(slug: str, db: AsyncSession) -> bool:
|
||||
return r.scalar_one_or_none() is not None
|
||||
|
||||
|
||||
async def _vk_post_exists(owner_id: int, post_id: int, db: AsyncSession) -> bool:
|
||||
source_url = f"https://vk.com/wall{owner_id}_{post_id}"
|
||||
r = await db.execute(select(Article.id).where(Article.source_url == source_url))
|
||||
return r.scalar_one_or_none() is not None
|
||||
|
||||
|
||||
async def _get_or_create_category(name: str, db: AsyncSession) -> Category:
|
||||
s = slugify(name)
|
||||
cat = (await db.execute(select(Category).where(Category.slug == s))).scalar_one_or_none()
|
||||
@@ -381,21 +387,39 @@ async def _get_or_create_tag(name: str, db: AsyncSession) -> Tag:
|
||||
|
||||
# ── Post processing ───────────────────────────────────────────────────────────
|
||||
|
||||
async def _process_post(post: dict, category_name: str, db: AsyncSession) -> bool:
|
||||
async def _process_post(post: dict, category_name: str, db: AsyncSession, strict_filter: bool = True) -> bool:
|
||||
owner_id = post["owner_id"]
|
||||
post_id = post["id"]
|
||||
vk_slug = _vk_slug(owner_id, post_id)
|
||||
|
||||
if await _slug_exists(vk_slug, db):
|
||||
if await _vk_post_exists(owner_id, post_id, db):
|
||||
return False
|
||||
|
||||
text = post.get("text", "")
|
||||
attachments = post.get("attachments", [])
|
||||
ts = post["date"]
|
||||
|
||||
# Фильтр: импортируем только посты с разрешёнными тегами
|
||||
# Обработка репостов (copy_history)
|
||||
copy_history = post.get("copy_history", [])
|
||||
if copy_history:
|
||||
original = copy_history[0]
|
||||
orig_gid = str(abs(original.get("owner_id", 0)))
|
||||
in_db = (await db.execute(
|
||||
select(VkSource.id).where(VkSource.group_id == orig_gid)
|
||||
)).scalar_one_or_none()
|
||||
if in_db:
|
||||
# Оригинал придёт сам из своей группы — пропускаем
|
||||
return False
|
||||
# Берём содержимое оригинала; комментарий репостера добавляем в начало
|
||||
orig_text = original.get("text", "")
|
||||
orig_atts = original.get("attachments", [])
|
||||
text = (text + "\n" + orig_text).strip() if text else orig_text
|
||||
attachments = orig_atts if orig_atts else attachments
|
||||
|
||||
tag_names = _detect_tags(text)
|
||||
if not tag_names:
|
||||
# strict_filter=True: пропускаем пост если нет тегов
|
||||
# strict_filter=False: импортируем всё, теги назначаем если найдены
|
||||
if strict_filter and not tag_names:
|
||||
return False
|
||||
|
||||
# Заранее генерируем UUID и slug — используются в путях MinIO
|
||||
@@ -648,7 +672,7 @@ async def run_import() -> dict:
|
||||
batch_new = 0
|
||||
for post in posts:
|
||||
try:
|
||||
ok = await _process_post(post, source.group_name, db)
|
||||
ok = await _process_post(post, source.group_name, db, source.strict_filter)
|
||||
if ok:
|
||||
imported += 1; grp_new += 1; batch_new += 1
|
||||
else:
|
||||
@@ -721,7 +745,7 @@ async def run_history_import(group_id: str | None = None, since_days: int = 365)
|
||||
hit_cutoff = True
|
||||
continue
|
||||
try:
|
||||
ok = await _process_post(post, source.group_name, db)
|
||||
ok = await _process_post(post, source.group_name, db, source.strict_filter)
|
||||
if ok:
|
||||
imported += 1; grp_new += 1; batch_new += 1
|
||||
else:
|
||||
|
||||
Reference in New Issue
Block a user