#!/usr/bin/env python3 """ Hermes Bridge - 将 Hermes AIAgent 封装为 HTTP API,供 Spring Boot 工作流引擎调用。 启动方式: python hermes_bridge.py --port 18731 python hermes_bridge.py --port 18731 --hermes-home /path/to/.hermes API: POST /run 运行一次 Agent 对话(SSE 流式响应) POST /skills/invalidate 清空技能索引缓存与 Agent 池(技能写时同步后调用) GET /health 健康检查 """ import argparse import json import logging import os import sys import threading import time import uuid from collections import OrderedDict, deque from pathlib import Path from typing import Optional # --------------------------------------------------------------------------- # 强制 stdout/stderr 使用 UTF-8(Windows 默认 cp936 会导致中文输出乱码) # 必须在任何 print/logging 之前执行 # --------------------------------------------------------------------------- for _stream in (sys.stdout, sys.stderr): try: _stream.reconfigure(encoding="utf-8", errors="replace") except Exception: pass # --------------------------------------------------------------------------- # 父进程守护:Java 后端被 kill 后,子进程自动退出,释放端口与资源 # --------------------------------------------------------------------------- def _start_parent_watcher(): """启动守护线程,当父进程退出时自杀。""" try: import psutil except ImportError: logging.warning("未安装 psutil,无法监听父进程状态;Java 退出后子进程可能残留") return try: parent = psutil.Process(os.getppid()) except Exception: return def _watch(): while True: time.sleep(2) try: if not parent.is_running() or parent.status() == psutil.STATUS_ZOMBIE: logging.info("父进程已退出,Hermes Bridge 自动终止") os._exit(0) except Exception: logging.info("父进程状态不可获取,Hermes Bridge 自动终止") os._exit(0) watcher = threading.Thread(target=_watch, daemon=True, name="parent-watcher") watcher.start() _start_parent_watcher() # --------------------------------------------------------------------------- # Hermes 项目路径注入(让 Python 能找到 hermes-agent 的源码) # --------------------------------------------------------------------------- BRIDGE_DIR = Path(__file__).resolve().parent HERMES_AGENT_DIR = BRIDGE_DIR.parent / "hermes-agent" if HERMES_AGENT_DIR.is_dir() and str(HERMES_AGENT_DIR) not in sys.path: sys.path.insert(0, str(HERMES_AGENT_DIR)) # --------------------------------------------------------------------------- # 日志 # --------------------------------------------------------------------------- logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", stream=sys.stderr, ) logger = logging.getLogger("hermes-bridge") # --------------------------------------------------------------------------- # FastAPI app # --------------------------------------------------------------------------- try: from fastapi import FastAPI, HTTPException from fastapi.responses import StreamingResponse from pydantic import BaseModel except ImportError: logger.error("请先安装 FastAPI: pip install fastapi uvicorn[standard]") sys.exit(1) app = FastAPI(title="Hermes Bridge", version="1.0.0") # --------------------------------------------------------------------------- # 共享密钥(X-Bridge-Token 头校验) # --------------------------------------------------------------------------- # 启动时从环境变量读取;为空则禁用认证(仅依赖 127.0.0.1 绑定) _BRIDGE_AUTH_TOKEN = os.environ.get("HERMES_BRIDGE_AUTH_TOKEN", "").strip() @app.middleware("http") async def _verify_token(request, call_next): """校验 X-Bridge-Token 头(/health 不校验)""" if _BRIDGE_AUTH_TOKEN and request.url.path != "/health": token = request.headers.get("X-Bridge-Token", "") if token != _BRIDGE_AUTH_TOKEN: from fastapi.responses import JSONResponse return JSONResponse(status_code=401, content={"detail": "invalid or missing X-Bridge-Token"}) return await call_next(request) # --------------------------------------------------------------------------- # 全局 Agent 池(按 session_id::hermes_home 复用,有界 LRU) # --------------------------------------------------------------------------- _agent_lock = threading.Lock() _agents: "OrderedDict[str, object]" = OrderedDict() _AGENT_CACHE_MAX = 8 # 最大缓存实例数;超出按 FIFO 淘汰,防止无界增长 # --------------------------------------------------------------------------- # 全局 clarify 会话映射:clarify_id → session_key # 用于 POST /clarify/{clarify_id}/answer 端点定位 _ClarifyEntry # (session_key 在 register 时与 clarify_id 绑定,可直接 resolve) # --------------------------------------------------------------------------- _clarify_lock = threading.Lock() _clarify_ids: "set[str]" = set() # 仅用于跟踪 Bridge 模式发起的 clarify_id # --------------------------------------------------------------------------- # 请求模型 # --------------------------------------------------------------------------- class RunRequest(BaseModel): """运行 Agent 对话的请求体""" system_prompt: Optional[str] = None # 系统提示(如 SKILL.md 内容) user_message: str # 用户消息 max_iterations: int = 100 # 最大工具调用轮次(兜底默认值,由后端配置覆盖) session_id: Optional[str] = None # 会话 ID(复用 Agent 实例) hermes_home: Optional[str] = None # 本次运行使用的 HERMES_HOME 路径(其下应有 skills 子目录) working_dir: Optional[str] = None # agent 工具调用的实际工作目录(终端、文件读写以它为根) model_id: Optional[str] = None # 节点指定的模型 ID base_url: Optional[str] = None # 节点指定的模型 base URL api_key: Optional[str] = None # 节点指定的模型 API Key model_name: Optional[str] = None # 节点指定的模型标识 # --------------------------------------------------------------------------- # Agent 工厂 # --------------------------------------------------------------------------- def _get_or_create_agent(session_id: str, max_iterations: int, hermes_home: Optional[str] = None, model_id: Optional[str] = None, base_url: Optional[str] = None, api_key: Optional[str] = None, model_name: Optional[str] = None): """获取或创建 AIAgent 实例 hermes_home 用于覆盖本次运行的 HERMES_HOME(其下应有 skills 子目录), 实现按工作流运行目录隔离的 Skill 加载。 model_id/base_url/api_key/model_name 支持按节点动态切换模型; 未提供时回退到环境变量默认配置。 """ # 缓存键同时绑定 hermes_home 与模型 ID,避免跨运行目录/模型复用 agent cache_key = f"{session_id}::{hermes_home or ''}::{model_id or 'default'}" if session_id else None if cache_key: with _agent_lock: if cache_key in _agents: _agents.move_to_end(cache_key) # LRU 提升到最新 return _agents[cache_key] # 在导入/创建 AIAgent 之前显式覆盖 HERMES_HOME,使 get_hermes_home() 返回正确路径 if hermes_home: os.environ["HERMES_HOME"] = str(hermes_home) logger.info("使用请求指定的 HERMES_HOME: %s", hermes_home) # 安装沙箱补丁(幂等),强制 hermes-agent 文件操作限制在 _SESSION_CWD 内 try: import sandbox_patch sandbox_patch.install_sandbox_patches() except Exception as e: logger.warning("沙箱补丁加载异常(继续但不保证隔离): %s", e) try: from run_agent import AIAgent except ImportError: raise ImportError( "无法导入 hermes-agent。请确保 hermes-agent 目录存在或已安装 hermes-agent 包。\n" f"尝试的路径: {HERMES_AGENT_DIR}" ) # 优先使用请求中传入的模型配置,其次回退到环境变量 effective_base_url = base_url or os.environ.get("HERMES_BRIDGE_BASE_URL", "") effective_api_key = api_key or os.environ.get("HERMES_BRIDGE_API_KEY", "") effective_model = model_name or os.environ.get("HERMES_BRIDGE_MODEL", "") # 尝试从 Hermes config.yaml 读取(如果未通过请求或环境变量指定) if not effective_base_url or not effective_api_key: _load_hermes_config() effective_base_url = effective_base_url or os.environ.get("HERMES_BRIDGE_BASE_URL", "") effective_api_key = effective_api_key or os.environ.get("HERMES_BRIDGE_API_KEY", "") effective_model = effective_model or os.environ.get("HERMES_BRIDGE_MODEL", "") if not effective_api_key: raise ValueError("未配置 API Key。请设置 HERMES_BRIDGE_API_KEY 环境变量或在 Hermes config.yaml 中配置。") logger.info( "创建 AIAgent: session=%s, model_id=%s, model=%s, base_url=%s", session_id, model_id or "默认", effective_model or "默认", effective_base_url or "默认" ) kwargs = dict( base_url=effective_base_url or None, api_key=effective_api_key, model=effective_model or None, max_iterations=max_iterations, quiet_mode=True, tool_progress_mode="off", skip_context_files=True, load_soul_identity=False, skip_memory=True, ) agent = AIAgent(**kwargs) if cache_key: with _agent_lock: _agents[cache_key] = agent # 超出上限按 FIFO 淘汰最旧实例 while len(_agents) > _AGENT_CACHE_MAX: _agents.popitem(last=False) return agent def _load_hermes_config(): """尝试从 Hermes 的 config.yaml 和 .env 加载配置""" try: hermes_home = Path(os.environ.get("HERMES_HOME", Path.home() / ".hermes")) # 加载 .env env_file = hermes_home / ".env" if env_file.exists(): for line in env_file.read_text(encoding="utf-8").splitlines(): line = line.strip() if not line or line.startswith("#") or "=" not in line: continue key, _, value = line.partition("=") key = key.strip() value = value.strip().strip('"').strip("'") # 只设置 Bridge 相关的环境变量(如果尚未设置) if key in ("OPENAI_API_KEY", "OPENROUTER_API_KEY") and not os.environ.get("HERMES_BRIDGE_API_KEY"): os.environ["HERMES_BRIDGE_API_KEY"] = value if key == "OPENAI_BASE_URL" and not os.environ.get("HERMES_BRIDGE_BASE_URL"): os.environ["HERMES_BRIDGE_BASE_URL"] = value # 加载 config.yaml config_file = hermes_home / "config.yaml" if config_file.exists(): import yaml with open(config_file, "r", encoding="utf-8") as f: config = yaml.safe_load(f) or {} model_cfg = config.get("model") or {} provider = config.get("provider") or {} if not os.environ.get("HERMES_BRIDGE_MODEL") and model_cfg.get("name"): os.environ["HERMES_BRIDGE_MODEL"] = model_cfg["name"] if not os.environ.get("HERMES_BRIDGE_BASE_URL") and provider.get("base_url"): os.environ["HERMES_BRIDGE_BASE_URL"] = provider["base_url"] if not os.environ.get("HERMES_BRIDGE_API_KEY") and provider.get("api_key"): os.environ["HERMES_BRIDGE_API_KEY"] = provider["api_key"] except Exception as e: logger.warning("加载 Hermes 配置失败: %s", e) # --------------------------------------------------------------------------- # SSE 事件流 # --------------------------------------------------------------------------- def _run_agent_stream(agent, system_prompt: str, user_message: str, run_id: str, working_dir: Optional[str] = None, session_id: Optional[str] = None): """运行 Agent 并生成 SSE 事件流 事件类型(v1.1 扩展): - text: 流式文本片段(LLM 输出) - tool_start: 工具调用开始(含工具名与参数) - tool_end: 工具调用结束(含截断后的结果摘要) - todo_update: todo 工具的完整 JSON 输出(不截断,独立事件) - step: 每轮 LLM 调用前进度回调 - status: 生命周期/警告类状态消息 - thinking_status: TUI 风格思考 spinner 文案(与 text 区分) - tool_progress: 工具细粒度进度(tool.started / reasoning.available 等) - clarify_request: 模型请求用户确认的问题(含 clarify_id 与选项) - done: 执行完成,包含最终回复 - error: 执行出错 """ # L-4:deque 的 popleft() 是 O(1);原 list.pop(0) 是 O(n),事件多时每次复制整个数组浪费 CPU queue = deque() done_event = threading.Event() def on_stream(delta: str): """接收流式文本片段""" queue.append(("text", {"content": delta})) def on_tool_start(tool_call_id, tool_name: str, args: dict): """工具开始。 注意:Hermes tool_start_callback 实际签名为 (tool_call_id, name, args), 必须接收 3 个位置参数,否则会触发 TypeError 被 hermes-agent 静默吞掉。 """ queue.append(("tool_start", {"tool": tool_name, "args": args})) def on_tool_end(tool_call_id, tool_name: str, args, result): """工具结束。 注意:Hermes tool_complete_callback 实际签名为 (tool_call_id, name, args, result), 必须接收 4 个位置参数。result 可能是 str / list / dict,统一转字符串处理。 todo 工具特殊处理:不截断、发独立 todo_update 事件。 """ # result 统一转字符串(OpenAI tool_result 可能是 list/dict) if isinstance(result, str): result_text = result elif result is None: result_text = "" else: try: result_text = json.dumps(result, ensure_ascii=False) except Exception: result_text = str(result) if tool_name == "todo": # 完整保留 todo 工具的 JSON 输出(用于前端 TODO List 渲染) queue.append(("todo_update", {"tool": tool_name, "result": result_text})) # 同时发标准 tool_end(保持向后兼容,前端旧逻辑仍能看到调用记录) summary = result_text[:500] queue.append(("tool_end", {"tool": tool_name, "result_summary": summary})) return summary = result_text[:500] queue.append(("tool_end", {"tool": tool_name, "result_summary": summary})) def on_step(api_call_count: int, prev_tools: list): """每轮 LLM 调用前的进度回调""" queue.append(("step", { "iteration": api_call_count, "prev_tools": list(prev_tools) if prev_tools else [], })) def on_status(category: str, message: str): """生命周期/警告类状态消息""" queue.append(("status", {"category": category, "message": message})) def on_thinking(content: str): """TUI 风格思考 spinner 文案(如 '🤔 思考中...') 与 text 事件区分:text 是 LLM 真实输出片段,thinking_status 是 UI 状态提示""" if content: # 忽略清空信号(content == ""),避免噪声 queue.append(("thinking_status", {"content": content})) def on_tool_progress(event_type: str, name: Optional[str] = None, preview: Optional[str] = None, args: Optional[dict] = None, **kwargs): """工具细粒度进度事件(tool.started / reasoning.available / _thinking 等)""" queue.append(("tool_progress", { "event_type": event_type, "name": name, "preview": preview, "args": args, })) def on_clarify(question: str, choices=None) -> str: """同步阻塞 clarify 回调 Bridge 模式下: 1. 生成 clarify_id,注册到 clarify_gateway 2. 发 clarify_request SSE 事件(Java 收到后会持久化 + 推给前端) 3. 阻塞等待 POST /clarify/{clarify_id}/answer 唤醒(无超时) 4. 返回用户答案给 agent """ import secrets from tools import clarify_gateway clarify_id = secrets.token_hex(8) session_key = f"bridge-{run_id}" clarify_gateway.register( clarify_id=clarify_id, session_key=session_key, question=question, choices=list(choices) if choices else None, ) with _clarify_lock: _clarify_ids.add(clarify_id) queue.append(("clarify_request", { "clarify_id": clarify_id, "run_id": run_id, "question": question, "choices": list(choices) if choices else None, })) logger.info("[%s] Clarify 等待用户回答: clarify_id=%s, question=%s", run_id, clarify_id, question[:80]) # 无超时等待;实际超时由 Java/HermesBridgeClient 的 24h 兜底控制 response = clarify_gateway.wait_for_response(clarify_id, timeout=None) with _clarify_lock: _clarify_ids.discard(clarify_id) if response is None: # 不应发生(None 表示超时,但 timeout=None 永不超时),降级返回空让 agent 自决策 logger.warning("[%s] Clarify 意外返回 None: clarify_id=%s", run_id, clarify_id) return "" logger.info("[%s] Clarify 收到回答: clarify_id=%s, response=%s", run_id, clarify_id, response[:80]) return response def _run(): # 通过 contextvar 设置本次运行的 agent 工作目录; # hermes-agent 内 resolve_agent_cwd() 会优先读取此值, # 影响 agent 工具(终端、文件读写)的当前目录。 cwd_token = None original_cwd = None effective_task_id = session_id or "default" if working_dir: wd = str(working_dir) # 同时切换进程 cwd,覆盖那些直接调用 os.getcwd() 的工具 try: original_cwd = os.getcwd() os.chdir(wd) logger.info("[%s] 切换进程 cwd: %s", run_id, wd) except Exception as e: logger.warning("[%s] 切换 cwd 失败: %s", run_id, e) # 设置 TERMINAL_CWD 作为 local environment 的 fallback cwd os.environ["TERMINAL_CWD"] = wd # 注册 task env override,强制更新(或创建)该 session 的 terminal environment cwd try: from tools.terminal_tool import register_task_env_overrides register_task_env_overrides(effective_task_id, {"cwd": wd}) logger.info("[%s] 注册 terminal cwd override for task %s: %s", run_id, effective_task_id, wd) except Exception as e: logger.warning("[%s] 注册 terminal cwd override 失败: %s", run_id, e) try: from agent.runtime_cwd import _SESSION_CWD cwd_token = _SESSION_CWD.set(wd) logger.info("[%s] 设置 agent 工作目录 contextvar: %s", run_id, wd) except Exception as e: logger.warning("[%s] 设置工作目录失败(agent 将使用进程 cwd): %s", run_id, e) try: agent.tool_start_callback = on_tool_start agent.tool_complete_callback = on_tool_end agent.step_callback = on_step agent.status_callback = on_status agent.thinking_callback = on_thinking agent.tool_progress_callback = on_tool_progress agent.clarify_callback = on_clarify result = agent.run_conversation( user_message=user_message, system_message=system_prompt, stream_callback=on_stream, task_id=effective_task_id, ) final = result.get("final_response", "") if isinstance(result, dict) else str(result) queue.append(("done", {"content": final})) except Exception as e: logger.error("Agent 执行异常: %s", e, exc_info=True) queue.append(("error", {"message": str(e)})) finally: # 清理本 run 未完成的 clarify(agent 异常退出时不留孤儿阻塞线程) # L-5:clear_session 失败时记录 warning,便于发现 clarify_gateway 内部异常 # (原 `except: pass` 完全静默,问题发生时无任何线索可查) try: from tools import clarify_gateway clarify_gateway.clear_session(f"bridge-{run_id}") except Exception as cleanup_err: logger.warning("[%s] 清理 clarify 会话失败: %s", run_id, cleanup_err) if cwd_token is not None: try: _SESSION_CWD.reset(cwd_token) except Exception: pass if original_cwd is not None: try: os.chdir(original_cwd) except Exception: pass done_event.set() thread = threading.Thread(target=_run, daemon=True) thread.start() while not done_event.is_set() or queue: while queue: event_type, data = queue.popleft() yield f"event: {event_type}\ndata: {json.dumps(data, ensure_ascii=False)}\n\n" time.sleep(0.05) # 刷出最后可能残留的事件 while queue: event_type, data = queue.popleft() yield f"event: {event_type}\ndata: {json.dumps(data, ensure_ascii=False)}\n\n" # --------------------------------------------------------------------------- # API 端点 # --------------------------------------------------------------------------- @app.post("/run") async def run_agent(req: RunRequest): """ 运行 Hermes Agent,返回 SSE 流式响应。 事件类型: - text: 流式文本片段 - tool_start: 工具调用开始 - tool_end: 工具调用结束 - done: 执行完成,包含最终回复 - error: 执行出错 """ if not req.user_message.strip(): raise HTTPException(status_code=400, detail="user_message 不能为空") session_id = req.session_id or "" run_id = uuid.uuid4().hex[:8] try: agent = _get_or_create_agent( session_id, req.max_iterations, req.hermes_home, model_id=req.model_id, base_url=req.base_url, api_key=req.api_key, model_name=req.model_name ) except (ImportError, ValueError) as e: raise HTTPException(status_code=503, detail=str(e)) return StreamingResponse( _run_agent_stream(agent, req.system_prompt, req.user_message, run_id, req.working_dir, session_id), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "X-Accel-Buffering": "no", "Connection": "keep-alive", }, ) @app.post("/clarify/{clarify_id}/answer") async def clarify_answer(clarify_id: str, payload: dict): """接收前端送回的 clarify 答案,唤醒阻塞中的 agent 线程。 Bridge 模式下,agent 调用 clarify 工具时会在 on_clarify 回调里阻塞等待。 Java/HermesBridgeClient 收到 clarify_request SSE 事件后持久化到 DB 并推给前端, 前端展示 Askfor UI;用户作答后通过 Java 转发到本端点。 请求体:{"answer": ""} 返回:{"ok": true, "resolved": true} 或 {"ok": true, "resolved": false}(已超时/不存在) """ from tools import clarify_gateway answer = (payload or {}).get("answer", "") if not isinstance(answer, str): answer = str(answer) if answer is not None else "" resolved = clarify_gateway.resolve_gateway_clarify(clarify_id, answer) if not resolved: with _clarify_lock: known = clarify_id in _clarify_ids if not known: logger.warning("收到未知 clarify_id 的回答: %s", clarify_id) else: logger.info("clarify_id %s 已被清理(agent 已退出)", clarify_id) else: logger.info("Clarify 回答已送达: clarify_id=%s", clarify_id) return {"ok": True, "resolved": resolved} @app.post("/skills/invalidate") async def skills_invalidate(): """ 清空技能相关缓存,使 Java 侧新同步到 HERMES_HOME/skills 的技能立即生效。 两层缓存: 1. agent.prompt_builder._SKILLS_PROMPT_CACHE —— 进程内 LRU,key 不含目录 mtime, 同一 skills_dir 下新增技能后不会自动失效,必须显式清空; 2. _agents —— Agent 池,系统提示词(含技能索引)在 Agent 创建时构建一次, 清空后下次请求重建 Agent 以拉取最新技能索引。 磁盘快照 .skills_prompt_snapshot.json 自带 mtime/size 校验,文件变更后自动失效,无需处理。 """ cleared_prompt_cache = False try: from agent.prompt_builder import clear_skills_system_prompt_cache clear_skills_system_prompt_cache() cleared_prompt_cache = True except Exception as e: logger.warning("清空 skills prompt 缓存失败(可能尚未加载 prompt_builder): %s", e) with _agent_lock: evicted = len(_agents) _agents.clear() logger.info("技能缓存已失效: prompt_cache=%s, evicted_agents=%d", cleared_prompt_cache, evicted) return {"ok": True, "promptCacheCleared": cleared_prompt_cache, "evictedAgents": evicted} @app.get("/health") async def health(): """健康检查""" hermes_home = os.environ.get("HERMES_HOME", str(Path.home() / ".hermes")) hermes_agent_found = HERMES_AGENT_DIR.is_dir() # 尝试导入验证 can_import = False try: import run_agent can_import = True except ImportError: pass return { "status": "ok" if (hermes_agent_found or can_import) else "degraded", "hermes_home": hermes_home, "hermes_agent_dir": str(HERMES_AGENT_DIR), "hermes_agent_found": hermes_agent_found, "can_import_agent": can_import, "active_sessions": len(_agents), } # --------------------------------------------------------------------------- # 启动 # --------------------------------------------------------------------------- def main(): parser = argparse.ArgumentParser(description="Hermes Bridge Server") parser.add_argument("--port", type=int, default=18731, help="监听端口 (默认 18731)") parser.add_argument("--host", default="127.0.0.1", help="监听地址 (默认 127.0.0.1)") parser.add_argument("--hermes-home", default=None, help="Hermes Home 目录") args = parser.parse_args() if args.hermes_home: os.environ["HERMES_HOME"] = args.hermes_home import uvicorn # M-7:必须单 worker 运行!_agents(OrderedDict)与 _clarify_ids(set)是进程级状态, # clarify 还依赖进程内阻塞线程 + set 内的 clarify_id 来路由答案。多 worker 会导致: # - Agent 缓存命中率骤降(每个 worker 独立 _agents,反复重建) # - clarify_answer 命中非阻塞 worker 时 resolve 返回 false,前端提交的回答丢失 # 这里直接传 app 实例(而非 import 字符串),uvicorn 会强制单 worker; # 若未来改为 uvicorn.run("hermes_bridge:app", workers=N) 会破坏该约束,需重新设计共享状态。 logger.info("启动 Hermes Bridge (单 worker): %s:%d", args.host, args.port) uvicorn.run(app, host=args.host, port=args.port, log_level="info", workers=1) if __name__ == "__main__": main()