"""设备自动化管理后台(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 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 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__) # 生产必须通过环境变量 WEB_SECRET_KEY 设置强密钥;缺省用开发密钥(不安全) app.config["SECRET_KEY"] = os.environ.get( "WEB_SECRET_KEY", "dev-secret-key-change-in-production") 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}) # ================== 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()