""" agent67.core.assistant — lambda 的大脑 🐂 (v2) v2 升级: - 使用 lambdagent 内置 21 工具(替代手写 6 工具) - ToolGateway 安全网关(替代手写 safety.py) - Hook 系统全局审计 - ContextManager 防状态爆炸 - 保留 macOS 专属工具(browser, app, system, screenshot) Lambda 表达式: lambda = Memory(Loop(Brain >> parse_action >> route_execute)) """ from __future__ import annotations import json import re import sys import os from pathlib import Path from typing import Any # 路径设置 PROJECT_ROOT = Path(__file__).resolve().parent.parent.parent.parent sys.path.insert(0, str(PROJECT_ROOT)) from lambdagent import Lam, Tool, Loop, Context from lambdagent.extensions import Memory, Guard from lambdagent.context_manager import ContextManager from lambdagent.hooks import HookRegistry # 内置工具 from lambdagent.builtin_tools.registry import BUILTIN_TOOLS, resolve_tools from lambdagent.builtin_tools.terminal_ui import TerminalUI from .claude_lam import ClaudeLam from .ollama_lam import OllamaLam from .config import BACKEND_OLLAMA, BACKEND_CLAUDE_CODE, BACKEND_API from .prompt import SYSTEM_PROMPT # macOS 专属工具(保留) from ..tools.browser_controller import control_browser from ..tools.app_controller import control_app from ..tools.system_info import query_system from ..tools.screenshot import take_screenshot from ..tools.research_workflow import research_workflow # ════════════════════════════════════════════════════════════ # 工具注册表 (v2: 内置工具 + macOS 工具) # ════════════════════════════════════════════════════════════ def _build_tool_registry() -> dict: """合并内置工具 + macOS 专属工具""" # 内置 21 个工具 tools = dict(BUILTIN_TOOLS) # macOS 专属(保留原有能力) tools["browser"] = Tool("browser", control_browser) tools["app"] = Tool("app", control_app) tools["system"] = Tool("system", query_system) tools["screenshot"] = Tool("screenshot", take_screenshot) # 科研流程 (research skill pack) tools["ResearchWorkflow"] = Tool("ResearchWorkflow", research_workflow) # 完成信号 tools["done"] = Tool("done", lambda x: x) return tools TOOL_REGISTRY = _build_tool_registry() # ════════════════════════════════════════════════════════════ # 解析 & 执行引擎 (v2) # ════════════════════════════════════════════════════════════ def _extract_json_blocks(text: str) -> list[str]: """Extract JSON code blocks, handling nested ``` in content. Strategy: find ```json markers, then find the matching ``` closer by looking for a ``` that is followed by a newline or end-of-string (not inside a JSON string value). """ blocks = [] i = 0 marker = "```json" while i < len(text): start = text.find(marker, i) if start == -1: break # Skip past the marker and optional newline content_start = start + len(marker) if content_start < len(text) and text[content_start] == '\n': content_start += 1 # Find the closing ``` — look for ``` at start of line or after newline # that's NOT inside a JSON string best_end = -1 j = content_start while j < len(text): pos = text.find("```", j) if pos == -1: break # Check if this ``` is a real closer (not inside JSON string content) # Heuristic: if the text between content_start and pos contains # a valid-looking JSON with "action", it's the closer candidate = text[content_start:pos].strip() if candidate and ('"action"' in candidate or '"tool"' in candidate): try: # Try to parse — if it works, this is the right closer _try_parse_json(candidate) best_end = pos break except Exception: pass # Even if parse fails, if it looks like JSON, take it if candidate.startswith("{") and candidate.rstrip().endswith("}"): best_end = pos break j = pos + 3 if best_end != -1: block = text[content_start:best_end].strip() if block: blocks.append(block) i = best_end + 3 else: # Fallback: take everything until the next ``` end = text.find("```", content_start) if end != -1: block = text[content_start:end].strip() if block: blocks.append(block) i = end + 3 else: break return blocks def _try_parse_json(json_str: str) -> dict | None: """Try multiple strategies to parse potentially malformed JSON from LLM output.""" # Strategy 1: direct parse try: return json.loads(json_str) except json.JSONDecodeError: pass # Strategy 2: fix newlines inside JSON string values # LLM outputs real newlines inside "content": "line1\nline2" # We need to escape them but ONLY inside string values try: fixed = _fix_json_newlines(json_str) return json.loads(fixed) except (json.JSONDecodeError, Exception): pass # Strategy 3: balanced brace extraction + newline fix try: depth = 0 start = json_str.index("{") in_string = False escape = False for i in range(start, len(json_str)): c = json_str[i] if escape: escape = False continue if c == '\\': escape = True continue if c == '"': in_string = not in_string continue if not in_string: if c == '{': depth += 1 elif c == '}': depth -= 1 if depth == 0: candidate = json_str[start:i+1] fixed = _fix_json_newlines(candidate) return json.loads(fixed) except (json.JSONDecodeError, ValueError): pass return None def _fix_json_newlines(s: str) -> str: """Replace real newlines inside JSON string values with \\n. Walks the string character by character, only replacing newlines when inside a quoted string value. """ result = [] in_string = False escape = False for c in s: if escape: result.append(c) escape = False continue if c == '\\': result.append(c) escape = True continue if c == '"': in_string = not in_string result.append(c) continue if c == '\n' and in_string: result.append('\\n') continue result.append(c) return ''.join(result) def parse_and_execute(llm_output: str, hooks: HookRegistry = None) -> tuple[str, bool, str]: """ 解析 LLM 输出,提取工具调用并执行。 Returns: (result_text, is_done, tool_name) """ matches = _extract_json_blocks(llm_output) if not matches: matches = re.findall(r'\{[^{}]*(?:"action"|"tool")\s*:.*?\}', llm_output, re.DOTALL) if not matches: return llm_output, False, "" # 尝试每个 match(从第一个开始),找到第一个包含 action/tool 的可解析 JSON data = None for candidate in matches: candidate = candidate.strip() parsed = _try_parse_json(candidate) if parsed and ("action" in parsed or "tool" in parsed): data = parsed break if data is None: return f"(无法解析 JSON: {matches[0][:100]})\n{llm_output}", False, "" # 统一格式: action/tool 都支持 tool_name = data.get("action", data.get("tool", "")) # 完成信号 if tool_name in ("done", "terminate"): # 兼容多种字段: summary / input.answer / args.answer / args.summary / 顶层 answer answer = "" for holder in (data.get("input"), data.get("args")): if isinstance(holder, dict): answer = holder.get("answer") or holder.get("summary") or "" if answer: break if not answer: answer = data.get("summary") or data.get("answer") or "" return (str(answer) if answer else "任务完成"), True, tool_name # 路由到工具 tool = TOOL_REGISTRY.get(tool_name) if not tool: return f"❌ 未知工具: {tool_name}. 可用: {', '.join(sorted(TOOL_REGISTRY.keys())[:15])}...", False, tool_name # 准备输入 tool_input = data.get("input", {}) if not tool_input: # 兼容旧格式: {"tool": "shell", "command": "ls"} → input = {"command": "ls"} tool_input = {k: v for k, v in data.items() if k not in ("action", "tool")} if isinstance(tool_input, dict): tool_input_str = json.dumps(tool_input, ensure_ascii=False) else: tool_input_str = str(tool_input) # Hook: pre_tool if hooks: hooks.fire("pre_tool", term=tool, input=tool_input_str, ctx=None) # 执行 import time t0 = time.time() try: if hasattr(tool, 'apply'): result = tool.apply(tool_input_str) else: result = tool(tool_input_str) except Exception as e: result = f"[ERROR] {e}" duration_ms = (time.time() - t0) * 1000 # Hook: post_tool if hooks: output_wrapper = {"value": result} hooks.fire("post_tool", term=tool, input=tool_input_str, output=output_wrapper, duration_ms=duration_ms, ctx=None) result = output_wrapper["value"] return str(result), False, tool_name # ════════════════════════════════════════════════════════════ # PersonalAssistant v2 # ════════════════════════════════════════════════════════════ class PersonalAssistant: """ lambda 个人助理 v2 = Memory(ReAct(Brain)) + 21 内置工具 + Hook 新增能力: - 21 个内置工具 (文件/代码/Shell/Git/Web/Notebook/Task) - ToolGateway 安全网关 - Hook 全局审计 - ContextManager 上下文压缩 - TerminalUI 流式显示 """ def __init__(self, model: str = "sonnet", use_api: bool = False, backend: str = "", verbose: bool = False): # Brain (LLM) if backend == BACKEND_OLLAMA: self.brain = OllamaLam( "lambda_v2", prompt=SYSTEM_PROMPT, model=model, temperature=0.3, max_tokens=4096, ) elif use_api or backend == BACKEND_API: self.brain = Lam( "lambda_v2", prompt=SYSTEM_PROMPT, model=model, temperature=0.3, max_tokens=4096, ) else: self.brain = ClaudeLam( "lambda_v2", prompt=SYSTEM_PROMPT, model=model, max_tokens=4096, inject_override=False, ) self.ctx = Context() self.conversation_history: list[dict] = [] self.max_history = 30 # v2 新增 self.hooks = HookRegistry() self.context_manager = ContextManager(max_tokens=100000, keep_recent=5) self.ui = TerminalUI(verbose=verbose) # 注册审计 hook self._total_tokens = 0 self._total_tool_calls = 0 self.hooks.register("post_tool", self._audit_tool) def _audit_tool(self, **kwargs): """全局工具调用审计""" self._total_tool_calls += 1 def _build_input(self, user_msg: str, observations: list[str] | None = None) -> str: """构建完整输入(历史 + 消息 + 工具结果)""" parts = [] # 对话历史 if self.conversation_history: parts.append("=== 对话历史 ===") for entry in self.conversation_history[-self.max_history:]: content = entry["content"] if len(content) > 300: content = content[:300] + "..." parts.append(f"[{entry['role']}] {content}") parts.append("=== 历史结束 ===\n") # 工具观察 (sliding window: summarize old, keep recent in full) if observations: _KEEP_RECENT = 5 parts.append("=== 工具执行结果 ===") if len(observations) > _KEEP_RECENT: # Summarize old observations parts.append("[之前的工具调用摘要]") for obs in observations[:-_KEEP_RECENT]: # Extract tool name and first line of result lines = obs.split("\n") tool_line = lines[0] if lines else "" result_line = lines[1][:120] if len(lines) > 1 else "" parts.append(f" {tool_line}: {result_line}...") parts.append("") # Recent observations in full for obs in observations[-_KEEP_RECENT:]: parts.append(obs) else: for obs in observations: parts.append(obs) parts.append("=== 结果结束 ===\n") parts.append( "请基于以上结果决定下一步。直接输出一个JSON代码块调用下一个工具。\n" "注意:工具名是 ReadFile(不是Read)、WriteFile(不是Write)、EditFile(不是Edit)、" "Bash(不是bash或shell)、ListFiles(不是Glob或ls)。" ) else: parts.append(f"[用户] {user_msg}") state = "\n".join(parts) # v2: 上下文压缩 if self.context_manager.should_compact(state): state = self.context_manager.compact(state) return state def chat(self, user_msg: str) -> str: """ 与 lambda 对话。完整 ReAct 循环 (v2 流式交互)。 """ self.conversation_history.append({"role": "用户", "content": user_msg}) observations = [] max_steps = 30 final_response = "" import time as _time for step in range(max_steps): # 显示当前步骤状态 _step_label = f"[步骤 {step + 1}/{max_steps}]" # 构建输入 if step == 0: llm_input = self._build_input(user_msg) else: llm_input = self._build_input(user_msg, observations) # 显示思考状态(如果不是流式模式,显示 spinner) is_streaming = hasattr(self.brain, 'stream') and self.brain.stream if not is_streaming: print(f" ⏳ {_step_label} 思考中...", end="\r", flush=True) # Brain 思考 (β-规约) — 流式模式下 ClaudeLam 会实时输出 t0 = _time.time() llm_output = self.brain.apply(llm_input, self.ctx) think_ms = (_time.time() - t0) * 1000 if not is_streaming: # 清除 spinner print(f" ✅ {_step_label} 思考完成 ({think_ms:.0f}ms) ") # 解析并执行 (v2: 支持内置工具 + hooks) result, is_done, tool_name = parse_and_execute(str(llm_output), self.hooks) if is_done: text_part = re.sub(r'```json.*?```', '', str(llm_output), flags=re.DOTALL).strip() final_response = text_part + ("\n" + result if result else "") print(f" ✅ {_step_label} 任务完成") break if result != str(llm_output): # 工具被执行 text_part = re.sub(r'```json.*?```', '', str(llm_output), flags=re.DOTALL).strip() # 显示思考内容 if text_part and not is_streaming: print(f" 💭 {text_part[:200]}") # 显示工具调用 print(f" 🔧 {_step_label} 调用工具: {tool_name or '?'}") # 工具结果预览 result_preview = result[:300].replace("\n", "\n ") print(f" 结果: {result_preview}") print() observations.append(f"[工具: {tool_name}]\n[执行结果] {result}") else: # 纯文本回复,没有工具调用 output_str = str(llm_output) completion_signals = [ "任务完成", "已完成", "完成了", "已保存", "已写入", "已创建", "总结", "结论", "综上", "done", "finished", "completed", ] looks_final = any( sig.lower() in output_str.lower() for sig in completion_signals ) used_tools = len(observations) > 0 can_continue = step < max_steps - 2 if used_tools and not looks_final and can_continue: # 已调过工具但只是中间叙述 → 推动继续 observations.append( f"[LLM回复] {output_str[:500]}\n" f"[系统提醒] 你还没给出最终答案。请选择:\n" f" (a) 继续调用工具推进任务;\n" f" (b) 如果信息已足够,用 done 工具把完整结果写在 answer 里," f"例如:\n" f' {{"tool": "done", "args": {{"answer": "用户最近3封未读邮件:\\n1. ...\\n2. ...\\n3. ..."}}}}\n' f" 注意: answer 必须是真正的摘要/结果内容,不要写占位符。" ) print(f" 💭 {_step_label} 中间叙述,推动继续...") else: # 真的是最终回复 (首轮纯问答 / 含完成信号 / 步数耗尽) final_response = output_str break else: final_response = f"(达到最大步数 {max_steps})\n最后的观察:\n" + ( observations[-1] if observations else "无" ) self.conversation_history.append({ "role": "lambda", "content": final_response[:500], }) return final_response def print_trace(self): """打印 β-规约追踪 (v2: 使用 TerminalUI)""" if self.ctx.trace: self.ui.trace(self.ctx.trace) else: print(" (无追踪记录)") def print_stats(self): """打印会话统计""" print(f"\n📊 会话统计:") print(f" β-规约步数: {len(self.ctx.trace)}") print(f" 工具调用数: {self._total_tool_calls}") print(f" 对话轮数: {len(self.conversation_history) // 2}")