#!/usr/bin/env python3 """ s08_context_compact.py - 上下文压缩 在调用 LLM 前插入四层压缩流水线: L1: snip_compact — 消息数量 > 50 时裁掉中间消息 L2: micro_compact — 用占位符替换旧的 tool_results L3: tool_result_budget — 把大型结果持久化到磁盘 L4: compact_history — LLM 完整摘要(1 次 API 调用) 应急:reactive_compact — 当 API 仍然返回 prompt_too_long 时触发 ┌─────────────────────────────────────────────────────────────┐ │ messages[] │ │ ↓ │ │ L3 budget ─→ L1 snip ─→ L2 micro ─→ [token > threshold?] │ │ ├─ 否 → LLM │ │ └─ 是 → L4 summary │ │ ↓ │ │ LLM 调用 │ │ [prompt_too_long?] │ │ └─ 是 → reactive │ └─────────────────────────────────────────────────────────────┘ 核心原则:先便宜,后昂贵。 执行顺序匹配 CC 源码:budget → snip → micro → auto。 基于 s07(技能加载)构建。用法: python s08_context_compact/code.py 需要: pip install anthropic python-dotenv + .env 中配置 ANTHROPIC_API_KEY """ import ast, json, os, subprocess, time from pathlib import Path try: import readline readline.parse_and_bind('set bind-tty-special-chars off') except ImportError: pass from anthropic import Anthropic from dotenv import load_dotenv load_dotenv(override=True) if os.getenv("ANTHROPIC_BASE_URL"): os.environ.pop("ANTHROPIC_AUTH_TOKEN", None) WORKDIR = Path.cwd() SKILLS_DIR = WORKDIR / "skills" TRANSCRIPT_DIR = WORKDIR / ".transcripts" TOOL_RESULTS_DIR = WORKDIR / ".task_outputs" / "tool-results" client = Anthropic(base_url=os.getenv("ANTHROPIC_BASE_URL")) MODEL = os.environ["MODEL_ID"] CURRENT_TODOS: list[dict] = [] # s07: 技能目录扫描 (继承自 s07) def _parse_frontmatter(text: str) -> tuple[dict, str]: if not text.startswith("---"): return {}, text parts = text.split("---", 2) if len(parts) < 3: return {}, text meta = {} for line in parts[1].strip().splitlines(): if ":" in line: k, v = line.split(":", 1) meta[k.strip()] = v.strip().strip('"').strip("'") return meta, parts[2].strip() SKILL_REGISTRY: dict[str, dict] = {} def _scan_skills(): if not SKILLS_DIR.exists(): return for d in sorted(SKILLS_DIR.iterdir()): if not d.is_dir(): continue manifest = d / "SKILL.md" if manifest.exists(): raw = manifest.read_text() meta, body = _parse_frontmatter(raw) name = meta.get("name", d.name) desc = meta.get("description", raw.split("\n")[0].lstrip("#").strip()) SKILL_REGISTRY[name] = {"name": name, "description": desc, "content": raw} _scan_skills() def list_skills() -> str: if not SKILL_REGISTRY: return "(未找到技能)" return "\n".join(f"- **{s['name']}**: {s['description']}" for s in SKILL_REGISTRY.values()) def load_skill(name: str) -> str: skill = SKILL_REGISTRY.get(name) if not skill: return f"未找到技能:{name}" return skill["content"] # s08: SYSTEM 包含技能目录 (继承自 s07 build_system) def build_system() -> str: catalog = list_skills() return ( f"你是位于 {WORKDIR}. " f"可用技能:\n{catalog}\n" "需要时使用 load_skill 获取完整详情。" ) SYSTEM = build_system() # s08: 子 Agent 使用自己的系统提示词 — 不压缩、不加载技能 SUB_SYSTEM = ( f"你是位于 {WORKDIR}. " "完成交给你的任务,然后返回简洁摘要。" "不要继续委派。" ) # ═══════════════════════════════════════════════════════════ # 来自 s02-s07 (未改动): 基础工具 # ═══════════════════════════════════════════════════════════ def safe_path(p: str) -> Path: path = (WORKDIR / p).resolve() if not path.is_relative_to(WORKDIR): raise ValueError(f"路径逃逸出工作区:{p}") return path def run_bash(command: str) -> str: try: r = subprocess.run(command, shell=True, cwd=WORKDIR, capture_output=True, text=True, timeout=120) out = (r.stdout + r.stderr).strip() return out[:50000] if out else "(无输出)" except subprocess.TimeoutExpired: return "错误:执行超时(120 秒)" def run_read(path: str, limit: int | None = None) -> str: try: lines = safe_path(path).read_text().splitlines() if limit and limit < len(lines): lines = lines[:limit] + [f"... ({len(lines) - limit} 行更多内容)"] return "\n".join(lines) except Exception as e: return f"错误:{e}" def run_write(path: str, content: str) -> str: try: file_path = safe_path(path); file_path.parent.mkdir(parents=True, exist_ok=True) file_path.write_text(content); return f"已写入 {len(content)} 字节到 {path}" except Exception as e: return f"错误:{e}" def run_edit(path: str, old_text: str, new_text: str) -> str: try: file_path = safe_path(path) text = file_path.read_text() if old_text not in text: return f"错误:在文件中未找到目标文本:{path}" file_path.write_text(text.replace(old_text, new_text, 1)) return f"已编辑 {path}" except Exception as e: return f"错误:{e}" def run_glob(pattern: str) -> str: import glob as g try: results = [] for match in g.glob(pattern, root_dir=WORKDIR): if (WORKDIR / match).resolve().is_relative_to(WORKDIR): results.append(match) return "\n".join(results) if results else "(无匹配)" except Exception as e: return f"错误:{e}" def _normalize_todos(todos): if isinstance(todos, str): try: 个待办 = json.loads(todos) except json.JSONDecodeError: try: 个待办 = ast.literal_eval(todos) except (SyntaxError, ValueError): return None, "错误:todos 必须是列表或 JSON 数组字符串" if not isinstance(todos, list): return None, "错误:todos 必须是列表" for i, t in enumerate(todos): if not isinstance(t, dict): return None, f"错误:todos[{i}] 必须是对象" if "content" not in t or "status" not in t: return None, f"错误:todos[{i}] 缺少 'content' 或 'status'" if t["status"] not in ("pending", "in_progress", "completed"): return None, f"错误:todos[{i}] 包含无效状态 '{t['status']}'" return 个待办, None def run_todo_write(todos: list) -> str: global CURRENT_TODOS 个待办, error = _normalize_todos(todos) if error: return error CURRENT_TODOS = 个待办 lines = ["\n\033[33m## 当前任务\033[0m"] for t in CURRENT_TODOS: icon = {"pending": " ", "in_progress": "\033[36m▸\033[0m", "completed": "\033[32m✓\033[0m"}[t["status"]] lines.append(f" [{icon}] {t['content']}") print("\n".join(lines)) return f"已更新 {len(CURRENT_TODOS)} 个任务" def extract_text(content) -> str: if not isinstance(content, list): return str(content) return "\n".join(getattr(b, "text", "") for b in content if getattr(b, "type", None) == "text") # ═══════════════════════════════════════════════════════════ # 来自 s06-s07 (未改动): 子 Agent # ═══════════════════════════════════════════════════════════ SUB_TOOLS = [ {"name": "bash", "description": "运行一条 shell 命令。", "input_schema": {"type": "object", "properties": {"command": {"type": "string"}}, "required": ["command"]}}, {"name": "read_file", "description": "读取文件内容。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}}, "required": ["path"]}}, {"name": "write_file", "description": "向文件写入内容。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "content": {"type": "string"}}, "required": ["path", "content"]}}, {"name": "edit_file", "description": "在文件中替换一次完全匹配的文本。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "old_text": {"type": "string"}, "new_text": {"type": "string"}}, "required": ["path", "old_text", "new_text"]}}, {"name": "glob", "description": "查找匹配 glob 模式的文件。", "input_schema": {"type": "object", "properties": {"pattern": {"type": "string"}}, "required": ["pattern"]}}, ] SUB_HANDLERS = {"bash": run_bash, "read_file": run_read, "write_file": run_write, "edit_file": run_edit, "glob": run_glob} def spawn_subagent(description: str) -> str: print(f"\n\033[35m[子 Agent 已启动]\033[0m") messages = [{"role": "user", "content": description}] for _ in range(30): response = client.messages.create(model=MODEL, system=SUB_SYSTEM, messages=messages, tools=SUB_TOOLS, max_tokens=8000) messages.append({"role": "assistant", "content": response.content}) if response.stop_reason != "tool_use": break results = [] for block in response.content: if block.type == "tool_use": blocked = trigger_hooks("PreToolUse", block) if blocked: results.append({"type": "tool_result", "tool_use_id": block.id, "content": str(blocked)}) continue handler = SUB_HANDLERS.get(block.name) output = handler(**block.input) if handler else f"未知工具:{block.name}" trigger_hooks("PostToolUse", block, output) print(f" \033[90m[sub] {block.name}: {str(output)[:100]}\033[0m") results.append({"type": "tool_result", "tool_use_id": block.id, "content": output}) messages.append({"role": "user", "content": results}) result = extract_text(messages[-1]["content"]) if not result: for msg in reversed(messages): if msg["role"] == "assistant": result = extract_text(msg["content"]) if result: break if not result: result = "子 Agent stopped 等待 30 turns without final answer." print(f"\033[35m[子 Agent 已完成]\033[0m") return result # ═══════════════════════════════════════════════════════════ # 新增于 s08: 四层压缩流水线 # ═══════════════════════════════════════════════════════════ CONTEXT_LIMIT = 50000 KEEP_RECENT = 3 PERSIST_THRESHOLD = 30000 def estimate_size(msgs): return len(str(msgs)) def _block_type(block): return block.get("type") if isinstance(block, dict) else getattr(block, "type", None) def _message_has_tool_use(msg): if msg.get("role") != "assistant": return False content = msg.get("content") if not isinstance(content, list): return False return any(_block_type(block) == "tool_use" for block in content) def _is_tool_result_message(msg): if msg.get("role") != "user": return False content = msg.get("content") if not isinstance(content, list): return False return any(isinstance(block, dict) and block.get("type") == "tool_result" for block in content) # L1: snipCompact — 裁剪中间消息 def snip_compact(messages, max_messages=50): if len(messages) <= max_messages: return messages keep_head, keep_tail = 3, max_messages - 3 head_end, tail_start = keep_head, len(messages) - keep_tail if head_end > 0 and _message_has_tool_use(messages[head_end - 1]): while head_end < len(messages) and _is_tool_result_message(messages[head_end]): head_end += 1 if (tail_start > 0 and tail_start < len(messages) and _is_tool_result_message(messages[tail_start]) and _message_has_tool_use(messages[tail_start - 1])): tail_start -= 1 if head_end >= tail_start: return messages snipped = tail_start - head_end return messages[:head_end] + [{"role": "user", "content": f"[snipped {snipped} messages]"}] + messages[tail_start:] # L2: microCompact — 旧结果占位符 def collect_tool_results(messages): blocks = [] for mi, msg in enumerate(messages): if msg.get("role") != "user" or not isinstance(msg.get("content"), list): continue for bi, block in enumerate(msg["content"]): if isinstance(block, dict) and block.get("type") == "tool_result": blocks.append((mi, bi, block)) return blocks def micro_compact(messages): tool_results = collect_tool_results(messages) if len(tool_results) <= KEEP_RECENT: return messages for _, _, block in tool_results[:-KEEP_RECENT]: if len(block.get("content", "")) > 120: block["content"] = "[早前工具结果已压缩。如有需要请重新运行。]" return messages # L3: toolResultBudget — 将大型结果持久化到磁盘 def persist_large_output(tool_use_id, output): if len(output) <= PERSIST_THRESHOLD: return output TOOL_RESULTS_DIR.mkdir(parents=True, exist_ok=True) path = TOOL_RESULTS_DIR / f"{tool_use_id}.txt" if not path.exists(): path.write_text(output) return f"\n完整输出:{path}\nPreview:\n{output[:2000]}\n" def tool_result_budget(messages, max_bytes=200_000): last = messages[-1] if messages else None if not last or last.get("role") != "user" or not isinstance(last.get("content"), list): return messages blocks = [(i, b) for i, b in enumerate(last["content"]) if isinstance(b, dict) and b.get("type") == "tool_result"] total = sum(len(str(b.get("content", ""))) for _, b in blocks) if total <= max_bytes: return messages ranked = sorted(blocks, key=lambda p: len(str(p[1].get("content", ""))), reverse=True) for _, block in ranked: if total <= max_bytes: break content = str(block.get("content", "")) if len(content) <= PERSIST_THRESHOLD: continue tid = block.get("tool_use_id", "unknown") block["content"] = persist_large_output(tid, content) total = sum(len(str(b.get("content", ""))) for _, b in blocks) return messages # L4: autoCompact — LLM 完整摘要 def write_transcript(messages): TRANSCRIPT_DIR.mkdir(parents=True, exist_ok=True) path = TRANSCRIPT_DIR / f"transcript_{int(time.time())}.jsonl" with path.open("w") as f: for msg in messages: f.write(json.dumps(msg, default=str) + "\n") return path def summarize_history(messages): conversation = json.dumps(messages, default=str)[:80000] prompt = ("总结这段编码 Agent 对话,以便继续工作。\n" "保留:1. 当前目标,2. 关键发现/决策,3. 已读/已改文件," "4. 剩余工作,5. 用户约束。\n保持简洁但具体。\n\n" + conversation) response = client.messages.create(model=MODEL, messages=[{"role": "user", "content": prompt}], max_tokens=2000) return "\n".join( getattr(block, "text", "") for block in response.content if getattr(block, "type", None) == "text").strip() or "(空摘要)" def compact_history(messages): transcript_path = write_transcript(messages) print(f"[对话记录已保存:{transcript_path}]") summary = summarize_history(messages) return [{"role": "user", "content": f"[Compacted]\n\n{summary}"}] # Emergency: reactiveCompact — API 错误时触发 def reactive_compact(messages): transcript = write_transcript(messages) tail_start = max(0, len(messages) - 5) if (tail_start > 0 and tail_start < len(messages) and _is_tool_result_message(messages[tail_start]) and _message_has_tool_use(messages[tail_start - 1])): tail_start -= 1 summary = summarize_history(messages[:tail_start]) return [{"role": "user", "content": f"[Reactive compact]\n\n{summary}"}, *messages[tail_start:]] # ═══════════════════════════════════════════════════════════ # 来自 s07: 工具定义 # ═══════════════════════════════════════════════════════════ TOOLS = [ {"name": "bash", "description": "运行一条 shell 命令。", "input_schema": {"type": "object", "properties": {"command": {"type": "string"}}, "required": ["command"]}}, {"name": "read_file", "description": "读取文件内容。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "limit": {"type": "integer"}}, "required": ["path"]}}, {"name": "write_file", "description": "向文件写入内容。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "content": {"type": "string"}}, "required": ["path", "content"]}}, {"name": "edit_file", "description": "在文件中替换一次完全匹配的文本。", "input_schema": {"type": "object", "properties": {"path": {"type": "string"}, "old_text": {"type": "string"}, "new_text": {"type": "string"}}, "required": ["path", "old_text", "new_text"]}}, {"name": "glob", "description": "查找匹配 glob 模式的文件。", "input_schema": {"type": "object", "properties": {"pattern": {"type": "string"}}, "required": ["pattern"]}}, {"name": "todo_write", "description": "为当前编码会话创建并维护任务清单。", "input_schema": {"type": "object", "properties": {"todos": {"type": "array", "items": {"type": "object", "properties": {"content": {"type": "string"}, "status": {"type": "string", "enum": ["pending", "in_progress", "completed"]}}, "required": ["content", "status"]}}}, "required": ["todos"]}}, {"name": "task", "description": "启动一个子 Agent 处理复杂子任务。只返回最终结论。", "input_schema": {"type": "object", "properties": {"description": {"type": "string"}}, "required": ["description"]}}, {"name": "load_skill", "description": "按名称加载某个技能的完整内容。", "input_schema": {"type": "object", "properties": {"name": {"type": "string"}}, "required": ["name"]}}, # s08 变化: 新的 compact 工具 — 触发 compact_history,而不是空操作 {"name": "compact", "description": "总结早前对话以释放上下文空间。", "input_schema": {"type": "object", "properties": {"focus": {"type": "string"}}}}, ] TOOL_HANDLERS = { "bash": run_bash, "read_file": run_read, "write_file": run_write, "edit_file": run_edit, "glob": run_glob, "todo_write": run_todo_write, "task": spawn_subagent, "load_skill": load_skill, } # 来自 s04 (未改动): Hooks HOOKS = {"PreToolUse": [], "PostToolUse": []} def trigger_hooks(event, *args): for cb in HOOKS[event]: r = cb(*args) if r is not None: return r return None DENY_LIST = ["rm -rf /", "sudo", "shutdown"] def permission_hook(block): if block.name == "bash": for p in DENY_LIST: if p in block.input.get("command", ""): return "权限被拒绝" return None def log_hook(block): print(f"\033[90m[HOOK] {block.name}\033[0m") return None HOOKS["PreToolUse"].append(permission_hook) HOOKS["PreToolUse"].append(log_hook) # ═══════════════════════════════════════════════════════════ # agent_loop — s08 核心:调用 LLM 前运行压缩流水线 # ═══════════════════════════════════════════════════════════ MAX_REACTIVE_RETRIES = 1 # 响应式压缩的重试上限 def agent_loop(messages: list): reactive_retries = 0 while True: # s08 变化: 三个预处理器(0 次 API 调用,便宜的优先) # 顺序匹配 CC 源码:budget → snip → micro messages[:] = tool_result_budget(messages) # L3: 先持久化大型结果 messages[:] = snip_compact(messages) # L1: 裁剪中间部分 messages[:] = micro_compact(messages) # L2: 旧结果占位符 # s08 变化: tokens 仍超过阈值 → LLM 摘要(1 次 API 调用) if estimate_size(messages) > CONTEXT_LIMIT: print("[自动压缩]") messages[:] = compact_history(messages) try: response = client.messages.create(model=MODEL, system=SYSTEM, messages=messages, tools=TOOLS, max_tokens=8000) reactive_retries = 0 # API 调用成功后重置 except Exception as e: if ("prompt_too_long" in str(e).lower() or "token 过多" in str(e).lower()) and reactive_retries < MAX_REACTIVE_RETRIES: print("[响应式压缩]") messages[:] = reactive_compact(messages) reactive_retries += 1 continue raise messages.append({"role": "assistant", "content": response.content}) if response.stop_reason != "tool_use": return results = [] for block in response.content: if block.type != "tool_use": continue print(f"\033[36m> {block.name}\033[0m") # s08: compact 工具触发 compact_history,而不是返回空操作字符串 if block.name == "compact": messages[:] = compact_history(messages) results.append({"type": "tool_result", "tool_use_id": block.id, "content": "[已压缩。对话历史已完成摘要。]"}) messages.append({"role": "user", "content": results}) break # 结束当前轮次,使用压缩后的上下文重新开始 blocked = trigger_hooks("PreToolUse", block) if blocked: results.append({"type": "tool_result", "tool_use_id": block.id, "content": str(blocked)}) continue handler = TOOL_HANDLERS.get(block.name) output = handler(**block.input) if handler else f"未知工具:{block.name}" trigger_hooks("PostToolUse", block, output) print(str(output)[:200]) results.append({"type": "tool_result", "tool_use_id": block.id, "content": str(output)}) else: # 正常路径:没有调用 compact messages.append({"role": "user", "content": results}) continue # 已调用 compact:结果已在上方追加 continue if __name__ == "__main__": print("s08: 上下文压缩 — 四层压缩流水线") print("输入问题,回车发送。输入 q 退出。\n") history = [] while True: try: query = input("\033[36ms08 >> \033[0m") except (EOFError, KeyboardInterrupt): break if query.strip().lower() in ("q", "exit", ""): break history.append({"role": "user", "content": query}) agent_loop(history) for block in history[-1]["content"]: if getattr(block, "type", None) == "text": print(block.text) print()