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"⛔ ارسال به کانال مقصد «{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.")