Agent学习记录-9

阅读学习时间:约5分钟

笔记

Agent的学习记录--9

来源Github的学习指南Agent Learning Hub的学习:

Stage 3——Harness

当我们在说Agent时,我们指的是Agent 产品 = 模型 + Harness

构建 Harness。 即编写代码,为模型提供一个可操作的环境。Harness 是 agent 在特定领域工作所需要的一切:

Harness = Tools + Knowledge + Observation + Action Interfaces + Permissions

Tools:          文件读写、Shell、网络、数据库、浏览器
Knowledge:      产品文档、领域资料、API 规范、风格指南
Observation:    git diff、错误日志、浏览器状态、传感器数据
Action:         CLI 命令、API 调用、UI 交互
Permissions:    沙箱隔离、审批流程、信任边界

一句话:模型做决策。Harness 执行。模型做推理。Harness 提供上下文。模型是驾驶者。Harness 是载具。对于编程 agent 的 harness 是它的 IDE、终端和文件系统。

对于一个Harness工程,我们所做的是:

  • 实现工具。 给 agent 一双手。文件读写、Shell 执行、API 调用、浏览器控制、数据库查询。每个工具都是 agent 在环境中可以采取的一个行动。设计它们时要原子化、可组合、描述清晰。
  • 策划知识。 给 agent 领域专长。产品文档、架构决策记录、风格指南、合规要求。按需加载,不要前置塞入。Agent 应该知道有什么可用,然后自己拉取所需。
  • 管理上下文。 子 Agent 把明确的工作留在另一份消息列表中;上下文压缩缩短较早的历史;任务系统让目标持久化到单次对话之外。
  • 控制权限。 给 agent 边界。沙箱化文件访问。对破坏性操作要求审批。在 agent 和外部系统之间实施信任边界。这是安全工程与 harness 工程的交汇点。
  • 收集任务过程数据。 Agent 在你的 harness 中执行的每一条行动序列都是训练信号。真实部署中的感知-推理-行动轨迹是微调下一代 agent 模型的原材料。你的 harness 不仅服务于 agent -- 它还可以帮助进化 agent。

拆开Claude Code:其包含以下结构:

Claude Code = 一个 agent loop
            + 工具 (bash, read, write, edit, glob, grep, browser...)
            + 按需 skill 加载
            + 上下文压缩
            + 子 agent 派生
            + 带依赖图的任务系统
            + 异步邮箱的团队协调
            + 任务绑定的 worktree 并行执行
            + 权限治理
            + hooks 扩展系统
            + memory 持久化
            + MCP 外部能力路由

对于一个Harness而言,其有五大核心部件:

  • Agent Loop

  • Tool Registry

  • Permission Gate

  • Session Store

  • Context Compaction

在之前的学习里,我们将所有逻辑(循环、工具字典、硬编码判断)挤在一个单一脚本里,但在实际上,当 Agent 需要执行终端指令、读写用户本地文件时,必须有标准的分层体系

┌─────────────────────────────────────────────────────────────┐
│                    User / Application                       │
└──────────────────────────────┬──────────────────────────────┘
                               │ (query)
                               ▼
┌─────────────────────────────────────────────────────────────┐
│                       Agent Harness                         │
│                                                             │
│   ┌───────────────────┐               ┌──────────────────┐  │
│   │ 5. Compaction     │               │ 4. Session Store │  │
│   │ (微压缩/防爆仓)     │               │ (存盘/轨迹跟踪)    │  │
│   └─────────▲─────────┘               └─────────▲────────┘  │
│             │                                   │           │
│   ┌─────────┴───────────────────────────────────┴────────┐  │
│   │                 1. Agent Loop                        │  │
│   │          (Observe -> Think -> Act)                   │  │
│   └─────────┬───────────────────────────────────┬────────┘  │
│             │                                   │           │
│             ▼                                   ▼           │
│   ┌───────────────────┐               ┌──────────────────┐  │
│   │ 2. Tool Registry  │               │ 3. Permission    │  │
│   │ (工具自省注册)      │ ────────────> │ (安全门禁网关)     │  │
│   └───────────────────┘               └──────────────────┘  │
└─────────────────────────────────────────────────────────────┘

Agent Loop

