diff --git a/mcp_agent/agent.py b/mcp_agent/agent.py index ead0a5b..3d7f28e 100644 --- a/mcp_agent/agent.py +++ b/mcp_agent/agent.py @@ -1,11 +1,13 @@ """Agent 编排层:OpenAI 兼容模型(DeepSeek 等)经 MCP 工具控制手机。 -工作流: - 1. 启动时从 MCP Server 拉工具列表 → 转 OpenAI function schema - 2. run(prompt):循环 chat/completions - - 模型返回 tool_calls → 依次执行(经 MCP)→ 结果回喂 - - de_screenshot 的返回图像转为 image_url 追加为下一轮 user 消息(多模态看图) - - 无 tool_calls → 返回最终文本 +支持两种运行模式: + - 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 回调带缩略 """ import base64 import json @@ -26,9 +28,10 @@ SYSTEM_PROMPT = """你是手机自动化控制助手。你通过工具实时操 1. 先 de_list_devices 确定目标设备(在线才可操作) 2. 观察屏幕:先 de_screenshot 获取截图(图像会随后给你),基于截图理解当前界面 3. 操作:de_tap/de_swipe 的坐标必须与最近一次 de_screenshot 图像一致(直接看图给坐标,服务器自动换算) -4. 每次关键操作后再次 de_screenshot 验证结果,直到完成用户目标 -5. 完成或失败时用中文总结:做了什么、当前状态、需要用户注意的事项 -6. 设备不可用/操作失败时如实报告错误,不要臆测成功 +4. 元素操作优先:能用 de_ui_tree/de_tap_element(text/id/desc 定位)就不用裸坐标 +5. 每次关键操作后再次 de_screenshot 验证结果,直到完成用户目标 +6. 完成或失败时用中文总结:做了什么、当前状态、需要用户注意的事项 +7. 设备不可用/操作失败时如实报告错误,不要臆测成功 可用工具清单将由系统提供。""" @@ -37,9 +40,13 @@ class Agent: def __init__(self, settings: AgentSettings = None): self.s = settings or S self.tools_schema = [] # OpenAI function schema - self._tool_exec = {} # name -> callable self.messages = [] - self.on_step = None # 可选回调 fn(step_dict),web 展示进度用 + # 回调(Web 展示用,均可选): + # on_delta(text, kind) kind: content | reasoning —— 流式文本增量 + # on_tool(step) step: {tool, args, result, image} —— 工具调用完成 + self.on_delta = None + self.on_tool = None + self._mcp = None # ---------- MCP 工具桥 ---------- async def _load_tools(self): @@ -52,37 +59,49 @@ class Agent: for t in tools: # MCP SDK v2 改名 input_schema,兼容新旧字段 schema = getattr(t, "input_schema", None) or getattr(t, "inputSchema", {}) - # fastmcp Tool 属性兼容:name/description/inputSchema name = getattr(t, "name", "") desc = getattr(t, "description", "") or "" self.tools_schema.append({ "type": "function", "function": {"name": name, "description": desc, "parameters": schema}}) - self._tool_exec[name] = t _log.info("MCP 工具已加载: %s", [s["function"]["name"] for s in self.tools_schema]) async def close(self): - if getattr(self, "_mcp", None): - await self._mcp.__aexit__(None, None, None) + if self._mcp: + try: + await self._mcp.__aexit__(None, None, None) + except Exception: + pass - # ---------- 模型调用 ---------- - async def _chat(self): - """调用 OpenAI 兼容 chat/completions,返回完整 response JSON。""" + # ---------- 模型调用(流式) ---------- + async def _chat_stream(self): + """流式 chat/completions:逐 chunk 产出 JSON(async generator)。""" body = { "model": self.s.model, "messages": self.messages, "tools": self.tools_schema if self.tools_schema else None, "max_tokens": 4096, + "stream": 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: - r = await client.post(f"{self.s.api_base.rstrip('/')}/chat/completions", - json=body, headers=headers) - if r.status_code != 200: - raise RuntimeError(f"模型 API HTTP {r.status_code}: {r.text[:300]}") - return r.json() + async with client.stream("POST", url, json=body, headers=headers) as r: + if r.status_code != 200: + text = (await r.aread()).decode(errors="replace") + 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 # ---------- 工具执行 ---------- async def _execute_tool(self, name, arguments): @@ -94,28 +113,33 @@ class Agent: data = getattr(result, "data", result) except Exception as e: return {"ok": False, "error": f"工具执行失败: {e}"}, None - # de_screenshot:图像分离(作为 image_url 追加给模型看) + # 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_step: + if self.on_tool: try: - self.on_step({"tool": name, "args": args, - "result": text_result if image_b64 else data, - "image_b64": image_b64}) + self.on_tool({"tool": name, "args": args, + "result": text_result, "image": image_b64}) except Exception: pass - if image_b64: - return text_result, image_b64 - return data, None + return text_result, image_b64 - # ---------- 主循环 ---------- - async def run(self, prompt: str, serial: str = "") -> str: - """执行一轮指令,返回最终回答文本。""" + # ---------- 主循环(流式) ---------- + async def run_stream(self, prompt: str, serial: str = "", + on_delta=None, on_tool=None): + """流式执行一轮指令,返回最终完整文本。 + + on_delta(text, kind):content/reasoning 文本增量(实时推给前端) + on_tool(step):工具调用完成(实时显示 MCP 步骤) + """ + self.on_delta = on_delta + self.on_tool = on_tool target = serial or self.s.default_serial sys_txt = SYSTEM_PROMPT if target: @@ -123,28 +147,55 @@ class Agent: self.messages = [{"role": "system", "content": sys_txt}, {"role": "user", "content": prompt}] - for step in range(self.s.max_steps): - resp = await self._chat() - choice = (resp.get("choices") or [{}])[0] - msg = choice.get("message") or {} - - # 1) 工具调用 - tool_calls = msg.get("tool_calls") - if tool_calls: - self.messages.append({ - "role": "assistant", - "content": msg.get("content") or "", - "tool_calls": tool_calls}) - for tc in tool_calls: + for _step in range(self.s.max_steps): + content_parts = [] + tool_acc = {} # index -> {id, name, args} + has_tool = False + async for chunk in self._chat_stream(): + 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 {} - name = fn.get("name", "") + if fn.get("name"): + acc["name"] += fn["name"] + if fn.get("arguments"): + acc["args"] += fn["arguments"] + + 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}) + for tc in tcs: + fn = tc["function"] text_result, image_b64 = await self._execute_tool( - name, fn.get("arguments", "{}")) + fn["name"], fn["arguments"]) self.messages.append({ - "role": "tool", - "tool_call_id": tc.get("id", ""), + "role": "tool", "tool_call_id": tc["id"], "content": json.dumps(text_result, ensure_ascii=False)[:4000]}) - # 截图图像:作为下一轮 user 图像内容(OpenAI 协议 tool 结果只能文本) if image_b64: self.messages.append({ "role": "user", @@ -155,7 +206,12 @@ class Agent: f"data:image/jpeg;base64,{image_b64}"}}]}) continue - # 2) 最终回答 - return msg.get("content") or "(模型无输出)" + # 无工具调用:本轮即最终回答 + return full_content return "(达到最大步骤数未完成,请检查操作是否卡在循环)" + + # ---------- 非流式(CLI) ---------- + async def run(self, prompt: str, serial: str = "") -> str: + """非流式执行,返回最终文本(CLI 用,内部走流式收集)。""" + return await self.run_stream(prompt, serial) diff --git a/static/admin/agent.js b/static/admin/agent.js new file mode 100644 index 0000000..4f5eb38 --- /dev/null +++ b/static/admin/agent.js @@ -0,0 +1,198 @@ +// AI 控制台(顶级 Tab):DeepSeek 风格聊天 + 流式输出 + 实时 MCP 步骤 +let _agentBusy = false; +let _agentStream = null; // EventSource +let _agentCfgLoaded = false; + +// ================== 初始化 ================== +function initAgentChat(){ + if(_agentCfgLoaded)return; + _agentCfgLoaded = true; + loadAgentConfig(); + const chat = document.getElementById('agent-chat'); + if(chat && !chat.children.length){ + chat.innerHTML = '
' + esc(d.tool||'') + ''
+ + '' + esc(d.args||'') + '';
+ if(d.image){
+ const img = document.createElement('img');
+ img.src = 'data:image/jpeg;base64,' + d.image;
+ card.appendChild(img);
+ }
+ cards.appendChild(card);
+ scrollChat();
+ });
+
+ es.addEventListener('done', ev=>{
+ const d = JSON.parse(ev.data);
+ if(d.answer){
+ if(!msgEl) msgEl = newAssistantMsg();
+ msgEl.querySelector('.agent-text').textContent = d.answer;
+ }
+ endRun();
+ });
+
+ es.addEventListener('error', ev=>{
+ let msg = '连接中断';
+ try{
+ if(ev.data) msg = JSON.parse(ev.data).message || msg;
+ }catch(e){}
+ if(!msgEl) msgEl = newAssistantMsg();
+ msgEl.querySelector('.agent-text').textContent = '⚠ ' + msg;
+ endRun();
+ });
+
+ es.onerror = ()=>{
+ // 事件流正常结束(done 后服务器关流)会触发一次 error——done 已处理则忽略
+ if(_agentBusy) endRun();
+ };
+}
+
+function endRun(){
+ _agentBusy = false;
+ _agentStream && _agentStream.close();
+ _agentStream = null;
+ document.getElementById('agent-running-tag').style.display = 'none';
+ document.getElementById('btn-agent-send').disabled = false;
+ scrollChat();
+}
diff --git a/static/admin/base.js b/static/admin/base.js
index 548f9a6..88ac53b 100644
--- a/static/admin/base.js
+++ b/static/admin/base.js
@@ -100,6 +100,7 @@ function showTab(name){
if(name==='tasks'){showSubTab('tasks',_activeSubs.tasks);loadTasks();loadCustomActions();}
if(name==='tools'){showSubTab('tools',_activeSubs.tools);loadToolsDevices();loadAdbDevices();loadTailscaleDevices();loadApks();}
if(name==='logs'){loadLogs();if(document.getElementById('log-auto').checked)_logTimer=setInterval(loadLogs,3000);}
+ if(name==='agent' && typeof initAgentChat==='function') initAgentChat();
if(name==='users')loadUsers();
}
@@ -114,7 +115,6 @@ function showSubTab(tabId, name){
tab.querySelectorAll('.sub-tab').forEach(b=>b.classList.toggle('active', b.dataset.sub===name));
tab.querySelectorAll('.sub-panel').forEach(p=>p.classList.toggle('active', p.id===tabId+'-sub-'+name));
if(name==='groups' && typeof loadGroups==='function') loadGroups();
- if(name==='agent' && typeof loadAgentConfig==='function') loadAgentConfig();
if(name==='devpool' && typeof loadDevPool==='function'){
loadDevPool();
if(typeof loadDiscovery==='function'){
diff --git a/static/admin/tools.js b/static/admin/tools.js
index c106461..d81aaf2 100644
--- a/static/admin/tools.js
+++ b/static/admin/tools.js
@@ -41,89 +41,6 @@ async function loadToolsDevices(force){
status.textContent = '共 '+_clipDevices.length+' 台设备';
}
-// ================== AI 控制台(模型配置 + 指令执行 + 进度轮询) ==================
-let _agentPoll = null;
-
-function loadAgentConfig(){
- apiGet('/api/agent/config').then(r=>{
- if(!r||!r.ok)return;
- const b=document.getElementById('agent-api-base');
- if(!b)return; // 面板未渲染
- b.value=r.api_base||'https://api.deepseek.com';
- document.getElementById('agent-model').value=r.model||'';
- document.getElementById('agent-default-serial').value=r.default_serial||'';
- const hint=document.getElementById('agent-key-hint');
- hint.textContent=r.api_key_masked?('已配置 '+r.api_key_masked):'未配置 Key';
- });
-}
-function saveAgentConfig(){
- const body={api_base:document.getElementById('agent-api-base').value.trim(),
- model:document.getElementById('agent-model').value.trim(),
- default_serial:document.getElementById('agent-default-serial').value.trim()};
- const key=document.getElementById('agent-api-key').value.trim();
- if(key)body.api_key=key;
- apiPost('/api/agent/config',body).then(r=>{
- if(r&&r.ok){showToast('配置已保存','success');loadAgentConfig();}
- else showToast('保存失败: '+((r&&r.error)||''),'error');
- });
-}
-function runAgent(){
- const prompt=document.getElementById('agent-prompt').value.trim();
- if(!prompt){showToast('请输入指令','error');return;}
- apiPost('/api/agent/run',{prompt}).then(r=>{
- if(r&&r.ok){
- showToast('Agent 已启动','success');
- document.getElementById('agent-status').textContent='运行中...';
- document.getElementById('btn-agent-run').disabled=true;
- document.getElementById('agent-answer-wrap').style.display='none';
- startAgentPoll();
- }else showToast('启动失败: '+((r&&r.error)||''),'error');
- });
-}
-function startAgentPoll(){
- if(_agentPoll)clearInterval(_agentPoll);
- _agentPoll=setInterval(pollAgent,2000);
- pollAgent();
-}
-function stopAgentPoll(){
- if(_agentPoll){clearInterval(_agentPoll);_agentPoll=null;}
-}
-function pollAgent(){
- apiGet('/api/agent/status').then(r=>{
- if(!r||!r.ok)return;
- const st=document.getElementById('agent-status');
- if(r.state==='running'){
- st.textContent='运行中... ('+(r.steps||[]).length+' 步)';
- }else{
- st.textContent=r.state==='done'?'完成':'失败';
- document.getElementById('btn-agent-run').disabled=false;
- stopAgentPoll();
- }
- // 步骤流(含截图缩略)
- const steps=document.getElementById('agent-steps');
- const html=(r.steps||[]).map(s=>{
- const args=esc(s.args||'');
- const img=s.image
- ?''+esc(s.tool||'')+' '
- +''+args+''+img+'