Agent Harness 深度拆解:当 AI Agent 决定「给自己套上缰绳」——从裸奔到可控的生产级运行时架构如何重新定义 AI 工程化的终极形态
2026年,AI 工程圈最火的一个词不是某个新模型,不是某个新框架,而是一个古老单词的新用法:Harness(马具)。从 Mitchell Hashimoto 在博客中首次定义,到 OpenAI 发布百万行代码的实验报告,再到 Martin Fowler 的深度分析——几周之内,这个术语成了讨论 AI Agent 开发绕不开的话题。本文将深入拆解 Agent Harness 的架构设计、核心组件与生产级实现,带你看懂这场正在重塑 AI 工程化的核心范式革命。
目录
- 引言:当 Agent 从 Demo 走向生产
- 核心概念:什么是 Agent Harness
- 四层架构模型深度解析
- 核心组件一:执行与编排引擎
- 核心组件二:上下文与轨迹管理
- 核心组件三:交互层与执行环境
- 核心组件四:安全护栏与可观测性
- 代码实战:从零构建一个生产级 Agent Harness
- 性能优化:Harness 层的关键调优策略
- 总结与展望
一、引言:当 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 工程化进入了一个新阶段:
从"能用"到"可靠":Harness 提供了生产级的运行时基础设施,让 Agent 从 Demo 走向真实业务。
从"自由"到"可控":通过权限矩阵、沙箱隔离、人类审批等机制,Harness 确保 Agent 不会越界。
从"黑盒"到"可审计":分布式追踪、成本归因、轨迹持久化让 Agent 的行为完全透明。
从"单次"到"长期":三层记忆架构、上下文压缩、状态机让 Agent 能够执行长周期复杂任务。
10.2 未来趋势
Harness as a Service:Harness 将从框架演变为云服务,类似于 Kubernetes 之于容器编排。
标准化协议:MCP、A2A 等协议将被整合进 Harness,形成统一的 Agent 通信标准。
自优化 Harness:Harness 本身也将由 AI 驱动,能够根据运行时数据自动优化配置。
多模态 Harness:随着多模态模型的成熟,Harness 将扩展到支持语音、视频、图像等模态。
10.3 开发者行动指南
如果你正在构建 AI Agent 系统,现在就应该开始关注 Harness:
评估你的 Agent:它们是否能在生产环境稳定运行?是否有完善的监控和安全机制?
引入 Harness 层:即使是最简单的 Harness(状态机 + 权限控制),也能大幅提升可靠性。
建立可观测性:从第一天起就追踪成本、延迟、错误率,这些数据将指导你的优化。
拥抱标准化:使用 MCP 等标准协议,避免被锁定在特定框架中。
保持人类在环路中:无论 Harness 多么自动化,关键决策仍需人类审批。
结语:Agent Harness 不是一个时髦的 buzzword,而是 AI 从实验室走向真实世界的必经之路。就像操作系统之于 CPU、数据库之于存储、容器编排之于微服务——Harness 之于 AI Agent,是让强大能力可靠落地的关键基础设施。2026年,每个认真对待 AI 的团队,都应该认真对待 Harness。
参考资源:
- Awesome-Agent-Harness — Agent Harness 资源合集
- OpenAI: Harness engineering — Codex Agent 实验报告
- Mitchell Hashimoto: Agent Harness — 概念首次定义
- Martin Fowler: Agent Harness Analysis — 深度分析
- MCP Specification — Model Context Protocol 规范
本文由程序员茄子发布,转载请注明出处。如有技术问题或建议,欢迎在评论区讨论。