diff --git a/src/telegram_exporter/__init__.py b/src/telegram_exporter/__init__.py new file mode 100644 index 0000000..bc5721d --- /dev/null +++ b/src/telegram_exporter/__init__.py @@ -0,0 +1,3 @@ +"""Export all media from a Telegram group to local disk, grouped by post.""" + +__version__ = "0.1.0" diff --git a/src/telegram_exporter/cli.py b/src/telegram_exporter/cli.py new file mode 100644 index 0000000..1571486 --- /dev/null +++ b/src/telegram_exporter/cli.py @@ -0,0 +1,253 @@ +"""Command line entry point: process rails, argument parsing, dispatch. + +The rails come first and in this order for reasons that are not stylistic: +umask before anything can create a file, logging policy before anything can log a +credential, the git-safety assertion before anything writes into a repo, and the +lock before the destructive .part sweep. +""" + +from __future__ import annotations + +import argparse +import asyncio +import fcntl +import logging +import os +import socket +import sys +from contextlib import contextmanager +from pathlib import Path + +from . import paths +from .downloader import ( + FATAL_ERRNO, + Config, + Fetcher, + run_download, + sweep_part_files, + write_title, +) +from .estimate import EstimateSink, resume_from, run_estimate +from .session import ( + Abort, + assert_not_committable, + connected_client, + default_session_path, + load_api_credentials, + now_iso, + resolve_entity, + with_session_suffix, +) +from .sidecar import Sidecar +from .state import State +from .traversal import TraversalError, build_filters, iter_posts, parse_size + +log = logging.getLogger("telegram_exporter.cli") + +LOCK_NAME = ".export.lock" +DEFAULT_MIN_FREE = "2GiB" + + +def configure_logging(verbose: bool) -> None: + """--verbose means *our* logger only. + + telethon's request logging carries api_id, phone_number and + phone_code_hash, so its floor is unconditional and there is deliberately no + --debug-telethon flag. The root level is never touched. + """ + logging.basicConfig(level=logging.INFO, + format="%(asctime)s %(levelname)s %(message)s") + logging.getLogger("telegram_exporter").setLevel( + logging.DEBUG if verbose else logging.INFO) + logging.getLogger("telethon").setLevel(logging.WARNING) + logging.getLogger("asyncio").setLevel(logging.WARNING) + + +@contextmanager +def exclusive(path: Path): + """Single-instance lock, per export root and per session. + + flock releases on process death, kill -9 included, so there is no stale-lock + cleanup to get wrong. Contention fails fast with exit 3 - no wait, no + --force. The message names pid, host and start time so that "did it hang?" + gets an answer instead of a second process. + """ + path.parent.mkdir(parents=True, exist_ok=True) + f = open(path, "a+") + try: + fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + f.seek(0) + holder = f.read().strip() or "?" + f.close() + raise Abort( + f"busy: another tg-export holds {path} [{holder}]\n" + f"if you meant to export a second group at the same time, give it its " + f"own --session too: one SQLite session cannot serve two clients.", 3 + ) from None + f.seek(0) + f.truncate() + f.write(f"pid={os.getpid()} host={socket.gethostname()} started={now_iso()}\n") + f.flush() + try: + yield + finally: + f.close() + + +def _non_negative(text: str) -> int: + """A negative --limit is the one input where the estimator and the real run + would disagree (one stops before counting, the other after one post), so it + is rejected at the boundary rather than handled twice.""" + value = int(text) + if value < 0: + raise argparse.ArgumentTypeError(f"must be 0 or greater, got {value}") + return value + + +def build_parser() -> argparse.ArgumentParser: + p = argparse.ArgumentParser( + prog="tg-export", + description="Export all media from a Telegram group, grouped by post.", + epilog="credentials: export TG_API_ID and TG_API_HASH " + "(create an app at https://my.telegram.org). " + "The login code and any 2FA password are prompted for, never read " + "from the environment.") + p.add_argument("--group", required=True, + help="numeric id, @username, or t.me link") + p.add_argument("--out", required=True, metavar="DIR", + help="parent directory; the export goes in /g") + p.add_argument("--dry-run", action="store_true", + help="report what a real run would fetch, write nothing") + p.add_argument("--since", metavar="YYYY-MM-DD", help="messages on or after (UTC)") + p.add_argument("--until", metavar="YYYY-MM-DD", help="messages on or before (UTC)") + p.add_argument("--types", metavar="LIST", + help="comma-separated: photo,video,document,audio,voice") + p.add_argument("--max-size", metavar="SIZE", + help="skip media larger than this, e.g. 100MB") + p.add_argument("--include-text", action="store_true", + help="also record text-only messages in messages.jsonl with an " + "empty files[]; part of the stored filter set") + p.add_argument("--limit", type=_non_negative, metavar="N", + help="stop after N posts that contain files (0 = connect only)") + p.add_argument("--session", metavar="PATH", + help=f"session file (default: {default_session_path()})") + p.add_argument("--reset-state", action="store_true", + help="discard the cursor and rotate messages.jsonl; keeps " + "downloaded files") + p.add_argument("--min-free", default=DEFAULT_MIN_FREE, metavar="SIZE", + help=f"disk kept free (default {DEFAULT_MIN_FREE})") + p.add_argument("--max-flood-wait", type=int, metavar="SECONDS", + help="exit resumably (code 6) instead of sleeping through a " + "flood wait longer than this; unset means sleep as long " + "as Telegram demands") + p.add_argument("--verbose", action="store_true", + help="debug logging for this tool only, never for telethon") + return p + + +async def _run(args) -> int: + api_id, api_hash = load_api_credentials() + filters = build_filters(types=args.types, since=args.since, until=args.until, + max_size=args.max_size, include_text=args.include_text) + min_free = parse_size(args.min_free) + session_path = with_session_suffix( + Path(args.session).expanduser() if args.session else default_session_path()) + + async with connected_client(session_path, api_id, api_hash) as client: + entity, peer_id = await resolve_entity(client, args.group) + title = (getattr(entity, "title", None) + or getattr(entity, "username", None) or str(peer_id)) + root = paths.export_root(Path(args.out).expanduser(), peer_id) + assert_not_committable(root, is_dir=True) + log.info("group %r (id %s) -> %s", title, peer_id, root) + + if args.dry_run: + return await _dry_run(client, entity, root=root, title=title, + peer_id=peer_id, filters=filters, + min_free=min_free, args=args) + return await _real_run(client, entity, root=root, title=title, + peer_id=peer_id, filters=filters, + min_free=min_free, session_path=session_path, + args=args) + + +async def _dry_run(client, entity, *, root: Path, title: str, peer_id: int, + filters, min_free: int, args) -> int: + """Takes no lock, creates no directories, mutates no state - so it can run + while a real export is in progress.""" + start = resume_from(root, filters) + posts = iter_posts(client, entity, after_id=start, filters=filters, + max_flood_wait=args.max_flood_wait) + return await run_estimate(posts, sink=EstimateSink(), root=root, title=title, + peer_id=peer_id, filters=filters, start_id=start, + min_free=min_free, limit=args.limit) + + +async def _real_run(client, entity, *, root: Path, title: str, peer_id: int, + filters, min_free: int, session_path: Path, args) -> int: + root.mkdir(parents=True, exist_ok=True) + session_lock = session_path.with_name(session_path.name + ".lock") + + # Both locks before the .part sweep: the sweep is the destructive operation, + # so locking after it would lock after the damage. Per-root rather than + # global, because exporting two different groups concurrently is legitimate. + with exclusive(root / LOCK_NAME), exclusive(session_lock): + sweep_part_files(root) + write_title(root, title) + + # State first: it refuses a filter or chat mismatch (exit 7) before the + # sidecar is touched. + state = State.open(root, chat_id=peer_id, chat_title=title, + filters=filters, reset=args.reset_state) + sidecar = Sidecar(root) + if args.reset_state: + sidecar.rotate() + + with sidecar: + posts = iter_posts(client, entity, after_id=state.cursor_id, + filters=filters, max_flood_wait=args.max_flood_wait) + totals = await run_download( + posts, + fetcher=Fetcher(client, entity, args.max_flood_wait), + state=state, sidecar=sidecar, root=root, + cfg=Config(min_free=min_free, disk_path=root, + max_flood_wait=args.max_flood_wait), + limit=args.limit) + log.info("done: %s", totals.summary()) + return 0 + + +def main() -> int: + os.umask(0o077) # FIRST statement: covers the session file and every sqlite + # sibling, the export tree, messages.jsonl, .part files, + # the state file and the lock - one line, no per-file chmod. + args = build_parser().parse_args() + configure_logging(args.verbose) + try: + return asyncio.run(_run(args)) + except Abort as a: + log.error("%s", a.reason) + return a.code + except TraversalError as e: + # An album split or a non-ascending sweep. Loud on purpose: the export on + # disk cannot be trusted, and the cursor is where the last good post was. + log.error("%s", e) + return 1 + except OSError as e: + # The download path maps these itself; this catches the same failures + # arriving from the sidecar or state writes, so disk exhaustion exits 3 + # there too instead of printing a traceback. + code = 3 if e.errno in FATAL_ERRNO else 1 + log.error("filesystem: %s", e) + return code + except ValueError as e: # bad --types/--since/--max-size value + log.error("%s", e) + return 1 + except KeyboardInterrupt: + log.info("interrupted - re-run the same command to resume") + return 130 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/src/telegram_exporter/downloader.py b/src/telegram_exporter/downloader.py new file mode 100644 index 0000000..1d4a788 --- /dev/null +++ b/src/telegram_exporter/downloader.py @@ -0,0 +1,443 @@ +"""The download path: fetch sequentially, write durably, dedupe, record. + +Governing principle: **the filesystem is the state; the state file is a hint; the +sidecar is the log.** Everything here follows from it. + +Telethon has no byte-level resume - download_media writes from offset 0 every +call - so resumability lives at *message* granularity, and completeness is proven +by arrival at the target path: + + per file: write -> flush -> fsync(fd) -> os.replace(tmp, target) + per post: fsync(post_dir) -> sidecar.append + fsync -> state.commit + +With fsync before os.replace, a file at the target path can only have arrived +fully downloaded. That is what licenses existence-based dedupe below. +""" + +from __future__ import annotations + +import asyncio +import errno +import logging +import os +import shutil +import time +from dataclasses import dataclass +from pathlib import Path + +from telethon.errors import ( + AuthKeyDuplicatedError, + AuthKeyUnregisteredError, + ChannelInvalidError, + ChannelPrivateError, + ChatForbiddenError, + FileIdInvalidError, + FileReferenceExpiredError, + FilerefUpgradeNeededError, + MediaEmptyError, + ServerError, + SessionRevokedError, + UserDeactivatedBanError, +) +from telethon.errors import common as telethon_common + +from . import paths +from .media import file_ext, has_media, media_kind, media_size +from .session import Abort, with_flood_retry +from .state import fsync_dir + +log = logging.getLogger("telegram_exporter.downloader") + +TRANSIENT_ERRNO = {errno.EAGAIN, errno.EINTR, errno.EIO, errno.ETIMEDOUT, + errno.ECONNRESET} +FATAL_ERRNO = {errno.ENOSPC, errno.EDQUOT, errno.EROFS, errno.EACCES} + +MAX_ATTEMPTS = 3 +PART_SUFFIX = ".part" +TITLE_NAME = "title.txt" + +# Commit the cursor at least this often even across a long filtered region, so a +# 100k-message sweep does not become 100k state writes. +COMMIT_INTERVAL_S = 5.0 + +DOWNLOADED, SKIPPED, FAILED = "DOWNLOADED", "SKIPPED", "FAILED" + +# Network trouble is retried; the unavailable-media errors are skipped; the +# access and auth errors end the run. FloodWaitError never appears here - it is +# owned by with_flood_retry, because a retry during an active wait is counted as +# a fresh violation and escalates the next wait. +# +# ServerError rather than only RpcCallFailError: the latter is one leaf of it, so +# a plain Telegram -500 would otherwise escape. The telethon.errors.common set +# derives from Exception and BufferError, not OSError - a corrupt packet or a +# checksum failure on a flaky link is transient, but nothing below OSError would +# have caught it, and an unhandled one ends a 20-hour unattended run with a +# traceback. +_NETWORK_ERRORS = (ConnectionError, TimeoutError, asyncio.TimeoutError, + asyncio.IncompleteReadError, ServerError, + telethon_common.InvalidBufferError, + telethon_common.InvalidChecksumError, + telethon_common.TypeNotFoundError, + telethon_common.SecurityError, + telethon_common.BadMessageError) +_AUTH_ERRORS = (AuthKeyUnregisteredError, AuthKeyDuplicatedError, + SessionRevokedError, UserDeactivatedBanError) +_ACCESS_ERRORS = (ChannelPrivateError, ChatForbiddenError, ChannelInvalidError) +_UNAVAILABLE_ERRORS = (MediaEmptyError, FileIdInvalidError, FilerefUpgradeNeededError) + + +def backoff(attempt: int) -> float: + return 5.0 * 2 ** (attempt - 1) # 5, 10, 20 s + + +@dataclass(frozen=True) +class Config: + """Everything download_one needs from the CLI, so it takes no globals.""" + + min_free: int # bytes held in reserve, default 2 GiB + disk_path: Path # filesystem to measure - the export root + # No max_flood_wait here: the ceiling belongs to the two flood primitives, + # which Fetcher owns. A copy on this object would read as though check_disk + # and download_one honored it. + + +@dataclass +class Result: + """What happened to one message's media, as the sidecar records it.""" + + message_id: int + status: str + path: Path | None = None + size: int | None = None + declared_size: int | None = None + name: str | None = None + mime: str | None = None + kind: str | None = None + error: str | None = None + + def as_record(self, root: Path) -> dict: + record = { + "path": str(self.path.relative_to(root)), + "size": self.size, + "name": self.name, + "mime": self.mime, + "kind": self.kind, + } + # declared_size appears only when it differs, which incidentally + # accumulates empirical data on photo-variant accuracy. + if self.declared_size is not None and self.declared_size != self.size: + record["declared_size"] = self.declared_size + return record + + +@dataclass +class Totals: + """Per-run, in-memory counters. Deliberately not persisted: counters in the + state file double-count on every resume, because a re-processed album + re-increments them.""" + + posts: int = 0 + downloaded: int = 0 + skipped: int = 0 + failed: int = 0 + bytes_downloaded: int = 0 + + def summary(self) -> str: + line = (f"{self.posts} posts, {self.downloaded} downloaded " + f"({human_bytes(self.bytes_downloaded)}), {self.skipped} already present") + if self.failed: + line += f", {self.failed} files failed (see messages.jsonl)" + return line + + +def human_bytes(n: int) -> str: + step = 1024.0 + value = float(n) + for unit in ("B", "KiB", "MiB", "GiB", "TiB"): + if value < step or unit == "TiB": + return f"{value:,.1f} {unit}" if unit != "B" else f"{int(value):,} B" + value /= step + raise AssertionError("unreachable: the loop returns at TiB") + + +# --------------------------------------------------------------------------- # +# Fetcher - the only part that touches the network, so tests can replace it +# --------------------------------------------------------------------------- # + +class Fetcher: + """Wraps the two client calls download_one needs, each behind flood retry.""" + + def __init__(self, client, entity, max_flood_wait: int | None = None) -> None: + self.client = client + self.entity = entity + self.max_flood_wait = max_flood_wait + + async def fetch(self, msg, fh) -> None: + await with_flood_retry(lambda: self.client.download_media(msg, file=fh), + max_wait_s=self.max_flood_wait) + + async def refresh(self, message_id: int): + """Re-fetch a message to renew an expired file reference.""" + return await with_flood_retry( + lambda: self.client.get_messages(self.entity, ids=message_id), + max_wait_s=self.max_flood_wait) + + +# --------------------------------------------------------------------------- # +# Per-file decisions +# --------------------------------------------------------------------------- # + +def check_disk(msg, cfg: Config) -> None: + """Refuse before the first byte, comparing against *this file's* size. + + One statvfs is microseconds, unmeasurable against a multi-second download at + concurrency 1. Polling every N files would optimize a cost that does not + exist while opening a blind spot: a single Telegram file reaches 2-4 GB, so + ~100 GB could be written between two polls. An unknown size counts as 0 and + is absorbed by the reserve, which is why the reserve defaults to 2 GiB - one + entire unknown-size file at the non-premium upload cap. + """ + need = (media_size(msg) or 0) + cfg.min_free + free = shutil.disk_usage(cfg.disk_path).free + if free < need: + raise Abort(f"insufficient disk on {cfg.disk_path}: need " + f"{human_bytes(need)} (file + {human_bytes(cfg.min_free)} " + f"reserve), have {human_bytes(free)}", 3) + + +def validate(tmp: Path, msg) -> bool: + """Advisory only - never raises. True means accept. + + Documents, video and audio carry an exact document.size, so a mismatch there + is real truncation. Photos are the only approximate case: .size describes the + size *variant* Telethon selected. An exact-equality gate would delete every + completed photo on resume and then raise, which is why this is advisory and + why dedupe below tests existence rather than size. + """ + actual, declared = tmp.stat().st_size, media_size(msg) + if actual == 0: + # Nothing arrived. download_media returns without writing when the media + # resolves to an empty variant, and the photo excuse below would + # otherwise rename a 0-byte file to the target - where the dedupe rule + # ("zero bytes is never a complete download") could never see it again, + # because a completed export does not revisit existing targets. + log.warning("message %s produced 0 bytes - retrying", msg.id) + return False + if declared is None or actual == declared: + return True + if media_kind(msg) == "photo": + log.debug("photo %s declared %s, got %s - variant size, accepted", + msg.id, declared, actual) + return True + log.warning("size mismatch on message %s: declared %s, got %s - retrying", + msg.id, declared, actual) + return False + + +def sweep_part_files(root: Path) -> int: + """Remove stray .part files left by a killed run. + + Destructive, so the caller must already hold the export lock: sweeping first + and locking second would lock after the damage. + """ + removed = 0 + for stray in root.rglob(f"*{PART_SUFFIX}"): + stray.unlink(missing_ok=True) + removed += 1 + if removed: + log.info("swept %d stray %s file(s)", removed, PART_SUFFIX) + return removed + + +async def download_one(fetcher, post_dir_path: Path, msg, *, cfg: Config) -> Result: + """Download one message's media into post_dir_path. + + Every exit is one of: DOWNLOADED, SKIPPED, FAILED (recorded, run continues), + or Abort (run ends with a correct cursor). There is no path that lets an + unhandled exception kill an unattended 20-hour export. + """ + try: + target = paths.safe_join(post_dir_path, paths.derive_filename(msg)) + except (ValueError, TypeError) as e: + # sanitize should make this unreachable, so reaching it means sanitize has + # a hole worth finding - hence ERROR, not a silent repair. But one + # hostile filename must not be able to end the run: without the cursor + # advancing past this post, every later run would die on the same + # message. Recorded in errors[] and skipped, like unavailable media. + log.error("cannot build a safe path for message %s: %s - skipping", msg.id, e) + return _result(msg, FAILED, error=f"unsafe filename: {e}") + + if target.exists(): + if target.stat().st_size > 0: + return _result(msg, SKIPPED, target) + # Zero bytes means ENOSPC junk, never a complete download - re-fetch. + target.unlink() + + check_disk(msg, cfg) + post_dir_path.mkdir(parents=True, exist_ok=True) + tmp = target.with_name(target.name + PART_SUFFIX) + refreshed = False + + for attempt in range(1, MAX_ATTEMPTS + 1): + try: + with open(tmp, "wb") as fh: + await fetcher.fetch(msg, fh) + fh.flush() + os.fsync(fh.fileno()) + if not validate(tmp, msg): + tmp.unlink(missing_ok=True) + await asyncio.sleep(backoff(attempt)) + continue + os.replace(tmp, target) + return _result(msg, DOWNLOADED, target) + + # ---- SKIP: this message is unavailable; the run continues ---------- # + except FileReferenceExpiredError: + # References expire in hours, so a 20-hour run *will* hit this. One + # refresh is the highest-value entry in this table. + tmp.unlink(missing_ok=True) + if refreshed: + return _result(msg, FAILED, error="file reference expired twice") + fresh = await fetcher.refresh(msg.id) + refreshed = True + if fresh is None or not has_media(fresh): + return _result(msg, FAILED, error="message deleted") + msg = fresh + continue + except _UNAVAILABLE_ERRORS as e: + tmp.unlink(missing_ok=True) + return _result(msg, FAILED, + error=f"media unavailable or expired: {type(e).__name__}") + + # ---- ABORT: the run cannot continue ------------------------------- # + except _AUTH_ERRORS: + tmp.unlink(missing_ok=True) + raise Abort("session invalidated - delete the session file and re-login", 4) + except _ACCESS_ERRORS: + tmp.unlink(missing_ok=True) + raise Abort("lost access to the group", 5) + + # ---- RETRY: network trouble --------------------------------------- # + # Before the OSError clause on purpose: ConnectionError *is* an OSError, + # and a ConnectionResetError raised by asyncio can carry errno=None, + # which the fatal/transient errno test below would fall through to + # "unexpected OS error" and abort a healthy run. + except _NETWORK_ERRORS as e: + tmp.unlink(missing_ok=True) + log.warning("network error on message %s (attempt %d/%d): %s", + msg.id, attempt, MAX_ATTEMPTS, e) + await asyncio.sleep(backoff(attempt)) + continue + + # ---- filesystem: fatal checked before transient ------------------- # + except OSError as e: + tmp.unlink(missing_ok=True) + if e.errno in FATAL_ERRNO: + raise Abort(f"filesystem: {e}", 3) from None + if e.errno in TRANSIENT_ERRNO: + await asyncio.sleep(backoff(attempt)) + continue + raise Abort(f"unexpected OS error: {e}", 1) from None + + return _result(msg, FAILED, error=f"exhausted {MAX_ATTEMPTS} attempts") + + +def _result(msg, status: str, target: Path | None = None, + error: str | None = None) -> Result: + f = getattr(msg, "file", None) + return Result( + message_id=msg.id, + status=status, + path=target, + size=target.stat().st_size if target is not None else None, + declared_size=media_size(msg), + name=getattr(f, "name", None) or f"{media_kind(msg)}{file_ext(msg)}", + mime=getattr(f, "mime_type", None), + kind=media_kind(msg), + error=error, + ) + + +# --------------------------------------------------------------------------- # +# The sink loop - Invariant 2 by construction +# --------------------------------------------------------------------------- # + +async def run_download(posts, *, fetcher, state, sidecar, root: Path, + cfg: Config, limit: int | None = None) -> Totals: + """Drive the shared traversal, downloading as it goes. + + state.commit is reached only after a post is fully handled, and there is no + commit on any error path - so the cursor can never point past in-flight work. + "Commit the cursor and exit" is deliberately not a concept here; it was the + mechanism by which the disk guard could have violated Invariant 2. + + Abort and KeyboardInterrupt propagate to cli.main, which owns exit codes. + The cursor is already correct when they do. + """ + totals = Totals() + last_commit = time.monotonic() + last_complete = state.cursor_id + if limit == 0: + return totals # connectivity check, no work requested + + async for post in posts: + media = [m for m in post.messages if has_media(m)] + results: list[Result] = [] + + if media: + target_dir = paths.post_dir(root, post.post_id) + for msg in media: + result = await download_one(fetcher, target_dir, msg, cfg=cfg) + results.append(result) + _tally(totals, result, post.post_id) + if target_dir.exists(): + fsync_dir(target_dir) + + # A post may be empty (everything filtered out) or text-only under + # --include-text; both still advance the cursor, and neither is counted. + if post.messages: + sidecar.append(post, results) + sidecar.fsync() + if media: + totals.posts += 1 + + # The post is now fully handled: files on disk, sidecar durable. + last_complete = post.max_message_id + wrote = any(r.status in (DOWNLOADED, FAILED) for r in results) + if wrote or (time.monotonic() - last_commit) > COMMIT_INTERVAL_S: + state.commit(last_complete) + last_commit = time.monotonic() + + if limit is not None and totals.posts >= limit: + log.info("--limit %d reached", limit) + state.commit(last_complete) + return totals + + # Flush the throttled tail before claiming the export finished. Without this, + # a run ending on a stretch of filtered or text-only posts shorter than the + # commit interval would record completed_at against a cursor still pointing + # at the last post that happened to contain a file - and the next run would + # re-sweep that whole tail to download nothing. + state.commit(last_complete) + state.mark_completed() + return totals + + +def _tally(totals: Totals, result: Result, post_id: int) -> None: + if result.status == DOWNLOADED: + totals.downloaded += 1 + totals.bytes_downloaded += result.size or 0 + log.info("%s/%s %s", post_id, result.path.name, human_bytes(result.size or 0)) + elif result.status == SKIPPED: + totals.skipped += 1 + log.debug("%s/%s already present", post_id, result.path.name) + else: + totals.failed += 1 + log.warning("%s message %s: %s", post_id, result.message_id, result.error) + + +def write_title(root: Path, title: str) -> None: + """The human-readable group name lives here, as data - never as a path + component. Written once; a later rename updates it without moving anything.""" + root.mkdir(parents=True, exist_ok=True) + (root / TITLE_NAME).write_text((title or "") + "\n", encoding="utf-8") diff --git a/src/telegram_exporter/estimate.py b/src/telegram_exporter/estimate.py new file mode 100644 index 0000000..3b8f558 --- /dev/null +++ b/src/telegram_exporter/estimate.py @@ -0,0 +1,179 @@ +"""--dry-run: sweep, download nothing, report what a real run would fetch *from +where it would actually start*. + +run_estimate drives the same iter_posts generator with the same Filters as +run_download. That shared generator is the anti-drift property: an estimator with +its own traversal would report numbers that do not match reality. There is no +Sink Protocol - two implementations with one call site and no shared behavior is +indirection describing a for loop - so both run modes own a plain `async for`, +and the duplication is two lines. +""" + +from __future__ import annotations + +import shutil +from dataclasses import dataclass, field +from pathlib import Path + +from .downloader import human_bytes +from .media import has_media, media_kind, media_size +from .state import State +from .traversal import Filters + +# Photo .size describes the size *variant* Telethon selected, so photo totals are +# approximate. Documents, video and audio carry an exact document.size. +_LABELS = {"photo": "photos", "video": "videos", "document": "documents", + "audio": "audio files", "voice": "voice notes"} + + +@dataclass +class EstimateSink: + """Counts posts *with files*, files, and bytes by kind. + + Empty posts (a fully filtered group, which iter_posts yields to keep the + cursor advancing) and text-only posts under --include-text are not counted, + on both sides of any comparison - so the dry-run and real-run reconciliation + in phase 7 compares like with like. + """ + + posts: int = 0 + files: int = 0 + by_kind: dict[str, list[int]] = field(default_factory=dict) # kind -> [count, bytes] + unknown_size_files: int = 0 + + def add(self, post) -> None: + media = [m for m in post.messages if has_media(m)] + if not media: + return + self.posts += 1 + for msg in media: + self.files += 1 + kind = media_kind(msg) # has_media guarantees this is set + row = self.by_kind.setdefault(kind, [0, 0]) + row[0] += 1 + size = media_size(msg) + if size is None: + # Never coerced to 0: an unbounded estimate error must be visible. + self.unknown_size_files += 1 + else: + row[1] += size + + @property + def total_bytes(self) -> int: + return sum(row[1] for row in self.by_kind.values()) + + +def _measurable_dir(p: Path) -> Path: + """Nearest existing ancestor - dry-run must not create the export root just + to ask how much space is left.""" + for candidate in [p, *p.parents]: + if candidate.is_dir(): + return candidate + return Path.cwd() + + +def describe_filters(filters: Filters) -> str: + # `is not None` throughout: --max-size 0 is an active filter that drops + # everything, and printing it as "(none)" would misdescribe the run. + parts = [] + if filters.types is not None: + parts.append("types=" + ",".join(sorted(filters.types))) + if filters.since is not None: + parts.append(f"since={filters.since.date()}") + if filters.until is not None: + parts.append(f"until={filters.until.date()}") + if filters.max_size is not None: + parts.append(f"max-size={human_bytes(filters.max_size)}") + if filters.include_text: + parts.append("include-text") + return " ".join(parts) or "(none)" + + +def stored_filters_differ(root: Path, filters: Filters) -> bool: + """Whether a real run would refuse (exit 7) rather than resume.""" + data = State.read(root) + return bool(data) and data.get("filters") != filters.to_state() + + +def resume_from(root: Path, filters: Filters) -> int: + """Start where the real run would. + + Without this, `tg-export --dry-run && tg-export` fails *closed* on every + resumed run: a 38 GB group with 30 GB already downloaded and 11 GB free would + re-measure the whole 38 GB, exit 2, and refuse to fetch the 8 GB that fits + comfortably - leaving the user no option but to bypass the guard they were + told to trust. + + Deliberately not scanning the export dir to subtract bytes already on disk: + with a cursor-based sweep everything before the cursor is already excluded, + and the residue is a handful of failed files and one partial post. + """ + data = State.read(root) + if not data: + return 0 + if data.get("filters") != filters.to_state(): + return 0 # different filters: the real run would refuse + return int(data.get("cursor_id", 0)) + + +def _mode_line(root: Path, filters: Filters, start_id: int) -> str: + """Say which run this describes. + + When the stored filters differ, the numbers below are honest for a *fresh* + export but misleading as "remaining work", because a real run would refuse + with exit 7 until --reset-state. Saying so beats printing a total the next + command cannot act on. + """ + if start_id: + return f"resume from message {start_id}" + if stored_filters_differ(root, filters): + return ("full history - stored filters differ, so a real run would refuse " + "(exit 7) until --reset-state") + return "full history (no prior state)" + + +async def run_estimate(posts, *, sink: EstimateSink, root: Path, + title: str, peer_id: int, filters: Filters, + start_id: int, min_free: int, + limit: int | None = None) -> int: + """Consume the shared traversal, print the report, return the exit code. + + The verdict is advisory. Phase 5's per-file check_disk is the sole + enforcement mechanism, so a SHORT verdict does not mean the run is unsafe - + it means the run will stop partway with a clean, resumable cursor. + """ + if limit != 0: + async for post in posts: + sink.add(post) + # Same stop rule as run_download, so `--dry-run --limit N` and + # `--limit N` from the same start describe the same work. + if limit is not None and sink.posts >= limit: + break + + free = shutil.disk_usage(_measurable_dir(root)).free + total = sink.total_bytes + need = total + min_free + + print(f"Group: {title or '(untitled)'} (id {peer_id})") + print(f"Filters: {describe_filters(filters)}") + print(f"Mode: {_mode_line(root, filters, start_id)}") + print() + print(f"Remaining: {sink.posts:,} posts / {sink.files:,} files") + for kind, (count, size) in sorted(sink.by_kind.items()): + note = " (approximate - size variant)" if kind == "photo" else "" + print(f" {_LABELS.get(kind, kind):<14}{count:>6,} files " + f"{human_bytes(size):>12}{note}") + if sink.unknown_size_files: + print(f" {'unknown size':<14}{sink.unknown_size_files:>6,} files " + f"{'(not counted)':>12}") + print() + print(f"Total: {human_bytes(total)}") + print(f"Free: {human_bytes(free)} Reserve: {human_bytes(min_free)}") + + if free < need: + print(f"Verdict: SHORT by {human_bytes(need - free)}") + print(f" advisory - the run would stop partway with a " + f"resumable cursor, not corrupt anything") + return 2 + print(f"Verdict: FITS ({human_bytes(free - need)} to spare)") + return 0 diff --git a/src/telegram_exporter/media.py b/src/telegram_exporter/media.py new file mode 100644 index 0000000..d448472 --- /dev/null +++ b/src/telegram_exporter/media.py @@ -0,0 +1,106 @@ +"""Media classification and size, in one place because three consumers need it. + +`--types` filtering, filename derivation, the download validator, the disk guard +and the estimator all have to agree on "what kind of media is this" and "how big +does the server say it is". Two copies of that answer would drift, and a drift +between the estimator and the downloader is exactly the lie the shared-traversal +design exists to prevent. + +Duck-typed on purpose: every accessor uses getattr with a default, so unit tests +can build a stub message with just the attributes a case needs, offline. +""" + +from __future__ import annotations + +from pathlib import PurePosixPath + +# The kinds `--types` accepts. "video" absorbs video notes and gifs; stickers +# fall through to "document" unless they carry a video attribute. +KINDS = ("photo", "video", "document", "audio", "voice") + + +def media_kind(msg) -> str | None: + """The media kind, or None when the message carries no downloadable file. + + A link preview is deliberately not media. Telethon's `Message.photo` and + `Message.document` also return the *web page preview's* photo or document, + so classifying with those properties alone would make every message + containing a URL look like an attachment - inflating counts and downloading + thumbnails nobody asked for. Checking web_preview first excludes them. + """ + if getattr(msg, "web_preview", None) is not None: + return None + if getattr(msg, "photo", None) is not None: + return "photo" + if getattr(msg, "voice", None) is not None: + return "voice" # a voice note is an audio document + if (getattr(msg, "video", None) is not None + or getattr(msg, "video_note", None) is not None + or getattr(msg, "gif", None) is not None): + return "video" + if getattr(msg, "audio", None) is not None: + return "audio" + if getattr(msg, "document", None) is not None: + return "document" + return None # text, poll, geo, contact, dice... + + +def has_media(msg) -> bool: + return media_kind(msg) is not None + + +def media_size(msg) -> int | None: + """Server-declared size in bytes, or None when the server did not say. + + Single source of truth. None is not zero, and each consumer has its own + stated rule for it: + + --max-size filter include, warn once (unknown-size media is + pathological and near-always small; a silent drop + would violate the no-silent-gaps ethos) + download validator accept - there is nothing to check against + estimator counted on a separate "unknown size" line, never + coerced to 0, because an unbounded estimate error + has to be visible + disk guard treated as 0, absorbed by the --min-free reserve + + Comparing `None <= limit` unguarded raises TypeError inside the sweep and + kills a 20-hour export. That is the failure this helper exists to prevent. + """ + f = getattr(msg, "file", None) + if f is None: + return None + size = getattr(f, "size", None) + return size if isinstance(size, int) else None + + +def file_stem(msg) -> str: + """Stem for the filename, before sanitization. Untrusted: the name comes + from DocumentAttributeFilename, which any group member controls.""" + name = getattr(getattr(msg, "file", None), "name", None) + if name: + stem = PurePosixPath(str(name)).stem + if stem: + return stem + return media_kind(msg) or "media" # photos and voice notes carry no name + + +def file_ext(msg) -> str: + """Extension including the dot, falling back to .bin. + + Prefers the extension of the declared filename, then Telethon's mime-derived + File.ext. Never empty, so a name is always recognizable on disk. + """ + f = getattr(msg, "file", None) + name = getattr(f, "name", None) + if name: + suffix = PurePosixPath(str(name)).suffix + # Byte-counted, like the stem truncation in paths.sanitize: 16 multibyte + # characters would be 48+ bytes, and "{id}_" + 200 stem bytes + that + # would overflow ext4's 255-byte component limit. + if suffix and len(suffix.encode()) <= 16: + return suffix + ext = getattr(f, "ext", None) + if ext: + return ext if str(ext).startswith(".") else f".{ext}" + return ".bin" diff --git a/src/telegram_exporter/paths.py b/src/telegram_exporter/paths.py new file mode 100644 index 0000000..b711c27 --- /dev/null +++ b/src/telegram_exporter/paths.py @@ -0,0 +1,167 @@ +"""All path construction, from the export root down to the filename leaf. + +Pure functions, no network, no I/O - which is why this is its own module rather +than a helper inside downloader.py, and why every rule below is table-tested +offline. + + The governing invariant: exactly one untrusted string ever becomes a path + component - the filename leaf. + +Every other component is an int, enforced by a type check, which makes +safe_join's "post_dir is trusted" precondition true by type rather than by +assumption. Two untrusted strings were in play before this module existed: + +1. DocumentAttributeFilename, supplied by arbitrary group members. + `../../../.ssh/authorized_keys` is a cheap, real attack. +2. The group title - server-supplied and editable by any admin. safe_join could + never have caught that one, because the escape happens in a parent component + before the join. It is fixed by never putting the title in a path at all. +""" + +from __future__ import annotations + +import os +import re +from pathlib import Path + +from .media import file_ext, file_stem, media_kind + +# Control characters are removed outright rather than replaced: a NUL truncates +# the filename at the syscall boundary, so leaving a placeholder would only hide +# what happened. +_CONTROL = {c: None for c in range(0x20)} +_CONTROL[0x7F] = None + +# Path separators and shell/Windows reserved characters become underscores. +_RESERVED = str.maketrans({ch: "_" for ch in '/\\:*?"<>|'}) + +_WHITESPACE_RUN = re.compile(r"\s+") + +# ext4 accepts 255 bytes per component; 200 leaves room for the "{id}_" prefix +# and the extension, both of which are byte-counted too. +MAX_NAME_BYTES = 200 +MAX_EXT_BYTES = 16 + + +def export_root(out_dir: Path, peer_id: int) -> Path: + """`/g` - the directory name is the chat id, never the title. + + peer_id comes from telethon.utils.get_peer_id(entity). An int cannot + traverse, cannot collapse into the parent when empty, cannot collide across + groups, and cannot strand an in-progress export when an admin renames the + group. The human-readable name is written to /title.txt as data. + + The `g` prefix only avoids a leading `-` in shell arguments; it is ergonomic, + not a security property. + """ + _require_int(peer_id, "peer_id") + return (Path(out_dir) / f"g{peer_id}").resolve() + + +def post_dir(root: Path, post_id: int) -> Path: + """One directory per logical post, named for the lowest message id in it.""" + _require_int(post_id, "post_id") + return root / str(post_id) + + +def _require_int(value: object, label: str) -> None: + # bool is an int subclass, and a bool here always means a caller bug. + if not isinstance(value, int) or isinstance(value, bool): + raise TypeError(f"{label} must be an int, got {type(value).__name__}: {value!r}") + + +def derive_filename(msg) -> str: + """`{message_id}_{sanitized stem}{ext}` - keyed on the message id. + + One Telegram message carries at most one media, so the message id makes + uniqueness inside a post folder structural rather than positional. + + A positional index is unstable under history mutation: delete one album + member between runs and every later index shifts, silently orphaning files + already on disk. It would also force Post to carry the full pre-filter + member list purely for naming, reintroducing the filter coupling that + Invariant 1 exists to remove. + + The redundancy in `1042/1042_photo.jpg` buys self-identifying files and lets + `grep 1043 messages.jsonl` join sidecar to disk with no index arithmetic. + """ + _require_int(getattr(msg, "id", None), "msg.id") + stem = sanitize(file_stem(msg), fallback=media_kind(msg) or "media") + return f"{msg.id}_{stem}{sanitize_ext(file_ext(msg))}" + + +def sanitize_ext(ext: str) -> str: + """The extension is part of the untrusted leaf and gets the same treatment. + + Sanitizing only the stem left a hole: the extension is derived from the same + attacker-supplied DocumentAttributeFilename, so `holiday.jpg` produced a + name containing a NUL, and the ValueError that Path.resolve then raised + escaped the download loop - halting the export on that message on every + subsequent run, since the cursor could not advance past it. Control + characters also meant an operator's terminal could be fed ANSI escapes from + a group member's filename. + """ + body = sanitize(str(ext).lstrip("."), fallback="bin", max_bytes=MAX_EXT_BYTES) + return f".{body}" + + +def sanitize(name: str, *, fallback: str, max_bytes: int = MAX_NAME_BYTES) -> str: + """Reduce an untrusted string to one safe path component. + + Order matters: control characters go before the leaf is taken, so a NUL + cannot confuse the split; the leaf is taken before separators are replaced, + so `../../../etc/passwd` loses its directories rather than becoming + `.._.._.._etc_passwd`. + """ + s = str(name).translate(_CONTROL) + s = _leaf(s) # discards every directory part, incl. .. + s = s.translate(_RESERVED) + s = _WHITESPACE_RUN.sub(" ", s).strip() + s = _strip_edges(s) + s = _truncate_utf8(s, max_bytes) + s = _strip_edges(s) # truncation can expose a new leading dot + return s or fallback + + +def _leaf(s: str) -> str: + """The final component only. Both separators, so a Windows-style path passed + to a Linux host still loses its directories.""" + return Path(s.replace("\\", "/")).name + + +def _strip_edges(s: str) -> str: + """Leading dots hide files; leading dashes make a filename look like a flag + to any tool the export is later piped through. Trailing dots and spaces are + stripped because they survive on Linux and silently vanish elsewhere.""" + return s.lstrip(". -\t").rstrip(". \t") + + +def _truncate_utf8(s: str, max_bytes: int) -> str: + """Truncate by UTF-8 bytes, preserving the extension. + + ext4's limit is 255 *bytes*, and a name of CJK or emoji characters overflows + it well before 255 characters - a character-counted truncation raises + OSError: File name too long partway through an export. + """ + if len(s.encode()) <= max_bytes: + return s + stem, ext = os.path.splitext(s) + ext_bytes = ext.encode() + if len(ext_bytes) > max_bytes: # pathological "extension"; drop it + stem, ext, ext_bytes = s, "", b"" + budget = max_bytes - len(ext_bytes) + # errors="ignore" discards a partial multi-byte character at the cut. + return stem.encode()[:budget].decode("utf-8", "ignore") + ext + + +def safe_join(post_dir: Path, filename: str) -> Path: + """Join and assert containment. Raises on escape - never silently repairs. + + sanitize should make this unreachable. Reaching it means sanitize has a hole, + so repairing here would convert a discovered bug into a silent one. + """ + base = Path(post_dir).resolve() + resolved = (base / filename).resolve() + if resolved == base or not resolved.is_relative_to(base): + raise ValueError(f"path escape attempt: {filename!r} against {post_dir}") + return resolved diff --git a/src/telegram_exporter/session.py b/src/telegram_exporter/session.py new file mode 100644 index 0000000..cdcde8f --- /dev/null +++ b/src/telegram_exporter/session.py @@ -0,0 +1,303 @@ +"""Credentials, session file preparation, the authenticated client, and the two +flood-wait primitives every network call routes through. + +Credential policy (the one authoritative copy): + + api_id, api_hash | env TG_API_ID / TG_API_HASH only | needed every run + phone | prompt; TG_PHONE optional | PII, not a secret + login code | prompt only | single-use, 5 min TTL + 2FA password | getpass only - no env, no file, no flag + +Nothing here reads a .env file. The shell already does that +(`set -a; . ./.env; set +a`). +""" + +from __future__ import annotations + +import asyncio +import logging +import os +import random +import stat +import subprocess +from contextlib import asynccontextmanager +from datetime import UTC, datetime, timedelta +from getpass import getpass +from pathlib import Path +from typing import AsyncIterator, Callable + +from telethon import TelegramClient, utils +from telethon.errors import FloodWaitError, SessionPasswordNeededError + +log = logging.getLogger("telegram_exporter.session") + +SUFFIX = ".session" + +# Routine short waits are absorbed inside Telethon rather than bubbling up. +FLOOD_SLEEP_THRESHOLD = 120 + + +class Abort(Exception): + """A run-ending condition that carries the process exit code. + + Raised instead of calling sys.exit deep in the call stack, so the caller + decides when to exit and the resume cursor is never written from an error + path. + """ + + def __init__(self, reason: str, code: int) -> None: + super().__init__(reason) + self.reason = reason + self.code = code + + +# --------------------------------------------------------------------------- # +# Session path +# --------------------------------------------------------------------------- # + +def default_session_path() -> Path: + """Outside the repo by default, so the safe location needs no opt-in and the + git check only bites a deliberately chosen path.""" + base = os.environ.get("XDG_STATE_HOME") or str(Path.home() / ".local" / "state") + return Path(base) / "tg-export" / f"default{SUFFIX}" + + +def with_session_suffix(p: Path) -> Path: + """Telethon appends '.session' when absent, so every consumer - the privacy + assertion, the git check, the lock file - has to agree on the real filename. + Pure: creates nothing, so callers can check a path before it exists.""" + return p if p.name.endswith(SUFFIX) else p.with_name(p.name + SUFFIX) + + +def prepare_session_path(p: Path) -> Path: + """Create the session file 0600 *before* Telethon can create it 0644. + + sqlite3.connect() on a fresh path yields 0644 under a default umask, and the + auth key is written during login - so a post-hoc chmod is too late. cli.main + sets umask(0o077) first; this is the belt to that braces. + """ + p = with_session_suffix(p) + p.parent.mkdir(parents=True, exist_ok=True) + try: + os.close(os.open(p, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)) + except FileExistsError: + os.chmod(p, 0o600) + return p + + +def assert_session_private(p: Path) -> None: + """Every sqlite sibling too: -journal, -wal and -shm hold the same secrets.""" + for f in sorted(p.parent.glob(p.name + "*")): + mode = stat.S_IMODE(f.stat().st_mode) + if mode & 0o077: + raise SystemExit(f"insecure mode on {f}: {oct(mode)}") + + +def _nearest_existing_dir(p: Path) -> Path | None: + """First existing ancestor directory, used as git's cwd. + + The target usually does not exist yet - an export root on a first run, or a + session path under a directory we have not created. Running git from a + nonexistent cwd raises FileNotFoundError before git can answer. + """ + for candidate in [p, *p.parents]: + if candidate.is_dir(): + return candidate + return None + + +def _check_ignore(target: str | Path, cwd: Path) -> int: + """git check-ignore rc: 0 = ignored, 1 = NOT ignored, 128 = not a work tree. + Returns 0 when git is unavailable, so a missing git never blocks a run.""" + try: + return subprocess.run(["git", "check-ignore", "-q", str(target)], + cwd=cwd, capture_output=True).returncode + except (FileNotFoundError, NotADirectoryError): + return 0 # no git installed, or the tree moved under us - not our problem + + +def assert_not_committable(p: Path, *, is_dir: bool = False) -> None: + """Refuse a session file or export root that git would happily commit. + + Adding .gitignore patterns cannot help when --out and --session point + anywhere, so this is checked at runtime against the actual path. + + is_dir must be True only for a path that will become a directory. The + trailing-slash retry below is what makes a directory-only pattern match + before the directory exists - and asking about a *file* that way would be + unsound in the other direction: a `.gitignore` line of `creds.session/` + would then excuse a file named `creds.session`, which git would commit. + """ + target = p.absolute() + cwd = _nearest_existing_dir(target) + if cwd is None: + return + rc = _check_ignore(target, cwd) + if rc == 1 and is_dir and not target.exists(): + # A directory-only pattern such as `exports/` does not match a path git + # cannot see is a directory, so an export root would be refused on the + # first run and accepted on every run after it. The trailing slash tells + # git to read it as a directory - passed as a string, because Path drops + # it. Verified empirically against git's exit codes. + rc = _check_ignore(f"{target}{os.sep}", cwd) + if rc == 1: + raise SystemExit( + f"refusing: {p} is inside a git repo and is not gitignored.\n" + f"add it to .gitignore, or choose a path outside the repo.") + + +# --------------------------------------------------------------------------- # +# Flood-wait primitives - two, because one cannot cover both shapes +# --------------------------------------------------------------------------- # + +async def _flood_sleep(e: FloodWaitError, max_wait_s: int | None, + where: str = "") -> None: + """Sleep out one flood wait, or abort if it exceeds an explicit ceiling. + + max_wait_s is None by default: sleep however long Telegram demands. Always + log the computed WAKE TIME, not just the duration - an unattended four-hour + sleep has to read as a sleep rather than as a hang. That log line is the + mitigation for the no-ceiling default. + """ + if max_wait_s is not None and e.seconds > max_wait_s: + raise Abort(f"flood wait {e.seconds}s exceeds --max-flood-wait {max_wait_s}s", 6) + wake = datetime.now(UTC) + timedelta(seconds=e.seconds) + log.warning("flood wait %ss%s - sleeping until %s", + e.seconds, where, wake.isoformat(timespec="seconds")) + await asyncio.sleep(e.seconds + random.uniform(1.0, 3.0)) + + +async def with_flood_retry(make_coro: Callable[[], object], *, + max_wait_s: int | None = None): + """Download path. Never tight-retry: a retry during an active wait is a + fresh violation that escalates the next wait.""" + while True: + try: + return await make_coro() + except FloodWaitError as e: + await _flood_sleep(e, max_wait_s) + + +async def aiter_with_flood_retry(make_agen, *, start_id: int, + max_wait_s: int | None = None): + """Sweep path. + + iter_messages is consumed with `async for`, and Telethon raises + FloodWaitError on the *next page fetch* - deep inside the loop, where a + wrapper guarding only construction never sees it. Rebuilding the generator + from the last id actually yielded is what keeps the sweep gap-free: with + reverse=True, offset_id is an exclusive lower bound, so no message is + skipped or replayed. + """ + last = start_id + while True: + try: + async for item in make_agen(last): + last = item.id + yield item + return + except FloodWaitError as e: + await _flood_sleep(e, max_wait_s, where=f" mid-sweep after id {last}") + + +# --------------------------------------------------------------------------- # +# Credentials and client +# --------------------------------------------------------------------------- # + +def load_api_credentials() -> tuple[int, str]: + api_id, api_hash = os.environ.get("TG_API_ID"), os.environ.get("TG_API_HASH") + if not api_id or not api_hash: + raise Abort( + "missing credentials: set TG_API_ID and TG_API_HASH.\n" + " create an app at https://my.telegram.org -> API development tools\n" + " then: export TG_API_ID=... TG_API_HASH=...", 1) + try: + return int(api_id), api_hash + except ValueError: + raise Abort(f"TG_API_ID must be an integer, got {api_id!r}", 1) from None + + +async def _login(client: TelegramClient) -> None: + """Interactive first login. Later runs return immediately.""" + if await client.is_user_authorized(): + return + phone = os.environ.get("TG_PHONE") or input("phone (+countrycode): ").strip() + if not phone: + raise Abort("no phone number given", 1) + await client.send_code_request(phone) + code = input("login code (sent via Telegram): ").strip() + try: + await client.sign_in(phone, code) + except SessionPasswordNeededError: + # 2FA password: getpass only. Never argv, never env, never a file. + await client.sign_in(password=getpass("2FA password: ")) + + +@asynccontextmanager +async def connected_client(session_path: Path, api_id: int, + api_hash: str) -> AsyncIterator[TelegramClient]: + """Connected, authorized client that always disconnects. + + Telethon does not flush SQLite session state without disconnect(), so this + is a context manager rather than a discipline. + """ + # Check before creating: a refused path should not be left holding an empty + # session file. + session_path = with_session_suffix(session_path) + assert_not_committable(session_path) + session_path = prepare_session_path(session_path) + client = TelegramClient(str(session_path), api_id, api_hash) + client.flood_sleep_threshold = FLOOD_SLEEP_THRESHOLD + await client.connect() + try: + await _login(client) + assert_session_private(session_path) + yield client + finally: + await client.disconnect() + + +async def resolve_entity(client: TelegramClient, spec: str): + """Accept a numeric id, @username, or t.me link. Returns (entity, peer_id). + + peer_id is telethon.utils.get_peer_id(entity) - the canonical int that + becomes the export directory name. + """ + if "/+" in spec or "joinchat/" in spec: + raise Abort( + f"{spec} is a private invite link, which cannot be resolved without " + f"joining.\njoin the group in a Telegram client first, then pass its " + f"numeric id or @username.", 5) + + target = spec.strip() + try: + entity = await client.get_entity(int(target)) + except ValueError as e: + # Not an int, or an int Telethon has never seen. Raw ids resolve only + # from the session cache, so fall back to a dialog scan before giving up. + try: + entity = await client.get_entity(target) + except ValueError: + entity = await _find_in_dialogs(client, target) + if entity is None: + raise Abort( + f"cannot resolve group {spec!r}: {e}\n" + f"confirm this account is a member, and try the @username or " + f"the numeric id.", 5) from None + return entity, utils.get_peer_id(entity) + + +async def _find_in_dialogs(client: TelegramClient, target: str): + """Last resort for a numeric id absent from the session cache.""" + try: + wanted = int(target) + except ValueError: + return None + async for dialog in client.iter_dialogs(): + if utils.get_peer_id(dialog.entity) == wanted: + return dialog.entity + return None + + +def now_iso() -> str: + return datetime.now(UTC).isoformat(timespec="seconds") diff --git a/src/telegram_exporter/sidecar.py b/src/telegram_exporter/sidecar.py new file mode 100644 index 0000000..7e9cf78 --- /dev/null +++ b/src/telegram_exporter/sidecar.py @@ -0,0 +1,149 @@ +"""messages.jsonl - the append-only log linking every file on disk to its message. + +A record is an event: "the files written for this message in this run". The +correct consumer rule is therefore the UNION of all records for a message_id, +not last-write-wins. Last-write-wins is actively wrong: after a filter change it +would report a narrower file list than what is actually on disk. +""" + +from __future__ import annotations + +import json +import logging +import os +from datetime import UTC, datetime +from pathlib import Path + +log = logging.getLogger("telegram_exporter.sidecar") + +SIDECAR_NAME = "messages.jsonl" + + +class Sidecar: + """Append-only JSONL writer. Opened for the life of the run.""" + + def __init__(self, root: Path) -> None: + self.root = root + self.path = root / SIDECAR_NAME + self._fh = None + + # ---- lifecycle ------------------------------------------------------ # + + def open(self) -> Sidecar: + self.root.mkdir(parents=True, exist_ok=True) + self.repair() + self._fh = open(self.path, "a", encoding="utf-8") + return self + + def close(self) -> None: + if self._fh is not None: + self._fh.close() + self._fh = None + + def __enter__(self) -> Sidecar: + return self.open() + + def __exit__(self, *exc) -> None: + self.close() + + def repair(self) -> None: + """Truncate a partial trailing line written when the host died mid-write. + + One unterminated line breaks every JSONL consumer, and the damage is + always confined to the tail, so scanning back from EOF for the last + newline is both sufficient and cheap on a 100k-line file. + """ + if not self.path.exists() or self.path.stat().st_size == 0: + return + with open(self.path, "rb+") as f: + f.seek(0, os.SEEK_END) + end = f.tell() + window = min(end, 1 << 20) + f.seek(end - window) + tail = f.read(window) + cut = tail.rfind(b"\n") + trailing = tail[cut + 1:] + if not trailing: + return # file ends on a newline: intact + try: + json.loads(trailing) + except (json.JSONDecodeError, UnicodeDecodeError): + offset = end - len(trailing) + f.truncate(offset) + log.warning("repaired %s: truncated %d trailing bytes of a partial " + "record at offset %d", self.path, len(trailing), offset) + return + # Valid JSON but no terminating newline - complete record, just add one. + f.write(b"\n") + + def rotate(self) -> Path | None: + """--reset-state rotates rather than deletes, so the current sidecar always + describes exactly one filter regime while the old one stays inspectable.""" + if not self.path.exists(): + return None + stamp = datetime.now(UTC).strftime("%Y%m%dT%H%M%SZ") + target = self.root / f"messages-{stamp}.jsonl" + os.replace(self.path, target) + log.info("rotated %s -> %s", self.path.name, target.name) + return target + + # ---- writing -------------------------------------------------------- # + + def append(self, post, results) -> None: + """One record per surviving message in the post. + + Text-only messages under --include-text produce the same shape with + "files": [] - consumers distinguish media from context by files being + empty, not by a separate record type. + """ + by_message = {} + for r in results: + by_message.setdefault(r.message_id, []).append(r) + + for msg in post.messages: + rs = by_message.get(msg.id, []) + record = { + "message_id": msg.id, + "post_id": post.post_id, + "grouped_id": getattr(msg, "grouped_id", None), + "date": _iso(getattr(msg, "date", None)), + "sender_id": getattr(msg, "sender_id", None), + "sender_name": _sender_name(msg), + "caption": (getattr(msg, "message", None) or None), + "reply_to": _reply_to(msg), + "files": [r.as_record(self.root) for r in rs if r.path is not None], + "errors": [r.error for r in rs if r.error is not None], + } + self._fh.write(json.dumps(record, ensure_ascii=False) + "\n") + + def fsync(self) -> None: + """Called before the cursor commits: the sidecar must be durable before + the state file claims the post is done.""" + self._fh.flush() + os.fsync(self._fh.fileno()) + + +def _iso(value) -> str | None: + return value.isoformat() if isinstance(value, datetime) else None + + +def _reply_to(msg) -> int | None: + header = getattr(msg, "reply_to", None) + return getattr(header, "reply_to_msg_id", None) if header is not None else None + + +def _sender_name(msg) -> str | None: + """Only from an already-cached sender. + + Never an extra get_entity call: on a 100k-message sweep that is 100k extra + RPCs and a flood ban, in exchange for a display string. + """ + sender = getattr(msg, "sender", None) + if sender is None: + return None + title = getattr(sender, "title", None) + if title: + return title + parts = [getattr(sender, "first_name", None), getattr(sender, "last_name", None)] + name = " ".join(p for p in parts if p) + return name or getattr(sender, "username", None) diff --git a/src/telegram_exporter/state.py b/src/telegram_exporter/state.py new file mode 100644 index 0000000..6d55fcc --- /dev/null +++ b/src/telegram_exporter/state.py @@ -0,0 +1,163 @@ +"""The resume cursor, and nothing else. + +state.py and sidecar.py are separate modules on purpose. The correctness-critical +part is the *ordering* between them - sidecar first, cursor second - and that +ordering belongs to the caller. Two objects in two modules makes the sequence +visible at one call site, and makes a "helpful" commit_and_append() impossible to +write without noticing that it destroys Invariant 2. +""" + +from __future__ import annotations + +import json +import logging +import os +from pathlib import Path + +from .session import Abort, now_iso +from .traversal import Filters + +log = logging.getLogger("telegram_exporter.state") + +STATE_NAME = ".export-state.json" +VERSION = 1 + + +def fsync_dir(path: Path) -> None: + """A rename is not durable until its directory entry is. + + Lives here rather than in downloader.py because durability is what this + module exists for; downloader.py imports it for the post-directory fsync that + has to happen before the cursor claims a post is done. + """ + fd = os.open(path, os.O_RDONLY) + try: + os.fsync(fd) + finally: + os.close(fd) + + +class State: + """cursor_id means: every post with max_message_id <= cursor_id has been + handled *under these filters*. + + One cursor field, not two. The sweep is a generator driven by the sink at + concurrency 1, so it can never run ahead of the downloader - a second + "swept_to" field would invent a divergence to manage. + """ + + def __init__(self, path: Path, data: dict) -> None: + self.path = path + self.data = data + + # ---- construction -------------------------------------------------- # + + @classmethod + def path_for(cls, root: Path) -> Path: + return root / STATE_NAME + + @classmethod + def read(cls, root: Path) -> dict | None: + """Read-only load for --dry-run. Never creates anything.""" + p = cls.path_for(root) + if not p.exists(): + return None + try: + return json.loads(p.read_text()) + except (json.JSONDecodeError, OSError) as e: + raise Abort(f"cannot read {p}: {e}\n" + f"inspect it, or use --reset-state to start over", 7) from None + + @classmethod + def open(cls, root: Path, *, chat_id: int, chat_title: str, + filters: Filters, reset: bool = False) -> State: + """Load and validate, or create. Refuses a mismatch rather than guessing.""" + path = cls.path_for(root) + existing = None if reset else cls.read(root) + + if existing is not None: + _assert_compatible(existing, chat_id=chat_id, filters=filters) + existing["chat_title"] = chat_title # titles change; that is fine + existing["completed_at"] = None # this run has not finished yet + existing["updated_at"] = now_iso() + state = cls(path, existing) + else: + state = cls(path, { + "version": VERSION, + "chat_id": chat_id, + "chat_title": chat_title, + "filters": filters.to_state(), + "cursor_id": 0, + "completed_at": None, + "started_at": now_iso(), + "updated_at": now_iso(), + }) + state.save() + return state + + # ---- accessors ------------------------------------------------------ # + + @property + def cursor_id(self) -> int: + return int(self.data["cursor_id"]) + + # ---- mutation ------------------------------------------------------- # + + def commit(self, cursor_id: int) -> None: + """Invariant 2: only ever called after a COMPLETE post, and unreachable + from every error exit, so the cursor can never point past in-flight work.""" + if cursor_id < self.cursor_id: + log.debug("ignoring cursor regression %s -> %s", self.cursor_id, cursor_id) + return + self.data["cursor_id"] = int(cursor_id) + self.save() + + def mark_completed(self) -> None: + """Set only on clean exhaustion of the sweep. The only way to answer + "did my 20-hour export finish?" without guessing.""" + self.data["completed_at"] = now_iso() + self.save() + + def save(self) -> None: + """Atomic: tmp -> fsync -> replace -> fsync(dir). A half-written state + file would be indistinguishable from a corrupted cursor.""" + self.data["updated_at"] = now_iso() + self.path.parent.mkdir(parents=True, exist_ok=True) + tmp = self.path.with_name(self.path.name + ".tmp") + with open(tmp, "w") as f: + json.dump(self.data, f, indent=2, sort_keys=True) + f.write("\n") + f.flush() + os.fsync(f.fileno()) + os.replace(tmp, self.path) + fsync_dir(self.path.parent) + + +def _assert_compatible(existing: dict, *, chat_id: int, filters: Filters) -> None: + """Invariant 3: filters are stored verbatim, and a mismatch names what changed. + + A sha256 fingerprint would be smaller to store and useless here: the whole + point is to *print the diff*, and a digest is one-way. + + --limit is deliberately absent from the filter set. It changes when we stop, + not what a post contains, so including it would let a `--limit 20` test run + poison every later full run. + """ + if int(existing.get("chat_id", 0)) != chat_id: + raise Abort( + f"chat mismatch: state file belongs to chat {existing.get('chat_id')}, " + f"not {chat_id}\nexport each group to its own --out directory", 7) + + was, now = existing.get("filters") or {}, filters.to_state() + changed = [k for k in sorted(set(was) | set(now)) if was.get(k) != now.get(k)] + if changed: + lines = "\n".join(f" {k+':':12} {_show(was.get(k))} -> {_show(now.get(k))}" + for k in changed) + raise Abort( + f"filter mismatch:\n{lines}\n" + f"the cursor means 'handled through here under these filters'.\n" + f"re-run with the original filters, or --reset-state to start over", 7) + + +def _show(v: object) -> str: + return "(none)" if v is None else repr(v) diff --git a/src/telegram_exporter/traversal.py b/src/telegram_exporter/traversal.py new file mode 100644 index 0000000..4426a7a --- /dev/null +++ b/src/telegram_exporter/traversal.py @@ -0,0 +1,239 @@ +"""The single traversal generator both run modes consume. + +Chronological sweep, album buffering, filter application. Yields one Post per +logical post a human would recognize. + +`run_download()` and `run_estimate()` drive *this* generator with the same +Filters. That shared generator - not any class hierarchy - is what guarantees the +dry-run estimate describes what a real run would download. A forked traversal +would drift and lie, so never fork it. +""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass, field +from datetime import UTC, datetime +from typing import Any + +from .media import KINDS, media_kind, media_size +from .session import aiter_with_flood_retry + +log = logging.getLogger("telegram_exporter.traversal") + +# An album is capped at 10 members by Telegram; 20 is slack for a protocol +# change, and exceeding it means the grouping logic is wrong, not that Telegram +# got generous. +MAX_ALBUM_BUFFER = 20 + +_SIZE_UNITS = {"": 1, "B": 1, "K": 10**3, "KB": 10**3, "KIB": 2**10, + "M": 10**6, "MB": 10**6, "MIB": 2**20, + "G": 10**9, "GB": 10**9, "GIB": 2**30, + "T": 10**12, "TB": 10**12, "TIB": 2**40} + + +class TraversalError(Exception): + """The sweep is not the shape the whole design assumes. Never repaired: a + wrong assumption here produces silent corruption that looks like success.""" + + +class AlbumSplitError(TraversalError): + """A grouped_id reopened after its run closed - the album was split across a + non-adjacent boundary, so folder identity is no longer trustworthy.""" + + +class SweepOrderError(TraversalError): + """The sweep is not strictly ascending, or an album exceeded its bound. + + A real exception rather than a bare assert: `python -O` strips assertions, + and these two tripwires guard the same class of silent corruption as + AlbumSplitError sitting beside them. + """ + + +@dataclass(frozen=True) +class Post: + """One logical post: a standalone message, or every member of one album. + + grouped_id is deliberately absent. The sidecar reads it off the message and + folder identity uses post_id, so a field here would be a third copy of a + fact with two owners already. + """ + + post_id: int # lowest message id in the FULL group (pre-filter) + messages: list[Any] # filter-surviving members, ascending by id + max_message_id: int # highest id in the FULL group - the cursor value + + +@dataclass(frozen=True) +class Filters: + """What a post *contains*. Stored verbatim in the state file, because the + cursor means "handled through here under these filters".""" + + types: frozenset[str] | None = None + since: datetime | None = None + until: datetime | None = None + max_size: int | None = None + include_text: bool = False + + # Set once, so the unknown-size warning does not repeat per file. + _warned: set[str] = field(default_factory=set, compare=False, repr=False) + + def to_state(self) -> dict: + """The stored form. Verbatim and comparable - no digest, because two + call sites have to *name* what changed, and a hash cannot. + + `is not None`, not truthiness: an empty set is a filter that drops + everything, and storing it as "no filter" would let a later unfiltered + run inherit an end-of-history cursor and report success having downloaded + nothing. build_filters refuses to construct one, so this is the second + line of defense on Invariant 3. + """ + return { + "types": sorted(self.types) if self.types is not None else None, + "since": self.since.isoformat() if self.since else None, + "until": self.until.isoformat() if self.until else None, + "max_size": self.max_size, + "include_text": self.include_text, + } + + +def parse_size(text: str) -> int: + """`100MB`, `2GiB`, `1500000`. Decimal units are powers of 10, `*iB` are + powers of 2 - the same convention `ls -h` and `df -h` use.""" + raw = str(text).strip().replace(" ", "") + i = 0 + while i < len(raw) and (raw[i].isdigit() or raw[i] == "."): + i += 1 + number, unit = raw[:i], raw[i:].upper() + if not number or unit not in _SIZE_UNITS: + raise ValueError(f"cannot parse size {text!r}; try 100MB, 2GiB or a byte count") + return int(float(number) * _SIZE_UNITS[unit]) + + +def parse_date(text: str) -> datetime: + """`--since`/`--until` as UTC-aware, because msg.date is UTC-aware and a + naive comparison raises TypeError mid-sweep.""" + value = datetime.fromisoformat(str(text).strip()) + return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC) + + +def build_filters(*, types: str | None = None, since: str | None = None, + until: str | None = None, max_size: str | None = None, + include_text: bool = False) -> Filters: + kinds = None + if types: + kinds = frozenset(t.strip().lower() for t in types.split(",") if t.strip()) + unknown = kinds - set(KINDS) + if unknown: + raise ValueError(f"unknown media types {sorted(unknown)}; " + f"choose from {', '.join(KINDS)}") + if not kinds: + # `--types " "` or `--types ,` - easy to produce from an unset shell + # variable. Left alone it is a filter that drops every media kind, + # which would sweep the whole history, download nothing, and commit + # an end-of-history cursor. + raise ValueError(f"--types {types!r} names no media kind; " + f"choose from {', '.join(KINDS)}, or omit the flag") + return Filters( + types=kinds, + since=parse_date(since) if since else None, + until=parse_date(until) if until else None, + max_size=parse_size(max_size) if max_size else None, + include_text=include_text, + ) + + +def keep(msg, filters: Filters) -> bool: + """Whether this message survives the filters. + + Text-only messages are dropped unless --include-text, in which case they + reach the sidecar with an empty files[] so captions and replies keep their + surrounding conversation. They never create a post folder and never reach + download_one. + """ + kind = media_kind(msg) + if kind is None: + return filters.include_text + if filters.types is not None and kind not in filters.types: + return False + date = getattr(msg, "date", None) + if filters.since is not None and (date is None or date < filters.since): + return False + if filters.until is not None and (date is None or date > filters.until): + return False + if filters.max_size is not None: + size = media_size(msg) + if size is None: + # Include and warn once. A silent drop violates the no-silent-gaps + # ethos, and unknown-size media is pathological and near-always small. + if "max_size_none" not in filters._warned: + filters._warned.add("max_size_none") + log.warning("some media reports no size; --max-size cannot apply " + "to it and it is being included") + elif size > filters.max_size: + return False + return True + + +async def iter_posts(client, entity, *, after_id: int = 0, + filters: Filters, max_flood_wait: int | None = None): + """Sweep oldest -> newest, yielding one Post per logical post. + + after_id is an exclusive lower bound: with reverse=True Telethon bumps + offset_id by one before the first request and again per page (verified + against telethon 1.44.0's _MessagesIter), so passing the last handled id + resumes without replaying it. The explicit `msg.id <= after_id` guard makes + the sweep correct even if that internal detail ever changes. + """ + buf: list[Any] = [] + closed: set[int] = set() + prev = 0 + + def agen(since: int): + return client.iter_messages(entity, reverse=True, offset_id=since) + + async for msg in aiter_with_flood_retry(agen, start_id=after_id, + max_wait_s=max_flood_wait): + if msg.id <= after_id: + continue + if msg.id <= prev: + raise SweepOrderError( + f"sweep not strictly ascending: {msg.id} after {prev}. " + f"Album grouping and the cursor both assume ascending order.") + prev = msg.id + + gid = getattr(msg, "grouped_id", None) + if buf and (gid is None or gid != buf[0].grouped_id): + yield _close(buf, closed, filters) + buf = [] + if gid is None: + yield _close([msg], closed, filters) + else: + buf.append(msg) + if len(buf) > MAX_ALBUM_BUFFER: + raise SweepOrderError( + f"album buffer exceeded {MAX_ALBUM_BUFFER} for grouped_id " + f"{gid} - grouping logic is wrong, refusing to keep buffering") + if buf: + yield _close(buf, closed, filters) + + +def _close(buf: list[Any], closed: set[int], filters: Filters) -> Post: + """Turn a finished run of messages into a Post. + + Invariant 1 - identity is filter-invariant. post_id and max_message_id come + from the complete grouped_id group, before filters drop any member, so + `--types photo` and `--types video` address the same folder. Phase 4's + message-id filenames close the same loop for files. + """ + gid = getattr(buf[0], "grouped_id", None) + if gid is not None: + if gid in closed: + raise AlbumSplitError( + f"grouped_id {gid} reopened at msg {buf[0].id} after being closed - " + f"album split across a non-adjacent boundary. Do not trust this export.") + closed.add(gid) + return Post(post_id=buf[0].id, + messages=[m for m in buf if keep(m, filters)], + max_message_id=buf[-1].id)