Files
copykar/services/publisher.py
T

229 lines
10 KiB
Python

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"⛔ <b>ارسال به کانال مقصد «{target.title}» ممکن نیست!</b>\n\n"
f"• 🆔 شناسه کانال: <code>{target.channel_id}</code>\n"
f"• ❗️ خطا: <code>{type(error).__name__}</code>\n\n"
"<b>ربات</b> در این کانال عضو یا ادمین با دسترسی ارسال پیام نیست.\n"
"لطفا ربات را در کانال <b>ادمین</b> کنید و دسترسی <b>ارسال پیام (Post Messages)</b> بدهید، "
"سپس پست را دوباره به صف بفرستید.\n\n"
"<i>تا رفع این مشکل، پست‌های این کانال ارسال نمی‌شوند و در وضعیت بررسی باقی می‌مانند.</i>"
)
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.")