对于Agent loop而言,其是操作系统内核与事件循环:负责管理生命周期、调用驱动(Tools)、权限沙箱(Permission Gate)、内存分页压缩(Context Compaction)、Token 熔断等。其有许多模式:

  • ReAct模式:即Observe -> Think -> Act
  • Plan-and-Solve模式:先在第一阶段生成一个完整的分解步骤列表(Plan),随后的 Loop 只负责按任务清单顺序调工具执行,并在出错时再触发重规划(Re-plan)。
  • Direct Tool Calling模式:通过大语言模型直接输出结构化的 tool_calls的json,不输出自然语言,这时候loop直接接受参数-> 执行工具 -> 喂回结果
  • Reflexion模式:Loop 中包含专门的评判器(Evaluator)和自我修正(Self-Reflection),循环不仅仅是为了调用工具获取外部数据,还为了评判当前方案是否达标并在 Memory 中沉淀反思教训。
  • State Machine/Workflow-driven Loop模式:循环的状态跳转是由预设的图结构(如 LangGraph)或条件分支决定的,模型的自由度受状态机严格限制。

Tool Registry

其具有管理 Agent 拥有的所有能力,对于我们使用OpenAI协议的,其利用 Python inspect 与 Type Annotations 自动生成标准 Function Calling JSON Schema,且标注工具安全元数据(只读、写操作、高危)和负责安全执行分发与统一异常封箱。

registry.py:

import inspect
import json
from typing import Callable, Dict, Any, List, Optional, get_type_hints

# Python 类型到 JSON Schema 类型的映射
TYPE_MAPPING = {
    str: "string",
    int: "integer",
    float: "number",
    bool: "boolean",
    list: "array",
    dict: "object",
} 

# 工具元数据与可执行体定义
class ToolDefinition:
    def __init__(
        self,
        name: str,
        func: Callable,
        description: str,
        parameters_schema: Dict[str, Any],
        is_destructive: bool = False,
        requires_permission: bool = False,
    ):
        self.name = name
        self.func = func
        self.description = description
        self.parameters_schema = parameters_schema
        self.is_destructive = is_destructive
        self.requires_permission = requires_permission

    def to_openai_schema(self) -> Dict[str, Any]:
        """生成 OpenAI 兼容的 tools schema 格式"""
        return {
            "type": "function",
            "function": {
                "name": self.name,
                "description": self.description,
                "parameters": self.parameters_schema,
            }
        }

# Harness 工具中心注册器
class ToolRegistry:
    def __init__(self):
        self._tools: Dict[str, ToolDefinition] = {}

    def register(
        self,
        name: Optional[str] = None,
        description: Optional[str] = None,
        is_destructive: bool = False,
        requires_permission: bool = False,
    ):
        """装饰器:注册一个 Python 函数为 Agent 工具"""
        def decorator(func: Callable):
            tool_name = name or func.__name__
            tool_desc = description or (func.__doc__ or "No description provided.").strip()
            schema = self._generate_schema_from_func(func)
            
            tool_def = ToolDefinition(
                name=tool_name,
                func=func,
                description=tool_desc,
                parameters_schema=schema,
                is_destructive=is_destructive,
                requires_permission=requires_permission,
            )
            self._tools[tool_name] = tool_def
            return func
        return decorator

    def _generate_schema_from_func(self, func: Callable) -> Dict[str, Any]:
        """将python函数转化为JSON Schema"""
        sig = inspect.signature(func)
        try:
            type_hints = get_type_hints(func)
        except Exception:
            type_hints = {}

        properties = {}
        required = []

        for param_name, param in sig.parameters.items():
            if param_name in ("self", "cls"):
                continue
            
            param_type = type_hints.get(param_name, str)
            json_type = TYPE_MAPPING.get(param_type, "string")
            
            properties[param_name] = {
                "type": json_type,
                "description": f"Parameter: {param_name}",
            }
            
            if param.default == inspect.Parameter.empty:
                required.append(param_name)

        return {
            "type": "object",
            "properties": properties,
            "required": required,
        }

    def get_tool(self, name: str) -> Optional[ToolDefinition]:
        return self._tools.get(name)

    def get_schemas(self) -> List[Dict[str, Any]]:
        """获取供大模型 API 调用的 tools 列表"""
        return [tool.to_openai_schema() for tool in self._tools.values()]

    def execute(self, name: str, arguments: Dict[str, Any]) -> str:
        """安全执行工具:自动捕获异常并转化为格式化自愈文本"""
        tool = self.get_tool(name)
        if not tool:
            return f"[Tool Error]: 未找到名为 '{name}' 的工具,可用工具有: {list(self._tools.keys())}"
        
        try:
            result = tool.func(**arguments)
            return str(result)
        except Exception as e:
            return f"[Tool Execution Error in '{name}']: {type(e).__name__}: {str(e)}"

