diff --git a/services/admin_bot.py b/services/admin_bot.py index cef4ee2..33fa738 100644 --- a/services/admin_bot.py +++ b/services/admin_bot.py @@ -1,5 +1,9 @@ import os import re +import io +import time +import json +import asyncio import logging from datetime import datetime, timezone from typing import Optional, List, Dict, Any, Tuple @@ -8,22 +12,64 @@ from telethon.tl.types import PeerChannel, PeerChat, PeerUser, Channel, Chat from db.models import Post, TargetChannel, SourceChannel from db.repository import Repository from core.queue import RedisQueue -from core.metrics import ADMIN_ACTIONS_TOTAL, TARGET_ACTIVITY_TOTAL +from core.llm import LLMClient +from core.metrics import ( + ADMIN_ACTIONS_TOTAL, TARGET_ACTIVITY_TOTAL, AUTO_ROUTED_POSTS_TOTAL, + ERRORS_RESOLVED_TOTAL, ERRORS_OPEN_GAUGE, ERRORS_OPEN_TOTAL_GAUGE, +) from core.proxy import get_telegram_proxy from core.error_logger import log_exception +from services.collector import MIN_FETCH_LIMIT, MAX_FETCH_LIMIT +from services.ai_processor import SUPPORTED_LANGUAGES +from services.metrics_reporter import get_instant_metrics_report, generate_metrics_graph logger = logging.getLogger(__name__) SESSION_DIR = os.getenv("SESSION_DIR", "/app/sessions" if os.path.exists("/app") else "/projects/telegram-bots/copykar/sessions") -def get_persian_main_menu(): +# Telegram hard limits. Exceeding either raises MediaCaptionTooLongError / MessageTooLongError. +MAX_MEDIA_CAPTION = 1024 +MAX_TEXT_MESSAGE = 4096 +TRUNCATION_MARK = "\n\n… (متن طولانی بود و کوتاه شد)" + + +def clamp_for_telegram(text: str, has_media: bool) -> str: + """Trim text to what Telegram will accept for this message type.""" + limit = MAX_MEDIA_CAPTION if has_media else MAX_TEXT_MESSAGE + if len(text) <= limit: + return text + return text[: limit - len(TRUNCATION_MARK)].rstrip() + TRUNCATION_MARK + + +def get_persian_main_menu(is_paused: bool = False): + pause_btn = "▶️ راه‌اندازی و ادامه سیستم" if is_paused else "🛑 توقف اضطراری سیستم" return [ [Button.text("📊 آمار و وضعیت ناوگان", resize=True), Button.text("🔑 درخواست کد لاگین", resize=True)], [Button.text("📡 کانال‌های مبدا", resize=True), Button.text("🎯 کانال‌های مقصد", resize=True)], + [Button.text("📂 دسته‌بندی کانال‌ها", resize=True), Button.text("🧠 تنظیمات و لاگ‌های AI", resize=True)], [Button.text("➕ افزودن کانال مبدا", resize=True), Button.text("➕ افزودن کانال مقصد", resize=True)], - [Button.text("❓ راهنمای سیستم", resize=True)] + [Button.text("📨 ارسال پست‌های بررسی‌نشده", resize=True), Button.text("⚠️ خطاهای سیستم", resize=True)], + [Button.text(pause_btn, resize=True), Button.text("❓ راهنمای سیستم", resize=True)] ] + +def get_source_fetch_buttons(channel_id: int, source_id: Optional[int] = None): + rows = [ + [ + Button.inline("📥 ۲۰ پست", data=f"hist:{channel_id}:20"), + Button.inline("📥 ۵۰ پست", data=f"hist:{channel_id}:50"), + Button.inline("🔢 تعداد دلخواه", data=f"histn:{channel_id}"), + ] + ] + if source_id is not None: + rows.append([ + Button.inline("📁 تعیین دسته‌بندی", data=f"src_cat:{source_id}"), + Button.inline("🗑 حذف این کانال مبدا", data=f"del_src:{source_id}") + ]) + return rows + + + def get_cancel_button(): return [[Button.inline("❌ انصراف و بازگشت به منو", data="cancel_state")]] @@ -56,6 +102,11 @@ class AdminBotService: self.preview_cache: Dict[str, str] = {} # User step-by-step state: {user_id: {"action": "...", "target_id": ...}} self.user_states: Dict[int, Dict[str, Any]] = {} + # Callback data is limited to 64 bytes, so error groups are referenced by + # index into this cache rather than by their service/type strings. + self.error_group_cache: Dict[int, Tuple[str, str]] = {} + self._metrics_task: Optional[Any] = None + self._running = False def set_collector(self, collector): self.collector = collector @@ -69,6 +120,33 @@ class AdminBotService: def is_admin(self, user_id: int) -> bool: return not self.admin_user_ids or user_id in self.admin_user_ids + async def get_menu(self): + is_paused = await self.repo.is_system_paused() + return get_persian_main_menu(is_paused=is_paused) + + async def on_ai_fallback_alert(self, from_profile: Any, to_profile: Any, error_msg: str, step: int): + clean_err = (error_msg or "")[:250].replace("<", "<").replace(">", ">") + text = ( + f"⚠️ هشدار جابجایی خودکار هوش مصنوعی (AI Fallback)\n\n" + f"❌ ارائه‌دهنده با خطا: «{getattr(from_profile, 'name', 'اصلی')}» ({getattr(from_profile, 'provider_type', '')} / {getattr(from_profile, 'model', '')})\n" + f"🔍 علت خطا: {clean_err}\n\n" + f"🔄 سوییچ خودکار به ارائه‌دهنده پشتیبان:\n" + f"✅ «{getattr(to_profile, 'name', 'پشتیبان')}» ({getattr(to_profile, 'provider_type', '')} / {getattr(to_profile, 'model', '')}) — (گام {step} در زنجیره)" + ) + await self.notify_admins(text) + + async def on_ai_chain_failure_alert(self, chain: List[Any], last_error_msg: str): + chain_names = " ➔ ".join([f"«{getattr(p, 'name', 'AI')}» ({getattr(p, 'model', '')})" for p in chain]) + clean_err = (last_error_msg or "")[:300].replace("<", "<").replace(">", ">") + text = ( + f"🚨 خطای بحرانی: قطعی تمامی ارائه‌دهندگان هوش مصنوعی (AI Chain Failure)\n\n" + f"تمامی {len(chain)} ارائه‌دهنده در زنجیره پشتیبان با خطا مواجه شدند:\n" + f"🔗 مسیر زنجیره: {chain_names}\n\n" + f"⚠️ آخرین خطا: {clean_err}\n\n" + f"لطفا وضعیت ارائه‌دهنده‌ها، اعتبار API یا پراکسی را در منوی «🧠 تنظیمات و لاگ‌های AI» بررسی کنید." + ) + await self.notify_admins(text) + async def notify_admins(self, text: str): """Broadcast Persian message to review channel and admin DMs.""" if self.review_channel_id: @@ -77,17 +155,21 @@ class AdminBotService: except Exception as e: logger.error(f"Failed to notify review channel: {e}") + menu = await self.get_menu() for admin_id in self.admin_user_ids: try: - await self.client.send_message(admin_id, text, parse_mode="html", buttons=get_persian_main_menu()) + await self.client.send_message(admin_id, text, parse_mode="html", buttons=menu) except Exception as e: logger.debug(f"Could not send DM to admin {admin_id}: {e}") + async def start(self): logger.info("Starting Admin Bot Service...") await self.client.start(bot_token=self.bot_token) logger.info("Admin Bot connected successfully.") self._register_handlers() + self._running = True + self._metrics_task = asyncio.create_task(self._error_metrics_loop()) # --- Helper to resolve channel from forwarded message or text --- async def _resolve_channel(self, event: events.NewMessage.Event) -> Tuple[Optional[int], Optional[str], Optional[str], str]: @@ -147,15 +229,38 @@ class AdminBotService: qsize = await self.queue.get_target_queue_size(target.id) if self.queue else 0 sleep_st = f"🌙 فعال ({target.sleep_start_hour}:00 تا {target.sleep_end_hour}:00)" if target.is_sleep_enabled else "☀️ غیرفعال" + auto_ids = list(target.auto_source_ids or []) + if auto_ids: + all_sources = {s.channel_id: s for s in await self.repo.get_active_sources()} + names = [(all_sources[cid].title or str(cid)) if cid in all_sources else str(cid) for cid in auto_ids] + auto_st = "🤖 فعال از: " + "، ".join(f"{n}" for n in names) + else: + auto_st = "🛑 غیرفعال (فقط با تایید دستی ادمین)" + + lang = getattr(target, "language", "fa") or "fa" + lang_label = SUPPORTED_LANGUAGES.get(lang, {}).get("label", "فارسی 🇮🇷") + + cat_name = "بدون دسته‌بندی" + if getattr(target, "category_id", None): + cat = await self.repo.get_category_by_id(target.category_id) + if cat: + cat_name = cat.name + + custom_p = getattr(target, "custom_prompt", "") or "" card_text = ( f"🎯 تنظیمات کانال مقصد: {target.title}\n\n" f"• 🆔 شناسه کانال: {target.channel_id} (ID: {target.id})\n" f"• 🔗 یوزرنیم: @{target.username or 'ندارد'}\n" + f"• 📁 دسته‌بندی: {cat_name}\n" + f"• 🌐 زبان کانال: {lang_label}\n" f"• ⏱ فاصله ارسال پست‌ها: هر {target.post_interval_min} دقیقه\n" f"• 🌙 وضعیت ساعت خواب: {sleep_st}\n" - f"• 📥 پست‌های منتظر در صف: {qsize} پست\n\n" + f"• 📥 پست‌های منتظر در صف: {qsize} پست\n" + f"• 🤖 ارسال خودکار: {auto_st}\n\n" f"🎭 شخصیت و لحن نگارش:\n" - f"{target.personality or 'پیش‌فرض (رسمی و روان)'}\n\n" + f"{target.personality or 'پیش‌فرض'}\n\n" + f"📝 دستورالعمل و فرامین اختصاصی (Custom Prompt):\n" + f"{custom_p or 'ثبت نشده (پیش‌فرض)'}\n\n" f"🏷 فوتر و هشتگ‌های اختصاصی:\n" f"{target.custom_footer or 'ندارد'}\n\n" "👇 برای ویرایش هر بخش روی دکمه مربوطه کلیک کنید:" @@ -164,11 +269,19 @@ class AdminBotService: buttons = [ [ Button.inline("🎭 تغییر لحن و استایل", data=f"st_pers:{target.id}"), + Button.inline("📝 پرامپت و فرامین اختصاصی", data=f"st_cprompt:{target.id}") + ], + [ + Button.inline("📁 تعیین دسته‌بندی", data=f"trg_cat:{target.id}"), Button.inline("🏷 تغییر فوتر و تگ‌ها", data=f"st_foot:{target.id}") ], [ - Button.inline("⏱ تغییر فاصله ارسال", data=f"st_intv:{target.id}"), - Button.inline("🌙 تنظیم ساعت خواب", data=f"st_slp:{target.id}") + Button.inline("🌐 تنظیم زبان", data=f"st_lang:{target.id}"), + Button.inline("⏱ تغییر فاصله ارسال", data=f"st_intv:{target.id}") + ], + [ + Button.inline("🌙 تنظیم ساعت خواب", data=f"st_slp:{target.id}"), + Button.inline("🤖 ارسال خودکار از مبدا", data=f"auto_src:{target.id}") ] ] if target.is_sleep_enabled: @@ -180,6 +293,242 @@ class AdminBotService: ]) return card_text, buttons + async def _render_target_language_menu(self, target_id: int) -> Tuple[str, List[List[Button]]]: + target = await self.repo.get_target_by_id(target_id) + if not target: + return "❌ کانال مقصد یافت نشد.", [] + + current_lang = getattr(target, "language", "fa") or "fa" + text = ( + f"🌐 انتخاب زبان برای کانال «{target.title}»\n\n" + "هوش مصنوعی پست‌ها را متناسب با این زبان بازنویسی و ترجمه می‌کند:\n" + f"زبان فعلی: {SUPPORTED_LANGUAGES.get(current_lang, {}).get('label', 'فارسی 🇮🇷')}" + ) + buttons = [] + for code, info in SUPPORTED_LANGUAGES.items(): + mark = " ✅" if code == current_lang else "" + buttons.append([Button.inline(f"{info['label']}{mark}", data=f"set_lang:{target_id}:{code}")]) + buttons.append([Button.inline("🔙 بازگشت به تنظیمات کانال", data=f"trg_view:{target_id}")]) + return text, buttons + + # --- AI Management & Multi-Provider Profiles --- + async def _render_ai_menu(self) -> Tuple[str, List[List[Button]]]: + active_prof = await self.repo.get_active_provider_profile() + if active_prof: + provider = active_prof.provider_type + model = active_prof.model + base_url = active_prof.base_url + api_key = active_prof.api_key + reasoning = active_prof.reasoning_effort + prof_name = active_prof.name + else: + provider = await self.repo.get_setting("ai_provider", os.getenv("AI_PROVIDER", "openai")) + model = await self.repo.get_setting("ai_model", os.getenv("AI_MODEL", "antigravity" if provider == "agy" else "orcarouter/auto")) + base_url = await self.repo.get_setting("ai_base_url", os.getenv("AI_BASE_URL", "")) + api_key = await self.repo.get_setting("ai_api_key", os.getenv("AI_API_KEY", "")) + reasoning = await self.repo.get_setting("ai_reasoning_effort", os.getenv("AI_REASONING_EFFORT", "")) + prof_name = "پیش‌فرض سیستم" + + double_check = await self.repo.get_setting("ai_double_check", os.getenv("AI_DOUBLE_CHECK", "false")) + dc_status = "✅ فعال (بررسی مجدد)" if double_check.lower() in ("true", "1", "yes") else "❌ غیرفعال" + reason_label = reasoning.upper() if reasoning else "پیش‌فرض" + masked_key = f"{api_key[:6]}...{api_key[-4:]}" if len(api_key) > 10 else ("تنظیم نشده" if not api_key else "******") + url_label = base_url if base_url else "پیش‌فرض (لوکال/داخلی)" + + text = ( + "🧠 تنظیمات هوش مصنوعی (Multi-Provider AI System):\n\n" + f"• 🔌 پروفایل فعال: {prof_name} ({provider})\n" + f"• 🏷 مدل فعال: {model}\n" + f"• 🌐 Base URL: {url_label}\n" + f"• 🔑 API Key: {masked_key}\n" + f"• 🔍 بررسی مجدد (Double Check): {dc_status}\n" + f"• 🧠 میزان تفکر: {reason_label}\n\n" + "👇 برای تغییر پروفایل فعال، ثبت ارائه‌دهنده جدید یا تست روی دکمه‌ها بزنید:" + ) + + buttons = [ + [ + Button.inline("🔌 لیست و تغییر سرویس‌دهنده‌ها", data="ai_prov_list"), + Button.inline("➕ ثبت سرویس‌دهنده جدید", data="ai_add_prov"), + ], + [ + Button.inline(f"🔍 بررسی مجدد ({'فعال' if double_check.lower() in ('true', '1', 'yes') else 'غیرفعال'})", data="ai_toggle_dc"), + Button.inline("🧪 تست سریع اتصال AI", data="ai_test"), + ], + [ + Button.inline("📜 مشاهده ۱۰ لاگ اخیر AI", data="ai_logs:10") + ] + ] + return text, buttons + + async def _render_ai_providers_list(self) -> Tuple[str, List[List[Button]]]: + profiles = await self.repo.get_provider_profiles() + text = ( + "🔌 لیست سرویس‌دهنده‌های هوش مصنوعی (AI Providers):\n\n" + "برای مشاهده جزئیات، ویرایش، تست یا فعال‌سازی روی هر گزینه کلیک کنید:\n" + ) + buttons = [] + for p in profiles: + st = "🟢 [فعال]" if p.is_active else "⚪️" + label = f"{st} {p.name} ({p.model})" + buttons.append([Button.inline(label, data=f"ai_pv_view:{p.id}")]) + + buttons.append([ + Button.inline("➕ ثبت سرویس‌دهنده جدید", data="ai_add_prov"), + Button.inline("🔙 بازگشت به منوی AI", data="ai_menu"), + ]) + return text, buttons + + async def _render_ai_provider_detail(self, profile_id: int) -> Tuple[str, List[List[Button]]]: + p = await self.repo.get_provider_profile_by_id(profile_id) + if not p: + return "❌ پروفایل سرویس‌دهنده یافت نشد.", [[Button.inline("🔙 بازگشت", data="ai_prov_list")]] + + st_text = "🟢 فعال (در حال استفاده)" if p.is_active else "⚪️ غیرفعال" + masked_key = f"{p.api_key[:6]}...{p.api_key[-4:]}" if len(p.api_key) > 10 else ("تنظیم نشده" if not p.api_key else "******") + url_label = p.base_url if p.base_url else "پیش‌فرض (لوکال)" + fallback_label = f"🔄 {p.fallback_provider_name}" if p.fallback_provider_name else "❌ بدون پشتیبان (پایان زنجیره)" + + text = ( + f"🔌 اطلاعات سرویس‌دهنده: {p.name}\n\n" + f"• 🏷 نوع سرویس (Type): {p.provider_type}\n" + f"• 🤖 مدل (Model): {p.model}\n" + f"• 🌐 Base URL: {url_label}\n" + f"• 🔑 API Key: {masked_key}\n" + f"• 🧠 تفکر (Reasoning): {p.reasoning_effort or 'پیش‌فرض'}\n" + f"• 🔄 پشتیبان دستی (Fallback): {fallback_label}\n" + f"• ⚡️ وضعیت: {st_text}" + ) + + buttons = [] + if not p.is_active: + buttons.append([Button.inline("⚡️ فعال‌سازی این سرویس‌دهنده", data=f"ai_pv_activate:{p.id}")]) + + buttons.append([ + Button.inline("🧪 تست اختصاصی این Provider", data=f"ai_pv_test:{p.id}"), + Button.inline("🔄 تنظیم ارائه‌دهنده پشتیبان", data=f"ai_pv_fb_menu:{p.id}"), + ]) + buttons.append([ + Button.inline("🧠 حالت تفکر و استدلال", data=f"ai_pv_effort:{p.id}"), + Button.inline("🏷 تغییر مدل", data=f"ai_pv_model:{p.id}"), + ]) + buttons.append([ + Button.inline("🌐 تغییر Base URL", data=f"ai_pv_url:{p.id}"), + Button.inline("🔑 تغییر API Key", data=f"ai_pv_key:{p.id}"), + ]) + if not p.is_active: + buttons.append([Button.inline("🗑 حذف این سرویس‌دهنده", data=f"ai_pv_del:{p.id}")]) + + buttons.append([Button.inline("🔙 بازگشت به لیست سرویس‌دهنده‌ها", data="ai_prov_list")]) + return text, buttons + + async def _render_ai_provider_fallback_menu(self, profile_id: int) -> Tuple[str, List[List[Button]]]: + p = await self.repo.get_provider_profile_by_id(profile_id) + if not p: + return "❌ سرویس‌دهنده یافت نشد.", [[Button.inline("🔙 بازگشت", data="ai_prov_list")]] + + all_providers = await self.repo.get_provider_profiles() + other_providers = [op for op in all_providers if op.id != profile_id] + + curr_fb = p.fallback_provider_name or "بدون پشتیبان" + text = ( + f"🔄 تنظیم دستی ارائه‌دهنده پشتیبان (Fallback)\n" + f"برای سرویس‌دهنده: {p.name}\n\n" + f"• وضعیت فعلی: {curr_fb}\n\n" + "در صورت بروز هرگونه خطا در این ارائه‌دهنده، درخواست هوش مصنوعی به کدام سرویس‌دهنده منتقل شود؟" + ) + + buttons: List[List[Button]] = [] + is_none_selected = not p.fallback_provider_id + buttons.append([ + Button.inline(f"{'✅ ' if is_none_selected else ''}❌ بدون پشتیبان (توقف زنجیره)", data=f"ai_pv_set_fb:{p.id}:0") + ]) + + for op in other_providers: + is_selected = (p.fallback_provider_id == op.id) + btn_text = f"{'✅ ' if is_selected else ''}🔄 {op.name} ({op.model})" + buttons.append([Button.inline(btn_text, data=f"ai_pv_set_fb:{p.id}:{op.id}")]) + + buttons.append([Button.inline("🔙 بازگشت به جزئیات", data=f"ai_pv_view:{p.id}")]) + return text, buttons + + async def _render_ai_provider_effort_menu(self, profile_id: int) -> Tuple[str, List[List[Button]]]: + p = await self.repo.get_provider_profile_by_id(profile_id) + if not p: + return "❌ سرویس‌دهنده یافت نشد.", [[Button.inline("🔙 بازگشت", data="ai_prov_list")]] + + current = (p.reasoning_effort or "").lower() or "none" + effort_map = { + "none": "خاموش / بدون استدلال (پیش‌فرض)", + "low": "سبک / سریع (Low)", + "medium": "متوسط و استاندارد (Medium)", + "high": "عمیق و دقیق (High)" + } + text = ( + f"🧠 تنظیم حالت تفکر و استدلال (Thinking Mode / Reasoning Effort)\n" + f"برای سرویس‌دهنده: {p.name} (مدل: {p.model})\n\n" + f"• وضعیت فعلی: {effort_map.get(current, current)}\n\n" + "میزان عمق تفکر و استدلال مدل قبل از تولید پاسخ را انتخاب کنید:" + ) + buttons = [ + [Button.inline(f"{'✅ ' if current in ('', 'none') else ''}⚡️ خاموش (Off)", data=f"ai_pv_set_effort:{p.id}:none")], + [Button.inline(f"{'✅ ' if current == 'low' else ''}🟢 سبک و سریع (Low)", data=f"ai_pv_set_effort:{p.id}:low")], + [Button.inline(f"{'✅ ' if current == 'medium' else ''}🟡 متوسط و استاندارد (Medium)", data=f"ai_pv_set_effort:{p.id}:medium")], + [Button.inline(f"{'✅ ' if current == 'high' else ''}🔴 عمیق و حداکثری (High)", data=f"ai_pv_set_effort:{p.id}:high")], + [Button.inline("🔙 بازگشت به جزئیات", data=f"ai_pv_view:{p.id}")] + ] + return text, buttons + + + async def _render_ai_logs(self, limit: int = 10) -> Tuple[str, List[List[Button]]]: + logs = await self.repo.get_recent_ai_logs(limit) + if not logs: + text = "📜 هیچ لاگی از درخواست‌های هوش مصنوعی یافت نشد." + buttons = [[Button.inline("🔙 بازگشت به تنظیمات AI", data="ai_menu")]] + return text, buttons + + text = ( + f"📜 آخرین {len(logs)} درخواست هوش مصنوعی (AI Logs):\n\n" + "برای مشاهده جزئیات کامل پرامپت و پاسخ، روی هر لاگ بزنید:" + ) + buttons = [] + for l in logs: + st = "✅" if l.status == "success" else "❌" + dur = f"{l.duration_sec:.1f}s" + title = f"{st} [{dur}] {l.provider} • {l.action_name} (#{l.id})" + buttons.append([Button.inline(title, data=f"ai_log_view:{l.id}")]) + buttons.append([Button.inline("🔙 بازگشت به تنظیمات AI", data="ai_menu")]) + return text, buttons + + async def _render_ai_log_detail(self, log_id: int) -> Tuple[str, List[List[Button]]]: + log = await self.repo.get_ai_log_by_id(log_id) + if not log: + return "❌ لاگ یافت نشد.", [[Button.inline("🔙 بازگشت", data="ai_logs:10")]] + + st_icon = "✅ موفق" if log.status == "success" else f"❌ خطا: {log.error_message or 'نامشخص'}" + prompt_preview = (log.prompt[:350] + "...") if len(log.prompt) > 350 else log.prompt + sys_preview = ((log.system_prompt[:200] + "...") if len(log.system_prompt) > 200 else log.system_prompt) if log.system_prompt else "ندارد" + resp_text = log.response_text or "ندارد" + resp_preview = (resp_text[:700] + "...") if len(resp_text) > 700 else resp_text + + text = ( + f"📜 جزئیات درخواست هوش مصنوعی #{log.id}\n\n" + f"• ⏰ زمان: {log.created_at}\n" + f"• 📌 عملیات: {log.action_name}\n" + f"• 🔌 سرویس‌دهنده: {log.provider}\n" + f"• 🏷 مدل: {log.model}\n" + f"• ⏱ مدت زمان: {log.duration_sec:.2f} ثانیه\n" + f"• 📊 وضعیت: {st_icon}\n\n" + f"⚙️ System Prompt:\n{sys_preview}\n\n" + f"📥 Prompt:\n{prompt_preview}\n\n" + f"📤 Response:\n{resp_preview}" + ) + buttons = [ + [Button.inline("🔙 بازگشت به لیست لاگ‌ها", data="ai_logs:10")], + [Button.inline("🧠 بازگشت به منوی AI", data="ai_menu")] + ] + return text, buttons + # --- Review Post UI Helpers --- def _build_raw_post_keyboard(self, post: Post, targets: List[TargetChannel]): buttons = [] @@ -216,15 +565,24 @@ class AdminBotService: published_lines += f" • {t_title}\n" published_lines += "➖➖➖➖➖➖➖➖➖➖\n\n" + rejection_line = "" + if post.status == "rejected" or post.rejection_reason: + rejection_line = ( + f"❌ وضعیت: رد شده توسط هوش مصنوعی (Rejected)\n" + f"📝 علت رد: {post.rejection_reason or 'طبق ارزیابی AI'}\n" + f"➖➖➖➖➖➖➖➖➖➖\n\n" + ) + caption = ( f"📥 پست جدید از مبدا ({post.source_channel_id}):\n\n" + f"{rejection_line}" f"{published_lines}" f"{post.raw_text or ''}\n\n" f"👇 کانال مقصد مورد نظر را برای بازنویسی هوشمند انتخاب کنید:" ) return caption - async def send_raw_review_post(self, post_id: int): + async def send_raw_review_post(self, post_id: int, refresh_only: bool = False): post = await self.repo.get_post_by_id(post_id) if not post or not self.review_channel_id or post.is_deleted: return @@ -233,8 +591,25 @@ class AdminBotService: keyboard = self._build_raw_post_keyboard(post, targets) caption = self._format_raw_post_caption(post) + if refresh_only: + # Update the card that already exists instead of posting a second one. + if not post.review_message_id: + return + try: + await self.client.edit_message( + self.review_channel_id, post.review_message_id, + clamp_for_telegram(caption, bool(post.media_path)), + parse_mode="html", buttons=keyboard, + ) + except Exception as e: + logger.debug(f"Could not refresh review card for post {post.id}: {e}") + return + + has_media = bool(post.media_path and os.path.exists(post.media_path)) + caption = clamp_for_telegram(caption, has_media) + try: - if post.media_path and os.path.exists(post.media_path): + if has_media: msg = await self.client.send_file( self.review_channel_id, file=post.media_path, @@ -251,7 +626,436 @@ class AdminBotService: ) await self.repo.update_review_message_id(post.id, msg.id) except Exception as e: - logger.error(f"Failed to send raw post {post.id} to review channel: {e}", exc_info=True) + await log_exception("admin_bot.send_review_post", e, {"post_id": post.id, "review_channel_id": self.review_channel_id}) + + async def refresh_error_metrics(self) -> int: + """Republish the open-error gauges from the database. Returns the open count.""" + groups = await self.repo.get_open_error_summary(limit=200) + ERRORS_OPEN_GAUGE.clear() + total = 0 + for g in groups: + ERRORS_OPEN_GAUGE.labels(service=g["service_name"], error_type=g["error_type"]).set(g["occurrences"]) + total += g["occurrences"] + ERRORS_OPEN_TOTAL_GAUGE.set(total) + return total + + async def _render_error_report(self): + groups = await self.repo.get_open_error_summary(limit=20) + await self.refresh_error_metrics() + + if not groups: + return "✅ هیچ خطای رفع‌نشده‌ای در سیستم وجود ندارد.", [] + + self.error_group_cache = { + i: (g["service_name"], g["error_type"]) for i, g in enumerate(groups) + } + + total = sum(g["occurrences"] for g in groups) + lines = [f"⚠️ خطاهای رفع‌نشده سیستم ({total} مورد در {len(groups)} گروه)\n"] + for i, g in enumerate(groups): + last_seen = g["last_seen"].strftime("%m-%d %H:%M") if g["last_seen"] else "?" + message = (g["last_message"] or "")[:110] + lines.append( + f"\n{i + 1}. {g['error_type']}{g['service_name']}\n" + f" • تعداد: {g['occurrences']} | آخرین: {last_seen}\n" + f" • {message}" + ) + + buttons = [ + [Button.inline(f"✅ رفع شد: {i + 1}. {g['error_type'][:22]}", data=f"errfix:{i}")] + for i, g in enumerate(groups) + ] + buttons.append([Button.inline("✅✅ علامت‌گذاری همه به‌عنوان رفع‌شده", data="errfix_all")]) + return "\n".join(lines), buttons + + async def _render_source_list(self): + sources = await self.repo.get_active_sources() + if not sources: + return "هیچ کانال مبدایی ثبت نشده است. روی ➕ افزودن کانال مبدا بزنید.", [] + + categories = await self.repo.get_categories() + if not categories: + # If no categories exist yet, show flat list + text = ( + f"📡 کانال‌های مبدا فعال ({len(sources)} کانال)\n\n" + "برای دیدن جزئیات و استخراج پست، روی نام هر کانال بزنید:" + ) + buttons = [[Button.inline(f"📢 {s.title or s.channel_id}", data=f"src_view:{s.id}")] for s in sources] + buttons.append([Button.inline("➕ ایجاد دسته‌بندی جدید", data="cat_add")]) + return text, buttons + + # Show categories view for sources + text = ( + f"📡 دسته‌بندی‌های کانال‌های مبدا ({len(sources)} کانال)\n\n" + "روی هر دسته‌بندی بزنید تا کانال‌های مبدای مربوط به آن نمایش داده شود:" + ) + buttons = [] + for cat in categories: + srcs = await self.repo.get_sources_by_category(cat.id) + buttons.append([Button.inline(f"📁 {cat.name} ({len(srcs)} کانال)", data=f"src_cat_view:{cat.id}")]) + + uncat_srcs = await self.repo.get_sources_by_category(None) + if uncat_srcs: + buttons.append([Button.inline(f"📁 کانال‌های بدون دسته‌بندی ({len(uncat_srcs)} کانال)", data="src_cat_view:0")]) + + buttons.append([Button.inline("🌐 نمایش همه کانال‌های مبدا یکجا", data="src_all")]) + buttons.append([Button.inline("📂 مدیریت دسته‌بندی‌ها", data="list_cat")]) + return text, buttons + + async def _render_source_channels_in_category(self, category_id: int): + if category_id == 0: + sources = await self.repo.get_sources_by_category(None) + cat_name = "کانال‌های بدون دسته‌بندی" + else: + cat = await self.repo.get_category_by_id(category_id) + if not cat: + return "❌ دسته‌بندی یافت نشد.", [[Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_src")]] + sources = await self.repo.get_sources_by_category(category_id) + cat_name = cat.name + + if not sources: + text = ( + f"📡 کانال‌های مبدا در دسته: «{cat_name}»\n\n" + "هیچ کانال مبدایی در این دسته‌بندی قرار ندارد." + ) + buttons = [[Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_src")]] + return text, buttons + + text = ( + f"📡 کانال‌های مبدا در دسته: «{cat_name}» ({len(sources)} کانال)\n\n" + "برای مشاهده جزئیات و استخراج پست، روی کانال مورد نظر بزنید:" + ) + buttons = [[Button.inline(f"📢 {s.title or s.channel_id}", data=f"src_view:{s.id}")] for s in sources] + buttons.append([Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_src")]) + return text, buttons + + async def _render_source_all_flat(self): + sources = await self.repo.get_active_sources() + if not sources: + return "هیچ کانال مبدایی ثبت نشده است. روی ➕ افزودن کانال مبدا بزنید.", [] + cats = {c.id: c.name for c in await self.repo.get_categories()} + text = ( + f"📡 همه کانال‌های مبدا فعال ({len(sources)} کانال)\n\n" + "برای دیدن جزئیات و استخراج پست، روی نام هر کانال بزنید:" + ) + buttons = [] + for s in sources: + cat_label = f" [{cats[s.category_id]}]" if s.category_id and s.category_id in cats else "" + buttons.append([Button.inline(f"📢 {s.title or s.channel_id}{cat_label}", data=f"src_view:{s.id}")]) + buttons.append([Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_src")]) + return text, buttons + + async def _render_target_list(self): + targets = await self.repo.get_active_targets() + if not targets: + return "هیچ کانال مقصدی ثبت نشده است. روی ➕ افزودن کانال مقصد بزنید.", [] + + categories = await self.repo.get_categories() + if not categories: + # If no categories exist yet, show flat list + text = ( + f"🎯 کانال‌های مقصد ({len(targets)} کانال)\n\n" + "برای دیدن و تغییر تنظیمات، روی نام هر کانال بزنید:" + ) + buttons = [] + for t in targets: + auto = " 🤖" if t.auto_source_ids else "" + buttons.append([Button.inline(f"🎯 {t.title or t.channel_id}{auto}".strip(), data=f"trg_view:{t.id}")]) + buttons.append([Button.inline("➕ ایجاد دسته‌بندی جدید", data="cat_add")]) + return text, buttons + + # Show categories view for targets + text = ( + f"🎯 دسته‌بندی‌های کانال‌های مقصد ({len(targets)} کانال)\n\n" + "روی هر دسته‌بندی بزنید تا کانال‌های مقصد مربوط به آن نمایش داده شود:" + ) + buttons = [] + for cat in categories: + trgs = await self.repo.get_targets_by_category(cat.id) + buttons.append([Button.inline(f"📁 {cat.name} ({len(trgs)} کانال)", data=f"trg_cat_view:{cat.id}")]) + + uncat_trgs = await self.repo.get_targets_by_category(None) + if uncat_trgs: + buttons.append([Button.inline(f"📁 کانال‌های بدون دسته‌بندی ({len(uncat_trgs)} کانال)", data="trg_cat_view:0")]) + + buttons.append([Button.inline("🌐 نمایش همه کانال‌های مقصد یکجا", data="trg_all")]) + buttons.append([Button.inline("📂 مدیریت دسته‌بندی‌ها", data="list_cat")]) + return text, buttons + + async def _render_target_channels_in_category(self, category_id: int): + if category_id == 0: + targets = await self.repo.get_targets_by_category(None) + cat_name = "کانال‌های بدون دسته‌بندی" + else: + cat = await self.repo.get_category_by_id(category_id) + if not cat: + return "❌ دسته‌بندی یافت نشد.", [[Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_trg")]] + targets = await self.repo.get_targets_by_category(category_id) + cat_name = cat.name + + if not targets: + text = ( + f"🎯 کانال‌های مقصد در دسته: «{cat_name}»\n\n" + "هیچ کانال مقصدی در این دسته‌بندی قرار ندارد." + ) + buttons = [[Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_trg")]] + return text, buttons + + text = ( + f"🎯 کانال‌های مقصد در دسته: «{cat_name}» ({len(targets)} کانال)\n\n" + "برای مشاهده و تغییر تنظیمات، روی کانال مورد نظر بزنید:" + ) + buttons = [] + for t in targets: + auto = " 🤖" if t.auto_source_ids else "" + buttons.append([Button.inline(f"🎯 {t.title or t.channel_id}{auto}".strip(), data=f"trg_view:{t.id}")]) + buttons.append([Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_trg")]) + return text, buttons + + async def _render_target_all_flat(self): + targets = await self.repo.get_active_targets() + if not targets: + return "هیچ کانال مقصدی ثبت نشده است. روی ➕ افزودن کانال مقصد بزنید.", [] + cats = {c.id: c.name for c in await self.repo.get_categories()} + text = ( + f"🎯 همه کانال‌های مقصد ({len(targets)} کانال)\n\n" + "برای دیدن و تغییر تنظیمات، روی نام هر کانال بزنید:" + ) + buttons = [] + for t in targets: + auto = " 🤖" if t.auto_source_ids else "" + cat_label = f" [{cats[t.category_id]}]" if t.category_id and t.category_id in cats else "" + buttons.append([Button.inline(f"🎯 {t.title or t.channel_id}{cat_label}{auto}".strip(), data=f"trg_view:{t.id}")]) + buttons.append([Button.inline("🔙 بازگشت به دسته‌بندی‌ها", data="list_trg")]) + return text, buttons + + + async def _render_source_config(self, source_id: int): + source = await self.repo.get_source_by_id(source_id) + if not source: + return "❌ کانال مبدا یافت نشد.", [] + + cat_name = "بدون دسته‌بندی" + if getattr(source, "category_id", None): + cat = await self.repo.get_category_by_id(source.category_id) + if cat: + cat_name = cat.name + + auto_targets = await self.repo.get_targets_auto_routed_from(source.channel_id) + auto_line = ( + "🤖 ارسال خودکار به: " + "، ".join(f"{t.title}" for t in auto_targets) + if auto_targets else "🛑 ارسال خودکار: غیرفعال" + ) + collected = await self.repo.count_posts_from_source(source.channel_id) + + card = ( + f"📢 {source.title or 'کانال مبدا'}\n\n" + f"• 🆔 شناسه کانال: {source.channel_id}\n" + f"• 🔗 یوزرنیم: @{source.username or 'ندارد'}\n" + f"• 📁 دسته‌بندی: {cat_name}\n" + f"• 📥 پست‌های دریافت‌شده: {collected}\n" + f"• {auto_line}\n\n" + "👇 برای دریافت پست‌های گذشته یکی از گزینه‌ها را انتخاب کنید:" + ) + buttons = get_source_fetch_buttons(source.channel_id, source.id) + buttons.append([Button.inline("🔙 بازگشت به لیست کانال‌های مبدا", data="list_src")]) + return card, buttons + + + async def _render_auto_sources(self, target_id: int): + """Toggle screen listing every source with its on/off state for this target.""" + target = await self.repo.get_target_by_id(target_id) + if not target: + return "❌ کانال مقصد یافت نشد.", [] + + sources = await self.repo.get_active_sources() + if not sources: + return ( + "📡 هنوز هیچ کانال مبدایی ثبت نشده است.\n" + "ابتدا از منوی اصلی یک کانال مبدا اضافه کنید.", + [[Button.inline("🔙 بازگشت", data=f"trg_view:{target_id}")]], + ) + + enabled = set(target.auto_source_ids or []) + text = ( + f"🤖 ارسال خودکار برای کانال «{target.title}»\n\n" + "هر کانال مبدایی که روشن باشد، پست‌هایش به‌صورت خودکار بازنویسی شده و " + "به صف ارسال این کانال مقصد اضافه می‌شود.\n\n" + "⚠️ پست‌ها همچنان به کانال ادمین هم ارسال می‌شوند تا در جریان باشید.\n\n" + "👇 برای روشن/خاموش کردن هر مبدا روی آن بزنید:" + ) + buttons = [ + [Button.inline( + f"{'✅' if s.channel_id in enabled else '⬜️'} {s.title or s.channel_id}", + data=f"auto_tg:{target_id}:{s.channel_id}", + )] + for s in sources + ] + buttons.append([Button.inline("🔙 بازگشت به تنظیمات کانال", data=f"trg_view:{target_id}")]) + return text, buttons + + async def _render_categories_menu(self) -> Tuple[str, List[List[Button]]]: + cats = await self.repo.get_categories() + uncat_srcs = await self.repo.get_sources_by_category(None) + uncat_trgs = await self.repo.get_targets_by_category(None) + + lines = [ + "📂 مدیریت دسته‌بندی کانال‌ها (Categories)\n", + "با ایجاد دسته‌بندی، می‌توانید کانال‌های مبدا و مقصد را گروه‌بندی کرده و منظم‌تر مدیریت کنید:\n" + ] + buttons = [] + + if cats: + for c in cats: + counts = await self.repo.get_category_channel_counts(c.id) + s_c = counts["sources"] + t_c = counts["targets"] + label = f"📁 {c.name} ({s_c} مبدا | {t_c} مقصد)" + buttons.append([Button.inline(label, data=f"cat_view:{c.id}")]) + lines.append(f"• {c.name}: {s_c} مبدا | {t_c} مقصد") + else: + lines.append("هنوز هیچ دسته‌بندی ساخته نشده است.") + + if uncat_srcs or uncat_trgs: + lines.append(f"\n• 📁 کانال‌های بدون دسته: {len(uncat_srcs)} مبدا | {len(uncat_trgs)} مقصد") + buttons.append([Button.inline(f"📁 کانال‌های بدون دسته ({len(uncat_srcs)} مبدا | {len(uncat_trgs)} مقصد)", data="cat_view:0")]) + + buttons.append([Button.inline("➕ ایجاد دسته‌بندی جدید", data="cat_add")]) + return "\n".join(lines), buttons + + async def _render_category_detail(self, cat_id: int) -> Tuple[str, List[List[Button]]]: + if cat_id == 0: + sources = await self.repo.get_sources_by_category(None) + targets = await self.repo.get_targets_by_category(None) + cat_name = "کانال‌های بدون دسته‌بندی" + desc_line = "کانال‌هایی که هنوز در هیچ دسته‌بندی قرار نگرفته‌اند." + else: + cat = await self.repo.get_category_by_id(cat_id) + if not cat: + return "❌ دسته‌بندی یافت نشد.", [[Button.inline("🔙 بازگشت به لیست دسته‌ها", data="list_cat")]] + sources = await self.repo.get_sources_by_category(cat_id) + targets = await self.repo.get_targets_by_category(cat_id) + cat_name = cat.name + desc_line = cat.description or "توضیحاتی ثبت نشده است." + + text = ( + f"📁 دسته‌بندی: {cat_name}\n\n" + f"• 📝 توضیحات: {desc_line}\n" + f"• 📡 کانال‌های مبدا ({len(sources)} کانال):\n" + f"• 🎯 کانال‌های مقصد ({len(targets)} کانال):\n\n" + "👇 برای مدیریت هر کانال یا ویرایش این دسته، روی دکمه مربوطه بزنید:" + ) + + buttons = [] + if sources: + for s in sources: + buttons.append([Button.inline(f"📡 {s.title or s.channel_id}", data=f"src_view:{s.id}")]) + if targets: + for t in targets: + buttons.append([Button.inline(f"🎯 {t.title or t.channel_id}", data=f"trg_view:{t.id}")]) + + if cat_id != 0: + buttons.append([ + Button.inline("✏️ تغییر نام دسته", data=f"cat_rename:{cat_id}"), + Button.inline("🗑 حذف این دسته‌بندی", data=f"cat_del:{cat_id}") + ]) + buttons.append([Button.inline("🔙 بازگشت به لیست دسته‌ها", data="list_cat")]) + return text, buttons + + async def _render_set_channel_category_menu(self, channel_type: str, channel_id: int) -> Tuple[str, List[List[Button]]]: + cats = await self.repo.get_categories() + if channel_type == "src": + ch = await self.repo.get_source_by_id(channel_id) + title = ch.title if ch else f"مبدا #{channel_id}" + back_cb = f"src_view:{channel_id}" + else: + ch = await self.repo.get_target_by_id(channel_id) + title = ch.title if ch else f"مقصد #{channel_id}" + back_cb = f"trg_view:{channel_id}" + + curr_cat_id = getattr(ch, "category_id", None) if ch else None + + text = ( + f"📁 انتخاب دسته‌بندی برای: «{title}»\n\n" + "دسته‌بندی مورد نظرتان را از لیست زیر انتخاب کنید:" + ) + buttons = [] + for c in cats: + mark = " ✅" if curr_cat_id == c.id else "" + buttons.append([Button.inline(f"📁 {c.name}{mark}", data=f"cat_set:{channel_type}:{channel_id}:{c.id}")]) + + buttons.append([Button.inline("❌ بدون دسته‌بندی (حذف دسته)", data=f"cat_set:{channel_type}:{channel_id}:0")]) + buttons.append([Button.inline("🔙 بازگشت", data=back_cb)]) + return text, buttons + + + async def auto_route_post(self, post_id: int) -> int: + """Queue a post to every target subscribed to its source channel. + + The admin review card is always sent as well; auto-routing is additive so an + operator can still intervene on anything that went out automatically. + """ + post = await self.repo.get_post_by_id(post_id) + if not post or post.is_deleted: + return 0 + + if os.getenv("AI_AUTO_POSTING_ENABLED", "true").lower() not in ("true", "1", "yes"): + logger.info(f"Auto-routing skipped for post {post_id}: AI_AUTO_POSTING_ENABLED is disabled") + return 0 + + targets = await self.repo.get_targets_auto_routed_from(post.source_channel_id) + if not targets: + return 0 + + routed = 0 + ai_rejected_reasons = [] + for target in targets: + try: + text = post.raw_text or "" + if self.ai_processor: + rewrite_res = await self.ai_processor.rewrite_for_target( + text, target, has_media=bool(post.media_path), image_path=post.media_path + ) + if getattr(rewrite_res, "is_rejected", False): + reason = getattr(rewrite_res, "rejection_reason", "رد شده توسط هوش مصنوعی") + logger.info(f"AI rejected post {post.id} for target {target.title}: {reason}") + ai_rejected_reasons.append(f"{target.title}: {reason}") + continue + text = str(rewrite_res) + + payload = { + "post_id": post.id, + "text": text, + "media_path": post.media_path, + "target_id": target.id, + "target_title": target.title, + } + if self.queue: + await self.queue.push_target_post(target.id, payload) + await self.repo.record_post_queued_to_target(post.id, target.id, target.title or "Target") + ADMIN_ACTIONS_TOTAL.labels(action="auto_routed").inc() + AUTO_ROUTED_POSTS_TOTAL.labels( + source_channel_id=str(post.source_channel_id), + target_title=target.title or "", + ).inc() + routed += 1 + logger.info(f"Auto-routed post {post.id} to target {target.title} (#{target.id})") + except Exception as e: + await log_exception("admin_bot.auto_route", e, {"post_id": post.id, "target_id": target.id}) + + if not routed and ai_rejected_reasons: + combined_reason = " | ".join(ai_rejected_reasons) + await self.repo.reject_post(post.id, rejection_reason=combined_reason) + await self.send_raw_review_post(post.id, refresh_only=True) + elif routed: + await self.send_raw_review_post(post.id, refresh_only=True) + return routed + + + async def handle_collected_post(self, post_id: int): + """Entry point used by the collector: review card first, then auto-routing.""" + await self.send_raw_review_post(post_id) + await self.auto_route_post(post_id) def _register_handlers(self): # --- Start / Menu --- @@ -274,27 +1078,36 @@ class AdminBotService: f"• 📋 شناسه کانال ادمین‌ها: {self.review_channel_id}\n\n" "برای مدیریت کانال‌ها از دکمه‌های زیر استفاده کنید:" ) - await event.reply(welcome_text, parse_mode="html", buttons=get_persian_main_menu()) + menu = await self.get_menu() + await event.reply(welcome_text, parse_mode="html", buttons=menu) - # --- Statistics --- - @self.client.on(events.NewMessage(pattern=r"(?i)^(/stats|📊 آمار و وضعیت ناوگان)$")) + # --- Statistics & Grafana Metrics Report --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/stats|/report|📊 آمار و وضعیت ناوگان|📊 گزارش و آمار سیستم)$")) async def cmd_stats(event: events.NewMessage.Event): if not self.is_admin(event.sender_id): return - pending_review = len(await self.repo.get_posts_by_status("pending_review", limit=5000)) - published = len(await self.repo.get_posts_by_status("published", limit=5000)) - rejected = len(await self.repo.get_posts_by_status("rejected", limit=5000)) - total_redis_q = await self.queue.get_total_queued_posts() if self.queue else 0 + report_text = await get_instant_metrics_report() + buttons = [ + [ + Button.inline("⏱ نمودار ۱۵ دقیقه اخیر", data="rep_gr:15m"), + Button.inline("⏱ نمودار ۳ ساعت اخیر", data="rep_gr:3h"), + ], + [ + Button.inline("⏱ نمودار ۲۴ ساعت اخیر", data="rep_gr:24h"), + Button.inline("🔄 بروزرسانی گزارش", data="rep_refresh"), + ] + ] + await event.reply(report_text, parse_mode="html", buttons=buttons) + + + # --- AI Settings & Logs --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/ai|🧠 تنظیمات و لاگ‌های AI)$")) + async def cmd_ai(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + text, buttons = await self._render_ai_menu() + await event.reply(text, parse_mode="html", buttons=buttons) - text = ( - "📊 آمار زنده سیستم کپی‌کار:\n\n" - f"• 📥 مجموع پست‌های در صف ارسال کانال‌های مقصد: {total_redis_q}\n" - f"• 📋 پست‌های در انتظار بررسی ادمین: {pending_review}\n" - f"• 🚀 پست‌های منتشر شده: {published}\n" - f"• ❌ پست‌های رد شده: {rejected}\n\n" - "📈 داشبورد مانیتورینگ گرانافا: http://localhost:3000" - ) - await event.reply(text, parse_mode="html", buttons=get_persian_main_menu()) # --- Userbot Authentication --- @self.client.on(events.NewMessage(pattern=r"(?i)^(/request_code|🔑 درخواست کد لاگین)$")) @@ -324,28 +1137,8 @@ class AdminBotService: async def cmd_sources(event: events.NewMessage.Event): if not self.is_admin(event.sender_id): return - sources = await self.repo.get_active_sources() - if not sources: - await event.reply("هیچ کانال مبدایی ثبت نشده است. روی ➕ افزودن کانال مبدا بزنید.", parse_mode="html") - return - - await event.reply(f"📡 کانال‌های مبدا فعال ({len(sources)} کانال):", parse_mode="html") - for s in sources: - card = ( - f"📢 {s.title or 'کانال مبدا'}\n" - f"• شناسه: {s.channel_id}\n" - f"• یوزرنیم: @{s.username or 'ندارد'}" - ) - buttons = [ - [ - Button.inline("📥 استخراج ۲۰ پست", data=f"hist:{s.channel_id}:20"), - Button.inline("📥 استخراج ۵۰ پست", data=f"hist:{s.channel_id}:50"), - ], - [ - Button.inline("🗑 حذف این کانال مبدا", data=f"del_src:{s.id}") - ] - ] - await event.reply(card, parse_mode="html", buttons=buttons) + text, buttons = await self._render_source_list() + await event.reply(text, parse_mode="html", buttons=buttons or None) @self.client.on(events.NewMessage(pattern=r"(?i)^(➕ افزودن کانال مبدا)$")) async def cmd_add_source_prompt(event: events.NewMessage.Event): @@ -364,15 +1157,16 @@ class AdminBotService: async def cmd_targets(event: events.NewMessage.Event): if not self.is_admin(event.sender_id): return - targets = await self.repo.get_active_targets() - if not targets: - await event.reply("هیچ کانال مقصدی ثبت نشده است. روی ➕ افزودن کانال مقصد بزنید.", parse_mode="html") - return + text, buttons = await self._render_target_list() + await event.reply(text, parse_mode="html", buttons=buttons or None) - await event.reply(f"🎯 کانال‌های مقصد برای انتشار ({len(targets)} کانال):\nبرای تنظیمات هر کانال، روی دکمه آن بزنید:", parse_mode="html") - for t in targets: - card, buttons = await self._render_target_config(t.id) - await event.reply(card, parse_mode="html", buttons=buttons) + # --- Categories Management --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/categories|/cats|📂 دسته‌بندی کانال‌ها)$")) + async def cmd_categories(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + text, buttons = await self._render_categories_menu() + await event.reply(text, parse_mode="html", buttons=buttons or None) @self.client.on(events.NewMessage(pattern=r"(?i)^(➕ افزودن کانال مقصد)$")) async def cmd_add_target_prompt(event: events.NewMessage.Event): @@ -384,8 +1178,69 @@ class AdminBotService: "ساده‌ترین روش: یک پیام از کانال مورد نظر به این ربات فوروارد (Forward) کنید!\n\n" "(یا می‌توانید آیدی، یوزرنیم یا لینک کانال مثل @channel_username را بفرستید)" ) + await event.reply(guide, parse_mode="html", buttons=get_cancel_button()) + # --- Re-send unreviewed posts --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/pending|📨 ارسال پست‌های بررسی‌نشده)$")) + async def cmd_pending(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + posts = await self.repo.get_unreviewed_posts(limit=50) + if not posts: + menu = await self.get_menu() + await event.reply("✅ هیچ پست بررسی‌نشده‌ای در صف نیست.", buttons=menu) + return + + msg = await event.reply( + f"⏳ در حال ارسال {len(posts)} پست بررسی‌نشده به کانال ادمین...", + parse_mode="html", + ) + for post in posts: + await self.send_raw_review_post(post.id) + await msg.edit( + f"✅ تعداد {len(posts)} پست بررسی‌نشده به کانال ادمین ارسال شد.", + parse_mode="html", + ) + + # --- System Pause / Emergency Stop & Resume --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/pause|/stop|🛑 توقف اضطراری سیستم)$")) + async def cmd_pause(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + await self.repo.set_system_paused(True) + text = ( + "🛑 سیستم به صورت اضطراری متوقف شد (System Paused).\n\n" + "• ⏸ دریافت خودکار پست‌های جدید از کانال‌های مبدا متوقف شد.\n" + "• ⏸ ارسال زمان‌بندی‌شده پست‌ها به کانال‌های مقصد متوقف شد.\n" + "• 🟢 پنل مدیریت، دستورات و پیام‌های ادمین همچنان کاملاً فعال هستند.\n\n" + "برای از سرگیری فعالیت سیستم، روی دکمه «▶️ راه‌اندازی و ادامه سیستم» بزنید." + ) + menu = await self.get_menu() + await event.reply(text, parse_mode="html", buttons=menu) + + @self.client.on(events.NewMessage(pattern=r"(?i)^(/resume|/start_fleet|▶️ راه‌اندازی و ادامه سیستم)$")) + async def cmd_resume(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + await self.repo.set_system_paused(False) + text = ( + "▶️ سیستم با موفقیت مجدداً فعال شد (System Resumed).\n\n" + "• 🟢 دریافت پست‌های جدید از کانال‌های مبدا از سر گرفته شد.\n" + "• 🟢 ارسال پست‌های در صف به کانال‌های مقصد طبق زمان‌بندی آغاز شد." + ) + menu = await self.get_menu() + await event.reply(text, parse_mode="html", buttons=menu) + + # --- Release Notes & Changes --- + @self.client.on(events.NewMessage(pattern=r"(?i)^(/changes|/notes|📝 گزارش تغییرات اخیر)$")) + async def cmd_changes(event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + notes = self.get_latest_change_notes() + menu = await self.get_menu() + await event.reply(notes, parse_mode="html", buttons=menu) + # --- Help --- @self.client.on(events.NewMessage(pattern=r"(?i)^(/help|❓ راهنمای سیستم)$")) async def cmd_help(event: events.NewMessage.Event): @@ -393,382 +1248,1253 @@ class AdminBotService: return help_text = ( "📖 راهنمای سیستم کپی‌کار:\n\n" + "• توقف و راه‌اندازی اضطراری: با دکمه 🛑 توقف اضطراری سیستم تمام فرآیندهای دریافت و ارسال متوقف می‌شوند در حالی که ربات همیشه به ادمین پاسخ می‌دهد.\n" + "• زنجیره پشتیبان هوش مصنوعی: در صورت بروز خطا در ارائه‌دهنده فعال، درخواست طبق تنظیمات دستی شما به ارائه‌دهنده پشتیبان تعیین‌شده منتقل می‌شود.\n" + "• گزارش تغییرات: با دستور /changes آخرین تغییرات و امکانات سیستم را مشاهده کنید.\n" "• افزودن آسان کانال‌ها: فقط کافیست یک پیام از کانال را به ربات فوروارد کنید تا شناسه و نام آن خودکار ثبت شود!\n" "• تنظیمات با دکمه: در بخش کانال‌های مقصد، روی هر کانال دکمه‌های تغییر لحن، تگ‌ها، فاصله زمانی و ساعت خواب وجود دارد.\n" - "• بازنویسی هوشمند: پست‌های مبدا به کانال ادمین می‌آیند و با لمس دکمه هر مقصد، متن با هوش مصنوعی و لحن همان کانال بازنویسی می‌شود." + "• بازنویسی هوشمند: پست‌های مبدا به کانال ادمین می‌آیند و با لمس دکمه هر مقصد، متن با هوش مصنوعی و لحن همان کانال بازنویسی می‌شود.\n" + "• پست‌های بررسی‌نشده: با دکمه 📨 ارسال پست‌های بررسی‌نشده هر پستی که هنوز به کانال ادمین نرفته، دوباره ارسال می‌شود.\n" + "• ارسال خودکار: در تنظیمات هر کانال مقصد، دکمه 🤖 ارسال خودکار از مبدا را بزنید تا پست‌های یک مبدا خودکار به آن مقصد بروند.\n" + "• خطاها: با دکمه ⚠️ خطاهای سیستم خطاهای رفع‌نشده را ببینید و بعد از رفع، علامت‌گذاری کنید." ) - await event.reply(help_text, parse_mode="html", buttons=get_persian_main_menu()) + menu = await self.get_menu() + await event.reply(help_text, parse_mode="html", buttons=menu) + + def get_latest_change_notes(self) -> str: + """Return the formatted change notes in Persian.""" + return ( + "📢 یادداشت تغییرات و قابلیت‌های جدید سیستم (Release Notes):\n\n" + "1️⃣ تنظیم دستی زنجیره پشتیبان هوش مصنوعی (Manual Fallback Chain):\n" + "• امکان تعیین ارائه‌دهنده پشتیبان اختصاصی برای تک‌تک سرویس‌دهنده‌ها در منوی هوش مصنوعی.\n" + "• هدایت دقیق زنجیره به ارائه‌دهنده مدنظر شما در صورت بروز خطا و ارسال اعلان لحظه‌ای به ادمین.\n\n" + "2️⃣ پشتیبانی کامل از تصاویر و مالتی‌مدال (Vision Multimodal):\n" + "• ارسال تصاویر پست‌های ورودی به مدل‌های هوش مصنوعی (AGY، OpenAI Vision، Gemini).\n" + "• درک و تحلیل متن درون عکس‌ها، نمودارها و اینفوگرافیک‌ها در بازنویسی.\n\n" + "3️⃣ پایش دیسک سیستم‌عامل سرور (Host OS Disk Metrics):\n" + "• نمایش ظرفیت واقعی، فضای مصرف‌شده و آزاد دیسک اصلی سرور در گرافانا و دستور /stats.\n\n" + "4️⃣ دکمه توقف و شروع اضطراری (Emergency Stop):\n" + "• دکمه توقف و ازسرگیری مستقیم از منوی اصلی، با فعال ماندن همیشگی پاسخگویی ربات به ادمین‌ها." + ) + + async def broadcast_change_notes(self): + """Send latest change notes to admins and review channel.""" + notes = self.get_latest_change_notes() + await self.notify_admins(notes) + # --- Step-by-Step Parameter & Text Input Message Handler --- @self.client.on(events.NewMessage) async def handle_user_input(event: events.NewMessage.Event): - if not self.is_admin(event.sender_id): - return - - text = (event.raw_text or "").strip() - - # Ignore main menu commands - if text in ["/start", "/menu", "منو", "/stats", "📊 آمار و وضعیت ناوگان", - "/sources", "📡 کانال‌های مبدا", "/targets", "🎯 کانال‌های مقصد", - "➕ افزودن کانال مبدا", "➕ افزودن کانال مقصد", "❓ راهنمای سیستم", - "/request_code", "🔑 درخواست کد لاگین"]: - return - - # Cancel command - if text in ["/cancel", "انصراف", "لغو"]: - self.user_states.pop(event.sender_id, None) - await event.reply("❌ عملیات لغو شد.", buttons=get_persian_main_menu()) - return - - state = self.user_states.get(event.sender_id) - if not state: - # Direct /code or /password handling - if text.startswith("/code "): - code = text.split("/code ", 1)[1].strip() - res = await self.collector.submit_code(code) - await event.reply(res, parse_mode="html", buttons=get_persian_main_menu()) - elif text.startswith("/password "): - pwd = text.split("/password ", 1)[1].strip() - res = await self.collector.submit_password(pwd) - await event.reply(res, parse_mode="html", buttons=get_persian_main_menu()) - return - - action = state.get("action") - - # 1. Waiting for Source Channel Forward or Text - if action == "wait_source_fwd": - ch_id, title, username, err = await self._resolve_channel(event) - if err or not ch_id: - await event.reply(err or "خطا در دریافت اطلاعات کانال. لطفا مجدد ارسال کنید:", buttons=get_cancel_button()) - return - - self.user_states.pop(event.sender_id, None) - await self.repo.add_source(channel_id=ch_id, title=title, username=username) - - buttons = [ - [ - Button.inline("📥 استخراج ۲۰ پست گذشته", data=f"hist:{ch_id}:20"), - Button.inline("📥 استخراج ۵۰ پست گذشته", data=f"hist:{ch_id}:50") - ] - ] - await event.reply( - f"✅ کانال مبدا {title} با شناسه {ch_id} با موفقیت افزوده شد!\n\n" - f"آیا مایلید پست‌های قبلی این کانال هم دریافت شود؟", - parse_mode="html", - buttons=buttons - ) - - # 2. Waiting for Target Channel Forward or Text - elif action == "wait_target_fwd": - ch_id, title, username, err = await self._resolve_channel(event) - if err or not ch_id: - await event.reply(err or "خطا در دریافت اطلاعات کانال. لطفا مجدد ارسال کنید:", buttons=get_cancel_button()) - return - - self.user_states.pop(event.sender_id, None) - target_id = await self.repo.add_target(channel_id=ch_id, title=title, username=username) - - card, buttons = await self._render_target_config(target_id) - await event.reply( - f"✅ کانال مقصد {title} افزوده شد!\n\nاکنون می‌توانید با دکمه‌های زیر لحن و زمان‌بندی آن را تنظیم کنید:", - parse_mode="html" - ) - await event.reply(card, parse_mode="html", buttons=buttons) - - # 3. Waiting for Target Personality - elif action == "wait_personality": - target_id = state.get("target_id") - self.user_states.pop(event.sender_id, None) - await self.repo.update_target_personality(target_id, text) - card, buttons = await self._render_target_config(target_id) - await event.reply("✅ لحن و شخصیت کانال با موفقیت به روز شد!", parse_mode="html") - await event.reply(card, parse_mode="html", buttons=buttons) - - # 4. Waiting for Target Custom Footer - elif action == "wait_footer": - target_id = state.get("target_id") - self.user_states.pop(event.sender_id, None) - await self.repo.update_target_footer(target_id, text) - card, buttons = await self._render_target_config(target_id) - await event.reply("✅ فوتر و تگ‌های اختصاصی کانال با موفقیت به روز شد!", parse_mode="html") - await event.reply(card, parse_mode="html", buttons=buttons) - - # 5. Waiting for Target Post Interval - elif action == "wait_interval": - target_id = state.get("target_id") - if not text.isdigit() or int(text) < 1: - await event.reply("⚠️ لطفا یک عدد معتبر (به دقیقه) ارسال کنید:", buttons=get_cancel_button()) - return - interval_min = int(text) - self.user_states.pop(event.sender_id, None) - await self.repo.update_target_schedule(target_id=target_id, post_interval_min=interval_min) - card, buttons = await self._render_target_config(target_id) - await event.reply(f"✅ فاصله ارسال به هر {interval_min} دقیقه تغییر یافت!", parse_mode="html") - await event.reply(card, parse_mode="html", buttons=buttons) - - # 6. Waiting for Sleep Window Hours - elif action == "wait_sleep": - target_id = state.get("target_id") - parts = re.findall(r"\d+", text) - if len(parts) < 2: - await event.reply("⚠️ لطفا ساعت شروع و پایان را به این شکل بفرستید: 23 8 (برای ۲۳:۰۰ تا ۰۸:۰۰)", parse_mode="html", buttons=get_cancel_button()) - return - start_h = int(parts[0]) % 24 - end_h = int(parts[1]) % 24 - self.user_states.pop(event.sender_id, None) - await self.repo.update_target_schedule(target_id=target_id, sleep_start_hour=start_h, sleep_end_hour=end_h, is_sleep_enabled=True) - card, buttons = await self._render_target_config(target_id) - await event.reply(f"🌙 ساعت خواب از {start_h}:00 تا {end_h}:00 فعال شد!", parse_mode="html") - await event.reply(card, parse_mode="html", buttons=buttons) - - # 7. Waiting for Login Code - elif action == "wait_login_code": - self.user_states.pop(event.sender_id, None) - res = await self.collector.submit_code(text) - await event.reply(res, parse_mode="html", buttons=get_persian_main_menu()) + # Telethon swallows handler exceptions into its own logger, so an unhandled + # error here looks to the admin like the bot simply ignored them. + try: + await self._handle_user_input(event) + except Exception as e: + await log_exception("admin_bot.user_input", e, { + "sender_id": event.sender_id, + "state": self.user_states.get(event.sender_id), + }) + await self._report_error_to_admin(event, e) # --- Inline Callback Queries --- @self.client.on(events.CallbackQuery) async def on_callback(event: events.CallbackQuery.Event): - if not self.is_admin(event.sender_id): - await event.answer("⛔ دسترسی غیرمجاز.", alert=True) + try: + await self._handle_callback(event) + except Exception as e: + await log_exception("admin_bot.callback", e, { + "sender_id": event.sender_id, + "data": event.data.decode("utf-8", errors="replace") if event.data else None, + }) + await self._report_error_to_admin(event, e) + + async def _report_error_to_admin(self, event, error: Exception) -> None: + """Surface a failure in the chat it happened in, so no action ever fails silently.""" + message = ( + "❌ خطایی در انجام این عملیات رخ داد.\n\n" + f"{type(error).__name__}: {error}\n\n" + "جزئیات کامل در جدول خطاهای سیستم ثبت شد." + ) + try: + await event.reply(message, parse_mode="html") + except Exception: + try: + await event.answer("❌ خطا در انجام عملیات. جزئیات در لاگ ثبت شد.", alert=True) + except Exception as notify_err: + logger.error(f"Could not report error to admin: {notify_err}") + + async def _handle_user_input(self, event: events.NewMessage.Event): + if not self.is_admin(event.sender_id): + return + + text = (event.raw_text or "").strip() + + # Ignore main menu commands + if text in ["/start", "/menu", "منو", "/stats", "📊 آمار و وضعیت ناوگان", + "/sources", "📡 کانال‌های مبدا", "/targets", "🎯 کانال‌های مقصد", + "/request_code", "🔑 درخواست کد لاگین", + "/pending", "📨 ارسال پست‌های بررسی‌نشده", + "/errors", "⚠️ خطاهای سیستم", + "/ai", "🧠 تنظیمات و لاگ‌های AI", + "/categories", "/cats", "📂 دسته‌بندی کانال‌ها", + "/pause", "/stop", "🛑 توقف اضطراری سیستم", + "/resume", "/start_fleet", "▶️ راه‌اندازی و ادامه سیستم"]: + return + + # Cancel command + if text in ["/cancel", "انصراف", "لغو"]: + self.user_states.pop(event.sender_id, None) + menu = await self.get_menu() + await event.reply("❌ عملیات لغو شد.", buttons=menu) + return + + state = self.user_states.get(event.sender_id) + if not state: + # Direct /code or /password handling + if text.startswith(("/code ", "/password ")) and not self.collector: + await event.reply("سرویس کالکتور متصل نیست.") + return + if text.startswith("/code "): + code = text.split("/code ", 1)[1].strip() + res = await self.collector.submit_code(code) + menu = await self.get_menu() + await event.reply(res, parse_mode="html", buttons=menu) + elif text.startswith("/password "): + pwd = text.split("/password ", 1)[1].strip() + res = await self.collector.submit_password(pwd) + menu = await self.get_menu() + await event.reply(res, parse_mode="html", buttons=menu) + return + + action = state.get("action") + + # 1. Waiting for Source Channel Forward or Text + if action == "wait_source_fwd": + ch_id, title, username, err = await self._resolve_channel(event) + if err or not ch_id: + await event.reply(err or "خطا در دریافت اطلاعات کانال. لطفا مجدد ارسال کنید:", buttons=get_cancel_button()) return - data = event.data.decode("utf-8") + self.user_states.pop(event.sender_id, None) + await self.repo.add_source(channel_id=ch_id, title=title, username=username) - # 0. Cancel Active State - if data == "cancel_state": - self.user_states.pop(event.sender_id, None) - await event.edit("❌ عملیات لغو شد.") - await event.answer("لغو شد.") + buttons = get_source_fetch_buttons(ch_id) + await event.reply( + f"✅ کانال مبدا {title} با شناسه {ch_id} با موفقیت افزوده شد!\n\n" + f"آیا مایلید پست‌های قبلی این کانال هم دریافت شود؟", + parse_mode="html", + buttons=buttons + ) - # --- Target Channel Config Buttons --- - elif data.startswith("st_pers:"): - target_id = int(data.split(":")[1]) - target = await self.repo.get_target_by_id(target_id) - self.user_states[event.sender_id] = {"action": "wait_personality", "target_id": target_id} + # 2. Waiting for Target Channel Forward or Text + elif action == "wait_target_fwd": + ch_id, title, username, err = await self._resolve_channel(event) + if err or not ch_id: + await event.reply(err or "خطا در دریافت اطلاعات کانال. لطفا مجدد ارسال کنید:", buttons=get_cancel_button()) + return + + self.user_states.pop(event.sender_id, None) + target_id = await self.repo.add_target(channel_id=ch_id, title=title, username=username) + + card, buttons = await self._render_target_config(target_id) + await event.reply( + f"✅ کانال مقصد {title} افزوده شد!\n\nاکنون می‌توانید با دکمه‌های زیر لحن و زمان‌بندی آن را تنظیم کنید:", + parse_mode="html" + ) + await event.reply(card, parse_mode="html", buttons=buttons) + + # 3. Waiting for Target Personality + elif action == "wait_personality": + target_id = state.get("target_id") + self.user_states.pop(event.sender_id, None) + await self.repo.update_target_personality(target_id, text) + card, buttons = await self._render_target_config(target_id) + await event.reply("✅ لحن و شخصیت کانال با موفقیت به روز شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 4. Waiting for Target Custom Footer + elif action == "wait_footer": + target_id = state.get("target_id") + self.user_states.pop(event.sender_id, None) + await self.repo.update_target_footer(target_id, text) + card, buttons = await self._render_target_config(target_id) + await event.reply("✅ فوتر و تگ‌های اختصاصی کانال با موفقیت به روز شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 4.1. Waiting for Target Custom Prompt + elif action == "wait_custom_prompt": + target_id = state.get("target_id") + self.user_states.pop(event.sender_id, None) + new_prompt = "" if text.strip() in ("/clear", "clear", "حذف", "پاک") else text.strip() + await self.repo.update_target_custom_prompt(target_id, new_prompt) + card, buttons = await self._render_target_config(target_id) + if new_prompt: + await event.reply("✅ دستورالعمل و فرامین اختصاصی کانال با موفقیت ذخیره شد!", parse_mode="html") + else: + await event.reply("✅ پرامپت اختصاصی حذف شد و به حالت پیش‌فرض بازگشت.", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 4.2. Waiting for Category Name Creation + elif action == "wait_cat_name": + self.user_states.pop(event.sender_id, None) + cat_name = text.strip() + if not cat_name: + menu = await self.get_menu() + await event.reply("⚠️ نام دسته‌بندی نمی‌تواند خالی باشد.", buttons=menu) + return + cat_id = await self.repo.create_category(name=cat_name) + card, buttons = await self._render_category_detail(cat_id) + await event.reply(f"✅ دسته‌بندی {cat_name} با موفقیت ایجاد شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 4.3. Waiting for Category Rename + elif action == "wait_cat_rename": + cat_id = state.get("cat_id") + self.user_states.pop(event.sender_id, None) + cat_name = text.strip() + if not cat_name: + menu = await self.get_menu() + await event.reply("⚠️ نام دسته‌بندی نمی‌تواند خالی باشد.", buttons=menu) + return + await self.repo.update_category(cat_id, name=cat_name) + card, buttons = await self._render_category_detail(cat_id) + await event.reply(f"✅ نام دسته‌بندی به {cat_name} تغییر یافت!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + + + # 5. Waiting for Target Post Interval + elif action == "wait_interval": + target_id = state.get("target_id") + if not text.isdigit() or int(text) < 1: + await event.reply("⚠️ لطفا یک عدد معتبر (به دقیقه) ارسال کنید:", buttons=get_cancel_button()) + return + interval_min = int(text) + self.user_states.pop(event.sender_id, None) + await self.repo.update_target_schedule(target_id=target_id, post_interval_min=interval_min) + card, buttons = await self._render_target_config(target_id) + await event.reply(f"✅ فاصله ارسال به هر {interval_min} دقیقه تغییر یافت!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 6. Waiting for Sleep Window Hours + elif action == "wait_sleep": + target_id = state.get("target_id") + parts = re.findall(r"\d+", text) + if len(parts) < 2: + await event.reply("⚠️ لطفا ساعت شروع و پایان را به این شکل بفرستید: 23 8 (برای ۲۳:۰۰ تا ۰۸:۰۰)", parse_mode="html", buttons=get_cancel_button()) + return + start_h = int(parts[0]) % 24 + end_h = int(parts[1]) % 24 + self.user_states.pop(event.sender_id, None) + await self.repo.update_target_schedule(target_id=target_id, sleep_start_hour=start_h, sleep_end_hour=end_h, is_sleep_enabled=True) + card, buttons = await self._render_target_config(target_id) + await event.reply(f"🌙 ساعت خواب از {start_h}:00 تا {end_h}:00 فعال شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + # 7. Waiting for a custom history-fetch count + elif action == "wait_fetch_count": + ch_id = state.get("channel_id") + digits = re.findall(r"\d+", text) + if not digits: await event.reply( - f"🎭 تغییر لحن و شخصیت کانال «{target.title}»:\n\n" - "لطفا در پیام بعدی، لحن و استایل نگارش مورد نظرتان را بفرستید:\n" - "(مثال: لحن دوستانه و پرانرژی همراه با ایموجی و تیترهای جذاب)", - parse_mode="html", - buttons=get_cancel_button() + f"⚠️ لطفا فقط یک عدد بین {MIN_FETCH_LIMIT} تا {MAX_FETCH_LIMIT} بفرستید:", + buttons=get_cancel_button(), ) - await event.answer() - - elif data.startswith("st_foot:"): - target_id = int(data.split(":")[1]) - target = await self.repo.get_target_by_id(target_id) - self.user_states[event.sender_id] = {"action": "wait_footer", "target_id": target_id} + return + count = int(digits[0]) + if not (MIN_FETCH_LIMIT <= count <= MAX_FETCH_LIMIT): await event.reply( - f"🏷 تغییر فوتر و تگ‌های اختصاصی «{target.title}»:\n\n" - "لطفا متن فوتر، آیدی یا هشتگ‌هایی که می‌خواهید انتهای هر پست قرار بگیرد را بفرستید:\n" - "(مثال: 🆔 @my_chan\n#فناوری #اخبار)", + f"⚠️ عدد باید بین {MIN_FETCH_LIMIT} و {MAX_FETCH_LIMIT} باشد. دوباره بفرستید:", parse_mode="html", - buttons=get_cancel_button() + buttons=get_cancel_button(), ) - await event.answer() + return - elif data.startswith("st_intv:"): - target_id = int(data.split(":")[1]) - target = await self.repo.get_target_by_id(target_id) - self.user_states[event.sender_id] = {"action": "wait_interval", "target_id": target_id} - await event.reply( - f"⏱ تغییر فاصله ارسال کانال «{target.title}»:\n\n" - "فاصله زمانی بین ارسال هر پست را به دقیقه بفرستید:\n" - "(مثال: 30 برای ارسال هر نیم ساعت یک پست)", - parse_mode="html", - buttons=get_cancel_button() - ) - await event.answer() + self.user_states.pop(event.sender_id, None) + if not self.collector: + await event.reply("سرویس کالکتور متصل نیست.") + return - elif data.startswith("st_slp:"): - target_id = int(data.split(":")[1]) - target = await self.repo.get_target_by_id(target_id) - self.user_states[event.sender_id] = {"action": "wait_sleep", "target_id": target_id} - await event.reply( - f"🌙 تنظیم ساعت خواب کانال «{target.title}»:\n\n" - "ساعت شروع و پایان خواب را با یک فاصله بفرستید:\n" - "(مثال: 23 8 برای عدم ارسال پست از ساعت ۲۳:۰۰ تا ۰۸:۰۰ صبح)", - parse_mode="html", - buttons=get_cancel_button() - ) - await event.answer() + msg = await event.reply( + f"⏳ در حال دریافت {count} پست گذشته از {ch_id}...", + parse_mode="html", + ) - elif data.startswith("dis_slp:"): - target_id = int(data.split(":")[1]) - await self.repo.update_target_schedule(target_id=target_id, is_sleep_enabled=False) - card, buttons = await self._render_target_config(target_id) - await event.edit(card, parse_mode="html", buttons=buttons) - await event.answer("☀️ ساعت خواب غیرفعال شد.") + async def progress_notify(txt: str): + await msg.edit(txt, parse_mode="html") - elif data.startswith("del_trg:"): - target_id = int(data.split(":")[1]) - await self.repo.delete_target(target_id) - await event.edit("🗑 کانال مقصد با موفقیت حذف شد.", parse_mode="html", buttons=None) - await event.answer("کانال مقصد حذف شد.") + await self.collector.scrape_channel_history( + channel_id=ch_id, limit=count, progress_callback=progress_notify + ) - elif data.startswith("del_src:"): - source_id = int(data.split(":")[1]) - await self.repo.delete_source(source_id) - await event.edit("🗑 کانال مبدا با موفقیت حذف شد.", parse_mode="html", buttons=None) - await event.answer("کانال مبدا حذف شد.") + # 8. Waiting for Login Code + elif action == "wait_login_code": + self.user_states.pop(event.sender_id, None) + if not self.collector: + await event.reply("سرویس کالکتور متصل نیست.") + return + res = await self.collector.submit_code(text) + menu = await self.get_menu() + await event.reply(res, parse_mode="html", buttons=menu) - elif data == "list_trg": - targets = await self.repo.get_active_targets() - if not targets: - await event.edit("هیچ کانال مقصدی ثبت نشده است.") - return - await event.edit("🎯 لیست کانال‌های مقصد به روز شد:") - for t in targets: - card, buttons = await self._render_target_config(t.id) - await event.reply(card, parse_mode="html", buttons=buttons) + # 9. Waiting for AI Model + elif action == "wait_ai_model": + self.user_states.pop(event.sender_id, None) + new_model = text.strip() + await self.repo.set_setting("ai_model", new_model) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_menu() + await event.reply(f"✅ مدل فعال به {new_model} تغییر یافت!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - # --- Review Channel Buttons --- - elif data.startswith("hist:"): - _, ch_id_str, limit_str = data.split(":") - ch_id = int(ch_id_str) - limit = int(limit_str) + # 10. Waiting for AI Base URL + elif action == "wait_ai_base_url": + self.user_states.pop(event.sender_id, None) + new_url = "" if text.strip().lower() in ("none", "null", "خالی", "default") else text.strip() + await self.repo.set_setting("ai_base_url", new_url) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_menu() + await event.reply("✅ Base URL با موفقیت ذخیره شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - if not self.collector: - await event.answer("کالکتور در دسترس نیست.", alert=True) - return + # 11. Waiting for AI API Key + elif action == "wait_ai_api_key": + self.user_states.pop(event.sender_id, None) + new_key = "" if text.strip().lower() in ("none", "null", "خالی") else text.strip() + await self.repo.set_setting("ai_api_key", new_key) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_menu() + await event.reply("✅ کلید API با موفقیت ذخیره شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - await event.edit(f"⏳ در حال دریافت {limit} پست گذشته از کانال {ch_id}...", parse_mode="html") - - async def progress_notify(txt: str): - await event.edit(txt, parse_mode="html") - - await self.collector.scrape_channel_history(channel_id=ch_id, limit=limit, progress_callback=progress_notify) - await event.answer("فرآیند دریافت آغاز شد.") - - elif data.startswith("sel_trg:"): - _, post_id_str, target_id_str = data.split(":") - post_id = int(post_id_str) - target_id = int(target_id_str) - - post = await self.repo.get_post_by_id(post_id) - target = await self.repo.get_target_by_id(target_id) - if not post or not target: - await event.answer("پست یا کانال مقصد یافت نشد.", alert=True) - return - - await event.answer(f"در حال بازنویسی برای {target.title}...") - - loading_caption = ( - f"🤖 در حال بازنویسی هوشمند برای کانال: {target.title}...\n" - f"(اعمال لحن اختصاصی و حذف تگ‌های مبدا)" - ) - try: - await event.edit(loading_caption, parse_mode="html", buttons=None) - except Exception: - pass - - rewritten_text = await self.ai_processor.rewrite_for_target(post.raw_text or "", target) - cache_key = f"{post_id}:{target_id}" - self.preview_cache[cache_key] = rewritten_text - - preview_caption = ( - f"🎯 پیش‌نمایش بازنویسی شده برای: {target.title}\n" - f"🎭 شخصیت و لحن: {target.personality or 'پیش‌فرض'}\n" - f"➖➖➖➖➖➖➖➖➖➖\n\n" - f"{rewritten_text}\n\n" - f"➖➖➖➖➖➖➖➖➖➖\n" - f"آیا این متن برای صف انتشار تایید است؟" - ) - - preview_buttons = [ - [ - Button.inline(f"✅ تایید و افزودن به صف {target.title}", data=f"pub:{post_id}:{target_id}"), - ], - [ - Button.inline("🔙 انصراف / بازگشت به پست اصلی", data=f"cancel:{post_id}") - ] + # 12. Waiting for Provider Profile Name + elif action == "wait_prov_name": + name = text.strip() + self.user_states[event.sender_id] = {"action": "wait_prov_type", "prov_name": name} + type_buttons = [ + [ + Button.inline("🤖 AGY (سرور داخلی/لوکال)", data="ai_pv_type_sel:agy"), + Button.inline("⚡️ OpenAI / OrcaRouter", data="ai_pv_type_sel:openai"), + ], + [ + Button.inline("🌟 Google Gemini", data="ai_pv_type_sel:gemini"), + ], + [ + Button.inline("❌ انصراف", data="cancel_state") ] - await event.edit(preview_caption, parse_mode="html", buttons=preview_buttons) + ] + await event.reply( + f"📌 نام سرویس‌دهنده: {name}\n\n" + "اکنون نوع سرویس (Provider Type) را انتخاب کنید:", + parse_mode="html", + buttons=type_buttons + ) - elif data.startswith("pub:"): - _, post_id_str, target_id_str = data.split(":") - post_id = int(post_id_str) - target_id = int(target_id_str) + # 13. Waiting for Provider Profile Model + elif action == "wait_prov_model": + model_name = text.strip() + state["prov_model"] = model_name + self.user_states[event.sender_id] = { + "action": "wait_prov_url", + "prov_name": state["prov_name"], + "prov_type": state["prov_type"], + "prov_model": model_name + } + default_url = "http://host.docker.internal:8000/v1" if state["prov_type"] == "agy" else ("https://api.orcarouter.ai/v1" if state["prov_type"] == "openai" else "https://generativelanguage.googleapis.com/v1beta") + await event.reply( + f"🤖 مدل: {model_name}\n\n" + f"اکنون آدرس Base URL را ارسال کنید:\n" + f"(پیش‌فرض: {default_url} - برای استفاده از پیش‌فرض کلمه default یا خالی بفرستید)", + parse_mode="html", + buttons=get_cancel_button() + ) - post = await self.repo.get_post_by_id(post_id) - target = await self.repo.get_target_by_id(target_id) - if not post or not target: - await event.answer("اطلاعات یافت نشد.", alert=True) - return + # 14. Waiting for Provider Profile Base URL + elif action == "wait_prov_url": + url_val = "" if text.strip().lower() in ("default", "none", "null", "خالی") else text.strip() + state["prov_url"] = url_val + self.user_states[event.sender_id] = { + "action": "wait_prov_key", + "prov_name": state["prov_name"], + "prov_type": state["prov_type"], + "prov_model": state["prov_model"], + "prov_url": url_val + } + await event.reply( + "🔑 اکنون کلید API (API Key) را ارسال کنید:\n" + "(اگر سرور داخلی بدون کلید است، کلمه none یا خالی بفرستید)", + parse_mode="html", + buttons=get_cancel_button() + ) - cache_key = f"{post_id}:{target_id}" - text_to_publish = self.preview_cache.get(cache_key) or post.raw_text or "" + # 15. Waiting for Provider Profile API Key + elif action == "wait_prov_key": + key_val = "" if text.strip().lower() in ("none", "null", "خالی") else text.strip() + name = state["prov_name"] + ptype = state["prov_type"] + pmodel = state["prov_model"] + purl = state.get("prov_url", "") + self.user_states.pop(event.sender_id, None) - payload = { - "post_id": post.id, - "text": text_to_publish, - "media_path": post.media_path, - "target_id": target.id, - "target_title": target.title - } - if self.queue: - await self.queue.push_target_post(target.id, payload) - ADMIN_ACTIONS_TOTAL.labels(action="approved").inc() + prof_id = await self.repo.add_provider_profile( + name=name, + provider_type=ptype, + model=pmodel, + base_url=purl, + api_key=key_val, + is_active=False + ) + card, buttons = await self._render_ai_provider_detail(prof_id) + await event.reply(f"✅ سرویس‌دهنده «{name}» با موفقیت اضافه شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - await self.repo.record_post_published_to_target(post_id, target.id, target.title or "Target") + # 16. Waiting to edit specific provider model + elif action == "wait_pv_edit_model": + prof_id = state["prof_id"] + self.user_states.pop(event.sender_id, None) + await self.repo.update_provider_profile(prof_id, model=text.strip()) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_provider_detail(prof_id) + await event.reply("✅ مدل سرویس‌دهنده بروز شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - updated_post = await self.repo.get_post_by_id(post_id) - targets = await self.repo.get_active_targets() - new_caption = self._format_raw_post_caption(updated_post) - new_buttons = self._build_raw_post_keyboard(updated_post, targets) + # 17. Waiting to edit specific provider URL + elif action == "wait_pv_edit_url": + prof_id = state["prof_id"] + self.user_states.pop(event.sender_id, None) + new_url = "" if text.strip().lower() in ("default", "none", "null", "خالی") else text.strip() + await self.repo.update_provider_profile(prof_id, base_url=new_url) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_provider_detail(prof_id) + await event.reply("✅ Base URL سرویس‌دهنده بروز شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) - await event.edit( + # 18. Waiting to edit specific provider Key + elif action == "wait_pv_edit_key": + prof_id = state["prof_id"] + self.user_states.pop(event.sender_id, None) + new_key = "" if text.strip().lower() in ("none", "null", "خالی") else text.strip() + await self.repo.update_provider_profile(prof_id, api_key=new_key) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + card, buttons = await self._render_ai_provider_detail(prof_id) + await event.reply("✅ کلید API سرویس‌دهنده بروز شد!", parse_mode="html") + await event.reply(card, parse_mode="html", buttons=buttons) + + + + + async def _handle_callback(self, event: events.CallbackQuery.Event): + if not self.is_admin(event.sender_id): + await event.answer("⛔ دسترسی غیرمجاز.", alert=True) + return + + data = event.data.decode("utf-8") + + # 0. Cancel Active State + if data == "cancel_state": + self.user_states.pop(event.sender_id, None) + await event.edit("❌ عملیات لغو شد.") + await event.answer("لغو شد.") + + # --- AI Settings & Logs Callbacks --- + elif data == "ai_menu": + text, buttons = await self._render_ai_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data == "ai_prov_menu": + text, buttons = await self._render_ai_provider_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_set_prov:"): + new_prov = data.split(":")[1] + await self.repo.set_setting("ai_provider", new_prov) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + text, buttons = await self._render_ai_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer(f"سرویس‌دهنده به {new_prov} تغییر یافت!", alert=True) + + elif data == "ai_set_model": + self.user_states[event.sender_id] = {"action": "wait_ai_model"} + await event.reply( + "🏷 نام مدل مورد نظر را ارسال کنید:\n" + "(مثال: antigravity یا google/gemini-3.5-flash)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data == "ai_set_url": + self.user_states[event.sender_id] = {"action": "wait_ai_base_url"} + await event.reply( + "🌐 آدرس Base URL را ارسال کنید:\n" + "(برای پیش‌فرض یا لوکال، کلمه default یا خالی بفرستید)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data == "ai_set_key": + self.user_states[event.sender_id] = {"action": "wait_ai_api_key"} + await event.reply( + "🔑 کلید API جدید را ارسال کنید:", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data == "ai_toggle_dc": + current = await self.repo.get_setting("ai_double_check", os.getenv("AI_DOUBLE_CHECK", "false")) + new_val = "false" if current.lower() in ("true", "1", "yes") else "true" + await self.repo.set_setting("ai_double_check", new_val) + text, buttons = await self._render_ai_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer(f"بررسی مجدد: {'فعال' if new_val == 'true' else 'غیرفعال'}") + + elif data == "ai_reason_menu": + text, buttons = await self._render_ai_reasoning_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_set_reason:"): + lvl = data.split(":")[1] + val = "" if lvl == "none" else lvl + await self.repo.set_setting("ai_reasoning_effort", val) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + text, buttons = await self._render_ai_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer(f"میزان تفکر روی {lvl} تنظیم شد!") + + elif data.startswith("ai_logs:"): + limit = int(data.split(":")[1]) + text, buttons = await self._render_ai_logs(limit) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_log_view:"): + log_id = int(data.split(":")[1]) + text, buttons = await self._render_ai_log_detail(log_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data == "ai_test": + await event.answer("⏳ در حال ارسال پرامپت تستی...", alert=False) + test_prompt = "پست آزمایشی: سیستم هوشمند کپی‌کار با معماری Multi-Agent فعال است." + test_system = "متن ورودی را بازنویسی کن و خروجی را فقط به صورت JSON با کلید 'rewritten_text' برگردان." + t0 = time.time() + try: + if self.ai_processor and self.ai_processor.llm: + res = await self.ai_processor.llm.generate_json( + prompt=test_prompt, + system_prompt=test_system, + action_name="test_connection" + ) + latency = time.time() - t0 + model_used = getattr(self.ai_processor.llm, "last_used_model", self.ai_processor.llm.model) + provider_used = self.ai_processor.llm.provider + res_text = json.dumps(res, ensure_ascii=False, indent=2) + card = ( + "🧪 نتیجه تست موفقیت‌آمیز اتصال هوش مصنوعی:\n\n" + f"• 🔌 Provider: {provider_used}\n" + f"• 🏷 Model: {model_used}\n" + f"• ⏱ مدت زمان پاسخ: {latency:.2f} ثانیه\n" + f"• 📊 وضعیت: ✅ موفق (Success)\n\n" + f"📥 Input Prompt:\n{test_prompt}\n\n" + f"📤 Output JSON:\n{res_text}" + ) + else: + card = "❌ سرویس هوش مصنوعی بارگذاری نشده است." + except Exception as e: + latency = time.time() - t0 + card = ( + "🧪 نتیجه تست اتصال هوش مصنوعی:\n\n" + f"• ⏱ مدت زمان: {latency:.2f}s\n" + f"• 📊 وضعیت: ❌ خطا (Failed)\n\n" + f"❌ علت خطا:\n{e}" + ) + buttons = [ + [Button.inline("🔄 اجرای مجدد تست", data="ai_test")], + [Button.inline("🔙 بازگشت به تنظیمات AI", data="ai_menu")] + ] + await event.edit(card, parse_mode="html", buttons=buttons) + + elif data == "ai_prov_list": + text, buttons = await self._render_ai_providers_list() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_pv_view:"): + prof_id = int(data.split(":")[1]) + text, buttons = await self._render_ai_provider_detail(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_pv_activate:"): + prof_id = int(data.split(":")[1]) + await self.repo.set_active_provider_profile(prof_id) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + text, buttons = await self._render_ai_provider_detail(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer("✅ این سرویس‌دهنده با موفقیت فعال شد!") + + elif data.startswith("ai_pv_test:"): + prof_id = int(data.split(":")[1]) + prof = await self.repo.get_provider_profile_by_id(prof_id) + if not prof: + await event.answer("سرویس‌دهنده یافت نشد.", alert=True) + return + await event.answer("⏳ در حال ارسال درخواست تست به ارائه‌دهنده...") + test_client = LLMClient( + provider=prof.provider_type, + api_key=prof.api_key, + model=prof.model, + base_url=prof.base_url, + reasoning_effort=prof.reasoning_effort, + repo=self.repo + ) + t0 = time.time() + test_prompt = "Say hello and return JSON with keys: status (ok), message (Hello from Copykar), and ping (pong)." + try: + res = await test_client.generate_json( + prompt=test_prompt, + system_prompt="You are a test agent. Always reply with valid JSON only.", + action_name="test_provider_profile" + ) + latency = time.time() - t0 + res_text = json.dumps(res, ensure_ascii=False, indent=2) + card = ( + f"🧪 نتیجه تست ارائه‌دهنده «{prof.name}»:\n\n" + f"• 🔌 Provider: {prof.provider_type}\n" + f"• 🏷 Model: {prof.model}\n" + f"• ⏱ مدت زمان: {latency:.2f}s\n" + f"• 📊 وضعیت: ✅ موفق (Success)\n\n" + f"📥 Input Prompt:\n{test_prompt}\n\n" + f"📤 Output JSON:\n{res_text}" + ) + except Exception as e: + latency = time.time() - t0 + card = ( + f"🧪 نتیجه تست ارائه‌دهنده «{prof.name}»:\n\n" + f"• 🔌 Provider: {prof.provider_type}\n" + f"• 🏷 Model: {prof.model}\n" + f"• ⏱ مدت زمان: {latency:.2f}s\n" + f"• 📊 وضعیت: ❌ خطا (Failed)\n\n" + f"❌ علت خطا:\n{e}" + ) + buttons = [ + [Button.inline("🔄 اجرای مجدد تست", data=f"ai_pv_test:{prof.id}")], + [Button.inline("🔙 بازگشت به جزئیات", data=f"ai_pv_view:{prof.id}")] + ] + await event.edit(card, parse_mode="html", buttons=buttons) + + elif data.startswith("ai_pv_fb_menu:"): + prof_id = int(data.split(":")[1]) + text, buttons = await self._render_ai_provider_fallback_menu(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_pv_set_fb:"): + _, prof_id_str, fb_id_str = data.split(":") + prof_id = int(prof_id_str) + fb_id = int(fb_id_str) + target_fb_id = fb_id if fb_id > 0 else None + await self.repo.update_provider_fallback(prof_id, target_fb_id) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + text, buttons = await self._render_ai_provider_detail(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer("✅ ارائه‌دهنده پشتیبان دستی ذخیره شد.") + + elif data.startswith("ai_pv_effort:"): + prof_id = int(data.split(":")[1]) + text, buttons = await self._render_ai_provider_effort_menu(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("ai_pv_set_effort:"): + _, prof_id_str, effort_val = data.split(":") + prof_id = int(prof_id_str) + effort = "" if effort_val == "none" else effort_val + await self.repo.update_provider_reasoning_effort(prof_id, effort) + if self.ai_processor and self.ai_processor.llm: + await self.ai_processor.llm.sync_config_from_repo() + text, buttons = await self._render_ai_provider_detail(prof_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer(f"✅ قدرت استدلال روی {effort or 'خاموش'} تنظیم شد!") + + elif data.startswith("ai_pv_model:"): + prof_id = int(data.split(":")[1]) + self.user_states[event.sender_id] = {"action": "wait_pv_edit_model", "prof_id": prof_id} + await event.reply("🏷 نام مدل جدید (Model Name) را وارد کنید:", parse_mode="html", buttons=get_cancel_button()) + await event.answer() + + elif data.startswith("ai_pv_url:"): + prof_id = int(data.split(":")[1]) + self.user_states[event.sender_id] = {"action": "wait_pv_edit_url", "prof_id": prof_id} + await event.reply("🌐 آدرس Base URL جدید را وارد کنید:\n(برای پیش‌فرض کلمه default یا خالی بفرستید)", parse_mode="html", buttons=get_cancel_button()) + await event.answer() + + elif data.startswith("ai_pv_key:"): + prof_id = int(data.split(":")[1]) + self.user_states[event.sender_id] = {"action": "wait_pv_edit_key", "prof_id": prof_id} + await event.reply("🔑 کلید API جدید را وارد کنید:\n(اگر سرور بدون کلید است، کلمه none یا خالی بفرستید)", parse_mode="html", buttons=get_cancel_button()) + await event.answer() + + elif data.startswith("ai_pv_del:"): + prof_id = int(data.split(":")[1]) + await self.repo.delete_provider_profile(prof_id) + text, buttons = await self._render_ai_providers_list() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer("سرویس‌دهنده حذف شد.") + + + elif data == "ai_add_prov": + self.user_states[event.sender_id] = {"action": "wait_prov_name"} + await event.reply( + "➕ ثبت سرویس‌دهنده جدید (New AI Provider):\n\n" + "لطفا یک نام دلخواه برای این سرویس‌دهنده بفرستید:\n" + "(مثال: سرور لوکال AGY یا OpenRouter سریع یا DeepSeek)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("ai_pv_type_sel:"): + sel_type = data.split(":")[1] + state = self.user_states.get(event.sender_id, {}) + state["prov_type"] = sel_type + state["action"] = "wait_prov_model" + default_model = "antigravity" if sel_type == "agy" else ("google/gemini-2.5-flash" if sel_type == "openai" else "gemini-1.5-flash") + await event.reply( + f"📌 نوع انتخابی: {sel_type}\n\n" + f"لطفا نام مدل (Model Name) را ارسال کنید:\n" + f"(پیشنهادی: {default_model})", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("rep_gr:"): + time_range = data.split(":")[1] + await event.answer(f"⏳ در حال رسم نمودار بازه {time_range}...") + try: + chart_jpg, filename = await generate_metrics_graph(time_range) + label_map = {"15m": "۱۵ دقیقه اخیر", "3h": "۳ ساعت اخیر", "24h": "۲۴ ساعت اخیر"} + caption = ( + f"📊 نمودار متریک‌های عملکرد سیستم ({label_map.get(time_range, time_range)}):\n\n" + f"📁 {filename}\n\n" + "• 📥 نرخ دریافت و انتشار پیام\n" + "• 🧠 نرخ درخواست هوش مصنوعی\n" + "• 📬 حجم صف پیام‌ها\n" + "• ⚠️ نرخ خطاهای سیستم" + ) + chart_file = io.BytesIO(chart_jpg) + chart_file.name = filename + await self.client.send_file( + event.chat_id, + file=chart_file, + caption=caption, + parse_mode="html" + ) + except Exception as e: + logger.error(f"Failed to generate metrics graph: {e}", exc_info=True) + await event.reply(f"❌ خطا در تولید نمودار گرافیکی: {e}") + + + elif data == "rep_refresh": + report_text = await get_instant_metrics_report() + buttons = [ + [ + Button.inline("⏱ نمودار ۱۵ دقیقه اخیر", data="rep_gr:15m"), + Button.inline("⏱ نمودار ۳ ساعت اخیر", data="rep_gr:3h"), + ], + [ + Button.inline("⏱ نمودار ۲۴ ساعت اخیر", data="rep_gr:24h"), + Button.inline("🔄 بروزرسانی گزارش", data="rep_refresh"), + ] + ] + await event.edit(report_text, parse_mode="html", buttons=buttons) + await event.answer("گزارش بروز شد.") + + # --- Target Channel Config Buttons --- + elif data.startswith("st_pers:"): + target_id = int(data.split(":")[1]) + target = await self.repo.get_target_by_id(target_id) + if not target: + await event.answer("کانال مقصد یافت نشد.", alert=True) + return + self.user_states[event.sender_id] = {"action": "wait_personality", "target_id": target_id} + await event.reply( + f"🎭 تغییر لحن و شخصیت کانال «{target.title}»:\n\n" + "لطفا در پیام بعدی، لحن و استایل نگارش مورد نظرتان را بفرستید:\n" + "(مثال: لحن دوستانه و پرانرژی همراه با ایموجی و تیترهای جذاب)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("st_foot:"): + target_id = int(data.split(":")[1]) + target = await self.repo.get_target_by_id(target_id) + if not target: + await event.answer("کانال مقصد یافت نشد.", alert=True) + return + self.user_states[event.sender_id] = {"action": "wait_footer", "target_id": target_id} + await event.reply( + f"🏷 تغییر فوتر و تگ‌های اختصاصی «{target.title}»:\n\n" + "لطفا متن فوتر، آیدی یا هشتگ‌هایی که می‌خواهید انتهای هر پست قرار بگیرد را بفرستید:\n" + "(مثال: 🆔 @my_chan\n#فناوری #اخبار)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("st_cprompt:"): + target_id = int(data.split(":")[1]) + target = await self.repo.get_target_by_id(target_id) + if not target: + await event.answer("کانال مقصد یافت نشد.", alert=True) + return + self.user_states[event.sender_id] = {"action": "wait_custom_prompt", "target_id": target_id} + current_p = f"\n\nپرامپت فعلی:\n{target.custom_prompt}" if target.custom_prompt else "" + await event.reply( + f"📝 تنظیم دستورالعمل و فرامین اختصاصی برای «{target.title}»:\n\n" + "در پیام بعدی، فرامین و قوانین اختصاصی مد نظرتان برای شکل‌دهی به پیام‌های این کانال را بفرستید.\n" + "(مثال: همیشه تیتر را بولد کن، خلاصه‌ای در ۳ خط بنویس، هشتگ‌های تخصصی بزن و...)\n\n" + "💡 برای حذف پرامپت و بازگشت به حالت پیش‌فرض، عبارت /clear را ارسال کنید." + f"{current_p}", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + + elif data.startswith("st_intv:"): + target_id = int(data.split(":")[1]) + target = await self.repo.get_target_by_id(target_id) + if not target: + await event.answer("کانال مقصد یافت نشد.", alert=True) + return + self.user_states[event.sender_id] = {"action": "wait_interval", "target_id": target_id} + await event.reply( + f"⏱ تغییر فاصله ارسال کانال «{target.title}»:\n\n" + "فاصله زمانی بین ارسال هر پست را به دقیقه بفرستید:\n" + "(مثال: 30 برای ارسال هر نیم ساعت یک پست)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("st_lang:"): + target_id = int(data.split(":")[1]) + text, buttons = await self._render_target_language_menu(target_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("set_lang:"): + _, target_id_str, lang_code = data.split(":") + target_id = int(target_id_str) + if lang_code in SUPPORTED_LANGUAGES: + await self.repo.update_target_language(target_id, lang_code) + card_text, buttons = await self._render_target_config(target_id) + lang_label = SUPPORTED_LANGUAGES[lang_code]["label"] + await event.edit(card_text, parse_mode="html", buttons=buttons) + await event.answer(f"زبان کانال روی {lang_label} تنظیم شد!", alert=True) + else: + await event.answer("زبان نامعتبر است.", alert=True) + + elif data.startswith("st_slp:"): + target_id = int(data.split(":")[1]) + target = await self.repo.get_target_by_id(target_id) + if not target: + await event.answer("کانال مقصد یافت نشد.", alert=True) + return + self.user_states[event.sender_id] = {"action": "wait_sleep", "target_id": target_id} + await event.reply( + f"🌙 تنظیم ساعت خواب کانال «{target.title}»:\n\n" + "ساعت شروع و پایان خواب را با یک فاصله بفرستید:\n" + "(مثال: 23 8 برای عدم ارسال پست از ساعت ۲۳:۰۰ تا ۰۸:۰۰ صبح)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("errfix:"): + index = int(data.split(":")[1]) + group = self.error_group_cache.get(index) + if not group: + await event.answer("این گزارش منقضی شده. دوباره /errors را بزنید.", alert=True) + return + service_name, error_type = group + fixed = await self.repo.resolve_errors( + service_name=service_name, error_type=error_type, note=f"marked fixed by admin {event.sender_id}" + ) + ERRORS_RESOLVED_TOTAL.labels(service=service_name, error_type=error_type).inc(fixed) + text, buttons = await self._render_error_report() + await event.edit(text, parse_mode="html", buttons=buttons or None) + await event.answer(f"✅ {fixed} خطای «{error_type}» رفع‌شده علامت خورد.") + + elif data == "errfix_all": + fixed = await self.repo.resolve_errors(note=f"bulk marked fixed by admin {event.sender_id}") + ERRORS_RESOLVED_TOTAL.labels(service="all", error_type="all").inc(fixed) + text, buttons = await self._render_error_report() + await event.edit(text, parse_mode="html", buttons=buttons or None) + await event.answer(f"✅ {fixed} خطا رفع‌شده علامت خورد.") + + elif data == "list_src": + text, buttons = await self._render_source_list() + await event.edit(text, parse_mode="html", buttons=buttons or None) + await event.answer() + + elif data.startswith("src_view:"): + source_id = int(data.split(":")[1]) + card, buttons = await self._render_source_config(source_id) + await event.edit(card, parse_mode="html", buttons=buttons or None) + await event.answer() + + elif data.startswith("trg_view:"): + target_id = int(data.split(":")[1]) + card, buttons = await self._render_target_config(target_id) + await event.edit(card, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("auto_src:"): + target_id = int(data.split(":")[1]) + text, buttons = await self._render_auto_sources(target_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("auto_tg:"): + _, target_id_str, source_channel_str = data.split(":") + target_id = int(target_id_str) + enabled = await self.repo.toggle_target_auto_source(target_id, int(source_channel_str)) + text, buttons = await self._render_auto_sources(target_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer("✅ ارسال خودکار روشن شد." if enabled else "🛑 ارسال خودکار خاموش شد.") + + + elif data.startswith("dis_slp:"): + target_id = int(data.split(":")[1]) + await self.repo.update_target_schedule(target_id=target_id, is_sleep_enabled=False) + card, buttons = await self._render_target_config(target_id) + await event.edit(card, parse_mode="html", buttons=buttons) + await event.answer("☀️ ساعت خواب غیرفعال شد.") + + elif data.startswith("del_trg:"): + target_id = int(data.split(":")[1]) + await self.repo.delete_target(target_id) + text, buttons = await self._render_target_list() + await event.edit(f"🗑 کانال مقصد حذف شد.\n\n{text}", parse_mode="html", buttons=buttons or None) + await event.answer("کانال مقصد حذف شد.") + + elif data.startswith("del_src:"): + source_id = int(data.split(":")[1]) + await self.repo.delete_source(source_id) + text, buttons = await self._render_source_list() + await event.edit(f"🗑 کانال مبدا حذف شد.\n\n{text}", parse_mode="html", buttons=buttons or None) + await event.answer("کانال مبدا حذف شد.") + + elif data == "list_trg": + text, buttons = await self._render_target_list() + await event.edit(text, parse_mode="html", buttons=buttons or None) + await event.answer() + + # --- Category Callbacks --- + elif data.startswith("src_cat_view:"): + cat_id = int(data.split(":")[1]) + text, buttons = await self._render_source_channels_in_category(cat_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("trg_cat_view:"): + cat_id = int(data.split(":")[1]) + text, buttons = await self._render_target_channels_in_category(cat_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data == "src_all": + text, buttons = await self._render_source_all_flat() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data == "trg_all": + text, buttons = await self._render_target_all_flat() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data == "list_cat": + text, buttons = await self._render_categories_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("cat_view:"): + cat_id = int(data.split(":")[1]) + text, buttons = await self._render_category_detail(cat_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + + elif data == "cat_add": + self.user_states[event.sender_id] = {"action": "wait_cat_name"} + await event.reply( + "➕ ایجاد دسته‌بندی جدید:\n\n" + "لطفا یک نام برای دسته‌بندی جدید ارسال کنید:\n" + "(مثال: اخبار فناوری، رمزارز و کریپتو، سینما و سرگرمی و...)", + parse_mode="html", + buttons=get_cancel_button() + ) + await event.answer() + + elif data.startswith("cat_rename:"): + cat_id = int(data.split(":")[1]) + self.user_states[event.sender_id] = {"action": "wait_cat_rename", "cat_id": cat_id} + await event.reply("✏️ نام جدید دسته‌بندی را ارسال کنید:", parse_mode="html", buttons=get_cancel_button()) + await event.answer() + + elif data.startswith("cat_del:"): + cat_id = int(data.split(":")[1]) + await self.repo.delete_category(cat_id) + text, buttons = await self._render_categories_menu() + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer("دسته‌بندی حذف شد.") + + elif data.startswith("src_cat:"): + src_id = int(data.split(":")[1]) + text, buttons = await self._render_set_channel_category_menu("src", src_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("trg_cat:"): + trg_id = int(data.split(":")[1]) + text, buttons = await self._render_set_channel_category_menu("trg", trg_id) + await event.edit(text, parse_mode="html", buttons=buttons) + await event.answer() + + elif data.startswith("cat_set:"): + _, ch_type, ch_id_str, cat_id_str = data.split(":") + ch_id = int(ch_id_str) + cat_id = None if cat_id_str == "0" else int(cat_id_str) + if ch_type == "src": + await self.repo.set_source_category(ch_id, cat_id) + card, buttons = await self._render_source_config(ch_id) + await event.edit(card, parse_mode="html", buttons=buttons) + else: + await self.repo.set_target_category(ch_id, cat_id) + card, buttons = await self._render_target_config(ch_id) + await event.edit(card, parse_mode="html", buttons=buttons) + await event.answer("✅ دسته‌بندی کانال با موفقیت به روز شد!") + + + # --- Review Channel Buttons --- + elif data.startswith("hist:"): + _, ch_id_str, limit_str = data.split(":") + ch_id = int(ch_id_str) + limit = int(limit_str) + + if not self.collector: + await event.answer("کالکتور در دسترس نیست.", alert=True) + return + + await event.edit(f"⏳ در حال دریافت {limit} پست گذشته از کانال {ch_id}...", parse_mode="html") + + async def progress_notify(txt: str): + await event.edit(txt, parse_mode="html") + + await self.collector.scrape_channel_history(channel_id=ch_id, limit=limit, progress_callback=progress_notify) + await event.answer("فرآیند دریافت پایان یافت.") + + elif data.startswith("histn:"): + ch_id = int(data.split(":")[1]) + self.user_states[event.sender_id] = {"action": "wait_fetch_count", "channel_id": ch_id} + await event.reply( + f"🔢 تعداد پست‌های مورد نظر برای کانال {ch_id}:\n\n" + f"یک عدد بین {MIN_FETCH_LIMIT} تا {MAX_FETCH_LIMIT} بفرستید.\n" + "(مثال: 100 برای دریافت ۱۰۰ پست آخر)", + parse_mode="html", + buttons=get_cancel_button(), + ) + await event.answer() + + elif data.startswith("sel_trg:"): + _, post_id_str, target_id_str = data.split(":") + post_id = int(post_id_str) + target_id = int(target_id_str) + + post = await self.repo.get_post_by_id(post_id) + target = await self.repo.get_target_by_id(target_id) + if not post or not target: + await event.answer("پست یا کانال مقصد یافت نشد.", alert=True) + return + + await event.answer(f"در حال بازنویسی برای {target.title}...") + + loading_caption = ( + f"🤖 در حال بازنویسی هوشمند برای کانال: {target.title}...\n" + f"(اعمال لحن اختصاصی و حذف تگ‌های مبدا)" + ) + try: + await event.edit(loading_caption, parse_mode="html", buttons=None) + except Exception: + pass + + rewrite_res = await self.ai_processor.rewrite_for_target( + post.raw_text or "", target, has_media=bool(post.media_path), image_path=post.media_path + ) + rewritten_text = str(rewrite_res) + cache_key = f"{post_id}:{target_id}" + self.preview_cache[cache_key] = rewritten_text + + ai_warning = "" + if getattr(rewrite_res, "is_rejected", False): + reason = getattr(rewrite_res, "rejection_reason", "نامناسب برای این کانال") + ai_warning = ( + f"⚠️ هشدار هوش مصنوعی: این پست توسط هوش مصنوعی رد شده است!\n" + f"📝 علت رد: {reason}\n" + f"➖➖➖➖➖➖➖➖➖➖\n\n" + ) + + model_name = getattr(self.ai_processor.llm, "last_used_model", getattr(self.ai_processor.llm, "model", "AI")) if self.ai_processor else "AI" + preview_caption = ( + f"🎯 پیش‌نمایش بازنویسی شده برای: {target.title}\n" + f"🤖 مدل استفاده‌شده: {model_name}\n" + f"🎭 شخصیت و لحن: {target.personality or 'پیش‌فرض'}\n" + f"➖➖➖➖➖➖➖➖➖➖\n" + f"{ai_warning}" + f"{rewritten_text}\n\n" + f"➖➖➖➖➖➖➖➖➖➖\n" + f"آیا این متن برای صف انتشار تایید است؟" + ) + + preview_buttons = [ + [ + Button.inline(f"✅ تایید و افزودن به صف {target.title}", data=f"pub:{post_id}:{target_id}"), + ], + [ + Button.inline("🔙 انصراف / بازگشت به پست اصلی", data=f"cancel:{post_id}") + ] + ] + await event.edit( + clamp_for_telegram(preview_caption, bool(post.media_path and os.path.exists(post.media_path))), + parse_mode="html", + buttons=preview_buttons, + ) + + elif data.startswith("pub:"): + _, post_id_str, target_id_str = data.split(":") + post_id = int(post_id_str) + target_id = int(target_id_str) + + post = await self.repo.get_post_by_id(post_id) + target = await self.repo.get_target_by_id(target_id) + if not post or not target: + await event.answer("اطلاعات یافت نشد.", alert=True) + return + + cache_key = f"{post_id}:{target_id}" + text_to_publish = self.preview_cache.get(cache_key) or post.raw_text or "" + + payload = { + "post_id": post.id, + "text": text_to_publish, + "media_path": post.media_path, + "target_id": target.id, + "target_title": target.title + } + if self.queue: + await self.queue.push_target_post(target.id, payload) + ADMIN_ACTIONS_TOTAL.labels(action="approved").inc() + + await self.repo.record_post_queued_to_target(post_id, target.id, target.title or "Target") + + updated_post = await self.repo.get_post_by_id(post_id) + targets = await self.repo.get_active_targets() + new_caption = self._format_raw_post_caption(updated_post) + new_buttons = self._build_raw_post_keyboard(updated_post, targets) + + await event.edit( + clamp_for_telegram( f"✅ به صف انتشار کانال {target.title} اضافه شد!\n\n{new_caption}", - parse_mode="html", - buttons=new_buttons - ) - await event.answer(f"به صف {target.title} افزوده شد!") + bool(updated_post.media_path), + ), + parse_mode="html", + buttons=new_buttons + ) + await event.answer(f"به صف {target.title} افزوده شد!") - elif data.startswith("cancel:"): - _, post_id_str = data.split(":") - post_id = int(post_id_str) + elif data.startswith("cancel:"): + _, post_id_str = data.split(":") + post_id = int(post_id_str) - post = await self.repo.get_post_by_id(post_id) - if not post: - return + post = await self.repo.get_post_by_id(post_id) + if not post: + return - targets = await self.repo.get_active_targets() - caption = self._format_raw_post_caption(post) - buttons = self._build_raw_post_keyboard(post, targets) + targets = await self.repo.get_active_targets() + caption = self._format_raw_post_caption(post) + buttons = self._build_raw_post_keyboard(post, targets) - await event.edit(caption, parse_mode="html", buttons=buttons) - await event.answer("پیش‌نمایش لغو شد.") + await event.edit(caption, parse_mode="html", buttons=buttons) + await event.answer("پیش‌نمایش لغو شد.") - elif data.startswith("rej:"): - _, post_id_str = data.split(":") - post_id = int(post_id_str) + elif data.startswith("rej:"): + _, post_id_str = data.split(":") + post_id = int(post_id_str) - await self.repo.reject_post(post_id) - ADMIN_ACTIONS_TOTAL.labels(action="rejected").inc() + await self.repo.reject_post(post_id) + ADMIN_ACTIONS_TOTAL.labels(action="rejected").inc() - await event.edit( - f"{event.text}\n\n❌ این پست توسط ادمین رد و بایگانی شد.", - parse_mode="html", - buttons=None - ) - await event.answer("پست بایگانی شد.") + try: + original = await event.get_message() + original_text = (original.text or "") if original else "" + except Exception: + original_text = "" - elif data.startswith("del_msg:"): - _, post_id_str = data.split(":") - post_id = int(post_id_str) + await event.edit( + clamp_for_telegram( + f"{original_text}\n\n❌ این پست توسط ادمین رد و بایگانی شد.", False + ), + parse_mode="html", + buttons=None + ) + await event.answer("پست بایگانی شد.") - await self.repo.soft_delete_post(post_id) - try: - await event.delete() - except Exception as e: - logger.error(f"Could not delete message from review channel: {e}") - await event.edit("🗑 این پیام حذف و بایگانی شد.", parse_mode="html", buttons=None) + elif data.startswith("del_msg:"): + _, post_id_str = data.split(":") + post_id = int(post_id_str) - await event.answer("🗑 پیام از کانال ادمین حذف شد.", alert=False) + await self.repo.soft_delete_post(post_id) + try: + await event.delete() + except Exception as e: + logger.error(f"Could not delete message from review channel: {e}") + await event.edit("🗑 این پیام حذف و بایگانی شد.", parse_mode="html", buttons=None) + + await event.answer("🗑 پیام از کانال ادمین حذف شد.", alert=False) + + async def _error_metrics_loop(self): + """Keep the open-error gauges fresh for Prometheus without an admin opening /errors.""" + while self._running: + try: + await self.refresh_error_metrics() + except Exception as e: + logger.debug(f"Could not refresh error metrics: {e}") + await asyncio.sleep(30) async def stop(self): + self._running = False + if self._metrics_task: + self._metrics_task.cancel() if self.client.is_connected(): await self.client.disconnect() logger.info("Admin Bot disconnected.") diff --git a/tests/test_admin_menus.py b/tests/test_admin_menus.py new file mode 100644 index 0000000..b456e91 --- /dev/null +++ b/tests/test_admin_menus.py @@ -0,0 +1,89 @@ +"""Channel lists render as one button menu, and each button opens that channel's card.""" +import asyncio +import sys + +sys.path.insert(0, "/app") + +from db.database import init_db, close_db_pool +from db.repository import Repository +from services.admin_bot import AdminBotService + +SRC = -1009999000030 +TRG = -1009999000031 + + +async def _cleanup(repo): + pool = await repo._get_pool() + async with pool.acquire() as conn: + await conn.execute("DELETE FROM posts WHERE source_channel_id = $1;", SRC) + await conn.execute("DELETE FROM sources WHERE channel_id = $1;", SRC) + await conn.execute("DELETE FROM targets WHERE channel_id = $1;", TRG) + + +def button_data(rows): + return [b.data.decode() for row in rows for b in row] + + +async def run_tests(): + await init_db() + repo = Repository() + await _cleanup(repo) + try: + await _run_assertions(repo) + finally: + await _cleanup(repo) + await close_db_pool() + print("All admin menu tests passed successfully!") + + +async def _run_assertions(repo): + bot = AdminBotService(repo=repo, review_channel_id=0, admin_user_ids=[1]) + + src_id = await repo.add_source(SRC, "Menu Source", "menusrc") + target_id = await repo.add_target(TRG, "Menu Target", None) + await repo.create_raw_post(SRC, 1, "a post", content_hash="menu_1") + + # Source list: one button per source, pointing at its detail view. + text, buttons = await bot._render_source_list() + assert "Menu Source" not in text, "the list message itself should stay a header, not a card dump" + data = button_data(buttons) + assert f"src_view:{src_id}" in data, f"source button missing: {data}" + + # Target list likewise, and it flags auto-routing. + text, buttons = await bot._render_target_list() + data = button_data(buttons) + assert f"trg_view:{target_id}" in data, f"target button missing: {data}" + def label_for(rows, want): + for row in rows: + for b in row: + if b.data.decode() == want: + return b.text + raise AssertionError(f"button {want} not found") + + assert "🤖" not in label_for(buttons, f"trg_view:{target_id}"), "no auto-route configured yet" + await repo.set_target_auto_sources(target_id, [SRC]) + _, buttons = await bot._render_target_list() + assert "🤖" in label_for(buttons, f"trg_view:{target_id}"), "auto-routed target should be flagged" + + # Source detail card carries fetch buttons and a way back. + card, buttons = await bot._render_source_config(src_id) + assert "Menu Source" in card and str(SRC) in card + assert "1" in card, f"collected count missing from card: {card}" + assert "Menu Target" in card, "auto-routed target should be listed on the source card" + data = button_data(buttons) + assert f"hist:{SRC}:20" in data and f"histn:{SRC}" in data, f"fetch buttons missing: {data}" + assert "list_src" in data, f"back button missing: {data}" + + # Target detail card can get back to its list too. + card, buttons = await bot._render_target_config(target_id) + data = button_data(buttons) + assert f"auto_src:{target_id}" in data, f"auto-route button missing: {data}" + assert "list_trg" in data, f"back button missing: {data}" + + # Missing ids degrade gracefully instead of raising. + card, buttons = await bot._render_source_config(999999) + assert "یافت نشد" in card and buttons == [] + + +if __name__ == "__main__": + asyncio.run(run_tests()) diff --git a/tests/test_error_tracking.py b/tests/test_error_tracking.py new file mode 100644 index 0000000..53831bd --- /dev/null +++ b/tests/test_error_tracking.py @@ -0,0 +1,91 @@ +"""Errors can be grouped, marked fixed, and reported as metrics.""" +import asyncio +import sys + +sys.path.insert(0, "/app") + +from db.database import init_db, close_db_pool +from db.repository import Repository +from services.admin_bot import AdminBotService, clamp_for_telegram, MAX_MEDIA_CAPTION, MAX_TEXT_MESSAGE + +SVC = "test.tracking_service" + + +async def _cleanup(repo): + pool = await repo._get_pool() + async with pool.acquire() as conn: + await conn.execute("DELETE FROM error_logs WHERE service_name LIKE 'test.tracking%';") + + +async def _log(repo, service, etype, msg): + pool = await repo._get_pool() + async with pool.acquire() as conn: + await conn.execute( + "INSERT INTO error_logs (service_name, error_type, error_message) VALUES ($1,$2,$3);", + service, etype, msg) + + +def test_clamping(): + # Short text is untouched. + assert clamp_for_telegram("hello", True) == "hello" + + long_text = "x" * 5000 + media = clamp_for_telegram(long_text, has_media=True) + assert len(media) <= MAX_MEDIA_CAPTION, f"media caption still {len(media)} chars" + assert media.endswith(""), "truncation marker missing" + + plain = clamp_for_telegram(long_text, has_media=False) + assert len(plain) <= MAX_TEXT_MESSAGE, f"text message still {len(plain)} chars" + + # A message that exactly fits must not be altered. + exact = "y" * MAX_MEDIA_CAPTION + assert clamp_for_telegram(exact, True) == exact + + +async def run_tests(): + test_clamping() + + await init_db() + repo = Repository() + await _cleanup(repo) + + await _log(repo, SVC, "ValueError", "bad value one") + await _log(repo, SVC, "ValueError", "bad value two") + await _log(repo, SVC, "KeyError", "missing key") + await _log(repo, SVC + "_b", "ValueError", "other service") + + groups = await repo.get_open_error_summary(limit=50) + mine = {(g["service_name"], g["error_type"]): g for g in groups if g["service_name"].startswith("test.tracking")} + assert len(mine) == 3, f"expected 3 groups, got {list(mine)}" + assert mine[(SVC, "ValueError")]["occurrences"] == 2 + assert mine[(SVC, "ValueError")]["last_message"] == "bad value two", "last_message should be the newest" + + # Resolving one group leaves the others untouched. + fixed = await repo.resolve_errors(service_name=SVC, error_type="ValueError", note="fixed in tests") + assert fixed == 2, f"expected 2 rows resolved, got {fixed}" + groups = await repo.get_open_error_summary(limit=50) + mine = {(g["service_name"], g["error_type"]) for g in groups if g["service_name"].startswith("test.tracking")} + assert (SVC, "ValueError") not in mine + assert (SVC, "KeyError") in mine and (SVC + "_b", "ValueError") in mine + + # Resolving again is a no-op, not a double count. + assert await repo.resolve_errors(service_name=SVC, error_type="ValueError") == 0 + + # The gauge reflects what is still open. + bot = AdminBotService(repo=repo, review_channel_id=0, admin_user_ids=[1]) + total = await bot.refresh_error_metrics() + assert total == await repo.count_open_errors(), "gauge total must match the open count" + + # The report renders and caches groups for the callback buttons. + text, buttons = await bot._render_error_report() + assert "KeyError" in text + assert bot.error_group_cache, "groups must be cached for errfix callbacks" + assert any(b.data.decode().startswith("errfix:") for row in buttons for b in row) + + await _cleanup(repo) + await close_db_pool() + print("All error-tracking tests passed successfully!") + + +if __name__ == "__main__": + asyncio.run(run_tests()) diff --git a/tests/test_fallback_and_stop.py b/tests/test_fallback_and_stop.py new file mode 100644 index 0000000..cc0d759 --- /dev/null +++ b/tests/test_fallback_and_stop.py @@ -0,0 +1,146 @@ +import asyncio +from unittest.mock import AsyncMock, patch, MagicMock +from db.models import AIProviderProfile +from core.llm import LLMClient +from services.admin_bot import get_persian_main_menu, AdminBotService + + +def test_main_menu_pause_toggle(): + menu_running = get_persian_main_menu(is_paused=False) + button_texts_running = [getattr(b.button, 'text', '') for row in menu_running for b in row] + assert "🛑 توقف اضطراری سیستم" in button_texts_running + assert "▶️ راه‌اندازی و ادامه سیستم" not in button_texts_running + + menu_paused = get_persian_main_menu(is_paused=True) + button_texts_paused = [getattr(b.button, 'text', '') for row in menu_paused for b in row] + assert "▶️ راه‌اندازی و ادامه سیستم" in button_texts_paused + assert "🛑 توقف اضطراری سیستم" not in button_texts_paused + + +async def test_llm_fallback_chain_success_on_fallback(): + # Setup 2 providers: Primary (fails) -> Fallback (succeeds) + profile_1 = AIProviderProfile( + id=1, + name="Primary AI", + provider_type="openai", + model="gpt-4o", + base_url="https://api.openai.com/v1", + api_key="key1", + is_active=True, + fallback_provider_id=2 + ) + profile_2 = AIProviderProfile( + id=2, + name="Backup AI", + provider_type="gemini", + model="gemini-1.5-flash", + base_url="", + api_key="key2", + is_active=False, + fallback_provider_id=None + ) + + repo_mock = AsyncMock() + repo_mock.get_active_provider_profile.return_value = profile_1 + repo_mock.get_provider_profiles.return_value = [profile_1, profile_2] + repo_mock.record_ai_log = AsyncMock() + + fallback_alerts = [] + + async def on_fallback(from_p, to_p, err, step): + fallback_alerts.append((from_p.name, to_p.name, err, step)) + + client = LLMClient( + repo=repo_mock, + on_fallback_alert=on_fallback, + ) + client.max_retries_per_model = 0 # fail fast for test + + # Mock _call_openai to fail and _call_gemini to succeed + with patch.object(client, "_call_openai", side_effect=RuntimeError("Primary Provider Timeout")): + with patch.object(client, "_call_gemini", return_value={"decision": "accept", "rewritten_text": "Success from backup"}): + res = await client.generate_json("Test prompt", system_prompt="Test sys") + assert res["decision"] == "accept" + assert res["rewritten_text"] == "Success from backup" + assert len(fallback_alerts) == 1 + assert fallback_alerts[0][0] == "Primary AI" + assert fallback_alerts[0][1] == "Backup AI" + assert "Primary Provider Timeout" in fallback_alerts[0][2] + assert fallback_alerts[0][3] == 1 + + +async def test_llm_fallback_chain_all_fail_alert(): + profile_1 = AIProviderProfile( + id=1, + name="Primary AI", + provider_type="openai", + model="gpt-4o", + is_active=True, + fallback_provider_id=2 + ) + profile_2 = AIProviderProfile( + id=2, + name="Backup AI 1", + provider_type="openai", + model="gpt-4o-mini", + is_active=False, + fallback_provider_id=3 + ) + profile_3 = AIProviderProfile( + id=3, + name="Backup AI 2", + provider_type="gemini", + model="gemini-1.5-flash", + is_active=False, + fallback_provider_id=None + ) + + repo_mock = AsyncMock() + repo_mock.get_active_provider_profile.return_value = profile_1 + repo_mock.get_provider_profiles.return_value = [profile_1, profile_2, profile_3] + repo_mock.record_ai_log = AsyncMock() + + fallback_alerts = [] + chain_failure_alerts = [] + + async def on_fallback(from_p, to_p, err, step): + fallback_alerts.append((from_p.name, to_p.name, err, step)) + + async def on_chain_failure(chain, err): + chain_failure_alerts.append((len(chain), err)) + + client = LLMClient( + repo=repo_mock, + on_fallback_alert=on_fallback, + on_chain_failure_alert=on_chain_failure, + ) + client.max_retries_per_model = 0 + + with patch.object(client, "_call_openai", side_effect=RuntimeError("OpenAI Error")): + with patch.object(client, "_call_gemini", side_effect=RuntimeError("Gemini Quota Exceeded")): + try: + await client.generate_json("Test prompt") + assert False, "Should have raised exception" + except Exception as e: + assert "Gemini Quota Exceeded" in str(e) or "OpenAI Error" in str(e) + + # Check fallbacks: 1 -> 2, 2 -> 3 + assert len(fallback_alerts) == 2 + assert fallback_alerts[0][0] == "Primary AI" + assert fallback_alerts[0][1] == "Backup AI 1" + assert fallback_alerts[1][0] == "Backup AI 1" + assert fallback_alerts[1][1] == "Backup AI 2" + + # Check chain failure alert triggered + assert len(chain_failure_alerts) == 1 + assert chain_failure_alerts[0][0] == 3 + + +async def main(): + test_main_menu_pause_toggle() + await test_llm_fallback_chain_success_on_fallback() + await test_llm_fallback_chain_all_fail_alert() + print("All system stop and multi-provider AI fallback chain tests passed successfully!") + +if __name__ == "__main__": + asyncio.run(main())