目 录CONTENT

文章目录

Agent 可观测性:如何调试一个"不可预测"的系统

PySuper
2025-11-22 / 0 评论 / 0 点赞 / 1 阅读 / 0 字
温馨提示:
本文最后更新于2026-05-20,若内容或图片失效,请留言反馈。 所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

Agent 系统的行为往往是非确定性的、多步骤的、涉及复杂工具调用的,这让传统的调试方法变得力不从心。本文探讨如何构建 Agent 系统的可观测性体系,让"不可预测"变得可控可追踪。

一、Agent 系统的可观测性挑战

1.1 为什么 Agent 难以调试?

传统软件的行为是确定性的:相同的输入经过相同的处理逻辑,必然产生相同的输出。但 Agent 系统打破了这一确定性神话。

┌─────────────────────────────────────────────────────────────────────────────┐
│                    Agent 系统的不确定性来源                                   │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│  ┌─────────────────────────────────────────────────────────────────────┐   │
│  │                       1. LLM 的固有不确定性                          │   │
│  │                                                                       │   │
│  │   同一 Prompt → 不同模型可能产生不同回复                               │   │
│  │   相同输入 → 即使同模型也可能有随机性(temperature > 0)              │   │
│  │   上下文变化 → 微小变化可能导致截然不同的决策                          │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                    ↓                                        │
│  ┌─────────────────────────────────────────────────────────────────────┐   │
│  │                       2. 多步推理链路                                 │   │
│  │                                                                       │   │
│  │   Agent 执行往往需要多轮交互:                                        │   │
│  │   用户请求 → 意图识别 → 工具选择 → 参数生成 → 工具执行 → 结果处理    │   │
│  │   → 可能触发新的工具调用 → 最终回复                                   │   │
│  │                                                                       │   │
│  │   任何一步出错都可能导致最终结果偏离预期                               │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                    ↓                                        │
│  ┌─────────────────────────────────────────────────────────────────────┐   │
│  │                       3. 外部工具的不确定性                            │   │
│  │                                                                       │   │
│  │   API 响应可能超时、返回错误数据                                       │   │
│  │   网络波动导致间歇性失败                                              │   │
│  │   第三方服务行为可能变化                                              │   │
│  │   工具返回格式可能与预期不符                                           │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                    ↓                                        │
│  ┌─────────────────────────────────────────────────────────────────────┐   │
│  │                       4. 组合爆炸                                     │   │
│  │                                                                       │   │
│  │   10 个工具 × 10 种参数组合 × 5 种状态 = 500 种可能路径               │   │
│  │   难以穷尽测试所有场景                                                │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

1.2 可观测性的三大支柱

可观测性(Observability)≠ 监控(Monitoring)。传统监控告诉你"系统坏了",可观测性告诉你"为什么坏了"。

维度

说明

Agent 场景应用

追踪(Traces)

请求在系统中的完整执行路径

每个 Agent 决策、工具调用的完整链路

指标(Metrics)

可聚合的数值数据

决策正确率、Token 消耗、工具调用成功率

日志(Logs)

离散的、带时间戳的事件

错误详情、调试信息、审计记录

┌─────────────────────────────────────────────────────────────────────────────┐
│                         可观测性三大支柱                                     │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│                              ┌───────────────┐                              │
│                              │    Traces     │                              │
│                              │   请求链路追踪  │                              │
│                              └───────────────┘                              │
│                                       │                                      │
│                                       │                                      │
│         ┌───────────────┐     ┌───────────────┐     ┌───────────────┐      │
│         │   Metrics     │ ←→ │   Traces      │ ←→ │     Logs       │      │
│         │  可聚合指标    │     │   完整链路    │     │  离散事件     │      │
│         └───────────────┘     └───────────────┘     └───────────────┘      │
│              ↑                      ↑                      ↑               │
│              │                      │                      │               │
│         ┌────────────┐        ┌────────────┐         ┌────────────┐         │
│         │  Prometheus │        │   Jaeger    │         │   ELK      │         │
│         │   Grafana   │        │  Zipkin     │         │  Loki      │         │
│         └────────────┘        └────────────┘         └────────────┘         │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

二、LangSmith 与 LangFuse 集成

2.1 LangSmith 集成

LangSmith 是 LangChain 官方推出的可观测性平台,深度集成 LangChain。

import os
from langchain_openai import ChatOpenAI
from langchain.prompts import ChatPromptTemplate
from langchain.schema import StrOutputParser
from langchain.callbacks.tracers.langsmith import LangSmithTracer
from langchain.callbacks.manager import CallbackManager
from langchain.callbacks.base import BaseCallbackHandler

# ============== LangSmith 配置 ==============

# 设置环境变量
os.environ["LANGCHAIN_TRACING_V2"] = "true"
os.environ["LANGCHAIN_API_KEY"] = "your-langsmith-api-key"
os.environ["LANGCHAIN_PROJECT"] = "agent-debugging"  # 项目名称

# 可选:添加数据集用于评估
os.environ["LANGCHAIN_DATASET"] = "agent-evaluation-dataset"


class AgentCallbackHandler(BaseCallbackHandler):
    """
    自定义回调处理器:记录 Agent 特定事件
    
    可以在这里:
    1. 记录自定义指标
    2. 添加业务相关的上下文
    3. 实现自定义告警逻辑
    """
    
    def __init__(self):
        self.agent_decisions = []
        self.tool_calls = []
        self.token_usage = {"prompt": 0, "completion": 0, "total": 0}
    
    def on_llm_start(self, serialized, prompts, **kwargs):
        """LLM 开始调用"""
        print(f"[Callback] LLM 开始处理,Prompt tokens: {len(prompts)}")
    
    def on_llm_end(self, response, **kwargs):
        """LLM 调用结束"""
        if hasattr(response, "usage_metadata"):
            usage = response.usage_metadata
            self.token_usage["prompt"] += usage.get("input_tokens", 0)
            self.token_usage["completion"] += usage.get("output_tokens", 0)
            self.token_usage["total"] += usage.get("total_tokens", 0)
    
    def on_chain_start(self, serialized, inputs, **kwargs):
        """Chain 开始执行"""
        chain_name = serialized.get("name", "unknown")
        print(f"[Callback] Chain 开始: {chain_name}")
    
    def on_chain_end(self, outputs, **kwargs):
        """Chain 执行结束"""
        print(f"[Callback] Chain 结束,输出: {str(outputs)[:200]}")
    
    def on_tool_start(self, serialized, input_str, **kwargs):
        """工具开始执行"""
        tool_name = serialized.get("name", "unknown")
        print(f"[Callback] 工具开始: {tool_name}, 输入: {input_str[:100]}")
        self.tool_calls.append({
            "tool": tool_name,
            "input": input_str,
            "status": "started"
        })
    
    def on_tool_end(self, output, **kwargs):
        """工具执行结束"""
        self.tool_calls[-1]["status"] = "completed"
        self.tool_calls[-1]["output"] = str(output)[:200]
        print(f"[Callback] 工具结束: {self.tool_calls[-1]['tool']}, 输出: {str(output)[:100]}")
    
    def on_tool_error(self, error, **kwargs):
        """工具执行出错"""
        self.tool_calls[-1]["status"] = "error"
        self.tool_calls[-1]["error"] = str(error)
        print(f"[Callback] 工具错误: {error}")


def create_langsmith_callback_manager(project_name: str = "agent-debugging"):
    """创建 LangSmith 回调管理器"""
    
    # 创建 LangSmith tracer
    tracer = LangSmithTracer(
        project_name=project_name,
        # 可选:添加 tags 用于过滤
        tags=["production", "v1.0"],
        # 可选:添加 metadata
        metadata={
            "environment": os.getenv("ENV", "development"),
            "version": os.getenv("APP_VERSION", "1.0.0")
        }
    )
    
    # 创建自定义 handler
    custom_handler = AgentCallbackHandler()
    
    # 创建 callback manager
    callback_manager = CallbackManager(
        handlers=[tracer, custom_handler]
    )
    
    return callback_manager, custom_handler


