mirror of
https://gitverse.ru/kpa39l/vesti.git
synced 2026-09-29 09:55:03 +00:00
Initial import: vesti.nixg.ru — новостной апрув-проект (web, crawler, classifier, publisher, openspec)
This commit is contained in:
@@ -0,0 +1,50 @@
|
||||
# VESTI — разовая интерактивная авторизация Telethon-сессии (MTProto).
|
||||
# Запуск: .venv/bin/python -m crawler.auth_telegram
|
||||
# После успеха сессия /opt/vesti/telegram/vesti.session сохраняется (644? -> 600).
|
||||
import asyncio
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
|
||||
from config import TG_API_ID, TG_API_HASH, TG_PROXY, TG_SESSION_DIR, ensure_dirs
|
||||
from crawler.telegram_crawler import make_client
|
||||
from telethon import errors
|
||||
|
||||
import logging
|
||||
logging.basicConfig(level=logging.INFO)
|
||||
|
||||
|
||||
async def main():
|
||||
ensure_dirs()
|
||||
TG_SESSION_DIR.mkdir(parents=True, exist_ok=True)
|
||||
client = make_client(TG_SESSION_DIR)
|
||||
await client.connect()
|
||||
if await client.is_user_authorized():
|
||||
me = await client.get_me()
|
||||
print(f"Уже авторизован: @{me.username or me.phone}")
|
||||
await client.disconnect()
|
||||
return
|
||||
phone = input("Телефон (с кодом страны, +7...): ").strip()
|
||||
if not phone:
|
||||
print("Нет телефона.")
|
||||
return
|
||||
try:
|
||||
await client.send_code_request(phone)
|
||||
code = input("Код из Telegram: ").strip()
|
||||
# 2FA пароль, если включён
|
||||
try:
|
||||
await client.sign_in(phone, code)
|
||||
except errors.SessionPasswordNeededError:
|
||||
pwd = input("2FA-пароль: ").strip()
|
||||
await client.sign_in(password=pwd)
|
||||
me = await client.get_me()
|
||||
print(f"Авторизован: @{me.username or me.phone}")
|
||||
except Exception as e:
|
||||
print(f"Ошибка: {e}")
|
||||
finally:
|
||||
await client.disconnect()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
@@ -0,0 +1,185 @@
|
||||
#!/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())
|
||||
@@ -0,0 +1,370 @@
|
||||
# VESTI — Telegram-краулер (Telethon / MTProto).
|
||||
# Читает публичные каналы из sources.yaml, инкрементально после last_post_id,
|
||||
# собирает метрики (views/reactions), пишет в SQLite vesti.db.
|
||||
# РАБОТАЕТ ТОЛЬКО через SOCKS5-туннель 127.0.0.1:1080 (из РФ TG заблокирован).
|
||||
#
|
||||
# Запуск по каждому источнику отдельно (задача источника):
|
||||
# python -m crawler.telegram_crawler --source linuxklub
|
||||
# python -m crawler.telegram_crawler --direction linux
|
||||
# python -m crawler.telegram_crawler --all
|
||||
import argparse
|
||||
import asyncio
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
|
||||
from config import TG_API_ID, TG_API_HASH, TG_PROXY, TG_SESSION_DIR, ensure_dirs
|
||||
from db.db import db
|
||||
from sources.sources import (
|
||||
get_enabled_sources, get_source, load_sources_yaml, sync_sources_to_db,
|
||||
start_run, finish_run,
|
||||
)
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s %(message)s")
|
||||
log = logging.getLogger("vesti.crawler")
|
||||
|
||||
MAX_POSTS_PER_CHANNEL = 20 # лимит на запуск (защита от 429/flood)
|
||||
DUPLICATE_KEYS = ("url", "sha256")
|
||||
|
||||
|
||||
def parse_socks(url):
|
||||
"""socks5://user:pass@host:port -> (host, port, user, pass)."""
|
||||
from urllib.parse import urlparse
|
||||
u = urlparse(url)
|
||||
return u.hostname, (u.port or 1080), (u.username or None), (u.password or None)
|
||||
|
||||
|
||||
def make_client(session_dir):
|
||||
"""Создаёт Telethon-клиент. Сессия — session_dir/vesti.session."""
|
||||
from telethon import TelegramClient
|
||||
from telethon.sessions import SQLiteSession
|
||||
|
||||
host, port, user, pwd = parse_socks(TG_PROXY)
|
||||
# Telethon: proxy=(proxy_type, addr, port[, username, password]).
|
||||
# proxy_type — это socks.SOCKS5 / socks.HTTP из pysocks (telethon использует socksio/pysocks).
|
||||
import socks as pysocks
|
||||
proxy = (pysocks.SOCKS5, host, port)
|
||||
if user:
|
||||
proxy = (pysocks.SOCKS5, host, port, user, pwd)
|
||||
client = TelegramClient(SQLiteSession(str(session_dir / "vesti")), TG_API_ID, TG_API_HASH, proxy=proxy)
|
||||
return client
|
||||
|
||||
|
||||
def sha256_text(text: str) -> str:
|
||||
return hashlib.sha256((text or "").encode("utf-8", "ignore")).hexdigest()
|
||||
|
||||
|
||||
def detect_media_type(msg) -> str:
|
||||
"""Определяет тип контента поста по media-полю Message."""
|
||||
if msg.photo:
|
||||
return "photo"
|
||||
if msg.video:
|
||||
return "video"
|
||||
if msg.document:
|
||||
# document может быть file/voice/sticker/video_note
|
||||
mime = (msg.document.mime_type or "") if msg.document else ""
|
||||
if "voice" in mime:
|
||||
return "voice"
|
||||
if "sticker" in mime:
|
||||
return "sticker"
|
||||
return "document"
|
||||
if msg.voice:
|
||||
return "voice"
|
||||
if msg.video_note:
|
||||
return "video_note"
|
||||
if msg.sticker:
|
||||
return "sticker"
|
||||
if msg.animation if hasattr(msg, "animation") else False:
|
||||
return "animation"
|
||||
return "text"
|
||||
|
||||
|
||||
def post_media_path(msg, channel) -> str:
|
||||
"""Строит относительный путь media/<channel>_<id>.<ext> (без скачивания здесь)."""
|
||||
if not channel:
|
||||
channel = "chan"
|
||||
ext = "jpg"
|
||||
mtype = detect_media_type(msg)
|
||||
if mtype == "video":
|
||||
ext = "mp4"
|
||||
elif mtype == "voice":
|
||||
ext = "ogg"
|
||||
elif mtype == "document" and msg.document:
|
||||
# имя файла в атрибутах документа (DocumentAttributesFilename)
|
||||
fname = None
|
||||
doc = msg.document
|
||||
if hasattr(doc, "file_name"):
|
||||
fname = doc.file_name
|
||||
else:
|
||||
for a in (getattr(doc, "attributes", None) or []):
|
||||
if hasattr(a, "file_name"):
|
||||
fname = a.file_name
|
||||
break
|
||||
if fname:
|
||||
fn = fname.split(".")
|
||||
ext = fn[-1] if len(fn) > 1 and len(fn[-1]) <= 5 else "bin"
|
||||
else:
|
||||
ext = "bin"
|
||||
return f"media/{channel}_{msg.id}.{ext}"
|
||||
|
||||
|
||||
async def download_media(client, msg, rel_path, source):
|
||||
"""Скачивает медиа поста в /opt/vesti/media/<basename> (MEDIA_DIR — уже /opt/vesti/media).
|
||||
|
||||
rel_path из post_media_path = 'media/<channel>_<id>.<ext>' (от корня проекта).
|
||||
Файл сохраняется в MEDIA_DIR/<basename> = /opt/vesti/media/<file> — ровно туда,
|
||||
куда указывает media_path в БД (media/<file>). БЕЗ вложенной media/media.
|
||||
"""
|
||||
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
|
||||
# определяем источник скачивания: photo/video/document/voice
|
||||
dl = await client.download_media(msg, file=full)
|
||||
if dl:
|
||||
return rel_path
|
||||
return None
|
||||
|
||||
|
||||
async def fetch_channel(client, source, state, src_ids=None, is_own_source=False):
|
||||
"""Читает новые посты канала, возвращает (list_of_raw, last_id).
|
||||
src_ids: set channel_id наших источников — для пропуска форвардов из своих каналов.
|
||||
is_own_source: источник own:true → медиа скачивается (контент пользователя)."""
|
||||
channel = source["channel"]
|
||||
try:
|
||||
entity = await client.get_entity(channel)
|
||||
except Exception as e:
|
||||
raise RuntimeError(f"get_entity({channel}) failed: {e}")
|
||||
|
||||
last_id = state.get("last_post_id") or 0
|
||||
posts = []
|
||||
limit = MAX_POSTS_PER_CHANNEL
|
||||
src_ids = src_ids or set()
|
||||
|
||||
# Берём историю с конца; фильтруем по id > last_id на лету.
|
||||
async for msg in client.iter_messages(entity, limit=limit * 3, reverse=False):
|
||||
if msg.id <= last_id:
|
||||
continue
|
||||
# --- форвард: нужен только исходный пост ---
|
||||
fw = getattr(msg, "fwd_from", None)
|
||||
fwd_channel_id = None
|
||||
fwd_post_id = None
|
||||
is_forward = False
|
||||
if fw is not None:
|
||||
from_id = getattr(fw, "from_id", None)
|
||||
if from_id is not None:
|
||||
# PeerChannel: у пересланных из канала
|
||||
if hasattr(from_id, "channel_id"):
|
||||
fwd_channel_id = from_id.channel_id
|
||||
fwd_post_id = getattr(fw, "channel_post", None)
|
||||
is_forward = True
|
||||
# PeerUser / другие — не обрабатываем как форвард канала
|
||||
# Пост-форвард из канала, который есть в наших источниках -> пропускаем (будет скачан в исходном)
|
||||
if is_forward and fwd_channel_id in src_ids:
|
||||
continue
|
||||
# Пост-форвард из чужого канала: сохраняем как есть, но БЕЗ скачивания медиа (не наше).
|
||||
# Для own-канала это требование тоже действует: чужой форвард = is_own=1 + fwd-поля, БЕЗ медиа.
|
||||
skip_media = is_forward
|
||||
|
||||
has_content = (
|
||||
msg.text is not None
|
||||
or msg.photo or msg.document or msg.video
|
||||
or msg.voice or msg.sticker
|
||||
or (msg.animation if hasattr(msg, "animation") else False)
|
||||
)
|
||||
if has_content:
|
||||
text = msg.text or ""
|
||||
ctype = detect_media_type(msg)
|
||||
media_rel = None
|
||||
if not skip_media:
|
||||
media_rel = post_media_path(msg, channel) if ctype != "text" else None
|
||||
# скачиваем медиа, если есть (в media/)
|
||||
if media_rel:
|
||||
try:
|
||||
saved = await download_media(client, msg, media_rel, source)
|
||||
media_rel = saved or media_rel
|
||||
except Exception as e:
|
||||
log.warning("media download failed %s: %s", channel, e)
|
||||
views = getattr(msg, "views", None) or 0
|
||||
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
|
||||
url = f"https://t.me/{channel}/{msg.id}"
|
||||
posts.append({
|
||||
"tg_channel": channel,
|
||||
"tg_post_id": msg.id,
|
||||
"url": url,
|
||||
"text": text,
|
||||
"content_type": ctype,
|
||||
"media_path": media_rel,
|
||||
"views": int(views or 0),
|
||||
"reactions": 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": fwd_channel_id,
|
||||
"fwd_from_post_id": fwd_post_id,
|
||||
"is_own": 1 if is_own_source else 0,
|
||||
})
|
||||
if len(posts) >= limit:
|
||||
break
|
||||
return posts, max([p["tg_post_id"] for p in posts] or [last_id])
|
||||
|
||||
|
||||
def store_posts(conn, source_id, posts, is_own_source=False):
|
||||
"""Добавляет посты в БД, возвращает количество НОВЫХ (не дублей).
|
||||
|
||||
own-семантика:
|
||||
- источник own:true → посты is_own=1 + is_own_canonical=1 (первичный экземпляр)
|
||||
- дубль того же текста (sha256) во внешнем канале → is_own=0 (упоминание),
|
||||
канонический (is_own_canonical=1) остаётся первичным
|
||||
"""
|
||||
new = 0
|
||||
for p in posts:
|
||||
sha = sha256_text(p["text"])
|
||||
# дедуп: уникальный (sha256) и (url)
|
||||
dup = conn.execute("SELECT id, media_path, fwd_from_channel_id, is_own, is_own_canonical FROM posts WHERE sha256=? OR url=?", (sha, p["url"])).fetchone()
|
||||
if dup:
|
||||
# дозаполняем типы/медиа, если появились
|
||||
if p.get("media_path") and not dup[1]:
|
||||
conn.execute(
|
||||
"UPDATE posts SET content_type=?, media_path=? WHERE id=?",
|
||||
(p.get("content_type", "text"), p["media_path"], dup[0]),
|
||||
)
|
||||
# дозаполняем fwd-поля, если пост-форвард и они не были записаны
|
||||
if p.get("fwd_from_channel_id") and not dup[2]:
|
||||
conn.execute(
|
||||
"UPDATE posts SET fwd_from_channel_id=?, fwd_from_post_id=? WHERE id=?",
|
||||
(p["fwd_from_channel_id"], p.get("fwd_from_post_id"), dup[0]),
|
||||
)
|
||||
# дубль того же контента:
|
||||
# - тот же sha (сам себя в том же канале) — не создаём
|
||||
# - внешний канал, но у нас уже есть канонический СВОЙ экземпляр → упоминание (is_own=0),
|
||||
# НО только если существующий пост канонический (иначе это просто дубль во внешних)
|
||||
if (p.get("is_own") and not dup[3] and dup[4]):
|
||||
# тот же контент в своём канале, а в БД уже есть канон → обновляем is_own у дубля как упоминание
|
||||
conn.execute(
|
||||
"UPDATE posts SET is_own=0, is_own_canonical=0 WHERE id=?",
|
||||
(dup[0],),
|
||||
)
|
||||
continue
|
||||
is_own = 1 if (is_own_source or p.get("is_own")) else 0
|
||||
is_own_canonical = 1 if is_own else 0
|
||||
conn.execute(
|
||||
"""INSERT INTO posts (sha256, source_id, tg_channel, tg_post_id, url, text, content_type,
|
||||
media_path, views, reactions, reactions_total, published_at, status,
|
||||
fwd_from_channel_id, fwd_from_post_id, is_own, is_own_canonical)
|
||||
VALUES (?,?,?,?,?,?,?,?,?,?,?,?, 'new',?,?,?,?)""",
|
||||
(sha, source_id, p["tg_channel"], p["tg_post_id"], p["url"], p["text"],
|
||||
p.get("content_type", "text"), p.get("media_path"),
|
||||
p["views"], p["reactions"], p["reactions_total"], p["published_at"],
|
||||
p.get("fwd_from_channel_id"), p.get("fwd_from_post_id"),
|
||||
is_own, is_own_canonical),
|
||||
)
|
||||
new += 1
|
||||
return new
|
||||
|
||||
|
||||
async def run_source(client, source, src_ids=None):
|
||||
"""Краулинг одного источника: инкрементально, с записью запуска в runs."""
|
||||
slug = source["slug"]
|
||||
is_own_source = int(source.get("own") or 0) == 1
|
||||
with db() as conn:
|
||||
state = conn.execute("SELECT * FROM tg_state WHERE channel_slug=?", (slug,)).fetchone()
|
||||
state = dict(state) if state else {}
|
||||
|
||||
run_id = start_run(source["id"], f"telegram:{slug}", "cron")
|
||||
t0 = time.monotonic()
|
||||
try:
|
||||
posts, last_id = await fetch_channel(client, source, state, src_ids, is_own_source)
|
||||
with db() as conn:
|
||||
new = store_posts(conn, source["id"], posts, is_own_source)
|
||||
# сохраняем состояние
|
||||
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, last_id),
|
||||
)
|
||||
conn.execute("UPDATE sources SET last_fetch=datetime('now'), status='alive', last_error=NULL WHERE id=?", (source["id"],))
|
||||
duration = int((time.monotonic() - t0) * 1000)
|
||||
finish_run(run_id, "ok", len(posts), new, None, duration)
|
||||
log.info("OK %s: fetched=%d new=%d last_id=%d (%.1fs)", slug, len(posts), new, last_id, duration / 1000)
|
||||
except Exception as e:
|
||||
duration = int((time.monotonic() - t0) * 1000)
|
||||
finish_run(run_id, "error", 0, 0, str(e)[:500], duration)
|
||||
log.error("ERROR %s: %s", slug, e)
|
||||
with db() as conn:
|
||||
conn.execute("UPDATE sources SET status='dead', last_error=? WHERE id=?", (str(e)[:500], source["id"]))
|
||||
raise
|
||||
|
||||
|
||||
async def resolve_source_ids(client, sources):
|
||||
"""Возвращает set channel_id всех наших telegram-источников (для дедупа форвардов)."""
|
||||
ids = set()
|
||||
for s in sources:
|
||||
try:
|
||||
ent = await client.get_entity(s["channel"])
|
||||
ids.add(getattr(ent, "id", None))
|
||||
except Exception as e:
|
||||
log.warning("resolve id %s failed: %s", s["channel"], e)
|
||||
return ids
|
||||
|
||||
|
||||
async def main():
|
||||
ap = argparse.ArgumentParser(description="VESTI Telegram crawler (per-source)")
|
||||
ap.add_argument("--source", help="slug источника (например linuxklub)")
|
||||
ap.add_argument("--direction", help="направление (например linux)")
|
||||
ap.add_argument("--all", action="store_true", help="все включённые telegram-источники")
|
||||
args = ap.parse_args()
|
||||
|
||||
ensure_dirs()
|
||||
TG_SESSION_DIR.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# синк реестра в БД
|
||||
data = load_sources_yaml()
|
||||
sync_sources_to_db(data)
|
||||
|
||||
if args.source:
|
||||
src = get_source(args.source)
|
||||
if not src:
|
||||
log.error("source %s not found", args.source)
|
||||
sys.exit(1)
|
||||
sources = [src]
|
||||
elif args.direction:
|
||||
sources = get_enabled_sources(crawler="telegram", direction=args.direction)
|
||||
else:
|
||||
sources = get_enabled_sources(crawler="telegram")
|
||||
|
||||
if not sources:
|
||||
log.info("no sources to crawl")
|
||||
return
|
||||
|
||||
client = make_client(TG_SESSION_DIR)
|
||||
try:
|
||||
await client.connect()
|
||||
if not await client.is_user_authorized():
|
||||
# Требуется код подтверждения — интерактивно (редко; сессия сохраняется)
|
||||
log.warning("Session not authorized. Проверьте, что сессия /opt/vesti/telegram/vesti.session существует.")
|
||||
log.warning("Для первой авторизации: python -m crawler.auth_telegram")
|
||||
await client.disconnect()
|
||||
sys.exit(2)
|
||||
src_ids = await resolve_source_ids(client, sources)
|
||||
log.info("source channel ids: %s", sorted(src_ids))
|
||||
for src in sources:
|
||||
await run_source(client, src, src_ids)
|
||||
finally:
|
||||
await client.disconnect()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user