| import os |
| import logging |
| import subprocess |
| import time |
| import threading |
| import requests |
| import json |
| import zipfile |
| import shutil |
| from pathlib import Path |
| from flask import Flask, request, jsonify |
| from huggingface_hub import HfApi |
|
|
| |
| app = Flask(__name__) |
| logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s') |
| logger = logging.getLogger(__name__) |
|
|
| |
| SPACE_ID = os.getenv('SPACE_ID') |
| HF_TOKEN = os.getenv('HF_TOKEN') |
| MEGA_ACCOUNTS_JSON = os.getenv('MEGA_ACCOUNTS_JSON', '[]') |
|
|
| |
| try: |
| MEGA_ACCOUNTS = json.loads(MEGA_ACCOUNTS_JSON) |
| logger.info(f"✅ {len(MEGA_ACCOUNTS)} Mega hesabı sisteme yüklendi.") |
| except Exception as e: |
| MEGA_ACCOUNTS = [] |
| logger.warning(f"⚠️ Hesaplar yüklenemedi: {str(e)}") |
|
|
| current_account_index = 0 |
| current_task = None |
| task_lock = threading.Lock() |
| DOWNLOAD_DIR = Path("/tmp/downloads") |
| DOWNLOAD_DIR.mkdir(exist_ok=True, parents=True) |
|
|
| |
| MEGACMD_TIMEOUT = 7200 |
|
|
| |
|
|
| def run_megacmd(command, timeout=30): |
| """MegaCMD komutlarını çalıştırır ve çıktıları analiz eder.""" |
| try: |
| if isinstance(command, str): |
| command = command.split() |
| |
| logger.info(f"🔧 Komut: {' '.join(command)}") |
| result = subprocess.run(command, capture_output=True, text=True, timeout=timeout) |
| |
| output = result.stdout.strip() |
| error = result.stderr.strip() |
| combined = (output + " " + error).lower() |
| |
| |
| if any(k in combined for k in ['api_eblocked', 'temporary ban', 'connection refused', 'blocked', 'bus error']): |
| logger.error(f"🚨 IP BAN ALGILANDI: {error}") |
| return {'success': False, 'critical': True, 'error': error} |
| |
| |
| if 'bandwidth quota exceeded' in combined or 'transfer quota exceeded' in combined: |
| logger.warning(f"📊 KOTA DOLDU: {error}") |
| return {'success': False, 'bandwidth_exceeded': True, 'error': error} |
| |
| if result.returncode != 0: |
| return {'success': False, 'error': error} |
| |
| return {'success': True, 'output': output} |
| |
| except subprocess.TimeoutExpired: |
| return {'success': False, 'timeout': True, 'error': f'Timeout {timeout}s'} |
| except Exception as e: |
| return {'success': False, 'error': str(e)} |
|
|
| def init_megacmd_server(): |
| """Mega sunucusunu arka planda başlatır.""" |
| try: |
| logger.info("🚀 MegaCMD server başlatılıyor...") |
| subprocess.Popen(['mega-cmd-server'], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) |
| time.sleep(7) |
| return True |
| except Exception as e: |
| logger.error(f"❌ MegaCMD server hatası: {str(e)}") |
| return False |
|
|
| def mega_login(email, password): |
| """Mega hesabına giriş yapar.""" |
| try: |
| logger.info(f"🔐 Giriş yapılıyor: {email}") |
| run_megacmd(['mega-logout'], timeout=10) |
| time.sleep(2) |
| result = run_megacmd(['mega-login', email, password], timeout=45) |
| if result['success']: |
| logger.info(f"✅ Giriş başarılı: {email}") |
| return True |
| return False |
| except: |
| return False |
|
|
| def get_next_account(): |
| """Sıradaki hesabı getirir (Zombie Mode).""" |
| global current_account_index |
| if not MEGA_ACCOUNTS: return None |
| account = MEGA_ACCOUNTS[current_account_index] |
| current_account_index = (current_account_index + 1) % len(MEGA_ACCOUNTS) |
| return account |
|
|
| def rotate_account(): |
| """Hesap değiştirir.""" |
| logger.warning("🔄 HESAP ROTASYONU BAŞLATILIYOR...") |
| account = get_next_account() |
| if not account: return False |
| return mega_login(account['email'], account['password']) |
|
|
| def phoenix_restart(): |
| """IP Ban yiyince Space'i yeniden başlatır (Phoenix Mode).""" |
| try: |
| if not SPACE_ID or not HF_TOKEN: |
| logger.error("❌ SPACE_ID veya HF_TOKEN eksik, restart atılamıyor!") |
| return False |
| |
| logger.warning("🔥 PHOENIX MODE: IP Banlandı, Space restart ediliyor...") |
| api = HfApi(token=HF_TOKEN) |
| api.restart_space(repo_id=SPACE_ID) |
| return True |
| except Exception as e: |
| logger.error(f"❌ Restart hatası: {str(e)}") |
| return False |
|
|
| def send_webhook(webhook_url, data): |
| """Render tarafına durum bildirir.""" |
| if not webhook_url: return |
| try: |
| requests.post(webhook_url, json=data, timeout=10) |
| except: |
| pass |
|
|
| |
|
|
| def process_task(task_data): |
| global current_task |
| task_id = task_data['task_id'] |
| mega_url = task_data['mega_url'] |
| webhook_url = task_data.get('webhook_url') |
| |
| try: |
| logger.info(f"🚀 İşlem Başladı: {task_id}") |
| send_webhook(webhook_url, {'status': 'started', 'progress': 0}) |
|
|
| |
| if MEGA_ACCOUNTS: |
| acc = MEGA_ACCOUNTS[0] |
| mega_login(acc['email'], acc['password']) |
| |
| |
| target_dir = DOWNLOAD_DIR / task_id |
| target_dir.mkdir(exist_ok=True) |
| |
| |
| max_retries = 10 |
| retry_count = 0 |
| success = False |
| |
| while retry_count < max_retries: |
| send_webhook(webhook_url, {'status': 'downloading', 'progress': 10 + (retry_count * 2)}) |
| |
| |
| result = run_megacmd(['mega-get', mega_url, str(target_dir)], timeout=MEGACMD_TIMEOUT) |
| |
| |
| if result.get('critical'): |
| send_webhook(webhook_url, {'status': 'error', 'error': 'IP Ban - Restarting System'}) |
| phoenix_restart() |
| return |
| |
| |
| if result.get('bandwidth_exceeded'): |
| logger.warning("📊 Kota doldu, hesap değiştiriliyor...") |
| if rotate_account(): |
| retry_count += 1 |
| time.sleep(3) |
| continue |
| else: |
| raise Exception("Hesap rotasyonu başarısız, hesap kalmadı!") |
| |
| |
| if result['success']: |
| success = True |
| break |
| |
| |
| logger.error(f"İndirme hatası: {result.get('error')}") |
| break |
| |
| if not success: |
| raise Exception("İndirme tamamlanamadı.") |
|
|
| |
| logger.info("📦 Zipleniyor...") |
| send_webhook(webhook_url, {'status': 'zipping', 'progress': 80}) |
| |
| zip_path = DOWNLOAD_DIR / f"{task_id}.zip" |
| with zipfile.ZipFile(zip_path, 'w', zipfile.ZIP_DEFLATED) as zipf: |
| for root, dirs, files in os.walk(target_dir): |
| for file in files: |
| file_path = os.path.join(root, file) |
| arcname = os.path.relpath(file_path, target_dir) |
| zipf.write(file_path, arcname) |
| |
| |
| |
| |
| upload_url = "https://dosya-yuklendi-varsayildi.com" |
| |
| send_webhook(webhook_url, { |
| 'status': 'completed', |
| 'progress': 100, |
| 'download_url': upload_url |
| }) |
| logger.info(f"✅ İşlem Tamam: {task_id}") |
|
|
| |
| shutil.rmtree(target_dir, ignore_errors=True) |
| if os.path.exists(zip_path): |
| os.remove(zip_path) |
|
|
| except Exception as e: |
| logger.error(f"🔥 Kritik Hata: {str(e)}") |
| send_webhook(webhook_url, {'status': 'failed', 'error': str(e)}) |
| finally: |
| with task_lock: |
| current_task = None |
|
|
| |
|
|
| @app.route('/health', methods=['GET']) |
| def health(): |
| """Render bu adrese ping atar, worker yaşıyor mu bakar.""" |
| with task_lock: |
| is_busy = current_task is not None |
| return jsonify({'status': 'healthy', 'busy': is_busy}), 200 |
|
|
| @app.route('/process', methods=['POST']) |
| def process(): |
| """Render'dan gelen iş emrini alır.""" |
| try: |
| with task_lock: |
| if current_task is not None: |
| return jsonify({'error': 'Worker şu an meşgul'}), 503 |
| |
| data = request.get_json() |
| task_id = f"task_{int(time.time())}" |
| |
| task_data = { |
| 'task_id': task_id, |
| 'mega_url': data['mega_url'], |
| 'webhook_url': data.get('webhook_url') |
| } |
| |
| with task_lock: |
| current_task = task_data |
| |
| |
| t = threading.Thread(target=process_task, args=(task_data,)) |
| t.daemon = True |
| t.start() |
| |
| return jsonify({'success': True, 'task_id': task_id}), 200 |
| |
| except Exception as e: |
| return jsonify({'error': str(e)}), 500 |
|
|
| if __name__ == '__main__': |
| |
| if not init_megacmd_server(): |
| exit(1) |
| |
| |
| app.run(host='0.0.0.0', port=7860) |