比如现在输入以下函数:

@registry.register(description="计算给定长宽的矩形面积") #装饰器注册
    def calculate_area(width: float, height: float) -> float:
        return width * height

其会自动转化为:

{
  "type": "object",
  "properties": {
    "width": {
      "type": "number",
      "description": "Parameter: width"
    },
    "height": {
      "type": "number",
      "description": "Parameter: height"
    }
  },
  "required": ["width", "height"]
}

Permission Gate

其职责包括:拦截每一次由大模型决策生成的工具调用请求。其根据输入的Code里的内容判定是否为ALLOW (免审放行)、ASK_USER (需人工授权/交互确认) 还是 DENY (高危直接驳回)。它防止破坏性文件写入、路径穿越逃逸、危险系统命令或未授权操作。

permission.py:

from enum import Enum
from typing import Dict, Any, Optional, Callable
import re

class PermissionLevel(Enum):
    ALLOW = "allow"          # 安全无害操作(如只读文件、提取公式),直接执行
    ASK_USER = "ask_user"    # 潜在副作用(如写入文件、执行新脚本),需要人机交互确认
    DENY = "deny"            # 高危破坏性行为(如删除根目录、格式化、非法穿越),强制阻断

class PermissionGate:
    
    # 高危关键字黑名单(直接 DENY)
    DENY_PATTERNS = [
        r"(?:^|\s)(?:rm|del|rmdir)\s+-[rf]+",  # 递归删除
        r"format\s+[a-zA-Z]:",                 # 磁盘格式化
        r":\(\)\s*\{\s*:\|\:&\s*\};:",         # Fork 炸弹
        r"drop\s+database",                    # 删库
        r"\.\.[\\/]\.\.[\\/]",                 # 连续向上目录穿越逃逸
    ]

    def __init__(
        self,
        interactive: bool = True,
        approval_callback: Optional[Callable[[str, Dict[str, Any]], bool]] = None
    ):
        """
        :param interactive: 是否在终端等待用户人工输入 [y/n] 授权
        :param approval_callback: 自定义审批回调函数(自动化测试或 WebUI 时使用)
        """
        self.interactive = interactive
        self.approval_callback = approval_callback

    def evaluate(self, tool_name: str, tool_args: Dict[str, Any], is_destructive: bool = False, requires_permission: bool = False) -> PermissionLevel: # return permission level (DENY, ASK_USER, ALLOW)
        """评估该次工具调用请求的权限级别"""
        args_str = " ".join(str(v) for v in tool_args.values())

        # 1. 检查是否存在高危黑名单特征 -> 强制 DENY
        for pattern in self.DENY_PATTERNS:
            if re.search(pattern, args_str, re.IGNORECASE):
                return PermissionLevel.DENY

        # 2. 检查工具元数据是否声明了破坏性或需要权限
        if is_destructive:
            return PermissionLevel.DENY
        if requires_permission:
            return PermissionLevel.ASK_USER

        # 3. 常见写操作启发式规则
        if any(keyword in tool_name.lower() for keyword in ["write", "save", "delete", "execute", "run", "bash"]):
            return PermissionLevel.ASK_USER

        return PermissionLevel.ALLOW

    def check_and_authorize(
        self,
        tool_name: str,
        tool_args: Dict[str, Any],
        is_destructive: bool = False,
        requires_permission: bool = False
    ) -> tuple[bool, str]:
        """
        执行门禁检查。
        返回: (是否放行, 原因或诊断信息)
        """
        level = self.evaluate(tool_name, tool_args, is_destructive, requires_permission)

        if level == PermissionLevel.DENY:
            reason = f"[Permission DENIED]: 工具 '{tool_name}' 包含破坏性行为或被安全规则禁止: {tool_args}"
            return False, reason

        if level == PermissionLevel.ALLOW:
            return True, "[Permission ALLOWED]: 只读或安全操作,自动放行。"

        # 此时 level == PermissionLevel.ASK_USER
        if self.approval_callback: # 用于自动化测试时使用
            approved = self.approval_callback(tool_name, tool_args)
            if approved:
                return True, "[Permission APPROVED]: 外部审批回调确认放行。"
            else:
                return False, f"[Permission REJECTED]: 外部审批回调拒绝了工具 '{tool_name}' 的执行。"

        if self.interactive: # 用于用户交互批准
            print(f"\n" + "!" * 55)
            print(f"【安全门禁提示】Agent 请求执行以下操作:")
            print(f"工具名称: {tool_name}")
            print(f"调用参数: {tool_args}")
            print("!" * 55)
            user_choice = input("是否批准执行该操作?[y/N]: ").strip().lower()
            if user_choice in ["y", "yes"]:
                return True, "[Permission APPROVED]: 用户在终端手动确认放行。"
            else:
                return False, f"[Permission REJECTED]: 用户在终端拒绝执行工具 '{tool_name}'。"

        # 如果非交互模式且未提供回调,对 ASK_USER 默认做安全拒绝
        return False, f"[Permission REJECTED]: 非交互模式下未提供审批回调,已安全阻断工具 '{tool_name}'。"

