281 lines
9.1 KiB
Python
281 lines
9.1 KiB
Python
|
|
"""
|
||
|
|
AUTARCH Local Agent
|
||
|
|
Lightweight client that runs on the user's machine and executes
|
||
|
|
commands locally, streaming results back to the AUTARCH server.
|
||
|
|
|
||
|
|
This makes the web UI feel like AUTARCH is running on the user's
|
||
|
|
machine — scans, tool output, file operations all happen locally.
|
||
|
|
|
||
|
|
Usage:
|
||
|
|
python3 autarch-agent.py --server https://10.0.0.1:8181
|
||
|
|
|
||
|
|
The agent:
|
||
|
|
1. Connects to the AUTARCH server via WebSocket
|
||
|
|
2. Authenticates with the user's token
|
||
|
|
3. Receives command requests from the server
|
||
|
|
4. Executes them locally (subprocess)
|
||
|
|
5. Streams stdout/stderr back in real-time
|
||
|
|
6. Reports system info (OS, arch, hostname, IP) on connect
|
||
|
|
|
||
|
|
Protocol:
|
||
|
|
Server -> Agent: {"type": "exec", "id": "xxx", "command": "nmap -sV 10.0.0.1"}
|
||
|
|
Agent -> Server: {"type": "output", "id": "xxx", "stream": "stdout", "data": "..."}
|
||
|
|
Agent -> Server: {"type": "output", "id": "xxx", "stream": "stderr", "data": "..."}
|
||
|
|
Agent -> Server: {"type": "done", "id": "xxx", "exit_code": 0}
|
||
|
|
Agent -> Server: {"type": "sysinfo", "hostname": "...", "os": "...", ...}
|
||
|
|
"""
|
||
|
|
|
||
|
|
import argparse
|
||
|
|
import asyncio
|
||
|
|
import json
|
||
|
|
import os
|
||
|
|
import platform
|
||
|
|
import shutil
|
||
|
|
import socket
|
||
|
|
import subprocess
|
||
|
|
import ssl
|
||
|
|
import sys
|
||
|
|
import time
|
||
|
|
from pathlib import Path
|
||
|
|
|
||
|
|
try:
|
||
|
|
import websockets
|
||
|
|
except ImportError:
|
||
|
|
print("Installing websockets...")
|
||
|
|
subprocess.check_call([sys.executable, "-m", "pip", "install", "websockets", "-q"])
|
||
|
|
import websockets
|
||
|
|
|
||
|
|
|
||
|
|
def get_system_info():
|
||
|
|
"""Collect local system information to send on connect."""
|
||
|
|
info = {
|
||
|
|
"hostname": socket.gethostname(),
|
||
|
|
"os": platform.system(),
|
||
|
|
"os_version": platform.version(),
|
||
|
|
"arch": platform.machine(),
|
||
|
|
"python": platform.python_version(),
|
||
|
|
"username": os.getenv("USER", os.getenv("USERNAME", "unknown")),
|
||
|
|
"home": str(Path.home()),
|
||
|
|
"cwd": os.getcwd(),
|
||
|
|
"pid": os.getpid(),
|
||
|
|
}
|
||
|
|
|
||
|
|
# Network interfaces
|
||
|
|
try:
|
||
|
|
info["ip"] = socket.gethostbyname(socket.gethostname())
|
||
|
|
except Exception:
|
||
|
|
info["ip"] = "unknown"
|
||
|
|
|
||
|
|
# Available tools
|
||
|
|
tools = {}
|
||
|
|
for tool in ["nmap", "tcpdump", "wireshark", "python3", "git", "curl",
|
||
|
|
"wget", "ssh", "netstat", "ss", "ip", "ifconfig", "arp",
|
||
|
|
"dig", "nslookup", "traceroute", "ping", "whois",
|
||
|
|
"openssl", "nikto", "sqlmap", "hydra", "john", "hashcat"]:
|
||
|
|
path = shutil.which(tool)
|
||
|
|
if path:
|
||
|
|
tools[tool] = path
|
||
|
|
info["tools"] = tools
|
||
|
|
|
||
|
|
# Disk usage
|
||
|
|
try:
|
||
|
|
usage = shutil.disk_usage("/")
|
||
|
|
info["disk_total"] = usage.total
|
||
|
|
info["disk_free"] = usage.free
|
||
|
|
except Exception:
|
||
|
|
pass
|
||
|
|
|
||
|
|
return info
|
||
|
|
|
||
|
|
|
||
|
|
async def execute_command(command, cmd_id, websocket):
|
||
|
|
"""Execute a command locally and stream output back."""
|
||
|
|
try:
|
||
|
|
process = await asyncio.create_subprocess_shell(
|
||
|
|
command,
|
||
|
|
stdout=asyncio.subprocess.PIPE,
|
||
|
|
stderr=asyncio.subprocess.PIPE,
|
||
|
|
)
|
||
|
|
|
||
|
|
# Stream stdout
|
||
|
|
async def read_stream(stream, name):
|
||
|
|
while True:
|
||
|
|
line = await stream.readline()
|
||
|
|
if not line:
|
||
|
|
break
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "output",
|
||
|
|
"id": cmd_id,
|
||
|
|
"stream": name,
|
||
|
|
"data": line.decode("utf-8", errors="replace"),
|
||
|
|
}))
|
||
|
|
|
||
|
|
# Run both streams concurrently
|
||
|
|
await asyncio.gather(
|
||
|
|
read_stream(process.stdout, "stdout"),
|
||
|
|
read_stream(process.stderr, "stderr"),
|
||
|
|
)
|
||
|
|
|
||
|
|
exit_code = await process.wait()
|
||
|
|
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "done",
|
||
|
|
"id": cmd_id,
|
||
|
|
"exit_code": exit_code,
|
||
|
|
}))
|
||
|
|
|
||
|
|
except Exception as e:
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "error",
|
||
|
|
"id": cmd_id,
|
||
|
|
"message": str(e),
|
||
|
|
}))
|
||
|
|
|
||
|
|
|
||
|
|
async def read_file(path, cmd_id, websocket):
|
||
|
|
"""Read a file and send its contents back."""
|
||
|
|
try:
|
||
|
|
p = Path(path).expanduser()
|
||
|
|
if not p.exists():
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "error", "id": cmd_id,
|
||
|
|
"message": f"File not found: {path}"
|
||
|
|
}))
|
||
|
|
return
|
||
|
|
|
||
|
|
if p.is_dir():
|
||
|
|
# List directory
|
||
|
|
entries = []
|
||
|
|
for entry in sorted(p.iterdir()):
|
||
|
|
stat = entry.stat()
|
||
|
|
entries.append({
|
||
|
|
"name": entry.name,
|
||
|
|
"is_dir": entry.is_dir(),
|
||
|
|
"size": stat.st_size,
|
||
|
|
"modified": stat.st_mtime,
|
||
|
|
})
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "dirlist", "id": cmd_id,
|
||
|
|
"path": str(p), "entries": entries,
|
||
|
|
}))
|
||
|
|
else:
|
||
|
|
# Read file (limit to 10MB)
|
||
|
|
size = p.stat().st_size
|
||
|
|
if size > 10 * 1024 * 1024:
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "error", "id": cmd_id,
|
||
|
|
"message": f"File too large: {size} bytes"
|
||
|
|
}))
|
||
|
|
return
|
||
|
|
|
||
|
|
content = p.read_text(encoding="utf-8", errors="replace")
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "file", "id": cmd_id,
|
||
|
|
"path": str(p), "content": content,
|
||
|
|
"size": size,
|
||
|
|
}))
|
||
|
|
|
||
|
|
except Exception as e:
|
||
|
|
await websocket.send(json.dumps({
|
||
|
|
"type": "error", "id": cmd_id,
|
||
|
|
"message": str(e),
|
||
|
|
}))
|
||
|
|
|
||
|
|
|
||
|
|
async def agent_loop(server_url, token):
|
||
|
|
"""Main agent loop - connect to server and process commands."""
|
||
|
|
# Allow self-signed certs
|
||
|
|
ssl_context = ssl.SSLContext(ssl.PROTOCOL_TLS_CLIENT)
|
||
|
|
ssl_context.check_hostname = False
|
||
|
|
ssl_context.verify_mode = ssl.CERT_NONE
|
||
|
|
|
||
|
|
ws_url = server_url.replace("https://", "wss://").replace("http://", "ws://")
|
||
|
|
ws_url = f"{ws_url}/ws/agent"
|
||
|
|
|
||
|
|
print(f"Connecting to {ws_url}...")
|
||
|
|
|
||
|
|
while True:
|
||
|
|
try:
|
||
|
|
async with websockets.connect(ws_url, ssl=ssl_context,
|
||
|
|
extra_headers={"Authorization": f"Bearer {token}"}) as ws:
|
||
|
|
print("Connected to AUTARCH server")
|
||
|
|
|
||
|
|
# Send system info on connect
|
||
|
|
sysinfo = get_system_info()
|
||
|
|
sysinfo["type"] = "sysinfo"
|
||
|
|
await ws.send(json.dumps(sysinfo))
|
||
|
|
|
||
|
|
# Send heartbeat periodically
|
||
|
|
async def heartbeat():
|
||
|
|
while True:
|
||
|
|
await asyncio.sleep(30)
|
||
|
|
try:
|
||
|
|
await ws.send(json.dumps({"type": "heartbeat", "time": time.time()}))
|
||
|
|
except Exception:
|
||
|
|
break
|
||
|
|
|
||
|
|
heartbeat_task = asyncio.create_task(heartbeat())
|
||
|
|
|
||
|
|
# Process commands
|
||
|
|
try:
|
||
|
|
async for message in ws:
|
||
|
|
try:
|
||
|
|
msg = json.loads(message)
|
||
|
|
msg_type = msg.get("type", "")
|
||
|
|
cmd_id = msg.get("id", "")
|
||
|
|
|
||
|
|
if msg_type == "exec":
|
||
|
|
command = msg.get("command", "")
|
||
|
|
if command:
|
||
|
|
print(f" Executing: {command[:80]}")
|
||
|
|
asyncio.create_task(execute_command(command, cmd_id, ws))
|
||
|
|
|
||
|
|
elif msg_type == "read":
|
||
|
|
path = msg.get("path", "")
|
||
|
|
if path:
|
||
|
|
asyncio.create_task(read_file(path, cmd_id, ws))
|
||
|
|
|
||
|
|
elif msg_type == "ping":
|
||
|
|
await ws.send(json.dumps({"type": "pong", "id": cmd_id}))
|
||
|
|
|
||
|
|
elif msg_type == "kill":
|
||
|
|
# Server wants us to disconnect
|
||
|
|
print("Server requested disconnect")
|
||
|
|
break
|
||
|
|
|
||
|
|
except json.JSONDecodeError:
|
||
|
|
pass
|
||
|
|
|
||
|
|
finally:
|
||
|
|
heartbeat_task.cancel()
|
||
|
|
|
||
|
|
except (websockets.exceptions.ConnectionClosed,
|
||
|
|
ConnectionRefusedError, OSError) as e:
|
||
|
|
print(f"Connection lost: {e}. Reconnecting in 5s...")
|
||
|
|
await asyncio.sleep(5)
|
||
|
|
|
||
|
|
except KeyboardInterrupt:
|
||
|
|
print("\nAgent stopped")
|
||
|
|
break
|
||
|
|
|
||
|
|
|
||
|
|
def main():
|
||
|
|
parser = argparse.ArgumentParser(description="AUTARCH Local Agent")
|
||
|
|
parser.add_argument("--server", "-s", required=True,
|
||
|
|
help="AUTARCH server URL (e.g. https://10.0.0.1:8181)")
|
||
|
|
parser.add_argument("--token", "-t", default="",
|
||
|
|
help="Authentication token")
|
||
|
|
args = parser.parse_args()
|
||
|
|
|
||
|
|
print("AUTARCH Local Agent")
|
||
|
|
print(f"Server: {args.server}")
|
||
|
|
print(f"System: {platform.system()} {platform.machine()}")
|
||
|
|
print(f"Host: {socket.gethostname()}")
|
||
|
|
print()
|
||
|
|
|
||
|
|
asyncio.run(agent_loop(args.server, args.token))
|
||
|
|
|
||
|
|
|
||
|
|
if __name__ == "__main__":
|
||
|
|
main()
|