Files
auto_control/web_server.py
T

1840 lines
73 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 secrets
import functools
import shlex
import re
import subprocess
import urllib.parse
import threading
from datetime import datetime
from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError as FuturesTimeout
from flask import (Flask, jsonify, request, redirect, url_for, render_template,
render_template_string, Response, make_response, session)
from markupsafe import escape as _esc
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, get_all_worker_status
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_connect, _ADB_LOCK, _adb)
from core import uiauto_helper, tailscale_client, stf_device_mgmt, ssh_client
from core.tailscale_client import TailscaleError
from core.stf_device_mgmt import StfDevError
from core.ssh_client import SSHError
from config import ADB_PATH, 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()
# 设备池(本地清单 + adb 在线状态):阶段 0 起作为 STF 的替代数据源,
# 首次启动自动从 STF 导入现有设备;阶段 3 摘除 STF 后此模块独立运行
from core import device_pool
device_pool.init_app(app)
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/<job_id>", 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/<job_id>", 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/<job_id>/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/<job_id>/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/<name>", 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/<name>", 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/<int:uid>", 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/<int:uid>", 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("/locate")
def locate_page():
"""设备端定位页(免登录):点亮后浏览器打开,大字显示设备 IP/序列号。
只展示 serial 文本,无任何平台数据,免登录可接受。
"""
serial = request.args.get("serial", "")
return render_template_string(
"""<!DOCTYPE html><html><head><meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>设备定位</title><style>
body{margin:0;background:#000;color:#fff;display:flex;flex-direction:column;
align-items:center;justify-content:center;height:100vh;font-family:monospace}
.serial{font-size:min(9vw,72px);font-weight:700;color:#ffd700;word-break:break-all;padding:0 20px;text-align:center}
.ip{font-size:min(5vw,36px);color:#7dd3fc;margin-top:24px;word-break:break-all;padding:0 20px;text-align:center}
.hint{font-size:min(3vw,16px);color:#64748b;margin-top:40px}
</style></head><body>
<div class="serial">{{ serial }}</div>
<div class="ip">{{ ip }}</div>
<div class="hint">定位完成 · 按返回键退出</div>
</body></html>""",
serial=_esc(serial), ip=_esc(serial.split(":")[0]))
# 常见浏览器包名(结束定位时 force-stop 用)
_BROWSER_PKGS = {"com.android.browser", "com.miui.browser", "com.android.chrome",
"com.brave.browser", "com.opera.browser", "org.mozilla.firefox",
"com.UCMobile", "com.tencent.mtt", "com.baidu.browser.apps"}
# 定位时启动的浏览器:serial -> 包名(结束定位时无论前后台都关闭它)
_locate_browsers = {}
_locate_lock = threading.Lock()
def _current_focus_pkg(serial):
"""取设备当前前台包名(dumpsys window 解析,失败返回空串)。"""
try:
r = subprocess.run([ADB_PATH, "-s", serial, "shell", "dumpsys", "window"],
capture_output=True, timeout=20)
out = (r.stdout or b"").decode("utf-8", errors="replace")
m = re.search(r"mCurrentFocus=.*?([\w.]+)/", out)
return m.group(1) if m else ""
except Exception:
return ""
@app.route("/api/device/locate/stop", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_device_locate_stop():
"""结束定位:关闭定位启动的浏览器(无论当前是否前台)。
优先 force-stop 定位时记录下的浏览器包名(切到后台也能关掉);
无记录时退回:前台是浏览器 → force-stop;否则按返回键(不误杀任务应用)。
"""
serial = str((request.json or {}).get("serial", "")).strip()
if not serial:
return jsonify({"ok": False, "error": "缺少设备 serial"}), 400
try:
if ":" in serial:
adb_connect(serial)
# 1. 优先关闭定位时启动的浏览器(无论前后台)
with _locate_lock:
started = _locate_browsers.pop(serial, None)
if started:
subprocess.run([ADB_PATH, "-s", serial, "shell", "am", "force-stop", started],
capture_output=True, timeout=20)
msg = f"已关闭定位浏览器 {started}"
else:
# 2. 退回:前台是浏览器则关闭,否则返回键轻量退出
pkg = _current_focus_pkg(serial)
if pkg in _BROWSER_PKGS:
subprocess.run([ADB_PATH, "-s", serial, "shell", "am", "force-stop", pkg],
capture_output=True, timeout=20)
msg = f"已关闭浏览器 {pkg}"
else:
subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "4"],
capture_output=True, timeout=15) # BACK 轻量退出
msg = "已按返回键退出" + (f"(前台 {pkg})" if pkg else "")
_log.info(f"结束定位 {serial}: {msg}")
return jsonify({"ok": True, "msg": msg})
except Exception as e:
return jsonify({"ok": False, "error": f"结束定位失败: {e}"}), 502
@app.route("/api/device/screen_all", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_device_screen_all():
"""批量亮屏/息屏所有在线设备(并发)。mode: on=亮屏解锁 / off=息屏。
息屏会中断运行中的任务(u2 无法操作),前端已有确认提示。
"""
data = request.json or {}
mode = data.get("mode", "on")
try:
serials = device_pool.list_online()
except Exception as e:
return jsonify({"ok": False, "error": f"获取设备列表失败: {e}"}), 502
if not serials:
return jsonify({"ok": False, "error": "无在线设备"}), 400
def _do(serial):
try:
if ":" in serial:
adb_connect(serial)
if mode == "off":
subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "26"],
capture_output=True, timeout=15) # KEYCODE_POWER 息屏
else:
subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "224"],
capture_output=True, timeout=15) # WAKEUP
subprocess.run([ADB_PATH, "-s", serial, "shell", "wm", "dismiss-keyguard"],
capture_output=True, timeout=15)
return True
except Exception:
return False
with ThreadPoolExecutor(max_workers=min(10, len(serials))) as pool:
results = list(pool.map(_do, serials))
ok_n = sum(1 for x in results if x)
_log.info(f"批量{'息屏' if mode == 'off' else '亮屏'}: 成功 {ok_n}/{len(serials)}")
return jsonify({"ok": True, "success": ok_n, "total": len(serials)})
@app.route("/api/device/locate", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_device_locate():
"""定位设备:点亮屏幕并解除锁屏,可选打开大字定位页(查找设备用)。
请求 {"serial": "...", "show": true}:show=true 时额外用浏览器打开
平台 /locate 定位页,全屏大字显示设备 IP(更醒目,但会切换前台,
任务运行中慎用)。点亮组合:WAKEUP → dismiss-keyguard → MENU 兜底。
绝不 disconnect,红线。
"""
data = request.json or {}
serial = str(data.get("serial", "")).strip()
if not serial:
return jsonify({"ok": False, "error": "缺少设备 serial"}), 400
try:
if ":" in serial:
adb_connect(serial)
for cmd in (["shell", "input", "keyevent", "224"],
["shell", "wm", "dismiss-keyguard"],
["shell", "input", "keyevent", "82"]):
try:
subprocess.run([ADB_PATH, "-s", serial, *cmd],
capture_output=True, timeout=15)
except Exception:
pass # 单步失败不阻塞,继续下一步
msgs = ["屏幕已点亮"]
if data.get("show"):
# 用设备浏览器打开定位页(设备走 tailnet 访问本机 100.100.10.2:18050)
locate_url = f"http://100.100.10.2:18050/locate?serial={urllib.parse.quote(serial)}"
r = subprocess.run([ADB_PATH, "-s", serial, "shell", "am", "start",
"-a", "android.intent.action.VIEW", "-d", locate_url],
capture_output=True, timeout=20)
out = ((r.stdout or b"") + (r.stderr or b"")).decode("utf-8", errors="replace")
if "error" in out.lower() or "exception" in out.lower():
msgs.append(f"打开定位页失败: {out.strip()[:80]}")
else:
msgs.append("已打开大字定位页(按返回键退出)")
# 记录启动的浏览器包名:结束定位时无论前后台都关闭它。
# 优先从 am start 输出解析 pkg=(浏览器可能尚未到前台,前台检测会扑空)
pkg = ""
m2 = re.search(r"pkg=([\w.]+)", out)
if m2 and m2.group(1) in _BROWSER_PKGS:
pkg = m2.group(1)
else:
time.sleep(1.5) # 等浏览器切到前台再查
pkg = _current_focus_pkg(serial)
if pkg in _BROWSER_PKGS:
with _locate_lock:
_locate_browsers[serial] = pkg
_log.info(f"定位设备 {serial}: {'; '.join(msgs)}")
return jsonify({"ok": True, "msg": f"{serial} {';'.join(msgs)}"})
except Exception as e:
return jsonify({"ok": False, "error": f"定位失败: {e}"}), 502
@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():
# 阶段 1:调度已不占 STF,此处仅清理迁移期遗留的 STF 占用记录
# (阶段 3 摘除 STF 后此接口改为清理本实例 _running/worker 状态)
try:
released, failed = stf.release_all_mine()
except Exception:
released, failed = [], []
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,前端用 <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/devices/<serial>/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/<action_id>", 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/<action_id>", 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 已连接的设备;设备池里在线但未本地连接的设备
在这里补全,并先做一次轻量 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:
online = set(device_pool.list_online())
pool = [s for s in device_pool.list_configured() if s in online]
except Exception:
pool = []
added = 0
for serial in pool:
if not serial or serial in seen:
continue
# 轻量连接(单次尝试,不重试不 kill-server);连不上也照样列出,
# 前端抓取时会报明确错误,不影响其他设备的连接
adb_connect_light(serial)
data.append({
"serial": serial,
"model": "",
"product": "",
"name": serial,
"status": "device",
"enabled": True,
})
seen.add(serial)
added += 1
if added:
_log.info(f"抓取元素设备列表补全 {added} 台在线设备: {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,前端用 <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:远程看屏(阶段 2) ==================
# MJPEG 实时画面流 + u2 触控注入。数据源 u2(atx-agent minicap 截图,单帧 ~200ms),
# 零新依赖(原计划 ws-scrcpy 不在 npm 分发,GitHub 下载在国内不可靠,改自建)。
# 坐标约定:前端以 <img> 显示尺寸归一化后映射到设备原生分辨率(naturalWidth/Height)。
_SCREEN_JPEG_QUALITY = 65
_SCREEN_FRAME_GAP = 0.05 # 帧间最小间隔(秒),防止空转烧 CPU
@app.route("/api/screen/stream")
@perm_required(PERM_DEVICES)
def api_screen_stream():
"""远程看屏:MJPEG 实时画面流(multipart/x-mixed-replace)。
?serial=xxx 设备 serial
浏览器 <img> 直接渲染;客户端断开(GeneratorExit)自动停止,不占资源。
设备离线/atx-agent 无响应时流自然结束,前端提示重新连接。
"""
serial = request.args.get("serial", "")
if not serial:
return jsonify({"ok": False, "error": "缺少 serial"}), 400
def generate():
import io
import uiautomator2 as u2
try:
d = u2.connect(serial)
except Exception as e:
_log.warning(f"远程看屏 {serial} u2 连接失败: {e}")
return
while True:
try:
img = d.screenshot()
if img is None:
break
if img.mode != "RGB":
img = img.convert("RGB")
buf = io.BytesIO()
img.save(buf, "JPEG", quality=_SCREEN_JPEG_QUALITY)
yield (b"--frame\r\nContent-Type: image/jpeg\r\n\r\n"
+ buf.getvalue() + b"\r\n")
except GeneratorExit:
break
except Exception as e:
_log.debug(f"远程看屏 {serial} 流中断: {e}")
break
time.sleep(_SCREEN_FRAME_GAP)
return Response(generate(),
mimetype="multipart/x-mixed-replace; boundary=frame")
def _screen_get_device(serial):
"""远程看屏用的 u2 连接(连接失败抛异常由调用方转 503)。"""
import uiautomator2 as u2
return u2.connect(serial)
@app.route("/api/screen/tap", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_screen_tap():
"""点击:{serial, x, y}(设备原生分辨率坐标)。"""
data = request.json or {}
serial, x, y = data.get("serial", ""), data.get("x"), data.get("y")
if not serial or x is None or y is None:
return jsonify({"ok": False, "error": "缺少 serial/x/y"}), 400
try:
_screen_get_device(serial).click(int(x), int(y))
return jsonify({"ok": True})
except Exception as e:
return jsonify({"ok": False, "error": f"点击失败: {e}"}), 503
@app.route("/api/screen/swipe", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_screen_swipe():
"""滑动:{serial, x1, y1, x2, y2, duration?}(设备原生分辨率坐标)。"""
data = request.json or {}
serial = data.get("serial", "")
x1, y1, x2, y2 = (data.get(k) for k in ("x1", "y1", "x2", "y2"))
if not serial or None in (x1, y1, x2, y2):
return jsonify({"ok": False, "error": "缺少 serial/x1/y1/x2/y2"}), 400
try:
_screen_get_device(serial).swipe(int(x1), int(y1), int(x2), int(y2),
duration=float(data.get("duration", 0.2)))
return jsonify({"ok": True})
except Exception as e:
return jsonify({"ok": False, "error": f"滑动失败: {e}"}), 503
_SCREEN_KEYS = {"back", "home", "recent", "menu", "power", "volume_up",
"volume_down", "enter", "delete", "search", "camera"}
@app.route("/api/screen/key", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_screen_key():
"""按键:{serial, key}(back/home/recent/menu/power 等,见 _SCREEN_KEYS)。"""
data = request.json or {}
serial, key = data.get("serial", ""), (data.get("key") or "").strip().lower()
if not serial or key not in _SCREEN_KEYS:
return jsonify({"ok": False, "error": "缺少 serial 或不支持的按键"}), 400
try:
_screen_get_device(serial).press(key)
return jsonify({"ok": True})
except Exception as e:
return jsonify({"ok": False, "error": f"按键失败: {e}"}), 503
@app.route("/api/screen/text", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_screen_text():
"""输入文字:{serial, text}(需焦点在输入框,u2 send_keys)。"""
data = request.json or {}
serial, text = data.get("serial", ""), (data.get("text") or "").strip()
if not serial or not text:
return jsonify({"ok": False, "error": "缺少 serial/text"}), 400
try:
_screen_get_device(serial).send_keys(text)
return jsonify({"ok": True})
except Exception as e:
return jsonify({"ok": False, "error": f"输入失败: {e}"}), 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/<apk_id>", 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():
"""可安装设备列表(应用管理):设备池在线设备 + 本机 adb 设备(含 USB 有线连接)。
source: pool=设备池 / usb=本机有线 adb / adb=本机网络 adb。
"""
devices, seen = [], set()
try:
online = set(device_pool.list_online())
for serial in device_pool.list_configured():
if serial in online:
devices.append({"serial": serial, "model": "", "source": "pool"})
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)
def _merged_device_list():
"""本地 adb(含 USB)+ 设备池的合并列表([{serial, state}])。
供维护终端/剪贴板注入/应用版本管理等设备选择场景共用。
状态:device/offline=本机 adb 实际状态;pool=设备池已配置但本机未连接;
stf_not_ready=STF agent 未就绪(迁移期保留,阶段 3 摘除 STF 后删除)。
"""
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 serial in device_pool.list_configured():
if serial not in states:
states[serial] = "pool"
except Exception:
pass
# 迁移期:STF 里 adb 可达但 agent 未就绪的设备仍列出,供维护终端修复 agent
try:
for d in stf.list_all_devices():
serial = d.get("serial", "")
if serial and d.get("present") and not d.get("ready") and serial not in states:
states[serial] = "stf_not_ready"
except Exception:
pass
return [{"serial": s, "state": st} for s, st in sorted(states.items())]
@app.route("/api/adb/devices")
@admin_required
def api_adb_devices():
"""维护终端设备列表(仅管理员):本地 adb 已连接 + STF 在线设备池
(含 agent 未就绪的待修复设备,state=stf_not_ready)。
供终端设备选择器使用——选中后自动附加 `-s <serial>`,
离线/未连接的设备也可选,配合"重连设备"按钮恢复。
"""
return jsonify({"ok": True, "devices": _merged_device_list()})
@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/tools/clipboard/set", methods=["POST"])
@admin_required
def api_tools_clipboard_set():
"""工具-剪贴板注入:把指定文字写入一台或多台设备的剪贴板。
请求: {"serials": ["100.100.10.11:5555", ...], "text": "要注入的文字"}
设备来源与维护终端一致:本地 adb 已连接设备(含 USB)+ STF 在线池。
实现:u2 jsonrpc.setClipboard(实测 cmd clipboard 在 MIUI 上不存在)。
"""
data = request.json or {}
serials = data.get("serials") or []
text = (data.get("text") or "").strip()
if not isinstance(serials, list) or not serials:
return jsonify({"ok": False, "error": "未选择设备"}), 400
if not text:
return jsonify({"ok": False, "error": "注入内容不能为空"}), 400
results = {}
for serial in serials:
ok, msg = _set_device_clipboard(serial, text)
results[serial] = {"ok": ok, "msg": msg}
ok_count = sum(1 for v in results.values() if v["ok"])
fail_count = len(results) - ok_count
return jsonify({"ok": True, "results": results,
"ok_count": ok_count, "fail_count": fail_count,
"error": None if ok_count == len(results) else f"{fail_count} 台设备注入失败"})
def _set_device_clipboard(serial, text):
"""向单台设备注入剪贴板文字(u2 jsonrpc.setClipboard)。
- IP:port 设备:先 adb connect(已连接自动跳过;绝不 disconnect,红线)
- USB 设备(无冒号):u2 按 serial 直连,首次会自动推送 atx-agent
- 连接/注入都带超时保护,避免 atx-agent 无响应时挂住请求
"""
import uiautomator2 as u2
if ":" in serial:
try:
adb_connect(serial)
except Exception as e:
return False, f"adb 连接失败: {e}"
d = None
try:
with ThreadPoolExecutor(max_workers=1) as pool:
d = pool.submit(u2.connect, serial).result(timeout=30)
with ThreadPoolExecutor(max_workers=1) as pool:
pool.submit(d.set_clipboard, text).result(timeout=15)
return True, "已注入"
except FuturesTimeout:
return False, "连接/注入超时(atx-agent 可能无响应)"
except Exception as e:
return False, f"{type(e).__name__}: {str(e)[:120]}"
@app.route("/api/stf/devmgmt/status")
@admin_required
def api_stf_devmgmt_status():
"""STF 设备管理:脚本配置的 IP + 220 adb 实际连接状态。"""
try:
st = stf_device_mgmt.status()
return jsonify({"ok": True, **st})
except StfDevError as e:
return jsonify({"ok": False, "error": str(e)}), 502
@app.route("/api/stf/devmgmt/add", methods=["POST"])
@admin_required
def api_stf_devmgmt_add():
"""STF 设备管理:添加设备(写入脚本 + 立即 connect)。
迁移期桥接:同步加入本地设备池(SQLite),保证新设备可被调度。
阶段 3 摘除 STF 后此接口随设备池管理面板退役。
"""
ip = (request.json or {}).get("ip", "").strip()
if not ip:
return jsonify({"ok": False, "error": "请输入设备 IP"}), 400
try:
msgs = stf_device_mgmt.add_device(ip)
except StfDevError as e:
return jsonify({"ok": False, "error": str(e)}), 502
# 桥接:同步本地设备池(serial 统一 IP:5555)
try:
device_pool.add_device(f"{ip}:5555" if ":" not in ip else ip)
except Exception as e:
_log.warning(f"STF 添加设备桥接设备池失败: {e}")
return jsonify({"ok": True, "msgs": msgs})
@app.route("/api/stf/devmgmt/remove", methods=["POST"])
@admin_required
def api_stf_devmgmt_remove():
"""STF 设备管理:移除设备(脚本删除 + disconnect)。
迁移期桥接:同步从本地设备池删除。阶段 3 此接口退役。
"""
ip = (request.json or {}).get("ip", "").strip()
if not ip:
return jsonify({"ok": False, "error": "请输入设备 IP"}), 400
try:
msgs = stf_device_mgmt.remove_device(ip)
except StfDevError as e:
return jsonify({"ok": False, "error": str(e)}), 502
try:
device_pool.remove_device(f"{ip}:5555" if ":" not in ip else ip)
except Exception:
pass
return jsonify({"ok": True, "msgs": msgs})
def _app_ver_on_device(serial, pkg):
"""单设备包版本查询:pm path 检查安装 → dumpsys 取 versionName/versionCode。
IP:port 设备先 adb connect(已连接自动跳过;绝不 disconnect,红线)。
查询是只读 shell 命令,无需 adb server 锁(锁只保护 connect/kill-server 类操作)。
"""
if ":" in serial:
try:
adb_connect(serial)
except Exception as e:
return {"installed": False, "error": f"adb 连接失败: {e}"}
try:
r = subprocess.run([ADB_PATH, "-s", serial, "shell", "pm", "path", pkg],
capture_output=True, timeout=15)
if b"package:" not in (r.stdout or b""):
return {"installed": False, "version_name": "", "version_code": ""}
r2 = subprocess.run([ADB_PATH, "-s", serial, "shell", "dumpsys", "package", pkg],
capture_output=True, timeout=20)
out2 = (r2.stdout or b"").decode("utf-8", errors="replace")
vm = re.search(r"versionName=(\S+)", out2)
vc = re.search(r"versionCode=(\d+)", out2)
return {"installed": True,
"version_name": vm.group(1) if vm else "",
"version_code": vc.group(1) if vc else ""}
except subprocess.TimeoutExpired:
return {"installed": False, "error": "查询超时"}
except Exception as e:
return {"installed": False, "error": str(e)}
@app.route("/api/tools/appver", methods=["POST"])
@admin_required
def api_tools_appver():
"""工具-应用版本管理:查询所有设备上指定包名的安装情况与版本号。
请求: {"pkg": "com.ss.android.ugc.aweme"}
设备来源:本地 adb(含 USB)+ STF 在线池,并发查询(最多 10 台同时)。
"""
pkg = (request.json or {}).get("pkg", "").strip()
if not re.match(r"^[A-Za-z0-9_.]+$", pkg or ""):
return jsonify({"ok": False, "error": "包名格式不正确(仅字母/数字/._)"}), 400
devices = _merged_device_list()
if not devices:
return jsonify({"ok": False, "error": "无可用设备"}), 404
results = {}
with ThreadPoolExecutor(max_workers=min(10, len(devices))) as pool:
futures = {pool.submit(_app_ver_on_device, d["serial"], pkg): d for d in devices}
for fut in as_completed(futures, timeout=90):
d = futures[fut]
try:
results[d["serial"]] = fut.result()
except Exception as e:
results[d["serial"]] = {"installed": False, "error": str(e)}
fail = sum(1 for v in results.values() if v.get("error"))
return jsonify({"ok": True, "results": results, "fail": fail,
"total": len(results)})
@app.route("/api/stf/devmgmt/delete", methods=["POST"])
@admin_required
def api_stf_devmgmt_delete():
"""STF 设备管理:彻底删除(脚本移除 + disconnect + 清 STF 池记录)。
迁移期桥接:同步从本地设备池删除。阶段 3 此接口退役。
"""
ip = (request.json or {}).get("ip", "").strip()
if not ip:
return jsonify({"ok": False, "error": "请输入设备 IP"}), 400
try:
msgs = stf_device_mgmt.delete_device(ip, stf=stf)
except StfDevError as e:
return jsonify({"ok": False, "error": str(e)}), 502
try:
device_pool.remove_device(f"{ip}:5555" if ":" not in ip else ip)
except Exception:
pass
return jsonify({"ok": True, "msgs": msgs})
@app.route("/api/stf/devmgmt/reconnect", methods=["POST"])
@admin_required
def api_stf_devmgmt_reconnect():
"""STF 设备管理:一键重连(手动运行 220 的 connect_devices.sh)。"""
try:
ok, summary = stf_device_mgmt.reconnect_all()
return jsonify({"ok": True, "output": summary})
except StfDevError as e:
return jsonify({"ok": False, "error": str(e)}), 502
# ================== API:设备池管理(本地 SQLite 清单,新增设备的规范化入口) ==================
@app.route("/api/devices/pool", methods=["GET"])
@perm_required(PERM_DEVICES)
def api_devices_pool_list():
"""设备池清单(SQLite,含实时在线状态)。"""
try:
online = set(device_pool.list_online())
except Exception:
online = set()
rows = device_pool.list_devices()
for r in rows:
r["online"] = r["serial"] in online
return jsonify({"ok": True, "devices": rows})
@app.route("/api/devices/pool/add", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_devices_pool_add():
"""添加/更新设备:{serial, name?, note?}。
IP:5555 设备添加后立即尝试 adb connect(不可达显示离线,不影响其他设备)。
"""
data = request.json or {}
serial = (data.get("serial") or "").strip()
if not serial:
return jsonify({"ok": False, "error": "请输入设备 serial(如 100.100.10.20:5555)"}), 400
is_new = device_pool.add_device(serial,
name=(data.get("name") or "").strip(),
note=(data.get("note") or "").strip())
msg = "已添加" if is_new else "已更新"
if ":" in serial:
try:
adb_connect(serial)
except Exception:
pass
_log.info(f"设备池管理: {msg} {serial}")
return jsonify({"ok": True, "msg": msg, "is_new": is_new})
@app.route("/api/devices/pool/remove", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_devices_pool_remove():
"""从设备池删除(不再参与调度;不影响其他系统)。"""
serial = (request.json or {}).get("serial", "").strip()
if not serial:
return jsonify({"ok": False, "error": "缺少 serial"}), 400
ok = device_pool.remove_device(serial)
if not ok:
return jsonify({"ok": False, "error": "设备不存在"}), 404
return jsonify({"ok": True, "msg": "已删除"})
@app.route("/api/devices/pool/toggle", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_devices_pool_toggle():
"""启用/停用设备(停用后不参与调度)。"""
data = request.json or {}
serial = (data.get("serial") or "").strip()
enabled = data.get("enabled")
if not serial or enabled is None:
return jsonify({"ok": False, "error": "缺少参数"}), 400
ok = device_pool.set_enabled(serial, bool(enabled))
if not ok:
return jsonify({"ok": False, "error": "设备不存在"}), 404
return jsonify({"ok": True, "msg": "已" + ("启用" if enabled else "停用")})
@app.route("/api/devices/pool/reconnect", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_devices_pool_reconnect():
"""一键重连:并发 adb connect 池内全部 IP:5555 设备(后台执行,不阻塞)。"""
def _run():
try:
serials = [s for s in device_pool.list_configured() if ":" in s]
except Exception:
return
if not serials:
return
from concurrent.futures import ThreadPoolExecutor, as_completed
with ThreadPoolExecutor(max_workers=10) as pool:
futures = {pool.submit(adb_connect, s): s for s in serials}
for _ in as_completed(futures):
pass
_log.info(f"设备池一键重连完成({len(serials)} 台)")
import threading
threading.Thread(target=_run, daemon=True).start()
return jsonify({"ok": True, "msg": "重连已启动(后台并发,约 10-20 秒)"})
@app.route("/api/stf/restart", methods=["POST"])
@admin_required
def api_stf_restart():
"""一键重启 STF Docker 容器(仅管理员)。
通过 SSH(config.STF_SSH_TARGET,默认 [email protected])执行
`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
try:
code, out, err = ssh_client.run(
f"docker restart {STF_DOCKER_CONTAINER} && "
f"docker ps --filter name={STF_DOCKER_CONTAINER} --format '{{{{.Names}}}}: {{{{.Status}}}}'",
timeout=60)
except SSHError as e:
return jsonify({"ok": False, "error": str(e)}), 502
if code != 0:
_log.error(f"STF 容器重启失败: {err or out}")
return jsonify({"ok": False,
"error": f"SSH 执行失败:{err or out}(请确认 SSH 配置(STF_SSH_TARGET / STF_SSH_PASSWORD)且容器名 {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/<device_id>", 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/<device_id>/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/<device_id>", 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()