Add telegram-archiver microservice (MVP)

This commit is contained in:
2026-02-19 21:41:52 +03:00
parent 7166f535d6
commit 6332b91ba4
13 changed files with 1635 additions and 0 deletions
+10
View File
@@ -0,0 +1,10 @@
"""
Telegram Archiver - Archive Telegram channels to Markdown bundles.
This package provides both CLI and REST API interfaces for archiving
Telegram channels with their media content to local filesystem.
"""
__version__ = "1.0.0"
__author__ = "Evgeny Storozhenko"
__email__ = "dedinit"
+374
View File
@@ -0,0 +1,374 @@
"""
Core archiver logic.
Handles message processing, markdown generation, and file organization.
"""
import asyncio
import logging
from datetime import datetime
from pathlib import Path
from typing import Optional
import yaml
from telethon.tl.types import Message
from app.models import PostData, MediaFile, MediaType
from app.telethon_client import TelethonArchiver
from config import Settings, get_settings
logger = logging.getLogger(__name__)
class ChannelArchiver:
"""
Handles archiving of a single Telegram channel.
Manages bundle creation, markdown generation, and media downloads.
"""
def __init__(
self,
telethon_client: TelethonArchiver,
settings: Optional[Settings] = None,
):
self.client = telethon_client
self.settings = settings or get_settings()
self.channel_path: Optional[Path] = None
self.log_file: Optional[Path] = None
# Statistics
self.stats = {
"posts_archived": 0,
"posts_skipped": 0,
"media_downloaded": 0,
"media_skipped_large": 0,
"large_files": [],
}
def _setup_logging(self, channel_username: str) -> None:
"""Setup file logging for this channel archive session."""
timestamp = datetime.now().strftime("%Y%m%d")
log_filename = f"{timestamp}-{channel_username}.log"
self.log_file = Path(__file__).parent.parent / log_filename
# Create file handler
file_handler = logging.FileHandler(self.log_file, encoding="utf-8")
file_handler.setLevel(self.settings.log_level)
# Create formatter
formatter = logging.Formatter(
"%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
file_handler.setFormatter(formatter)
# Add handler to logger
logger.addHandler(file_handler)
logger.info(f"=== Archive session started for {channel_username} ===")
def _get_channel_path(self, channel_username: str) -> Path:
"""Get the output directory path for this channel."""
# Remove @ prefix if present
clean_username = channel_username.lstrip("@")
base_path = self.settings.output_dir / clean_username
base_path.mkdir(parents=True, exist_ok=True)
return base_path
async def archive_channel(
self,
channel: str,
limit: Optional[int] = None,
from_message_id: Optional[int] = None,
force: bool = False,
) -> dict:
"""
Archive all messages from a channel.
Args:
channel: Channel username or ID
limit: Maximum number of messages to archive
from_message_id: Start from this message ID
force: Force re-download of already archived posts
Returns:
Statistics dictionary
"""
start_time = datetime.now()
# Setup
self._setup_logging(channel)
self.channel_path = self._get_channel_path(channel)
logger.info(f"Starting archive of {channel}")
logger.info(f"Output directory: {self.channel_path}")
if limit:
logger.info(f"Limit: {limit} messages")
if from_message_id:
logger.info(f"Starting from message ID: {from_message_id}")
# Get channel info
entity = await self.client.get_channel_info(channel)
channel_title = getattr(entity, "title", channel)
# Reset stats
self.stats = {
"posts_archived": 0,
"posts_skipped": 0,
"media_downloaded": 0,
"media_skipped_large": 0,
"large_files": [],
}
# Process messages
async for message in self.client.fetch_messages(
channel, limit=limit, from_message_id=from_message_id
):
await self._process_message(message, force=force)
# Write large files report
large_files_report = await self._write_large_files_report()
# Calculate duration
duration = (datetime.now() - start_time).total_seconds()
logger.info(f"=== Archive completed in {duration:.2f}s ===")
logger.info(f"Posts archived: {self.stats['posts_archived']}")
logger.info(f"Posts skipped: {self.stats['posts_skipped']}")
logger.info(f"Media downloaded: {self.stats['media_downloaded']}")
logger.info(f"Media skipped (large): {self.stats['media_skipped_large']}")
return {
"channel": channel,
"channel_title": channel_title,
"posts_archived": self.stats["posts_archived"],
"posts_skipped": self.stats["posts_skipped"],
"media_downloaded": self.stats["media_downloaded"],
"media_skipped_large": self.stats["media_skipped_large"],
"output_path": str(self.channel_path),
"log_file": str(self.log_file),
"large_files_report": large_files_report,
"duration_seconds": duration,
}
async def _process_message(
self, message: Message, force: bool = False
) -> None:
"""
Process a single message: create bundle, download media, generate markdown.
Args:
message: Telegram message to process
force: Force re-processing if bundle exists
"""
message_id = message.id
bundle_dir = self.channel_path / str(message_id)
# Check if already archived (skip if exists and not force)
if bundle_dir.exists() and not force:
logger.debug(f"Skipping message {message_id}: already archived")
self.stats["posts_skipped"] += 1
return
logger.debug(f"Processing message {message_id}...")
# Create bundle directory
bundle_dir.mkdir(parents=True, exist_ok=True)
# Extract post data
post_data = await self._extract_post_data(message)
# Download media
if message.media:
media_info = await self.client.download_media(
message, bundle_dir, self.settings.max_file_size
)
# Update post_data with media info
for info in media_info:
media_file = MediaFile(
filename=info["filename"],
type=info["type"],
caption=info.get("caption"),
size=info.get("size"),
is_too_large=info.get("is_too_large", False),
)
post_data.media_files.append(media_file)
if info.get("is_too_large"):
self.stats["media_skipped_large"] += 1
self.stats["large_files"].append(
{
"message_id": message_id,
"filename": info["filename"],
"size": info.get("size"),
}
)
else:
self.stats["media_downloaded"] += 1
# Generate markdown
markdown_content = self._generate_markdown(post_data)
# Write index.md
index_path = bundle_dir / "index.md"
index_path.write_text(markdown_content, encoding="utf-8")
self.stats["posts_archived"] += 1
logger.debug(f"Message {message_id} archived successfully")
async def _extract_post_data(self, message: Message) -> PostData:
"""
Extract structured data from a Telegram message.
Args:
message: Telegram message
Returns:
PostData object with all metadata
"""
# Get channel info
channel = await message.get_chat()
channel_username = getattr(channel, "username", "unknown")
channel_title = getattr(channel, "title", channel_username)
# Convert message text to markdown
text = message.message or ""
if message.message:
# Use Telethon's built-in markdown converter
text = message.get_message_text(markdown=True)
# Handle reply-to
reply_to = None
if message.reply_to and message.reply_to.reply_to_msg_id:
reply_to = message.reply_to.reply_to_msg_id
# Handle repost/forward
repost_from = None
repost_channel = None
if message.fwd_from:
repost_from = message.fwd_from.channel_post
if message.fwd_from.from_id:
# Try to get original channel info
try:
fwd_channel = await self.client.client.get_entity(
message.fwd_from.from_id
)
repost_channel = getattr(fwd_channel, "title", None)
except Exception:
pass
return PostData(
message_id=message.id,
date=message.date,
text=text,
author=channel_title,
channel_username=channel_username,
reply_to=reply_to,
repost_from=repost_from,
repost_channel=repost_channel,
views=getattr(message, "views", None),
)
def _generate_markdown(self, post: PostData) -> str:
"""
Generate markdown content with front-matter for Hugo.
Args:
post: PostData object
Returns:
Markdown string with YAML front-matter
"""
# Build front-matter
frontmatter = {
"message_id": post.message_id,
"date": post.date.isoformat(),
"author": post.author,
}
# Add optional fields
if post.reply_to:
# Relative link to the replied post's index.md
frontmatter["reply_to"] = f"../{post.reply_to}/index.md"
if post.repost_from:
frontmatter["repost_from"] = post.repost_from
if post.repost_channel:
frontmatter["repost_channel"] = post.repost_channel
if post.media_files:
frontmatter["media_files"] = [
{
"filename": mf.filename,
"type": mf.type,
"caption": mf.caption,
"size": mf.size,
"is_too_large": mf.is_too_large,
}
for mf in post.media_files
]
# Build markdown content
content = []
# Add media embeds for downloaded files
for mf in post.media_files:
if not mf.is_too_large:
if mf.type == "photo":
content.append(f"![{mf.caption or ''}]({mf.filename})")
elif mf.type == "video":
content.append(f"@[video]({mf.filename})")
elif mf.type == "audio":
content.append(f"@[audio]({mf.filename})")
elif mf.type == "document":
content.append(f"📎 [{mf.filename}]({mf.filename})")
# Add text content
if post.text:
content.append(post.text)
# Mark reposts visually
if post.repost_from:
content.insert(0, f"> *Repost from {post.repost_channel or 'Unknown'}*")
content.insert(1, "")
# Combine front-matter and content
yaml_frontmatter = yaml.dump(
frontmatter,
allow_unicode=True,
default_flow_style=False,
sort_keys=False,
)
return f"---\n{yaml_frontmatter}---\n\n{''.join(content)}"
async def _write_large_files_report(self) -> Optional[str]:
"""
Write report of files that were too large to download.
Returns:
Path to report file or None if no large files
"""
if not self.stats["large_files"]:
return None
report_path = self.channel_path / "2big2get.md"
content = ["# Files Too Large to Download\n\n"]
content.append(
f"Generated: {datetime.now().isoformat()}\n\n"
)
content.append(
f"Total files skipped: {len(self.stats['large_files'])}\n\n"
)
content.append("| Message ID | Filename | Size (bytes) |\n")
content.append("|------------|----------|-------------|\n")
for file_info in self.stats["large_files"]:
content.append(
f"| {file_info['message_id']} | {file_info['filename']} | "
f"{file_info['size']} |\n"
)
report_path.write_text("".join(content), encoding="utf-8")
logger.info(f"Large files report written to {report_path}")
return str(report_path)
+87
View File
@@ -0,0 +1,87 @@
"""
Logging configuration for Telegram Archiver.
Supports both console and file output with OpenTelemetry integration.
"""
import logging
import sys
from pathlib import Path
from typing import Optional
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import ConsoleSpanExporter, SimpleSpanProcessor
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
def setup_logging(
log_level: str = "INFO",
log_file: Optional[Path] = None,
enable_otel: bool = False,
) -> None:
"""
Configure logging for the application.
Args:
log_level: Logging level (DEBUG, INFO, WARNING, ERROR)
log_file: Optional file path for log output
enable_otel: Enable OpenTelemetry tracing
"""
# Create root logger
root_logger = logging.getLogger()
root_logger.setLevel(getattr(logging, log_level.upper()))
# Clear existing handlers
root_logger.handlers.clear()
# Console handler
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setLevel(getattr(logging, log_level.upper()))
console_formatter = logging.Formatter(
"%(asctime)s - %(name)s - %(levelname)s - %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
console_handler.setFormatter(console_formatter)
root_logger.addHandler(console_handler)
# File handler (if specified)
if log_file:
file_handler = logging.FileHandler(log_file, encoding="utf-8")
file_handler.setLevel(getattr(logging, log_level.upper()))
file_formatter = logging.Formatter(
"%(asctime)s - %(name)s - %(levelname)s - %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
file_handler.setFormatter(file_formatter)
root_logger.addHandler(file_handler)
# OpenTelemetry configuration
if enable_otel:
setup_opentelemetry()
def setup_opentelemetry() -> None:
"""Configure OpenTelemetry tracing."""
# Set up tracer provider
trace.set_tracer_provider(TracerProvider())
# Add console exporter (for debugging)
trace.get_tracer_provider().add_span_processor(
SimpleSpanProcessor(ConsoleSpanExporter())
)
logger = logging.getLogger(__name__)
logger.info("OpenTelemetry tracing enabled")
def get_logger(name: str) -> logging.Logger:
"""
Get a logger instance with the specified name.
Args:
name: Logger name (usually __name__)
Returns:
Configured logger instance
"""
return logging.getLogger(name)
+237
View File
@@ -0,0 +1,237 @@
"""
Telegram Archiver - Main Application Entry Point.
Provides both FastAPI REST API and CLI interfaces for archiving
Telegram channels to local filesystem bundles.
"""
import asyncio
import time
from contextlib import asynccontextmanager
from pathlib import Path
from typing import AsyncGenerator
from fastapi import FastAPI, HTTPException, BackgroundTasks
from fastapi.responses import JSONResponse
from app.models import (
ArchiveRequest,
ArchiveResponse,
HealthCheck,
)
from app.telethon_client import TelethonArchiver
from app.archiver import ChannelArchiver
from app.logger import setup_logging, get_logger
from config import get_settings, Settings
# Initialize logger
logger = get_logger(__name__)
# Global archiver client (shared across requests)
archiver_client: TelethonArchiver | None = None
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
"""Application lifespan manager."""
global archiver_client
# Startup
logger.info("Starting Telegram Archiver...")
settings = get_settings()
setup_logging(
log_level=settings.log_level,
enable_otel=settings.otel_enabled,
)
archiver_client = TelethonArchiver(settings)
await archiver_client.connect()
logger.info("Telegram Archiver started successfully")
yield
# Shutdown
logger.info("Shutting down Telegram Archiver...")
if archiver_client:
await archiver_client.disconnect()
logger.info("Telegram Archiver shutdown complete")
# Create FastAPI app
app = FastAPI(
title="Telegram Archiver",
description="Archive Telegram channels to local filesystem with media",
version="1.0.0",
lifespan=lifespan,
)
@app.get("/health", response_model=HealthCheck)
async def health_check() -> HealthCheck:
"""Health check endpoint."""
is_connected = archiver_client is not None and archiver_client._connected
return HealthCheck(
status="healthy" if is_connected else "disconnected",
version="1.0.0",
session_active=is_connected,
)
@app.post("/archive", response_model=ArchiveResponse)
async def archive_channel(request: ArchiveRequest) -> ArchiveResponse:
"""
Archive a Telegram channel to local filesystem.
Creates a bundle directory for each message with:
- index.md: Message content in Markdown with front-matter
- Media files: Photos, videos, documents (if under size limit)
Large files (>MAX_FILE_SIZE) are skipped and logged in 2big2get.md
"""
if not archiver_client:
raise HTTPException(status_code=503, detail="Archiver not initialized")
try:
logger.info(f"Archive request for channel: {request.channel}")
# Get settings and override output_dir if specified
settings = get_settings()
if request.output_dir:
settings.output_dir = Path(request.output_dir)
# Create archiver for this request
channel_archiver = ChannelArchiver(archiver_client, settings)
# Run archive process
result = await channel_archiver.archive_channel(
channel=request.channel,
limit=request.limit,
from_message_id=request.from_message_id,
force=request.force,
)
logger.info(f"Archive completed: {result['posts_archived']} posts")
return ArchiveResponse(**result)
except Exception as e:
logger.error(f"Archive failed: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))
@app.get("/")
async def root() -> dict:
"""Root endpoint with API info."""
return {
"name": "Telegram Archiver",
"version": "1.0.0",
"description": "Archive Telegram channels to local filesystem",
"endpoints": {
"health": "GET /health",
"archive": "POST /archive",
"docs": "GET /docs",
},
}
# CLI entry point
def cli_main() -> None:
"""Command-line interface entry point."""
import click
@click.command()
@click.option(
"--channel",
"-c",
required=True,
help="Telegram channel username (with or without @)",
)
@click.option(
"--output-dir",
"-o",
type=click.Path(),
default=None,
help="Output directory for archived channel",
)
@click.option(
"--limit",
"-l",
type=int,
default=None,
help="Limit number of messages to archive",
)
@click.option(
"--from-message-id",
"-f",
type=int,
default=None,
help="Start archiving from this message ID",
)
@click.option(
"--force",
is_flag=True,
help="Force re-download of already archived posts",
)
@click.option(
"--env-file",
type=click.Path(exists=True),
default=".env",
help="Path to .env file",
)
def archive(
channel: str,
output_dir: str | None,
limit: int | None,
from_message_id: int | None,
force: bool,
env_file: str,
) -> None:
"""Archive a Telegram channel to local filesystem."""
# Load settings
settings = get_settings()
if output_dir:
settings.output_dir = Path(output_dir)
# Setup logging
setup_logging(log_level=settings.log_level)
# Run archive
async def run_archive() -> None:
async with TelethonArchiver(settings) as client:
archiver = ChannelArchiver(client, settings)
result = await archiver.archive_channel(
channel=channel,
limit=limit,
from_message_id=from_message_id,
force=force,
)
# Print summary
click.echo("\n" + "=" * 50)
click.echo("Archive Summary")
click.echo("=" * 50)
click.echo(f"Channel: {result['channel_title']} ({result['channel']})")
click.echo(f"Posts archived: {result['posts_archived']}")
click.echo(f"Posts skipped: {result['posts_skipped']}")
click.echo(f"Media downloaded: {result['media_downloaded']}")
click.echo(f"Media skipped (large): {result['media_skipped_large']}")
click.echo(f"Output path: {result['output_path']}")
click.echo(f"Log file: {result['log_file']}")
if result['large_files_report']:
click.echo(f"Large files report: {result['large_files_report']}")
click.echo(f"Duration: {result['duration_seconds']:.2f}s")
click.echo("=" * 50)
asyncio.run(run_archive())
archive()
if __name__ == "__main__":
# For development: run with uvicorn
# uvicorn app.main:app --reload
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
+138
View File
@@ -0,0 +1,138 @@
"""
Pydantic models for Telegram Archiver.
Used for data validation and serialization.
"""
from pydantic import BaseModel, Field
from datetime import datetime
from typing import Optional, Literal
from enum import Enum
class MediaType(str, Enum):
"""Types of media files supported by the archiver."""
PHOTO = "photo"
VIDEO = "video"
DOCUMENT = "document"
AUDIO = "audio"
VOICE = "voice"
VIDEO_NOTE = "video_note"
STICKER = "sticker"
ANIMATION = "animation" # GIF
class MediaFile(BaseModel):
"""Represents a media file attached to a Telegram message."""
filename: str = Field(..., description="Original filename")
type: MediaType = Field(..., description="Type of media")
caption: Optional[str] = Field(None, description="Caption text for the media")
size: Optional[int] = Field(None, description="File size in bytes")
is_too_large: bool = Field(
default=False,
description="True if file exceeds MAX_FILE_SIZE limit"
)
download_url: Optional[str] = Field(
None,
description="Telegram URL for files that weren't downloaded"
)
model_config = {"use_enum_values": True}
class PostData(BaseModel):
"""
Represents a single Telegram post with all metadata.
This is the main data structure passed to the markdown generator.
"""
message_id: int = Field(..., description="Telegram message ID")
date: datetime = Field(..., description="Message timestamp")
text: str = Field(default="", description="Message text in Markdown format")
author: str = Field(..., description="Channel name or author")
channel_username: str = Field(..., description="Channel username (for output path)")
# Relationships
reply_to: Optional[int] = Field(
None,
description="Message ID this post is replying to"
)
repost_from: Optional[int] = Field(
None,
description="Original message ID if this is a forwarded post"
)
repost_channel: Optional[str] = Field(
None,
description="Original channel name if this is a forwarded post"
)
# Media attachments
media_files: list[MediaFile] = Field(
default_factory=list,
description="List of media files attached to the post"
)
# Metadata
views: Optional[int] = Field(None, description="View count (if available)")
has_large_files: bool = Field(
default=False,
description="True if some files were not downloaded due to size limit"
)
@property
def bundle_dir(self) -> str:
"""Return the directory name for this post's bundle."""
return str(self.message_id)
@property
def index_path(self) -> str:
"""Return relative path to index.md within bundle."""
return f"{self.message_id}/index.md"
def has_media(self) -> bool:
"""Check if post has any media files."""
return len(self.media_files) > 0
def has_downloaded_media(self) -> bool:
"""Check if post has any downloaded (not too large) media files."""
return any(not f.is_too_large for f in self.media_files)
class ArchiveRequest(BaseModel):
"""Request model for FastAPI endpoint."""
channel: str = Field(..., description="Telegram channel username or ID")
output_dir: Optional[str] = Field(
None,
description="Override default output directory"
)
limit: Optional[int] = Field(
None,
description="Limit number of posts to archive (for testing)"
)
from_message_id: Optional[int] = Field(
None,
description="Start archiving from this message ID"
)
force: bool = Field(
default=False,
description="Force re-download of already archived posts"
)
class ArchiveResponse(BaseModel):
"""Response model for FastAPI endpoint."""
channel: str
channel_title: str
posts_archived: int
posts_skipped: int
media_downloaded: int
media_skipped_large: int
output_path: str
log_file: str
large_files_report: Optional[str] = None
duration_seconds: float
class HealthCheck(BaseModel):
"""Health check response."""
status: str
version: str
session_active: bool
+347
View File
@@ -0,0 +1,347 @@
"""
Telethon client wrapper for Telegram API interaction.
Handles connection, authentication, and message fetching.
"""
import asyncio
from typing import AsyncGenerator, Optional
from pathlib import Path
import logging
from telethon import TelegramClient
from telethon.sessions import StringSession
from telethon.tl.types import (
Message,
Channel,
Chat,
User,
MessageMediaPhoto,
MessageMediaDocument,
MessageMediaWebPage,
MessageMediaGeo,
MessageMediaPoll,
MessageMediaDice,
MessageMediaInvoice,
MessageMediaContact,
MessageMediaGame,
MessageMediaGeoLive,
MessageMediaVenue,
MessageMediaGiveaway,
MessageMediaUnsupported,
DocumentAttributeFilename,
)
from telethon.errors import SessionPasswordNeededError, FloodWaitError
from config import Settings, get_settings
logger = logging.getLogger(__name__)
class TelethonArchiver:
"""
Wrapper around Telethon client for archiving purposes.
Manages connection, authentication, and message iteration.
"""
def __init__(self, settings: Optional[Settings] = None):
self.settings = settings or get_settings()
self.client: Optional[TelegramClient] = None
self._connected = False
async def connect(self) -> None:
"""
Initialize and connect the Telegram client.
Handles phone number verification and 2FA if needed.
"""
if self._connected:
logger.debug("Already connected to Telegram")
return
logger.info(f"Connecting to Telegram as {self.settings.phone}...")
# Create client with session file
self.client = TelegramClient(
session=self.settings.session_path,
api_id=self.settings.api_id,
api_hash=self.settings.api_hash,
device_model="Telegram Archiver",
app_version="1.0.0",
lang_code="en",
system_lang_code="en",
)
await self.client.connect()
# Check if authorized
if not await self.client.is_user_authorized():
logger.info("Not authorized. Starting authentication...")
await self._authenticate()
self._connected = True
logger.info("Connected to Telegram successfully")
async def _authenticate(self) -> None:
"""Handle the authentication flow."""
try:
# Send code request
await self.client.send_code_request(self.settings.phone)
logger.info(f"Code sent to {self.settings.phone}")
# Get code from user (in MVP, we'll use input())
code = input("Enter the code you received: ")
try:
await self.client.sign_in(
phone=self.settings.phone,
code=code
)
except SessionPasswordNeededError:
# 2FA is enabled
logger.info("2FA enabled. Enter password:")
password = input("2FA Password: ")
await self.client.sign_in(password=password)
except FloodWaitError as e:
logger.error(f"Flood wait: must wait {e.seconds} seconds")
raise
except Exception as e:
logger.error(f"Authentication failed: {e}")
raise
async def disconnect(self) -> None:
"""Disconnect from Telegram."""
if self.client:
await self.client.disconnect()
self._connected = False
logger.info("Disconnected from Telegram")
async def get_channel_info(self, channel: str) -> Channel | Chat | User:
"""
Get channel/chat/user information.
Args:
channel: Channel username (with or without @) or ID
Returns:
Channel, Chat, or User object
"""
if not self._connected or not self.client:
raise RuntimeError("Not connected to Telegram")
entity = await self.client.get_entity(channel)
logger.info(f"Found channel: {getattr(entity, 'title', 'Unknown')}")
return entity
async def fetch_messages(
self,
channel: str,
limit: Optional[int] = None,
from_message_id: Optional[int] = None,
) -> AsyncGenerator[Message, None]:
"""
Fetch messages from a channel.
Args:
channel: Channel username or ID
limit: Maximum number of messages to fetch (None = all)
from_message_id: Start from this message ID (None = latest)
Yields:
Message objects from newest to oldest
"""
if not self._connected or not self.client:
raise RuntimeError("Not connected to Telegram")
entity = await self.get_channel_info(channel)
# Get total message count for progress
total = entity.messages_count if hasattr(entity, 'messages_count') else None
if limit:
total = min(total, limit) if total else limit
logger.info(f"Fetching messages from {getattr(entity, 'title', channel)}...")
if from_message_id:
logger.info(f"Starting from message ID: {from_message_id}")
# Iterate messages (newest first)
message_count = 0
async for message in self.client.iter_messages(
entity,
limit=limit,
min_id=from_message_id, # Messages with ID > from_message_id
):
yield message
message_count += 1
if message_count % 100 == 0:
logger.debug(f"Fetched {message_count} messages...")
logger.info(f"Finished fetching. Total messages: {message_count}")
async def download_media(
self,
message: Message,
output_path: Path,
max_size: Optional[int] = None,
) -> list[dict]:
"""
Download all media from a message.
Args:
message: Telegram message with media
output_path: Directory to save media files
max_size: Maximum file size to download (bytes)
Returns:
List of dicts with file info (filename, type, size, path, is_too_large)
"""
if not message.media:
return []
max_size = max_size or self.settings.max_file_size
media_info = []
# Ensure output directory exists
output_path.mkdir(parents=True, exist_ok=True)
# Handle different media types
if isinstance(message.media, MessageMediaPhoto):
# Photo
info = await self._download_photo(message, output_path, max_size)
if info:
media_info.append(info)
elif isinstance(message.media, MessageMediaDocument):
# Document, video, audio, etc.
info = await self._download_document(message, output_path, max_size)
if info:
media_info.append(info)
# Note: Other media types (geo, poll, etc.) are not downloadable
return media_info
async def _download_photo(
self,
message: Message,
output_path: Path,
max_size: int,
) -> Optional[dict]:
"""Download photo from message."""
try:
# Get file size
photo = message.media.photo
# Get highest resolution
size = max(photo.sizes, key=lambda s: getattr(s, 'w', 0) * getattr(s, 'h', 0))
file_size = getattr(size, 'size', 0)
# Generate filename
timestamp = message.date.strftime("%Y%m%d_%H%M%S")
filename = f"photo_{timestamp}.jpg"
filepath = output_path / filename
info = {
"filename": filename,
"type": "photo",
"size": file_size,
"path": str(filepath),
"is_too_large": file_size > max_size,
"caption": message.message,
}
if file_size > max_size:
logger.debug(f"Photo too large: {file_size} > {max_size}")
return info
# Download
await self.client.download_media(message.photo, filepath)
logger.debug(f"Downloaded photo: {filename}")
return info
except Exception as e:
logger.error(f"Failed to download photo: {e}")
return None
async def _download_document(
self,
message: Message,
output_path: Path,
max_size: int,
) -> Optional[dict]:
"""Download document/video/audio from message."""
try:
doc = message.media.document
file_size = doc.size
# Get original filename
filename = None
for attr in doc.attributes:
if isinstance(attr, DocumentAttributeFilename):
filename = attr.file_name
break
if not filename:
# Generate filename based on type
ext = self._get_extension(doc)
timestamp = message.date.strftime("%Y%m%d_%H%M%S")
filename = f"file_{timestamp}{ext}"
filepath = output_path / filename
# Determine media type
media_type = self._get_media_type(doc)
info = {
"filename": filename,
"type": media_type,
"size": file_size,
"path": str(filepath),
"is_too_large": file_size > max_size,
"caption": message.message,
}
if file_size > max_size:
logger.debug(f"File too large: {file_size} > {max_size}")
return info
# Download
await self.client.download_media(doc, filepath)
logger.debug(f"Downloaded {media_type}: {filename}")
return info
except Exception as e:
logger.error(f"Failed to download document: {e}")
return None
def _get_extension(self, doc) -> str:
"""Get file extension from document attributes."""
for attr in doc.attributes:
if isinstance(attr, DocumentAttributeFilename):
name = attr.file_name
if "." in name:
return "." + name.rsplit(".", 1)[-1]
return ".bin"
def _get_media_type(self, doc) -> str:
"""Determine media type from document attributes."""
mime_type = getattr(doc, 'mime_type', '')
if mime_type.startswith('video/'):
return 'video'
elif mime_type.startswith('audio/'):
return 'audio'
elif mime_type.startswith('image/'):
return 'photo'
elif mime_type == 'application/x-tgsticker':
return 'sticker'
else:
return 'document'
async def __aenter__(self):
"""Async context manager entry."""
await self.connect()
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
"""Async context manager exit."""
await self.disconnect()