diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..fd11181 --- /dev/null +++ b/.env.example @@ -0,0 +1,9 @@ +# Required +BOT_TOKEN=your_telegram_bot_token_here +SUPER_ADMIN_ID=your_telegram_id_here + +# Optional +PAYMENT_CARD=your_card_number +FORCE_CHANNEL=@your_channel +WEB_PORT=10674 +BASE_URL=http://your-server:10674 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..3ff1462 --- /dev/null +++ b/.gitignore @@ -0,0 +1,12 @@ +__pycache__/ +*.pyc +*.pyo +.env +config.json +bot.db +*.sqlite +uploads/ +logs/ +backups/ +temp/ +.yt-dlp/ diff --git a/README.md b/README.md new file mode 100644 index 0000000..9eb7b9a --- /dev/null +++ b/README.md @@ -0,0 +1,79 @@ +# Telegram Bot v4.0 — Modular Architecture + +A feature-rich Telegram bot with VPN config distribution, media downloading, proxy service, file-to-link, VIP subscriptions, and full admin panel. + +## Features + +- **VPN Config Distribution** — Free and VIP tiers with cooldown and usage tracking +- **Media Download** — YouTube, Instagram, TikTok, Twitter video/audio download with queue system +- **Proxy Service** — Aggregates free HTTP/HTTPS/SOCKS4/SOCKS5 proxies from multiple sources, auto-refreshes, REST API endpoint, live proxy checker +- **Speed Test** — One-click download speed test using Cloudflare and Hetzner endpoints +- **File to Link** — Upload any file, get a direct download link (7-day expiry) +- **VIP System** — Silver/Gold/Diamond plans with payment receipts and wallet purchases +- **Wallet** — Balance management, referral bonuses, VIP purchases +- **Referral System** — Invite friends, earn wallet credit +- **Admin Panel** — User management, broadcast, stats, config management, discount codes, wallet charge, weekly analytics dashboard +- **HTTP Streaming Server** — Direct file downloads via HTTP, proxy list API + +## Project Structure + +``` +bot/ +├── __init__.py # Package init +├── config.py # Configuration (env vars, config.json, constants) +├── database.py # Async SQLite database layer (aiosqlite) +├── helpers.py # Utilities (formatting, QR, rate limiting, validation) +├── keyboards.py # Inline keyboard builders +├── decorators.py # Guard decorator (auth, rate limit, ban check) +├── server.py # HTTP streaming server + proxy API +├── main.py # App entry point, handler registration +└── handlers/ + ├── start.py # /start command + ├── main_menu.py # Profile, help, referral, claim, VIP, wallet + ├── youtube.py # Media download (YouTube, Instagram, TikTok) + ├── file.py # File-to-link + ├── proxy.py # Proxy service (fetch, check, list, manage) + └── admin.py # Admin panel (configs, users, broadcast, stats) +``` + +## Setup + +1. Clone the repo +2. Install dependencies: + ```bash + pip install -r requirements.txt + ``` +3. Copy `.env.example` to `.env` and fill in your values: + ```bash + cp .env.example .env + ``` +4. Run: + ```bash + python run.py + ``` + +## Configuration + +Edit `config.json` (auto-generated on first run) to customize: +- VIP plan pricing and durations +- Download limits and queue sizes +- Rate limiting +- Proxy update intervals +- Referral bonuses + +## Proxy API + +The bot exposes a REST API at `/api/proxies` for programmatic proxy access: +``` +GET /api/proxies?protocol=http&limit=50 +``` + +## Architecture Highlights + +- **Async database** via `aiosqlite` — non-blocking DB operations +- **Persistent DB connection** with WAL mode and 8 MB cache +- **Separated concerns** — each module handles one domain +- **Guard decorator** — centralized auth, rate limiting, ban checking +- **Queue system** for media downloads with progress tracking +- **Concurrent updates** enabled for better throughput +- **Periodic jobs** — proxy refresh, daily reports, file cleanup diff --git a/bot/__init__.py b/bot/__init__.py new file mode 100644 index 0000000..128e6d1 --- /dev/null +++ b/bot/__init__.py @@ -0,0 +1,6 @@ +""" +Telegram Bot - Restructured v4.0 +Modular architecture with separated DB, proxy support, and optimized performance. +""" + +__version__ = "4.0.0" diff --git a/bot/config.py b/bot/config.py new file mode 100644 index 0000000..4702820 --- /dev/null +++ b/bot/config.py @@ -0,0 +1,138 @@ +""" +Centralized configuration loader. +Reads from .env and config.json, merges defaults, and exposes typed constants. +""" + +import os +import sys +import json +from pathlib import Path + +from dotenv import load_dotenv + +load_dotenv() + +# ── Telegram ────────────────────────────────────── +BOT_TOKEN: str = os.getenv("BOT_TOKEN", "") +SUPER_ADMIN_ID: int = int(os.getenv("SUPER_ADMIN_ID", "0")) +PAYMENT_CARD: str = os.getenv("PAYMENT_CARD", "") +FORCE_CHANNEL: str = os.getenv("FORCE_CHANNEL", "") + +if not BOT_TOKEN: + print("BOT_TOKEN is not set in .env") + sys.exit(1) + +# ── Web / streaming server ──────────────────────── +WEB_SERVER_HOST: str = "0.0.0.0" +WEB_SERVER_PORT: int = int(os.getenv("WEB_PORT", "10674")) +BASE_DOWNLOAD_URL: str = os.getenv("BASE_URL", f"http://localhost:{WEB_SERVER_PORT}") +COOKIE_FILE: str = os.path.join(os.path.expanduser("~"), ".yt-dlp", "cookies.txt") + +# ── Directories ─────────────────────────────────── +for _d in ("uploads", "logs", "backups", "temp"): + Path(_d).mkdir(exist_ok=True) +Path(COOKIE_FILE).parent.mkdir(exist_ok=True) + +# ── Config JSON ─────────────────────────────────── +DEFAULT_CONFIG: dict = { + "database": "bot.db", + "claim_cooldown_hours": 6, + "max_yt_size_mb": 500, + "max_file_size_mb": 2000, + "yt_sleep_interval": 2, + "max_queue_size": 50, + "referral_bonus": 5000, + "enable_qr_code": True, + "welcome_message": "به ربات خوش اومدی!", + "rate_limit_per_minute": 25, + "max_broadcast_delay": 0.05, + "daily_report_hour": 8, + "proxy_update_interval_minutes": 30, + "proxy_check_timeout": 8, + "proxy_max_results": 30, + "vip_levels": { + "silver": { + "name": "Silver", + "price": 50000, + "days": 30, + "max_yt_quality": "720", + "max_configs": 3, + "wallet_bonus": 0, + }, + "gold": { + "name": "Gold", + "price": 120000, + "days": 90, + "max_yt_quality": "1080", + "max_configs": 5, + "wallet_bonus": 10000, + }, + "diamond": { + "name": "Diamond", + "price": 250000, + "days": 180, + "max_yt_quality": "1080", + "max_configs": 10, + "wallet_bonus": 30000, + }, + }, + "discount_codes": {}, +} + +CONFIG_PATH = Path("config.json") +if not CONFIG_PATH.exists(): + CONFIG_PATH.write_text(json.dumps(DEFAULT_CONFIG, indent=4, ensure_ascii=False), encoding="utf-8") + +with open(CONFIG_PATH, "r", encoding="utf-8") as _f: + CONFIG: dict = json.load(_f) + +for _k, _v in DEFAULT_CONFIG.items(): + CONFIG.setdefault(_k, _v) + +# ── Typed shortcuts ─────────────────────────────── +DB_NAME: str = CONFIG["database"] +CLAIM_COOLDOWN: int = CONFIG["claim_cooldown_hours"] +MAX_YT_SIZE_MB: int = CONFIG["max_yt_size_mb"] +MAX_FILE_SIZE: int = CONFIG["max_file_size_mb"] +YT_SLEEP: int = CONFIG["yt_sleep_interval"] +MAX_QUEUE: int = CONFIG["max_queue_size"] +REFERRAL_BONUS: int = CONFIG["referral_bonus"] +ENABLE_QR: bool = CONFIG["enable_qr_code"] +WELCOME_MSG: str = CONFIG["welcome_message"] +VIP_LEVELS: dict = CONFIG["vip_levels"] +RATE_LIMIT: int = CONFIG.get("rate_limit_per_minute", 25) +DISCOUNT_CODES: dict = CONFIG.get("discount_codes", {}) +PROXY_UPDATE_INTERVAL: int = CONFIG.get("proxy_update_interval_minutes", 30) +PROXY_CHECK_TIMEOUT: int = CONFIG.get("proxy_check_timeout", 8) +PROXY_MAX_RESULTS: int = CONFIG.get("proxy_max_results", 30) + +VALID_PERMS = frozenset({ + "can_manage_configs", + "can_manage_payments", + "can_broadcast", + "can_block_users", + "can_give_vip", +}) + +# ── Conversation states ─────────────────────────── +( + STATE_NONE, + STATE_WAITING_RECEIPT, + STATE_ADDING_CONFIG, + STATE_DELETING_CONFIG, + STATE_BROADCASTING, + STATE_WAITING_FILE, + STATE_WAITING_YT_URL, + STATE_MANAGE_ADMIN, + STATE_SET_COOKIE, + STATE_WAITING_CONFIG_TEXT, + STATE_SEARCHING_USER, + STATE_GIVE_VIP, + STATE_BAN_USER, + STATE_UNBAN_USER, + STATE_WAITING_DISCOUNT, + STATE_ADD_DISCOUNT, + STATE_WALLET_CHARGE, + STATE_CHECKING_PROXY, + STATE_ADMIN_WALLET_CHARGE, +) = range(19) diff --git a/bot/database.py b/bot/database.py new file mode 100644 index 0000000..b36a62a --- /dev/null +++ b/bot/database.py @@ -0,0 +1,846 @@ +""" +Async database layer using aiosqlite. +All DB operations are non-blocking; connection pooling via a persistent connection. +""" + +import shutil +import logging +import asyncio +from datetime import datetime, timedelta +from typing import Optional, List, Tuple + +import aiosqlite + +from bot.config import ( + DB_NAME, + SUPER_ADMIN_ID, + CLAIM_COOLDOWN, + VIP_LEVELS, + REFERRAL_BONUS, + VALID_PERMS, +) +from bot.helpers import generate_secure_code + +logger = logging.getLogger("BOT.db") + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS users ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + telegram_id INTEGER UNIQUE NOT NULL, + username TEXT DEFAULT '', + full_name TEXT DEFAULT '', + role TEXT DEFAULT 'user', + is_vip INTEGER DEFAULT 0, + vip_level TEXT DEFAULT 'silver', + vip_expiry TEXT, + last_claim TEXT, + total_claims INTEGER DEFAULT 0, + balance INTEGER DEFAULT 0, + total_spent INTEGER DEFAULT 0, + joined_at TEXT DEFAULT (datetime('now')), + banned INTEGER DEFAULT 0, + ban_reason TEXT, + referral_code TEXT UNIQUE, + referred_by INTEGER, + total_referrals INTEGER DEFAULT 0, + last_seen TEXT DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS admin_permissions ( + admin_id INTEGER PRIMARY KEY, + can_manage_configs INTEGER DEFAULT 0, + can_manage_payments INTEGER DEFAULT 0, + can_broadcast INTEGER DEFAULT 0, + can_block_users INTEGER DEFAULT 0, + can_give_vip INTEGER DEFAULT 0 +); + +CREATE TABLE IF NOT EXISTS configs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + config_text TEXT NOT NULL, + category TEXT DEFAULT 'free', + protocol TEXT DEFAULT 'unknown', + remark TEXT DEFAULT '', + created_at TEXT DEFAULT (datetime('now')), + expires_at TEXT, + usage_count INTEGER DEFAULT 0, + active INTEGER DEFAULT 1 +); + +CREATE TABLE IF NOT EXISTS file_links ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + file_unique_id TEXT UNIQUE NOT NULL, + file_id TEXT NOT NULL, + file_type TEXT NOT NULL, + file_name TEXT DEFAULT 'file', + file_size INTEGER DEFAULT 0, + mime_type TEXT DEFAULT 'application/octet-stream', + uploader_id INTEGER, + created_at TEXT DEFAULT (datetime('now')), + expires_at TEXT, + download_count INTEGER DEFAULT 0 +); + +CREATE TABLE IF NOT EXISTS payments ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER NOT NULL, + amount INTEGER NOT NULL, + plan TEXT NOT NULL, + status TEXT DEFAULT 'pending', + receipt_file TEXT, + discount_code TEXT, + created_at TEXT DEFAULT (datetime('now')), + reviewed_at TEXT, + confirmed_by INTEGER +); + +CREATE TABLE IF NOT EXISTS wallet_transactions ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER NOT NULL, + amount INTEGER NOT NULL, + type TEXT NOT NULL, + description TEXT DEFAULT '', + created_at TEXT DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS discount_codes ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + code TEXT UNIQUE NOT NULL, + percent INTEGER NOT NULL, + max_uses INTEGER DEFAULT 0, + used_count INTEGER DEFAULT 0, + expires_at TEXT, + created_at TEXT DEFAULT (datetime('now')), + active INTEGER DEFAULT 1 +); + +CREATE TABLE IF NOT EXISTS bot_stats ( + key TEXT PRIMARY KEY, + value INTEGER DEFAULT 0 +); + +CREATE TABLE IF NOT EXISTS user_logs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER, + action TEXT, + details TEXT, + created_at TEXT DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS broadcasts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + admin_id INTEGER, + message TEXT, + total_sent INTEGER DEFAULT 0, + total_failed INTEGER DEFAULT 0, + created_at TEXT DEFAULT (datetime('now')) +); + +CREATE TABLE IF NOT EXISTS proxies ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ip TEXT NOT NULL, + port INTEGER NOT NULL, + protocol TEXT DEFAULT 'http', + country TEXT DEFAULT '', + speed_ms INTEGER DEFAULT 0, + anonymity TEXT DEFAULT 'unknown', + alive INTEGER DEFAULT 1, + last_check TEXT DEFAULT (datetime('now')), + created_at TEXT DEFAULT (datetime('now')), + UNIQUE(ip, port, protocol) +); + +CREATE INDEX IF NOT EXISTS idx_users_tg ON users(telegram_id); +CREATE INDEX IF NOT EXISTS idx_configs_cat ON configs(category, active); +CREATE INDEX IF NOT EXISTS idx_payments_status ON payments(status); +CREATE INDEX IF NOT EXISTS idx_files_uid ON file_links(file_unique_id); +CREATE INDEX IF NOT EXISTS idx_proxies_alive ON proxies(alive, protocol); +""" + +_STAT_KEYS = ( + "total_users", + "total_downloads", + "total_claims", + "total_payments", + "total_revenue", +) + + +class Database: + """Async SQLite database with connection pooling via a persistent connection.""" + + def __init__(self) -> None: + self._db: Optional[aiosqlite.Connection] = None + self._lock = asyncio.Lock() + + async def connect(self) -> aiosqlite.Connection: + if self._db is None: + self._db = await aiosqlite.connect(DB_NAME, timeout=30) + self._db.row_factory = aiosqlite.Row + await self._db.execute("PRAGMA journal_mode=WAL") + await self._db.execute("PRAGMA foreign_keys=ON") + await self._db.execute("PRAGMA synchronous=NORMAL") + await self._db.execute("PRAGMA cache_size=-8000") # 8 MB cache + return self._db + + async def init_db(self) -> None: + db = await self.connect() + await db.executescript(_SCHEMA) + + # migrations + for col, default in [ + ("vip_level", "'silver'"), + ("total_spent", "0"), + ]: + try: + await db.execute(f"ALTER TABLE users ADD COLUMN {col} TEXT DEFAULT {default}") + except Exception: + pass + try: + await db.execute("ALTER TABLE payments ADD COLUMN discount_code TEXT") + except Exception: + pass + + for key in _STAT_KEYS: + await db.execute("INSERT OR IGNORE INTO bot_stats (key, value) VALUES (?, 0)", (key,)) + + if SUPER_ADMIN_ID: + await db.execute( + "INSERT OR IGNORE INTO users (telegram_id, role, referral_code) VALUES (?, 'super_admin', ?)", + (SUPER_ADMIN_ID, generate_secure_code()), + ) + await db.execute( + "INSERT OR IGNORE INTO admin_permissions " + "(admin_id, can_manage_configs, can_manage_payments, can_broadcast, can_block_users, can_give_vip) " + "VALUES (?, 1, 1, 1, 1, 1)", + (SUPER_ADMIN_ID,), + ) + await db.commit() + + async def close(self) -> None: + if self._db: + await self._db.close() + self._db = None + + # ── helpers ──────────────────────────────────── + + async def _fetchone(self, sql: str, params: tuple = ()) -> Optional[dict]: + db = await self.connect() + async with db.execute(sql, params) as cur: + row = await cur.fetchone() + return dict(row) if row else None + + async def _fetchall(self, sql: str, params: tuple = ()) -> List[dict]: + db = await self.connect() + async with db.execute(sql, params) as cur: + return [dict(r) for r in await cur.fetchall()] + + async def _execute(self, sql: str, params: tuple = ()) -> int: + db = await self.connect() + cur = await db.execute(sql, params) + await db.commit() + return cur.rowcount + + async def _insert(self, sql: str, params: tuple = ()) -> int: + db = await self.connect() + cur = await db.execute(sql, params) + await db.commit() + return cur.lastrowid + + # ── Users ───────────────────────────────────── + + async def get_user(self, tg_id: int) -> Optional[dict]: + return await self._fetchone("SELECT * FROM users WHERE telegram_id=?", (tg_id,)) + + async def create_user( + self, tg_id: int, username: str, full_name: str, ref_code: str = None + ) -> Tuple[dict, Optional[int]]: + db = await self.connect() + async with self._lock: + existing = await self._fetchone("SELECT 1 FROM users WHERE telegram_id=?", (tg_id,)) + if existing: + await self._execute( + "UPDATE users SET username=?, full_name=?, last_seen=datetime('now') WHERE telegram_id=?", + (username or "", full_name or "", tg_id), + ) + return await self.get_user(tg_id), None + + code = generate_secure_code() + referred_by = None + + if ref_code and not ref_code.startswith("dl_"): + ref_row = await self._fetchone( + "SELECT telegram_id FROM users WHERE referral_code=?", (ref_code,) + ) + if ref_row and ref_row["telegram_id"] != tg_id: + referred_by = ref_row["telegram_id"] + await db.execute( + "UPDATE users SET balance=balance+?, total_referrals=total_referrals+1 WHERE telegram_id=?", + (REFERRAL_BONUS, referred_by), + ) + await db.execute( + "INSERT INTO wallet_transactions (user_id, amount, type, description) VALUES (?, ?, ?, ?)", + (referred_by, REFERRAL_BONUS, "credit", "Referral bonus"), + ) + + role = "super_admin" if tg_id == SUPER_ADMIN_ID else "user" + await db.execute( + "INSERT INTO users (telegram_id, username, full_name, referral_code, referred_by, role) " + "VALUES (?, ?, ?, ?, ?, ?)", + (tg_id, username or "", full_name or "", code, referred_by, role), + ) + await db.execute("UPDATE bot_stats SET value=value+1 WHERE key='total_users'") + await db.commit() + return await self.get_user(tg_id), referred_by + + async def update_last_seen(self, tg_id: int) -> None: + await self._execute( + "UPDATE users SET last_seen=datetime('now') WHERE telegram_id=?", (tg_id,) + ) + + async def is_vip(self, tg_id: int) -> bool: + user = await self.get_user(tg_id) + if not user or not user["is_vip"]: + return False + if user["vip_expiry"]: + try: + if datetime.now() > datetime.fromisoformat(user["vip_expiry"]): + await self._execute( + "UPDATE users SET is_vip=0, vip_expiry=NULL WHERE telegram_id=?", + (tg_id,), + ) + return False + except ValueError: + return False + return True + + async def get_vip_level(self, tg_id: int) -> Optional[str]: + if not await self.is_vip(tg_id): + return None + user = await self.get_user(tg_id) + return user.get("vip_level", "silver") if user else None + + async def can_claim(self, tg_id: int) -> Tuple[bool, int]: + user = await self.get_user(tg_id) + if not user or not user["last_claim"]: + return True, 0 + try: + last = datetime.fromisoformat(user["last_claim"]) + except ValueError: + return True, 0 + passed = (datetime.now() - last).total_seconds() / 3600 + if passed >= CLAIM_COOLDOWN: + return True, 0 + remaining = int((CLAIM_COOLDOWN - passed) * 3600) + return False, remaining + + async def update_claim(self, tg_id: int) -> None: + db = await self.connect() + await db.execute( + "UPDATE users SET last_claim=datetime('now'), total_claims=total_claims+1 WHERE telegram_id=?", + (tg_id,), + ) + await db.execute("UPDATE bot_stats SET value=value+1 WHERE key='total_claims'") + await db.commit() + + async def has_perm(self, tg_id: int, perm: str) -> bool: + if perm not in VALID_PERMS: + return False + if tg_id == SUPER_ADMIN_ID: + return True + row = await self._fetchone( + f"SELECT {perm} FROM admin_permissions WHERE admin_id=?", (tg_id,) + ) + return bool(row and row[perm]) + + async def get_role(self, tg_id: int) -> str: + user = await self.get_user(tg_id) + return user["role"] if user else "user" + + async def ban_user(self, tg_id: int, reason: str = "") -> None: + await self._execute( + "UPDATE users SET banned=1, ban_reason=? WHERE telegram_id=?", (reason, tg_id) + ) + + async def unban_user(self, tg_id: int) -> bool: + return ( + await self._execute( + "UPDATE users SET banned=0, ban_reason=NULL WHERE telegram_id=?", (tg_id,) + ) + > 0 + ) + + async def get_banned_users(self, limit: int = 50) -> List[dict]: + return await self._fetchall( + "SELECT * FROM users WHERE banned=1 ORDER BY id DESC LIMIT ?", (limit,) + ) + + async def search_user(self, query: str) -> List[dict]: + if query.isdigit(): + return await self._fetchall("SELECT * FROM users WHERE telegram_id=?", (int(query),)) + pattern = f"%{query}%" + return await self._fetchall( + "SELECT * FROM users WHERE username LIKE ? OR full_name LIKE ? LIMIT 10", + (pattern, pattern), + ) + + async def get_all_users(self) -> List[int]: + rows = await self._fetchall("SELECT telegram_id FROM users WHERE banned=0") + return [r["telegram_id"] for r in rows] + + # ── Wallet ──────────────────────────────────── + + async def add_balance(self, tg_id: int, amount: int, description: str = "") -> int: + db = await self.connect() + await db.execute("UPDATE users SET balance=balance+? WHERE telegram_id=?", (amount, tg_id)) + await db.execute( + "INSERT INTO wallet_transactions (user_id, amount, type, description) VALUES (?, ?, ?, ?)", + (tg_id, amount, "credit", description), + ) + await db.commit() + row = await self._fetchone("SELECT balance FROM users WHERE telegram_id=?", (tg_id,)) + return row["balance"] if row else 0 + + async def deduct_balance(self, tg_id: int, amount: int, description: str = "") -> bool: + db = await self.connect() + async with self._lock: + row = await self._fetchone("SELECT balance FROM users WHERE telegram_id=?", (tg_id,)) + if not row or row["balance"] < amount: + return False + await db.execute( + "UPDATE users SET balance=balance-?, total_spent=total_spent+? WHERE telegram_id=?", + (amount, amount, tg_id), + ) + await db.execute( + "INSERT INTO wallet_transactions (user_id, amount, type, description) VALUES (?, ?, ?, ?)", + (tg_id, -amount, "debit", description), + ) + await db.commit() + return True + + async def get_wallet_history(self, tg_id: int, limit: int = 10) -> List[dict]: + return await self._fetchall( + "SELECT * FROM wallet_transactions WHERE user_id=? ORDER BY created_at DESC LIMIT ?", + (tg_id, limit), + ) + + # ── Discount codes ──────────────────────────── + + async def add_discount_code( + self, code: str, percent: int, max_uses: int = 0, days_valid: int = 30 + ) -> bool: + expires = (datetime.now() + timedelta(days=days_valid)).isoformat() + try: + await self._insert( + "INSERT INTO discount_codes (code, percent, max_uses, expires_at) VALUES (?, ?, ?, ?)", + (code.upper(), percent, max_uses, expires), + ) + return True + except Exception: + return False + + async def validate_discount(self, code: str) -> Optional[dict]: + d = await self._fetchone( + "SELECT * FROM discount_codes WHERE code=? AND active=1", (code.upper(),) + ) + if not d: + return None + if d["expires_at"] and datetime.now() > datetime.fromisoformat(d["expires_at"]): + return None + if d["max_uses"] > 0 and d["used_count"] >= d["max_uses"]: + return None + return d + + async def use_discount(self, code: str) -> None: + await self._execute( + "UPDATE discount_codes SET used_count=used_count+1 WHERE code=?", (code.upper(),) + ) + + async def list_discount_codes(self) -> List[dict]: + return await self._fetchall( + "SELECT * FROM discount_codes ORDER BY created_at DESC LIMIT 20" + ) + + # ── Configs ─────────────────────────────────── + + async def add_config(self, text: str, category: str = "free", remark: str = "") -> int: + proto = text.split("://")[0] if "://" in text else "unknown" + expiry = (datetime.now() + timedelta(days=30)).isoformat() + return await self._insert( + "INSERT INTO configs (config_text, category, protocol, remark, expires_at) VALUES (?, ?, ?, ?, ?)", + (text, category, proto, remark, expiry), + ) + + async def delete_config(self, cfg_id: int) -> bool: + return await self._execute("DELETE FROM configs WHERE id=?", (cfg_id,)) > 0 + + async def get_all_configs( + self, limit: int = 20, offset: int = 0, category: str = None + ) -> List[dict]: + if category: + return await self._fetchall( + "SELECT * FROM configs WHERE category=? ORDER BY id DESC LIMIT ? OFFSET ?", + (category, limit, offset), + ) + return await self._fetchall( + "SELECT * FROM configs ORDER BY id DESC LIMIT ? OFFSET ?", (limit, offset) + ) + + async def get_config(self, cfg_id: int) -> Optional[dict]: + return await self._fetchone("SELECT * FROM configs WHERE id=?", (cfg_id,)) + + async def get_active_configs(self, category: str, limit: int = 3) -> List[dict]: + now = datetime.now().isoformat() + return await self._fetchall( + "SELECT * FROM configs WHERE active=1 AND category=? AND expires_at>? " + "ORDER BY usage_count ASC, id DESC LIMIT ?", + (category, now, limit), + ) + + async def count_configs(self, category: str = None) -> int: + if category: + row = await self._fetchone( + "SELECT COUNT(*) AS cnt FROM configs WHERE category=?", (category,) + ) + else: + row = await self._fetchone("SELECT COUNT(*) AS cnt FROM configs") + return row["cnt"] if row else 0 + + async def increment_usage(self, cfg_id: int) -> None: + await self._execute("UPDATE configs SET usage_count=usage_count+1 WHERE id=?", (cfg_id,)) + + # ── Files ───────────────────────────────────── + + async def save_file( + self, unique_id: str, file_id: str, ftype: str, + name: str, size: int, mime: str, uploader_id: int = None + ) -> bool: + expires = (datetime.now() + timedelta(days=7)).isoformat() + await self._execute( + "INSERT OR IGNORE INTO file_links " + "(file_unique_id, file_id, file_type, file_name, file_size, mime_type, uploader_id, expires_at) " + "VALUES (?, ?, ?, ?, ?, ?, ?, ?)", + (unique_id, file_id, ftype, name, size, mime, uploader_id, expires), + ) + return True + + async def get_file(self, unique_id: str) -> Optional[dict]: + return await self._fetchone( + "SELECT * FROM file_links WHERE file_unique_id=?", (unique_id,) + ) + + async def increment_file_downloads(self, unique_id: str) -> None: + await self._execute( + "UPDATE file_links SET download_count=download_count+1 WHERE file_unique_id=?", + (unique_id,), + ) + + async def clean_expired_files(self) -> None: + await self._execute("DELETE FROM file_links WHERE expires_at < datetime('now')") + + # ── Payments ────────────────────────────────── + + async def create_payment( + self, user_id: int, amount: int, plan: str, receipt: str, discount_code: str = None + ) -> int: + return await self._insert( + "INSERT INTO payments (user_id, amount, plan, receipt_file, discount_code) VALUES (?, ?, ?, ?, ?)", + (user_id, amount, plan, receipt, discount_code), + ) + + async def get_pending_payments(self) -> List[dict]: + return await self._fetchall( + "SELECT p.*, u.username, u.full_name " + "FROM payments p JOIN users u ON p.user_id=u.telegram_id " + "WHERE p.status='pending' ORDER BY p.created_at DESC" + ) + + async def get_payment(self, pay_id: int) -> Optional[dict]: + return await self._fetchone("SELECT * FROM payments WHERE id=?", (pay_id,)) + + async def confirm_payment(self, pay_id: int, admin_id: int) -> Optional[dict]: + db = await self.connect() + async with self._lock: + pay = await self._fetchone( + "SELECT * FROM payments WHERE id=? AND status='pending'", (pay_id,) + ) + if not pay: + return None + plan_info = VIP_LEVELS.get(pay["plan"], {}) + days = plan_info.get("days", 30) + wallet_bonus = plan_info.get("wallet_bonus", 0) + + user = await self._fetchone( + "SELECT vip_expiry, is_vip FROM users WHERE telegram_id=?", (pay["user_id"],) + ) + if user and user["is_vip"] and user["vip_expiry"]: + try: + base = max(datetime.fromisoformat(user["vip_expiry"]), datetime.now()) + except Exception: + base = datetime.now() + else: + base = datetime.now() + expiry = (base + timedelta(days=days)).isoformat() + + await db.execute( + "UPDATE payments SET status='confirmed', confirmed_by=?, reviewed_at=datetime('now') WHERE id=?", + (admin_id, pay_id), + ) + await db.execute( + "UPDATE users SET is_vip=1, vip_level=?, vip_expiry=?, total_spent=total_spent+? WHERE telegram_id=?", + (pay["plan"], expiry, pay["amount"], pay["user_id"]), + ) + if wallet_bonus > 0: + await db.execute( + "UPDATE users SET balance=balance+? WHERE telegram_id=?", + (wallet_bonus, pay["user_id"]), + ) + await db.execute( + "INSERT INTO wallet_transactions (user_id, amount, type, description) VALUES (?, ?, ?, ?)", + (pay["user_id"], wallet_bonus, "credit", f"VIP purchase bonus ({pay['plan']})"), + ) + await db.execute("UPDATE bot_stats SET value=value+1 WHERE key='total_payments'") + await db.execute( + "UPDATE bot_stats SET value=value+? WHERE key='total_revenue'", + (pay["amount"],), + ) + await db.commit() + + result = dict(pay) + result["wallet_bonus"] = wallet_bonus + result["vip_expiry"] = expiry + result["plan_info"] = plan_info + return result + + async def reject_payment(self, pay_id: int, admin_id: int) -> Optional[dict]: + pay = await self._fetchone( + "SELECT * FROM payments WHERE id=? AND status='pending'", (pay_id,) + ) + if not pay: + return None + await self._execute( + "UPDATE payments SET status='rejected', confirmed_by=?, reviewed_at=datetime('now') WHERE id=?", + (admin_id, pay_id), + ) + return dict(pay) + + # ── Stats ───────────────────────────────────── + + async def get_stats(self) -> dict: + today = datetime.now().strftime("%Y-%m-%d") + db = await self.connect() + + async def _count(sql, params=()): + async with db.execute(sql, params) as c: + r = await c.fetchone() + return r[0] if r else 0 + + return { + "total_users": await _count("SELECT COUNT(*) FROM users"), + "today_users": await _count( + "SELECT COUNT(*) FROM users WHERE joined_at LIKE ?", (f"{today}%",) + ), + "total_vips": await _count("SELECT COUNT(*) FROM users WHERE is_vip=1"), + "total_banned": await _count("SELECT COUNT(*) FROM users WHERE banned=1"), + "total_configs_free": await _count( + "SELECT COUNT(*) FROM configs WHERE category='free' AND active=1" + ), + "total_configs_vip": await _count( + "SELECT COUNT(*) FROM configs WHERE category='vip' AND active=1" + ), + "total_downloads": await _count( + "SELECT COALESCE(value,0) FROM bot_stats WHERE key='total_downloads'" + ), + "total_claims": await _count( + "SELECT COALESCE(value,0) FROM bot_stats WHERE key='total_claims'" + ), + "total_revenue": await _count( + "SELECT COALESCE(value,0) FROM bot_stats WHERE key='total_revenue'" + ), + "total_payments": await _count( + "SELECT COALESCE(value,0) FROM bot_stats WHERE key='total_payments'" + ), + "pending_payments": await _count( + "SELECT COUNT(*) FROM payments WHERE status='pending'" + ), + "total_wallet_balance": await _count( + "SELECT COALESCE(SUM(balance),0) FROM users" + ), + } + + async def get_daily_stats(self) -> dict: + today = datetime.now().strftime("%Y-%m-%d") + yesterday = (datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d") + db = await self.connect() + + async def _count(sql, params=()): + async with db.execute(sql, params) as c: + r = await c.fetchone() + return r[0] if r else 0 + + return { + "new_users_today": await _count( + "SELECT COUNT(*) FROM users WHERE joined_at LIKE ?", (f"{today}%",) + ), + "new_users_yesterday": await _count( + "SELECT COUNT(*) FROM users WHERE joined_at LIKE ?", (f"{yesterday}%",) + ), + "payments_today": await _count( + "SELECT COUNT(*) FROM payments WHERE status='confirmed' AND reviewed_at LIKE ?", + (f"{today}%",), + ), + "revenue_today": await _count( + "SELECT COALESCE(SUM(amount),0) FROM payments WHERE status='confirmed' AND reviewed_at LIKE ?", + (f"{today}%",), + ), + "downloads_today": await _count( + "SELECT COUNT(*) FROM user_logs WHERE action='yt_download' AND created_at LIKE ?", + (f"{today}%",), + ), + "claims_today": await _count( + "SELECT COUNT(*) FROM user_logs WHERE action='claim_config' AND created_at LIKE ?", + (f"{today}%",), + ), + "active_vips": await _count("SELECT COUNT(*) FROM users WHERE is_vip=1"), + "total_users": await _count("SELECT COUNT(*) FROM users"), + } + + async def get_weekly_stats(self) -> dict: + db = await self.connect() + days = [] + for i in range(7): + d = (datetime.now() - timedelta(days=i)).strftime("%Y-%m-%d") + async with db.execute( + "SELECT COUNT(*) FROM users WHERE joined_at LIKE ?", (f"{d}%",) + ) as c: + r = await c.fetchone() + days.append({"date": d, "new_users": r[0] if r else 0}) + + async def _count(sql, params=()): + async with db.execute(sql, params) as c: + r = await c.fetchone() + return r[0] if r else 0 + + top_actions = await self._fetchall( + "SELECT action, COUNT(*) as cnt FROM user_logs " + "WHERE created_at > datetime('now', '-7 days') " + "GROUP BY action ORDER BY cnt DESC LIMIT 5" + ) + return { + "daily_breakdown": days, + "week_revenue": await _count( + "SELECT COALESCE(SUM(amount),0) FROM payments " + "WHERE status='confirmed' AND reviewed_at > datetime('now', '-7 days')" + ), + "week_downloads": await _count( + "SELECT COUNT(*) FROM user_logs WHERE action='yt_download' AND created_at > datetime('now', '-7 days')" + ), + "week_claims": await _count( + "SELECT COUNT(*) FROM user_logs WHERE action='claim_config' AND created_at > datetime('now', '-7 days')" + ), + "top_actions": top_actions, + } + + # ── Logs ────────────────────────────────────── + + async def add_log(self, user_id: int, action: str, details: str = "") -> None: + await self._insert( + "INSERT INTO user_logs (user_id, action, details) VALUES (?, ?, ?)", + (user_id, action, details), + ) + + async def add_broadcast(self, admin_id: int, msg: str, sent: int, failed: int) -> None: + await self._insert( + "INSERT INTO broadcasts (admin_id, message, total_sent, total_failed) VALUES (?, ?, ?, ?)", + (admin_id, msg, sent, failed), + ) + + # ── VIP ─────────────────────────────────────── + + async def give_vip(self, tg_id: int, days: int, level: str = "silver") -> None: + user = await self._fetchone( + "SELECT vip_expiry, is_vip FROM users WHERE telegram_id=?", (tg_id,) + ) + if user and user["is_vip"] and user["vip_expiry"]: + try: + base = max(datetime.fromisoformat(user["vip_expiry"]), datetime.now()) + except Exception: + base = datetime.now() + else: + base = datetime.now() + expiry = (base + timedelta(days=days)).isoformat() + await self._execute( + "UPDATE users SET is_vip=1, vip_level=?, vip_expiry=? WHERE telegram_id=?", + (level, expiry, tg_id), + ) + + async def remove_admin(self, tg_id: int) -> bool: + count = await self._execute("DELETE FROM admin_permissions WHERE admin_id=?", (tg_id,)) + await self._execute("UPDATE users SET role='user' WHERE telegram_id=?", (tg_id,)) + return count > 0 + + # ── Proxies ─────────────────────────────────── + + async def upsert_proxies(self, proxy_list: List[dict]) -> int: + db = await self.connect() + count = 0 + for p in proxy_list: + try: + await db.execute( + "INSERT INTO proxies (ip, port, protocol, country, speed_ms, anonymity, alive, last_check) " + "VALUES (?, ?, ?, ?, ?, ?, 1, datetime('now')) " + "ON CONFLICT(ip, port, protocol) DO UPDATE SET " + "alive=1, speed_ms=excluded.speed_ms, last_check=datetime('now'), country=excluded.country, anonymity=excluded.anonymity", + (p["ip"], p["port"], p.get("protocol", "http"), + p.get("country", ""), p.get("speed_ms", 0), p.get("anonymity", "unknown")), + ) + count += 1 + except Exception: + pass + await db.commit() + return count + + async def get_alive_proxies( + self, protocol: str = None, limit: int = 30 + ) -> List[dict]: + if protocol: + return await self._fetchall( + "SELECT * FROM proxies WHERE alive=1 AND protocol=? ORDER BY speed_ms ASC LIMIT ?", + (protocol, limit), + ) + return await self._fetchall( + "SELECT * FROM proxies WHERE alive=1 ORDER BY speed_ms ASC LIMIT ?", (limit,) + ) + + async def mark_proxy_dead(self, proxy_id: int) -> None: + await self._execute("UPDATE proxies SET alive=0 WHERE id=?", (proxy_id,)) + + async def get_proxy_stats(self) -> dict: + db = await self.connect() + + async def _count(sql, params=()): + async with db.execute(sql, params) as c: + r = await c.fetchone() + return r[0] if r else 0 + + return { + "total": await _count("SELECT COUNT(*) FROM proxies"), + "alive": await _count("SELECT COUNT(*) FROM proxies WHERE alive=1"), + "http": await _count("SELECT COUNT(*) FROM proxies WHERE alive=1 AND protocol='http'"), + "https": await _count("SELECT COUNT(*) FROM proxies WHERE alive=1 AND protocol='https'"), + "socks4": await _count("SELECT COUNT(*) FROM proxies WHERE alive=1 AND protocol='socks4'"), + "socks5": await _count("SELECT COUNT(*) FROM proxies WHERE alive=1 AND protocol='socks5'"), + } + + async def clean_dead_proxies(self, older_than_hours: int = 24) -> int: + cutoff = (datetime.now() - timedelta(hours=older_than_hours)).isoformat() + return await self._execute( + "DELETE FROM proxies WHERE alive=0 AND last_check < ?", (cutoff,) + ) + + # ── Backup ──────────────────────────────────── + + def backup(self) -> str: + dst = f"backups/db_{datetime.now():%Y%m%d_%H%M%S}.sqlite" + shutil.copy2(DB_NAME, dst) + return dst diff --git a/bot/decorators.py b/bot/decorators.py new file mode 100644 index 0000000..c52707a --- /dev/null +++ b/bot/decorators.py @@ -0,0 +1,74 @@ +""" +Guard decorator: ban check, rate limit, forced channel membership, permission check. +""" + +from functools import wraps + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.ext import ContextTypes + +from bot.config import FORCE_CHANNEL, SUPER_ADMIN_ID +from bot.helpers import rate_limit_check + + +def guard(perm: str = None): + def decorator(func): + @wraps(func) + async def wrapper(update: Update, context: ContextTypes.DEFAULT_TYPE, *args, **kwargs): + uid = update.effective_user.id if update.effective_user else 0 + db = context.bot_data["db"] + + user = await db.get_user(uid) + if user and user.get("banned"): + text = "Your account is suspended." + if update.callback_query: + await update.callback_query.answer(text, show_alert=True) + else: + await update.message.reply_text(text) + return + + if rate_limit_check(uid): + text = "Too many requests. Please wait." + if update.callback_query: + await update.callback_query.answer(text, show_alert=True) + else: + await update.message.reply_text(text) + return + + if FORCE_CHANNEL and uid != SUPER_ADMIN_ID: + role = await db.get_role(uid) + if role not in ("admin", "super_admin"): + try: + member = await context.bot.get_chat_member(FORCE_CHANNEL, uid) + if member.status in ("left", "kicked"): + raise Exception + except Exception: + kb = InlineKeyboardMarkup([[ + InlineKeyboardButton( + "Join Channel", + url=f"https://t.me/{FORCE_CHANNEL.lstrip('@')}", + ), + InlineKeyboardButton("I Joined", callback_data="check_join"), + ]]) + text = f"Please join {FORCE_CHANNEL} first." + if update.callback_query: + await update.callback_query.answer( + "Join the channel first!", show_alert=True + ) + else: + await update.message.reply_text(text, reply_markup=kb) + return + + if perm and not await db.has_perm(uid, perm): + text = "Permission denied." + if update.callback_query: + await update.callback_query.answer(text, show_alert=True) + else: + await update.message.reply_text(text) + return + + await db.update_last_seen(uid) + return await func(update, context, *args, **kwargs) + + return wrapper + return decorator diff --git a/bot/handlers/__init__.py b/bot/handlers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/bot/handlers/admin.py b/bot/handlers/admin.py new file mode 100644 index 0000000..81fc939 --- /dev/null +++ b/bot/handlers/admin.py @@ -0,0 +1,833 @@ +""" +Admin panel handlers: config management, user management, broadcast, +VIP grants, discount codes, backup, cookie, stats. +""" + +import os +import asyncio +from datetime import datetime, timedelta + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.constants import ParseMode +from telegram.error import BadRequest +from telegram.ext import ContextTypes, ConversationHandler + +from bot.config import ( + SUPER_ADMIN_ID, + COOKIE_FILE, + VIP_LEVELS, + CONFIG, + STATE_WAITING_CONFIG_TEXT, + STATE_DELETING_CONFIG, + STATE_BROADCASTING, + STATE_SET_COOKIE, + STATE_BAN_USER, + STATE_UNBAN_USER, + STATE_SEARCHING_USER, + STATE_GIVE_VIP, + STATE_MANAGE_ADMIN, + STATE_ADD_DISCOUNT, + STATE_ADMIN_WALLET_CHARGE, +) +from bot.decorators import guard +from bot.helpers import safe_edit, sanitize_text, is_valid_tg_id, fmt_time, make_progress_bar +from bot.keyboards import kb_admin, kb_cancel, kb_back_main + + +@guard() +async def cb_admin(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + uid = q.from_user.id if q else update.effective_user.id + db = context.bot_data["db"] + role = await db.get_role(uid) + if role not in ("admin", "super_admin") and uid != SUPER_ADMIN_ID: + if q: + await q.answer("Permission denied.", show_alert=True) + return + if q: + await q.answer() + stats = await db.get_stats() + text = ( + "*Admin Panel*\n" + "--------------------\n" + f"Users: {stats['total_users']:,} | Today: {stats['today_users']:,}\n" + f"VIP: {stats['total_vips']:,} | Banned: {stats['total_banned']:,}\n" + f"Revenue: {stats['total_revenue']:,}\n" + f"Pending payments: {stats['pending_payments']:,}\n" + f"Free configs: {stats['total_configs_free']:,}\n" + f"VIP configs: {stats['total_configs_vip']:,}" + ) + if q: + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + else: + await update.message.reply_text(text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + + +async def cmd_admin(update: Update, context: ContextTypes.DEFAULT_TYPE): + db = context.bot_data["db"] + user = await db.get_user(update.effective_user.id) + if user and user.get("banned"): + return + await cb_admin(update, context) + + +# ── Config management ───────────────────────────── + +@guard(perm="can_manage_configs") +async def cb_admin_configs(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + stats = await db.get_stats() + text = ( + "*Config Management*\n\n" + f"Free: {stats['total_configs_free']}\n" + f"VIP: {stats['total_configs_vip']}" + ) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=InlineKeyboardMarkup([ + [InlineKeyboardButton("Add", callback_data="add_config")], + [InlineKeyboardButton("Free List", callback_data="list_configs_free"), + InlineKeyboardButton("VIP List", callback_data="list_configs_vip")], + [InlineKeyboardButton("Delete", callback_data="del_config")], + [InlineKeyboardButton("Back", callback_data="admin")], + ])) + + +@guard(perm="can_manage_configs") +async def cb_add_config_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit( + q, + "*Add Config*\n\n" + "Format:\n`config_text | category | remark`\n\n" + "Example:\n`vmess://base64 | vip | US server`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_WAITING_CONFIG_TEXT + + +async def handle_add_config(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_manage_configs"): + return ConversationHandler.END + raw = sanitize_text(update.message.text, 5000) + parts = [p.strip() for p in raw.split("|")] + config_text = parts[0] + category = parts[1].lower() if len(parts) > 1 and parts[1].lower() in ("free", "vip") else "free" + remark = parts[2] if len(parts) > 2 else "" + + cfg_id = await db.add_config(config_text, category, remark) + await db.add_log(uid, "add_config", f"id={cfg_id} cat={category}") + await update.message.reply_text( + f"Config #{cfg_id} added ({category}).", + reply_markup=kb_admin(uid), + ) + return ConversationHandler.END + + +@guard(perm="can_manage_configs") +async def cb_del_config_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Enter config ID to delete:", reply_markup=kb_cancel()) + return STATE_DELETING_CONFIG + + +async def handle_del_config(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_manage_configs"): + return ConversationHandler.END + raw = sanitize_text(update.message.text, 20) + if not raw.isdigit(): + await update.message.reply_text("Enter a valid numeric ID.", reply_markup=kb_cancel()) + return STATE_DELETING_CONFIG + cfg_id = int(raw) + ok = await db.delete_config(cfg_id) + if ok: + await db.add_log(uid, "del_config", f"id={cfg_id}") + await update.message.reply_text(f"Config #{cfg_id} deleted.", reply_markup=kb_admin(uid)) + else: + await update.message.reply_text("Config not found.", reply_markup=kb_admin(uid)) + return ConversationHandler.END + + +@guard(perm="can_manage_configs") +async def cb_list_configs(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + cat = "free" if "free" in q.data else "vip" + configs = await db.get_all_configs(limit=10, category=cat) + if not configs: + await safe_edit(q, f"No {cat} configs.", reply_markup=kb_admin(q.from_user.id)) + return + lines = [f"*{cat.upper()} Configs:*\n"] + for c in configs: + lines.append(f"#{c['id']} | {c['protocol']} | used: {c['usage_count']} | {c['remark'][:30]}") + await safe_edit(q, "\n".join(lines)[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(q.from_user.id)) + + +# ── User management ────────────────────────────── + +@guard() +async def cb_admin_users(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Search User", callback_data="search_user")], + [InlineKeyboardButton("Ban User", callback_data="ban_user"), + InlineKeyboardButton("Unban User", callback_data="unban_user")], + [InlineKeyboardButton("Banned List", callback_data="banned_list")], + [InlineKeyboardButton("Back", callback_data="admin")], + ]) + await safe_edit(q, "*User Management*", parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard(perm="can_block_users") +async def cb_ban_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Enter user ID to ban:", reply_markup=kb_cancel()) + return STATE_BAN_USER + + +async def handle_ban_user(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_block_users"): + return ConversationHandler.END + raw = sanitize_text(update.message.text, 60) + parts = raw.split(maxsplit=1) + if not parts[0].isdigit(): + await update.message.reply_text("Enter a valid user ID.", reply_markup=kb_cancel()) + return STATE_BAN_USER + target = int(parts[0]) + reason = parts[1] if len(parts) > 1 else "" + await db.ban_user(target, reason) + await db.add_log(uid, "ban_user", f"target={target} reason={reason}") + await update.message.reply_text(f"User `{target}` banned.", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + return ConversationHandler.END + + +@guard(perm="can_block_users") +async def cb_unban_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Enter user ID to unban:", reply_markup=kb_cancel()) + return STATE_UNBAN_USER + + +async def handle_unban_user(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_block_users"): + return ConversationHandler.END + raw = sanitize_text(update.message.text, 20).strip() + if not raw.isdigit(): + await update.message.reply_text("Enter a valid user ID.", reply_markup=kb_cancel()) + return STATE_UNBAN_USER + target = int(raw) + ok = await db.unban_user(target) + if ok: + await db.add_log(uid, "unban_user", f"target={target}") + await update.message.reply_text(f"User `{target}` unbanned.", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + else: + await update.message.reply_text("User not found.", reply_markup=kb_admin(uid)) + return ConversationHandler.END + + +@guard() +async def cb_search_user_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Enter user ID or username to search:", reply_markup=kb_cancel()) + return STATE_SEARCHING_USER + + +async def handle_search_user(update: Update, context: ContextTypes.DEFAULT_TYPE): + db = context.bot_data["db"] + uid = update.effective_user.id + raw = sanitize_text(update.message.text, 60) + users = await db.search_user(raw) + if not users: + await update.message.reply_text("No users found.", reply_markup=kb_admin(uid)) + return ConversationHandler.END + lines = ["*Search Results:*\n"] + for u in users: + lines.append( + f"ID: `{u['telegram_id']}` | @{u['username'] or '--'}\n" + f"Name: {u['full_name'] or '--'} | Role: {u['role']}\n" + f"VIP: {'Yes' if u['is_vip'] else 'No'} | Balance: {u['balance']:,}\n" + ) + await update.message.reply_text("\n".join(lines)[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + return ConversationHandler.END + + +@guard() +async def cb_banned_list(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + banned = await db.get_banned_users(limit=20) + if not banned: + await safe_edit(q, "No banned users.", reply_markup=kb_admin(q.from_user.id)) + return + lines = ["*Banned Users:*\n"] + for u in banned: + lines.append(f"`{u['telegram_id']}` @{u['username'] or '--'} | {u.get('ban_reason', '')[:30]}") + await safe_edit(q, "\n".join(lines)[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(q.from_user.id)) + + +# ── VIP grant ──────────────────────────────────── + +@guard() +async def cb_admin_give_vip(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit( + q, + "*Give VIP*\n\nFormat: `user_id days [level]`\n\nExample: `123456 30 gold`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_GIVE_VIP + + +async def handle_give_vip(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_give_vip"): + return ConversationHandler.END + raw = sanitize_text(update.message.text, 60) + parts = raw.split() + if len(parts) < 2 or not parts[0].isdigit() or not parts[1].isdigit(): + await update.message.reply_text("Format: `user_id days [level]`", parse_mode=ParseMode.MARKDOWN) + return STATE_GIVE_VIP + target, days = int(parts[0]), int(parts[1]) + level = parts[2].lower() if len(parts) > 2 and parts[2].lower() in VIP_LEVELS else "silver" + if days <= 0 or days > 3650: + await update.message.reply_text("Days must be between 1 and 3650.") + return STATE_GIVE_VIP + if not await db.get_user(target): + await update.message.reply_text(f"User {target} not found.") + return ConversationHandler.END + await db.give_vip(target, days, level) + await db.add_log(uid, "give_vip", f"target={target} days={days} level={level}") + level_info = VIP_LEVELS.get(level, {}) + try: + await context.bot.send_message( + target, + f"*Congratulations!*\n\nYou received {days} days of VIP {level_info.get('name', level)}!", + parse_mode=ParseMode.MARKDOWN, + ) + except Exception: + pass + await update.message.reply_text( + f"Granted {days} days VIP {level_info.get('name', level)} to `{target}`.", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_admin(uid), + ) + return ConversationHandler.END + + +# ── Broadcast ──────────────────────────────────── + +@guard(perm="can_broadcast") +async def cb_admin_broadcast(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "*Broadcast*\n\nSend your message (text, photo, video):", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_cancel()) + return STATE_BROADCASTING + + +async def handle_broadcast(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_broadcast"): + return ConversationHandler.END + users = await db.get_all_users() + total = len(users) + if not total: + await update.message.reply_text("No users in database.") + return ConversationHandler.END + msg = update.message + status = await msg.reply_text(f"Broadcasting to {total:,} users...") + success = failed = 0 + for i, user_id in enumerate(users): + try: + if msg.text: + await context.bot.send_message(user_id, msg.text, parse_mode=ParseMode.HTML) + elif msg.photo: + await context.bot.send_photo(user_id, msg.photo[-1].file_id, caption=msg.caption or "", parse_mode=ParseMode.HTML) + elif msg.video: + await context.bot.send_video(user_id, msg.video.file_id, caption=msg.caption or "", parse_mode=ParseMode.HTML) + elif msg.document: + await context.bot.send_document(user_id, msg.document.file_id, caption=msg.caption or "", parse_mode=ParseMode.HTML) + success += 1 + except Exception: + failed += 1 + if (i + 1) % 50 == 0: + try: + pct = int((i + 1) / total * 100) + bar = make_progress_bar(pct) + await safe_edit(status, f"{bar} {pct}%\nSent: {success} | Failed: {failed}") + except Exception: + pass + await asyncio.sleep(CONFIG.get("max_broadcast_delay", 0.05)) + await safe_edit( + status, + f"*Broadcast complete!*\n\n" + f"Total: {total:,}\nSent: {success:,}\nFailed: {failed:,}\n" + f"Success rate: {success / total * 100:.1f}%", + parse_mode=ParseMode.MARKDOWN, + ) + await db.add_broadcast(uid, msg.text or msg.caption or "", success, failed) + return ConversationHandler.END + + +# ── Payments ───────────────────────────────────── + +@guard(perm="can_manage_payments") +async def cb_admin_payments(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + uid = q.from_user.id + db = context.bot_data["db"] + pays = await db.get_pending_payments() + if not pays: + await safe_edit(q, "No pending payments.", reply_markup=kb_admin(uid)) + return + lines = [f"*Pending Payments ({len(pays)}):*\n--------------------"] + for p in pays[:10]: + uname = f"@{p['username']}" if p["username"] else p["full_name"] or str(p["user_id"]) + lines.append( + f"#{p['id']} | {uname}\n" + f"{VIP_LEVELS.get(p['plan'], {}).get('name', p['plan'])} | " + f"{p['amount']:,}" + + (f" | code: {p.get('discount_code', '')}" if p.get("discount_code") else "") + + f"\n/confirm\\_{p['id']} /reject\\_{p['id']}" + ) + await safe_edit(q, "\n\n".join(lines)[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + + +@guard(perm="can_manage_payments") +async def cb_pay_action(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + uid = q.from_user.id + db = context.bot_data["db"] + parts = q.data.split("_") + action = parts[1] + pay_id = int(parts[2]) + + if action == "confirm": + pay = await db.confirm_payment(pay_id, uid) + if not pay: + await q.answer("Not found or already reviewed.", show_alert=True) + return + try: + await context.bot.send_message(pay["user_id"], "Payment confirmed! VIP activated!") + except Exception: + pass + await q.answer(f"Payment #{pay_id} confirmed.") + else: + pay = await db.reject_payment(pay_id, uid) + if not pay: + await q.answer("Not found.", show_alert=True) + return + try: + await context.bot.send_message(pay["user_id"], "Payment rejected.") + except Exception: + pass + await q.answer(f"Payment #{pay_id} rejected.") + + await cb_admin_payments(update, context) + + +# ── Discount codes ─────────────────────────────── + +@guard() +async def cb_admin_discount(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + codes = await db.list_discount_codes() + lines = ["*Discount Codes:*\n"] + if codes: + for c in codes: + lines.append(f"`{c['code']}` | {c['percent']}% | used: {c['used_count']}/{c['max_uses'] or 'unlimited'}") + else: + lines.append("No discount codes.") + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Add Code", callback_data="add_discount_code")], + [InlineKeyboardButton("Back", callback_data="admin")], + ]) + await safe_edit(q, "\n".join(lines), parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard() +async def cb_add_discount_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit( + q, + "*Add Discount Code*\n\nFormat: `CODE PERCENT MAX_USES DAYS`\n\nExample: `SAVE20 20 100 30`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_ADD_DISCOUNT + + +async def handle_add_discount(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + raw = sanitize_text(update.message.text, 100) + parts = raw.split() + if len(parts) < 2: + await update.message.reply_text("Format: `CODE PERCENT [MAX_USES] [DAYS]`", parse_mode=ParseMode.MARKDOWN) + return STATE_ADD_DISCOUNT + code = parts[0].upper() + try: + percent = int(parts[1]) + max_uses = int(parts[2]) if len(parts) > 2 else 0 + days = int(parts[3]) if len(parts) > 3 else 30 + except ValueError: + await update.message.reply_text("Invalid numbers.", reply_markup=kb_cancel()) + return STATE_ADD_DISCOUNT + if percent < 1 or percent > 100: + await update.message.reply_text("Percent must be 1-100.") + return STATE_ADD_DISCOUNT + ok = await db.add_discount_code(code, percent, max_uses, days) + if ok: + await db.add_log(uid, "add_discount", f"code={code} percent={percent}") + await update.message.reply_text(f"Discount `{code}` ({percent}%) created.", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + else: + await update.message.reply_text("Code already exists.", reply_markup=kb_admin(uid)) + return ConversationHandler.END + + +# ── Super admin ────────────────────────────────── + +@guard() +async def cb_super_admin(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + uid = q.from_user.id + if uid != SUPER_ADMIN_ID: + await q.answer("Super admin only!", show_alert=True) + return + await q.answer() + db = context.bot_data["db"] + conn = await db.connect() + async with conn.execute( + "SELECT u.telegram_id, u.username, u.full_name FROM users u " + "JOIN admin_permissions ap ON u.telegram_id=ap.admin_id WHERE u.role='admin'" + ) as cur: + admins = [dict(r) for r in await cur.fetchall()] + lines = ["*Admin Management*\n--------------------"] + if admins: + for a in admins: + lines.append(f"`{a['telegram_id']}` @{a['username'] or '--'} | {a['full_name'] or '--'}") + else: + lines.append("No admins.") + await safe_edit(q, "\n".join(lines), parse_mode=ParseMode.MARKDOWN, reply_markup=InlineKeyboardMarkup([ + [InlineKeyboardButton("Add Admin", callback_data="add_admin")], + [InlineKeyboardButton("Remove Admin", callback_data="remove_admin")], + [InlineKeyboardButton("Back", callback_data="admin")], + ])) + + +@guard() +async def cb_add_admin_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + if q.from_user.id != SUPER_ADMIN_ID: + await q.answer("Not allowed.", show_alert=True) + return + context.user_data["admin_action"] = "add_admin" + await safe_edit(q, "*Add Admin*\n\nEnter numeric user ID:", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_cancel()) + return STATE_MANAGE_ADMIN + + +@guard() +async def cb_remove_admin_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + if q.from_user.id != SUPER_ADMIN_ID: + await q.answer("Not allowed.", show_alert=True) + return + context.user_data["admin_action"] = "remove_admin" + await safe_edit(q, "*Remove Admin*\n\nEnter numeric user ID:", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_cancel()) + return STATE_MANAGE_ADMIN + + +async def handle_admin_action(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + if uid != SUPER_ADMIN_ID: + return ConversationHandler.END + db = context.bot_data["db"] + raw = sanitize_text(update.message.text, 50) + if not is_valid_tg_id(raw): + await update.message.reply_text("Enter a valid ID.", reply_markup=kb_cancel()) + return STATE_MANAGE_ADMIN + target = int(raw.strip()) + action = context.user_data.get("admin_action", "add_admin") + if action == "add_admin": + if not await db.get_user(target): + await update.message.reply_text(f"User {target} not registered.") + return STATE_MANAGE_ADMIN + conn = await db.connect() + await conn.execute("UPDATE users SET role='admin' WHERE telegram_id=?", (target,)) + await conn.execute( + "INSERT OR IGNORE INTO admin_permissions " + "(admin_id, can_manage_configs, can_manage_payments, can_broadcast, can_block_users, can_give_vip) " + "VALUES (?, 0, 0, 0, 0, 0)", + (target,), + ) + await conn.commit() + try: + await context.bot.send_message(target, "You have been promoted to admin!") + except Exception: + pass + await update.message.reply_text(f"`{target}` is now admin.", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + else: + ok = await db.remove_admin(target) + if not ok: + await update.message.reply_text(f"Admin {target} not found.") + else: + try: + await context.bot.send_message(target, "Your admin access has been revoked.") + except Exception: + pass + await update.message.reply_text(f"Admin `{target}` removed.", parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(uid)) + context.user_data.pop("admin_action", None) + return ConversationHandler.END + + +# ── Full stats ─────────────────────────────────── + +@guard() +async def cb_full_stats(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + if q.from_user.id != SUPER_ADMIN_ID: + await q.answer("Super admin only!", show_alert=True) + return + await q.answer() + db = context.bot_data["db"] + stats = await db.get_stats() + weekly = await db.get_weekly_stats() + conn = await db.connect() + + async with conn.execute("SELECT id, protocol, remark, usage_count FROM configs ORDER BY usage_count DESC LIMIT 3") as c: + top_configs = [dict(r) for r in await c.fetchall()] + async with conn.execute("SELECT telegram_id, username, total_referrals FROM users ORDER BY total_referrals DESC LIMIT 3") as c: + top_referrers = [dict(r) for r in await c.fetchall()] + + proxy_stats = await db.get_proxy_stats() + + # Weekly user chart (simple text bar) + chart_lines = [] + for day in reversed(weekly["daily_breakdown"]): + bar = make_progress_bar(min(day["new_users"] * 10, 100), 8) if day["new_users"] else "░" * 8 + chart_lines.append(f"`{day['date'][5:]}` {bar} {day['new_users']}") + + text = ( + f"*Full Statistics*\n" + f"--------------------\n" + f"Users: {stats['total_users']:,} | Today: {stats['today_users']:,}\n" + f"VIP: {stats['total_vips']:,} | Banned: {stats['total_banned']:,}\n\n" + f"*Weekly Overview:*\n" + f"Revenue: {weekly['week_revenue']:,}\n" + f"Downloads: {weekly['week_downloads']:,}\n" + f"Config claims: {weekly['week_claims']:,}\n\n" + f"*New users (7 days):*\n" + ) + text += "\n".join(chart_lines) + text += ( + f"\n\n*Lifetime:*\n" + f"Downloads: {stats['total_downloads']:,}\n" + f"Revenue: {stats['total_revenue']:,}\n" + f"Payments: {stats['total_payments']:,} | Pending: {stats['pending_payments']:,}\n" + f"Wallet pool: {stats['total_wallet_balance']:,}\n" + f"Proxies: {proxy_stats['alive']} alive / {proxy_stats['total']} total\n\n" + f"*Top configs:*\n" + ) + for i, c in enumerate(top_configs, 1): + text += f"{i}. #{c['id']} {c['protocol'].upper()} -- {c['usage_count']} uses\n" + text += "\n*Top referrers:*\n" + for i, r in enumerate(top_referrers, 1): + text += f"{i}. @{r['username'] or r['telegram_id']} -- {r['total_referrals']} referrals\n" + + if weekly["top_actions"]: + text += "\n*Top actions (week):*\n" + for a in weekly["top_actions"]: + text += f" {a['action']}: {a['cnt']}\n" + + await safe_edit(q, text[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb_admin(q.from_user.id)) + + +# ── Backup ─────────────────────────────────────── + +@guard() +async def cb_admin_backup(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + if q.from_user.id != SUPER_ADMIN_ID: + await q.answer("Super admin only!", show_alert=True) + return + await q.answer("Creating backup...") + db = context.bot_data["db"] + dst = db.backup() + await db.add_log(q.from_user.id, "db_backup", dst) + with open(dst, "rb") as f: + await context.bot.send_document( + q.from_user.id, document=f, filename=os.path.basename(dst), + caption=f"DB Backup | {datetime.now():%Y/%m/%d %H:%M}", + ) + await safe_edit(q, "Backup sent!", reply_markup=kb_admin(q.from_user.id)) + + +# ── Cookie ─────────────────────────────────────── + +@guard() +async def cb_set_cookie(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + if q.from_user.id != SUPER_ADMIN_ID: + await q.answer("Super admin only!", show_alert=True) + return + await safe_edit( + q, + "*Set YouTube Cookie*\n\nSend `cookies.txt` file (Netscape format).", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_SET_COOKIE + + +async def handle_cookie(update: Update, context: ContextTypes.DEFAULT_TYPE): + if update.effective_user.id != SUPER_ADMIN_ID: + return ConversationHandler.END + db = context.bot_data["db"] + if not update.message.document: + await update.message.reply_text("Send a .txt file.") + return STATE_SET_COOKIE + if not update.message.document.file_name.lower().endswith(".txt"): + await update.message.reply_text("Only .txt files accepted.") + return STATE_SET_COOKIE + tg_file = await context.bot.get_file(update.message.document.file_id) + await tg_file.download_to_drive(COOKIE_FILE) + await db.add_log(update.effective_user.id, "set_cookie", COOKIE_FILE) + await update.message.reply_text("Cookie saved!", reply_markup=kb_admin(update.effective_user.id)) + return ConversationHandler.END + + +# ── Payment commands ───────────────────────────── + +async def cmd_confirm_payment(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_manage_payments"): + await update.message.reply_text("Permission denied.") + return + try: + pay_id = int(update.message.text.split("_")[1]) + pay = await db.confirm_payment(pay_id, uid) + if not pay: + await update.message.reply_text("Not found or already reviewed.") + return + try: + await context.bot.send_message(pay["user_id"], "Payment confirmed! VIP activated!") + except Exception: + pass + await update.message.reply_text(f"Payment #{pay_id} confirmed.") + except (IndexError, ValueError): + await update.message.reply_text("Invalid format. Use /confirm_ID") + + +async def cmd_reject_payment(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_manage_payments"): + await update.message.reply_text("Permission denied.") + return + try: + pay_id = int(update.message.text.split("_")[1]) + pay = await db.reject_payment(pay_id, uid) + if not pay: + await update.message.reply_text("Not found.") + return + try: + await context.bot.send_message(pay["user_id"], "Payment rejected.") + except Exception: + pass + await update.message.reply_text(f"Payment #{pay_id} rejected.") + except (IndexError, ValueError): + await update.message.reply_text("Invalid format. Use /reject_ID") + + +# ── Admin wallet charge ────────────────────────── + +@guard(perm="can_manage_payments") +async def cb_admin_wallet_charge(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + await safe_edit( + q, + "*Charge User Wallet*\n\n" + "Send in format:\n" + "`user_id amount reason`\n\n" + "Example: `123456789 50000 Manual top-up`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_ADMIN_WALLET_CHARGE + + +async def handle_wallet_charge(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not await db.has_perm(uid, "can_manage_payments"): + return ConversationHandler.END + + raw = sanitize_text(update.message.text, 200) + parts = raw.split(None, 2) + if len(parts) < 2: + await update.message.reply_text( + "Invalid format. Use: `user_id amount reason`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_ADMIN_WALLET_CHARGE + + target_id_str, amount_str = parts[0], parts[1] + reason = parts[2] if len(parts) > 2 else "Admin charge" + + if not is_valid_tg_id(target_id_str): + await update.message.reply_text("Invalid user ID.", reply_markup=kb_cancel()) + return STATE_ADMIN_WALLET_CHARGE + + try: + amount = int(amount_str) + if amount <= 0 or amount > 10_000_000: + raise ValueError + except ValueError: + await update.message.reply_text("Invalid amount (1 to 10,000,000).", reply_markup=kb_cancel()) + return STATE_ADMIN_WALLET_CHARGE + + target_id = int(target_id_str) + target = await db.get_user(target_id) + if not target: + await update.message.reply_text("User not found.", reply_markup=kb_cancel()) + return STATE_ADMIN_WALLET_CHARGE + + new_balance = await db.add_balance(target_id, amount, f"Admin: {reason}") + await db.add_log(uid, "admin_wallet_charge", f"target={target_id} amount={amount} reason={reason}") + + await update.message.reply_text( + f"Charged *{amount:,}* to user `{target_id}`\n" + f"New balance: *{new_balance:,}*\n" + f"Reason: {reason}", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_admin(uid), + ) + try: + await context.bot.send_message( + target_id, + f"*{amount:,}* added to your wallet!\n" + f"Reason: {reason}\n" + f"New balance: *{new_balance:,}*", + parse_mode=ParseMode.MARKDOWN, + ) + except Exception: + pass + return ConversationHandler.END diff --git a/bot/handlers/file.py b/bot/handlers/file.py new file mode 100644 index 0000000..39e0921 --- /dev/null +++ b/bot/handlers/file.py @@ -0,0 +1,87 @@ +""" +File-to-link handler. +""" + +from telegram import Update +from telegram.constants import ParseMode +from telegram.ext import ContextTypes, ConversationHandler + +from bot.config import MAX_FILE_SIZE, BASE_DOWNLOAD_URL, STATE_WAITING_FILE +from bot.decorators import guard +from bot.helpers import safe_edit, fmt_size +from bot.keyboards import kb_back_main, kb_cancel + + +@guard() +async def cb_file(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + await safe_edit( + q, + f"*File to Link*\n\nSend a file (max {MAX_FILE_SIZE} MB):", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_WAITING_FILE + + +async def handle_file(update: Update, context: ContextTypes.DEFAULT_TYPE): + msg = update.message + uid = msg.from_user.id + db = context.bot_data["db"] + + user = await db.get_user(uid) + if user and user.get("banned"): + return ConversationHandler.END + + fid = uid_unique = ftype = name = mime = None + size = 0 + + if msg.document: + d = msg.document + fid, uid_unique, ftype = d.file_id, d.file_unique_id, "document" + name, size, mime = d.file_name or "file", d.file_size or 0, d.mime_type or "application/octet-stream" + elif msg.video: + v = msg.video + fid, uid_unique, ftype = v.file_id, v.file_unique_id, "video" + name, size, mime = v.file_name or "video.mp4", v.file_size or 0, v.mime_type or "video/mp4" + elif msg.audio: + a = msg.audio + fid, uid_unique, ftype = a.file_id, a.file_unique_id, "audio" + name, size, mime = a.file_name or "audio.mp3", a.file_size or 0, a.mime_type or "audio/mpeg" + elif msg.photo: + p = msg.photo[-1] + fid, uid_unique, ftype = p.file_id, p.file_unique_id, "photo" + name, size, mime = "photo.jpg", p.file_size or 0, "image/jpeg" + elif msg.voice: + v2 = msg.voice + fid, uid_unique, ftype = v2.file_id, v2.file_unique_id, "voice" + name, size, mime = "voice.ogg", v2.file_size or 0, "audio/ogg" + else: + await msg.reply_text("Unsupported file type.", reply_markup=kb_cancel()) + return STATE_WAITING_FILE + + if size and size > MAX_FILE_SIZE * 1024 * 1024: + await msg.reply_text( + f"File too large ({fmt_size(size)} > {MAX_FILE_SIZE} MB).", + reply_markup=kb_cancel(), + ) + return STATE_WAITING_FILE + + await db.save_file(uid_unique, fid, ftype, name, size, mime, uploader_id=uid) + await db.add_log(uid, "file_upload", f"{name} ({fmt_size(size)})") + + bot_info = await context.bot.get_me() + bot_link = f"https://t.me/{bot_info.username}?start=dl_{uid_unique}" + stream_link = f"{BASE_DOWNLOAD_URL}/stream/{uid_unique}" + + text = ( + f"*File saved!*\n\n" + f"Name: `{name}`\n" + f"Size: {fmt_size(size)}\n" + f"Expiry: 7 days\n\n" + f"Telegram link:\n`{bot_link}`\n\n" + f"Direct link:\n`{stream_link}`" + ) + await msg.reply_text(text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_back_main()) + return ConversationHandler.END diff --git a/bot/handlers/main_menu.py b/bot/handlers/main_menu.py new file mode 100644 index 0000000..9bc4882 --- /dev/null +++ b/bot/handlers/main_menu.py @@ -0,0 +1,349 @@ +""" +Main menu callback handlers: profile, help, referral, claim, VIP, wallet. +""" + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.constants import ParseMode +from telegram.ext import ContextTypes + +from bot.config import ( + SUPER_ADMIN_ID, + PAYMENT_CARD, + REFERRAL_BONUS, + VIP_LEVELS, + WELCOME_MSG, + STATE_WAITING_RECEIPT, + STATE_WAITING_DISCOUNT, + STATE_WALLET_CHARGE, +) +from bot.decorators import guard +from bot.helpers import safe_edit, fmt_time, vip_level_info, generate_qr +from bot.keyboards import kb_main, kb_back_main, kb_cancel + + +@guard() +async def cb_main(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + uid = q.from_user.id + db = context.bot_data["db"] + role = await db.get_role(uid) + await safe_edit( + q, + "*Main Menu*\nSelect an option:", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_main(uid, role), + ) + + +@guard() +async def cb_check_join(update: Update, context: ContextTypes.DEFAULT_TYPE): + from bot.config import FORCE_CHANNEL + q = update.callback_query + uid = q.from_user.id + db = context.bot_data["db"] + try: + member = await context.bot.get_chat_member(FORCE_CHANNEL, uid) + if member.status in ("left", "kicked"): + await q.answer("You haven't joined yet!", show_alert=True) + return + except Exception: + await q.answer("Error checking membership.", show_alert=True) + return + await q.answer("Membership confirmed!") + role = await db.get_role(uid) + await safe_edit(q, "Membership confirmed! You can now use the bot.", reply_markup=kb_main(uid, role)) + + +@guard() +async def cb_help(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + text = ( + "*Bot Guide*\n" + "--------------------\n\n" + "*Config:* Get a free VPN config every 6 hours\n\n" + "*Download:* Send a YouTube/Instagram/TikTok link to download\n\n" + "*File2Link:* Send any file, get a download link (7 day expiry)\n\n" + "*VIP:* Premium configs, priority downloads, and more\n\n" + "*Wallet:* Balance for VIP purchases\n\n" + "*Referral:* Earn rewards for inviting friends\n\n" + "*Proxy:* Get free HTTP/SOCKS proxies\n\n" + "--------------------\n" + "`/cancel` - Cancel any operation" + ) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_back_main()) + + +@guard() +async def cb_profile(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + uid = q.from_user.id + db = context.bot_data["db"] + user = await db.get_user(uid) + if not user: + await safe_edit(q, "User not found.", reply_markup=kb_back_main()) + return + role_map = {"super_admin": "Super Admin", "admin": "Admin", "user": "User"} + role_text = role_map.get(user["role"], "User") + vip_text = vip_level_info(user) if user["is_vip"] else "None" + text = ( + f"*Your Profile*\n" + f"--------------------\n" + f"ID: `{uid}`\n" + f"Name: {user['full_name'] or '--'}\n" + f"Role: {role_text}\n" + f"VIP: {vip_text}\n" + f"Balance: *{user['balance']:,}*\n" + f"Configs claimed: {user['total_claims']}\n" + f"Referrals: {user['total_referrals']}\n" + f"Total spent: {user.get('total_spent', 0):,}\n" + f"Joined: {fmt_time(user['joined_at'])}\n" + f"Last seen: {fmt_time(user['last_seen'])}" + ) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_back_main()) + + +@guard() +async def cb_referral(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + uid = q.from_user.id + db = context.bot_data["db"] + user = await db.get_user(uid) + bot_info = await context.bot.get_me() + link = f"https://t.me/{bot_info.username}?start={user['referral_code']}" + text = ( + f"*Referral System*\n" + f"--------------------\n\n" + f"For each friend who joins via your link:\n" + f"*{REFERRAL_BONUS:,}* is added to your wallet\n\n" + f"Your link:\n`{link}`\n\n" + f"Successful referrals: *{user['total_referrals']}*\n" + f"Wallet balance: *{user['balance']:,}*" + ) + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Share Link", switch_inline_query=f"Check out this bot: {link}")], + [InlineKeyboardButton("Home", callback_data="main")], + ]) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard() +async def cb_claim(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + uid = q.from_user.id + db = context.bot_data["db"] + can, remaining = await db.can_claim(uid) + if not can: + h, m = divmod(remaining, 3600) + await q.answer(f"Wait {h}h {m // 60}m", show_alert=True) + return + + is_vip = await db.is_vip(uid) + category = "vip" if is_vip else "free" + configs = await db.get_active_configs(category, limit=3) + + if not configs: + if is_vip: + configs = await db.get_active_configs("free", limit=3) + if not configs: + await safe_edit(q, "No configs available right now.", reply_markup=kb_back_main()) + return + + await db.update_claim(uid) + lines = ["*Your Configs:*\n--------------------\n"] + for c in configs: + await db.increment_usage(c["id"]) + await db.add_log(uid, "claim_config", f"config_id={c['id']}") + lines.append(f"`{c['config_text']}`\n") + + await safe_edit(q, "\n".join(lines), parse_mode=ParseMode.MARKDOWN, reply_markup=kb_back_main()) + + +@guard() +async def cb_vip(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + lines = ["*VIP Plans*\n--------------------\n"] + btns = [] + for key, info in VIP_LEVELS.items(): + lines.append( + f"*{info['name']}*\n" + f"Price: {info['price']:,} | Duration: {info['days']} days\n" + f"Max quality: {info['max_yt_quality']}p | Configs: {info['max_configs']}\n" + ) + btns.append(InlineKeyboardButton(f"Buy {info['name']}", callback_data=f"buy_{key}")) + + kb_rows = [btns[i : i + 2] for i in range(0, len(btns), 2)] + kb_rows.append([InlineKeyboardButton("Home", callback_data="main")]) + await safe_edit( + q, "\n".join(lines), parse_mode=ParseMode.MARKDOWN, + reply_markup=InlineKeyboardMarkup(kb_rows), + ) + + +@guard() +async def cb_buy_vip(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + plan = q.data.replace("buy_", "") + info = VIP_LEVELS.get(plan) + if not info: + await safe_edit(q, "Invalid plan.", reply_markup=kb_back_main()) + return + context.user_data["vip_plan"] = plan + context.user_data["vip_price"] = info["price"] + text = ( + f"*{info['name']}*\n" + f"--------------------\n" + f"Price: *{info['price']:,}*\n" + f"Duration: {info['days']} days\n\n" + f"Card: `{PAYMENT_CARD}`\n\n" + f"After payment, send your receipt photo." + ) + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Send Receipt", callback_data="send_receipt")], + [InlineKeyboardButton("Apply Discount", callback_data="apply_discount")], + [InlineKeyboardButton("Home", callback_data="main")], + ]) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard() +async def cb_receipt_prompt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Send your payment receipt (photo):", reply_markup=kb_cancel()) + return STATE_WAITING_RECEIPT + + +async def handle_receipt(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + if not update.message.photo: + await update.message.reply_text("Please send a photo.", reply_markup=kb_cancel()) + return STATE_WAITING_RECEIPT + plan = context.user_data.get("vip_plan", "silver") + price = context.user_data.get("vip_price", VIP_LEVELS.get(plan, {}).get("price", 0)) + discount = context.user_data.get("discount_code") + receipt_id = update.message.photo[-1].file_id + pay_id = await db.create_payment(uid, price, plan, receipt_id, discount) + if discount: + await db.use_discount(discount) + await db.add_log(uid, "payment_submit", f"plan={plan} amount={price} id={pay_id}") + await update.message.reply_text( + f"Receipt submitted (#{pay_id}). Awaiting admin review.", + reply_markup=kb_back_main(), + ) + context.user_data.pop("vip_plan", None) + context.user_data.pop("vip_price", None) + context.user_data.pop("discount_code", None) + from telegram.ext import ConversationHandler + return ConversationHandler.END + + +@guard() +async def cb_apply_discount(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit(q, "Enter your discount code:", reply_markup=kb_cancel()) + return STATE_WAITING_DISCOUNT + + +async def handle_discount_code(update: Update, context: ContextTypes.DEFAULT_TYPE): + from bot.helpers import sanitize_text + from telegram.ext import ConversationHandler + + uid = update.effective_user.id + db = context.bot_data["db"] + code = sanitize_text(update.message.text, 30).strip().upper() + disc = await db.validate_discount(code) + if not disc: + await update.message.reply_text("Invalid or expired code.", reply_markup=kb_cancel()) + return STATE_WAITING_DISCOUNT + plan = context.user_data.get("vip_plan", "silver") + original = VIP_LEVELS.get(plan, {}).get("price", 0) + new_price = int(original * (100 - disc["percent"]) / 100) + context.user_data["vip_price"] = new_price + context.user_data["discount_code"] = code + await update.message.reply_text( + f"Discount applied! {disc['percent']}% off\n" + f"Original: {original:,} -> New: *{new_price:,}*\n\n" + f"Now send your receipt photo.", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_WAITING_RECEIPT + + +@guard() +async def cb_wallet(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + uid = q.from_user.id + db = context.bot_data["db"] + user = await db.get_user(uid) + history = await db.get_wallet_history(uid, limit=5) + lines = [ + f"*Wallet*\n--------------------\n", + f"Balance: *{user['balance']:,}*\n", + ] + if history: + lines.append("\nRecent transactions:") + for tx in history: + icon = "+" if tx["type"] == "credit" else "-" + lines.append(f" {icon}{abs(tx['amount']):,} | {tx['description'][:30]}") + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Buy VIP with Wallet", callback_data="wallet_buy_vip")], + [InlineKeyboardButton("Home", callback_data="main")], + ]) + await safe_edit(q, "\n".join(lines), parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard() +async def cb_wallet_buy_vip(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + btns = [] + for key, info in VIP_LEVELS.items(): + btns.append( + InlineKeyboardButton( + f"{info['name']} ({info['price']:,})", + callback_data=f"wallet_pay_{key}", + ) + ) + kb_rows = [btns[i : i + 2] for i in range(0, len(btns), 2)] + kb_rows.append([InlineKeyboardButton("Back", callback_data="wallet")]) + await safe_edit( + q, "*Buy VIP with wallet balance:*", + parse_mode=ParseMode.MARKDOWN, + reply_markup=InlineKeyboardMarkup(kb_rows), + ) + + +@guard() +async def cb_wallet_pay(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + uid = q.from_user.id + db = context.bot_data["db"] + plan = q.data.replace("wallet_pay_", "") + info = VIP_LEVELS.get(plan) + if not info: + await q.answer("Invalid plan.", show_alert=True) + return + ok = await db.deduct_balance(uid, info["price"], f"VIP purchase ({info['name']})") + if not ok: + await q.answer("Insufficient balance!", show_alert=True) + return + await db.give_vip(uid, info["days"], plan) + await db.add_log(uid, "wallet_vip", f"plan={plan} price={info['price']}") + await q.answer(f"VIP {info['name']} activated!") + await safe_edit( + q, + f"*VIP Activated!*\n\n" + f"Plan: {info['name']}\n" + f"Duration: {info['days']} days", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_back_main(), + ) diff --git a/bot/handlers/proxy.py b/bot/handlers/proxy.py new file mode 100644 index 0000000..3acf194 --- /dev/null +++ b/bot/handlers/proxy.py @@ -0,0 +1,408 @@ +""" +Proxy feature: fetch free proxy lists, check alive, serve to users. +Supports HTTP, HTTPS, SOCKS4, SOCKS5. +""" + +import asyncio +import logging +import time +from typing import List + +import aiohttp + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.constants import ParseMode +from telegram.ext import ContextTypes + +from telegram.ext import ConversationHandler + +from bot.config import PROXY_CHECK_TIMEOUT, PROXY_MAX_RESULTS, STATE_CHECKING_PROXY +from bot.decorators import guard +from bot.helpers import safe_edit +from bot.keyboards import kb_back_main, kb_proxy_menu + +logger = logging.getLogger("BOT.proxy") + +# Public proxy API endpoints +PROXY_SOURCES = [ + { + "url": "https://api.proxyscrape.com/v2/?request=displayproxies&protocol=http&timeout=10000&country=all&ssl=all&anonymity=all", + "protocol": "http", + }, + { + "url": "https://api.proxyscrape.com/v2/?request=displayproxies&protocol=socks5&timeout=10000&country=all", + "protocol": "socks5", + }, + { + "url": "https://api.proxyscrape.com/v2/?request=displayproxies&protocol=socks4&timeout=10000&country=all", + "protocol": "socks4", + }, + { + "url": "https://www.proxy-list.download/api/v1/get?type=http", + "protocol": "http", + }, + { + "url": "https://www.proxy-list.download/api/v1/get?type=https", + "protocol": "https", + }, + { + "url": "https://www.proxy-list.download/api/v1/get?type=socks5", + "protocol": "socks5", + }, + { + "url": "https://www.proxy-list.download/api/v1/get?type=socks4", + "protocol": "socks4", + }, +] + + +async def _fetch_proxy_list(session: aiohttp.ClientSession, source: dict) -> List[dict]: + """Fetch proxies from a single API source.""" + proxies = [] + try: + async with session.get(source["url"], timeout=aiohttp.ClientTimeout(total=15)) as resp: + if resp.status != 200: + return [] + text = await resp.text() + for line in text.strip().splitlines(): + line = line.strip() + if ":" in line: + parts = line.split(":") + if len(parts) == 2: + ip, port_str = parts + try: + port = int(port_str) + proxies.append({ + "ip": ip.strip(), + "port": port, + "protocol": source["protocol"], + }) + except ValueError: + continue + except Exception as e: + logger.debug(f"Failed to fetch from {source['url']}: {e}") + return proxies + + +async def fetch_all_proxies() -> List[dict]: + """Fetch proxies from all sources concurrently.""" + all_proxies = [] + async with aiohttp.ClientSession() as session: + tasks = [_fetch_proxy_list(session, src) for src in PROXY_SOURCES] + results = await asyncio.gather(*tasks, return_exceptions=True) + for result in results: + if isinstance(result, list): + all_proxies.extend(result) + + # deduplicate + seen = set() + unique = [] + for p in all_proxies: + key = (p["ip"], p["port"], p["protocol"]) + if key not in seen: + seen.add(key) + unique.append(p) + return unique + + +async def check_proxy(ip: str, port: int, protocol: str, timeout: int = None) -> dict: + """Check if a proxy is alive and measure speed.""" + timeout = timeout or PROXY_CHECK_TIMEOUT + proxy_url = f"{protocol}://{ip}:{port}" + start = time.monotonic() + try: + async with aiohttp.ClientSession() as session: + async with session.get( + "http://httpbin.org/ip", + proxy=proxy_url if protocol in ("http", "https") else None, + timeout=aiohttp.ClientTimeout(total=timeout), + ) as resp: + elapsed = int((time.monotonic() - start) * 1000) + if resp.status == 200: + return {"alive": True, "speed_ms": elapsed} + except Exception: + pass + return {"alive": False, "speed_ms": 0} + + +async def update_proxy_database(db) -> int: + """Fetch proxies from all sources and upsert into DB.""" + proxies = await fetch_all_proxies() + if not proxies: + return 0 + return await db.upsert_proxies(proxies) + + +# ── Handlers ────────────────────────────────────── + +@guard() +async def cb_proxy_menu(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + stats = await db.get_proxy_stats() + text = ( + "*Proxy Service*\n" + "--------------------\n\n" + f"Available proxies: *{stats['alive']}*\n" + f"HTTP: {stats['http']} | HTTPS: {stats['https']}\n" + f"SOCKS4: {stats['socks4']} | SOCKS5: {stats['socks5']}\n\n" + "Select a proxy type:" + ) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_proxy_menu()) + + +@guard() +async def cb_proxy_list(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + + proto_map = { + "proxy_http": "http", + "proxy_socks5": "socks5", + "proxy_socks4": "socks4", + "proxy_all": None, + } + protocol = proto_map.get(q.data) + label = protocol.upper() if protocol else "ALL" + + proxies = await db.get_alive_proxies(protocol=protocol, limit=PROXY_MAX_RESULTS) + if not proxies: + await safe_edit( + q, + f"No {label} proxies available.\nTry refreshing or check back later.", + reply_markup=InlineKeyboardMarkup([ + [InlineKeyboardButton("Refresh Proxies", callback_data="proxy_refresh")], + [InlineKeyboardButton("Back", callback_data="proxy_menu")], + ]), + ) + return + + lines = [f"*{label} Proxies* ({len(proxies)} found)\n--------------------\n"] + for i, p in enumerate(proxies[:20], 1): + speed = f"{p['speed_ms']}ms" if p["speed_ms"] else "N/A" + country = p.get("country") or "??" + lines.append(f"`{p['ip']}:{p['port']}` | {p['protocol']} | {speed} | {country}") + + if len(proxies) > 20: + lines.append(f"\n... and {len(proxies) - 20} more") + + text = "\n".join(lines) + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Copy as Text", callback_data=f"proxy_copy_{protocol or 'all'}")], + [InlineKeyboardButton("Refresh", callback_data="proxy_refresh")], + [InlineKeyboardButton("Back", callback_data="proxy_menu")], + ]) + await safe_edit(q, text[:4000], parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard() +async def cb_proxy_copy(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + proto = q.data.replace("proxy_copy_", "") + protocol = proto if proto != "all" else None + + proxies = await db.get_alive_proxies(protocol=protocol, limit=PROXY_MAX_RESULTS) + if not proxies: + await q.answer("No proxies available.", show_alert=True) + return + + text_list = "\n".join(f"{p['ip']}:{p['port']}" for p in proxies) + await context.bot.send_message( + q.from_user.id, + f"```\n{text_list}\n```", + parse_mode=ParseMode.MARKDOWN, + ) + + +@guard() +async def cb_proxy_stats(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + stats = await db.get_proxy_stats() + text = ( + "*Proxy Statistics*\n" + "--------------------\n\n" + f"Total in DB: {stats['total']}\n" + f"Alive: *{stats['alive']}*\n" + f"Dead: {stats['total'] - stats['alive']}\n\n" + f"HTTP: {stats['http']}\n" + f"HTTPS: {stats['https']}\n" + f"SOCKS4: {stats['socks4']}\n" + f"SOCKS5: {stats['socks5']}" + ) + await safe_edit( + q, text, parse_mode=ParseMode.MARKDOWN, + reply_markup=InlineKeyboardMarkup([ + [InlineKeyboardButton("Refresh Proxies", callback_data="proxy_refresh")], + [InlineKeyboardButton("Back", callback_data="proxy_menu")], + ]), + ) + + +@guard() +async def cb_proxy_refresh(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer("Refreshing proxy list...") + db = context.bot_data["db"] + + await safe_edit(q, "Fetching proxies from sources...", reply_markup=None) + count = await update_proxy_database(db) + await db.clean_dead_proxies(older_than_hours=12) + + stats = await db.get_proxy_stats() + text = ( + f"*Proxy list updated!*\n\n" + f"Fetched: {count} proxies\n" + f"Alive: {stats['alive']}" + ) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_proxy_menu()) + + +@guard() +async def cb_proxy_check(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + await safe_edit( + q, + "*Proxy Checker*\n\n" + "Send a proxy in format:\n" + "`ip:port` or `protocol://ip:port`\n\n" + "Example: `1.2.3.4:8080` or `socks5://1.2.3.4:1080`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_back_main(), + ) + return STATE_CHECKING_PROXY + + +async def handle_proxy_check(update: Update, context: ContextTypes.DEFAULT_TYPE): + """User sent a proxy string to validate.""" + text = update.message.text.strip() + protocol = "http" + addr = text + + if "://" in text: + protocol, addr = text.split("://", 1) + protocol = protocol.lower() + + if ":" not in addr: + await update.message.reply_text( + "Invalid format. Use `ip:port` or `protocol://ip:port`", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_back_main(), + ) + return ConversationHandler.END + + parts = addr.split(":") + ip = parts[0].strip() + try: + port = int(parts[1].strip()) + except (ValueError, IndexError): + await update.message.reply_text( + "Invalid port number.", + reply_markup=kb_back_main(), + ) + return ConversationHandler.END + + await update.message.reply_text(f"Checking `{protocol}://{ip}:{port}` ...") + + result = await check_proxy(ip, port, protocol) + if result["alive"]: + text = ( + f"*Proxy is ALIVE*\n\n" + f"Address: `{ip}:{port}`\n" + f"Protocol: {protocol.upper()}\n" + f"Response time: *{result['speed_ms']}ms*" + ) + else: + text = ( + f"*Proxy is DEAD*\n\n" + f"Address: `{ip}:{port}`\n" + f"Protocol: {protocol.upper()}\n" + f"Could not connect within {PROXY_CHECK_TIMEOUT}s" + ) + await update.message.reply_text( + text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_proxy_menu() + ) + return ConversationHandler.END + + +# Admin proxy management +@guard(perm="can_manage_configs") +async def cb_admin_proxy(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + stats = await db.get_proxy_stats() + text = ( + "*Admin: Proxy Management*\n" + "--------------------\n\n" + f"Total: {stats['total']} | Alive: {stats['alive']}\n" + f"HTTP: {stats['http']} | SOCKS5: {stats['socks5']}" + ) + kb = InlineKeyboardMarkup([ + [InlineKeyboardButton("Fetch New Proxies", callback_data="proxy_refresh")], + [InlineKeyboardButton("Clean Dead Proxies", callback_data="admin_proxy_clean")], + [InlineKeyboardButton("Back", callback_data="admin")], + ]) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb) + + +@guard(perm="can_manage_configs") +async def cb_admin_proxy_clean(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + removed = await db.clean_dead_proxies(older_than_hours=1) + await safe_edit( + q, + f"Cleaned {removed} dead proxies.", + reply_markup=InlineKeyboardMarkup([ + [InlineKeyboardButton("Back", callback_data="admin_proxy")], + ]), + ) + + +@guard() +async def cb_speed_test(update: Update, context: ContextTypes.DEFAULT_TYPE): + """Run a quick download speed test.""" + q = update.callback_query + await q.answer() + await safe_edit(q, "Running speed test...\nPlease wait 5-10 seconds.", reply_markup=None) + + test_urls = [ + ("Cloudflare", "https://speed.cloudflare.com/__down?bytes=5000000"), + ("Hetzner", "https://speed.hetzner.de/1MB.bin"), + ] + results = [] + async with aiohttp.ClientSession() as session: + for name, url in test_urls: + try: + start = time.monotonic() + total = 0 + async with session.get(url, timeout=aiohttp.ClientTimeout(total=15)) as resp: + async for chunk in resp.content.iter_chunked(65536): + total += len(chunk) + elapsed = time.monotonic() - start + speed_mbps = (total * 8) / (elapsed * 1_000_000) + results.append(f"{name}: *{speed_mbps:.1f} Mbps* ({total / 1024:.0f} KB in {elapsed:.1f}s)") + except Exception: + results.append(f"{name}: Failed") + + text = "*Speed Test Results*\n--------------------\n\n" + "\n".join(results) + await safe_edit(q, text, parse_mode=ParseMode.MARKDOWN, reply_markup=kb_back_main()) + + +async def scheduled_proxy_update(context: ContextTypes.DEFAULT_TYPE): + """Periodic job to refresh proxy list.""" + db = context.bot_data["db"] + try: + count = await update_proxy_database(db) + await db.clean_dead_proxies(older_than_hours=24) + logger.info(f"Scheduled proxy update: {count} proxies fetched") + except Exception as e: + logger.error(f"Scheduled proxy update failed: {e}") diff --git a/bot/handlers/start.py b/bot/handlers/start.py new file mode 100644 index 0000000..c5449dd --- /dev/null +++ b/bot/handlers/start.py @@ -0,0 +1,86 @@ +""" +/start command handler. +""" + +from datetime import datetime + +from telegram import Update +from telegram.constants import ParseMode +from telegram.ext import ContextTypes + +from bot.config import WELCOME_MSG, REFERRAL_BONUS +from bot.keyboards import kb_main +from bot.helpers import fmt_size + + +async def cmd_start(update: Update, context: ContextTypes.DEFAULT_TYPE): + user = update.effective_user + if not user: + return + db = context.bot_data["db"] + + args = context.args or [] + ref_code = args[0] if args else None + + u, referred_by = await db.create_user( + user.id, user.username or "", user.full_name or "", ref_code + ) + + if referred_by: + referrer = await db.get_user(referred_by) + try: + await context.bot.send_message( + referred_by, + f"*Referral success!*\n\n" + f"*{user.full_name or user.username or 'New user'}* joined via your link!\n" + f"*{REFERRAL_BONUS:,}* added to your wallet\n" + f"Balance: *{referrer['balance']:,}*", + parse_mode=ParseMode.MARKDOWN, + ) + except Exception: + pass + + if ref_code and ref_code.startswith("dl_"): + file_uid = ref_code[3:] + file_data = await db.get_file(file_uid) + if file_data: + if file_data.get("expires_at") and datetime.now() > datetime.fromisoformat( + file_data["expires_at"] + ): + await update.message.reply_text("This link has expired.") + return + await db.increment_file_downloads(file_uid) + try: + send_map = { + "document": context.bot.send_document, + "video": context.bot.send_video, + "audio": context.bot.send_audio, + "photo": context.bot.send_photo, + "voice": context.bot.send_voice, + } + fn = send_map.get(file_data["file_type"], context.bot.send_document) + await fn( + user.id, + file_data["file_id"], + caption=( + f"*{file_data['file_name']}*\n" + f"Size: {fmt_size(file_data['file_size'])}" + ), + parse_mode=ParseMode.MARKDOWN, + ) + except Exception: + await update.message.reply_text("Failed to send file.") + return + + role = u["role"] if u else "user" + is_new = referred_by is not None + text = ( + f"{'Welcome!' if is_new else 'Hello!'} *{user.full_name or user.username or 'Friend'}*\n\n" + f"{WELCOME_MSG}\n\n" + f"{'Your referrer got a bonus!' if is_new else ''}" + ) + await update.message.reply_text( + text.strip(), + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_main(user.id, role), + ) diff --git a/bot/handlers/youtube.py b/bot/handlers/youtube.py new file mode 100644 index 0000000..a0ada76 --- /dev/null +++ b/bot/handlers/youtube.py @@ -0,0 +1,435 @@ +""" +YouTube / social media download handlers. +""" + +import os +import io +import re +import asyncio +import tempfile +import glob as _glob +import logging +from typing import List + +import yt_dlp + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.constants import ParseMode +from telegram.ext import ContextTypes, ConversationHandler + +from bot.config import ( + COOKIE_FILE, + MAX_YT_SIZE_MB, + MAX_QUEUE, + YT_SLEEP, + VIP_LEVELS, + STATE_WAITING_YT_URL, +) +from bot.decorators import guard +from bot.helpers import ( + safe_edit, + sanitize_text, + is_yt_url, + is_social_url, + fmt_duration, + make_progress_bar, +) +from bot.keyboards import kb_back_main, kb_cancel + +logger = logging.getLogger("BOT.youtube") + +youtube_queue: List[dict] = [] +is_processing = False + + +def _cookie_args_cmd() -> list: + if os.path.exists(COOKIE_FILE) and os.path.getsize(COOKIE_FILE) > 0: + return ["--cookies", COOKIE_FILE] + return ["--extractor-args", "youtube:player_client=android,ios,web"] + + +def _cookie_args_opts() -> dict: + if os.path.exists(COOKIE_FILE) and os.path.getsize(COOKIE_FILE) > 0: + return {"cookiefile": COOKIE_FILE} + return {"extractor_args": {"youtube": {"player_client": ["android", "ios", "web"]}}} + + +async def get_media_info(url: str) -> dict: + loop = asyncio.get_event_loop() + opts = { + "quiet": True, + "no_warnings": True, + "skip_download": True, + "socket_timeout": 30, + **_cookie_args_opts(), + } + + def _extract(): + with yt_dlp.YoutubeDL(opts) as ydl: + return ydl.extract_info(url, download=False) + + info = await loop.run_in_executor(None, _extract) + if not info: + raise ValueError("Could not fetch video info") + + filesize = 0 + for f in sorted(info.get("formats", []), key=lambda x: x.get("height") or 0, reverse=True): + if (f.get("height") or 999) <= 720: + fs = f.get("filesize") or f.get("filesize_approx") or 0 + if fs: + filesize = fs + break + if not filesize: + filesize = info.get("filesize") or info.get("filesize_approx") or 0 + + return { + "title": (info.get("title") or "video")[:80], + "duration": info.get("duration") or 0, + "filesize": filesize, + "uploader": info.get("uploader") or "Unknown", + "view_count": info.get("view_count") or 0, + "thumbnail": info.get("thumbnail") or "", + "is_yt": is_yt_url(url), + } + + +async def download_media(url: str, audio_only: bool = False, quality: str = "720") -> bytes: + if audio_only: + fmt = "bestaudio[ext=m4a]/bestaudio[ext=webm]/bestaudio/best" + else: + h = int(quality) + fmt = "/".join([ + f"bestvideo[ext=mp4][height<={h}]+bestaudio[ext=m4a]", + f"bestvideo[ext=mp4][height<={h}]+bestaudio", + f"bestvideo[height<={h}]+bestaudio[ext=m4a]", + f"bestvideo[height<={h}]+bestaudio", + f"best[height<={h}][ext=mp4]", + f"best[height<={h}]", + "bestvideo[ext=mp4]+bestaudio[ext=m4a]", + "bestvideo+bestaudio", + "best[ext=mp4]", + "best", + ]) + + with tempfile.TemporaryDirectory(prefix="ytdl_") as tmpdir: + out_tmpl = os.path.join(tmpdir, "video.%(ext)s") + cmd = ["yt-dlp", "-f", fmt] + if audio_only: + cmd += ["-x", "--audio-format", "mp3", "--audio-quality", "0"] + else: + cmd += ["--merge-output-format", "mp4"] + cmd += _cookie_args_cmd() + cmd += [ + "--no-playlist", "--no-part", + "--retries", "3", "--fragment-retries", "3", + "--socket-timeout", "60", + "-o", out_tmpl, url, + ] + + proc = await asyncio.create_subprocess_exec( + *cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE + ) + try: + _, stderr = await asyncio.wait_for(proc.communicate(), timeout=360) + except asyncio.TimeoutError: + proc.kill() + raise RuntimeError("Download timed out.") + + if proc.returncode != 0: + err = stderr.decode("utf-8", errors="ignore") + logger.warning(f"yt-dlp stderr: {err[:600]}") + el = err.lower() + if "sign in" in el or "login" in el or "age" in el: + raise RuntimeError("Video requires login. Admin must set cookies.") + if "private" in el: + raise RuntimeError("This video is private.") + if "not available" in el or "unavailable" in el: + raise RuntimeError("Video is unavailable or region-restricted.") + if "copyright" in el: + raise RuntimeError("Video cannot be downloaded due to copyright.") + if "live" in el and "stream" in el: + raise RuntimeError("Live streams cannot be downloaded.") + raise RuntimeError("Download failed. The video may have restrictions.") + + files = _glob.glob(os.path.join(tmpdir, "video.*")) + if not files: + raise RuntimeError("No file was downloaded.") + with open(files[0], "rb") as fh: + data = fh.read() + + if not data: + raise RuntimeError("Downloaded file is empty.") + return data + + +async def process_queue(context: ContextTypes.DEFAULT_TYPE): + global youtube_queue, is_processing + if is_processing: + return + is_processing = True + db = context.bot_data["db"] + try: + while youtube_queue: + job = youtube_queue.pop(0) + uid = job["uid"] + url = job["url"] + title = job["title"] + audio = job.get("audio", False) + quality = job.get("quality", "720") + try: + icon = "audio" if audio else "video" + short_title = title[:40] + ("..." if len(title) > 40 else "") + prog_msg = await context.bot.send_message( + uid, + f"*Downloading {icon}*\n" + f"--------------------\n" + f"`{short_title}`\n\n" + f"`{make_progress_bar(0, 12)}` 0%\n" + f"Connecting...", + parse_mode=ParseMode.MARKDOWN, + ) + + async def _update_progress(pct: int, status: str): + try: + bar = make_progress_bar(pct, 12) + await safe_edit( + prog_msg, + f"*Downloading {icon}*\n" + f"--------------------\n" + f"`{short_title}`\n\n" + f"`{bar}` {pct}%\n" + f"{status}", + parse_mode=ParseMode.MARKDOWN, + ) + except Exception: + pass + + await _update_progress(10, "Connecting to server...") + download_task = asyncio.create_task( + download_media(url, audio_only=audio, quality=quality) + ) + + stages = [ + (25, "Fetching info..."), + (50, "Downloading..."), + (75, "Processing..."), + (90, "Preparing to send..."), + ] + for pct, status in stages: + await asyncio.sleep(4) + if download_task.done(): + break + await _update_progress(pct, status) + + data = await download_task + await _update_progress(100, "Ready!") + + size_mb = len(data) / (1024 * 1024) + if size_mb > MAX_YT_SIZE_MB: + await context.bot.send_message( + uid, + f"File too large ({size_mb:.1f} MB > {MAX_YT_SIZE_MB} MB limit).", + ) + continue + + file_buf = io.BytesIO(data) + safe = re.sub(r"[^\w\s\-]", "", title)[:50] + if audio: + file_buf.name = f"{safe}.mp3" + await context.bot.send_audio( + uid, audio=file_buf, title=title, + caption=f"*{title[:60]}*\n{size_mb:.1f} MB", + parse_mode=ParseMode.MARKDOWN, + ) + else: + file_buf.name = f"{safe}.mp4" + await context.bot.send_video( + uid, video=file_buf, + caption=( + f"*{title[:60]}*\n" + f"--------------------\n" + f"{size_mb:.1f} MB | {quality}p | Ready" + ), + parse_mode=ParseMode.MARKDOWN, + supports_streaming=True, + ) + + try: + await prog_msg.delete() + except Exception: + pass + + await db.add_log(uid, "yt_download", f"title={title[:50]} size={size_mb:.1f}MB quality={quality}") + async with (await db.connect()).execute( + "UPDATE bot_stats SET value=value+1 WHERE key='total_downloads'" + ): + pass + await (await db.connect()).commit() + + except Exception as e: + logger.error(f"Queue error: {e}") + try: + await context.bot.send_message(uid, f"Error: {str(e)[:200]}") + except Exception: + pass + + await asyncio.sleep(YT_SLEEP) + finally: + is_processing = False + + +@guard() +async def cb_yt(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await safe_edit( + q, + "*Video Download*\n\n" + "Send a link from:\n" + "YouTube | Instagram | TikTok | Twitter\n\n" + "Supported formats: video + audio", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_cancel(), + ) + return STATE_WAITING_YT_URL + + +async def handle_yt_url(update: Update, context: ContextTypes.DEFAULT_TYPE): + uid = update.effective_user.id + db = context.bot_data["db"] + url = sanitize_text(update.message.text, 500).strip() + + if not is_yt_url(url) and not is_social_url(url): + await update.message.reply_text("Invalid URL. Send a valid video link.", reply_markup=kb_cancel()) + return STATE_WAITING_YT_URL + + if len(youtube_queue) >= MAX_QUEUE: + await update.message.reply_text("Queue is full. Try again later.", reply_markup=kb_back_main()) + return ConversationHandler.END + + status_msg = await update.message.reply_text("Fetching video info...") + try: + info = await asyncio.wait_for(get_media_info(url), timeout=45) + size_mb = info["filesize"] / (1024 * 1024) if info["filesize"] else 0 + context.user_data["yt_url"] = url + context.user_data["yt_info"] = info + + size_text = f"{size_mb:.1f} MB" if size_mb else "Unknown" + views = f"{info['view_count']:,}" if info["view_count"] else "--" + text = ( + f"*Video Info*\n\n" + f"*{info['title']}*\n\n" + f"Channel: {info['uploader']}\n" + f"Duration: `{fmt_duration(info['duration'])}`\n" + f"Size: `{size_text}`\n" + f"Views: {views}\n\n" + f"*Select download quality:*" + ) + + is_vip = await db.is_vip(uid) + vip_level = await db.get_vip_level(uid) or "silver" + max_q = VIP_LEVELS.get(vip_level, {}).get("max_yt_quality", "720") if is_vip else "720" + + quality_btns = [] + if info.get("is_yt"): + available = [("360", "360p"), ("480", "480p"), ("720", "720p HD")] + if is_vip and max_q == "1080": + available.append(("1080", "1080p FHD")) + quality_btns = [ + InlineKeyboardButton(label, callback_data=f"yt_q_{q_val}") + for q_val, label in available + ] + kb_rows = [] + if quality_btns: + kb_rows.append(quality_btns[:2]) + if len(quality_btns) > 2: + kb_rows.append(quality_btns[2:]) + kb_rows.append([InlineKeyboardButton("Audio (MP3)", callback_data="yt_dl_audio")]) + if not info.get("is_yt"): + kb_rows = [ + [InlineKeyboardButton("Download Video", callback_data="yt_dl_video")], + [InlineKeyboardButton("Download Audio", callback_data="yt_dl_audio")], + ] + kb_rows.append([InlineKeyboardButton("Cancel", callback_data="main")]) + + await safe_edit(status_msg, text, parse_mode=ParseMode.MARKDOWN, reply_markup=InlineKeyboardMarkup(kb_rows)) + except asyncio.TimeoutError: + await safe_edit(status_msg, "Timed out fetching info. Try again.") + except Exception as e: + await safe_edit(status_msg, f"Error: {str(e)[:200]}") + return ConversationHandler.END + + return ConversationHandler.END + + +@guard() +async def cb_yt_quality(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + quality = q.data.replace("yt_q_", "") + url = context.user_data.get("yt_url") + info = context.user_data.get("yt_info", {}) + if not url: + await safe_edit(q, "Session expired.", reply_markup=kb_back_main()) + return + uid = q.from_user.id + title = info.get("title", "video") + job = {"uid": uid, "url": url, "title": title, "audio": False, "quality": quality} + if await db.is_vip(uid): + youtube_queue.insert(0, job) + pos_text = "VIP priority!" + else: + youtube_queue.append(job) + pos_text = f"Queue position: {len(youtube_queue)}" + await safe_edit( + q, + f"*Added to queue*\n" + f"--------------------\n" + f"`{title[:50]}`\n" + f"Quality: {quality}p\n\n" + f"{pos_text}", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_back_main(), + ) + asyncio.create_task(process_queue(context)) + context.user_data.pop("yt_url", None) + context.user_data.pop("yt_info", None) + + +@guard() +async def cb_yt_dl(update: Update, context: ContextTypes.DEFAULT_TYPE): + q = update.callback_query + await q.answer() + db = context.bot_data["db"] + uid = q.from_user.id + audio = q.data == "yt_dl_audio" + url = context.user_data.get("yt_url") + info = context.user_data.get("yt_info", {}) + if not url: + await safe_edit(q, "Session expired.", reply_markup=kb_back_main()) + return + title = info.get("title", "video") + quality = "720" + if await db.is_vip(uid): + lvl = await db.get_vip_level(uid) or "silver" + quality = VIP_LEVELS.get(lvl, {}).get("max_yt_quality", "720") + job = {"uid": uid, "url": url, "title": title, "audio": audio, "quality": quality} + if await db.is_vip(uid): + youtube_queue.insert(0, job) + pos_text = "VIP priority!" + else: + youtube_queue.append(job) + pos_text = f"Queue position: {len(youtube_queue)}" + icon = "Audio" if audio else "Video" + await safe_edit( + q, + f"*Added to queue*\n" + f"--------------------\n" + f"{icon}: `{title[:50]}`\n\n" + f"{pos_text}", + parse_mode=ParseMode.MARKDOWN, + reply_markup=kb_back_main(), + ) + asyncio.create_task(process_queue(context)) + context.user_data.pop("yt_url", None) + context.user_data.pop("yt_info", None) diff --git a/bot/helpers.py b/bot/helpers.py new file mode 100644 index 0000000..030ace2 --- /dev/null +++ b/bot/helpers.py @@ -0,0 +1,156 @@ +""" +Shared utility functions: formatting, security checks, QR codes, etc. +""" + +import io +import re +import time +import string +import secrets +import logging +from collections import defaultdict +from datetime import datetime +from typing import Dict, List + +import qrcode +from telegram.error import BadRequest + +from bot.config import RATE_LIMIT, VIP_LEVELS + +logger = logging.getLogger("BOT.helpers") + +_rate_tracker: Dict[int, List[float]] = defaultdict(list) + + +# ── Safe edit wrappers ──────────────────────────── + +async def safe_edit(msg_or_query, text: str, **kwargs): + try: + if hasattr(msg_or_query, "edit_message_text"): + await msg_or_query.edit_message_text(text, **kwargs) + else: + await msg_or_query.edit_text(text, **kwargs) + except BadRequest as e: + if "message is not modified" not in str(e).lower(): + raise + + +async def safe_edit_caption(q, caption: str, **kwargs): + try: + await q.edit_message_caption(caption=caption, **kwargs) + except BadRequest as e: + if "message is not modified" not in str(e).lower(): + raise + + +# ── Security ────────────────────────────────────── + +def rate_limit_check(uid: int) -> bool: + now = time.time() + calls = _rate_tracker[uid] + calls[:] = [t for t in calls if now - t < 60.0] + if len(calls) >= RATE_LIMIT: + return True + calls.append(now) + return False + + +def sanitize_text(text: str, max_len: int = 500) -> str: + if not text: + return "" + text = text.strip() + text = re.sub(r"[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]", "", text) + return text[:max_len] + + +def is_valid_tg_id(val: str) -> bool: + return bool(re.match(r"^\d{5,15}$", val.strip())) + + +def generate_secure_code(length: int = 10) -> str: + alphabet = string.ascii_uppercase + string.digits + return "".join(secrets.choice(alphabet) for _ in range(length)) + + +def is_social_url(url: str) -> bool: + return bool( + re.search( + r"(instagram\.com|instagr\.am|tiktok\.com|vm\.tiktok\.com" + r"|twitter\.com|x\.com|facebook\.com|fb\.watch)", + url, + ) + ) + + +def is_yt_url(url: str) -> bool: + return bool(re.search(r"(youtube\.com|youtu\.be|yt\.be)", url)) + + +# ── Formatting ──────────────────────────────────── + +def fmt_time(iso: str) -> str: + if not iso: + return "N/A" + try: + return datetime.fromisoformat(iso).strftime("%Y/%m/%d %H:%M") + except ValueError: + return "N/A" + + +def fmt_size(b: int) -> str: + if b < 1024: + return f"{b} B" + if b < 1024 ** 2: + return f"{b / 1024:.1f} KB" + if b < 1024 ** 3: + return f"{b / 1024 ** 2:.1f} MB" + return f"{b / 1024 ** 3:.1f} GB" + + +def fmt_duration(secs: int) -> str: + if not secs: + return "N/A" + h, rem = divmod(int(secs), 3600) + m, s = divmod(rem, 60) + if h: + return f"{h}:{m:02d}:{s:02d}" + return f"{m}:{s:02d}" + + +def vip_level_info(user: dict) -> str: + if not user.get("is_vip"): + return "--" + level = user.get("vip_level", "silver") + info = VIP_LEVELS.get(level, {}) + name = info.get("name", level) + if user.get("vip_expiry"): + try: + exp = datetime.fromisoformat(user["vip_expiry"]) + days_left = (exp - datetime.now()).days + return f"{name} | {days_left} days left" + except ValueError: + pass + return name + + +def make_progress_bar(pct: int, width: int = 10) -> str: + filled = int(width * pct / 100) + return "\u2588" * filled + "\u2591" * (width - filled) + + +# ── QR Code ─────────────────────────────────────── + +def generate_qr(text: str) -> io.BytesIO: + qr = qrcode.QRCode( + version=1, + box_size=7, + border=3, + error_correction=qrcode.constants.ERROR_CORRECT_M, + ) + qr.add_data(text) + qr.make(fit=True) + img = qr.make_image(fill_color="#1a1a2e", back_color="white") + bio = io.BytesIO() + img.save(bio, format="PNG") + bio.seek(0) + return bio diff --git a/bot/keyboards.py b/bot/keyboards.py new file mode 100644 index 0000000..c532142 --- /dev/null +++ b/bot/keyboards.py @@ -0,0 +1,103 @@ +""" +Inline keyboard builders. +""" + +from telegram import InlineKeyboardButton, InlineKeyboardMarkup + +from bot.config import SUPER_ADMIN_ID + + +def is_admin_check(tg_id: int, db_role: str) -> bool: + return tg_id == SUPER_ADMIN_ID or db_role in ("admin", "super_admin") + + +def kb_main(uid: int, role: str = "user") -> InlineKeyboardMarkup: + rows = [ + [ + InlineKeyboardButton("Config", callback_data="claim"), + InlineKeyboardButton("Download", callback_data="yt"), + ], + [ + InlineKeyboardButton("File2Link", callback_data="file"), + InlineKeyboardButton("VIP", callback_data="vip"), + ], + [ + InlineKeyboardButton("Wallet", callback_data="wallet"), + InlineKeyboardButton("Profile", callback_data="profile"), + ], + [ + InlineKeyboardButton("Proxy", callback_data="proxy_menu"), + InlineKeyboardButton("Speed Test", callback_data="speed_test"), + ], + [ + InlineKeyboardButton("Referral", callback_data="referral"), + InlineKeyboardButton("Help", callback_data="help"), + ], + ] + if is_admin_check(uid, role): + rows.append([InlineKeyboardButton("Admin Panel", callback_data="admin")]) + return InlineKeyboardMarkup(rows) + + +def kb_back_main() -> InlineKeyboardMarkup: + return InlineKeyboardMarkup( + [[InlineKeyboardButton("Home", callback_data="main")]] + ) + + +def kb_cancel() -> InlineKeyboardMarkup: + return InlineKeyboardMarkup( + [[InlineKeyboardButton("Cancel", callback_data="main")]] + ) + + +def kb_admin(uid: int) -> InlineKeyboardMarkup: + rows = [ + [ + InlineKeyboardButton("Configs", callback_data="admin_configs"), + InlineKeyboardButton("Payments", callback_data="admin_payments"), + ], + [ + InlineKeyboardButton("Users", callback_data="admin_users"), + InlineKeyboardButton("Broadcast", callback_data="admin_broadcast"), + ], + [ + InlineKeyboardButton("Give VIP", callback_data="admin_give_vip"), + InlineKeyboardButton("Discount", callback_data="admin_discount"), + ], + [ + InlineKeyboardButton("Wallet Charge", callback_data="admin_wallet_charge"), + InlineKeyboardButton("Proxies", callback_data="admin_proxy"), + ], + ] + if uid == SUPER_ADMIN_ID: + rows += [ + [ + InlineKeyboardButton("Admins", callback_data="super_admin"), + InlineKeyboardButton("Full Stats", callback_data="full_stats"), + ], + [ + InlineKeyboardButton("DB Backup", callback_data="admin_backup"), + InlineKeyboardButton("YT Cookie", callback_data="set_cookie"), + ], + ] + rows.append([InlineKeyboardButton("Home", callback_data="main")]) + return InlineKeyboardMarkup(rows) + + +def kb_proxy_menu() -> InlineKeyboardMarkup: + return InlineKeyboardMarkup([ + [ + InlineKeyboardButton("HTTP/S Proxies", callback_data="proxy_http"), + InlineKeyboardButton("SOCKS5 Proxies", callback_data="proxy_socks5"), + ], + [ + InlineKeyboardButton("SOCKS4 Proxies", callback_data="proxy_socks4"), + InlineKeyboardButton("All Proxies", callback_data="proxy_all"), + ], + [ + InlineKeyboardButton("Proxy Stats", callback_data="proxy_stats"), + InlineKeyboardButton("Check Proxy", callback_data="proxy_check"), + ], + [InlineKeyboardButton("Home", callback_data="main")], + ]) diff --git a/bot/main.py b/bot/main.py new file mode 100644 index 0000000..8d21dcc --- /dev/null +++ b/bot/main.py @@ -0,0 +1,327 @@ +""" +Bot entry point: initializes DB, registers handlers, starts polling. +""" + +import logging +import sys +from datetime import datetime + +from telegram import Update +from telegram.error import BadRequest +from telegram.ext import ( + ApplicationBuilder, + CommandHandler, + CallbackQueryHandler, + MessageHandler, + ConversationHandler, + ContextTypes, + filters, +) + +from bot.config import ( + BOT_TOKEN, + SUPER_ADMIN_ID, + CONFIG, + PROXY_UPDATE_INTERVAL, + STATE_WAITING_YT_URL, + STATE_WAITING_FILE, + STATE_WAITING_RECEIPT, + STATE_WAITING_DISCOUNT, + STATE_WAITING_CONFIG_TEXT, + STATE_DELETING_CONFIG, + STATE_SET_COOKIE, + STATE_BAN_USER, + STATE_UNBAN_USER, + STATE_SEARCHING_USER, + STATE_BROADCASTING, + STATE_GIVE_VIP, + STATE_MANAGE_ADMIN, + STATE_ADD_DISCOUNT, + STATE_CHECKING_PROXY, + STATE_ADMIN_WALLET_CHARGE, +) +from bot.database import Database +from bot.keyboards import kb_main +from bot.server import start_web_server + +from bot.handlers.start import cmd_start +from bot.handlers.main_menu import ( + cb_main, + cb_check_join, + cb_help, + cb_profile, + cb_referral, + cb_claim, + cb_vip, + cb_buy_vip, + cb_receipt_prompt, + handle_receipt, + cb_apply_discount, + handle_discount_code, + cb_wallet, + cb_wallet_buy_vip, + cb_wallet_pay, +) +from bot.handlers.youtube import ( + cb_yt, + handle_yt_url, + cb_yt_quality, + cb_yt_dl, +) +from bot.handlers.file import cb_file, handle_file +from bot.handlers.proxy import ( + cb_proxy_menu, + cb_proxy_list, + cb_proxy_copy, + cb_proxy_stats, + cb_proxy_refresh, + cb_proxy_check, + handle_proxy_check, + cb_speed_test, + cb_admin_proxy, + cb_admin_proxy_clean, + scheduled_proxy_update, +) +from bot.handlers.admin import ( + cb_admin, + cmd_admin, + cb_admin_configs, + cb_add_config_prompt, + handle_add_config, + cb_del_config_prompt, + handle_del_config, + cb_list_configs, + cb_admin_users, + cb_ban_prompt, + handle_ban_user, + cb_unban_prompt, + handle_unban_user, + cb_search_user_prompt, + handle_search_user, + cb_banned_list, + cb_admin_give_vip, + handle_give_vip, + cb_admin_broadcast, + handle_broadcast, + cb_admin_payments, + cb_pay_action, + cb_admin_discount, + cb_add_discount_prompt, + handle_add_discount, + cb_super_admin, + cb_add_admin_prompt, + cb_remove_admin_prompt, + handle_admin_action, + cb_full_stats, + cb_admin_backup, + cb_set_cookie, + handle_cookie, + cmd_confirm_payment, + cmd_reject_payment, + cb_admin_wallet_charge, + handle_wallet_charge, +) + +# ── Logging ─────────────────────────────────────── +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", + handlers=[ + logging.FileHandler( + f"logs/bot_{datetime.now():%Y%m%d}.log", encoding="utf-8" + ), + logging.StreamHandler(sys.stdout), + ], +) +logger = logging.getLogger("BOT") + + +# ── Daily report ───────────────────────────────── +async def send_daily_report(context: ContextTypes.DEFAULT_TYPE): + if not SUPER_ADMIN_ID: + return + db = context.bot_data["db"] + try: + from telegram.constants import ParseMode + stats = await db.get_daily_stats() + text = ( + f"*Daily Report -- {datetime.now().strftime('%Y/%m/%d')}*\n" + f"--------------------\n" + f"New users today: *{stats['new_users_today']}*\n" + f"New users yesterday: {stats['new_users_yesterday']}\n" + f"Active VIPs: {stats['active_vips']}\n" + f"Payments today: {stats['payments_today']} ({stats['revenue_today']:,})\n" + f"Downloads today: {stats['downloads_today']}\n" + f"Configs today: {stats['claims_today']}\n" + f"Total users: {stats['total_users']:,}" + ) + await context.bot.send_message(SUPER_ADMIN_ID, text, parse_mode=ParseMode.MARKDOWN) + except Exception as e: + logger.error(f"Daily report error: {e}") + + +# ── Error handler ──────────────────────────────── +async def error_handler(update: object, context: ContextTypes.DEFAULT_TYPE): + logger.error(f"Unhandled error: {context.error}", exc_info=context.error) + if isinstance(update, Update) and update.effective_message: + try: + await update.effective_message.reply_text("An error occurred. Please try again.") + except Exception: + pass + + +# ── Cancel ─────────────────────────────────────── +async def cmd_cancel(update: Update, context: ContextTypes.DEFAULT_TYPE): + context.user_data.clear() + db = context.bot_data["db"] + uid = update.effective_user.id if update.effective_user else 0 + role = await db.get_role(uid) if uid else "user" + kb = kb_main(uid, role) if update.effective_user else None + if update.message: + await update.message.reply_text("Cancelled.", reply_markup=kb) + elif update.callback_query: + await update.callback_query.answer() + try: + await update.callback_query.edit_message_text("Cancelled.", reply_markup=kb) + except BadRequest: + pass + return ConversationHandler.END + + +# ── Post-init ──────────────────────────────────── +async def post_init(app): + db = Database() + await db.init_db() + app.bot_data["db"] = db + + await db.clean_expired_files() + await start_web_server(app.bot, db) + + # daily report + hour = CONFIG.get("daily_report_hour", 8) + app.job_queue.run_daily( + send_daily_report, + time=datetime.now().replace(hour=hour, minute=0, second=0).time(), + ) + + # periodic proxy update + app.job_queue.run_repeating( + scheduled_proxy_update, + interval=PROXY_UPDATE_INTERVAL * 60, + first=10, + ) + + logger.info("Bot initialized!") + + +def main(): + app = ( + ApplicationBuilder() + .token(BOT_TOKEN) + .post_init(post_init) + .concurrent_updates(True) + .build() + ) + app.add_error_handler(error_handler) + + # ── Commands ────────────────────────────────── + app.add_handler(CommandHandler("start", cmd_start)) + app.add_handler(CommandHandler("admin", cmd_admin)) + app.add_handler(CommandHandler("cancel", cmd_cancel)) + app.add_handler(MessageHandler(filters.Regex(r"^/confirm_\d+$"), cmd_confirm_payment)) + app.add_handler(MessageHandler(filters.Regex(r"^/reject_\d+$"), cmd_reject_payment)) + + # ── Callbacks ───────────────────────────────── + app.add_handler(CallbackQueryHandler(cb_main, pattern="^main$")) + app.add_handler(CallbackQueryHandler(cb_check_join, pattern="^check_join$")) + app.add_handler(CallbackQueryHandler(cb_help, pattern="^help$")) + app.add_handler(CallbackQueryHandler(cb_profile, pattern="^profile$")) + app.add_handler(CallbackQueryHandler(cb_referral, pattern="^referral$")) + app.add_handler(CallbackQueryHandler(cb_claim, pattern="^claim$")) + app.add_handler(CallbackQueryHandler(cb_vip, pattern="^vip$")) + app.add_handler(CallbackQueryHandler(cb_buy_vip, pattern="^buy_")) + app.add_handler(CallbackQueryHandler(cb_wallet, pattern="^wallet$")) + app.add_handler(CallbackQueryHandler(cb_wallet_buy_vip, pattern="^wallet_buy_vip$")) + app.add_handler(CallbackQueryHandler(cb_wallet_pay, pattern="^wallet_pay_")) + app.add_handler(CallbackQueryHandler(cb_admin, pattern="^admin$")) + app.add_handler(CallbackQueryHandler(cb_admin_configs, pattern="^admin_configs$")) + app.add_handler(CallbackQueryHandler(cb_list_configs, pattern="^list_configs_(free|vip)$")) + app.add_handler(CallbackQueryHandler(cb_admin_payments, pattern="^admin_payments$")) + app.add_handler(CallbackQueryHandler(cb_admin_users, pattern="^admin_users$")) + app.add_handler(CallbackQueryHandler(cb_banned_list, pattern="^banned_list$")) + app.add_handler(CallbackQueryHandler(cb_admin_give_vip, pattern="^admin_give_vip$")) + app.add_handler(CallbackQueryHandler(cb_admin_discount, pattern="^admin_discount$")) + app.add_handler(CallbackQueryHandler(cb_add_discount_prompt, pattern="^add_discount_code$")) + app.add_handler(CallbackQueryHandler(cb_super_admin, pattern="^super_admin$")) + app.add_handler(CallbackQueryHandler(cb_full_stats, pattern="^full_stats$")) + app.add_handler(CallbackQueryHandler(cb_admin_backup, pattern="^admin_backup$")) + app.add_handler(CallbackQueryHandler(cb_yt_quality, pattern="^yt_q_")) + app.add_handler(CallbackQueryHandler(cb_yt_dl, pattern="^yt_dl_(video|audio)$")) + app.add_handler(CallbackQueryHandler(cb_pay_action, pattern=r"^pay_(confirm|reject)_\d+$")) + + # ── Proxy callbacks ─────────────────────────── + app.add_handler(CallbackQueryHandler(cb_proxy_menu, pattern="^proxy_menu$")) + app.add_handler(CallbackQueryHandler(cb_proxy_list, pattern="^proxy_(http|socks4|socks5|all)$")) + app.add_handler(CallbackQueryHandler(cb_proxy_copy, pattern="^proxy_copy_")) + app.add_handler(CallbackQueryHandler(cb_proxy_stats, pattern="^proxy_stats$")) + app.add_handler(CallbackQueryHandler(cb_proxy_refresh, pattern="^proxy_refresh$")) + app.add_handler(CallbackQueryHandler(cb_speed_test, pattern="^speed_test$")) + app.add_handler(CallbackQueryHandler(cb_admin_proxy, pattern="^admin_proxy$")) + app.add_handler(CallbackQueryHandler(cb_admin_proxy_clean, pattern="^admin_proxy_clean$")) + + # ── Conversation handler ────────────────────── + conv = ConversationHandler( + entry_points=[ + CallbackQueryHandler(cb_yt, pattern="^yt$"), + CallbackQueryHandler(cb_file, pattern="^file$"), + CallbackQueryHandler(cb_receipt_prompt, pattern="^send_receipt$"), + CallbackQueryHandler(cb_apply_discount, pattern="^apply_discount$"), + CallbackQueryHandler(cb_add_config_prompt, pattern="^add_config$"), + CallbackQueryHandler(cb_del_config_prompt, pattern="^del_config$"), + CallbackQueryHandler(cb_set_cookie, pattern="^set_cookie$"), + CallbackQueryHandler(cb_ban_prompt, pattern="^ban_user$"), + CallbackQueryHandler(cb_unban_prompt, pattern="^unban_user$"), + CallbackQueryHandler(cb_search_user_prompt, pattern="^search_user$"), + CallbackQueryHandler(cb_admin_broadcast, pattern="^admin_broadcast$"), + CallbackQueryHandler(cb_admin_give_vip, pattern="^admin_give_vip$"), + CallbackQueryHandler(cb_add_admin_prompt, pattern="^add_admin$"), + CallbackQueryHandler(cb_remove_admin_prompt, pattern="^remove_admin$"), + CallbackQueryHandler(cb_add_discount_prompt, pattern="^add_discount_code$"), + CallbackQueryHandler(cb_proxy_check, pattern="^proxy_check$"), + CallbackQueryHandler(cb_admin_wallet_charge, pattern="^admin_wallet_charge$"), + ], + states={ + STATE_WAITING_YT_URL: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_yt_url)], + STATE_WAITING_FILE: [MessageHandler(filters.ALL & ~filters.COMMAND, handle_file)], + STATE_WAITING_RECEIPT: [MessageHandler(filters.PHOTO, handle_receipt)], + STATE_WAITING_DISCOUNT: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_discount_code)], + STATE_WAITING_CONFIG_TEXT: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_add_config)], + STATE_DELETING_CONFIG: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_del_config)], + STATE_SET_COOKIE: [MessageHandler(filters.Document.ALL, handle_cookie)], + STATE_BAN_USER: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_ban_user)], + STATE_UNBAN_USER: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_unban_user)], + STATE_SEARCHING_USER: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_search_user)], + STATE_BROADCASTING: [MessageHandler(filters.ALL & ~filters.COMMAND, handle_broadcast)], + STATE_GIVE_VIP: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_give_vip)], + STATE_MANAGE_ADMIN: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_admin_action)], + STATE_ADD_DISCOUNT: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_add_discount)], + STATE_CHECKING_PROXY: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_proxy_check)], + STATE_ADMIN_WALLET_CHARGE: [MessageHandler(filters.TEXT & ~filters.COMMAND, handle_wallet_charge)], + }, + fallbacks=[ + CommandHandler("cancel", cmd_cancel), + CallbackQueryHandler(cmd_cancel, pattern="^main$"), + ], + allow_reentry=True, + per_user=True, + per_chat=True, + ) + app.add_handler(conv) + + logger.info("Bot v4.0 running!") + logger.info(f"Super Admin: {SUPER_ADMIN_ID}") + app.run_polling(drop_pending_updates=True) + + +if __name__ == "__main__": + main() diff --git a/bot/server.py b/bot/server.py new file mode 100644 index 0000000..1154835 --- /dev/null +++ b/bot/server.py @@ -0,0 +1,88 @@ +""" +HTTP streaming server for file downloads. +""" + +import re +import logging +from datetime import datetime + +import aiohttp +from aiohttp import web + +from bot.config import BOT_TOKEN, WEB_SERVER_HOST, WEB_SERVER_PORT + +logger = logging.getLogger("BOT.server") + + +async def stream_handler(request: web.Request) -> web.Response: + f_uid = request.match_info.get("f_uid", "") + if not re.match(r"^[\w\-]{10,}$", f_uid): + return web.Response(status=400, text="Bad Request") + + db = request.app["db"] + file_data = await db.get_file(f_uid) + if not file_data: + return web.Response(status=404, text="File not found") + + if file_data.get("expires_at") and datetime.now() > datetime.fromisoformat( + file_data["expires_at"] + ): + return web.Response(status=410, text="Link expired") + + await db.increment_file_downloads(f_uid) + bot = request.app["bot"] + try: + tg_file = await bot.get_file(file_data["file_id"]) + file_url = f"https://api.telegram.org/file/bot{BOT_TOKEN}/{tg_file.file_path}" + resp = web.StreamResponse( + status=200, + headers={ + "Content-Type": file_data.get("mime_type", "application/octet-stream"), + "Content-Disposition": f'attachment; filename="{file_data["file_name"]}"', + }, + ) + await resp.prepare(request) + async with request.app["http_client"].get(file_url) as r: + async for chunk in r.content.iter_chunked(512 * 1024): + await resp.write(chunk) + return resp + except Exception as e: + logger.error(f"Stream error: {e}") + return web.Response(status=500, text="Internal error") + + +async def health_handler(request: web.Request) -> web.Response: + return web.Response(text="OK") + + +async def proxy_api_handler(request: web.Request) -> web.Response: + """REST API endpoint for proxy list.""" + db = request.app["db"] + protocol = request.query.get("protocol") + try: + limit = min(int(request.query.get("limit", "50")), 200) + except (ValueError, TypeError): + limit = 50 + proxies = await db.get_alive_proxies(protocol=protocol, limit=limit) + lines = [f"{p['ip']}:{p['port']}" for p in proxies] + return web.Response(text="\n".join(lines), content_type="text/plain") + + +async def start_web_server(bot, db): + app = web.Application() + app["bot"] = bot + app["db"] = db + app["http_client"] = aiohttp.ClientSession() + + async def close_session(app): + await app["http_client"].close() + + app.on_cleanup.append(close_session) + app.router.add_get("/stream/{f_uid}", stream_handler) + app.router.add_get("/health", health_handler) + app.router.add_get("/api/proxies", proxy_api_handler) + runner = web.AppRunner(app) + await runner.setup() + site = web.TCPSite(runner, WEB_SERVER_HOST, WEB_SERVER_PORT) + await site.start() + logger.info(f"HTTP server: {WEB_SERVER_HOST}:{WEB_SERVER_PORT}") diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..9e93222 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,7 @@ +python-telegram-bot[job-queue]>=20.7 +python-dotenv>=1.0.0 +aiohttp>=3.9.0 +aiosqlite>=0.19.0 +yt-dlp>=2024.1.0 +qrcode[pil]>=7.4.0 +Pillow>=10.0.0 diff --git a/run.py b/run.py new file mode 100644 index 0000000..07b3448 --- /dev/null +++ b/run.py @@ -0,0 +1,6 @@ +#!/usr/bin/env python3 +"""Entry point.""" +from bot.main import main + +if __name__ == "__main__": + main()