import os import asyncio import logging from datetime import datetime, timedelta, timezone from typing import Optional, List, Awaitable, Callable, Set, Tuple, Any from telethon import TelegramClient from telethon.errors import ( ChannelPrivateError, ChatAdminRequiredError, ChatWriteForbiddenError, PeerIdInvalidError, UserBannedInChannelError, ) from db.models import TargetChannel from db.repository import Repository from core.queue import RedisQueue from core.metrics import TARGET_ACTIVITY_TOTAL, QUEUE_POSTS_GAUGE from core.proxy import get_telegram_proxy from core.error_logger import log_exception logger = logging.getLogger(__name__) SESSION_DIR = os.getenv("SESSION_DIR", "/app/sessions" if os.path.exists("/app") else "/projects/telegram-bots/copykar/sessions") # Sleep windows are configured in the operator's wall-clock time, not UTC. TIMEZONE_OFFSET_HOURS = float(os.getenv("TIMEZONE_OFFSET_HOURS", "3.5")) # Failures that will never resolve by retrying: the account simply cannot post there. # Re-queueing these would spin the same post through the loop forever. PERMANENT_DELIVERY_ERRORS = ( ChatAdminRequiredError, ChatWriteForbiddenError, ChannelPrivateError, UserBannedInChannelError, PeerIdInvalidError, ) async def send_styled_message(client: TelegramClient, entity: Any, text: str, file_path: Optional[str] = None): """Deliver message or media trying Markdown first, then HTML, then raw plain text.""" if not text and not file_path: return None # 1. Try standard Markdown first (AI's natural formatting) try: if file_path and os.path.exists(file_path): return await client.send_file(entity, file=file_path, caption=text, parse_mode="md") return await client.send_message(entity, text, parse_mode="md") except Exception as md_err: logger.debug(f"Markdown send failed ({md_err}), trying HTML...") # 2. Try HTML format fallback try: if file_path and os.path.exists(file_path): return await client.send_file(entity, file=file_path, caption=text, parse_mode="html") return await client.send_message(entity, text, parse_mode="html") except Exception as html_err: logger.debug(f"HTML send failed ({html_err}), falling back to plain text...") # 3. Plain text safety fallback if file_path and os.path.exists(file_path): return await client.send_file(entity, file=file_path, caption=text, parse_mode=None) return await client.send_message(entity, text, parse_mode=None) class PublisherService: def __init__( self, repo: Repository, queue: RedisQueue, client: Optional[TelegramClient] = None, api_id: Optional[int] = None, api_hash: Optional[str] = None, bot_token: Optional[str] = None, session_name: Optional[str] = None, notify_fn: Optional[Callable[[str], Awaitable[None]]] = None, ): self.repo = repo self.queue = queue self.client = client self.notify_fn = notify_fn # (target_id, error type) pairs already reported, so a broken target is # announced once instead of every polling cycle. self._reported_failures: Set[Tuple[int, str]] = set() self.api_id = api_id or int(os.getenv("API_ID", "0")) self.api_hash = api_hash or os.getenv("API_HASH", "") self.bot_token = bot_token or os.getenv("BOT_TOKEN") self.session_name = session_name or os.path.join(SESSION_DIR, "publisher.session") os.makedirs(os.path.dirname(self.session_name), exist_ok=True) if not self.client: self.client = TelegramClient(self.session_name, self.api_id, self.api_hash, proxy=get_telegram_proxy()) self._running = False self._task: Optional[asyncio.Task] = None async def start(self): logger.info("Starting Paced Target Publisher Service...") if not self.client.is_connected(): # Only owns the connection when it built its own client; a shared client # (the admin bot's) is already connected by its owner. if self.bot_token: await self.client.start(bot_token=self.bot_token) else: await self.client.start() logger.info("Paced Target Publisher Service connected (delivering as the bot account).") self._running = True self._task = asyncio.create_task(self._publisher_loop()) def _is_in_sleep_window(self, target: TargetChannel, current_hour: int) -> bool: if not target.is_sleep_enabled: return False start = target.sleep_start_hour end = target.sleep_end_hour if start == end: return False if start < end: return start <= current_hour < end else: # Overnight sleep (e.g. 23:00 to 08:00) return current_hour >= start or current_hour < end async def _publisher_loop(self): while self._running: try: if await self.repo.is_system_paused(): logger.debug("[publisher] System is paused by admin. Skipping queue processing.") else: await self._process_all_target_queues() except Exception as e: logger.error(f"Error in target publisher loop: {e}", exc_info=True) await asyncio.sleep(15) async def _process_all_target_queues(self): targets = await self.repo.get_active_targets() now = datetime.now(timezone.utc) local_now = now + timedelta(hours=TIMEZONE_OFFSET_HOURS) current_hour_local = local_now.hour for target in targets: qsize = await self.queue.get_target_queue_size(target.id) QUEUE_POSTS_GAUGE.labels(status=f"target_{target.id}").set(qsize) if qsize == 0: continue # 1. Check Sleep Window if self._is_in_sleep_window(target, current_hour_local): logger.debug(f"Target #{target.id} ({target.title}) in sleep window ({target.sleep_start_hour}:00-{target.sleep_end_hour}:00). Skipping.") continue # 2. Check Cooldown Interval if target.last_post_time: last_post = target.last_post_time if last_post.tzinfo is None: last_post = last_post.replace(tzinfo=timezone.utc) diff_minutes = (now - last_post).total_seconds() / 60.0 if diff_minutes < target.post_interval_min: continue # 3. Pop next post payload for this target disp_order = getattr(target, "dispatch_order", "order") or "order" payload = await self.queue.pop_target_post(target.id, dispatch_order=disp_order) if not payload: continue post_id = payload.get("post_id") text = payload.get("text", "") media_path = payload.get("media_path") try: await send_styled_message(self.client, target.channel_id, text, file_path=media_path) # Record publication in database await self.repo.record_post_published_to_target(post_id, target.id, target.title or "Target") await self.repo.update_target_last_post(target.id) TARGET_ACTIVITY_TOTAL.labels(channel_id=str(target.channel_id), title=target.title or '').inc() logger.info(f"Published post ID {post_id} to Target {target.title} ({target.channel_id})") except Exception as e: permanent = isinstance(e, PERMANENT_DELIVERY_ERRORS) if permanent: # Dropping the payload is deliberate: the post stays 'pending_review' # so an admin can re-send it once the permission problem is resolved. await self._report_broken_target(target, e) else: # The payload was already popped; putting it back keeps the post from # being silently lost on a transient Telegram failure. try: await self.queue.push_target_post(target.id, payload) except Exception as requeue_err: logger.critical(f"Failed to requeue post {post_id} for target {target.id}: {requeue_err}") await log_exception("publisher.publish", e, { "post_id": post_id, "target_id": target.id, "channel_id": target.channel_id, "permanent": permanent, }) async def _report_broken_target(self, target: TargetChannel, error: Exception) -> None: """Tell the admins once that a target channel is unreachable for this account.""" key = (target.id, type(error).__name__) if key in self._reported_failures: return self._reported_failures.add(key) logger.error( f"Target #{target.id} ({target.title}) rejected delivery permanently: " f"{type(error).__name__}. Queue drained for this target until it is fixed." ) if not self.notify_fn: return try: await self.notify_fn( f"⛔ ارسال به کانال مقصد «{target.title}» ممکن نیست!\n\n" f"• 🆔 شناسه کانال: {target.channel_id}\n" f"• ❗️ خطا: {type(error).__name__}\n\n" "ربات در این کانال عضو یا ادمین با دسترسی ارسال پیام نیست.\n" "لطفا ربات را در کانال ادمین کنید و دسترسی ارسال پیام (Post Messages) بدهید، " "سپس پست را دوباره به صف بفرستید.\n\n" "تا رفع این مشکل، پست‌های این کانال ارسال نمی‌شوند و در وضعیت بررسی باقی می‌مانند." ) except Exception as notify_err: logger.error(f"Could not notify admins about broken target {target.id}: {notify_err}") async def stop(self): self._running = False if self._task: self._task.cancel() logger.info("Publisher Service stopped.")