| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524 |
- #!/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"<persisted-output>\n完整输出:{path}\nPreview:\n{output[:2000]}\n</persisted-output>"
- 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()
|