feat: 新增通用步骤任务+元素抓取,完善全部项目文档
- 新增 generic_steps 通用步骤任务(可视化步骤编辑器编排流程,支持 open_app/click/swipe/input_text/wait/loop/group) - 新增 core/uiauto_helper.py 封装 uiautodev 元素抓取客户端 - web_server 新增元素抓取/截图/设备列表等 API - monitor.html 新增步骤编辑器、元素抓取模态框、独立关闭逻辑 - 新增 README.md 项目总览(快速上手/架构/配置/FAQ) - 新增 doc/ARCHITECTURE.md 架构详解、doc/DEPLOY.md 部署指南、doc/API.md 接口文档 - 修复 doc/TASK_DEV.md:移除已删除的 comment 引用,补充 generic 包,更新注册示例 - .gitignore 忽略 .claude/ 工具产物
This commit is contained in:
+96
-71
@@ -302,6 +302,7 @@ class TaskManager:
|
||||
self.groups = {} # name -> DeviceGroup(内存业务对象)
|
||||
self.jobs = {} # id -> TaskJob(内存业务对象)
|
||||
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._load()
|
||||
@@ -557,88 +558,112 @@ class TaskManager:
|
||||
异常分类:
|
||||
DeviceOfflineError — 设备掉线,立即放弃不重试(换设备也没用)
|
||||
其他异常 — 按 max_attempts 重试
|
||||
用户停止(stop_device)会加入 _stop_requested,阻止任何后续重试。
|
||||
"""
|
||||
for attempt in range(1, max_attempts + 1):
|
||||
# 同一 serial 同时只能一个 worker
|
||||
try:
|
||||
for attempt in range(1, max_attempts + 1):
|
||||
# 用户已请求停止 → 不再启动新 attempt
|
||||
with self._lock:
|
||||
if serial in self._stop_requested:
|
||||
_log.info(f"{serial} 用户已请求停止,取消重试 (job={job.name})")
|
||||
return
|
||||
# 同一 serial 同时只能一个 worker
|
||||
with self._lock:
|
||||
if serial in self._running:
|
||||
_log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})")
|
||||
return
|
||||
self._running[serial] = {"job_id": job.id, "started_at": time.time(),
|
||||
"attempt": attempt, "task_type": job.task_type}
|
||||
_update_status(serial, task_job=job.name, attempt=attempt,
|
||||
max_attempts=max_attempts)
|
||||
|
||||
worker = None
|
||||
try:
|
||||
worker = task.create_worker(self.stf, serial, job.params)
|
||||
with self._lock:
|
||||
self._running[serial]["worker"] = worker
|
||||
_log.info(f"{serial} 开始任务 {job.name} (第{attempt}/{max_attempts}次)")
|
||||
worker.start()
|
||||
worker.join() # 等待 worker 结束
|
||||
# worker 正常结束(done 或被 stop)
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
# 判断是否成功:看 status
|
||||
with _WORKERS_LOCK:
|
||||
st = _WORKERS.get(serial, {}).get("status")
|
||||
# 清除 task_job 标记,避免前端误判"运行中"
|
||||
_update_status(serial, task_job="")
|
||||
if st == "done":
|
||||
_log.info(f"{serial} 任务 {job.name} 成功完成")
|
||||
return
|
||||
# 用户请求停止(无论 attempt 第几次、status 是什么)→ 不重试
|
||||
with self._lock:
|
||||
if serial in self._stop_requested:
|
||||
_log.info(f"{serial} 用户已请求停止,不再重试 (job={job.name})")
|
||||
return
|
||||
_log.warning(f"{serial} 任务未成功(status={st})")
|
||||
except DeviceOfflineError as e:
|
||||
# 设备掉线,立即放弃,不重试
|
||||
_log.error(f"{serial} 设备离线,放弃任务 {job.name}: {e}")
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
_update_status(serial, status="failed", task_job="",
|
||||
last_error=f"设备离线: {e}")
|
||||
return
|
||||
except Exception as e:
|
||||
_log.error(f"{serial} 执行异常: {e}", exc_info=True)
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
_update_status(serial, task_job="")
|
||||
|
||||
if attempt < max_attempts:
|
||||
# 用户在 sleep 期间点停止也能中断
|
||||
with self._lock:
|
||||
if serial in self._stop_requested:
|
||||
_log.info(f"{serial} 用户已请求停止,取消重试 (job={job.name})")
|
||||
return
|
||||
# 检查是否是端口耗尽类临时错误,需要更长退避等端口释放
|
||||
with _WORKERS_LOCK:
|
||||
err = _WORKERS.get(serial, {}).get("last_error", "")
|
||||
if err.startswith("[transient]"):
|
||||
# Windows TCP 端口耗尽,TIME_WAIT 默认 2-4 分钟,等 120 秒
|
||||
extra_delay = max(delay, 120)
|
||||
_log.info(f"{serial} ADB 连接临时错误(端口耗尽),{extra_delay}s 后重试 ({attempt+1}/{max_attempts})")
|
||||
time.sleep(extra_delay)
|
||||
else:
|
||||
_log.info(f"{serial} {delay}s 后重试 ({attempt+1}/{max_attempts})")
|
||||
time.sleep(delay)
|
||||
|
||||
_log.error(f"{serial} 任务 {job.name} 重试耗尽,放弃")
|
||||
_update_status(serial, status="failed", last_error=f"{job.name} 重试{max_attempts}次失败")
|
||||
finally:
|
||||
# 清除停止标志:整个重试循环结束(成功/失败/停止)后允许下次任务
|
||||
with self._lock:
|
||||
if serial in self._running:
|
||||
_log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})")
|
||||
return
|
||||
self._running[serial] = {"job_id": job.id, "started_at": time.time(),
|
||||
"attempt": attempt, "task_type": job.task_type}
|
||||
_update_status(serial, task_job=job.name, attempt=attempt,
|
||||
max_attempts=max_attempts)
|
||||
|
||||
worker = None
|
||||
try:
|
||||
worker = task.create_worker(self.stf, serial, job.params)
|
||||
with self._lock:
|
||||
self._running[serial]["worker"] = worker
|
||||
_log.info(f"{serial} 开始任务 {job.name} (第{attempt}/{max_attempts}次)")
|
||||
worker.start()
|
||||
worker.join() # 等待 worker 结束
|
||||
# worker 正常结束(done 或被 stop)
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
# 判断是否成功:看 status
|
||||
with _WORKERS_LOCK:
|
||||
st = _WORKERS.get(serial, {}).get("status")
|
||||
# 清除 task_job 标记,避免前端误判"运行中"
|
||||
_update_status(serial, task_job="")
|
||||
if st == "done":
|
||||
_log.info(f"{serial} 任务 {job.name} 成功完成")
|
||||
return
|
||||
if st == "released" and attempt == 1:
|
||||
# 被手动停止,不重试
|
||||
return
|
||||
_log.warning(f"{serial} 任务未成功(status={st})")
|
||||
except DeviceOfflineError as e:
|
||||
# 设备掉线,立即放弃,不重试
|
||||
_log.error(f"{serial} 设备离线,放弃任务 {job.name}: {e}")
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
_update_status(serial, status="failed", task_job="",
|
||||
last_error=f"设备离线: {e}")
|
||||
return
|
||||
except Exception as e:
|
||||
_log.error(f"{serial} 执行异常: {e}", exc_info=True)
|
||||
with self._lock:
|
||||
self._running.pop(serial, None)
|
||||
_update_status(serial, task_job="")
|
||||
|
||||
if attempt < max_attempts:
|
||||
# 检查是否是端口耗尽类临时错误,需要更长退避等端口释放
|
||||
with _WORKERS_LOCK:
|
||||
err = _WORKERS.get(serial, {}).get("last_error", "")
|
||||
if err.startswith("[transient]"):
|
||||
# Windows TCP 端口耗尽,TIME_WAIT 默认 2-4 分钟,等 120 秒
|
||||
extra_delay = max(delay, 120)
|
||||
_log.info(f"{serial} ADB 连接临时错误(端口耗尽),{extra_delay}s 后重试 ({attempt+1}/{max_attempts})")
|
||||
time.sleep(extra_delay)
|
||||
else:
|
||||
_log.info(f"{serial} {delay}s 后重试 ({attempt+1}/{max_attempts})")
|
||||
time.sleep(delay)
|
||||
|
||||
_log.error(f"{serial} 任务 {job.name} 重试耗尽,放弃")
|
||||
_update_status(serial, status="failed", last_error=f"{job.name} 重试{max_attempts}次失败")
|
||||
self._stop_requested.discard(serial)
|
||||
|
||||
# ---- 运行控制 ----
|
||||
def stop_device(self, serial):
|
||||
"""停止指定设备的 worker。"""
|
||||
"""停止指定设备的 worker,并阻止后续重试。
|
||||
|
||||
无论 worker 当前在运行还是在重试 sleep 中,都会阻止下一次重试。
|
||||
"""
|
||||
with self._lock:
|
||||
self._stop_requested.add(serial)
|
||||
info = self._running.get(serial)
|
||||
if not info:
|
||||
return False
|
||||
w = info.get("worker")
|
||||
if w and w.is_alive():
|
||||
w.stop()
|
||||
return True
|
||||
return False
|
||||
if info:
|
||||
w = info.get("worker")
|
||||
if w and w.is_alive():
|
||||
w.stop()
|
||||
return True
|
||||
# worker 已结束但重试循环可能还在 sleep —— 仍然返回 True 表示已阻止重试
|
||||
return True
|
||||
|
||||
def stop_all(self):
|
||||
"""停止所有运行中的 worker。"""
|
||||
"""停止所有运行中的 worker,并阻止后续重试。"""
|
||||
with self._lock:
|
||||
items = list(self._running.items())
|
||||
for s, _ in items:
|
||||
self._stop_requested.add(s)
|
||||
stopped = []
|
||||
for serial, info in items:
|
||||
w = info.get("worker")
|
||||
|
||||
Reference in New Issue
Block a user