Context Compaction

这个环节是上下文压缩与生命周期管理,其职责在于:

  • 防止多轮长任务中 context window 溢出导致模型报错崩溃
  • Micro-Compaction (微压缩):就地轻量裁剪历史轮次中过长冗余的 Tool Result,保留头尾与核心诊断。
  • Auto-Compaction (自动摘要压缩):当轮次或 Token 预算超限时,自动将前半段历史压缩为快照摘要,保留最近几轮完整会话。

一个很自然的想法就是:直接让模型总结整段历史可以明显缩短上下文,但需要注意的是,这样会丢失细节的同时,还需要多产生一次模型调用。除此之外,还存在一些问题:

  • 模型缺乏自身的“全局 Token 计量感知”: 大模型本身在推理时并不知道自己当前对话到底占了 80k 还是 180k tokens,也没有物理内存计量的概念。
  • 存在致命“死锁与崩溃风险(Fatal Out of Window): 如果把压缩责任推给模型,一旦模型陷入死循环、幻觉,或者连续执行了几次大输出工具(如一次 git diff 吐出 3 万行),上下文瞬间撑爆。下一轮 API 直接返回 400 ContextWindowExceeded 报错,进程直接炸死,模型根本没有机会去调用工具自救
  • 污染模型的行动空间(Action Space): 给模型暴露不相关的系统运维类 Tool,会增加模型选错工具的概率,模型可能会在仅执行两轮任务时就莫名其妙调用了压缩工具。

所以,实际在工程化时,会采用分层压缩架构:

级别 机制 解决的痛点 状态
L1 (入口级) 大输出落盘 + 路径引用 阻止偶发的大体积 Tool Output 突刺 无损(外存化)
L2 (结构级) 消息归档落盘 + 墓碑指针裁剪 控制会话消息的总轮数,建立历史索引 无损(可回溯)
L3 (内存级) 在场消息微压缩 (Micro-compact) 抹去历史思维链、修剪残余冗余格式 低损(仅去噪)
L4 (语义级) 语义快照摘要 (Auto-compact) 极限情况下的终极防爆,由小模型重新提炼主线 强损(但保主线)

compaction.py

第一步我们先做一个大结果的落盘

import os
import json
import hashlib
from datetime import datetime
from typing import List, Dict, Any, Optional, Tuple

try:
    from openai import OpenAI
except ImportError:
    OpenAI = Any
    