# 使用示例
async def demo_langsmith():
    """LangSmith 集成示例"""
    
    # 创建带回调的 LLM
    callback_manager, handler = create_langsmith_callback_manager()
    
    llm = ChatOpenAI(
        model="gpt-4o",
        temperature=0,
        callback_manager=callback_manager
    )
    
    # 创建简单的 chain
    prompt = ChatPromptTemplate.from_template(
        "用一句话解释 {topic},用中文回答"
    )
    
    chain = prompt | llm | StrOutputParser()
    
    # 执行
    result = await chain.ainvoke({"topic": "量子计算"})
    
    print(f"\n最终结果: {result}")
    print(f"Token 消耗: {handler.token_usage}")
    print(f"工具调用: {handler.tool_calls}")

2.2 LangFuse 集成

LangFuse 是开源的可观测性方案,支持自部署,数据完全私有。

"""
LangFuse 集成:开源可观测性方案
"""

from langfuse import Langfuse
from langfuse.callback import CallbackHandler
from langchain_openai import ChatOpenAI
from langchain.schema import HumanMessage
import os

# ============== LangFuse 配置 ==============

# 方式1:使用云服务
# os.environ["LANGFUSE_PUBLIC_KEY"] = "pk-xxx"
# os.environ["LANGFUSE_SECRET_KEY"] = "sk-xxx"
# langfuse = Langfuse()

# 方式2:自部署
langfuse = Langfuse(
    public_key=os.getenv("LANGFUSE_PUBLIC_KEY", ""),
    secret_key=os.getenv("LANGFUSE_SECRET_KEY", ""),
    host=os.getenv("LANGFUSE_HOST", "http://localhost:3000")  # 自部署地址
)


def create_langfuse_callback():
    """创建 LangFuse 回调处理器"""
    
    return CallbackHandler(
        langfuse=langfuse,
        # 会话相关
        session_id="session-123",  # 用于追踪同一会话
        user_id="user-456",       # 用于追踪用户
        
        # 追踪元数据
        tags=["agent", "production"],
        metadata={
            "agent_version": "1.0.0",
            "environment": "production"
        }
    )


# ============== 追踪 Agent 执行 ==============

class AgentTracer:
    """
    Agent 追踪器:封装 LangFuse 功能
    
    功能:
    1. 追踪每个 Agent 决策
    2. 记录工具调用链路
    3. 收集性能指标
    """
    
    def __init__(self, agent_name: str):
        self.agent_name = agent_name
        self.langfuse = langfuse
        self.trace = None
        self.spans = []
    
    def start_trace(self, name: str, input_data: dict) -> str:
        """
        开始追踪
        
        Returns:
            trace_id: 用于关联后续操作
        """
        self.trace = self.langfuse.trace(
            name=f"{self.agent_name}.{name}",
            input=input_data,
            metadata={
                "agent": self.agent_name,
                "type": "agent_execution"
            }
        )
        return self.trace.id
    
    def create_span(self, name: str, metadata: dict = None):
        """创建子跨度"""
        
        class SpanContext:
            def __init__(self, tracer, name, metadata):
                self.tracer = tracer
                self.name = name
                self.metadata = metadata or {}
                self.start_time = None
                self.end_time = None
                self.observations = []
            
            def __enter__(self):
                import time
                self.start_time = time.time()
                return self
            
            def __exit__(self, exc_type, exc_val, exc_tb):
                import time
                self.end_time = time.time()
                duration = self.end_time - self.start_time
                
                # 记录到 LangFuse
                if self.tracer.trace:
                    self.tracer.trace.span(
                        name=self.name,
                        input=self.metadata.get("input"),
                        output=self.metadata.get("output"),
                        metadata={
                            **self.metadata,
                            "duration_seconds": duration,
                            "status": "success" if exc_type is None else "error"
                        }
                    )
                
                return False
        
        return SpanContext(self, name, metadata)
    
    def record_decision(
        self,
        step: str,
        reasoning: str,
        action: str,
        confidence: float
    ) -> None:
        """记录 Agent 决策点"""
        
        if not self.trace:
            return
        
        self.trace.generation(
            name=f"decision.{step}",
            input={"step": step, "reasoning": reasoning},
            output={
                "action": action,
                "confidence": confidence
            },
            metadata={
                "type": "agent_decision",
                "step": step
            }
        )
    
    def record_tool_call(
        self,
        tool_name: str,
        input_args: dict,
        output: Any,
        error: str = None
    ) -> None:
        """记录工具调用"""
        
        if not self.trace:
            return
        
        self.trace.span(
            name=f"tool.{tool_name}",
            input=input_args,
            output=output if not error else None,
            metadata={
                "type": "tool_call",
                "tool_name": tool_name,
                "error": error
            }
        )
    
    def end_trace(self, output: Any, metadata: dict = None) -> None:
        """结束追踪"""
        
        if not self.trace:
            return
        
        self.trace.update(
            output=output,
            metadata=metadata or {}
        )
        self.trace = None


# 使用示例
async def demo_langfuse():
    """LangFuse 使用示例"""
    
    tracer = AgentTracer("weather_assistant")
    
    # 开始追踪
    trace_id = tracer.start_trace(
        "get_weather_and_remind",
        {"user_input": "明天北京天气怎么样?"}
    )
    
    # 记录决策
    tracer.record_decision(
        step="intent_detection",
        reasoning="用户询问天气,需要调用天气查询工具",
        action="call_get_weather",
        confidence=0.95
    )
    
    # 记录工具调用
    with tracer.create_span("get_weather", {"city": "北京", "date": "明天"}):
        weather_result = {"temp": 15, "condition": "晴"}
        tracer.record_tool_call(
            "get_weather",
            {"city": "北京", "date": "明天"},
            weather_result
        )
    
    # 结束追踪
    tracer.end_trace(
        output={"answer": "北京明天天气晴朗,温度15°C"},
        metadata={"tokens_used": 500}
    )

三、Trace 设计:每步决策可追踪

3.1 完整 Trace 架构

┌─────────────────────────────────────────────────────────────────────────────┐
│                         Agent Trace 完整架构                                 │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│  用户输入: "帮我查下张三的订单,如果已发货就设置个提醒"                         │
│                                                                             │
│  ┌─────────────────────────────────────────────────────────────────────┐   │
│  │                        Trace Root                                  │   │
│  │                      trace_id: abc123                              │   │
│  │                      start_time: 2024-01-15 10:00:00               │   │
│  │                      duration: 2.5s                                │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                    │                                        │
│         ┌──────────────────────────┼──────────────────────────┐              │
│         ↓                          ↓                          ↓              │
│  ┌─────────────┐           ┌─────────────┐           ┌─────────────┐        │
│  │   Span 1    │           │   Span 2    │           │   Span 3    │        │
│  │  意图识别   │           │  工具调用   │           │  工具调用   │        │
│  │ intent.detect│          │ query.order │           │ set.reminder│        │
│  └─────────────┘           └─────────────┘           └─────────────┘        │
│         │                          │                          │            │
│         ↓                          ↓                          ↓            │
│  ┌─────────────┐           ┌─────────────┐           ┌─────────────┐        │
│  │  Generation  │           │    Event    │           │    Event    │        │
│  │  决策理由    │           │  工具输入   │           │  工具输出   │        │
│  │ 置信度0.92   │           │ order_id    │           │ 已发货     │        │
│  └─────────────┘           └─────────────┘           └─────────────┘        │
│                                                                             │
│  嵌套关系:                                                                 │
│  Trace                                                                        │
│    └── Span: 意图识别                                                         │
│          └── Generation: 推理过程                                            │
│    └── Span: 订单查询                                                        │
│          └── Event: API调用详情                                              │
│    └── Span: 设置提醒                                                        │
│          └── Event: 提醒创建成功                                             │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

3.2 追踪实现代码

"""
Agent 追踪系统完整实现
支持:LangSmith、LangFuse、自建追踪
"""

