mirror of
https://gitverse.ru/kpa39l/chronicle.nixg.ru.git
synced 2026-09-29 09:55:08 +00:00
Phase 0 MVP: Telegram Archiver fully functional (992 posts tested)
This commit is contained in:
@@ -5,6 +5,7 @@ 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
|
||||
@@ -19,6 +20,54 @@ 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.
|
||||
@@ -167,52 +216,62 @@ class ChannelArchiver:
|
||||
self.stats["posts_skipped"] += 1
|
||||
return
|
||||
|
||||
logger.debug(f"Processing message {message_id}...")
|
||||
try:
|
||||
logger.debug(f"Processing message {message_id}...")
|
||||
|
||||
# Create bundle directory
|
||||
bundle_dir.mkdir(parents=True, exist_ok=True)
|
||||
# Create bundle directory
|
||||
bundle_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
# Extract post data
|
||||
post_data = await self._extract_post_data(message)
|
||||
# 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),
|
||||
# Download media
|
||||
if message.media:
|
||||
media_info = await self.client.download_media(
|
||||
message, bundle_dir, self.settings.max_file_size, timeout=30
|
||||
)
|
||||
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"),
|
||||
}
|
||||
# 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),
|
||||
)
|
||||
else:
|
||||
self.stats["media_downloaded"] += 1
|
||||
post_data.media_files.append(media_file)
|
||||
|
||||
# Generate markdown
|
||||
markdown_content = self._generate_markdown(post_data)
|
||||
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
|
||||
|
||||
# Write index.md
|
||||
index_path = bundle_dir / "index.md"
|
||||
index_path.write_text(markdown_content, encoding="utf-8")
|
||||
# Generate markdown
|
||||
markdown_content = self._generate_markdown(post_data)
|
||||
|
||||
self.stats["posts_archived"] += 1
|
||||
logger.debug(f"Message {message_id} archived successfully")
|
||||
# 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:
|
||||
"""
|
||||
@@ -230,10 +289,18 @@ class ChannelArchiver:
|
||||
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)
|
||||
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
|
||||
@@ -258,18 +325,20 @@ class ChannelArchiver:
|
||||
return PostData(
|
||||
message_id=message.id,
|
||||
date=message.date,
|
||||
text=text,
|
||||
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
|
||||
@@ -284,6 +353,10 @@ class ChannelArchiver:
|
||||
"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
|
||||
@@ -309,21 +382,31 @@ class ChannelArchiver:
|
||||
# Build markdown content
|
||||
content = []
|
||||
|
||||
# Add media embeds for downloaded files
|
||||
# Add media embeds using Hugo shortcodes
|
||||
for mf in post.media_files:
|
||||
if not mf.is_too_large:
|
||||
if mf.type == "photo":
|
||||
content.append(f"")
|
||||
# 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":
|
||||
content.append(f"@[video]({mf.filename})")
|
||||
# Hugo video shortcode
|
||||
content.append(f'{{{{< video src="{mf.filename}" >}}}}')
|
||||
elif mf.type == "audio":
|
||||
content.append(f"@[audio]({mf.filename})")
|
||||
# Hugo audio shortcode
|
||||
content.append(f'{{{{< audio src="{mf.filename}" >}}}}')
|
||||
elif mf.type == "document":
|
||||
content.append(f"📎 [{mf.filename}]({mf.filename})")
|
||||
# Download link for documents
|
||||
content.append(f'📎 [{mf.filename}]({mf.filename})')
|
||||
|
||||
# Add text content
|
||||
# Add text content (only if not already in media caption)
|
||||
if post.text:
|
||||
content.append(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:
|
||||
|
||||
@@ -8,10 +8,15 @@ 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
|
||||
# OpenTelemetry is optional - skip if not installed
|
||||
try:
|
||||
from opentelemetry import trace
|
||||
from opentelemetry.sdk.trace import TracerProvider
|
||||
from opentelemetry.sdk.trace.export import ConsoleSpanExporter, SimpleSpanProcessor
|
||||
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
|
||||
OTEL_AVAILABLE = True
|
||||
except ImportError:
|
||||
OTEL_AVAILABLE = False
|
||||
|
||||
|
||||
def setup_logging(
|
||||
@@ -62,6 +67,11 @@ def setup_logging(
|
||||
|
||||
def setup_opentelemetry() -> None:
|
||||
"""Configure OpenTelemetry tracing."""
|
||||
if not OTEL_AVAILABLE:
|
||||
logger = logging.getLogger(__name__)
|
||||
logger.warning("OpenTelemetry not available - tracing disabled")
|
||||
return
|
||||
|
||||
# Set up tracer provider
|
||||
trace.set_tracer_provider(TracerProvider())
|
||||
|
||||
|
||||
@@ -19,6 +19,7 @@ class MediaType(str, Enum):
|
||||
VIDEO_NOTE = "video_note"
|
||||
STICKER = "sticker"
|
||||
ANIMATION = "animation" # GIF
|
||||
TIMEOUT = "timeout" # File skipped due to timeout
|
||||
|
||||
|
||||
class MediaFile(BaseModel):
|
||||
@@ -76,6 +77,10 @@ class PostData(BaseModel):
|
||||
default=False,
|
||||
description="True if some files were not downloaded due to size limit"
|
||||
)
|
||||
tags: list[str] = Field(
|
||||
default_factory=list,
|
||||
description="Hashtags extracted from the post"
|
||||
)
|
||||
|
||||
@property
|
||||
def bundle_dir(self) -> str:
|
||||
|
||||
@@ -64,10 +64,11 @@ class TelethonArchiver:
|
||||
session=self.settings.session_path,
|
||||
api_id=self.settings.api_id,
|
||||
api_hash=self.settings.api_hash,
|
||||
device_model="Telegram Archiver",
|
||||
device_model="Desktop (Ubuntu)", # Helps with code delivery
|
||||
system_version="Ubuntu 22.04",
|
||||
app_version="1.0.0",
|
||||
lang_code="en",
|
||||
system_lang_code="en",
|
||||
lang_code="ru", # Russian language
|
||||
system_lang_code="ru",
|
||||
)
|
||||
|
||||
await self.client.connect()
|
||||
@@ -83,26 +84,58 @@ class TelethonArchiver:
|
||||
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}")
|
||||
# Try app code first
|
||||
logger.info("Requesting auth code via Telegram app...")
|
||||
sent_code = await self.client.send_code_request(
|
||||
self.settings.phone,
|
||||
force_sms=False
|
||||
)
|
||||
logger.info(f"Code request sent. Type: {sent_code.type}")
|
||||
|
||||
# Get code from user (in MVP, we'll use input())
|
||||
code = input("Enter the code you received: ")
|
||||
# Get code from user
|
||||
print("\n" + "=" * 50)
|
||||
print("AUTHENTICATION REQUIRED")
|
||||
print("=" * 50)
|
||||
print("Check your Telegram app (NOT SMS) for the code.")
|
||||
print("Look for a message from 'Telegram' with the code.")
|
||||
print("=" * 50)
|
||||
print("\nOptions:")
|
||||
print(" 1. Enter the code from Telegram app")
|
||||
print(" 2. Type 'sms' to receive code via SMS")
|
||||
print(" 3. Press Ctrl+C to cancel")
|
||||
print("=" * 50)
|
||||
|
||||
user_input = input("\nEnter code (or 'sms'): ").strip()
|
||||
|
||||
# If user requests SMS
|
||||
if user_input.lower() == 'sms':
|
||||
logger.info("Requesting SMS code...")
|
||||
print("\n📱 Requesting SMS code...")
|
||||
sent_code = await self.client.send_code_request(
|
||||
self.settings.phone,
|
||||
force_sms=True
|
||||
)
|
||||
print("SMS sent! Enter the code:")
|
||||
user_input = input("SMS Code: ").strip()
|
||||
|
||||
code = user_input
|
||||
|
||||
try:
|
||||
await self.client.sign_in(
|
||||
phone=self.settings.phone,
|
||||
code=code
|
||||
code=code,
|
||||
phone_code_hash=sent_code.phone_code_hash
|
||||
)
|
||||
except SessionPasswordNeededError:
|
||||
# 2FA is enabled
|
||||
logger.info("2FA enabled. Enter password:")
|
||||
print("\n🔐 2FA enabled")
|
||||
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")
|
||||
print(f"\n⚠️ You must wait {e.seconds} seconds before trying again.")
|
||||
raise
|
||||
except Exception as e:
|
||||
logger.error(f"Authentication failed: {e}")
|
||||
@@ -165,24 +198,33 @@ class TelethonArchiver:
|
||||
|
||||
# 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
|
||||
error_count = 0
|
||||
|
||||
# Build iter_messages kwargs
|
||||
iter_kwargs = {"limit": limit}
|
||||
if from_message_id:
|
||||
iter_kwargs["min_id"] = from_message_id
|
||||
|
||||
async for message in self.client.iter_messages(entity, **iter_kwargs):
|
||||
try:
|
||||
yield message
|
||||
message_count += 1
|
||||
|
||||
if message_count % 100 == 0:
|
||||
logger.debug(f"Fetched {message_count} messages...")
|
||||
if message_count % 50 == 0:
|
||||
logger.info(f"Progress: {message_count} messages fetched...")
|
||||
except Exception as e:
|
||||
error_count += 1
|
||||
logger.error(f"Error processing message {message.id}: {e}")
|
||||
continue # Skip problematic messages
|
||||
|
||||
logger.info(f"Finished fetching. Total messages: {message_count}")
|
||||
logger.info(f"Finished fetching. Total: {message_count}, Errors: {error_count}")
|
||||
|
||||
async def download_media(
|
||||
self,
|
||||
message: Message,
|
||||
output_path: Path,
|
||||
max_size: Optional[int] = None,
|
||||
timeout: int = 60, # Timeout in seconds
|
||||
) -> list[dict]:
|
||||
"""
|
||||
Download all media from a message.
|
||||
@@ -191,6 +233,7 @@ class TelethonArchiver:
|
||||
message: Telegram message with media
|
||||
output_path: Directory to save media files
|
||||
max_size: Maximum file size to download (bytes)
|
||||
timeout: Timeout for each file download in seconds
|
||||
|
||||
Returns:
|
||||
List of dicts with file info (filename, type, size, path, is_too_large)
|
||||
@@ -204,20 +247,40 @@ class TelethonArchiver:
|
||||
# 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)
|
||||
try:
|
||||
# Handle different media types with timeout
|
||||
if isinstance(message.media, MessageMediaPhoto):
|
||||
# Photo
|
||||
info = await asyncio.wait_for(
|
||||
self._download_photo(message, output_path, max_size),
|
||||
timeout=timeout
|
||||
)
|
||||
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)
|
||||
elif isinstance(message.media, MessageMediaDocument):
|
||||
# Document, video, audio, etc.
|
||||
info = await asyncio.wait_for(
|
||||
self._download_document(message, output_path, max_size),
|
||||
timeout=timeout
|
||||
)
|
||||
if info:
|
||||
media_info.append(info)
|
||||
|
||||
# Note: Other media types (geo, poll, etc.) are not downloadable
|
||||
except asyncio.TimeoutError:
|
||||
logger.warning(f"Timeout downloading media for message {message.id}")
|
||||
# Add info about skipped file
|
||||
media_info.append({
|
||||
"filename": "timeout_skipped.bin",
|
||||
"type": "timeout",
|
||||
"caption": None,
|
||||
"size": 0,
|
||||
"path": None,
|
||||
"is_too_large": False,
|
||||
"timeout": True,
|
||||
})
|
||||
except Exception as e:
|
||||
logger.error(f"Error downloading media for message {message.id}: {e}")
|
||||
|
||||
return media_info
|
||||
|
||||
|
||||
Reference in New Issue
Block a user