mirror of
https://gitverse.ru/kpa39l/vesti.git
synced 2026-09-29 18:05:03 +00:00
185 lines
8.2 KiB
Python
185 lines
8.2 KiB
Python
#!/usr/bin/env python3
|
|
"""VESTI — бэкфилл своего канала @dedinit (own-content-hub, задача 2.3).
|
|
|
|
Итерирует ВСЕ посты канала dedinit (с конца, без лимита MAX_POSTS_PER_CHANNEL),
|
|
складывает в БД через store_posts (own-семантика: is_own=1, канон, дедуп по sha256/url),
|
|
скачивает медиа своего канала, пишет tg_state.last_post_id = max (без перечитывания
|
|
при инкрементальных запусках).
|
|
|
|
Запуск:
|
|
cd /opt/vesti
|
|
.venv/bin/python -m crawler.backfill_dedinit [--limit N] [--no-media] [--max-errors N]
|
|
|
|
--limit: сколько постов обработать (для теста; по умолчанию — все).
|
|
--no-media: не скачивать медиа (текст-бэкфилл; позже докачать отдельно).
|
|
--max-errors: сколько подряд идущих ошибок выдержать (по умолчанию 20).
|
|
"""
|
|
import argparse
|
|
import asyncio
|
|
import logging
|
|
import sys
|
|
import time
|
|
from pathlib import Path
|
|
|
|
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
|
|
|
from config import TG_PROXY, TG_SESSION_DIR, ensure_dirs
|
|
from db.db import db
|
|
from sources.sources import get_source, start_run, finish_run
|
|
from crawler.telegram_crawler import (
|
|
make_client, store_posts, detect_media_type, post_media_path,
|
|
sha256_text,
|
|
)
|
|
|
|
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s %(message)s")
|
|
log = logging.getLogger("vesti.backfill")
|
|
|
|
|
|
async def download_media_safe(client, msg, rel_path):
|
|
"""Скачивает медиа, возвращает (rel_path|None). Безопасная версия для бэкфилла.
|
|
|
|
Сохраняет в MEDIA_DIR/<basename> (= /opt/vesti/media/<file>), как и download_media
|
|
в telegram_crawler — чтобы media_path из БД (media/<file>) совпадал с фактическим файлом.
|
|
"""
|
|
import os as _os
|
|
from config import MEDIA_DIR
|
|
full = _os.path.join(str(MEDIA_DIR), _os.path.basename(rel_path))
|
|
_os.makedirs(str(MEDIA_DIR), exist_ok=True)
|
|
if _os.path.exists(full):
|
|
return rel_path
|
|
try:
|
|
dl = await client.download_media(msg, file=full)
|
|
except Exception as e:
|
|
log.warning("media dl %s: %s", rel_path, e)
|
|
return None
|
|
return rel_path if dl else None
|
|
|
|
|
|
async def backfill(client, source, limit=None, no_media=False, max_errors=20):
|
|
slug = source["slug"]
|
|
channel = source["channel"]
|
|
with db() as conn:
|
|
state = conn.execute("SELECT * FROM tg_state WHERE channel_slug=?", (slug,)).fetchone()
|
|
state = dict(state) if state else {}
|
|
last_id = state.get("last_post_id") or 0
|
|
|
|
run_id = start_run(source["id"], f"telegram:{slug}", "backfill")
|
|
t0 = time.monotonic()
|
|
fetched = 0
|
|
errors = 0
|
|
min_id_seen = None
|
|
|
|
try:
|
|
entity = await client.get_entity(channel)
|
|
except Exception as e:
|
|
finish_run(run_id, "error", 0, 0, f"get_entity: {e}")
|
|
raise
|
|
|
|
async for msg in client.iter_messages(entity, reverse=False):
|
|
if limit and fetched >= limit:
|
|
break
|
|
# Бэкфилл НЕ пропускает по last_post_id (частично могли забежать вперёд):
|
|
# store_posts сам отсеет дубликаты по sha256/url. Просто идём от свежего к старому.
|
|
min_id_seen = msg.id if min_id_seen is None else min(min_id_seen, msg.id)
|
|
|
|
# Дубликаты-форварды: если пост — пересланный из нашего же канала или из канала-источника,
|
|
# то store_posts сам разрулит (is_own=0 упоминание). Для фулл-бэкфилла пропускаем чужие форварды:
|
|
fw = getattr(msg, "fwd_from", None)
|
|
is_forward = fw is not None and getattr(getattr(fw, "from_id", None), "channel_id", None) is not None
|
|
|
|
try:
|
|
ctype = detect_media_type(msg)
|
|
text = msg.text or ""
|
|
url = f"https://t.me/{channel}/{msg.id}"
|
|
media_rel = None
|
|
if not (is_forward or no_media) and ctype != "text":
|
|
media_rel = post_media_path(msg, channel)
|
|
saved = await download_media_safe(client, msg, media_rel)
|
|
media_rel = saved or media_rel
|
|
|
|
reactions = {}
|
|
rr = getattr(msg, "reactions", None)
|
|
if rr and getattr(rr, "results", None):
|
|
for r in rr.results:
|
|
if getattr(r, "reaction", None) and getattr(r, "count", None):
|
|
key = getattr(r.reaction, "emoji", None) or str(r.reaction)
|
|
reactions[key] = r.count
|
|
|
|
post = {
|
|
"tg_channel": channel,
|
|
"tg_post_id": msg.id,
|
|
"url": url,
|
|
"text": text,
|
|
"content_type": ctype,
|
|
"media_path": media_rel,
|
|
"views": int(getattr(msg, "views", None) or 0),
|
|
"reactions": __import__("json").dumps(reactions, ensure_ascii=False),
|
|
"reactions_total": sum(reactions.values()),
|
|
"published_at": msg.date.isoformat() if msg.date else None,
|
|
"fwd_from_channel_id": getattr(getattr(fw, "from_id", None), "channel_id", None) if is_forward else None,
|
|
"fwd_from_post_id": getattr(fw, "channel_post", None) if is_forward else None,
|
|
"is_own": 1,
|
|
}
|
|
with db() as conn:
|
|
new = store_posts(conn, source["id"], [post], is_own_source=True)
|
|
fetched += 1
|
|
if new:
|
|
log.info(" + #%d (new=%d) [%s]", msg.id, new, ctype)
|
|
else:
|
|
log.info(" = #%d dup [%s]", msg.id, ctype)
|
|
errors = 0
|
|
except Exception as e:
|
|
errors += 1
|
|
log.warning("ERROR #%d: %s", msg.id, e)
|
|
if errors >= max_errors:
|
|
log.error("Слишком много ошибок подряд (%d) — прерываю", errors)
|
|
break
|
|
|
|
# состояние: last_post_id = реальный max tg_post_id канала (исключая тестовые 9000+,
|
|
# которые были вставлены в прошлых сессиях как искусственные id)
|
|
max_id_seen = None
|
|
with db() as conn:
|
|
r = conn.execute(
|
|
"SELECT tg_post_id FROM posts WHERE tg_channel=? AND tg_post_id < 9000 ORDER BY tg_post_id DESC LIMIT 1",
|
|
(channel,),
|
|
).fetchone()
|
|
max_id_seen = r[0] if r else None
|
|
if max_id_seen:
|
|
with db() as conn:
|
|
conn.execute(
|
|
"INSERT INTO tg_state (channel_slug, last_post_id, last_ts) VALUES (?,?,datetime('now'))"
|
|
" ON CONFLICT(channel_slug) DO UPDATE SET last_post_id=excluded.last_post_id, last_ts=datetime('now')",
|
|
(slug, max_id_seen),
|
|
)
|
|
duration = int((time.monotonic() - t0) * 1000)
|
|
finish_run(run_id, "ok", fetched, fetched, None, duration)
|
|
log.info("DONE %s: fetched=%d (%.1fs), max_post_id=%s", slug, fetched, duration / 1000, max_id_seen)
|
|
|
|
|
|
async def main():
|
|
ap = argparse.ArgumentParser(description="VESTI backfill @dedinit")
|
|
ap.add_argument("--limit", type=int, default=None, help="сколько постов обработать (тест)")
|
|
ap.add_argument("--no-media", action="store_true", help="без скачивания медиа")
|
|
ap.add_argument("--max-errors", type=int, default=20)
|
|
args = ap.parse_args()
|
|
|
|
ensure_dirs()
|
|
TG_SESSION_DIR.mkdir(parents=True, exist_ok=True)
|
|
src = get_source("dedinit")
|
|
if not src:
|
|
log.error("source dedinit not found")
|
|
sys.exit(1)
|
|
|
|
client = make_client(TG_SESSION_DIR)
|
|
try:
|
|
await client.connect()
|
|
if not await client.is_user_authorized():
|
|
log.error("Session not authorized. /opt/vesti/telegram/vesti.session?")
|
|
sys.exit(2)
|
|
await backfill(client, src, limit=args.limit, no_media=args.no_media, max_errors=args.max_errors)
|
|
finally:
|
|
await client.disconnect()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main()) |