import json
import time
import uuid
from typing import Any, Dict, List, Optional, Callable
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from contextvars import ContextVar
from collections import defaultdict

# 追踪上下文
_current_trace: ContextVar[Optional["TraceContext"]] = ContextVar(
    "current_trace", 
    default=None
)


class SpanKind(Enum):
    """跨度类型"""
    INTERNAL = "internal"           # 内部操作
    LLM = "llm"                     # LLM 调用
    TOOL = "tool"                   # 工具调用
    CHAIN = "chain"                 # Chain 执行
    RETRIEVER = "retriever"         # 检索操作


class SpanStatus(Enum):
    """跨度状态"""
    OK = "ok"
    ERROR = "error"
    UNSET = "unset"


@dataclass
class Span:
    """追踪跨度"""
    span_id: str
    name: str
    kind: SpanKind
    trace_id: str
    parent_span_id: Optional[str]
    
    start_time: datetime = field(default_factory=datetime.now)
    end_time: Optional[datetime] = None
    
    input_data: Any = None
    output_data: Any = None
    error: Optional[str] = None
    
    metadata: Dict[str, Any] = field(default_factory=dict)
    attributes: Dict[str, Any] = field(default_factory=dict)
    
    # 关联数据
    tags: List[str] = field(default_factory=list)
    
    @property
    def duration_ms(self) -> float:
        """计算持续时间(毫秒)"""
        if self.end_time:
            return (self.end_time - self.start_time).total_seconds() * 1000
        return 0.0
    
    def to_dict(self) -> Dict[str, Any]:
        """转换为字典"""
        return {
            "span_id": self.span_id,
            "name": self.name,
            "kind": self.kind.value,
            "trace_id": self.trace_id,
            "parent_span_id": self.parent_span_id,
            "start_time": self.start_time.isoformat(),
            "end_time": self.end_time.isoformat() if self.end_time else None,
            "duration_ms": self.duration_ms,
            "input": self._truncate(self.input_data),
            "output": self._truncate(self.output_data),
            "error": self.error,
            "metadata": self.metadata,
            "attributes": self.attributes,
            "tags": self.tags
        }
    
    def _truncate(self, data: Any, max_len: int = 1000) -> Any:
        """截断数据"""
        if data is None:
            return None
        s = str(data)
        if len(s) > max_len:
            return s[:max_len] + "..."
        return data


@dataclass
class TraceContext:
    """追踪上下文"""
    trace_id: str
    name: str
    start_time: datetime = field(default_factory=datetime.now)
    end_time: Optional[datetime] = None
    
    spans: List[Span] = field(default_factory=list)
    _span_stack: List[str] = field(default_factory=list)  # 用于嵌套追踪
    
    metadata: Dict[str, Any] = field(default_factory=dict)
    tags: List[str] = field(default_factory=list)
    
    # 指标收集
    metrics: Dict[str, Any] = field(default_factory=lambda: defaultdict(int))
    
    @property
    def duration_ms(self) -> float:
        """总持续时间"""
        if self.end_time:
            return (self.end_time - self.start_time).total_seconds() * 1000
        return (datetime.now() - self.start_time).total_seconds() * 1000
    
    def add_span(self, span: Span) -> None:
        """添加跨度"""
        self.spans.append(span)
    
    def get_current_span(self) -> Optional[Span]:
        """获取当前活跃跨度"""
        if not self._span_stack:
            return None
        span_id = self._span_stack[-1]
        for span in reversed(self.spans):
            if span.span_id == span_id:
                return span
        return None
    
    def to_dict(self) -> Dict[str, Any]:
        """转换为字典"""
        return {
            "trace_id": self.trace_id,
            "name": self.name,
            "start_time": self.start_time.isoformat(),
            "end_time": self.end_time.isoformat() if self.end_time else None,
            "duration_ms": self.duration_ms,
            "spans": [s.to_dict() for s in self.spans],
            "metadata": self.metadata,
            "tags": self.tags,
            "metrics": dict(self.metrics)
        }


class AgentTracer:
    """
    Agent 追踪器:完整的追踪系统实现
    
    功能:
    1. 自动追踪 LLM 调用
    2. 自动追踪工具执行
    3. 支持嵌套跨度
    4. 收集关键指标
    5. 支持多后端导出(LangSmith、LangFuse、自建)
    """
    
    def __init__(
        self,
        service_name: str,
        export_backends: Optional[List[Callable]] = None
    ):
        self.service_name = service_name
        self.export_backends = export_backends or []
    
    def start_trace(
        self,
        name: str,
        input_data: Any = None,
        metadata: Optional[Dict] = None,
        tags: Optional[List[str]] = None
    ) -> TraceContext:
        """
        开始追踪会话
        
        Args:
            name: 追踪名称
            input_data: 输入数据
            metadata: 元数据
            tags: 标签
        
        Returns:
            TraceContext: 追踪上下文
        """
        trace = TraceContext(
            trace_id=str(uuid.uuid4())[:16],
            name=name,
            metadata=metadata or {},
            tags=tags or []
        )
        
        # 存入上下文
        _current_trace.set(trace)
        
        # 创建根跨度
        root_span = Span(
            span_id=str(uuid.uuid4())[:16],
            name=name,
            kind=SpanKind.INTERNAL,
            trace_id=trace.trace_id,
            parent_span_id=None,
            input_data=input_data
        )
        trace.add_span(root_span)
        trace._span_stack.append(root_span.span_id)
        
        return trace
    
    def end_trace(
        self,
        output_data: Any = None,
        error: Optional[str] = None,
        metadata: Optional[Dict] = None
    ) -> Optional[TraceContext]:
        """结束追踪"""
        trace = _current_trace.get()
        
        if not trace:
            return None
        
        trace.end_time = datetime.now()
        
        # 更新根跨度
        root_span = trace.spans[0]
        root_span.end_time = trace.end_time
        root_span.output_data = output_data
        root_span.error = error
        
        # 更新元数据
        if metadata:
            trace.metadata.update(metadata)
        
        # 导出到后端
        for backend in self.export_backends:
            try:
                backend(trace)
            except Exception as e:
                print(f"导出追踪数据失败: {e}")
        
        # 清理上下文
        _current_trace.set(None)
        
        return trace
    
    def create_span(
        self,
        name: str,
        kind: SpanKind,
        input_data: Any = None,
        metadata: Optional[Dict] = None
    ) -> Span:
        """创建新的跨度"""
        trace = _current_trace.get()
        
        if not trace:
            # 没有活跃追踪,创建一个匿名的
            trace = self.start_trace(f"implicit.{name}")
        
        # 获取父跨度
        parent_span = trace.get_current_span()
        parent_span_id = parent_span.span_id if parent_span else None
        
        # 创建新跨度
        span = Span(
            span_id=str(uuid.uuid4())[:16],
            name=name,
            kind=kind,
            trace_id=trace.trace_id,
            parent_span_id=parent_span_id,
            input_data=input_data,
            metadata=metadata or {}
        )
        
        trace.add_span(span)
        trace._span_stack.append(span.span_id)
        
        return span
    
    def end_span(
        self,
        span: Span,
        output_data: Any = None,
        error: Optional[str] = None,
        attributes: Optional[Dict] = None
    ) -> None:
        """结束跨度"""
        span.end_time = datetime.now()
        span.output_data = output_data
        span.error = error
        
        if attributes:
            span.attributes.update(attributes)
        
        # 更新追踪上下文
        trace = _current_trace.get()
        if trace and span.span_id in trace._span_stack:
            trace._span_stack.remove(span.span_id)
        
        # 更新指标
        if span.kind == SpanKind.LLM and output_data:
            trace.metrics["llm_calls"] += 1
            if hasattr(output_data, "usage_metadata"):
                usage = output_data.usage_metadata
                trace.metrics["prompt_tokens"] += usage.get("input_tokens", 0)
                trace.metrics["completion_tokens"] += usage.get("output_tokens", 0)
        
        if span.kind == SpanKind.TOOL:
            trace.metrics["tool_calls"] += 1
    
    def span(self, name: str, kind: SpanKind = SpanKind.INTERNAL):
        """
        上下文管理器:自动管理跨度生命周期
        
        用法:
        with tracer.span("query_order", SpanKind.TOOL) as span:
            result = query_order(order_id)
            span.output_data = result
        """
        
        class SpanContext:
            def __init__(ctx, tracer, name, kind):
                ctx.tracer = tracer
                ctx.name = name
                ctx.kind = kind
                ctx.span = None
            
            def __enter__(ctx):
                ctx.span = ctx.tracer.create_span(
                    name=ctx.name,
                    kind=ctx.kind
                )
                return ctx.span
            
            def __exit__(ctx, exc_type, exc_val, exc_tb):
                if exc_type:
                    ctx.span.error = str(exc_val)
                ctx.tracer.end_span(ctx.span)
                return False
        
        return SpanContext(self, name, kind)
    
    def trace_llm(
        self,
        model_name: str,
        prompt: Any,
        metadata: Optional[Dict] = None
    ):
        """
        LLM 调用追踪装饰器
        
        用法:
        @tracer.trace_llm("gpt-4", {"temperature": 0})
        async def call_llm(messages):
            return await llm.agenerate(messages)
        """
        
        def decorator(func: Callable) -> Callable:
            async def async_wrapper(*args, **kwargs):
                with self.span(f"llm.{model_name}", SpanKind.LLM) as span:
                    span.attributes["model"] = model_name
                    span.attributes["prompt_length"] = len(str(prompt))
                    
                    result = await func(*args, **kwargs)
                    
                    if hasattr(result, "usage_metadata"):
                        span.attributes["usage"] = result.usage_metadata
                    
                    return result
            
            def sync_wrapper(*args, **kwargs):
                with self.span(f"llm.{model_name}", SpanKind.LLM) as span:
                    span.attributes["model"] = model_name
                    result = func(*args, **kwargs)
                    return result
            
            import asyncio
            if asyncio.iscoroutinefunction(func):
                return async_wrapper
            return sync_wrapper
        
        return decorator
    
    def trace_tool(self, tool_name: str):
        """
        工具调用追踪装饰器
        
        用法:
        @tracer.trace_tool("get_weather")
        def get_weather(city, date):
            return weather_api.query(city, date)
        """
        
        def decorator(func: Callable) -> Callable:
            async def async_wrapper(*args, **kwargs):
                with self.span(f"tool.{tool_name}", SpanKind.TOOL) as span:
                    span.attributes["tool_name"] = tool_name
                    span.attributes["input"] = kwargs
                    
                    try:
                        result = await func(*args, **kwargs)
                        span.attributes["success"] = True
                        return result
                    except Exception as e:
                        span.attributes["success"] = False
                        span.attributes["error"] = str(e)
                        raise
            
            def sync_wrapper(*args, **kwargs):
                with self.span(f"tool.{tool_name}", SpanKind.TOOL) as span:
                    span.attributes["tool_name"] = tool_name
                    try:
                        result = func(*args, **kwargs)
                        span.attributes["success"] = True
                        return result
                    except Exception as e:
                        span.attributes["success"] = False
                        span.attributes["error"] = str(e)
                        raise
            
            import asyncio
            if asyncio.iscoroutinefunction(func):
                return async_wrapper
            return sync_wrapper
        
        return decorator


