用户报的"探索完无法点击创建任务"真因:一个 run 的事件原先只有**一条** queue.Queue, 聊天页与建任务页同时开着时两个 EventSource 会**瓜分**它——建任务页的回放卡在中间、 `done` 被聊天页取走 → 永远等不到草稿,页面上自然没有可点的"创建"。 一、修(根因 + 表现) - `web/agent_api.py` 新增 `_Fanout`:**每个订阅者一个专属队列**,多开页面各看各的, 还带单轮事件缓冲(晚订阅/刷新重连也能补齐回放,终止事件一定送达)。 实测两路订阅者收到完全一致的 1039 条事件(含 done)。 - `static/admin/agent.js`:断线重连的兜底订阅也按 `mode` 让开(此前漏了这一处)。 二、补齐上一批的三项 - **「直接创建」**:`POST /api/agent/task_draft/create`(草稿体只在服务端、创建前再校验一次、 成功后清草稿避免重复建)+ 草稿预览里的「✓ 直接创建任务」按钮 + 「探索完直接创建任务」勾选框。 - **草稿沉淀经验/动作**:designer 轮次也走 `_distill_experience/_distill_actions`, 但**只在草稿通过校验时**(没走通的试错不入库,免得把误点当经验)。 - **MCP `de_snapshot`**(第 20 个工具):截图+元素树一次取齐(省一次来回、不会因界面在动而错位), 附带 `screen_state`/`unstable`;两套提示词都改为优先用它。 平台侧 `/api/uiauto/snapshot` 随之多返回 `screen_state`。 三、文档 - AI_TASK_GEN §10:§10.3 记两个 bug 的真因与修法、§10.4 三项标完成、§10.5 剩余项。 - AI_CONSOLE(扇出语义、多页面同时看一轮)、API(task_draft/create、snapshot 字段)、 MCP/MCP_DESIGN/staffdeck/README/ARCHITECTURE:工具数 19→20 + de_snapshot 条目。 - backlog:记一条新发现的缺陷——`mcp_server/platform_client._login()` 会把"登录页 200" 当成登录成功(现场进程缺 `MCP_PLATFORM_PASS` 时表现为含糊的 platform_unavailable)。 自测:真机浏览器端到端(勾上"探索完直接创建")→ 探索 12 步 → 草稿 → 自动建任务成功; `de_snapshot` 直连真机校验;校验器 21 条用例、扇出单元用例、本地工具契约用例全绿。 自测产生的任务/草稿已全部清理(未碰用户既有数据)。
565 lines
32 KiB
Python
565 lines
32 KiB
Python
"""Agent 编排层:OpenAI 兼容模型(DeepSeek 等)经 MCP 工具控制手机。
|
||
|
||
支持两种运行模式:
|
||
- run_stream():流式(SSE 逐 token + 工具调用实时回调)——Web AI 控制台用
|
||
- run():非流式收集结果——CLI 用(内部调 run_stream)
|
||
|
||
流式细节(OpenAI 兼容):
|
||
- content/reasoning_content 增量逐 chunk 回调(kind 区分)
|
||
- tool_calls 分片累积(arguments 按 index 拼接),流结束后统一执行
|
||
- 截图(de_screenshot)图像转 image_url 追加下一轮,同时 on_tool 回调带缩略
|
||
- token 用量:请求带 stream_options.include_usage,按「每次模型调用」累计,
|
||
on_usage 回调吐出累计值(Web 控制台展示)
|
||
"""
|
||
import base64
|
||
import json
|
||
import logging
|
||
|
||
import httpx
|
||
from fastmcp import Client
|
||
|
||
from mcp_agent.config import AgentSettings
|
||
|
||
_log = logging.getLogger("agent")
|
||
|
||
S = AgentSettings()
|
||
|
||
|
||
# 模型单次回复的长度上限。建任务模式要一次性输出整份任务 JSON(几十个步骤),
|
||
# 4096 很容易被截断 → 工具参数变成半截 JSON(见 _execute_tool 的容错)。
|
||
_DEFAULT_MAX_TOKENS = 4096
|
||
_DESIGNER_MAX_TOKENS = 8192
|
||
|
||
|
||
class _UsageUnsupported(RuntimeError):
|
||
"""模型/网关不认 stream_options.include_usage(400/422 或报错点名该字段)——降级重试用。"""
|
||
|
||
CHAT_SYSTEM_PROMPT = """你是手机自动化控制助手。你通过工具实时操作 Android 手机。
|
||
|
||
工作规范:
|
||
1. 先 de_list_devices 确定目标设备(在线才可操作);设备有 name(名称)与 serial(地址),
|
||
**给用户汇报时用名称**(同型号多台靠它区分),调工具时仍用 serial
|
||
2. 观察屏幕:优先 `de_snapshot`(截图+元素树一次取齐,图像会随后给你)——比
|
||
de_screenshot + de_ui_tree 两次取数更省步骤,界面在动时也不会出现"元素位置和
|
||
截图对不上";需要单独看画面时才用 de_screenshot
|
||
3. 点击定位分优先级(不要自己推算像素坐标,那是精度最差的方式):
|
||
a) 目标有可见文字(按钮/菜单/列表标题/标签/输入框提示)→ de_tap_text 直接给文字,
|
||
一次完成「找到并点击」,原生控件与 WebView/图片渲染文字都支持
|
||
b) 文字有歧义或 de_tap_text 未命中 → de_ui_tree(limit=80) 看可点元素后
|
||
用 de_tap_element(text/text_contains 匹配)
|
||
c) 只有纯图形目标(视频画面/无文字图标且树里没有)才用 de_tap 给坐标——
|
||
坐标只需大致对准目标中心,服务端会自动吸附到该处可点击元素中心,无需精算
|
||
4. de_tap 点击后若返回 snapped=true 表示已吸附命中元素(可核对 label);
|
||
截图判断界面变化=点击成功,无变化=未命中
|
||
5. 输入文字:先 de_tap_text 或 de_tap 点中输入框,再 de_type_text 输入
|
||
6. 每次关键操作后再次 de_screenshot 验证结果,直到完成用户目标
|
||
7. 若点击后截图无任何变化:不要重复点同一坐标,换 de_tap_text/de_tap_element
|
||
重新定位,或先 de_ui_tree 确认元素文案再试
|
||
8. 完成或失败时用中文总结:做了什么、当前状态、需要用户注意的事项
|
||
9. 设备不可用/操作失败时如实报告错误,不要臆测成功
|
||
10. 效率:界面未变化时不要重复截图/点击同一位置;每步都要推进目标;
|
||
若连续 6 步无进展(截图内容未变/操作无效),停止并总结原因,不要空转
|
||
|
||
可用工具清单将由系统提供。"""
|
||
|
||
# 兼容旧名(CLI / 其它调用方仍可能 import SYSTEM_PROMPT)
|
||
SYSTEM_PROMPT = CHAT_SYSTEM_PROMPT
|
||
|
||
|
||
# ================== 建任务模式(designer)==================
|
||
# 与聊天模式的根本区别:**产出物是一条任务,不是"替用户做完这件事"**。
|
||
# 所以它要"看懂了就写下来",而不是"一路点到底";有副作用的动作只许核对、不许真做。
|
||
DESIGNER_SYSTEM_PROMPT = """你是手机自动化**任务设计师**。你的产出物是**一条可调度、可在步骤编辑器里继续修改的任务**(不是替用户把这件事做完)。
|
||
|
||
工作方式:在真机上探索够用即止 → 把探索到的元素与操作写成任务步骤 → 调 submit_task 提交草稿 → 人工确认后才会真正入库执行。
|
||
|
||
## 探索规范
|
||
1. 只用本次指定的设备 serial;给用户汇报时用设备名。设备离线/异常/被任务占用时**立即停止并说明**,不要硬试。
|
||
2. **先看再动**:优先 `de_snapshot(limit=80~150)`——它**一次取齐**截图与元素树(元素字段只有 text/id/desc/class/clickable/bounds,**没有 xpath**),比"先 de_screenshot 再 de_ui_tree"少一次来回、且不会因为界面在动而错位;返回 `unstable=true` 时说明取的时候界面在变,等界面静下来再取一次。同一屏只取一次,不要反复截图。
|
||
3. 不知道包名先 `de_list_apps(keyword=…)`,再 `de_open_app(package)`;用 `de_foreground_app` 确认前台。
|
||
4. **定位优先级**:`text` / `description` / `resourceId` > `xpath` > 坐标。**禁止坐标**(绝不产 `click_xy`)。
|
||
包含匹配只能 `descriptionContains`,或 `xpath` 里的 `//*[contains(@text,"…")]`(没有 textContains 这个类型)。
|
||
5. XPath 写法:`//*[@resource-id="包名:id/xxx"]`、`//*[@text="…"]`、`//*[@text="…" and @resource-id="…"]`。
|
||
**禁止** `//*[@id="x"][3]` 这种位置谓词(那是"父节点内第 3 个",同类元素一多就全失配)。
|
||
6. 验证分两档:
|
||
- **无害导航类**(tab、返回、搜索框、列表项、设置项)→ 可以 `de_tap_element` 真点一次确认能命中,记 `evidence.verified="tapped"`;
|
||
- **有副作用类**(点赞/关注/评论/发送/转发/购买/删除/退出登录)→ **只核对元素存在**(`verified="present"`),**绝不真点**;写进任务时 `probability ≤ 30`,并在 `notes` 里写明"请人工复核"。
|
||
7. 没验证命中的元素**一律不写进任务**;拿不准的进 `notes`。宁可少写一步,也不要编一个选择器。
|
||
|
||
## 步骤规范
|
||
8. 只能用这 18 种步骤类型:open_app / stop_app / screen_on / screen_off / keep_screen / key_event / swipe / swipe_until / click / long_click / wait_el / input_text / clipboard / wait / loop / group / if_el(click_xy 禁用)。
|
||
9. 输入文字必须"先 click 输入框,再 input_text";滚动用 `swipe{direction}`,不用像素。
|
||
10. 时长语义必须映射对:
|
||
- "跑 2 小时" → `params.max_duration = 7200` + 一个顶层 `loop{loop_mode:"forever"}`;
|
||
- "每天 8-9 点" → `schedule{mode:"cron_stop", cron:"0 8 * * *", stop_cron:"0 9 * * *"}`;
|
||
- **长任务必须给 max_duration**,否则等于无限跑。
|
||
11. 随机化三件套:随机等待 `wait{min,max}`、随机文案 `input_text{mode:"random",texts:"a\\nb"}`、随机触发 `group{probability:<100}` 包住子步骤(概率对任何类型都生效)。
|
||
12. 首步建议 `screen_on`(必要时 `keep_screen{mode:"on"}`);**末步用 `key_event{key:"home"}` 把设备还给用户**(不要默认息屏)。
|
||
13. 规模控制:工具调用 ≤ 25 步、steps ≤ 30 个节点、notes ≤ 5 条、evidence ≤ 10 条。结构要精简(去掉冗余的等待/滑动)。
|
||
|
||
## 收尾
|
||
14. 探索够了就调 `submit_task` 提交草稿。**校验失败时按返回的 errors 逐条修正后重提(最多 3 次)**,不要重复提交同一份草稿。
|
||
15. 提交成功后**不要再调用任何工具**,用一句中文总结:"建了什么任务、哪些步骤需要人工复核"。
|
||
16. 若 system 里注入了"可复用动作",优先复用其中的定位,跳过重复探索。
|
||
|
||
可用工具清单将由系统提供。"""
|
||
|
||
|
||
# 平台级(本地)工具:不经 MCP server,由 Agent 直接分派到调用方注册的 handler。
|
||
# 为什么不放进 mcp_server:这些工具要读写平台自身的任务/草稿(依赖 Flask app context 与
|
||
# 平台权限模型),放进 MCP 层等于给外部客户端开一个写任务的后门,且要连带改
|
||
# MCP.md/MCP_DESIGN.md 与审计——收益为零(见 doc/AI_TASK_GEN.md §9.2)。
|
||
LOCAL_TOOL_SPECS = {
|
||
"submit_task": {
|
||
"name": "submit_task",
|
||
"description": "提交任务草稿(**结束性调用**:提交成功后就不要再调任何工具,直接总结)。"
|
||
"服务端会校验草稿(步骤类型/必填参数/嵌套深度/调度格式),"
|
||
"校验失败返回 errors 数组,请逐条修正后重新调用本工具。"
|
||
"常见被打回的原因:步骤缺 selector_value;loop 的 loop_mode 不是 "
|
||
"rounds/time/forever(别写 count);按时间循环缺 loop_duration;"
|
||
"cron 不是 5 段;steps 为空。提交前请自己先按这些自查一遍。",
|
||
"parameters": {
|
||
"type": "object",
|
||
"properties": {
|
||
"summary": {"type": "string",
|
||
"description": "一句话说明这条任务做什么(≤200 字)"},
|
||
"task": {
|
||
"type": "object",
|
||
"description": "任务信封,字段同平台任务:name/target/schedule/retry/enabled/params",
|
||
"properties": {
|
||
"name": {"type": "string", "description": "任务名(≤40 字)"},
|
||
"target": {"type": "object",
|
||
"description": '{"mode":"all"|"group"|"serial", "serial":"…", "group_name":"…"}'},
|
||
"schedule": {"type": "object",
|
||
"description": '{"mode":"once"} 或 {"mode":"cron","cron":"分 时 日 月 周"} 或 {"mode":"cron_stop","cron":"…","stop_cron":"…"}'},
|
||
"retry": {"type": "object",
|
||
"description": '{"max_attempts":1-10,"delay":10-3600}'},
|
||
"enabled": {"type": "boolean"},
|
||
"params": {
|
||
"type": "object",
|
||
"properties": {
|
||
"max_duration": {"type": "integer",
|
||
"description": "单次运行时长上限(秒),0=不限时"},
|
||
"steps": {"type": "array",
|
||
"description": "步骤树:[{type,label,params}],容器用 params.children / params.then / params.else",
|
||
"items": {"type": "object"}},
|
||
},
|
||
"required": ["steps"],
|
||
},
|
||
},
|
||
"required": ["params"],
|
||
},
|
||
"notes": {"type": "array", "items": {"type": "string"},
|
||
"description": "需要人工复核/拿不准的点(≤5 条)"},
|
||
"evidence": {"type": "array", "items": {"type": "object"},
|
||
"description": "每个选择器的探索依据:{screen,element,selector_type,selector_value,verified}"},
|
||
},
|
||
"required": ["summary", "task"],
|
||
},
|
||
},
|
||
}
|
||
|
||
|
||
class Agent:
|
||
def __init__(self, settings: AgentSettings = None):
|
||
self.s = settings or S
|
||
self.tools_schema = [] # OpenAI function schema
|
||
self.messages = []
|
||
# 回调(Web 展示用,均可选):
|
||
# on_delta(text, kind) kind: content | reasoning —— 流式文本增量
|
||
# on_tool(step) step: {tool, args, result, image} —— 工具调用完成
|
||
# on_usage(usage) usage: {prompt_tokens, completion_tokens,
|
||
# total_tokens, calls} —— 累计 token 用量
|
||
self.on_delta = None
|
||
self.on_tool = None
|
||
self.on_usage = None
|
||
# 本轮累计用量(run_stream 开始时重置)
|
||
self.usage = {"prompt_tokens": 0, "completion_tokens": 0,
|
||
"total_tokens": 0, "calls": 0}
|
||
self._include_usage = True # 模型不认 stream_options 时自动置 False
|
||
self._mcp = None
|
||
# ---- 建任务模式(designer)相关 ----
|
||
self.mode = "chat"
|
||
# 平台级本地工具:{名字: async handler(args) -> dict},不走 MCP
|
||
self._local_handlers = {}
|
||
# 本轮每个工具调用了几次(防"原地打转":同一个工具反复调不推进目标)
|
||
self._tool_counts = {}
|
||
# 单个工具本轮最多调用次数(防"原地打转")。要**大于**正常重试次数:
|
||
# 探索里 submit_task 反复被校验驳回调几次是正常的,卡太死会把整轮探索截断
|
||
# (实测 25 时正好把一轮设计任务耗在修字段上)。总步数由 max_steps 兜底。
|
||
self.tool_call_limit = 40
|
||
self.max_tokens = None # 模型单次回复上限(None=用 _DEFAULT_MAX_TOKENS)
|
||
|
||
# ---------- 本地(平台级)工具 ----------
|
||
def register_local_tool(self, name, handler):
|
||
"""注册平台级本地工具。**必须在 `_load_tools()` 之前调用**(schema 在加载时组装)。
|
||
|
||
handler: `async def(args: dict) -> dict`,返回值会作为工具结果回给模型并推给前端。
|
||
"""
|
||
if name not in LOCAL_TOOL_SPECS:
|
||
raise KeyError(f"未定义的本地工具: {name}(需先在 LOCAL_TOOL_SPECS 里声明 schema)")
|
||
self._local_handlers[name] = handler
|
||
|
||
def _count_tool_call(self, name):
|
||
n = self._tool_counts.get(name, 0) + 1
|
||
self._tool_counts[name] = n
|
||
return n
|
||
|
||
# ---------- MCP 工具桥 ----------
|
||
async def _load_tools(self):
|
||
"""从 MCP Server 拉工具,转 OpenAI function schema。"""
|
||
# 短连接超时:MCP 不可用时快速失败(默认会无限重试卡死线程)
|
||
self._mcp = Client(self.s.mcp_url, timeout=10.0, init_timeout=10.0)
|
||
await self._mcp.__aenter__()
|
||
tools = await self._mcp.list_tools()
|
||
self.tools_schema = []
|
||
for t in tools:
|
||
# MCP SDK v2 改名 input_schema,兼容新旧字段
|
||
schema = getattr(t, "input_schema", None) or getattr(t, "inputSchema", {})
|
||
name = getattr(t, "name", "")
|
||
desc = getattr(t, "description", "") or ""
|
||
self.tools_schema.append({
|
||
"type": "function",
|
||
"function": {"name": name, "description": desc,
|
||
"parameters": schema}})
|
||
# 平台级本地工具(如 submit_task)随 MCP 工具一起暴露给模型,
|
||
# 执行时由 _execute_tool 本地分派(不发给 MCP server)
|
||
for name in self._local_handlers:
|
||
spec = LOCAL_TOOL_SPECS.get(name)
|
||
if spec:
|
||
self.tools_schema.append({"type": "function", "function": spec})
|
||
_log.info("MCP 工具已加载: %s", [s["function"]["name"] for s in self.tools_schema])
|
||
|
||
async def close(self):
|
||
if self._mcp:
|
||
try:
|
||
await self._mcp.__aexit__(None, None, None)
|
||
except Exception:
|
||
pass
|
||
|
||
# ---------- 模型调用(流式) ----------
|
||
async def _chat_stream(self):
|
||
"""流式 chat/completions:逐 chunk 产出 JSON(async generator)。
|
||
|
||
默认要求服务端在末尾 chunk 带 usage(token 统计);个别网关不认
|
||
`stream_options` 会直接 400,此时自动降级重试一次(不影响主流程)。
|
||
"""
|
||
if self._include_usage:
|
||
try:
|
||
async for chunk in self._chat_stream_once(True):
|
||
yield chunk
|
||
return
|
||
except _UsageUnsupported as e:
|
||
_log.warning("模型不支持 stream_options.include_usage,降级重试:%s", e)
|
||
self._include_usage = False
|
||
async for chunk in self._chat_stream_once(False):
|
||
yield chunk
|
||
|
||
async def _chat_stream_once(self, with_usage):
|
||
body = {
|
||
"model": self.s.model,
|
||
"messages": self.messages,
|
||
"tools": self.tools_schema if self.tools_schema else None,
|
||
"max_tokens": self.max_tokens or _DEFAULT_MAX_TOKENS,
|
||
"stream": True,
|
||
}
|
||
if with_usage:
|
||
body["stream_options"] = {"include_usage": True}
|
||
headers = {"Authorization": f"Bearer {self.s.api_key}",
|
||
"Content-Type": "application/json"}
|
||
url = f"{self.s.api_base.rstrip('/')}/chat/completions"
|
||
async with httpx.AsyncClient(timeout=self.s.request_timeout) as client:
|
||
async with client.stream("POST", url, json=body, headers=headers) as r:
|
||
if r.status_code != 200:
|
||
text = (await r.aread()).decode(errors="replace")
|
||
# 本次开了 include_usage 却被打回(400/422,或报错里点名这个字段)
|
||
# → 视作网关不支持,交给上层降级重试(401/余额等真错误照常抛出)
|
||
if with_usage and (r.status_code in (400, 422)
|
||
or "stream_options" in text):
|
||
raise _UsageUnsupported(text[:200])
|
||
raise RuntimeError(f"模型 API HTTP {r.status_code}: {text[:300]}")
|
||
async for line in r.aiter_lines():
|
||
if not line.startswith("data:"):
|
||
continue
|
||
data = line[5:].strip()
|
||
if data == "[DONE]":
|
||
break
|
||
try:
|
||
yield json.loads(data)
|
||
except json.JSONDecodeError:
|
||
continue
|
||
|
||
# ---------- token 用量 ----------
|
||
@staticmethod
|
||
def _read_usage(raw):
|
||
"""把一次模型调用返回的 usage 规整为 {prompt, completion, total};无效返回 None。"""
|
||
if not isinstance(raw, dict):
|
||
return None
|
||
try:
|
||
pt = int(raw.get("prompt_tokens") or 0)
|
||
ct = int(raw.get("completion_tokens") or 0)
|
||
tt = int(raw.get("total_tokens") or 0) or (pt + ct)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
if not (pt or ct or tt):
|
||
return None
|
||
return {"prompt_tokens": pt, "completion_tokens": ct, "total_tokens": tt}
|
||
|
||
def _accumulate_usage(self, raw):
|
||
"""把一次模型调用的 usage 累加进本轮总量,并回调 on_usage(累计值)。"""
|
||
u = self._read_usage(raw)
|
||
if not u:
|
||
return
|
||
self.usage["prompt_tokens"] += u["prompt_tokens"]
|
||
self.usage["completion_tokens"] += u["completion_tokens"]
|
||
self.usage["total_tokens"] += u["total_tokens"]
|
||
self.usage["calls"] += 1
|
||
if self.on_usage:
|
||
try:
|
||
self.on_usage(dict(self.usage))
|
||
except Exception:
|
||
pass
|
||
|
||
# ---------- 工具执行 ----------
|
||
async def _execute_tool(self, name, arguments):
|
||
"""执行工具(MCP 或平台级本地工具),返回 (文本结果, image_data_or_None)。"""
|
||
try:
|
||
args = json.loads(arguments) if isinstance(arguments, str) else (arguments or {})
|
||
except (json.JSONDecodeError, TypeError) as e:
|
||
# 长参数被输出长度截断时 arguments 是半截 JSON。**不能整轮报错**——
|
||
# 把"坏了"告诉模型,让它精简后重发(designer 的 steps 可能很长)。
|
||
_log.warning("工具 %s 参数不是合法 JSON: %s", name, e)
|
||
return {"ok": False,
|
||
"error": f"参数不是合法 JSON({e})——常见原因是这次输出过长被截断,"
|
||
"请精简要提交的内容后重新调用一次"}, None
|
||
if not isinstance(args, dict):
|
||
return {"ok": False, "error": "工具参数必须是 JSON 对象"}, None
|
||
_log.info("执行工具 %s %s", name, args)
|
||
|
||
# 同一个工具被反复调用(原地打转)时给模型一个明确的刹车
|
||
if self._count_tool_call(name) > self.tool_call_limit:
|
||
return {"ok": False,
|
||
"error": f"{name} 本轮调用次数已达上限({self.tool_call_limit} 次),"
|
||
"请换一种做法推进,或直接总结当前进展"}, None
|
||
|
||
# 平台级本地工具:不发给 MCP server
|
||
handler = self._local_handlers.get(name)
|
||
if handler is not None:
|
||
# 平台级工具自己会推更贴切的提示卡(📝 任务草稿 / 草稿被拦下 / 未存下),
|
||
# 这里**不再重复推一张工具卡**——只在 handler 意外抛异常(自己没来得及推)时补一张
|
||
try:
|
||
result = await handler(args)
|
||
except Exception as e:
|
||
_log.exception("本地工具 %s 执行失败", name)
|
||
result = {"ok": False, "error": f"工具执行失败: {e}"}
|
||
if self.on_tool:
|
||
try:
|
||
self.on_tool({"tool": name, "args": args,
|
||
"result": result, "image": None})
|
||
except Exception:
|
||
pass
|
||
return result, None
|
||
|
||
try:
|
||
result = await self._mcp.call_tool(name, args)
|
||
data = getattr(result, "data", result)
|
||
except Exception as e:
|
||
return {"ok": False, "error": f"工具执行失败: {e}"}, None
|
||
# de_screenshot:图像分离(作为 image_url 追加给模型看 + on_tool 缩略展示)
|
||
image_b64 = None
|
||
text_result = data
|
||
if name == "de_screenshot" and isinstance(data, dict) and data.get("ok"):
|
||
img = (data.get("data") or {}).get("image") or {}
|
||
if img.get("data"):
|
||
text_result = {k: v for k, v in (data.get("data") or {}).items()
|
||
if k != "image"}
|
||
image_b64 = img["data"]
|
||
if self.on_tool:
|
||
try:
|
||
self.on_tool({"tool": name, "args": args,
|
||
"result": text_result, "image": image_b64})
|
||
except Exception:
|
||
pass
|
||
return text_result, image_b64
|
||
|
||
def _repair_tool_messages(self):
|
||
"""修复 tool_calls 配对不完整:从尾部移除「assistant 带 tool_calls 但其后
|
||
tool 回应不足」的消息段(流中断可能丢失分片,400 重试前自愈)。"""
|
||
for i in range(len(self.messages) - 1, -1, -1):
|
||
m = self.messages[i]
|
||
if m.get("role") == "assistant" and m.get("tool_calls"):
|
||
# 只数**紧跟着的连续** tool 消息:模型侧要求工具回应连续排列,
|
||
# 中间夹一条 user(如截图图像)就会被判成"回应不足"。
|
||
need = len(m["tool_calls"])
|
||
have = 0
|
||
for x in self.messages[i + 1:]:
|
||
if x.get("role") != "tool":
|
||
break
|
||
have += 1
|
||
if have < need:
|
||
_log.warning("修复不完整 tool_calls 段(need=%d have=%d),回退 %d 条消息",
|
||
need, have, len(self.messages) - i)
|
||
self.messages = self.messages[:i]
|
||
return
|
||
|
||
# ---------- 主循环(流式) ----------
|
||
async def run_stream(self, prompt: str, serial: str = "",
|
||
history=None, on_delta=None, on_tool=None,
|
||
should_stop=None, extra_context=None, on_usage=None,
|
||
mode: str = "chat"):
|
||
"""流式执行一轮指令,返回最终完整文本。
|
||
|
||
history:上一轮的 [{"role": "user"|"assistant", "content": 文本}] 列表,
|
||
用于多轮对话保持上下文(截图/工具消息不入历史,控制 token)。
|
||
on_delta(text, kind):content/reasoning 文本增量(实时推给前端)
|
||
on_tool(step):工具调用完成(实时显示 MCP 步骤)
|
||
should_stop:可调用 fn() -> bool,每轮模型调用前检查(用户中断用)
|
||
extra_context:附加文本(经验记忆注入,放在 system prompt 末尾)
|
||
on_usage(usage):每完成一次模型调用回调一次(累计值,见 self.usage)
|
||
mode:`chat`(默认,AI 控制台聊天)或 `designer`(AI 建任务:改系统提示词、
|
||
放宽输出长度上限、注册的本地工具生效)
|
||
|
||
本轮累计 token 用量同时留在 self.usage(调用方可直接读)。
|
||
"""
|
||
self.on_delta = on_delta
|
||
self.on_tool = on_tool
|
||
self.on_usage = on_usage
|
||
self.mode = mode
|
||
self.max_tokens = _DESIGNER_MAX_TOKENS if mode == "designer" else _DEFAULT_MAX_TOKENS
|
||
self._tool_counts = {}
|
||
self.usage = {"prompt_tokens": 0, "completion_tokens": 0,
|
||
"total_tokens": 0, "calls": 0}
|
||
target = serial or self.s.default_serial
|
||
sys_txt = DESIGNER_SYSTEM_PROMPT if mode == "designer" else CHAT_SYSTEM_PROMPT
|
||
if target:
|
||
sys_txt += f"\n\n本次默认目标设备 serial:{target}(未指定设备时用它)。"
|
||
if extra_context:
|
||
sys_txt += f"\n\n## 过往成功经验参考(同类任务,可参考其中的操作套路,但要根据当前界面灵活调整)\n{extra_context}"
|
||
self.messages = [{"role": "system", "content": sys_txt}]
|
||
for h in (history or []):
|
||
if h.get("role") in ("user", "assistant") and h.get("content"):
|
||
self.messages.append({"role": h["role"], "content": h["content"]})
|
||
self.messages.append({"role": "user", "content": prompt})
|
||
|
||
for _step in range(self.s.max_steps):
|
||
if should_stop and should_stop():
|
||
_log.info("Agent 被用户中断")
|
||
return "(已按用户要求停止操作)"
|
||
content_parts = []
|
||
tool_acc = {} # index -> {id, name, args}
|
||
has_tool = False
|
||
retried = False
|
||
call_usage = None # 本次模型调用的 usage(末尾 chunk 带)
|
||
while True:
|
||
call_usage = None # 重试时丢弃上一次(未完成)的用量
|
||
try:
|
||
async for chunk in self._chat_stream():
|
||
if chunk.get("usage"):
|
||
call_usage = chunk["usage"]
|
||
choice = (chunk.get("choices") or [{}])[0]
|
||
delta = choice.get("delta") or {}
|
||
text = delta.get("content")
|
||
if text:
|
||
content_parts.append(text)
|
||
if on_delta:
|
||
on_delta(text, "content")
|
||
rtext = delta.get("reasoning_content")
|
||
if rtext:
|
||
if on_delta:
|
||
on_delta(rtext, "reasoning")
|
||
for tc in delta.get("tool_calls") or []:
|
||
has_tool = True
|
||
idx = tc.get("index", 0)
|
||
acc = tool_acc.setdefault(idx, {"id": "", "name": "", "args": ""})
|
||
if tc.get("id"):
|
||
acc["id"] = tc["id"]
|
||
fn = tc.get("function") or {}
|
||
if fn.get("name"):
|
||
acc["name"] += fn["name"]
|
||
if fn.get("arguments"):
|
||
acc["args"] += fn["arguments"]
|
||
break
|
||
except RuntimeError as e:
|
||
# 流中断导致 tool_calls 分片丢失:修复后重试一次
|
||
if ("tool_calls" in str(e) or "must be followed" in str(e)) and not retried:
|
||
_log.warning("tool_calls 消息不完整,自愈重试")
|
||
self._repair_tool_messages()
|
||
retried = True
|
||
continue
|
||
raise
|
||
|
||
self._accumulate_usage(call_usage)
|
||
full_content = "".join(content_parts)
|
||
|
||
if has_tool:
|
||
# 组装 assistant 消息(含 tool_calls)并执行工具
|
||
tcs = []
|
||
for idx in sorted(tool_acc):
|
||
acc = tool_acc[idx]
|
||
tcs.append({"id": acc["id"] or f"call_{idx}",
|
||
"type": "function",
|
||
"function": {"name": acc["name"],
|
||
"arguments": acc["args"]}})
|
||
self.messages.append({"role": "assistant",
|
||
"content": full_content,
|
||
"tool_calls": tcs})
|
||
# 工具结果必须**连续**跟在带 tool_calls 的 assistant 消息后面:
|
||
# 中间插任何消息都会被模型侧判成"工具回应不足"而 400
|
||
# (An assistant message with 'tool_calls' must be followed by tool
|
||
# messages responding to each 'tool_call_id')。
|
||
# 截图图像因此先攒着,等本轮所有 tool 消息都发完,再作为一条 user 消息附上。
|
||
pending_images = []
|
||
for tc in tcs:
|
||
fn = tc["function"]
|
||
text_result, image_b64 = await self._execute_tool(
|
||
fn["name"], fn["arguments"])
|
||
self.messages.append({
|
||
"role": "tool", "tool_call_id": tc["id"],
|
||
"content": json.dumps(text_result, ensure_ascii=False)[:4000]})
|
||
if image_b64:
|
||
pending_images.append(image_b64)
|
||
if pending_images:
|
||
content = [{"type": "text",
|
||
"text": "这是最新屏幕截图,请基于它继续判断"}]
|
||
for img in pending_images:
|
||
content.append({"type": "image_url",
|
||
"image_url": {"url":
|
||
f"data:image/jpeg;base64,{img}"}})
|
||
self.messages.append({"role": "user", "content": content})
|
||
continue
|
||
|
||
# 无工具调用:本轮即最终回答
|
||
return full_content
|
||
|
||
# 步骤超限:不带工具让模型做最终总结(避免机械提示,给用户有意义的结论)
|
||
try:
|
||
_log.warning("达到最大步骤数,请求模型收尾总结")
|
||
saved_tools = self.tools_schema
|
||
self.tools_schema = []
|
||
self.messages.append({"role": "user",
|
||
"content": "已达最大操作步骤数,请立即用中文总结:"
|
||
"已完成的部分、当前设备状态、未能完成的原因与下一步建议。"
|
||
"不要调用任何工具。"})
|
||
parts = []
|
||
call_usage = None
|
||
async for chunk in self._chat_stream():
|
||
if chunk.get("usage"):
|
||
call_usage = chunk["usage"]
|
||
delta = (chunk.get("choices") or [{}])[0].get("delta") or {}
|
||
text = delta.get("content")
|
||
if text:
|
||
parts.append(text)
|
||
if on_delta:
|
||
on_delta(text, "content")
|
||
self._accumulate_usage(call_usage)
|
||
self.tools_schema = saved_tools
|
||
summary = "".join(parts)
|
||
return summary or "(已达步骤上限,模型未能生成总结)"
|
||
except Exception as e:
|
||
return f"(已达最大步骤数,且收尾总结失败: {e})"
|
||
|
||
# ---------- 非流式(CLI) ----------
|
||
async def run(self, prompt: str, serial: str = "") -> str:
|
||
"""非流式执行,返回最终文本(CLI 用,内部走流式收集)。"""
|
||
return await self.run_stream(prompt, serial)
|