Files
auto_control/web_server.py
T
butubb 3a38d9a57f fix: adb 命令超时 + STF 释放失败检查,杜绝设备卡死/占用悬空
- adb 所有命令加 30s 超时:connect 到不可达地址不再无限挂起(worker 不再卡 connecting)
- STF release 检查响应状态并重试:504 等失败不再被静默吞掉,释放失败会真实上报
- release_all_mine 返回 (released, failed);web /api/release 透出失败列表
2026-08-08 16:30:42 +08:00

633 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""设备自动化管理后台(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 subprocess
from datetime import datetime
from flask import Flask, jsonify, request, redirect, url_for, render_template, Response, make_response
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.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
from core import uiauto_helper
from tasks import list_task_types, get_task_class
_log = get_logger("web")
app = Flask(__name__)
app.config["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))
# ================== 页面路由 ==================
@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)
_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/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"])
@login_required
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 ==================
@app.route("/api/jobs")
@login_required
def api_jobs_list():
return jsonify({"ok": True, "jobs": [j.to_dict() for j in mgr.jobs.values()],
"task_types": list_task_types()})
@app.route("/api/jobs", methods=["POST"])
@login_required
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))
return jsonify({"ok": True, "msg": "任务已创建", "job": job.to_dict()})
@app.route("/api/jobs/<job_id>", methods=["PUT"])
@login_required
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)
return jsonify({"ok": True, "msg": "任务已更新", "job": job.to_dict()})
@app.route("/api/jobs/<job_id>", methods=["DELETE"])
@login_required
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/<job_id>/run", methods=["POST"])
@login_required
def api_jobs_run(job_id):
return jsonify(mgr.run_job_now(job_id))
@app.route("/api/jobs/<job_id>/toggle", methods=["POST"])
@login_required
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"])
@login_required
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/<name>", methods=["PUT"])
@login_required
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/<name>", methods=["DELETE"])
@login_required
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")
@login_required
def api_users_list():
users = [{"id": u.id, "username": u.username, "is_admin": u.is_admin} for u in User.query.all()]
return jsonify({"ok": True, "users": users})
@app.route("/api/users", methods=["POST"])
@login_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
u = User(username=username, is_admin=data.get("is_admin", True))
u.set_password(password)
db.session.add(u)
db.session.commit()
_log.info(f"创建用户 {username}")
return jsonify({"ok": True, "msg": "用户已创建"})
@app.route("/api/users/<int:uid>", methods=["PUT"])
@login_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:
u.is_admin = bool(data["is_admin"])
db.session.commit()
_log.info(f"更新用户 {u.username}")
return jsonify({"ok": True, "msg": "用户已更新"})
@app.route("/api/users/<int:uid>", methods=["DELETE"])
@login_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
db.session.delete(u)
db.session.commit()
_log.info(f"删除用户 {u.username}")
return jsonify({"ok": True, "msg": "用户已删除"})
# ================== API:日志查看 ==================
@app.route("/api/logs")
@login_required
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"])
@login_required
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"])
@login_required
def api_stop_all():
stopped = mgr.stop_all()
return jsonify({"ok": True, "stopped": stopped})
@app.route("/api/release", methods=["POST"])
@login_required
def api_release():
released, failed = stf.release_all_mine()
return jsonify({"ok": True, "released": released, "failed": failed})
@app.route("/api/device/screenshot")
@login_required
def api_device_screenshot():
"""获取设备当前画面截图(PNG)。
用 adb exec-out screencap -p,只读操作,不抢占 u2 的 atx-agent 通道,
任务运行中调用安全。直接返回 image/png,前端用 <img> 加载。
?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/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"])
@login_required
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/<action_id>", methods=["PUT"])
@login_required
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/<action_id>", methods=["DELETE"])
@login_required
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")
@login_required
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/uiauto/devices")
@login_required
def api_uiauto_devices():
"""获取 uiauto2 已连接的设备列表(供抓取元素时选择设备)。"""
ok, data = uiauto_helper.list_devices()
if ok:
return jsonify({"ok": True, "devices": data})
return jsonify({"ok": False, "error": data}), 503
@app.route("/api/uiauto/screenshot")
@login_required
def api_uiauto_screenshot():
"""通过 uiauto2 获取设备截图(供抓取元素时显示设备画面)。
?serial=xxx 设备 serial
直接返回 image/jpeg,前端用 <img> 加载。
"""
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"])
@login_required
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/<apk_id>", methods=["DELETE"])
@login_required
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", methods=["POST"])
@login_required
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})
# ================== uiautodev 自动启动 ==================
_uiauto_proc = None
def _ensure_uiauto_running():
"""确保 uiautodev 本地服务在运行(元素抓取功能依赖,端口 20242)。
已在运行则跳过;未运行则用子进程启动 `uiautodev server --no-browser`。
启动失败只记日志,不影响主服务。
"""
global _uiauto_proc
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})")
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
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
# 自动启动 uiautodev 服务(元素抓取功能依赖,端口 20242)
_ensure_uiauto_running()
_log.info(f"管理后台: http://localhost:{WEB_PORT}/ (admin/admin123)")
try:
_run_server(WEB_HOST, WEB_PORT)
finally:
mgr.shutdown()
_stop_uiauto()