From 01bd139af78911090cc725609724ade58b64ecf2 Mon Sep 17 00:00:00 2001 From: drjones Date: Thu, 13 Aug 2026 16:41:45 -0700 Subject: [PATCH] android-fleet-godhead v1: fleet control node + ProxyFly engine + dashboard --- .gitignore | 3 + android-fleet.service | 15 + android-harvest.service | 15 + app.py | 807 ++++++++++++++++++++++++++++++++++++++++ config.json.example | 10 + docs/MASTER_PROMPT.md | 92 +++++ docs/README.md | 73 ++++ onboard.py | 313 ++++++++++++++++ proxy_harvest.py | 215 +++++++++++ static/index.html | 457 +++++++++++++++++++++++ 10 files changed, 2000 insertions(+) create mode 100644 .gitignore create mode 100644 android-fleet.service create mode 100644 android-harvest.service create mode 100644 app.py create mode 100644 config.json.example create mode 100644 docs/MASTER_PROMPT.md create mode 100644 docs/README.md create mode 100644 onboard.py create mode 100644 proxy_harvest.py create mode 100644 static/index.html diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..4eb85ae --- /dev/null +++ b/.gitignore @@ -0,0 +1,3 @@ +config.json +__pycache__/ +*.pyc diff --git a/android-fleet.service b/android-fleet.service new file mode 100644 index 0000000..696c293 --- /dev/null +++ b/android-fleet.service @@ -0,0 +1,15 @@ +[Unit] +Description=Android Fleet Godhead — fleet control node +After=network.target + +[Service] +Type=simple +WorkingDirectory=/opt/android-fleet +ExecStart=/usr/bin/python3 /opt/android-fleet/app.py +Restart=always +RestartSec=5 +StandardOutput=journal +StandardError=journal + +[Install] +WantedBy=multi-user.target diff --git a/android-harvest.service b/android-harvest.service new file mode 100644 index 0000000..6659a13 --- /dev/null +++ b/android-harvest.service @@ -0,0 +1,15 @@ +[Unit] +Description=Android Fleet — ProxyFly harvest engine +After=network.target + +[Service] +Type=simple +WorkingDirectory=/opt/android-fleet +ExecStart=/usr/bin/python3 /opt/android-fleet/proxy_harvest.py +Restart=always +RestartSec=30 +StandardOutput=journal +StandardError=journal + +[Install] +WantedBy=multi-user.target diff --git a/app.py b/app.py new file mode 100644 index 0000000..abecb23 --- /dev/null +++ b/app.py @@ -0,0 +1,807 @@ +#!/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 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) + +# ---------------- 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")]: + 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") + 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 + +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: + 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)) + con.execute("UPDATE proxies SET health=?, latency_ms=?, last_check=? WHERE id=?", + (health, lat, time.time(), 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 = [] + try: + for _ in range(count): + vmid = None + 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')}") + pve("PUT", f"/nodes/pve/qemu/{vmid}/config", {"args": args, "memory": 3072, "cores": 2}) + pve("POST", f"/nodes/pve/qemu/{vmid}/status/start", {}) + vmids.append(vmid) + log("SPAWN", f"spawned {name} (vmid {vmid}, display {display})") + 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=? WHERE id=?", (json.dumps(vmids), job_id)) + con.commit() + con.close() + log("SPAWN_DONE", f"job {job_id}: {len(vmids)} vms up") + except Exception as e: + con = db() + con.execute("UPDATE spawn_jobs SET status='error', error=? WHERE id=?", (str(e), job_id)) + con.commit() + con.close() + log("SPAWN_ERROR", f"job {job_id}: {e}") + +@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") + +@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() + +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): + 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) diff --git a/config.json.example b/config.json.example new file mode 100644 index 0000000..8781fd8 --- /dev/null +++ b/config.json.example @@ -0,0 +1,10 @@ +{ + "proxmox_host": "10.30.20.85", + "proxmox_token": "root@pam!androidcloud=REPLACE_ME", + "template_vmid": 1301, + "lan_cidr": "10.30.20.0/24", + "adb_port": 5555, + "next_display": 62, + "ssh_host": "10.30.20.85", + "api_token": "REPLACE_ME" +} \ No newline at end of file diff --git a/docs/MASTER_PROMPT.md b/docs/MASTER_PROMPT.md new file mode 100644 index 0000000..2a4c9b5 --- /dev/null +++ b/docs/MASTER_PROMPT.md @@ -0,0 +1,92 @@ +# MASTER PROMPT — Android Fleet Godhead (complete rebuild spec) + +You are an expert systems architect, network engineer, and full-stack developer. +Rebuild the ANDROID FLEET GODHEAD: an autonomous Android VM fleet controller for Proxmox. + +## REAL INFRASTRUCTURE (do not invent alternatives) + +- Proxmox host 10.30.20.85 (SSH root, key auth; API https://10.30.20.85:8006/api2/json, token `PVEAPIToken root@pam!androidcloud=d324351f-5c70-45a1-855f-3552bd290bc1`, verify_ssl=False). Node name: pve. +- SINGLE bridge `vmbr0` = LAN 10.30.20.0/24, Xfinity router at 10.30.20.1 handles DHCP (dynamic pool .100-.254). THERE IS NO vmbr1, NO dnsmasq — do NOT build one. Android VMs DHCP on the existing LAN. +- Android template: VM 1301, Android-x86 9.0-r2, 4 cores / 6GB / 12GB qcow2 (storage `local`, content types iso,vztmpl,backup,rootdir,images). + Boot args (per-clone, unique -vnc display XX): + `-kernel /var/lib/vz/template/kernel -initrd /var/lib/vz/template/initrd.img -append "quiet root=/dev/ram0 androidboot.hardware=android_x86_64 SRC=/android-9.0-r2 androidboot.selinux=permissive SETUPWIZARD=0 ip=dhcp" -vnc 0.0.0.0:59XX,password=off` +- Template is backdoored: `ro.adb.secure=0` + `service.adb.tcp.port=5555` in build.prop AND an /system/etc/init.sh hook that re-issues setprop + restarts adbd at boot. system.sfs inside the ext2 disk at /android-9.0-r2/system.sfs — when rebuilding it: mksquashfs MUST use `-comp gzip -b 131072` (xz/256K makes this kernel hang at "Detecting Android-x86"). The sfs wraps a single `system.img` (ext4) — edit build.prop/init.sh by loop-mounting it. +- Each clone gets a random MAC -> random DHCP IP -> adb open on 5555 (insecure, root). +- Control node: LXC CT 704 `android-fleet-godhead` @ 10.30.20.31:8080, Debian 12, Python 3.11, Flask + flask-sock + SQLite + adb + requests + pysocks. systemd unit `android-fleet.service`, WorkingDirectory /opt/android-fleet. SSH key from CT704 → Proxmox authorized_keys (for `ip neigh` ARP). +- Proxy pool: local NordVPN SOCKS5 containers running gost `-L socks5://0.0.0.0:1080` — CT680 nord-tunnel 10.30.20.154, CT681 nord-london 10.30.20.71, CT682 nord-sydney 10.30.20.189. Each device gets ONE proxy → unique egress IP. Assign via `adb shell settings put global http_proxy host:port` + `settings put global global_http_proxy host:port`. +- AI automation: Hermes agents on local GPU Ollama ONLY — light_reaper 10.30.20.186:11434 (granite4.1:3b agentic) or shadow-death 10.30.20.128:11434 (gemma4:26b heavy). NEVER MacBook models, NEVER run LLMs inside Proxmox CTs. LLM calls: POST /api/chat with {"model", "keep_alive": "30m", "messages", "stream": false, "options": {"num_predict": ...}} — big models need num_predict >= 2048 (thinking tokens are uncapped), never instruct them to "hide reasoning" (infinite loop). + +## ARCHITECTURE + +``` +Proxmox (10.30.20.85) ──clone/start/destroy──▶ CT704 android-fleet-godhead (10.30.20.31:8080) + │ + vmbr0 LAN 10.30.20.0/24 ──────────────────────┤ + │ + android clones (VMIDs 1310+, DHCP IPs, adb :5555) ◀─── scan + adb + metrics + nord CTs (680/681/682 gost :1080) ◀─── proxy health + per-device assignment +``` + +## COMPONENTS + +### 1. Discovery (every 10s) +- Threaded port scan of 10.30.20.0/24 for open :5555 (socket timeout 1.0s, 64-256 threads). +- For each open IP: `adb connect ip:5555` then `adb shell cat /sys/class/net/eth0/address` = MAC (primary). Fallback: `ssh root@10.30.20.85 ip neigh show dev vmbr0` (host ARP only shows IPs the HOST has talked to — never rely on it alone). +- vmid resolution: match MAC against Proxmox `net0` configs of qemu VMs. +- Register new MAC → DEVICE_JOINED event; mark offline when port closes. + +### 2. ADB orchestrator (every 15s) +- `adb connect` retry; 3 consecutive failures → offline (proxy released after 5 min). +- Collect: ro.product.model, ro.build.version.release, /proc/meminfo, /proc/stat (cpu delta), dumpsys battery, /proc/net/dev eth0 rx/tx deltas. + +### 3. Proxy switchboard (every 60s) +- Health: GET https://api.ipify.org through each proxy (requests + socks5:// via PySocks), 8s timeout → healthy/dead + latency + public egress IP. +- Auto-assign healthy unassigned proxy to any online device without one; push via `settings put global http_proxy`. + +### 4. Spawn / kill +- POST /api/spawn {count, name_prefix} → background job: free vmid scan (1310+), POST /nodes/pve/qemu/1301/clone {newid, name, full:1}, PUT config {args: unique 59XX display}, POST status/start. 10 max per job, ~4s stagger. +- Kill: qm stop → DELETE VM → release proxy → drop device row. + +### 5. Dashboard + WS +- GET / → static/index.html (Tailwind CDN + Chart.js CDN, dark godhead theme). +- WS /ws pushes: {type:"devices"} every 3s, {type:"metrics"} every 5s, {type:"log"} on events. + +## REST CONTRACT + +| Method | Path | Body/Notes | +|---|---|---| +| GET | /api/status | {devices_online, devices_total, proxies_healthy, proxies_total, uptime} | +| GET/PATCH | /api/devices | PATCH {mac, label} | +| GET/POST | /api/proxies | POST {name, host, port, proto=socks5, user, pass} | +| DELETE | /api/proxies/ | releases assignment | +| POST | /api/proxies//assign | {mac} | +| POST | /api/spawn | {count, name_prefix} → 202 {job_id} | +| GET | /api/spawn/jobs | last 10 jobs | +| POST | /api/device//kill | destroy VM + release proxy | +| POST | /api/device//shell | {cmd} — WHITELIST regex only (getprop, dumpsys, cat /proc, pm list, settings get/list, ls, df, whoami, logcat -d, ip addr, wm size, id, uname) | +| POST | /api/device//screenshot | returns PNG | +| GET | /api/device//screenshot/latest | cached PNG | +| POST | /api/device//reboot | adb reboot | +| GET | /api/logs | last 200 events | + +## DB SCHEMA (SQLite /opt/android-fleet/fleet.db, WAL) +- devices(mac PK, vmid, ip, model, android, battery, cpu_pct, mem_used_mb, mem_total_mb, adb_state, proxy_name, public_ip, last_seen, first_seen, label) +- proxies(id PK, name, host, port, proto, user, pass, health, latency_ms, assigned_mac, last_check) +- spawn_jobs(id PK, status, count, done, vmid_list JSON, error, created) + +## HARD LESSONS (baked into this spec) +1. system.sfs rebuild: gzip + 128K blocks ONLY (xz hangs boot). +2. PVE API status/current may return {"data": null} — never .get() on it blindly. +3. requests needs pysocks for socks5:// proxies. +4. Host ARP table is unreliable for discovery — get MAC from the device over adb. +5. Scan timeout must be >= 1s for emulated NICs. +6. adb over IPv6 link-local needs zone %eth0 and may hang — use IPv4 whenever present. +7. qm clone full copies ~12GB → free vmid scan must probe existence safely. +8. Template base images are immutable (chattr +i) — update a template by cloning from a fixed running VM, never by editing the base. + +## ACCEPTANCE TESTS +1. Spawn 1 → device appears in /api/devices with ip/mac/vmid within 5 min, adb_state online. +2. Proxy auto-assigns → device.proxy_name set, public_ip ≠ LAN IP. +3. Screenshot endpoint returns a valid PNG. +4. Kill → VM destroyed in Proxmox, proxy released, row gone. +5. Dashboard shows device card with live cpu/mem bars + log events, no page refresh. diff --git a/docs/README.md b/docs/README.md new file mode 100644 index 0000000..346f93e --- /dev/null +++ b/docs/README.md @@ -0,0 +1,73 @@ +# Android Fleet Godhead — Operator Manual + +Dashboard: http://10.30.20.31:8080 · CT704 · systemd: android-fleet + +## What it is +Master control node for the Android VM fleet. Clones Android-x86 VMs from template 1301 on demand. Each clone gets a random MAC → random DHCP IP from the LAN router → auto-registers via adb (:5555) → gets a unique egress proxy from the NordVPN bank. Full live telemetry in the dashboard. + +## Quick commands + +```bash +# spawn 3 fleet members +curl -X POST http://10.30.20.31:8080/api/spawn -H "Content-Type: application/json" -d '{"count":3,"name_prefix":"fleet"}' + +# watch the job +curl -s http://10.30.20.31:8080/api/spawn/jobs + +# list devices +curl -s http://10.30.20.31:8080/api/devices + +# shell (whitelisted cmds) +curl -X POST http://10.30.20.31:8080/api/device//shell -H "Content-Type: application/json" -d '{"cmd":"getprop ro.build.fingerprint"}' + +# screenshot +curl -X POST http://10.30.20.31:8080/api/device//screenshot -o shot.png + +# reboot / kill (kill = destroy VM + release proxy) +curl -X POST http://10.30.20.31:8080/api/device//reboot +curl -X POST http://10.30.20.31:8080/api/device//kill + +# add proxy +curl -X POST http://10.30.20.31:8080/api/proxies -H "Content-Type: application/json" -d '{"name":"nord-tokyo","host":"10.30.20.154","port":1080,"proto":"socks5"}' + +# assign proxy to device +curl -X POST http://10.30.20.31:8080/api/proxies/4/assign -H "Content-Type: application/json" -d '{"mac":"bc:24:11:cb:6a:f7"}' + +# direct adb (from CT704) +pct exec 704 -- adb connect 10.30.20.195:5555 +pct exec 704 -- adb -s 10.30.20.195:5555 shell "pm list packages" +``` + +## How pieces connect +``` +Proxmox (clone/start/destroy via PVE API token) + ↕ +CT704 godhead (Flask :8080 + SQLite + adb + ws) + ├─ scan 10.30.20.0/24:5555 every 10s ─▶ new devices auto-register + ├─ adb metrics every 15s (cpu/mem/battery/net) + ├─ proxy health every 60s → auto-assign → settings put global http_proxy on device + └─ WS push → dashboard live cards + charts +``` + +## How AI agents drive it +- Agents = Hermes cron jobs on local GPU Ollama (light_reaper granite4.1:3b / shadow-death gemma4:26b). NEVER MacBook models, never LLMs on CTs. +- Pattern: agent reads /api/devices → picks device(s) → calls shell/screenshot/spawn endpoints → acts (e.g. install APK, run app, tap via `input` — extend the whitelist regex in app.py for new cmds). +- MCP: add a remote HTTP MCP in Hermes config (url http://10.30.20.31:8080/mcp) wrapping the same endpoints when needed. + +## Troubleshooting +| Symptom | Fix | +|---|---| +| Device offline | `pct exec 704 -- adb connect :5555` manually; check VM status `qm status `; reboot via VM | +| Never appears | DHCP lease? check router; ensure adb port open: `nc -z -w2 5555` | +| Spawn job error | `curl /api/spawn/jobs` — free-vmid scan 1310+; check Proxmox storage space (local dir, full clone = 12GB each) | +| Proxy dead | `pct exec 680 -- gost` is the process; check nordvpn: `pct exec 680 -- nordvpn status`; then health loop auto-recovers | +| No egress IP shown | proxy assigned but check failed — verify from CT704: `curl -x socks5h://10.30.20.71:1080 https://api.ipify.org` | +| Boot hang at "Detecting Android-x86" | system.sfs rebuilt with wrong compression — MUST be gzip/128K blocks | + +## Template surgery (when updating the image) +1. Pick a known-good running clone. 2. `qm stop `. 3. qemu-nbd map its disk (NOT the base image — templates are immutable). 4. Extract /android-9.0-r2/system.sfs → unsquashfs → loop-mount system.img → edit build.prop / etc/init.sh → umount → `mksquashfs -comp gzip -b 131072`. 5. Write back, unmount, disconnect. 6. `qm destroy 1301 --purge` → `qm clone 1301 --full` → set args (ip=dhcp, -vnc 5961) → `qm template 1301`. + +## Backups / state +- /opt/android-fleet/{fleet.db, config.json} on CT704 — copy both before any surgery. +- Template base: /var/lib/vz/images/1301/base-1301-disk-1.qcow2 (immutable, don't hand-edit). +- Docs: Obsidian Hermes/android-fleet-godhead.md · fleet db mirrors to Teable if wired. diff --git a/onboard.py b/onboard.py new file mode 100644 index 0000000..b11e37a --- /dev/null +++ b/onboard.py @@ -0,0 +1,313 @@ +#!/usr/bin/env python3 +""" +Android fleet onboarding: turn a bare Android-x86 clone into a fully-set-up +phone with a fresh Instagram account. +Steps: basics -> Instagram install -> account creation (device proxy) -> +IMAP email verification -> app login (best-effort UI) -> DB record. +""" +import imaplib +import json +import os +import re +import sqlite3 +import subprocess +import sys +import time + +DB_PATH = "/opt/android-fleet/fleet.db" +IG_APK = "/opt/android-fleet/ig_instagram.apk" +SESSIONS_DIR = "/opt/android-fleet/sessions" +GMAIL = "jones.henry666q@gmail.com" +APP_PASS = "hkdbecbdkncmvshg" +PASSWORD = "aaaa1111" + +os.makedirs(SESSIONS_DIR, exist_ok=True) + + +def log(msg): + print(f"[onboard] {msg}", flush=True) + + +def dbcon(): + con = sqlite3.connect(DB_PATH, timeout=15) + con.row_factory = sqlite3.Row + return con + + +def get_device(mac): + con = dbcon() + row = con.execute("SELECT * FROM devices WHERE mac=?", (mac,)).fetchone() + con.close() + return dict(row) if row else None + + +def adb_shell(serial, cmd, timeout=45): + r = subprocess.run(["adb", "-s", serial, "shell", cmd], + capture_output=True, text=True, timeout=timeout) + return r.stdout.strip() + + +def adb_cmd(serial, args, timeout=60): + return subprocess.run(["adb", "-s", serial] + args, + capture_output=True, text=True, timeout=timeout) + + +def next_counter(): + con = dbcon() + row = con.execute("SELECT COUNT(*) AS n FROM accounts").fetchone() + con.close() + return row["n"] + 1 + + +def setup_basics(serial): + log("basics: timezone, screen, animations") + adb_shell(serial, "settings put global stay_on_while_plugged_in 3") + adb_shell(serial, "settings put system screen_off_timeout 2147483647") + adb_shell(serial, "settings put global window_animation_scale 0") + adb_shell(serial, "settings put global transition_animation_scale 0") + adb_shell(serial, "settings put global animator_duration_scale 0") + adb_shell(serial, "setprop persist.sys.timezone America/Los_Angeles") + adb_shell(serial, "settings put global airplane_mode_on 0") + log("basics done") + + +def install_instagram(serial): + """Try the APK; arm64-only builds won't install on x86 — non-fatal.""" + log("attempting Instagram APK install (arm64 builds fail on x86 — that's fine)...") + try: + r = adb_cmd(serial, ["install", "-r", IG_APK], timeout=300) + out = (r.stdout or "") + (r.stderr or "") + log(f"install result: {out.strip()[:120]}") + if "Success" in out: + return "installed" + if "NO_MATCHING_ABIS" in out: + log("APK is arm64-only — skipping app install, using API-level account") + return "skipped" + return "failed" + except Exception as e: + log(f"install error: {e}") + return "failed" + + +def fetch_verification(timeout=240): + """Poll Gmail IMAP for IG confirm: 6-digit code OR confirmation link.""" + end = time.time() + timeout + while time.time() < end: + try: + M = imaplib.IMAP4_SSL("imap.gmail.com") + M.login(GMAIL, APP_PASS) + M.select("INBOX") + typ, data = M.search(None, '(FROM "Instagram")') + nums = data[0].split()[-6:] + for num in reversed(nums): + typ, msg = M.fetch(num, "(BODY.PEEK[])") + raw = msg[0][1] + body = raw.decode(errors="replace") + m = re.search(r"\b(\d{6})\b", body) + if m: + M.logout() + return {"code": m.group(1)} + links = re.findall(r'https?://[^\s"<>]+(?:confirm|verify|challenge)[^\s"<>]*', body) + if links: + M.logout() + return {"link": links[0]} + M.logout() + except Exception as e: + log(f"imap poll: {e}") + time.sleep(8) + return None + + +def create_instagram_account(serial, proxy_url, counter): + log(f"creating IG account #{counter} via device proxy") + from instagrapi import Client + cl = Client() + if proxy_url: + cl.set_proxy(proxy_url) + email = f"igbot{counter}+{GMAIL}" + username = f"fleet.jones{counter}" + status = "created" + try: + cl.signup(username=username, password=PASSWORD, email=email, full_name="Android User") + log("signup ok, no verification needed") + except Exception as e: + msg = str(e).lower() + log(f"signup challenge: {str(e)[:100]}") + if "checkpoint" in msg or "challenge" in msg: + return {"status": "checkpoint", "email": email, "username": username} + if "sent" in msg or "verification" in msg or "confirm" in msg: + ver = fetch_verification() + if not ver: + return {"status": "no_code", "email": email, "username": username} + if "code" in ver: + try: + cl.signup(username=username, password=PASSWORD, email=email, + full_name="Android User", verification_code=ver["code"]) + log("signup ok after email code") + except Exception as e2: + log(f"verify signup failed: {str(e2)[:100]}") + return {"status": f"verify_failed: {str(e2)[:60]}", + "email": email, "username": username} + elif "link" in ver: + log(f"clicking confirm link through device proxy: {ver['link'][:60]}...") + try: + import requests as _rq + proxies = None + if proxy_url: + proxies = {"http": proxy_url, "https": proxy_url} + r = _rq.get(ver["link"], proxies=proxies, timeout=20, + headers={"User-Agent": "Mozilla/5.0 (Linux; Android 9) AppleWebKit/537.36"}) + log(f"link click: HTTP {r.status_code}") + time.sleep(5) + try: + cl.login(username, PASSWORD) + log("login after link confirm OK") + except Exception as le: + log(f"login after link: {str(le)[:80]}") + status = "created_unverified" + except Exception as ce: + log(f"link click failed: {ce}") + status = "created_unverified" + else: + return {"status": f"signup_failed: {str(e)[:60]}", "email": email, "username": username} + sess = os.path.join(SESSIONS_DIR, f"{username}.json") + cl.dump_settings(sess) + return {"status": status, "email": email, "username": username, + "session_file": sess, "password": PASSWORD} + + +def login_app(serial, username): + """Best-effort: launch IG and log in via UI taps (uiautomator-based).""" + log("launching IG app for login") + adb_shell(serial, "am start -n com.instagram.android/.activity.MainTabActivity", 15) + time.sleep(8) + def find_node(text): + dump = adb_shell(serial, "uiautomator dump /sdcard/ui.xml && cat /sdcard/ui.xml", 20) + m = re.search(r'text="[^"]*' + re.escape(text) + r'[^"]*"[^>]*bounds="\[(\d+),(\d+)\]\[(\d+),(\d+)\]"', dump) + if m: + x = (int(m.group(1)) + int(m.group(3))) // 2 + y = (int(m.group(2)) + int(m.group(4))) // 2 + return x, y + return None + try: + p = find_node("Log in") + if p: + adb_shell(serial, f"input tap {p[0]} {p[1]}", 10) + time.sleep(4) + p = find_node("Username") + if p: + adb_shell(serial, f"input tap {p[0]} {p[1]}", 10) + adb_shell(serial, f"input text {username}", 10) + time.sleep(1) + p = find_node("Password") + if p: + adb_shell(serial, f"input tap {p[0]} {p[1]}", 10) + adb_shell(serial, f"input text {PASSWORD}", 10) + time.sleep(1) + p = find_node("Log In") or find_node("Log in") + if p: + adb_shell(serial, f"input tap {p[0]} {p[1]}", 10) + log("login tapped") + return True + except Exception as e: + log(f"app login UI best-effort failed: {e}") + return False + + +def login_existing(mac, username, password): + """Login an existing IG account through the device's proxy. Session saved.""" + dev = get_device(mac) + if not dev or not dev["ip"]: + return + con = dbcon() + pr = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (mac,)).fetchone() + con.close() + from instagrapi import Client + cl = Client() + if pr: + cred = f"{pr['user']}:{pr['pass']}@" if pr["user"] else "" + cl.set_proxy(f"socks5://{cred}{pr['host']}:{pr['port']}") + try: + cl.login(username, password) + sess = os.path.join(SESSIONS_DIR, f"{username}.json") + cl.dump_settings(sess) + con = dbcon() + con.execute("""INSERT OR REPLACE INTO accounts(mac, email, username, password, status, session_file, created) + VALUES(?,?,?,?,?,?,?)""", + (mac, None, username, password, "logged_in", sess, time.time())) + con.commit() + con.close() + log(f"{mac}: login OK for {username} — session saved") + except Exception as e: + log(f"{mac}: login failed for {username}: {str(e)[:120]}") + con = dbcon() + con.execute("""INSERT OR REPLACE INTO accounts(mac, email, username, password, status, session_file, created) + VALUES(?,?,?,?,?,?,?)""", + (mac, None, username, password, f"login_failed: {str(e)[:60]}", None, time.time())) + con.commit() + con.close() + + +def onboard(mac): + dev = get_device(mac) + if not dev or not dev["ip"]: + log(f"device {mac} offline — aborting") + return + serial = f"{dev['ip']}:5555" + log(f"onboarding {mac} @ {serial}") + try: + adb_cmd(serial, ["connect", serial], timeout=15) + except Exception: + pass + counter = next_counter() + con = dbcon() + con.execute("INSERT OR REPLACE INTO accounts(mac, status, created) VALUES(?,?,?)", + (mac, "onboarding", time.time())) + con.commit() + con.close() + + setup_basics(serial) + try: + adb_cmd(serial, ["root"], timeout=15) + time.sleep(2) + adb_cmd(serial, ["connect", serial], timeout=15) + except Exception: + pass + apk_state = install_instagram(serial) + + # proxy for this device + con = dbcon() + pr = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (mac,)).fetchone() + con.close() + proxy_url = None + if pr: + cred = f"{pr['user']}:{pr['pass']}@" if pr["user"] else "" + proxy_url = f"socks5://{cred}{pr['host']}:{pr['port']}" + + result = create_instagram_account(serial, proxy_url, counter) + result["apk"] = apk_state + if result.get("session_file"): + login_app(serial, result["username"]) + _set_status(mac, result["status"], extra=result) + + +def _set_status(mac, status, extra=None): + con = dbcon() + if extra: + con.execute("""UPDATE accounts SET status=?, email=?, username=?, password=?, + session_file=?, created=? WHERE mac=?""", + (status, extra.get("email"), extra.get("username"), + extra.get("password", PASSWORD), extra.get("session_file"), + time.time(), mac)) + else: + con.execute("UPDATE accounts SET status=? WHERE mac=?", (status, mac)) + con.commit() + con.close() + log(f"account status: {status}") + + +if __name__ == "__main__": + if len(sys.argv) < 2: + print("usage: onboard.py ") + sys.exit(1) + onboard(sys.argv[1]) diff --git a/proxy_harvest.py b/proxy_harvest.py new file mode 100644 index 0000000..d88e0a7 --- /dev/null +++ b/proxy_harvest.py @@ -0,0 +1,215 @@ +#!/usr/bin/env python3 +""" +ProxyFly integration: harvest free proxies -> live-test -> cleanliness-grade +(hosting/proxy flags + Spamhaus DNSBL) -> feed the fleet proxy bank. +Every machine gets a verified egress IP, 100% backend, zero-touch. +""" +import json +import os +import socket +import sqlite3 +import threading +import time +import urllib.request +from concurrent.futures import ThreadPoolExecutor + +DB_PATH = "/opt/android-fleet/fleet.db" +SOURCES = [ + "https://cdn.jsdelivr.net/gh/proxifly/free-proxy-list@main/proxies/all/data.json", +] +MAX_FETCH = 200 # candidates per cycle +MAX_TEST = 120 # live-test at most this many per cycle +MAX_POOL = 40 # cap the proxies table (perf) +TEST_TIMEOUT = 7 +HARVEST_INTERVAL = 1500 # 25 min + +def log(msg): + print(f"[harvest] {msg}", flush=True) + +def dbcon(): + con = sqlite3.connect(DB_PATH, timeout=30) + con.row_factory = sqlite3.Row + return con + +def fetch_candidates(): + cands = [] + for src in SOURCES: + try: + with urllib.request.urlopen(src, timeout=45) as r: + data = json.loads(r.read()) + for p in data: + proto = p.get("protocol", "") + if proto not in ("socks5", "https", "http", "socks4"): + continue + if p.get("anonymity") == "transparent": + continue # leaks real IP — useless for the fleet + cands.append({ + "host": p["ip"], "port": p["port"], "proto": proto, + "country": (p.get("geolocation") or {}).get("country", "?"), + }) + except Exception as e: + log(f"source error {src}: {e}") + # dedupe by host:port + seen = set() + out = [] + for c in cands: + k = f"{c['host']}:{c['port']}" + if k not in seen: + seen.add(k) + out.append(c) + return out[:MAX_FETCH] + +def test_proxy(c): + """Returns (c, egress_ip, latency_ms) or (c, None, None).""" + proto_map = {"socks5": "socks5h://", "https": "https://", "http": "http://", + "socks4": "socks4://"} + url = proto_map[c["proto"]] + f"{c['host']}:{c['port']}" + import urllib.request as ur + opener = ur.build_opener(ur.ProxyHandler({"http": url, "https": url})) + t0 = time.time() + try: + with opener.open("https://api.ipify.org", timeout=TEST_TIMEOUT) as r: + ip = r.read().decode().strip() + lat = int((time.time() - t0) * 1000) + return (c, ip, lat) + except Exception: + return (c, None, None) + +def check_clean(ips): + """ip-api batch: flags for each ip. Returns {ip: dict}.""" + out = {} + ips = [i for i in ips if i] + for i in range(0, len(ips), 15): + batch = ips[i:i + 15] + try: + req = urllib.request.Request( + "http://ip-api.com/batch?fields=status,country,isp,as,hosting,proxy,mobile,query", + data=json.dumps(batch).encode(), + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=20) as r: + res = json.loads(r.read()) + for e in res: + if e.get("status") == "success": + out[e["query"]] = e + except Exception as e: + log(f"ip-api batch error: {e}") + time.sleep(4) # 45 req/min limit — throttle hard + return out + +def dnsbl_listed(ip): + """Spamhaus ZEN: True if listed.""" + try: + rev = ".".join(reversed(ip.split("."))) + ".zen.spamhaus.org" + socket.gethostbyname(rev) + return True + except socket.gaierror: + return False + except Exception: + return False + +def prune_pool(): + con = dbcon() + rows = con.execute("SELECT id, egress_ip, source FROM proxies ORDER BY id").fetchall() + # remove duplicates by egress ip (keep lowest id) + seen = set() + for r in rows: + if r["egress_ip"] and r["egress_ip"] in seen: + con.execute("DELETE FROM proxies WHERE id=?", (r["id"],)) + elif r["egress_ip"]: + seen.add(r["egress_ip"]) + # cap pool size (drop oldest proxifly entries beyond cap) + over = [r for r in rows if r["source"] == "proxifly"] + if len(over) > MAX_POOL: + for r in over[:-MAX_POOL]: + con.execute("DELETE FROM proxies WHERE id=?", (r["id"],)) + con.commit() + con.close() + +def harvest(): + log(f"cycle start — fetching {len(SOURCES)} source(s)") + cands = fetch_candidates() + log(f"{len(cands)} candidates, testing up to {MAX_TEST}") + tested = [] + with ThreadPoolExecutor(max_workers=20) as ex: + for c, ip, lat in ex.map(test_proxy, cands[:MAX_TEST]): + if ip: + c["egress_ip"] = ip + c["latency_ms"] = lat + tested.append(c) + log(f"{len(tested)} proxies alive; grading cleanliness") + flags = check_clean([c["egress_ip"] for c in tested]) + dns_flags = {} + with ThreadPoolExecutor(max_workers=20) as ex: + for c in tested: + dns_flags[c["egress_ip"]] = ex.submit(dnsbl_listed, c["egress_ip"]) + for k, f in dns_flags.items(): + dns_flags[k] = f.result() + con = dbcon() + added = 0 + for c in tested: + f = flags.get(c["egress_ip"], {}) + dns = dns_flags.get(c["egress_ip"], False) + clean = 1 if (not f.get("hosting") and not f.get("proxy") and not dns) else 0 + mobile = 1 if f.get("mobile") else 0 + # keep only clean ones + a few fallbacks, skip dirty datacenter spam + if clean == 0: + why = "hosting" if f.get("hosting") else ("proxy" if f.get("proxy") else ("dnsbl" if dns else "unknown")) + log(f"rejected {c['egress_ip']} ({why})") + continue + name = f"proxifly-{c['country']}-{c['egress_ip'].split('.')[-2]}" + con.execute("""INSERT INTO proxies(name, host, port, proto, health, latency_ms, + assigned_mac, last_check, clean, country, egress_ip, source) + VALUES(?,?,?,?,'healthy',?,NULL,?,?,?,?,?) + ON CONFLICT(egress_ip) DO UPDATE SET health='healthy', latency_ms=?, + last_check=?, clean=?""", + (name, c["host"], c["port"], c["proto"], c["latency_ms"], time.time(), + clean, c["country"], c["egress_ip"], "proxifly", + c["latency_ms"], time.time(), clean)) + added += 1 + con.commit() + con.close() + prune_pool() + log(f"cycle done: +{added} clean proxies into the bank") + assign_best() + log("auto-assignment sweep complete") + +def assign_best(): + """Give every online device without a proxy the best clean one available.""" + con = dbcon() + con.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_prox_egress ON proxies(egress_ip)") + devs = con.execute("SELECT * FROM devices WHERE adb_state='online' AND ip IS NOT NULL").fetchall() + for d in devs: + have = con.execute("SELECT * FROM proxies WHERE assigned_mac=?", (d["mac"],)).fetchone() + if have: + continue # already assigned — keep stable + best = con.execute("""SELECT * FROM proxies + WHERE clean=1 AND health='healthy' AND assigned_mac IS NULL + ORDER BY latency_ms ASC LIMIT 1""").fetchone() + if not best: + best = con.execute("""SELECT * FROM proxies + WHERE health='healthy' AND assigned_mac IS NULL + ORDER BY latency_ms ASC LIMIT 1""").fetchone() + if not best: + continue + con.execute("UPDATE proxies SET assigned_mac=? WHERE id=?", (d["mac"], best["id"])) + con.execute("UPDATE devices SET proxy_name=?, public_ip=? WHERE mac=?", + (best["name"], best["egress_ip"], d["mac"])) + log(f"assigned {best['name']} (egress {best['egress_ip']}, clean={best['clean']}) -> {d['mac'][:14]}") + con.commit() + con.close() + +def loop(): + while True: + time.sleep(HARVEST_INTERVAL) + try: + harvest() + except Exception as e: + log(f"cycle error: {e}") + +if __name__ == "__main__": + log("proxyfly harvest engine starting") + try: + harvest() + except Exception as e: + log(f"startup harvest error: {e}") + loop() diff --git a/static/index.html b/static/index.html new file mode 100644 index 0000000..1088400 --- /dev/null +++ b/static/index.html @@ -0,0 +1,457 @@ + + + + + +ANDROID FLEET — GODHEAD + + + + + + + + +
+ + +
+
+

ANDROID FLEET — GODHEAD

+

autonomous android spawner · proxy switchboard · live telemetry

+
+
+ --:--:-- + + WS +
+
+ + +
+ + +
+ SPAWN + + + +
+
discovery: LAN :5555 scan + ARP · every 10s
+
+ + +
+
FLEET BANDWIDTH (KB/s)
+
CPU % / MEM % (fleet avg)
+
ONLINE DEVICES
+
PROXY LATENCY (ms)
+
+ + +
+
+
📱
+
NO DEVICES — awaiting spawn
+
spawn a fleet member and it will auto-register here within seconds
+
+ + +
+
+
+ PROXY BANK + +
+
+
+
+ PROXYFLY ENGINE +
auto-harvest · test · clean-check (ip-api + Spamhaus) · assign — every 45 min, zero touch
+
+
+ IG ACCOUNTS +
+
+
+ EVENT LOG +
+
+
+
+ + + + + + + + + +