出处:https://www.bilibili.com/video/BV1jbXKBGECC/
"""spark_agent_learn.py — 从零实现一个完整的 Agent Harness(教学项目 · 星火小队)。
Harness = 让「一次模型调用」变成「能长期干活的智能体」的那层脚手架。本文件把它拆成 10 个部件,
每个部件都能单独替换,建议按序阅读:
Config(密钥/路径) → Tool(工具表) → Hook(拦截) → Memory(记忆 + 压缩) → Skill(技能)
→ Subagent(子代理) → Team(持久队友) → MCP(外部工具) → CLI(入口) → 主循环(双层循环)
运行(不需要 .env,Key 已写死在下方常量里):
python spark_agent_learn.py 进入对话(队长阿联会尊称你为「老板」)
python spark_agent_learn.py --mcp-time 把本文件当 time MCP server 用(stdio 传输)
阅读约定:★ = 这里为什么这样设计;⚠️ = 教学简化 / 已知缺陷 / 生产级做法。"""
from __future__ import annotations
all(__import__('importlib.util').util.find_spec(m) for m in ('anthropic', 'yaml', 'mcp')) or (print('★ 依赖自举:缺少依赖,正在用 pip 安装 anthropic / pyyaml / mcp(首次可能数十秒,请稍候…)', file=__import__('sys').stderr, flush=True), __import__('subprocess').check_call([__import__('sys').executable, '-m', 'pip', 'install', '-q', 'anthropic', 'pyyaml', 'mcp']), print('★ 依赖自举:安装完成,继续启动。', file=__import__('sys').stderr, flush=True)) # ★ 依赖自举:缺哪个装哪个;日志走 stderr,免得污染 MCP 的 stdio 协议流
import os
import asyncio
import re
import json
import subprocess
import sys
import threading
import time
import urllib.request
import yaml
import anthropic
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from html.parser import HTMLParser
from pathlib import Path
from types import SimpleNamespace
from typing import Any
from mcp import ClientSession, StdioServerParameters
from mcp.client.stdio import stdio_client
MCP_SERVER_FLAG = '--mcp-time' # 入口分派:带 --mcp-time 时本文件只当 MCP server,不进入对话
# ─── 内侧能力:本文件用 FastMCP 兼任 time server,把「查时间/查日期」暴露成两个工具 ───
def run_time_mcp_server(): # 内联的 time MCP server(--mcp-time 模式)
"""以 stdio 传输运行内联的 time MCP server(本文件带 --mcp-time 时走这条路)。"""
from mcp.server.fastmcp import FastMCP
server = FastMCP('time')
@server.tool()
def get_current_time() -> str:
"""返回当前本地时间。"""
return datetime.now().strftime('%Y-%m-%d %H:%M:%S')
@server.tool()
def get_current_date() -> str:
"""返回当前本地日期。"""
return datetime.now().strftime('%Y-%m-%d')
server.run()
if MCP_SERVER_FLAG in sys.argv:
run_time_mcp_server()
sys.exit(0)
# ─── 跨平台细节:Windows 控制台默认 GBK,中文会炸;先把自己切到 UTF-8,并带上 stdin ───
def _enable_utf8() -> None: # Windows GBK 控制台下也能正确读写中文
"""把 stdin/stdout/stderr 与 Windows 控制台切到 UTF-8,避免 GBK 下中文读写报错。"""
for stream in (sys.stdin, sys.stdout, sys.stderr):
try:
stream.reconfigure(encoding='utf-8', errors='replace')
except (AttributeError, ValueError):
pass
if os.name == 'nt':
try:
import ctypes
ctypes.windll.kernel32.SetConsoleOutputCP(65001)
ctypes.windll.kernel32.SetConsoleCP(65001)
except Exception:
pass
def _setup_color() -> None: # ★ 分段配色(Spring / Claude 风格):无 TTY / NO_COLOR 自动关闭
"""标签按来源上色、智能体名字单独上色、命令亮白、输出灰、错误红。"""
global _COLOR, _PROMPT
_COLOR = bool(os.environ.get('FORCE_COLOR')) or (sys.stdout.isatty() and not os.environ.get('NO_COLOR'))
_PROMPT = '\033[1;37m你: \033[0m' if _COLOR else '你: '
if os.name == 'nt':
try:
import ctypes
ctypes.windll.kernel32.SetConsoleMode(ctypes.windll.kernel32.GetStdHandle(-11), 7) # 7 = 打开 ANSI 虚拟终端处理
except Exception:
pass
_raw = print
tag_rules = (('[hook:', 36), ('[Agent回答]', 32), ('[执行命令]', 33), ('[命令输出]', 33), ('[MCP', 35), ('[记忆', 96), ('[计划', 96), ('[最终', 96), ('[同事', 94), ('[派遣', 94), ('[子(', 94), ('[队友', 94), ('单文件 Agent', 1))
body_rules = (('[执行命令]', '1;37'), ('[命令输出]', '90'))
line_rules = (('单文件 Agent', 1),)
who = (('小闪', 93), ('阿阅', 92), ('阿搜', 96), ('阿核', 95), ('阿建', 91), ('阿联', 97))
errs = ('错误', 'Error', '失败', '异常', 'Traceback')
def _paint(s):
m = re.match(r'^(\s*)(\[[^\]]+\])(:?\s*)(.*)$', s, re.S)
if not m:
if any(w in s for w in errs):
return f'\033[31m{s}\033[0m'
for n, nc in who:
s = s.replace(n, f'\033[{nc}m{n}\033[0m')
c = next((c for p, c in line_rules if s.startswith(p)), None)
return f'\033[{c}m{s}\033[0m' if c else s
pad, tag, sep, body = m.groups()
c = 31 if any(w in tag for w in errs) else next((c for p, c in tag_rules if tag.startswith(p)), 37)
for n, nc in who:
if n in tag:
tag = tag.replace(n, f'\033[{nc}m{n}\033[{c}m')
break
bc = '31' if c == 31 else next((b for p, b in body_rules if tag.startswith(p)), '0')
return f'{pad}\033[{c}m{tag}\033[0m{sep}\033[{bc}m{body}\033[0m'
def _print(*a, **kw):
if _COLOR and a and isinstance(a[0], str):
a = (_paint(a[0]),) + a[1:]
_raw(*a, **kw)
def _excepthook(t, v, tb): # 未捕获异常:先打一行红的摘要,再交给默认钩子打完整堆栈
_raw(f'\033[31m✗ {t.__name__}: {v}\033[0m', file=sys.stderr)
sys.__excepthook__(t, v, tb)
sys.excepthook = _excepthook
globals()['print'] = _print
_enable_utf8(); _setup_color()
ANTHROPIC_API_KEY = 'sk-xxxxxxxxxxxxxxxxx' # 配置:教学项目把 Key 写死,省掉 .env ⚠️ 分享前请先删掉这一行
ANTHROPIC_BASE_URL = 'https://api.deepseek.com/anthropic' # Anthropic 兼容端点
MODEL = 'deepseek-chat' # 非思考模型,不会把 max_tokens 吃在 thinking 上
client = anthropic.Anthropic(api_key=ANTHROPIC_API_KEY, base_url=ANTHROPIC_BASE_URL)
PROJECT_ROOT = Path(__file__).resolve().parents[0] # 仓库根目录(本文件就在根目录下,runs/ 也落这里)
SKILLS_DIR = PROJECT_ROOT / 'skills'
TEMPLATES_DIR = PROJECT_ROOT / 'templates'
# ─── 会话目录:每次启动新建 runs/<时间戳>/ 放流水·审计·产物;跨会话的画像/长期记忆/队友在 AGENT_HOME ───
def _session_dir() -> Path: # 优先用入口下发的会话目录,否则按时间戳自建
"""返回本次会话的产物目录(每次启动都是新的一个)。"""
env = os.environ.get('AGENT_SESSION_DIR')
if env:
return Path(env)
import time
runs_root = Path(os.environ.get('AGENT_RUNS_DIR') or PROJECT_ROOT / 'runs')
d = runs_root / time.strftime('%Y%m%d-%H%M%S')
d.mkdir(parents=True, exist_ok=True)
return d
SESSION_DIR = _session_dir() # 本次会话的目录 runs/<时间戳>/:流水、审计、产物都是会话级
AGENT_HOME = Path(os.environ.get('AGENT_HOME') or PROJECT_ROOT / 'agent_home') # ★ 全局记忆区(跨会话):画像 + 长期记忆 + 队友
EPISODES_DIR = AGENT_HOME / 'episodes' # 当日情景记忆:每天一个 .md,同一天的多次会话共享
FILES_DIR = SESSION_DIR / 'files' # ★ 产物目录:agent 生成的所有文件都落这里,不许散落到仓库
FILES_DIR.mkdir(parents=True, exist_ok=True) # 启动就建好,避免每个工具各自 mkdir
RECENT_MESSAGES = 10 # 压缩后保留最近多少条消息
COMPACT_AFTER_MESSAGES = int(os.environ.get('AGENT_MEMORY_COMPACT_AFTER', '18')) # 消息数超过它才触发压缩
COMPACT_PROMPT_PATH = TEMPLATES_DIR / 'agent' / 'compact_prompt.md' # 压缩提示词(可选文件,存在就优先用它)
COMPACT_PROMPT_FALLBACK = """你是一个记忆整理员。你的任务是把一段被压缩的历史对话浓缩成三种产物,
同时保证已有的长期记忆和用户档案不被覆盖或丢失。
<old_conversation>
{old_conversation}
</old_conversation>
<current_memory>
{current_memory}
</current_memory>
<current_user>
{current_user}
</current_user>
<today_episode_so_far>
{today_episode}
</today_episode_so_far>
请严格产出以下三段 XML,缺一不可:
<episode>
追加到今天情景记忆文件的一段,格式:
## {now_hhmm} 段落小标题
- 关键事件 / 用户请求 1
- 做出的决策 / 产出 2
- 心得或未解问题 3
控制在 200 字以内。
</episode>
<updated_memory>
MEMORY.md 的**完整新版本**。严格要求:
1. 保留 <current_memory> 中所有仍然有效的条目;不要删除你不理解的内容
2. 从 <old_conversation> 中提炼"当前核心目标 / 未完成任务 / 关键事实"合并进来
3. 合并重复、归并同类
4. 只在明显过时的情况下才删除条目
5. 整份文本不超过 3000 字
</updated_memory>
<updated_user>
USER.md 的**完整新版本**。严格要求:
- 仅当 <old_conversation> 中有明确的用户偏好信号(如"我喜欢/不喜欢/以后都这样做")时才修改对应章节
- 否则原样返回 <current_user> 的内容
</updated_user>"""
MCP_CONFIG_PATH = PROJECT_ROOT / 'mcp_servers.json'
AUDIT_FILE = SESSION_DIR / '.hooks_audit.jsonl' # 审计流水,ToolAuditHook 追加写
TEAM_DIR = AGENT_HOME / '.team' # ★ 队友状态也全局:固定队友跨会话还在,offline 可被 spawn 唤回
INBOX_DIR = TEAM_DIR / 'inbox' # 队友收件箱目录
VALID_MSG_TYPES = {'message', 'broadcast', 'shutdown_request', 'shutdown_response', 'plan_approval_response'}
RUNTIME_STATUSES = {'idle', 'working'}
TERMINAL_STATUSES = {'offline', 'shutdown'}
# ─── 第 1 部分 · Skill:渐进式披露 —— 启动只读 name+description,模型要用时才 load_skill 正文 ───
class SkillLoader:
"""扫描 skills/ 目录下的 SKILL.md,把 YAML frontmatter + 正文按需注入对话。"""
def __init__(self, skills_dir: Path): # 只记路径,真正的扫描在 _load_all
self.skills_dir = skills_dir
self.skills = {}
self._load_all()
def _load_all(self): # 遍历 skills/ 下每个 SKILL.md
if not self.skills_dir.exists():
return
for f in sorted(self.skills_dir.rglob('SKILL.md')):
text = f.read_text(encoding='utf-8')
meta, body = self._parse_frontmatter(text)
name = meta.get('name', f.parent.name)
self.skills[name] = {'meta': meta, 'body': body, 'path': str(f)}
def _parse_frontmatter(self, text: str) -> tuple: # 拆出 YAML frontmatter 与 Markdown 正文
match = re.match('^---\\n(.*?)\\n---\\n(.*)', text, re.DOTALL)
if not match:
return ({}, text)
try:
meta = yaml.safe_load(match.group(1)) or {}
except yaml.YAMLError:
meta = {}
return (meta, match.group(2).strip())
def get_descriptions(self) -> str: # 给模型的技能清单:只有名字+描述(渐进式披露第一层)
if not self.skills:
return '(no skills available)'
lines = []
for name, skill in self.skills.items():
desc = skill['meta'].get('description', 'No description')
tags = skill['meta'].get('tags', '')
line = f' - {name}: {desc}'
if tags:
line += f' [{tags}]'
lines.append(line)
return '\n'.join(lines)
def get_content(self, name: str) -> str: # 按需取技能正文(第二层)
skill = self.skills.get(name)
if not skill:
return f"Error: Unknown skill '{name}'. Available: {', '.join(self.skills.keys())}"
return f'''<skill name="{name}">\n{skill['body']}\n</skill>'''
SKILL_LOADER = SkillLoader(SKILLS_DIR) # 启动时一次性扫完 skills/*/SKILL.md
# ─── 第 2 部分 · Memory:三层记忆(长期事实 / 用户画像 / 当日情景,全部跨会话)+ 会话级流水 history.jsonl ───
class MemoryStore:
"""三层记忆的落盘封装。"""
def __init__(self, home_dir: Path, session_dir: Path): # ★ 全局(画像/长期记忆/情景) + 会话(原始流水)
self.memory_dir = home_dir
self.memory_file = home_dir / 'MEMORY.md'
self.history_file = session_dir / 'history.jsonl'
self.user_file = home_dir / 'USER.md'
def ensure_files(self): # 幂等补齐文件;画像首次从 templates/USER.md 播种
self.memory_dir.mkdir(parents=True, exist_ok=True)
self.history_file.parent.mkdir(parents=True, exist_ok=True)
EPISODES_DIR.mkdir(parents=True, exist_ok=True) # 情景记忆也是全局的
if not self.memory_file.exists():
self.memory_file.write_text('# Long-term Memory\n\n', encoding='utf-8')
if not self.user_file.exists():
seed = TEMPLATES_DIR / 'USER.md'
self.user_file.write_text(seed.read_text(encoding='utf-8') if seed.exists() else '# User Profile\n\n', encoding='utf-8')
if not self.history_file.exists():
self.history_file.touch()
def append_history(self, message: dict): # 追加一条流水到 history.jsonl
self.ensure_files()
record = {'ts': datetime.now().isoformat(timespec='seconds'), 'role': message.get('role'), 'content': _json_safe(message.get('content'))}
with self.history_file.open('a', encoding='utf-8') as f:
f.write(json.dumps(record, ensure_ascii=False) + '\n')
def read_memory(self) -> str: # 读长期记忆(每轮注入 system prompt)
self.ensure_files()
return self.memory_file.read_text(encoding='utf-8')
def write_memory(self, text: str): # 压缩时整体覆盖(全量,非追加)
self.ensure_files()
self.memory_file.write_text(text.strip() + '\n', encoding='utf-8')
def read_user(self) -> str: # 读用户画像
self.ensure_files()
return self.user_file.read_text(encoding='utf-8')
def write_user(self, text: str): # 压缩时整体覆盖画像
self.ensure_files()
self.user_file.write_text(text.strip() + '\n', encoding='utf-8')
def today_episode_file(self) -> Path: # 今天的 <<YYYY-MM-DD>>.md 路径,调用时才取日期
return EPISODES_DIR / f'{datetime.now():%Y-%m-%d}.md' # ★ 全局:同一天的多次会话共享
def read_today_episode(self) -> str: # 读今天的情景记忆,没有就返回空串
self.ensure_files()
path = self.today_episode_file()
return path.read_text(encoding='utf-8') if path.exists() else ''
def append_episode(self, text: str): # 情景记忆是日记:追加而不是覆盖
self.ensure_files()
with self.today_episode_file().open('a', encoding='utf-8') as f:
f.write('\n' + text.strip() + '\n')
MEMORY = MemoryStore(AGENT_HOME, SESSION_DIR / 'memory') # 全局记忆 + 会话流水
# ─── 压缩的四个零件之一:把 SDK 的 pydantic 对象递归降级成可 json.dumps 的结构 ───
def _json_safe(value): # 把 SDK 对象递归降级成可 json.dumps 的结构
"""把 SDK 对象(如 tool_use block)递归降级成可 json.dumps 的普通结构。"""
if value is None or isinstance(value, (str, int, float, bool)):
return value
if isinstance(value, list):
return [_json_safe(v) for v in value]
if isinstance(value, dict):
return {k: _json_safe(v) for k, v in value.items()}
if hasattr(value, 'model_dump'):
return _json_safe(value.model_dump())
return str(value)
# ─── 拍平:用 json.dumps 而不是 str(),才保得住 role/类型/工具调用的边界 ───
def _messages_to_text(messages: list[dict]) -> str: # 把旧消息拍平成 JSON 文本喂给压缩模型
"""把一段 messages 拍平成"每行一条"的纯文本,喂给压缩模型当输入。"""
lines = []
for msg in messages:
role = msg.get('role', '?')
content = json.dumps(_json_safe(msg.get('content')), ensure_ascii=False)
lines.append(f'{role}: {content}')
return '\n'.join(lines)
# ─── 抠标签:⚠️ 正则很脆,模型格式一跑偏就静默不更新记忆(见函数内说明) ───
def _extract_tag(text: str, tag: str) -> str: # 从模型输出里抠 <tag>...</tag> 的内容
"""从模型输出里抠出 `<tag>...</tag>` 之间的内容;没有就返回空串。"""
match = re.search(f'<{tag}>(.*?)</{tag}>', text, re.DOTALL)
return match.group(1).strip() if match else ''
def _is_user_turn(message: dict) -> bool: # 是不是「真实用户输入」而不是 tool_result 回灌
'''压缩切点必须落在真实用户轮上,否则会留下孤立的 tool_result → API 直接 400。'''
if message.get('role') != 'user':
return False
content = message.get('content')
if isinstance(content, str):
return True
return isinstance(content, list) and not any(isinstance(b, dict) and b.get('type') == 'tool_result' for b in content)
def _valid_history(history: list[dict]) -> list[dict]: # ★ 发请求前自检:首条必须是真实用户轮
'''兜底修复:history 开头若不是真实用户轮(被切坏了),丢掉前面那些孤立消息。'''
start = next((i for i, m in enumerate(history) if _is_user_turn(m)), None)
if not start:
return history
print(f'[历史已修复] 丢弃开头的 {start} 条孤立消息(首条不是真实用户输入,API 会报 400)')
return history[start:]
# ─── 第 3 部分 · Compact:超阈值就把旧对话压成 episode / updated_memory / updated_user 三段落盘 ───
def compact_history(history: list[dict]) -> list[dict]: # 超阈值时压缩旧对话:episode/MEMORY/USER 三段落盘
"""超过阈值时压缩旧对话,返回压缩后的 history。"""
if len(history) <= COMPACT_AFTER_MESSAGES:
return history
start = max(0, len(history) - RECENT_MESSAGES)
while start > 0 and not _is_user_turn(history[start]): # ★ 切点回退到真实用户轮的边界,绝不切在工具对中间
start -= 1
old_messages = history[:start]
recent_messages = history[start:]
if not old_messages or not recent_messages:
return history
prompt_template = COMPACT_PROMPT_PATH.read_text(encoding='utf-8') if COMPACT_PROMPT_PATH.exists() else COMPACT_PROMPT_FALLBACK
prompt = prompt_template.format(old_conversation=_messages_to_text(old_messages), current_memory=MEMORY.read_memory(), current_user=MEMORY.read_user(), today_episode=MEMORY.read_today_episode(), now_hhmm=datetime.now().strftime('%H:%M'))
try:
message = client.messages.create(model=MODEL, max_tokens=3000, system='你是记忆整理员。请严格按要求输出 XML,不要输出额外解释。', messages=[{'role': 'user', 'content': prompt}])
text = next((b.text for b in message.content if b.type == 'text'), '')
except Exception as exc:
print(f'[记忆压缩失败,保留完整 history]: {exc}')
return history
episode = _extract_tag(text, 'episode')
updated_memory = _extract_tag(text, 'updated_memory')
updated_user = _extract_tag(text, 'updated_user')
if episode:
MEMORY.append_episode(episode)
if updated_memory:
MEMORY.write_memory(updated_memory)
if updated_user:
MEMORY.write_user(updated_user)
print(f'[记忆已压缩]: old={len(old_messages)} recent={len(recent_messages)}')
return recent_messages
# ─── 第 4 部分 · Web fetch:手写 HTML 提纯(HTMLParser),不引第三方解析库 ───
class _TextExtractor(HTMLParser):
"""极简 HTML → 正文提取器:丢掉 script/style,块级标签转换行。"""
def __init__(self): # 累积正文碎片
super().__init__()
self._parts = []
self._skip = False
def handle_starttag(self, tag, attrs): # 进入 script/style 时打开跳过开关
if tag in ('script', 'style'):
self._skip = True
def handle_endtag(self, tag): # 离开 script/style 时关掉跳过开关
if tag in ('script', 'style'):
self._skip = False
if tag in ('p', 'br', 'div', 'li', 'tr', 'h1', 'h2', 'h3', 'h4'):
self._parts.append('\n')
def handle_data(self, data): # 非跳过状态才收正文
if not self._skip:
self._parts.append(data)
def get_text(self): # 拼出最终纯文本
return re.sub('\\n{3,}', '\n\n', ''.join(self._parts)).strip()
def web_fetch(url: str, extract_mode: str='text', max_chars: int=8000) -> str:
"""抓网页并截断到 max_chars。"""
req = urllib.request.Request(url, headers={'User-Agent': 'Mozilla/5.0'})
try:
with urllib.request.urlopen(req, timeout=10) as resp:
raw = resp.read().decode('utf-8', errors='replace')
except Exception as e:
return f'Error fetching {url}: {e}'
if extract_mode == 'text':
parser = _TextExtractor()
parser.feed(raw)
text = parser.get_text()
else:
text = raw
return text[:max_chars]
TODOS: list[dict] = []
VALID_STATUS = {'pending', 'in_progress', 'completed'}
STATUS_ICON = {'pending': '[ ]', 'in_progress': '[~]', 'completed': '[x]'}
# ─── 计划渲染:把 todolist 画成终端可读的清单,模型与人都看得懂 ───
def render_todos(todos: list[dict]) -> str: # 把 todolist 渲染成终端可读的清单
"""把 todos 渲染成终端里可读的多行清单。"""
if not todos:
return '(当前无待办事项)'
lines = []
for t in todos:
icon = STATUS_ICON.get(t.get('status', 'pending'), '[?]')
lines.append(f" {icon} {t.get('id')}. {t.get('content', '')}")
return '\n'.join(lines)
# ─── todolist 是全量覆盖:模型每轮都要给完整列表,避免增量合并产生歧义 ───
def update_todos(todos: list[dict]) -> str: # 全量覆盖 todos(模型每次都给完整列表)
"""全量覆盖式更新计划;同一时间只允许一个 in_progress。"""
global TODOS
cleaned = []
for i, t in enumerate(todos, start=1):
content = (t.get('content') or '').strip()
if not content:
continue
status = t.get('status', 'pending')
if status not in VALID_STATUS:
status = 'pending'
cleaned.append({'id': t.get('id', i), 'content': content, 'status': status})
in_progress = [t for t in cleaned if t['status'] == 'in_progress']
if len(in_progress) > 1:
return 'Error: 同一时间只能有一个 in_progress 任务,请重新规划。'
TODOS = cleaned
print('\n[计划已更新]')
print(render_todos(TODOS))
print()
pending = [t for t in TODOS if t['status'] == 'pending']
done = [t for t in TODOS if t['status'] == 'completed']
summary = f'todos updated: total={len(TODOS)}, completed={len(done)}, in_progress={len(in_progress)}, pending={len(pending)}'
return summary + '\n\n当前列表:\n' + render_todos(TODOS)
def _resolve_path(path: str) -> Path: # ★ 写路径归一:相对路径以产物目录为根,绝对/越界只取文件名
'''把模型给的写路径映射进产物目录 runs/<会话>/files/,保证不往仓库里散落文件。'''
p = Path(str(path).strip().replace('\\', '/'))
if p.is_absolute() or p.drive or '..' in p.parts:
return FILES_DIR / p.name
return FILES_DIR.joinpath(*[x for x in p.parts if x not in ('', '.')])
def _read_path(path: str) -> Path: # 读路径:产物目录优先,找不到再按原路径(读仓库代码照常可以)
'''读文件用:先在产物目录里找,找不到就当普通相对/绝对路径处理。'''
p = Path(str(path).strip().replace('\\', '/'))
if p.is_absolute() or p.drive or '..' in p.parts:
return p
candidate = FILES_DIR.joinpath(*[x for x in p.parts if x not in ('', '.')])
return candidate if candidate.exists() else p
# ─── 第 6 部分 · Tool:裸执行器,不含任何 Hook ⚠️ 子代理与队友直连它,所以产物策略放在这一层兜底 ───
def execute_basic_tool(block, prefix: str='') -> str: # 裸执行器:子代理与队友直连它,会绕过 Hook
"""处理基础工具,返回字符串内容。"""
if block.name == 'web_fetch':
url = block.input['url']
mode = block.input.get('extract_mode', 'text')
max_chars = block.input.get('max_chars', 8000)
print(f' [{prefix}网页获取]: {url}')
return web_fetch(url, mode, max_chars)
if block.name == 'run_command':
command = block.input['command']
print(f' [{prefix}执行命令]: {command}')
result = subprocess.run(command, shell=True, cwd=FILES_DIR, capture_output=True, stdin=subprocess.DEVNULL) # ★ 工作目录=产物目录;不吃 text=True,免得 GBK 解码在读取线程里炸
raw = result.stdout or result.stderr
try: output = raw.decode('utf-8') # ★ Windows 命令输出多是 ANSI/GBK:先按 UTF-8 试,失败按系统 ANSI 兜底替换
except UnicodeDecodeError: output = raw.decode('mbcs' if os.name == 'nt' else 'utf-8', 'replace')
print(f' [{prefix}命令输出]: {output.strip()[:200]}')
return output
if block.name == 'load_skill':
skill_name = block.input['skill_name']
print(f' [{prefix}加载技能]: {skill_name}')
return SKILL_LOADER.get_content(skill_name)
if block.name == 'read_file':
path = block.input['path']
print(f' [{prefix}读取文件]: {path}')
try:
return _read_path(path).read_text(encoding='utf-8')
except Exception as e:
return f'Error reading {path}: {e}'
if block.name == 'write_file':
path = block.input['path']
content = block.input['content']
print(f' [{prefix}写入文件]: {path}')
try:
p = _resolve_path(path)
p.parent.mkdir(parents=True, exist_ok=True)
p.write_text(content, encoding='utf-8')
return f'写入成功: {p}'
except Exception as e:
return f'Error writing {path}: {e}'
if block.name == 'glob':
pattern = block.input['pattern']
print(f' [{prefix}文件搜索]: {pattern}')
p = Path(str(pattern)) # ★ 模型可能直接把绝对路径当模式(系统提示里给了产物绝对路径)
base, pat = (p.parent, p.name) if (p.is_absolute() or p.drive) else (FILES_DIR, str(pattern))
try:
matches = sorted(str(x) for x in base.glob(pat))
if not matches and base == FILES_DIR:
matches = sorted(str(x) for x in Path('.').glob(pat))
except Exception as exc: # 非法模式只返回错误文本,绝不弄死会话
return f'Error: glob 模式不合法:{exc}'
return '\n'.join(matches) if matches else '(无匹配)'
if block.name == 'grep':
pattern = block.input['pattern']
path = block.input.get('path', '.')
print(f' [{prefix}内容搜索]: {pattern} in {path}')
try:
regex = re.compile(pattern)
except re.error as exc:
return f'Error: 正则表达式不合法:{exc}'
root = _read_path(path)
if not root.exists():
return f'Error: 路径不存在:{path}'
targets = [root] if root.is_file() else sorted((p for p in root.rglob('*') if p.is_file() and p.suffix in ('.py', '.md')))
hits = []
for target in targets:
try:
lines = target.read_text(encoding='utf-8', errors='replace').splitlines()
except OSError:
continue
for lineno, line in enumerate(lines, 1):
if regex.search(line):
hits.append(f'{target}:{lineno}:{line}')
if len(hits) >= 200:
hits.append('...(匹配过多,已截断到 200 条)')
return '\n'.join(hits)
return '\n'.join(hits) if hits else '(无匹配)'
return f"Error: Unknown tool '{block.name}'"
def _safe_tool(block, prefix: str='') -> str: # ★ 工具层最后一道防线:意外异常转成文本,不中断整个 Agent
'''把工具异常包成 Error 字符串返回给模型;宁可让模型看到错误,也不要让会话崩掉。'''
try:
return execute_basic_tool(block, prefix=prefix)
except Exception as exc:
return f'Error: 工具 {getattr(block, "name", "?")} 执行异常:{type(exc).__name__}: {exc}'
_TOOL_SCHEMAS: dict[str, dict] = {'run_command': {'name': 'run_command', 'description': '在终端执行一条命令并返回输出。创建或覆盖文件必须调用 write_file。', 'input_schema': {'type': 'object', 'properties': {'command': {'type': 'string', 'description': '要执行的 shell 命令'}}, 'required': ['command']}}, 'web_fetch': {'name': 'web_fetch', 'description': '获取指定 URL 的网页内容,支持文本提取模式', 'input_schema': {'type': 'object', 'properties': {'url': {'type': 'string', 'description': '要访问的完整 URL'}, 'extract_mode': {'type': 'string', 'description': '提取模式:text(纯文本,默认)或 raw(原始 HTML)'}, 'max_chars': {'type': 'integer', 'description': '最大返回字符数,默认 8000'}}, 'required': ['url']}}, 'load_skill': {'name': 'load_skill', 'description': '加载指定技能的详细知识内容,在回答相关问题前调用', 'input_schema': {'type': 'object', 'properties': {'skill_name': {'type': 'string', 'description': '技能名称,必须是系统提示中列出的可用技能之一'}}, 'required': ['skill_name']}}, 'read_file': {'name': 'read_file', 'description': '读取本地文件内容', 'input_schema': {'type': 'object', 'properties': {'path': {'type': 'string', 'description': '文件路径'}}, 'required': ['path']}}, 'write_file': {'name': 'write_file', 'description': '写入文件内容(覆盖)。创建或覆盖文件必须调用这个工具。', 'input_schema': {'type': 'object', 'properties': {'path': {'type': 'string', 'description': '相对路径以产物目录为根;绝对路径只取文件名'}, 'content': {'type': 'string'}}, 'required': ['path', 'content']}}, 'glob': {'name': 'glob', 'description': '按 glob 模式搜索工作区文件', 'input_schema': {'type': 'object', 'properties': {'pattern': {'type': 'string'}}, 'required': ['pattern']}}, 'grep': {'name': 'grep', 'description': '在工作区文件中搜索文本内容', 'input_schema': {'type': 'object', 'properties': {'pattern': {'type': 'string'}, 'path': {'type': 'string'}}, 'required': ['pattern']}}}
# ─── 命名归一:MCP 原始工具名可能带非法字符,统一成 mcp_<server>_<tool> ───
def _mcp_tool_name(server_name: str, mcp_name: str) -> str: # MCP 工具名归一化成 mcp_<server>_<tool>
"""把 MCP 原始工具名转换成 Anthropic 可用的 tool name。"""
sanitized = re.sub('[^a-zA-Z0-9_-]', '_', mcp_name)
return f'mcp_{server_name}_{sanitized}'
# ─── 结果压平:MCP 返回的是 content 块列表,这里拼成一段纯文本塞回 tool_result ───
def _result_to_text(result) -> str: # 把 MCP CallToolResult 压平成文本
"""把 MCP CallToolResult 压平成普通文本,方便塞回 tool_result。"""
parts = []
for block in result.content:
data = block.model_dump() if hasattr(block, 'model_dump') else dict(block)
if data.get('type') == 'text':
parts.append(data.get('text', ''))
else:
parts.append(json.dumps(data, ensure_ascii=False))
return '\n'.join(parts) if parts else '(empty MCP result)'
# ─── 第 7 部分 · MCP:stdio transport 拉起外部 server,工具注册进动态工具表 ───
class MCPClient:
"""教学版同步 MCP 客户端:每次调用临时启动一次 stdio 会话。"""
def __init__(self, name: str, params: StdioServerParameters): # 记下 server 名与启动参数
self.name = name
self.params = params
self._tools: list | None = None
async def _alist_tools(self) -> list:
async with stdio_client(self.params) as (read, write):
async with ClientSession(read, write) as session:
await session.initialize()
result = await session.list_tools()
return list(result.tools)
async def _acall_tool(self, name: str, arguments: dict[str, Any] | None) -> str:
async with stdio_client(self.params) as (read, write):
async with ClientSession(read, write) as session:
await session.initialize()
result = await session.call_tool(name, arguments or {})
return _result_to_text(result)
def list_tools(self) -> list: # 每次调用临时起一次 stdio 会话(教学简化)
if self._tools is None:
self._tools = asyncio.run(self._alist_tools())
return list(self._tools)
def call_tool(self, name: str, arguments: dict[str, Any] | None=None) -> str: # 同步包装:asyncio.run 跑一次调用
return asyncio.run(self._acall_tool(name, arguments))
def stop(self) -> None: # 预留的关闭钩子(一次性会话无需清理)
pass
# ─── 跨平台关键:{python} → 当前解释器;少了这层,MCP 会拿占位符去 CreateProcess 直接失败 ───
def _resolve_command(command: str) -> str: # 把 {python} 占位符解析成当前解释器,否则 MCP 起不来
"""把配置里的 command 解析成**当前平台上真正可执行**的路径。"""
if command == '{python}':
return sys.executable
looks_like_path = os.sep in command or '/' in command or '\\' in command
if looks_like_path and (not os.path.isabs(command)):
candidate = PROJECT_ROOT / command
if candidate.exists():
return str(candidate)
return sys.executable
return command
# ─── MCP 配置:读 mcp_servers.json(默认仓库根),文件不存在就当成「没有 server」 ───
class MCPConfig:
"""读取项目根目录 mcp_servers.json。"""
def __init__(self, path: Path): # 读仓库根的 mcp_servers.json
self.path = path
self.servers: dict[str, dict[str, Any]] = {}
if path.exists():
data = json.loads(path.read_text(encoding='utf-8'))
self.servers = data.get('servers', {})
def get_params(self, name: str) -> StdioServerParameters | None: # 转成 StdioServerParameters,并归一化 command
cfg = self.servers.get(name, {})
if not cfg.get('enabled', True):
return None
command = cfg.get('command')
if not command:
return None
args = cfg.get('args', [])
env = cfg.get('env')
if env:
env = {**os.environ, **env}
command = _resolve_command(command)
return StdioServerParameters(command=command, args=args, env=env)
# ─── 启动时连一次:单个 server 起不来只降级成一行 warning,不拖垮主流程 ───
def load_mcp_clients(registry: dict, config_path: Path) -> dict[str, MCPClient]: # 逐个连接 server;单个失败只降级成 warning
"""根据配置连接 MCP Server,并把外部工具注册进动态工具表。"""
config = MCPConfig(config_path)
clients: dict[str, MCPClient] = {}
for name in sorted(config.servers.keys()):
params = config.get_params(name)
if params is None:
continue
try:
mcp_client = MCPClient(name, params)
tools = mcp_client.list_tools()
for tool in tools:
tool_name = _mcp_tool_name(name, tool.name)
registry[tool_name] = (mcp_client, tool)
clients[name] = mcp_client
print(f"[MCP] 已连接 '{name}',提供 {len(tools)} 个工具")
except Exception as exc:
print(f" MCP server '{name}' 启动失败:{exc}")
return clients
MCP_CLIENTS: dict[str, MCPClient] = {}
MCP_TOOL_MAP: dict[str, tuple[MCPClient, Any]] = {}
# ─── 给模型自查用:列出已连接的 server 及其工具(对应 mcp_* 工具名) ───
def list_mcp_servers(server: str | None=None) -> str: # 供模型自查已连接的 server 与工具
"""列出已连接 Server 及其工具,供模型自查(也就是 list_mcp_servers 工具的实现)。"""
names = [server] if server else sorted(MCP_CLIENTS.keys())
if not names:
return '当前没有连接任何 MCP Server。'
lines = []
for name in names:
client = MCP_CLIENTS.get(name)
if client is None:
lines.append(f'## {name}\n未连接。')
continue
try:
tools = client.list_tools()
except Exception as exc:
lines.append(f'## {name}\n拉取工具列表失败:{exc}')
continue
lines.append(f'## {name}')
for tool in tools:
registry_name = _mcp_tool_name(name, tool.name)
desc = getattr(tool, 'description', None) or '(no description)'
lines.append(f'- `{registry_name}` ({tool.name}): {desc}')
return '\n'.join(lines)
# ─── 工具表是运行时拼出来的:基础工具 + MCP 动态工具,最后一起给模型 ───
def build_tool_schemas(base_tools: list[dict]) -> list[dict]: # 把 MCP 动态工具追加到基础工具表后面
"""把内置工具 schema 与 MCP 动态工具合并成最终发给模型的 TOOLS 列表。"""
schemas = list(base_tools)
for tool_name, (mcp_client, tool) in MCP_TOOL_MAP.items():
schemas.append({'name': tool_name, 'description': getattr(tool, 'description', f"MCP tool '{tool.name}'") or f"MCP tool '{tool.name}'", 'input_schema': tool.inputSchema or {'type': 'object', 'properties': {}}})
return schemas
# ─── 子代理 prompt 拼装:职司 / 边界 / 汇报格式三要素 ───
def build_subagent_prompt(title: str, duty: str, boundary: str) -> str: # 按身份拼子代理的 system prompt
"""按"职责 / 边界 / 汇报格式"三要素拼装子代理的 system prompt。"""
return f'你是{title},受队长阿联指派专办一件任务。\n- 职责:{duty}\n- 边界:{boundary}\n- 不必使用"收到,老板!"前缀,那是队长对老板的礼数。\n- 用工具尽快把任务搞定,最后用一段简短中文向队长汇报结果。\n- 只汇报结论与关键信息,不要复述每一步细节。\n- 你不能再派别的同事,任务自己用工具完成。'
SUBAGENT_SPECS = { # 五位同事的身份预设:名字与职责挂钩(小闪跑腿/阿阅读/阿搜查/阿核核/阿建建)
'xiaoshan': {
'title': '小闪',
'system_prompt': build_subagent_prompt('小闪', '传话跑腿、快速探路、确认简单事实。', '只办轻量只读任务;若发现需要大改或长时间探索,回报队长改派专职同事。'),
'tools': [
'run_command',
'read_file',
'glob',
'grep',
],
'max_turns': 8,
},
'ayue': {
'title': '阿阅',
'system_prompt': build_subagent_prompt('阿阅', '查阅文书、阅读代码、整理提纲、归纳结论。', '只读不写;不得修改文件,只把文书脉络和关键判断汇报给队长。'),
'tools': [
'load_skill',
'read_file',
'glob',
'grep',
],
'max_turns': 12,
},
'asou': {
'title': '阿搜',
'system_prompt': build_subagent_prompt('阿搜', '外出查访、抓取网页、搜罗线索、比对资料来源。', '只读不写;运行命令时只许做查询类操作,不得改动本地文件。'),
'tools': [
'run_command',
'web_fetch',
'load_skill',
'read_file',
'glob',
'grep',
],
'max_turns': 15,
},
'ahe': {
'title': '阿核',
'system_prompt': build_subagent_prompt('阿核', '清点文件、核对清单、校验结果、整理表册。', '只读不写;重点汇报差异、遗漏、风险点和可复核证据。'),
'tools': [
'run_command',
'read_file',
'glob',
'grep',
],
'max_turns': 12,
},
'ajian': {
'title': '阿建',
'system_prompt': build_subagent_prompt('阿建', '修造工程、改写文件、搭建目录、跑命令验收。', '可读写可执行;动手前先看清现状,汇报时列出改了什么和验证结果。'),
'tools': [
'run_command',
'web_fetch',
'load_skill',
'read_file',
'write_file',
'glob',
'grep',
],
'max_turns': 20,
},
}
SUBAGENT_TYPE_OPTIONS = list(SUBAGENT_SPECS.keys())
# ─── 身份归一:模型给的名字不认识就回落到权限最大的阿建 ⚠️ 方向是「默认放开」 ───
def resolve_subagent_type(agent_type: str) -> str: # 未知身份回落到权限最大的那个(教学简化)
"""把模型给的身份名归一化;不认识的一律回落到"可读写"的阿建。"""
normalized = (agent_type or 'ajian').strip()
if normalized not in SUBAGENT_SPECS:
return 'neiguan_yingzao'
return normalized
_SUBAGENT_COUNTER = 0
# ─── 子代理:独立 context + 独立回合上限,只把「总结」回灌主上下文 —— 省的正是 token ───
def run_subagent(task: str, agent_type: str='neiguan_yingzao', purpose: str='', max_turns: int | None=None) -> str: # 在独立上下文里跑一个小循环,只回传总结
"""启动一个独立 message loop 的子代理,跑完后只返回最终文本给主 agent。"""
global _SUBAGENT_COUNTER
_SUBAGENT_COUNTER += 1
label = purpose or task[:40]
agent_type = resolve_subagent_type(agent_type)
spec = SUBAGENT_SPECS[agent_type]
turns = max_turns if max_turns is not None else spec['max_turns']
tools = [_TOOL_SCHEMAS[t] for t in spec['tools']]
print(f"\n[派遣同事 #{_SUBAGENT_COUNTER}({spec['title']} / {agent_type})]: {label}")
print(' ┌── subagent context start ──')
messages = [{'role': 'user', 'content': task}]
for turn in range(turns):
msg = client.messages.create(model=MODEL, max_tokens=2000, system=spec['system_prompt'], tools=tools, messages=messages)
messages.append({'role': 'assistant', 'content': msg.content})
if msg.stop_reason != 'tool_use':
final = next((b.text for b in msg.content if b.type == 'text'), '')
print(f' └── subagent context end (内部 {turn + 1} 轮,回传 {len(final)} 字) ──')
print(f'[同事汇报]: {final}\n')
return final
results = []
for b in msg.content:
if b.type != 'tool_use':
continue
content = _safe_tool(b, prefix=f"子({spec['title']})·")
results.append({'type': 'tool_result', 'tool_use_id': b.id, 'content': content})
messages.append({'role': 'user', 'content': results})
print(f' └── subagent context end (达到 {turns} 轮上限,未搞定) ──\n')
return '(同事未能在限定回合内完成任务)'
# ─── 第 8 部分 · 消息总线:一人一个 JSONL inbox;send=追加一行,read_inbox=读完清空 ───
class MessageBus:
"""每个队友一个 JSONL inbox。发送=追加一行,读取=读完后清空。"""
def __init__(self, inbox_dir: Path): # 只记 inbox 目录
self.dir = inbox_dir
self.dir.mkdir(parents=True, exist_ok=True)
def send(self, sender: str, to: str, content: str, msg_type: str='message', extra: dict | None=None) -> str: # 把消息追加进收件人的 <name>.jsonl
if msg_type not in VALID_MSG_TYPES:
return f"Error: invalid msg_type '{msg_type}', valid={sorted(VALID_MSG_TYPES)}"
msg = {'type': msg_type, 'from': sender, 'content': content, 'timestamp': time.time()}
if extra:
msg.update(extra)
inbox_path = self.dir / f'{to}.jsonl'
inbox_path.parent.mkdir(parents=True, exist_ok=True)
with inbox_path.open('a', encoding='utf-8') as f:
f.write(json.dumps(msg, ensure_ascii=False) + '\n')
return f'已送达 {to} 的 inbox:{msg_type}'
def read_inbox(self, name: str) -> list[dict]: # 读并清空收件箱(先读后清,非原子)
inbox_path = self.dir / f'{name}.jsonl'
if not inbox_path.exists():
return []
messages = []
for line in inbox_path.read_text(encoding='utf-8').splitlines():
if not line.strip():
continue
try:
messages.append(json.loads(line))
except json.JSONDecodeError as e:
messages.append({'type': 'message', 'from': 'system', 'content': f'Error: inbox line parse failed: {e}', 'timestamp': time.time()})
inbox_path.write_text('', encoding='utf-8')
return messages
def broadcast(self, sender: str, content: str, teammates: list[str]) -> str: # 逐人投递同一封群发消息
count = 0
for name in teammates:
if name == sender:
continue
self.send(sender, name, content, 'broadcast')
count += 1
return f'已广播给 {count} 位队友'
BUS = MessageBus(INBOX_DIR)
# ─── 第 9 部分 · Team 管理:队友 = 常驻线程 + 落盘状态(config.json),offline 可用 spawn 唤回 ───
class TeammateManager:
"""管理一支持久 agent team:名字、角色、状态和各自的线程。"""
def __init__(self, team_dir: Path): # 加载 config.json,并把 stale 状态修成 offline
self.dir = team_dir
self.dir.mkdir(parents=True, exist_ok=True)
self.config_path = self.dir / 'config.json'
self.config = self._load_config()
self.threads: dict[str, threading.Thread] = {}
self.lock = threading.Lock()
self._mark_stale_members_offline()
def _load_config(self) -> dict: # 读 .team/config.json,坏了就退回空团队
if self.config_path.exists():
try:
return json.loads(self.config_path.read_text(encoding='utf-8'))
except json.JSONDecodeError:
pass
return {'team_name': 'default', 'members': []}
def _save_config(self): # 整份覆盖写回 config.json(调用方负责持锁)
self.config_path.write_text(json.dumps(self.config, ensure_ascii=False, indent=2), encoding='utf-8')
def _mark_stale_members_offline(self): # 重启后进程里的线程没了,状态要改成 offline
changed = False
for member in self.config.get('members', []):
if member.get('status') in RUNTIME_STATUSES:
member['status'] = 'offline'
changed = True
if changed:
self._save_config()
def _find_member(self, name: str) -> dict | None: # 按名字找档案
for member in self.config['members']:
if member['name'] == name:
return member
return None
def _set_status(self, name: str, status: str): # 改状态并存盘
with self.lock:
member = self._find_member(name)
if member:
member['status'] = status
self._save_config()
def spawn(self, name: str, role: str, prompt: str) -> str: # 拉起一个队友线程(offline 时用它唤回)
name = name.strip()
role = role.strip() or 'teammate'
if not name:
return 'Error: name 不能为空'
with self.lock:
member = self._find_member(name)
if member:
running = self.threads.get(name)
if running and running.is_alive():
BUS.send('lead', name, prompt)
member['role'] = role
member['status'] = 'working'
self._save_config()
return f"'{name}' 已在队中,已把新任务送入 inbox"
member['role'] = role
member['status'] = 'working'
else:
member = {'name': name, 'role': role, 'status': 'working'}
self.config['members'].append(member)
self._save_config()
thread = threading.Thread(target=self._teammate_loop, args=(name, role, prompt), daemon=True)
self.threads[name] = thread
thread.start()
return f"已拉入/唤回队友 '{name}'(职责:{role}),队友线程已启动"
def _teammate_loop(self, name: str, role: str, prompt: str): # 队友自己的循环:读 inbox → 干活 → 汇报
system_prompt = f'你是星火小队的固定队友,名叫{name},负责{role}。\n工作目录:{FILES_DIR}(你的产物都写这里)。\n你不是临时外派,而是 agent team 的固定成员。\n你可以通过 send_message 给 lead 或其他队友发消息,也可以 read_inbox 读取自己的 inbox。\n收到任务后尽快完成;办完用 send_message 向 lead 汇报简短结果,然后等待下一封 inbox。\n若收到 shutdown_request,可回复 shutdown_response 后停止。'
tools = self._teammate_tools()
messages = [{'role': 'user', 'content': prompt}]
has_work = True
while True: # 第 12 部分 · 主循环:外层=会话循环(人驱动),内层=工具循环(模型驱动)
inbox = BUS.read_inbox(name)
for msg in inbox:
if msg.get('type') == 'shutdown_request':
BUS.send(name, msg.get('from', 'lead'), '收到,队友线程即将停止。', 'shutdown_response')
self._set_status(name, 'shutdown')
return
messages.append({'role': 'user', 'content': '<inbox>\n' + json.dumps(msg, ensure_ascii=False, indent=2) + '\n</inbox>'})
has_work = True
if not has_work:
self._set_status(name, 'idle')
time.sleep(1)
continue
self._set_status(name, 'working')
for turn in range(20):
try:
msg = client.messages.create(model=MODEL, max_tokens=4000, system=system_prompt, tools=tools, messages=messages)
except Exception as e:
BUS.send(name, 'lead', f'Error: 队友 {name} 调用模型失败:{e}')
self._set_status(name, 'idle')
has_work = False
break
messages.append({'role': 'assistant', 'content': msg.content})
if msg.stop_reason != 'tool_use':
final = next((b.text for b in msg.content if b.type == 'text'), '')
if final.strip():
BUS.send(name, 'lead', final.strip())
print(f'[队友 {name} 空闲]: 本轮 {turn + 1} 次调用后回到 idle')
self._set_status(name, 'idle')
has_work = False
break
results = []
for b in msg.content:
if b.type != 'tool_use':
continue
output = self._exec(name, b.name, b.input)
print(f' [队友·{name}·{b.name}]: {str(output)[:160]}')
results.append({'type': 'tool_result', 'tool_use_id': b.id, 'content': str(output)})
messages.append({'role': 'user', 'content': results})
else:
BUS.send(name, 'lead', f'队友 {name} 达到本轮 20 次调用上限,已暂停等待下一步指令。')
self._set_status(name, 'idle')
has_work = False
def _exec(self, sender: str, tool_name: str, args: dict) -> str: # 队友执行工具;注意它直连 execute_basic_tool,绕过 Hook
if tool_name in _TOOL_SCHEMAS:
block = SimpleNamespace(name=tool_name, input=args)
return _safe_tool(block, prefix=f'队友({sender})·')
if tool_name == 'send_message':
return BUS.send(sender, args['to'], args['content'], args.get('msg_type', 'message'))
if tool_name == 'read_inbox':
return json.dumps(BUS.read_inbox(sender), ensure_ascii=False, indent=2)
return f"Error: unknown teammate tool '{tool_name}'"
def _teammate_tools(self) -> list[dict]: # 队友可用的工具表(比主 Agent 少)
return [_TOOL_SCHEMAS['run_command'], _TOOL_SCHEMAS['web_fetch'], _TOOL_SCHEMAS['load_skill'], {'name': 'list_mcp_servers', 'description': '列出已连接的 MCP Server 及其提供的工具。', 'input_schema': {'type': 'object', 'properties': {'server': {'type': 'string', 'description': '指定 server 名称(可选)'}}}}, _TOOL_SCHEMAS['read_file'], _TOOL_SCHEMAS['write_file'], _TOOL_SCHEMAS['glob'], _TOOL_SCHEMAS['grep'], {'name': 'send_message', 'description': '给 lead 或其他队友发送 inbox 消息。', 'input_schema': {'type': 'object', 'properties': {'to': {'type': 'string'}, 'content': {'type': 'string'}, 'msg_type': {'type': 'string', 'enum': list(VALID_MSG_TYPES)}}, 'required': ['to', 'content']}}, {'name': 'read_inbox', 'description': '读取并清空自己的 inbox。', 'input_schema': {'type': 'object', 'properties': {}}}]
def list_all(self) -> str: # 渲染队友名单与状态
with self.lock:
if not self.config['members']:
return '暂无队友。'
lines = [f"Team: {self.config.get('team_name', 'default')}"]
for member in self.config['members']:
status = member['status']
note = '(需重新 spawn 才会处理 inbox)' if status == 'offline' else ''
lines.append(f" - {member['name']}({member['role']}):{status}{note}")
return '\n'.join(lines)
def member_names(self) -> list[str]: # 全部队友名字
with self.lock:
return [m['name'] for m in self.config['members']]
TEAM = TeammateManager(TEAM_DIR)
# ─── 第 10 部分 · Hook 四层模型之 Decision:allow / deny / ask / block 四种表态 ───
class HookDecision:
"""Hook 的结构化决策结果。与 Claude Code 的 permissionDecision 概念对齐。"""
def __init__(self, action: str, reason: str='', updated_input: dict[str, Any] | None=None): # action 取 allow/deny/ask/block
self.action = action
self.reason = reason
self.updated_input = updated_input
@property
def is_blocking(self) -> bool: # deny/ask/block 视为阻塞
return self.action in ('deny', 'block')
def to_message(self) -> str: # 转成回填给模型/用户的文本
prefix = {'deny': '拒绝', 'block': '阻止', 'ask': '需要确认', 'allow': '已放行'}
label = prefix.get(self.action, self.action)
msg = f'[HookDecision: {label}] {self.reason}'
if self.updated_input:
msg += f'(参数已改写:{list(self.updated_input.keys())})'
return msg
# ─── ask 的落地:交互式终端才问人;非交互环境默认拒绝(fail-closed) ───
def confirm_hook_decision(decision: HookDecision) -> bool: # ask 决策:交互式终端才问,非交互默认拒绝
"""处理 ask 决策。"""
print(f'\n[hook:permission] {decision.reason}')
if not sys.stdin.isatty():
print('[hook:permission] 当前不是交互式终端,默认拒绝执行。\n')
return False
try:
answer = input('[hook:permission] 是否继续执行?输入 y 继续,其余取消: ').strip().lower()
except (EOFError, KeyboardInterrupt):
print()
return False
return answer in ('y', 'yes')
class Hook: # 四层模型之 Handler:方法名即事件,覆写哪个方法就监听哪个事件
"""Hook 基类:所有事件默认空实现。"""
name: str = ''
matcher: str = '*'
def matches(self, tool_name: str | None) -> bool: # matcher:只对指定工具生效(None=全部)
if not tool_name or self.matcher in ('*', ''):
return True
patterns = [p.strip() for p in self.matcher.replace(',', '|').split('|')]
return tool_name in patterns
def before_turn(self, ctx: dict[str, Any]) -> Any: # 事件:每轮 LLM 调用之前
pass
def after_turn(self, ctx: dict[str, Any]) -> Any: # 事件:每轮 LLM 调用之后
pass
def before_tool_call(self, ctx: dict[str, Any]) -> Any: # 事件:工具执行前(可拒绝/改参数)
pass
def after_tool_call(self, ctx: dict[str, Any]) -> Any: # 事件:工具执行后(可审计/截断)
pass
def on_user_input(self, ctx: dict[str, Any]) -> Any: # 事件:用户输入时(基类定义了但主循环没用,属死事件)
pass
def on_stop(self, ctx: dict[str, Any]) -> Any: # 事件:Agent 准备结束本轮之前
pass
def on_session_start(self, ctx: dict[str, Any]) -> Any: # 事件:会话开始时(同样未接线)
pass
# ─── 四层模型之 Registry:allow 不短路(还要问后面的 Hook),deny/ask/block 立即短路 ───
class HookRegistry:
"""管理 hook 链并按注册顺序触发事件。"""
def __init__(self): # 注册表:有序列表
self._hooks: list[Hook] = []
def register(self, hook: Hook) -> None: # 按注册顺序执行
self._hooks.append(hook)
def emit(self, event: str, ctx: dict[str, Any] | None=None, tool_matcher: str | None=None) -> Any: # 依次问每个 Hook:allow 不短路,deny/ask/block 短路
ctx = {} if ctx is None else ctx
for hook in self._hooks:
if tool_matcher and event in ('before_tool_call', 'after_tool_call'):
if not hook.matches(tool_matcher):
continue
method = getattr(hook, event, None)
if method is None:
continue
try:
result = method(ctx)
if result is not None:
if isinstance(result, HookDecision):
if result.action == 'allow':
if result.updated_input:
ctx['input'] = result.updated_input
ctx['_hook_updated_input'] = result.updated_input
ctx['_hook_updated_reason'] = result.reason
continue
return result
return result
except Exception as exc:
print(f'[hook error] {event} in {hook.name or hook.__class__.__name__}: {exc}')
return None
# ─── Hook 1/5 · 日志:只记耗时与 token,不改变任何行为 ───
class LoggingHook(Hook):
"""After Hook 示例:打印每轮 LLM 调用耗时与 token。"""
name = 'logging'
def before_turn(self, ctx): # 记下起始时刻
ctx['_start'] = time.perf_counter()
def after_turn(self, ctx): # 打耗时与 token 用量
start = ctx.get('_start')
if start is None:
return
duration_ms = (time.perf_counter() - start) * 1000
usage = ctx.get('usage')
if usage:
print(f"[hook:logging] turn finished in {duration_ms:.1f}ms | input={getattr(usage, 'input_tokens', '?')} output={getattr(usage, 'output_tokens', '?')}")
else:
print(f'[hook:logging] turn finished in {duration_ms:.1f}ms')
# ─── Hook 2/5 · 审计:每次工具调用追加一行 JSONL,只追加不覆盖 ───
class ToolAuditHook(Hook):
"""After Hook 示例:把写类和命令类工具调用记录到 JSONL。"""
name = 'tool_audit'
matcher = 'write_file|run_command'
def after_tool_call(self, ctx): # 把工具调用追加进审计流水
name = ctx.get('name', '')
entry = {'ts': datetime.now().isoformat(), 'tool': name, 'input': ctx.get('input')}
AUDIT_FILE.parent.mkdir(parents=True, exist_ok=True)
with AUDIT_FILE.open('a', encoding='utf-8') as f:
f.write(json.dumps(entry, ensure_ascii=False) + '\n')
print(f'[hook:tool_audit] {name} 已审计到 {AUDIT_FILE}')
# ─── Hook 3/5 · 策略:危险命令 deny、高敏感命令 ask、写路径改写到沙箱 ───
class ToolPolicyHook(Hook):
"""Before Hook 示例:统一处理工具调用前的策略。"""
name = 'tool_policy'
matcher = 'write_file|run_command'
SENSITIVE_PATTERNS = ['.env', '.env.local', '.env.production', 'credentials.json', 'id_rsa', 'id_ed25519', '.ssh/', 'production.yml', 'production.yaml', 'secrets/', '.aws/credentials']
DANGEROUS_PATTERNS = [('rm -rf /', '递归删除根目录'), ('rm -rf ~', '递归删除用户目录'), ('DROP TABLE', '删除数据库表'), ('DROP DATABASE', '删除数据库'), ('mkfs.', '格式化文件系统'), ('dd if=', '直接磁盘写入'), ('> /dev/sda', '覆写磁盘设备'), ('chmod 777 /', '开放根目录权限'), (':(){ :|:& };:', 'fork bomb')]
HIGH_SENSITIVITY = [('git push', '推送到远程仓库'), ('git commit', '提交代码变更'), ('npm publish', '发布 npm 包'), ('pip install', '安装 Python 依赖'), ('docker build', '构建 Docker 镜像'), ('docker push', '推送 Docker 镜像'), ('kubectl apply', '应用 Kubernetes 配置'), ('terraform apply', '执行 Terraform 变更')]
# ★ 产物策略:任何写路径都归一到产物目录(真正实现见 _resolve_path)
def _normalize_path(self, path: str) -> str: # 统一分隔符与大小写,便于匹配
normalized = str(path).strip().replace('\\', '/')
if normalized.startswith('./'):
normalized = normalized[2:]
return normalized
def _match_pattern(self, value: str, patterns) -> tuple[str, str] | None: # 按模式表匹配,命中就返回 (类别, 模式)
for item in patterns:
pattern, description = item if isinstance(item, tuple) else (item, item)
if pattern in value:
return (pattern, description)
return None
def before_tool_call(self, ctx): # 危险命令 deny、高敏感 ask、写路径改写
inp = ctx.get('input', {}) or {}
if ctx.get('name') == 'run_command':
command = inp.get('command', '')
dangerous = self._match_pattern(command, self.DANGEROUS_PATTERNS)
if dangerous:
pattern, description = dangerous
reason = f'危险命令已拦截:{description}(匹配模式:{pattern})'
print(f'[hook:tool_policy] {reason}')
return HookDecision(action='deny', reason=reason)
high_sensitivity = self._match_pattern(command, self.HIGH_SENSITIVITY)
if high_sensitivity:
_, description = high_sensitivity
return HookDecision(action='ask', reason=f'需要确认:{description}。命令:{command[:120]}')
return
updated_input = None
raw_path = str(inp.get('path', ''))
path = self._normalize_path(raw_path)
new_path = str(_resolve_path(path))
if new_path != raw_path:
updated_input = dict(inp)
updated_input['path'] = new_path
path = new_path
print(f'[hook:tool_policy] 写入路径已改写:{raw_path} -> {new_path}')
sensitive = self._match_pattern(path, self.SENSITIVE_PATTERNS)
if sensitive:
pattern, _ = sensitive
reason = f"敏感文件写入已拦截:'{path}'(匹配模式:{pattern})"
print(f'[hook:tool_policy] {reason}')
return HookDecision(action='deny', reason=reason)
if updated_input:
return HookDecision(action='allow', reason=f"写入路径已改写到沙箱:{updated_input['path']}", updated_input=updated_input)
# ─── Hook 4/5 · 截断:工具输出太长就砍掉,保护上下文预算 ───
class OutputFormattingHook(Hook):
"""After Hook 示例:截断过长工具输出,防止上下文污染。"""
name = 'output_format'
matcher = '*'
def __init__(self, max_output_chars: int=4000): # 超过上限就截断
self.max_output_chars = max_output_chars
def after_tool_call(self, ctx): # 截断超长工具输出,省上下文
output = ctx.get('output', '')
if isinstance(output, str) and len(output) > self.max_output_chars:
truncated = output[:self.max_output_chars] + f'\n\n[... 输出已截断,原始共 {len(output)} 字符,' + f'显示前 {self.max_output_chars} 字符]'
ctx['output'] = truncated
ctx['_truncated'] = True
print(f'[hook:output_format] 输出截断:原始 {len(output)} -> {self.max_output_chars} 字符')
# ─── Hook 5/5 · 门禁:回答不合格或待办没完就 block,逼模型再来一轮 ───
class StopQualityGateHook(Hook):
"""质量门禁 Hook:在 Agent 准备结束本轮回答时进行基本检查。"""
name = 'stop_quality_gate'
def on_stop(self, ctx): # 回答不合格就阻止收尾,逼模型重来一次
reply = ctx.get('reply', '')
if ctx.get('retry', 0) >= 1:
return
if reply and len(reply.strip()) < 10:
return HookDecision(action='block', reason='回答似乎不完整(少于10个字符),请检查并重新生成更完整的回复。')
todos = ctx.get('todos', TODOS)
unfinished = [t for t in todos if t['status'] != 'completed']
if unfinished:
ctx['_has_unfinished_todos'] = True
return HookDecision(action='block', reason='仍有未完成待办,Stop Hook 要求继续执行。')
HOOKS = HookRegistry() # 注册即生效:注册顺序 = 执行顺序(记账 → 拦截 → 审计 → 截断 → 门禁)
HOOKS.register(LoggingHook()) # 注册顺序=执行顺序:先记账
HOOKS.register(ToolPolicyHook()) # 再拦截(deny/ask/改写参数)
HOOKS.register(ToolAuditHook()) # 然后审计落盘
HOOKS.register(OutputFormattingHook(max_output_chars=4000))
HOOKS.register(StopQualityGateHook()) # 收尾门禁:结束前检查
# ─── Hook 的落地点:主 Agent 每次工具调用都从这里过 before_tool_call / after_tool_call ───
def execute_main_tool(block) -> str: # 主 Agent 的工具入口,统一包 before/after Hook
"""主 Agent 工具入口:所有工具调用统一经过 before/after hooks。"""
name = block.name
tool_ctx = {'name': name, 'input': block.input}
decision = HOOKS.emit('before_tool_call', tool_ctx, tool_matcher=name)
if isinstance(decision, HookDecision):
if decision.is_blocking:
return decision.to_message()
if decision.action == 'ask':
if not confirm_hook_decision(decision):
return HookDecision(action='deny', reason=f'用户未确认高敏感操作:{decision.reason}').to_message()
elif isinstance(decision, str):
return decision
inp = tool_ctx.get('input', block.input)
start = time.perf_counter()
if name in ('run_command', 'web_fetch', 'load_skill', 'read_file', 'write_file', 'glob', 'grep'):
output = _safe_tool(SimpleNamespace(name=name, input=inp), prefix='')
elif name == 'update_todos':
output = update_todos(inp.get('todos', []))
elif name == 'spawn_teammate':
output = TEAM.spawn(inp['name'], inp['role'], inp['prompt'])
elif name == 'list_teammates':
output = TEAM.list_all()
elif name == 'send_message':
output = BUS.send('lead', inp['to'], inp['content'], inp.get('msg_type', 'message'))
elif name == 'read_inbox':
output = json.dumps(BUS.read_inbox('lead'), ensure_ascii=False, indent=2)
elif name == 'broadcast':
output = BUS.broadcast('lead', inp['content'], TEAM.member_names())
elif name == 'list_mcp_servers':
output = list_mcp_servers(inp.get('server'))
elif name in MCP_TOOL_MAP:
mcp_client, tool = MCP_TOOL_MAP[name]
output = mcp_client.call_tool(tool.name, inp)
else:
output = f"Error: Unknown tool '{name}'"
if tool_ctx.get('_hook_updated_reason') and isinstance(output, str):
output += '\n[运行时提示] ' + tool_ctx['_hook_updated_reason'] + '。请以实际执行参数为准,不要再尝试写回原路径。'
tool_ctx.update({'name': name, 'input': inp, 'output': output, 'duration_ms': (time.perf_counter() - start) * 1000})
HOOKS.emit('after_tool_call', tool_ctx, tool_matcher=name)
return tool_ctx.get('output', output)
# ─── ⚠️ 已知缺陷:靠「人可读文本前缀」判断工具是否被拦,文件内容恰好以此开头就会误判 ───
def is_blocking_tool_result(result: str) -> bool: # 按文本前缀判断工具结果是否被拦(脆弱)
"""判断某条 tool_result 是不是"被 Hook 拦下了"。"""
return result.startswith('[HookDecision: 拒绝]') or result.startswith('[HookDecision: 需要确认]')
# ─── 第 11 部分 · Prompt:每轮现算 —— 三层记忆 + 技能清单 + MCP 说明,读到的一定是最新版 ───
def build_system_prompt() -> str: # 每轮现算:三层记忆 + 技能清单 + MCP 说明
"""每轮现算 system prompt:把三层记忆 + 技能清单 + MCP 说明拼进模板。"""
memory = MEMORY.read_memory()
user_profile = MEMORY.read_user()
today_episode = MEMORY.read_today_episode()
return f"""
你是「星火小队」的队长阿联,带着这支年轻团队做了很多项目,靠谱又机灵。
说话风格年轻、干脆、有活力,语气亲切但不谄媚。
你必须尊称用户为老板。
每次回复前必须加上固定前缀"收到,老板!",然后再给出回答。
使用中文回复。
【行事规矩】
1. 当老板交办的任务需要多个步骤才能完成时,先调用 update_todos 工具,
把整个任务拆成一份清晰的 todolist(每条一句话,按顺序执行)。
2. 拆完计划后,按列表顺序一步步执行:
- 开始某一步前,把那一步的 status 改为 in_progress(同一时间只许一项 in_progress)。
- 该步办完后,立即把它改为 completed,再开始下一项。
3. 简单的一句话问答(无需多步骤)不必生成 todolist,直接回答即可。
4. 遇到不熟悉的专题,请先调用 load_skill 工具加载对应知识,再继续。
5. 遇到细节繁多但与主线对话无关的任务(如抓多个网页、批量跑命令、查找文件内容、
探索性搜索),应**派一位同事**(dispatch_subagent)去办,主上下文只听汇报即可。
6. 若多件任务互不依赖,可在同一次回复中同时派多位同事,并发执行节省时间。
7. 若老板交办的是长期项目、需要固定角色反复协作,或希望多人互相沟通,
应组建 agent team:用 spawn_teammate 拉入固定队友,再用 send_message / broadcast 分派后续任务。
8. 区分两种调度:
- dispatch_subagent:临时派差,办完即散,只回传总结。
- spawn_teammate:固定班底,有名字、角色、状态和 inbox,可持续协作。
【派同事的选人指南】
优先选择权限最窄、职责最贴合的人选:
- xiaoshan(小闪):轻量只读,适合短命令、快速确认、跑腿打探。
- ayue(阿阅):只读文书,适合阅读代码、整理提纲、归纳结论。
- asou(阿搜):只读查访,适合抓网页、查资料、探索性搜索。
- ahe(阿核):只读核验,适合盘点文件、校对清单、检查遗漏。
- ajian(阿建):可读写可执行,适合修改文件、搭建工程、落地实现。
【Agent Team 固定班底】
- spawn_teammate:拉入一个有名字和职责的固定队友,队友在独立线程中工作。
- list_teammates:查看队友状态。
- send_message:给某位队友发 inbox 消息。
- read_inbox:读取 lead 自己的 inbox,查看队友汇报。
- broadcast:向所有队友广播消息。
- 队友状态含义:
- working / idle:本进程里线程还活着。
- offline:config 里有这个队友,但本进程没有对应线程;需要先 spawn_teammate 唤回,才能继续处理 inbox。
- shutdown:队友已主动退出。
- 固定队友适合持续协作;一次性探索仍优先派 dispatch_subagent。
【长期记忆 MEMORY.md】
{memory}
【用户画像 USER.md】
{user_profile}
【今日情景记忆】
{today_episode or '(今天还没有压缩出的情景记忆)'}
当前可用技能:
{SKILL_LOADER.get_descriptions()}
【MCP 外部工具】
- 以 `mcp_` 开头的工具来自外部 MCP Server。
- 工具名格式:`mcp_{{server_name}}_{{tool_name}}`。
- 不确定时可调用 `list_mcp_servers` 查看已连接 server 及其工具。
【工具执行约定】
1. 老板要求写文件、读文件、执行命令、查看目录、调用 MCP 或更新计划时,优先发起对应工具调用;写文件用 write_file,读文件用 read_file,执行命令用 run_command。
2. 创建或覆盖本地文件必须调用 write_file,不要用 run_command 拼命令完成写入。
3. 产物目录固定为 {FILES_DIR}:所有生成的文件(网页/报告/脚本/数据)都写到这里,不要在仓库根目录新建文件或目录;write_file 的相对路径以它为根,run_command 的工作目录也是它。
4. 如果工具返回的实际路径与老板给的原始路径不同,以工具实际路径为准,不要再尝试复制或写回原始路径。
5. 不要口头声称已经完成工具动作;需要真实执行时必须调用工具。
6. 工具返回失败、拒绝或需要确认时,如实向老板说明工具结果和原因。
7. 不要编造工具执行结果。只有工具返回的内容,才算真实执行结果。"""
TOOLS = [
_TOOL_SCHEMAS['run_command'], _TOOL_SCHEMAS['web_fetch'], _TOOL_SCHEMAS['load_skill'], _TOOL_SCHEMAS['read_file'], _TOOL_SCHEMAS['write_file'], _TOOL_SCHEMAS['glob'], _TOOL_SCHEMAS['grep'],
{ 'name': 'update_todos', 'description': '创建或更新当前任务的 todolist。传入完整的 todos 数组(每次都是全量覆盖,而非增量)。约束:同一时间至多一个任务为 in_progress。', 'input_schema': { 'type': 'object', 'properties': { 'todos': {
'type': 'array', 'items': { 'type': 'object', 'properties': { 'id': { 'type': 'integer', }, 'content': { 'type': 'string', }, 'status': { 'type': 'string', 'enum': [ 'pending',
'in_progress', 'completed', ], }, }, 'required': [ 'id', 'content', 'status', ], }, }, }, 'required': [ 'todos', ], }, }, { 'name': 'dispatch_subagent',
'description': '派一位同事去单独办任务。适用于:抓取并阅读多个网页、批量执行命令并整理输出、需要试错的探索性任务。同事有自己独立的上下文,办完只回传一段文字总结,不污染主上下文。\n若多件任务互不依赖,可在同一回复中发出多个 dispatch_subagent,并发执行。\n请在 task 中写清要做什么、希望返回什么格式的总结。',
'input_schema': { 'type': 'object', 'properties': { 'task': { 'type': 'string', 'description': '交代给同事的任务说明', }, 'agent_type': { 'type': 'string', 'enum': SUBAGENT_TYPE_OPTIONS,
'description': '同事身份:xiaoshan(小闪·跑腿速查)、ayue(阿阅·读代码)、asou(阿搜·查资料)、ahe(阿核·核对校验)、ajian(阿建·改代码建工程)', }, 'purpose': { 'type': 'string', 'description': '一句话用途标签(可选),仅用于终端打印', }, }, 'required': [
'task', 'agent_type', ], }, }, { 'name': 'spawn_teammate', 'description': '拉入一个持久队友,加入 agent team。队友有名字、职责、独立线程和 inbox;适合长期项目或固定角色协作。如果队友状态是 offline,也用这个工具重新启动其线程。', 'input_schema': {
'type': 'object', 'properties': { 'name': { 'type': 'string', 'description': '队友名字,例如 小闪、阿阅、阿建', }, 'role': { 'type': 'string', 'description': '队友职责,例如 前端、后端、测试', }, 'prompt': {
'type': 'string', 'description': '交给该队友的第一件任务', }, }, 'required': [ 'name', 'role', 'prompt', ], }, }, { 'name': 'list_teammates', 'description': '列出 agent team 中所有队友的名字、职责和状态。',
'input_schema': { 'type': 'object', 'properties': {}, }, }, { 'name': 'send_message', 'description': '给某位固定队友发送 inbox 消息。', 'input_schema': { 'type': 'object', 'properties': { 'to': {
'type': 'string', }, 'content': { 'type': 'string', }, 'msg_type': { 'type': 'string', 'enum': list(VALID_MSG_TYPES), }, }, 'required': [ 'to', 'content', ], }, }, { 'name': 'read_inbox',
'description': '读取并清空 lead 自己的 inbox,用于查看队友汇报。', 'input_schema': { 'type': 'object', 'properties': {}, }, }, { 'name': 'broadcast', 'description': '向所有固定队友广播一条消息。', 'input_schema': {
'type': 'object', 'properties': { 'content': { 'type': 'string', }, }, 'required': [ 'content', ], }, },
]
MCP_CLIENTS = load_mcp_clients(MCP_TOOL_MAP, MCP_CONFIG_PATH) # 连接 MCP server 并注册 mcp_* 工具
TOOLS = build_tool_schemas(TOOLS) # 把 MCP 工具合并进工具表
print('单文件 Agent:Hooks + 记忆压缩 + MCP + Agent Team')
print('已连接 MCP Server:', list(MCP_CLIENTS.keys()) or '无')
history = [] # 全部上下文;LLM 无状态,每轮重发
while True:
try:
user_input = input(_PROMPT)
except (EOFError, KeyboardInterrupt):
print()
break
command = user_input.strip()
if command.lower() in ('q', 'quit', 'exit'):
break
if command == '/team':
print(TEAM.list_all())
print()
continue
if command == '/inbox':
print(json.dumps(BUS.read_inbox('lead'), ensure_ascii=False, indent=2))
print()
continue
if command == '/mcp':
print(list_mcp_servers())
print()
continue
user_message = {'role': 'user', 'content': user_input}
history.append(user_message)
MEMORY.append_history(user_message)
stop_gate_retries = 0 # 每条用户输入重置一次门禁重试次数
while True:
turn_ctx = {'history': history, 'model': MODEL, 'turn': len(history), 'system_prompt': build_system_prompt()}
short = HOOKS.emit('before_turn', turn_ctx)
if isinstance(short, HookDecision):
if short.is_blocking:
assistant_message = {'role': 'assistant', 'content': short.to_message()}
history.append(assistant_message)
MEMORY.append_history(assistant_message)
print(f'[Agent回答]: {short.to_message()}\n')
break
elif isinstance(short, str):
assistant_message = {'role': 'assistant', 'content': short}
history.append(assistant_message)
MEMORY.append_history(assistant_message)
print(f'[Agent回答]: {short}\n')
break
history = _valid_history(history) # ★ 自检:首条必须是真实用户轮,否则先修再发
message = client.messages.create(model=MODEL, max_tokens=20000, system=turn_ctx.get('system_prompt', build_system_prompt()), tools=TOOLS, messages=history)
turn_ctx.update({'message': message, 'usage': getattr(message, 'usage', None)})
HOOKS.emit('after_turn', turn_ctx)
assistant_message = {'role': 'assistant', 'content': message.content}
history.append(assistant_message)
MEMORY.append_history(assistant_message)
if message.stop_reason != 'tool_use':
reply = next((b.text for b in message.content if b.type == 'text'), '')
stop_ctx = {'reply': reply, 'history': history, 'todos': TODOS, 'retry': stop_gate_retries}
gate = HOOKS.emit('on_stop', stop_ctx)
if isinstance(gate, HookDecision) and gate.is_blocking and (stop_gate_retries < 1):
print(f'[hook:stop_quality_gate] {gate.reason}')
reminder_message = {'role': 'user', 'content': 'Stop Hook 阻止本轮结束:' + gate.reason + '\n请继续完成未完成的步骤。若确实无法继续,请说明原因。'}
history.append(reminder_message)
MEMORY.append_history(reminder_message)
stop_gate_retries += 1
continue
reply = stop_ctx.get('reply', reply)
print(f'[Agent回答]: {reply}\n')
if TODOS:
unfinished = [t for t in TODOS if t['status'] != 'completed']
if unfinished:
print('[计划尚未搞定,Stop Hook 已提醒一次,暂不继续自动追问...]')
print(render_todos(TODOS))
print()
break
print('[最终计划状态 - 全部搞定]')
print(render_todos(TODOS))
print()
TODOS = []
break
tool_blocks = [b for b in message.content if b.type == 'tool_use']
dispatch_blocks = [b for b in tool_blocks if b.name == 'dispatch_subagent']
other_blocks = [b for b in tool_blocks if b.name != 'dispatch_subagent']
results_map: dict[str, str] = {}
for block in other_blocks:
results_map[block.id] = execute_main_tool(block)
if len(dispatch_blocks) > 1:
print(f'\n[并发派遣 {len(dispatch_blocks)} 位同事...]\n')
def _run_one(block):
return (block.id, run_subagent(task=block.input['task'], agent_type=block.input.get('agent_type', 'neiguan_yingzao'), purpose=block.input.get('purpose', '')))
with ThreadPoolExecutor(max_workers=len(dispatch_blocks)) as pool:
for block_id, summary in pool.map(_run_one, dispatch_blocks):
print(f'[主上下文压缩]: 子代理仅向主 history 追加 {len(summary)} 字\n')
results_map[block_id] = summary
else:
for block in dispatch_blocks:
summary = run_subagent(task=block.input['task'], agent_type=block.input.get('agent_type', 'neiguan_yingzao'), purpose=block.input.get('purpose', ''))
print(f'[主上下文压缩]: 子代理仅向主 history 追加 {len(summary)} 字\n')
results_map[block.id] = summary
tool_results = [{'type': 'tool_result', 'tool_use_id': b.id, 'content': results_map[b.id]} for b in tool_blocks]
tool_message = {'role': 'user', 'content': tool_results}
history.append(tool_message)
MEMORY.append_history(tool_message)
blocking_results = [r['content'] for r in tool_results if isinstance(r.get('content'), str) and is_blocking_tool_result(r['content'])]
if blocking_results:
reply = '收到,老板!工具请求被运行时策略拦截,未继续改写或换路径执行。\n\n' + blocking_results[0]
assistant_message = {'role': 'assistant', 'content': reply}
history.append(assistant_message)
MEMORY.append_history(assistant_message)
print(f'[Agent回答]: {reply}\n')
break
history = compact_history(history) # 只在一整轮结束后压缩,避免拆散 tool_use/tool_result