🏭 模块1:工业级轻量化综合 AI Agent 项目

综合实战 ⏱ 预计 4-5 小时

项目概述:AetherAgent

AetherAgent 是本教程的"毕业设计"——整合五大阶段所有核心能力:

能力来源实现方式
记忆系统第二阶段模块1 + 第四阶段模块2短期内存 + Chroma 向量长期记忆
工具调用第二阶段模块2ToolRegistry 注册中心
任务规划第二阶段模块3TaskPlanner + 动态调整
Prompt 工程第二阶段模块4分层 Prompt 模板
RAG 检索第三阶段项目1向量化 + 语义检索
多轮对话第三阶段项目3意图识别 + 状态管理
自我反思第四阶段模块1Self-Refine 循环
并行执行第四阶段模块2ThreadPoolExecutor
性能监控第四阶段模块3TokenMonitor + 缓存

系统架构

┌────────────────────────────────────────────────────────────────┐ │ AetherAgent 五层架构 │ │ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ │ 感知层: 输入解析 → 意图识别 → 信息提取 → 上下文组装 │ │ │ └──────────────────────────────────────────────────────────┘ │ │ │ │ │ ▼ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ │ 规划层: TaskPlanner → 任务拆解 → 依赖分析 → 动态调整 │ │ │ └──────────────────────────────────────────────────────────┘ │ │ │ │ │ ▼ │ │ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────┐ │ │ │ 工具层 │ │ 记忆层 │ │ 反思层 │ │ │ │ ToolRegistry │ │ ShortTerm │ │ SelfRefine │ │ │ │ 并行执行器 │ │ Chroma长期 │ │ 事实核查 │ │ │ └─────────────┘ └─────────────┘ └─────────────────────┘ │ │ │ │ │ │ │ └─────────────────┼──────────────────┘ │ │ ▼ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ │ 输出层: 结果汇总 → 质量验证 → 格式化 → 用户呈现 │ │ │ └──────────────────────────────────────────────────────────┘ │ │ │ │ ┌──────────────────────────────────────────────────────────┐ │ │ │ 基础设施: Token监控 | 日志系统 | 健康检查 | 容错重试 │ │ │ └──────────────────────────────────────────────────────────┘ │ └────────────────────────────────────────────────────────────────┘

完整代码(精简核心版)

#!/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()

项目拓展方向

✏️ 毕业设计

  1. 完整运行:运行 AetherAgent,完成至少10轮不同类型的任务交互。
  2. 功能增强:选择上述5个拓展方向之一,动手实现。
  3. 性能报告:使用 TokenMonitor 生成完整的 Token 消耗报告。
  4. 架构文档:为自己开发的 AetherAgent 撰写一份架构文档(含设计决策、技术选型理由)。
← 阶段首页 下一模块:知识图谱 →