Baseline md2vk: docs, audit log, docker 8420, openspec, deploy

This commit is contained in:
estorozhenko
2026-09-18 22:05:01 +00:00
commit 09e960a3a9
39 changed files with 3716 additions and 0 deletions
View File
+121
View File
@@ -0,0 +1,121 @@
"""Аудит-лог авторизации и запросов API (JSONL, ротация по дням).
Пишет строку JSON на каждый запрос /api/v1/*:
ts, ip, method, path, api_key_prefix, user_id, status, success, latency_ms, error.
Параметры из env:
- AUDIT_LOG_DIR — каталог для логов (по умолчанию "logs").
"""
from __future__ import annotations
import json
import logging
import os
import time
from datetime import date, datetime, timezone
from pathlib import Path
from starlette.middleware.base import BaseHTTPMiddleware
from starlette.requests import Request
logger = logging.getLogger("md2vk.audit")
class AuditMiddleware(BaseHTTPMiddleware):
"""Логирует каждый запрос /api/v1/* в JSONL с ротацией по дням."""
def __init__(self, app, log_dir: str | None = None):
super().__init__(app)
self.log_dir = Path(log_dir or os.getenv("AUDIT_LOG_DIR", "logs"))
self.log_dir.mkdir(parents=True, exist_ok=True)
self._fh = None
self._fh_date: date | None = None
def _ensure_file(self) -> None:
today = date.today()
if self._fh is None or self._fh_date != today:
if self._fh is not None:
try:
self._fh.close()
except Exception:
pass
path = self.log_dir / f"access.{today.isoformat()}.log"
self._fh = open(path, "a", encoding="utf-8")
self._fh_date = today
def _write(self, record: dict) -> None:
try:
self._ensure_file()
self._fh.write(json.dumps(record, ensure_ascii=False, default=str) + "\n")
self._fh.flush()
except Exception:
# Логгер не должен ронять API
logger.exception("audit write failed")
async def dispatch(self, request: Request, call_next):
start = time.monotonic()
response = None
error = None
try:
response = await call_next(request)
return response
except Exception as exc: # noqa: BLE001
error = str(exc)
raise
finally:
path = request.url.path
if path.startswith("/api/v1"):
try:
latency_ms = round((time.monotonic() - start) * 1000, 1)
status = response.status_code if response is not None else 500
auth = request.headers.get("authorization", "")
# api_key может быть в теле (POST) — пытаемся достать
api_key_prefix = ""
api_key_hash_short = ""
user_id = None
if auth.startswith("Bearer "):
api_key_prefix = auth[len("Bearer "):][:12]
elif request.method == "POST":
api_key_prefix = self._api_key_from_body(request)
if "md2vk_" in api_key_prefix:
# вычислим короткий хэш для привязки к user (без хранения ключа)
import hashlib
api_key_hash_short = hashlib.sha256(
api_key_prefix.encode()
).hexdigest()[:12]
record = {
"ts": datetime.now(timezone.utc).isoformat(timespec="seconds"),
"ip": (request.client.host if request.client else ""),
"method": request.method,
"path": path,
"query": str(request.url.query) or "",
"api_key_prefix": api_key_prefix,
"api_key_hash_short": api_key_hash_short,
"user_id": user_id,
"status": status,
"success": status < 400,
"latency_ms": latency_ms,
"error": error,
}
self._write(record)
except Exception: # noqa: BLE001
logger.exception("audit dispatch failed")
@staticmethod
def _api_key_from_body(request: Request) -> str:
"""Достаёт api_key из JSON-тела, не ломая повторное чтение."""
try:
# starlette кэширует _body — повторное чтение в роутере безопасно
body = getattr(request, "_body", None)
if body is None:
body = request.body() if hasattr(request, "body") else b""
if isinstance(body, bytes) and body:
data = json.loads(body)
key = data.get("api_key", "")
return str(key)[:12]
except Exception:
return ""
return ""
+64
View File
@@ -0,0 +1,64 @@
"""FastAPI-зависимости: аутентификация по API-ключу."""
from __future__ import annotations
import hashlib
from fastapi import Depends, HTTPException, status
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import get_db
from app.models import User
security_scheme = HTTPBearer(auto_error=False)
def hash_api_key(api_key: str) -> str:
return hashlib.sha256(api_key.encode()).hexdigest()
async def get_current_user_from_header(
credentials: HTTPAuthorizationCredentials | None = Depends(security_scheme),
db: AsyncSession = Depends(get_db),
) -> User:
"""Аутентификация по заголовку Authorization: Bearer <api_key>."""
if credentials is None:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Missing Authorization header. Use: Authorization: Bearer <api_key>",
)
api_key = credentials.credentials
if not api_key.startswith("md2vk_"):
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED,
detail="Invalid API key format",
)
api_key_hash = hash_api_key(api_key)
result = await db.execute(
select(User).where(User.api_key_hash == api_key_hash, User.is_active == True)
)
user = result.scalar_one_or_none()
if user is None:
raise HTTPException(status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid API key")
return user
async def get_user_by_api_key(api_key: str, db: AsyncSession) -> User:
"""Проверяет API-ключ из тела запроса. Используется в POST-эндпоинтах."""
if not api_key.startswith("md2vk_"):
raise HTTPException(status_code=401, detail="Invalid API key format")
api_key_hash = hash_api_key(api_key)
result = await db.execute(
select(User).where(User.api_key_hash == api_key_hash, User.is_active == True)
)
user = result.scalar_one_or_none()
if user is None:
raise HTTPException(status_code=401, detail="Invalid API key")
return user
+110
View File
@@ -0,0 +1,110 @@
"""Pydantic-схемы для API-запросов и ответов."""
from __future__ import annotations
from datetime import datetime
from typing import Optional
from pydantic import BaseModel, Field
# ─── Аутентификация ────────────────────────────────────────────────────────
class ApiKeyRequest(BaseModel):
api_key: str = Field(..., description="API-ключ пользователя")
# ─── Аккаунты ──────────────────────────────────────────────────────────────
class AccountCreateRequest(BaseModel):
api_key: str
vk_user_id: int = Field(..., description="VK owner_id (положительный — пользователь, отрицательный — сообщество)")
display_name: str = Field(..., max_length=255, description="Отображаемое имя аккаунта")
access_token: str = Field(..., description="VK OAuth-токен с правами wall")
token_type: str = Field(default="user", pattern="^(user|group)$", description="user | group")
class AccountResponse(BaseModel):
id: int
vk_user_id: int
display_name: str
token_type: str
is_active: bool
created_at: datetime
class AccountListResponse(BaseModel):
accounts: list[AccountResponse]
# ─── Публикация ────────────────────────────────────────────────────────────
class PublishRequest(BaseModel):
api_key: str
vk_account_id: int = Field(..., description="ID VK-аккаунта из /api/v1/accounts")
message_md: str = Field(..., min_length=1, description="Текст поста в Markdown")
publish_date: Optional[datetime] = Field(None, description="ISO datetime для отложенной публикации")
friends_only: bool = False
attachments: Optional[str] = Field(None, description="VK-вложения через запятую (photo123_456, ...)")
signed: Optional[bool] = Field(None, description="Подпись автора (для групп)")
class PublishResponse(BaseModel):
success: bool
publication_id: int
vk_post_id: Optional[int] = None
vk_owner_id: Optional[int] = None
url: Optional[str] = None
error: Optional[str] = None
# ─── Конвертация ───────────────────────────────────────────────────────────
class ConvertRequest(BaseModel):
message_md: str = Field(..., min_length=1, description="Markdown-текст для конвертации")
class FormatItemResponse(BaseModel):
type: str
offset: int
length: int
url: Optional[str] = None
class ConvertResponse(BaseModel):
text: str
format_data: Optional[dict] = None
# ─── Публикации (архив) ────────────────────────────────────────────────────
class PublicationFilterRequest(BaseModel):
api_key: str
vk_account_id: Optional[int] = None
status: Optional[str] = Field(None, pattern="^(draft|scheduled|published|error)$")
limit: int = Field(default=20, ge=1, le=100)
offset: int = Field(default=0, ge=0)
class PublicationItem(BaseModel):
id: int
vk_account_id: int
status: str
markdown_original: str
vk_post_id: Optional[int] = None
vk_owner_id: Optional[int] = None
scheduled_at: Optional[datetime] = None
published_at: Optional[datetime] = None
error_message: Optional[str] = None
created_at: datetime
class PublicationListResponse(BaseModel):
publications: list[PublicationItem]
total: int
# ─── Ошибки ────────────────────────────────────────────────────────────────
class ErrorResponse(BaseModel):
detail: str
+325
View File
@@ -0,0 +1,325 @@
"""API v1: эндпоинты публикации, конвертации, управления аккаунтами."""
from __future__ import annotations
import json
from datetime import datetime
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy import select, func
from sqlalchemy.ext.asyncio import AsyncSession
from app.database import get_db
from app.models import User, VkAccount, Publication
from app.security import decrypt_token, encrypt_token
from app.vk_client import VkClient, VkApiError
from app.converters.markdown_to_vk import markdown_to_vk, format_data_json
from app.api.schemas import (
AccountCreateRequest,
AccountResponse,
AccountListResponse,
PublishRequest,
PublishResponse,
ConvertRequest,
ConvertResponse,
PublicationFilterRequest,
PublicationItem,
PublicationListResponse,
)
from app.api.deps import get_current_user_from_header, get_user_by_api_key
router = APIRouter(prefix="/api/v1", tags=["v1"])
# ─── Health ─────────────────────────────────────────────────────────────────
@router.get("/health")
async def health():
return {"status": "ok"}
# ─── Аккаунты ──────────────────────────────────────────────────────────────
@router.post("/accounts", response_model=AccountResponse)
async def create_account(
req: AccountCreateRequest,
db: AsyncSession = Depends(get_db),
):
"""Добавить VK-аккаунт. Проверяет токен через VK API перед сохранением."""
user = await get_user_by_api_key(req.api_key, db)
# Проверяем токен через VK API
vk = VkClient(req.access_token)
try:
is_valid, display_name = await vk.check_token()
if not is_valid:
raise HTTPException(status_code=400, detail=f"VK token invalid: {display_name}")
except VkApiError as e:
raise HTTPException(status_code=400, detail=f"VK API error: {e}")
finally:
await vk.close()
# Шифруем токен перед сохранением
encrypted = encrypt_token(req.access_token)
name = req.display_name or display_name or f"VK-{req.vk_user_id}"
# Проверяем, нет ли уже такого VK-аккаунта у пользователя
existing = await db.execute(
select(VkAccount).where(
VkAccount.user_id == user.id,
VkAccount.vk_user_id == req.vk_user_id,
VkAccount.is_active == True,
)
)
if existing.scalar_one_or_none():
raise HTTPException(status_code=409, detail="This VK account is already registered")
account = VkAccount(
user_id=user.id,
vk_user_id=req.vk_user_id,
display_name=name,
access_token_enc=encrypted,
token_type=req.token_type,
)
db.add(account)
await db.flush()
await db.refresh(account)
return AccountResponse(
id=account.id,
vk_user_id=account.vk_user_id,
display_name=account.display_name,
token_type=account.token_type,
is_active=account.is_active,
created_at=account.created_at,
)
@router.get("/accounts", response_model=AccountListResponse)
async def list_accounts(
user: User = Depends(get_current_user_from_header),
db: AsyncSession = Depends(get_db),
):
"""Список VK-аккаунтов текущего пользователя."""
result = await db.execute(
select(VkAccount).where(VkAccount.user_id == user.id, VkAccount.is_active == True)
)
accounts = result.scalars().all()
return AccountListResponse(
accounts=[
AccountResponse(
id=a.id,
vk_user_id=a.vk_user_id,
display_name=a.display_name,
token_type=a.token_type,
is_active=a.is_active,
created_at=a.created_at,
)
for a in accounts
]
)
@router.delete("/accounts/{account_id}", status_code=204)
async def delete_account(
account_id: int,
user: User = Depends(get_current_user_from_header),
db: AsyncSession = Depends(get_db),
):
"""Удалить VK-аккаунт (soft delete)."""
result = await db.execute(
select(VkAccount).where(
VkAccount.id == account_id,
VkAccount.user_id == user.id,
)
)
account = result.scalar_one_or_none()
if account is None:
raise HTTPException(status_code=404, detail="Account not found")
account.is_active = False
await db.flush()
# ─── Публикация ────────────────────────────────────────────────────────────
@router.post("/publish", response_model=PublishResponse)
async def publish(
req: PublishRequest,
db: AsyncSession = Depends(get_db),
):
"""Опубликовать пост на стене VK."""
# Аутентификация по api_key из тела запроса
user = await get_user_by_api_key(req.api_key, db)
# Получаем VK-аккаунт
result = await db.execute(
select(VkAccount).where(
VkAccount.id == req.vk_account_id,
VkAccount.user_id == user.id,
VkAccount.is_active == True,
)
)
account = result.scalar_one_or_none()
if account is None:
raise HTTPException(status_code=404, detail="VK account not found")
# Расшифровываем токен (только в памяти!)
token = decrypt_token(account.access_token_enc)
if token is None:
raise HTTPException(status_code=500, detail="Failed to decrypt VK token (key mismatch?)")
# Конвертируем Markdown
chunks = markdown_to_vk(req.message_md)
if not chunks or not chunks[0].text.strip():
raise HTTPException(status_code=400, detail="Empty message after markdown conversion")
chunk = chunks[0]
fd_json = format_data_json(chunk.items) if chunk.items else None
# Создаём запись о публикации
publication = Publication(
vk_account_id=account.id,
status="draft",
markdown_original=req.message_md,
vk_text=chunk.text,
vk_format_data=fd_json,
scheduled_at=req.publish_date,
)
db.add(publication)
await db.flush()
# Если отложенная — сохраняем и выходим
if req.publish_date:
publication.status = "scheduled"
await db.flush()
return PublishResponse(
success=True,
publication_id=publication.id,
)
# Публикуем через VK API
vk = VkClient(token)
try:
result = await vk.wall_post(
message=chunk.text,
owner_id=account.vk_user_id if account.token_type == "group" else None,
from_group=(account.token_type == "group"),
friends_only=req.friends_only,
publish_date=int(req.publish_date.timestamp()) if req.publish_date else None,
attachments=req.attachments,
signed=req.signed,
format_data=chunk.items if chunk.items else None,
)
post_id = result.get("post_id")
owner_id = result.get("owner_id") or account.vk_user_id
publication.status = "published"
publication.vk_post_id = post_id
publication.vk_owner_id = owner_id
publication.published_at = datetime.utcnow()
await db.flush()
url = f"https://vk.com/wall{owner_id}_{post_id}"
return PublishResponse(
success=True,
publication_id=publication.id,
vk_post_id=post_id,
vk_owner_id=owner_id,
url=url,
)
except VkApiError as e:
publication.status = "error"
publication.error_message = str(e)
await db.flush()
return PublishResponse(
success=False,
publication_id=publication.id,
error=str(e),
)
finally:
await vk.close()
@router.post("/convert", response_model=ConvertResponse)
async def convert(req: ConvertRequest):
"""Конвертировать Markdown в VK format_data (без публикации)."""
chunks = markdown_to_vk(req.message_md)
if not chunks:
return ConvertResponse(text="", format_data={"version": 1, "items": []})
chunk = chunks[0]
items = []
for i in chunk.items:
items.append({"type": i.type, "offset": i.offset, "length": i.length, "url": i.url})
return ConvertResponse(
text=chunk.text,
format_data={"version": 1, "items": items},
)
# ─── Архив публикаций ──────────────────────────────────────────────────────
@router.post("/publications", response_model=PublicationListResponse)
async def list_publications(
req: PublicationFilterRequest,
db: AsyncSession = Depends(get_db),
):
"""Архив публикаций с фильтрацией."""
user = await get_user_by_api_key(req.api_key, db)
# Базовый запрос — только публикации пользователя
base_filter = VkAccount.user_id == user.id
query = (
select(Publication)
.join(VkAccount)
.where(base_filter)
)
count_query = (
select(func.count(Publication.id))
.join(VkAccount)
.where(base_filter)
)
if req.vk_account_id is not None:
query = query.where(Publication.vk_account_id == req.vk_account_id)
count_query = count_query.where(Publication.vk_account_id == req.vk_account_id)
if req.status is not None:
query = query.where(Publication.status == req.status)
count_query = count_query.where(Publication.status == req.status)
total_result = await db.execute(count_query)
total = total_result.scalar() or 0
query = query.order_by(Publication.created_at.desc()).offset(req.offset).limit(req.limit)
result = await db.execute(query)
publications = result.scalars().all()
return PublicationListResponse(
publications=[
PublicationItem(
id=p.id,
vk_account_id=p.vk_account_id,
status=p.status,
markdown_original=p.markdown_original,
vk_post_id=p.vk_post_id,
vk_owner_id=p.vk_owner_id,
scheduled_at=p.scheduled_at,
published_at=p.published_at,
error_message=p.error_message,
created_at=p.created_at,
)
for p in publications
],
total=total,
)