refactor: web_server 模块化拆分(web/ 蓝图包 auth/monitor/tasks/admin/tools/devices/apks/tailscale + context/common,web_server 1796→193 行仅留装配与启动)+ 大屏优化(网格 5s 轮询、弹窗独占缩略图、流支持 q/fps)+ 三修复(MIUI 屏幕状态检测改 dumpsys power、补 ADB_PATH 导入修复一键亮屏、renderGrid 增量渲染不重载缩略图)

This commit is contained in:
2026-08-19 14:55:29 +08:00
parent e07b339a51
commit ce47a5b83d
16 changed files with 1922 additions and 1780 deletions
+12 -1
View File
@@ -73,7 +73,18 @@ python web_server.py
```
platform-tools/
├── config.py # 根配置(部署配置统一从 .env 读,模板见 .env.example)
├── web_server.py # Flask 入口(API + 登录 + 单页应用)
├── web_server.py # Flask 入口(app 装配 + 蓝图注册 + 启动,193 行)
├── web/ # Web 层蓝图包(按功能域拆分,路由都在这里)
│ ├── auth.py # 登录/CSRF/权限装饰器/页面路由(/、/wall)
│ ├── monitor.py # 状态/运行控制/设备操作/远程看屏
│ ├── tasks_api.py # 任务计划/分组/自定义动作/元素抓取
│ ├── admin_api.py # 用户管理/日志
│ ├── tools_api.py # adb 终端/剪贴板注入/应用版本
│ ├── devices_api.py # 设备池管理
│ ├── apks_api.py # 应用管理
│ ├── tailscale_api.py # Tailscale 管理
│ ├── common.py # 跨模块共享工具(合并设备列表/屏幕状态)
│ └── context.py # 共享对象注入(mgr/apk_mgr/device_pool)
├── requirements.txt # Python 依赖清单
│
├── core/ # 核心基础设施层
+30 -30
View File
@@ -10,8 +10,8 @@
| 层 | 路径 | 职责 |
|----|------|------|
| 配置层 | `config.py` | 项目根配置:STF 地址、adb 路径、web 端口等基础设施。**不放任务参数** |
| 核心层 | `core/` | 框架运行时:日志、STF 客户端、adb 操作、Worker 基类、任务管理器、u2 辅助、Action 基类 |
| 配置层 | `config.py` | 项目根配置:adb 路径、web 端口、USB 远程 adb server 等基础设施。**不放任务参数** |
| 核心层 | `core/` | 框架运行时:日志、设备池、adb 操作、Worker 基类、任务管理器、u2 辅助、Action 基类 |
| 任务层 | `tasks/` | 每个 App 一个子包,自包含 `task.py` + `actions/`,互不依赖 |
| 前端层 | `templates/admin/` | 单页应用(纯 HTML+CSS+JS,无框架) |
| 数据层 | `data/` | SQLite 持久化 |
@@ -39,7 +39,7 @@
│ JSON API │ ┌─────────────────┐
▼ ┌──────────────┐ │ 设备 adb │
┌──────────────┐ │ 看门狗监控 │ │ (IP:5555 直连 │
│ web_server │ └──────────────┘ │ / USB 远程) │
│ web/ 蓝图包 │ └──────────────┘ │ / USB 远程) │
│ (Flask API) │ └─────────────────┘
└──────────────┘
│
@@ -137,7 +137,7 @@ run() 主循环(不要重写):
**并发控制**:同一 serial 同时只允许一个 worker,避免冲突。
**状态缓存**:`get_status()` 带 5 秒缓存,STF 请求慢时不阻塞前端。
**状态缓存**:`get_status()` 带 5 秒缓存,避免每次 /api/status 都查库/adb 阻塞前端。
### 3.4 前台 App 扫描器(`_ForegroundScanner`)
@@ -145,12 +145,11 @@ run() 主循环(不要重写):
| 设备状态 | 处理方式 | 是否打扰 |
|---------|---------|---------|
| worker 运行中 | 复用已有 ADB 连接查询 | 否 |
| 自己占用无 worker | STF remoteConnect 隧道查询 | 否 |
| 完全空闲 | 返回"空闲"(不 adb connect,避免 STF 误判离线) | 否 |
| 他人占用 | 标记"(他人占用)" | 否 |
| worker 运行中(IP:5555) | 复用已有 ADB 连接查询 | 否 |
| worker 运行中(USB) | 经 220 远程 adb server 查询 | 否 |
| 完全空闲 | 返回"空闲"(不主动 connect) | 否 |
> **为什么不扫描空闲设备的前台 App**:STF provider 内部通过 IP:5555 维持 adb 连接。外部 adb connect/disconnect 会让 adb server 断开该 transport,连带 STF provider 的连接断开,STF 误判设备 offline 并触发重连。
> **为什么不扫描空闲设备的前台 App**:IP:5555 的 adb transport 是共享的(历史与 STF provider 共用),外部 connect/disconnect 会扰动共享连接,遵守既有红线。
### 3.5 ADB 操作(`core/adb_helper.py`)
@@ -158,8 +157,8 @@ run() 主循环(不要重写):
**铁律:绝不 kill-server**:
- `adb kill-server` 会断开所有设备的 adb transport
- 导致 STF provider 对全部设备误判离线并触发重连
- 影响所有运行中的任务
- 全部设备连接被重建,影响所有运行中的任务
- 同理绝不 disconnect IP:5555(共享 transport 红线,见 DEVELOPMENT.md)
| 函数 | 说明 |
|------|------|
@@ -191,9 +190,9 @@ SQLAlchemy 模型,存于 `data/users.db`:
### 权限模型
- `User.perms` 存业务权限位 JSON 数组(`tasks`/`devices`/`apks`/`logs`),`is_admin=true` 拥有全部权限(`has_perm` 短路)
- 后端统一用 `@perm_required(PERM_X)` / `@admin_required` 装饰器拦截(web_server.py),无权限返回 403;
查看类 GET 接口只要求登录;用户管理、维护接口(STF 重启 `/api/stf/restart`、adb 终端 `/api/adb/cmd`)强制 `admin_required`
- adb 终端安全红线:拒绝 `kill-server` / `disconnect`(STF provider 共享 adb transport,断开会误判全设备离线)
- 后端统一用 `@perm_required(PERM_X)` / `@admin_required` 装饰器拦截(web/auth.py),无权限返回 403;
查看类 GET 接口只要求登录;用户管理、adb 终端(`/api/adb/cmd`)强制 `admin_required`
- adb 终端安全红线:拒绝 `kill-server` / `disconnect`(共享 adb transport)
- 前端 `loadMe()` 拉取 `/api/me`,用 `data-perm` 属性隐藏无权限的 tab/按钮,行内按钮用 `_can(perm)` 判断
- 安全兜底:**前端隐藏只是 UX,权限强制在后端**;新增路由时按"写操作必须带权限装饰器"的约定
@@ -201,10 +200,10 @@ SQLAlchemy 模型,存于 `data/users.db`:
APK 上传/解析/批量安装。
**设备连接策略(直连,绕过 STF)**:
**设备连接策略(直连)**:
- 直接 `adb connect serial`(serial 是 IP:5555)
- 不经过 STF occupy/release,避免 STF release 触发 agent 清理卸载 app
- 安装后不主动 disconnect(STF provider 共享该 adb transport)
- 不经过占用/释放(单实例互斥由调度器内存锁保证)
- 安装后不主动 disconnect(共享 adb transport 红线)
**安装流程**:
1. 上传 APK → 保存到 `data/apks/` → pyaxmlparser 解析包名/版本 → 入库
@@ -232,23 +231,24 @@ APK 上传/解析/批量安装。
- 引擎懒加载单例 + 并发加锁(识别约 0.2-0.5s/次)
- 返回坐标约定:像素、原点左上
### 3.10 STF 设备管理(`core/stf_device_mgmt.py`)
### 3.10 Web 层蓝图包(`web/`)
工具页"STF 设备管理":维护 220 上 OpenSTF 设备池。
路由按功能域拆分(模块化开发底线,便于定位问题):
背景:220 的 adb 跑在 Docker 容器(`STF_ADB_CONTAINER`,默认 `adb`)里,设备池由
`/mnt/data/openstf/connect_devices.sh`(`STF_SCRIPT_PATH`)的 `DEVICES` 数组维护,cron 每 5 分钟补连。
| 函数 | 说明 |
| 模块 | 职责 |
|------|------|
| `status()` | 脚本配置 IP + 220 adb 实际连接状态 |
| `add_device(ip)` | 用 220 的 python3 精确改脚本 DEVICES 块 + `docker exec adb adb connect` |
| `remove_device(ip)` | 脚本删条目 + `adb disconnect` |
| `auth.py` | 登录/登出/CSRF/权限装饰器/页面路由(/、/wall) |
| `monitor.py` | 状态/运行控制/设备操作/远程看屏(流+缩略图+触控) |
| `tasks_api.py` | 任务计划/分组/自定义动作/元素抓取/步骤测试 |
| `admin_api.py` | 用户管理/日志 |
| `tools_api.py` | adb 终端/剪贴板注入/应用版本 |
| `devices_api.py` | 设备池管理(增删停用/重连/型号采集) |
| `apks_api.py` | 应用管理 |
| `tailscale_api.py` | Tailscale 管理 |
| `common.py` | 跨模块共享(合并设备列表/屏幕状态) |
| `context.py` | 共享对象注入(mgr/apk_mgr/device_pool) |
关键点:
- 脚本编辑用 python3 行级增删(sed 处理多行数组易误伤)
- `adb connect` 到不可达 IP 会挂 40s+,220 侧用 `timeout 15` 兜底,超时保留在脚本由 cron 重试
- 断开的是 220(STF provider 侧)的 adb 连接,与本机任务直连的 adb 相互独立
`web_server.py` 只保留:app 创建、数据库初始化、蓝图注册、uiautodev 生命周期、启动(193 行)。
---
+73 -34
View File
@@ -239,50 +239,89 @@ function fmtProgress(p){
}
function shortIp(serial){return String(serial||'').replace(/:5555$/,'');}
function cardHtml(d, i){
const cls=cardCls(d);
const serial=esc(d.serial);
const name=esc(d.device_name||d.model||'设备');
const ip=esc(shortIp(d.serial));
const st=badgeText(d);
const stTag=_screenState[d.serial]==='off'?'<div class="screen-tag off">🌙</div>':'';
const action=d.current_action
?'<div class="action">'+(d.task_job?'<b>'+esc(d.task_job)+'</b> · ':'')+esc(d.current_action)+'</div>'
:(d.task_job?'<div class="action"><b>'+esc(d.task_job)+'</b></div>':'<div class="action"></div>');
return '<div class="card '+cls+'" data-serial="'+serial+'" style="animation-delay:'+(i*30)+'ms" onclick="openCtrl(\''+serial+'\')">'+
'<div class="thumb">'+
'<img data-thumb="'+serial+'" alt="">'+
'<div class="ph"><span class="ico">📱</span><span>'+serial+'</span></div>'+
'<div class="badge"><span class="dot"></span>'+st+'</div>'+stTag+
'</div>'+
'<div class="info">'+
'<div class="name">'+name+'</div>'+
'<div class="ip">'+ip+'</div>'+
'<div class="model">'+esc(d.model||'')+'</div>'+
action+fmtProgress(d.progress)+
'<div class="hint">点击进入设备操作</div>'+
'</div>'+
'</div>';
}
// 增量渲染:已有卡片只更新信息区(不重建 img,缩略图不闪不重载);增删设备才动 DOM
function renderGrid(devs){
const grid=document.getElementById('wall-grid');
if(!devs.length){
grid.innerHTML='<div class="empty-wall">设备池为空,请到「工具 → 设备池管理」添加设备</div>';
return;
}
grid.innerHTML=devs.map((d,i)=>{
const cls=cardCls(d);
const serial=esc(d.serial);
const name=esc(d.device_name||d.model||'设备');
const ip=esc(shortIp(d.serial));
const st=badgeText(d);
const stTag=_screenState[d.serial]==='off'?'<div class="screen-tag off">🌙</div>':'';
const action=d.current_action
?'<div class="action">'+(d.task_job?'<b>'+esc(d.task_job)+'</b> · ':'')+esc(d.current_action)+'</div>'
:(d.task_job?'<div class="action"><b>'+esc(d.task_job)+'</b></div>':'<div class="action"></div>');
return '<div class="card '+cls+'" data-serial="'+serial+'" style="animation-delay:'+(i*30)+'ms" onclick="openCtrl(\''+serial+'\')">'+
'<div class="thumb">'+
'<img data-thumb="'+serial+'" alt="">'+
'<div class="ph"><span class="ico">📱</span><span>'+serial+'</span></div>'+
'<div class="badge"><span class="dot"></span>'+st+'</div>'+stTag+
'</div>'+
'<div class="info">'+
'<div class="name">'+name+'</div>'+
'<div class="ip">'+ip+'</div>'+
'<div class="model">'+esc(d.model||'')+'</div>'+
action+fmtProgress(d.progress)+
'<div class="hint">点击进入设备操作</div>'+
'</div>'+
'</div>';
}).join('');
// 离线卡片显示占位
grid.querySelectorAll('.card').forEach(c=>{
const d=devs.find(x=>x.serial===c.dataset.serial);
const img=c.querySelector('img[data-thumb]');
const ph=c.querySelector('.ph');
if(d&&d.present&&img){img.style.display='block';ph.style.display='none';}
else if(ph){ph.style.display='flex';}
const existing={};
grid.querySelectorAll('.card').forEach(c=>{existing[c.dataset.serial]=c;});
// 1. 移除已不存在的卡片
Object.keys(existing).forEach(s=>{
if(!devs.find(d=>d.serial===s))existing[s].remove();
});
// 2. 更新/新建
devs.forEach((d,i)=>{
const serial=d.serial;
const card=existing[serial];
if(card){
// 已有卡片:只更新状态类信息(badge/class/信息区),不动 img
card.className='card '+cardCls(d);
const badge=card.querySelector('.badge');
const bt=badgeText(d);
if(badge&&badge.textContent!==bt)badge.innerHTML='<span class="dot"></span>'+bt;
// 熄屏角标
const tag=card.querySelector('.screen-tag');
if(_screenState[serial]==='off'&&!tag)card.querySelector('.thumb').insertAdjacentHTML('beforeend','<div class="screen-tag off">🌙</div>');
else if(_screenState[serial]!=='off'&&tag)tag.remove();
// 信息区
card.querySelector('.name').textContent=d.device_name||d.model||'设备';
card.querySelector('.ip').textContent=shortIp(serial);
card.querySelector('.model').textContent=d.model||'';
const action=d.current_action
?(d.task_job?'<b>'+esc(d.task_job)+'</b> · ':'')+esc(d.current_action)
:(d.task_job?'<b>'+esc(d.task_job)+'</b>':'');
card.querySelector('.action').innerHTML=action;
// 进度
const progWrap=card.querySelector('.prog-wrap');
if(progWrap)progWrap.innerHTML=fmtProgress(d.progress);
else if(d.progress&&d.progress.total)card.querySelector('.info').insertAdjacentHTML('beforeend','<div class="prog-wrap">'+fmtProgress(d.progress)+'</div>');
// 离线占位切换
const img=card.querySelector('img[data-thumb]');
const ph=card.querySelector('.ph');
if(d.present){if(img)img.style.display='block';if(ph)ph.style.display='none';}
else{if(img)img.style.display='none';if(ph)ph.style.display='flex';}
}else{
grid.insertAdjacentHTML('beforeend',cardHtml(d,i));
const nc=grid.lastElementChild;
if(!d.present)nc.querySelector('.ph').style.display='flex';
}
});
refreshThumbs();
}
// 缩略图轮询(2.5s):fetch 读取 X-Screen-State 头 → blob 喂 img
async function refreshThumbs(){
// 进入设备操作弹窗后:停止全部缩略图获取,资源全部给当前设备的高清流
if(_ctrlSerial)return;
if(document.visibilityState!=='visible')return;
const now=Date.now();
const imgs=[...document.querySelectorAll('img[data-thumb]')];
@@ -343,7 +382,7 @@ function closeCtrl(){
function startCtrlStream(){
if(!_ctrlSerial)return;
const img=document.getElementById('ctrl-img');
img.src='/api/screen/stream?serial='+encodeURIComponent(_ctrlSerial)+'&t='+Date.now();
img.src='/api/screen/stream?serial='+encodeURIComponent(_ctrlSerial)+'&q=85&fps=10&t='+Date.now(); // 弹窗高清高帧
}
function updateCtrlHead(){
const d=_devices.find(x=>x.serial===_ctrlSerial);
@@ -414,7 +453,7 @@ function ctrlSendText(){
})();
// 定时器
setInterval(()=>{if(document.visibilityState==='visible')refreshThumbs();},2500);
setInterval(()=>{if(document.visibilityState==='visible')refreshThumbs();},5000); // 网格低帧率(省资源)
setInterval(()=>{if(document.visibilityState==='visible')refreshStatus();},5000);
// 初始化
+28
View File
@@ -0,0 +1,28 @@
"""Web 层蓝图包:web_server.py 只保留 app 初始化与装配,路由按功能域拆分。
模块划分(按功能定位问题):
auth — 登录/登出/CSRF/权限装饰器/页面路由(/、/wall)
monitor — 状态/运行控制/设备操作/远程看屏
tasks — 任务计划/分组/自定义动作/元素抓取/步骤测试
admin — 用户管理/日志
tools — adb 终端/剪贴板注入/应用版本
devices — 设备池管理
apks — 应用管理
tailscale — Tailscale 管理
"""
from flask import Blueprint
def register_blueprints(app):
"""注册全部蓝图(url_prefix 为空,路由路径保持原样)。"""
from .auth import bp as auth_bp
from .monitor import bp as monitor_bp
from .tasks_api import bp as tasks_bp
from .admin_api import bp as admin_bp
from .tools_api import bp as tools_bp
from .devices_api import bp as devices_bp
from .apks_api import bp as apks_bp
from .tailscale_api import bp as tailscale_bp
for bp in (auth_bp, monitor_bp, tasks_bp, admin_bp, tools_bp,
devices_bp, apks_bp, tailscale_bp):
app.register_blueprint(bp)
+101
View File
@@ -0,0 +1,101 @@
"""管理域 API:用户管理/日志查看。"""
import os
import time
from flask import Blueprint, jsonify, request
from flask_login import current_user
from flask import session
from core.logger import _LOG_DIR, _MODULE_FILES
from core.models import db, User
from web import context
from web.auth import (admin_required, perm_required, _validate_perms,
PERM_LOGS)
bp = Blueprint("admin", __name__)
@bp.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})
@bp.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": "用户已创建"})
@bp.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": "用户已更新"})
@bp.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:日志查看 ==================
@bp.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:运行控制 ==================
+123
View File
@@ -0,0 +1,123 @@
"""应用管理 API(APK 上传/安装/删除/设备已装应用)。"""
import re
import shlex
import subprocess
from flask import Blueprint, jsonify, request
from flask_login import login_required
from core import device_pool
from core.adb_helper import adb_connect_light, _adb
from web import context
from web.auth import perm_required, PERM_APKS
from web.common import _merged_device_list
bp = Blueprint("apks", __name__)
@bp.route("/api/apks")
@login_required
def api_apks_list():
"""列出所有已上传的 APK。"""
return jsonify({"ok": True, "apks": context.apk_mgr.list_all()})
@bp.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 = context.apk_mgr.upload(file)
if info:
return jsonify({"ok": True, "apk": info,
"msg": f"上传成功: {info['display_name']}"})
return jsonify({"ok": False, "error": "上传失败,请检查文件格式"}), 500
@bp.route("/api/apks/<apk_id>", methods=["DELETE"])
@perm_required(PERM_APKS)
def api_apks_delete(apk_id):
"""删除 APK 文件和记录。"""
ok, msg = context.apk_mgr.delete(apk_id)
if ok:
return jsonify({"ok": True, "msg": msg})
return jsonify({"ok": False, "error": msg}), 400
@bp.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})
@bp.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 = context.apk_mgr.install(apk_id, serials)
if ok:
return jsonify({"ok": True, "msg": msg})
return jsonify({"ok": False, "error": msg}), 400
@bp.route("/api/apks/install/status")
@login_required
def api_apks_install_status():
"""获取安装任务实时状态。"""
status = context.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
return [{"serial": s, "state": st} for s, st in sorted(states.items())]
+157
View File
@@ -0,0 +1,157 @@
"""认证与页面路由:登录/登出/CSRF/权限装饰器/页面(/、/wall)。
CSRF 的 before_request 由 web_server 注册(app 级),本模块提供实现。
"""
import functools
import secrets
from flask import (Blueprint, jsonify, request, redirect, url_for,
render_template, make_response, session)
from flask_login import (login_user, logout_user, login_required, current_user)
from core.logger import get_logger
from core.models import User
_log = get_logger("web")
bp = Blueprint("auth", __name__)
# ================== 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"]
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
@bp.route("/api/csrf")
@login_required
def api_csrf():
"""获取 CSRF token(前端非 GET 请求需携带 X-CSRF-Token 头)。"""
return jsonify({"ok": True, "token": _csrf_token()})
# ================== 权限控制 ==================
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
@bp.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(),
}})
# ================== 页面路由 ==================
@bp.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
@bp.route("/wall")
@login_required
def wall():
"""监控大屏(独立全屏页面,供挂墙/电视展示)。
深色控制室风格;设备缩略图按需轮询(每设备 ~5s 一帧),
20 台设备整体开销约 0.1 核 CPU + 50KB/s,普通电脑无压力。
"""
resp = make_response(render_template("admin/wall.html"))
resp.headers["Cache-Control"] = "no-store, no-cache, must-revalidate, max-age=0"
return resp
@bp.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("auth.index"))
return render_template("admin/login.html", error="用户名或密码错误")
return render_template("admin/login.html", error=None)
@bp.route("/logout")
@login_required
def logout():
_log.info(f"用户 {current_user.username} 登出")
logout_user()
return redirect(url_for("auth.login"))
+46
View File
@@ -0,0 +1,46 @@
"""跨模块共享工具(web 层)。"""
from core import device_pool
from core.adb_helper import _adb
from core.logger import get_logger
_log = get_logger("web")
def _merged_device_list():
"""本地 adb(含 USB)+ 设备池的合并列表([{serial, state}])。
供维护终端/剪贴板注入/应用版本管理等设备选择场景共用。
状态:device/offline=本机 adb 实际状态;pool=设备池已配置但本机未连接。
"""
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
return [{"serial": s, "state": st} for s, st in sorted(states.items())]
def _device_screen_state(serial):
"""设备屏幕状态:on/off/unknown(dumpsys power 只读查询)。
用 mWakefulness 判断(Asleep=熄屏):MIUI 的 dumpsys display 没有
mScreenState 字段,power 字段全机型存在。USB 设备查询失败返回 unknown。
"""
try:
if ":" not in serial:
return "unknown"
out = _adb("-s", serial, "shell", "dumpsys", "power") or ""
import re
m = re.search(r"mWakefulness=(\w+)", out)
if m:
return "off" if m.group(1).lower() == "asleep" else "on"
except Exception:
pass
return "unknown"
+12
View File
@@ -0,0 +1,12 @@
"""共享上下文:web_server 启动时注入全局对象,各蓝图模块从这里取。
避免蓝图模块与 web_server 循环 import。
"""
mgr = None # core.task_manager.TaskManager
apk_mgr = None # core.apk_manager.ApkManager
device_pool = None # core.device_pool
def init(mgr_, apk_mgr_, device_pool_):
global mgr, apk_mgr, device_pool
mgr, apk_mgr, device_pool = mgr_, apk_mgr_, device_pool_
+119
View File
@@ -0,0 +1,119 @@
"""设备池管理 API(本地 SQLite 清单)。"""
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
from flask import Blueprint, jsonify, request
from core import device_pool
from core.adb_helper import adb_connect
from core.logger import get_logger
from web.auth import perm_required, PERM_DEVICES
_log = get_logger("web")
bp = Blueprint("devices", __name__)
@bp.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})
@bp.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
# 后台采集型号(不阻塞添加响应)
import threading
threading.Thread(target=device_pool.refresh_model, args=(serial,),
daemon=True).start()
_log.info(f"设备池管理: {msg} {serial}")
return jsonify({"ok": True, "msg": msg, "is_new": is_new})
@bp.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": "已删除"})
@bp.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 "停用")})
@bp.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)} 台)")
# 重连后顺手刷新型号
try:
device_pool.refresh_all_models()
except Exception:
pass
import threading
threading.Thread(target=_run, daemon=True).start()
return jsonify({"ok": True, "msg": "重连已启动(后台并发,约 10-20 秒)"})
@bp.route("/api/devices/pool/refresh_models", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_devices_pool_refresh_models():
"""批量采集池内在线设备的型号(后台执行,不阻塞)。"""
import threading
threading.Thread(target=device_pool.refresh_all_models, daemon=True).start()
return jsonify({"ok": True, "msg": "型号采集已启动(后台并发,约 10 秒)"})
# ================== API:Tailscale 管理(仅管理员) ==================
# 通过 Tailscale 官方 API v2 管理 tailnet 设备(列表/改名/授权/密钥不过期/删除/生成 auth key)。
# 设备 IP 由 tailnet 自动分配,API 无法修改,列表只读展示。
-94
View File
@@ -1,94 +0,0 @@
<!DOCTYPE html>
<html lang="zh-CN"><head><meta charset="UTF-8"><title>设备监控大屏</title>
<style>
*{box-sizing:border-box} body{font-family:-apple-system,"Segoe UI",Roboto,sans-serif;margin:0;background:#f0f2f5;color:#2c3e50}
.header{background:#fff;padding:14px 24px;box-shadow:0 1px 4px rgba(0,0,0,.08);display:flex;align-items:center;gap:16px}
.header h1{font-size:18px;margin:0} .header .ts{color:#95a5a6;font-size:12px;margin-left:auto}
.header a{color:#1890ff;text-decoration:none;font-size:13px}
.header a:hover{text-decoration:underline}
.content{padding:20px 24px;max-width:1400px;margin:0 auto}
.bar{display:flex;gap:8px;margin-bottom:16px;align-items:center;flex-wrap:wrap}
button{padding:7px 16px;border:1px solid #d9d9d9;border-radius:4px;cursor:pointer;font-size:13px;background:#fff;color:#595959}
button:hover{border-color:#40a9ff;color:#40a9ff} button.primary{background:#1890ff;color:#fff;border-color:#1890ff}
button.primary:hover{background:#40a9ff} button.danger{color:#ff4d4f;border-color:#ff4d4f}
button.danger:hover{background:#ff4d4f;color:#fff} button.sm{padding:3px 10px;font-size:12px}
table{width:100%;border-collapse:collapse;background:#fff;border-radius:6px;overflow:hidden;box-shadow:0 1px 3px rgba(0,0,0,.08)}
th,td{padding:9px 10px;text-align:left;border-bottom:1px solid #f0f0f0;font-size:13px}
th{background:#fafafa;font-weight:600;color:#595959} tr:hover td{background:#fafafa}
.badge{padding:2px 8px;border-radius:10px;font-size:11px;color:#fff}
.b-running{background:#52c41a}.b-connecting{background:#faad14}.b-error{background:#ff4d4f}
.b-done{background:#1890ff}.b-idle{background:#bfbfbf;color:#595959}.b-released{background:#8c8c8c}.b-failed{background:#ff4d4f}
.ok{color:#52c41a}.no{color:#ff4d4f}.muted{color:#8c8c8c;font-size:12px}
.err{color:#ff4d4f;font-size:12px;max-width:180px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap}
.progress{width:70px;height:5px;background:#f0f0f0;border-radius:3px;overflow:hidden;display:inline-block;vertical-align:middle}
.progress>div{height:100%;background:#1890ff}
.tag{display:inline-block;background:#e6f7ff;color:#1890ff;padding:1px 8px;border-radius:3px;font-size:11px;margin:1px}
.kv{display:inline-block;background:#f5f5f5;border-radius:3px;padding:0 6px;margin:1px;font-size:11px;color:#595959}
.stats{display:flex;gap:16px;margin-bottom:16px}
.stat-card{background:#fff;padding:14px 20px;border-radius:6px;box-shadow:0 1px 3px rgba(0,0,0,.08);min-width:120px}
.stat-card .num{font-size:24px;font-weight:600;color:#1890ff}
.stat-card .lbl{font-size:12px;color:#8c8c8c;margin-top:2px}
</style></head><body>
<div class="header">
<h1>设备监控大屏</h1>
<span class="muted" id="devCount"></span>
<span class="ts" id="srvTime"></span>
<a href="/admin" target="_blank">管理后台 →</a>
</div>
<div class="content">
<div class="stats" id="stats"></div>
<div class="bar">
<button class="primary" onclick="doPost('/api/jobs/_quick/run')" title="快速启动:对全部空闲设备跑默认抖音养号">快速启动全部</button>
<button class="danger" onclick="doPost('/api/stop_all')">停止全部</button>
<button onclick="doPost('/api/release')">强制释放占用</button>
<button onclick="loadStatus()">刷新</button>
<span class="muted">每 5 秒自动刷新 · 分组/任务管理请到 <a href="/admin" target="_blank">管理后台</a></span>
</div>
<table><thead><tr>
<th>设备</th><th>型号</th><th>STF</th><th>任务状态</th><th>抖音</th><th>进度</th>
<th>任务名称</th><th>当前动作</th><th>操作计数</th><th>最近错误</th><th>操作</th>
</tr></thead><tbody id="tb-devices"></tbody></table>
</div>
<script>
const ACTION_LABELS={like:'点赞',comment:'评论',follow:'关注',share:'分享'};
function doPost(url,body){return fetch(url,{method:'POST',headers:{'Content-Type':'application/json'},body:body?JSON.stringify(body):'{}'}).then(r=>r.json()).then(d=>{alert(JSON.stringify(d,null,1));loadStatus();}).catch(e=>alert(e));}
function fmtActionCounts(c){
if(!c||!Object.keys(c).length)return '<span class="muted">-</span>';
const parts=Object.entries(c).filter(([_,v])=>v>0).map(([k,v])=>'<span class="kv">'+(ACTION_LABELS[k]||k)+v+'</span>');
return parts.length?parts.join(''):'<span class="muted">-</span>';
}
async function loadStatus(){
try{
const r=await fetch('/api/status');const d=await r.json();if(!d.ok)return;
document.getElementById('devCount').textContent='设备 '+d.devices.length;
document.getElementById('srvTime').textContent=new Date(d.server_time*1000).toLocaleTimeString();
// 统计卡片
const running=d.devices.filter(x=>x.worker_status==='running').length;
const idle=d.devices.filter(x=>x.worker_status==='idle').length;
const error=d.devices.filter(x=>['error','failed'].includes(x.worker_status)).length;
const occupied=d.devices.filter(x=>x.stf_occupied).length;
document.getElementById('stats').innerHTML=
'<div class="stat-card"><div class="num" style="color:#52c41a">'+running+'</div><div class="lbl">运行中</div></div>'+
'<div class="stat-card"><div class="num" style="color:#1890ff">'+occupied+'</div><div class="lbl">STF占用</div></div>'+
'<div class="stat-card"><div class="num" style="color:#bfbfbf">'+idle+'</div><div class="lbl">空闲</div></div>'+
'<div class="stat-card"><div class="num" style="color:#ff4d4f">'+error+'</div><div class="lbl">异常</div></div>';
// 设备表格
document.getElementById('tb-devices').innerHTML=d.devices.map(dev=>{
const stf=dev.stf_occupied?'<span class="badge b-running">占用</span>':'<span class="badge b-idle">空闲</span>';
if(!dev.present)stf+=' <span class="no">离线</span>';
const ws=dev.worker_status;
const wb='<span class="badge b-'+ws+'">'+({idle:'未启动',connecting:'连接中',running:'运行中',done:'完成',error:'异常',released:'已释放',failed:'失败'}[ws]||ws)+'</span>';
const dy=dev.douyin_running?'<span class="ok">●</span>':'<span class="no">○</span>';
const pct=Math.round(dev.videos_watched/80*100);
const prog='<span class="progress"><div style="width:'+pct+'%"></div></span> '+dev.videos_watched;
const tname=dev.task_job?'<span class="tag">'+dev.task_job+'</span>':'<span class="muted">-</span>';
const retry=dev.attempt>0?' <span class="muted">'+dev.attempt+'次</span>':'';
const btn=dev.worker_status==='running'||dev.worker_status==='connecting'?'<button class="danger sm" onclick="doPost(\'/api/stop_device\',{serial:\''+dev.serial+'\'})">停止</button>':'';
return '<tr><td>'+dev.serial+'</td><td>'+dev.model+'</td><td>'+stf+'</td><td>'+wb+'</td><td>'+dy+'</td><td>'+prog+'</td><td>'+tname+retry+'</td><td class="muted">'+(dev.current_action||'-')+'</td><td>'+fmtActionCounts(dev.action_counts)+'</td><td class="err" title="'+(dev.last_error||'')+'">'+(dev.last_error||'-')+'</td><td>'+btn+'</td></tr>';
}).join('');
}catch(e){console.error(e);}
}
loadStatus();
setInterval(loadStatus,5000);
</script>
</body></html>
+533
View File
@@ -0,0 +1,533 @@
"""监控域 API:状态/运行控制/设备操作/远程看屏。"""
import time
import threading
import subprocess
import shlex
from concurrent.futures import ThreadPoolExecutor, as_completed
from flask import (Blueprint, jsonify, request, Response, render_template)
from flask_login import login_required
from core import device_pool
from core.adb_helper import (screenshot, adb_connect, adb_connect_light,
_ADB_LOCK, _adb)
from web import context
from web.auth import perm_required, PERM_DEVICES
from web.common import _merged_device_list, _device_screen_state
from config import ADB_PATH
from core.logger import get_logger
_log = get_logger("web")
bp = Blueprint("monitor", __name__)
# 远程看屏参数
_SCREEN_JPEG_QUALITY = 65
_SCREEN_FRAME_GAP = 0.05
@bp.route("/api/health")
def api_health():
"""轻量健康检查:返回进程/设备/任务摘要,不暴露敏感信息,供运维探活。"""
try:
devices, err = context.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")),
"jobs": len(context.mgr.jobs),
})
except Exception as e:
return jsonify({"ok": False, "status": "down", "error": str(e)}), 500
# ================== API:失败任务汇总 ==================
@bp.route("/api/summary")
@login_required
def api_summary():
"""失败/异常任务汇总:各状态计数 + 异常设备列表(last_error / 重试次数)。
供监控页"异常汇总"面板使用,便于长期运行观察设备健康度。
"""
try:
sts = get_all_worker_status()
# 已删除/不在当前 STF 设备池的设备,其陈旧失败记录不应展示。
# 复用 get_status 缓存(顺带触发上面的状态清理),STF 异常时不隐藏任何记录。
try:
devices, derr = context.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:状态(监控大屏用)==================
@bp.route("/api/status")
@login_required
def api_status():
devices, err = context.mgr.get_status()
if err:
return jsonify({"ok": False, "error": err}), 500
return jsonify({
"ok": True, "devices": devices, "server_time": time.time(),
"fg_scanning": context.mgr._fg_scanner.is_scanning,
"fg_last_scan": context.mgr._fg_scanner.last_scan_time,
})
@bp.route("/api/scan_foreground", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_scan_foreground():
"""手动触发前台 App 扫描(不打扰设备)。"""
started = context.mgr._fg_scanner.scan_once()
if started:
return jsonify({"ok": True, "msg": "扫描已启动"})
return jsonify({"ok": False, "error": "已有扫描在进行中"})
@bp.route("/api/stop_device", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_stop_device():
serial = (request.json or {}).get("serial", "")
if context.mgr.stop_device(serial):
return jsonify({"ok": True, "msg": f"已发送停止信号给 {serial}"})
return jsonify({"ok": False, "error": f"{serial} 没有运行中的任务"}), 400
@bp.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 ""
@bp.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
@bp.route("/api/device/screen_all", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_device_screen_all():
"""批量亮屏/息屏(并发)。mode: on=亮屏解锁 / off=息屏。
请求可带 serials 指定设备(前端勾选);不带则作用于全部在线设备。
息屏会中断运行中的任务(u2 无法操作),前端已有确认提示。
"""
data = request.json or {}
mode = data.get("mode", "on")
try:
online = set(device_pool.list_online())
except Exception as e:
return jsonify({"ok": False, "error": f"获取设备列表失败: {e}"}), 502
serials = data.get("serials") or []
if serials:
# 指定设备:过滤掉不在线的(离线设备操作无意义,直接跳过)
serials = [s for s in serials if s in online]
else:
serials = list(online)
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)})
@bp.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
@bp.route("/api/stop_all", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_stop_all():
stopped = context.mgr.stop_all()
return jsonify({"ok": True, "stopped": stopped})
@bp.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 context.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} 的异常状态"})
@bp.route("/api/device/clear_all_errors", methods=["POST"])
@perm_required(PERM_DEVICES)
def api_device_clear_all_errors():
"""一键清除所有异常/失败设备(跳过正在运行/等待重试的)。"""
running = set(context.mgr.get_running())
cleared = clear_all_worker_errors(exclude=running)
_log.info(f"一键清除异常:共清除 {cleared} 台")
return jsonify({"ok": True, "cleared": cleared,
"msg": f"已清除 {cleared} 台设备的异常状态"})
@bp.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
@bp.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
# ================== 自定义动作(步骤打包) ==================
@bp.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
# 参数:q=JPEG 质量(大屏弹窗用 85 更清晰);fps=目标帧率(0=默认节流)
try:
quality = max(30, min(95, int(request.args.get("q", _SCREEN_JPEG_QUALITY))))
fps = float(request.args.get("fps", 0) or 0)
except ValueError:
quality, fps = _SCREEN_JPEG_QUALITY, 0
gap = (1.0 / fps) if fps > 0 else _SCREEN_FRAME_GAP
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=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(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)
@bp.route("/api/screen/thumb")
@perm_required(PERM_DEVICES)
def api_screen_thumb():
"""大屏缩略图:单张 JPEG(360px 宽,质量 55),供监控大屏轮询。
一次性请求(非流),与任务并发安全(与任务截图同走 u2 minicap)。
大屏按需轮询:每设备约 2.5s 一帧,20 台 ≈ 0.2 核 CPU + 100KB/s。
响应头 X-Screen-State: on/off/unknown —— 大屏据此显示「亮屏中/熄屏中」。
"""
serial = request.args.get("serial", "")
if not serial:
return jsonify({"ok": False, "error": "缺少 serial"}), 400
import io
import uiautomator2 as u2
try:
d = u2.connect(serial)
img = d.screenshot()
if img is None:
return jsonify({"ok": False, "error": "截图失败"}), 503
if img.width > 360:
ratio = 360 / img.width
img = img.resize((360, int(img.height * ratio)))
buf = io.BytesIO()
img.convert("RGB").save(buf, "JPEG", quality=55)
return Response(buf.getvalue(), mimetype="image/jpeg",
headers={"Cache-Control": "no-store",
"X-Screen-State": _device_screen_state(serial)})
except Exception as e:
return jsonify({"ok": False, "error": str(e)}), 503
@bp.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
@bp.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"}
@bp.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
@bp.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 上传/安装)==================
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 ""
def _screen_get_device(serial):
"""远程看屏用的 u2 连接(连接失败抛异常由调用方转 503)。"""
import uiautomator2 as u2
return u2.connect(serial)
+98
View File
@@ -0,0 +1,98 @@
"""Tailscale 管理 API(仅管理员)。"""
import re
from flask import Blueprint, jsonify, request
from core import tailscale_client
from core.tailscale_client import TailscaleError
from core.logger import get_logger
_log = get_logger("web")
from web.auth import admin_required
bp = Blueprint("tailscale", __name__)
@bp.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 ""})
@bp.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
@bp.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
@bp.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
@bp.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
@bp.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
+336
View File
@@ -0,0 +1,336 @@
"""任务域 API:任务计划/分组/自定义动作/元素抓取/步骤测试。"""
import re
import time
import uuid
import threading
from datetime import datetime
from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError as FuturesTimeout
from flask import (Blueprint, jsonify, request, Response)
from flask_login import login_required
from core import device_pool, uiauto_helper
from core.adb_helper import adb_connect_light
from core.models import db, DeviceGroup, TaskJob, CustomAction
from web import context
from web.auth import perm_required, PERM_DEVICES, PERM_TASKS
from tasks import list_task_types, get_task_class
from web.common import _merged_device_list
bp = Blueprint("tasks", __name__)
@bp.route("/api/task_types")
@login_required
def api_task_types():
return jsonify({"ok": True, "task_types": list_task_types()})
@bp.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})
@bp.route("/api/devices")
@login_required
def api_devices():
"""返回所有在线设备 serial(供分组表单勾选用)。"""
try:
return jsonify({"ok": True, "devices": context.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 = context.mgr.next_run_of(job)
return nr.strftime("%Y-%m-%d %H:%M") if nr else None
@bp.route("/api/jobs")
@login_required
def api_jobs_list():
jobs = []
for j in context.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()})
@bp.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 = context.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})
@bp.route("/api/jobs/<job_id>", methods=["PUT"])
@perm_required(PERM_TASKS)
def api_jobs_update(job_id):
job = context.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
context.mgr.update_job(job_id, **fields)
d = job.to_dict()
d["next_run"] = _job_next_run(job)
return jsonify({"ok": True, "msg": "任务已更新", "job": d})
@bp.route("/api/jobs/<job_id>", methods=["DELETE"])
@perm_required(PERM_TASKS)
def api_jobs_delete(job_id):
if context.mgr.delete_job(job_id):
return jsonify({"ok": True, "msg": "任务已删除"})
return jsonify({"ok": False, "error": "任务不存在"}), 404
@bp.route("/api/jobs/<job_id>/run", methods=["POST"])
@perm_required(PERM_TASKS)
def api_jobs_run(job_id):
return jsonify(context.mgr.run_job_now(job_id))
@bp.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 = context.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 ==================
@bp.route("/api/groups")
@login_required
def api_groups_list():
return jsonify({"ok": True, "groups": [g.to_dict() for g in context.mgr.groups.values()]})
@bp.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 context.mgr.groups:
return jsonify({"ok": False, "error": "分组名已存在"}), 400
context.mgr.add_group(name, data.get("serials", []), data.get("description", ""))
return jsonify({"ok": True, "msg": "分组已创建"})
@bp.route("/api/groups/<name>", methods=["PUT"])
@perm_required(PERM_TASKS)
def api_groups_update(name):
if name not in context.mgr.groups:
return jsonify({"ok": False, "error": "分组不存在"}), 404
data = request.json or {}
context.mgr.update_group(name,
serials=data.get("serials"),
description=data.get("description"))
return jsonify({"ok": True, "msg": "分组已更新"})
@bp.route("/api/groups/<name>", methods=["DELETE"])
@perm_required(PERM_TASKS)
def api_groups_delete(name):
if context.mgr.delete_group(name):
return jsonify({"ok": True, "msg": "分组已删除"})
return jsonify({"ok": False, "error": "分组不存在"}), 404
# ================== API:用户管理 CRUD ==================
@bp.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]})
@bp.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()})
@bp.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()})
@bp.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": "已删除"})
@bp.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
@bp.route("/api/uiauto/status")
@login_required
def api_uiauto_status():
"""探测 uiauto2 本地服务是否运行(前端按钮禁启用)。"""
return jsonify({"ok": True, "running": uiauto_helper.is_running()})
@bp.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),与运行中任务互不干扰。
返回 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})
@bp.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})
@bp.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
def _job_next_run(job):
"""任务下次执行时间,格式化为 "YYYY-MM-DD HH:MM"(Flask 序列化 aware datetime 会变 GMT 格式)。"""
nr = context.mgr.next_run_of(job)
return nr.strftime("%Y-%m-%d %H:%M") if nr else None
+236
View File
@@ -0,0 +1,236 @@
"""工具域 API:adb 终端/剪贴板注入/应用版本。"""
import re
import subprocess
import shlex
from concurrent.futures import ThreadPoolExecutor, as_completed
from flask import Blueprint, jsonify, request
from core import device_pool
from core.adb_helper import adb_connect, adb_connect_light, _ADB_LOCK, _adb
from core.logger import get_logger
from web import context
from web.auth import admin_required, perm_required
from web.common import _merged_device_list
_log = get_logger("web")
bp = Blueprint("tools", __name__)
@bp.route("/api/adb/devices")
@admin_required
def api_adb_devices():
"""维护终端设备列表(仅管理员):本地 adb 已连接 + 设备池(SQLite)。
状态:device/offline=本机 adb 实际状态;pool=设备池已配置但本机未连接。
供终端设备选择器使用——选中后自动附加 `-s <serial>`,
离线/未连接的设备也可选,配合"重连设备"按钮恢复。
"""
return jsonify({"ok": True, "devices": _merged_device_list()})
@bp.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
@bp.route("/api/tools/clipboard/set", methods=["POST"])
@admin_required
def api_tools_clipboard_set():
"""工具-剪贴板注入:把指定文字写入一台或多台设备的剪贴板。
请求: {"serials": ["100.100.10.11:5555", ...], "text": "要注入的文字"}
设备来源与维护终端一致:本地 adb 已连接设备(含 USB)+ 设备池。
实现: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]}"
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_light(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)}
@bp.route("/api/tools/appver", methods=["POST"])
@admin_required
def api_tools_appver():
"""工具-应用版本管理:查询所有设备上指定包名的安装情况与版本号。
请求: {"pkg": "com.ss.android.ugc.aweme"}
设备来源:本地 adb(含 USB)+ 设备池,并发查询(最多 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
# 只查在线设备:池内离线条目(如陈旧记录)不连,避免 connect 重试拖慢查询
try:
online = set(device_pool.list_online())
devices = [d for d in _merged_device_list() if d["serial"] in online]
except Exception:
devices = []
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)})
def _blocked_adb_cmd(cmd):
"""命中红线的 adb 命令(kill-server / disconnect)直接拒绝。"""
low = cmd.lower()
return any(p in low for p in _ADB_BLOCKED_PATTERNS)
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]}"
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_light(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)}
+18 -1621
View File
File diff suppressed because it is too large Load Diff