# ============== 导出后端示例 ==============

def export_to_langsmith(trace: TraceContext) -> None:
    """导出到 LangSmith"""
    # 实现 LangSmith 导出逻辑
    pass


def export_to_langfuse(trace: TraceContext) -> None:
    """导出到 LangFuse"""
    # 实现 LangFuse 导出逻辑
    pass


def export_to_file(trace: TraceContext, filepath: str = "traces.jsonl") -> None:
    """导出到文件"""
    import json
    
    with open(filepath, "a") as f:
        f.write(json.dumps(trace.to_dict(), ensure_ascii=False) + "\n")


def export_to_stdout(trace: TraceContext) -> None:
    """导出到标准输出(调试用)"""
    import json
    print(json.dumps(trace.to_dict(), indent=2, ensure_ascii=False))

四、关键指标体系

4.1 Agent 核心指标

# Agent 可观测性指标配置
metrics_config:
  # 核心业务指标
  business_metrics:
    - name: "agent.task.success_rate"
      description: "任务成功率"
      type: "gauge"
      aggregation: "ratio"
      
    - name: "agent.task.avg_duration_ms"
      description: "平均任务耗时"
      type: "histogram"
      buckets: [100, 500, 1000, 2000, 5000, 10000, 30000]
      
    - name: "agent.task.user_satisfaction"
      description: "用户满意度评分"
      type: "gauge"
      aggregation: "avg"
  
  # LLM 调用指标
  llm_metrics:
    - name: "llm.calls.total"
      description: "LLM 调用总次数"
      type: "counter"
      
    - name: "llm.tokens.prompt"
      description: "Prompt Token 消耗"
      type: "counter"
      
    - name: "llm.tokens.completion"
      description: "Completion Token 消耗"
      type: "counter"
      
    - name: "llm.latency.p50_ms"
      description: "LLM 延迟 P50"
      type: "gauge"
      
    - name: "llm.latency.p99_ms"
      description: "LLM 延迟 P99"
      type: "gauge"
  
  # 工具调用指标
  tool_metrics:
    - name: "tool.calls.total"
      description: "工具调用总次数"
      type: "counter"
      labels: ["tool_name"]
      
    - name: "tool.calls.success_rate"
      description: "工具调用成功率"
      type: "gauge"
      labels: ["tool_name"]
      
    - name: "tool.latency.avg_ms"
      description: "工具平均延迟"
      type: "gauge"
      labels: ["tool_name"]
      
    - name: "tool.errors.total"
      description: "工具错误总次数"
      type: "counter"
      labels: ["tool_name", "error_type"]
  
  # Agent 行为指标
  agent_metrics:
    - name: "agent.steps.avg"
      description: "平均决策步数"
      type: "histogram"
      buckets: [1, 2, 3, 5, 7, 10, 15, 20]
      
    - name: "agent.loops.detected"
      description: "检测到的循环次数"
      type: "counter"
      
    - name: "agent.escalations.human"
      description: "转人工次数"
      type: "counter"

4.2 指标收集器实现

"""
Agent 指标收集器
"""

from typing import Dict, List, Optional, Any
from dataclasses import dataclass, field
from datetime import datetime
from collections import defaultdict
import time


@dataclass
class MetricPoint:
    """指标数据点"""
    name: str
    value: float
    timestamp: datetime = field(default_factory=datetime.now)
    labels: Dict[str, str] = field(default_factory=dict)


