diff --git a/mcp_server/mcp_server.py b/mcp_server/mcp_server.py index 2abbbda..2728a90 100644 --- a/mcp_server/mcp_server.py +++ b/mcp_server/mcp_server.py @@ -52,6 +52,32 @@ def _check_write(): raise PlatformError("write_disabled", "写操作未启用(MCP_ALLOW_WRITE=1 开启)") +# 设备任务占用锁:serial -> (ts, worker_status, task_job)。AI 写操作前检查, +# 任务 running/connecting 的设备拒绝操作(AI 不与任务抢设备)。5s 缓存。 +_busy_cache = {} + + +def _ensure_device_free(serial): + """写操作前检查设备是否有任务在跑(worker running/connecting → 拒绝)。""" + import time + now = time.time() + c = _busy_cache.get(serial) + if not c or now - c[0] > 5: + try: + devs = platform().list_devices() + info = next((d for d in devs if d.get("serial") == serial), {}) + c = (now, info.get("worker_status") or "idle", + info.get("task_job") or "") + _busy_cache[serial] = c + except Exception: + return # 状态查询失败不阻塞(操作失败会另行报错) + if c[1] in ("running", "connecting"): + raise PlatformError( + "device_busy", + f"设备正在执行任务「{c[2] or '未知'}」——AI 不与任务抢设备," + f"任务结束后才能操作(可在平台任务页先停止任务)") + + def _to_native(serial, x, y): """截图坐标 → 设备原生坐标(按最近一次截图的比例换算)。""" c = _coord.get(serial) @@ -138,6 +164,7 @@ def de_tap(serial: str, x: int, y: int) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if x < 0 or y < 0: raise PlatformError("invalid_param", "坐标不能为负") nx, ny = _to_native(serial, x, y) @@ -160,6 +187,7 @@ def de_swipe(serial: str, x1: int, y1: int, x2: int, y2: int, try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) nx1, ny1 = _to_native(serial, x1, y1) nx2, ny2 = _to_native(serial, x2, y2) platform().swipe(serial, nx1, ny1, nx2, ny2, duration) @@ -217,6 +245,7 @@ def de_tap_element(serial: str, by: str, value: str, index: int = 1) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if by not in ("text", "id", "desc", "text_contains", "desc_contains"): raise PlatformError("invalid_param", "by 可选 text/id/desc/text_contains/desc_contains") @@ -270,6 +299,7 @@ def de_wake(serial: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) platform().wake(serial) except PlatformError as e: return _err(e) @@ -283,6 +313,7 @@ def de_press_key(serial: str, key: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if key not in _KEYS: raise PlatformError("invalid_param", f"不支持的按键: {key}(可选 {_KEYS})") platform().press_key(serial, key) @@ -299,6 +330,7 @@ def de_open_app(serial: str, package: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if not package: raise PlatformError("invalid_param", "缺少包名") from mcp_server import direct_ops @@ -317,6 +349,7 @@ def de_stop_app(serial: str, package: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if not package: raise PlatformError("invalid_param", "缺少包名") from mcp_server import direct_ops @@ -350,6 +383,7 @@ def de_type_text(serial: str, text: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if not text: raise PlatformError("invalid_param", "内容为空") from mcp_server import direct_ops @@ -370,6 +404,7 @@ def de_set_clipboard(serial: str, text: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if not text: raise PlatformError("invalid_param", "内容为空") from mcp_server import direct_ops @@ -390,6 +425,7 @@ def de_sleep(serial: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) platform().sleep(serial) except PlatformError as e: return _err(e) @@ -429,6 +465,7 @@ def de_tap_text(serial: str, text: str) -> dict: try: _check_write() serial = _check_serial(serial) + _ensure_device_free(serial) if not text or len(text) > 100: raise PlatformError("invalid_param", "text 不能为空且 ≤100 字符") res = platform().tap_text(serial, text) diff --git a/static/admin/agent.js b/static/admin/agent.js index 27bf38c..82c5727 100644 --- a/static/admin/agent.js +++ b/static/admin/agent.js @@ -4,6 +4,7 @@ let _agentStream = null; // EventSource let _agentCfgLoaded = false; // ================== 初始化 ================== +let _agentDefaultSerial = ''; function initAgentChat(){ if(_agentCfgLoaded)return; _agentCfgLoaded = true; @@ -12,10 +13,12 @@ function initAgentChat(){ if(chat && !chat.children.length){ chat.innerHTML = '
👋 给 AI 下达指令,它将通过截图观察手机并执行操作。
' + '例如:「打开抖音搜索奚学东,告诉我第一个视频的标题」
' - + '点击右上角 ⚙ 配置模型与 API Key。
'; + + '先选「🎯 目标设备」(AI 只操作你选定的设备),点击右上角 ⚙ 配置模型与 API Key。'; } bindChatScroll(); loadLiveDevices(); + loadAgentTargetDevices(); + startRunPoll(); // 输入框快捷键 const inp = document.getElementById('agent-input'); inp.addEventListener('keydown', ev=>{ @@ -32,6 +35,58 @@ function loadAgentConfig(){ if(!r||!r.ok)return; const tag = document.getElementById('agent-model-tag'); if(tag && r.model) tag.textContent = r.model; + _agentDefaultSerial = r.default_serial || ''; + }); +} + +// ================== 目标设备选择(AI 只操作选定设备) ================== +function loadAgentTargetDevices(){ + apiGet('/api/agent/devices').then(r=>{ + if(!r||!r.ok)return; + const sel = document.getElementById('agent-target-select'); + if(!sel)return; + const cur = sel.value; + const devs = (r.devices||[]).filter(x=>x.online); + sel.innerHTML = '' + + devs.map(d=>{ + const label = d.serial + (d.model ? ' · ' + d.model : '') + + (d.busy ? ' ⛔ 任务中:' + d.task_job : ''); + return ''; + }).join(''); + // 保留当前选择;否则预选配置的默认设备 + if(cur && [...sel.options].some(o=>o.value===cur)) sel.value = cur; + else if(_agentDefaultSerial && [...sel.options].some(o=>o.value===_agentDefaultSerial)) + sel.value = _agentDefaultSerial; + const hint = document.getElementById('agent-target-hint'); + if(hint && sel.value) hint.textContent = '将操作:' + sel.value; + }); +} + +// ================== 跨窗口运行状态(多人/多窗口可见并可停止) ================== +let _runPoll = null; +function startRunPoll(){ + if(_runPoll) return; + _runPoll = setInterval(pollRunState, 8000); +} +function pollRunState(){ + if(_agentStream) return; // 本窗口正在跑(事件流驱动),不轮询 + apiGet('/api/agent/run').then(r=>{ + if(!r||!r.ok)return; + const running = r.state === 'running'; + document.getElementById('agent-running-tag').style.display = running ? 'inline' : 'none'; + document.getElementById('btn-agent-stop').style.display = running ? 'inline-block' : 'none'; + const hint = document.getElementById('agent-target-hint'); + if(hint){ + hint.textContent = running + ? ('⏳ 运行中' + (r.started ? ' ' + r.started + ' 起' : '') + + (r.serial ? ' · ' + r.serial : '') + + ':' + (r.prompt||'').slice(0,70)) + : (document.getElementById('agent-target-select').value + ? '将操作:' + document.getElementById('agent-target-select').value : ''); + } + document.getElementById('btn-agent-send').disabled = running; + if(!running && _agentBusy){ setRunning(false); } // 流异常丢失时复位 }); } function openAgentConfig(){ @@ -64,6 +119,7 @@ function saveAgentConfig(){ showToast('配置已保存','success'); closeAgentConfig(); loadAgentConfig(); + loadAgentTargetDevices(); }else showToast('保存失败: ' + ((r&&r.error)||''),'error'); }); } @@ -129,6 +185,7 @@ function setRunning(on){ document.getElementById('btn-agent-send').disabled = on; } function stopAgent(){ + if(!confirm('确定停止当前 AI 运行任务?\n(其他人发起的任务也会被停止)')) return; apiPost('/api/agent/stop',{}).then(r=>{ if(r&&r.ok) showToast('已请求停止,正在中断…','success'); else showToast((r&&r.error)||'停止失败','error'); @@ -136,13 +193,20 @@ function stopAgent(){ } function sendAgentMsg(){ if(_agentBusy){showToast('上一轮还在运行','error');return;} + const sel = document.getElementById('agent-target-select'); + const serial = sel ? sel.value : ''; + if(!serial){ + showToast('请先选择目标设备(AI 只操作你选定的设备;任务中的设备不可选)','error'); + if(sel) sel.focus(); + return; + } const inp = document.getElementById('agent-input'); const text = inp.value.trim(); if(!text){showToast('请输入指令','error');return;} inp.value = ''; addUserMsg(text); - apiPost('/api/agent/run', {prompt: text}).then(r=>{ + apiPost('/api/agent/run', {prompt: text, serial: serial}).then(r=>{ if(!r||!r.ok){ const msg = r ? (r.error||'启动失败') : '请求失败'; const div = newAssistantMsg(); diff --git a/templates/admin/monitor.html b/templates/admin/monitor.html index d5c4b5f..31e37a8 100644 --- a/templates/admin/monitor.html +++ b/templates/admin/monitor.html @@ -467,7 +467,6 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
🤖 AI 控制台
- @@ -487,10 +486,22 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14 onchange="watchDevice(this.value)">
-
- - +
+
+ + + + + +
+
+ + +
diff --git a/web/agent_api.py b/web/agent_api.py index e76adbd..b6f7ffe 100644 --- a/web/agent_api.py +++ b/web/agent_api.py @@ -498,23 +498,85 @@ def agent_run(): 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 锁兜底) with _lock: if _run["state"] == "running": - return jsonify({"ok": False, "error": "已有 Agent 运行中,请等待完成"}), 409 + 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=(data.get("serial") or "").strip(), + serial=serial, + started=_dt.now().strftime("%H:%M:%S"), answer="", error="") # history 保留(同会话多轮对话),由前端「清空对话」调用 clear 重置 _queues[run_id] = queue.Queue() _stop_events[run_id] = threading.Event() - _log.info(f"Agent 启动: {prompt[:60]}") + _log.info(f"Agent 启动: {prompt[:60]} @ {serial or 'default'}") threading.Thread(target=_agent_thread, args=(run_id, prompt, serial, cfg), daemon=True).start() return jsonify({"ok": True, "run_id": run_id}) +@bp.route("/api/agent/run", methods=["GET"]) +@admin_required +def agent_run_status(): + """当前 Agent 运行状态(多窗口/多人可见):idle / running + 任务摘要。 + + 前端轮询它同步「运行中」状态(换浏览器/他人启动的任务也能看到并停止)。 + """ + with _lock: + return jsonify({"ok": True, + "state": _run.get("state", "idle"), + "prompt": _run.get("prompt", ""), + "serial": _run.get("serial", ""), + "started": _run.get("started", "")}) + + +@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"), + "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():