编程 Agent Harness 深度拆解:当 AI Agent 决定给自己套上缰绳——从裸奔到可控的生产级运行时架构如何重新定义 AI 工程化的终极形态

2026-08-05 06:17:27 +0800 CST views 11

Agent Harness 深度拆解:当 AI Agent 决定「给自己套上缰绳」——从裸奔到可控的生产级运行时架构如何重新定义 AI 工程化的终极形态

2026年,AI 工程圈最火的一个词不是某个新模型,不是某个新框架,而是一个古老单词的新用法:Harness(马具)。从 Mitchell Hashimoto 在博客中首次定义,到 OpenAI 发布百万行代码的实验报告,再到 Martin Fowler 的深度分析——几周之内,这个术语成了讨论 AI Agent 开发绕不开的话题。本文将深入拆解 Agent Harness 的架构设计、核心组件与生产级实现,带你看懂这场正在重塑 AI 工程化的核心范式革命。

目录

  1. 引言:当 Agent 从 Demo 走向生产
  2. 核心概念:什么是 Agent Harness
  3. 四层架构模型深度解析
  4. 核心组件一:执行与编排引擎
  5. 核心组件二:上下文与轨迹管理
  6. 核心组件三:交互层与执行环境
  7. 核心组件四:安全护栏与可观测性
  8. 代码实战:从零构建一个生产级 Agent Harness
  9. 性能优化:Harness 层的关键调优策略
  10. 总结与展望

一、引言:当 Agent 从 Demo 走向生产

1.1 一个让硅谷震动的实验

2026年2月,HashiCorp 联合创始人 Mitchell Hashimoto 发表了一篇博客,提出了一个看似简单却深刻的问题:

"我们能让 AI Agent 从零开始搭建一个完整的真实应用吗?"

OpenAI 随后用 Codex Agent 基于这个思路进行了大规模实验。结果令人震惊——Agent 确实能写出能运行的代码,但几乎无法通过生产级的代码审查。问题不在于模型能力,而在于缺乏一套可靠的"运行时基础设施"来约束、引导和验证 Agent 的行为。

这个发现催生了 Agent Harness(智能体驾驭层)这一全新概念。

1.2 裸奔的 Agent:为什么需要缰绳?

让我们直面现实:2025年到2026年初,大量企业在部署 AI Agent 时遇到了同一个痛点——"能做 Demo,但无法在企业环境稳定落地"

典型症状包括:

# 你可能遇到过这些场景:
# 场景1:Agent 无限循环调用工具
while true; do
    agent.call("search", query=last_query)  # 永远找不到结果
done

# 场景2:Agent 忘记之前的决策
agent.execute("create_user", params=...)  # 第一次:成功
agent.execute("create_user", params=...)  # 第二次:重复创建!

# 场景3:Agent 越权操作
agent.execute("delete_table", table="production_data")  # ???

这些问题的根源在于:我们给了 Agent 太多自由,却没有给它足够的约束

1.3 Harness 的隐喻:马、马具与骑手

Agent Harness 的核心隐喻来自骑马:

隐喻Agent 系统对应
马(Horse)AI 模型(LLM)—— 力量强大但难以预测
马具(Harness)运行时基础设施 —— 约束、引导、保护
骑手(Rider)人类开发者 —— 定义目标、监督过程、处理异常

OpenAI 的表述更加工程化:

Coding Agent = AI Model + Harness

Harness 不是单次对话的包装,而是一个能驱动工具调用、状态流转、事件流、客户端交互的长期运行系统


二、核心概念:什么是 Agent Harness

2.1 定义

Agent Harness 是围绕 AI Agent 构建的生产级运行时基础设施与工程化范式。它解决的核心问题是:如何让 Agent 在复杂、长周期任务中保持稳定、可控、可审计

Salesforce 的定义更直白:

Agent Harness 是一层 operational software layer,负责管理 AI 的 tools、memory、safety,从而让 autonomous task execution 更可靠。

2.2 核心公式

Agent Harness = 执行引擎 + 上下文管理 + 交互层 + 安全护栏 + 可观测性

更具体地:

Agent Harness = {
    ExecutionEngine:     模型调用 + 工具执行 + 状态机
    ContextManager:      记忆系统 + 轨迹持久化 + 上下文压缩
    InteractionSurface:  API网关 + UI适配 + 事件总线
    SafetyGuardrails:    权限控制 + 沙箱隔离 + 人类审批
    Observability:       日志追踪 + 性能监控 + 成本核算
}

2.3 Agent Harness vs 传统 Agent 框架

很多人会问:Agent Harness 和 LangChain、AutoGPT 这些框架有什么区别?

维度传统 Agent 框架Agent Harness
定位Agent 的"大脑"Agent 的"身体"
关注点如何思考(推理策略)如何执行(运行时基础设施)
抽象层级单次任务循环长期运行的生产系统
状态管理简单的内存状态持久化、可恢复的轨迹
安全机制基础的输入过滤多层沙箱 + 权限矩阵
可观测性基础日志分布式追踪 + 成本归因

关键洞察:Harness 不是替代框架,而是在框架之上提供一层统一的管控平面


三、四层架构模型深度解析

基于 Awesome-Agent-Harness 仓库中的调研论文,Agent Harness 可以被组织为一个四层架构栈

