diff --git a/android-mcp.service b/android-mcp.service new file mode 100644 index 0000000..26b16dc --- /dev/null +++ b/android-mcp.service @@ -0,0 +1,13 @@ +[Unit] +Description=Android Fleet MCP server (Hermes agent control) +After=network.target android-fleet.service + +[Service] +Type=simple +WorkingDirectory=/opt/android-fleet +ExecStart=/usr/bin/python3 /opt/android-fleet/mcp_server.py +Restart=always +RestartSec=5 + +[Install] +WantedBy=multi-user.target diff --git a/app.py b/app.py index ca5b7a9..a378305 100644 --- a/app.py +++ b/app.py @@ -6,6 +6,7 @@ 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 @@ -46,6 +47,11 @@ 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() @@ -315,6 +321,10 @@ def find_vmid_for_mac(mac): 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: @@ -328,6 +338,14 @@ def adb_loop(): 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: @@ -874,13 +892,168 @@ def ws_handler(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): + 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) diff --git a/mcp_server.py b/mcp_server.py new file mode 100644 index 0000000..c479008 --- /dev/null +++ b/mcp_server.py @@ -0,0 +1,136 @@ +#!/usr/bin/env python3 +""" +Android Fleet MCP server — lets Hermes agents (and any MCP client) drive the +fleet: see screens, tap, type, run commands, trigger onboarding, broadcast. +HTTP streamable-http on :8095/mcp. All calls proxy to the godhead REST API. +""" +import json +import urllib.request + +from mcp.server.fastmcp import FastMCP + +GODHEAD = "http://127.0.0.1:8080" + +mcp = FastMCP("android-fleet", host="0.0.0.0", port=8095) + +def _api(path, method="GET", payload=None, timeout=120): + data = json.dumps(payload).encode() if payload is not None else None + req = urllib.request.Request(GODHEAD + path, data=data, method=method, + headers={"Content-Type": "application/json"}) + with urllib.request.urlopen(req, timeout=timeout) as r: + return json.loads(r.read()) + +def _err(e): + return {"error": str(e)[:200]} + +@mcp.tool() +def fleet_devices() -> dict: + """List all fleet devices with status, IP, proxy and IG account state.""" + try: + devs = _api("/api/devices") + accts = {a["mac"]: a.get("status") for a in _api("/api/accounts")} + for d in devs: + d["ig_status"] = accts.get(d["mac"], "none") + return {"count": len(devs), + "online": sum(1 for d in devs if d["adb_state"] == "online"), + "devices": devs} + except Exception as e: + return _err(e) + +@mcp.tool() +def device_shell(mac: str, cmd: str) -> dict: + """Run an adb shell command on one device (mac = device MAC address).""" + try: + return _api(f"/api/device/{mac}/exec", "POST", {"cmd": cmd}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_vision(mac: str, instruction: str) -> dict: + """Screenshot the device and ask the vision model (qwen2.5vl:7b on GPU) + about what's on screen. Use before tapping to find coordinates.""" + try: + return _api(f"/api/device/{mac}/vision", "POST", {"instruction": instruction}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_tap(mac: str, x: int, y: int) -> dict: + """Tap screen coordinates on a device.""" + try: + return _api(f"/api/device/{mac}/tap", "POST", {"x": x, "y": y}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_swipe(mac: str, x1: int, y1: int, x2: int, y2: int, ms: int = 300) -> dict: + """Swipe on a device (e.g. scroll the feed).""" + try: + return _api(f"/api/device/{mac}/swipe", "POST", + {"x1": x1, "y1": y1, "x2": x2, "y2": y2, "ms": ms}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_type(mac: str, text: str) -> dict: + """Type text on a device (into the focused field).""" + try: + return _api(f"/api/device/{mac}/text", "POST", {"text": text}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_launch(mac: str, package: str) -> dict: + """Launch an app by package name on a device.""" + try: + return _api(f"/api/device/{mac}/launch", "POST", {"package": package}) + except Exception as e: + return _err(e) + +@mcp.tool() +def device_screenshot(mac: str) -> dict: + """Take a screenshot; returns a vision description + the PNG is cached + at /api/device//screenshot/latest.""" + try: + shot = _api(f"/api/device/{mac}/screenshot", "POST", timeout=60) + return {"saved": shot.get("path"), "view": f"/api/device/{mac}/screenshot/latest"} + except Exception as e: + return _err(e) + +@mcp.tool() +def fleet_broadcast_tap(macs: list[str], x: int, y: int) -> dict: + """Tap the same coordinates on many devices at once.""" + try: + return _api("/api/fleet/action", "POST", + {"macs": macs, "action": {"kind": "tap", "params": {"x": x, "y": y}}}) + except Exception as e: + return _err(e) + +@mcp.tool() +def onboard_device(mac: str) -> dict: + """Run the full phone+IG onboarding on a device (creates IG account when + Instagram allows it; rate limits are reported honestly).""" + try: + return _api(f"/api/device/{mac}/onboard", "POST", {}) + except Exception as e: + return _err(e) + +@mcp.tool() +def spawn_devices(count: int) -> dict: + """Clone new Android VMs from the template; they auto-register.""" + try: + return _api("/api/spawn", "POST", {"count": count}) + except Exception as e: + return _err(e) + +@mcp.tool() +def fleet_message(text: str) -> dict: + """Broadcast a message via saved IG sessions (requires created accounts).""" + try: + return _api("/api/fleet/message", "POST", {"text": text}) + except Exception as e: + return _err(e) + +if __name__ == "__main__": + import uvicorn + uvicorn.run(mcp.streamable_http_app(), host="0.0.0.0", port=8095, log_level="warning") diff --git a/onboard.py b/onboard.py index b11e37a..708484d 100644 --- a/onboard.py +++ b/onboard.py @@ -12,6 +12,7 @@ import re import sqlite3 import subprocess import sys +import threading import time DB_PATH = "/opt/android-fleet/fleet.db" @@ -52,11 +53,21 @@ def adb_cmd(serial, args, timeout=60): capture_output=True, text=True, timeout=timeout) +NORD_FALLBACK = "http://10.30.20.154:1081" # nord-tunnel HTTP CONNECT (instagrapi's socks client is broken) +COUNTER_FILE = "/opt/android-fleet/ig_counter" +COUNTER_LOCK = threading.Lock() + def next_counter(): - con = dbcon() - row = con.execute("SELECT COUNT(*) AS n FROM accounts").fetchone() - con.close() - return row["n"] + 1 + """Atomic account counter — never reuse an email alias.""" + with COUNTER_LOCK: + n = 1 + if os.path.exists(COUNTER_FILE): + try: + n = int(open(COUNTER_FILE).read().strip() or 0) + 1 + except ValueError: + n = 1 + open(COUNTER_FILE, "w").write(str(n)) + return n def setup_basics(serial): @@ -133,6 +144,9 @@ def create_instagram_account(serial, proxy_url, counter): except Exception as e: msg = str(e).lower() log(f"signup challenge: {str(e)[:100]}") + if "429" in str(e) or "rate" in msg: + return {"status": "rate_limited_ip", "email": email, "username": username, + "session_file": None} 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: @@ -275,16 +289,15 @@ def onboard(mac): 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) + # SIGNUP ALWAYS USES THE STABLE NORD TUNNEL — free-pool device proxies are + # too flaky for the account-creation handshake (dead SOCKS = instant fail). + result = None + try: + result = create_instagram_account(serial, NORD_FALLBACK, counter) + except Exception as e: + log(f"onboard crash: {str(e)[:100]}") + result = {"status": f"onboard_error: {str(e)[:60]}", "email": None, + "username": None, "session_file": None} result["apk"] = apk_state if result.get("session_file"): login_app(serial, result["username"]) diff --git a/static/index.html b/static/index.html index f2c7264..9d7c627 100644 --- a/static/index.html +++ b/static/index.html @@ -239,8 +239,8 @@ function renderFleet() {
+ - @@ -248,6 +248,10 @@ function renderFleet() {
+
+ +
▶ full screen
+
`; }).join(""); } @@ -459,6 +463,12 @@ function humanize(sec) { document.getElementById("modal").addEventListener("click", (e) => { if (e.target.id === "modal") closeModal(); }); setInterval(() => { document.getElementById("clock").textContent = new Date().toLocaleTimeString(); }, 1000); setInterval(loadProxies, 15000); +// keep every card's mini-screen fresh (~4s; backend caches a new shot every 8s) +setInterval(() => { + document.querySelectorAll("img.shot").forEach(img => { + img.src = "/api/device/" + img.dataset.mac + "/screenshot/latest?t=" + Date.now(); + }); +}, 4000); initCharts(); connect(); loadProxies();