class MetricsCollector:
    """
    指标收集器
    
    功能:
    1. 计数指标
    2. 数值指标
    3. 直方图指标
    4. 聚合计算
    """
    
    def __init__(self):
        self._counters: Dict[str, float] = defaultdict(float)
        self._gauges: Dict[str, float] = {}
        self._histograms: Dict[str, List[float]] = defaultdict(list)
        self._last_values: Dict[str, float] = {}
    
    def increment(
        self,
        name: str,
        value: float = 1.0,
        labels: Optional[Dict[str, str]] = None
    ) -> None:
        """增加计数"""
        key = self._make_key(name, labels)
        self._counters[key] += value
    
    def gauge(
        self,
        name: str,
        value: float,
        labels: Optional[Dict[str, str]] = None
    ) -> None:
        """设置仪表值"""
        key = self._make_key(name, labels)
        self._gauges[key] = value
        self._last_values[key] = value
    
    def histogram(
        self,
        name: str,
        value: float,
        labels: Optional[Dict[str, str]] = None
    ) -> None:
        """记录直方图值"""
        key = self._make_key(name, labels)
        self._histograms[key].append(value)
    
    def _make_key(self, name: str, labels: Optional[Dict[str, str]]) -> str:
        """生成带标签的键"""
        if not labels:
            return name
        label_str = ",".join(f"{k}={v}" for k, v in sorted(labels.items()))
        return f"{name}{{{label_str}}}"
    
    def get_prometheus_metrics(self) -> str:
        """生成 Prometheus 格式的指标"""
        lines = []
        
        # 计数器
        for key, value in self._counters.items():
            name = key.split("{")[0]
            lines.append(f"# TYPE {name} counter")
            lines.append(f"{key} {value}")
        
        # 仪表
        for key, value in self._gauges.items():
            name = key.split("{")[0]
            lines.append(f"# TYPE {name} gauge")
            lines.append(f"{key} {value}")
        
        # 直方图
        for key, values in self._histograms.items():
            name = key.split("{")[0]
            lines.append(f"# TYPE {name} histogram")
            
            # 计算分位数
            sorted_values = sorted(values)
            total = len(sorted_values)
            
            for quantile in [0.5, 0.9, 0.95, 0.99]:
                idx = int(total * quantile)
                value = sorted_values[idx] if idx < total else sorted_values[-1]
                labels = key.split("{")[1].rstrip("}") if "{" in key else ""
                q_key = f'{name}_quantile{{{labels},quantile="{quantile}"}}'
                lines.append(f"{q_key} {value}")
            
            # 总和与计数
            sum_key = f'{name}_sum{key[key.index("{"):] if "{" in key else ""}'
            count_key = f'{name}_count{key[key.index("{"):] if "{" in key else ""}'
            lines.append(f"{sum_key} {sum(values)}")
            lines.append(f"{count_key} {total}")
        
        return "\n".join(lines)
    
    def get_summary(self) -> Dict[str, Any]:
        """获取指标摘要"""
        return {
            "counters": dict(self._counters),
            "gauges": dict(self._gauges),
            "histograms": {
                k: {
                    "count": len(v),
                    "mean": sum(v) / len(v),
                    "min": min(v),
                    "max": max(v),
                    "p50": self._percentile(v, 0.5),
                    "p95": self._percentile(v, 0.95),
                    "p99": self._percentile(v, 0.99)
                }
                for k, v in self._histograms.items()
            }
        }
    
    def _percentile(self, values: List[float], p: float) -> float:
        """计算百分位数"""
        if not values:
            return 0.0
        sorted_values = sorted(values)
        idx = int(len(sorted_values) * p)
        return sorted_values[min(idx, len(sorted_values) - 1)]


# 全局指标收集器
_global_metrics = MetricsCollector()


def get_metrics_collector() -> MetricsCollector:
    """获取全局指标收集器"""
    return _global_metrics


# ============== Agent 指标追踪装饰器 ==============

def track_agent_metrics(operation: str):
    """追踪 Agent 操作指标的装饰器"""
    
    def decorator(func):
        async def async_wrapper(*args, **kwargs):
            start_time = time.time()
            metrics = get_metrics_collector()
            
            try:
                result = await func(*args, **kwargs)
                duration_ms = (time.time() - start_time) * 1000
                
                # 记录指标
                metrics.increment("agent.operations.total", labels={"operation": operation})
                metrics.histogram("agent.operations.duration_ms", duration_ms, {"operation": operation})
                
                return result
                
            except Exception as e:
                metrics.increment("agent.operations.errors", labels={
                    "operation": operation,
                    "error_type": type(e).__name__
                })
                raise
        
        def sync_wrapper(*args, **kwargs):
            start_time = time.time()
            metrics = get_metrics_collector()
            
            try:
                result = func(*args, **kwargs)
                duration_ms = (time.time() - start_time) * 1000
                
                metrics.increment("agent.operations.total", labels={"operation": operation})
                metrics.histogram("agent.operations.duration_ms", duration_ms, {"operation": operation})
                
                return result
                
            except Exception as e:
                metrics.increment("agent.operations.errors", labels={
                    "operation": operation,
                    "error_type": type(e).__name__
                })
                raise
        
        import asyncio
        if asyncio.iscoroutinefunction(func):
            return async_wrapper
        return sync_wrapper
    
    return decorator

五、调试技巧:回放与断点

5.1 回放系统

"""
Agent 执行回放系统
支持:记录、重放、调试
"""

import json
import pickle
from typing import Any, Dict, List, Optional
from dataclasses import dataclass, field
from datetime import datetime
from pathlib import Path

from .tracing import TraceContext


@dataclass
class ReplayEvent:
    """回放事件"""
    timestamp: datetime
    event_type: str  # "llm_call", "tool_call", "decision", "error"
    data: Dict[str, Any]


@dataclass
class ReplaySession:
    """回放会话"""
    session_id: str
    trace_id: str
    start_time: datetime
    end_time: Optional[datetime] = None
    
    events: List[ReplayEvent] = field(default_factory=list)
    
    # 原始输入输出
    input_data: Any = None
    output_data: Any = None
    
    # 评估结果
    evaluation: Dict[str, Any] = field(default_factory=dict)


class ReplayRecorder:
    """
    回放记录器:记录 Agent 执行过程
    
    用于:
    1. 生产问题排查
    2. 测试用例生成
    3. 模型效果评估
    """
    
    def __init__(self, storage_path: str = "./replays"):
        self.storage_path = Path(storage_path)
        self.storage_path.mkdir(parents=True, exist_ok=True)
        
        self._current_session: Optional[ReplaySession] = None
        self._sessions: Dict[str, ReplaySession] = {}
    
    def start_recording(
        self,
        session_id: str,
        trace_id: str,
        input_data: Any
    ) -> ReplaySession:
        """开始记录"""
        
        session = ReplaySession(
            session_id=session_id,
            trace_id=trace_id,
            start_time=datetime.now(),
            input_data=input_data
        )
        
        self._current_session = session
        self._sessions[session_id] = session
        
        return session
    
    def record_event(
        self,
        event_type: str,
        data: Dict[str, Any]
    ) -> None:
        """记录事件"""
        
        if not self._current_session:
            return
        
        event = ReplayEvent(
            timestamp=datetime.now(),
            event_type=event_type,
            data=data
        )
        
        self._current_session.events.append(event)
    
    def end_recording(
        self,
        output_data: Any,
        evaluation: Optional[Dict[str, Any]] = None
    ) -> ReplaySession:
        """结束记录"""
        
        if not self._current_session:
            raise RuntimeError("没有活跃的记录会话")
        
        self._current_session.end_time = datetime.now()
        self._current_session.output_data = output_data
        
        if evaluation:
            self._current_session.evaluation = evaluation
        
        # 持久化
        self._persist_session(self._current_session)
        
        session = self._current_session
        self._current_session = None
        
        return session
    
    def _persist_session(self, session: ReplaySession) -> None:
        """持久化会话"""
        
        filepath = self.storage_path / f"{session.session_id}.json"
        
        data = {
            "session_id": session.session_id,
            "trace_id": session.trace_id,
            "start_time": session.start_time.isoformat(),
            "end_time": session.end_time.isoformat() if session.end_time else None,
            "input": self._serialize(session.input_data),
            "output": self._serialize(session.output_data),
            "evaluation": session.evaluation,
            "events": [
                {
                    "timestamp": e.timestamp.isoformat(),
                    "event_type": e.event_type,
                    "data": e.data
                }
                for e in session.events
            ]
        }
        
        with open(filepath, "w", encoding="utf-8") as f:
            json.dump(data, f, ensure_ascii=False, indent=2)
    
    def _serialize(self, data: Any) -> Any:
        """序列化数据"""
        try:
            return json.dumps(data, ensure_ascii=False, default=str)
        except:
            return str(data)
    
    def load_session(self, session_id: str) -> Optional[ReplaySession]:
        """加载会话"""
        
        filepath = self.storage_path / f"{session_id}.json"
        
        if not filepath.exists():
            return None
        
        with open(filepath, "r", encoding="utf-8") as f:
            data = json.load(f)
        
        session = ReplaySession(
            session_id=data["session_id"],
            trace_id=data["trace_id"],
            start_time=datetime.fromisoformat(data["start_time"]),
            end_time=datetime.fromisoformat(data["end_time"]) if data.get("end_time") else None,
            input_data=data.get("input"),
            output_data=data.get("output"),
            evaluation=data.get("evaluation", {})
        )
        
        session.events = [
            ReplayEvent(
                timestamp=datetime.fromisoformat(e["timestamp"]),
                event_type=e["event_type"],
                data=e["data"]
            )
            for e in data.get("events", [])
        ]
        
        return session
    
    def list_sessions(
        self,
        limit: int = 100,
        offset: int = 0
    ) -> List[Dict[str, Any]]:
        """列出会话"""
        
        sessions = []
        
        for filepath in sorted(self.storage_path.glob("*.json"), reverse=True):
            try:
                with open(filepath, "r", encoding="utf-8") as f:
                    data = json.load(f)
                    sessions.append({
                        "session_id": data["session_id"],
                        "trace_id": data["trace_id"],
                        "start_time": data["start_time"],
                        "event_count": len(data.get("events", []))
                    })
            except:
                continue
        
        return sessions[offset:offset + limit]