┌─────────────────────────────────────────────────┐
│  Layer 4: 安全与治理 (Safety & Governance)       │
│  权限控制 · 沙箱隔离 · 合规审计 · 人类审批       │
├─────────────────────────────────────────────────┤
│  Layer 3: 交互层 (Interaction Surface)           │
│  API网关 · UI适配 · 事件总线 · 多模态接口         │
├─────────────────────────────────────────────────┤
│  Layer 2: 上下文管理 (Context & Trajectory)      │
│  记忆系统 · 轨迹持久化 · 状态压缩 · 可观测性     │
├─────────────────────────────────────────────────┤
│  Layer 1: 执行引擎 (Execution & Orchestration)   │
│  模型路由 · 工具调用 · 状态机 · 多Agent编排       │
└─────────────────────────────────────────────────┘

Layer 1:执行引擎 —— 时间引擎

这是 Harness 的心脏,驱动自主执行循环:

class ExecutionEngine:
    """Layer 1: 执行与编排引擎"""
    
    def __init__(self, config: HarnessConfig):
        self.model_router = ModelRouter(config.models)
        self.tool_registry = ToolRegistry(config.tools)
        self.state_machine = StateMachine(config.states)
        self.orchestrator = MultiAgentOrchestrator(config.agents)
    
    async def run(self, task: Task) -> ExecutionResult:
        """主执行循环"""
        state = self.state_machine.initialize(task)
        
        while not state.is_terminal():
            # 1. 路由到合适的模型
            model = self.model_router.select(state.context)
            
            # 2. 调用模型获取决策
            decision = await model.generate(
                messages=state.messages,
                tools=self.tool_registry.get_available(state),
                constraints=state.constraints
            )
            
            # 3. 执行工具调用
            if decision.has_tool_calls():
                for tool_call in decision.tool_calls:
                    result = await self.execute_tool(tool_call)
                    state.append_result(tool_call, result)
            
            # 4. 更新状态机
            state = self.state_machine.transition(state, decision)
            
            # 5. 安全检查
            if not self.safety_check(state):
                return ExecutionResult.blocked("Safety violation")
        
        return state.to_result()

Layer 2:上下文管理 —— 认知层

Agent 最大的挑战之一是"记忆":如何在长周期任务中保持连贯性?

class ContextManager:
    """Layer 2: 上下文与轨迹管理"""
    
    def __init__(self, storage: StorageBackend):
        self.storage = storage
        self.compressor = ContextCompressor()
        self.trajectory = TrajectoryTracker()
    
    async def get_context(self, agent_id: str, window: int = 100) -> Context:
        """获取Agent的上下文"""
        # 1. 加载持久化轨迹
        trajectory = await self.trajectory.load(agent_id)
        
        # 2. 加载短期记忆
        short_term = await self.storage.get_short_term(agent_id, window)
        
        # 3. 加载长期记忆(摘要)
        long_term = await self.storage.get_long_term(agent_id)
        
        # 4. 压缩上下文
        compressed = self.compressor.compress(
            trajectory=trajectory,
            short_term=short_term,
            long_term=long_term
        )
        
        return Context(
            messages=compressed.messages,
            state=compressed.state,
            metadata=compressed.metadata
        )
    
    async def save_trajectory(self, agent_id: str, event: Event):
        """保存执行轨迹"""
        await self.trajectory.append(agent_id, event)
        
        # 定期压缩旧轨迹
        if self.trajectory.should_compress(agent_id):
            compressed = self.compressor.compress_trajectory(
                await self.trajectory.get_old(agent_id)
            )
            await self.storage.archive_compressed(agent_id, compressed)

Layer 3:交互层 —— 接口层

Agent 需要与外部世界交互——用户、API、其他服务:

class InteractionSurface:
    """Layer 3: 交互层与执行环境"""
    
    def __init__(self, config: InteractionConfig):
        self.api_gateway = APIGateway(config.api)
        self.event_bus = EventBus(config.events)
        self.ui_adapter = UIAdapter(config.ui)
        self.mcp_client = MCPClient(config.mcp)
    
    async def handle_request(self, request: Request) -> Response:
        """处理外部请求"""
        # 1. 认证与授权
        auth_result = await self.authenticate(request)
        if not auth_result.ok:
            return Response.unauthorized(auth_result.reason)
        
        # 2. 限流检查
        if not self.rate_limiter.allow(auth_result.user_id):
            return Response.too_many_requests()
        
        # 3. 路由到正确的处理链
        handler = self.router.match(request.path)
        
        # 4. 执行并返回
        result = await handler.execute(
            request=request,
            context=auth_result.context
        )
        
        # 5. 发布事件
        await self.event_bus.publish(
            "request.completed",
            {"request_id": request.id, "result": result}
        )
        
        return result

Layer 4:安全与治理 —— 护栏

这是 Harness 中最关键的一层——确保 Agent 不会越界:

class SafetyGuardrails:
    """Layer 4: 安全护栏与可观测性"""
    
    def __init__(self, config: SafetyConfig):
        self.permission_matrix = PermissionMatrix(config.permissions)
        self.sandbox = Sandbox(config.sandbox)
        self.audit_log = AuditLog(config.audit)
        self.human_review = HumanReview(config.review)
    
    async def check(self, action: Action, context: Context) -> SafetyResult:
        """安全检查"""
        # 1. 权限检查
        if not self.permission_matrix.allows(
            agent_id=context.agent_id,
            action=action.type,
            resource=action.target
        ):
            return SafetyResult.blocked("Permission denied")
        
        # 2. 沙箱检查
        if self.sandbox.is_dangerous(action):
            return SafetyResult.need_review(
                reason="Dangerous operation requires human approval",
                action=action
            )
        
        # 3. 合规检查
        compliance = await self.check_compliance(action, context)
        if not compliance.ok:
            return SafetyResult.blocked(compliance.reason)
        
        # 4. 记录审计日志
        await self.audit_log.record(
            agent_id=context.agent_id,
            action=action,
            result="approved"
        )
        
        return SafetyResult.approved()

