Files
RMM-UI-BE/client_agent_prototype.py

229 lines
10 KiB
Python

import urllib.request
import json
import time
import subprocess
import psutil
CLIENT_ID = "client_1"
import os
# Using the central production backend URL
CENTRAL_API_URL = f"http://rmm-backend.seekright.com/api/get-command?client_id={CLIENT_ID}"
CENTRAL_TELEMETRY_URL = f"http://rmm-backend.seekright.com/api/telemetry?client_id={CLIENT_ID}"
platform_name = "nt" if os.name == "nt" else "posix"
CENTRAL_COMMANDS_URL = f"http://rmm-backend.seekright.com/api/get-commands?platform={platform_name}"
# SECURITY MECHANISM: The "Whitelist"
# The agent NEVER executes arbitrary strings from the internet.
# It dynamically pulls safe command arrays from the central FastAPI server configuration.
# An in-memory cache is used as a resilient fallback if the central server goes offline.
CACHED_COMMANDS = {}
def load_allowed_commands():
global CACHED_COMMANDS
try:
req = urllib.request.Request(CENTRAL_COMMANDS_URL)
with urllib.request.urlopen(req, timeout=3) as response:
commands = json.loads(response.read().decode())
if commands:
CACHED_COMMANDS = commands
return CACHED_COMMANDS
except Exception as e:
print(f" -> [!] Failed to load commands from central API: {e}")
if CACHED_COMMANDS:
print(" -> Returning last successfully cached whitelist.")
return CACHED_COMMANDS
# Return fallback safe command in case the server is offline on startup
return {"whoami": ["whoami"]}
ROCKET_CHAT_WEBHOOK = "https://text.seekright.com/hooks/6a0eb41ea88367bc6e101f0c/gvXfD8zSNRti43kEabqdE9tMCvNi9zAAoGfbao7cxcKbiW6R"
def send_rocket_chat_notification(title, message, color="#764FA5"):
try:
if not message:
message = "No output provided."
payload = {
"alias": f"Agent: {CLIENT_ID}",
"text": "Execution Report",
"attachments": [{
"title": title,
"text": message,
"color": color
}]
}
data = json.dumps(payload).encode('utf-8')
req = urllib.request.Request(ROCKET_CHAT_WEBHOOK, data=data, headers={'Content-Type': 'application/json'}, method='POST')
urllib.request.urlopen(req, timeout=3)
except Exception as e:
print(f" -> [!] Failed to send webhook: {e}")
def collect_and_send_telemetry():
try:
# 1. Gather comprehensive GPU stats
gpus = []
try:
# Query multiple specific fields (including detailed power and fan speed metrics)
result = subprocess.run(
["nvidia-smi", "--query-gpu=name,utilization.gpu,temperature.gpu,memory.used,memory.total,power.draw,power.limit,fan.speed", "--format=csv,noheader"],
capture_output=True, text=True
)
if result.returncode == 0:
for line in result.stdout.strip().split('\n'):
if line:
parts = [p.strip() for p in line.split(',')]
if len(parts) >= 8:
gpus.append({
"name": parts[0],
"utilization": parts[1],
"temp": parts[2] + "C",
"memory_used": parts[3],
"memory_total": parts[4],
"power_draw": parts[5],
"power_limit": parts[6],
"fan_speed": parts[7]
})
except FileNotFoundError:
pass # nvidia-smi not installed, ignore
# 2. Gather ALL Disk Partitions (Filtering snaps, loops, and virtual mounts)
disks = []
for part in psutil.disk_partitions(all=False):
# Ignore loop devices, snaps, system, and WSL virtual mountpoints
if 'loop' in part.device or 'loop' in part.fstype or part.fstype == '':
continue
if part.mountpoint.startswith('/snap') or part.fstype == 'squashfs':
continue
if part.mountpoint.startswith('/boot') or part.mountpoint.startswith('/mnt/wsl'):
continue
try:
usage = psutil.disk_usage(part.mountpoint)
disks.append({
"mount": part.mountpoint,
"total_gb": round(usage.total / (1024**3), 2),
"free_gb": round(usage.free / (1024**3), 2),
"percent": usage.percent
})
except PermissionError:
continue
# 3. Assemble Payload
payload = {
"cpu_percent": psutil.cpu_percent(interval=0.5),
"memory_percent": psutil.virtual_memory().percent,
"memory_total_gb": round(psutil.virtual_memory().total / (1024**3), 2),
"memory_free_gb": round(psutil.virtual_memory().available / (1024**3), 2),
"disks": disks,
"gpus": gpus
}
data = json.dumps(payload).encode('utf-8')
req = urllib.request.Request(CENTRAL_TELEMETRY_URL, data=data, headers={'Content-Type': 'application/json'}, method='POST')
urllib.request.urlopen(req, timeout=3)
print("[*] Telemetry pushed successfully.")
except Exception as e:
print(f"[!] Failed to push telemetry: {e}")
def send_command_result_to_server(command: str, returncode: int, stdout: str, stderr: str):
try:
base_url = CENTRAL_API_URL.split("/api/")[0]
url = f"{base_url}/api/command-result?client_id={CLIENT_ID}"
payload = {
"command": command,
"returncode": returncode,
"stdout": stdout,
"stderr": stderr
}
data = json.dumps(payload).encode('utf-8')
req = urllib.request.Request(
url,
data=data,
headers={'Content-Type': 'application/json'},
method='POST'
)
with urllib.request.urlopen(req, timeout=3) as response:
pass
print(" -> Sent command execution output to Central Server.")
except Exception as e:
print(f" -> [!] Failed to send command result to server: {e}")
def poll_server():
try:
print(f"[*] Checking for commands...")
# 1. Fetch data from Central API
req = urllib.request.Request(CENTRAL_API_URL)
with urllib.request.urlopen(req, timeout=3) as response:
data = json.loads(response.read().decode())
command_key = data.get("command", "none")
# 2. Process the command
if command_key == "none":
print(" -> No pending commands.")
else:
# Load whitelisted commands dynamically from commands.json
allowed_commands = load_allowed_commands()
if command_key in allowed_commands:
print(f" -> Executing approved command: '{command_key}'")
safe_cmd_list = allowed_commands[command_key]
# Special check: If rebooting or restarting agent, send pre-execution notification
# to ensure the network buffer has time to transmit the status before connection drops!
is_shutdown_or_restart = any(x in command_key.lower() for x in ["reboot", "restart", "kill_python"])
if is_shutdown_or_restart:
print(" -> Shutdown/restart command detected. Sending pre-execution success notification...")
send_rocket_chat_notification(
f"🔄 Executing: {command_key}",
f"System/agent is performing a scheduled action: {' '.join(safe_cmd_list)}\nConnection may drop temporarily.",
"#FFA500"
)
time.sleep(3) # Give 3 seconds for the HTTP payload to hit the Rocket.Chat server!
try:
# Execute securely. capture_output prevents it from hanging.
result = subprocess.run(safe_cmd_list, capture_output=True, text=True)
output = result.stdout.strip()
error_output = result.stderr.strip() if result.stderr else ""
# Post execution results to Central Server
send_command_result_to_server(command_key, result.returncode, output, error_output)
if not is_shutdown_or_restart:
if result.returncode == 0:
print(f" -> Output: {output}")
send_rocket_chat_notification(f"✅ Success: {command_key}", output, "#00FF00")
else:
print(f" -> Error: {error_output}")
send_rocket_chat_notification(f"❌ Failed: {command_key}", error_output, "#FF0000")
except Exception as ex:
print(f" -> [!] Execution Error: {ex}")
send_command_result_to_server(command_key, -1, "", str(ex))
if not is_shutdown_or_restart:
send_rocket_chat_notification(f"❌ Execution Error: {command_key}", str(ex), "#FF0000")
else:
# If a hacker injects "rm -rf /" into the database, it hits this block and fails securely.
warning_msg = f"Server requested unknown/malicious command: '{command_key}'. Ignored."
print(f" -> [!] SECURITY WARNING: {warning_msg}")
send_rocket_chat_notification("⚠️ SECURITY WARNING", warning_msg, "#FFA500")
except Exception as e:
print(f" -> [!] Failed to connect to central server: {e}")
if __name__ == "__main__":
print(f"Starting Agent for {CLIENT_ID}...")
print("[*] Running in continuous background polling mode (every 60 seconds).")
print("[*] Press Ctrl+C to stop the agent.")
try:
while True:
collect_and_send_telemetry()
poll_server()
time.sleep(10) # Poll every 10 seconds
except KeyboardInterrupt:
print("\n[*] Agent stopped by user.")