class ReplayDebugger:
    """
    回放调试器:重放历史执行,支持断点
    """
    
    def __init__(self, recorder: ReplayRecorder):
        self.recorder = recorder
        self._breakpoints: Dict[str, bool] = {}
    
    def set_breakpoint(self, event_type: str) -> None:
        """设置断点"""
        self._breakpoints[event_type] = True
    
    def clear_breakpoint(self, event_type: str) -> None:
        """清除断点"""
        self._breakpoints.pop(event_type, None)
    
    def replay(
        self,
        session_id: str,
        callback: Optional[callable] = None
    ) -> None:
        """
        重放会话
        
        Args:
            session_id: 会话 ID
            callback: 每个事件的回调函数,可用于调试
        """
        
        session = self.recorder.load_session(session_id)
        
        if not session:
            raise ValueError(f"会话不存在: {session_id}")
        
        print(f"开始重放会话: {session_id}")
        print(f"开始时间: {session.start_time}")
        print(f"事件数: {len(session.events)}")
        print(f"输入: {session.input_data}")
        print("-" * 80)
        
        for i, event in enumerate(session.events):
            # 检查断点
            if self._breakpoints.get(event.event_type):
                print(f"\n[断点] 事件类型: {event.event_type}")
                input("按 Enter 继续...")
            
            print(f"\n[{i+1}] {event.timestamp} - {event.event_type}")
            print(json.dumps(event.data, indent=2, ensure_ascii=False, default=str))
            
            # 执行回调
            if callback:
                callback(event)
        
        print("\n" + "-" * 80)
        print(f"重放结束")
        print(f"输出: {session.output_data}")
        print(f"评估: {session.evaluation}")

5.2 调试面板

"""
Agent 调试面板:实时查看 Agent 执行状态
"""

import asyncio
from typing import Dict, List, Any, Optional
from dataclasses import dataclass
from datetime import datetime
import streamlit as st
from tracelite import TraceLite


@dataclass
class DebugState:
    """调试状态"""
    current_step: int = 0
    total_steps: int = 0
    current_span: Optional[str] = None
    decision_reasoning: Optional[str] = None
    last_tool_call: Optional[Dict] = None
    error: Optional[str] = None


class AgentDebugPanel:
    """
    Agent 调试面板
    
    功能:
    1. 实时显示执行状态
    2. 查看历史追踪
    3. 设置条件断点
    4. 分析性能瓶颈
    """
    
    def __init__(self):
        self.tracer = TraceLite()
        self.state = DebugState()
    
    def render(self):
        """渲染调试面板"""
        
        st.set_page_config(page_title="Agent Debug Panel")
        
        # 侧边栏:会话列表
        with st.sidebar:
            st.title("调试会话")
            
            # 会话选择
            sessions = self.tracer.list_sessions()
            session_options = {s["session_id"]: s for s in sessions}
            
            selected = st.selectbox(
                "选择会话",
                options=list(session_options.keys()),
                format_func=lambda x: f"{x} ({session_options[x]['event_count']} 事件)"
            )
            
            if selected:
                self._render_session_detail(session_options[selected])
        
        # 主区域:执行可视化
        st.title("Agent 执行追踪")
        
        # 追踪树
        self._render_trace_tree()
        
        # 指标面板
        self._render_metrics_panel()
        
        # 日志流
        self._render_log_stream()
    
    def _render_session_detail(self, session: Dict):
        """渲染会话详情"""
        
        st.subheader("会话信息")
        
        col1, col2 = st.columns(2)
        
        with col1:
            st.metric("事件数", session.get("event_count", 0))
        
        with col2:
            st.metric("开始时间", session.get("start_time", "N/A"))
        
        # 加载会话
        full_session = self.tracer.load_session(session["session_id"])
        
        if full_session:
            st.json({
                "input": full_session.input_data,
                "output": full_session.output_data,
                "evaluation": full_session.evaluation
            })
    
    def _render_trace_tree(self):
        """渲染追踪树"""
        
        st.subheader("执行追踪")
        
        # 模拟追踪数据
        trace_data = {
            "name": "agent.query_order",
            "duration_ms": 2500,
            "spans": [
                {
                    "id": "span1",
                    "name": "intent_detection",
                    "kind": "llm",
                    "duration_ms": 500,
                    "status": "ok"
                },
                {
                    "id": "span2",
                    "name": "query_order",
                    "kind": "tool",
                    "duration_ms": 1200,
                    "status": "ok",
                    "children": [
                        {
                            "id": "span2.1",
                            "name": "api_call",
                            "kind": "internal",
                            "duration_ms": 1000,
                            "status": "ok"
                        }
                    ]
                },
                {
                    "id": "span3",
                    "name": "format_response",
                    "kind": "llm",
                    "duration_ms": 800,
                    "status": "ok"
                }
            ]
        }
        
        # 使用树形图展示
        self._render_span_tree(trace_data["spans"], level=0)
    
    def _render_span_tree(self, spans: List[Dict], level: int):
        """递归渲染跨度树"""
        
        for span in spans:
            indent = "  " * level
            status_emoji = "✅" if span["status"] == "ok" else "❌"
            
            # 展开/折叠
            with st.expander(f"{indent}{status_emoji} {span['name']} ({span['duration_ms']}ms)"):
                st.json(span)
                
                # 递归渲染子跨度
                if "children" in span:
                    self._render_span_tree(span["children"], level + 1)
    
    def _render_metrics_panel(self):
        """渲染指标面板"""
        
        st.subheader("关键指标")
        
        col1, col2, col3, col4 = st.columns(4)
        
        with col1:
            st.metric("LLM 调用", 3, delta="-1 vs 上次")
        
        with col2:
            st.metric("Token 消耗", "1.2K", delta="-15%")
        
        with col3:
            st.metric("工具调用", 1, delta="0 vs 上次")
        
        with col4:
            st.metric("执行时间", "2.5s", delta="-0.3s")
    
    def _render_log_stream(self):
        """渲染日志流"""
        
        st.subheader("实时日志")
        
        # 模拟日志
        logs = [
            {"time": "10:00:01", "level": "INFO", "message": "开始处理请求"},
            {"time": "10:00:01", "level": "DEBUG", "message": "意图识别: query_order"},
            {"time": "10:00:02", "level": "INFO", "message": "调用工具: query_order"},
            {"time": "10:00:03", "level": "INFO", "message": "工具执行成功"},
            {"time": "10:00:03", "level": "DEBUG", "message": "生成响应..."},
            {"time": "10:00:04", "level": "INFO", "message": "请求处理完成"},
        ]
        
        for log in logs:
            color = {
                "INFO": "white",
                "DEBUG": "gray",
                "WARNING": "yellow",
                "ERROR": "red"
            }.get(log["level"], "white")
            
            st.text(f"[{log['time']}] {log['level']}: {log['message']}")

六、告警规则与自动告警

6.1 告警规则配置