四、核心组件一:执行与编排引擎

4.1 模型路由器

Harness 的第一个核心组件是模型路由器——决定每次调用使用哪个模型:

class ModelRouter:
    """智能模型路由"""
    
    def __init__(self, models: List[ModelConfig]):
        self.models = {m.name: m for m in models}
        self.cost_tracker = CostTracker()
        self.latency_tracker = LatencyTracker()
    
    def select(self, context: TaskContext) -> Model:
        """根据任务特征选择最优模型"""
        candidates = []
        
        for name, config in self.models.items():
            # 1. 能力匹配
            if not config.capabilities.match(context.required_capabilities):
                continue
            
            # 2. 成本约束
            estimated_cost = self.estimate_cost(config, context)
            if estimated_cost > context.budget:
                continue
            
            # 3. 延迟约束
            estimated_latency = self.latency_tracker.get_p95(name)
            if estimated_latency > context.latency_budget:
                continue
            
            # 4. 综合评分
            score = self.score(
                config=config,
                cost=estimated_cost,
                latency=estimated_latency,
                quality=historical_quality[name]
            )
            
            candidates.append((score, name, config))
        
        if not candidates:
            raise NoModelAvailable("No model meets the constraints")
        
        # 选择得分最高的模型
        candidates.sort(key=lambda x: x[0], reverse=True)
        return self.models[candidates[0][1]]
    
    def score(self, config, cost, latency, quality) -> float:
        """综合评分函数"""
        # 权重可配置
        weights = {
            "quality": 0.5,
            "cost": 0.3,
            "latency": 0.2
        }
        
        return (
            weights["quality"] * quality +
            weights["cost"] * (1 - cost / MAX_COST) +
            weights["latency"] * (1 - latency / MAX_LATENCY)
        )

4.2 状态机引擎

Agent 的执行过程可以被建模为一个有限状态机

from enum import Enum
from dataclasses import dataclass
from typing import Optional, Dict, Any

class AgentState(Enum):
    """Agent 状态定义"""
    IDLE = "idle"
    PLANNING = "planning"
    EXECUTING = "executing"
    WAITING_TOOL = "waiting_tool"
    WAITING_HUMAN = "waiting_human"
    RECOVERING = "recovering"
    COMPLETED = "completed"
    FAILED = "failed"
    BLOCKED = "blocked"

@dataclass
class StateTransition:
    """状态转移规则"""
    from_state: AgentState
    to_state: AgentState
    trigger: str
    guard: Optional[str] = None  # 守卫条件

class StateMachine:
    """Agent 执行状态机"""
    
    def __init__(self):
        self.transitions = self._define_transitions()
        self.current_state = AgentState.IDLE
        self.history = []
    
    def _define_transitions(self) -> Dict[tuple, StateTransition]:
        """定义状态转移规则"""
        return {
            (AgentState.IDLE, "task_received"): StateTransition(
                from_state=AgentState.IDLE,
                to_state=AgentState.PLANNING,
                trigger="task_received"
            ),
            (AgentState.PLANNING, "plan_ready"): StateTransition(
                from_state=AgentState.PLANNING,
                to_state=AgentState.EXECUTING,
                trigger="plan_ready"
            ),
            (AgentState.EXECUTING, "tool_call"): StateTransition(
                from_state=AgentState.EXECUTING,
                to_state=AgentState.WAITING_TOOL,
                trigger="tool_call"
            ),
            (AgentState.WAITING_TOOL, "tool_result"): StateTransition(
                from_state=AgentState.WAITING_TOOL,
                to_state=AgentState.EXECUTING,
                trigger="tool_result"
            ),
            (AgentState.EXECUTING, "need_approval"): StateTransition(
                from_state=AgentState.EXECUTING,
                to_state=AgentState.WAITING_HUMAN,
                trigger="need_approval"
            ),
            (AgentState.WAITING_HUMAN, "approved"): StateTransition(
                from_state=AgentState.WAITING_HUMAN,
                to_state=AgentState.EXECUTING,
                trigger="approved"
            ),
            (AgentState.EXECUTING, "completed"): StateTransition(
                from_state=AgentState.EXECUTING,
                to_state=AgentState.COMPLETED,
                trigger="completed"
            ),
            (AgentState.EXECUTING, "error"): StateTransition(
                from_state=AgentState.EXECUTING,
                to_state=AgentState.RECOVERING,
                trigger="error"
            ),
            (AgentState.RECOVERING, "recovered"): StateTransition(
                from_state=AgentState.RECOVERING,
                to_state=AgentState.EXECUTING,
                trigger="recovered"
            ),
            (AgentState.RECOVERING, "max_retries"): StateTransition(
                from_state=AgentState.RECOVERING,
                to_state=AgentState.FAILED,
                trigger="max_retries"
            ),
        }
    
    def transition(self, trigger: str, context: Dict[str, Any] = None) -> AgentState:
        """执行状态转移"""
        key = (self.current_state, trigger)
        
        if key not in self.transitions:
            raise InvalidTransition(
                f"Cannot transition from {self.current_state} with trigger {trigger}"
            )
        
        transition = self.transitions[key]
        
        # 执行守卫条件检查
        if transition.guard and not self.evaluate_guard(transition.guard, context):
            raise GuardConditionFailed(f"Guard condition failed: {transition.guard}")
        
        # 记录历史
        self.history.append({
            "from": self.current_state,
            "to": transition.to_state,
            "trigger": trigger,
            "timestamp": time.time()
        })
        
        # 执行转移
        old_state = self.current_state
        self.current_state = transition.to_state
        
        # 发布状态变更事件
        self.emit("state.changed", {
            "from": old_state,
            "to": transition.to_state,
            "trigger": trigger
        })
        
        return self.current_state

