"""设备自动化管理后台(Flask + Flask-Login,单页应用)。 启动:python web_server.py 访问:http://localhost:5000 页面结构: / — 单页应用(监控/任务/分组/日志/用户 Tab 切换,需登录) /login /logout — 用户登录/登出 /api/* — JSON API 用户系统: 首次启动自动创建默认管理员 admin/admin123(请及时改密码)。 用户数据存 data/users.db(SQLite)。 """ import os import sys import time import uuid import atexit import secrets import functools import shlex import re import subprocess from datetime import datetime from flask import (Flask, jsonify, request, redirect, url_for, render_template, Response, make_response, session) from flask_login import (LoginManager, login_user, logout_user, login_required, current_user) from core.stf_client import STFClient from core.task_manager import TaskManager from core.device_worker import clear_worker_error, clear_all_worker_errors from core.logger import get_logger, _LOG_DIR, _MODULE_FILES from core.models import db, init_db, User, DeviceGroup, TaskJob, CustomAction from core.apk_manager import ApkManager from core.adb_helper import screenshot, list_installed_apps, adb_connect_light, _ADB_LOCK, _adb from core import uiauto_helper, tailscale_client from core.tailscale_client import TailscaleError from config import ADB_PATH, STF_SSH_TARGET, STF_DOCKER_CONTAINER from tasks import list_task_types, get_task_class _log = get_logger("web") app = Flask(__name__) # 会话密钥:优先 .env 的 WEB_SECRET_KEY;未配置则随机生成(重启后登录态失效,生产务必配置固定值) _web_secret = os.environ.get("WEB_SECRET_KEY") or secrets.token_hex(32) if not os.environ.get("WEB_SECRET_KEY"): _log.warning("未配置 WEB_SECRET_KEY,已随机生成会话密钥(重启后登录态失效)") app.config["SECRET_KEY"] = _web_secret # debug=False 时模板默认不自动重载,开发期改 HTML 需重启才生效;这里显式开启 app.config["TEMPLATES_AUTO_RELOAD"] = True app.config["SQLALCHEMY_DATABASE_URI"] = "sqlite:///" + os.path.join( os.path.dirname(os.path.abspath(__file__)), "data", "users.db") app.config["SQLALCHEMY_TRACK_MODIFICATIONS"] = False login_manager = LoginManager(app) login_manager.login_view = "login" # 先初始化数据库(含旧 JSON 迁移),再创建 TaskManager(需要 app context 读写 DB) init_db(app) stf = STFClient() mgr = TaskManager(stf, app=app) apk_mgr = ApkManager(stf, app=app) @login_manager.user_loader def load_user(user_id): return User.query.get(int(user_id)) # ================== CSRF 防护 ================== def _csrf_token(): """获取或生成当前会话的 CSRF token(非 GET 请求需在 X-CSRF-Token 头携带)。""" if "csrf_token" not in session: session["csrf_token"] = secrets.token_hex(16) return session["csrf_token"] @app.before_request def _csrf_protect(): """非安全方法(POST/PUT/DELETE/PATCH)校验 X-CSRF-Token 请求头。""" if request.method not in ("POST", "PUT", "DELETE", "PATCH"): return # 登录表单(未登录,尚无 token)和静态文件跳过 if request.endpoint in ("login", "static"): return sess = session.get("csrf_token", "") header = request.headers.get("X-CSRF-Token", "") if not sess or header != sess: return jsonify({"ok": False, "error": "CSRF 校验失败"}), 403 @app.route("/api/csrf") @login_required def api_csrf(): """获取 CSRF token(前端非 GET 请求需携带 X-CSRF-Token 头)。""" return jsonify({"ok": True, "token": _csrf_token()}) # ================== 权限控制 ================== # 业务权限位(is_admin 管理员拥有全部权限,不受 perms 限制): # tasks — 任务管理:任务/自定义动作/分组的创建、修改、删除、启停、执行 # devices — 设备控制:停止设备、释放占用、清除异常、前台扫描、元素抓取 # apks — 应用管理:APK 上传、安装、删除 # logs — 日志查看 # 查看类接口(GET 状态/列表)所有登录用户可用;用户管理仅管理员可用(admin_required)。 PERM_TASKS = "tasks" PERM_DEVICES = "devices" PERM_APKS = "apks" PERM_LOGS = "logs" ALL_PERMS = (PERM_TASKS, PERM_DEVICES, PERM_APKS, PERM_LOGS) _PERM_LABELS = {PERM_TASKS: "任务管理", PERM_DEVICES: "设备控制", PERM_APKS: "应用管理", PERM_LOGS: "日志查看"} def _has_perm(perm): """当前用户是否拥有指定权限。管理员恒为 True。""" u = current_user return bool(u and (u.is_admin or u.has_perm(perm))) def _validate_perms(raw): """校验权限位列表:只保留合法值、去重。返回合法列表。""" valid = set(ALL_PERMS) out = [] for p in raw or []: if p in valid and p not in out: out.append(p) return out def perm_required(perm): """路由装饰器:要求登录且拥有指定业务权限,否则 403。""" def deco(fn): @functools.wraps(fn) @login_required def wrapper(*args, **kwargs): if not _has_perm(perm): return jsonify({"ok": False, "error": f"无权限执行此操作(需要权限: {_PERM_LABELS.get(perm, perm)})"}), 403 return fn(*args, **kwargs) return wrapper return deco def admin_required(fn): """路由装饰器:仅管理员可用(用户管理类接口),否则 403。""" @functools.wraps(fn) @login_required def wrapper(*args, **kwargs): if not current_user.is_admin: return jsonify({"ok": False, "error": "仅管理员可执行此操作"}), 403 return fn(*args, **kwargs) return wrapper @app.route("/api/me") @login_required def api_me(): """当前登录用户信息(含权限位),前端据此隐藏无权限的功能入口。""" u = current_user return jsonify({"ok": True, "user": { "id": u.id, "username": u.username, "is_admin": bool(u.is_admin), # 管理员返回全部权限位,前端统一用 perms 判断 "perms": list(ALL_PERMS) if u.is_admin else u.get_perms(), }}) # ================== 页面路由 ================== @app.route("/") @login_required def index(): """单页应用首页。""" resp = make_response(render_template("admin/monitor.html")) resp.headers["Cache-Control"] = "no-store, no-cache, must-revalidate, max-age=0" resp.headers["Pragma"] = "no-cache" return resp @app.route("/login", methods=["GET", "POST"]) def login(): if request.method == "POST": username = request.form.get("username", "") password = request.form.get("password", "") user = User.query.filter_by(username=username).first() if user and user.check_password(password): login_user(user) _csrf_token() # 建立 CSRF token,前端通过 /api/csrf 获取 _log.info(f"用户 {username} 登录") return redirect(request.args.get("next") or url_for("index")) return render_template("admin/login.html", error="用户名或密码错误") return render_template("admin/login.html", error=None) @app.route("/logout") @login_required def logout(): _log.info(f"用户 {current_user.username} 登出") logout_user() return redirect(url_for("login")) # ================== API:健康检查(无需登录,供探活) ================== @app.route("/api/health") def api_health(): """轻量健康检查:返回进程/设备/任务摘要,不暴露敏感信息,供运维探活。""" try: devices, err = mgr.get_status() devs = devices or [] return jsonify({ "ok": True, "status": "up", "time": time.time(), "device_total": len(devs), "device_online": sum(1 for d in devs if d.get("present")), "device_running": sum(1 for d in devs if d.get("worker_status") in ("running", "connecting")), "device_error": sum(1 for d in devs if d.get("worker_status") in ("error", "failed")), "occupied": sum(1 for d in devs if d.get("stf_occupied")), "jobs": len(mgr.jobs), }) except Exception as e: return jsonify({"ok": False, "status": "down", "error": str(e)}), 500 # ================== API:失败任务汇总 ================== @app.route("/api/summary") @login_required def api_summary(): """失败/异常任务汇总:各状态计数 + 异常设备列表(last_error / 重试次数)。 供监控页"异常汇总"面板使用,便于长期运行观察设备健康度。 """ try: sts = get_all_worker_status() # 已删除/不在当前 STF 设备池的设备,其陈旧失败记录不应展示。 # 复用 get_status 缓存(顺带触发上面的状态清理),STF 异常时不隐藏任何记录。 try: devices, derr = mgr.get_status() if not derr: present = {d.get("serial") for d in devices} sts = [s for s in sts if s.get("serial") in present] except Exception: pass counts = {"total": 0, "running": 0, "done": 0, "error": 0, "failed": 0, "idle": 0} errors = [] for s in sts: st = s.get("status", "idle") counts["total"] += 1 if st in counts: counts[st] += 1 if st in ("error", "failed") and s.get("last_error"): errors.append({ "serial": s.get("serial"), "model": s.get("model", ""), "status": st, "last_error": s.get("last_error", ""), "task": s.get("task_job", ""), "attempt": s.get("attempt", 0), "updated": s.get("last_heartbeat", 0), }) errors.sort(key=lambda x: x.get("updated", 0), reverse=True) return jsonify({"ok": True, "counts": counts, "errors": errors[:50]}) except Exception as e: return jsonify({"ok": False, "error": str(e)}), 500 # ================== API:状态(监控大屏用)================== @app.route("/api/status") @login_required def api_status(): devices, err = mgr.get_status() if err: return jsonify({"ok": False, "error": err}), 500 return jsonify({ "ok": True, "devices": devices, "server_time": time.time(), "fg_scanning": mgr._fg_scanner.is_scanning, "fg_last_scan": mgr._fg_scanner.last_scan_time, }) @app.route("/api/scan_foreground", methods=["POST"]) @perm_required(PERM_DEVICES) def api_scan_foreground(): """手动触发前台 App 扫描(不打扰设备)。""" started = mgr._fg_scanner.scan_once() if started: return jsonify({"ok": True, "msg": "扫描已启动"}) return jsonify({"ok": False, "error": "已有扫描在进行中"}) @app.route("/api/task_types") @login_required def api_task_types(): return jsonify({"ok": True, "task_types": list_task_types()}) @app.route("/api/actions") @login_required def api_actions(): """返回指定任务类型支持的专属操作。""" task_type = request.args.get("task_type", "") task_cls = get_task_class(task_type) if not task_cls: return jsonify({"ok": False, "error": "未知任务类型", "actions": []}) actions = task_cls.list_action_types() return jsonify({"ok": True, "actions": actions}) @app.route("/api/devices") @login_required def api_devices(): """返回所有在线设备 serial(供分组表单勾选用)。""" try: return jsonify({"ok": True, "devices": mgr.list_all_serials()}) except Exception as e: return jsonify({"ok": False, "error": str(e)}), 500 # ================== API:任务计划 CRUD ================== def _job_next_run(job): """任务下次执行时间,格式化为 "YYYY-MM-DD HH:MM"(Flask 序列化 aware datetime 会变 GMT 格式)。""" nr = mgr.next_run_of(job) return nr.strftime("%Y-%m-%d %H:%M") if nr else None @app.route("/api/jobs") @login_required def api_jobs_list(): jobs = [] for j in mgr.jobs.values(): d = j.to_dict() d["next_run"] = _job_next_run(j) jobs.append(d) return jsonify({"ok": True, "jobs": jobs, "task_types": list_task_types()}) @app.route("/api/jobs", methods=["POST"]) @perm_required(PERM_TASKS) def api_jobs_create(): data = request.json or {} name = (data.get("name") or "").strip() if not name: return jsonify({"ok": False, "error": "任务名不能为空"}), 400 task_type = data.get("task_type", "douyin_nurture") if not get_task_class(task_type): return jsonify({"ok": False, "error": f"未知任务类型: {task_type}"}), 400 job = mgr.add_job( name=name, task_type=task_type, target=data.get("target", {"mode": "all"}), params=data.get("params", {}), schedule=data.get("schedule", {"mode": "once"}), retry=data.get("retry", {"max_attempts": 1, "delay": 60}), enabled=data.get("enabled", True)) d = job.to_dict() d["next_run"] = _job_next_run(job) return jsonify({"ok": True, "msg": "任务已创建", "job": d}) @app.route("/api/jobs/", methods=["PUT"]) @perm_required(PERM_TASKS) def api_jobs_update(job_id): job = mgr.jobs.get(job_id) if not job: return jsonify({"ok": False, "error": "任务不存在"}), 404 data = request.json or {} fields = {} for k in ("name", "task_type", "target", "params", "schedule", "retry", "enabled"): if k in data: fields[k] = data[k] if "task_type" in fields and not get_task_class(fields["task_type"]): return jsonify({"ok": False, "error": f"未知任务类型: {fields['task_type']}"}), 400 mgr.update_job(job_id, **fields) d = job.to_dict() d["next_run"] = _job_next_run(job) return jsonify({"ok": True, "msg": "任务已更新", "job": d}) @app.route("/api/jobs/", methods=["DELETE"]) @perm_required(PERM_TASKS) def api_jobs_delete(job_id): if mgr.delete_job(job_id): return jsonify({"ok": True, "msg": "任务已删除"}) return jsonify({"ok": False, "error": "任务不存在"}), 404 @app.route("/api/jobs//run", methods=["POST"]) @perm_required(PERM_TASKS) def api_jobs_run(job_id): return jsonify(mgr.run_job_now(job_id)) @app.route("/api/jobs//toggle", methods=["POST"]) @perm_required(PERM_TASKS) def api_jobs_toggle(job_id): enabled = (request.json or {}).get("enabled", True) job = mgr.toggle_job(job_id, enabled) if not job: return jsonify({"ok": False, "error": "任务不存在"}), 404 return jsonify({"ok": True, "msg": f"任务已{'启用' if enabled else '停用'}"}) # ================== API:设备分组 CRUD ================== @app.route("/api/groups") @login_required def api_groups_list(): return jsonify({"ok": True, "groups": [g.to_dict() for g in mgr.groups.values()]}) @app.route("/api/groups", methods=["POST"]) @perm_required(PERM_TASKS) def api_groups_create(): data = request.json or {} name = (data.get("name") or "").strip() if not name: return jsonify({"ok": False, "error": "分组名不能为空"}), 400 if name in mgr.groups: return jsonify({"ok": False, "error": "分组名已存在"}), 400 mgr.add_group(name, data.get("serials", []), data.get("description", "")) return jsonify({"ok": True, "msg": "分组已创建"}) @app.route("/api/groups/", methods=["PUT"]) @perm_required(PERM_TASKS) def api_groups_update(name): if name not in mgr.groups: return jsonify({"ok": False, "error": "分组不存在"}), 404 data = request.json or {} mgr.update_group(name, serials=data.get("serials"), description=data.get("description")) return jsonify({"ok": True, "msg": "分组已更新"}) @app.route("/api/groups/", methods=["DELETE"]) @perm_required(PERM_TASKS) def api_groups_delete(name): if mgr.delete_group(name): return jsonify({"ok": True, "msg": "分组已删除"}) return jsonify({"ok": False, "error": "分组不存在"}), 404 # ================== API:用户管理 CRUD ================== @app.route("/api/users") @admin_required def api_users_list(): users = [{"id": u.id, "username": u.username, "is_admin": u.is_admin, "perms": u.get_perms()} for u in User.query.all()] return jsonify({"ok": True, "users": users}) @app.route("/api/users", methods=["POST"]) @admin_required def api_users_create(): data = request.json or {} username = (data.get("username") or "").strip() password = data.get("password", "") if not username or not password: return jsonify({"ok": False, "error": "用户名和密码不能为空"}), 400 if User.query.filter_by(username=username).first(): return jsonify({"ok": False, "error": "用户名已存在"}), 400 # 默认普通用户(最小权限);管理员由 is_admin 决定,权限位照常保存(降级后生效) u = User(username=username, is_admin=bool(data.get("is_admin", False))) u.set_perms(_validate_perms(data.get("perms") or [])) u.set_password(password) db.session.add(u) db.session.commit() _log.info(f"创建用户 {username} (admin={u.is_admin}, perms={u.get_perms()})") return jsonify({"ok": True, "msg": "用户已创建"}) @app.route("/api/users/", methods=["PUT"]) @admin_required def api_users_update(uid): u = User.query.get(uid) if not u: return jsonify({"ok": False, "error": "用户不存在"}), 404 data = request.json or {} if "password" in data and data["password"]: u.set_password(data["password"]) if "is_admin" in data: new_admin = bool(data["is_admin"]) if u.is_admin and not new_admin and User.query.filter_by(is_admin=True).count() <= 1: return jsonify({"ok": False, "error": "不能取消最后一个管理员"}), 400 u.is_admin = new_admin if "perms" in data: u.set_perms(_validate_perms(data["perms"])) db.session.commit() _log.info(f"更新用户 {u.username} (admin={u.is_admin}, perms={u.get_perms()})") return jsonify({"ok": True, "msg": "用户已更新"}) @app.route("/api/users/", methods=["DELETE"]) @admin_required def api_users_delete(uid): u = User.query.get(uid) if not u: return jsonify({"ok": False, "error": "用户不存在"}), 404 if u.username == "admin": return jsonify({"ok": False, "error": "不能删除默认管理员"}), 400 if u.id == current_user.id: return jsonify({"ok": False, "error": "不能删除当前登录用户"}), 400 if u.is_admin and User.query.filter_by(is_admin=True).count() <= 1: return jsonify({"ok": False, "error": "不能删除最后一个管理员"}), 400 db.session.delete(u) db.session.commit() _log.info(f"删除用户 {u.username}") return jsonify({"ok": True, "msg": "用户已删除"}) # ================== API:日志查看 ================== @app.route("/api/logs") @perm_required(PERM_LOGS) def api_logs(): files = list(_MODULE_FILES.values()) current = request.args.get("file", "core.log") lines = int(request.args.get("lines", 300)) content = "" path = os.path.join(_LOG_DIR, current) if os.path.exists(path): try: with open(path, encoding="utf-8") as f: content = "".join(f.readlines()[-lines:]) except Exception as e: content = f"读取失败: {e}" return jsonify({"ok": True, "content": content, "file": current, "files": files}) # ================== API:运行控制 ================== @app.route("/api/stop_device", methods=["POST"]) @perm_required(PERM_DEVICES) def api_stop_device(): serial = (request.json or {}).get("serial", "") if mgr.stop_device(serial): return jsonify({"ok": True, "msg": f"已发送停止信号给 {serial}"}) return jsonify({"ok": False, "error": f"{serial} 没有运行中的任务"}), 400 @app.route("/api/stop_all", methods=["POST"]) @perm_required(PERM_DEVICES) def api_stop_all(): stopped = mgr.stop_all() return jsonify({"ok": True, "stopped": stopped}) @app.route("/api/release", methods=["POST"]) @perm_required(PERM_DEVICES) def api_release(): released, failed = stf.release_all_mine() return jsonify({"ok": True, "released": released, "failed": failed}) @app.route("/api/device/clear_error", methods=["POST"]) @perm_required(PERM_DEVICES) def api_device_clear_error(): """清除单台设备的异常状态(error/failed → idle),供设备列表"清除异常"按钮使用。""" serial = (request.json or {}).get("serial", "") if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 if serial in mgr.get_running(): return jsonify({"ok": False, "error": "设备正在运行或等待重试,无法清除"}), 400 if not clear_worker_error(serial): return jsonify({"ok": False, "error": "设备正在运行任务,无法清除"}), 400 _log.info(f"清除设备 {serial} 的异常状态") return jsonify({"ok": True, "msg": f"已清除 {serial} 的异常状态"}) @app.route("/api/device/clear_all_errors", methods=["POST"]) @perm_required(PERM_DEVICES) def api_device_clear_all_errors(): """一键清除所有异常/失败设备(跳过正在运行/等待重试的)。""" running = set(mgr.get_running()) cleared = clear_all_worker_errors(exclude=running) _log.info(f"一键清除异常:共清除 {cleared} 台") return jsonify({"ok": True, "cleared": cleared, "msg": f"已清除 {cleared} 台设备的异常状态"}) @app.route("/api/device/screenshot") @login_required def api_device_screenshot(): """获取设备当前画面截图(PNG)。 用 adb exec-out screencap -p,只读操作,不抢占 u2 的 atx-agent 通道, 任务运行中调用安全。直接返回 image/png,前端用 加载。 ?serial=xxx 设备 serial ?t=123 时间戳,避免浏览器缓存(前端自动加) """ serial = request.args.get("serial", "") if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 ok, data = screenshot(serial) if ok: return Response(data, mimetype="image/png", headers={"Cache-Control": "no-store"}) return jsonify({"ok": False, "error": data}), 500 @app.route("/api/devices//apps") @login_required def api_device_apps(serial): """获取指定设备上已安装的应用列表(包名 + versionCode + versionName + 路径)。""" ok, data = list_installed_apps(serial) if ok: return jsonify({"ok": True, "apps": data}) return jsonify({"ok": False, "error": data}), 500 # ================== 自定义动作(步骤打包) ================== @app.route("/api/custom_actions") @login_required def api_custom_actions_list(): rows = CustomAction.query.order_by(CustomAction.created_at.desc()).all() return jsonify({"ok": True, "actions": [r.to_dict() for r in rows]}) @app.route("/api/custom_actions", methods=["POST"]) @perm_required(PERM_TASKS) def api_custom_actions_create(): data = request.json or {} name = (data.get("name") or "").strip() if not name: return jsonify({"ok": False, "error": "动作名不能为空"}), 400 steps = data.get("steps", []) if not steps: return jsonify({"ok": False, "error": "至少需要1个步骤"}), 400 row = CustomAction(id=uuid.uuid4().hex[:8], name=name, icon=data.get("icon") or "📦", created_at=datetime.now().strftime("%Y-%m-%d %H:%M")) row.set_steps(steps) db.session.add(row) db.session.commit() return jsonify({"ok": True, "msg": "动作已保存", "action": row.to_dict()}) @app.route("/api/custom_actions/", methods=["PUT"]) @perm_required(PERM_TASKS) def api_custom_actions_update(action_id): row = CustomAction.query.get(action_id) if not row: return jsonify({"ok": False, "error": "动作不存在"}), 404 data = request.json or {} if "name" in data: name = (data["name"] or "").strip() if not name: return jsonify({"ok": False, "error": "动作名不能为空"}), 400 row.name = name if "icon" in data: row.icon = data["icon"] or "📦" if "steps" in data: if not data["steps"]: return jsonify({"ok": False, "error": "至少需要1个步骤"}), 400 row.set_steps(data["steps"]) db.session.commit() return jsonify({"ok": True, "msg": "已更新", "action": row.to_dict()}) @app.route("/api/custom_actions/", methods=["DELETE"]) @perm_required(PERM_TASKS) def api_custom_actions_delete(action_id): row = CustomAction.query.get(action_id) if not row: return jsonify({"ok": False, "error": "动作不存在"}), 404 db.session.delete(row) db.session.commit() return jsonify({"ok": True, "msg": "已删除"}) @app.route("/api/uiauto/elements") @perm_required(PERM_DEVICES) def api_uiauto_elements(): """获取设备当前 UI 元素树(供步骤编辑器"抓取元素"用)。 依赖本地运行的 uiautodev 服务(端口 20242)。 ?serial=xxx 设备 serial 返回: 200 — {ok:true, elements:[...]} 503 — {ok:false, error:"..."}(uiauto2 未启动) """ serial = request.args.get("serial", "") if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 ok, data = uiauto_helper.get_elements(serial) if ok: return jsonify({"ok": True, "elements": data}) return jsonify({"ok": False, "error": data}), 503 @app.route("/api/uiauto/status") @login_required def api_uiauto_status(): """探测 uiauto2 本地服务是否运行(前端按钮禁启用)。""" return jsonify({"ok": True, "running": uiauto_helper.is_running()}) @app.route("/api/steps/test", methods=["POST"]) @perm_required(PERM_DEVICES) def api_steps_test(): """测试单个步骤:在指定设备上试执行,验证选择器是否命中(步骤编辑器"测试此步骤")。 请求: {"serial": "100.100.10.11:5555", "step": {"type": "click", "params": {...}}} 只读连接(adb connect + u2),不占用/释放 STF,与运行中任务互不干扰。 返回 result: "命中" / "未找到" / "已执行"。 """ data = request.json or {} serial = (data.get("serial") or "").strip() step = data.get("step") or {} if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 if not step.get("type"): return jsonify({"ok": False, "error": "缺少步骤类型"}), 400 from tasks.generic.task import test_step ok, msg, result = test_step(serial, step) res_txt = {True: "命中", False: "未找到"}.get(result, "已执行") _log.info(f"测试步骤 {step.get('type')} @ {serial}: {res_txt} ({msg})") return jsonify({"ok": ok, "msg": msg, "result": res_txt}) @app.route("/api/uiauto/devices") @perm_required(PERM_DEVICES) def api_uiauto_devices(): """获取可选设备列表(抓取元素用)。 uiautodev 只认识本地 adb 已连接的设备;平台设备池(STF)里在线但未 本地连接的设备(如 100.100.10.10)在这里补全,并先做一次轻量 adb connect(单次尝试,绝不 kill-server),让它们可被 uiautodev 抓取。 """ ok, data = uiauto_helper.list_devices() if not ok: return jsonify({"ok": False, "error": data}), 503 seen = {d.get("serial") for d in data} try: pool = stf.list_all_devices() except Exception: pool = [] added = 0 for d in pool: serial = d.get("serial", "") if not serial or serial in seen: continue if not (d.get("present") and d.get("ready")): continue # 轻量连接(单次尝试,不重试不 kill-server);连不上也照样列出, # 前端抓取时会报明确错误,不影响其他设备的连接 adb_connect_light(serial) data.append({ "serial": serial, "model": d.get("model") or d.get("product") or "", "product": d.get("product") or "", "name": serial, "status": "device", "enabled": True, }) seen.add(serial) added += 1 if added: _log.info(f"抓取元素设备列表补全 {added} 台 STF 在线设备: {list(seen)}") return jsonify({"ok": True, "devices": data}) @app.route("/api/uiauto/screenshot") @perm_required(PERM_DEVICES) def api_uiauto_screenshot(): """通过 uiauto2 获取设备截图(供抓取元素时显示设备画面)。 ?serial=xxx 设备 serial 直接返回 image/jpeg,前端用 加载。 """ serial = request.args.get("serial", "") if not serial: return jsonify({"ok": False, "error": "缺少 serial"}), 400 ok, data = uiauto_helper.get_screenshot(serial) if ok: return Response(data, mimetype="image/jpeg", headers={"Cache-Control": "no-store"}) return jsonify({"ok": False, "error": data}), 503 # ================== API:应用管理(APK 上传/安装)================== @app.route("/api/apks") @login_required def api_apks_list(): """列出所有已上传的 APK。""" return jsonify({"ok": True, "apks": apk_mgr.list_all()}) @app.route("/api/apks/upload", methods=["POST"]) @perm_required(PERM_APKS) def api_apks_upload(): """上传 APK 文件(multipart/form-data,字段名 file)。""" file = request.files.get("file") if not file or not file.filename: return jsonify({"ok": False, "error": "未选择文件"}), 400 info = apk_mgr.upload(file) if info: return jsonify({"ok": True, "apk": info, "msg": f"上传成功: {info['display_name']}"}) return jsonify({"ok": False, "error": "上传失败,请检查文件格式"}), 500 @app.route("/api/apks/", methods=["DELETE"]) @perm_required(PERM_APKS) def api_apks_delete(apk_id): """删除 APK 文件和记录。""" ok, msg = apk_mgr.delete(apk_id) if ok: return jsonify({"ok": True, "msg": msg}) return jsonify({"ok": False, "error": msg}), 400 @app.route("/api/apks/install/devices") @perm_required(PERM_APKS) def api_apks_install_devices(): """可安装设备列表(应用管理):STF 在线设备池 + 本机 adb 设备(含 USB 有线连接)。 source: stf=STF 池 / usb=本机有线 adb / adb=本机网络 adb。 """ devices, seen = [], set() try: for d in stf.list_all_devices(): serial = d.get("serial", "") if serial and d.get("present") and d.get("ready"): devices.append({"serial": serial, "model": d.get("model") or d.get("manufacturer") or "", "source": "stf"}) seen.add(serial) except Exception: pass out = _adb("devices") for line in (out or "").splitlines()[1:]: parts = line.split() if len(parts) >= 2 and parts[0] and parts[0] not in seen: devices.append({"serial": parts[0], "model": "", "source": "usb" if ":" not in parts[0] else "adb"}) seen.add(parts[0]) return jsonify({"ok": True, "devices": devices}) @app.route("/api/apks/install", methods=["POST"]) @perm_required(PERM_APKS) def api_apks_install(): """批量安装 APK 到指定设备。参数: {apk_id, serials:[]}""" data = request.json or {} apk_id = data.get("apk_id", "") serials = data.get("serials", []) ok, msg = apk_mgr.install(apk_id, serials) if ok: return jsonify({"ok": True, "msg": msg}) return jsonify({"ok": False, "error": msg}), 400 @app.route("/api/apks/install/status") @login_required def api_apks_install_status(): """获取安装任务实时状态。""" status = apk_mgr.get_install_status() return jsonify({"ok": True, "status": status}) # ================== API:维护(仅管理员) ================== # 维护功能:STF 容器一键重启 + adb 远程终端。两个接口都仅管理员可用。 # adb 终端安全红线(技术约束):绝不 kill-server / 绝不 disconnect IP:5555—— # STF provider 共享该 adb transport,断开会让 STF 误判全部设备离线并触发重连。 _ADB_BLOCKED_PATTERNS = ("kill-server", "disconnect") def _blocked_adb_cmd(cmd): """命中红线的 adb 命令(kill-server / disconnect)直接拒绝。""" low = cmd.lower() return any(p in low for p in _ADB_BLOCKED_PATTERNS) @app.route("/api/adb/devices") @admin_required def api_adb_devices(): """维护终端设备列表(仅管理员):本地 adb 已连接 + STF 在线设备池。 供终端设备选择器使用——选中后自动附加 `-s `, 离线/未连接的设备也可选,配合"重连设备"按钮恢复。 """ states = {} out = _adb("devices") for line in (out or "").splitlines()[1:]: parts = line.split() if len(parts) >= 2 and parts[0] and not parts[0].startswith("*"): states[parts[0]] = parts[1] try: for d in stf.list_all_devices(): serial = d.get("serial", "") if serial and d.get("present") and d.get("ready") and serial not in states: states[serial] = "stf" except Exception: pass devices = [{"serial": s, "state": st} for s, st in sorted(states.items())] return jsonify({"ok": True, "devices": devices}) @app.route("/api/adb/cmd", methods=["POST"]) @admin_required def api_adb_cmd(): """维护终端:执行 adb 命令(仅管理员)。 请求: {"cmd": "adb -s 100.100.10.11:5555 shell ls /sdcard"} 用平台 adb 二进制执行(config.ADB_PATH),带 20s 超时。 """ data = request.json or {} cmd = (data.get("cmd") or "").strip() if not cmd: return jsonify({"ok": False, "error": "命令不能为空"}), 400 if _blocked_adb_cmd(cmd): return jsonify({"ok": False, "error": "禁止执行 kill-server / disconnect(会中断 STF 设备监控,影响所有任务)"}), 400 tokens = shlex.split(cmd) if tokens and tokens[0] == "adb": tokens = tokens[1:] if not tokens: return jsonify({"ok": False, "error": "命令不能为空"}), 400 try: # 与 worker 的 adb 调用共用锁,避免并发操作同一 adb server with _ADB_LOCK: r = subprocess.run([ADB_PATH, *tokens], capture_output=True, timeout=20) out = (r.stdout or b"").decode("utf-8", errors="replace") err = (r.stderr or b"").decode("utf-8", errors="replace") _log.info(f"维护终端执行 adb {' '.join(tokens)} -> exit {r.returncode}") return jsonify({"ok": True, "stdout": out, "stderr": err, "code": r.returncode}) except subprocess.TimeoutExpired: return jsonify({"ok": False, "error": "命令执行超时(20s)"}), 408 except Exception as e: _log.warning(f"维护终端 adb 执行异常: {e}") return jsonify({"ok": False, "error": f"执行失败: {e}"}), 500 @app.route("/api/stf/restart", methods=["POST"]) @admin_required def api_stf_restart(): """一键重启 STF Docker 容器(仅管理员)。 通过 SSH(config.STF_SSH_TARGET,默认 stf@192.168.20.220)执行 `docker restart <容器>`(config.STF_DOCKER_CONTAINER,默认 stf), 随后轮询 STF API 确认恢复。 影响:STF 约 30s 不可用,期间占用状态重置、任务中断——有任务在跑时拒绝执行。 """ running = mgr.get_running() if running: return jsonify({"ok": False, "error": f"有 {len(running)} 台设备正在运行任务,请先停止再维护"}), 409 ssh_cmd = ["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=5", "-o", "StrictHostKeyChecking=no", STF_SSH_TARGET, f"docker restart {STF_DOCKER_CONTAINER} && " f"docker ps --filter name={STF_DOCKER_CONTAINER} --format '{{{{.Names}}}}: {{{{.Status}}}}'"] try: r = subprocess.run(ssh_cmd, capture_output=True, timeout=60) except subprocess.TimeoutExpired: return jsonify({"ok": False, "error": "SSH 执行超时(60s)"}), 408 except FileNotFoundError: return jsonify({"ok": False, "error": "本机未安装 ssh,无法远程操作"}), 500 out = (r.stdout or b"").decode("utf-8", errors="replace") err = (r.stderr or b"").decode("utf-8", errors="replace") if r.returncode != 0: _log.error(f"STF 容器重启失败: {err or out}") return jsonify({"ok": False, "error": f"SSH 执行失败:{err or out}(请确认本机可免密 SSH 到 {STF_SSH_TARGET},且容器名 {STF_DOCKER_CONTAINER} 正确)"}), 502 # 轮询 STF API 确认恢复(最长 30s) alive = False for _ in range(30): try: stf.list_all_devices() alive = True break except Exception: time.sleep(1) _log.info(f"STF 容器 {STF_DOCKER_CONTAINER} 已重启,API 恢复={alive}") return jsonify({"ok": True, "msg": "STF 容器已重启" + ("" if alive else "(STF API 尚未恢复,请稍后刷新)"), "detail": out.strip(), "api_alive": alive}) # ================== API:Tailscale 管理(仅管理员) ================== # 通过 Tailscale 官方 API v2 管理 tailnet 设备(列表/改名/授权/密钥不过期/删除/生成 auth key)。 # 设备 IP 由 tailnet 自动分配,API 无法修改,列表只读展示。 @app.route("/api/tailscale/status") @admin_required def api_tailscale_status(): """Tailscale 管理配置状态(API key / tailnet 是否已配置)。""" c = tailscale_client.tailscale return jsonify({"ok": True, "configured": c.is_configured(), "hint": c.config_hint() if not c.is_configured() else ""}) @app.route("/api/tailscale/devices") @admin_required def api_tailscale_devices(): """列出 tailnet 全部设备。""" try: return jsonify({"ok": True, "devices": tailscale_client.tailscale.list_devices()}) except TailscaleError as e: return jsonify({"ok": False, "error": str(e)}), 502 @app.route("/api/tailscale/devices/", methods=["POST"]) @admin_required def api_tailscale_device_update(device_id): """更新设备:按传入字段分发到专属端点(改名 /name、授权 /authorized、密钥 /key)。""" data = request.json or {} c = tailscale_client.tailscale try: if "name" in data: c.set_device_name(device_id, data.get("name")) if "authorized" in data: c.set_device_authorized(device_id, data.get("authorized")) if "key_expiry_disabled" in data: c.set_device_key_expiry(device_id, data.get("key_expiry_disabled")) _log.info(f"Tailscale 设备更新 {device_id}: {data}") return jsonify({"ok": True, "msg": "设备已更新"}) except TailscaleError as e: return jsonify({"ok": False, "error": str(e)}), 502 @app.route("/api/tailscale/devices//ip", methods=["POST"]) @admin_required def api_tailscale_device_ip(device_id): """设置设备 IPv4 地址(未公开端点,实测可用)。 ⚠ 改 IP 会断开设备当前 tailscale 会话,且平台设备池 serial 随之变化, 改完需同步更新 STF 设备池(分组/任务目标里的旧 IP 会失效)。 """ data = request.json or {} ipv4 = (data.get("ipv4") or "").strip() if not ipv4: return jsonify({"ok": False, "error": "缺少 ipv4"}), 400 try: tailscale_client.tailscale.set_device_ip(device_id, ipv4) _log.info(f"Tailscale 设备 {device_id} IP 已设置为 {ipv4}") return jsonify({"ok": True, "msg": f"设备 IP 已设置为 {ipv4}"}) except TailscaleError as e: return jsonify({"ok": False, "error": str(e)}), 502 @app.route("/api/tailscale/devices/", methods=["DELETE"]) @admin_required def api_tailscale_device_delete(device_id): """从 tailnet 移除设备。""" try: tailscale_client.tailscale.delete_device(device_id) _log.info(f"Tailscale 设备已移除: {device_id}") return jsonify({"ok": True, "msg": "设备已从 tailnet 移除"}) except TailscaleError as e: return jsonify({"ok": False, "error": str(e)}), 502 @app.route("/api/tailscale/authkey", methods=["POST"]) @admin_required def api_tailscale_authkey(): """生成设备接入 auth key(key 只显示一次)。""" data = request.json or {} try: # Tailscale 要求 description 仅 ASCII(中文会 400) desc = re.sub(r"[^\x20-\x7e]", "", (data.get("description") or "").strip()) key = tailscale_client.tailscale.create_auth_key( description=desc or "auto_control", reusable=bool(data.get("reusable", False)), ephemeral=bool(data.get("ephemeral", False)), preauthorized=data.get("preauthorized", True), expiry_seconds=int(data.get("expiry_seconds", 3600))) _log.info(f"Tailscale auth key 已生成: {key['id']}") return jsonify({"ok": True, "key": key["key"], "id": key["id"], "expires": key["expires"]}) except TailscaleError as e: return jsonify({"ok": False, "error": str(e)}), 502 # ================== uiautodev 自动启动 ================== _uiauto_proc = None # PID 文件:web_server 异常退出(kill -9/崩溃)时 uiautodev 成孤儿残留, # 用 PID 文件在下次启动时识别并清理,避免多次重启后堆积进程。 _UIAUTO_PID_FILE = os.path.join( os.path.dirname(os.path.abspath(__file__)), "data", "uiauto.pid") def _kill_stale_uiauto(): """清理上次 web_server 残留的 uiautodev 进程(读取 PID 文件)。""" if not os.path.exists(_UIAUTO_PID_FILE): return try: with open(_UIAUTO_PID_FILE) as f: pid = int(f.read().strip()) if pid > 0 and pid != os.getpid(): try: os.kill(pid, 0) # 探测进程是否存活 os.kill(pid, 15) # 终止残留 _log.warning(f"已清理残留的 uiautodev 进程 (PID={pid})") except ProcessLookupError: pass except PermissionError: _log.warning(f"无权限清理残留 uiautodev (PID={pid}),忽略") except Exception: pass finally: try: os.remove(_UIAUTO_PID_FILE) except OSError: pass def _ensure_uiauto_running(): """确保 uiautodev 本地服务在运行(元素抓取功能依赖,端口 20242)。 已在运行则跳过;未运行则用子进程启动 `uiautodev server --no-browser`。 启动失败只记日志,不影响主服务。 """ global _uiauto_proc _kill_stale_uiauto() if uiauto_helper.is_running(): _log.info("uiauto2 服务已在运行 (端口 20242)") return try: _log.info("正在自动启动 uiautodev 服务 (端口 20242)...") kwargs = dict(stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) if sys.platform == "win32": kwargs["creationflags"] = subprocess.CREATE_NO_WINDOW _uiauto_proc = subprocess.Popen( [sys.executable, "-m", "uiautodev", "server", "--no-browser"], **kwargs, ) _log.info(f"uiautodev 服务已启动 (PID={_uiauto_proc.pid})") try: with open(_UIAUTO_PID_FILE, "w") as f: f.write(str(_uiauto_proc.pid)) except OSError: pass atexit.register(_stop_uiauto) except FileNotFoundError: _log.warning("uiautodev 未安装,元素抓取功能不可用。请运行: pip install uiautodev") except Exception as e: _log.warning(f"自动启动 uiautodev 失败: {e}(元素抓取功能不可用)") def _stop_uiauto(): """停止自动启动的 uiautodev 子进程。""" global _uiauto_proc try: if os.path.exists(_UIAUTO_PID_FILE): os.remove(_UIAUTO_PID_FILE) except OSError: pass if not _uiauto_proc: return try: _uiauto_proc.terminate() _uiauto_proc.wait(timeout=5) _log.info("uiautodev 服务已停止") except Exception: try: _uiauto_proc.kill() except Exception: pass _uiauto_proc = None def _run_server(host, port): """启动 Flask,自动处理端口冲突和权限问题。 常见失败原因: WinError 10013 — 端口在 Windows 动态端口范围内被出站连接占用, 或非管理员绑定 0.0.0.0。前者换端口,后者降级 127.0.0.1。 WinError 10048 — 端口已被其他进程监听。换端口重试。 策略:原端口失败 → 试 127.0.0.1:原端口 → 试 127.0.0.1:原端口+1..+5 """ import socket # 先探测原端口是否可用(避免 flask 报错打日志难看) candidates = [(host, port)] if host == "0.0.0.0": candidates.append(("127.0.0.1", port)) for i in range(1, 6): candidates.append(("127.0.0.1", port + i)) for h, p in candidates: try: with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) s.bind((h, p)) # 探测成功,启动 flask if (h, p) != (host, port): _log.warning(f"原端口 {host}:{port} 不可用,改用 {h}:{p}") _log.info(f"启动服务: http://localhost:{p}/") app.run(host=h, port=p, debug=False, threaded=True) return except OSError as e: _log.warning(f"绑定 {h}:{p} 失败: {e}") # 全部失败 raise OSError(f"无法绑定任何候选端口({host}:{port} 及 127.0.0.1:{port}~{port+5})") if __name__ == "__main__": from config import WEB_HOST, WEB_PORT, AUTO_RELEASE_STALE_OCCUPY # 自动启动 uiautodev 服务(元素抓取功能依赖,端口 20242) _ensure_uiauto_running() # 启动时处理可能残留的 STF 占用(上次异常退出/强杀遗留)。 # AUTO_RELEASE_STALE_OCCUPY=True(单实例)时自动释放;否则仅提示。 # 多实例共用一个 STF 账户时不要开启,误放会中断另一实例的任务。 try: _mine = stf.list_my_devices() if _mine: _serials = ", ".join(d["serial"] for d in _mine) if AUTO_RELEASE_STALE_OCCUPY: _log.warning(f"启动时发现 {len(_mine)} 台设备残留占用,自动释放(AUTO_RELEASE_STALE_OCCUPY=True):{_serials}") _released, _failed = stf.release_all_mine() _log.info(f"自动释放完成:成功 {len(_released)} 台,失败 {len(_failed)} 台") else: _log.warning(f"启动时发现 {len(_mine)} 台设备仍被本账户占用(可能异常退出残留):{_serials}。" f"单实例无人值守可设 AUTO_RELEASE_STALE_OCCUPY=True 自动清理") except Exception: pass _log.info(f"管理后台: http://localhost:{WEB_PORT}/ (admin/admin123)") try: _run_server(WEB_HOST, WEB_PORT) finally: mgr.shutdown() _stop_uiauto()