Files
butubb ed9e8bacb1 feat(AI 建任务): 直接创建 + 草稿沉淀 + MCP de_snapshot;修「建任务页收不到 done」
用户报的"探索完无法点击创建任务"真因:一个 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 条用例、扇出单元用例、本地工具契约用例全绿。
自测产生的任务/草稿已全部清理(未碰用户既有数据)。
2026-09-14 08:14:11 +08:00

565 lines
32 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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)