4.3 多 Agent 编排器

在复杂任务中,往往需要多个 Agent 协作:

class MultiAgentOrchestrator:
    """多Agent编排器"""
    
    def __init__(self, config: OrchestratorConfig):
        self.agents = {a.name: a for a in config.agents}
        self.workflow_engine = WorkflowEngine(config.workflow)
        self.communication = AgentCommunication(config.comm)
    
    async def execute_workflow(self, workflow: Workflow) -> WorkflowResult:
        """执行多Agent工作流"""
        context = WorkflowContext(workflow)
        
        for step in workflow.steps:
            # 1. 确定执行的Agent
            agent = self.select_agent(step, context)
            
            # 2. 准备输入
            input_data = self.prepare_input(step, context)
            
            # 3. 执行
            result = await agent.execute(input_data)
            
            # 4. 处理Agent间通信
            if step.communication:
                await self.communication.send(
                    sender=agent.name,
                    receivers=step.communication.receivers,
                    message=result
                )
            
            # 5. 更新上下文
            context.update(step.name, result)
            
            # 6. 检查是否需要人工干预
            if result.need_human_review:
                review = await self.request_human_review(result)
                if not review.approved:
                    return WorkflowResult.blocked(reason=review.reason)
        
        return context.to_result()
    
    def select_agent(self, step: WorkflowStep, context: WorkflowContext) -> Agent:
        """根据步骤特征选择Agent"""
        # 基于能力匹配 + 负载均衡
        candidates = [
            agent for agent in self.agents.values()
            if agent.can_handle(step.task_type)
        ]
        
        if not candidates:
            raise NoAgentAvailable(f"No agent can handle {step.task_type}")
        
        # 选择负载最低的Agent
        return min(candidates, key=lambda a: a.current_load)

五、核心组件二:上下文与轨迹管理

5.1 三层记忆架构

Agent 的记忆系统是 Harness 中最复杂的部分之一:

┌─────────────────────────────────────────────┐
│           长期记忆 (Long-term)              │
│  ┌─────────────────────────────────────┐    │
│  │  摘要存储 · 知识图谱 · 向量数据库    │    │
│  │  生命周期: 永久 · 压缩率: 高         │    │
│  └─────────────────────────────────────┘    │
├─────────────────────────────────────────────┤
│           短期记忆 (Short-term)             │
│  ┌─────────────────────────────────────┐    │
│  │  最近N轮对话 · 工具调用结果          │    │
│  │  生命周期: 当前任务 · 压缩率: 中     │    │
│  └─────────────────────────────────────┘    │
├─────────────────────────────────────────────┤
│           工作记忆 (Working)                │
│  ┌─────────────────────────────────────┐    │
│  │  当前推理状态 · 临时变量 · 缓存      │    │
│  │  生命周期: 即时 · 压缩率: 无         │    │
│  └─────────────────────────────────────┘    │
└─────────────────────────────────────────────┘

5.2 轨迹持久化

每个 Agent 的执行轨迹都应该被持久化:

@dataclass
class TrajectoryEvent:
    """轨迹事件"""
    event_id: str
    timestamp: float
    agent_id: str
    event_type: str  # "message", "tool_call", "tool_result", "state_change"
    content: Dict[str, Any]
    metadata: Dict[str, Any]
    
    def to_vector(self) -> List[float]:
        """转换为向量表示(用于语义搜索)"""
        text = f"{self.event_type}: {json.dumps(self.content)}"
        return embedding_model.encode(text)

class TrajectoryTracker:
    """轨迹追踪器"""
    
    def __init__(self, storage: StorageBackend, embedding_model):
        self.storage = storage
        self.embedding_model = embedding_model
        self.vector_store = VectorStore(dimension=768)
    
    async def append(self, agent_id: str, event: TrajectoryEvent):
        """追加轨迹事件"""
        # 1. 存储原始事件
        await self.storage.append(agent_id, event)
        
        # 2. 生成向量并存储
        vector = event.to_vector()
        await self.vector_store.upsert(
            id=event.event_id,
            vector=vector,
            metadata={"agent_id": agent_id, "type": event.event_type}
        )
    
    async def search(self, agent_id: str, query: str, top_k: int = 10) -> List[TrajectoryEvent]:
        """语义搜索轨迹"""
        query_vector = self.embedding_model.encode(query)
        
        results = await self.vector_store.search(
            vector=query_vector,
            filter={"agent_id": agent_id},
            top_k=top_k
        )
        
        return [await self.storage.get_event(r.id) for r in results]
    
    async def get_summary(self, agent_id: str) -> str:
        """获取轨迹摘要"""
        recent = await self.storage.get_recent(agent_id, limit=50)
        
        # 使用LLM生成摘要
        summary = await llm.generate(
            messages=[{
                "role": "system",
                "content": "你是一个轨迹摘要助手。请用简洁的语言总结以下Agent执行轨迹的关键信息。"
            }, {
                "role": "user",
                "content": f"轨迹事件:\n{format_events(recent)}"
            }]
        )
        
        return summary

5.3 上下文压缩策略

当轨迹过长时,需要压缩:

class ContextCompressor:
    """上下文压缩器"""
    
    def __init__(self, config: CompressionConfig):
        self.strategies = {
            "sliding_window": SlidingWindowStrategy(config.window_size),
            "summarization": SummarizationStrategy(config.llm),
            "importance_based": ImportanceBasedStrategy(config.threshold),
            "hierarchical": HierarchicalStrategy(config.levels)
        }
    
    def compress(self, trajectory, short_term, long_term) -> CompressedContext:
        """智能压缩上下文"""
        total_tokens = count_tokens(trajectory + short_term + long_term)
        
        if total_tokens <= self.max_tokens:
            # 无需压缩
            return CompressedContext(
                messages=trajectory + short_term,
                summary=long_term
            )
        
        # 选择压缩策略
        strategy = self.select_strategy(total_tokens)
        
        if strategy == "sliding_window":
            compressed = self.strategies["sliding_window"].compress(
                trajectory, short_term, long_term
            )
        elif strategy == "summarization":
            compressed = await self.strategies["summarization"].compress(
                trajectory, short_term, long_term
            )
        elif strategy == "importance_based":
            compressed = self.strategies["importance_based"].compress(
                trajectory, short_term, long_term
            )
        else:
            compressed = await self.strategies["hierarchical"].compress(
                trajectory, short_term, long_term
            )
        
        return compressed
    
    def select_strategy(self, total_tokens: int) -> str:
        """根据上下文大小选择策略"""
        ratio = total_tokens / self.max_tokens
        
        if ratio < 1.5:
            return "sliding_window"
        elif ratio < 3.0:
            return "importance_based"
        elif ratio < 5.0:
            return "summarization"
        else:
            return "hierarchical"

六、核心组件三:交互层与执行环境

6.1 MCP 协议集成

Agent Harness 需要与外部工具交互,而 MCP(Model Context Protocol)是2026年的标准协议:

class MCPClient:
    """MCP协议客户端"""
    
    def __init__(self, config: MCPConfig):
        self.servers = {}
        self.transport = StreamableHTTPTransport(config.transport)
        self.auth = OAuth2Auth(config.auth)
    
    async def connect(self, server_url: str):
        """连接MCP服务器"""
        # 1. 建立连接
        connection = await self.transport.connect(server_url)
        
        # 2. 认证
        token = await self.auth.get_token(server_url)
        await connection.authenticate(token)
        
        # 3. 获取可用工具列表
        tools = await connection.list_tools()
        
        # 4. 注册工具
        self.servers[server_url] = {
            "connection": connection,
            "tools": tools,
            "capabilities": await connection.get_capabilities()
        }
    
    async def call_tool(self, server_url: str, tool_name: str, params: Dict) -> Any:
        """调用MCP工具"""
        server = self.servers.get(server_url)
        if not server:
            raise ServerNotConnected(server_url)
        
        # 检查工具是否存在
        tool = next(
            (t for t in server["tools"] if t.name == tool_name),
            None
        )
        if not tool:
            raise ToolNotFound(tool_name)
        
        # 检查参数
        if not tool.validate_params(params):
            raise InvalidParams(f"Invalid params for {tool_name}")
        
        # 调用工具
        result = await server["connection"].call_tool(tool_name, params)
        
        return result

6.2 沙箱执行环境

Agent 的代码执行需要在隔离的沙箱中:

class SandboxExecutor:
    """沙箱执行器"""
    
    def __init__(self, config: SandboxConfig):
        self.runtime = config.runtime  # "docker", "wasm", "v8"
        self.limits = config.limits
    
    async def execute(self, code: str, language: str) -> ExecutionResult:
        """在沙箱中执行代码"""
        # 1. 创建沙箱
        sandbox = await self.create_sandbox()
        
        try:
            # 2. 设置资源限制
            await sandbox.set_limits(
                memory=self.limits.memory,
                cpu=self.limits.cpu,
                timeout=self.limits.timeout,
                network=self.limits.network
            )
            
            # 3. 执行代码
            result = await sandbox.execute(code, language)
            
            # 4. 收集输出
            stdout = await sandbox.get_stdout()
            stderr = await sandbox.get_stderr()
            
            return ExecutionResult(
                success=result.exit_code == 0,
                stdout=stdout,
                stderr=stderr,
                exit_code=result.exit_code
            )
        
        finally:
            # 5. 清理沙箱
            await sandbox.cleanup()
    
    async def create_sandbox(self) -> Sandbox:
        """创建沙箱实例"""
        if self.runtime == "docker":
            return DockerSandbox(self.config.docker)
        elif self.runtime == "wasm":
            return WASMSandbox(self.config.wasm)
        elif self.runtime == "v8":
            return V8Sandbox(self.config.v8)
        else:
            raise UnsupportedRuntime(self.runtime)

七、核心组件四:安全护栏与可观测性

7.1 权限矩阵

class PermissionMatrix:
    """权限矩阵"""
    
    def __init__(self, config: PermissionConfig):
        self.rules = config.rules
        self.cache = PermissionCache()
    
    def allows(self, agent_id: str, action: str, resource: str) -> bool:
        """检查权限"""
        # 1. 检查缓存
        cache_key = f"{agent_id}:{action}:{resource}"
        cached = self.cache.get(cache_key)
        if cached is not None:
            return cached
        
        # 2. 检查规则
        for rule in self.rules:
            if rule.matches(agent_id, action, resource):
                result = rule.effect == "allow"
                self.cache.set(cache_key, result, ttl=300)
                return result
        
        # 3. 默认拒绝
        return False
    
    def matches(self, rule: PermissionRule, agent_id: str, action: str, resource: str) -> bool:
        """规则匹配"""
        return (
            (rule.agent_pattern == "*" or re.match(rule.agent_pattern, agent_id)) and
            (rule.action_pattern == "*" or re.match(rule.action_pattern, action)) and
            (rule.resource_pattern == "*" or re.match(rule.resource_pattern, resource))
        )