# Agent 告警规则配置
alerting_config:
  # 告警通道
  channels:
    - name: "slack"
      type: "webhook"
      url: "https://hooks.slack.com/services/xxx"
      enabled: true
    
    - name: "email"
      type: "smtp"
      recipients: ["oncall@company.com"]
      enabled: false
    
    - name: "pagerduty"
      type: "pagerduty"
      integration_key: "xxx"
      enabled: true
  
  # 告警规则
  rules:
    # 无限循环检测
    - name: "infinite_loop_detection"
      description: "检测 Agent 进入无限循环"
      condition: >
        agent.loops.detected > 0 
        and agent.steps.current > 20
      severity: "critical"
      cooldown: 5m
      channels: ["slack", "pagerduty"]
    
    # 工具调用异常
    - name: "tool_failure_rate_high"
      description: "工具调用失败率过高"
      condition: >
        rate(tool.calls.errors) / rate(tool.calls.total) > 0.1
      severity: "warning"
      cooldown: 10m
      channels: ["slack"]
    
    # 延迟过高
    - name: "latency_high"
      description: "LLM 响应延迟过高"
      condition: >
        llm.latency.p99 > 10000  # 10秒
      severity: "warning"
      cooldown: 5m
      channels: ["slack"]
    
    # Token 消耗异常
    - name: "token_consumption_anomaly"
      description: "Token 消耗异常"
      condition: >
        agent.tokens.current > agent.tokens.avg * 3
      severity: "warning"
      cooldown: 15m
      channels: ["slack"]
    
    # 错误率过高
    - name: "error_rate_critical"
      description: "Agent 错误率超过阈值"
      condition: >
        rate(agent.operations.errors) / rate(agent.operations.total) > 0.05
      severity: "critical"
      cooldown: 5m
      channels: ["slack", "pagerduty"]
    
    # 转人工率过高
    - name: "human_escalation_high"
      description: "转人工频率过高"
      condition: >
        rate(agent.escalations.human) / rate(agent.operations.total) > 0.1
      severity: "warning"
      cooldown: 30m
      channels: ["slack"]
    
    # 成本超限
    - name: "cost_over_budget"
      description: "API 成本超过预算"
      condition: >
        agent.cost.daily > 1000  # 美元
      severity: "warning"
      cooldown: 1h
      channels: ["slack", "email"]

6.2 告警系统实现

"""
Agent 告警系统
"""

import asyncio
from typing import Dict, List, Optional, Callable, Any
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from enum import Enum
import logging

logger = logging.getLogger(__name__)


class AlertSeverity(Enum):
    """告警级别"""
    INFO = "info"
    WARNING = "warning"
    CRITICAL = "critical"


@dataclass
class Alert:
    """告警"""
    name: str
    message: str
    severity: AlertSeverity
    
    timestamp: datetime = field(default_factory=datetime.now)
    labels: Dict[str, str] = field(default_factory=dict)
    
    # 关联数据
    trace_id: Optional[str] = None
    session_id: Optional[str] = None
    
    # 上下文
    context: Dict[str, Any] = field(default_factory=dict)


@dataclass
class AlertRule:
    """告警规则"""
    name: str
    description: str
    condition: Callable[[Dict], bool]  # 返回 True 触发告警
    severity: AlertSeverity
    
    # 抑制配置
    cooldown: timedelta = field(default_factory=lambda: timedelta(minutes=5))
    last_triggered: Optional[datetime] = None
    
    # 通道
    channels: List[str] = field(default_factory=list)
    
    # 启用状态
    enabled: bool = True


class AlertManager:
    """
    告警管理器
    
    功能:
    1. 定义和评估告警规则
    2. 告警聚合和抑制
    3. 多通道通知
    4. 告警历史记录
    """
    
    def __init__(self):
        self.rules: Dict[str, AlertRule] = {}
        self.channels: Dict[str, Callable] = {}
        self.alert_history: List[Alert] = []
        
        # 初始化默认规则
        self._init_default_rules()
    
    def _init_default_rules(self):
        """初始化默认告警规则"""
        
        # 无限循环检测
        self.add_rule(AlertRule(
            name="infinite_loop_detection",
            description="检测 Agent 进入无限循环",
            condition=lambda m: m.get("agent.loops.detected", 0) > 0 and 
                               m.get("agent.steps.current", 0) > 20,
            severity=AlertSeverity.CRITICAL,
            channels=["slack", "pagerduty"]
        ))
        
        # 工具失败率
        self.add_rule(AlertRule(
            name="tool_failure_rate_high",
            description="工具调用失败率超过 10%",
            condition=lambda m: (
                m.get("tool.calls.total", 0) > 0 and
                m.get("tool.calls.errors", 0) / m.get("tool.calls.total", 1) > 0.1
            ),
            severity=AlertSeverity.WARNING,
            channels=["slack"]
        ))
        
        # 延迟过高
        self.add_rule(AlertRule(
            name="latency_high",
            description="LLM P99 延迟超过 10 秒",
            condition=lambda m: m.get("llm.latency.p99", 0) > 10000,
            severity=AlertSeverity.WARNING,
            channels=["slack"]
        ))
    
    def add_rule(self, rule: AlertRule) -> None:
        """添加告警规则"""
        self.rules[rule.name] = rule
    
    def add_channel(self, name: str, sender: Callable) -> None:
        """添加告警通道"""
        self.channels[name] = sender
    
    def evaluate(self, metrics: Dict[str, Any]) -> List[Alert]:
        """
        评估所有规则
        
        Args:
            metrics: 当前指标数据
        
        Returns:
            触发的告警列表
        """
        triggered_alerts = []
        
        for rule in self.rules.values():
            if not rule.enabled:
                continue
            
            # 检查冷却期
            if rule.last_triggered:
                if datetime.now() - rule.last_triggered < rule.cooldown:
                    continue
            
            # 评估条件
            try:
                if rule.condition(metrics):
                    alert = Alert(
                        name=rule.name,
                        message=rule.description,
                        severity=rule.severity,
                        context={"metrics": metrics}
                    )
                    
                    triggered_alerts.append(alert)
                    rule.last_triggered = datetime.now()
                    
            except Exception as e:
                logger.error(f"规则 {rule.name} 评估失败: {e}")
        
        return triggered_alerts
    
    async def send_alerts(self, alerts: List[Alert]) -> None:
        """发送告警"""
        
        for alert in alerts:
            # 记录历史
            self.alert_history.append(alert)
            
            # 发送到各通道
            for rule in self.rules.values():
                if rule.name == alert.name:
                    for channel_name in rule.channels:
                        if channel_name in self.channels:
                            try:
                                await self.channels[channel_name](alert)
                            except Exception as e:
                                logger.error(f"发送告警到 {channel_name} 失败: {e}")
    
    def get_active_alerts(self) -> List[Alert]:
        """获取活跃告警(最近 1 小时内)"""
        
        cutoff = datetime.now() - timedelta(hours=1)
        
        return [
            a for a in self.alert_history
            if a.timestamp > cutoff
        ]
    
    def get_alert_history(
        self,
        limit: int = 100,
        severity: Optional[AlertSeverity] = None
    ) -> List[Alert]:
        """获取告警历史"""
        
        alerts = self.alert_history[-limit:]
        
        if severity:
            alerts = [a for a in alerts if a.severity == severity]
        
        return alerts


# ============== 告警通道实现 ==============

async def slack_sender(webhook_url: str, alert: Alert) -> None:
    """发送 Slack 告警"""
    
    import aiohttp
    
    severity_emoji = {
        AlertSeverity.INFO: "ℹ️",
        AlertSeverity.WARNING: "⚠️",
        AlertSeverity.CRITICAL: "🚨"
    }
    
    payload = {
        "text": f"{severity_emoji[alert.severity]} *{alert.name}*",
        "blocks": [
            {
                "type": "header",
                "text": {
                    "type": "plain_text",
                    "text": f"Agent 告警: {alert.name}"
                }
            },
            {
                "type": "section",
                "fields": [
                    {"type": "mrkdwn", "text": f"*级别:*\n{alert.severity.value}"},
                    {"type": "mrkdwn", "text": f"*时间:*\n{alert.timestamp.isoformat()}"}
                ]
            },
            {
                "type": "section",
                "text": {
                    "type": "mrkdwn",
                    "text": f"*描述:*\n{alert.message}"
                }
            }
        ]
    }
    
    async with aiohttp.ClientSession() as session:
        await session.post(webhook_url, json=payload)


