Fix duplicate-post and timeout bugs in TG/MAX posting pipeline

- Stop blindly retrying send_photo/send_video/send_media_group and MAX
  send_message on ambiguous network timeouts - a timeout doesn't prove
  the message wasn't delivered, and retrying risked posting duplicates
  (observed live: a video posted 3-5x after repeated timeout retries).
- Raise/rework timeouts that were too short for real large-file transfer
  speeds: TG media upload timeout, and yt-dlp download now uses stall
  detection (killed only on true silence) instead of a flat ceiling that
  was cutting off legitimately slow-but-successful video downloads.
- Fix sendRichMessage: ok=true with an unparseable message_id no longer
  triggers a fallback send (Telegram already created the message).
- MAX: videos over the documented 250MB cap are now sent as a "watch via
  link" note instead of silently failing the upload.
- Add PRAGMA busy_timeout to all DB connections.
- Remove unused MAX_MEDIA_CHANNEL_ID (MAX's /uploads returns a portable
  token directly, no staging channel needed, unlike Telegram).

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-15 12:17:47 +05:00
parent be01558ddb
commit ec65aaf57d
9 changed files with 276 additions and 108 deletions
+6
View File
@@ -21,6 +21,9 @@ TG_CHAT_ID=-1001234567890
TG_MEDIA_CHANNEL_ID= TG_MEDIA_CHANNEL_ID=
# Admin Telegram ID(s) for notifications & reports (single ID or comma-separated "123456,789012") # Admin Telegram ID(s) for notifications & reports (single ID or comma-separated "123456,789012")
TG_ADMIN_IDS=123456789 TG_ADMIN_IDS=123456789
# Timeout (seconds) for send_photo/send_video/send_media_group. Must comfortably
# cover your largest expected video at real upload speed (default 1800 = 30min).
TG_MEDIA_UPLOAD_TIMEOUT_SEC=1800
# ========================================== # ==========================================
# Local Telegram Bot API (Optional) # Local Telegram Bot API (Optional)
@@ -39,6 +42,9 @@ MAX_BOT_TOKEN=your_max_bot_token_here
MAX_CHAT_ID=123456 MAX_CHAT_ID=123456
# MAX API Base URL # MAX API Base URL
MAX_API_BASE_URL=https://platform-api2.max.ru MAX_API_BASE_URL=https://platform-api2.max.ru
# MAX's documented hard cap for a single video attachment. Videos over this are
# sent as a "watch via link" text note instead of failing the upload.
MAX_VIDEO_LIMIT_MB=250
# ========================================== # ==========================================
# Polling, Reporting & Maintenance # Polling, Reporting & Maintenance
+3 -2
View File
@@ -2,9 +2,8 @@ from __future__ import annotations
from loguru import logger from loguru import logger
import asyncio import asyncio
import os
import time import time
from pathlib import Path from typing import Optional
try: try:
from .config import settings from .config import settings
from .media_processor import is_path_active from .media_processor import is_path_active
@@ -17,6 +16,8 @@ async def cleanup_stale_cache(max_age_minutes: Optional[int] = None) -> int:
""" """
Deletes files from the cache directory that are older than max_age_minutes, Deletes files from the cache directory that are older than max_age_minutes,
strictly skipping any files that are currently active in downloads/processing. strictly skipping any files that are currently active in downloads/processing.
Safety net for the immediate per-post cleanup in main.py's process_new_post:
catches anything left behind if the process died mid-post.
""" """
if max_age_minutes is None: if max_age_minutes is None:
max_age_minutes = settings.cache_max_age_minutes max_age_minutes = settings.cache_max_age_minutes
+14 -1
View File
@@ -26,6 +26,10 @@ class Settings(BaseSettings):
tg_media_channel_id: str = "" # Optional storage channel tg_media_channel_id: str = "" # Optional storage channel
tg_admin_ids: str = "" # Comma-separated admin IDs for reports, e.g. "123456,789012" tg_admin_ids: str = "" # Comma-separated admin IDs for reports, e.g. "123456,789012"
local_bot_api_url: str = "" # e.g., "http://127.0.0.1:8081" local_bot_api_url: str = "" # e.g., "http://127.0.0.1:8081"
# Ceiling for send_photo/send_video/send_media_group requests. Must comfortably
# cover a 2GB upload at realistic throughput, not just the observed happy path -
# measured ~3.5MB/s for a 207MB video means 2GB alone can take ~10min.
tg_media_upload_timeout_sec: int = 1800
# NOTE: TELEGRAM_API_ID / TELEGRAM_API_HASH are intentionally not modeled here - # NOTE: TELEGRAM_API_ID / TELEGRAM_API_HASH are intentionally not modeled here -
# they're only consumed by docker-entrypoint.sh (raw env) to start the local # they're only consumed by docker-entrypoint.sh (raw env) to start the local
# telegram-bot-api binary, never read from Python. # telegram-bot-api binary, never read from Python.
@@ -34,6 +38,9 @@ class Settings(BaseSettings):
max_bot_token: str = "" max_bot_token: str = ""
max_chat_id: str = "" # Destination chat ID in MAX max_chat_id: str = "" # Destination chat ID in MAX
max_api_base_url: str = "https://platform-api2.max.ru" max_api_base_url: str = "https://platform-api2.max.ru"
# MAX's documented hard cap for a single video attachment (dev.max.ru/docs-api).
# Videos over this are sent as a text link instead of failing the whole post.
max_video_limit_mb: int = 250
max_video_ready_attempts: int = 6 max_video_ready_attempts: int = 6
max_video_ready_delay_sec: float = 8.0 max_video_ready_delay_sec: float = 8.0
@@ -57,7 +64,13 @@ class Settings(BaseSettings):
video_max_duration_sec: int = 7200 video_max_duration_sec: int = 7200
video_max_height: int = 720 video_max_height: int = 720
media_download_timeout_sec: int = 120 media_download_timeout_sec: int = 120
yt_dlp_timeout_sec: int = 600 # Absolute backstop for the whole download, regardless of progress (guards
# against a pathological slow-trickle that never actually stalls).
yt_dlp_timeout_sec: int = 3600
# Killed only if yt-dlp produces NO progress output for this long - a real
# stall, not just a big/slow file. This is the timeout that actually matters
# day to day; yt_dlp_timeout_sec above is just the outer safety net.
yt_dlp_stall_timeout_sec: int = 120
# Text Styling & Decoration # Text Styling & Decoration
header_text: str = "" header_text: str = ""
+25 -13
View File
@@ -2,8 +2,8 @@ from __future__ import annotations
from loguru import logger from loguru import logger
import json import json
from datetime import datetime from contextlib import asynccontextmanager
from typing import Any, Optional from typing import Any, AsyncIterator, Optional
import aiosqlite import aiosqlite
try: try:
from .config import settings from .config import settings
@@ -11,13 +11,26 @@ except (ImportError, ValueError):
from config import settings from config import settings
# Each method opens its own short-lived connection (no pooling). busy_timeout
# means a second connection opened concurrently (e.g. inspecting the DB by
# hand with sqlite3 while the service runs) waits instead of immediately
# failing with "database is locked".
_BUSY_TIMEOUT_MS = 5000
class Database: class Database:
def __init__(self, db_path: Optional[str] = None) -> None: def __init__(self, db_path: Optional[str] = None) -> None:
self.db_path = str(settings.db_path if db_path is None else db_path) self.db_path = str(settings.db_path if db_path is None else db_path)
@asynccontextmanager
async def _connect(self) -> AsyncIterator[aiosqlite.Connection]:
async with aiosqlite.connect(self.db_path) as db:
await db.execute(f"PRAGMA busy_timeout = {_BUSY_TIMEOUT_MS};")
yield db
async def init(self) -> None: async def init(self) -> None:
logger.info("Initializing database at {}", self.db_path) logger.info("Initializing database at {}", self.db_path)
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
await db.execute("PRAGMA journal_mode=WAL;") await db.execute("PRAGMA journal_mode=WAL;")
await db.execute( await db.execute(
""" """
@@ -58,7 +71,7 @@ class Database:
await db.commit() await db.commit()
async def has_any_posts(self, owner_id: int) -> bool: async def has_any_posts(self, owner_id: int) -> bool:
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
cursor = await db.execute( cursor = await db.execute(
"SELECT 1 FROM posts WHERE vk_owner_id = ? LIMIT 1", "SELECT 1 FROM posts WHERE vk_owner_id = ? LIMIT 1",
(owner_id,), (owner_id,),
@@ -76,7 +89,7 @@ class Database:
reason: str = "bootstrap_initial_skip", reason: str = "bootstrap_initial_skip",
) -> None: ) -> None:
raw_json = json.dumps(raw_data, ensure_ascii=False) raw_json = json.dumps(raw_data, ensure_ascii=False)
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
await db.execute( await db.execute(
""" """
INSERT INTO posts (vk_owner_id, vk_post_id, posted_at, text, raw_json, tg_status, max_status, tg_error, max_error) INSERT INTO posts (vk_owner_id, vk_post_id, posted_at, text, raw_json, tg_status, max_status, tg_error, max_error)
@@ -90,23 +103,22 @@ class Database:
await db.commit() await db.commit()
async def is_post_processed(self, owner_id: int, post_id: int) -> bool: async def is_post_processed(self, owner_id: int, post_id: int) -> bool:
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
db.row_factory = aiosqlite.Row db.row_factory = aiosqlite.Row
cursor = await db.execute( cursor = await db.execute(
"SELECT id, tg_status, max_status FROM posts WHERE vk_owner_id = ? AND vk_post_id = ?", "SELECT tg_status, max_status FROM posts WHERE vk_owner_id = ? AND vk_post_id = ?",
(owner_id, post_id), (owner_id, post_id),
) )
row = await cursor.fetchone() row = await cursor.fetchone()
if not row: if not row:
return False return False
# If already published on both or marked skipped, it's processed
return bool( return bool(
row["tg_status"] in ("published", "skipped") row["tg_status"] in ("published", "skipped")
and row["max_status"] in ("published", "skipped") and row["max_status"] in ("published", "skipped")
) )
async def get_post(self, owner_id: int, post_id: int) -> Optional[dict[str, Any]]: async def get_post(self, owner_id: int, post_id: int) -> Optional[dict[str, Any]]:
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
db.row_factory = aiosqlite.Row db.row_factory = aiosqlite.Row
cursor = await db.execute( cursor = await db.execute(
"SELECT * FROM posts WHERE vk_owner_id = ? AND vk_post_id = ?", "SELECT * FROM posts WHERE vk_owner_id = ? AND vk_post_id = ?",
@@ -124,7 +136,7 @@ class Database:
raw_data: dict[str, Any], raw_data: dict[str, Any],
) -> int: ) -> int:
raw_json = json.dumps(raw_data, ensure_ascii=False) raw_json = json.dumps(raw_data, ensure_ascii=False)
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
cursor = await db.execute( cursor = await db.execute(
""" """
INSERT INTO posts (vk_owner_id, vk_post_id, posted_at, text, raw_json) INSERT INTO posts (vk_owner_id, vk_post_id, posted_at, text, raw_json)
@@ -149,7 +161,7 @@ class Database:
error: Optional[str] = None, error: Optional[str] = None,
) -> None: ) -> None:
msg_str = ",".join(str(m) for m in message_ids) if message_ids else None msg_str = ",".join(str(m) for m in message_ids) if message_ids else None
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
await db.execute( await db.execute(
""" """
UPDATE posts UPDATE posts
@@ -173,7 +185,7 @@ class Database:
error: Optional[str] = None, error: Optional[str] = None,
) -> None: ) -> None:
msg_str = ",".join(message_ids) if message_ids else None msg_str = ",".join(message_ids) if message_ids else None
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
await db.execute( await db.execute(
""" """
UPDATE posts UPDATE posts
@@ -196,7 +208,7 @@ class Database:
status: str = "ok", status: str = "ok",
error: Optional[str] = None, error: Optional[str] = None,
) -> None: ) -> None:
async with aiosqlite.connect(self.db_path) as db: async with self._connect() as db:
await db.execute( await db.execute(
""" """
INSERT INTO publication_runs (found_count, published_tg_count, published_max_count, status, error) INSERT INTO publication_runs (found_count, published_tg_count, published_max_count, status, error)
+10 -19
View File
@@ -1,7 +1,6 @@
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import os
import signal import signal
import sys import sys
from typing import Any, Optional from typing import Any, Optional
@@ -47,15 +46,12 @@ class ServiceApp:
) )
logger.info("Initializing VK to TG & MAX Poster Service...") logger.info("Initializing VK to TG & MAX Poster Service...")
# Initialize SQLite DB
await self.db.init() await self.db.init()
# Initialize Telegram Poster
await self.tg_poster.init() await self.tg_poster.init()
if self.tg_poster.bot: if self.tg_poster.bot:
self.admin_notifier = AdminNotifier(self.tg_poster.bot) self.admin_notifier = AdminNotifier(self.tg_poster.bot)
# Resolve VK Group
if not settings.vk_source: if not settings.vk_source:
raise ValueError("VK_SOURCE is not set in configuration") raise ValueError("VK_SOURCE is not set in configuration")
@@ -66,14 +62,12 @@ class ServiceApp:
self.vk_group_url = f"https://vk.com/{screen_name}" self.vk_group_url = f"https://vk.com/{screen_name}"
logger.info("Resolved VK Group: '{}' (owner_id: {}, url: {})", name, owner_id, self.vk_group_url) logger.info("Resolved VK Group: '{}' (owner_id: {}, url: {})", name, owner_id, self.vk_group_url)
# Start periodic background cache cleaner
self.cleaner_task = asyncio.create_task(run_cleaner_loop(interval_minutes=15)) self.cleaner_task = asyncio.create_task(run_cleaner_loop(interval_minutes=15))
async def process_new_post(self, post: VKPost) -> dict[str, Any]: async def process_new_post(self, post: VKPost) -> dict[str, Any]:
vk_url = f"https://vk.com/wall{post.owner_id}_{post.post_id}" vk_url = f"https://vk.com/wall{post.owner_id}_{post.post_id}"
logger.info("Processing post #{} from {}", post.post_id, vk_url) logger.info("Processing post #{} from {}", post.post_id, vk_url)
# Save to database
post_db_id = await self.db.save_or_update_post( post_db_id = await self.db.save_or_update_post(
owner_id=post.owner_id, owner_id=post.owner_id,
post_id=post.post_id, post_id=post.post_id,
@@ -100,12 +94,13 @@ class ServiceApp:
} }
try: try:
# 1. Download/extract media (skip entirely if both platforms are already done) # Download media once, shared by both platforms (skip entirely if
# both are already done - nothing left to attach).
if post.media and not (tg_done and max_done): if post.media and not (tg_done and max_done):
logger.info("Downloading {} media items for post #{}...", len(post.media), post.post_id) logger.info("Downloading {} media items for post #{}...", len(post.media), post.post_id)
processed_media = await media_processor.process_media_items(post.media) processed_media = await media_processor.process_media_items(post.media)
# 2. Publish to Telegram (only if not already published/skipped for this post) # Telegram - independent of MAX, so a MAX failure never blocks/retries this.
if tg_done: if tg_done:
logger.debug("Post #{} already resolved for Telegram ({}), skipping resend.", post.post_id, existing.get("tg_status")) logger.debug("Post #{} already resolved for Telegram ({}), skipping resend.", post.post_id, existing.get("tg_status"))
else: else:
@@ -130,7 +125,7 @@ class ServiceApp:
result_summary["tg_status"] = "failed" result_summary["tg_status"] = "failed"
result_summary["tg_error"] = err result_summary["tg_error"] = err
# 3. Publish to MAX Messenger (only if not already published/skipped for this post) # MAX - independent of Telegram.
if not (settings.max_bot_token and settings.max_chat_id): if not (settings.max_bot_token and settings.max_chat_id):
if not max_done: if not max_done:
await self.db.update_max_result(post_db_id=post_db_id, status="skipped") await self.db.update_max_result(post_db_id=post_db_id, status="skipped")
@@ -160,7 +155,8 @@ class ServiceApp:
result_summary["max_error"] = err result_summary["max_error"] = err
finally: finally:
# Immediate cleanup of temporary media files # Immediate cleanup of temporary media files. If the process dies
# before this runs, cleaner.py's age-based sweep catches it later.
if processed_media: if processed_media:
await media_processor.cleanup(processed_media) await media_processor.cleanup(processed_media)
@@ -182,7 +178,6 @@ class ServiceApp:
count=settings.vk_check_count, count=settings.vk_check_count,
) )
# Check if this is the very first run on an empty database
is_initial_start = not await self.db.has_any_posts(self.vk_group_owner_id) is_initial_start = not await self.db.has_any_posts(self.vk_group_owner_id)
if is_initial_start and latest_posts: if is_initial_start and latest_posts:
logger.info( logger.info(
@@ -202,7 +197,6 @@ class ServiceApp:
logger.info("Marked {} existing posts as already known. Only new future posts will be published.", len(latest_posts)) logger.info("Marked {} existing posts as already known. Only new future posts will be published.", len(latest_posts))
latest_posts = [] latest_posts = []
elif settings.bootstrap_mode == "publish_latest_one": elif settings.bootstrap_mode == "publish_latest_one":
# Mark all except the single latest post as skipped
newest = max(latest_posts, key=lambda p: p.date) newest = max(latest_posts, key=lambda p: p.date)
for p in latest_posts: for p in latest_posts:
if p.post_id != newest.post_id: if p.post_id != newest.post_id:
@@ -217,7 +211,6 @@ class ServiceApp:
latest_posts = [newest] latest_posts = [newest]
logger.info("Bootstrap mode: keeping only the single newest post #{}", newest.post_id) logger.info("Bootstrap mode: keeping only the single newest post #{}", newest.post_id)
# Filter out processed posts and sort oldest -> newest
for p in latest_posts: for p in latest_posts:
if p.is_repost: if p.is_repost:
logger.debug("Skipping repost #{}", p.post_id) logger.debug("Skipping repost #{}", p.post_id)
@@ -226,7 +219,6 @@ class ServiceApp:
if not is_done: if not is_done:
posts_to_process.append(p) posts_to_process.append(p)
# Sort chronological (oldest to newest)
posts_to_process.sort(key=lambda p: p.date) posts_to_process.sort(key=lambda p: p.date)
if posts_to_process: if posts_to_process:
@@ -234,7 +226,7 @@ class ServiceApp:
for post in posts_to_process: for post in posts_to_process:
report = await self.process_new_post(post) report = await self.process_new_post(post)
published_reports.append(report) published_reports.append(report)
await asyncio.sleep(2.0) # Pause between posts await asyncio.sleep(2.0) # gentle pacing between posts
else: else:
logger.info("No new posts found.") logger.info("No new posts found.")
@@ -258,7 +250,9 @@ class ServiceApp:
error=cycle_error, error=cycle_error,
) )
# Notify admins # Report straight from this cycle's own results/error - no separate
# task reading the DB back, so there's no way for the two to drift
# out of sync (that was the actual bug the worker-split version had).
if self.admin_notifier: if self.admin_notifier:
await self.admin_notifier.notify_cycle_result( await self.admin_notifier.notify_cycle_result(
group_name=self.vk_group_name, group_name=self.vk_group_name,
@@ -278,7 +272,6 @@ class ServiceApp:
except Exception as exc: except Exception as exc:
logger.exception("Unexpected error in main loop: {}", exc) logger.exception("Unexpected error in main loop: {}", exc)
# Sleep between cycles
sleep_seconds = settings.check_interval_minutes * 60 sleep_seconds = settings.check_interval_minutes * 60
logger.debug("Sleeping for {} seconds until next check...", sleep_seconds) logger.debug("Sleeping for {} seconds until next check...", sleep_seconds)
for _ in range(sleep_seconds): for _ in range(sleep_seconds):
@@ -292,7 +285,6 @@ class ServiceApp:
if self.cleaner_task: if self.cleaner_task:
self.cleaner_task.cancel() self.cleaner_task.cancel()
await self.tg_poster.close() await self.tg_poster.close()
# Clean any remaining stale files on shutdown
await cleanup_stale_cache(max_age_minutes=0) await cleanup_stale_cache(max_age_minutes=0)
logger.info("Service stopped cleanly.") logger.info("Service stopped cleanly.")
@@ -309,7 +301,6 @@ async def main() -> None:
try: try:
loop.add_signal_handler(sig, handle_signal) loop.add_signal_handler(sig, handle_signal)
except NotImplementedError: except NotImplementedError:
# Signal handlers not implemented on Windows event loop for non-main threads
pass pass
try: try:
+91 -34
View File
@@ -11,11 +11,11 @@ import aiohttp
try: try:
from .config import settings from .config import settings
from .media_processor import ProcessedMedia from .media_processor import ProcessedMedia
from .text_formatter import build_media_unavailable_note, format_post_text, split_message_chunks from .text_formatter import build_media_unavailable_note, build_oversized_video_note, format_post_text, split_message_chunks
except (ImportError, ValueError): except (ImportError, ValueError):
from config import settings from config import settings
from media_processor import ProcessedMedia from media_processor import ProcessedMedia
from text_formatter import build_media_unavailable_note, format_post_text, split_message_chunks from text_formatter import build_media_unavailable_note, build_oversized_video_note, format_post_text, split_message_chunks
MAX_MESSAGE_LIMIT = 4000 MAX_MESSAGE_LIMIT = 4000
MAX_MEDIA_ITEMS = 10 MAX_MEDIA_ITEMS = 10
@@ -55,7 +55,15 @@ class MAXAPIClient:
if self.session: if self.session:
await self.session.close() await self.session.close()
async def request(self, method: str, path: str, **kwargs: Any) -> dict[str, Any]: async def request(
self, method: str, path: str, retry_on_exception: bool = True, **kwargs: Any
) -> dict[str, Any]:
"""retry_on_exception=False for calls where a timeout/connection error is
ambiguous about whether the server actually processed it (e.g. send_message):
a network exception there does NOT prove nothing was sent, so blindly
retrying risks posting duplicates. HTTP-response-based retries (429,
attachment.not.ready) stay safe regardless - those only fire once we
know the server explicitly rejected the request."""
if not self.session: if not self.session:
raise RuntimeError("MAX session is not initialized") raise RuntimeError("MAX session is not initialized")
url = f"{self.base_url}{path}" url = f"{self.base_url}{path}"
@@ -84,9 +92,11 @@ class MAXAPIClient:
continue continue
if resp.status < 500: if resp.status < 500:
raise MAXAPIError(last_error, resp.status, str(code or "")) raise MAXAPIError(last_error, resp.status, str(code or ""))
except MAXAPIError:
raise
except Exception as exc: except Exception as exc:
last_error = str(exc) last_error = str(exc)
if attempt >= self.max_attempts: if not retry_on_exception or attempt >= self.max_attempts:
raise raise
await asyncio.sleep(min(2 ** attempt, self.retry_backoff_max_sec)) await asyncio.sleep(min(2 ** attempt, self.retry_backoff_max_sec))
raise RuntimeError(last_error or "MAX API request failed") raise RuntimeError(last_error or "MAX API request failed")
@@ -100,7 +110,9 @@ class MAXAPIClient:
payload: dict[str, Any] = {"text": text, "notify": True, "format": "html"} payload: dict[str, Any] = {"text": text, "notify": True, "format": "html"}
if attachments: if attachments:
payload["attachments"] = attachments payload["attachments"] = attachments
return await self.request("POST", f"/messages?chat_id={quote(str(chat_id))}", json=payload) return await self.request(
"POST", f"/messages?chat_id={quote(str(chat_id))}", retry_on_exception=False, json=payload
)
async def get_message(self, message_id: str) -> dict[str, Any]: async def get_message(self, message_id: str) -> dict[str, Any]:
return await self.request("GET", f"/messages/{quote(message_id, safe='')}") return await self.request("GET", f"/messages/{quote(message_id, safe='')}")
@@ -228,27 +240,59 @@ class MAXPoster:
return {"type": media_type, "payload": payload} return {"type": media_type, "payload": payload}
async def upload_media_group( async def upload_media(
self, client: MAXAPIClient, items: list[ProcessedMedia] self, media_items: list[ProcessedMedia]
) -> list[dict[str, Any]]: ) -> tuple[list[dict[str, Any]], list[ProcessedMedia]]:
attachments: list[dict[str, Any]] = [] """Upload stage: push every media item to MAX and poll until all video
for item in items: attachments report ready, so send_post() never blocks on readiness.
try:
att = await self.upload_media_item(client, item)
if att:
attachments.append(att)
except Exception as exc:
logger.warning("MAX media upload failed for {}: {}", item.attachment_id, exc)
if attachments:
await self.wait_for_videos(client, attachments)
return attachments
async def post_to_max( Returns (attachments, oversized_videos): videos over MAX's documented
250MB cap are never even attempted (MAX would just reject them) - they're
returned separately so send_post() can add a "watch via link" note
instead of silently dropping them.
"""
if not self.token or not self.chat_id:
return [], []
valid_media = [m for m in media_items if not m.is_link_only and m.local_path]
if not valid_media:
return [], []
max_video_bytes = settings.max_video_limit_mb * 1024 * 1024
oversized: list[ProcessedMedia] = []
uploadable: list[ProcessedMedia] = []
for item in valid_media:
if item.media_type == "video" and item.size_bytes > max_video_bytes:
logger.info(
"Video {} ({} MB) exceeds MAX's {}MB limit, sending as a link instead",
item.attachment_id, round(item.size_bytes / (1024 * 1024), 1), settings.max_video_limit_mb,
)
oversized.append(item)
else:
uploadable.append(item)
attachments: list[dict[str, Any]] = []
async with MAXAPIClient(self.token, self.api_base_url) as client:
for item in uploadable:
try:
att = await self.upload_media_item(client, item)
if att:
attachments.append(att)
except Exception as exc:
logger.warning("MAX media upload failed for {}: {}", item.attachment_id, exc)
if attachments:
await self.wait_for_videos(client, attachments)
return attachments, oversized
async def send_post(
self, self,
raw_text: str, raw_text: str,
media_items: list[ProcessedMedia], attachments: list[dict[str, Any]],
vk_url: Optional[str] = None, vk_url: Optional[str] = None,
link_only_media: Optional[list[ProcessedMedia]] = None,
oversized_videos: Optional[list[ProcessedMedia]] = None,
) -> tuple[list[str], Optional[str]]: ) -> tuple[list[str], Optional[str]]:
"""Post stage - assumes upload_media() already ran and every video
attachment is confirmed ready."""
if not self.token or not self.chat_id: if not self.token or not self.chat_id:
logger.warning("MAX bot token or chat ID is not set; skipping MAX post.") logger.warning("MAX bot token or chat ID is not set; skipping MAX post.")
return [], None return [], None
@@ -260,18 +304,18 @@ class MAXPoster:
vk_url=vk_url, vk_url=vk_url,
) )
link_only = [m for m in media_items if m.is_link_only] note = build_media_unavailable_note(link_only_media or [], parse_mode="html")
note = build_media_unavailable_note(link_only, parse_mode="html") video_note = build_oversized_video_note(oversized_videos or [], parse_mode="html")
if note: for extra in (note, video_note):
formatted_text = f"{formatted_text}\n\n{note}" if formatted_text else note if extra:
formatted_text = f"{formatted_text}\n\n{extra}" if formatted_text else extra
chunks = split_message_chunks(formatted_text, self.message_limit) chunks = split_message_chunks(formatted_text, self.message_limit)
valid_media = [m for m in media_items if not m.is_link_only and m.local_path]
# Chunk into groups of MAX_MEDIA_ITEMS instead of silently dropping the excess: # Chunk into groups of MAX_MEDIA_ITEMS instead of silently dropping the excess:
# the first group rides with the text message, extra groups go out as follow-ups. # the first group rides with the text message, extra groups go out as follow-ups.
media_groups = ( media_groups = (
[valid_media[i : i + MAX_MEDIA_ITEMS] for i in range(0, len(valid_media), MAX_MEDIA_ITEMS)] [attachments[i : i + MAX_MEDIA_ITEMS] for i in range(0, len(attachments), MAX_MEDIA_ITEMS)]
if valid_media if attachments
else [[]] else [[]]
) )
@@ -279,10 +323,8 @@ class MAXPoster:
first_url: Optional[str] = None first_url: Optional[str] = None
async with MAXAPIClient(self.token, self.api_base_url) as client: async with MAXAPIClient(self.token, self.api_base_url) as client:
first_attachments = await self.upload_media_group(client, media_groups[0])
first_text = chunks[0] if chunks else "" first_text = chunks[0] if chunks else ""
res = await self.send_message_waiting_for_media(client, first_text, first_attachments) res = await self.send_message_waiting_for_media(client, first_text, media_groups[0])
first_mid = self.message_id_from_response(res) first_mid = self.message_id_from_response(res)
if first_mid: if first_mid:
@@ -291,11 +333,10 @@ class MAXPoster:
first_url = self.message_url_from_response(msg_obj) first_url = self.message_url_from_response(msg_obj)
for group in media_groups[1:]: for group in media_groups[1:]:
attachments = await self.upload_media_group(client, group) if not group:
if not attachments:
continue continue
await asyncio.sleep(1.0) await asyncio.sleep(1.0)
sub_res = await self.send_message_waiting_for_media(client, "", attachments) sub_res = await self.send_message_waiting_for_media(client, "", group)
sub_mid = self.message_id_from_response(sub_res) sub_mid = self.message_id_from_response(sub_res)
if sub_mid: if sub_mid:
message_ids.append(sub_mid) message_ids.append(sub_mid)
@@ -309,3 +350,19 @@ class MAXPoster:
logger.info("Sent MAX Message: {}", message_ids) logger.info("Sent MAX Message: {}", message_ids)
return message_ids, first_url return message_ids, first_url
async def post_to_max(
self,
raw_text: str,
media_items: list[ProcessedMedia],
vk_url: Optional[str] = None,
) -> tuple[list[str], Optional[str]]:
"""Back-compat convenience wrapper: upload_media() + send_post() in one call."""
if not self.token or not self.chat_id:
logger.warning("MAX bot token or chat ID is not set; skipping MAX post.")
return [], None
link_only = [m for m in media_items if m.is_link_only]
attachments, oversized = await self.upload_media(media_items)
return await self.send_post(
raw_text, attachments, vk_url=vk_url, link_only_media=link_only, oversized_videos=oversized
)
+34 -16
View File
@@ -132,36 +132,54 @@ class MediaProcessor:
f"/bestvideo[height<={settings.video_max_height}][filesize<{max_size}]+bestaudio" f"/bestvideo[height<={settings.video_max_height}][filesize<{max_size}]+bestaudio"
f"/bestvideo[height<={settings.video_max_height}]+bestaudio" f"/bestvideo[height<={settings.video_max_height}]+bestaudio"
), ),
"--quiet", "--no-warnings", # No --quiet: we need yt-dlp's own progress lines to tell a slow-but-alive
# download (fine, however long it takes) apart from a genuinely hung one -
# a flat wall-clock timeout can't tell those apart and was killing large
# videos that just needed more time (same class of bug as the TG upload
# timeout). --newline makes each progress update its own line instead of
# overwriting via \r, so we can read it with readline().
"--newline", "--no-warnings",
]) ])
proc = None proc = None
stderr = b"" output_lines: list[str] = []
deadline = asyncio.get_running_loop().time() + settings.yt_dlp_timeout_sec
stall_reason: Optional[str] = None
try: try:
proc = await asyncio.create_subprocess_exec( proc = await asyncio.create_subprocess_exec(
*cmd, *cmd,
stdout=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT,
) )
_, stderr = await asyncio.wait_for( while True:
proc.communicate(), timeout=settings.yt_dlp_timeout_sec remaining_total = deadline - asyncio.get_running_loop().time()
) if remaining_total <= 0:
except asyncio.TimeoutError: stall_reason = "video download exceeded overall timeout"
if proc: break
wait_for = min(settings.yt_dlp_stall_timeout_sec, remaining_total)
try: try:
proc.kill() line = await asyncio.wait_for(proc.stdout.readline(), timeout=wait_for)
await proc.communicate() except asyncio.TimeoutError:
except Exception: stall_reason = "video download stalled (no progress from yt-dlp)"
pass break
output_path.unlink(missing_ok=True) if not line:
await unregister_active_path(output_path) break # stdout closed - process is finishing up
return None, "video download timeout", False output_lines.append(line.decode("utf-8", errors="ignore"))
if stall_reason:
proc.kill()
await proc.communicate()
output_path.unlink(missing_ok=True)
await unregister_active_path(output_path)
return None, stall_reason, False
await proc.wait()
finally: finally:
if netrc_path: if netrc_path:
Path(netrc_path).unlink(missing_ok=True) Path(netrc_path).unlink(missing_ok=True)
if proc.returncode != 0 or not output_path.exists() or output_path.stat().st_size == 0: if proc.returncode != 0 or not output_path.exists() or output_path.stat().st_size == 0:
err_msg = (stderr or b"").decode("utf-8", errors="ignore").lower() err_msg = "".join(output_lines).lower()
is_permanent = any( is_permanent = any(
m in err_msg for m in ( m in err_msg for m in (
"removed", "unavailable", "private", "access denied", "does not pass filter", "sign in" "removed", "unavailable", "private", "access denied", "does not pass filter", "sign in"
+23
View File
@@ -176,6 +176,29 @@ def build_media_unavailable_note(link_only_items: list, parse_mode: str = "html"
return header + "\n" + "\n".join(lines) return header + "\n" + "\n".join(lines)
def build_oversized_video_note(items: list, parse_mode: str = "html") -> str:
"""Like build_media_unavailable_note, but for videos that were deliberately
skipped for exceeding a platform's size limit (e.g. MAX's 250MB video cap) -
friendlier tone since this isn't a failure, just a known platform limit."""
if not items:
return ""
lines = []
for item in items:
url = str(getattr(item, "original_url", "") or "").strip()
if not url:
continue
if parse_mode == "html":
lines.append(f'- <a href="{html.escape(url, quote=True)}">видео по ссылке</a>')
else:
lines.append(f"- видео по ссылке: {url}")
if not lines:
return ""
header = "Видео слишком большое для платформы, посмотреть можно здесь:"
if parse_mode == "html":
header = f"<i>{html.escape(header)}</i>"
return header + "\n" + "\n".join(lines)
def format_post_text( def format_post_text(
raw_text: str, raw_text: str,
*, *,
+70 -23
View File
@@ -73,7 +73,8 @@ class TelegramPoster:
if local_url: if local_url:
try: try:
session = AiohttpSession( session = AiohttpSession(
api=TelegramAPIServer.from_base(local_url, is_local=True) api=TelegramAPIServer.from_base(local_url, is_local=True),
timeout=settings.tg_media_upload_timeout_sec,
) )
test_bot = Bot(token=settings.tg_bot_token, session=session) test_bot = Bot(token=settings.tg_bot_token, session=session)
me = await test_bot.get_me() me = await test_bot.get_me()
@@ -86,10 +87,10 @@ class TelegramPoster:
local_url, local_url,
exc, exc,
) )
self.bot = Bot(token=settings.tg_bot_token) self.bot = Bot(token=settings.tg_bot_token, session=AiohttpSession(timeout=settings.tg_media_upload_timeout_sec))
self.is_local_api = False self.is_local_api = False
else: else:
self.bot = Bot(token=settings.tg_bot_token) self.bot = Bot(token=settings.tg_bot_token, session=AiohttpSession(timeout=settings.tg_media_upload_timeout_sec))
self.is_local_api = False self.is_local_api = False
async def close(self) -> None: async def close(self) -> None:
@@ -119,6 +120,21 @@ class TelegramPoster:
await asyncio.sleep(2 ** attempt) await asyncio.sleep(2 ** attempt)
raise RuntimeError("Telegram retries exhausted") raise RuntimeError("Telegram retries exhausted")
async def tg_retry_media(self, fn):
"""Like tg_retry, but for calls that upload actual file bytes (send_photo/
send_video/send_media_group). A network timeout there is NOT proof the
message wasn't delivered - large videos can finish server-side after the
client gives up waiting - so blindly resending risks posting duplicates.
Only flood control (an explicit, safe-to-retry signal) gets a retry;
anything else propagates immediately."""
try:
return await fn()
except TelegramRetryAfter as exc:
delay = float(exc.retry_after) + 0.5
logger.warning("Telegram flood control: retry after {}s", delay)
await asyncio.sleep(delay)
return await fn()
async def set_reaction(self, message_id: int) -> None: async def set_reaction(self, message_id: int) -> None:
if not (settings.tg_auto_reaction_enabled and settings.tg_auto_reaction and self.bot): if not (settings.tg_auto_reaction_enabled and settings.tg_auto_reaction and self.bot):
return return
@@ -158,19 +174,19 @@ class TelegramPoster:
item = valid_items[0] item = valid_items[0]
fs = FSInputFile(str(item.local_path)) fs = FSInputFile(str(item.local_path))
if item.media_type == "photo": if item.media_type == "photo":
msg = await self.tg_retry( msg = await self.tg_retry_media(
lambda: self.bot.send_photo(photo=fs, **self.chat_kwargs(use_storage)) lambda: self.bot.send_photo(photo=fs, **self.chat_kwargs(use_storage))
) )
if msg.photo: if msg.photo:
file_ids[item.attachment_id] = msg.photo[-1].file_id file_ids[item.attachment_id] = msg.photo[-1].file_id
else: else:
msg = await self.tg_retry( msg = await self.tg_retry_media(
lambda: self.bot.send_video(video=fs, **self.chat_kwargs(use_storage)) lambda: self.bot.send_video(video=fs, **self.chat_kwargs(use_storage))
) )
if msg.video: if msg.video:
file_ids[item.attachment_id] = msg.video.file_id file_ids[item.attachment_id] = msg.video.file_id
else: else:
msgs = await self.tg_retry( msgs = await self.tg_retry_media(
lambda: self.bot.send_media_group(media=group, **self.chat_kwargs(use_storage)) lambda: self.bot.send_media_group(media=group, **self.chat_kwargs(use_storage))
) )
for item, msg in zip(valid_items, msgs): for item, msg in zip(valid_items, msgs):
@@ -251,7 +267,14 @@ class TelegramPoster:
mid = res.get("message_id") mid = res.get("message_id")
if mid: if mid:
return [int(mid)] return [int(mid)]
raise RichMessageUnavailable("sendRichMessage returned no message_id") # ok=true means Telegram DID create the message, even though we
# couldn't parse its id from this response. Must NOT raise
# RichMessageUnavailable here - that would trigger send_post's
# fallback and post the same content a second time via the
# standard path, landing one "broken" (unparsed-id) message next
# to a normal one.
logger.warning("sendRichMessage returned ok=true but no parseable message_id: {}", payload)
return []
description = str(payload.get("description") or f"HTTP {resp.status}") description = str(payload.get("description") or f"HTTP {resp.status}")
if "Too Many Requests" in description and isinstance(payload.get("parameters"), dict): if "Too Many Requests" in description and isinstance(payload.get("parameters"), dict):
@@ -300,7 +323,7 @@ class TelegramPoster:
if len(first_group) == 1: if len(first_group) == 1:
item = first_group[0] item = first_group[0]
if item["type"] == "photo": if item["type"] == "photo":
msg = await self.tg_retry( msg = await self.tg_retry_media(
lambda: self.bot.send_photo( lambda: self.bot.send_photo(
photo=item["media"], photo=item["media"],
caption=first_caption or None, caption=first_caption or None,
@@ -309,7 +332,7 @@ class TelegramPoster:
) )
) )
else: else:
msg = await self.tg_retry( msg = await self.tg_retry_media(
lambda: self.bot.send_video( lambda: self.bot.send_video(
video=item["media"], video=item["media"],
caption=first_caption or None, caption=first_caption or None,
@@ -326,7 +349,7 @@ class TelegramPoster:
group.append(InputMediaPhoto(media=item["media"], caption=cap, parse_mode="HTML")) group.append(InputMediaPhoto(media=item["media"], caption=cap, parse_mode="HTML"))
else: else:
group.append(InputMediaVideo(media=item["media"], caption=cap, parse_mode="HTML")) group.append(InputMediaVideo(media=item["media"], caption=cap, parse_mode="HTML"))
msgs = await self.tg_retry( msgs = await self.tg_retry_media(
lambda: self.bot.send_media_group(media=group, **self.chat_kwargs()) lambda: self.bot.send_media_group(media=group, **self.chat_kwargs())
) )
message_ids.extend(int(m.message_id) for m in msgs) message_ids.extend(int(m.message_id) for m in msgs)
@@ -341,7 +364,7 @@ class TelegramPoster:
else InputMediaVideo(media=item["media"]) else InputMediaVideo(media=item["media"])
for item in chunk for item in chunk
] ]
msgs = await self.tg_retry( msgs = await self.tg_retry_media(
lambda: self.bot.send_media_group(media=g, **self.chat_kwargs()) lambda: self.bot.send_media_group(media=g, **self.chat_kwargs())
) )
message_ids.extend(int(m.message_id) for m in msgs) message_ids.extend(int(m.message_id) for m in msgs)
@@ -377,19 +400,36 @@ class TelegramPoster:
mids.append(int(msg.message_id)) mids.append(int(msg.message_id))
return mids return mids
async def post_to_telegram( async def upload_media(self, media_items: list[ProcessedMedia]) -> dict[str, str]:
"""Upload stage: mint reusable file_ids from the storage channel, if configured.
Returns {} when no storage channel is set up - send_post() then falls back to
sending local files directly, same as the pre-split behavior.
"""
valid_media = [m for m in media_items if not m.is_link_only and m.local_path]
if not valid_media or not self.storage_chat_id:
return {}
return await self.upload_media_for_file_ids(valid_media)
async def send_post(
self, self,
raw_text: str, raw_text: str,
media_items: list[ProcessedMedia], media_items: list[ProcessedMedia],
file_ids: dict[str, str],
vk_url: Optional[str] = None, vk_url: Optional[str] = None,
) -> tuple[list[int], Optional[str]]: ) -> tuple[list[int], Optional[str]]:
""" """
Main Telegram posting routine: Post stage - assumes upload_media() already ran (file_ids may be empty if no
storage channel is configured, in which case local files are sent directly):
1. Formats text for HTML parse mode and notes any media that couldn't be attached. 1. Formats text for HTML parse mode and notes any media that couldn't be attached.
2. Uploads media to the storage channel (if configured) to obtain reusable file_ids. 2. If the text fits the platform's normal limit (1024 caption chars with media,
3. Tries sendRichMessage (Bot API 10.1+) for a proper collage + rich text. 4096 message chars without), sends it via the standard, well-tested
4. Falls back to standard aiogram calls (send_photo/send_video/send_media_group) send_photo/send_video/send_media_group/send_message calls directly.
if rich message is unavailable (no storage channel, old Bot API server, etc). 3. Only reaches for sendRichMessage - a non-standard endpoint - when the text is
too long to fit as a single caption, since that's the one thing the standard
path can't do in a single message (it would otherwise split into a media
message + separate follow-up text messages).
4. Falls back to the standard path if sendRichMessage is unavailable/fails.
""" """
formatted_text = format_post_text( formatted_text = format_post_text(
raw_text, raw_text,
@@ -404,13 +444,10 @@ class TelegramPoster:
if note: if note:
formatted_text = f"{formatted_text}\n\n{note}" if formatted_text else note formatted_text = f"{formatted_text}\n\n{note}" if formatted_text else note
# Obtain file_ids from the storage channel, if configured, to avoid re-uploading normal_limit = self.caption_limit if valid_media else self.message_limit
# (also a prerequisite for rich messages, which reference media by file_id). fits_normal_limit = len(formatted_text) <= normal_limit
file_ids: dict[str, str] = {}
if valid_media and self.storage_chat_id:
file_ids = await self.upload_media_for_file_ids(valid_media)
if valid_media and file_ids: if valid_media and file_ids and not fits_normal_limit:
rich_msg = self.build_rich_message(formatted_text, valid_media, file_ids) rich_msg = self.build_rich_message(formatted_text, valid_media, file_ids)
if rich_msg: if rich_msg:
try: try:
@@ -429,3 +466,13 @@ class TelegramPoster:
if mids: if mids:
await self.set_reaction(mids[0]) await self.set_reaction(mids[0])
return mids, url return mids, url
async def post_to_telegram(
self,
raw_text: str,
media_items: list[ProcessedMedia],
vk_url: Optional[str] = None,
) -> tuple[list[int], Optional[str]]:
"""Back-compat convenience wrapper: upload_media() + send_post() in one call."""
file_ids = await self.upload_media(media_items)
return await self.send_post(raw_text, media_items, file_ids, vk_url=vk_url)