# 示例权限配置
PERMISSION_CONFIG = {
    "rules": [
        # Agent A 可以读取任何资源
        {"agent_pattern": "agent-a", "action_pattern": "read.*", "resource_pattern": "*", "effect": "allow"},
        
        # Agent A 只能写入自己的资源
        {"agent_pattern": "agent-a", "action_pattern": "write.*", "resource_pattern": "agent-a/.*", "effect": "allow"},
        
        # 所有Agent都可以调用搜索工具
        {"agent_pattern": "*", "action_pattern": "tool:search", "resource_pattern": "*", "effect": "allow"},
        
        # 没有任何Agent可以删除生产数据
        {"agent_pattern": "*", "action_pattern": "delete", "resource_pattern": "production/.*", "effect": "deny"},
    ]
}

7.2 人类审批机制

class HumanReview:
    """人类审批机制"""
    
    def __init__(self, config: ReviewConfig):
        self.channels = config.channels  # slack, email, dashboard
        self.timeout = config.timeout
        self.auto_approve_rules = config.auto_approve_rules
    
    async def request_review(self, action: Action, context: Context) -> ReviewResult:
        """请求人类审批"""
        # 1. 检查是否自动批准
        if self.should_auto_approve(action, context):
            return ReviewResult(approved=True, reviewer="auto")
        
        # 2. 创建审批请求
        request = ReviewRequest(
            id=generate_id(),
            action=action,
            context=context,
            created_at=time.time()
        )
        
        # 3. 通知审批人
        await self.notify_reviewers(request)
        
        # 4. 等待审批
        result = await self.wait_for_review(request)
        
        # 5. 记录审批结果
        await self.log_review(request, result)
        
        return result
    
    def should_auto_approve(self, action: Action, context: Context) -> bool:
        """检查是否应该自动批准"""
        for rule in self.auto_approve_rules:
            if rule.matches(action, context):
                return True
        return False
    
    async def wait_for_review(self, request: ReviewRequest) -> ReviewResult:
        """等待审批结果"""
        start_time = time.time()
        
        while time.time() - start_time < self.timeout:
            result = await self.check_review_status(request.id)
            
            if result is not None:
                return result
            
            await asyncio.sleep(1)
        
        # 超时处理
        return ReviewResult(
            approved=False,
            reviewer="timeout",
            reason="Review timed out"
        )

7.3 可观测性:分布式追踪

class ObservabilityLayer:
    """可观测性层"""
    
    def __init__(self, config: ObservabilityConfig):
        self.tracer = Tracer(config.tracing)
        self.metrics = MetricsCollector(config.metrics)
        self.logger = StructuredLogger(config.logging)
        self.cost_tracker = CostTracker(config.costing)
    
    async def trace(self, operation: str, func: Callable) -> Any:
        """分布式追踪"""
        span = self.tracer.start_span(operation)
        
        try:
            # 记录开始
            span.set_attribute("start_time", time.time())
            
            # 执行操作
            result = await func()
            
            # 记录成功
            span.set_attribute("status", "success")
            span.set_attribute("result_type", type(result).__name__)
            
            return result
        
        except Exception as e:
            # 记录失败
            span.set_attribute("status", "error")
            span.set_attribute("error.type", type(e).__name__)
            span.set_attribute("error.message", str(e))
            raise
        
        finally:
            # 关闭span
            span.end()
    
    def record_metric(self, name: str, value: float, tags: Dict[str, str] = None):
        """记录指标"""
        self.metrics.record(name, value, tags)
        
        # 同时记录到日志
        self.logger.info(f"Metric: {name}={value}", extra={"tags": tags})
    
    def track_cost(self, agent_id: str, model: str, tokens: int, operation: str):
        """跟踪成本"""
        cost = self.cost_tracker.calculate(model, tokens)
        
        self.record_metric(
            "agent.cost.total",
            cost,
            {"agent_id": agent_id, "model": model, "operation": operation}
        )
        
        # 检查成本预算
        daily_cost = self.cost_tracker.get_daily_cost(agent_id)
        if daily_cost > self.cost_tracker.get_budget(agent_id):
            self.logger.warning(f"Agent {agent_id} exceeded daily cost budget")

八、代码实战:从零构建一个生产级 Agent Harness

8.1 项目结构

agent-harness/
├── src/
│   ├── core/
│   │   ├── __init__.py
│   │   ├── engine.py          # 执行引擎
│   │   ├── state.py           # 状态机
│   │   └── orchestrator.py    # 编排器
│   ├── context/
│   │   ├── __init__.py
│   │   ├── manager.py         # 上下文管理
│   │   ├── trajectory.py      # 轨迹追踪
│   │   └── compressor.py      # 压缩器
│   ├── interaction/
│   │   ├── __init__.py
│   │   ├── gateway.py         # API网关
│   │   ├── mcp.py             # MCP客户端
│   │   └── sandbox.py         # 沙箱执行
│   ├── safety/
│   │   ├── __init__.py
│   │   ├── permissions.py     # 权限矩阵
│   │   ├── guardrails.py      # 安全护栏
│   │   └── review.py          # 人类审批
│   └── observability/
│       ├── __init__.py
│       ├── tracing.py         # 分布式追踪
│       ├── metrics.py         # 指标收集
│       └── cost.py            # 成本跟踪
├── config/
│   ├── harness.yaml           # 主配置
│   ├── permissions.yaml       # 权限配置
│   └── models.yaml            # 模型配置
├── tests/
├── docker-compose.yml
└── README.md

8.2 主入口:Agent Harness

# src/harness.py
from .core.engine import ExecutionEngine
from .context.manager import ContextManager
from .interaction.gateway import APIGateway
from .safety.guardrails import SafetyGuardrails
from .observability.tracing import ObservabilityLayer