# ============== 与 Agent 集成 ==============

class AlertingMiddleware:
    """
    告警中间件:集成到 Agent 执行链路
    """
    
    def __init__(self, alert_manager: AlertManager):
        self.alert_manager = alert_manager
        self.metrics_collector = get_metrics_collector()
    
    async def check_alerts(self) -> None:
        """检查告警规则"""
        
        metrics = self.metrics_collector.get_summary()
        
        # 评估规则
        alerts = self.alert_manager.evaluate(metrics)
        
        # 发送告警
        if alerts:
            await self.alert_manager.send_alerts(alerts)

七、完整可观测性集成代码

"""
Agent 可观测性完整集成示例
使用 LangChain + LangSmith/LangFuse + 自建追踪系统
"""

import os
from typing import Dict, Any, List, Optional
from langchain_openai import ChatOpenAI
from langchain.prompts import ChatPromptTemplate
from langchain.schema import StrOutputParser
from langchain.tools import tool
from langchain.agents import AgentExecutor, create_openai_functions_agent
from langchain.callbacks.manager import CallbackManager

from agent_tracing import AgentTracer, SpanKind
from agent_metrics import get_metrics_collector, track_agent_metrics
from agent_alerting import AlertManager, AlertingMiddleware


class ObservableAgent:
    """
    可观测 Agent:集成完整的可观测性能力
    
    功能:
    1. 全链路追踪
    2. 指标收集
    3. 自动告警
    4. 回放支持
    """
    
    def __init__(
        self,
        llm: ChatOpenAI,
        tools: List[Any],
        project_name: str = "agent-debugging"
    ):
        self.llm = llm
        self.tools = tools
        self.project_name = project_name
        
        # 初始化组件
        self._init_tracer()
        self._init_metrics()
        self._init_alerting()
        self._init_replay()
    
    def _init_tracer(self):
        """初始化追踪器"""
        
        # 导出后端列表
        export_backends = []
        
        # 如果配置了 LangSmith
        if os.getenv("LANGCHAIN_TRACING_V2"):
            from langchain.callbacks.tracers.langsmith import LangSmithTracer
            export_backends.append(
                lambda trace: LangSmithTracer(project_name=self.project_name).on_chain_end(
                    {"trace": trace.to_dict()}
                )
            )
        
        self.tracer = AgentTracer(
            service_name="observable_agent",
            export_backends=export_backends
        )
    
    def _init_metrics(self):
        """初始化指标收集"""
        
        self.metrics = get_metrics_collector()
    
    def _init_alerting(self):
        """初始化告警"""
        
        self.alert_manager = AlertManager()
        
        # 配置告警通道
        if os.getenv("SLACK_WEBHOOK_URL"):
            from agent_alerting import slack_sender
            self.alert_manager.add_channel(
                "slack",
                lambda alert: slack_sender(os.getenv("SLACK_WEBHOOK_URL"), alert)
            )
        
        self.alerting_middleware = AlertingMiddleware(self.alert_manager)
    
    def _init_replay(self):
        """初始化回放系统"""
        
        from agent_replay import ReplayRecorder, ReplayDebugger
        
        self.replay_recorder = ReplayRecorder()
        self.replay_debugger = ReplayDebugger(self.replay_recorder)
    
    @track_agent_metrics("agent.run")
    async def run(self, user_input: str, session_id: Optional[str] = None) -> Dict[str, Any]:
        """
        运行 Agent
        
        Args:
            user_input: 用户输入
            session_id: 会话 ID(用于回放)
        
        Returns:
            {
                "response": str,           # Agent 回复
                "trace_id": str,           # 追踪 ID
                "metrics": Dict,           # 指标快照
                "success": bool            # 是否成功
            }
        """
        
        import uuid
        trace_id = str(uuid.uuid4())[:16]
        session_id = session_id or trace_id
        
        # 开始记录
        self.replay_recorder.start_recording(
            session_id=session_id,
            trace_id=trace_id,
            input_data=user_input
        )
        
        # 开始追踪
        self.tracer.start_trace(
            name=f"agent.{session_id}",
            input_data={"user_input": user_input},
            metadata={"session_id": session_id}
        )
        
        try:
            # 构建 Agent
            prompt = ChatPromptTemplate.from_messages([
                ("system", "你是一个智能助手,使用工具来回答用户问题。"),
                ("human", "{input}"),
                ("ai", "{agent_scratchpad}")
            ])
            
            agent = create_openai_functions_agent(
                llm=self.llm,
                tools=self.tools,
                prompt=prompt
            )
            
            agent_executor = AgentExecutor.from_agent_and_tools(
                agent=agent,
                tools=self.tools,
                verbose=True,
                handle_parsing_errors=True
            )
            
            # 执行
            with self.tracer.span("agent_execution", SpanKind.CHAIN):
                result = await agent_executor.ainvoke({"input": user_input})
            
            # 结束追踪
            output = result.get("output", "")
            self.tracer.end_trace(output_data=output)
            
            # 结束记录
            self.replay_recorder.end_recording(
                output_data=output,
                evaluation={"success": True}
            )
            
            # 检查告警
            await self.alerting_middleware.check_alerts()
            
            return {
                "response": output,
                "trace_id": trace_id,
                "session_id": session_id,
                "metrics": self.metrics.get_summary(),
                "success": True
            }
            
        except Exception as e:
            # 错误处理
            self.tracer.end_trace(error=str(e))
            self.replay_recorder.end_recording(
                output_data=None,
                evaluation={"success": False, "error": str(e)}
            )
            
            return {
                "response": f"抱歉,处理您的请求时遇到错误: {str(e)}",
                "trace_id": trace_id,
                "session_id": session_id,
                "error": str(e),
                "success": False
            }
    
    def debug_session(self, session_id: str):
        """调试历史会话"""
        
        self.replay_debugger.replay(session_id)


# ============== 使用示例 ==============

async def demo_observable_agent():
    """可观测 Agent 使用示例"""
    
    # 初始化 LLM
    llm = ChatOpenAI(model="gpt-4o", temperature=0)
    
    # 定义工具
    @tool
    def get_weather(city: str) -> str:
        """获取城市天气"""
        weathers = {"北京": "晴 15°C", "上海": "多云 22°C"}
        return weathers.get(city, "未知")
    
    tools = [get_weather]
    
    # 创建可观测 Agent
    agent = ObservableAgent(llm=llm, tools=tools)
    
    # 运行
    result = await agent.run("北京天气怎么样?")
    
    print(f"回复: {result['response']}")
    print(f"追踪 ID: {result['trace_id']}")
    
    # 调试历史会话
    # agent.debug_session(result['session_id'])


if __name__ == "__main__":
    import asyncio
    asyncio.run(demo_observable_agent())

八、总结

Agent 系统的可观测性是构建可靠生产系统的基石。通过追踪、指标、日志三大支柱,我们可以让"不可预测"的 Agent 变得可控可追踪。

核心要点回顾:

  1. 追踪是基础:完整记录每个决策点、工具调用的输入输出

  2. 指标要全面:从业务、LLM、工具、Agent 行为多个维度收集

  3. 告警要精准:避免告警疲劳,聚焦真正需要人工介入的问题

  4. 回放是调试利器:生产问题的完美复现离不开完整的执行记录

技术栈推荐:

维度

推荐方案

追踪

LangSmith(开箱即用)、LangFuse(开源自部署)、自建 TraceLite

指标

Prometheus + Grafana

日志

ELK Stack / Loki

告警

AlertManager + 业务告警系统

可观测性不是事后补救,而是从设计之初就需要考虑的架构问题。投入建设的可观测性基础设施,终将在生产问题排查时得到丰厚回报。

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区