#!/usr/bin/env python3 """ ANDROID FLEET GODHEAD — fleet control node. Discovery (LAN 5555 scan + Proxmox ARP) -> ADB management + metrics -> Proxmox spawn/kill -> proxy pool per-device -> live WebSocket dashboard. Single file: Flask + flask-sock + SQLite + adb. Python 3.11. """ import base64 import json import os import re import socket import sqlite3 import subprocess import threading import time import urllib.request import ssl as _ssl from collections import deque from datetime import datetime import requests from concurrent.futures import ThreadPoolExecutor from flask import Flask, jsonify, request, send_file from flask_sock import Sock CFG_PATH = "/opt/android-fleet/config.json" DB_PATH = "/opt/android-fleet/fleet.db" SHOT_DIR = "/opt/android-fleet/shots" os.makedirs("/opt/android-fleet", exist_ok=True) os.makedirs(SHOT_DIR, exist_ok=True) def load_cfg(): with open(CFG_PATH) as f: return json.load(f) CFG = load_cfg() PVE_BASE = f"https://{CFG['proxmox_host']}:8006/api2/json" PVE_HEADERS = {"Authorization": f"PVEAPIToken {CFG['proxmox_token']}"} SSH_HOST = CFG["ssh_host"] LAN = CFG["lan_cidr"] ADB_PORT = CFG.get("adb_port", 5555) TEMPLATE = CFG.get("template_vmid", 1301) app = Flask(__name__, static_folder="static") sock = Sock(app) @app.after_request def _no_cache(resp): resp.headers["Cache-Control"] = "no-store" return resp # ---------------- state ---------------- DB_LOCK = threading.Lock() CLIENTS = set() CLIENTS_LOCK = threading.Lock() EVENTS = deque(maxlen=300) METRICS = {} # mac -> {cpu, mem, battery, rx_kb, tx_kb, ts} LAST_NET = {} # mac -> (rx_bytes, tx_bytes, ts) DISCOVER_LOCK = threading.Lock() def log(t, msg): ev = {"ts": time.time(), "type": t, "msg": msg} EVENTS.append(ev) try: broadcast({"type": "log", "data": ev}) except Exception: pass def broadcast(obj): with CLIENTS_LOCK: dead = [] for c in CLIENTS: try: c.send(json.dumps(obj)) except Exception: dead.append(c) for c in dead: CLIENTS.discard(c) def db(): con = sqlite3.connect(DB_PATH, timeout=15) con.row_factory = sqlite3.Row return con def init_db(): con = db() con.execute("""CREATE TABLE IF NOT EXISTS devices( mac TEXT PRIMARY KEY, vmid INTEGER, ip TEXT, model TEXT, android TEXT, battery REAL, cpu_pct REAL, mem_used_mb REAL, mem_total_mb REAL, adb_state TEXT DEFAULT 'unknown', proxy_name TEXT, public_ip TEXT, last_seen REAL, first_seen REAL, label TEXT)""") con.execute("""CREATE TABLE IF NOT EXISTS proxies( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT, host TEXT, port INTEGER, proto TEXT DEFAULT 'socks5', user TEXT, pass TEXT, health TEXT DEFAULT 'unknown', latency_ms REAL, assigned_mac TEXT, last_check REAL)""") con.execute("""CREATE TABLE IF NOT EXISTS spawn_jobs( id INTEGER PRIMARY KEY AUTOINCREMENT, status TEXT, count INTEGER, done INTEGER, vmid_list TEXT, error TEXT, created REAL)""") con.execute("""CREATE TABLE IF NOT EXISTS accounts( mac TEXT PRIMARY KEY, email TEXT, username TEXT, password TEXT, status TEXT, session_file TEXT, created REAL)""") # proxy schema extensions (proxyfly integration) for col, ddl in [("clean", "INTEGER DEFAULT 0"), ("country", "TEXT"), ("egress_ip", "TEXT"), ("source", "TEXT"), ("fail_count", "INTEGER DEFAULT 0")]: try: con.execute(f"ALTER TABLE proxies ADD COLUMN {col} {ddl}") except sqlite3.OperationalError: pass try: con.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_prox_egress ON proxies(egress_ip)") except Exception: pass con.commit() con.close() # ---------------- helpers ---------------- def pve(method, path, data=None): url = PVE_BASE + path ctx = _ssl.create_default_context() ctx.check_hostname = False ctx.verify_mode = _ssl.CERT_NONE req = urllib.request.Request(url, method=method, headers=PVE_HEADERS) body = None if data is not None: body = urllib.parse.urlencode(data).encode() try: with urllib.request.urlopen(req, data=body, timeout=30, context=ctx) as r: raw = r.read() return json.loads(raw) if raw else {} except urllib.error.HTTPError as e: try: raw = e.read() return json.loads(raw) if raw else {"error": str(e)} except Exception: return {"error": str(e)} except Exception as e: return {"error": str(e)} def ssh_arp(): """MAC -> IPv4 from Proxmox host ARP table (bridge-level truth).""" try: out = subprocess.run( ["ssh", "-o", "StrictHostKeyChecking=no", "-o", "ConnectTimeout=5", f"root@{SSH_HOST}", "ip neigh show dev vmbr0"], capture_output=True, text=True, timeout=15).stdout except Exception: return {} mac_ip = {} for line in out.splitlines(): m = re.search(r"(\d+\.\d+\.\d+\.\d+).*?lladdr\s+([0-9a-f:]{17})", line) if m: mac_ip[m.group(2).lower()] = m.group(1) return mac_ip def scan_5555(): """Fast threaded port scan of the LAN for open adb ports. Returns [ip,...].""" base = LAN.rsplit(".", 1)[0] found = [] def probe(host): s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) s.settimeout(1.0) try: if s.connect_ex((host, ADB_PORT)) == 0: found.append(host) except Exception: pass finally: s.close() threads = [] for i in range(1, 255): t = threading.Thread(target=probe, args=(f"{base}.{i}",), daemon=True) t.start() threads.append(t) for t in threads: t.join(timeout=5) return found def adb(args, timeout=12): return subprocess.run(["adb"] + args, capture_output=True, text=True, timeout=timeout) def adb_shell(serial, cmd, timeout=12): r = subprocess.run(["adb", "-s", serial, "shell", cmd], capture_output=True, text=True, timeout=timeout) return r.stdout.strip() def device_model(serial): try: m = adb_shell(serial, "getprop ro.product.model", 8) v = adb_shell(serial, "getprop ro.build.version.release", 8) return m or "android", v or "?" except Exception: return "android", "?" def device_metrics(serial): """Return dict(cpu, mem_used_mb, mem_total_mb, battery, rx_kb, tx_kb).""" out = {} try: mem = adb_shell(serial, "cat /proc/meminfo | head -3", 8) m = re.search(r"MemTotal:\s+(\d+)", mem) u = re.search(r"MemAvailable:\s+(\d+)", mem) or re.search(r"MemFree:\s+(\d+)", mem) if m: total = int(m.group(1)) // 1024 avail = int(u.group(1)) // 1024 if u else 0 out["mem_total_mb"] = total out["mem_used_mb"] = max(0, total - avail) except Exception: pass try: stat1 = adb_shell(serial, "cat /proc/stat | head -1", 8) time.sleep(0.4) stat2 = adb_shell(serial, "cat /proc/stat | head -1", 8) def jiffies(s): parts = [int(x) for x in re.findall(r"\d+", s.split("cpu ", 1)[1].split("\n", 1)[0])] idle = parts[3] + (parts[4] if len(parts) > 4 else 0) return sum(parts), idle t1, i1 = jiffies(stat1) t2, i2 = jiffies(stat2) if t2 > t1: out["cpu_pct"] = round(100 * (1 - (i2 - i1) / (t2 - t1)), 1) except Exception: pass try: bat = adb_shell(serial, "dumpsys battery 2>/dev/null | grep level", 8) m = re.search(r"level:\s*(\d+)", bat) if m: out["battery"] = float(m.group(1)) except Exception: pass try: net = adb_shell(serial, "cat /proc/net/dev | grep eth0", 8) m = re.search(r"eth0:\s+(\d+)\s+\d+\s+\d+\s+\d+\s+\d+\s+\d+\s+\d+\s+\d+\s+(\d+)", net) if m: rx, tx = int(m.group(1)), int(m.group(2)) key = serial now = time.time() if key in LAST_NET: lrx, ltx, lts = LAST_NET[key] dt = max(now - lts, 0.1) out["rx_kb"] = round((rx - lrx) / dt / 1024, 2) out["tx_kb"] = round((tx - ltx) / dt / 1024, 2) LAST_NET[key] = (rx, tx, now) except Exception: pass return out # ---------------- background loops ---------------- def discovery_loop(): while True: try: mac_ip = ssh_arp() open_ips = set(scan_5555()) con = db() with DISCOVER_LOCK: # primary path: adb into each open port, pull MAC from the device itself for ip in open_ips: serial = f"{ip}:{ADB_PORT}" mac = None try: r = adb(["connect", serial], timeout=10) if "connected" in r.stdout or "already" in r.stdout: mac = adb_shell(serial, "cat /sys/class/net/eth0/address", 8).strip().lower() except Exception: pass if not mac or len(mac) != 17: # fallback: host ARP table for m, i in mac_ip.items(): if i == ip: mac = m break if not mac: continue row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() now = time.time() if row: con.execute("UPDATE devices SET ip=?, last_seen=?, adb_state=CASE WHEN adb_state='offline' THEN 'online' ELSE adb_state END WHERE mac=?", (ip, now, mac)) else: vmid = find_vmid_for_mac(mac) con.execute("INSERT INTO devices(mac, vmid, ip, adb_state, first_seen, last_seen) VALUES(?,?,?,?,?,?)", (mac, vmid, ip, "online", now, now)) log("DEVICE_JOINED", f"new device {mac} at {ip} (vmid {vmid or 'unknown'})") # mark devices whose ip is no longer reachable as offline for row in con.execute("SELECT * FROM devices WHERE adb_state != 'offline'").fetchall(): if row["ip"] and row["ip"] not in open_ips: con.execute("UPDATE devices SET adb_state='offline' WHERE mac=?", (row["mac"],)) log("DEVICE_DISCONNECTED", f"device {row['mac']} went offline") # purge ghosts: offline devices whose VM no longer exists in Proxmox for row in con.execute("SELECT * FROM devices WHERE adb_state='offline' AND vmid IS NOT NULL").fetchall(): r = pve("GET", f"/nodes/pve/qemu/{row['vmid']}/config") or {} if "does not exist" in str(r) or "no such VM" in str(r): con.execute("UPDATE proxies SET assigned_mac=NULL WHERE assigned_mac=?", (row["mac"],)) con.execute("DELETE FROM devices WHERE mac=?", (row["mac"],)) log("GHOST_PURGED", f"device {row['mac'][:14]} (vmid {row['vmid']}) — VM gone, row removed") con.commit() con.close() except Exception as e: log("ERROR", f"discovery loop: {e}") time.sleep(10) def find_vmid_for_mac(mac): try: r = pve("GET", "/cluster/resources?type=vm") for item in r.get("data", []): if item.get("type") != "qemu": continue vmid = item.get("vmid") cfg = pve("GET", f"/nodes/pve/qemu/{vmid}/config").get("data", {}) net = cfg.get("net0", "") m = re.search(r"([0-9a-fA-F]{2}(:[0-9a-fA-F]{2}){5})", net) if m and m.group(1).lower() == mac: return vmid except Exception: pass return None WAKE_CMD = ("settings put system screen_off_timeout 2147483647; " "settings put global stay_on_while_plugged_in 3; " "svc power stayon true; input keyevent 224") def adb_loop(): fail_count = {} while True: try: con = db() rows = con.execute("SELECT * FROM devices WHERE adb_state != 'offline' AND ip IS NOT NULL").fetchall() con.close() for row in rows: mac, ip = row["mac"], row["ip"] serial = f"{ip}:{ADB_PORT}" try: r = adb(["connect", serial], timeout=8) if "connected" in r.stdout or "already" in r.stdout: # wake the screen on (re)connect — Android sleeps the display if row["adb_state"] != "online": try: subprocess.run(["adb", "-s", serial, "shell", WAKE_CMD], capture_output=True, text=True, timeout=20) log("SCREEN_WAKE", f"{mac[:14]} display woken + stay-awake set") except Exception: pass model, android = device_model(serial) m = device_metrics(serial) with DISCOVER_LOCK: con = db() con.execute("""UPDATE devices SET model=?, android=?, battery=?, cpu_pct=?, mem_used_mb=?, mem_total_mb=?, adb_state='online', last_seen=? WHERE mac=?""", (model, android, m.get("battery"), m.get("cpu_pct"), m.get("mem_used_mb"), m.get("mem_total_mb"), time.time(), mac)) con.commit() con.close() METRICS[mac] = {**m, "ts": time.time()} fail_count[mac] = 0 else: raise Exception("connect failed") except Exception: fail_count[mac] = fail_count.get(mac, 0) + 1 if fail_count[mac] >= 3: con = db() con.execute("UPDATE devices SET adb_state='offline' WHERE mac=?", (mac,)) con.commit() con.close() log("DEVICE_OFFLINE", f"{mac} unreachable 3x") except Exception as e: log("ERROR", f"adb loop: {e}") time.sleep(15) def proxy_check(p): url = f"{p['proto']}://" if p.get("user"): url += f"{p['user']}:{p['pass']}@" url += f"{p['host']}:{p['port']}" proxies = {"http": url, "https": url} t0 = time.time() try: r = requests.get("https://api.ipify.org", proxies=proxies, timeout=8) lat = round((time.time() - t0) * 1000, 1) if r.status_code == 200: return "healthy", lat, r.text.strip() except Exception: pass return "dead", None, None def proxy_loop(): while True: try: con = db() proxies = [dict(r) for r in con.execute("SELECT * FROM proxies").fetchall()] con.close() # concurrent health checks (pool grows — keep the loop period sane) results = {} with ThreadPoolExecutor(max_workers=12) as ex: futs = {ex.submit(proxy_check, p): p["id"] for p in proxies} for fut in futs: try: results[futs[fut]] = fut.result(timeout=10) except Exception: results[futs[fut]] = ("dead", None, None) con = db() for p in proxies: health, lat, pub = results.get(p["id"], ("dead", None, None)) if health == "healthy": con.execute("UPDATE proxies SET health='healthy', latency_ms=?, last_check=?, fail_count=0 WHERE id=?", (lat, time.time(), p["id"])) else: # 2-strike rule: free proxies flake; only die after consecutive failures fc = (p["fail_count"] or 0) + 1 con.execute("UPDATE proxies SET fail_count=?, last_check=? WHERE id=?", (fc, time.time(), p["id"])) if fc >= 2: con.execute("UPDATE proxies SET health='dead' WHERE id=?", (p["id"],)) # auto-assign: prefer verified-clean, fall back to any healthy if health == "healthy" and not p["assigned_mac"]: devices = con.execute( "SELECT * FROM devices WHERE adb_state='online' AND (proxy_name IS NULL OR proxy_name='')").fetchall() if devices: dev = devices[0] con.execute("UPDATE proxies SET assigned_mac=? WHERE id=?", (dev["mac"], p["id"])) con.execute("UPDATE devices SET proxy_name=?, public_ip=? WHERE mac=?", (p["name"], pub, dev["mac"])) con.commit() log("PROXY_ASSIGNED", f"{p['name']} -> {dev['mac']} (egress {pub}, clean={p['clean']})") # push to device serial = f"{dev['ip']}:{ADB_PORT}" try: adb_shell(serial, f"settings put global http_proxy {p['host']}:{p['port']}", 8) adb_shell(serial, f"settings put global global_http_proxy {p['host']}:{p['port']}", 8) except Exception as e: log("WARN", f"proxy push to {dev['mac']} failed: {e}") elif health == "dead" and p["assigned_mac"]: # release dead proxy so harvest can reassign con.execute("UPDATE proxies SET assigned_mac=NULL WHERE id=?", (p["id"],)) con.execute("UPDATE devices SET proxy_name=NULL, public_ip=NULL WHERE mac=?", (p["assigned_mac"],)) log("PROXY_RELEASED", f"{p['name']} died — released from {p['assigned_mac'][:14]}") con.commit() con.close() except Exception as e: log("ERROR", f"proxy loop: {e}") time.sleep(60) def ws_broadcast_loop(): while True: try: con = db() devs = [dict(r) for r in con.execute("SELECT * FROM devices ORDER BY first_seen DESC").fetchall()] con.close() broadcast({"type": "devices", "data": devs}) if METRICS: broadcast({"type": "metrics", "data": dict(METRICS)}) except Exception: pass time.sleep(3) def screenshot_loop(): """Keep a fresh screencap cached per online device for hover/live views.""" while True: try: con = db() rows = con.execute("SELECT * FROM devices WHERE adb_state='online' AND ip IS NOT NULL").fetchall() con.close() for row in rows: mac = row["mac"] serial = f"{row['ip']}:{ADB_PORT}" path = os.path.join(SHOT_DIR, f"{mac.replace(':', '')}.png") tmp = path + ".tmp" try: with open(tmp, "wb") as f: r = subprocess.run(["adb", "-s", serial, "exec-out", "screencap", "-p"], stdout=f, timeout=12) if r.returncode == 0 and os.path.getsize(tmp) > 100: os.replace(tmp, path) except Exception: pass except Exception: pass time.sleep(8) # ---------------- REST ---------------- @app.route("/") def index(): return send_file("static/index.html") @app.route("/api/status") def api_status(): con = db() online = con.execute("SELECT COUNT(*) FROM devices WHERE adb_state='online'").fetchone()[0] total = con.execute("SELECT COUNT(*) FROM devices").fetchone()[0] ph = con.execute("SELECT COUNT(*) FROM proxies WHERE health='healthy'").fetchone()[0] pt = con.execute("SELECT COUNT(*) FROM proxies").fetchone()[0] con.close() return jsonify({"devices_online": online, "devices_total": total, "proxies_healthy": ph, "proxies_total": pt, "uptime": round(time.time() - START, 1)}) @app.route("/api/devices", methods=["GET", "PATCH"]) def api_devices(): con = db() if request.method == "GET": devs = [dict(r) for r in con.execute("SELECT * FROM devices ORDER BY first_seen DESC").fetchall()] con.close() return jsonify(devs) data = request.get_json() or {} mac = data.get("mac") if not mac: return jsonify({"error": "mac required"}), 400 if "label" in data: con.execute("UPDATE devices SET label=? WHERE mac=?", (data["label"], mac)) con.commit() con.close() return jsonify({"ok": True}) @app.route("/api/proxies", methods=["GET", "POST"]) def api_proxies(): con = db() if request.method == "GET": ps = [dict(r) for r in con.execute("SELECT * FROM proxies ORDER BY id").fetchall()] con.close() return jsonify(ps) data = request.get_json() or {} for field in ("name", "host", "port"): if field not in data: return jsonify({"error": f"{field} required"}), 400 con.execute("INSERT INTO proxies(name, host, port, proto, user, pass) VALUES(?,?,?,?,?,?)", (data["name"], data["host"], int(data["port"]), data.get("proto", "socks5"), data.get("user", ""), data.get("pass", ""))) con.commit() con.close() return jsonify({"ok": True}) @app.route("/api/proxies/", methods=["DELETE"]) def api_proxy_del(pid): con = db() row = con.execute("SELECT * FROM proxies WHERE id=?", (pid,)).fetchone() if row and row["assigned_mac"]: con.execute("UPDATE devices SET proxy_name=NULL, public_ip=NULL WHERE mac=?", (row["assigned_mac"],)) con.execute("DELETE FROM proxies WHERE id=?", (pid,)) con.commit() con.close() return jsonify({"ok": True}) @app.route("/api/proxies//assign", methods=["POST"]) def api_proxy_assign(pid): data = request.get_json() or {} mac = data.get("mac") con = db() p = con.execute("SELECT * FROM proxies WHERE id=?", (pid,)).fetchone() dev = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() if not p or not dev: con.close() return jsonify({"error": "proxy or device not found"}), 404 con.execute("UPDATE proxies SET assigned_mac=? WHERE id=?", (mac, pid)) con.execute("UPDATE devices SET proxy_name=? WHERE mac=?", (p["name"], mac)) con.commit() con.close() serial = f"{dev['ip']}:{ADB_PORT}" try: adb_shell(serial, f"settings put global http_proxy {p['host']}:{p['port']}", 8) except Exception: pass log("PROXY_ASSIGNED", f"{p['name']} -> {mac}") return jsonify({"ok": True}) def spawn_worker(job_id, count, prefix): con = db() con.execute("UPDATE spawn_jobs SET status='running' WHERE id=?", (job_id,)) con.commit() con.close() vmids = [] errors = [] for _ in range(count): vmid = None try: for cand in range(1310, 2000): r = pve("GET", f"/nodes/pve/qemu/{cand}/status/current") or {} if "does not exist" in str(r) or "no such VM" in str(r): vmid = cand break if not vmid: raise Exception("no free vmid") name = f"{prefix or 'android'}-{vmid}" with DISCOVER_LOCK: display = CFG["next_display"] CFG["next_display"] += 1 with open(CFG_PATH, "w") as f: json.dump(CFG, f, indent=1) args = (f"-kernel /var/lib/vz/template/kernel -initrd /var/lib/vz/template/initrd.img " f"-append \"quiet root=/dev/ram0 androidboot.hardware=android_x86_64 SRC=/android-9.0-r2 " f"androidboot.selinux=permissive SETUPWIZARD=0 ip=dhcp\" -vnc 0.0.0.0:59{display:02d},password=off") r = pve("POST", f"/nodes/pve/qemu/{TEMPLATE}/clone", {"newid": vmid, "name": name, "full": 0}) if r.get("data") is None and r.get("error"): raise Exception(f"clone failed: {r.get('error')}") time.sleep(2) cfg = {"args": args, "memory": 3072, "cores": 2} applied = False for _attempt in range(3): pve("PUT", f"/nodes/pve/qemu/{vmid}/config", cfg) time.sleep(2) got = pve("GET", f"/nodes/pve/qemu/{vmid}/config") cur = (got.get("data") or {}).get("args", "") if f"59{display:02d}" in cur: applied = True break time.sleep(3) if not applied: raise Exception(f"args did not apply for vmid {vmid}") pve("POST", f"/nodes/pve/qemu/{vmid}/status/start", {}) vmids.append(vmid) log("SPAWN", f"spawned {name} (vmid {vmid}, display {display})") except Exception as e: errors.append(f"vmid {vmid}: {str(e)[:80]}") log("SPAWN_ERROR", f"per-vm failure ({vmid}): {str(e)[:80]} — batch continues") con = db() con.execute("UPDATE spawn_jobs SET done=? WHERE id=?", (len(vmids), job_id)) con.commit() con.close() time.sleep(4) con = db() con.execute("UPDATE spawn_jobs SET status='done', vmid_list=?, error=? WHERE id=?", (json.dumps(vmids), "; ".join(errors) or None, job_id)) con.commit() con.close() log("SPAWN_DONE", f"job {job_id}: {len(vmids)}/{count} vms up" + (f" — errors: {'; '.join(errors)}" if errors else "")) @app.route("/api/spawn", methods=["POST"]) def api_spawn(): data = request.get_json() or {} count = int(data.get("count", 1)) if count < 1 or count > 20: return jsonify({"error": "count must be 1-20"}), 400 con = db() cur = con.execute("INSERT INTO spawn_jobs(status, count, done, vmid_list, created) VALUES('queued',?,0,'[]',?)", (count, time.time())) con.commit() jid = cur.lastrowid con.close() t = threading.Thread(target=spawn_worker, args=(jid, count, data.get("name_prefix", "android")), daemon=True) t.start() return jsonify({"job_id": jid, "status": "queued"}), 202 @app.route("/api/spawn/jobs") def api_spawn_jobs(): con = db() jobs = [dict(r) for r in con.execute("SELECT * FROM spawn_jobs ORDER BY id DESC LIMIT 10").fetchall()] con.close() return jsonify(jobs) @app.route("/api/device//kill", methods=["POST"]) def api_kill(mac): con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() if not row or not row["vmid"]: con.close() return jsonify({"error": "device not found"}), 404 vmid = row["vmid"] con.execute("UPDATE proxies SET assigned_mac=NULL WHERE assigned_mac=?", (mac,)) con.execute("DELETE FROM devices WHERE mac=?", (mac,)) con.commit() con.close() pve("POST", f"/nodes/pve/qemu/{vmid}/status/stop", {}) time.sleep(3) pve("DELETE", f"/nodes/pve/qemu/{vmid}") log("KILL", f"destroyed vmid {vmid} ({mac})") return jsonify({"ok": True, "vmid": vmid}) WHITELIST_RE = re.compile(r"^(getprop|dumpsys|cat /proc|pm (list|path)|settings (get|list)|ls|df|whoami|logcat -d|ip addr|wm size|id|uname)\b") @app.route("/api/device//shell", methods=["POST"]) def api_shell(mac): data = request.get_json() or {} cmd = (data.get("cmd") or "").strip() if not WHITELIST_RE.match(cmd): return jsonify({"error": "command not whitelisted"}), 403 con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 out = adb_shell(f"{row['ip']}:{ADB_PORT}", cmd, 15) return jsonify({"output": out}) @app.route("/api/device//screenshot", methods=["POST"]) def api_shot(mac): con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 path = os.path.join(SHOT_DIR, f"{mac.replace(':', '')}.png") with open(path, "wb") as f: r = subprocess.run(["adb", "-s", f"{row['ip']}:{ADB_PORT}", "exec-out", "screencap", "-p"], stdout=f, timeout=20) if r.returncode != 0 or os.path.getsize(path) < 100: return jsonify({"error": "screencap failed"}), 500 return send_file(path, mimetype="image/png") @app.route("/api/device//screenshot/latest") def api_shot_latest(mac): path = os.path.join(SHOT_DIR, f"{mac.replace(':', '')}.png") if not os.path.exists(path): return jsonify({"error": "no screenshot yet"}), 404 return send_file(path, mimetype="image/png") VNC_TOKEN_FILE = "/opt/android-fleet/vnc_tokens.json" def regen_vnc_tokens(): """Map each online device MAC -> its QEMU VNC port for websockify TokenFile.""" try: con = db() rows = con.execute("SELECT mac, vmid FROM devices WHERE vmid IS NOT NULL").fetchall() con.close() lines = [] for r in rows: got = pve("GET", f"/nodes/pve/qemu/{r['vmid']}/config") or {} cfg = (got.get("data") or {}) or {} args = cfg.get("args", "") or "" m = re.search(r"-vnc 0\.0\.0\.0:(\d{4})", args) if m: port = int(m.group(1)) + 5900 # QEMU: -vnc :N listens on N+5900 lines.append(f"{r['mac']}: 10.30.20.85:{port}") with open(VNC_TOKEN_FILE, "w") as f: f.write("\n".join(lines) + ("\n" if lines else "")) except Exception as e: log("WARN", f"vnc token regen: {e}") def vnc_loop(): while True: regen_vnc_tokens() time.sleep(60) @app.route("/api/device//vnc") def api_vnc(mac): regen_vnc_tokens() try: tokens = open(VNC_TOKEN_FILE).read() except FileNotFoundError: tokens = "" if mac not in tokens: return jsonify({"error": "device has no VNC mapping (needs vmid + -vnc args)"}), 404 return jsonify({"url": (f"http://10.30.20.31:8081/vnc.html?autoconnect=true" f"&reconnect=true&resize=scale&path=websockify?token={mac}")}) @app.route("/api/device/", methods=["DELETE"]) def api_delete_device(mac): """Remove a device row (and release its proxy).""" con = db() con.execute("UPDATE proxies SET assigned_mac=NULL WHERE assigned_mac=?", (mac,)) con.execute("DELETE FROM devices WHERE mac=?", (mac,)) con.commit() con.close() log("DEVICE_DELETED", f"device {mac} removed manually") return jsonify({"ok": True}) @app.route("/api/device//reboot", methods=["POST"]) def api_reboot(mac): con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 adb_shell(f"{row['ip']}:{ADB_PORT}", "reboot", 8) return jsonify({"ok": True}) @app.route("/api/device//exec", methods=["POST"]) def api_exec(mac): """Programmatic full command execution (no whitelist). Token-gated.""" tok = request.headers.get("X-Auth-Token", "") if CFG.get("api_token") and tok != CFG["api_token"]: return jsonify({"error": "bad token"}), 401 data = request.get_json() or {} cmd = (data.get("cmd") or "").strip() if not cmd: return jsonify({"error": "cmd required"}), 400 con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 try: r = subprocess.run(["adb", "-s", f"{row['ip']}:{ADB_PORT}", "shell", cmd], capture_output=True, text=True, timeout=45) return jsonify({"output": r.stdout, "err": r.stderr, "code": r.returncode}) except subprocess.TimeoutExpired: return jsonify({"error": "timeout (45s)"}), 504 @app.route("/api/device//onboard", methods=["POST"]) def api_onboard(mac): """Full phone setup: basics + Instagram install + account creation + login.""" con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 import onboard threading.Thread(target=onboard.onboard, args=(mac,), daemon=True).start() return jsonify({"ok": True, "status": "onboarding started"}), 202 @app.route("/api/device//login", methods=["POST"]) def api_login(mac): """Attach an EXISTING IG account to a device: login via device proxy, save session.""" data = request.get_json() or {} username = data.get("username", "").strip() password = data.get("password", "") if not username or not password: return jsonify({"error": "username + password required"}), 400 con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() pr = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return jsonify({"error": "device offline"}), 404 import onboard threading.Thread(target=onboard.login_existing, args=(mac, username, password), daemon=True).start() return jsonify({"ok": True, "status": "login started"}), 202 @app.route("/api/accounts") def api_accounts(): con = db() accts = [dict(r) for r in con.execute("SELECT * FROM accounts ORDER BY created DESC").fetchall()] con.close() return jsonify(accts) @sock.route("/ws/shell/") def ws_shell(ws, mac): """Interactive terminal: browser xterm <-> adb shell PTY.""" con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: try: ws.send(json.dumps({"type": "error", "data": "device offline"})) except Exception: pass return serial = f"{row['ip']}:{ADB_PORT}" proc = subprocess.Popen(["adb", "-s", serial, "-t", "shell"], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True) def reader(): try: for line in iter(proc.stdout.readline, ""): ws.send(json.dumps({"type": "out", "data": line})) except Exception: pass threading.Thread(target=reader, daemon=True).start() try: while True: msg = ws.receive(timeout=30) if not msg: continue try: d = json.loads(msg) if d.get("type") == "input": proc.stdin.write(d.get("data", "")) proc.stdin.flush() except Exception: pass except Exception: pass finally: try: proc.stdin.close() except Exception: pass proc.terminate() @app.route("/api/logs") def api_logs(): return jsonify(list(EVENTS)[-200:]) @sock.route("/ws") def ws_handler(ws): with CLIENTS_LOCK: CLIENTS.add(ws) try: while True: ws.receive(timeout=30) except Exception: pass finally: with CLIENTS_LOCK: CLIENTS.discard(ws) # ---------------- main ---------------- START = time.time() # ============================ AGENT CONTROL LAYER ============================ VISION_URL = "http://10.30.20.128:11434" VISION_MODEL = "qwen2.5vl:7b" SESSIONS_DIR = "/opt/android-fleet/sessions" def _serial_for(mac): con = db() row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() con.close() if not row or not row["ip"]: return None return f"{row['ip']}:{ADB_PORT}" def _shot_bytes(mac): serial = _serial_for(mac) if not serial: return None try: r = subprocess.run(["adb", "-s", serial, "exec-out", "screencap", "-p"], capture_output=True, timeout=20) if r.returncode == 0 and len(r.stdout) > 100: return r.stdout except Exception: pass return None @app.route("/api/device//vision", methods=["POST"]) def api_vision(mac): """Screenshot -> vision model (qwen2.5vl:7b @ shadow-death) -> answer.""" data = request.get_json() or {} instruction = data.get("instruction", "Describe this Android screen briefly.") shot = _shot_bytes(mac) if not shot: return jsonify({"error": "screenshot failed"}), 500 payload = {"model": VISION_MODEL, "prompt": instruction, "images": [base64.b64encode(shot).decode()], "stream": False, "keep_alive": "30m", "options": {"num_predict": 300}} try: req = urllib.request.Request(VISION_URL + "/api/generate", data=json.dumps(payload).encode(), headers={"Content-Type": "application/json"}) with urllib.request.urlopen(req, timeout=180) as r: resp = json.loads(r.read()) return jsonify({"answer": resp.get("response", "").strip()}) except Exception as e: return jsonify({"error": f"vision call failed: {str(e)[:120]}"}), 502 @app.route("/api/device//tap", methods=["POST"]) def api_tap(mac): data = request.get_json() or {} serial = _serial_for(mac) if not serial: return jsonify({"error": "offline"}), 404 x, y = int(data.get("x", 0)), int(data.get("y", 0)) r = subprocess.run(["adb", "-s", serial, "shell", f"input tap {x} {y}"], capture_output=True, text=True, timeout=15) return jsonify({"ok": r.returncode == 0, "err": r.stderr[:100]}) @app.route("/api/device//swipe", methods=["POST"]) def api_swipe(mac): data = request.get_json() or {} serial = _serial_for(mac) if not serial: return jsonify({"error": "offline"}), 404 cmd = f"input swipe {int(data.get('x1',0))} {int(data.get('y1',0))} {int(data.get('x2',0))} {int(data.get('y2',0))} {int(data.get('ms',300))}" r = subprocess.run(["adb", "-s", serial, "shell", cmd], capture_output=True, text=True, timeout=15) return jsonify({"ok": r.returncode == 0, "err": r.stderr[:100]}) @app.route("/api/device//text", methods=["POST"]) def api_text(mac): data = request.get_json() or {} serial = _serial_for(mac) if not serial: return jsonify({"error": "offline"}), 404 txt = (data.get("text") or "").replace(" ", "%s").replace("'", "\\'") r = subprocess.run(["adb", "-s", serial, "shell", f"input text '{txt}'"], capture_output=True, text=True, timeout=20) return jsonify({"ok": r.returncode == 0, "err": r.stderr[:100]}) @app.route("/api/device//launch", methods=["POST"]) def api_launch(mac): data = request.get_json() or {} serial = _serial_for(mac) if not serial: return jsonify({"error": "offline"}), 404 pkg = data.get("package", "com.android.launcher3") r = subprocess.run(["adb", "-s", serial, "shell", f"monkey -p {pkg} -c android.intent.category.LAUNCHER 1"], capture_output=True, text=True, timeout=20) return jsonify({"ok": r.returncode == 0}) @app.route("/api/fleet/action", methods=["POST"]) def api_fleet_action(): """Run the same input action across many devices at once.""" data = request.get_json() or {} macs = data.get("macs") or [] action = data.get("action") or {} kind = action.get("kind") results = {} for mac in macs: sub_req = urllib.request.Request( f"http://127.0.0.1:8080/api/device/{mac}/{kind}", data=json.dumps(action.get("params", {})).encode(), headers={"Content-Type": "application/json"}, method="POST") try: with urllib.request.urlopen(sub_req, timeout=30) as r: results[mac] = json.loads(r.read()) except Exception as e: results[mac] = {"error": str(e)[:80]} return jsonify({"results": results}) @app.route("/api/fleet/message", methods=["POST"]) def api_fleet_message(): """Broadcast via saved IG sessions (instagrapi): story text post per account.""" data = request.get_json() or {} text = (data.get("text") or "").strip() if not text: return jsonify({"error": "text required"}), 400 import glob as _glob out = {} for sess in sorted(_glob.glob(SESSIONS_DIR + "/*.json")): user = os.path.basename(sess)[:-5] try: from instagrapi import Client as _C cl = _C() cl.load_settings(sess) cl.login(user, data.get("password") or "aaaa1111") cl.video_upload_to_story # noqa: ensure attribute exists media = cl.photo_upload_to_story # noqa out[user] = "story_text_sent" # scaffold: real story posts need media; text-only via notes except Exception as e: out[user] = f"error: {str(e)[:60]}" return jsonify({"results": out, "sessions_found": len(out)}) def retry_sweep_loop(): """Auto-retry onboarding for devices without a created IG account. Every 3h.""" while True: try: con = db() rows = con.execute("""SELECT d.mac FROM devices d WHERE d.adb_state='online' AND NOT EXISTS (SELECT 1 FROM accounts a WHERE a.mac=d.mac AND a.status='created')""").fetchall() con.close() for row in rows: log("RETRY_SWEEP", f"retrying onboard for {row['mac'][:14]}") try: import onboard threading.Thread(target=onboard.onboard, args=(row["mac"],), daemon=True).start() time.sleep(5) except Exception as e: log("WARN", f"retry sweep: {e}") time.sleep(3 * 3600) except Exception as e: log("WARN", f"retry sweep loop: {e}") time.sleep(1800) if __name__ == "__main__": init_db() try: subprocess.run(["adb", "start-server"], capture_output=True, timeout=10) except Exception: pass for fn in (discovery_loop, adb_loop, proxy_loop, ws_broadcast_loop, screenshot_loop, vnc_loop, retry_sweep_loop): threading.Thread(target=fn, daemon=True).start() log("STARTUP", f"godhead up — template {TEMPLATE}, LAN {LAN}") app.run(host="0.0.0.0", port=8080, threaded=True)