class ContextCompactor:
    def __init__(
        self,
        max_tool_output_chars: int = 600,       #历史 tool 返回内容允许的最大字符数
        keep_recent_tools: int = 2,				#保留最近几次 tool 调用的原始输出(豁免微压缩)
        auto_compact_message_limit: int = 14,   #触发消息换页与快照压缩的消息条数上限 k
        tool_spill_threshold_chars: int = 2000, #工具单次超长输出直接落盘的外存绝对阈值
        turn_tool_budget_chars: int = 4000,     #单轮中所有并发工具输出的总预算字符数上限
        max_context_chars: int = 16000,         #全局上下文总物理字符预算上限
        head_keep_count: int = 3,				#换页归档时头部保留的锚点消息数
        storage_dir: str = "context_storage",	#物理落盘与历史归档存储目录
    ):
        self.max_tool_output_chars = max_tool_output_chars
        self.keep_recent_tools = keep_recent_tools
        self.auto_compact_message_limit = auto_compact_message_limit
        self.tool_spill_threshold_chars = tool_spill_threshold_chars
        self.turn_tool_budget_chars = turn_tool_budget_chars
        self.max_context_chars = max_context_chars
        self.head_keep_count = head_keep_count
        self.storage_dir = storage_dir
        self.spills_dir = os.path.join(self.storage_dir, "spills")
        self.archives_dir = os.path.join(self.storage_dir, "archives")
        
    # 将超长内容物理写入磁盘,并格式化返回首尾预览与文件路径指针
    def _write_spill_file(
        self,
        content: str,
        tool_name: str = "tool",
        preview_chars: int = 400
    ) -> Tuple[str, str]:
        os.makedirs(self.spills_dir, exist_ok=True)
        content_hash = hashlib.md5(content.encode("utf-8")).hexdigest()[:8]
        timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
        filename = f"{tool_name}_{timestamp}_{content_hash}.txt"
        file_path = os.path.abspath(os.path.join(self.spills_dir, filename))

        with open(file_path, "w", encoding="utf-8") as f:
            f.write(content)

        # 自适应控制首尾预览长度,确保落盘后格式化文本长度必然显著小于原始文本
        actual_preview = min(preview_chars, max(60, int(len(content) * 0.4)))
        half = max(25, actual_preview // 2)
        head_snippet = content[:half]
        tail_snippet = content[-half:] if half > 0 else ""
        omitted_chars = max(0, len(content) - (len(head_snippet) + len(tail_snippet)))

        spilled_message = (
            f"[Harness Warning: 工具 '{tool_name}' 输出过长 ({len(content)} 字符),已物理落盘保存]\n"
            f"[完整输出归档路径]: {file_path}\n"
            f"--- [首部内容预览 ({len(head_snippet)} 字符)] ---\n"
            f"{head_snippet}\n"
            f"... [中间省略 {omitted_chars} 字符,若需查看完整细节可调用文件读取或搜索工具查阅上述路径] ...\n"
            f"--- [尾部内容预览 ({len(tail_snippet)} 字符)] ---\n"
            f"{tail_snippet}"
        )
        return spilled_message, file_path
        
    # 对于模型返回单个或多个Tool的信息的时候,需要进行信息剪切,以处理上下文。对于多个Tool返回的时候,我们会判断最大的结果,并且从最大的结果开始落盘处理
    def process_tool_results(
        self,
        results: List[Dict[str, Any]],
        preview_chars: int = 400,
        turn_budget: Optional[int] = None,
        single_threshold: Optional[int] = None,
    ) -> List[Dict[str, Any]]:
        budget = turn_budget or self.turn_tool_budget_chars
        threshold = single_threshold or self.tool_spill_threshold_chars
        
        # 初始化提取结果到 items
        items=[]
        for r in results:
            item = dict(r)
            item["_was_spilled"] = False
            item["_spill_path"] = None
            item["_original_chars"] = len(str(item.get("content") or ""))
            items.append(item)
        
        total_chars = sum(item["_original_chars"] for item in items)
        
        while TRUE:
            unspilled = [
                (i,len(str(items[i].get("content") or "")))
                for i in range(len(items))
                if not item[i]["_was_spilled"]
            ]
            if not unspilled:
                break
            
            unspilled.sort(key = lambda x: x[1], reverse = TRUE) #倒序排列
            max_idx, max_len = unspilled[0]
            
            if max_len <= threshold and total_chars <= budget:
                break
            
            target = items[max_idx]
            raw_content = str(target.get("content") or "")
            tool_name = str(target.get("name") or "tool")
            
            spilled_text, file_path = self._write_spill_file(
                raw_content, tool_name=tool_name, preview_chars=preview_chars
            )
            saved = len(raw_content) - len(spilled_text)
            target["content"] = spilled_text
            target["_was_spilled"] = TRUE
            target["_spill_path"] = file_path
            target["_spill_reason"] = "single_item_exceeded" if max_len > threshold else "batch_greedy_spill"
            total_chars -= saved
            
        return items
        
    # 在真实调用时,使用一个轻量适配器
    def spill_tool_output(
        self,
        content: str,
        tool_name: str = "tool",
        preview_chars: int = 400
    ) -> Tuple[str, bool, Optional[str]]:
        res = self.process_tool_results(
            [{"name": tool_name, "content": content}],
            preview_chars=preview_chars
        )[0]
        return res["content"], res["_was_spilled"], res["_spill_path"]

在这里我们完成了最入门的一级,下面进行第二步,做消息裁切,即当消息过多的时候(比方为k条),那么我们只保留前3条和后面的k-3-1条,还有1条用于归档标记,其中写明删去了多少条消息,以及完整记录保存在哪里。这一步控制消息数量,但保留下来的旧消息仍可能包含很长的工具结果。

class ContextCompactor:
    ...
    
    def archive_and_page_messages(
        self,
        messages: List[Dict[str, Any]],
        session_id: Optional[str] = None
    ) -> Tuple[List[Dict[str, Any]], bool, Optional[str]]:
        if len(messages) <= self.auto_compact_message_limit:
            return messages, FALSE, None
        
        os.makedirs(self.archives_dir, exist_ok = TRUE)
        timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
        sess = session_id or "session"
        filename = f"history_archive_{sess}_{timestamp}_{len(messages)}msgs.txt
        file_path = os.path.abspath(os.path.join(self.archives_dir, filename))
        
        with open(file_path, "w", encoding = "utf-8") as f:
            f.write("=== Agent Session History Full Archive ===\n")
            f.write(f"Timestamp: {datetime.now().isoformat()}\n")
            f.write(f"Total Messages: {len(messages)}\n\n")
            for idx, m in emunated(messages):
                f.write(f"--- [Message #{idx} | Role: {m.get('role')}] ---\n")
                if "tool_calls" in m:
                    f.write(f"Tool Calls: {json.dumps(m['tool_calls'], ensure_ascii=False)}\n")
                if "tool_call_id" in m:
                    f.write(f"Tool Call ID: {m['tool_call_id']} | Name: {m.get('name')}\n")
                content = str(m.get("content") or "")
                f.write(f"content \n{content}\n\n")
                
        tail_count = max(2, self.auto_compact_message_limit - self.head_keep_count - 1)
        
        head_end = min(self.head_keep_count, len(messages))
        tail_start = max(head_end, len(messages) - tail_count)
        
        # 如果 head_end 恰好切在含 tool_calls 的 assistant 之后,必须把跟随的 tool 响应拉入头部
        while head_end < tail_start:
            curr_msg = messages[head_end]
            if curr_msg.get("role") == "tool":
                head_end += 1
            else:
                break
        
        # 如果 tail_start 恰好切在 tool 消息上,说明前置的 assistant(tool_calls) 被切入了墓碑,必须将前面的assistant纳入,避免 API 400 报错
        while tail_start > head_end and messages[tail_start].get("role") == "tool":
            tail_start -= 1
        
        # 说明不需要切
        if tail_start <= head_end:
            return messages, False, None
        
        archived_count = tail_start - head_end
        tombstone_msg = {
            "role": "user",
            "content": (
                f"【Harness 上下文历史物理换页墓碑】\n"
                f"由于会话轮次超过阈值 ({len(messages)} > {self.auto_compact_message_limit}),"
                f"系统已将中间阶段的 {archived_count} 条历史交互完整归档至外存。\n"
                f"物理归档快照文件: {file_path}\n"
                f"(注:前序初始目标与最近活跃状态已完整保留。如在后续推理中必须核实被归档轮次的特定数据,"
                f"可使用文件检索工具精准读取上述路径。)"
            )
        }
        
        paged_messages = messages[:head_end] + tombstone_msg + messages[tail_start:]
        return paged_messages, True, file_path

在这两步后,我们还是会估计剩余上下文的大小,现在只有超过字数限制的时候才会执行micro_compact。对于模型已经读取过的结果,它会豁免最近保留的 tool 调用。旧结果被替换前会先完整落盘,因此每个占位都带有可恢复路径:我们现在做第三步:即微压缩

class ContextCompactor:
    ...
    
    def micro_compact(self, messages: List[Dict[str, Any]]) -> Tuple[List[Dict[str, Any]], int]:
        tool_indices = [i for i, m in enumerate(messages) if m.get("role") == "tool"]
        exempt_indices = set(tool_indices[-self.keep_recent_tools:]) if self.keep_recent_tools > 0 else set()
        compacted_messages = []
        total_char_saved = 0
        
        for i, msg in enumerate(messages):
            new_msg = dict(msg)
            if msg.get("role") == "tool" and i not in exempt_indices:
                content = str(msg.get("content") or "")
                if len(content) > self.max_tool_output_chars:
                    head_len = self.max_tool_output_chars // 2 - 40
                    tail_len = self.max_tool_output_chars // 2 - 40
                    omitted_len = len(content) - (head_len + tail_len)
                    
                    compressed_content = (
                    	content[:head_len] +
                        f"\n...[Context Micro-Compacted: {omitted_len} characters truncated by Harness] ...\n" +
                        content[-tail_len:]
                    )
                    total_chars_saved += (len(content) - len(compressed_content))
                    new_msg["content"] = compressed_content
            compacted_messages.append(new_msg)
        return compacted_messages, total_chars_saved

第四步,当我们前三步做完了,但还是超过上下文窗口限制的时候,我们可以拉起辅助模型来提取主线进展与核心结论,即在这里我们要调用API了

class ContextCompactor:
    ...
    
    def auto_compact(
        self,
        messages: List[Dict[str, Any]],
        client: Optional[OpenAI] = None,
        model_name: Optional[str] = None
    ) -> Tuple[List[Dict[str, Any]], bool]:
        if len(messages) <= self.auto_compact_message_limit:
            return messages, False
        
        # 保留最前方的 System Prompt,如果没有则创造一条
        system_msgs = [m for m in messages if m.get("role") == "system"]
        sys_msg = system_msgs[0] if system_msgs else {"role": "system", "content": "You are a helpful assistant."}
        # 保留最近4条消息
        recent_msgs = messages[-4:]
        middle_msgs = [m for m in messages if m not in system_msgs and m not in recent_msgs]
        
        if not middle_msgs:
            return messages, False
        
        # 实际上这样会丢失很多信息,还是使用API调用
        middle_text = ""
        for m in middle_msgs:
            role = m.get("role")
            content = m.get("content") or ""
            middle_text += f"[{role}]: {content[:200]}\n"
        
        if client and model_name:
            try:
                resp = client.chat.completions.create(
                	model = model_name,
                    messages = [
                        {"role": "system", "content": "请用一段话简明扼要总结以下对话和工具调用的核心进展与关键事实:"},
                        {"role": "user", "content": middle_text[:2000]}
                    ],
                    temperature = 0.0
                )
                summary_content = f"【Harness 上下文历史自动快照】\n" + (resp.choices[0].message.content or "").strip()
            except Exception:
                pass
        
        compacted_messages = [
            sys_msg,
            {"role": "user", "content": "summary_content"},
            {"role": "assistant", "content": "收到,我已知晓前序工作进展,我们将基于此快照继续执行。"}
        ] + recent_msgs
        
        return compacted_messages, True

最后我们把这四个压缩上下文函数进行组装:

class ContextCompactor:
    ...
    
    def calculate_total_chars(self, messages: List[Dict[str, Any]]) -> int:
        """计算消息列表的全局物理字符总数"""
        return sum(len(str(m.get("content") or "")) for m in messages)

    def is_overflow(self, messages: List[Dict[str, Any]]) -> bool:
        """
        评估上下文是否突破安全水位:
        1. 消息总轮数超过上限 k (auto_compact_message_limit);
        2. 或全局总字符数超过安全预算 (max_context_chars)。
        只要两者皆未超标,上下文即处于健康安全状态。
        """
        too_many_msgs = len(messages) > self.auto_compact_message_limit
        too_many_chars = self.calculate_total_chars(messages) > self.max_context_chars
        return too_many_msgs or too_many_chars
    
    def compact_pipeline(
    	self,
        messages: List[Dict[str, Any]],
        client: Optional[OpenAI] = None,
        model_name: Optional[str] = None,
        session_id: Optional[str] = None,
    ) -> Tuple[List[Dict[str, Any]], Dict[str, Any]]:
        stats: Dict[str, Any] = {
            "spilled_tools": 0,
            "paged": False,
            "archive_path": None,
            "saved_chars": 0,
            "auto_compacted": False,
            "short_circuited_at": None
        }
        current_msgs = [dict(m) for m in messages]
        # 0. 判断有无单项超长或者上下文累积填满
        has_giant_tool = any(
            m.get("role") == "tool"
            and not m.get("_was_spilled")
            and len(str(m.get("content") or "")) > self.tool_spill_threshold_chars
            and "[完整输出归档路径]:" not in str(m.get("content") or "")
            for m in current_msgs
        )

        if not has_giant_tool and not self.is_overflow(current_msgs):
            stats["short_circuited_at"] = "healthy_initial"
            return current_msgs, stats
        # 1. 大工具落盘
        for m in current_msgs:
            if m.get("role") == "tool" and not m.get("_was_spilled"):
                content_str = str(m.get("content") or "")
                if "[完整输出归档路径]:" not in content_str:
                    spilled_text, was_spilled, _ = self.spill_tool_output(
                        content_str, tool_name=str(m.get("name") or "tool")
                    )
                    if was_spilled:
                        m["content"] = spilled_text
                        m["_was_spilled"] = True
                        stats["spilled_tools"] += 1
        if not self.is_overflow(current_msgs):
            stats["short_circuited_at"] = "level_1_spill"
            return current_msgs, stats
        # 2. 对话切片
        paged_msgs, paged, archive_path = self.archive_and_page_messages(current_msgs, session_id=session_id)
        stats["paged"] = paged
        stats["archive_path"] = archive_path
        current_msgs = paged_msgs
        if not self.is_overflow(current_msgs):
            stats["short_circuited_at"] = "level_2_paging"
            return current_msgs, stats
        # 3. 微压缩
        compacted_msgs, saved_chars = self.micro_compact(current_msgs)
        stats["saved_chars"] = saved_chars
        current_msgs = compacted_msgs
        if not self.is_overflow(current_msgs):
            stats["short_circuited_at"] = "level_3_micro"
            return current_msgs, stats
        # 4. 语义对照
        final_msgs, auto_compacted = self.auto_compact(compacted_msgs, client=client, model_name=model_name)
        stats["auto_compacted"] = auto_compacted
        stats["short_circuited_at"] = "level_4_auto"
        return final_msgs, stats

Session store

这个环节主要是会话存储和轨迹的持久化,其主要用于支持持久化存储到本地磁盘和会话重载与断点续跑。

session.py

import os
import json
import time
from typing import List, Dict, Any, Optional

class SessionStore:
    def __init__(self, session_id: str, storage_dir: str = "sessions"):
        self.session_id = session_id
        self.storage_dir = storage_dir
        os.makedirs(self.storage_dir, exist_ok=True)
        self.session_file = os.path.join(self.storage_dir, f"{session_id}.json")
        self.trace_log_file = os.path.join(self.storage_dir, f"{session_id}_trace.jsonl")
        
    def save_state(self, messages: List[Dict[str, Any]], metadata: Optional[Dict[str, Any]] = None):
        data = {
            "session_id": self.session_id,
            "updated_at": time.strftime("%Y-%m-%d %H:%M:%S"),
            "messages_count": len(messages),
            "metadata": metadata or {},
            "messages": messages,
        }
        with open(self.session_file, "w", encoding = "utf-8") as f:
            json.dump(data, f, ensure_ascii=False, indent=2)
    
    def load_state(self) -> Optional[Dict[str, Any]]:
        if not os.path.exists(elf.session_file):
            return None
        try: 
            with open(self.session_file, "r",encoding = "utf-8") as f:
                return json.load(f)
        except Exception as e:
            print(f"[SessionStore Warning]: 加载会话快照失败: {e}")
            return None	
       
    def append_trace_event(self, event_type: str, payload: Dict[str, Any]):
        record = {
            "timestamp": time.strftime("%Y-%m-%d %H:%M:%S"),
            "session_id": self.session_id,
            "event": event_type,
            "payload": payload,
        }
        with open(self.trace_log_file, "w", encoding = "utf-8") as f:
            f.write(json.dumps(record, ensure_ascii = False) + "\n")