阶段3: 摘除 STF——删除 stf_client/stf_device_mgmt 与全部引用(STF 容器重启/设备管理/强制释放占用/agent 按钮/占用列),调度与设备管理完全基于本地设备池+adb;错误类迁入 device_worker,create_worker 去 stf 参数
This commit is contained in:
@@ -2,7 +2,6 @@
|
||||
|
||||
config — 全局配置常量(STF、adb 路径、养号参数)
|
||||
adb_helper — adb 命令封装(并发安全)
|
||||
stf_client — OpenSTF REST API 客户端
|
||||
device_worker — 设备生命周期 + 养号 worker + 全局状态注册表
|
||||
task_manager — 任务管理框架(分组/计划/调度/重试/持久化)
|
||||
"""
|
||||
|
||||
+3
-8
@@ -29,7 +29,6 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from config import APK_DIR, ADB_PATH
|
||||
from core.logger import get_logger
|
||||
from core.models import db, ApkFile as ApkRow
|
||||
from .stf_client import STFClient
|
||||
from .device_worker import get_all_worker_status
|
||||
|
||||
_log = get_logger("core.apk")
|
||||
@@ -43,8 +42,7 @@ _INSTALL_CONCURRENCY = 5
|
||||
class ApkManager:
|
||||
"""APK 文件管理 + 批量安装。"""
|
||||
|
||||
def __init__(self, stf_client=None, app=None):
|
||||
self.stf = stf_client or STFClient()
|
||||
def __init__(self, app=None):
|
||||
self.app = app # Flask app,用于 db context
|
||||
os.makedirs(APK_DIR, exist_ok=True)
|
||||
self._lock = threading.Lock()
|
||||
@@ -194,11 +192,8 @@ class ApkManager:
|
||||
if not serials:
|
||||
return False, "未选择设备"
|
||||
|
||||
# 获取设备信息(型号等)
|
||||
try:
|
||||
all_devices = {d["serial"]: d for d in self.stf.list_all_devices()}
|
||||
except Exception:
|
||||
all_devices = {}
|
||||
# 设备信息(型号从 worker 状态取,池内无型号信息)
|
||||
all_devices = {}
|
||||
|
||||
# 获取 worker 状态(判断设备是否运行中)
|
||||
worker_status = {w["serial"]: w for w in get_all_worker_status()}
|
||||
|
||||
+5
-51
@@ -1,22 +1,17 @@
|
||||
"""设备池(STF 替代数据源):SQLite devices 表 = 设备清单,本地 adb = 在线状态。
|
||||
|
||||
阶段 0:纯新增模块,不改现有行为。task_manager / web_server 的 STF 调用
|
||||
在阶段 1 逐个切换到本模块,阶段 3 摘除 STF 后删除 seed 相关代码。
|
||||
"""设备池:SQLite devices 表 = 设备清单,本地 adb = 在线状态(STF 替代数据源)。
|
||||
|
||||
设计:
|
||||
- list_configured() — devices 表里 enabled 的设备(清单,管理页维护)
|
||||
- list_online() — 本机 `adb devices` 里 state=device 的设备(实时)
|
||||
- list_ready() — configured ∩ online(调度用,取代 STF list_free_devices)
|
||||
- list_online() — 本机 `adb devices` 里 state=device 的设备(实时;
|
||||
池内有 USB 设备时合并 220 远程 adb server 状态)
|
||||
- list_ready() — configured ∩ online(任务调度用)
|
||||
- CRUD — add/update/remove/set_enabled(设备池管理页用)
|
||||
- 首次启动自动从 STF 导入现有设备(一次性,app_meta 标记防重复)
|
||||
|
||||
红线:全模块不 connect / 不 kill-server / 不 disconnect,遵守既有技术约束。
|
||||
后台线程(task_manager 等)调用时由本模块自行推 app context,调用方无需关心。
|
||||
"""
|
||||
import time
|
||||
|
||||
from sqlalchemy import text
|
||||
|
||||
from config import USB_ADB_HOST, USB_ADB_PORT
|
||||
from core.adb_helper import _adb, _adb_remote
|
||||
from core.logger import get_logger
|
||||
@@ -28,13 +23,9 @@ _app = None
|
||||
|
||||
|
||||
def init_app(app):
|
||||
"""web_server 启动时调用:绑定 Flask app,并做首次设备导入。"""
|
||||
"""web_server 启动时调用:绑定 Flask app(供后台线程推 db context)。"""
|
||||
global _app
|
||||
_app = app
|
||||
try:
|
||||
_seed_from_stf_if_empty()
|
||||
except Exception as e:
|
||||
_log.warning(f"设备池初始化异常(不影响启动): {e}")
|
||||
|
||||
|
||||
def _ctx():
|
||||
@@ -142,40 +133,3 @@ def set_enabled(serial, enabled):
|
||||
d.enabled = enabled
|
||||
db.session.commit()
|
||||
return True
|
||||
|
||||
|
||||
# ================== 首次导入(阶段 3 摘除 STF 后删除) ==================
|
||||
def _seed_from_stf_if_empty():
|
||||
"""首次启动把 STF 现有设备导入本地清单。
|
||||
|
||||
条件:device 表为空 且 未导入过(app_meta 标记)。STF 不可达时跳过,
|
||||
不写标记,下次启动重试。摘除 STF 后此函数随 stf_client 一起删除。
|
||||
"""
|
||||
try:
|
||||
from core.stf_client import STFClient
|
||||
except Exception:
|
||||
return
|
||||
with _ctx():
|
||||
if Device.query.count() > 0:
|
||||
return
|
||||
seeded = db.session.execute(text(
|
||||
"SELECT value FROM app_meta WHERE key='devices_seeded_from_stf'")).scalar()
|
||||
if seeded:
|
||||
return
|
||||
try:
|
||||
# 导入全部 serial(不限于 present):STF 池 = 配置的舰队,
|
||||
# present=false 的是掉线/陈旧记录,导入后显示离线,管理页可删
|
||||
serials = [d["serial"] for d in STFClient().list_all_devices()]
|
||||
except Exception as e:
|
||||
_log.warning(f"从 STF 导入设备清单失败(可稍后手动添加): {e}")
|
||||
return
|
||||
if not serials:
|
||||
_log.warning("STF 无在线设备,跳过设备池首次导入")
|
||||
return
|
||||
now = time.strftime("%Y-%m-%d %H:%M")
|
||||
for s in serials:
|
||||
db.session.add(Device(serial=s, created_at=now))
|
||||
db.session.execute(text(
|
||||
"INSERT OR REPLACE INTO app_meta(key,value) VALUES('devices_seeded_from_stf','1')"))
|
||||
db.session.commit()
|
||||
_log.info(f"设备池首次导入 {len(serials)} 台设备(来源 STF)")
|
||||
|
||||
+19
-9
@@ -18,12 +18,23 @@ import uiautomator2 as u2
|
||||
|
||||
from config import USB_ADB_HOST, USB_ADB_PORT
|
||||
from core.logger import get_logger
|
||||
from core.stf_client import DeviceOfflineError, STFError
|
||||
from core import device_pool
|
||||
from .adb_helper import adb_connect, _adb_remote
|
||||
|
||||
_log = get_logger("core.worker")
|
||||
|
||||
|
||||
class STFError(Exception):
|
||||
"""设备操作错误(历史名称保留兼容)。code: offline/conflict/..."""
|
||||
|
||||
def __init__(self, message, code=""):
|
||||
super().__init__(message)
|
||||
self.code = code
|
||||
|
||||
|
||||
class DeviceOfflineError(STFError):
|
||||
"""设备离线/不可用:任务不重试(换设备也没用)。"""
|
||||
|
||||
# 心跳超时阈值(秒)。worker 超过这个时间没更新心跳,判定为卡死
|
||||
_HEARTBEAT_TIMEOUT = 120
|
||||
# 看门狗检查间隔
|
||||
@@ -201,11 +212,11 @@ def stop_watchdog():
|
||||
class BaseWorker(threading.Thread):
|
||||
"""通用 worker 基类。
|
||||
|
||||
自动处理:STF 占用/释放、u2 连接、状态上报、异常捕获、stop 信号、心跳。
|
||||
自动处理:设备获取、u2 连接、状态上报、异常捕获、stop 信号、心跳。
|
||||
子类只需实现 run_task(d) 方法,专注业务逻辑。
|
||||
|
||||
生命周期(基类 run() 已封装,不要重写):
|
||||
1. acquire 设备(STF 占用 + adb 连接)
|
||||
1. acquire 设备(直连 IP:5555 / USB 走远程 adb server)
|
||||
2. u2.connect 拿到 d
|
||||
3. 调用 setup(d) ← 子类可选钩子,做初始化
|
||||
4. 调用子类 run_task(d) ← 业务逻辑
|
||||
@@ -246,9 +257,8 @@ class BaseWorker(threading.Thread):
|
||||
其他异常 — 可重试
|
||||
"""
|
||||
|
||||
def __init__(self, stf_client, serial, params=None, daemon=True):
|
||||
def __init__(self, serial, params=None, daemon=True):
|
||||
super().__init__(daemon=daemon)
|
||||
self.stf = stf_client
|
||||
self.serial = serial
|
||||
self.params = params or {}
|
||||
self._stop_flag = threading.Event()
|
||||
@@ -347,7 +357,7 @@ class BaseWorker(threading.Thread):
|
||||
def run(self):
|
||||
"""基类主循环:不要重写。子类实现 run_task + 可选钩子。"""
|
||||
device = STFDevice(serial=self.serial)
|
||||
_update_status(self.serial, status="connecting", stf_occupied=False,
|
||||
_update_status(self.serial, status="connecting",
|
||||
last_error="", remote_adb_url="", usb_server="", model="",
|
||||
current_action="")
|
||||
try:
|
||||
@@ -370,7 +380,7 @@ class BaseWorker(threading.Thread):
|
||||
model = info.get("productName") or ""
|
||||
except Exception:
|
||||
pass
|
||||
_update_status(self.serial, status="running", stf_occupied=True,
|
||||
_update_status(self.serial, status="running",
|
||||
remote_adb_url=remote or "", usb_server=usb_server or "",
|
||||
model=model)
|
||||
|
||||
@@ -425,9 +435,9 @@ class BaseWorker(threading.Thread):
|
||||
with _WORKERS_LOCK:
|
||||
cur = _WORKERS.get(self.serial, {}).get("status")
|
||||
if cur in ("running", "connecting"):
|
||||
_update_status(self.serial, stf_occupied=False, status="released")
|
||||
_update_status(self.serial, status="released")
|
||||
else:
|
||||
_update_status(self.serial, stf_occupied=False)
|
||||
pass
|
||||
|
||||
def _u2_connect_with_timeout(self, remote):
|
||||
"""u2.connect 带超时保护,避免 atx-agent 无响应时永久 hang。"""
|
||||
|
||||
@@ -1,247 +0,0 @@
|
||||
"""OpenSTF REST API 客户端封装。
|
||||
|
||||
负责:设备列表查询、占用、远程 ADB 隧道建立、释放。
|
||||
不负责实际 UI 操作(那是 DeviceWorker 的事)。
|
||||
|
||||
健壮性设计:
|
||||
- 所有请求带 timeout,避免 STF 网关 504 时长时间卡死
|
||||
- remoteConnect 带重试,STF provider 偶尔慢启动
|
||||
- 友好错误分类:设备离线 / 网络不通 / STF 内部错误,便于上层决策
|
||||
- 占用冲突自动重试(STF 偶发 400 "device already in use")
|
||||
"""
|
||||
import time
|
||||
|
||||
import requests
|
||||
|
||||
from config import STF_URL, STF_TOKEN
|
||||
from core.logger import get_logger
|
||||
|
||||
_log = get_logger("core.stf")
|
||||
|
||||
# 请求超时(秒)。connect 超时 + read 超时
|
||||
_TIMEOUT = (5, 15)
|
||||
# remoteConnect 重试次数(STF provider 慢启动时需要重试)
|
||||
_RC_RETRIES = 3
|
||||
# remoteConnect 重试间隔
|
||||
_RC_DELAY = 2
|
||||
|
||||
|
||||
class STFError(Exception):
|
||||
"""STF 相关错误基类。"""
|
||||
|
||||
def __init__(self, message, code=""):
|
||||
super().__init__(message)
|
||||
self.code = code # 错误码:offline / network / conflict / server / unknown
|
||||
|
||||
|
||||
class DeviceOfflineError(STFError):
|
||||
"""设备离线/不可用。"""
|
||||
|
||||
|
||||
class STFNetworkError(STFError):
|
||||
"""STF 服务不可达。"""
|
||||
|
||||
|
||||
class DeviceConflictError(STFError):
|
||||
"""设备已被占用。"""
|
||||
|
||||
|
||||
def _headers():
|
||||
return {"Authorization": f"Bearer {STF_TOKEN}"}
|
||||
|
||||
|
||||
class STFClient:
|
||||
"""OpenSTF REST API 封装。"""
|
||||
|
||||
def __init__(self, base_url=STF_URL, token=STF_TOKEN):
|
||||
self.base_url = base_url
|
||||
self.headers = _headers()
|
||||
|
||||
# ================== 设备查询 ==================
|
||||
def list_all_devices(self):
|
||||
"""返回 STF 上所有设备(含状态)。"""
|
||||
try:
|
||||
resp = requests.get(f"{self.base_url}/api/v1/devices",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
resp.raise_for_status()
|
||||
return resp.json().get("devices", [])
|
||||
except requests.exceptions.ConnectionError as e:
|
||||
raise STFNetworkError(f"STF 服务不可达: {e}", "network") from e
|
||||
except requests.exceptions.Timeout as e:
|
||||
raise STFNetworkError(f"STF 请求超时: {e}", "network") from e
|
||||
except requests.exceptions.RequestException as e:
|
||||
_log.error(f"获取设备列表失败: {e}")
|
||||
raise STFError(f"获取设备列表失败: {e}", "server") from e
|
||||
|
||||
def list_free_devices(self):
|
||||
"""返回可占用的空闲设备。"""
|
||||
return [d for d in self.list_all_devices()
|
||||
if d.get("present") and d.get("ready")
|
||||
and not d.get("using") and d.get("owner") is None]
|
||||
|
||||
def list_my_devices(self):
|
||||
"""返回当前账户已占用的设备。"""
|
||||
try:
|
||||
resp = requests.get(f"{self.base_url}/api/v1/user/devices",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
resp.raise_for_status()
|
||||
return resp.json().get("devices", [])
|
||||
except requests.exceptions.RequestException as e:
|
||||
_log.error(f"获取已占用设备失败: {e}")
|
||||
return []
|
||||
|
||||
# ================== 占用/释放 ==================
|
||||
def occupy(self, serial, retries=2):
|
||||
"""占用设备。返回 (ok, msg)。
|
||||
|
||||
冲突(已被占用)时自动重试,偶发 400 "already in use" 可能是脏状态。
|
||||
"""
|
||||
for attempt in range(1, retries + 1):
|
||||
try:
|
||||
r = requests.post(f"{self.base_url}/api/v1/user/devices/{serial}",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
if r.status_code in (200, 201):
|
||||
return True, r.text[:200]
|
||||
if r.status_code == 400 and "already" in r.text.lower():
|
||||
# 已被占用(可能是自己之前占用没释放干净),尝试先释放再占
|
||||
_log.warning(f"{serial} 占用冲突,尝试清理后重试 ({attempt}/{retries})")
|
||||
try:
|
||||
requests.delete(f"{self.base_url}/api/v1/user/devices/{serial}",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
except Exception:
|
||||
pass
|
||||
time.sleep(1)
|
||||
continue
|
||||
return False, f"HTTP {r.status_code}: {r.text[:200]}"
|
||||
except requests.exceptions.RequestException as e:
|
||||
_log.warning(f"{serial} 占用请求异常 ({attempt}/{retries}): {e}")
|
||||
if attempt < retries:
|
||||
time.sleep(1)
|
||||
return False, str(e)
|
||||
return False, "占用冲突重试耗尽"
|
||||
|
||||
def release(self, serial, retries=3):
|
||||
"""释放设备占用。返回 (ok, msg)。
|
||||
|
||||
STF 释放可能因设备响应超时返回 504(如设备过载/卡死),此时设备实际未释放。
|
||||
必须检查响应状态并重试,不能静默吞掉——否则占用会一直悬着。
|
||||
"""
|
||||
self.remote_disconnect(serial)
|
||||
last_status = None
|
||||
last_body = ""
|
||||
for attempt in range(1, retries + 1):
|
||||
try:
|
||||
r = requests.delete(f"{self.base_url}/api/v1/user/devices/{serial}",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
if r.status_code in (200, 201, 202, 204):
|
||||
return True, "released"
|
||||
last_status = r.status_code
|
||||
last_body = r.text[:150]
|
||||
_log.warning(f"释放 {serial} 失败 HTTP {r.status_code}: {last_body} ({attempt}/{retries})")
|
||||
except requests.exceptions.RequestException as e:
|
||||
_log.warning(f"释放 {serial} 请求异常 ({attempt}/{retries}): {e}")
|
||||
last_status = -1
|
||||
last_body = str(e)
|
||||
if attempt < retries:
|
||||
time.sleep(2)
|
||||
return False, f"STF 释放失败(HTTP {last_status}): {last_body}"
|
||||
|
||||
def release_all_mine(self):
|
||||
"""释放当前账户占用的所有设备(清理用)。返回 (released, failed)。"""
|
||||
released, failed = [], []
|
||||
for d in self.list_my_devices():
|
||||
serial = d["serial"]
|
||||
ok, msg = self.release(serial)
|
||||
if ok:
|
||||
released.append(serial)
|
||||
_log.info(f"已释放 {serial}")
|
||||
else:
|
||||
failed.append(serial)
|
||||
_log.error(f"释放 {serial} 失败: {msg}")
|
||||
return released, failed
|
||||
|
||||
def delete_device(self, serial):
|
||||
"""从 STF 池彻底删除设备记录(清理离线幽灵设备用)。
|
||||
|
||||
DELETE /api/v1/devices/{serial}(实测可用):断开后 STF 数据库里残留的
|
||||
present=False 陈旧记录会一直显示在设备池,用此端点彻底移除。
|
||||
"""
|
||||
import urllib.parse
|
||||
try:
|
||||
r = requests.delete(
|
||||
f"{self.base_url}/api/v1/devices/{urllib.parse.quote(serial, safe='')}",
|
||||
headers=self.headers, timeout=_TIMEOUT)
|
||||
if r.status_code in (200, 202, 204):
|
||||
_log.info(f"已从 STF 池删除设备 {serial}")
|
||||
return True, "已从 STF 池删除"
|
||||
return False, f"HTTP {r.status_code}: {r.text[:150]}"
|
||||
except requests.exceptions.RequestException as e:
|
||||
_log.warning(f"删除 STF 设备 {serial} 请求异常: {e}")
|
||||
return False, str(e)
|
||||
|
||||
# ================== 远程 ADB 隧道 ==================
|
||||
def remote_connect(self, serial):
|
||||
"""建立远程 ADB 隧道,返回 remoteConnectUrl。
|
||||
|
||||
STF provider 慢启动或设备掉线时会 504,这里带重试。
|
||||
设备真离线时快速失败,不长时间卡住 worker。
|
||||
"""
|
||||
last_err = None
|
||||
for attempt in range(1, _RC_RETRIES + 1):
|
||||
try:
|
||||
conn = requests.post(
|
||||
f"{self.base_url}/api/v1/user/devices/{serial}/remoteConnect",
|
||||
headers=self.headers, timeout=_TIMEOUT,
|
||||
)
|
||||
# 504 = STF 网关等 provider 响应超时,通常是设备掉线或 provider 卡死
|
||||
if conn.status_code == 504:
|
||||
last_err = f"STF 网关超时(504),设备可能掉线"
|
||||
_log.warning(f"{serial} remoteConnect 504 ({attempt}/{_RC_RETRIES})")
|
||||
if attempt < _RC_RETRIES:
|
||||
time.sleep(_RC_DELAY)
|
||||
continue
|
||||
raise DeviceOfflineError(
|
||||
f"{serial} remoteConnect 超时,设备可能掉线或 provider 卡死",
|
||||
"offline",
|
||||
)
|
||||
conn.raise_for_status()
|
||||
data = conn.json()
|
||||
if not data.get("success"):
|
||||
raise STFError(f"远程连接失败: {data.get('description')}", "server")
|
||||
url = (data.get("remoteConnectUrl") or data.get("remoteAdbUrl")
|
||||
or data.get("remote_adb_url") or data.get("adbUrl")
|
||||
or data.get("url"))
|
||||
if not url:
|
||||
raise STFError("STF 返回成功但无 remoteConnectUrl", "server")
|
||||
return url
|
||||
except DeviceOfflineError:
|
||||
raise
|
||||
except requests.exceptions.ConnectionError as e:
|
||||
last_err = str(e)
|
||||
_log.warning(f"{serial} remoteConnect 网络异常 ({attempt}/{_RC_RETRIES}): {e}")
|
||||
if attempt < _RC_RETRIES:
|
||||
time.sleep(_RC_DELAY)
|
||||
except requests.exceptions.Timeout as e:
|
||||
last_err = str(e)
|
||||
_log.warning(f"{serial} remoteConnect 超时 ({attempt}/{_RC_RETRIES}): {e}")
|
||||
if attempt < _RC_RETRIES:
|
||||
time.sleep(_RC_DELAY)
|
||||
except requests.exceptions.RequestException as e:
|
||||
last_err = str(e)
|
||||
_log.error(f"{serial} remoteConnect 请求异常: {e}")
|
||||
raise STFError(f"远程连接请求失败: {e}", "server") from e
|
||||
# 重试耗尽
|
||||
raise DeviceOfflineError(
|
||||
f"{serial} remoteConnect 重试 {_RC_RETRIES} 次失败: {last_err}",
|
||||
"offline",
|
||||
)
|
||||
|
||||
def remote_disconnect(self, serial):
|
||||
"""断开远程 ADB 隧道。"""
|
||||
try:
|
||||
requests.post(
|
||||
f"{self.base_url}/api/v1/user/devices/{serial}/remoteDisconnect",
|
||||
headers=self.headers, timeout=_TIMEOUT,
|
||||
)
|
||||
except requests.exceptions.RequestException:
|
||||
pass # 断开失败不影响主流程
|
||||
@@ -1,158 +0,0 @@
|
||||
"""OpenSTF 设备池管理(工具页"STF 设备管理"用)。
|
||||
|
||||
背景:220 上 STF 的 adb 跑在 Docker 容器里,设备池由
|
||||
/mnt/data/openstf/connect_devices.sh 的 DEVICES 数组维护(cron 每 5 分钟补连)。
|
||||
本模块把"改脚本 + 实时 adb connect/disconnect"封装成 API:
|
||||
- status 已配置 IP 列表 + 220 adb 实际连接状态
|
||||
- add_device 写入脚本 DEVICES + 立即 docker exec adb adb connect
|
||||
- remove_device 从脚本 DEVICES 移除 + adb disconnect
|
||||
|
||||
脚本编辑用 220 上的 python3 做精确的行级增删(sed 处理多行数组容易误伤)。
|
||||
注意:这里的 adb 是 220 上的(STF provider 的),与本机 Mac 直连的 adb 相互独立;
|
||||
移除设备会断开 STF 池连接(设备从 OpenSTF 下线),不影响本机已有的任务 adb 连接。
|
||||
"""
|
||||
import re
|
||||
import shlex
|
||||
|
||||
from config import STF_ADB_CONTAINER, STF_SCRIPT_PATH
|
||||
from core.logger import get_logger
|
||||
from core.ssh_client import run as _ssh_run, SSHError
|
||||
|
||||
_log = get_logger("core.stfdev")
|
||||
|
||||
|
||||
class StfDevError(Exception):
|
||||
"""STF 设备管理错误(SSH/脚本/命令失败)。"""
|
||||
|
||||
|
||||
def _ssh(cmd, timeout=25):
|
||||
"""在 220 上执行一条 shell 命令(认证见 core.ssh_client)。返回 (code, stdout, stderr)。"""
|
||||
try:
|
||||
return _ssh_run(cmd, timeout=timeout)
|
||||
except SSHError as e:
|
||||
raise StfDevError(str(e))
|
||||
|
||||
|
||||
def _adb(args, remote_timeout=15):
|
||||
"""在 220 的 adb 容器里执行 adb 命令。返回 (code, stdout, stderr)。
|
||||
|
||||
remote_timeout 是 220 侧的 timeout 秒数:adb connect 到不可达 IP 会长时间
|
||||
挂起重试(实测 ≥40s),必须用 host 的 timeout 兜底,否则 API 请求会卡死。
|
||||
"""
|
||||
return _ssh(f"timeout {remote_timeout} docker exec {shlex.quote(STF_ADB_CONTAINER)} adb "
|
||||
+ " ".join(shlex.quote(a) for a in args))
|
||||
|
||||
|
||||
def _norm_ip(ip):
|
||||
"""规范化 IP:去空白、去 :5555 后缀。"""
|
||||
ip = str(ip or "").strip()
|
||||
if ip.endswith(":5555"):
|
||||
ip = ip[:-5]
|
||||
return ip
|
||||
|
||||
|
||||
# ================== 查询 ==================
|
||||
def configured_ips():
|
||||
"""解析脚本 DEVICES 数组里的 IP 列表。"""
|
||||
code, out, err = _ssh(f"cat {shlex.quote(STF_SCRIPT_PATH)}")
|
||||
if code != 0:
|
||||
raise StfDevError(f"读取脚本失败: {err.strip()}")
|
||||
m = re.search(r"DEVICES=\((.*?)\)", out, re.S)
|
||||
if not m:
|
||||
raise StfDevError("脚本中未找到 DEVICES 列表")
|
||||
return re.findall(r'"([^"]+)"', m.group(1))
|
||||
|
||||
|
||||
def connected_devices():
|
||||
"""220 adb 实际连接状态 [{serial, state}]。"""
|
||||
code, out, err = _adb(["devices"])
|
||||
if code != 0:
|
||||
raise StfDevError(f"adb devices 失败: {err.strip()}")
|
||||
devs = []
|
||||
for line in out.splitlines()[1:]:
|
||||
parts = line.split()
|
||||
if len(parts) >= 2 and parts[0] and not parts[0].startswith("*"):
|
||||
devs.append({"serial": parts[0], "state": parts[1]})
|
||||
return devs
|
||||
|
||||
|
||||
def status():
|
||||
"""综合状态:脚本配置的 IP + 实际连接。返回 {configured, connected}。"""
|
||||
return {"configured": configured_ips(), "connected": connected_devices()}
|
||||
|
||||
|
||||
# ================== 增删设备 ==================
|
||||
def _edit_script(python_code):
|
||||
"""用 220 的 python3 执行脚本编辑代码(代码内请用双引号)。"""
|
||||
cmd = f"python3 -c {shlex.quote(python_code)}"
|
||||
code, out, err = _ssh(cmd)
|
||||
if code != 0:
|
||||
raise StfDevError(f"脚本编辑失败: {err.strip() or out.strip()}")
|
||||
|
||||
|
||||
def add_device(ip):
|
||||
"""添加设备:写入脚本 DEVICES + 立即 connect。返回操作消息列表。"""
|
||||
ip = _norm_ip(ip)
|
||||
if not re.match(r"^\d{1,3}(\.\d{1,3}){3}$", ip):
|
||||
raise StfDevError(f"IP 格式不正确: {ip}")
|
||||
msgs = []
|
||||
if ip in configured_ips():
|
||||
msgs.append(f"{ip} 已在脚本配置中")
|
||||
else:
|
||||
# 在 DEVICES=( 行后插入一行 IP(python 精确行级插入)
|
||||
_edit_script(
|
||||
"p='" + STF_SCRIPT_PATH + "';s=open(p).read();"
|
||||
"s=s.replace('DEVICES=(\\n','DEVICES=(\\n \\\"" + ip + "\\\"\\n',1);open(p,'w').write(s)")
|
||||
msgs.append(f"已写入脚本 {STF_SCRIPT_PATH}")
|
||||
code, out, err = _adb(["connect", f"{ip}:5555"])
|
||||
if code == 0:
|
||||
msgs.append(out.strip() or "adb connect 成功")
|
||||
elif code == 124:
|
||||
msgs.append("连接超时(设备当前不可达),已保留在脚本中,cron 会每 5 分钟自动重试")
|
||||
else:
|
||||
msgs.append(f"连接失败: {out.strip() or err.strip() or '未知错误'}")
|
||||
_log.info(f"STF 添加设备 {ip}: {msgs}")
|
||||
return msgs
|
||||
|
||||
|
||||
def remove_device(ip):
|
||||
"""移除设备:从脚本 DEVICES 删除 + adb disconnect。返回操作消息列表。"""
|
||||
ip = _norm_ip(ip)
|
||||
msgs = []
|
||||
if ip in configured_ips():
|
||||
# 删除 DEVICES 块内该 IP 的整行
|
||||
_edit_script(
|
||||
"p='" + STF_SCRIPT_PATH + "';s=open(p).read();"
|
||||
"s=s.replace(' \\\"" + ip + "\\\"\\n','');open(p,'w').write(s)")
|
||||
msgs.append("已从脚本移除")
|
||||
else:
|
||||
msgs.append(f"{ip} 不在脚本配置中")
|
||||
code, out, err = _adb(["disconnect", f"{ip}:5555"])
|
||||
msgs.append(out.strip() or (f"adb disconnect 执行完成(exit {code})" if code == 0 else "adb disconnect 无输出"))
|
||||
_log.info(f"STF 移除设备 {ip}: {msgs}")
|
||||
return msgs
|
||||
|
||||
|
||||
def delete_device(ip, stf=None):
|
||||
"""彻底删除设备:脚本移除 + disconnect + 清 STF 池记录(幽灵设备)。"""
|
||||
msgs = remove_device(ip)
|
||||
if stf is not None:
|
||||
ok, msg = stf.delete_device(f"{_norm_ip(ip)}:5555")
|
||||
msgs.append(msg)
|
||||
_log.info(f"STF 彻底删除设备 {ip}: {msgs}")
|
||||
return msgs
|
||||
|
||||
|
||||
def reconnect_all():
|
||||
"""一键重连:后台运行 220 的 connect_devices.sh(cron 之外的立即补连)。
|
||||
|
||||
脚本对未连接设备逐个 connect(不可达设备可能挂 40s+),全程可能 1-3 分钟,
|
||||
因此 nohup 后台执行立即返回,不阻塞请求;结果用"刷新"看设备连接状态。
|
||||
"""
|
||||
cmd = (f"nohup bash {shlex.quote(STF_SCRIPT_PATH)} >/tmp/connect_devices.out 2>&1 &")
|
||||
code, out, err = _ssh(cmd, timeout=15)
|
||||
if code != 0:
|
||||
raise StfDevError(f"启动重连失败: {err.strip() or out.strip()}")
|
||||
_log.info("STF 一键重连已在 220 后台启动")
|
||||
return True, ("重连已在 220 后台执行(对未连接设备逐个 connect,约 1-3 分钟),"
|
||||
"稍后点\"刷新\"查看设备连接状态")
|
||||
@@ -31,9 +31,9 @@ from config import DATA_DIR
|
||||
from core.logger import get_logger
|
||||
from core.models import db, DeviceGroup as GroupRow, TaskJob as JobRow
|
||||
from core import device_pool
|
||||
from .stf_client import STFClient, DeviceOfflineError
|
||||
from .adb_helper import get_foreground_app, get_foreground_app_remote
|
||||
from .device_worker import (
|
||||
DeviceOfflineError,
|
||||
get_all_worker_status, _update_status, _remove_worker,
|
||||
_WORKERS, _WORKERS_LOCK,
|
||||
start_watchdog, stop_watchdog,
|
||||
@@ -195,9 +195,9 @@ class _ForegroundScanner:
|
||||
|
||||
策略(按设备状态区分获取逻辑):
|
||||
- worker 运行中设备:用已有 remote_adb_url 直接查询(已有 adb 连接,无额外开销)
|
||||
- 空闲设备:直接 adb connect <serial> → 查询 → adb disconnect(绕过 STF,不打扰设备)
|
||||
- 被别人占用的设备:标记 "他人占用"
|
||||
- STF occupy/release 会打扰设备(可能退回桌面),绝不使用
|
||||
- USB 运行中设备:经远程 adb server(220)查询
|
||||
- 空闲设备:不主动连接(返回"空闲")——IP:5555 的 adb transport 与 220
|
||||
共享,外部 connect/disconnect 会扰动共享连接
|
||||
|
||||
serial 格式为 IP:5555(设备本身的 adb 网络地址),可直接 adb connect。
|
||||
adb connect/disconnect 只建立/断开调试连接,不影响设备 UI。
|
||||
@@ -334,8 +334,7 @@ class _ForegroundScanner:
|
||||
|
||||
# ================== 任务管理器 ==================
|
||||
class TaskManager:
|
||||
def __init__(self, stf_client=None, app=None):
|
||||
self.stf = stf_client or STFClient()
|
||||
def __init__(self, app=None):
|
||||
self.app = app # Flask app,用于 db context
|
||||
self.scheduler = BackgroundScheduler(timezone="Asia/Shanghai")
|
||||
self.scheduler.start()
|
||||
@@ -683,7 +682,7 @@ class TaskManager:
|
||||
|
||||
worker = None
|
||||
try:
|
||||
worker = task.create_worker(self.stf, serial, job.params)
|
||||
worker = task.create_worker(serial, job.params)
|
||||
with self._lock:
|
||||
self._running[serial]["worker"] = worker
|
||||
_log.info(f"{serial} 开始任务 {job.name} (第{attempt}/{max_attempts}次)")
|
||||
@@ -850,7 +849,6 @@ class TaskManager:
|
||||
"device_name": names.get(serial, ""),
|
||||
"present": is_online,
|
||||
"ready": is_online, # 阶段 1:ready 概念并入在线状态
|
||||
"stf_occupied": False, # 阶段 1:已无 STF 占用(阶段 3 删字段)
|
||||
"owner": "",
|
||||
"worker_status": w.get("status", "idle"),
|
||||
"foreground_app": self._fg_scanner.get(serial),
|
||||
|
||||
Reference in New Issue
Block a user