Files
auto_control/web/agent_api.py
T
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

1904 lines
86 KiB
Python
Raw 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.
"""AI 控制台 API:模型/Key 前端配置 + Agent 流式执行(SSE 事件流)。
流程:
POST /api/agent/run {prompt} → 启动 Agent 线程,返回 run_id
GET /api/agent/stream?run_id= → SSE 事件流(EventSource 订阅):
event: delta {text, kind: content|reasoning} 流式文本增量
event: step {tool, args, image?} 工具调用完成(MCP 步骤)
event: usage {prompt_tokens, completion_tokens,
total_tokens, calls} 本轮累计 token 用量
event: done {answer, usage} 完成
event: error {message} 失败
GET/POST /api/agent/config → 配置读写(key 打码回显)
单实例:同时只允许一个 Agent 运行。
"""
import asyncio
import base64
import io
import json
import os
import queue
import re
import sys
import threading
import uuid
from flask import Blueprint, Response, jsonify, request
from core.logger import get_logger
from core.models import db
from web.auth import admin_required
_log = get_logger("web.agent")
bp = Blueprint("agent", __name__)
# ---------- 配置键(app_meta) ----------
_CFG_KEYS = {"api_base": "agent_api_base",
"model": "agent_model",
"api_key": "agent_api_key",
"default_serial": "agent_default_serial",
"max_steps": "agent_max_steps"}
# 推理链落库上限(字符):只留够回看的量,避免会话消息无限膨胀
_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 -> _Fanout(SSE 消费者各自订阅,见下)
_stop_events = {} # run_id -> threading.Event(用户中断)
_lock = threading.Lock()
_flask_app = None # web_server 注册时注入(后台线程 db 操作需 app context)
def set_app(app):
global _flask_app
_flask_app = app
class _Fanout:
"""一个 Agent 轮次的事件,扇出给**所有**订阅者(每个订阅者一个独立队列)。
历史实现是"一个 run 一个 queue.Queue":第二个页面一订阅,两个 EventSource 就
开始**瓜分**同一条队列——谁先取到算谁的。2026-09-14 实测踩到:同一个浏览器里
聊天页与建任务页同时开着,建任务页的回放卡到一半就不动了、`done` 事件被另一个
页面取走 → 永远等不到草稿(表现为"探索完没法创建任务")。
扇出之后,多开几个页面都各看各的,谁都不丢事件。
"""
# 单轮事件缓冲上限:新订阅者从这里补齐"订阅之前"的回放(刷新/重连不丢进度)。
# 只留最近 _HISTORY_KEEP 条,避免长任务把内存拖大;终止事件(done/error)永远保留。
_HISTORY_KEEP = 400
def __init__(self):
self._subs = []
self._history = []
self._closed = False
self._lock = threading.Lock()
def subscribe(self):
"""新订阅者拿一个专属队列(**预填已发生的事件**);轮次已结束则返回 None。"""
with self._lock:
if self._closed:
return None
q = queue.Queue()
for item in self._history:
q.put(item)
self._subs.append(q)
return q
def unsubscribe(self, q):
with self._lock:
if q in self._subs:
self._subs.remove(q)
def put(self, item):
with self._lock:
subs = list(self._subs)
self._history.append(item)
# 只留最近 N 条。终止事件(done/error 及其后的哨兵)永远在尾部,
# 按这个规则裁剪不会把它们丢掉。
if len(self._history) > self._HISTORY_KEEP:
del self._history[:len(self._history) - self._HISTORY_KEEP]
for q in subs:
q.put(item)
def close(self):
"""结束:给每个订阅者发终止哨兵(None)。"""
with self._lock:
self._closed = True
subs, self._subs = self._subs, []
for q in subs:
q.put(None)
def _meta_get(key):
# 走方言中立助手(app_meta 的列名 key 在 MySQL 里是保留字,裸 SQL 会语法错)
from core.db_config import meta_get
return meta_get(key) or ""
def _meta_put(key, value):
from core.db_config import meta_set
meta_set(key, str(value))
def _read_cfg():
return {k: _meta_get(v) for k, v in _CFG_KEYS.items()}
# ================== 经验记忆(自进化) ==================
# agent_experience / experience_audit / agent_action / agent_conversation 四张表
# 已在 core/models.py 里定义为 ORM 模型(唯一真相),建表统一由 db.create_all() 完成。
# 以前这里是各写一段 SQLite 的 CREATE TABLE + try/except: pass —— 换 MySQL 后
# AUTOINCREMENT / TEXT DEFAULT 全都建不出来,而异常被吞掉,表现为"功能静默失灵"。""
# 巡检评审提示词:模型只判「保留/建议删除」+ 评分 + 理由,不直接删。
_AUDIT_PROMPT = """你是「手机自动化经验库」质检员。经验库存放 AI 成功操作手机后提炼的
操作套路,下次相似任务会自动注入参考。请判断下面这条经验是否值得继续保留。
判断标准(全部满足才保留):
1. 具体可执行:步骤是明确的工具操作(de_* 或清晰的中文步骤),不是「de_tap / de_input」
这类空泛占位,也没有编造不存在的工具名
2. 不依赖易过时信息:不含写死的像素坐标(x,y 数字)这类界面一变就失效的步骤
3. 是独立任务的操作套路(如「打开X搜索Y并点赞」),不是对话追问/抱怨/单次性内容
4. 配方与下方实际操作序列大致一致,不是模型自由发挥编造
经验内容:
任务描述:{prompt}
操作配方:{recipe}
实际操作序列:{tool_seq}
被引用次数:{hits}
只输出 JSON:{{"keep": true或false, "score": 0到10的整数, "reason": "一句话中文理由"}}"""
def _ensure_tables():
"""确保 agent 系四张表存在(幂等)。create_all 只建缺失的表,开销可以忽略。
失败必须落日志:以前是 except: pass,MySQL 下建表失败会被完全吞掉,
表不存在 → 后续查询全报错,而日志里什么都看不到。
"""
try:
db.create_all()
except Exception as e:
_log.error(f"agent 表创建失败: {e}")
def _ensure_exp_table():
_ensure_tables()
def _bigrams(text):
"""中文/英文文本 bigram 集合(无空格分词,粗粒度相似度)。"""
t = "".join(c for c in (text or "").lower() if c.isalnum() or "\u4e00" <= c <= "\u9fff")
return {t[i:i + 2] for i in range(len(t) - 1)}
# ---------- 经验质量门槛 ----------
# 自进化经验只应保存「独立操作任务」的成功套路。多轮对话里用户的短句
# (质疑/纠正/催促,如「继续啊」「你确定我是卡1吗」「还没有完成啊」)
# 不是新任务——把执行出错被纠正的轮次存成经验会教坏后续任务。
# 用启发式过滤(零成本,可解释),配方层再校验工具名真实性。
# 任务性动词:命中任一视为有明确操作诉求(疑问/纠错短句一般不含它们)
_TASK_VERBS = ("打开", "搜索", "查看", "找到", "截图", "输入", "点击", "点开",
"发送", "安装", "卸载", "下载", "启动", "停止", "关闭", "退出",
"登录", "切换", "设置", "删除", "清理", "复制", "粘贴", "读取",
"剪贴板", "长按", "滑动", "播放", "发布", "检查", "看看",
"帮我", "给我", "请", "拍张", "查一下")
# 对话续语开头:几乎只出现在承接上一轮(「继续啊」「还有吗」)
_CONTINUE_PREFIXES = ("继续", "还有", "然后呢", "再来", "快点", "刚才",
"接着", "上一步")
# 强质疑/纠错信号(不含「为什么」——「查一下为什么」是正当任务)
_DOUBT_MARKS = ("你确定", "是不是", "不是吧", "不是吗", "怎么都", "怎么还",
"还没有", "没看到", "我说的是", "你听我说", "不对吧", "你又",
"重新来", "错了", "你说得", "你回答")
# 全部真实 MCP 工具(配方里出现不存在的 de_* 说明蒸馏模型在编造,弃存)
_KNOWN_TOOLS = frozenset({
"de_list_devices", "de_screenshot", "de_tap", "de_swipe", "de_ui_tree",
"de_tap_element", "de_read_clipboard", "de_wake", "de_press_key",
"de_open_app", "de_stop_app", "de_foreground_app", "de_type_text",
"de_set_clipboard", "de_sleep", "de_ocr", "de_tap_text", "de_list_apps",
"de_list_tasks"})
def _qualify_experience(prompt, recipe):
"""经验入库前质量门槛,返回 True=值得保存。
1) prompt 太短 / 纯续语开头 / 质疑纠错 → 非独立任务,弃
2) 疑问短句(≤30 字、以 吗/? 结尾)且无任务动词 → 追问/反问,弃
3) 配方含不存在的 de_* 工具(蒸馏模型自由发挥)→ 弃
"""
t = "".join(c for c in (prompt or "") if not c.isspace())
if len(t) < 6:
_log.info("经验弃存:prompt 过短「%s」", t[:20])
return False
if any(t.startswith(p) for p in _CONTINUE_PREFIXES):
_log.info("经验弃存:对话续语开头「%s」", t[:20])
return False
if any(m in t for m in _DOUBT_MARKS):
_log.info("经验弃存:质疑/纠错语气「%s」", t[:20])
return False
if (t.endswith("吗") or t.endswith("?") or t.endswith("?")) \
and len(t) <= 30 and not any(v in t for v in _TASK_VERBS):
_log.info("经验弃存:无操作诉求的追问「%s」", t[:20])
return False
for name in re.findall(r"de_[a-z_]+", recipe or ""):
if name not in _KNOWN_TOOLS:
_log.info("经验弃存:配方含不存在的工具 %s", name)
return False
return True
def _find_experiences(prompt, limit=2, threshold=0.10):
"""按 bigram 重叠检索相似历史经验(prompt 与任务描述的字符相似度)。
返回 (注入文本, 命中的配方列表);命中的经验 hits+1(回写),让被反复
参考的有效经验浮到前面。配方列表供调用方展示「命中 N 条」与摘要。
"""
try:
if _flask_app is None:
return "", []
with _flask_app.app_context():
_ensure_exp_table()
rows = db.session.execute(db.text(
"SELECT id, task_prompt, recipe, hits FROM agent_experience "
"WHERE recipe != '' ORDER BY hits DESC, id DESC LIMIT 50")).fetchall()
except Exception:
return "", []
if not rows:
return "", []
cur = _bigrams(prompt)
if not cur:
return "", []
scored = []
for row_id, task_prompt, recipe, hits in rows:
sim = len(cur & _bigrams(task_prompt)) / len(cur)
if sim >= threshold:
scored.append((sim, hits or 0, recipe, row_id))
scored.sort(key=lambda x: (-x[0], -x[1]))
if scored:
# hits 回写(尽力而为,失败不影响检索)
try:
with _flask_app.app_context():
for _sim, _hits, _recipe, row_id in scored:
db.session.execute(db.text(
"UPDATE agent_experience SET hits=hits+1 WHERE id=:i"),
{"i": row_id})
db.session.commit()
except Exception:
pass
picked = [recipe for _sim, _hits, recipe, _row_id in scored[:limit]]
parts = [f"- {r[:600]}" for r in picked]
return "\n".join(parts), picked
def _msg_text(message, allow_reasoning=False):
"""取模型回复正文。
推理模型(DeepSeek 等)**偶发**把 token 预算全花在 reasoning 上、content 为空
(finish_reason=length),直接读 content 会当空处理 → 经验/动作被静默丢弃
(2026-09-10 实测)。默认**不回退 reasoning**(那是思考草稿,做"配方"会污染);
仅在能结构化解析的场景(动作 JSON 提取)才允许回退。
"""
m = message or {}
text = (m.get("content") or "").strip()
if text:
return text
if allow_reasoning:
return (m.get("reasoning_content") or "").strip()
return ""
_RECIPE_JSON_RE = re.compile(r'"recipe"\s*:\s*"((?:[^"\\]|\\.)*)"', re.S)
def _recipe_ok(recipe):
"""配方质量门槛:太短、或只有序号+省略号(模型照抄占位)视为无效。
2026-09-10 实测:提示词里写了示例 JSON,模型会把占位符 `1. …\\n2. …`
原样当成配方存下来 → 这里拦掉,宁可重试/不存,也不污染经验库。
"""
t = (recipe or "").strip()
if len(t) < 15:
return False
if "…" in t and len(t) < 30:
return False
return True
def _extract_recipe(text):
"""从(可能被截断的)模型输出里取配方正文。
蒸馏提示词要求输出 {"recipe": "…"};推理模型会把预算烧在 reasoning 上并
截断(finish_reason=length),严格 json.loads 会失败,故用正则直接抠字段。
"""
if not text:
return ""
m = _RECIPE_JSON_RE.search(text)
if m:
try:
return json.loads('"' + m.group(1) + '"').strip()
except Exception:
return m.group(1).replace("\\n", "\n").strip()
t = text.strip()
if t.startswith("{") or t.startswith("["):
return "" # JSON 外壳但没抠到 recipe:视为无效,避免把草稿当配方
return t
def _clean_recipe(recipe):
"""清洗蒸馏输出:去掉模型可能加的「操作配方:」之类前缀标签与包裹引号。
旧逻辑用 `"配方" not in recipe[:50]` 直接拒绝——而蒸馏提示词本就要求以
「操作配方」作答,模型一旦照做就会被整条丢弃(静默),导致经验很难存下。
改为清洗前缀而非丢弃。
"""
t = (recipe or "").strip()
if not t:
return ""
t = re.sub(r'^[^\n]{0,12}?配方[::]\s*', '', t, count=1).strip()
t = re.sub(r'^(操作步骤|步骤|流程)[::]\s*', '', t).strip()
return t.strip('"“”\'').strip()
def _distill_experience(cfg, prompt, tool_seq):
"""任务完成后用模型把操作序列提炼为可复用配方(失败静默,不阻塞)。"""
try:
import httpx
body = {
"model": cfg.get("model") or "deepseek-v4-flash-vision-exp",
# 关推理:蒸馏是"给定轨迹写配方"的确定性任务,开启 thinking 会占满
# 预算导致 content 空/被截断(实测);该代理支持 thinking.type=disabled
"thinking": {"type": "disabled"},
"messages": [{"role": "user",
"content": "以下是一次成功的手机自动化操作记录。请提炼成简洁的"
"「操作配方」(2-6 步,每步:目标 → 用哪个工具),"
"供下次同类任务参考。不要解释,直接输出配方正文。\n"
f"任务:{prompt[:300]}\n操作序列:{tool_seq[:800]}"}],
"max_tokens": 800,
}
headers = {"Authorization": f"Bearer {cfg.get('api_key', '')}",
"Content-Type": "application/json"}
url = f"{(cfg.get('api_base') or 'https://api.deepseek.com').rstrip('/')}/chat/completions"
# content 为空(推理吃满预算)时重试一次——实测同样提示第二次常能出正文
for _attempt in (1, 2):
r = httpx.post(url, json=body, headers=headers, timeout=25)
if r.status_code != 200:
_log.warning(f"经验提炼: 模型返回 HTTP {r.status_code}")
continue
msg = (r.json().get("choices") or [{}])[0].get("message") or {}
# 正文优先;为空时从 reasoning 里抠(可能含清晰步骤),都要过质量门槛
for cand in (_msg_text(msg), _extract_recipe(msg.get("reasoning_content"))):
if _recipe_ok(cand):
return cand.strip()[:1500]
if cand:
_log.info(f"经验提炼: 候选配方质量不足({len(cand)} 字符)被丢弃")
_log.info("经验提炼: 两次均未产出可用配方(推理占满预算/输出被截断/仅占位)")
return ""
except Exception as e:
_log.warning(f"经验提炼失败: {e}")
return ""
def _save_experience(prompt, recipe, tool_seq):
"""保存经验(后台线程调用,包 app context)。返回是否保存成功。"""
try:
if _flask_app is None:
return False
with _flask_app.app_context():
_ensure_exp_table()
from datetime import datetime
db.session.execute(db.text(
"INSERT INTO agent_experience(task_prompt, recipe, tool_seq, hits, created_at) "
"VALUES(:p, :r, :t, 0, :c)"),
{"p": prompt[:500], "r": recipe, "t": tool_seq[:1000],
"c": datetime.now().strftime("%Y-%m-%d %H:%M")})
db.session.commit()
_log.info("经验已保存(配方 %d 字符)", len(recipe))
return True
except Exception as e:
_log.warning(f"经验保存失败: {e}")
return False
# ---------- 动作经验库(agent_action)----------
# 与「任务级配方」(agent_experience) 互补:动作 = 有语义名的可复用单元,可含 1~N 步,
# steps 直接用编辑器 schema(open_app/click/input_text…),带**元素定位**
# (selector_type/selector_value),**不含坐标**(分辨率/旋转/改版即失效)。
# 复用:执行前按 name/别名/App 召回并注入 system prompt,模型可跳过重新探索。
# 建表见 core/models.py 的 AgentAction
# 允许沉淀的步骤类型(编辑器 STEP_TYPES 子集;显式排除 click_xy 等坐标类)
_ACTION_STEP_TYPES = {
"open_app", "stop_app", "screen_on", "screen_off", "key_event",
"swipe", "swipe_until", "click", "long_click", "wait_el",
"input_text", "clipboard", "wait", "loop", "group", "if_el"}
_ACTION_REQUIRED = { # 类型 → 必需的 params 键(缺则丢弃该动作)
"open_app": ("package",), "stop_app": ("package",),
"click": ("selector_value",), "long_click": ("selector_value",),
"wait_el": ("selector_value",), "swipe_until": ("selector_value",),
"if_el": ("selector_value", "then"),
"group": ("children",), "loop": ("children",)}
# 可写入动作的关键工具(探索类 de_screenshot/de_ui_tree/de_ocr 不沉淀)
_ACTION_TOOLS = {
"de_open_app", "de_stop_app", "de_tap_text", "de_tap_element", "de_type_text",
"de_set_clipboard", "de_swipe", "de_press_key", "de_wake", "de_sleep"}
def _ensure_action_table():
_ensure_tables()
def _brief_result(result):
"""工具结果的精简摘要(判成败 + 供提炼模型参考)。"""
try:
s = json.dumps(result, ensure_ascii=False) if not isinstance(result, str) else result
except Exception:
s = str(result)
return (s or "")[:160]
def _tool_ok(tool, result):
"""启发式判断一次工具调用是否成功(只沉淀成功动作,避免把误点当经验)。"""
if result is None:
return False
if isinstance(result, dict):
if result.get("error") or result.get("ok") is False or result.get("success") is False:
return False
if tool == "de_tap_text":
return bool(result.get("matched")) or result.get("method") in ("ui", "ocr")
if result.get("matched") is False:
return False
s = str(result).lower()
return not any(k in s for k in ("not_found", "error", "failed", "occupied"))
def _normalize_action_item(a):
"""把模型可能的「单动作」形状({action/tool/type, params})归一成标准动作。
实测蒸馏模型常不按提示词的 name/steps 输出,而是照抄输入行给出
{"action":"de_open_app","params":{...}}——这里做兼容映射,避免全被丢弃。
"""
if not isinstance(a, dict):
return None
if isinstance(a.get("steps"), list): # 已是标准形状
return a
tool = str(a.get("action") or a.get("tool") or a.get("type") or "").strip()
p = a.get("params") if isinstance(a.get("params"), dict) else {}
if not tool:
return None
name = str(a.get("name") or "").strip()
t, sp = None, {}
if tool in ("de_open_app", "open_app"):
t, sp = "open_app", {"package": p.get("package", "")}
name = name or f"打开应用 {p.get('package', '')}"
elif tool in ("de_stop_app", "stop_app"):
t, sp = "stop_app", {"package": p.get("package", "")}
name = name or f"关闭应用 {p.get('package', '')}"
elif tool == "de_tap_text":
t, sp = "click", {"selector_type": "text", "selector_value": p.get("text", "")}
name = name or f"点击「{p.get('text', '')}」"
elif tool == "de_tap_element":
by = {"text": "text", "id": "resourceId", "desc": "description",
"text_contains": "text", "desc_contains": "descriptionContains"}.get(
p.get("by"), "text")
t, sp = "click", {"selector_type": by, "selector_value": p.get("value", "")}
name = name or f"点击元素 {p.get('value', '')}"
elif tool == "de_type_text":
t, sp = "input_text", {"mode": "fixed", "fixed_text": p.get("text", "")}
name = name or f"输入「{p.get('text', '')}」"
elif tool == "de_set_clipboard":
t, sp, name = "clipboard", {"text": p.get("text", "")}, name or "写入剪贴板"
elif tool == "de_swipe":
t, sp = "swipe", {"direction": p.get("direction") or "up"}
name = name or "滑动"
elif tool == "de_press_key":
t, sp = "key_event", {"key": p.get("key") or "back"}
name = name or f"按键 {p.get('key', '')}"
elif tool == "de_wake":
t, name = "screen_on", name or "亮屏解锁"
elif tool == "de_sleep":
t, name = "screen_off", name or "息屏"
if not t:
return None
return {"name": name, "app": a.get("app", ""), "aliases": a.get("aliases") or [],
"params": [], "preconditions": a.get("preconditions", ""),
"steps": [{"type": t, "params": sp}]}
def _sanitize_actions(actions):
"""校验/清洗模型产出的动作:名字非空、步骤白名单+必填、显式禁坐标。"""
out = []
for a in (actions or []):
a = _normalize_action_item(a)
if not isinstance(a, dict):
continue
name = str(a.get("name") or "").strip()[:40]
steps = a.get("steps")
if not name or not isinstance(steps, list) or not steps:
continue
good = []
for st in steps:
if not isinstance(st, dict):
continue
t = st.get("type")
if t not in _ACTION_STEP_TYPES: # click_xy 等坐标类在此被剔除
continue
p = st.get("params") if isinstance(st.get("params"), dict) else {}
if any(not p.get(k) for k in _ACTION_REQUIRED.get(t, ())):
continue
good.append({"type": t, "label": str(st.get("label") or "")[:20], "params": p})
if not good:
continue
out.append({
"name": name,
"app": str(a.get("app") or "")[:80],
"aliases": [str(x)[:20] for x in (a.get("aliases") or []) if str(x).strip()][:5],
"params": [str(x)[:20] for x in (a.get("params") or []) if str(x).strip()][:5],
"preconditions": str(a.get("preconditions") or "")[:100],
"steps": good})
return out
def _iter_json_objects(raw):
"""从文本里按大括号配对切出顶层 JSON 对象(容忍坏片段,逐条抢救)。"""
depth = 0
start = None
in_str = esc = False
for i, c in enumerate(raw):
if in_str:
if esc:
esc = False
elif c == "\\":
esc = True
elif c == '"':
in_str = False
continue
if c == '"':
in_str = True
elif c == "{":
if depth == 0:
start = i
depth += 1
elif c == "}":
depth -= 1
if depth == 0 and start is not None:
yield raw[start:i + 1]
start = None
def _loads_lenient(text):
"""尽量解析模型输出的 JSON 数组:容忍 markdown 围栏、尾逗号、中文引号、坏对象。"""
s = (text or "").strip()
s = re.sub(r"^```[a-zA-Z]*\s*|\s*```$", "", s).strip()
m = re.search(r"\[[\s\S]*\]", s)
raw = m.group(0) if m else s
candidates = [raw,
re.sub(r",\s*([\]}])", r"\1", raw), # 去尾逗号
re.sub(r"[“”]", '"', raw), # 中文引号 → 英文
re.sub(r"[“”]", '"', re.sub(r",\s*([\]}])", r"\1", raw))]
for cand in candidates:
try:
v = json.loads(cand)
if isinstance(v, list):
return v
except Exception:
continue
# 逐对象抢救(顶层大括号配对),坏的跳过
out = []
for obj in _iter_json_objects(raw):
try:
out.append(json.loads(obj))
except Exception:
continue
return out
def _distill_actions(cfg, prompt, trace):
"""从**成功**的工具轨迹提炼命名动作(JSON 数组)。失败返回 []。"""
ok_ops = [t for t in (trace or [])
if t.get("tool") in _ACTION_TOOLS and _tool_ok(t.get("tool"), t.get("result"))]
_log.info(f"动作提炼: 轨迹 {len(trace or [])} 步, 成功可沉淀 {len(ok_ops)} 步")
if not ok_ops:
return []
# 输入截断到前 10 步:步数过多会让模型输出变长被截断(实测 10 步 → 2002 字符无有效 JSON)
ok_ops = ok_ops[:10]
lines = [f"- {t['tool']} 参数={t['args']} 结果={_brief_result(t['result'])}" for t in ok_ops]
instruction = (
"以下是一次成功的手机自动化操作的**成功步骤**。请提炼为**至多 3 个**「动作」"
"(每个动作 = 一个有语义名的可复用单元,**至多 4 步**,字段尽量简短)。"
"只输出 JSON 数组,不要解释,不要思考过程。\n"
"字段:name(动作名,如「打开抖音」「搜索关键词」);app(包名,未知则空串);"
"aliases(别名数组);params(参数名数组,如[\"关键词\"]);steps(步骤数组)。\n"
"steps 每步:type + params,type 取值:open_app{package} / click{selector_type,"
"selector_value,wait_timeout?} / input_text{mode,fixed_text} / swipe{direction} /"
"wait{min,max} / key_event{key} / group{children}。\n"
"**定位必须用元素定位**:selector_type 取 xpath/text/resourceId/description/"
"descriptionContains,值用上面步骤里出现的真实文字或 id;**禁止坐标**。"
"若某步只能用坐标定位,就不要产出该动作。\n"
"输出示例(**顶层字段必须是 name/steps,禁止用 action/tool 当顶层字段**):\n"
'[{"name":"打开抖音","app":"com.ss.android.ugc.aweme","aliases":["启动抖音"],'
'"params":[],"steps":[{"type":"open_app","params":{"package":"com.ss.android.ugc.aweme"}}]},'
'{"name":"搜索关键词","app":"","aliases":["点搜索"],"params":["关键词"],'
'"steps":[{"type":"click","params":{"selector_type":"text","selector_value":"搜索"}}]}]\n'
f"任务:{prompt[:200]}\n成功步骤:\n" + "\n".join(lines)[:1500])
try:
import httpx
body = {"model": cfg.get("model") or "deepseek-v4-flash-vision-exp",
"thinking": {"type": "disabled"}, # 同 _distill_experience:关推理
"messages": [{"role": "user", "content": instruction}],
"max_tokens": 900}
headers = {"Authorization": f"Bearer {cfg.get('api_key', '')}",
"Content-Type": "application/json"}
url = f"{(cfg.get('api_base') or 'https://api.deepseek.com').rstrip('/')}/chat/completions"
content = ""
# content 空(推理吃满预算)→ 重试一次;仍空则回退 reasoning_content
# (这里能结构化提取 JSON 数组,草稿里也常含可用 JSON,故允许回退)
for _attempt in (1, 2):
r = httpx.post(url, json=body, headers=headers, timeout=30)
if r.status_code != 200:
_log.warning(f"动作提炼: 模型返回 HTTP {r.status_code}")
continue
msg = (r.json().get("choices") or [{}])[0].get("message") or {}
content = _msg_text(msg, allow_reasoning=True)
if content:
break
acts = _sanitize_actions(_loads_lenient(content))
if not acts:
_log.info(f"动作提炼: 解析后无有效动作(原始输出 {len(content)} 字符)"
f" 样本={content[:160]!r}")
return acts
except Exception as e:
_log.warning(f"动作提炼失败: {e}")
return []
def _save_actions(prompt, actions):
"""按 (name, app) upsert 保存动作。返回保存条数。"""
if not actions or _flask_app is None:
return 0
n = 0
try:
from datetime import datetime
now = datetime.now().strftime("%Y-%m-%d %H:%M")
with _flask_app.app_context():
_ensure_action_table()
for a in actions:
row = db.session.execute(db.text(
"SELECT id FROM agent_action WHERE name=:n AND app=:a"),
{"n": a["name"], "a": a["app"]}).fetchone()
if row:
db.session.execute(db.text(
"UPDATE agent_action SET aliases=:al, params=:p, steps=:s, "
"preconditions=:pc, updated_at=:t WHERE id=:i"),
{"al": json.dumps(a["aliases"], ensure_ascii=False),
"p": json.dumps(a["params"], ensure_ascii=False),
"s": json.dumps(a["steps"], ensure_ascii=False),
"pc": a["preconditions"], "t": now, "i": row[0]})
else:
db.session.execute(db.text(
"INSERT INTO agent_action(name, app, aliases, params, steps, "
"preconditions, hits, source_prompt, created_at, updated_at) "
"VALUES(:n,:a,:al,:p,:s,:pc,0,:sp,:t,:t)"),
{"n": a["name"], "a": a["app"],
"al": json.dumps(a["aliases"], ensure_ascii=False),
"p": json.dumps(a["params"], ensure_ascii=False),
"s": json.dumps(a["steps"], ensure_ascii=False),
"pc": a["preconditions"], "sp": prompt[:200], "t": now})
n += 1
db.session.commit()
_log.info(f"动作经验已保存 {n} 条")
except Exception as e:
_log.warning(f"动作保存失败: {e}")
return n
def _find_actions(prompt, limit=3):
"""召回可复用动作:动作名/别名命中 prompt,或与来源提示够相似。返回 (文本, 名字列表)。"""
try:
if _flask_app is None:
return "", []
with _flask_app.app_context():
_ensure_action_table()
rows = db.session.execute(db.text(
"SELECT id, name, app, aliases, params, steps, hits FROM agent_action "
"ORDER BY hits DESC, id DESC LIMIT 100")).fetchall()
except Exception:
return "", []
if not rows:
return "", []
p_norm = "".join(c for c in (prompt or "").lower()
if c.isalnum() or "一" <= c <= "鿿")
cur = _bigrams(prompt)
scored = []
for rid, name, app, aliases, params, steps, hits in rows:
try:
keys = [name] + [str(x) for x in (json.loads(aliases) if aliases else [])]
except Exception:
keys = [name]
hit = any(k and k.lower() in p_norm for k in keys if k)
sim = (len(cur & _bigrams(name)) / len(cur)) if cur else 0.0
if hit or sim >= 0.34:
scored.append((1 if hit else 0, sim, hits or 0, rid, name, app, params, steps))
if not scored:
return "", []
scored.sort(key=lambda x: (-x[0], -x[1], -x[2]))
picked = scored[:limit]
try:
with _flask_app.app_context():
for item in picked:
db.session.execute(db.text("UPDATE agent_action SET hits=hits+1 WHERE id=:i"),
{"i": item[3]})
db.session.commit()
except Exception:
pass
items, parts = [], []
for item in picked:
_hit, _sim, _hts, _rid, name, app, params, steps = item
try:
st = json.loads(steps) if steps else []
except Exception:
st = []
items.append(name)
brief = " → ".join(
str(s.get("params", {}).get("selector_value")
or s.get("params", {}).get("package")
or s.get("params", {}).get("fixed_text")
or s.get("type")) for s in st[:8])
parts.append(f"- 「{name}」" + (f"(app={app})" if app else "")
+ (f" 参数:{params}" if params else "") + f":{brief}")
return "\n".join(parts), items
# ================== 经验巡检(AI 质检,删除需人工确认) ==================
_audit_state = {"running": False, "last": "", "last_summary": ""}
def _ensure_audit_table():
_ensure_tables()
def _review_one_exp(cfg, prompt, recipe, tool_seq, hits):
"""调模型评审单条经验,返回 (verdict, score, reason)。
模型偶尔输出解释文字不带 JSON——重试一次(强调只输出 JSON),
仍失败则保守返回保留。
"""
import httpx
base_prompt = _AUDIT_PROMPT.format(prompt=(prompt or "")[:400],
recipe=(recipe or "")[:600],
tool_seq=(tool_seq or "")[:600],
hits=hits or 0)
for attempt in (1, 2):
try:
content = base_prompt if attempt == 1 else (
base_prompt + "\n注意:只输出一个 JSON 对象,不要输出任何其它文字、解释或代码块标记。")
body = {
"model": cfg.get("model") or "deepseek-v4-flash-vision-exp",
"messages": [{"role": "user", "content": content}],
"max_tokens": 200,
}
headers = {"Authorization": f"Bearer {cfg.get('api_key', '')}",
"Content-Type": "application/json"}
r = httpx.post(f"{(cfg.get('api_base') or 'https://api.deepseek.com').rstrip('/')}/chat/completions",
json=body, headers=headers, timeout=20)
if r.status_code != 200:
raise RuntimeError(f"HTTP {r.status_code}")
text = _msg_text(((r.json() or {}).get("choices") or [{}])[0]
.get("message") or {})
m = re.search(r"\{.*\}", text, re.S)
if not m:
raise RuntimeError("响应无 JSON")
j = json.loads(m.group(0))
keep = bool(j.get("keep", True))
score = max(0, min(10, int(j.get("score", 5) or 5)))
reason = str(j.get("reason", ""))[:200]
return ("keep" if keep else "delete"), score, reason
except Exception as e:
_log.warning(f"经验评审第 {attempt} 次失败: {e}")
return "keep", 0, "评审失败自动保留(两次尝试均无有效 JSON)"
def run_experience_audit():
"""巡检一轮:逐条评审经验 → 写 experience_audit。
规则:
- 最新已 kept / deleted 的经验不再重复评审(人工拍板过的尊重)
- 评审 keep → 直接 action=kept;delete → action=pending(等人工确认)
- 绝不自动删除;无 api_key 配置时跳过并记日志
由 web_server 每日定时调度或 API 手动触发(后台线程)。
"""
if _audit_state["running"]:
_log.info("经验巡检已在运行,跳过本次触发")
return
_audit_state["running"] = True
try:
if _flask_app is None:
return
with _flask_app.app_context():
_ensure_exp_table()
_ensure_audit_table()
cfg = _read_cfg()
if not cfg.get("api_key"):
_log.info("经验巡检跳过:未配置 API Key")
return
rows = db.session.execute(db.text(
"SELECT id, task_prompt, recipe, tool_seq, hits FROM agent_experience "
"ORDER BY id")).fetchall()
if not rows:
_log.info("经验巡检:经验库为空")
return
from datetime import datetime
now = datetime.now().strftime("%Y-%m-%d %H:%M")
reviewed = suggested = 0
for exp_id, prompt, recipe, tool_seq, hits in rows:
# 最新一次人工/评审结论:kept / deleted 的不再打扰
last = db.session.execute(db.text(
"SELECT action FROM experience_audit WHERE exp_id=:e "
"ORDER BY id DESC LIMIT 1"), {"e": exp_id}).scalar()
if last in ("kept", "deleted"):
continue
verdict, score, reason = _review_one_exp(
cfg, prompt, recipe, tool_seq, hits)
action = "kept" if verdict == "keep" else "pending"
db.session.execute(db.text(
"INSERT INTO experience_audit"
"(exp_id, verdict, score, reason, hits, action, audited_at) "
"VALUES(:e, :v, :s, :r, :h, :a, :t)"),
{"e": exp_id, "v": verdict, "s": score, "r": reason,
"h": hits or 0, "a": action, "t": now})
reviewed += 1
suggested += 1 if action == "pending" else 0
db.session.commit()
summary = f"评审 {reviewed} 条,建议删除 {suggested} 条(待人工确认)"
_audit_state.update(last=now, last_summary=summary)
_log.info("经验巡检完成: %s", summary)
except Exception as e:
_log.warning(f"经验巡检异常: {e}")
finally:
_audit_state["running"] = False
# ================== 经验库管理 API ==================
@bp.route("/api/agent/experience")
@admin_required
def agent_experience_list():
"""经验库列表(含最近一次巡检结论),供 AI 控制台「🧠 经验库」面板管理。"""
_ensure_exp_table()
_ensure_audit_table()
rows = db.session.execute(db.text(
"SELECT e.id, e.task_prompt, e.recipe, e.tool_seq, e.hits, e.created_at, "
"a.verdict, a.score, a.reason, a.action, a.audited_at "
"FROM agent_experience e LEFT JOIN experience_audit a ON a.id = "
"(SELECT id FROM experience_audit WHERE exp_id = e.id ORDER BY id DESC LIMIT 1) "
"ORDER BY e.id DESC")).fetchall()
exps = []
for r in rows:
(eid, task_prompt, recipe, tool_seq, hits, created_at,
verdict, score, reason, action, audited_at) = r
audit = None
if action:
audit = {"verdict": verdict, "score": round(score or 0, 1),
"reason": reason, "action": action, "at": audited_at}
exps.append({"id": eid, "task_prompt": task_prompt or "",
"recipe": recipe or "", "tool_seq": tool_seq or "",
"hits": hits or 0, "created_at": created_at or "",
"audit": audit})
return jsonify({"ok": True, "running": _audit_state["running"],
"last": _audit_state["last"],
"last_summary": _audit_state["last_summary"],
"experiences": exps})
@bp.route("/api/agent/experience/delete", methods=["POST"])
@admin_required
def agent_experience_delete():
"""人工确认删除经验(真删;前端有二次确认,巡检绝不自动调它)。"""
data = request.json or {}
exp_id = int(data.get("id") or 0)
if exp_id <= 0:
return jsonify({"ok": False, "error": "缺少 id"}), 400
db.session.execute(db.text(
"DELETE FROM agent_experience WHERE id=:i"), {"i": exp_id})
db.session.execute(db.text(
"UPDATE experience_audit SET action='deleted' "
"WHERE exp_id=:i AND action='pending'"), {"i": exp_id})
db.session.commit()
_log.info("人工确认删除经验 #%d", exp_id)
return jsonify({"ok": True, "msg": f"经验 #{exp_id} 已删除"})
@bp.route("/api/agent/experience/audit", methods=["POST"])
@admin_required
def agent_experience_audit_now():
"""手动触发一轮经验巡检(后台线程;进行中返回 409)。"""
if _audit_state["running"]:
return jsonify({"ok": False, "error": "巡检正在进行中"}), 409
threading.Thread(target=run_experience_audit, daemon=True).start()
return jsonify({"ok": True, "msg": "巡检已启动,完成后刷新列表查看建议"})
@bp.route("/api/agent/experience/keep", methods=["POST"])
@admin_required
def agent_experience_keep():
"""人工保留:撤销建议删除(该条后续巡检不再重复建议)。"""
data = request.json or {}
exp_id = int(data.get("id") or 0)
if exp_id <= 0:
return jsonify({"ok": False, "error": "缺少 id"}), 400
db.session.execute(db.text(
"UPDATE experience_audit SET action='kept' "
"WHERE exp_id=:i AND action='pending'"), {"i": exp_id})
db.session.commit()
return jsonify({"ok": True, "msg": f"经验 #{exp_id} 已保留"})
# ================== 动作经验库(agent_action:查看/删除/编辑) ==================
def _action_row_to_dict(r):
def _j(s, d):
try:
return json.loads(s) if s else d
except Exception:
return d
return {"id": r[0], "name": r[1], "app": r[2], "aliases": _j(r[3], []),
"params": _j(r[4], []), "steps": _j(r[5], []), "preconditions": r[6],
"hits": r[7] or 0, "updated_at": r[8] or ""}
@bp.route("/api/agent/actions")
@admin_required
def agent_actions_list():
"""动作经验库列表(命名动作 + 编辑器 schema 步骤 + 元素定位,禁坐标)。"""
with _flask_app.app_context():
_ensure_action_table()
rows = db.session.execute(db.text(
"SELECT id, name, app, aliases, params, steps, preconditions, hits, updated_at "
"FROM agent_action ORDER BY hits DESC, id DESC")).fetchall()
return jsonify({"ok": True, "actions": [_action_row_to_dict(r) for r in rows]})
@bp.route("/api/agent/actions/delete", methods=["POST"])
@admin_required
def agent_actions_delete():
aid = int((request.json or {}).get("id") or 0)
if aid <= 0:
return jsonify({"ok": False, "error": "缺少 id"}), 400
with _flask_app.app_context():
_ensure_action_table()
db.session.execute(db.text("DELETE FROM agent_action WHERE id=:i"), {"i": aid})
db.session.commit()
return jsonify({"ok": True, "msg": f"动作 #{aid} 已删除"})
@bp.route("/api/agent/actions/save", methods=["POST"])
@admin_required
def agent_actions_save():
"""新增/编辑动作:{id?, name, app?, aliases?, params?, steps, preconditions?}。
steps 可为数组或 JSON 字符串;经 _sanitize_actions 校验(白名单/必填/禁坐标)。
"""
from datetime import datetime
d = request.json or {}
steps = d.get("steps")
if isinstance(steps, str):
try:
steps = json.loads(steps)
except Exception:
return jsonify({"ok": False, "error": "steps 不是合法 JSON"}), 400
item = _sanitize_actions([{"name": d.get("name"), "app": d.get("app", ""),
"aliases": d.get("aliases") or [],
"params": d.get("params") or [],
"preconditions": d.get("preconditions", ""),
"steps": steps if isinstance(steps, list) else []}])
if not item:
return jsonify({"ok": False, "error": "动作无效:需要 name,且 steps 需为受支持的步骤"
"(click/open_app/input_text… 且带元素定位;不接受坐标)"}), 400
a = item[0]
now = datetime.now().strftime("%Y-%m-%d %H:%M")
with _flask_app.app_context():
_ensure_action_table()
aid = int(d.get("id") or 0)
if aid > 0:
db.session.execute(db.text(
"UPDATE agent_action SET name=:n, app=:a, aliases=:al, params=:p, steps=:s, "
"preconditions=:pc, updated_at=:t WHERE id=:i"),
{"n": a["name"], "a": a["app"],
"al": json.dumps(a["aliases"], ensure_ascii=False),
"p": json.dumps(a["params"], ensure_ascii=False),
"s": json.dumps(a["steps"], ensure_ascii=False),
"pc": a["preconditions"], "t": now, "i": aid})
else:
db.session.execute(db.text(
"INSERT INTO agent_action(name, app, aliases, params, steps, preconditions, "
"hits, source_prompt, created_at, updated_at) "
"VALUES(:n,:a,:al,:p,:s,:pc,0,'(人工新增)',:t,:t)"),
{"n": a["name"], "a": a["app"],
"al": json.dumps(a["aliases"], ensure_ascii=False),
"p": json.dumps(a["params"], ensure_ascii=False),
"s": json.dumps(a["steps"], ensure_ascii=False),
"pc": a["preconditions"], "t": now})
db.session.commit()
return jsonify({"ok": True, "msg": "动作已保存"})
# ================== 历史会话(DeepSeek 式:多会话持久化) ==================
# 会话 = 消息序列(JSON),一轮 run 的 user/assistant 文本完成后追加落库。
# 单表 JSON 存储(每会话几十 KB 内,无需独立消息表)。建表见 core/models.py。
def _ensure_conv_table():
_ensure_tables()
def _conv_msgs(conv_id):
"""读会话消息列表 [{role, content}]。会话不存在返回 None。"""
try:
row = db.session.execute(db.text(
"SELECT messages FROM agent_conversation WHERE id=:i"),
{"i": conv_id}).fetchone()
except Exception:
return None
if not row:
return None
try:
return json.loads(row[0] or "[]")
except Exception:
return []
def _conv_save_messages(conv_id, messages, title=None):
"""整段写回会话消息 + 时间戳(尽力而为,失败不影响主流程)。"""
try:
from datetime import datetime
now = datetime.now().strftime("%Y-%m-%d %H:%M")
msgs_json = json.dumps(messages[-60:], ensure_ascii=False)
if title:
db.session.execute(db.text(
"UPDATE agent_conversation SET messages=:m, title=:t, "
"updated_at=:u WHERE id=:i"),
{"m": msgs_json, "t": title[:100], "u": now, "i": conv_id})
else:
db.session.execute(db.text(
"UPDATE agent_conversation SET messages=:m, updated_at=:u "
"WHERE id=:i"),
{"m": msgs_json, "u": now, "i": conv_id})
db.session.commit()
return True
except Exception as e:
_log.warning(f"会话保存失败: {e}")
return False
@bp.route("/api/agent/conversations", methods=["GET"])
@admin_required
def agent_conversations_list():
"""会话列表(按最近更新倒序):{id, title, updated_at, count}。"""
_ensure_conv_table()
rows = db.session.execute(db.text(
"SELECT id, title, messages, updated_at FROM agent_conversation "
"ORDER BY updated_at DESC, created_at DESC")).fetchall()
out = []
for cid, title, messages, updated_at in rows:
try:
count = len(json.loads(messages or "[]")) // 2
except Exception:
count = 0
out.append({"id": cid, "title": title or "新会话",
"updated_at": updated_at or "", "count": count})
return jsonify({"ok": True, "conversations": out})
@bp.route("/api/agent/conversations", methods=["POST"])
@admin_required
def agent_conversations_create():
"""新建会话(空消息)。返回 {id}。"""
_ensure_conv_table()
from datetime import datetime
cid = uuid.uuid4().hex[:10]
now = datetime.now().strftime("%Y-%m-%d %H:%M")
db.session.execute(db.text(
"INSERT INTO agent_conversation(id, title, messages, created_at, updated_at) "
"VALUES(:i, '', '[]', :c, :c)"), {"i": cid, "c": now})
db.session.commit()
return jsonify({"ok": True, "id": cid, "title": "新会话"})
@bp.route("/api/agent/conversations/<conv_id>", methods=["GET"])
@admin_required
def agent_conversations_detail(conv_id):
"""会话详情(全部消息文本)。"""
_ensure_conv_table()
msgs = _conv_msgs(conv_id)
if msgs is None:
return jsonify({"ok": False, "error": "会话不存在"}), 404
row = db.session.execute(db.text(
"SELECT title, created_at, updated_at FROM agent_conversation "
"WHERE id=:i"), {"i": conv_id}).fetchone()
return jsonify({"ok": True, "id": conv_id,
"title": (row[0] if row else "") or "新会话",
"messages": msgs,
"created_at": row[1] if row else "",
"updated_at": row[2] if row else ""})
@bp.route("/api/agent/conversations/<conv_id>", methods=["DELETE"])
@admin_required
def agent_conversations_delete(conv_id):
"""删除会话(消息一并删除,不可恢复)。"""
_ensure_conv_table()
db.session.execute(db.text(
"DELETE FROM agent_conversation WHERE id=:i"), {"i": conv_id})
db.session.commit()
_log.info(f"删除会话 {conv_id}")
return jsonify({"ok": True, "msg": "会话已删除"})
@bp.route("/api/agent/conversations/<conv_id>/rename", methods=["POST"])
@admin_required
def agent_conversations_rename(conv_id):
"""重命名会话。"""
data = request.json or {}
title = (data.get("title") or "").strip()[:100]
if not title:
return jsonify({"ok": False, "error": "标题不能为空"}), 400
_ensure_conv_table()
db.session.execute(db.text(
"UPDATE agent_conversation SET title=:t WHERE id=:i"),
{"t": title, "i": conv_id})
db.session.commit()
return jsonify({"ok": True, "msg": "已重命名"})
# ================== 运行 ==================
@bp.route("/api/agent/run", methods=["POST"])
@admin_required
def agent_run():
"""启动 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
if not cfg.get("model"):
return jsonify({"ok": False, "error": "请先填写模型名"}), 400
if not serial and not (cfg.get("default_serial") or "").strip():
return jsonify({"ok": False, "error": "请先选择目标设备(AI 只操作你指定的设备)"}), 400
# 目标设备校验:必须在池、在线、且无任务运行(AI 不与任务抢设备;
# MCP 层另有 busy 锁兜底外部客户端)
if serial:
try:
from web import context
devs, err = context.mgr.get_status()
dev = next((x for x in (devs or []) if x.get("serial") == serial), None)
if err or dev is None:
return jsonify({"ok": False, "error": f"设备 {serial} 不在设备池"}), 400
if not dev.get("present"):
return jsonify({"ok": False, "error": f"设备 {serial} 当前离线,请稍后再试"}), 400
if dev.get("worker_status") in ("running", "connecting"):
return jsonify({"ok": False, "error":
f"设备 {serial} 正在执行任务「{dev.get('task_job') or ''}」——"
f"AI 不与任务抢设备,任务结束后才能操作(或在任务页先停止)"}), 409
except Exception:
pass # 状态服务异常不阻塞(MCP busy 锁兜底)
# 会话校验:不存在则自动新建(标题=首条消息截断)
if conv_id:
_ensure_conv_table()
if _conv_msgs(conv_id) is None:
try:
from datetime import datetime as _dt2
now = _dt2.now().strftime("%Y-%m-%d %H:%M")
# 上面刚确认过会话不存在,普通 INSERT 即可;
# 并发下撞主键就回滚忽略(原来用 SQLite 专有的 INSERT OR IGNORE)
db.session.execute(db.text(
"INSERT INTO agent_conversation(id, messages, "
"created_at, updated_at) VALUES(:i, '[]', :c, :c)"),
{"i": conv_id, "c": now})
db.session.commit()
except Exception:
db.session.rollback()
with _lock:
if _run["state"] == "running":
return jsonify({"ok": False, "error":
f"已有 Agent 运行中({_run.get('prompt', '')[:40]}…),"
f"请等待完成或先停止"}), 409
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, mode=mode,
started=_dt.now().strftime("%H:%M:%S"),
answer="", error="", usage={},
draft=None, draft_error="", warnings=[])
# history 保留(多轮上下文),由会话/「新建会话」管理
_queues[run_id] = _Fanout()
_stop_events[run_id] = threading.Event()
_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, mode, settings),
daemon=True).start()
return jsonify({"ok": True, "run_id": run_id})
@bp.route("/api/agent/config", methods=["GET"])
@admin_required
def agent_config_get():
"""读 Agent 配置(key 打码返回)。"""
cfg = _read_cfg()
if cfg["api_key"]:
k = cfg["api_key"]
cfg["api_key_masked"] = k[:6] + "***" + k[-4:]
return jsonify({"ok": True, **cfg})
@bp.route("/api/agent/config", methods=["POST"])
@admin_required
def agent_config_save():
"""保存 Agent 配置:{api_base?, model?, api_key?, default_serial?} 部分更新。"""
data = request.json or {}
for key, meta_key in _CFG_KEYS.items():
if key in data and data[key] is not None:
_meta_put(meta_key, str(data[key]).strip())
db.session.commit()
return jsonify({"ok": True, "msg": "已保存"})
@bp.route("/api/agent/run", methods=["GET"])
@admin_required
def agent_run_status():
"""当前 Agent 运行状态(多窗口/页面刷新恢复用):
idle/running/done + run_id + 最终答案 + 会话历史。
前端刷新后据此:渲染历史轮次、running 时重新订阅事件流(服务端队列
保留积压事件,重连后补发)、done 时展示最终答案。
"""
with _lock:
hist = [{"role": h.get("role"), "content": (h.get("content") or "")[:4000]}
for h in (_run.get("history") or [])]
return jsonify({"ok": True,
"state": _run.get("state", "idle"),
"run_id": _run.get("id") or "",
"prompt": _run.get("prompt", ""),
"serial": _run.get("serial", ""),
"started": _run.get("started", ""),
"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:]})
@bp.route("/api/agent/devices")
@admin_required
def agent_devices():
"""AI 可用设备列表:在线状态 + 是否有任务运行(前端选择器 busy 设备禁选)。
worker 状态实时(内存);busy = 任务 running/connecting。
"""
try:
from web import context
devs, err = context.mgr.get_status()
except Exception as e:
return jsonify({"ok": False, "error": str(e)[:120]}), 503
out = []
for d in (devs or []):
busy = d.get("worker_status") in ("running", "connecting")
out.append({"serial": d.get("serial"),
"name": d.get("device_name") or "", # 设备名称(设备池里的身份标识)
"model": d.get("model") or "",
"online": bool(d.get("present")),
"busy": busy,
"worker_status": d.get("worker_status") or "idle",
"task_job": d.get("task_job") or ""})
return jsonify({"ok": True, "devices": out})
@bp.route("/api/agent/stream")
@admin_required
def agent_stream():
"""SSE 事件流(EventSource):delta/step/done/error。"""
run_id = request.args.get("run_id", "")
with _lock:
if run_id != _run["id"]:
return jsonify({"ok": False, "error": "run_id 不存在"}), 404
fan = _queues.get(run_id)
# 每个订阅者一个专属队列(多开页面互不抢事件);轮次已结束 → 404
q = fan.subscribe() if fan is not None else None
if q is None:
return jsonify({"ok": False, "error": "本轮已结束"}), 404
def gen():
try:
while True:
try:
evt = q.get(timeout=15)
except queue.Empty:
yield ": keepalive\n\n" # 心跳防超时
continue
if evt is None:
break
kind, payload = evt
yield f"event: {kind}\ndata: {json.dumps(payload, ensure_ascii=False)}\n\n"
if kind in ("done", "error"):
break
finally:
# 页面关掉/断线时摘掉自己的队列,别让扇出往里堆没人读的事件
fan.unsubscribe(q)
return Response(gen(), mimetype="text/event-stream",
headers={"Cache-Control": "no-cache",
"X-Accel-Buffering": "no"})
@bp.route("/api/agent/stop", methods=["POST"])
@admin_required
def agent_stop():
"""中断当前运行的 Agent(下一个检查点生效,通常在数秒内)。"""
with _lock:
if _run["state"] != "running":
return jsonify({"ok": False, "error": "当前没有运行中的任务"}), 400
evt = _stop_events.get(_run["id"])
if evt:
evt.set()
_log.info("用户请求中断 Agent")
return jsonify({"ok": True, "msg": "已请求停止"})
@bp.route("/api/agent/clear", methods=["POST"])
@admin_required
def agent_clear():
"""清空对话历史。"""
with _lock:
_run["history"] = []
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/create", methods=["POST"])
@admin_required
def agent_task_draft_create():
"""把暂存的草稿**直接创建**成任务(AI 建任务页的「直接创建任务」)。
这是"AI 直接建任务"的落点,两条设计约束:
- **草稿体只在服务端**(app_meta):客户端只能传 overrides,不能把任务内容
再提交一遍——少一个可被篡改的入口;
- 创建前**再校验一次**(草稿可能是人工改过的),通过后走与「新建任务」
**完全相同**的落库路径 `context.mgr.add_job`。
创建成功后清掉草稿(同一份草稿不该被重复建成多个任务)。
"""
from core.task_draft import validate_draft
from web import context
saved = _load_draft_meta()
if not saved or not saved.get("draft"):
return jsonify({"ok": False, "error": "没有可创建的草稿,请先探索一次"}), 400
data = request.json or {}
overrides = data.get("overrides") if isinstance(data.get("overrides"), dict) else {}
groups, pool = _draft_env()
res = validate_draft(saved["draft"], groups=groups, pool=pool,
default_serial=(data.get("serial") or "").strip(),
overrides=overrides)
if not res["ok"]:
return jsonify({"ok": False, "error": "草稿校验未通过(请在编辑器里核对)",
"errors": res["errors"]}), 400
task = res["draft"]["task"]
params = task["params"]
params.setdefault("skip_offline", True) # 与编辑器保存时一致
try:
job = context.mgr.add_job(name=task["name"], task_type=task["task_type"],
target=task["target"], params=params,
schedule=task["schedule"], retry=task["retry"],
enabled=bool(task.get("enabled", True)))
except Exception as e:
_log.exception("AI 草稿创建任务失败")
return jsonify({"ok": False, "error": f"创建失败: {e}"}), 500
# 用掉就清掉,避免重复创建
_meta_put(_DRAFT_META_KEY, "")
try:
db.session.commit()
except Exception:
db.session.rollback()
with _lock:
_run["draft"] = None
_run["warnings"] = []
d = job.to_dict()
nr = context.mgr.next_run_of(job)
d["next_run"] = nr.strftime("%Y-%m-%d %H:%M") if nr else None
_log.info(f"AI 草稿已创建任务: {job.name}({job.id})"
f"{',' + str(len(res['warnings'])) + ' 条提醒' if res['warnings'] else ''}")
return jsonify({"ok": True, "msg": f"任务「{job.name}」已创建",
"job": d, "warnings": res["warnings"]})
@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:
from PIL import Image
img = Image.open(io.BytesIO(base64.b64decode(b64)))
if img.width > width:
img = img.resize((width, int(img.height * width / img.width)))
buf = io.BytesIO()
img.convert("RGB").save(buf, "JPEG", quality=quality)
return base64.b64encode(buf.getvalue()).decode()
except Exception:
return b64
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__))))
from mcp_agent.agent import Agent
# 推理链(reasoning_content)累计——单纯流式展示会随页面刷新丢失,
# 落库后可在会话里折叠回看(只存文本,不回灌给模型)
reason_parts = []
reason_len = 0
def on_delta(text, kind):
nonlocal reason_len
if kind == "reasoning" and reason_len < _REASONING_KEEP:
reason_parts.append(text)
reason_len += len(text)
q.put(("delta", {"text": text, "kind": kind}))
def on_usage(usage):
"""每次模型调用完成 → 推累计用量(前端实时刷新 token 计数)。"""
q.put(("usage", dict(usage)))
tool_seq = [] # 本轮工具序列(任务级配方提炼用)
tool_trace = [] # 结构化轨迹:{tool, args, result}(动作经验提炼用,判成败)
def on_tool(step):
# args 保留对象(json.dumps 序列化)——前端要解析 serial 做画面跟随;
# 之前 str() 成 Python repr(单引号)导致前端 JSON.parse 失败、跟随失效
rec = {"tool": step.get("tool"),
"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:
args = step.get("args") or {}
brief = {k: v for k, v in args.items() if k != "serial"}
tool_seq.append(f"{step.get('tool')}({str(brief)[:60]})")
res = step.get("result")
# 可沉淀工具保留原始 result(_tool_ok 需按字段判成败);其余只存摘要
tool_trace.append({"tool": step.get("tool"), "args": brief,
"result": res if step.get("tool") in _ACTION_TOOLS
else _brief_result(res)})
except Exception:
pass
agent = Agent()
agent.s.api_base = cfg.get("api_base") or agent.s.api_base
agent.s.model = cfg.get("model") or agent.s.model
agent.s.api_key = cfg.get("api_key") or agent.s.api_key
agent.s.default_serial = cfg.get("default_serial") or agent.s.default_serial
try:
agent.s.max_steps = max(1, min(200, int(cfg.get("max_steps") or 40)))
except (TypeError, ValueError):
pass # 非法值用默认 40
with _lock:
conv_id = _run.get("conv_id") or ""
# 历史来源:绑定了会话 → 从会话读(多轮上下文延续,DeepSeek 式);
# 无会话(CLI/兼容)→ 用运行态 history
if conv_id:
# 后台线程需自行推 app context(与 _find_experiences 同思路)
with _flask_app.app_context():
msgs = _conv_msgs(conv_id)
history = [{"role": m.get("role"), "content": m.get("content", "")[:4000]}
for m in (msgs or []) if m.get("role") in ("user", "assistant")]
history = history[-24:]
with _lock:
_run["history"] = list(history)
else:
with _lock:
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)
if exp_ctx:
_log.info("命中历史经验,注入参考配方")
brief = ";".join("".join(x.split())[:16] for x in exp_items[:2])
q.put(("step", {"tool": "🧠 经验记忆",
"args": f"命中 {len(exp_items)} 条同类历史经验,已注入参考"
+ (f":{brief}" if brief else ""),
"image": None}))
# 动作经验:可复用的命名动作(带元素定位),优先复用可跳过重新探索
act_ctx, act_items = _find_actions(prompt)
if act_ctx:
_log.info(f"命中可复用动作 {len(act_items)} 个,注入参考")
q.put(("step", {"tool": "🧠 动作经验",
"args": f"命中 {len(act_items)} 个可复用动作,已注入参考"
+ (f":{'、'.join(act_items[:3])}" if act_items else ""),
"image": None}))
recall_ctx = exp_ctx
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", "")
def _mcp_unreachable(err):
"""MCP 客户端在工具服务不可达/返回非 MCP 响应时只抛含糊字样,
识别出来以便给用户明确指引(默认端口 8033)。"""
m = str(err).lower()
return any(k in m for k in (
"server returned an error response", "connecterror",
"connection refused", "all connection attempts failed",
"connection error", "server disconnected",
"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:
if _mcp_unreachable(e):
raise RuntimeError(
f"MCP server(8033) 不可达:无法加载设备工具({mcp_url})。"
f"请确认 MCP server 已启动(本机可运行 "
f"python -m mcp_server.mcp_server,默认端口 8033)") from e
raise
try:
return await agent.run_stream(prompt, target,
history=history,
on_delta=on_delta, on_tool=on_tool,
should_stop=lambda: bool(
stop_evt and stop_evt.is_set()),
extra_context=recall_ctx,
on_usage=on_usage,
mode=mode)
except Exception as e:
if _mcp_unreachable(e):
raise RuntimeError(
f"MCP server(8033) 连接中断:AI 无法继续操作设备({mcp_url})") from e
raise
# 整体超时保护:卡死时结束,释放单实例
answer = asyncio.run(asyncio.wait_for(_execute(), timeout=900))
usage = dict(getattr(agent, "usage", None) or {})
with _lock:
_run["state"] = "done"
_run["answer"] = answer
_run["usage"] = usage
# 追加本轮进历史(多轮连续性;上限 12 轮防 token 膨胀)。
# usage/reasoning 仅用于前端展示与落库,不进模型上下文(读回时只取 role/content)
hist = _run.setdefault("history", [])
hist.append({"role": "user", "content": prompt[:2000]})
turn = {"role": "assistant", "content": (answer or "")[:4000]}
if usage:
turn["usage"] = usage
reason = "".join(reason_parts).strip()
if reason:
turn["reasoning"] = reason[:_REASONING_KEEP]
hist.append(turn)
_run["history"] = hist[-24:]
# 会话落库:本轮追加写回(新会话自动用首条消息作标题)。
# 后台线程 db 访问需 app context。
if conv_id and _flask_app is not None:
try:
with _flask_app.app_context():
saved = list(_run["history"])
title = None
if not db.session.execute(db.text(
"SELECT title FROM agent_conversation WHERE id=:i"),
{"i": conv_id}).scalar():
title = "".join(prompt.split())[:30] or "新会话"
_conv_save_messages(conv_id, saved, title=title)
_log.info(f"会话已保存: {conv_id}({len(saved)} 条消息)")
except Exception as e:
_log.warning(f"会话落库失败: {e}")
# 自进化:成功执行过工具则提炼配方写入经验。必须在 done 之前完成——
# done 发出后 SSE 关流,用户就看不到「已写入经验」的提示了。
# 提炼/保存失败静默(不阻塞、不影响结果),只在成功时推送 🧠 卡片。
# 建任务(designer)模式也沉淀,但**只在草稿通过校验时**:那说明这轮探索
# 确实走通了一条路。没有草稿的探索(试错、半途而废)不入库,免得把误点当经验。
with _lock:
draft_ok = bool(_run.get("draft"))
if tool_seq and (not designer or draft_ok):
try:
recipe = _clean_recipe(_distill_experience(cfg, prompt, " -> ".join(tool_seq)))
if recipe:
if _qualify_experience(prompt, recipe):
if _save_experience(prompt, recipe, " -> ".join(tool_seq)):
_log.info("经验已写入记忆库,随事件流提示")
q.put(("step", {"tool": "🧠 经验记忆",
"args": "本轮操作已提炼为经验并写入记忆库"
"(下次相似任务会自动参考)",
"image": None}))
# 动作经验:把本轮**成功**步骤沉淀为命名动作(带元素定位,禁坐标)
acts = _distill_actions(cfg, prompt, tool_trace)
n_act = _save_actions(prompt, acts)
if n_act:
q.put(("step", {"tool": "🧠 动作经验",
"args": f"已沉淀 {n_act} 个可复用动作"
"(含元素定位,下次同类任务可直接复用)",
"image": None}))
else:
q.put(("step", {"tool": "🧠 动作经验",
"args": "本轮未沉淀动作(步骤以坐标定位为主,"
"缺少可复用的元素定位信息)",
"image": None}))
except Exception as e:
_log.warning(f"经验保存异常: {e}")
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 消息配对检查)
try:
roles = [m.get("role", "?") for m in agent.messages]
tcs = sum(1 for m in agent.messages
if m.get("tool_calls") and isinstance(m.get("tool_calls"), list))
tools_msg = sum(1 for m in agent.messages if m.get("role") == "tool")
_log.warning(f"诊断 messages: {len(agent.messages)} 条 roles={roles[-8:]} "
f"tool_calls消息={tcs} tool回应={tools_msg}")
except Exception:
pass
with _lock:
_run["state"] = "error"
_run["error"] = f"{type(e).__name__}: {str(e)[:200]}"
# 失败也保留已花费的 token(前端仍能展示本轮用量;agent 可能未建出来)
try:
_run["usage"] = dict(agent.usage or {})
except Exception:
pass
q.put(("error", {"message": str(e)[:200]}))
finally:
# 关闭 SSE:扇出把终止哨兵发给**每个**订阅者
if isinstance(q, _Fanout):
q.close()
else:
q.put(None)
_queues.pop(run_id, None)
_stop_events.pop(run_id, None)