| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525 |
- """
- 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}")
|