diff --git a/core/device_worker.py b/core/device_worker.py index 5990141..f48ad38 100644 --- a/core/device_worker.py +++ b/core/device_worker.py @@ -18,6 +18,7 @@ import uiautomator2 as u2 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_disconnect _log = get_logger("core.worker") @@ -32,13 +33,15 @@ _U2_INFO_TIMEOUT = 10 class STFDevice: - """单设备生命周期:acquire(占用+连接)→ release(断开+释放)。 + """单设备生命周期:acquire(直连)→ release(无操作)。 - 用 try/finally 保证设备一定被释放,即使养号中途报错。 + 阶段 1:已摘除 STF occupy/release——单实例互斥由 TaskManager._running 保证, + 设备健康由 device_pool(本地 adb 在线状态)判断。保留 STF 引用仅为 + 非 IP:port(USB 序列号)设备的 remoteConnect 桥接,阶段 2 换成 220 adb server。 """ - def __init__(self, stf_client, serial=None): - self.stf = stf_client + def __init__(self, stf_client=None, serial=None): + self.stf = stf_client # 仅 USB 桥接用(阶段 2 移除) self.serial = serial self.remote_adb_url = None @@ -46,22 +49,15 @@ class STFDevice: self.serial = serial or self.serial or self._pick_free() _log.info(f"[{self.serial}] 选定设备") - ok, msg = self.stf.occupy(self.serial) - if not ok: - raise STFError(f"占用设备失败: {msg}", "conflict") - _log.info(f"[{self.serial}] STF 占用成功") - # 直连优先:serial 本身是 IP:5555(如 Tailscale 网络设备)时直接 adb connect。 - # 实测(本机 100.100.10.2,设备 100.100.10.x): - # - 本机在 tailnet 内时直连可靠且快(<1s),且只 connect、绝不 disconnect/kill-server, - # 不会影响 STF provider 共享的 adb 连接(STF 端设备保持 present/ready)。 - # - STF remoteConnect 桥接在此环境反而有 adb key 认证问题(隧道显示 unauthorized, - # 即使设备已授权本机 key),直连才是可靠路径。 - # 仅当 serial 不是 IP:port(如 USB 序列号)时才走 STF 桥接。 + # 只 connect、绝不 disconnect/kill-server(红线),不影响 STF provider 的连接。 + # 仅当 serial 不是 IP:port(如 USB 序列号)时才走 STF 桥接(阶段 2 替换)。 if ":" in self.serial: self.remote_adb_url = self.serial _log.info(f"[{self.serial}] 直连设备: {self.remote_adb_url}") else: + if self.stf is None: + raise STFError(f"USB 设备 {self.serial} 暂不支持(未启用 STF 桥接)", "offline") self.remote_adb_url = self.stf.remote_connect(self.serial) _log.info(f"[{self.serial}] STF 桥接: {self.remote_adb_url}") @@ -70,27 +66,21 @@ class STFDevice: time.sleep(2) def _pick_free(self): - free = self.stf.list_free_devices() + free = device_pool.list_ready() if not free: raise STFError("没有可用的在线空闲设备", "offline") - return free[0]["serial"] + return free[0] def release(self): if not self.serial: return - # 直连模式(remote_adb_url == serial):不 disconnect - # STF provider 共享该 IP:5555 的 adb transport,disconnect 会让 STF 误判离线并重连 + # 直连模式(remote_adb_url == serial):不 disconnect(红线,STF 共享 transport) if self.remote_adb_url and self.remote_adb_url != self.serial: adb_disconnect(self.remote_adb_url) try: self.stf.remote_disconnect(self.serial) except Exception: pass - ok, msg = self.stf.release(self.serial) - if ok: - _log.info(f"[{self.serial}] 已释放") - else: - _log.warning(f"[{self.serial}] STF 释放失败: {msg}") # ================== 全局 worker 状态注册表(供 web_server 读取) ================== diff --git a/core/task_manager.py b/core/task_manager.py index c1c7014..bf59e13 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -30,6 +30,7 @@ from apscheduler.triggers.cron import CronTrigger 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, adb_connect_light, adb_disconnect from .device_worker import ( @@ -41,8 +42,8 @@ from tasks import list_task_types, get_task_class _log = get_logger("core.tm") -# 错峰启动间隔(秒):批量触发时设备逐个开始占用/连接,避免 adb 全局锁串行堆积 -# 和 STF occupy 并发风暴。100 台 × 0.2s = 20s 全部开始启动。 +# 错峰启动间隔(秒):批量触发时设备逐个开始连接,避免 adb 全局锁串行堆积。 +# 100 台 × 0.2s = 20s 全部开始启动。 _START_STAGGER_SEC = 0.2 @@ -143,7 +144,10 @@ class TaskJob: d.get("params"), d.get("schedule"), d.get("retry"), d.get("enabled", True)) def resolve_serials(self, manager): - """根据 target 解析出实际要跑的 serial 列表。""" + """根据 target 解析出实际要跑的 serial 列表。 + + 数据源:core.device_pool(本地清单 + adb 在线状态),不再查 STF。 + """ mode = self.target.get("mode", "all") if mode == "serial": serials = [self.target["serial"]] @@ -151,37 +155,31 @@ class TaskJob: g = manager.groups.get(self.target.get("group_name")) serials = list(g.serials) if g else [] else: - # all:默认返回所有空闲设备(天然只含在线设备); - # 抢占模式返回全部在线就绪设备(含被占用,执行时抢占) + # all:默认返回池内在线设备(单实例互斥由 _run_with_retry 的 _running 保证); + # 抢占模式返回全部在线设备(含运行中,执行时抢占) if self.params.get("preempt"): try: - return [d["serial"] for d in manager.stf.list_all_devices() - if d.get("present") and d.get("ready")] + return device_pool.list_online() except Exception: return [] try: - return [d["serial"] for d in manager.stf.list_free_devices()] + return device_pool.list_ready() except Exception: return [] - # 离线自动跳过(serial/group 模式):跳过 STF 池里不在线(present=False)或 - # 未就绪(ready=False,provider 刚接入还在初始化)的设备, - # 避免对离线/未就绪设备反复尝试占用后报"设备离线"。默认开启,可在任务编辑器取消勾选。 + # 离线自动跳过(serial/group 模式):跳过本机 adb 不可达的设备, + # 避免反复尝试连接后报"设备离线"。默认开启,可在任务编辑器取消勾选。 if serials and self.params.get("skip_offline", True): try: - devs = manager.stf.list_all_devices() + online = set(device_pool.list_online()) except Exception: - devs = None # STF 查询失败时不过滤,维持原行为 - if devs is not None: - present = {d["serial"] for d in devs if d.get("present")} - ready = {d["serial"] for d in devs if d.get("present") and d.get("ready")} + online = None # 查询失败时不过滤,维持原行为 + if online is not None: kept, skipped = [], [] for s in serials: - if s not in present: - skipped.append((s, "离线")) - elif s not in ready: - skipped.append((s, "未就绪(provider初始化中)")) - else: + if s in online: kept.append(s) + else: + skipped.append((s, "离线(未连接)")) if skipped: detail = ", ".join(f"{s}({why})" for s, why in skipped) _log.info(f"任务 {self.name} 跳过 {len(skipped)} 台设备: {detail}") @@ -206,8 +204,7 @@ class _ForegroundScanner: dumpsys window 只读取窗口状态,不执行任何操作。 """ - def __init__(self, stf): - self.stf = stf + def __init__(self): self._cache = {} # serial -> app_name self._cache_lock = threading.Lock() self._scanning = threading.Event() # 标记是否正在扫描 @@ -245,61 +242,40 @@ class _ForegroundScanner: def _scan_all(self): """扫描所有在线设备(后台线程执行)。 - 按设备归属分四类处理: + 按设备归属分两类处理: 1. worker 运行中:用已有 remote_adb_url 直接查询(无额外开销) - 2. 自己账户占用但无 worker:调用 STF remoteConnect 获取隧道查询 - (不 occupy/release,不打扰设备 UI) - 3. 完全空闲设备(using=False):尝试轻量 adb connect serial - (单次尝试,不 kill-server,不影响其他 worker) - 4. 被他人占用:标记 "(他人占用)" + 2. 空闲设备:返回"空闲"——不主动连接。IP:5555 的 adb transport + 与 STF provider 共享,外部 connect/disconnect 会让 STF 误判 + 设备离线并触发重连(迁移期仍保留此约束,摘除 STF 后可放开) """ self._scanning.set() _log.info("前台 App 扫描已启动") try: - all_devices = self.stf.list_all_devices() + online = set(device_pool.list_online()) except Exception as e: - _log.error("前台 App 扫描: 获取设备列表失败: %s", e) + _log.error("前台 App 扫描: 获取在线设备失败: %s", e) self._scanning.clear() return - # 自己账户已占用的设备列表(用于区分"自己占用"vs"他人占用") - try: - my_serials = {d["serial"] for d in self.stf.list_my_devices()} - except Exception: - my_serials = set() - worker_status = {w["serial"]: w for w in get_all_worker_status()} # 分类设备 have_conn = {} # serial -> remote_adb_url(worker 运行中,已有 adb 连接) - my_owned = [] # 自己账户占用但无 running worker,需 remoteConnect 获取隧道 - free_serials = [] # 完全空闲设备,尝试轻量 adb connect - skip_serials = [] # 被他人占用 + free_serials = [] # 空闲设备 - for dev in all_devices: - if not dev.get("present"): - continue - serial = dev.get("serial", "") - if not serial: - continue + for serial in online: w = worker_status.get(serial, {}) url = w.get("remote_adb_url") if url and w.get("status") in ("running", "connecting"): # 1. 设备正在执行任务,已有 adb 连接,直接查询 have_conn[serial] = url - elif serial in my_serials: - # 2. 自己账户占用但无 worker,可安全 remoteConnect(不打扰设备) - my_owned.append(serial) - elif dev.get("using"): - # 4. 被他人占用 - skip_serials.append(serial) else: - # 3. 完全空闲,尝试轻量 adb connect(不 kill-server) + # 2. 空闲设备 free_serials.append(serial) results = {} - _log.info("前台 App 扫描分类: 运行中=%d, 自己占用=%d, 空闲=%d, 他人占用=%d", - len(have_conn), len(my_owned), len(free_serials), len(skip_serials)) + _log.info("前台 App 扫描分类: 运行中=%d, 空闲=%d", + len(have_conn), len(free_serials)) # 1. worker 运行中设备:用已有 remote_adb_url 查询(并发 10) if have_conn: @@ -313,33 +289,9 @@ class _ForegroundScanner: except Exception: results[s] = None - # 2. 自己占用设备:remoteConnect 获取隧道 → 查询 → 断开(并发 5,不打扰设备) - if my_owned: - with ThreadPoolExecutor(max_workers=min(5, len(my_owned))) as pool: - futures = {pool.submit(self._scan_my_owned, s): s - for s in my_owned} - for fut in as_completed(futures, timeout=30): - s = futures[fut] - try: - results[s] = fut.result() - except Exception: - results[s] = None - - # 3. 空闲设备:轻量 adb connect serial → 查询 → disconnect(并发 10) - if free_serials: - with ThreadPoolExecutor(max_workers=min(10, len(free_serials))) as pool: - futures = {pool.submit(self._scan_free, s): s - for s in free_serials} - for fut in as_completed(futures, timeout=20): - s = futures[fut] - try: - results[s] = fut.result() - except Exception: - results[s] = None - - # 4. 被别人占用的设备 - for s in skip_serials: - results[s] = "(他人占用)" + # 2. 空闲设备:不打扰,直接返回"空闲" + for s in free_serials: + results[s] = "空闲" # 更新缓存 with self._cache_lock: @@ -351,37 +303,6 @@ class _ForegroundScanner: self._scanning.clear() _log.info("前台 App 扫描完成: %d 台设备", len(results)) - def _scan_my_owned(self, serial): - """自己账户占用的设备:通过 STF remoteConnect 获取隧道查询。 - - 不调用 occupy/release,只建立/断开 ADB 隧道,不打扰设备 UI。 - """ - try: - url = self.stf.remote_connect(serial) - if not url: - return None - try: - if not adb_connect_light(url): - return None - return get_foreground_app(url) - finally: - adb_disconnect(url) - self.stf.remote_disconnect(serial) - except Exception: - return None - - def _scan_free(self, serial): - """空闲设备:不扫描前台 App,直接返回"空闲"。 - - 原因:STF provider 内部通过 IP:5555 维持 adb 连接监控设备。 - 外部 adb connect/disconnect IP:5555 会让 adb server 断开该地址的 - transport,连带 STF provider 的连接一起断,STF 误判设备 offline - 并触发重连——表现就是"一扫描前台 App 设备就离线、需要重连"。 - 因此空闲设备绝不主动 adb connect/disconnect,只对运行中设备 - (走 STF 隧道端口,不动 5555)和自占设备(remoteConnect 隧道)扫描。 - """ - return "空闲" - # ================== 任务管理器 ================== class TaskManager: @@ -396,7 +317,7 @@ class TaskManager: self._running = {} # serial -> {"worker", "job_id", "started_at", "attempt"} self._stop_requested = set() # serial 集合:用户请求停止,阻止后续重试 self._lock = threading.Lock() - self._fg_scanner = _ForegroundScanner(self.stf) + self._fg_scanner = _ForegroundScanner() self._load() def _db(self): @@ -837,15 +758,15 @@ class TaskManager: with self._lock: return {s: dict(v) for s, v in self._running.items()} - # ---- 状态查询(带缓存,避免 STF 请求阻塞前端)---- + # ---- 状态查询(带缓存,避免设备列表查询阻塞前端)---- _status_cache = None # (timestamp, data, error) _status_cache_lock = threading.Lock() _STATUS_CACHE_TTL = 5.0 # 缓存 5 秒,前端 5 秒刷新刚好命中 def get_status(self): - """综合状态:STF 设备池 + 本地 worker + 运行中的任务。 + """综合状态:设备池(清单+在线)+ 本地 worker + 运行中的任务。 - 带 5 秒缓存:STF 请求慢时避免每次 /api/status 都打 STF 阻塞 Flask。 + 带 5 秒缓存:避免每次 /api/status 都查库/adb 阻塞 Flask。 worker 状态实时读(内存,无 IO),不受缓存影响。 """ with self._status_cache_lock: @@ -856,46 +777,53 @@ class TaskManager: return None, err # 用缓存的设备列表 + 实时 worker 状态重新组装 return self._merge_status(cached), None - # 缓存过期或不存在,重新拉 STF + # 缓存过期或不存在,重新拉设备池 try: - all_devices = self.stf.list_all_devices() + configured = device_pool.list_configured() except Exception as e: with self._status_cache_lock: self._status_cache = (time.time(), None, str(e)) return None, f"获取设备列表失败: {e}" with self._status_cache_lock: - self._status_cache = (time.time(), all_devices, None) - return self._merge_status(all_devices), None + self._status_cache = (time.time(), configured, None) + return self._merge_status(configured), None - def _merge_status(self, all_devices): - """用 STF 设备列表 + 实时 worker 状态 + 前台 App 组装返回结果。""" + def _merge_status(self, configured_serials): + """用设备池清单 + 实时 worker 状态 + 前台 App 组装返回结果。""" worker_status = {w["serial"]: w for w in get_all_worker_status()} running = self.get_running() + try: + online = set(device_pool.list_online()) + except Exception: + online = set() + try: + names = {d["serial"]: d["name"] for d in device_pool.list_devices()} + except Exception: + names = {} - # 清理陈旧状态:serial 已不在 STF 设备池、且没有在跑 worker 的条目, + # 清理陈旧状态:serial 已不在设备池、且没有在跑 worker 的条目, # 避免设备被删除后其失败记录仍残留在"异常汇总"里 - present = {dev.get("serial", "") for dev in all_devices} + configured = set(configured_serials) for serial, w in list(worker_status.items()): - if serial in present or serial in running: + if serial in configured or serial in running: continue if w.get("status") in ("running", "connecting"): continue _remove_worker(serial) result = [] - for dev in all_devices: - serial = dev.get("serial", "") - owner = dev.get("owner") + for serial in configured_serials: + is_online = serial in online w = worker_status.get(serial, {}) r = running.get(serial, {}) result.append({ "serial": serial, - "model": dev.get("model") or w.get("model", "") or dev.get("product", ""), - "device_name": dev.get("name") or "", - "present": dev.get("present", False), - "ready": dev.get("ready", False), - "stf_occupied": dev.get("using", False), - "owner": owner.get("name", "") if owner else "", + "model": w.get("model", ""), + "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), # 通用进度字段(任意 app 通用,前端统一解析展示) @@ -912,14 +840,11 @@ class TaskManager: return result def list_all_serials(self): - """返回 STF 上所有在线设备的 serial 列表(供分组表单勾选用)。 - - 复用 get_status 缓存,避免开页面时阻塞。 - """ - devices, err = self.get_status() - if err or not devices: + """返回设备池在线设备的 serial 列表(供分组表单勾选用)。""" + try: + return device_pool.list_ready() + except Exception: return [] - return [d["serial"] for d in devices if d.get("present")] def shutdown(self): self.stop_all() diff --git a/doc/STF_REMOVAL.md b/doc/STF_REMOVAL.md index 1c61517..0c3f446 100644 --- a/doc/STF_REMOVAL.md +++ b/doc/STF_REMOVAL.md @@ -90,8 +90,8 @@ STF 当前仅提供:occupy/release 互斥、present+ready 健康信号、设 |---|---|---| | `task_manager.resolve_serials` (145-183) | `stf.list_free_devices()` / `list_all_devices()` | `device_pool.list_ready()`;preempt 模式 = 全部 list_online | | `resolve_serials` skip_offline 过滤 (166-183) | STF present/ready 判定 | `device_pool.is_online()`;跳过原因文案"离线/未连接" | -| `device_worker.STFDevice.acquire` (45-70) | `stf.occupy()` + 直连/桥接 | 去掉 occupy 与桥接分支(IP:port 直连保留,USB 序列号仅提示不支持);保留 adb_connect 重试 + 2s 等待 | -| `STFDevice.release` (78-93) | `stf.release()` | 置空(互斥由 `_running[serial]` 负责,本来就是内存锁) | +| `device_worker.STFDevice.acquire` (45-70) | `stf.occupy()` + 直连/桥接 | 去掉 occupy(互斥由 `_running[serial]` 负责);**USB 桥接保留到阶段 2**(避免迁移期 USB 断档,阶段 2 换成 220 adb server);保留 adb_connect 重试 + 2s 等待 | +| `STFDevice.release` (78-93) | `stf.release()` | 置空(仅保留 USB 隧道断开) | | `STFDevice._pick_free` (72-76) | `list_free_devices` | `device_pool.list_ready()` | | `_ForegroundScanner` (234-386) | `list_all_devices` / `list_my_devices` / remote_connect | `device_pool.list_online()`;桥接分支删除 | | `web_server.api_device_screen_all` (646) | STF present 列表 | `device_pool.list_online()` | diff --git a/static/admin/apps.js b/static/admin/apps.js index ce730e0..49e4067 100644 --- a/static/admin/apps.js +++ b/static/admin/apps.js @@ -167,7 +167,7 @@ async function openInstallModal(apkId){ '