# VESTI Web — локальный интерфейс управления (FastAPI + Jinja2 + Bootstrap 5.3 + HTMX). # Спека vesti-web: 127.0.0.1:8400, авторизация (пароль в .env), списки кандидатов, # подтверждение/отклонение черновиков (→ tg-publisher + банк), опубликованные с метриками. import json import os import secrets import sqlite3 from datetime import datetime from pathlib import Path from urllib.parse import quote import markdown as md_lib from fastapi import FastAPI, Form, HTTPException, Request, Response from fastapi.responses import FileResponse, HTMLResponse, RedirectResponse from fastapi.staticfiles import StaticFiles from jinja2 import Environment, FileSystemLoader from markupsafe import Markup, escape try: # bleach — санитайзер HTML из RSS (allowlist); см. requirements.txt import bleach except ImportError: bleach = None # CRUD источников: работа с sources.yaml (источник истины) + БД синк from sources.sources import ( db_source_to_yaml, add_or_update_source, delete_source, set_source_enabled, ) BASE_DIR = Path(__file__).resolve().parent.parent DB_PATH = BASE_DIR / "db" / "vesti.db" TPL_DIR = Path(__file__).resolve().parent / "templates" STATIC_DIR = Path(__file__).resolve().parent / "static" # Медиа постов: свежие — media/; старые (исторический баг) — media/media/ MEDIA_DIRS = [ BASE_DIR / "media", BASE_DIR / "media" / "media", ] # Авторизация: одна учётка, пароль из .env (VESTI_WEB_PASSWORD или ADMIN_PASSWORD — legacy имя) WEB_PASSWORD = os.getenv("VESTI_WEB_PASSWORD") or os.getenv("ADMIN_PASSWORD") or "admin" # TODO: в проде только из .env SESSION_COOKIE = "vesti_session" app = FastAPI(title="VESTI Web") app.mount("/static", StaticFiles(directory=str(STATIC_DIR)), name="static") tpl = Environment(loader=FileSystemLoader(str(TPL_DIR))) tpl.filters["from_json"] = lambda s: json.loads(s or "[]") def md_filter(text: str) -> Markup: """Markdown → безопасный HTML (экранируем HTML до парсинга, XSS-safe).""" if not text: return Markup("") html = md_lib.markdown(escape(text), extensions=["nl2br", "sane_lists"]) return Markup(html) # Теги/атрибуты, разрешённые в текстах из RSS (безопасное подмножество HTML). _SAFE_TAGS = [ "p", "br", "b", "strong", "i", "em", "u", "s", "strike", "del", "a", "ul", "ol", "li", "blockquote", "pre", "code", "h1", "h2", "h3", "h4", "h5", "h6", "table", "thead", "tbody", "tr", "th", "td", ] _SAFE_ATTRS = {"a": ["href", "title", "rel"]} _SAFE_PROTOCOLS = ["http", "https", "mailto"] def safe_html_filter(text: str) -> Markup: """Исходный текст из RSS → безопасный HTML (allowlist-санитайзер). Сохраняет «заложенное форматирование» (жирный, ссылки, списки, абзацы), вырезает всё опасное (script/style/on*/iframe/svg/...) и делает ссылки кликабельными (linkify). Для markdown-пересказов/комментариев НЕ используется — они идут через md_filter. """ if not text: return Markup("") if bleach is None: # запасной вариант без bleach — экранируем (безопасно, но без разметки) return Markup(md_lib.markdown(escape(text), extensions=["nl2br", "sane_lists"])) try: cleaned = bleach.clean( text, tags=_SAFE_TAGS, attributes=_SAFE_ATTRS, protocols=_SAFE_PROTOCOLS, strip=True, ) cleaned = bleach.linkify(cleaned, parse_email=True) return Markup(cleaned) except Exception: return Markup(escape(text)) def dt_filter(value) -> str: """2026-08-15T15:53 → 15:53 15.08.2026; None/мусор → ''.""" if not value: return "" try: if isinstance(value, str): value = datetime.fromisoformat(value.replace("Z", "+00:00")) return value.strftime("%H:%M %d.%m.%Y") except (ValueError, TypeError): return str(value)[:16] if value else "" tpl.filters["markdown"] = md_filter tpl.filters["safe_html"] = safe_html_filter tpl.filters["dt"] = dt_filter def _db() -> sqlite3.Connection: conn = sqlite3.connect(DB_PATH, timeout=10) conn.row_factory = sqlite3.Row conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA busy_timeout=8000") return conn def _migrate(): """Идемпотентные миграции при старте веб-сервиса.""" conn = sqlite3.connect(DB_PATH, timeout=10) try: cols = [r[1] for r in conn.execute("PRAGMA table_info(posts)").fetchall()] if "is_read" not in cols: conn.execute("ALTER TABLE posts ADD COLUMN is_read INTEGER DEFAULT 0") conn.commit() finally: conn.close() _migrate() def _session_ok(request: Request) -> bool: return request.cookies.get(SESSION_COOKIE) == os.getenv("VESTI_WEB_SESSION", "") def _require_auth(request: Request): if not _session_ok(request): raise HTTPException(status_code=303, headers={"Location": "/login"}) @app.get("/", response_class=HTMLResponse) def index(request: Request): _require_auth(request) return RedirectResponse(url="/candidates", status_code=302) @app.get("/login", response_class=HTMLResponse) def login_page(request: Request): if _session_ok(request): return RedirectResponse(url="/candidates", status_code=302) return HTMLResponse(tpl.get_template("login.html").render(request=request)) @app.post("/login") def login(request: Request, password: str = Form(...)): if secrets.compare_digest(password, WEB_PASSWORD): sess = secrets.token_hex(16) # сохраняем сессию (в проде — в env; для прототипа — атрибут) os.environ["VESTI_WEB_SESSION"] = sess resp = RedirectResponse(url="/candidates", status_code=302) resp.set_cookie(SESSION_COOKIE, sess, httponly=True) return resp return RedirectResponse(url="/login?error=1", status_code=302) DIRECTIONS = ["linux", "tech", "politics", "games", "electronics", "llm"] class _Row(dict): """sqlite3.Row → dict с доступом через точку (для шаблонов).""" def __getattr__(self, k): try: return self[k] except KeyError: raise AttributeError(k) def _row(r): return _Row(dict(r)) if r is not None else None def _cand_back(group_by: str = "", direction: str = "", own: str = "", q: str = "", src: str = "candidates", **extra): """URL возврата в список (src: candidates|selected) с сохранением группировки и фильтров.""" parts = [] if group_by: parts.append(f"group_by={quote(group_by)}") if direction: parts.append(f"direction={quote(direction)}") if own in ("0", "1"): parts.append(f"own={own}") if q: parts.append(f"q={quote(q)}") for k, v in extra.items(): if v: parts.append(f"{k}={quote(str(v))}") return (f"/{src}?" + "&".join(parts)) if parts else f"/{src}" STATUS_LABELS = { "new": "💎 Новые", "selected": "⭐ Отобранные", "rejected": "🗑 Отклонённые", "published": "✅ Опубликованные", "": "Все", } STATUS_COLORS = {"new": "success", "selected": "warning", "rejected": "danger", "published": "primary"} def _group_label(group_by: str, key: str) -> str: """Человекочитаемая подпись группы для списка кандидатов.""" if group_by == "source": return key or "без источника" if group_by == "date": if key == "today": return "Сегодня" if key == "yesterday": return "Вчера" if key == "week": return "Ранее на этой неделе" return f"Ранее · {key}" # status return STATUS_LABELS.get(key, key or "?") def _fetch_candidates(conn, direction: str, status: str, own: str, q: str, group_by: str, limit: int = 200): """Выборка кандидатов для левой панели. Возвращает (grouped, counts, flat). flat — сгруппированный список; counts — {статус: число} для фильтра-табов.""" where = [] params = [] # Кандидаты = только внешние (is_own=0) посты в статусе new. # Свои посты канала (is_own=1) не являются кандидатами — это контент # собственного канала, управляется отдельно (fan-out); published/rejected — # уже решённые посты, им не место в очереди кандидатов. where.append("p.is_own=0") if not status: # по умолчанию — только кандидаты (new); явный статус ниже where.append("p.status='new'") else: # явный статус (new/rejected/published) — единственный фильтр статуса where.append("p.status=?") params.append(status) if direction: where.append("p.direction=?") params.append(direction) if own in ("1", "0"): where.append("p.is_own=?") params.append(int(own)) if q: where.append("(p.text LIKE ? OR p.summary LIKE ?)") params += [f"%{q}%", f"%{q}%"] w = ("WHERE " + " AND ".join(where)) if where else "" sql = f"""SELECT p.*, s.name source_name, s.slug source_slug, s.lang source_lang, s.direction AS source_direction FROM posts p LEFT JOIN sources s ON s.id=p.source_id {w}""" rows = conn.execute(sql + " ORDER BY COALESCE(p.published_at,p.fetched_at) DESC LIMIT ?", params + [limit]).fetchall() # counts для табов статуса: сколько кандидатов (внешних new) и сколько # решённых (rejected/published) доступно для просмотра. counts = {r["status"]: r["c"] for r in conn.execute( "SELECT p.status, COUNT(*) c FROM posts p WHERE p.is_own=0 GROUP BY p.status", )} counts[""] = sum(counts.values()) # статус по умолчанию для табов: если в выборке нет ничего "нового" — смотрим rejected if status == "new" and not rows: for s in ("rejected", "published"): if counts.get(s): status = s break flat = [] for r in rows: d = dict(r) d["source_name"] = d.get("source_name") or d.get("tg_channel") or "?" flat.append(_Row(d)) # группировка import itertools def _key_source(p): return p["source_name"] or "?" def _key_date(p): from datetime import date, timedelta today = date.today() ts = p.get("published_at") or p.get("fetched_at") or "" d = None try: d = datetime.fromisoformat(str(ts)[:19]).date() except ValueError: pass if d is None: return "other" if d == today: return "today" if d == today - timedelta(days=1): return "yesterday" if (today - d).days <= 7: return "week" return f"{d:%m.%Y}" def _key_status(p): return p["status"] keyf = {"source": _key_source, "date": _key_date, "status": _key_status}[group_by] def _group_status(items): """Статус группы: приоритет new > rejected > published, иначе — первый. При группировке по дате первый пост группы — свежайший и может быть published/rejected, хотя вся группа — новые кандидаты; статус группы должен отражать состав, а не первый пост.""" st = [p["status"] for p in items] for s in ("new", "rejected", "published"): if s in st: return s return items[0]["status"] grouped = [] for k, it in itertools.groupby(sorted(flat, key=keyf), key=keyf): items = list(it) grouped.append({ "key": k, "label": _group_label(group_by, k), "status": _group_status(items), "posts": items, }) # порядок групп: new → rejected → published (для status), иначе по ключу if group_by == "status": grouped.sort(key=lambda g: {"new": 0, "rejected": 1, "published": 2, "": 3}.get(g["key"], 9)) return grouped, counts, status @app.get("/candidates", response_class=HTMLResponse) def candidates(request: Request, direction: str = "", status: str = "", own: str = "", q: str = "", group_by: str = "date", selected: int = 0, error: str = ""): _require_auth(request) conn = _db() grouped, counts, status = _fetch_candidates(conn, direction, status, own, q, group_by) # выбор по умолчанию: первый кандидат списка (или указанный selected) sel = None if grouped: sel = grouped[0]["posts"][0] if selected: for g in grouped: for it in g["posts"]: if it["id"] == selected: sel = it break else: continue break # открытие поста = прочтение (как в почте): сразу помечаем, чтобы карточка перестала быть жирной if sel is not None and sel["id"]: conn.execute("UPDATE posts SET is_read=1 WHERE id=? AND is_read=0", (sel["id"],)) sel["is_read"] = 1 conn.commit() conn.close() # строка параметров фильтров для ссылок (выбор поста, смена группировки) qs = [] for k, v in (("direction", direction), ("status", status), ("own", own), ("q", q)): if v: qs.append(f"{k}={quote(v)}") if not selected: selected = sel["id"] if sel else 0 filters_qs = "&".join(qs) html = tpl.get_template("candidates.html").render( posts=grouped, counts=counts, status=status, direction=direction, own=own, q=q, group_by=group_by, selected=sel, error=error, filters_qs=filters_qs, directions=DIRECTIONS, STATUS_LABELS=STATUS_LABELS, STATUS_COLORS=STATUS_COLORS, ) return HTMLResponse(html) @app.get("/selected", response_class=HTMLResponse) def selected(request: Request, direction: str = "", own: str = "", q: str = "", group_by: str = "date", selected: int = 0, error: str = ""): """Отобранные: только status='selected' (внешние). Полная обработка постов.""" _require_auth(request) conn = _db() # отобранные = внешние посты со статусом selected grouped, counts, _ = _fetch_candidates(conn, direction, "selected", own, q, group_by) counts = {} sel = None if grouped: sel = grouped[0]["posts"][0] if selected: for g in grouped: for it in g["posts"]: if it["id"] == selected: sel = it break else: continue break if sel is not None and sel["id"]: conn.execute("UPDATE posts SET is_read=1 WHERE id=? AND is_read=0", (sel["id"],)) sel["is_read"] = 1 conn.commit() conn.close() qs = [] for k, v in (("direction", direction), ("own", own), ("q", q)): if v: qs.append(f"{k}={quote(v)}") if not selected: selected = sel["id"] if sel else 0 filters_qs = "&".join(qs) html = tpl.get_template("selected.html").render( posts=grouped, counts=counts, status="selected", direction=direction, own=own, q=q, group_by=group_by, selected=sel, error=error, filters_qs=filters_qs, directions=DIRECTIONS, STATUS_LABELS=STATUS_LABELS, STATUS_COLORS=STATUS_COLORS, ) return HTMLResponse(html) @app.post("/posts/{post_id}/select") def post_select(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form("")): """Отобрать кандидата (new → selected). Остаёмся на странице кандидатов.""" _require_auth(request) conn = _db() conn.execute("UPDATE posts SET status='selected' WHERE id=? AND status='new' AND is_own=0", (post_id,)) conn.commit() conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src="candidates"), status_code=302) @app.post("/candidates/select") def candidates_select(request: Request, id: int = Form(...), direction: str = Form(""), status: str = Form(""), own: str = Form(""), q: str = Form(""), group_by: str = Form("date")): """Выбор кандидата левой панели (обычная форма без JS).""" _require_auth(request) url = f"/candidates?selected={id}&group_by={group_by}" for k, v in (("direction", direction), ("status", status), ("own", own), ("q", q)): if v: url += f"&{k}={v}" return RedirectResponse(url=url, status_code=302) @app.post("/candidates/bulk") def candidates_bulk(request: Request, action: str = Form(...), ids: list[int] = Form(default_factory=list), direction: str = Form(""), status: str = Form(""), own: str = Form(""), q: str = Form(""), group_by: str = Form("date"), src: str = Form("candidates")): """Bulk-операции: select/reject/reject-old/process/approve по отмеченным ids. src — откуда вызывались (candidates|selected), туда и возвращаемся.""" _require_auth(request) back = f"/{src}?group_by={group_by}" for k, v in (("direction", direction), ("status", status), ("own", own), ("q", q)): if v: back += f"&{k}={quote(v)}" if not ids: return RedirectResponse(url=back + "&error=no_selection", status_code=302) conn = _db() ph = ", ".join("?" * len(ids)) rows = conn.execute(f"SELECT * FROM posts WHERE id IN ({ph})", ids).fetchall() if action == "select": # отбор кандидатов (new → selected); остаёмся на странице кандидатов conn.execute("UPDATE posts SET status='selected' WHERE id IN ({}) AND status='new' AND is_own=0".format(ph), ids) conn.commit() conn.close() return RedirectResponse(url=back, status_code=302) if action == "reject": conn.execute(f"UPDATE posts SET status='rejected' WHERE id IN ({ph})", ids) conn.commit() conn.close() # остаёмся в текущем списке (кандидаты/отобранные), не уходим в отклонённые return RedirectResponse(url=back, status_code=302) if action == "reject-old": # все выбранные, кроме самых свежих (по 1 от источника) — «почистить старьё» keep = {} for r in rows: key = r["source_id"] keep.setdefault(key, r["id"]) kill = [r["id"] for r in rows if r["id"] not in keep.values()] if kill: ph2 = ", ".join("?" * len(kill)) conn.execute(f"UPDATE posts SET status='rejected' WHERE id IN ({ph2})", kill) conn.commit() conn.close() return RedirectResponse(url=back, status_code=302) conn.close() if action == "process": # обработка локальной моделью (переклассификация) по выбранным errs = [] done = 0 for r in rows: try: impl_reclassify(int(r["id"])) done += 1 except Exception as e: errs.append(f"{r['id']}({e})") if errs: return RedirectResponse(url=f"/{src}?error=process_partial:{','.join(map(str,errs))[:120]}&group_by={group_by}", status_code=302) return RedirectResponse(url=back + "&error=processed:" + str(done), status_code=302) if action == "rewrite": # сгенерировать черновик своего поста для каждого выбранного errs = [] done = 0 for r in rows: try: from classifier.classify import call_ollama d = dict(r) text_ = (d.get("text") or "")[:2000] prompt = ( "Ты — редактор новостей. Перепиши новость СВОИМИ словами от первого лица\n" "(как автор канала «Дед в АйТи»), без копипасты и без канцелярита. Сохрани смысл,\n" "факты и ссылку на оригинал. Верни ТОЛЬКО текст пересказа, без пояснений.\n\n" f"Оригинал:\n{text_}" ) res = call_ollama(prompt, 0, 0, str(d.get("direction") or "")) draft = (res.get("text") or res.get("content") or "").strip() or res.get("summary") or "" if draft: conn = _db() conn.execute("UPDATE posts SET rewritten_text=? WHERE id=?", (draft, int(r["id"]))) conn.commit() conn.close() done += 1 else: errs.append(f"{r['id']}(empty)") except Exception as e: errs.append(f"{r['id']}({e})") errq = ("&error=rewrite_partial:" + ",".join(map(str, errs))[:120]) if errs else "" return RedirectResponse(url=back + "&error=rewritten:" + str(done) + errq, status_code=302) # approve — каждый через существующую логику (публикация + бандл) errs = [] for r in rows: if r["status"] == "published": continue err = impl_approve(int(r["id"])) if err: errs.append(f"{r['id']}({err})") if errs: return RedirectResponse(url=f"/{src}?error=bulk_partial:{','.join(map(str,errs))[:120]}&group_by={group_by}", status_code=302) return RedirectResponse(url="/published", status_code=302) @app.post("/posts/{post_id}/rewrite-save") async def post_rewrite_save(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): """Сохраняет отредактированный пользователем пересказ (rewritten_text).""" _require_auth(request) conn = _db() p = conn.execute("SELECT id FROM posts WHERE id=?", (post_id,)).fetchone() if not p: conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, error="notfound"), status_code=302) form = await request.form() draft = str(form.get("rewritten_text") or "").strip() conn.execute("UPDATE posts SET rewritten_text=? WHERE id=?", (draft or None, post_id)) conn.commit() conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id), status_code=302) @app.post("/posts/{post_id}/comment") def post_comment(post_id: int, request: Request, comment: str = Form(""), group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): _require_auth(request) conn = _db() conn.execute("UPDATE posts SET comment=? WHERE id=?", (comment.strip(), post_id)) conn.commit() conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id), status_code=302) @app.post("/posts/{post_id}/reclassify") def post_reclassify(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): """Переклассификация локальной моделью (qwen).""" _require_auth(request) try: impl_reclassify(post_id) except Exception as e: return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id, error=f"reclassify:{e}"), status_code=302) return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id), status_code=302) def impl_reclassify(post_id: int): """Перезапуск классификации локальной моделью. Бросает исключение при неудаче.""" from classifier.classify import classify_text conn = _db() p = conn.execute("SELECT * FROM posts WHERE id=?", (post_id,)).fetchone() if not p: conn.close() raise ValueError("notfound") res = classify_text(dict(p)) conn.execute( "UPDATE posts SET direction=?, relevance=?, interest=?, summary=?, classified=? WHERE id=?", (res["direction"], res["relevance"], res["interest"], res["summary"], 1 if res["classified"] else 0, post_id), ) conn.commit() conn.close() @app.post("/posts/{post_id}/rewrite") def post_rewrite(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): _require_auth(request) from classifier.classify import call_ollama conn = _db() p = conn.execute("SELECT * FROM posts WHERE id=?", (post_id,)).fetchone() if not p: conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, error="notfound"), status_code=302) d = dict(p) text_ = (d.get("text") or "")[:2000] prompt = ( "Ты — редактор новостей. Перепиши новость СВОИМИ словами от первого лица\n" "(как автор канала «Дед в АйТи»), без копипасты и без канцелярита. Сохрани смысл,\n" "факты и ссылку на оригинал. Верни ТОЛЬКО текст пересказа, без пояснений.\n\n" f"Оригинал:\n{text_}" ) try: r = call_ollama(prompt, 0, 0, str(d.get("direction") or "")) draft = (r.get("text") or r.get("content") or "").strip() if not draft: draft = r.get("summary") or "" except Exception as e: conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id, error=f"rewrite:{e}"), status_code=302) conn.execute("UPDATE posts SET rewritten_text=? WHERE id=?", (draft, post_id)) conn.commit() conn.close() return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id), status_code=302) @app.get("/published", response_class=HTMLResponse) def published(request: Request): _require_auth(request) conn = _db() rows = conn.execute( """SELECT p.*, s.name source_name, pub.tg_message_id, pub.views pub_views, pub.bundle_path, pub.distributed_dirs FROM posts p LEFT JOIN sources s ON s.id=p.source_id LEFT JOIN published pub ON pub.post_id=p.id WHERE p.status='published' ORDER BY p.published_at DESC LIMIT 100""" ).fetchall() conn.close() html = tpl.get_template("published.html").render(posts=rows) return HTMLResponse(html) def impl_approve(post_id: int, comment: str = "", dirs_selected: list[str] | None = None) -> str | None: """Публикация поста: publisher + бандл + статус published. Возвращает None при успехе, строку ошибки при неудаче (для bulk — собираем ошибки, для POST — редирект с error).""" from web.publisher_client import publish as http_publish from publisher.card import make_card from web.store import create_bundle conn = _db() post = conn.execute("SELECT p.*, s.name source_name, s.lang FROM posts p LEFT JOIN sources s ON s.id=p.source_id WHERE p.id=?", (post_id,)).fetchone() if not post: conn.close() return "notfound" post = dict(post) is_own = int(post.get("is_own") or 0) == 1 dirn = post.get("direction") or "linux" lang = post.get("lang") or "ru" if dirs_selected is None: # по умолчанию: направления классификации (classifications) или направление поста cls = conn.execute("SELECT direction FROM classifications WHERE post_id=? ORDER BY id", (post_id,)).fetchall() dirs_selected = [r["direction"] for r in cls] if cls else [dirn] dirs_selected = list(dict.fromkeys([d for d in dirs_selected if d])) card = make_card(post, comment) card["direction"] = dirn card["lang"] = lang try: res = http_publish(card) except RuntimeError as e: conn.close() return f"publisher:{e}" _res = res.get("results", {}) # свой контент: ключи — направления; внешний: ключи — каналы if is_own: tg_message_id = next(iter((r or {}).get("message_id") or 0 for r in _res.values()), 0) if res else 0 views_by_dir = {d: (r.get("views") or 0) for d, r in _res.items()} else: tg_message_id = next((r.get("message_id") or 0 for r in _res.values() if r.get("message_id")), 0) if _res else 0 views_by_dir = {d: (r.get("views") or 0) for d, r in _res.items()} # банк статей (markdown-бандл + медиа) — по каждому направлению fan-out try: bundle = create_bundle(post, bundles_dir=BASE_DIR / "bundles", directions=dirs_selected) except Exception as e: bundle = {"error": str(e)} import json as _json conn.execute("UPDATE posts SET status='published', comment=? WHERE id=?", (comment.strip(), post_id)) distributed = _json.dumps(dirs_selected, ensure_ascii=False) if is_own else _json.dumps([dirn]) conn.execute( "INSERT INTO published (post_id, bundle_path, tg_message_id, distributed_dirs, views) VALUES (?,?,?,?,?)", (post_id, bundle.get("path", ""), tg_message_id, distributed, sum(views_by_dir.values())), ) conn.commit() conn.close() return None @app.post("/posts/{post_id}/approve") async def approve(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): """Подтверждение черновика → публикация в TG (fan-out по направлениям) + создание бандла + статус published.""" _require_auth(request) # направления fan-out: из формы (чекбоксы) — по умолчанию направления классификации или направление поста dirs_selected: list[str] = [] comment = "" try: form = await request.form() dirs_selected = [str(d) for d in form.getlist("dirs")] comment = str(form.get("comment") or "").strip() except Exception: pass err = impl_approve(post_id, comment, dirs_selected) if err: return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, selected=post_id, error=err), status_code=302) # после публикации остаёмся в своей группировке, смотрим опубликованные return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src, status="published"), status_code=302) @app.post("/posts/{post_id}/reject") def reject(post_id: int, request: Request, group_by: str = Form(""), direction: str = Form(""), own: str = Form(""), q: str = Form(""), src: str = Form("candidates")): _require_auth(request) conn = _db() conn.execute("UPDATE posts SET status='rejected' WHERE id=?", (post_id,)) conn.commit() conn.close() # после отклонения остаёмся на странице, откуда пришли (кандидаты/отобранные), # в той же группировке/фильтрах; список отклонённых — через фильтр «🗑 Отклонённые» return RedirectResponse(url=_cand_back(group_by, direction, own, q, src=src), status_code=302) @app.get("/metrics", response_class=HTMLResponse) def metrics(request: Request): _require_auth(request) conn = _db() rows = conn.execute( """SELECT p.direction, COUNT(*) cnt, SUM(p.views) views, AVG(p.reactions_total) avg_reactions FROM posts p GROUP BY p.direction""" ).fetchall() runs = conn.execute("SELECT * FROM runs ORDER BY started_at DESC LIMIT 10").fetchall() conn.close() html = tpl.get_template("metrics.html").render(dirs=rows, runs=runs) return HTMLResponse(html) @app.get("/crawlers", response_class=HTMLResponse) def crawlers(request: Request): """Страница статуса краулеров: источники (alive/dead, last_fetch) + размер очереди.""" _require_auth(request) conn = _db() # источники + rss_state (для RSS — etag/modified/last_build_date/ошибки) sources = conn.execute( """SELECT s.id, s.slug, s.name, s.crawler, s.status, s.last_fetch, s.last_error, s.priority, s.enabled, COALESCE(rs.etag, '') AS etag, COALESCE(rs.modified, '') AS modified, COALESCE(rs.last_build_date, '') AS last_build_date, COALESCE(rs.error_count, 0) AS error_count, COALESCE(s.last_error, rs.last_error) AS last_error FROM sources s LEFT JOIN rss_state rs ON rs.slug=s.slug ORDER BY s.enabled DESC, s.priority, s.slug""" ).fetchall() # очередь (постов), карточки pending = conn.execute( "SELECT COUNT(*) FROM posts WHERE classified IS NULL OR classified=0" ).fetchone()[0] processing = conn.execute( "SELECT COUNT(*) FROM posts WHERE classified=-1" ).fetchone()[0] done_today = conn.execute( "SELECT COUNT(*) FROM posts WHERE classified=1 AND fetched_at >= date('now')" ).fetchone()[0] # статусы источников для бейджей alive = conn.execute("SELECT COUNT(*) FROM sources WHERE enabled=1 AND status='alive'").fetchone()[0] dead = conn.execute("SELECT COUNT(*) FROM sources WHERE enabled=1 AND status='dead'").fetchone()[0] total_sources = conn.execute("SELECT COUNT(*) FROM sources WHERE enabled=1").fetchone()[0] # последние запуски runs = conn.execute( """SELECT r.*, s.slug, s.name AS source_name FROM runs r LEFT JOIN sources s ON s.id=r.source_id ORDER BY r.id DESC LIMIT 12""" ).fetchall() conn.close() html = tpl.get_template("crawlers.html").render( request=request, sources=sources, pending=pending, processing=processing, done_today=done_today, alive=alive, dead=dead, total_sources=total_sources, runs=runs, ) return HTMLResponse(html) @app.post("/crawlers/reset") def crawlers_reset(request: Request): """Сброс dead-источников в alive (после починки).""" _require_auth(request) conn = _db() cur = conn.execute( "UPDATE sources SET status='alive', error_count=0, last_error=NULL WHERE status='dead'" ) conn.commit() conn.close() return RedirectResponse(url="/crawlers?reset=" + str(cur.rowcount), status_code=302) @app.post("/crawlers/run") async def crawlers_run(request: Request): """Запуск воркера (фаза 3): параллельная обработка очереди кандидатов. Запускается в фоне (subprocess.Popen, не блокирует веб); страница обновляется, очередь processing покажет прогресс. Результат — в runs/posts. """ _require_auth(request) workers, limit = 8, 200 try: form = await request.form() workers = int(str(form.get("workers") or 8)) limit = int(str(form.get("limit") or 200)) except Exception: pass try: import subprocess py = str(Path(__file__).resolve().parent.parent / ".venv" / "bin" / "python") log_path = Path(__file__).resolve().parent.parent / "logs" / "worker.log" log_path.parent.mkdir(exist_ok=True) with open(log_path, "a") as f: subprocess.Popen( [py, "-m", "crawler.worker", "--workers", str(workers), "--limit", str(limit)], cwd=str(Path(__file__).resolve().parent.parent), stdout=f, stderr=subprocess.STDOUT, ) except Exception as e: return RedirectResponse(url=f"/crawlers?run=error&msg={e}", status_code=302) return RedirectResponse(url=f"/crawlers?run=started&workers={workers}&limit={limit}", status_code=302) @app.get("/sources", response_class=HTMLResponse) def sources_page(request: Request): """Страница управления реестром источников: таблица + добавление/редактирование.""" _require_auth(request) conn = _db() sources = conn.execute( """SELECT s.*, COALESCE(rs.error_count, 0) AS error_count, COALESCE(s.last_error, rs.last_error) AS last_error, COALESCE(rs.etag, '') AS etag, COALESCE(rs.modified, '') AS modified FROM sources s LEFT JOIN rss_state rs ON rs.slug=s.slug ORDER BY s.enabled DESC, s.priority, s.slug""" ).fetchall() conn.close() html = tpl.get_template("sources.html").render( request=request, sources=sources, did=request.query_params.get("did", ""), ) return HTMLResponse(html) def _form_source(form) -> dict: """Собирает источник из формы (str поля, None для пустых).""" def s(k, default=None): v = form.get(k) return (v or "").strip() if v else default return { "slug": s("slug"), "name": s("name"), "url": s("url"), "channel": s("channel"), "crawler": s("crawler", "telegram"), "feed_url": s("feed_url"), "direction": s("direction"), "lang": s("lang", "ru"), "priority": s("priority", "P1"), "enabled": bool(form.get("enabled")), "own": bool(form.get("own")), } @app.post("/sources/add") async def sources_add(request: Request): """Добавление источника: валидация → sources.yaml → полный ре-синк в БД.""" _require_auth(request) form = await request.form() src = _form_source(form) ok, res = add_or_update_source(src) if not ok: return RedirectResponse(url=f"/sources?did=error:{res}", status_code=302) return RedirectResponse(url=f"/sources?did=added:{res}", status_code=302) @app.post("/sources/{source_id}/edit") async def sources_edit(source_id: int, request: Request): """Редактирование источника: обновляем yaml (slug не меняем) + полный ре-синк.""" _require_auth(request) conn = _db() row = conn.execute("SELECT * FROM sources WHERE id=?", (source_id,)).fetchone() conn.close() if not row: return RedirectResponse(url="/sources?did=error:не найден", status_code=302) form = await request.form() src = db_source_to_yaml(dict(row)) # текстовые поля for k in ("name", "url", "channel", "crawler", "feed_url", "direction", "lang", "priority"): if k in form: v = str(form.get(k) or "").strip() src[k] = v or None # чекбоксы src["enabled"] = bool(form.get("enabled")) src["own"] = bool(form.get("own")) ok, res = add_or_update_source(src) if not ok: return RedirectResponse(url=f"/sources?did=error:{res}", status_code=302) return RedirectResponse(url=f"/sources?did=updated:{res}", status_code=302) @app.post("/sources/{source_id}/toggle") async def sources_toggle(source_id: int, request: Request): """Пауза/возобновление: flip enabled в yaml + БД.""" _require_auth(request) conn = _db() row = conn.execute("SELECT * FROM sources WHERE id=?", (source_id,)).fetchone() conn.close() if not row: return RedirectResponse(url="/sources?did=error:не найден", status_code=302) if not set_source_enabled(row["slug"], not bool(row["enabled"])): return RedirectResponse(url="/sources?did=error:не найден", status_code=302) state = "paused" if row["enabled"] else "resumed" return RedirectResponse(url=f"/sources?did={state}:{row['slug']}", status_code=302) @app.post("/sources/{source_id}/delete") async def sources_delete(source_id: int, request: Request): """Удаление источника: yaml → БД, каскад (посты/раны остаются с source_id=NULL).""" _require_auth(request) conn = _db() row = conn.execute("SELECT * FROM sources WHERE id=?", (source_id,)).fetchone() conn.close() if not row: return RedirectResponse(url="/sources?did=error:не найден", status_code=302) if not delete_source(row["slug"]): return RedirectResponse(url=f"/sources?did=error:нет в yaml:{row['slug']}", status_code=302) return RedirectResponse(url=f"/sources?did=deleted:{row['slug']}", status_code=302) @app.get("/logout") def logout(): resp = RedirectResponse(url="/login", status_code=302) resp.delete_cookie(SESSION_COOKIE) return resp @app.get("/bundle/{bundle_path:path}") def bundle(request: Request, bundle_path: str): """Показывает markdown-бандл как текст (ссылка из published).""" _require_auth(request) full = (BASE_DIR / "bundles" / bundle_path).resolve() # защита от path traversal if not str(full).startswith(str((BASE_DIR / "bundles").resolve())): return HTMLResponse("bad path", status_code=400) if not full.exists(): return HTMLResponse("not found", status_code=404) return Response(full.read_text(encoding="utf-8"), media_type="text/markdown") @app.get("/media/{filename}") def media(request: Request, filename: str): """Отдаёт медиа-файл поста (из media/ или media/media/). Авторизация.""" _require_auth(request) name = os.path.basename(filename) # защита от path traversal if not name: return HTMLResponse("bad filename", status_code=400) for d in MEDIA_DIRS: f = (d / name).resolve() if f.exists() and f.is_file(): return FileResponse(f) return HTMLResponse("not found", status_code=404)