feat(AI 建任务): AI 自己在真机探索 → 写出可调度任务 → 人工确认入库

AI 控制台下新增子分栏「🧭 AI 建任务」:描述要做什么(例:建一个跑 2 小时的任务、自动刷
某 App、随机点赞),AI 用 de_* 工具自己在设备上探索(看屏/读元素树/点按验证),把走通的
路径写成一条任务草稿,经服务端校验后交人在步骤编辑器里核对/手改/试跑,保存才入库。

链路:POST /api/agent/run{mode:"designer", settings}
  → Agent 自探 → 本地工具 submit_task(draft)
  → core/task_draft 校验(失败把 errors 回灌模型让它改)
  → 只暂存(运行态 + app_meta.agent_task_draft,**不落库**)
  → SSE done{mode,draft,warnings} → 页面草稿预览 → openTaskModal(null, prefill) 预填

关键实现
- core/task_draft.py(新):把执行器的"静默跳过点"(未知 type/空 selector/空 children/
  嵌套>5/节点>60/cron 非法/必填缺失)前移成显式 error——POST /api/jobs 对 params 是盲存的,
  执行器又静默跳过错误步骤,没有这道闸门就是"任务建好了、跑起来什么都没做"。
  归一化兜底任务名/target/schedule/retry/时长;页面填的设置以 overrides 优先于模型。
  故意**不比执行器更严**:loop_mode 近义值归一(count→rounds)、缺 max_iterations 补默认 10
  (执行器本来就默认)——实测卡太死会把一轮探索耗在改字段上。
  有副作用的步骤(评论/发送/购买/删除…)只警告并把触发概率压到 30%(编辑器可改回)。
- mcp_agent/agent.py:双系统提示词(CHAT/DESIGNER)+ 平台级本地工具
  (LOCAL_TOOL_SPECS,不进 MCP)+ 每工具调用上限 40 + designer 输出上限 8192 +
  **json.loads 容错**(草稿被截断时给模型可读错误,而不是整轮失败)。
- web/agent_api.py:mode/settings 透传、submit_task 处理器(app_context 内校验+暂存)、
  done 带 draft、GET /api/agent/task_draft{,+POST,/clear}(草稿走 app_meta,不新建表)。
- 前端:static/admin/taskgen.js + #agent-sub-taskgen 子面板(showSubTab 机制);
  tasks.js 的 openTaskModal(jobId, prefill) + 信封归一化 + 唯一 draftKey;
  agent.js 按 mode 门控(一个 run 只有一个事件队列,两个 EventSource 会互相瓜分事件)。

顺带修掉一个 chat 也踩的协议 bug:一轮里同时调 de_screenshot 与别的工具时,截图图像会被
插在两条 tool 消息之间 → 模型侧判"工具回应不足"直接 400。改为本轮 tool 消息发完再附图像,
_repair_tool_messages 也改成只数**连续**的 tool 消息。

真机实测(Redmi 22120RN86C,设置页):8 步探索(含 tap_text 验证)→ submit_task 一次通过 →
草稿 8 个顶层步骤(screen_on/open_app/wait_el/click/wait/key_event…)、max_duration 1800、
无 click_xy、3 条 evidence;页面恢复草稿 + 预填编辑器 + 提示块渲染均正常,无 JS 报错。
自测数据已清理(草稿已丢弃、未创建任何任务)。

