from __future__ import annotations import asyncio import json import base64 import hashlib import mimetypes import re import secrets import time from datetime import date, datetime, timedelta, timezone from pathlib import Path from typing import Any from urllib.parse import urlencode from zoneinfo import ZoneInfo import aiohttp from fastapi import FastAPI, File, Form, Request, UploadFile, status from fastapi.responses import HTMLResponse, RedirectResponse from fastapi.staticfiles import StaticFiles from fastapi.templating import Jinja2Templates from loguru import logger from .config import settings from .constants import PLATFORM_VK from .db import fetch_setting, get_pool from .security import hash_password, new_token, token_hash, verify_password from .text_utils import normalize_hash_tag, parse_categories from .vk_api import VKAPIClient, normalize_vk_source from .workers.ai_qualifier import normalize_model, response_usage COOKIE_NAME = "vk_parser_admin" VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier" VK_OAUTH_STATE_COOKIE = "vk_oauth_state" BASE_DIR = Path(__file__).resolve().parent templates = Jinja2Templates(directory=str(BASE_DIR / "templates")) app = FastAPI(title="VK Parser Admin") LOCAL_TZ = ZoneInfo("Asia/Yekaterinburg") MODEL_CACHE: dict[str, Any] = {"key": "", "at": 0.0, "models": []} UPLOAD_ROOT = Path("uploads").resolve() EDITOR_MEDIA_DIR = UPLOAD_ROOT / "editor_media" EDITOR_MEDIA_DIR.mkdir(parents=True, exist_ok=True) app.mount("/uploads", StaticFiles(directory=str(UPLOAD_ROOT)), name="uploads") PROVIDER_OPTIONS = [ {"value": "openrouter", "label": "OpenRouter"}, {"value": "openai", "label": "OpenAI"}, {"value": "anthropic", "label": "Anthropic"}, {"value": "gemini", "label": "Google Gemini"}, {"value": "openai_compatible", "label": "OpenAI-compatible API"}, ] PROMPT_HINTS = { "ai_qualifier_prompt": { "title": "Свободная инструкция квалификатора", "body": ( "Здесь можно спокойно менять смысл: роль, тон, критерии пригодности, примеры хороших и плохих постов, " "логику оценки по шкале 1-10." ), "safe": ( "Можно менять: критерии, примеры, тон причины, пороговые объяснения. " "Не добавляй сюда JSON-схему: для неё есть отдельный технический контракт ниже." ), "input": "На вход уходит JSON-массив постов: id, source, original_url, media_count, media_types, text.", }, "ai_qualifier_contract": { "title": "Технический контракт квалификатора", "body": ( "Эта часть защищает парсер от сломанного ответа. Меняй её только если осознанно меняешь формат ответа " "и одновременно готов править валидатор в коде." ), "safe": "Обычно не трогаем. Здесь живут JSON-схема, допустимые decision/reject_tag и правило вернуть результат на каждый id.", "input": "Контракт склеивается после свободной инструкции и отправляется как system prompt.", "contract": '{"results":[{"id":123,"score":8,"decision":"accepted","reason":"до 10 слов на русском","reject_tag":null}]}', }, "ai_writer_prompt": { "title": "Свободная инструкция райтера", "body": ( "Здесь можно менять редакторскую часть: стиль, длину, голос канала, запреты на выдумки, примеры хорошего текста." ), "safe": ( "Можно менять: tone of voice, правила фактов, длину, примеры фраз. " "Хэштеги руками писать не надо: модель выбирает category, а код сам соберёт #category и #producer_tag." ), "input": "На вход уходит JSON-объект: categories и posts. В каждом post есть id, producer_name, producer_tag, original_url, qualification_score, media_count, media_types, text.", }, "ai_writer_contract": { "title": "Технический контракт райтера", "body": ( "Эта часть фиксирует формат JSON, запрет хэштегов в тексте и выбор категории. " "Менять её стоит только вместе с валидатором райтера." ), "safe": "Обычно не трогаем. Если промпт райтера переписывается, свободную часть меняем выше, контракт оставляем стабильным.", "input": "Контракт склеивается после свободной инструкции и отправляется как system prompt.", "contract": '{"rewrites":[{"id":123,"category":"разгрузка","text":"готовый текст без хэштегов","notes":"короткая заметка для редактора"}]}', }, } MODEL_FALLBACKS = { "openrouter": [ "openrouter/openai/gpt-4.1-mini", "openrouter/anthropic/claude-3.5-sonnet", "openrouter/google/gemini-2.5-flash", ], "openai": ["gpt-4.1-mini", "gpt-4o-mini", "o4-mini"], "anthropic": ["anthropic/claude-haiku-4-5-20251001", "anthropic/claude-sonnet-4-20250514"], "gemini": ["gemini/gemini-2.5-flash", "gemini/gemini-2.5-pro"], "openai_compatible": [], } PREFERRED_MODELS = { "openrouter": [ "openrouter/anthropic/claude-sonnet-4", "openrouter/anthropic/claude-3.5-sonnet", "openrouter/openai/gpt-4.1-mini", "openrouter/google/gemini-2.5-flash", ], "anthropic": [ "anthropic/claude-haiku-4-5-20251001", "anthropic/claude-4-sonnet-20250514", "anthropic/claude-sonnet-4-20250514", "anthropic/claude-3-7-sonnet-20250219", "anthropic/claude-3-5-sonnet-20241022", "anthropic/claude-3-5-haiku-20241022", ], "openai": ["gpt-4.1-mini", "gpt-4o-mini", "o4-mini"], "gemini": ["gemini/gemini-2.5-flash", "gemini/gemini-2.5-pro"], } CATEGORY_TITLES = { "AI Qualifier": "AI-квалификатор", "AI Writer": "AI-райтер", "Parser": "Парсер", "Uploader": "Аплоадер", "VK": "VK API", "General": "Общие", } CATEGORY_ORDER = { "AI Qualifier": 10, "AI Writer": 20, "Parser": 30, "VK": 40, "Uploader": 50, "General": 100, } SETTING_ORDER = { "AI Qualifier": [ "ai_qualifier_enabled", "ai_qualifier_provider", "ai_qualifier_model", "ai_qualifier_api_key", "ai_qualifier_api_base", "ai_qualifier_prompt", "ai_qualifier_contract", "ai_qualifier_batch_size", "ai_qualifier_min_score", "ai_qualifier_max_text_chars", "ai_qualifier_temperature", "ai_qualifier_timeout_sec", "ai_qualifier_interval_sec", ], "AI Writer": [ "ai_writer_enabled", "ai_writer_provider", "ai_writer_model", "ai_writer_api_key", "ai_writer_api_base", "ai_writer_prompt", "ai_writer_contract", "ai_writer_categories", "ai_writer_batch_size", "ai_writer_max_text_chars", "ai_writer_temperature", "ai_writer_timeout_sec", "ai_writer_interval_sec", ], "Parser": [ "parser_interval_sec", "parser_new_source_lookback_days", "parser_reparse_overlap_minutes", "parser_min_text_length", "parser_skip_empty_text", "parser_skip_no_media", "parser_skip_text_too_short", "parser_skip_reposts", "parser_store_skipped_posts", "parser_dedupe_content_hash", "parser_source_pause_sec", ], "VK": [ "vk_requests_per_second", "vk_wall_page_size", "vk_api_timeout_total_sec", "vk_api_timeout_connect_sec", "vk_rate_limit_sleep_sec", "vk_api_retry_attempts", "vk_api_retry_min_delay_sec", "vk_api_retry_max_delay_sec", ], "Uploader": [ "tg_media_channel_id", "local_bot_api_url", "uploader_interval_sec", "uploader_download_timeout_sec", "media_group_max_items", "media_upload_delay_sec", "media_post_job_pause_sec", "max_media_attempts", "tg_retry_attempts", "tg_retry_backoff_max_sec", "video_max_size_mb", "video_max_duration_sec", "uploader_yt_dlp_timeout_sec", ], } def now_utc() -> datetime: return datetime.now(timezone.utc) def redirect(path: str) -> RedirectResponse: return RedirectResponse(path, status_code=status.HTTP_303_SEE_OTHER) @app.middleware("http") async def add_security_headers(request: Request, call_next): response = await call_next(request) response.headers.setdefault("X-Frame-Options", "DENY") response.headers.setdefault("X-Content-Type-Options", "nosniff") response.headers.setdefault("Referrer-Policy", "same-origin") response.headers.setdefault("Permissions-Policy", "geolocation=(), microphone=(), camera=()") return response def pkce_challenge(verifier: str) -> str: digest = hashlib.sha256(verifier.encode("ascii")).digest() return base64.urlsafe_b64encode(digest).decode("ascii").rstrip("=") async def bootstrap_admin() -> None: if not settings.admin_bootstrap_login or not settings.admin_bootstrap_password: return pool = await get_pool() exists = await pool.fetchval("SELECT 1 FROM admin_users WHERE login=$1", settings.admin_bootstrap_login) if exists: return await pool.execute( """ INSERT INTO admin_users(login, password_hash, role, is_active) VALUES($1, $2, 'admin', TRUE) """, settings.admin_bootstrap_login.strip().lower(), hash_password(settings.admin_bootstrap_password), ) logger.warning("Bootstrap admin user created: {}", settings.admin_bootstrap_login) async def get_current_user(request: Request) -> dict | None: token = request.cookies.get(COOKIE_NAME) if not token: return None pool = await get_pool() row = await pool.fetchrow( """ SELECT s.csrf_token, u.id, u.login, u.role, u.is_active FROM admin_sessions s JOIN admin_users u ON u.id=s.user_id WHERE s.token_hash=$1 AND s.expires_at > NOW() AND u.is_active=TRUE """, token_hash(token, settings.app_secret_key), ) return dict(row) if row else None def require_csrf(user: dict, csrf_token: str) -> None: if not user or csrf_token != user["csrf_token"]: raise PermissionError("bad csrf") def base_context(request: Request, user: dict | None, **extra: Any) -> dict[str, Any]: ctx = {"request": request, "user": user, "app_env": settings.app_env} ctx.update(extra) return ctx def client_ip(request: Request) -> str: forwarded = request.headers.get("x-forwarded-for", "") if forwarded: return forwarded.split(",", 1)[0].strip()[:80] real_ip = request.headers.get("x-real-ip", "") if real_ip: return real_ip.strip()[:80] return (request.client.host if request.client else "unknown")[:80] def preserved_query(request: Request, exclude: set[str] | None = None) -> list[dict[str, str]]: exclude = exclude or {"page"} items = [] for key, value in request.query_params.multi_items(): if key not in exclude: items.append({"key": key, "value": value}) return items def query_path(request: Request, **replace: str) -> str: values: list[tuple[str, str]] = [] excluded = set(replace) | {"page"} for key, value in request.query_params.multi_items(): if key not in excluded: values.append((key, value)) for key, value in replace.items(): if value != "": values.append((key, value)) encoded = urlencode(values) return f"{request.url.path}?{encoded}" if encoded else request.url.path def raw_sort_headers(request: Request, current_sort: str) -> dict[str, dict[str, str | bool]]: columns = { "id": ("id_asc", "id_desc"), "source": ("source_asc", "source_desc"), "stage": ("stage_asc", "stage_desc"), "post": ("text_asc", "text_desc"), "ai": ("score_asc", "score_desc"), } headers: dict[str, dict[str, str | bool]] = {} for key, (asc, desc) in columns.items(): next_sort = desc if current_sort == asc else asc headers[key] = { "url": query_path(request, sort=next_sort), "active": current_sort in {asc, desc}, "direction": "asc" if current_sort == asc else "desc" if current_sort == desc else "", } return headers def parse_date_filter(value: str) -> date | None: value = (value or "").strip() if not value: return None try: return datetime.strptime(value, "%Y-%m-%d").date() except ValueError: return None def format_dt(value: datetime | None) -> str: if not value: return "" if value.tzinfo is None: value = value.replace(tzinfo=timezone.utc) return value.astimezone(LOCAL_TZ).strftime("%d.%m.%y %H:%M") def clamp_int(value: int, min_value: int, max_value: int) -> int: return max(min_value, min(max_value, int(value))) def pagination(page: int, per_page: int, total: int) -> dict[str, int | bool]: page = clamp_int(page, 1, 1_000_000) per_page = clamp_int(per_page, 10, 200) pages = max(1, (int(total) + per_page - 1) // per_page) page = min(page, pages) return { "page": page, "per_page": per_page, "total": int(total), "pages": pages, "offset": (page - 1) * per_page, "has_prev": page > 1, "has_next": page < pages, } FILTER_OPS = { "text": [ {"value": "contains", "label": "содержит"}, {"value": "eq", "label": "равно"}, {"value": "neq", "label": "не равно"}, ], "enum": [ {"value": "eq", "label": "равно"}, {"value": "neq", "label": "не равно"}, ], "number": [ {"value": "eq", "label": "="}, {"value": "gte", "label": ">="}, {"value": "lte", "label": "<="}, ], "date": [ {"value": "gte", "label": "от"}, {"value": "lte", "label": "до"}, {"value": "eq", "label": "день"}, ], } FILTER_OP_OPTIONS = [ {"value": "contains", "label": "содержит"}, {"value": "eq", "label": "равно / день"}, {"value": "neq", "label": "не равно"}, {"value": "gte", "label": ">= / от"}, {"value": "lte", "label": "<= / до"}, ] SOURCE_FIELD_SPECS = { "id": {"label": "ID", "expr": "s.id", "type": "number"}, "platform": {"label": "Площадка", "expr": "s.platform", "type": "enum"}, "name": {"label": "Название", "expr": "s.name", "type": "text"}, "tag": {"label": "Тэг", "expr": "COALESCE(s.tag,'')", "type": "text"}, "url": {"label": "Ссылка", "expr": "s.url", "type": "text"}, "external_id": {"label": "VK id", "expr": "COALESCE(s.external_id,'')", "type": "text"}, "active": {"label": "Включён", "expr": "s.active", "type": "enum"}, "status": {"label": "Статус", "expr": "s.status", "type": "enum"}, "posts_count": { "label": "Постов", "expr": "(SELECT COUNT(*) FROM raw_posts rp2 WHERE rp2.source_id=s.id)", "type": "number", }, "posts_24h": { "label": "Постов за 24ч", "expr": "(SELECT COUNT(*) FROM raw_posts rp3 WHERE rp3.source_id=s.id AND rp3.created_at > NOW() - INTERVAL '24 hours')", "type": "number", }, "last_parsed_at": {"label": "Последний парсинг", "expr": "timezone('Asia/Yekaterinburg', s.last_parsed_at)::date", "sort_expr": "s.last_parsed_at", "type": "date"}, "created_at": {"label": "Создан", "expr": "timezone('Asia/Yekaterinburg', s.created_at)::date", "sort_expr": "s.created_at", "type": "date"}, } RAW_FIELD_SPECS = { "id": {"label": "ID", "expr": "rp.id", "type": "number"}, "source_name": {"label": "Источник", "expr": "s.name", "type": "text"}, "source_tag": {"label": "Тэг источника", "expr": "COALESCE(s.tag,'')", "type": "text"}, "platform": {"label": "Площадка", "expr": "rp.platform", "type": "enum"}, "status": {"label": "Raw статус", "expr": "rp.status", "type": "enum"}, "original_url": {"label": "Оригинал", "expr": "rp.original_url", "type": "text"}, "text": {"label": "Текст", "expr": "rp.raw_text", "type": "text"}, "media_count": { "label": "Медиа", "expr": "(SELECT COUNT(*) FROM raw_post_media rmf WHERE rmf.raw_post_id=rp.id)", "type": "number", }, "posted_at": {"label": "Дата VK", "expr": "timezone('Asia/Yekaterinburg', rp.posted_at)::date", "sort_expr": "rp.posted_at", "type": "date"}, "created_at": {"label": "Дата загрузки", "expr": "timezone('Asia/Yekaterinburg', rp.created_at)::date", "sort_expr": "rp.created_at", "type": "date"}, "qualification_status": {"label": "Итог AI", "expr": "COALESCE(rp.qualification_status,'pending')", "type": "enum"}, "qualification_model_decision": {"label": "Решение модели", "expr": "COALESCE(rp.qualification_model_decision, rp.qualification_decision, '')", "type": "enum"}, "qualification_score": {"label": "Оценка", "expr": "rp.qualification_score", "type": "number"}, "qualification_reject_tag": {"label": "Reject tag", "expr": "COALESCE(rp.qualification_reject_tag,'')", "type": "enum"}, "rewrite_status": {"label": "Статус рерайта", "expr": "COALESCE(rp.rewrite_status,'pending')", "type": "enum"}, "rewrite_category": {"label": "Категория", "expr": "COALESCE(rp.rewrite_category,'')", "type": "enum"}, "rewrite_source_tag": {"label": "Хэштег источника", "expr": "COALESCE(rp.rewrite_source_tag,'')", "type": "text"}, } EDITOR_STATUS_OPTIONS = [ {"value": "review", "label": "На проверке"}, {"value": "edited", "label": "Отредактировано"}, {"value": "accepted", "label": "Принято"}, {"value": "rejected", "label": "Отклонено"}, ] def field_options(specs: dict[str, dict[str, str]]) -> list[dict[str, str]]: return [{"value": key, "label": spec["label"], "type": spec["type"]} for key, spec in specs.items()] def parse_table_filters(request: Request, specs: dict[str, dict[str, str]], slots: int = 4) -> list[dict[str, str]]: qp = request.query_params fields = qp.getlist("f_field") ops = qp.getlist("f_op") values = qp.getlist("f_value") rows: list[dict[str, str]] = [] for i in range(max(slots, len(fields))): field = fields[i] if i < len(fields) else "" op = ops[i] if i < len(ops) else "" value = values[i] if i < len(values) else "" if field not in specs: field = "" field_type = specs[field]["type"] if field else "text" allowed_ops = {item["value"] for item in FILTER_OPS[field_type]} if op not in allowed_ops: op = "contains" if field_type == "text" else "eq" rows.append({"field": field, "op": op, "value": value}) return rows def parse_table_sorts(request: Request, specs: dict[str, dict[str, str]], slots: int = 3) -> list[dict[str, str]]: qp = request.query_params fields = qp.getlist("sort_field") dirs = qp.getlist("sort_dir") rows: list[dict[str, str]] = [] for i in range(max(slots, len(fields))): field = fields[i] if i < len(fields) else "" direction = dirs[i] if i < len(dirs) else "desc" if field not in specs: field = "" if direction not in {"asc", "desc"}: direction = "desc" rows.append({"field": field, "dir": direction}) return rows def build_filter_sql(filters: list[dict[str, str]], specs: dict[str, dict[str, str]], start_index: int = 1) -> tuple[list[str], list[Any]]: clauses: list[str] = [] args: list[Any] = [] index = start_index for item in filters: field = item.get("field") or "" value = str(item.get("value") or "").strip() if not field or field not in specs or value == "": continue spec = specs[field] expr = spec["expr"] field_type = spec["type"] op = item.get("op") or ("contains" if field_type == "text" else "eq") if field_type == "text": if op == "eq": clauses.append(f"{expr} = ${index}") args.append(value) elif op == "neq": clauses.append(f"{expr} <> ${index}") args.append(value) else: clauses.append(f"{expr} ILIKE '%' || ${index} || '%'") args.append(value) elif field_type == "enum": if op == "neq": clauses.append(f"{expr} <> ${index}") else: clauses.append(f"{expr} = ${index}") if value.lower() in {"true", "false"}: args.append(value.lower() == "true") else: args.append(value) elif field_type == "number": try: number_value = int(value) except ValueError: continue sign = {"gte": ">=", "lte": "<=", "neq": "<>", "eq": "="}.get(op, "=") clauses.append(f"{expr} {sign} ${index}") args.append(number_value) elif field_type == "date": sign = {"gte": ">=", "lte": "<=", "eq": "="}.get(op, "=") clauses.append(f"{expr} {sign} ${index}::date") args.append(value) index += 1 return clauses, args def build_sort_sql(sorts: list[dict[str, str]], specs: dict[str, dict[str, str]], default_sql: str) -> str: parts: list[str] = [] seen: set[str] = set() for item in sorts: field = item.get("field") or "" if not field or field not in specs or field in seen: continue spec = specs[field] expr = spec.get("sort_expr") or spec["expr"] direction = "ASC" if item.get("dir") == "asc" else "DESC" nulls = " NULLS LAST" if direction == "DESC" else " NULLS FIRST" parts.append(f"{expr} {direction}{nulls}") seen.add(field) if parts: if "id" not in seen and "id" in specs: parts.append(f"{specs['id'].get('sort_expr') or specs['id']['expr']} DESC") return ", ".join(part for part in parts if part) return default_sql def selected_values(request: Request, name: str) -> list[str]: return [str(value).strip() for value in request.query_params.getlist(name) if str(value).strip()] def selected_int_values(request: Request, name: str) -> list[int]: values: list[int] = [] for value in selected_values(request, name): try: values.append(int(value)) except ValueError: continue return values def add_where(clauses: list[str], args: list[Any], sql_template: str, value: Any) -> None: args.append(value) clauses.append(sql_template.format(i=len(args))) async def raw_filter_facets(pool) -> dict[str, Any]: source_rows = await pool.fetch( """ SELECT s.id, s.name, COALESCE(s.tag, '') AS tag, COUNT(rp.id) AS count FROM raw_posts rp JOIN sources s ON s.id=rp.source_id GROUP BY s.id, s.name, s.tag ORDER BY COUNT(rp.id) DESC, s.name ASC """ ) category_rows = await pool.fetch( """ SELECT COALESCE(final_category, rewrite_category, '') AS value, COUNT(*) AS count FROM raw_posts WHERE COALESCE(final_category, rewrite_category, '') <> '' GROUP BY COALESCE(final_category, rewrite_category, '') ORDER BY COUNT(*) DESC, value ASC """ ) status_rows = await pool.fetch("SELECT status AS value, COUNT(*) AS count FROM raw_posts GROUP BY status ORDER BY count DESC, value") qualification_rows = await pool.fetch( """ SELECT COALESCE(qualification_status, 'pending') AS value, COUNT(*) AS count FROM raw_posts GROUP BY COALESCE(qualification_status, 'pending') ORDER BY count DESC, value """ ) rewrite_rows = await pool.fetch( """ SELECT COALESCE(rewrite_status, 'pending') AS value, COUNT(*) AS count FROM raw_posts GROUP BY COALESCE(rewrite_status, 'pending') ORDER BY count DESC, value """ ) return { "sources": [dict(row) for row in source_rows], "categories": [dict(row) for row in category_rows], "statuses": [dict(row) for row in status_rows], "qualification_statuses": [dict(row) for row in qualification_rows], "rewrite_statuses": [dict(row) for row in rewrite_rows], } async def editor_filter_facets(pool) -> dict[str, Any]: source_rows = await pool.fetch( """ SELECT s.id, s.name, COALESCE(s.tag, '') AS tag, COUNT(rp.id) AS count FROM raw_posts rp JOIN sources s ON s.id=rp.source_id WHERE rp.rewrite_status='ready' GROUP BY s.id, s.name, s.tag ORDER BY COUNT(rp.id) DESC, s.name ASC """ ) category_rows = await pool.fetch( """ SELECT COALESCE(final_category, rewrite_category, '') AS value, COUNT(*) AS count FROM raw_posts WHERE rewrite_status='ready' AND COALESCE(final_category, rewrite_category, '') <> '' GROUP BY COALESCE(final_category, rewrite_category, '') ORDER BY COUNT(*) DESC, value ASC """ ) status_rows = await pool.fetch( """ SELECT COALESCE(editorial_status, 'review') AS value, COUNT(*) AS count FROM raw_posts WHERE rewrite_status='ready' GROUP BY COALESCE(editorial_status, 'review') ORDER BY count DESC, value """ ) return { "sources": [dict(row) for row in source_rows], "categories": [dict(row) for row in category_rows], "statuses": [dict(row) for row in status_rows], } def parse_source_line(line: str) -> dict[str, str]: raw = (line or "").strip() if not raw: return {"name": "", "tag": "", "url": "", "error": "Пустая строка"} parts = re.split(r"\s+", raw) url_index = next((i for i, part in enumerate(parts) if "vk.com/" in part or part.startswith("club")), -1) if url_index < 0: if len(parts) == 1: value = parts[0].strip() url = value if value.startswith("http") else f"https://vk.com/{value}" return {"name": "", "tag": "", "url": url, "error": ""} return {"name": "", "tag": "", "url": "", "error": "Не нашёл ссылку VK"} url = parts[url_index].strip() before_url = parts[:url_index] if len(before_url) >= 2: tag = before_url[-1].strip() name = " ".join(before_url[:-1]).strip() elif len(before_url) == 1: name = before_url[0].strip() tag = before_url[0].strip() else: name = "" tag = "" if url.startswith("vk.com/"): url = f"https://{url}" return {"name": name, "tag": normalize_hash_tag(tag, ""), "url": url, "error": ""} async def build_sources_preview(lines_text: str, default_active: bool = True) -> list[dict[str, Any]]: lines = [line for line in (lines_text or "").splitlines() if line.strip()] parsed = [parse_source_line(line) for line in lines] normalized_values = [normalize_vk_source(item["url"]) for item in parsed if item.get("url")] pool = await get_pool() existing_rows = await pool.fetch( """ SELECT lower(external_id) AS external_id, lower(url) AS url, lower(COALESCE(tag,'')) AS tag FROM sources WHERE archived_at IS NULL AND (lower(external_id)=ANY($1::text[]) OR lower(url)=ANY($2::text[]) OR lower(COALESCE(tag,''))=ANY($3::text[])) """, [v.lower() for v in normalized_values if v], [str(item.get("url") or "").lower() for item in parsed], [normalize_hash_tag(str(item.get("tag") or ""), "").lower() for item in parsed if item.get("tag")], ) existing_ids = {str(row["external_id"] or "").lower() for row in existing_rows} existing_urls = {str(row["url"] or "").lower() for row in existing_rows} existing_tags = {str(row["tag"] or "").lower() for row in existing_rows} seen: set[str] = set() seen_tags: set[str] = set() async with VKAPIClient(rps=2, timeout_total_sec=12, timeout_connect_sec=5, retry_attempts=2) as client: preview = [] for idx, item in enumerate(parsed, start=1): url = item.get("url") or "" external_id = normalize_vk_source(url) if url else "" row = { "line_no": idx, "platform": PLATFORM_VK, "name": item.get("name") or "", "tag": normalize_hash_tag(item.get("tag") or external_id, external_id or "source"), "url": url, "external_id": external_id, "external_owner_id": None, "active": default_active, "ok": False, "error": item.get("error") or "", } key = external_id.lower() if not row["error"] and key in seen: row["error"] = "Дубль в этом списке" if not row["error"] and row["tag"].lower() in seen_tags: row["error"] = "Дубль тэга в этом списке" if not row["error"] and (key in existing_ids or url.lower() in existing_urls): row["error"] = "Уже есть в источниках" if not row["error"] and row["tag"].lower() in existing_tags: row["error"] = "Такой тэг уже есть" if not row["error"] and not external_id: row["error"] = "Не удалось разобрать VK-ссылку" if not row["error"]: try: screen_name, owner_id, resolved_name = await client.resolve_group(url) row["external_id"] = screen_name row["external_owner_id"] = owner_id row["name"] = row["name"] or resolved_name row["tag"] = normalize_hash_tag(row["tag"] or screen_name, screen_name) row["url"] = f"https://vk.com/{screen_name}" row["ok"] = True seen.add(key) seen_tags.add(row["tag"].lower()) except Exception as exc: row["error"] = str(exc) preview.append(row) return preview def pipeline_stage(post: dict[str, Any]) -> dict[str, str]: status = str(post.get("status") or "") qualification_status = str(post.get("qualification_status") or "").strip().lower() if status == "failed": return {"key": "failed", "label": "Ошибка", "class": "bad"} if status == "skipped": return {"key": "skipped", "label": "Пропущен", "class": "muted"} if status in {"raw_saved", "storage_pending"}: return {"key": status, "label": "Медиа", "class": ""} if status != "storage_ready": return {"key": status, "label": status or "Raw", "class": ""} rewrite_status = str(post.get("rewrite_status") or "").strip().lower() if rewrite_status == "ready": return {"key": "rewrite_ready", "label": "Рерайт готов", "class": "ok"} if rewrite_status == "processing": return {"key": "rewrite_processing", "label": "AI пишет", "class": "warn"} if rewrite_status == "failed": return {"key": "rewrite_failed", "label": "Рерайт ошибка", "class": "bad"} if qualification_status in {"accepted"}: return {"key": "qualified_accepted", "label": "AI принят", "class": "ok"} if qualification_status in {"rejected"}: return {"key": "qualified_rejected", "label": "AI отклонён", "class": "bad"} if qualification_status == "processing": return {"key": "qualification_processing", "label": "AI проверяет", "class": "warn"} if qualification_status == "failed": return {"key": "qualification_failed", "label": "AI ошибка", "class": "bad"} return {"key": "qualification_pending", "label": "Ждёт AI", "class": ""} def ai_badge(post: dict[str, Any]) -> dict[str, str]: status = str(post.get("qualification_status") or "pending").strip().lower() score = post.get("qualification_score") decision = str(post.get("qualification_decision") or status or "pending") reason = str(post.get("qualification_reason") or "") if isinstance(score, int): if score >= 8: css = "ok" elif score >= 5: css = "warn" else: css = "bad" label = f"{score}/10" sublabel = decision elif status == "processing": css = "warn" label = "..." sublabel = "processing" elif status == "failed": css = "bad" label = "!" sublabel = "failed" else: css = "" label = "AI" sublabel = "pending" return {"class": css, "label": label, "sublabel": sublabel, "reason": reason} def prepare_raw_post(row: Any) -> dict[str, Any]: post = dict(row) media_items = post.get("media_items") or [] if isinstance(media_items, str): try: media_items = json.loads(media_items) except json.JSONDecodeError: media_items = [] post["media_items"] = media_items post["posted_at_fmt"] = format_dt(post.get("posted_at")) post["created_at_fmt"] = format_dt(post.get("created_at")) post["qualified_at_fmt"] = format_dt(post.get("qualified_at")) post["stage"] = pipeline_stage(post) post["ai_badge"] = ai_badge(post) return post def prepare_editor_post(row: Any) -> dict[str, Any]: post = prepare_raw_post(row) post["review_text"] = post.get("final_text") or post.get("rewritten_text") or "" post["review_category"] = post.get("final_category") or post.get("rewrite_category") or "" post["review_category_tag"] = post.get("final_category_tag") or post.get("rewrite_category_tag") or "" post["review_source_tag"] = post.get("final_source_tag") or post.get("rewrite_source_tag") or post.get("source_tag") or "" post["editorial_status_label"] = next( (item["label"] for item in EDITOR_STATUS_OPTIONS if item["value"] == (post.get("editorial_status") or "review")), post.get("editorial_status") or "review", ) post["reviewed_at_fmt"] = format_dt(post.get("reviewed_at")) post["edited_at_fmt"] = format_dt(post.get("edited_at")) return post def uploaded_editor_media_path(url: str) -> Path | None: raw = str(url or "").strip() prefix = "/uploads/editor_media/" if not raw.startswith(prefix): return None name = Path(raw.removeprefix(prefix)).name if not name: return None path = (EDITOR_MEDIA_DIR / name).resolve() try: path.relative_to(EDITOR_MEDIA_DIR.resolve()) except ValueError: return None return path async def writer_categories() -> list[str]: pool = await get_pool() rows = await pool.fetch( """ SELECT name FROM content_categories WHERE is_active=TRUE ORDER BY sort_order, name """ ) categories = [str(row["name"]) for row in rows] if categories: return categories legacy_categories = parse_categories(await fetch_setting("ai_writer_categories", [])) return legacy_categories or [ "защита", "одежда", "разгрузка", "рюкзаки", "airsoft", "патчи", "электроника", "аксессуары", "производство", ] async def category_rows() -> list[dict[str, Any]]: pool = await get_pool() rows = await pool.fetch( """ SELECT id, name, tag, is_active, sort_order, created_at, updated_at FROM content_categories ORDER BY is_active DESC, sort_order, name """ ) result = [] for row in rows: item = dict(row) item["created_at_fmt"] = format_dt(item.get("created_at")) item["updated_at_fmt"] = format_dt(item.get("updated_at")) result.append(item) return result PROMPT_TEST_STATUS_OPTIONS = [ {"value": "all", "label": "Все"}, {"value": "qualified", "label": "Оцененные"}, {"value": "accepted", "label": "Принятые"}, {"value": "maybe", "label": "Спорные"}, {"value": "rejected", "label": "Отклоненные"}, {"value": "pending", "label": "Ждут AI"}, ] def normalize_prompt_test_status(value: str | None) -> str: value = str(value or "all").strip().lower() allowed = {item["value"] for item in PROMPT_TEST_STATUS_OPTIONS} return value if value in allowed else "all" def prompt_test_status_where(status_filter: str) -> str: if status_filter == "qualified": return "AND rp.qualification_status IN ('accepted', 'maybe', 'rejected')" if status_filter in {"accepted", "maybe", "rejected", "pending"}: return f"AND rp.qualification_status = '{status_filter}'" return "" async def prompt_test_posts( selected_post_id: int | None = None, status_filter: str = "all", ) -> tuple[list[dict[str, Any]], dict[str, Any] | None]: pool = await get_pool() status_filter = normalize_prompt_test_status(status_filter) where_sql = prompt_test_status_where(status_filter) rows = await pool.fetch( f""" SELECT rp.id, rp.raw_text, rp.original_url, rp.posted_at, rp.created_at, rp.qualification_status, rp.qualification_score, s.name AS source_name, s.tag AS source_tag, rp.qualification_reason, (SELECT COUNT(*) FROM raw_post_media m WHERE m.raw_post_id=rp.id) AS media_count, (SELECT ARRAY_AGG(DISTINCT m.media_type ORDER BY m.media_type) FROM raw_post_media m WHERE m.raw_post_id=rp.id) AS media_types FROM raw_posts rp JOIN sources s ON s.id=rp.source_id WHERE TRUE {where_sql} ORDER BY rp.created_at DESC, rp.id DESC LIMIT 150 """ ) posts = [] for row in rows: item = dict(row) text = str(item.get("raw_text") or "").strip() item["created_at_fmt"] = format_dt(item.get("created_at")) item["posted_at_fmt"] = format_dt(item.get("posted_at")) item["label"] = f"#{item['id']} · {item.get('source_name') or 'source'} · {item['created_at_fmt']}" item["snippet"] = text[:180] + ("..." if len(text) > 180 else "") posts.append(item) selected = next((item for item in posts if int(item["id"]) == int(selected_post_id or 0)), None) if not selected and selected_post_id: row = await pool.fetchrow( """ SELECT rp.id, rp.raw_text, rp.original_url, rp.posted_at, rp.created_at, rp.qualification_status, rp.qualification_score, s.name AS source_name, s.tag AS source_tag, rp.qualification_reason, (SELECT COUNT(*) FROM raw_post_media m WHERE m.raw_post_id=rp.id) AS media_count, (SELECT ARRAY_AGG(DISTINCT m.media_type ORDER BY m.media_type) FROM raw_post_media m WHERE m.raw_post_id=rp.id) AS media_types FROM raw_posts rp JOIN sources s ON s.id=rp.source_id WHERE rp.id=$1 """, selected_post_id, ) if row: selected = dict(row) text = str(selected.get("raw_text") or "").strip() selected["created_at_fmt"] = format_dt(selected.get("created_at")) selected["posted_at_fmt"] = format_dt(selected.get("posted_at")) selected["label"] = f"#{selected['id']} · {selected.get('source_name') or 'source'} · {selected['created_at_fmt']}" selected["snippet"] = text[:180] + ("..." if len(text) > 180 else "") posts.insert(0, selected) if not selected and posts: selected = posts[0] return posts, selected def prompt_test_payload(post: dict[str, Any], max_text_chars: int) -> dict[str, Any]: return { "task": "rewrite_accepted_posts", "categories": [], "posts": [ { "id": int(post["id"]), "producer_name": post.get("source_name") or "", "producer_tag": normalize_hash_tag(post.get("source_tag") or post.get("source_name") or "", "source"), "original_url": post.get("original_url") or "", "qualification_score": post.get("qualification_score"), "qualification_reason": post.get("qualification_reason") or "", "media_count": int(post.get("media_count") or 0), "media_types": list(post.get("media_types") or []), "text": str(post.get("raw_text") or "")[:max_text_chars], } ], } def prompt_test_first_rewrite(parsed: Any) -> dict[str, Any] | None: if not isinstance(parsed, dict): return None rewrites = parsed.get("rewrites") if not isinstance(rewrites, list) and isinstance(parsed.get("data"), dict): rewrites = parsed["data"].get("rewrites") if not isinstance(rewrites, list) or not rewrites: return None first = rewrites[0] return first if isinstance(first, dict) else None def prompt_test_preview_text(parsed: Any) -> str: rewrite = prompt_test_first_rewrite(parsed) return str((rewrite or {}).get("text") or "").strip() def prompt_test_preview_category(parsed: Any) -> str: rewrite = prompt_test_first_rewrite(parsed) return str((rewrite or {}).get("category") or "").strip() def prompt_test_preview_notes(parsed: Any) -> str: rewrite = prompt_test_first_rewrite(parsed) notes = (rewrite or {}).get("notes") return "" if notes in {None, ""} else str(notes).strip() def setting_value(settings_rows: list[dict], key: str, default: Any = None) -> Any: for row in settings_rows: if row.get("key") == key: return row.get("value_json", default) return default def setting_sort_key(item: dict) -> tuple[int, int, str]: category = str(item.get("category") or "General") key = str(item.get("key") or "") keys = SETTING_ORDER.get(category, []) pos = keys.index(key) if key in keys else 1000 return (CATEGORY_ORDER.get(category, 999), pos, key) def model_label(model_id: str, name: str | None = None, context: int | None = None) -> str: label = name or model_id if context: label = f"{label} · ctx {context}" return label async def fetch_openrouter_models() -> list[dict[str, str]]: async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=8)) as session: async with session.get("https://openrouter.ai/api/v1/models") as resp: resp.raise_for_status() data = await resp.json() models = [] for item in data.get("data", []): model_id = str(item.get("id") or "").strip() if not model_id: continue context = item.get("context_length") models.append( { "value": f"openrouter/{model_id}", "label": model_label(model_id, item.get("name"), context if isinstance(context, int) else None), } ) return models async def fetch_anthropic_models(api_key: str) -> list[dict[str, str]]: if not api_key: return [] headers = { "x-api-key": api_key, "anthropic-version": "2023-06-01", } models: list[dict[str, str]] = [] url = "https://api.anthropic.com/v1/models" async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session: for _ in range(5): async with session.get(url, headers=headers) as resp: resp.raise_for_status() data = await resp.json() for item in data.get("data", []): model_id = str(item.get("id") or "").strip() if not model_id: continue display_name = item.get("display_name") or model_id models.append({"value": f"anthropic/{model_id}", "label": f"{display_name} · live Anthropic"}) if not data.get("has_more"): break last_id = data.get("last_id") if not last_id: break url = f"https://api.anthropic.com/v1/models?after_id={last_id}" return models async def fetch_openai_models(api_key: str) -> list[dict[str, str]]: if not api_key: return [] headers = {"Authorization": f"Bearer {api_key}"} async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session: async with session.get("https://api.openai.com/v1/models", headers=headers) as resp: resp.raise_for_status() data = await resp.json() models = [] for item in data.get("data", []): model_id = str(item.get("id") or "").strip() if not model_id: continue if model_id.startswith(("gpt-", "o1", "o3", "o4")): models.append({"value": model_id, "label": f"{model_id} · live OpenAI"}) return sorted(models, key=lambda item: item["value"]) def litellm_catalog_models(provider: str) -> list[dict[str, str]]: try: import litellm except Exception: return [] catalog = getattr(litellm, "model_cost", {}) or {} models: list[dict[str, str]] = [] for model_id, meta in catalog.items(): if not isinstance(model_id, str): continue normalized = model_id if provider == "openai": if "/" in model_id or not model_id.startswith(("gpt-", "o1", "o3", "o4")): continue elif provider == "anthropic": if model_id.startswith("anthropic/"): normalized = model_id elif model_id.startswith("claude"): normalized = f"anthropic/{model_id}" else: continue elif provider == "gemini": if model_id.startswith("gemini/"): normalized = model_id elif model_id.startswith("gemini"): normalized = f"gemini/{model_id}" else: continue else: continue context = meta.get("max_input_tokens") if isinstance(meta, dict) else None models.append( { "value": normalized, "label": f"{model_label(normalized, context=context if isinstance(context, int) else None)} · LiteLLM catalog", } ) unique = {m["value"]: m for m in models} preferred = PREFERRED_MODELS.get(provider, []) def sort_key(item: dict[str, str]) -> tuple[int, str]: value = item["value"] if value in preferred: return (preferred.index(value), "") return (1000, item["label"].lower()) return sorted(unique.values(), key=sort_key) def sort_model_options(provider: str, models: list[dict[str, str]]) -> list[dict[str, str]]: preferred = PREFERRED_MODELS.get(provider, []) def sort_key(item: dict[str, str]) -> tuple[int, str]: value = item["value"] if value in preferred: return (preferred.index(value), "") return (1000, item["label"].lower()) unique = {m["value"]: m for m in models} return sorted(unique.values(), key=sort_key) async def model_options_for_provider(provider: str, api_key: str = "") -> list[dict[str, str]]: provider = (provider or "openrouter").strip().lower() now = time.time() cache_key = f"{provider}:{hashlib.sha256((api_key or '').encode('utf-8')).hexdigest()[:10]}" if MODEL_CACHE["key"] == cache_key and now - float(MODEL_CACHE["at"]) < 3600: return list(MODEL_CACHE["models"]) models: list[dict[str, str]] = [] if provider == "openrouter": try: models = await fetch_openrouter_models() except Exception as exc: logger.warning("OpenRouter model list failed: {}", exc) elif provider == "anthropic": try: models = await fetch_anthropic_models(api_key) except Exception as exc: logger.warning("Anthropic model list failed: {}", exc) elif provider == "openai": try: models = await fetch_openai_models(api_key) except Exception as exc: logger.warning("OpenAI model list failed: {}", exc) if not models: models = litellm_catalog_models(provider) if not models: models = [{"value": m, "label": f"{m} · fallback"} for m in MODEL_FALLBACKS.get(provider, [])] models = sort_model_options(provider, models) MODEL_CACHE.update({"key": cache_key, "at": now, "models": models}) return list(models) async def audit(actor_id: int | None, action: str, entity_type: str, entity_id: int | None, after: Any = None) -> None: pool = await get_pool() await pool.execute( """ INSERT INTO audit_log(actor_id, action, entity_type, entity_id, after_json) VALUES($1, $2, $3, $4, $5::jsonb) """, actor_id, action, entity_type, entity_id, json.dumps(after, ensure_ascii=False) if after is not None else None, ) @app.on_event("startup") async def startup() -> None: await bootstrap_admin() @app.get("/health") async def health() -> dict[str, str]: return {"ok": "true"} @app.get("/vk/oauth/callback", response_class=HTMLResponse) async def vk_oauth_callback(request: Request, code: str = "", error: str = "", error_description: str = ""): return templates.TemplateResponse( "vk_oauth_callback.html", base_context( request, None, code=code, code_verifier=request.cookies.get(VK_OAUTH_VERIFIER_COOKIE, ""), device_id=request.query_params.get("device_id", ""), state=request.query_params.get("state", ""), expected_state=request.cookies.get(VK_OAUTH_STATE_COOKIE, ""), error=error, error_description=error_description, ), ) @app.get("/vk/oauth/start") async def vk_oauth_start() -> RedirectResponse: verifier = secrets.token_urlsafe(64) state = secrets.token_urlsafe(24) redirect_uri = "https://sw.exostring.xyz/vk/oauth/callback" params = { "client_id": "54635120", "redirect_uri": redirect_uri, "response_type": "code", "scope": "wall photos video groups offline", "state": state, "code_challenge": pkce_challenge(verifier), "code_challenge_method": "s256", "origin": "https://sw.exostring.xyz", "v": "5.199", } response = RedirectResponse(f"https://id.vk.ru/authorize?{urlencode(params)}") response.set_cookie(VK_OAUTH_VERIFIER_COOKIE, verifier, httponly=True, samesite="lax", max_age=15 * 60) response.set_cookie(VK_OAUTH_STATE_COOKIE, state, httponly=True, samesite="lax", max_age=15 * 60) return response @app.get("/vk/group-oauth/start") async def vk_group_oauth_start() -> RedirectResponse: group_id = abs(int(settings.vk_storage_group_id)) params = { "client_id": "54635120", "group_ids": str(group_id), "display": "page", "redirect_uri": "https://sw.exostring.xyz/vk/oauth/callback", "scope": "manage,photos,docs", "response_type": "token", "state": secrets.token_urlsafe(24), "v": "5.199", } return RedirectResponse(f"https://oauth.vk.com/authorize?{urlencode(params)}") @app.get("/", response_class=HTMLResponse) async def index(request: Request): user = await get_current_user(request) if not user: return redirect("/login") return redirect("/sources") @app.get("/login", response_class=HTMLResponse) async def login_get(request: Request): return templates.TemplateResponse("login.html", base_context(request, None, error="")) @app.post("/login", response_class=HTMLResponse) async def login_post(request: Request, login: str = Form(...), password: str = Form(...)): pool = await get_pool() login_value = login.strip().lower() ip = client_ip(request) await pool.execute("DELETE FROM admin_login_attempts WHERE created_at < NOW() - INTERVAL '7 days'") ip_failures = int( await pool.fetchval( """ SELECT COUNT(*) FROM admin_login_attempts WHERE ip=$1 AND success=FALSE AND created_at > NOW() - INTERVAL '15 minutes' """, ip, ) or 0 ) login_failures = int( await pool.fetchval( """ SELECT COUNT(*) FROM admin_login_attempts WHERE login=$1 AND success=FALSE AND created_at > NOW() - INTERVAL '15 minutes' """, login_value, ) or 0 ) if ip_failures >= 10 or login_failures >= 5: await pool.execute( "INSERT INTO admin_login_attempts(ip, login, success) VALUES($1, $2, FALSE)", ip, login_value, ) return templates.TemplateResponse( "login.html", base_context(request, None, error="Слишком много попыток. Подожди 15 минут и попробуй снова."), status_code=status.HTTP_429_TOO_MANY_REQUESTS, ) row = await pool.fetchrow( "SELECT id, login, password_hash, role, is_active FROM admin_users WHERE login=$1", login_value, ) if not row or not row["is_active"] or not verify_password(password, row["password_hash"]): await pool.execute( "INSERT INTO admin_login_attempts(ip, login, success) VALUES($1, $2, FALSE)", ip, login_value, ) return templates.TemplateResponse("login.html", base_context(request, None, error="Неверный логин или пароль")) await pool.execute("INSERT INTO admin_login_attempts(ip, login, success) VALUES($1, $2, TRUE)", ip, login_value) token = new_token() csrf = new_token() await pool.execute( """ INSERT INTO admin_sessions(token_hash, user_id, csrf_token, expires_at) VALUES($1, $2, $3, $4) """, token_hash(token, settings.app_secret_key), row["id"], csrf, now_utc() + timedelta(days=14), ) response = redirect("/sources") response.set_cookie(COOKIE_NAME, token, httponly=True, secure=True, samesite="lax", max_age=14 * 24 * 3600) return response @app.post("/logout") async def logout(request: Request): token = request.cookies.get(COOKIE_NAME) if token: pool = await get_pool() await pool.execute("DELETE FROM admin_sessions WHERE token_hash=$1", token_hash(token, settings.app_secret_key)) response = redirect("/login") response.delete_cookie(COOKIE_NAME) return response @app.get("/sources", response_class=HTMLResponse) async def sources_list(request: Request, q: str = "", status_filter: str = "", page: int = 1, per_page: int = 50): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() table_filters = parse_table_filters(request, SOURCE_FIELD_SPECS) table_sorts = parse_table_sorts(request, SOURCE_FIELD_SPECS) filter_clauses, filter_args = build_filter_sql(table_filters, SOURCE_FIELD_SPECS, 3) where_parts = [ "s.archived_at IS NULL", "($1='' OR s.name ILIKE '%' || $1 || '%' OR COALESCE(s.tag,'') ILIKE '%' || $1 || '%' OR s.url ILIKE '%' || $1 || '%' OR COALESCE(s.external_id,'') ILIKE '%' || $1 || '%')", "($2='' OR s.status=$2)", *filter_clauses, ] where_sql = " AND ".join(where_parts) base_args = [q.strip(), status_filter.strip(), *filter_args] total = int( await pool.fetchval( f""" SELECT COUNT(*) FROM sources s WHERE {where_sql} """, *base_args, ) or 0 ) pager = pagination(page, per_page, total) limit_index = len(base_args) + 1 offset_index = len(base_args) + 2 order_sql = build_sort_sql(table_sorts, SOURCE_FIELD_SPECS, "s.active DESC, s.priority ASC, s.id DESC") rows = await pool.fetch( f""" SELECT s.*, COUNT(rp.id) AS posts_count, COUNT(rp.id) FILTER (WHERE rp.created_at > NOW() - INTERVAL '24 hours') AS posts_24h FROM sources s LEFT JOIN raw_posts rp ON rp.source_id=s.id WHERE {where_sql} GROUP BY s.id ORDER BY {order_sql} LIMIT ${limit_index} OFFSET ${offset_index} """, *base_args, pager["per_page"], pager["offset"], ) sources = [] for row in rows: item = dict(row) item["last_parsed_at_fmt"] = format_dt(item.get("last_parsed_at")) item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at")) item["created_at_fmt"] = format_dt(item.get("created_at")) sources.append(item) return templates.TemplateResponse( "sources.html", base_context( request, user, sources=sources, q=q, status_filter=status_filter, pagination=pager, table_filters=table_filters, table_sorts=table_sorts, field_options=field_options(SOURCE_FIELD_SPECS), filter_ops=FILTER_OPS, filter_op_options=FILTER_OP_OPTIONS, source_preview=[], bulk_text="", bulk_active=True, ), ) @app.post("/sources/preview", response_class=HTMLResponse) async def sources_preview( request: Request, csrf_token: str = Form(...), bulk_text: str = Form(""), bulk_active: str = Form("off"), q: str = Form(""), status_filter: str = Form(""), per_page: int = Form(50), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) preview = await build_sources_preview(bulk_text, default_active=bulk_active == "on") pool = await get_pool() table_filters = parse_table_filters(request, SOURCE_FIELD_SPECS) table_sorts = parse_table_sorts(request, SOURCE_FIELD_SPECS) filter_clauses, filter_args = build_filter_sql(table_filters, SOURCE_FIELD_SPECS, 3) where_parts = [ "s.archived_at IS NULL", "($1='' OR s.name ILIKE '%' || $1 || '%' OR COALESCE(s.tag,'') ILIKE '%' || $1 || '%' OR s.url ILIKE '%' || $1 || '%' OR COALESCE(s.external_id,'') ILIKE '%' || $1 || '%')", "($2='' OR s.status=$2)", *filter_clauses, ] where_sql = " AND ".join(where_parts) base_args = [q.strip(), status_filter.strip(), *filter_args] total = int(await pool.fetchval(f"SELECT COUNT(*) FROM sources s WHERE {where_sql}", *base_args) or 0) pager = pagination(1, per_page, total) limit_index = len(base_args) + 1 offset_index = len(base_args) + 2 order_sql = build_sort_sql(table_sorts, SOURCE_FIELD_SPECS, "s.active DESC, s.priority ASC, s.id DESC") rows = await pool.fetch( f""" SELECT s.*, COUNT(rp.id) AS posts_count, COUNT(rp.id) FILTER (WHERE rp.created_at > NOW() - INTERVAL '24 hours') AS posts_24h FROM sources s LEFT JOIN raw_posts rp ON rp.source_id=s.id WHERE {where_sql} GROUP BY s.id ORDER BY {order_sql} LIMIT ${limit_index} OFFSET ${offset_index} """, *base_args, pager["per_page"], pager["offset"], ) sources = [] for row in rows: item = dict(row) item["last_parsed_at_fmt"] = format_dt(item.get("last_parsed_at")) item["last_checked_at_fmt"] = format_dt(item.get("last_checked_at")) item["created_at_fmt"] = format_dt(item.get("created_at")) sources.append(item) return templates.TemplateResponse( "sources.html", base_context( request, user, sources=sources, q=q, status_filter=status_filter, pagination=pager, table_filters=table_filters, table_sorts=table_sorts, field_options=field_options(SOURCE_FIELD_SPECS), filter_ops=FILTER_OPS, filter_op_options=FILTER_OP_OPTIONS, source_preview=preview, bulk_text=bulk_text, bulk_active=bulk_active == "on", ), ) @app.post("/sources/bulk") async def sources_bulk_create( request: Request, csrf_token: str = Form(...), bulk_payload: str = Form("[]"), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) try: items = json.loads(bulk_payload) except json.JSONDecodeError: items = [] pool = await get_pool() created = 0 for item in items: if not isinstance(item, dict) or not item.get("ok"): continue try: row = await pool.fetchrow( """ INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, created_by) VALUES($1, $2, $3, $4, $5, $6, $7, 100, $8) ON CONFLICT DO NOTHING RETURNING id """, PLATFORM_VK, str(item.get("name") or item.get("external_id") or "").strip(), normalize_hash_tag(str(item.get("tag") or item.get("external_id") or ""), str(item.get("external_id") or "source")), str(item.get("url") or "").strip(), str(item.get("external_id") or "").strip(), item.get("external_owner_id"), bool(item.get("active", True)), user["id"], ) if row: created += 1 await audit(user["id"], "source.create", "source", int(row["id"]), {"url": item.get("url"), "bulk": True}) except Exception: logger.exception("Bulk source insert failed: {}", item) return redirect(f"/sources?q=&status_filter=&created={created}") @app.get("/sources/new", response_class=HTMLResponse) async def source_new(request: Request): user = await get_current_user(request) if not user: return redirect("/login") return templates.TemplateResponse( "source_form.html", base_context(request, user, source=None, action="/sources/new", title="Новый источник"), ) @app.post("/sources/new") async def source_create( request: Request, csrf_token: str = Form(...), platform: str = Form(PLATFORM_VK), name: str = Form(...), tag: str = Form(""), url: str = Form(...), active: str = Form("off"), priority: int = Form(100), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) platform = platform.strip().lower() or PLATFORM_VK external_id = normalize_vk_source(url) if platform == PLATFORM_VK else "" external_owner_id = None resolved_name = name.strip() resolved_url = url.strip() status_value = "new" status_msg = None if platform == PLATFORM_VK and external_id: try: async with VKAPIClient(rps=2, timeout_total_sec=12, timeout_connect_sec=5, retry_attempts=2) as client: screen_name, owner_id, vk_name = await client.resolve_group(url) external_id = screen_name external_owner_id = owner_id resolved_name = resolved_name or vk_name resolved_url = f"https://vk.com/{screen_name}" except Exception as exc: resolved_name = resolved_name or external_id status_value = "error" status_msg = str(exc)[:500] resolved_tag = normalize_hash_tag(tag or external_id or resolved_name, external_id or "source") pool = await get_pool() row = await pool.fetchrow( """ INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, status, status_msg, created_by) VALUES($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) ON CONFLICT DO NOTHING RETURNING id """, platform, resolved_name, resolved_tag, resolved_url, external_id, external_owner_id, active == "on", priority, status_value, status_msg, user["id"], ) if row: await audit(user["id"], "source.create", "source", int(row["id"]), {"url": url, "platform": platform}) return redirect("/sources") @app.get("/sources/{source_id}/edit", response_class=HTMLResponse) async def source_edit(request: Request, source_id: int): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() source = await pool.fetchrow("SELECT * FROM sources WHERE id=$1 AND archived_at IS NULL", source_id) if not source: return redirect("/sources") return templates.TemplateResponse( "source_form.html", base_context( request, user, source=dict(source), action=f"/sources/{source_id}/edit", title=f"Источник #{source_id}", ), ) @app.post("/sources/{source_id}/edit") async def source_update( request: Request, source_id: int, csrf_token: str = Form(...), platform: str = Form(PLATFORM_VK), name: str = Form(...), tag: str = Form(""), url: str = Form(...), active: str = Form("off"), priority: int = Form(100), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) platform = platform.strip().lower() or PLATFORM_VK external_id = normalize_vk_source(url) if platform == PLATFORM_VK else "" pool = await get_pool() await pool.execute( """ UPDATE sources SET platform=$2, name=$3, tag=$4, url=$5, external_id=$6, active=$7, priority=$8, updated_at=NOW() WHERE id=$1 """, source_id, platform, name.strip(), normalize_hash_tag(tag or external_id or name, external_id or "source"), url.strip(), external_id, active == "on", priority, ) await audit(user["id"], "source.update", "source", source_id, {"url": url, "platform": platform}) return redirect("/sources") @app.post("/sources/{source_id}/toggle") async def source_toggle(request: Request, source_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() row = await pool.fetchrow( """ UPDATE sources SET active=NOT active, status=CASE WHEN active THEN 'paused' ELSE 'new' END, updated_at=NOW() WHERE id=$1 RETURNING active """, source_id, ) await audit(user["id"], "source.toggle", "source", source_id, {"active": bool(row["active"]) if row else None}) return redirect("/sources") @app.post("/sources/{source_id}/delete") async def source_delete(request: Request, source_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() await pool.execute("UPDATE sources SET archived_at=NOW(), active=FALSE, updated_at=NOW() WHERE id=$1", source_id) await audit(user["id"], "source.archive", "source", source_id) return redirect("/sources") @app.get("/raw", response_class=HTMLResponse) async def raw_posts( request: Request, status_filter: str = "", date_from: str = "", date_to: str = "", score_min: str = "", score_max: str = "", sort: str = "created_desc", page: int = 1, per_page: int = 50, ): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() order_map = { "id_desc": "rp.id DESC", "id_asc": "rp.id ASC", "created_desc": "rp.created_at DESC", "created_asc": "rp.created_at ASC", "posted_desc": "rp.posted_at DESC NULLS LAST, rp.created_at DESC", "posted_asc": "rp.posted_at ASC NULLS LAST, rp.created_at ASC", "source_asc": "lower(s.name) ASC, rp.created_at DESC", "source_desc": "lower(s.name) DESC, rp.created_at DESC", "stage_asc": "rp.status ASC, COALESCE(rp.qualification_status, 'pending') ASC, COALESCE(rp.rewrite_status, 'pending') ASC, rp.created_at DESC", "stage_desc": "rp.status DESC, COALESCE(rp.qualification_status, 'pending') DESC, COALESCE(rp.rewrite_status, 'pending') DESC, rp.created_at DESC", "text_asc": "lower(rp.raw_text) ASC, rp.created_at DESC", "text_desc": "lower(rp.raw_text) DESC, rp.created_at DESC", "score_desc": "rp.qualification_score DESC NULLS LAST, rp.created_at DESC", "score_asc": "rp.qualification_score ASC NULLS LAST, rp.created_at DESC", } raw_statuses = selected_values(request, "raw_status") if status_filter.strip() and not raw_statuses: raw_statuses = [status_filter.strip()] source_ids = selected_int_values(request, "source_id") categories_selected = selected_values(request, "category") qualification_statuses = selected_values(request, "qualification_status") rewrite_statuses = selected_values(request, "rewrite_status") where_parts: list[str] = [] base_args: list[Any] = [] if raw_statuses: add_where(where_parts, base_args, "rp.status = ANY(${i}::text[])", raw_statuses) if source_ids: add_where(where_parts, base_args, "rp.source_id = ANY(${i}::bigint[])", source_ids) if categories_selected: add_where(where_parts, base_args, "COALESCE(rp.final_category, rp.rewrite_category, '') = ANY(${i}::text[])", categories_selected) if qualification_statuses: add_where(where_parts, base_args, "COALESCE(rp.qualification_status, 'pending') = ANY(${i}::text[])", qualification_statuses) if rewrite_statuses: add_where(where_parts, base_args, "COALESCE(rp.rewrite_status, 'pending') = ANY(${i}::text[])", rewrite_statuses) parsed_date_from = parse_date_filter(date_from) parsed_date_to = parse_date_filter(date_to) if parsed_date_from: add_where(where_parts, base_args, "timezone('Asia/Yekaterinburg', COALESCE(rp.posted_at, rp.created_at))::date >= ${i}::date", parsed_date_from) if parsed_date_to: add_where(where_parts, base_args, "timezone('Asia/Yekaterinburg', COALESCE(rp.posted_at, rp.created_at))::date <= ${i}::date", parsed_date_to) if score_min.strip().isdigit(): add_where(where_parts, base_args, "rp.qualification_score >= ${i}::int", int(score_min.strip())) if score_max.strip().isdigit(): add_where(where_parts, base_args, "rp.qualification_score <= ${i}::int", int(score_max.strip())) where_sql = " AND ".join(where_parts) if where_parts else "TRUE" order_sql = order_map.get(sort, order_map["created_desc"]) total = int( await pool.fetchval( f""" SELECT COUNT(*) FROM raw_posts rp JOIN sources s ON s.id=rp.source_id WHERE {where_sql} """, *base_args, ) or 0 ) pager = pagination(page, per_page, total) limit_index = len(base_args) + 1 offset_index = len(base_args) + 2 rows = await pool.fetch( f""" SELECT rp.*, s.name AS source_name, s.tag AS source_tag, s.platform AS source_platform, COUNT(rpm.id) AS media_count, COALESCE( jsonb_agg( jsonb_build_object( 'id', rpm.id, 'type', rpm.media_type, 'url', rpm.original_url, 'status', rpm.status, 'error', rpm.error, 'duration_sec', rpm.duration_sec ) ORDER BY rpm.sort_order ASC, rpm.id ASC ) FILTER (WHERE rpm.id IS NOT NULL), '[]'::jsonb ) AS media_items FROM raw_posts rp JOIN sources s ON s.id=rp.source_id LEFT JOIN raw_post_media rpm ON rpm.raw_post_id=rp.id AND COALESCE(rpm.editor_added, FALSE)=FALSE WHERE {where_sql} GROUP BY rp.id, s.name, s.tag, s.platform ORDER BY {order_sql} LIMIT ${limit_index} OFFSET ${offset_index} """, *base_args, pager["per_page"], pager["offset"], ) posts = [prepare_raw_post(row) for row in rows] return templates.TemplateResponse( "raw_posts.html", base_context( request, user, posts=posts, raw_statuses=raw_statuses, source_ids=source_ids, categories_selected=categories_selected, qualification_statuses=qualification_statuses, rewrite_statuses=rewrite_statuses, date_from=date_from, date_to=date_to, score_min=score_min, score_max=score_max, sort=sort, facets=await raw_filter_facets(pool), preserved_query=preserved_query(request, {"page", "per_page"}), sort_headers=raw_sort_headers(request, sort), pagination=pager, ), ) @app.get("/editor", response_class=HTMLResponse) async def editor_feed( request: Request, status_filter: str = "review", category_filter: str = "", date_from: str = "", date_to: str = "", score_min: str = "", score_max: str = "", sort: str = "rewritten_desc", page: int = 1, per_page: int = 25, ): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() editorial_statuses = selected_values(request, "editorial_status") if status_filter.strip() and not editorial_statuses: editorial_statuses = [status_filter.strip()] categories_selected = selected_values(request, "category") if category_filter.strip() and not categories_selected: categories_selected = [category_filter.strip()] source_ids = selected_int_values(request, "source_id") sort_map = { "rewritten_desc": "rp.rewritten_at DESC NULLS LAST, rp.id DESC", "rewritten_asc": "rp.rewritten_at ASC NULLS LAST, rp.id ASC", "score_desc": "rp.qualification_score DESC NULLS LAST, rp.rewritten_at DESC NULLS LAST", "score_asc": "rp.qualification_score ASC NULLS LAST, rp.rewritten_at DESC NULLS LAST", "posted_desc": "rp.posted_at DESC NULLS LAST, rp.rewritten_at DESC NULLS LAST", "posted_asc": "rp.posted_at ASC NULLS LAST, rp.rewritten_at DESC NULLS LAST", } where_parts = ["rp.rewrite_status='ready'"] base_args: list[Any] = [] if editorial_statuses: add_where(where_parts, base_args, "COALESCE(rp.editorial_status, 'review') = ANY(${i}::text[])", editorial_statuses) if categories_selected: add_where(where_parts, base_args, "COALESCE(rp.final_category, rp.rewrite_category, '') = ANY(${i}::text[])", categories_selected) if source_ids: add_where(where_parts, base_args, "rp.source_id = ANY(${i}::bigint[])", source_ids) parsed_date_from = parse_date_filter(date_from) parsed_date_to = parse_date_filter(date_to) if parsed_date_from: add_where(where_parts, base_args, "timezone('Asia/Yekaterinburg', COALESCE(rp.posted_at, rp.created_at))::date >= ${i}::date", parsed_date_from) if parsed_date_to: add_where(where_parts, base_args, "timezone('Asia/Yekaterinburg', COALESCE(rp.posted_at, rp.created_at))::date <= ${i}::date", parsed_date_to) if score_min.strip().isdigit(): add_where(where_parts, base_args, "rp.qualification_score >= ${i}::int", int(score_min.strip())) if score_max.strip().isdigit(): add_where(where_parts, base_args, "rp.qualification_score <= ${i}::int", int(score_max.strip())) where_sql = " AND ".join(where_parts) total = int( await pool.fetchval( f""" SELECT COUNT(*) FROM raw_posts rp JOIN sources s ON s.id=rp.source_id WHERE {where_sql} """, *base_args, ) or 0 ) pager = pagination(page, per_page, total) limit_index = len(base_args) + 1 offset_index = len(base_args) + 2 rows = await pool.fetch( f""" SELECT rp.*, s.name AS source_name, s.tag AS source_tag, s.platform AS source_platform, s.url AS source_url, u.login AS reviewed_by_login, COUNT(rpm.id) AS media_count, COALESCE( jsonb_agg( jsonb_build_object( 'id', rpm.id, 'type', rpm.media_type, 'url', rpm.original_url, 'status', rpm.status, 'error', rpm.error, 'duration_sec', rpm.duration_sec ) ORDER BY rpm.sort_order ASC, rpm.id ASC ) FILTER (WHERE rpm.id IS NOT NULL), '[]'::jsonb ) AS media_items FROM raw_posts rp JOIN sources s ON s.id=rp.source_id LEFT JOIN admin_users u ON u.id=rp.reviewed_by LEFT JOIN raw_post_media rpm ON rpm.raw_post_id=rp.id AND COALESCE(rpm.editor_hidden, FALSE)=FALSE WHERE {where_sql} GROUP BY rp.id, s.name, s.tag, s.platform, s.url, u.login ORDER BY CASE COALESCE(rp.editorial_status, 'review') WHEN 'review' THEN 0 WHEN 'edited' THEN 1 WHEN 'regenerating' THEN 2 WHEN 'rejected' THEN 3 WHEN 'accepted' THEN 4 ELSE 5 END, {sort_map.get(sort, sort_map["rewritten_desc"])} LIMIT ${limit_index} OFFSET ${offset_index} """, *base_args, pager["per_page"], pager["offset"], ) posts = [prepare_editor_post(row) for row in rows] categories = await writer_categories() counts_rows = await pool.fetch( """ SELECT COALESCE(editorial_status, 'review') AS status, COUNT(*) AS count FROM raw_posts WHERE rewrite_status='ready' GROUP BY COALESCE(editorial_status, 'review') """ ) counts = {str(row["status"]): int(row["count"]) for row in counts_rows} return templates.TemplateResponse( "editor.html", base_context( request, user, posts=posts, categories=categories, status_options=EDITOR_STATUS_OPTIONS, status_filter=editorial_statuses[0] if len(editorial_statuses) == 1 else "", editorial_statuses=editorial_statuses, categories_selected=categories_selected, source_ids=source_ids, date_from=date_from, date_to=date_to, score_min=score_min, score_max=score_max, sort=sort, facets=await editor_filter_facets(pool), counts=counts, preserved_query=preserved_query(request, {"page", "per_page"}), pagination=pager, ), ) @app.post("/editor/{post_id}/accept") async def editor_accept(request: Request, post_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() await pool.execute( """ UPDATE raw_posts SET editorial_status='accepted', final_text=COALESCE(final_text, rewritten_text), final_category=COALESCE(final_category, rewrite_category), final_category_tag=COALESCE(final_category_tag, rewrite_category_tag), final_source_tag=COALESCE(final_source_tag, rewrite_source_tag), reviewed_by=$2, reviewed_at=NOW(), updated_at=NOW() WHERE id=$1 AND rewrite_status='ready' """, post_id, user["id"], ) await audit(user["id"], "editor.accept", "raw_post", post_id) return redirect("/editor") @app.post("/editor/{post_id}/reject") async def editor_reject(request: Request, post_id: int, csrf_token: str = Form(...), editor_notes: str = Form("")): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() await pool.execute( """ UPDATE raw_posts SET editorial_status='rejected', editor_notes=$2, reviewed_by=$3, reviewed_at=NOW(), updated_at=NOW() WHERE id=$1 AND rewrite_status='ready' """, post_id, editor_notes.strip()[:1000], user["id"], ) await audit(user["id"], "editor.reject", "raw_post", post_id) return redirect("/editor") @app.post("/editor/{post_id}/save") async def editor_save( request: Request, post_id: int, csrf_token: str = Form(...), final_text: str = Form(...), final_category: str = Form(...), final_source_tag: str = Form(""), editor_notes: str = Form(""), action: str = Form("save"), delete_media_ids: list[int] | None = Form(default=None), media_files: list[UploadFile] | None = File(default=None), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) final_category = final_category.strip() final_category_tag = normalize_hash_tag(final_category, "category") final_source_tag = normalize_hash_tag(final_source_tag, "source") status_value = "accepted" if action == "accept" else "edited" pool = await get_pool() await pool.execute( """ UPDATE raw_posts SET editorial_status=$2, final_text=$3, final_category=$4, final_category_tag=$5, final_source_tag=$6, editor_notes=$7, reviewed_by=CASE WHEN $2='accepted' THEN $8 ELSE reviewed_by END, reviewed_at=CASE WHEN $2='accepted' THEN NOW() ELSE reviewed_at END, edited_at=NOW(), updated_at=NOW() WHERE id=$1 AND rewrite_status='ready' """, post_id, status_value, final_text.strip(), final_category, final_category_tag, final_source_tag, editor_notes.strip()[:1000], user["id"], ) delete_media_ids = delete_media_ids or [] media_files = media_files or [] if delete_media_ids: local_media_paths = [ uploaded_editor_media_path(str(row["original_url"] or "")) for row in await pool.fetch( """ SELECT original_url FROM raw_post_media WHERE raw_post_id=$1 AND id=ANY($2::bigint[]) AND editor_added=TRUE """, post_id, delete_media_ids, ) ] await pool.execute( """ UPDATE raw_post_media SET editor_hidden=TRUE, updated_at=NOW() WHERE raw_post_id=$1 AND id=ANY($2::bigint[]) """, post_id, delete_media_ids, ) for path in local_media_paths: if path and path.is_file(): try: path.unlink(missing_ok=True) except OSError: logger.warning("Could not delete editor media file {}", path) for media_file in media_files: if not media_file.filename: continue content = await media_file.read() if not content or len(content) > 50 * 1024 * 1024: continue original_name = Path(media_file.filename or "media").name suffix = Path(original_name).suffix.lower() if not suffix: guessed = mimetypes.guess_extension(media_file.content_type or "") suffix = guessed or ".bin" safe_name = f"{post_id}_{int(time.time())}_{secrets.token_hex(6)}{suffix}" target = EDITOR_MEDIA_DIR / safe_name target.write_bytes(content) public_url = f"/uploads/editor_media/{safe_name}" content_type = media_file.content_type or mimetypes.guess_type(original_name)[0] or "" if content_type.startswith("image/"): media_type = "photo" elif content_type.startswith("video/"): media_type = "video" else: media_type = "doc" sort_order = int( await pool.fetchval("SELECT COALESCE(MAX(sort_order), -1) + 1 FROM raw_post_media WHERE raw_post_id=$1", post_id) or 0 ) await pool.execute( """ INSERT INTO raw_post_media( raw_post_id, platform, media_type, original_url, preview_url, sort_order, status, editor_added, editor_uploaded_by, editor_uploaded_at ) VALUES($1, 'manual', $2, $3, $3, $4, 'ready', TRUE, $5, NOW()) """, post_id, media_type, public_url, sort_order, user["id"], ) await audit(user["id"], "editor.save" if status_value == "edited" else "editor.save_accept", "raw_post", post_id) return redirect("/editor") @app.get("/raw/{post_id}", response_class=HTMLResponse) async def raw_post_detail(request: Request, post_id: int): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() row = await pool.fetchrow( """ SELECT rp.*, s.name AS source_name, s.tag AS source_tag, s.platform AS source_platform, s.url AS source_url, COUNT(rpm.id) AS media_count, COALESCE( jsonb_agg( jsonb_build_object( 'id', rpm.id, 'type', rpm.media_type, 'url', rpm.original_url, 'status', rpm.status, 'error', rpm.error, 'duration_sec', rpm.duration_sec ) ORDER BY rpm.sort_order ASC, rpm.id ASC ) FILTER (WHERE rpm.id IS NOT NULL), '[]'::jsonb ) AS media_items, b.provider AS qualification_provider, b.model AS batch_model, b.status AS batch_status, b.posts_count AS batch_posts_count, b.accepted_count AS batch_accepted_count, b.rejected_count AS batch_rejected_count, b.maybe_count AS batch_maybe_count, b.error AS batch_error, b.prompt_text AS batch_prompt_text, b.created_at AS batch_created_at, b.completed_at AS batch_completed_at, b.response_json AS batch_response_json, b.prompt_tokens AS batch_prompt_tokens, b.completion_tokens AS batch_completion_tokens, b.total_tokens AS batch_total_tokens, b.estimated_cost_usd AS batch_estimated_cost_usd, wb.provider AS rewrite_provider, wb.model AS writer_batch_model, wb.status AS writer_batch_status, wb.posts_count AS writer_batch_posts_count, wb.ready_count AS writer_batch_ready_count, wb.failed_count AS writer_batch_failed_count, wb.error AS writer_batch_error, wb.prompt_text AS writer_batch_prompt_text, wb.created_at AS writer_batch_created_at, wb.completed_at AS writer_batch_completed_at, wb.response_json AS writer_batch_response_json, wb.prompt_tokens AS writer_prompt_tokens, wb.completion_tokens AS writer_completion_tokens, wb.total_tokens AS writer_total_tokens, wb.estimated_cost_usd AS writer_estimated_cost_usd FROM raw_posts rp JOIN sources s ON s.id=rp.source_id LEFT JOIN raw_post_media rpm ON rpm.raw_post_id=rp.id LEFT JOIN ai_qualification_batches b ON b.id=rp.qualification_batch_id LEFT JOIN ai_writer_batches wb ON wb.id=rp.rewrite_batch_id WHERE rp.id=$1 GROUP BY rp.id, s.name, s.tag, s.platform, s.url, b.id, wb.id """, post_id, ) if not row: return redirect("/raw") post = prepare_raw_post(row) post["batch_created_at_fmt"] = format_dt(post.get("batch_created_at")) post["batch_completed_at_fmt"] = format_dt(post.get("batch_completed_at")) response_json = post.get("batch_response_json") or {} if isinstance(response_json, str): try: response_json = json.loads(response_json) except json.JSONDecodeError: pass post["batch_response_pretty"] = json.dumps(response_json, ensure_ascii=False, indent=2) if not post.get("batch_prompt_text") and post.get("qualification_batch_id"): post["batch_prompt_text"] = ( "Точный prompt для этого старого batch ещё не сохранялся.\n" f"Prompt hash: {post.get('qualification_prompt_hash') or '—'}" ) else: post["batch_prompt_text"] = post.get("batch_prompt_text") or "" writer_response_json = post.get("writer_batch_response_json") or {} if isinstance(writer_response_json, str): try: writer_response_json = json.loads(writer_response_json) except json.JSONDecodeError: pass post["writer_batch_created_at_fmt"] = format_dt(post.get("writer_batch_created_at")) post["writer_batch_completed_at_fmt"] = format_dt(post.get("writer_batch_completed_at")) post["rewritten_at_fmt"] = format_dt(post.get("rewritten_at")) post["writer_batch_response_pretty"] = json.dumps(writer_response_json, ensure_ascii=False, indent=2) if not post.get("writer_batch_prompt_text") and post.get("rewrite_batch_id"): post["writer_batch_prompt_text"] = ( "Точный prompt для этого старого writer batch ещё не сохранялся.\n" f"Prompt hash: {post.get('rewrite_prompt_hash') or '—'}" ) else: post["writer_batch_prompt_text"] = post.get("writer_batch_prompt_text") or "" return templates.TemplateResponse("raw_post_detail.html", base_context(request, user, post=post)) @app.get("/prompt-test", response_class=HTMLResponse) async def prompt_test(request: Request, post_id: int | None = None, q_status: str = "all"): user = await get_current_user(request) if not user: return redirect("/login") q_status = normalize_prompt_test_status(q_status) posts, selected = await prompt_test_posts(post_id, q_status) categories = await writer_categories() max_text_chars = max(200, int(await fetch_setting("ai_writer_max_text_chars", 3500) or 3500)) prompt = str(await fetch_setting("ai_writer_prompt", "") or "") provider = str(await fetch_setting("ai_writer_provider", "anthropic") or "anthropic") model = str(await fetch_setting("ai_writer_model", "") or "") payload = prompt_test_payload(selected, max_text_chars) if selected else [] if payload: payload["categories"] = categories return templates.TemplateResponse( "prompt_test.html", base_context( request, user, posts=posts, selected=selected, prompt=prompt, provider=provider, model=model, q_status=q_status, status_options=PROMPT_TEST_STATUS_OPTIONS, max_text_chars=max_text_chars, payload_json=json.dumps(payload, ensure_ascii=False, indent=2), result=None, error=None, ), ) @app.post("/prompt-test", response_class=HTMLResponse) async def prompt_test_run( request: Request, csrf_token: str = Form(...), post_id: int = Form(...), prompt: str = Form(...), q_status: str = Form("all"), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) q_status = normalize_prompt_test_status(q_status) posts, selected = await prompt_test_posts(post_id, q_status) categories = await writer_categories() max_text_chars = max(200, int(await fetch_setting("ai_writer_max_text_chars", 3500) or 3500)) provider = str(await fetch_setting("ai_writer_provider", "anthropic") or "anthropic") model = str(await fetch_setting("ai_writer_model", "") or "").strip() api_key = str(await fetch_setting("ai_writer_api_key", "") or "").strip() api_base = str(await fetch_setting("ai_writer_api_base", "") or "").strip() temperature = max(0.0, float(await fetch_setting("ai_writer_temperature", 0.4) or 0.4)) timeout = max(10, int(await fetch_setting("ai_writer_timeout_sec", 180) or 180)) normalized_prompt = prompt.replace("\\r\\n", "\n").replace("\\n", "\n").strip() payload = prompt_test_payload(selected, max_text_chars) if selected else [] if payload: payload["categories"] = categories result = None error = None if not selected: error = "Пост не найден." elif not model or not api_key: error = "Не заполнены модель или API-ключ райтера в настройках." elif not normalized_prompt: error = "Промпт пустой." else: try: import litellm kwargs: dict[str, Any] = { "model": normalize_model(provider, model), "messages": [ {"role": "system", "content": normalized_prompt}, {"role": "user", "content": json.dumps(payload, ensure_ascii=False)}, ], "temperature": temperature, "timeout": timeout, } if api_key: kwargs["api_key"] = api_key if api_base: kwargs["api_base"] = api_base response = await asyncio.to_thread(litellm.completion, **kwargs) content = response.choices[0].message.content parsed = None raw_json = str(content or "").strip() try: if raw_json.startswith("```"): raw_json = raw_json.strip("`") if raw_json.lower().startswith("json"): raw_json = raw_json[4:].strip() parsed = json.loads(raw_json) except Exception: parsed = None result = { "content": content, "parsed_json": json.dumps(parsed, ensure_ascii=False, indent=2) if parsed is not None else "", "preview_text": prompt_test_preview_text(parsed), "preview_category": prompt_test_preview_category(parsed), "preview_notes": prompt_test_preview_notes(parsed), "usage": response_usage(response), "model": kwargs["model"], } await audit(user["id"], "writer_prompt_test.run", "raw_post", int(selected["id"]), {"provider": provider, "model": model}) except Exception as exc: logger.exception("Prompt test failed") error = str(exc) return templates.TemplateResponse( "prompt_test.html", base_context( request, user, posts=posts, selected=selected, prompt=normalized_prompt, provider=provider, model=model, q_status=q_status, status_options=PROMPT_TEST_STATUS_OPTIONS, max_text_chars=max_text_chars, payload_json=json.dumps(payload, ensure_ascii=False, indent=2), result=result, error=error, ), ) @app.get("/workers", response_class=HTMLResponse) async def workers(request: Request): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() rows = await pool.fetch( """ SELECT wc.name, wc.enabled, wc.settings_json, wh.heartbeat_at, wh.status, wh.current_job_id, wh.meta_json FROM worker_controls wc LEFT JOIN worker_heartbeats wh ON wh.name=wc.name ORDER BY wc.name """ ) settings_rows = await pool.fetch("SELECT * FROM app_settings ORDER BY category, key") settings = [] for row in settings_rows: item = dict(row) if isinstance(item.get("value_json"), str): try: item["value_json"] = json.loads(item["value_json"]) except json.JSONDecodeError: pass if item.get("key") == "ai_writer_categories": continue settings.append(item) settings.sort(key=setting_sort_key) ai_provider = str(setting_value(settings, "ai_qualifier_provider", "openrouter") or "openrouter") ai_model = str(setting_value(settings, "ai_qualifier_model", "") or "") ai_api_key = str(setting_value(settings, "ai_qualifier_api_key", "") or "") ai_model_options = await model_options_for_provider(ai_provider, ai_api_key) writer_provider = str(setting_value(settings, "ai_writer_provider", "anthropic") or "anthropic") writer_model = str(setting_value(settings, "ai_writer_model", "") or "") writer_api_key = str(setting_value(settings, "ai_writer_api_key", "") or "") writer_model_options = await model_options_for_provider(writer_provider, writer_api_key) return templates.TemplateResponse( "workers.html", base_context( request, user, workers=[dict(r) for r in rows], settings=settings, provider_options=PROVIDER_OPTIONS, ai_provider=ai_provider, ai_model=ai_model, ai_model_options=ai_model_options, writer_provider=writer_provider, writer_model=writer_model, writer_model_options=writer_model_options, categories=await category_rows(), category_titles=CATEGORY_TITLES, prompt_hints=PROMPT_HINTS, ), ) @app.post("/categories/create") async def category_create( request: Request, csrf_token: str = Form(...), name: str = Form(...), tag: str = Form(""), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) name = name.strip().strip("#") if not name: return redirect("/workers") tag = normalize_hash_tag(tag or name, "category") pool = await get_pool() sort_order = int(await pool.fetchval("SELECT COALESCE(MAX(sort_order), 0) + 10 FROM content_categories") or 10) await pool.execute( """ INSERT INTO content_categories(name, tag, sort_order) VALUES($1, $2, $3) ON CONFLICT (name) DO UPDATE SET tag=$2, is_active=TRUE, updated_at=NOW() """, name, tag, sort_order, ) await audit(user["id"], "category.create", "content_category", None, {"name": name, "tag": tag}) return redirect("/workers") @app.post("/categories/{category_id}/update") async def category_update( request: Request, category_id: int, csrf_token: str = Form(...), name: str = Form(...), tag: str = Form(""), sort_order: int = Form(0), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) name = name.strip().strip("#") if not name: return redirect("/workers") tag = normalize_hash_tag(tag or name, "category") pool = await get_pool() await pool.execute( """ UPDATE content_categories SET name=$2, tag=$3, sort_order=$4, updated_at=NOW() WHERE id=$1 """, category_id, name, tag, sort_order, ) await audit(user["id"], "category.update", "content_category", category_id, {"name": name, "tag": tag}) return redirect("/workers") @app.post("/categories/{category_id}/toggle") async def category_toggle(request: Request, category_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() row = await pool.fetchrow( """ UPDATE content_categories SET is_active=NOT is_active, updated_at=NOW() WHERE id=$1 RETURNING name, is_active """, category_id, ) await audit( user["id"], "category.toggle", "content_category", category_id, {"name": row["name"] if row else None, "is_active": bool(row["is_active"]) if row else None}, ) return redirect("/workers") @app.post("/categories/{category_id}/delete") async def category_delete(request: Request, category_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() row = await pool.fetchrow( """ UPDATE content_categories SET is_active=FALSE, updated_at=NOW() WHERE id=$1 RETURNING name """, category_id, ) await audit(user["id"], "category.archive", "content_category", category_id, {"name": row["name"] if row else None}) return redirect("/workers") @app.post("/workers/{worker_name}/toggle") async def worker_toggle(request: Request, worker_name: str, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() row = await pool.fetchrow( """ UPDATE worker_controls SET enabled=NOT enabled, updated_by=$2, updated_at=NOW() WHERE name=$1 RETURNING enabled """, worker_name, user["id"], ) await audit(user["id"], "worker.toggle", "worker", None, {"name": worker_name, "enabled": bool(row["enabled"])}) return redirect("/workers") @app.post("/settings/save") async def settings_save(request: Request, csrf_token: str = Form(...), key: str = Form(...), value: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) pool = await get_pool() row = await pool.fetchrow("SELECT value_type FROM app_settings WHERE key=$1", key) if not row: return redirect("/workers") value_type = str(row["value_type"]) if value_type == "int": parsed: Any = int(value) elif value_type == "float": parsed = float(value) elif value_type == "bool": parsed = str(value).lower() in {"1", "true", "yes", "on"} elif value_type == "secret": if not value.strip(): return redirect("/workers") parsed = value.strip() else: parsed = value if key in {"ai_qualifier_prompt", "ai_qualifier_contract", "ai_writer_prompt", "ai_writer_contract"}: parsed = parsed.replace("\\r\\n", "\n").replace("\\n", "\n") await pool.execute( """ UPDATE app_settings SET value_json=$2::jsonb, updated_by=$3, updated_at=NOW() WHERE key=$1 """, key, json.dumps(parsed, ensure_ascii=False), user["id"], ) await audit(user["id"], "setting.update", "setting", None, {"key": key, "value": parsed}) return redirect("/workers") @app.get("/users", response_class=HTMLResponse) async def users(request: Request, page: int = 1, per_page: int = 50): user = await get_current_user(request) if not user: return redirect("/login") pool = await get_pool() total = int(await pool.fetchval("SELECT COUNT(*) FROM admin_users") or 0) pager = pagination(page, per_page, total) rows = await pool.fetch( "SELECT id, login, role, is_active, created_at FROM admin_users ORDER BY id LIMIT $1 OFFSET $2", pager["per_page"], pager["offset"], ) users_rows = [] for row in rows: item = dict(row) item["created_at_fmt"] = format_dt(item.get("created_at")) users_rows.append(item) return templates.TemplateResponse("users.html", base_context(request, user, users=users_rows, pagination=pager)) @app.post("/users/create") async def user_create( request: Request, csrf_token: str = Form(...), login: str = Form(...), password: str = Form(...), role: str = Form("editor"), ): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) if user["role"] != "admin": return redirect("/users") role = role if role in {"admin", "editor", "viewer"} else "editor" pool = await get_pool() await pool.execute( "INSERT INTO admin_users(login, password_hash, role) VALUES($1, $2, $3)", login.strip().lower(), hash_password(password), role, ) await audit(user["id"], "user.create", "user", None, {"login": login, "role": role}) return redirect("/users") @app.post("/users/{user_id}/toggle") async def user_toggle(request: Request, user_id: int, csrf_token: str = Form(...)): user = await get_current_user(request) if not user: return redirect("/login") require_csrf(user, csrf_token) if user["role"] != "admin" or int(user["id"]) == user_id: return redirect("/users") pool = await get_pool() await pool.execute("UPDATE admin_users SET is_active=NOT is_active, updated_at=NOW() WHERE id=$1", user_id) await audit(user["id"], "user.toggle", "user", user_id) return redirect("/users")