mirror of
https://github.com/tiennm99/telegram-exporter.git
synced 2026-10-11 03:13:49 +00:00
feat: add tg-export, a resumable Telegram group media exporter
Exports all media from a group to disk, one folder per logical post with albums collapsed, alongside a messages.jsonl sidecar linking every file to its message. Three properties the design is built around: - Identity is filter-invariant. A post's folder is the lowest message id in the full grouped_id group and filenames key on message id, so --types photo and --types video address the same files. A positional index would shift whenever a member message was deleted. - The cursor commits only after a complete post and is unreachable from every error path, so it can never point past in-flight work. - Filters are stored verbatim; a mismatch refuses the run and names what changed, because the cursor means 'handled through here under these filters'. Exactly one untrusted string becomes a path component: the filename leaf, sanitized stem and extension both. Every other component is an int enforced by type. The export directory is the chat id rather than the group title, which makes traversal via a renamed group structurally impossible and stops a mid-export rename from orphaning the download. Downloads are sequential. Telethon does not parallelize transfers, and concurrency mainly accelerates flood-wait escalation from seconds to hours.
This commit is contained in:
1 parent
65f8abdf18
commit
275ec47576
10 files changed
+2005
No files matched your search
@@ -0,0 +1,3 @@
|
||||
"""Export all media from a Telegram group to local disk, grouped by post."""
|
||||
|
||||
__version__ = "0.1.0"
|
||||
@@ -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 <DIR>/g<chat_id>")
|
||||
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())
|
||||
@@ -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")
|
||||
@@ -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
|
||||
@@ -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"
|
||||
@@ -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:
|
||||
"""`<out>/g<chat_id>` - 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 <root>/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.<NUL>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
|
||||
@@ -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")
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
Reference in new issue
Block a user