文档:AI_TASK_GEN.md 状态改「P0 已实现」+ §10 实现记录(差异/护栏/未做项)、AI_CONSOLE.md
(子分栏、designer 分支、SSE done 负载、app_meta 键)、API.md、DATA_MODEL.md、ARCHITECTURE.md、
DEVELOPMENT.md(自测入口)、README.md 索引、backlog 勾掉 P0。
This commit is contained in:
2026-09-13 23:08:59 +08:00
parent 0bc713137d
commit 46e6ea1f37
16 changed files with 1707 additions and 47 deletions
+258 -11
View File
@@ -43,10 +43,20 @@ _CFG_KEYS = {"api_base": "agent_api_base",
# 推理链落库上限(字符):只留够回看的量,避免会话消息无限膨胀
_REASONING_KEEP = 6000
# AI 建任务:草稿落库的 app_meta 键(**不新建表** —— 只存最近一份,
# 刷新/重进页面能拿回来;见 doc/AI_TASK_GEN.md)
_DRAFT_META_KEY = "agent_task_draft"
# 草稿 JSON 体积上限(app_meta.value 是 TEXT):超了截断 notes/evidence,不报错
_DRAFT_MAX_CHARS = 60000
# ---------- 运行状态(单实例 + 事件队列) ----------
_run = {"id": None, "state": "idle", "prompt": "", "serial": "",
"answer": "", "error": "", "usage": {},
"mode": "chat", # chat(聊天)| designer(AI 建任务)
"draft": None, # designer:submit_task 校验通过的任务草稿
"draft_error": "", # designer:草稿被拦下的原因(errors 拼接)
"warnings": [], # designer:草稿的提醒项(不拦,前端醒目展示)
"history": []} # 多轮对话历史 [{role: user|assistant, content}]
_queues = {} # run_id -> queue.Queue(SSE 消费者读取)
_stop_events = {} # run_id -> threading.Event(用户中断)
@@ -1101,13 +1111,26 @@ def agent_conversations_rename(conv_id):
@bp.route("/api/agent/run", methods=["POST"])
@admin_required
def agent_run():
"""启动 Agent:{prompt, serial?, conversation_id?}。运行中返回 409。"""
"""启动 Agent:{prompt, serial?, conversation_id?, mode?, settings?}。运行中返回 409。
mode=chat(默认):AI 控制台聊天(现状不变)
mode=designer:AI 建任务 —— 探索设备 → 用 submit_task 提交任务草稿(不落库)
settings:建任务页上用户填的任务设置(任务名/目标/调度),作为草稿的 overrides
"""
data = request.json or {}
prompt = (data.get("prompt") or "").strip()
if not prompt:
return jsonify({"ok": False, "error": "请输入指令"}), 400
serial = (data.get("serial") or "").strip()
conv_id = (data.get("conversation_id") or "").strip()
mode = (data.get("mode") or "chat").strip()
if mode not in ("chat", "designer"):
return jsonify({"ok": False, "error": f"未知模式: {mode}"}), 400
settings = data.get("settings") if isinstance(data.get("settings"), dict) else {}
# 建任务模式是单轮的:不绑会话(避免把聊天历史灌进设计师上下文,
# 也避免探索过程污染聊天会话)
if mode == "designer":
conv_id = ""
cfg = _read_cfg()
if not cfg.get("api_key"):
return jsonify({"ok": False, "error": "请先在配置区填写 API Key"}), 400
@@ -1156,15 +1179,17 @@ def agent_run():
run_id = uuid.uuid4().hex[:8]
from datetime import datetime as _dt
_run.update(id=run_id, state="running", prompt=prompt,
serial=serial, conv_id=conv_id,
serial=serial, conv_id=conv_id, mode=mode,
started=_dt.now().strftime("%H:%M:%S"),
answer="", error="", usage={})
answer="", error="", usage={},
draft=None, draft_error="", warnings=[])
# history 保留(多轮上下文),由会话/「新建会话」管理
_queues[run_id] = queue.Queue()
_stop_events[run_id] = threading.Event()
_log.info(f"Agent 启动: {prompt[:60]} @ {serial or 'default'} conv={conv_id or '-'}")
_log.info(f"Agent 启动[{mode}]: {prompt[:60]} @ {serial or 'default'} "
f"conv={conv_id or '-'}")
threading.Thread(target=_agent_thread,
args=(run_id, prompt, serial, cfg),
args=(run_id, prompt, serial, cfg, mode, settings),
daemon=True).start()
return jsonify({"ok": True, "run_id": run_id})
@bp.route("/api/agent/config", methods=["GET"])
@@ -1211,6 +1236,10 @@ def agent_run_status():
"answer": (_run.get("answer") or "")[:4000],
"error": (_run.get("error") or "")[:400],
"usage": _run.get("usage") or {},
"mode": _run.get("mode") or "chat",
"draft": _run.get("draft"),
"draft_error": _run.get("draft_error") or "",
"warnings": _run.get("warnings") or [],
"history": hist[-16:]})
@@ -1291,6 +1320,110 @@ def agent_clear():
return jsonify({"ok": True, "msg": "已清空"})
# ================== AI 建任务:草稿 ==================
# 草稿只存**最近一份**(app_meta 单键),用途是"刷新/重进页面能拿回来"与审计;
# 它不是任务 —— 真正入库仍要用户在步骤编辑器里确认后走 POST /api/jobs。
def _load_draft_meta():
"""读回最近一份草稿(无则 None)。"""
raw = _meta_get(_DRAFT_META_KEY)
if not raw:
return None
try:
obj = json.loads(raw)
return obj if isinstance(obj, dict) else None
except (json.JSONDecodeError, TypeError):
_log.warning("草稿 JSON 解析失败,忽略")
return None
def _draft_env():
"""校验草稿要用的环境信息:现有分组名 + 设备池(取不到就跳过对应校验)。"""
groups = pool = None
try:
from core import device_pool
from web import context
groups = set(context.mgr.groups.keys())
pool = set(device_pool.list_configured())
except Exception as e:
_log.warning(f"分组/设备池读取失败,跳过对应校验: {e}")
return groups, pool
def _persist_draft(draft, warnings, prompt=""):
"""把草稿存进 app_meta(体积超限时逐级瘦身,不报错)。"""
from datetime import datetime as _dt
saved = {"draft": draft, "warnings": list(warnings or []),
"prompt": (prompt or "")[:500],
"created": _dt.now().strftime("%Y-%m-%d %H:%M")}
raw = json.dumps(saved, ensure_ascii=False)
if len(raw) > _DRAFT_MAX_CHARS:
saved["draft"]["evidence"] = []
raw = json.dumps(saved, ensure_ascii=False)
if len(raw) > _DRAFT_MAX_CHARS:
saved["draft"]["notes"] = list(saved["draft"].get("notes") or [])[:2]
raw = json.dumps(saved, ensure_ascii=False)
if len(raw) > _DRAFT_MAX_CHARS:
_log.warning(f"草稿体积仍超限({len(raw)} 字符),照原样存入")
_meta_put(_DRAFT_META_KEY, raw)
try:
db.session.commit()
except Exception:
db.session.rollback()
return saved
@bp.route("/api/agent/task_draft")
@admin_required
def agent_task_draft_get():
"""读回最近一次 AI 建任务的草稿(页面刷新/重进时恢复用)。"""
saved = _load_draft_meta()
with _lock:
cur = _run.get("draft")
st = _run.get("state")
mode = _run.get("mode") or "chat"
run_id = _run.get("id") or "" if st == "running" else ""
return jsonify({"ok": True, "saved": saved,
"running": st == "running", "mode": mode, "run_id": run_id,
"draft": cur, "draft_error": _run.get("draft_error") or ""})
@bp.route("/api/agent/task_draft", methods=["POST"])
@admin_required
def agent_task_draft_save():
"""回存一份草稿(人工在页面上改过之后)。会重新校验一次。"""
from core.task_draft import validate_draft
data = request.json or {}
groups, pool = _draft_env()
res = validate_draft(data.get("draft"), groups=groups, pool=pool,
default_serial=(data.get("serial") or "").strip())
if not res["ok"]:
return jsonify({"ok": False, "error": "草稿校验未通过",
"errors": res["errors"]}), 400
saved = _persist_draft(res["draft"], res["warnings"],
prompt=data.get("prompt") or "")
with _lock:
_run["draft"] = res["draft"]
_run["warnings"] = res["warnings"]
_run["draft_error"] = ""
return jsonify({"ok": True, "msg": "草稿已保存", "saved": saved})
@bp.route("/api/agent/task_draft/clear", methods=["POST"])
@admin_required
def agent_task_draft_clear():
"""丢弃当前草稿。"""
_meta_put(_DRAFT_META_KEY, "")
try:
db.session.commit()
except Exception:
db.session.rollback()
with _lock:
_run["draft"] = None
_run["warnings"] = []
_run["draft_error"] = ""
return jsonify({"ok": True, "msg": "已丢弃草稿"})
def _shrink_image(b64, width=220, quality=50):
"""截图降采样(SSE step 事件用,控制传输体积)。失败原样返回。"""
try:
@@ -1305,9 +1438,15 @@ def _shrink_image(b64, width=220, quality=50):
return b64
def _agent_thread(run_id, prompt, serial, cfg):
"""后台线程:Agent 流式执行,事件推入队列供 SSE 消费。"""
def _agent_thread(run_id, prompt, serial, cfg, mode="chat", settings=None):
"""后台线程:Agent 流式执行,事件推入队列供 SSE 消费。
mode=designer 时额外注册平台级工具 `submit_task`:它把 AI 探索出的任务草稿
交给 `core/task_draft` 校验,**通过也只暂存、不落库**(人工确认后才入库)。
"""
q = _queues.get(run_id)
settings = settings if isinstance(settings, dict) else {}
designer = (mode == "designer")
try:
# 平台进程 cwd=/app(含 mcp_agent 包),进程内 import
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
@@ -1339,6 +1478,12 @@ def _agent_thread(run_id, prompt, serial, cfg):
"args": step.get("args") or {}}
if step.get("image"):
rec["image"] = _shrink_image(step["image"])
# 工具报错时把原因带出去:前端卡片直接显示"这步为什么失败"
# (模型看得到 errors,人也要看得到)
r = step.get("result")
if isinstance(r, dict) and r.get("ok") is False:
rec["error"] = str(r.get("error")
or ";".join(r.get("errors") or []) or "失败")[:300]
q.put(("step", rec))
# 记录精简工具序列 + 结构化轨迹(去 serial、结果截断)
try:
@@ -1381,6 +1526,76 @@ def _agent_thread(run_id, prompt, serial, cfg):
history = list(_run.get("history") or [])
target = serial or cfg.get("default_serial") or ""
stop_evt = _stop_events.get(run_id)
if designer:
# 建任务模式是单轮的:不吃聊天历史(避免把聊天上下文灌进设计师)
with _lock:
_run["history"] = []
history = []
async def submit_task(args):
"""平台级工具:校验 AI 提交的任务草稿。
失败 → 把 errors 原样回给模型,让它逐条修正后重提(这是"AI 写坏任务"
的唯一闸门:POST /api/jobs 对 params 是盲存的,执行器又静默跳过错误步骤)。
成功 → 只暂存在运行态 + app_meta,等用户在步骤编辑器里确认后才入库。
"""
from core.task_draft import validate_draft
def _work():
"""校验 + 暂存(读分组/设备池、写 app_meta 都要 app context)。"""
groups, pool = _draft_env()
r = validate_draft(args, groups=groups, pool=pool,
default_serial=target, overrides=settings)
if r["ok"]:
_persist_draft(r["draft"], r["warnings"], prompt=prompt)
return r
try:
if _flask_app is not None:
with _flask_app.app_context():
res = _work()
else:
res = _work()
except Exception as e:
# 后台线程里任何异常都不能冒泡成"工具执行失败"这种含糊文案——
# 模型与用户都需要知道草稿到底存没存下
_log.exception("[designer] 草稿处理失败")
with _lock:
_run["draft_error"] = f"草稿处理失败: {e}"
q.put(("step", {"tool": "📝 草稿未存下",
"args": f"服务端处理草稿时出错:{type(e).__name__}: {str(e)[:160]}",
"image": None}))
return {"ok": False,
"error": f"服务端处理草稿时出错({type(e).__name__}: {e}),"
"草稿没有保存。请把 steps 精简后重试一次;"
"若仍失败,说明是平台问题,先向用户说明。"}
if not res["ok"]:
msg = ";".join(res["errors"][:6])
with _lock:
_run["draft_error"] = msg
q.put(("step", {"tool": "📝 草稿被拦下",
"args": f"{len(res['errors'])} 处问题需要修正:{msg}",
"image": None}))
return {"ok": False,
"error": "草稿未通过校验,请按下面的 errors 逐条修正后"
"**重新调用一次 submit_task**",
"errors": res["errors"], "warnings": res["warnings"]}
draft = res["draft"]
n_steps = len(((draft.get("task") or {}).get("params") or {}).get("steps") or [])
with _lock:
_run["draft"] = draft
_run["warnings"] = res["warnings"]
_run["draft_error"] = ""
q.put(("step", {"tool": "📝 任务草稿",
"args": f"已提交草稿「{draft['task'].get('name', '')}」:"
f"{n_steps} 个顶层步骤"
+ (f",{len(res['warnings'])} 条待复核提醒"
if res["warnings"] else ""),
"image": None}))
return {"ok": True,
"msg": "草稿已收到并通过校验,已交给用户确认。请立即停止调用工具,"
"用一句中文总结这条任务做什么、哪些步骤需要人工复核。",
"warnings": res["warnings"]}
# 经验检索:相似历史任务的操作配方注入 system(自进化记忆)
exp_ctx, exp_items = _find_experiences(prompt)
@@ -1403,6 +1618,25 @@ def _agent_thread(run_id, prompt, serial, cfg):
if act_ctx:
recall_ctx += ("\n\n## 可复用动作(优先按其中的元素定位操作;"
"若与当前界面不符,再自行截图确认)\n" + act_ctx)
if designer and settings:
# 用户在页面上填的任务设置:必须落到 draft 里,别让模型自己编
hints = []
if settings.get("name"):
hints.append(f"- 任务名:{settings['name']}")
mode_ = settings.get("target_mode")
if mode_ == "group" and settings.get("group_name"):
hints.append(f"- 目标:分组「{settings['group_name']}」")
elif mode_ == "all":
hints.append("- 目标:全部空闲设备")
elif mode_ == "serial":
hints.append(f"- 目标:指定设备 {settings.get('serial') or target}")
if settings.get("schedule_hint"):
hints.append(f"- 调度:{settings['schedule_hint']}")
if settings.get("max_duration"):
hints.append(f"- 运行时长上限:{settings['max_duration']} 秒")
if hints:
recall_ctx += ("\n\n## 用户已指定的任务设置(必须原样写进 draft.task)\n"
+ "\n".join(hints))
mcp_url = getattr(getattr(agent, "s", None), "mcp_url", "")
@@ -1417,6 +1651,9 @@ def _agent_thread(run_id, prompt, serial, cfg):
"network is unreachable", "connect timeout"))
async def _execute():
# 建任务模式:先注册平台级工具(必须在 _load_tools 之前,schema 在加载时组装)
if designer:
agent.register_local_tool("submit_task", submit_task)
try:
await agent._load_tools()
except Exception as e:
@@ -1433,7 +1670,8 @@ def _agent_thread(run_id, prompt, serial, cfg):
should_stop=lambda: bool(
stop_evt and stop_evt.is_set()),
extra_context=recall_ctx,
on_usage=on_usage)
on_usage=on_usage,
mode=mode)
except Exception as e:
if _mcp_unreachable(e):
raise RuntimeError(
@@ -1475,10 +1713,13 @@ def _agent_thread(run_id, prompt, serial, cfg):
except Exception as e:
_log.warning(f"会话落库失败: {e}")
# 自进化:成功执行过工具则提炼配方写入经验。必须在 done 之前完成——
# 自进化:AI 建任务(designer)**不沉淀**——探索轨迹是为"写出一条任务"服务的,
# 不是一次成功操作套路;沉淀它会把探索期的误点/试错当成经验,污染记忆库。
# (把"草稿→经验/动作"作为独立里程碑,见 doc/AI_TASK_GEN.md §6 P1)
# 必须在 done 之前完成——
# done 发出后 SSE 关流,用户就看不到「已写入经验」的提示了。
# 提炼/保存失败静默(不阻塞、不影响结果),只在成功时推送 🧠 卡片。
if tool_seq:
if tool_seq and not designer:
try:
recipe = _clean_recipe(_distill_experience(cfg, prompt, " -> ".join(tool_seq)))
if recipe:
@@ -1504,7 +1745,13 @@ def _agent_thread(run_id, prompt, serial, cfg):
"image": None}))
except Exception as e:
_log.warning(f"经验保存异常: {e}")
q.put(("done", {"answer": answer, "usage": usage}))
with _lock:
done_evt = {"answer": answer, "usage": usage, "mode": mode}
if designer:
done_evt["draft"] = _run.get("draft")
done_evt["warnings"] = _run.get("warnings") or []
done_evt["draft_error"] = _run.get("draft_error") or ""
q.put(("done", done_evt))
except Exception as e:
_log.warning(f"Agent 运行异常: {e}")
# 诊断:打印消息结构(tool_calls 与 tool 消息配对检查)