Files
copykar/services/publisher.py
T

213 lines
9.3 KiB
Python

import os
import asyncio
import logging
from datetime import datetime, timedelta, timezone
from typing import Optional, List, Awaitable, Callable, Set, Tuple
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,
)
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
payload = await self.queue.pop_target_post(target.id)
if not payload:
continue
post_id = payload.get("post_id")
text = payload.get("text", "")
media_path = payload.get("media_path")
try:
if media_path and os.path.exists(media_path):
await self.client.send_file(
target.channel_id,
file=media_path,
caption=text,
parse_mode="html"
)
else:
await self.client.send_message(
target.channel_id,
text,
parse_mode="html"
)
# 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.")