"""任务域 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 core.logger import get_logger _log = get_logger("web") 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 def _job_coverage(job, configured): """任务覆盖的设备(监控页「任务运行概况」展示用)。 与 TaskJob.resolve_serials 的差别:这里按 **target 定义**给出「这台任务会覆盖 哪些设备」,不因设备当前是否空闲而变(在线/运行中由前端拿实时状态标注), 也不做离线过滤、不写调度日志——它是每 5s 轮询的展示接口。 """ tgt = job.target or {} mode = tgt.get("mode", "all") try: if mode == "serial": serials = [tgt.get("serial") or ""] elif mode == "group": g = context.mgr.groups.get(tgt.get("group_name")) serials = [s for s in (list(g.serials) if g else []) if s in configured] else: serials = sorted(configured) serials = [s for s in serials if s] return {"mode": mode, "serials": serials, "total": len(serials)} except Exception as e: _log.warning(f"任务「{job.name}」覆盖设备解析失败: {e}") return {"mode": mode, "serials": [], "total": 0, "error": str(e)[:80]} def _configured_set(): """设备池清单(coverage 计算用,取一次给所有任务复用)。""" try: return set(device_pool.list_configured()) except Exception: return set() @bp.route("/api/jobs") @login_required def api_jobs_list(): configured = _configured_set() jobs = [] for j in context.mgr.jobs.values(): d = j.to_dict() d["next_run"] = _job_next_run(j) d["coverage"] = _job_coverage(j, configured) 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", "generic_steps") 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) d["coverage"] = _job_coverage(job, _configured_set()) return jsonify({"ok": True, "msg": "任务已创建", "job": d}) @bp.route("/api/jobs/", 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) d["coverage"] = _job_coverage(job, _configured_set()) return jsonify({"ok": True, "msg": "任务已更新", "job": d}) @bp.route("/api/jobs/", 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//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//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/", 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/", 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/", 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/", 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,前端用 加载。 """ 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 下载在国内不可靠,改自建)。 # 坐标约定:前端以 显示尺寸归一化后映射到设备原生分辨率(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