diff --git a/.env.example b/.env.example
index db94143..3ddeab8 100644
--- a/.env.example
+++ b/.env.example
@@ -10,5 +10,21 @@ MINIO_PUBLIC_URL=http://localhost:9000
DOCS_USERNAME=admin
DOCS_PASSWORD=admin
+# CORS: разрешённые домены через запятую
+ALLOWED_ORIGINS=http://localhost,http://localhost:80
+
ADMIN_USERNAME=admin
ADMIN_PASSWORD=changeme
+
+# VK импорт (необязательно)
+# Сервисный токен приложения: vk.com/apps → Настройки → Сервисный ключ доступа
+VK_ACCESS_TOKEN=
+# ID групп через запятую (без минуса): 12345,67890
+VK_GROUP_IDS=
+# Статус импортированных постов: published или draft
+VK_IMPORT_STATUS=published
+# Сколько последних постов проверять за раз
+VK_POSTS_PER_RUN=20
+# Время запуска импорта (UTC)
+VK_IMPORT_HOUR=6
+VK_IMPORT_MINUTE=0
diff --git a/.gitignore b/.gitignore
index e6ac574..c5a1982 100644
--- a/.gitignore
+++ b/.gitignore
@@ -17,10 +17,14 @@ __pycache__/
.ruff_cache/
uv.lock
-# ── Node / React editor UI ─────────────────────────────────────────────────
+# ── Node / React ───────────────────────────────────────────────────────────
+node_modules/
web/editor-ui/node_modules/
web/editor-ui/dist/
web/editor-ui/.vite/
+web-react/node_modules/
+web-react/dist/
+web-react/.vite/
# ── IDE / OS ───────────────────────────────────────────────────────────────
.DS_Store
diff --git a/api/bd/models.py b/api/bd/models.py
index 461df0c..7802c7c 100644
--- a/api/bd/models.py
+++ b/api/bd/models.py
@@ -1,6 +1,6 @@
import uuid
from datetime import datetime
-from sqlalchemy import String, Text, DateTime, ForeignKey, Table, Column, Integer, Enum as SAEnum
+from sqlalchemy import String, Text, DateTime, ForeignKey, Table, Column, Integer, Boolean, Enum as SAEnum
from sqlalchemy.orm import mapped_column, Mapped, relationship
from sqlalchemy.dialects.postgresql import UUID
import enum
@@ -44,6 +44,7 @@ class Article(Base):
excerpt: Mapped[str] = mapped_column(String(1000), default="")
cover_url: Mapped[str | None] = mapped_column(String(1000), nullable=True)
font_family: Mapped[str] = mapped_column(String(100), default="Merriweather")
+ source_url: Mapped[str | None] = mapped_column(String(1000), nullable=True)
status: Mapped[ArticleStatus] = mapped_column(SAEnum(ArticleStatus), default=ArticleStatus.draft)
view_count: Mapped[int] = mapped_column(Integer, default=0)
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
@@ -71,3 +72,13 @@ class Media(Base):
media_type: Mapped[str] = mapped_column(String(20)) # image | video
size_bytes: Mapped[int] = mapped_column(Integer, default=0)
created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
+
+
+class VkSource(Base):
+ __tablename__ = "vk_sources"
+ id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
+ 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)
+ last_run: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
+ created_at: Mapped[datetime] = mapped_column(DateTime, default=datetime.utcnow)
diff --git a/api/main.py b/api/main.py
index 468d5ab..0ebda36 100644
--- a/api/main.py
+++ b/api/main.py
@@ -4,6 +4,11 @@ from contextlib import asynccontextmanager
from fastapi import FastAPI, Depends, HTTPException, status
from fastapi.security import HTTPBasic, HTTPBasicCredentials
from fastapi.middleware.cors import CORSMiddleware
+from fastapi.responses import JSONResponse
+from fastapi.openapi.docs import get_swagger_ui_html, get_redoc_html
+from slowapi import _rate_limit_exceeded_handler
+from slowapi.errors import RateLimitExceeded
+from rate_limiter import limiter
from bd.database import wait_for_db, engine, Base
from route.public import router as public_router
@@ -11,12 +16,14 @@ from route.admin_auth import router as auth_router, ensure_default_admin
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
+import scheduler
DOCS_USER = os.getenv("DOCS_USERNAME", "admin")
DOCS_PASS = os.getenv("DOCS_PASSWORD", "admin")
_basic = HTTPBasic()
-
+# Rate limiter — ключ по IP-адресу
def _verify_docs(creds: HTTPBasicCredentials = Depends(_basic)):
ok_u = secrets.compare_digest(creds.username.encode(), DOCS_USER.encode())
ok_p = secrets.compare_digest(creds.password.encode(), DOCS_PASS.encode())
@@ -31,7 +38,9 @@ async def lifespan(app: FastAPI):
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
await ensure_default_admin()
+ scheduler.start()
yield
+ scheduler.shutdown()
app = FastAPI(
@@ -43,17 +52,18 @@ app = FastAPI(
openapi_url=None,
)
+app.state.limiter = limiter
+app.add_exception_handler(RateLimitExceeded, _rate_limit_exceeded_handler)
+
+# CORS: разрешаем только наш домен (меняй на продакшн-домен)
+_ALLOWED_ORIGINS = [o.strip() for o in os.getenv("ALLOWED_ORIGINS", "http://localhost,http://localhost:80").split(",") if o.strip()]
app.add_middleware(
CORSMiddleware,
- allow_origins=["*"],
- allow_methods=["*"],
- allow_headers=["*"],
+ allow_origins=_ALLOWED_ORIGINS,
+ allow_methods=["GET", "POST", "PUT", "PATCH", "DELETE"],
+ allow_headers=["Authorization", "Content-Type"],
)
-# Docs protected by HTTP Basic Auth
-from fastapi.openapi.docs import get_swagger_ui_html, get_redoc_html
-from fastapi.responses import JSONResponse
-
@app.get("/docs", include_in_schema=False)
async def docs(username: str = Depends(_verify_docs)):
@@ -75,11 +85,12 @@ async def health():
return {"ok": True}
-app.include_router(public_router, prefix="/news", tags=["public"])
-app.include_router(auth_router, prefix="/admin", tags=["admin-auth"])
-app.include_router(articles_router, prefix="/admin/articles", tags=["admin-articles"])
-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(public_router, prefix="/news", tags=["public"])
+app.include_router(auth_router, prefix="/admin", tags=["admin-auth"])
+app.include_router(articles_router, prefix="/admin/articles", tags=["admin-articles"])
+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"])
if __name__ == "__main__":
diff --git a/api/rate_limiter.py b/api/rate_limiter.py
new file mode 100644
index 0000000..38404a8
--- /dev/null
+++ b/api/rate_limiter.py
@@ -0,0 +1,4 @@
+from slowapi import Limiter
+from slowapi.util import get_remote_address
+
+limiter = Limiter(key_func=get_remote_address)
diff --git a/api/requirements.txt b/api/requirements.txt
index a657c84..0119fd3 100644
--- a/api/requirements.txt
+++ b/api/requirements.txt
@@ -10,3 +10,11 @@ python-multipart==0.0.12
minio==7.2.11
python-slugify==8.0.4
pillow==11.0.0
+httpx==0.27.2
+apscheduler==3.10.4
+bleach==6.1.0
+tinycss2==1.4.0
+beautifulsoup4==4.12.3
+lxml==5.3.0
+slowapi==0.1.9
+yt-dlp>=2024.12.3
diff --git a/api/route/admin_articles.py b/api/route/admin_articles.py
index 097d111..bfc7726 100644
--- a/api/route/admin_articles.py
+++ b/api/route/admin_articles.py
@@ -21,6 +21,7 @@ class ArticleIn(BaseModel):
content: str = ""
excerpt: str = ""
cover_url: Optional[str] = None
+ source_url: Optional[str] = None
font_family: str = "Merriweather"
category_id: Optional[str] = None
tag_names: list[str] = []
@@ -35,6 +36,7 @@ def _article_dict(a: Article) -> dict:
"content": a.content,
"excerpt": a.excerpt,
"cover_url": a.cover_url,
+ "source_url": a.source_url,
"font_family": a.font_family,
"status": a.status.value,
"view_count": a.view_count,
@@ -73,14 +75,19 @@ async def list_articles(
offset: int = Query(0, ge=0),
limit: int = Query(50, ge=1, le=200),
status: Optional[ArticleStatus] = Query(None),
+ q: Optional[str] = Query(None, max_length=200),
db: AsyncSession = Depends(get_db),
_: str = Depends(require_admin),
):
- q = select(Article).options(selectinload(Article.category), selectinload(Article.tags))
+ from sqlalchemy import or_
+ stmt = select(Article).options(selectinload(Article.category), selectinload(Article.tags))
if status:
- q = q.where(Article.status == status)
- total = (await db.execute(select(func.count()).select_from(q.subquery()))).scalar_one()
- articles = (await db.execute(q.order_by(Article.created_at.desc()).offset(offset).limit(limit))).scalars().all()
+ stmt = stmt.where(Article.status == status)
+ if q:
+ pattern = f"%{q}%"
+ stmt = stmt.where(or_(Article.title.ilike(pattern), Article.excerpt.ilike(pattern)))
+ total = (await db.execute(select(func.count()).select_from(stmt.subquery()))).scalar_one()
+ articles = (await db.execute(stmt.order_by(Article.created_at.desc()).offset(offset).limit(limit))).scalars().all()
return {
"total": total,
"offset": offset,
@@ -115,6 +122,7 @@ async def create_article(
content=data.content,
excerpt=data.excerpt,
cover_url=data.cover_url,
+ source_url=data.source_url,
font_family=data.font_family,
status=data.status,
category=cat,
@@ -173,7 +181,8 @@ async def update_article(
article.slug = slug
article.content = data.content
article.excerpt = data.excerpt
- article.cover_url = data.cover_url
+ article.cover_url = data.cover_url
+ article.source_url = data.source_url
article.font_family = data.font_family
article.status = data.status
article.category = cat
diff --git a/api/route/admin_auth.py b/api/route/admin_auth.py
index 33ad2f6..a0fe41e 100644
--- a/api/route/admin_auth.py
+++ b/api/route/admin_auth.py
@@ -1,6 +1,6 @@
import os
-from datetime import datetime, timedelta
-from fastapi import APIRouter, Depends, HTTPException, status
+from datetime import datetime, timedelta, timezone
+from fastapi import APIRouter, Depends, HTTPException, Request, status
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select
@@ -9,17 +9,24 @@ from passlib.context import CryptContext
from bd.database import get_db, AsyncSessionLocal
from bd.models import Admin
+from rate_limiter import limiter
router = APIRouter()
pwd_context = CryptContext(schemes=["bcrypt"], deprecated="auto")
-JWT_SECRET = os.getenv("JWT_SECRET", "changeme")
-JWT_ALGO = "HS256"
-JWT_EXP_HOURS = 24
+_raw_secret = os.getenv("JWT_SECRET", "")
+if not _raw_secret:
+ import secrets as _s
+ _raw_secret = _s.token_hex(32)
+ print("[WARN] JWT_SECRET не задан — сгенерирован временный секрет. Установи JWT_SECRET в .env!")
+
+JWT_SECRET = _raw_secret
+JWT_ALGO = "HS256"
+JWT_EXP_HOURS = 8
def _make_token(username: str) -> str:
- exp = datetime.utcnow() + timedelta(hours=JWT_EXP_HOURS)
+ exp = datetime.now(timezone.utc) + timedelta(hours=JWT_EXP_HOURS)
return jwt.encode({"sub": username, "exp": exp}, JWT_SECRET, algorithm=JWT_ALGO)
@@ -40,7 +47,8 @@ class LoginRequest(BaseModel):
@router.post("/login")
-async def login(data: LoginRequest, db: AsyncSession = Depends(get_db)):
+@limiter.limit("5/minute")
+async def login(request: Request, data: LoginRequest, db: AsyncSession = Depends(get_db)):
admin = (await db.execute(select(Admin).where(Admin.username == data.username))).scalar_one_or_none()
if not admin or not pwd_context.verify(data.password, admin.password_hash):
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="Неверный логин или пароль")
@@ -48,12 +56,7 @@ async def login(data: LoginRequest, db: AsyncSession = Depends(get_db)):
@router.post("/change-password")
-async def change_password(
- data: dict,
- db: AsyncSession = Depends(get_db),
-):
- from route.deps import require_admin
- # handled with deps in main
+async def change_password(data: dict, db: AsyncSession = Depends(get_db)):
admin = (await db.execute(select(Admin).where(Admin.username == data["username"]))).scalar_one_or_none()
if not admin or not pwd_context.verify(data["current_password"], admin.password_hash):
raise HTTPException(status_code=400, detail="Неверный текущий пароль")
diff --git a/api/route/admin_media.py b/api/route/admin_media.py
index 650fde7..7c32268 100644
--- a/api/route/admin_media.py
+++ b/api/route/admin_media.py
@@ -99,20 +99,25 @@ async def upload_media(
@router.get("")
async def list_media(
offset: int = Query(0, ge=0),
- limit: int = Query(50, ge=1, le=200),
+ limit: int = Query(40, ge=1, le=200),
media_type: str | None = Query(None),
db: AsyncSession = Depends(get_db),
_: str = Depends(require_admin),
):
- q = select(Media)
+ from sqlalchemy import func
+ base = select(Media)
if media_type:
- q = q.where(Media.media_type == media_type)
- items = (await db.execute(q.order_by(Media.created_at.desc()).offset(offset).limit(limit))).scalars().all()
- return [
- {"id": str(m.id), "url": m.url, "filename": m.filename, "media_type": m.media_type,
- "size_bytes": m.size_bytes, "created_at": m.created_at.isoformat()}
- for m in items
- ]
+ base = base.where(Media.media_type == media_type)
+ total = (await db.execute(select(func.count()).select_from(base.subquery()))).scalar_one()
+ items = (await db.execute(base.order_by(Media.created_at.desc()).offset(offset).limit(limit))).scalars().all()
+ return {
+ "total": total,
+ "items": [
+ {"id": str(m.id), "url": m.url, "filename": m.filename, "media_type": m.media_type,
+ "size_bytes": m.size_bytes, "created_at": m.created_at.isoformat()}
+ for m in items
+ ],
+ }
@router.delete("/{media_id}")
diff --git a/api/route/admin_vk.py b/api/route/admin_vk.py
new file mode 100644
index 0000000..9e0a7e4
--- /dev/null
+++ b/api/route/admin_vk.py
@@ -0,0 +1,113 @@
+from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException
+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(background_tasks: BackgroundTasks, _: str = Depends(require_admin)):
+ background_tasks.add_task(run_import)
+ return {"ok": True, "message": "Импорт последних постов запущен в фоне"}
+
+
+@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}» запущен в фоне"}
diff --git a/api/route/public.py b/api/route/public.py
index 4da6fda..7626896 100644
--- a/api/route/public.py
+++ b/api/route/public.py
@@ -1,12 +1,13 @@
import json
+from datetime import datetime
from fastapi import APIRouter, Depends, HTTPException, Query
from sqlalchemy.ext.asyncio import AsyncSession
-from sqlalchemy import select, func
+from sqlalchemy import select, func, or_
from sqlalchemy.orm import selectinload
import redis.asyncio as aioredis
from bd.database import get_db, get_redis
-from bd.models import Article, Category, Tag, ArticleStatus
+from bd.models import Article, Category, Tag, ArticleStatus, article_tag
router = APIRouter()
@@ -22,6 +23,7 @@ def _article_card(a: Article) -> dict:
"slug": a.slug,
"excerpt": a.excerpt,
"cover_url": a.cover_url,
+ "source_url": a.source_url,
"published_at": a.published_at.isoformat() if a.published_at else None,
"view_count": a.view_count,
"category": _cat_dict(a.category) if a.category else None,
@@ -42,10 +44,13 @@ async def list_news(
limit: int = Query(20, ge=1, le=100),
category: str | None = Query(None),
tag: str | None = Query(None),
+ date_from: str | None = Query(None),
+ date_to: str | None = Query(None),
+ q: str | None = Query(None, max_length=200),
db: AsyncSession = Depends(get_db),
redis: aioredis.Redis = Depends(get_redis),
):
- cache_key = f"news:list:{offset}:{limit}:{category or ''}:{tag or ''}"
+ cache_key = f"news:list:{offset}:{limit}:{category or ''}:{tag or ''}:{date_from or ''}:{date_to or ''}:{q or ''}"
cached = await redis.get(cache_key)
if cached:
return json.loads(cached)
@@ -68,6 +73,29 @@ async def list_news(
if t:
base = base.where(Article.tags.any(Tag.id == t.id))
+ if date_from:
+ try:
+ base = base.where(Article.published_at >= datetime.fromisoformat(date_from))
+ except ValueError:
+ pass
+
+ if date_to:
+ try:
+ base = base.where(Article.published_at <= datetime.fromisoformat(date_to))
+ except ValueError:
+ pass
+
+ if q:
+ pattern = f"%{q}%"
+ base = base.where(
+ or_(
+ Article.title.ilike(pattern),
+ Article.excerpt.ilike(pattern),
+ Article.content.ilike(pattern),
+ Article.tags.any(Tag.name.ilike(pattern)),
+ )
+ )
+
count_q = select(func.count()).select_from(base.subquery())
total: int = (await db.execute(count_q)).scalar_one()
@@ -82,7 +110,7 @@ async def list_news(
"has_more": (offset + limit) < total,
"items": [_article_card(a) for a in articles],
}
- await redis.setex(cache_key, 30, json.dumps(result, default=str))
+ await redis.setex(cache_key, 30 if q else 120, json.dumps(result, default=str))
return result
@@ -96,7 +124,36 @@ async def list_categories(
return json.loads(cached)
cats = (await db.execute(select(Category).order_by(Category.name))).scalars().all()
result = [_cat_dict(c) for c in cats]
- await redis.setex("news:categories", 60, json.dumps(result))
+ await redis.setex("news:categories", 300, json.dumps(result))
+ return result
+
+
+@router.get("/tags")
+async def list_tags(
+ limit: int = Query(30, ge=1, le=100),
+ db: AsyncSession = Depends(get_db),
+ redis: aioredis.Redis = Depends(get_redis),
+):
+ cache_key = f"news:tags:{limit}"
+ cached = await redis.get(cache_key)
+ if cached:
+ return json.loads(cached)
+
+ q = (
+ select(Tag, func.count(article_tag.c.article_id).label("cnt"))
+ .join(article_tag, Tag.id == article_tag.c.tag_id)
+ .join(Article, article_tag.c.article_id == Article.id)
+ .where(Article.status == ArticleStatus.published)
+ .group_by(Tag.id)
+ .order_by(func.count(article_tag.c.article_id).desc())
+ .limit(limit)
+ )
+ rows = (await db.execute(q)).all()
+ result = [
+ {"id": str(r.Tag.id), "name": r.Tag.name, "slug": r.Tag.slug, "count": r.cnt}
+ for r in rows
+ ]
+ await redis.setex(cache_key, 300, json.dumps(result))
return result
diff --git a/api/scheduler.py b/api/scheduler.py
new file mode 100644
index 0000000..cf87a6f
--- /dev/null
+++ b/api/scheduler.py
@@ -0,0 +1,26 @@
+import os
+from apscheduler.schedulers.asyncio import AsyncIOScheduler
+
+_scheduler = AsyncIOScheduler(timezone="UTC")
+
+VK_HOUR = int(os.getenv("VK_IMPORT_HOUR", "6"))
+VK_MINUTE = int(os.getenv("VK_IMPORT_MINUTE", "0"))
+
+
+def start():
+ from vk_parser import run_import
+ _scheduler.add_job(
+ run_import,
+ trigger="cron",
+ hour=VK_HOUR,
+ minute=VK_MINUTE,
+ id="vk_daily_import",
+ replace_existing=True,
+ )
+ _scheduler.start()
+ print(f"[Scheduler] VK-импорт запланирован на {VK_HOUR:02d}:{VK_MINUTE:02d} UTC каждый день")
+
+
+def shutdown():
+ if _scheduler.running:
+ _scheduler.shutdown(wait=False)
diff --git a/api/vk_parser.py b/api/vk_parser.py
new file mode 100644
index 0000000..f1b494a
--- /dev/null
+++ b/api/vk_parser.py
@@ -0,0 +1,725 @@
+import io
+import os
+import re
+import uuid
+from datetime import datetime, timezone
+
+import asyncio
+import bleach
+from bleach.css_sanitizer import CSSSanitizer
+import httpx
+from minio import Minio
+from minio.error import S3Error
+from slugify import slugify
+from sqlalchemy import select
+from sqlalchemy.ext.asyncio import AsyncSession
+
+from bd.database import AsyncSessionLocal
+from bd.models import Article, ArticleStatus, Category, Tag, VkSource
+
+VK_API = "https://api.vk.com/method"
+VK_VERSION = "5.199"
+VK_TOKEN = os.getenv("VK_ACCESS_TOKEN", "")
+VK_IMPORT_AS = os.getenv("VK_IMPORT_STATUS", "published") # published | draft
+VK_PER_RUN = int(os.getenv("VK_POSTS_PER_RUN", "100")) # макс 100 за запрос
+
+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")
+
+_URL_RE = re.compile(r"(https?://[^\s<>\"']+)")
+
+# Разрешённые HTML-теги и атрибуты после санитизации
+_ALLOWED_TAGS = [
+ "p", "br", "strong", "em", "a", "img",
+ "ul", "ol", "li", "blockquote",
+ "div", "iframe", "span", "video",
+]
+_ALLOWED_ATTRS = {
+ "a": ["href", "target", "rel"],
+ "img": ["src", "style", "alt"],
+ "iframe": ["src", "style", "frameborder", "allowfullscreen"],
+ "video": ["src", "controls", "style", "poster"],
+ "div": ["style"],
+ "span": ["style"],
+}
+
+
+_CSS = CSSSanitizer(allowed_css_properties=[
+ "max-width", "width", "height", "position", "top", "left",
+ "padding-bottom", "overflow", "margin", "border-radius",
+])
+
+
+def _sanitize(html: str) -> str:
+ return bleach.clean(html, tags=_ALLOWED_TAGS, attributes=_ALLOWED_ATTRS,
+ css_sanitizer=_CSS, strip=True)
+
+# ── Теги по ключевым словам ───────────────────────────────────────────────────
+# Ключ = название тега, значение = список подстрок (регистронезависимо)
+TAG_KEYWORDS: dict[str, list[str]] = {
+ # ── Учёба и поступление ───────────────────────────────────────────────────
+ "экзамены": ["экзамен", "зачёт", "зачет", "сессия", "огэ", "егэ",
+ "гиа", "защит", "диплом", "аттестац", "промежуточн"],
+ "абитуриентам": ["абитуриент", "поступ", "приёмн", "приемн", "зачислен",
+ "набор студент", "подай документ", "приглашаем поступ",
+ "бюджетн мест", "целевой приём", "контрольные цифры"],
+ "день открытых дверей": ["день открытых дверей", "открытые двери", "день открытых",
+ "экскурси", "познакомиться с колледж"],
+ "практика": ["практик", "стажировк", "производственн", "учебно-производ",
+ "на предприяти", "работодател", "наставник"],
+ "достижения": ["победител", "призёр", "призер", "награда", "грамот",
+ "диплом лауреат", "1 место", "2 место", "3 место",
+ "лучший студент", "гордост", "поздравляем", "медал"],
+ "студенческая жизнь": ["студсовет", "студенческ жизн", "общежити", "капустник",
+ "квн", "студенч", "первокурсник", "посвящение",
+ "студент год", "актив"],
+ "профессионалитет": ["профессионалитет", "фп профессионалитет",
+ "кластер", "федеральн проект"],
+
+ # ── Направления колледжей ─────────────────────────────────────────────────
+ "медицина": ["медицин", "сестринск", "фельдшер", "фармацевт", "лечебн",
+ "анатоми", "патологи", "здравоохранен", "санитар",
+ "первая помощ", "реанимац", "клиническ"],
+ "педагогика": ["педагог", "воспитател", "учител", "начальн класс",
+ "дошкольн", "детский сад", "логопед", "коррекцион",
+ "инклюзив", "тьютор"],
+ "юриспруденция": ["юрист", "юридическ", "правоохранительн", "право",
+ "законодательств", "суд", "прокурат", "полиц",
+ "юриспруденц", "правовед"],
+ "IT и технологии": ["программирован", "it-", "ит-", "информационн технолог",
+ "цифров", "кибербезопасн", "веб-разраб", "1с",
+ "компьютерн", "разработчик", "python", "frontend",
+ "backend", "хакатон", "ворлдскиллс"],
+ "кулинария и торговля": ["кулинар", "повар", "кондитер", "гастроном", "блюд",
+ "рецепт", "ресторан", "кафе", "торговл", "продавец",
+ "мерчандайзинг", "товаровед", "общепит"],
+ "строительство": ["строительств", "архитектур", "проектирован", "чертёж",
+ "чертеж", "монтаж", "сварк", "электромонтаж",
+ "сантехник", "отделочн"],
+ "техника и механика": ["механик", "двигател", "автомобил", "станок",
+ "металлообраб", "токар", "слесар", "техническ обслуж",
+ "ремонт оборудован", "машиностроен", "технолог производ"],
+ "сельское хозяйство": ["агроном", "агроинженер", "сельск хоз", "животновод",
+ "агротехник", "землепользован", "агро"],
+ "экономика и бухучёт": ["бухгалтер", "экономик", "финансов", "налог",
+ "аудит", "бизнес", "предпринимател", "менеджмент",
+ "маркетинг", "логистик"],
+ "социальная работа": ["социальн работ", "соцработник", "психолог",
+ "реабилитац", "инвалид", "ограниченн возможн",
+ "социальн помощ", "опека"],
+
+ # ── Общественные события ──────────────────────────────────────────────────
+ "спорт": ["спорт", "соревнован", "олимпиад", "футбол", "баскетбол",
+ "волейбол", "чемпион", "кубок", "турнир", "атлет",
+ "фитнес", "зарядк", "ворлдскиллс хайскул"],
+ "военка": ["военн", "армия", "призыв", "нво", "сво", "оборон",
+ "патриот", "защитник", "военно-патриот", "стрельб",
+ "зарниц", "юнармия", "допризывн"],
+ "воздушная опасность": ["воздушн тревог", "воздушн опасност", "ракет",
+ "бпла", "дрон", "обстрел", "укрытие", "сирен",
+ "эвакуац", "бомбоубежищ"],
+ "общественная деятельность": ["волонтёр", "волонтер", "благотвор", "субботник",
+ "экологическ акц", "посадк дерев", "уборк территор",
+ "донорств", "гуманитарн", "помощ фронт"],
+ "праздники и события": ["праздник", "концерт", "фестиваль", "выставк",
+ "торжеств", "церемони", "день колледж", "юбилей",
+ "8 марта", "23 февраля", "новый год", "масленица",
+ "день знаний", "последний звонок"],
+ "конкурсы и олимпиады": ["конкурс", "олимпиад", "чемпионат профессий",
+ "worldskills", "ворлдскиллс", "хакатон", "викторин",
+ "интеллектуальн", "акселератор"],
+}
+
+
+def _detect_tags(text: str) -> list[str]:
+ lower = text.lower()
+ found = []
+ for tag_name, keywords in TAG_KEYWORDS.items():
+ if any(kw in lower for kw in keywords):
+ found.append(tag_name)
+ return found
+
+
+# ── MinIO helpers ─────────────────────────────────────────────────────────────
+
+def _minio_client() -> Minio:
+ return Minio(MINIO_ENDPOINT, access_key=MINIO_ACCESS_KEY, secret_key=MINIO_SECRET_KEY, secure=False)
+
+
+def _ensure_bucket(mc: Minio):
+ try:
+ if not mc.bucket_exists(MINIO_BUCKET):
+ mc.make_bucket(MINIO_BUCKET)
+ policy = (
+ '{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*",'
+ f'"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::{MINIO_BUCKET}/*"]'
+ '}]}'
+ )
+ mc.set_bucket_policy(MINIO_BUCKET, policy)
+ except S3Error:
+ pass
+
+
+async def _upload_from_url(url: str, db: AsyncSession | None = None) -> str | None:
+ try:
+ async with httpx.AsyncClient(timeout=30, follow_redirects=True) as client:
+ r = await client.get(url)
+ if r.status_code != 200:
+ return None
+ data = r.content
+ ct = r.headers.get("content-type", "image/jpeg").split(";")[0].strip()
+ ext = ct.split("/")[-1].replace("jpeg", "jpg")
+ object_name = f"vk/images/{uuid.uuid4().hex}.{ext}"
+ mc = _minio_client()
+ _ensure_bucket(mc)
+ mc.put_object(MINIO_BUCKET, object_name, io.BytesIO(data), len(data), content_type=ct)
+ public_url = f"{MINIO_PUBLIC_URL}/{MINIO_BUCKET}/{object_name}"
+ if db is not None:
+ from bd.models import Media
+ db.add(Media(
+ filename=object_name.split("/")[-1],
+ url=public_url,
+ media_type="image",
+ size_bytes=len(data),
+ ))
+ return public_url
+ except Exception as exc:
+ print(f"[VK] Не удалось загрузить фото {url}: {exc}")
+ return None
+
+
+# ── Video download ────────────────────────────────────────────────────────────
+
+async def _fetch_video_info(owner_id: int, video_id: int) -> dict:
+ """Получить info о видео через video.get.
+ Возвращает thumb_url (лучший кадр превью).
+ files.mp4_* и player возвращаются только если токен имеет scope=video;
+ без scope они None — в этом случае скачивание идёт через yt-dlp.
+ """
+ if not VK_TOKEN:
+ return {}
+ params = {
+ "videos": f"{owner_id}_{video_id}",
+ "access_token": VK_TOKEN,
+ "v": VK_VERSION,
+ }
+ try:
+ async with httpx.AsyncClient(timeout=15) as client:
+ r = await client.get(f"{VK_API}/video.get", params=params)
+ data = r.json()
+ if "error" in data:
+ print(f"[VK video.get] {data['error'].get('error_msg')}")
+ return {}
+ items = data.get("response", {}).get("items", [])
+ if not items:
+ return {}
+ v = items[0]
+ result = {}
+ # лучший кадр превью (выше качеством чем в wall.get)
+ for src_key in ("image", "first_frame"):
+ frames = v.get(src_key, [])
+ if frames:
+ best = max(frames, key=lambda f: f.get("width", 0))
+ result["thumb_url"] = best.get("url")
+ break
+ return result
+ except Exception as exc:
+ print(f"[VK] video.get {owner_id}_{video_id}: {exc}")
+ return {}
+
+
+def _ytdlp_extract_url(vk_url: str, max_mb: int = 150) -> str | None:
+ """Синхронно: через yt-dlp достать прямую mp4-ссылку без скачивания файла."""
+ try:
+ import yt_dlp
+ except ImportError:
+ return None
+ max_bytes = max_mb * 1024 * 1024
+ ydl_opts = {
+ "format": "best[ext=mp4][filesize{}]/best[ext=mp4]/best".format(max_bytes),
+ "quiet": True,
+ "no_warnings": True,
+ }
+ try:
+ with yt_dlp.YoutubeDL(ydl_opts) as ydl:
+ info = ydl.extract_info(vk_url, download=False)
+ url = info.get("url") if info else None
+ if url:
+ print(f"[VK] yt-dlp: найден mp4 для {vk_url}")
+ return url
+ except Exception as exc:
+ print(f"[VK] yt-dlp: не удалось получить URL {vk_url}: {exc}")
+ return None
+
+
+async def _download_vk_video(owner_id: int, video_id: int,
+ db: AsyncSession | None = None,
+ max_mb: int = 150) -> str | None:
+ """Скачать VK-видео в MinIO и вернуть публичный URL.
+ Использует yt-dlp для получения прямой mp4-ссылки.
+ Стриминг — не держит весь файл в памяти.
+ """
+ vk_url = f"https://vk.com/video{owner_id}_{video_id}"
+ max_bytes = max_mb * 1024 * 1024
+
+ # yt-dlp работает синхронно — запускаем в пуле потоков
+ loop = asyncio.get_event_loop()
+ mp4_url = await loop.run_in_executor(None, _ytdlp_extract_url, vk_url, max_mb)
+ if not mp4_url:
+ return None
+
+ _chrome_ua = (
+ "Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
+ "AppleWebKit/537.36 (KHTML, like Gecko) "
+ "Chrome/124.0.0.0 Safari/537.36"
+ )
+ try:
+ async with httpx.AsyncClient(
+ timeout=300,
+ follow_redirects=True,
+ headers={"User-Agent": _chrome_ua, "Referer": "https://vk.com/"},
+ ) as client:
+ async with client.stream("GET", mp4_url) as r:
+ if r.status_code != 200:
+ print(f"[VK] Видео HTTP {r.status_code}: {mp4_url[:80]}")
+ return None
+ ct = r.headers.get("content-type", "video/mp4").split(";")[0].strip()
+ # Если content-type не видео — значит получили HTML, не файл
+ if "video" not in ct and "octet" not in ct:
+ print(f"[VK] Неверный content-type видео: {ct}")
+ return None
+ declared = int(r.headers.get("content-length", 0))
+ if declared and declared > max_bytes:
+ print(f"[VK] Видео слишком большое: {declared/1024/1024:.0f}MB > {max_mb}MB, пропускаем")
+ return None
+ chunks, total = [], 0
+ async for chunk in r.aiter_bytes(1024 * 256):
+ total += len(chunk)
+ if total > max_bytes:
+ print(f"[VK] Видео превысило {max_mb}MB при скачивании, пропускаем")
+ return None
+ chunks.append(chunk)
+ data = b"".join(chunks)
+ object_name = f"vk/videos/{uuid.uuid4().hex}.mp4"
+ mc = _minio_client()
+ _ensure_bucket(mc)
+ mc.put_object(MINIO_BUCKET, object_name, io.BytesIO(data), len(data), content_type="video/mp4")
+ public_url = f"{MINIO_PUBLIC_URL}/{MINIO_BUCKET}/{object_name}"
+ print(f"[VK] Видео загружено в MinIO: {object_name} ({total/1024/1024:.1f}MB)")
+ if db is not None:
+ from bd.models import Media
+ db.add(Media(
+ filename=object_name.split("/")[-1],
+ url=public_url,
+ media_type="video",
+ size_bytes=len(data),
+ ))
+ return public_url
+ except Exception as exc:
+ print(f"[VK] Ошибка скачивания видео {vk_url}: {exc}")
+ return None
+
+
+# ── VK content helpers ────────────────────────────────────────────────────────
+
+def _best_photo(sizes: list) -> str | None:
+ for t in ("w", "z", "y", "x", "m", "s"):
+ for s in sizes:
+ if s.get("type") == t:
+ return s["url"]
+ return sizes[-1]["url"] if sizes else None
+
+
+def _text_to_html(text: str) -> str:
+ parts = []
+ for line in text.split("\n"):
+ line = line.strip()
+ if not line:
+ continue
+ line = _URL_RE.sub(r'\1', line)
+ parts.append(f"
{line}
")
+ return "\n".join(parts)
+
+
+def _extract_title(text: str, ts: int) -> str:
+ for line in text.split("\n"):
+ line = line.strip()
+ if line:
+ return line[:200]
+ return "ВКонтакте " + datetime.utcfromtimestamp(ts).strftime("%d.%m.%Y")
+
+
+def _vk_slug(owner_id: int, post_id: int) -> str:
+ return f"vk-{abs(owner_id)}-{post_id}"
+
+
+# ── DB helpers ────────────────────────────────────────────────────────────────
+
+async def _slug_exists(slug: str, db: AsyncSession) -> bool:
+ r = await db.execute(select(Article.id).where(Article.slug == slug))
+ 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()
+ if not cat:
+ cat = Category(name=name, slug=s)
+ db.add(cat)
+ await db.flush()
+ return cat
+
+
+async def _get_or_create_tag(name: str, db: AsyncSession) -> Tag:
+ s = slugify(name)
+ tag = (await db.execute(select(Tag).where(Tag.slug == s))).scalar_one_or_none()
+ if not tag:
+ tag = Tag(name=name, slug=s)
+ db.add(tag)
+ await db.flush()
+ return tag
+
+
+# ── Post processing ───────────────────────────────────────────────────────────
+
+async def _process_post(post: dict, category_name: str, db: AsyncSession) -> 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):
+ return False
+
+ text = post.get("text", "")
+ attachments = post.get("attachments", [])
+ ts = post["date"]
+
+ # VK Clips и некоторые видео-посты: post.text пустой, описание лежит в video.description
+ if not text:
+ for att in attachments:
+ if att.get("type") == "video":
+ desc = att["video"].get("description", "").strip()
+ if desc:
+ text = desc
+ break
+
+ cover_url = None
+ content_parts = []
+
+ if text:
+ content_parts.append(_text_to_html(text))
+
+ for att in attachments:
+ att_type = att.get("type")
+
+ if att_type == "photo":
+ url = _best_photo(att["photo"].get("sizes", []))
+ if url:
+ minio_url = await _upload_from_url(url, db=db)
+ final_url = minio_url or url # fallback на оригинальный URL если MinIO недоступен
+ if cover_url is None:
+ cover_url = final_url
+ else:
+ content_parts.append(
+ f'
'
+ )
+
+ elif att_type == "video":
+ vid = att["video"]
+ player = vid.get("player")
+ vtitle = vid.get("title", "Видео")
+ vid_id = vid.get("id")
+ vid_oid = vid.get("owner_id")
+ vk_video_url = (
+ f"https://vk.com/video{vid_oid}_{vid_id}"
+ if vid_id and vid_oid else None
+ )
+
+ # Кадр превью из wall.get (запасной)
+ thumb_url = None
+ for src_key in ("image", "first_frame"):
+ frames = vid.get(src_key, [])
+ if frames:
+ best = max(frames, key=lambda f: f.get("width", 0))
+ thumb_url = best.get("url")
+ break
+
+ # Улучшенный кадр превью из video.get (+ может вернуть player при scope=video)
+ if vid_id and vid_oid:
+ vinfo = await _fetch_video_info(vid_oid, vid_id)
+ if vinfo.get("thumb_url"):
+ thumb_url = vinfo["thumb_url"]
+
+ # Всегда скачиваем превью в MinIO
+ minio_thumb = None
+ if thumb_url:
+ minio_thumb = await _upload_from_url(thumb_url, db=db)
+ final_thumb = minio_thumb or thumb_url
+ if cover_url is None:
+ cover_url = final_thumb
+ else:
+ final_thumb = None
+
+ # Скачиваем само видео через yt-dlp → MinIO
+ minio_video_url = None
+ if vid_id and vid_oid:
+ minio_video_url = await _download_vk_video(vid_oid, vid_id, db=db)
+
+ if minio_video_url:
+ # 1. Видео лежит у нас — нативный