class AgentHarness:
    """Agent Harness 主类"""
    
    def __init__(self, config_path: str):
        # 加载配置
        self.config = self.load_config(config_path)
        
        # 初始化各层
        self.engine = ExecutionEngine(self.config.engine)
        self.context = ContextManager(self.config.context)
        self.gateway = APIGateway(self.config.gateway)
        self.safety = SafetyGuardrails(self.config.safety)
        self.observability = ObservabilityLayer(self.config.observability)
        
        # 状态
        self.agents = {}
        self.running = False
    
    async def start(self):
        """启动Harness"""
        self.running = True
        
        # 启动API网关
        await self.gateway.start()
        
        # 启动事件监听
        await self.start_event_listeners()
        
        self.observability.record_metric("harness.started", 1)
        print("Agent Harness started successfully")
    
    async def stop(self):
        """停止Harness"""
        self.running = False
        
        # 停止所有Agent
        for agent_id, agent in self.agents.items():
            await agent.stop()
        
        # 停止网关
        await self.gateway.stop()
        
        self.observability.record_metric("harness.stopped", 1)
        print("Agent Harness stopped")
    
    async def create_agent(self, agent_config: dict) -> str:
        """创建新Agent"""
        agent_id = generate_id()
        
        # 安全检查
        safety_result = await self.safety.check_create_agent(agent_config)
        if not safety_result.approved:
            raise SafetyViolation(safety_result.reason)
        
        # 创建Agent
        agent = Agent(
            id=agent_id,
            config=agent_config,
            engine=self.engine,
            context=self.context,
            safety=self.safety,
            observability=self.observability
        )
        
        self.agents[agent_id] = agent
        
        self.observability.record_metric("agent.created", 1, {"agent_id": agent_id})
        
        return agent_id
    
    async def execute_task(self, agent_id: str, task: dict) -> dict:
        """执行任务"""
        agent = self.agents.get(agent_id)
        if not agent:
            raise AgentNotFound(agent_id)
        
        # 分布式追踪
        async with self.observability.trace("execute_task") as span:
            span.set_attribute("agent_id", agent_id)
            span.set_attribute("task.type", task.get("type"))
            
            # 获取上下文
            context = await self.context.get_context(agent_id)
            
            # 安全检查
            safety_result = await self.safety.check_task(agent_id, task)
            if not safety_result.approved:
                return {"error": safety_result.reason}
            
            # 执行任务
            result = await agent.execute(task, context)
            
            # 保存轨迹
            await self.context.save_trajectory(agent_id, {
                "task": task,
                "result": result,
                "timestamp": time.time()
            })
            
            # 记录指标
            self.observability.record_metric(
                "task.completed",
                1,
                {"agent_id": agent_id, "task_type": task.get("type")}
            )
            
            return result
    
    async def start_event_listeners(self):
        """启动事件监听器"""
        # 监听Agent状态变更
        self.engine.on("state.changed", self.handle_state_change)
        
        # 监听安全事件
        self.safety.on("safety.violation", self.handle_safety_violation)
        
        # 监听成本事件
        self.observability.on("cost.threshold", self.handle_cost_threshold)
    
    def handle_state_change(self, event):
        """处理状态变更"""
        self.observability.record_metric(
            "agent.state_change",
            1,
            {"agent_id": event.agent_id, "from": event.from_state, "to": event.to_state}
        )
    
    def handle_safety_violation(self, event):
        """处理安全违规"""
        self.observability.record_metric(
            "safety.violation",
            1,
            {"agent_id": event.agent_id, "reason": event.reason}
        )
        
        # 发送告警
        self.send_alert(f"Safety violation by agent {event.agent_id}: {event.reason}")
    
    def handle_cost_threshold(self, event):
        """处理成本阈值"""
        self.observability.record_metric(
            "cost.threshold_exceeded",
            1,
            {"agent_id": event.agent_id, "cost": event.cost}
        )
        
        # 发送告警
        self.send_alert(f"Agent {event.agent_id} exceeded cost threshold: ${event.cost}")

8.3 使用示例

# main.py
import asyncio
from src.harness import AgentHarness

async def main():
    # 1. 初始化Harness
    harness = AgentHarness("config/harness.yaml")
    await harness.start()
    
    # 2. 创建Agent
    agent_id = await harness.create_agent({
        "name": "code-review-agent",
        "model": "gpt-4o",
        "capabilities": ["code_review", "security_scan"],
        "permissions": {
            "tools": ["read_file", "analyze_code"],
            "max_cost_per_task": 0.50,
            "max_tokens_per_task": 100000
        }
    })
    
    # 3. 执行任务
    result = await harness.execute_task(agent_id, {
        "type": "code_review",
        "repo": "myorg/myapp",
        "pr_number": 123,
        "focus": ["security", "performance"]
    })
    
    print(f"Review result: {result}")
    
    # 4. 停止Harness
    await harness.stop()

if __name__ == "__main__":
    asyncio.run(main())

九、性能优化:Harness 层的关键调优策略

9.1 模型调用优化

class ModelCallOptimizer:
    """模型调用优化器"""
    
    def __init__(self, config: OptimizerConfig):
        self.cache = ResponseCache(config.cache)
        self.batcher = RequestBatcher(config.batch)
        self.retry = RetryStrategy(config.retry)
    
    async def optimize_call(self, request: ModelRequest) -> ModelResponse:
        """优化模型调用"""
        # 1. 检查缓存
        cached = await self.cache.get(request)
        if cached:
            return cached
        
        # 2. 批量处理
        if self.batcher.should_batch(request):
            batch_result = await self.batcher.add(request)
            return batch_result
        
        # 3. 带重试的调用
        response = await self.retry.execute(
            func=lambda: self.call_model(request),
            max_retries=3,
            backoff=exponential_backoff
        )
        
        # 4. 缓存结果
        await self.cache.set(request, response, ttl=3600)
        
        return response
    
    async def call_model(self, request: ModelRequest) -> ModelResponse:
        """实际调用模型"""
        start_time = time.time()
        
        response = await self.model_client.generate(
            model=request.model,
            messages=request.messages,
            tools=request.tools,
            max_tokens=request.max_tokens
        )
        
        latency = time.time() - start_time
        
        # 记录指标
        self.record_metrics(request, response, latency)
        
        return response

