综合实战 ⏱ 预计 4-5 小时
AetherAgent 是本教程的"毕业设计"——整合五大阶段所有核心能力:
| 能力 | 来源 | 实现方式 |
|---|---|---|
| 记忆系统 | 第二阶段模块1 + 第四阶段模块2 | 短期内存 + Chroma 向量长期记忆 |
| 工具调用 | 第二阶段模块2 | ToolRegistry 注册中心 |
| 任务规划 | 第二阶段模块3 | TaskPlanner + 动态调整 |
| Prompt 工程 | 第二阶段模块4 | 分层 Prompt 模板 |
| RAG 检索 | 第三阶段项目1 | 向量化 + 语义检索 |
| 多轮对话 | 第三阶段项目3 | 意图识别 + 状态管理 |
| 自我反思 | 第四阶段模块1 | Self-Refine 循环 |
| 并行执行 | 第四阶段模块2 | ThreadPoolExecutor |
| 性能监控 | 第四阶段模块3 | TokenMonitor + 缓存 |
#!/usr/bin/env python3
"""
AetherAgent — 工业级轻量化综合 AI Agent
==========================================
整合:记忆 + 工具 + 规划 + RAG + 反思 + 并行 + 监控
依赖: pip install anthropic python-dotenv chromadb sentence-transformers
运行: python3 aether_agent.py
"""
import json, os, time, hashlib
from datetime import datetime
from pathlib import Path
from typing import Any, Callable
from concurrent.futures import ThreadPoolExecutor, as_completed
from dotenv import load_dotenv
from anthropic import Anthropic
load_dotenv()
# ================================================================
# 配置
# ================================================================
class Config:
MODEL = "claude-sonnet-4-6"
MAX_RETRIES = 3
PARALLEL_WORKERS = 4
MEMORY_MAX_TURNS = 10
CACHE_SIZE = 100
REFINE_ITERATIONS = 2
# ================================================================
# Token 监控器
# ================================================================
class TokenMonitor:
def __init__(self):
self.total_in = 0; self.total_out = 0; self.calls = 0
def record(self, inp: int, out: int):
self.total_in += inp; self.total_out += out; self.calls += 1
def report(self):
total = self.total_in + self.total_out
return f"📊 Token: {total:,} ({self.calls}次调用) | 预估费用: ${self.total_in/1e6*3 + self.total_out/1e6*15:.4f}"
# ================================================================
# 工具注册中心
# ================================================================
class ToolRegistry:
def __init__(self):
self._tools: dict[str, dict] = {}
def register(self, name: str, desc: str, params: dict = None):
def deco(fn):
self._tools[name] = {"fn": fn, "desc": desc, "params": params or {}}
return fn
return deco
def execute(self, name: str, args: dict) -> str:
t = self._tools.get(name)
if not t: return f"❌ 工具不存在: {name}"
try:
first_arg = list(args.values())[0] if args else ""
return t["fn"](first_arg)
except Exception as e:
return f"❌ 工具执行失败: {e}"
def describe(self) -> str:
lines = ["## 可用工具"]
for name, t in self._tools.items():
params = ", ".join(f"{k}:{v}" for k,v in t["params"].items())
lines.append(f"- **{name}**({params}): {t['desc']}")
return "\n".join(lines)
def list_names(self) -> list[str]:
return list(self._tools.keys())
# 全局注册中心
tools = ToolRegistry()
@tools.register("get_time", "获取当前日期和时间")
def _(a=None): return datetime.now().strftime("%Y-%m-%d %H:%M:%S %A")
@tools.register("calculator", "安全计算数学表达式", {"expression": "数学表达式"})
def _(expr: str):
allowed = set("0123456789+-*/().% ")
if not all(c in allowed for c in expr): return "❌ 非法字符"
try: return str(eval(expr, {"__builtins__":{}}, {}))
except Exception as e: return f"❌ {e}"
@tools.register("read_file", "读取文件内容", {"path": "文件路径"})
def _(path: str):
p = Path(path)
if not p.exists(): return f"❌ 文件不存在: {path}"
try: return p.read_text("utf-8")[:3000]
except: return f"❌ 读取失败"
@tools.register("write_file", "写入文件", {"path_content": "路径|内容"})
def _(arg: str):
parts = arg.split("|", 1)
if len(parts) < 2: return "❌ 格式: 路径|内容"
Path(parts[0]).write_text(parts[1], encoding="utf-8")
return f"✅ 已写入 {parts[0]}"
@tools.register("web_search", "搜索网页内容(模拟)", {"query": "搜索词"})
def _(query: str):
return f"🔍 关于'{query}'的搜索结果: [开发中,请使用具体知识库]"
# ================================================================
# 记忆系统
# ================================================================
class MemorySystem:
def __init__(self):
self.short_term: list[dict] = []
self.max_turns = Config.MEMORY_MAX_TURNS
# 长期记忆用简单 JSON 文件(Chroma 版本见第四阶段模块2)
self.long_term_file = Path("./aether_memory.json")
self.long_term: dict = self._load_long()
def _load_long(self):
if self.long_term_file.exists():
return json.loads(self.long_term_file.read_text("utf-8"))
return {"facts": [], "preferences": {}, "stats": {}}
def remember_short(self, role: str, content: str):
self.short_term.append({"role": role, "content": content, "time": datetime.now().isoformat()})
if len(self.short_term) > self.max_turns * 2:
self.short_term = self.short_term[-(self.max_turns * 2):]
def get_context(self) -> list[dict]:
return [{"role": t["role"], "content": t["content"]} for t in self.short_term]
def remember_fact(self, fact: str):
self.long_term["facts"].append({"content": fact, "time": datetime.now().isoformat()})
self._save()
def recall_facts(self, keyword: str = "") -> list:
if not keyword: return self.long_term["facts"]
return [f for f in self.long_term["facts"] if keyword.lower() in f["content"].lower()]
def _save(self):
self.long_term_file.write_text(json.dumps(self.long_term, ensure_ascii=False, indent=2))
# ================================================================
# 自我反思器
# ================================================================
class SelfRefiner:
def __init__(self, client: Anthropic):
self.client = client
def refine(self, content: str, task: str, max_iter: int = None) -> str:
max_iter = max_iter or Config.REFINE_ITERATIONS
for i in range(max_iter):
feedback = self._evaluate(content, task)
if feedback.get("score", 0) >= 8:
break
content = self._improve(content, feedback, task)
return content
def _evaluate(self, content: str, task: str) -> dict:
resp = self.client.messages.create(
model=Config.MODEL, max_tokens=300,
messages=[{"role":"user","content":
f"评价以下内容(1-10分):\n任务:{task}\n内容:{content[:1500]}\n输出JSON: {{'score':int,'issues':[]}}"}],
temperature=0.1)
try:
text = resp.content[0].text.strip()
if "```" in text: text = text.split("```")[1].split("```")[0]
return json.loads(text)
except: return {"score": 8, "issues": []}
def _improve(self, content: str, feedback: dict, task: str) -> str:
issues = "\n".join(f"- {i}" for i in feedback.get("issues", []))
resp = self.client.messages.create(
model=Config.MODEL, max_tokens=1000,
messages=[{"role":"user","content":
f"改进以下内容。任务:{task}\n问题:{issues}\n内容:{content[:1500]}"}],
temperature=0.3)
return resp.content[0].text.strip()
# ================================================================
# AetherAgent 主类
# ================================================================
class AetherAgent:
"""工业级综合 AI Agent"""
def __init__(self, name: str = "Aether"):
self.name = name
self.client = Anthropic(api_key=os.getenv("ANTHROPIC_API_KEY"))
self.memory = MemorySystem()
self.refiner = SelfRefiner(self.client)
self.monitor = TokenMonitor()
self.cache: dict[str, dict] = {}
print(f"✨ {self.name}Agent 初始化完成")
print(f"🧠 模型: {Config.MODEL}")
print(f"🔧 工具: {tools.list_names()}")
print(f"💾 记忆: {len(self.memory.long_term['facts'])} 条长期记忆")
def _call_llm(self, system: str, user: str, max_tokens: int = 1000,
temperature: float = 0.3, use_cache: bool = True) -> str:
"""统一的 LLM 调用(带缓存和监控)"""
# 缓存检查
cache_key = hashlib.md5(f"{system}:{user}:{temperature}".encode()).hexdigest()
if use_cache and cache_key in self.cache:
entry = self.cache[cache_key]
if time.time() - entry["time"] < 300: # 5分钟
return entry["response"]
# API 调用
for attempt in range(Config.MAX_RETRIES):
try:
resp = self.client.messages.create(
model=Config.MODEL, max_tokens=max_tokens,
system=system, messages=[{"role": "user", "content": user}],
temperature=temperature)
# 记录 Token(Anthropic SDK 响应中包含 usage 信息)
if hasattr(resp, 'usage'):
self.monitor.record(resp.usage.input_tokens, resp.usage.output_tokens)
result = resp.content[0].text.strip()
# 缓存
self.cache[cache_key] = {"response": result, "time": time.time()}
if len(self.cache) > Config.CACHE_SIZE:
oldest = min(self.cache.items(), key=lambda x: x[1]["time"])
del self.cache[oldest[0]]
return result
except Exception as e:
if attempt == Config.MAX_RETRIES - 1:
return f"❌ LLM调用失败(重试{Config.MAX_RETRIES}次): {e}"
time.sleep(1 * (attempt + 1))
return "❌ 未知错误"
def run(self, user_input: str) -> str:
"""主运行入口"""
self.memory.remember_short("user", user_input)
# 构建系统提示词(整合全部能力)
system = f"""你是 {self.name},一个全能 AI Agent。
## 能力
- 理解并执行各种任务
- 调用工具完成实际操作
- 记住对话历史和重要信息
- 自我反思和改进输出质量
{tools.describe()}
## 行为规范
- 优先使用工具完成任务
- 工具调用格式: {{"tool":"工具名","args":{{"参数":"值"}}}}
- 输出高质量的结果,不确定时标明
- 使用中文回复
## 长期记忆上下文
{self._format_long_memory()}
"""
# 调用 LLM
response = self._call_llm(system, user_input)
# 处理工具调用
result = self._process_tool_calls(response)
# 自我反思优化(对重要输出)
if len(result) > 200 and "❌" not in result:
result = self.refiner.refine(result, user_input)
self.memory.remember_short("assistant", result)
return result
def _process_tool_calls(self, response: str) -> str:
"""解析并执行工具调用"""
try:
json_str = response
for marker in ["```json", "```"]:
if marker in response:
json_str = response.split(marker)[1].split("```")[0]
break
parsed = json.loads(json_str.strip())
if "tool" in parsed:
name = parsed["tool"]
args = parsed.get("args", {})
print(f"🔧 调用: {name}({args})")
return tools.execute(name, args)
except (json.JSONDecodeError, KeyError):
pass
return response
def _format_long_memory(self) -> str:
facts = self.memory.recall_facts()
if not facts: return "(暂无长期记忆)"
return "\n".join(f"- {f['content']}" for f in facts[-5:])
def run_parallel(self, tasks: list[str]) -> list[str]:
"""并行执行多个任务"""
results = [None] * len(tasks)
with ThreadPoolExecutor(max_workers=Config.PARALLEL_WORKERS) as executor:
futures = {executor.submit(self.run, task): i for i, task in enumerate(tasks)}
for future in as_completed(futures):
idx = futures[future]
try: results[idx] = future.result(timeout=60)
except Exception as e: results[idx] = f"❌ {e}"
return results
def remember(self, fact: str):
"""手动记忆事实"""
self.memory.remember_fact(fact)
return f"✅ 已记住: {fact}"
def status(self):
print(f"\n{'='*50}")
print(f"✨ {self.name}Agent 状态")
print(f"{'='*50}")
print(f"短期记忆: {len(self.memory.short_term)//2} 轮对话")
print(f"长期记忆: {len(self.memory.long_term['facts'])} 条事实")
print(f"缓存条目: {len(self.cache)}")
print(self.monitor.report())
# ================================================================
# 交互入口
# ================================================================
def main():
print("=" * 55)
print("✨ AetherAgent — 工业级综合 AI Agent")
print("=" * 55)
print("命令: 直接输入任务 | remember <事实> | parallel <任务1>;<任务2>")
print(" status | quit")
print("=" * 55)
if not os.getenv("ANTHROPIC_API_KEY"):
print("❌ 请配置 ANTHROPIC_API_KEY")
return
agent = AetherAgent()
while True:
try:
ui = input("\n👤 你: ").strip()
if not ui: continue
if ui.lower() in ("quit","exit","q"): break
if ui.lower().startswith("remember "):
print(f"🤖 Agent: {agent.remember(ui[9:])}")
elif ui.lower().startswith("parallel "):
tasks = [t.strip() for t in ui[9:].split(";") if t.strip()]
print(f"⚡ 并行执行 {len(tasks)} 个任务...")
results = agent.run_parallel(tasks)
for i, (t, r) in enumerate(zip(tasks, results), 1):
print(f"\n 任务{i}: {t[:50]}")
print(f" 结果: {r[:200]}")
elif ui.lower() == "status":
agent.status()
else:
result = agent.run(ui)
print(f"\n🤖 {agent.name}:\n{result}")
except KeyboardInterrupt: break
agent.status()
print(f"\n👋 {agent.name}Agent 已退出。")
if __name__ == "__main__":
main()