# 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/_. (без скачивания здесь).""" 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/ (MEDIA_DIR — уже /opt/vesti/media). rel_path из post_media_path = 'media/_.' (от корня проекта). Файл сохраняется в MEDIA_DIR/ = /opt/vesti/media/ — ровно туда, куда указывает media_path в БД (media/). БЕЗ вложенной 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())