9.2 上下文压缩优化

class AdaptiveCompressor:
    """自适应上下文压缩器"""
    
    def __init__(self, config: CompressorConfig):
        self.strategies = self.load_strategies(config)
        self.metrics = CompressionMetrics()
    
    async def compress(self, context: Context, target_tokens: int) -> CompressedContext:
        """自适应压缩"""
        current_tokens = context.token_count
        
        if current_tokens <= target_tokens:
            return context
        
        # 计算压缩比
        compression_ratio = target_tokens / current_tokens
        
        # 选择策略
        if compression_ratio > 0.8:
            strategy = "sliding_window"
        elif compression_ratio > 0.5:
            strategy = "importance_based"
        elif compression_ratio > 0.2:
            strategy = "summarization"
        else:
            strategy = "hierarchical"
        
        # 执行压缩
        compressed = await self.strategies[strategy].compress(
            context, target_tokens
        )
        
        # 记录压缩指标
        self.metrics.record(
            strategy=strategy,
            original_tokens=current_tokens,
            compressed_tokens=compressed.token_count,
            compression_ratio=compression_ratio
        )
        
        return compressed

9.3 监控仪表盘

class HarnessDashboard:
    """Harness监控仪表盘"""
    
    def __init__(self, harness: AgentHarness):
        self.harness = harness
        self.metrics = MetricsCollector()
    
    async def get_dashboard_data(self) -> dict:
        """获取仪表盘数据"""
        return {
            "agents": {
                "total": len(self.harness.agents),
                "active": sum(1 for a in self.harness.agents.values() if a.is_active()),
                "idle": sum(1 for a in self.harness.agents.values() if not a.is_active())
            },
            "tasks": {
                "completed_today": await self.metrics.get("tasks.completed.today"),
                "failed_today": await self.metrics.get("tasks.failed.today"),
                "avg_duration": await self.metrics.get("tasks.avg_duration")
            },
            "costs": {
                "total_today": await self.metrics.get("costs.total.today"),
                "by_model": await self.metrics.get_grouped("costs.by_model"),
                "budget_remaining": await self.metrics.get("costs.budget_remaining")
            },
            "performance": {
                "avg_latency_p50": await self.metrics.get("latency.p50"),
                "avg_latency_p95": await self.metrics.get("latency.p95"),
                "avg_latency_p99": await self.metrics.get("latency.p99"),
                "error_rate": await self.metrics.get("error.rate")
            },
            "safety": {
                "violations_today": await self.metrics.get("safety.violations.today"),
                "reviews_pending": await self.metrics.get("safety.reviews.pending"),
                "auto_approved": await self.metrics.get("safety.auto_approved")
            }
        }

十、总结与展望

10.1 Agent Harness 的核心价值

Agent Harness 的出现标志着 AI 工程化进入了一个新阶段:

  1. 从"能用"到"可靠":Harness 提供了生产级的运行时基础设施,让 Agent 从 Demo 走向真实业务。

  2. 从"自由"到"可控":通过权限矩阵、沙箱隔离、人类审批等机制,Harness 确保 Agent 不会越界。

  3. 从"黑盒"到"可审计":分布式追踪、成本归因、轨迹持久化让 Agent 的行为完全透明。

  4. 从"单次"到"长期":三层记忆架构、上下文压缩、状态机让 Agent 能够执行长周期复杂任务。

10.2 未来趋势

  1. Harness as a Service:Harness 将从框架演变为云服务,类似于 Kubernetes 之于容器编排。

  2. 标准化协议:MCP、A2A 等协议将被整合进 Harness,形成统一的 Agent 通信标准。

  3. 自优化 Harness:Harness 本身也将由 AI 驱动,能够根据运行时数据自动优化配置。

  4. 多模态 Harness:随着多模态模型的成熟,Harness 将扩展到支持语音、视频、图像等模态。

10.3 开发者行动指南

如果你正在构建 AI Agent 系统,现在就应该开始关注 Harness:

  1. 评估你的 Agent:它们是否能在生产环境稳定运行?是否有完善的监控和安全机制?

  2. 引入 Harness 层:即使是最简单的 Harness(状态机 + 权限控制),也能大幅提升可靠性。

  3. 建立可观测性:从第一天起就追踪成本、延迟、错误率,这些数据将指导你的优化。

  4. 拥抱标准化:使用 MCP 等标准协议,避免被锁定在特定框架中。

  5. 保持人类在环路中:无论 Harness 多么自动化,关键决策仍需人类审批。


结语:Agent Harness 不是一个时髦的 buzzword,而是 AI 从实验室走向真实世界的必经之路。就像操作系统之于 CPU、数据库之于存储、容器编排之于微服务——Harness 之于 AI Agent,是让强大能力可靠落地的关键基础设施。2026年,每个认真对待 AI 的团队,都应该认真对待 Harness。

参考资源


本文由程序员茄子发布,转载请注明出处。如有技术问题或建议,欢迎在评论区讨论。

推荐文章

PostgreSQL日常运维命令总结分享
2024-11-18 06:58:22 +0800 CST
全栈利器 H3 框架来了!
2025-07-07 17:48:01 +0800 CST
程序员茄子在线接单