Files

458 lines
15 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
Core archiver logic.
Handles message processing, markdown generation, and file organization.
"""
import asyncio
import logging
import re
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__)
def extract_hashtags(text: str) -> list[str]:
"""
Extract hashtags from text.
Matches: #tag, #тег, #tag123, #тэг_с_подчёркиванием
Supports: Latin, Cyrillic, digits, underscores
Args:
text: Message text
Returns:
List of unique hashtags (without #)
"""
if not text:
return []
# Regex for hashtags: # followed by word chars (including Cyrillic)
pattern = r'#([\wа-яА-ЯёЁ\d_]+)'
matches = re.findall(pattern, text)
# Return unique tags, lowercase
return list(set(tag.lower() for tag in matches))
def remove_hashtags(text: str) -> str:
"""
Remove hashtags from text.
Args:
text: Message text with hashtags
Returns:
Text without hashtags (extra whitespace cleaned)
"""
if not text:
return text
# Remove hashtags
pattern = r'#\w+[\wа-яА-ЯёЁ\d_]*'
text = re.sub(pattern, '', text)
# Clean up multiple spaces/newlines
text = re.sub(r'\n\s*\n', '\n\n', text)
text = re.sub(r' {2,}', ' ', text)
return text.strip()
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
try:
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, timeout=30
)
# 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
# Log progress every 50 posts
if self.stats["posts_archived"] % 50 == 0:
logger.info(f"Progress: {self.stats['posts_archived']} posts archived...")
logger.debug(f"Message {message_id} archived successfully")
except Exception as e:
logger.error(f"Failed to process message {message_id}: {e}")
self.stats["posts_skipped"] += 1
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
original_text = message.message or ""
if message.message and message.entities:
# Use Telethon's markdown parser with entities
text = message.text # Plain text
else:
text = original_text
# Extract hashtags from text
tags = extract_hashtags(text)
# Remove hashtags from text (clean up)
clean_text = remove_hashtags(text)
# 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=clean_text, # Use cleaned text (without hashtags)
author=channel_title,
channel_username=channel_username,
reply_to=reply_to,
repost_from=repost_from,
repost_channel=repost_channel,
views=getattr(message, "views", None),
tags=tags, # Extracted hashtags
)
def _generate_markdown(self, post: PostData) -> str:
"""
Generate markdown content with front-matter for Hugo.
Uses Hugo shortcodes for media embedding.
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 tags if present
if post.tags:
frontmatter["tags"] = post.tags
# 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 using Hugo shortcodes
for mf in post.media_files:
if not mf.is_too_large:
if mf.type == "photo":
# Hugo figure shortcode
if mf.caption:
content.append(f'{{{{< figure src="{mf.filename}" alt="{mf.caption}" title="{mf.caption}" >}}}}')
else:
content.append(f'{{{{< figure src="{mf.filename}" >}}}}')
elif mf.type == "video":
# Hugo video shortcode
content.append(f'{{{{< video src="{mf.filename}" >}}}}')
elif mf.type == "audio":
# Hugo audio shortcode
content.append(f'{{{{< audio src="{mf.filename}" >}}}}')
elif mf.type == "document":
# Download link for documents
content.append(f'📎 [{mf.filename}]({mf.filename})')
# Add text content (only if not already in media caption)
if post.text:
# Check if text is already used as caption in media
captions = [mf.caption for mf in post.media_files if mf.caption]
if post.text not in captions:
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)