对标 Claude Code 的生产级实现,逐模块诊断 lambdagentpaas 最亟待改善的工程短板。 按紧急程度分 P0 (必须) / P1 (应该) / P2 (建议) 三档。
lambdagentpaas 拥有同类项目中最好的理论基础 (Lambda 演算 isomorphism, 60/60 Church 编码验证),但在工程鲁棒性上与 Claude Code 存在代际差距:
| 维度 | Claude Code | lambdagentpaas | 差距等级 |
|---|---|---|---|
| 流式架构 | 全链路 Generator streaming | 全链路同步阻塞 | 致命 |
| 超时控制 | 每个工具/API 调用独立超时 | 几乎无超时 | 致命 |
| 错误恢复 | 重试+降级+fallback model | 单次调用,失败即终止 | 严重 |
| 并发安全 | isConcurrencySafe + AbortController 层级 | Par 是假并行,AsyncPar 有线程安全问题 | 严重 |
| 上下文管理 | Auto-compact + History snip + Microcompact | Sliding window (仅 ReAct) | 严重 |
| 取消机制 | AbortController 层级,WeakRef 防泄漏 | 无取消机制 | 严重 |
| 文件隔离 | Git Worktree 独立目录+分支 per agent | 无任何文件隔离,所有 Agent 共享 CWD | 严重 |
| 内存管理 | LRU cache + 磁盘持久化 + eviction | 无上限,trace 无限增长 | 中等 |
| 工具生态 | 60+ 内置工具,每个有输入校验/权限/渲染 | shell + MCP wrapper,无校验 | 中等 |
| 可观测性 | GrowthBook A/B + Telemetry + 结构化日志 | 基础 print trace | 中等 |
现状问题:
executor.py 的 reduce() 是同步函数,LLM 调用阻塞整个线程:
# executor.py - 当前实现
def _reduce_lam(self, lam, input_val, ctx):
result = self.llm_adapter.call(...) # 阻塞 30-60 秒
return result # 必须等完整响应
后果:
Claude Code 做法:
// query.ts - Generator-based streaming
async *query(params): AsyncGenerator<Message | StreamEvent> {
for await (const event of queryWithModel(...)) {
// content_block_delta → 立即 yield 给 UI
// tool_use block 还没写完就开始准备执行
yield event
}
}
改善方案:
# 目标:全链路 async generator
class Executor:
async def reduce(self, term, input_val, ctx) -> AsyncGenerator[Event, None]:
if isinstance(term, Lam):
async for token in self.llm_adapter.stream(model, prompt, input_text):
yield StreamEvent(type='token', data=token)
# 工具调用在流中解析,不等完整响应
elif isinstance(term, Compose):
for stage in term.stages:
async for event in self.reduce(stage, result, ctx):
yield event
# LLMAdapter 添加流式接口
class LLMAdapter:
async def stream(self, model, system, user, **kwargs) -> AsyncGenerator[str, None]:
if 'claude' in model:
async with client.messages.stream(...) as stream:
async for text in stream.text_stream:
yield text
工作量:大 (需重构 Executor + LLMAdapter + ReActEngine + 所有 Term.apply) 优先级:P0 — 不做流式,用户体验和服务吞吐都不可接受
现状问题:
几乎所有 I/O 操作都没有超时和重试:
# llm_adapter.py - 无超时
response = self.client.messages.create(...) # 可以挂起无限期
# mcp_client.py - 10 秒硬编码超时,无断路器
req = urllib.request.Request(url, data=payload)
resp = urllib.request.urlopen(req, timeout=10) # MCP 服务器慢就失败
# primitives.py - Lam._call_llm() 对 DashScope/Ollama 有 120 秒硬编码
urllib.request.urlopen(req, timeout=120) # 不可配置
后果:
Claude Code 做法:
// claude.ts
categorizeRetryableAPIError(error) // 区分 transient vs permanent
withRetry(fn, { maxRetries: 3, baseDelay: 1000 }) // 指数退避
FallbackTriggeredError → 自动切换 fallback model
checkQuotaStatus() → 速率限制
改善方案:
import asyncio
from dataclasses import dataclass
@dataclass
class RetryPolicy:
max_attempts: int = 3
base_delay: float = 1.0
max_delay: float = 30.0
jitter: bool = True
retryable_errors: tuple = (TimeoutError, ConnectionError, HTTPError)
async def with_retry(fn, policy: RetryPolicy):
for attempt in range(policy.max_attempts):
try:
return await asyncio.wait_for(fn(), timeout=timeout)
except policy.retryable_errors as e:
if attempt == policy.max_attempts - 1:
raise
delay = min(policy.base_delay * (2 ** attempt), policy.max_delay)
if policy.jitter:
delay *= (0.5 + random.random())
await asyncio.sleep(delay)
class CircuitBreaker:
"""MCP/LLM 调用断路器"""
def __init__(self, failure_threshold=5, reset_timeout=60):
self.failures = 0
self.state = 'closed' # closed → open → half_open
self.last_failure_time = 0
async def call(self, fn):
if self.state == 'open':
if time.time() - self.last_failure_time > self.reset_timeout:
self.state = 'half_open'
else:
raise CircuitOpenError("Circuit breaker is open")
try:
result = await fn()
self._record_success()
return result
except Exception as e:
self._record_failure()
raise
工作量:中 优先级:P0 — 没有重试和超时,任何网络问题都导致任务失败
现状问题:
一旦 Agent 开始执行,无法取消:
# 当前:GroupChat 开始后无法中断
chat = GroupChat(agents, max_rounds=100)
result = chat.apply(input, ctx) # 如果每轮 30 秒,可能运行 50 分钟
# 用户只能 Ctrl+C 杀进程
Claude Code 做法:
// 层级化取消
Main AbortController
├── Agent1 AbortController (child, WeakRef 防泄漏)
│ └── Tool AbortController (child)
│ └── Bash error → siblingAbortController.abort()
└── Agent2 AbortController (child)
// 用户按 ESC → parent.abort() → 所有子任务优雅终止
改善方案:
import asyncio
class CancellationToken:
"""层级化取消 token,等价于 AbortController"""
def __init__(self, parent: 'CancellationToken' = None):
self._cancelled = False
self._reason = None
self._children: list[weakref.ref] = []
self._callbacks: list[callable] = []
if parent:
parent._children.append(weakref.ref(self))
def cancel(self, reason: str = 'user_cancelled'):
if self._cancelled:
return
self._cancelled = True
self._reason = reason
for cb in self._callbacks:
cb(reason)
for child_ref in self._children:
child = child_ref()
if child:
child.cancel(reason)
@property
def is_cancelled(self) -> bool:
return self._cancelled
def check(self):
"""在每个 reduction 步骤前检查"""
if self._cancelled:
raise CancelledError(self._reason)
# 在 Executor 中使用
class Executor:
async def reduce(self, term, input_val, ctx, cancel_token: CancellationToken):
cancel_token.check() # 每步检查
if isinstance(term, Compose):
for stage in term.stages:
cancel_token.check()
result = await self.reduce(stage, result, ctx, cancel_token)
工作量:中 优先级:P0 — 长任务无法取消 = 不可控
现状问题:
lambdagentpaas 的所有 Agent 共享同一个工作目录 (CWD),没有任何文件系统隔离:
# shell_tool.py - 所有 Agent 的 shell 命令在同一目录执行
def shell_tool(command: str) -> str:
result = subprocess.run(command, shell=True, capture_output=True, text=True)
return result.stdout
# multiagent.py - GroupChat 中 3 个 Agent 同时写文件
agent_coder = Lam("coder", "Write code to main.py...")
agent_tester = Lam("tester", "Write tests to test_main.py...")
agent_reviewer = Lam("reviewer", "Refactor main.py...")
# coder 和 reviewer 同时修改 main.py → 文件冲突/覆盖/损坏!
后果:
main.py 的改动被 Agent B 覆盖 → 静默数据丢失pip install 不同版本 → 依赖冲突rm -rf temp/ 影响其他 Agent 的临时文件Claude Code 做法:
Claude Code 使用 Git Worktree 实现真正的文件隔离:
// worktree.ts — 每个 Agent 获得独立目录和分支
async function createAgentWorktree(slug: string) {
// 1. 校验 slug 防止路径穿越
validateWorktreeSlug(slug)
// 2. 在 .claude/worktrees/<slug>/ 创建独立目录
// 等价于完整 repo 副本
const { worktreePath, worktreeBranch } = await getOrCreateWorktree(gitRoot, slug)
// worktreeBranch = "claude-<slug>-<timestamp>" — 独立分支
// 3. 大目录 symlink 避免磁盘浪费
await symlinkDirectories(gitRoot, worktreePath, ['node_modules', '.venv'])
// 4. Agent 的 CWD 切换到 worktree 目录
return { worktreePath, worktreeBranch, headCommit, gitRoot }
}
// 检测 Agent 是否产生了变更
async function hasWorktreeChanges(worktreePath, gitRoot): boolean
// Agent 完成后清理
async function removeAgentWorktree(worktreePath, worktreeBranch, gitRoot)
隔离效果:
.claude/worktrees/<agent-slug>/node_modules 等通过 symlink 共享,不浪费磁盘git worktree remove改善方案:
lambdagentpaas 需要实现三级��件隔离,适应不同场景:
# lambdagent/isolation.py — 智能体文件隔离模块
import os
import shutil
import subprocess
import tempfile
from dataclasses import dataclass, field
from enum import Enum
from pathlib import Path
from typing import Optional
import weakref
class IsolationLevel(Enum):
"""三级隔离策略"""
NONE = "none" # 共享 CWD (仅单 Agent 场景)
DIRECTORY = "directory" # 临时目录隔离 (无 Git 项目)
WORKTREE = "worktree" # Git Worktree 隔离 (Git 项目, 推荐)
@dataclass
class IsolatedWorkspace:
"""Agent 的隔离工作空间"""
agent_id: str
workspace_path: str # Agent 的工作目录
isolation_level: IsolationLevel
branch_name: Optional[str] = None # Worktree 模式下的分支名
original_cwd: Optional[str] = None # 原始 CWD (用于回退)
git_root: Optional[str] = None
_cleanup_registered: bool = field(default=False, repr=False)
def has_changes(self) -> bool:
"""检测 Agent 是否产生了文件变更"""
if self.isolation_level == IsolationLevel.WORKTREE:
result = subprocess.run(
['git', 'status', '--porcelain'],
cwd=self.workspace_path,
capture_output=True, text=True
)
return bool(result.stdout.strip())
elif self.isolation_level == IsolationLevel.DIRECTORY:
# 对比原始目录的文件列表
return True # 临时目录默认有变更
return False
def get_diff(self) -> str:
"""获取 Agent 产生的所有变更"""
if self.isolation_level == IsolationLevel.WORKTREE:
result = subprocess.run(
['git', 'diff', 'HEAD'],
cwd=self.workspace_path,
capture_output=True, text=True
)
return result.stdout
return ""
def commit(self, message: str) -> Optional[str]:
"""在隔离分支上提交变更"""
if self.isolation_level != IsolationLevel.WORKTREE:
return None
subprocess.run(['git', 'add', '-A'], cwd=self.workspace_path, check=True)
result = subprocess.run(
['git', 'commit', '-m', message],
cwd=self.workspace_path,
capture_output=True, text=True
)
if result.returncode == 0:
# 返回 commit hash
hash_result = subprocess.run(
['git', 'rev-parse', 'HEAD'],
cwd=self.workspace_path,
capture_output=True, text=True
)
return hash_result.stdout.strip()
return None
def merge_back(self, target_branch: str = 'main') -> bool:
"""将 Agent 的变更合并回主分支"""
if self.isolation_level != IsolationLevel.WORKTREE:
return False
try:
subprocess.run(
['git', 'checkout', target_branch],
cwd=self.git_root, check=True
)
subprocess.run(
['git', 'merge', self.branch_name, '--no-ff',
'-m', f'Merge agent {self.agent_id} changes'],
cwd=self.git_root, check=True
)
return True
except subprocess.CalledProcessError:
return False
class WorkspaceManager:
"""管理所有 Agent 的隔离工作空间"""
WORKTREE_BASE = '.lambdagent/worktrees'
SYMLINK_DIRS = ['node_modules', '.venv', '__pycache__', '.tox']
def __init__(self, base_dir: str = '.'):
self.base_dir = os.path.abspath(base_dir)
self._workspaces: dict[str, IsolatedWorkspace] = {}
self._git_root = self._find_git_root()
def _find_git_root(self) -> Optional[str]:
try:
result = subprocess.run(
['git', 'rev-parse', '--show-toplevel'],
cwd=self.base_dir,
capture_output=True, text=True
)
return result.stdout.strip() if result.returncode == 0 else None
except FileNotFoundError:
return None
def _validate_slug(self, slug: str):
"""防止路径穿越攻击"""
if '..' in slug or '/' in slug or '\\' in slug:
raise ValueError(f"Invalid workspace slug: {slug}")
if len(slug) > 64:
raise ValueError(f"Slug too long: {slug}")
def create(
self,
agent_id: str,
level: IsolationLevel = IsolationLevel.WORKTREE,
slug: Optional[str] = None,
) -> IsolatedWorkspace:
"""为 Agent 创建隔离工作空间"""
slug = slug or agent_id.replace(' ', '-').lower()[:32]
self._validate_slug(slug)
if level == IsolationLevel.WORKTREE:
return self._create_worktree(agent_id, slug)
elif level == IsolationLevel.DIRECTORY:
return self._create_temp_directory(agent_id, slug)
else:
return IsolatedWorkspace(
agent_id=agent_id,
workspace_path=self.base_dir,
isolation_level=IsolationLevel.NONE,
original_cwd=self.base_dir
)
def _create_worktree(self, agent_id: str, slug: str) -> IsolatedWorkspace:
"""Git Worktree 隔离 (推荐)"""
if not self._git_root:
# 无 Git 仓库,降级为目录隔离
return self._create_temp_directory(agent_id, slug)
import time
branch_name = f"agent-{slug}-{int(time.time())}"
worktree_dir = os.path.join(self._git_root, self.WORKTREE_BASE, slug)
os.makedirs(os.path.dirname(worktree_dir), exist_ok=True)
# 创建 worktree + 新分支
subprocess.run(
['git', 'worktree', 'add', '-b', branch_name, worktree_dir],
cwd=self._git_root,
check=True,
capture_output=True
)
# Symlink 大目录避免磁盘浪费
for dirname in self.SYMLINK_DIRS:
src = os.path.join(self._git_root, dirname)
dst = os.path.join(worktree_dir, dirname)
if os.path.exists(src) and not os.path.exists(dst):
os.symlink(src, dst)
workspace = IsolatedWorkspace(
agent_id=agent_id,
workspace_path=worktree_dir,
isolation_level=IsolationLevel.WORKTREE,
branch_name=branch_name,
original_cwd=self.base_dir,
git_root=self._git_root
)
self._workspaces[agent_id] = workspace
return workspace
def _create_temp_directory(self, agent_id: str, slug: str) -> IsolatedWorkspace:
"""临时目录隔离 (无 Git 场景)"""
temp_base = os.path.join(self.base_dir, '.lambdagent', 'workspaces')
os.makedirs(temp_base, exist_ok=True)
workspace_dir = os.path.join(temp_base, slug)
if not os.path.exists(workspace_dir):
# 浅拷贝:只复制文件,不复制大目录
shutil.copytree(
self.base_dir, workspace_dir,
ignore=shutil.ignore_patterns(
'node_modules', '.venv', '__pycache__',
'.git', '.lambdagent'
)
)
# Symlink 大目录
for dirname in self.SYMLINK_DIRS:
src = os.path.join(self.base_dir, dirname)
dst = os.path.join(workspace_dir, dirname)
if os.path.exists(src) and not os.path.exists(dst):
os.symlink(src, dst)
workspace = IsolatedWorkspace(
agent_id=agent_id,
workspace_path=workspace_dir,
isolation_level=IsolationLevel.DIRECTORY,
original_cwd=self.base_dir
)
self._workspaces[agent_id] = workspace
return workspace
def cleanup(self, agent_id: str, force: bool = False):
"""清理 Agent 的工作空间"""
ws = self._workspaces.pop(agent_id, None)
if not ws:
return
if ws.isolation_level == IsolationLevel.WORKTREE:
if ws.has_changes() and not force:
raise RuntimeError(
f"Agent {agent_id} has uncommitted changes. "
f"Use force=True or commit first."
)
try:
subprocess.run(
['git', 'worktree', 'remove', '--force', ws.workspace_path],
cwd=ws.git_root, check=True, capture_output=True
)
# 删除临时分支 (如果没有合并)
subprocess.run(
['git', 'branch', '-D', ws.branch_name],
cwd=ws.git_root, capture_output=True
)
except subprocess.CalledProcessError:
pass # worktree 可能已被手动删除
elif ws.isolation_level == IsolationLevel.DIRECTORY:
if os.path.exists(ws.workspace_path):
shutil.rmtree(ws.workspace_path, ignore_errors=True)
def cleanup_all(self, force: bool = False):
"""清理所有工作空间"""
for agent_id in list(self._workspaces.keys()):
try:
self.cleanup(agent_id, force=force)
except RuntimeError:
pass # 跳过有未提交变更的
def get(self, agent_id: str) -> Optional[IsolatedWorkspace]:
return self._workspaces.get(agent_id)
def list_active(self) -> list[IsolatedWorkspace]:
return list(self._workspaces.values())
与现有模块集成:
# 1. 在 Executor 中自动创建隔离工作空间
class Executor:
def __init__(self, config, workspace_manager: WorkspaceManager = None):
self.workspace_manager = workspace_manager or WorkspaceManager()
async def reduce(self, term, input_val, ctx, cancel_token):
# 多 Agent 场景自动隔离
if isinstance(term, (Par, AsyncPar, GroupChat)):
return await self._reduce_with_isolation(term, input_val, ctx, cancel_token)
return await self._reduce_normal(term, input_val, ctx, cancel_token)
async def _reduce_with_isolation(self, term, input_val, ctx, cancel_token):
"""为每个子 Agent 创建隔离工作空间"""
agents = term.agents
workspaces = []
try:
for i, agent in enumerate(agents):
agent_id = f"{getattr(agent, 'name', f'agent-{i}')}"
ws = self.workspace_manager.create(agent_id, IsolationLevel.WORKTREE)
workspaces.append(ws)
# 并行执行,每个 Agent 在自己的 workspace 中
results = await asyncio.gather(*[
self._reduce_in_workspace(agent, input_val, ctx.fork(), cancel_token.child(), ws)
for agent, ws in zip(agents, workspaces)
])
return tuple(results)
finally:
for ws in workspaces:
self.workspace_manager.cleanup(ws.agent_id, force=True)
async def _reduce_in_workspace(self, term, input_val, ctx, cancel_token, workspace):
"""在隔离工作空间中执行 Agent"""
original_cwd = os.getcwd()
try:
os.chdir(workspace.workspace_path)
ctx.bindings['__workspace__'] = workspace # 让 Agent 能感知自己的 workspace
return await self.reduce(term, input_val, ctx, cancel_token)
finally:
os.chdir(original_cwd)
# 2. shell_tool 自动使用 workspace CWD
class IsolatedShellTool(Term):
def apply(self, input, ctx):
workspace = ctx.bindings.get('__workspace__')
cwd = workspace.workspace_path if workspace else '.'
result = subprocess.run(
input, shell=True, capture_output=True, text=True,
cwd=cwd, timeout=30
)
return result.stdout or result.stderr
# 3. GroupChat 中每个 Agent 独立 workspace
class IsolatedGroupChat(GroupChat):
async def apply(self, input, ctx, cancel_token):
workspaces = {}
try:
for agent in self.agents:
name = getattr(agent, 'name', str(id(agent)))
workspaces[name] = self.workspace_manager.create(
name, IsolationLevel.WORKTREE
)
# 每轮讨论,每个 agent 在自己的 workspace 执行
# 讨论结束后,coordinator 决定哪些变更合并回主分支
...
finally:
for ws in workspaces.values():
self.workspace_manager.cleanup(ws.agent_id, force=True)
# 4. 变更合并策略
class MergeStrategy(Enum):
AUTO_MERGE = "auto" # 自动合并无冲突的变更
COORDINATOR = "coordinator" # Coordinator Agent 决定合并
MANUAL = "manual" # 人工确认
class ChangeConsolidator:
"""合并多个 Agent 的文件变更"""
def consolidate(
self,
workspaces: list[IsolatedWorkspace],
strategy: MergeStrategy = MergeStrategy.AUTO_MERGE
) -> dict:
results = {}
for ws in workspaces:
if not ws.has_changes():
results[ws.agent_id] = {'status': 'no_changes'}
continue
ws.commit(f"Agent {ws.agent_id} changes")
if strategy == MergeStrategy.AUTO_MERGE:
success = ws.merge_back()
results[ws.agent_id] = {
'status': 'merged' if success else 'conflict',
'diff': ws.get_diff()
}
return results
YAML 配置集成:
# agent-config.yml 新增 isolation 字段
agentId: multi-coder
type: parallel
isolation:
level: worktree # none | directory | worktree
symlink: # 共享的大目录 (symlink)
- node_modules
- .venv
merge_strategy: auto # auto | coordinator | manual
cleanup: on_complete # on_complete | manual | never
agents:
- name: frontend-coder
type: react
systemPrompt: "You write frontend code..."
- name: backend-coder
type: react
systemPrompt: "You write backend code..."
工作量:中-大 (核心 WorkspaceManager ~300 行,集成 ~200 行,配置支持 ~100 行) 优先级:P0 — 多 Agent 写文件场景下,无隔离 = 数据损坏
问题 1:Par 是假并行
# extensions.py - Par.apply()
# 这不是并行!这是顺序执行!
return tuple(a.apply(input, ctx) for a in self.agents)
问题 2:AsyncPar 有线程安全问题
# multiagent.py - AsyncPar
# ctx 在多线程间共享,bindings dict 非线程安全
with ThreadPoolExecutor(max_workers=self.max_workers) as pool:
futures = {pool.submit(a.apply, input, ctx): a for a in self.agents}
# ^^^ 共享的 ctx!
问题 3:SharedMemory 无原子更新
# multiagent.py - SharedMemory.write()
def write(self, key, value):
with self._lock:
self._store[key] = value # OK: 加锁了
self._history.append(...) # OK
# 但 read + write 组合不是原子的:
# Agent A: read(x) → 10
# Agent B: read(x) → 10
# Agent A: write(x, 11)
# Agent B: write(x, 11) ← 丢失更新!
改善方案:
# Par: 用 asyncio.gather 实现真并行
class Par(Term):
async def apply(self, input, ctx, cancel_token):
tasks = [
self.executor.reduce(a, input, ctx.fork(), cancel_token.child())
for a in self.agents
]
return tuple(await asyncio.gather(*tasks))
# SharedMemory: 添加 compare-and-swap
class SharedMemory:
def compare_and_swap(self, key, expected, new_value) -> bool:
with self._lock:
if self._store.get(key) == expected:
self._store[key] = new_value
return True
return False
工作量:中 优先级:P1
现状问题:
# react_engine.py - state 每步增长,永不压缩
state = state + f"\n\nObservation: {obs}" # 线性增长!
# 100 步 × 2000 字符 = 200K 字符 → 超出任何模型上下文窗口
Sliding window 只在 from_config 编译器中实现,核心 react_engine.py 中未使用。
Claude Code 做法:
// 三层压缩机制
// 1. Auto-compact: 到达上下文 80% 时自动摘要旧消息
// 2. History snip: 裁剪最旧的消息,保留摘要
// 3. Microcompact: API 层面的消息合并
buildPostCompactMessages(messages, tokenBudget)
改善方案:
class ContextManager:
def __init__(self, max_tokens: int, model: str):
self.max_tokens = max_tokens
self.model = model
def should_compact(self, state: str) -> bool:
estimated_tokens = len(state) // 4 # 粗估
return estimated_tokens > self.max_tokens * 0.8
async def compact(self, state: str, llm: LLMAdapter) -> str:
"""保留最近 N 步原文,旧步骤摘要"""
steps = self._parse_steps(state)
if len(steps) <= 5:
return state
old_steps = steps[:-3]
recent_steps = steps[-3:]
summary = await llm.call(
model=self.model,
system="Summarize these agent steps concisely, preserving key findings and decisions.",
user="\n".join(old_steps),
max_tokens=500
)
return f"[Summary of steps 1-{len(old_steps)}]\n{summary.text}\n\n" + "\n".join(recent_steps)
工作量:中 优先级:P1 — 长 ReAct loop 必然 OOM 或超出上下文
现状问题:
# Tool 没有任何输入校验
class Tool(Term):
def apply(self, input, ctx):
return self.fn(input) # 直接执行,无校验
LLM 可能输出任何格式的 tool_input,直接传给 fn 可能导致崩溃或安全问题。
Claude Code 做法:
// 每个工具有 Zod schema
const inputSchema = z.object({
command: z.string(),
timeout: z.number().optional(),
description: z.string()
})
// 执行前校验
const parsed = tool.inputSchema.safeParse(input)
if (!parsed.success) {
return { error: `Invalid input: ${parsed.error.message}` }
}
改善方案:
from pydantic import BaseModel, ValidationError
class ToolSchema(BaseModel):
"""工具输入 schema 基类"""
class Config:
extra = 'forbid' # 拒绝未知字段
class ShellToolInput(ToolSchema):
command: str
timeout: int = 30
working_dir: str = "."
class ValidatedTool(Term):
def __init__(self, name, fn, schema: type[ToolSchema] = None):
self.name = name
self.fn = fn
self.schema = schema
def apply(self, input, ctx):
if self.schema:
try:
validated = self.schema.model_validate_json(input) if isinstance(input, str) else self.schema(**input)
return self.fn(validated.model_dump())
except ValidationError as e:
return f"Tool input validation failed: {e}"
return self.fn(input)
工作量:小 优先级:P1
现状问题:
# GroupChat 的 conversation history 每轮完整传给每个 Agent
speaker_input = f"Discussion so far:\n{state}\n\nYour turn..."
# state 包含所有前轮所有 agent 的完整输出
# 10 轮 × 3 agents × 1000 token/回复 = 30K tokens/轮
# 到第 10 轮,每次 API 调用消耗 30K tokens
改善方案:
class GroupChat(Term):
def _build_speaker_input(self, agent_name, history, round_num):
"""智能上下文窗口,避免状态爆炸"""
if len(history) <= 6:
return "\n".join(history)
# 保留:第一轮(主题) + 最近 3 轮 + 当前 agent 的所有历史发言
first_round = history[:len(self.agents)]
recent = history[-3 * len(self.agents):]
my_messages = [h for h in history if h.startswith(f"[{agent_name}]")]
return (
"[Topic]\n" + "\n".join(first_round) +
"\n\n[Your previous messages]\n" + "\n".join(my_messages[-3:]) +
"\n\n[Recent discussion]\n" + "\n".join(recent)
)
工作量:中 优先级:P1
现状问题:
# checkpoint.py - 只能保存 Context,不能保存执行位置
@dataclass
class Checkpoint:
context: dict # bindings + trace + memory
shared_memories: dict
metadata: dict
last_input: str
step_count: int # 记录了步数,但 Loop/GroupChat 不用它
如果 Agent 在 Loop 第 47 步崩溃,恢复后从第 0 步重新开始。
改善方案:
@dataclass
class ExecutionCheckpoint(Checkpoint):
"""支持中间状态恢复的增强 Checkpoint"""
execution_stack: list[StackFrame] # 执行栈快照
current_term: str # 当前执行的 Term 序列化
loop_state: dict # Loop 的 step_count + last_result
groupchat_state: dict # GroupChat 的 round + history
@dataclass
class StackFrame:
term_type: str # 'Loop', 'Compose', 'GroupChat'
term_id: str
step_index: int # 在 Compose 中的位置 / Loop 的当前步
local_state: dict # 该 frame 的局部状态
# Loop 支持恢复
class Loop(Term):
async def apply(self, input, ctx, cancel_token, resume_from: int = 0):
for step in range(resume_from, self.max_steps):
cancel_token.check()
result = await self.body.apply(input, ctx)
# 每步保存 checkpoint
if step % 5 == 0: # 每 5 步
save_checkpoint(ctx, loop_state={'step': step, 'result': result})
if self.condition(result, step):
break
return result
工作量:大 优先级:P1 — 长任务(100+ 步)没有中间恢复 = 浪费大量计算
现状:ctx.print_trace() 输出到 stdout,无结构化。
目标:
# OpenTelemetry 集成
from opentelemetry import trace as otel_trace
class InstrumentedExecutor(Executor):
async def reduce(self, term, input_val, ctx, cancel_token):
tracer = otel_trace.get_tracer("lambdagent")
with tracer.start_as_current_span(
f"reduce.{type(term).__name__}",
attributes={
"term.name": getattr(term, 'name', ''),
"input.length": len(str(input_val)),
}
) as span:
result = await super().reduce(term, input_val, ctx, cancel_token)
span.set_attribute("output.length", len(str(result)))
span.set_attribute("tokens.used", ctx.trace[-1].tokens_used)
return result
class TokenBudget:
def __init__(self, max_tokens: int, model: str):
self.max_tokens = max_tokens
self.used = 0
self.model = model
def estimate_cost(self, prompt: str) -> int:
return len(prompt) // 4 # 粗估,可接 tiktoken
def can_afford(self, prompt: str) -> bool:
return self.used + self.estimate_cost(prompt) < self.max_tokens
def record(self, input_tokens: int, output_tokens: int):
self.used += input_tokens + output_tokens
def remaining(self) -> int:
return self.max_tokens - self.used
借鉴 Claude Code 的 isConcurrencySafe:
class Tool(Term):
def __init__(self, name, fn, concurrent_safe: bool = False):
self.name = name
self.fn = fn
self.concurrent_safe = concurrent_safe
# 只读工具可并发
search = Tool("search", search_fn, concurrent_safe=True)
read_file = Tool("read", read_fn, concurrent_safe=True)
# 写工具不可并发
write_file = Tool("write", write_fn, concurrent_safe=False)
# Executor 根据标记决定并发策略
class Executor:
async def reduce_route(self, route, input_val, ctx, cancel_token):
tools = self._resolve_tools(route)
safe = [t for t in tools if t.concurrent_safe]
unsafe = [t for t in tools if not t.concurrent_safe]
# safe 并行执行,unsafe 顺序执行
import asyncio
import time
class RateLimiter:
"""令牌桶速率限制"""
def __init__(self, requests_per_minute: int):
self.rpm = requests_per_minute
self.tokens = requests_per_minute
self.last_refill = time.monotonic()
self._lock = asyncio.Lock()
async def acquire(self):
async with self._lock:
now = time.monotonic()
elapsed = now - self.last_refill
self.tokens = min(self.rpm, self.tokens + elapsed * (self.rpm / 60))
self.last_refill = now
if self.tokens < 1:
wait = (1 - self.tokens) / (self.rpm / 60)
await asyncio.sleep(wait)
self.tokens = 0
else:
self.tokens -= 1
class ResilientMCPClient:
def __init__(self, url, pool_size=5):
self.url = url
self.circuit = CircuitBreaker(failure_threshold=3, reset_timeout=30)
self._session = aiohttp.ClientSession(
connector=aiohttp.TCPConnector(limit=pool_size)
)
self._tool_cache = {} # 工具发现缓存
async def call_tool(self, name, arguments, timeout=30):
return await self.circuit.call(
lambda: with_retry(
lambda: self._do_call(name, arguments, timeout),
RetryPolicy(max_attempts=2, retryable_errors=(TimeoutError,))
)
)
Phase 1 (P0): 生存基础 [6-8 周]
├── P0-1: 全链路 async/await 重构
├── P0-1: LLMAdapter 流式接口
├── P0-2: 统一超时 + 重试 + 断路器
├── P0-3: CancellationToken 层级取消
└── P0-4: 智能体文件隔离 (WorkspaceManager + Git Worktree)
Phase 2 (P1): 可靠运行 [4-6 周]
├── P1-1: Par/AsyncPar 真并行 (asyncio.gather)
├── P1-2: Context 上下文压缩 (auto-compact)
├── P1-3: Pydantic 工具输入校验
├── P1-4: GroupChat 防状态爆炸
└── P1-5: 中间状态 Checkpoint
Phase 3 (P2): 竞争力提升 [4-6 周]
├── P2-1: OpenTelemetry 集成
├── P2-2: Token 预算管理
├── P2-3: 工具并发安全声明
├── P2-4: 速率限制器
└── P2-5: MCP 连接池 + 断路器
完整任务清单,按执行顺序排列。每项标注预估工作量和前置依赖。
[x] T01: 实现 lambdagent/async_core.py — 异步 Term 基类 2d
AsyncTerm 基类,async def apply() 签名Term.apply() 作为兼容层[x] T02: 实现 lambdagent/agentruntime/llm_adapter.py 流式接口 [3d]
async def stream() 方法 → AsyncGenerator[str, None]client.messages.stream() + text_streamclient.chat.completions.create(stream=True)stream=True 参数call() 作为兼容[x] T03: 重构 lambdagent/agentruntime/executor.py 为 async [3d]
reduce() → async def reduce()_reduce_* 方法异步化[x] T04: 重构 lambdagent/agentruntime/react_engine.py 为 async + streaming [3d]
asyncio.wait_for(tool(input), timeout=tool_timeout)[x] T05: 实现 lambdagent/retry.py — 统一重试/超时/断路器 [2d]
RetryPolicy dataclass (max_attempts, base_delay, jitter, retryable_errors)async def with_retry(fn, policy) — 指数退避 + jitterCircuitBreaker class (closed/open/half_open 三状态)[x] T06: 为所有 I/O 调用添加可配置超时 [1d]
LLMAdapter.call(): timeout 参数 (默认 120s)MCPClient.call_tool(): timeout 参数 (默认 30s)Lam._call_llm(): 移除硬编码 120s,使用 RuntimeConfig.timeoutshell_tool(): subprocess.run timeout (默认 30s)[x] T07: 实现 lambdagent/cancellation.py — CancellationToken [2d]
CancellationToken(parent=None): 层级化取消cancel(reason): 递归取消子 token (WeakRef 防泄漏)check(): 检查并抛出 CancelledErrorchild(): 创建子 token[x] T08: 实现 lambdagent/isolation.py — 智能体文件隔离 [4d]
IsolationLevel enum: NONE / DIRECTORY / WORKTREEIsolatedWorkspace dataclass: workspace_path, branch, has_changes(), get_diff(), commit(), merge_back()WorkspaceManager: create(), cleanup(), cleanup_all()[x] T09: WorkspaceManager 集成到 Executor [2d]
__workspace__ binding[x] T10: 变更合并策略 (ChangeConsolidator) [2d]
MergeStrategy: AUTO_MERGE / COORDINATOR / MANUAL[x] T11: from_config 编译器支持 isolation 字段 [1d]
isolation.level, isolation.symlink, isolation.merge_strategy[x] T12: 集成测试 — Phase 1 全部功能 [2d]
[x] T13: Par/AsyncPar 真并行 [2d]
asyncio.gather() 替代顺序 generator comprehensionctx.fork() 独立上下文[x] T14: Context 上下文压缩 (ContextManager) [3d]
should_compact(state): token 估算超过 80% 阈值compact(state, llm): 旧步骤摘要 + 保留最近 N 步[x] T15: Pydantic 工具输入校验 [2d]
ToolSchema(BaseModel) 基类ValidatedTool(name, fn, schema): apply 时自动校验[x] T16: GroupChat 防状态爆炸 [2d]
_build_speaker_input(): 智能上下文窗口[x] T17: 中间状态 Checkpoint [3d]
ExecutionCheckpoint: execution_stack + loop_state + groupchat_stateStackFrame: term_type, step_index, local_state[x] T18: OpenTelemetry 集成 [2d]
InstrumentedExecutor: 每个 reduce 步骤生成 span[x] T19: Token 预算管理 (TokenBudget) [1d]
[x] T20: 工具并发安全声明 [1d]
Tool(name, fn, concurrent_safe=False) 参数[x] T21: 速率限制器 (RateLimiter) [1d]
[x] T22: MCP 连接池 + 断路器 [2d]
ResilientMCPClient: aiohttp 连接池 + CircuitBreaker以下安全项来源于两份 SECURITYSPEC.md 中尚未完成的 TODO,按优先级排列。 标注
LSEC-xx来自 lambdagent,SEC-xx来自 agentpaas。
P1 — 高优先级安全 (阻塞生产部署)
[x] S01: 安全 lint 规则 (S001-S003) 1d
fromconfig/lint.py 新增: S001 无界 maxSteps, S002 无 Guard 的工具 Agent, S003 过宽 SandboxPolicy[x] S02: maxSteps 上限校验 0.5d
maxSteps <= 1000,超出报错[x] S03: 工具进程级超时执行 1d
[x] S04: 工具引用白名单 1d
from_config() 编译时工具引用必须来自预定义白名单,禁止任意 import[x] S05: MCP policy.mode 运行时执行 1d
mcp.policy.mode (auto/force/intelligence/disable) 当前声明但未读取,需真正执行[x] S06: PaaS 速率限制中间件 2d
TenantContext.rate_limit 的请求限流,每 API key 独立 (令牌桶)[x] S07: Secrets 注入方式修复 1d
inject.py 写入 os.environ (全局可见),改为 Context dict 注入,仅当前 Agent 可见[x] S08: API 输入大小校验 1d
RunRequest.input 添加 max_length (默认 100KB)%, _[x] S09: 移除加密 fallback 0.5d
cryptography[x] S10: HTTPS 强制中间件 0.5d
[x] S11: ToolGateway 审计日志接入 PaaS 1d
gateway.audit.entries 暴露到 /api/v1/traces/{run_id} 端点[x] S12: HIGH-risk 确认回调 webhook 2d
highRiskConfirmation: true 时通过 API/webhook 通知 tenant,暂停等待确认P2 — 中优先级安全 (纵深防御)
[ ] S13: Checkpoint 加密 1d
[ ] S14: Memory 按 Agent/Tenant 隔离 2d
scope_key = f"{tenant_id}:{agent_id}" 前缀,Redis/SQLite 表隔离[ ] S15: Channel 访问控制 1d
[x] S16: Prompt 注入防御 1d
[ ] S17: Token 用量硬限制 1d
[ ] S18: ToolGateway 路径 ACL 2d
allowed_read_paths / allowed_write_paths[ ] S19: ToolGateway 网络 ACL 1d
SandboxPolicy.network=False 时阻止工具出站网络[x] S20: API Key 哈希升级 PBKDF2 1d
[ ] S21: 安全审计日志 2d
[ ] S22: WeChat 凭证加密 0.5d
~/.agentpaas/wechat/credentials.json 用 master key 加密[x] S23: 安全响应头 0.5d
[x] S24: 请求 ID 追踪 0.5d
[ ] S25: Per-Tenant Gateway Policy 覆盖 1d
[ ] S26: ToolGateway 审计接入 OpenTelemetry 1d
P3 — 低优先级安全 (持续改进)
[ ] S27: L2 容器沙箱 5d
[ ] S28: YAML 解析器 Fuzz 测试 2d
from_config() 模糊测试[ ] S29: 多 Agent 安全审计 2d
[ ] S30: macOS sandbox-exec(1) 内核级隔离 3d
[ ] S31: 依赖漏洞扫描 0.5d
pip-audit / safety 到 CI[ ] S32: Per-Agent 工具速率限制 1d
[ ] S33: API Key 格式升级 0.5d
secrets.token_hex(16) → secrets.token_urlsafe(32)[ ] S34: WeChat token 轮换策略 1d
对标 Claude Code 60+ 内置工具,构建完整的编程助手 / 个人电脑助手能力。 每个工具是一个 Lambda 项
λx. tool(x),通过 ValidatedTool + ToolGateway 接入。
Tier 1 — 基础工具层 (MVP)
[x] A01: 结构化文件读取工具 ReadFile [1.5d]
PyPDF2 或 pdfplumber,按页范围)ReadFileSchema(file_path, offset?, limit?, pages?)[x] A02: 精确文件编辑工具 EditFile [2d]
replace_all 模式: 替换所有匹配EditFileSchema(file_path, old_string, new_string, replace_all?)[x] A03: 文件写入工具 WriteFile [0.5d]
os.makedirs)WriteFileSchema(file_path, content)[x] A04: 文件搜索工具 ListFiles (Glob) [1d]
**/*.py, src/**/*.ts)ListFilesSchema(pattern, path?, max_results?)[x] A05: 内容搜索工具 SearchContent (Grep) [1.5d]
ripgrep (rg),fallback 到纯 Python reSearchContentSchema(pattern, path?, glob?, context?, output_mode?)[x] A06: 增强 Shell 执行 Bash [2d]
run_in_background=True)BashSchema(command, timeout?, run_in_background?, working_dir?)[x] A07: Git 工作流工具集 [2d]
GitStatus: 查看工作区状态 + 未追踪文件GitDiff: 查看 staged/unstaged 变更GitLog: 查看提交历史 (--oneline, -n)GitCommit: 暂存 + 提交 (分析变更自动生成 message)GitBranch: 列出/创建/切换分支[x] A08: 内置工具注册表 + from_config 集成 [1d]
lambdagent/builtin_tools/__init__.py: 统一导出所有内置工具BUILTIN_TOOLS 注册表 (Dict[str, ValidatedTool])from_config 编译器: mcp.localTools 支持内置工具名自动解析localTools: [ReadFile, EditFile, WriteFile, Bash, terminate][x] A09: 流式终端 UI (基础版) [3d]
rich 库的终端渲染[x] A10: 集成测试 — MVP 全部工具 [2d]
Tier 2 — 智能层 (编程助手 V1)
[x] A11: 代码搜索引擎 CodeSearch [3d]
class Foo)、函数定义 (def/func/function)CodeSearchSchema(query, language?, type?, path?)[x] A12: 项目结构概览 ProjectMap [2d]
ProjectMapSchema(path?, depth?, include_summary?)[x] A13: 测试运行器 RunTests [2d]
RunTestsSchema(target?, framework?, verbose?)[x] A14: 任务管理器 TaskManager [1d]
TaskCreate(subject, description) → task_idTaskUpdate(task_id, status): pending → in_progress → completedTaskList(): 列出所有任务.lambdagent/tasks.json[x] A15: 权限审批 CLI 交互 [1.5d]
confirm_callback 实现: 终端 input("Allow? [y/n/always]")[x] A16: 项目级配置 .lambdagent.md [0.5d]
.lambdagent.mdfrom_config 编译器 + REPL 两种入口Tier 3 — 扩展层 (完整助手 V2)
[x] A17: Hook 系统 [2d]
pre_tool_call(tool_name, input) → 可修改或阻止post_tool_call(tool_name, input, output) → 可修改输出pre_llm_call(model, prompt) → prompt 注入、成本控制post_llm_call(model, prompt, response) → 输出过滤on_error(error, context) → 自定义恢复hooks.pre_tool_call: "python script.py"[x] A18: Notebook 编辑工具 NotebookEdit [3d]
jupyter_client 连接 kernel (可选)NotebookEditSchema(path, cell_index, action, content?)[x] A19: Web 搜索工具 WebSearch [1d]
WebSearchSchema(query, max_results?)[x] A20: Web 页面获取 WebFetch [1.5d]
markdownify 或 html2text)WebFetchSchema(url, selector?, max_length?)[x] A21: 模型自动 Fallback [1.5d]
[claude-sonnet, gpt-4o, qwen-max]model.fallback: [openai/gpt-4o, dashscope/qwen-max][x] A22: 流式终端 UI (完整版) [5d]
[x] A23: 集成测试 — V2 全部功能 [2d]
Tier 4 — 对标 Claude Code 关键差异能力
来源: Claude Code 33 个内置工具分析 (CC工具.md)。 以下是 A01-A23 未覆盖但对编程助手体验至关重要的能力。
[ ] A24: LSP 代码智能工具 [5d] ★★★ 关键
goto_definition(file, line, col) → 跳转到符号定义find_references(file, line, col) → 查找所有引用hover(file, line, col) → 获取类型/文档信息diagnostics(file) → 获取类型错误和警告workspace_symbols(query) → 工作区符号搜索LSPSchema(action, file_path, line?, column?, query?)[ ] A25: 用户交互工具 AskUser [1.5d] ★★ 重要
AskUser("你希望用哪种方案?")AskUser("选择框架", options=["React", "Vue", "Svelte"])AskUser("选择功能", options=[...], multi=True)AskUser("确认删除 main.py?", confirm=True)input() 交互 + rich 渲染选项列表AskUserSchema(question, options?, multi?, confirm?)[ ] A26: 规划模式 PlanMode [2d] ★★ 重要
EnterPlanMode: 切换到只读模式(禁止 Edit/Write/Bash),只允许 Read/Glob/Grep/LSP.lambdagent/plan.md 文件ExitPlanMode: 提交计划,展示给用户审批GatewayPolicy.allowed_tools 仅保留只读工具[ ] A27: Worktree 工具 (用户级) [1d]
EnterWorktree(slug) → 创建 Git Worktree,切换 CWDExitWorktree(action="keep"|"remove") → 退出,选择保留或清理WorkspaceManager.create() / cleanup()[ ] A28: 工具延迟加载 ToolSearch [2d] ★★ 重要
ToolSearch 元工具ToolSearch(query, max_results=5) → 返回匹配的工具定义 (name + schema + description)_compile_tools() 中,工具数超阈值时将 MCP 工具放入 deferred 池ToolSearchSchema(query, max_results?)[ ] A29: 定时调度 Cron [2d]
CronCreate(schedule, prompt, recurring=True) → 创建定时 promptCronDelete(cron_id) → 取消CronList() → 列出所有asyncio.create_task + 内存调度器[ ] A30: Agent Team 工具封装 [2d]
AgentSpawn(config, input, background=False) → 启动子 Agentbackground=True: 后台运行,返回 agent_idisolation="worktree": 自动创建隔离工作区AgentMessage(agent_id, message) → 向运行中的子 Agent 发消息AgentStatus(agent_id) → 查看子 Agent 状态AgentSpawnSchema(config, input, background?, isolation?)[ ] A31: Skill 工具封装 [1d]
SkillRun(skill_name, args) → 执行已注册的 Skill.lambdagent/skills/*.md 目录下的 Skill 文件lambdagent.skills.SkillRegistry[ ] A32: MCP 资源工具 [1d]
ListMcpResources(server?) → 列出 MCP Server 暴露的资源ReadMcpResource(uri) → 读取特定 MCP 资源内容MCPServer.read_resource(uri)ListMcpResourcesSchema(server?), ReadMcpResourceSchema(uri)[ ] A33: 集成测试 — 对标 CC 全部关键能力 [2d]
Tier 5 — 知识库 + OCR + 文档生成
基于已有 RAG 模块(rag.py),补充文档处理全链路能力。
[x] A34: 文档分块器 ChunkSplitter [1d]
split(text, strategy="paragraph", chunk_size=500, overlap=50) → List[str]create_rag(): create_rag(documents, chunk=True, chunk_size=500)[x] A35: OCR 工具 OCRTool [1d]
OCRTool(backend="auto"): 自动检测可用后端ocr({"file_path": "scan.png"}) → 提取的文字ocr({"file_path": "scan.pdf", "pages": "1-3"}) → PDF 逐页 OCRReadFile: 扫描型 PDF (无文本层) 自动降级到 OCROCRSchema(file_path, pages?, language?)[x] A36: 文档生成工具 DocGen [1.5d]
doc_gen({"source": "report.md", "format": "pdf", "output": "report.pdf"})DocGenSchema(source, format, output, template?)[x] A37: 知识库管理 CLI + 工具 [1.5d]
KBCreate(name, description) → 创建知识库(目录 + 索引)KBAdd(kb_name, file_path) → 添加文档(自动分块 + 索引)KBAdd(kb_name, directory, pattern="**/*.md") → 批量导入KBSearch(kb_name, query, top_k=5) → 检索KBList() → 列出所有知识库KBDelete(kb_name) → 删除.lambdagent/knowledge_bases/{name}/ 目录,JSON 索引 + SimpleVectorStore 序列化[x] A38: 集成测试 — 知识库全链路 [1d]
将 JARVIS 的核心能力抽象为 PaaS 服务,任何 Agent 通过 YAML 配置即可消费。 Agent 不需要改代码,只需
services.memory: enabled。
[x] P01: Memory Service — 持久化记忆服务 [3d] ★★★★★
POST /api/v1/memory/{agent_id}/store, GET .../recall, DELETE .../forgetMemoryStore, MemoryRecall, MemoryProfile[x] P02: Scheduler Service — 持久化调度服务 [2d] ★★★★
POST /api/v1/schedules, GET, DELETE, PATCHagent_id + input + notify_channel)~/.lambdagent/schedules.json 或 DBScheduleCreate, ScheduleList, ScheduleDelete[x] P03: Event Bus — 系统事件总线 [3d] ★★★★
POST /api/v1/events/subscribe, POST .../publishevents: [system.file.created:~/Desktop/*][x] P04: Notification Service — 多渠道通知 [2d] ★★★★
POST /api/v1/notifyNotify(message, channel, priority)[x] P05: Profile Service — 用户画像自动构建 [1.5d] ★★★
GET /api/v1/profile, PUT~/.lambdagent/user_profile.json 或 DB[x] P06: Persona Engine — 可配置人格 [2d] ★★★
persona.name, persona.style, persona.proactivity, persona.wake_word[x] P07: Multimodal Gateway — 多模态网关 [5d] ★★★
POST /api/v1/multimodal/stt, POST .../tts, POST .../visionListen, Speak, See[x] P08: Learning Service — 自适应学习 [2d] ★★
POST /api/v1/learning/feedback, GET .../strategies| 类别 | 数量 | 已完成 | 预估总工作量 |
|---|---|---|---|
| 工程改善 (T01-T22) | 22 项 | 22 ✅ | ~42 人天 |
| 安全加固 (S01-S34) | 34 项 | 16 ✅ | ~38 人天 |
| 助手工具 Tier 1-3 (A01-A23) | 23 项 | 23 ✅ | ~40 人天 |
| 助手工具 Tier 4 (A24-A33) | 10 项 | 0 | ~20 人天 |
| 知识库 Tier 5 (A34-A38) | 5 项 | 5 ✅ | ~6 人天 |
| PaaS 服务 Phase 6 (P01-P08) | 8 项 | 8 ✅ | ~20 人天 |
| Provider 统一 (L01-L05) | 5 项 | 0 | ~7 人天 |
| 双引擎切换 Phase 6.5 (E01-E07) | 7 项 | 0 | ~7 人天 |
| 集成层 Phase 7 (I01-I17) | 17 项 | 0 | ~32 人天 |
| 总计 | 131 项 | 74 ✅ | ~212 人天 |
T01 (async Term)
├── T02 (LLM streaming)
│ ├── T03 (async Executor)
│ │ ├── T04 (async ReAct)
│ │ │ ├── T12 (集成测试)
│ │ │ ├── T14 (Context 压缩)
│ │ │ └── T17 (中间 Checkpoint)
│ │ ├── T09 (Workspace 集成)
│ │ │ └── T10 (变更合并)
│ │ ├── T13 (真并行)
│ │ └── T18 (OpenTelemetry)
│ ├── T19 (Token 预算)
│ └── T21 (速率限制)
├── T05 (重试/断路器)
│ ├��─ T06 (统一超时)
│ └── T22 (MCP 连接池)
└── T07 (CancellationToken)
T08 (文件隔离) ← 独立,可与 T01-T07 并行开发
├── T09 (集成 Executor)
├── T10 (变更合并)
└── T11 (YAML 配置)
T15 (工具校验) ← 独立
T16 (GroupChat 优化) ← 依赖 T14
T20 (并发声明) ← 依赖 T13
A01-A07 (内置工具) ← 独立,可并行开发
├── A08 (工具注册表)
│ └── A10 (MVP 集成测试)
└── A09 (流式 UI) ← 依赖 T01, T02
A11 (代码搜索) ← 依赖 A05
A12 (项目概览) ← 依赖 A04
A13 (测试运行) ← 依赖 A06
A14 (任务管理) ← 独立
A15 (权限 UI) ← 独立
A17 (Hook 系统) ← 依赖 T03
A21 (模型 Fallback) ← 依赖 T05
A22 (完整 UI) ← 依赖 A09
建议并行路线:
I01 (MCP Server) ← 无依赖,可今天开始 ├── I02 (PyPI 发布) └── I03 (成本异常检测)
I04 (GitHub Action) ← 无依赖,可与 I01 并行 └── I05 (Marketplace 发布)
I06 (提取器接口) ← 无依赖 ├── I07 (LangChain 提取器) ├── I08 (CrewAI 提取器) ← 三个可并行 └── I09 (AutoGen 提取器)
I10 (守卫核心) ← 依赖 I06 ├── I11 (LangChain 守卫) ← 依赖 I07 ├── I12 (CrewAI 守卫) ← 三个可并行,依赖各自提取器 ├── I13 (AutoGen 守卫) ← 依赖 I09 └── I14 (PyPI 发布) ← 依赖 I11-I13
I15 (REST API) ← 无依赖,独立 I16 (VS Code) ← 无依赖,独立 └── I17 (Marketplace 发布)
[ ] L01: 实现 lambdagent/providers/base.py — LLMProvider 接口 [1d]
chat(messages: list[dict]) -> strLam._call_llm() 解耦[ ] L02: 实现 4 个 Provider [2d]
AnthropicProvider: client.messages.create(messages=...)OpenAICompatProvider: 覆盖 OpenAI / Ollama / DashScope / DeepSeek / MoonshotClaudeCodeProvider: 优先 --resume,fallback 到 messages 拼接[ ] L03: 实现 lambdagent/conversation.py — ConversationLam [2d]
messages: list[dict] 对话历史maxHistoryTokens 控制reset_session() 清空历史[ ] L04: from_config 编译器集成 [1d]
_compile_lam(): 根据 provider 创建对应 Provider + ConversationLammodel.conversation: true (默认开启), model.maxHistoryTokens: 80000_compress_state 依赖(ConversationLam 内部管理历史)[ ] L05: 端到端验证 + PaaS 集成 [1d]
agentpaas chat 统一走 from_config + ConversationLam目标: 将 Python 调用栈执行(Recursive)和 Agent CEK Machine 统一在同一个
Engine抽象下,通过配置切换。简单 agent 用 Recursive(低开销),复杂/高成本 agent 用 CEK(可暂停/恢复/成本熔断/精确追踪)。原则: Recursive 是默认引擎(向后兼容),CEK 是可选升级。双引擎对 125 个测试用例必须产生完全一致的结果。
[ ] E01: 实现 lambdagent/agentruntime/engine.py — Engine 抽象接口 [0.5d]
EngineMode 枚举: RECURSIVE | CEKEngine ABC: execute(term, input, **opts) -> EngineResult, execute_async(...)EngineResult 统一输出: value, trace: List[TraceRecord], cost: CostVector, steps: intfinal_state: Optional[CEKState], transitions: Optional[List[Transition]][ ] E02: 实现 RecursiveEngine — 封装现有 Executor [1d]
Executor.reduce(term, input, ctx) 为 Engine.execute() 接口term.aapply(input, ctx) 为 Engine.execute_async() 接口ctx.trace 转换为统一 TraceRecord 格式CostVector(tokens, latency, money)final_state 和 transitions 输出为 None[ ] E03: 实现 CEKEngine — 封装现有 CEKMachine [2d]
CEKMachine.run(term, input) 为 Engine.execute() 接口execute_async(): Yield 点映射为 await(LLM/Tool 异步调用)state.cost.money 与 cost_budgetInfiniteLoopDetectedTransition 列表转换为统一 TraceRecord 格式_format_kont(kont) → "compK(g, loopK(b,c,7, halt))"[ ] E04: 统一 TraceRecord 格式与转换 [1d]
TraceRecord dataclass: step, term_name, action, input_summary, output_summary, duration_ms, cost, cumulative_costcontinuation: Optional[str](K 栈描述), yield_type: Optional[str]("llm"|"tool"|None)RecursiveEngine._convert_trace(): ctx.trace → List[TraceRecord]CEKEngine._convert_transitions(): List[Transition] → List[TraceRecord](跳过 τ 静默转移或可选保留)[ ] E05: from_config 编译器支持 runtime.engine 字段 [0.5d]
YAML schema 新增:
runtime:
engine: recursive | cek # 默认 recursive
costBudget: 5.00 # CEK 独有: 逐步成本熔断
maxSteps: 10000 # CEK 独有: 最大转移步数
enablePauseResume: true # CEK 独有: 允许暂停/恢复
Runtime.__init__() 根据 engine 字段选择 RecursiveEngine 或 CEKEngineruntime.engine 字段时默认 RecursiveEngine依赖: E02, E03
[ ] E06: 实现 AdaptiveEngine — 自动选择引擎 [1d]
分析 term 结构评估复杂度: _assess_complexity(term) -> Complexity
决策规则:
type: simple 或管道 ≤5 步 → Recursive
maxSteps > 10 或含 Pair/AsyncPar 或 costBudget 已设 → CEK
runtime.engine: auto 时使用 AdaptiveEngine
依赖: E02, E03
[ ] E07: 双引擎一致性测试 [1d]
全部 125 个语义忠实性测试: RecursiveEngine 和 CEKEngine 各跑一遍,比较最终值完全一致
全部 60 个 Church 编码测试: 双引擎结果一致
成本向量对比: CEK 精确累积 vs Recursive 事后汇总,差异 < 1%
并行场景 (Pair/AsyncPar): 双引擎结果分布一致
Guard 重试场景: 双引擎重试次数和最终值一致
性能对比: 记录 Recursive vs CEK 的 overhead(预期 CEK 慢 5-15%,可接受)
依赖: E04, E05
核心策略: 不要求用户换框架。用户继续用 LangChain/CrewAI/AutoGen,lambdagentpaas 作为静态分析层介入,输出类型错误/成本预测/死循环风险/并行冲突。详见
docs/integration-strategy.md。真实案例支撑: 38+ 个 GitHub issue 验证了每一项静态检查的价值,包括 Claude Code #38029 ($342 单 session)、#34629 (10-20x 成本 28 天)、AutoGen #108 (空消息死循环)、LangChain #10997 (Dict vs Str 类型崩溃)。详见
REAL_WORLD_DEFECTS.md。
lambdagent/mcp_server.py — MCP Server 入口 [2d]lint_agent_config, estimate_agent_cost, check_agent_types, check_parallel_safetylint_agent_config: 调用已有 lint_config(), 支持 LangChain/CrewAI/AutoGen/Dify/通用 YAMLestimate_agent_cost: 调用已有 compute_grade(), 返回 (p, t, l, m) 分级上界 + 逐阶段分解check_agent_types: 调用已有 type_check(), 返回 T-Compose 错误列表 + 修复建议check_parallel_safety: 调用已有 check_store_independence(), 返回 store 冲突列表依赖: 无 (所有底层函数已存在)
[ ] I02: MCP Server 发布为独立 PyPI 包 lambdagent-mcp-server [1d]
pyproject.toml 配置 [project.scripts] 入口: lambdagent-mcp-server = "lambdagent.mcp_server:main"
支持 uvx lambdagent-mcp-server 一键启动
编写 README: 一行配置接入 Claude Code / Cursor / Windsurf
依赖: I01
[ ] I03: 成本异常运行时检测工具 [1d]
新 MCP 工具 monitor_agent_cost: 对比运行时实际 token 与分级类型预期值
偏差超 2x 时返回 COST_ANOMALY 警告 (对应 Claude Code #34629 场景)
依赖: I01
lambdagent/agent-lint-action — GitHub Action [2d]lambdagent lint --format jsonpaths, fail-on (error|warn|info), frameworks (auto|指定), cost-threshold, type-check, parallel-check**/agent*.yml, **/crew*.yml, **/autogen*.json 等依赖: CLI lambdagent lint (已存在)
[ ] I05: 发布到 GitHub Marketplace [1d]
action.yml ��数据 + Docker 构建
Marketplace listing: 名称 lambdagent/agent-lint-action
README: 示例 workflow 文件
依赖: I04
lambdagent/extractors/base.py — 提取器接口 [0.5d]FrameworkExtractor ABC: detect(obj) -> bool, extract(obj) -> dictextract_config(obj) 自动检测框架并返回归一化配置依赖: 无
[ ] I07: 实现 lambdagent/extractors/langchain_extractor.py [2d]
从 AgentExecutor 提取: model name/temperature, tools, max_iterations, prompt template
映射到 lambdagent YAML 格式 (type: react)
检测 terminate 等效物 (AgentFinish, early_stopping_method)
处理 LangChain LCEL Runnable 链 (RunnableSequence → type: chain)
依赖: I06
[ ] I08: 实现 lambdagent/extractors/crewai_extractor.py [2d]
从 Crew 提取: agents (role/goal/backstory → systemPrompt), tasks, process type
映射: sequential → type: chain, hierarchical → type: router, 其他 → type: parallel
提取 max_iter, tools, LLM config
依赖: I06
[ ] I09: 实现 lambdagent/extractors/autogen_extractor.py [2d]
从 GroupChatManager 提取: agents (system_message, llm_config), max_round, is_termination_msg
映射到 type: parallel + multiagent.maxRounds
检测终止条件类型 (字符串匹配 vs 函数 — 对应 AutoGen #108 风险)
依赖: I06
lambdagent_guard/core.py — 守卫核心 [2d]GuardedExecutor: 编译时检查 (type_check + cost_budget + store_independence) + 运行时 hookCostBudgetExceeded, InfiniteLoopDetected, StoreConflictError 异常类依赖: I06
[ ] I11: 实现 lambdagent_guard/langchain.py — LangChain 守卫 [2d]
guard_langchain(executor, cost_budget, type_check, loop_detection, cost_alert)
Hook 进 AgentExecutor 的 _take_next_step 或 agent.plan 方法
非侵入式: 用户代码只加 2 行
对应真实缺陷: LangChain #10997 (类型崩溃), #2495 (死循环), #26019 (tool 重复调用)
依赖: I07, I10
[ ] I12: 实现 lambdagent_guard/crewai.py — CrewAI 守卫 [2d]
guard_crewai(crew, cost_budget, parallel_safety, terminate_check)
Hook 进 Crew.kickoff() 执行流
对应真实缺陷: CrewAI #3847 (max_iter 无效), #737 (tool 结果不回传), #1355 (10x 调用)
依赖: I08, I10
[ ] I13: 实现 lambdagent_guard/autogen.py — AutoGen 守卫 [2d]
guard_autogen(manager, cost_budget, empty_message_detection, terminate_robustness)
Hook 进 GroupChatManager 消息循环
空消息检测: 连续 N 条消息长度 < 10 → 告警 (AutoGen #108)
终止鲁棒性: 不依赖精确字符串匹配,改用 substring + case-insensitive (AutoGen #391)
依赖: I09, I10
[ ] I14: 发布 lambdagent-guard PyPI 包 [1d]
包含 lambdagent_guard/ 下所有模块
框架依赖为 optional: pip install lambdagent-guard[langchain], [crewai], [autogen]
README: 每种框架的 2 行接入示例
依赖: I11, I12, I13
/api/v1/analyze/* 端点 [1d]POST /api/v1/analyze/lint: 接受 YAML 字符串或 dict, 返回 lint 结果POST /api/v1/analyze/type-check: 返回 T-Compose 类型检查结果POST /api/v1/analyze/cost: 返回分级类型成本预测POST /api/v1/analyze/parallel-safety: 返回 store independence 检查POST /api/v1/analyze/full: 合并以上四项, 一次调用返回完整分析framework 参数 (auto|langchain|crewai|autogen|dify)lambdagent.agent-lint [3d]λA: 2 errors | Est. cost: $6.00/run | Success: 0.16%lambdagent lint --format json依赖: CLI (已存在)
[ ] I17: 发布到 VS Code Marketplace + JetBrains Marketplace [1d]
VS Code: vsce package + vsce publish
JetBrains: 同功能的 IntelliJ 插件 (基于 External Annotator API)
依赖: I16
```
Phase 1 (P0): 生存基础 [6-8 周] ├── 路线 A: 全链路 async + 流式 + 超时重试 + 取消 ├── 路线 B: 智能体文件隔离 (WorkspaceManager) └── 路线 D 前半: S01-S10 安全独立项 (并行)
Phase 2 (P1): 可靠运行 [4-6 周] ✅ 已完成 ├── T13-T17: 真并行 + 上下文压缩 + 校验 + Checkpoint └── S11-S12, S13-S25: 安全纵深防御 (部分完成)
Phase 3 (P2+P3): 竞争力 + 持续安全 [4-6 周] ✅ 已完成 ├── T18-T22: OTel + Token预算 + 速率限制 + MCP连接池 └── S26-S34: 容器沙箱 + Fuzz + 漏洞扫描 + 审计 (部分 deferred)
Phase 4 (安全): 安全加固 [2-3 周] ✅ 部分完成 ├── S01-S12: P1 高优安全 (已完成) └── S13-S34: P2/P3 中低优安全 (部分 deferred)
Phase 5 (助手): 个人助手 & 编程助手 [8-12 周] ├── Tier 1 MVP: A01-A10 (文件工具 + Shell + Git + UI) ~15 人天 ├── Tier 2 V1: A11-A16 (代码搜索 + 测试 + 任务 + 权限) ~10 人天 ├── Tier 3 V2: A17-A23 (Hook + Notebook + Web + Fallback) ~15 人天 └── Tier 4 CC: A24-A33 (LSP + 规划模式 + Agent Team + 延迟加载) ~20 人天
Phase 6.5 (双引擎): Recursive / CEK 可切换执行引擎 [1-2 周] ├── E01: Engine 抽象接口 + EngineResult ~0.5 人天 ├── E02+E03: RecursiveEngine + CEKEngine (并行开发) ~3 人天 ├── E04: 统一 TraceRecord 格式 ~1 人天 ├── E05: from_config 支持 runtime.engine 字段 ~0.5 人天 ├── E06: AdaptiveEngine 自动选择 ~1 人天 └── E07: 双引擎 125+60 测试一致性验证 ~1 人天
Phase 7 (集成层): lambdagent 作为外部框架插件 [4-6 周] ├── 7a MCP Server: I01-I03 (AI IDE 集成, 最高优先级) ~4 人天 ├── 7b GitHub Action: I04-I05 (CI/CD 守门, 与 7a 并行) ~3 人天 ├── 7c 提取器: I06-I09 (LangChain/CrewAI/AutoGen 配置提取) ~6.5 人天 ├── 7d 运行时守卫: I10-I14 (Python 中间件, 依赖 7c) ~9 人天 ├── 7e REST API: I15 (分析端点, 独立) ~1 人天 └── 7f IDE 扩展: I16-I17 (VS Code/JetBrains, 独立) ~4 人天 ```
总计: 131 项 TODO (22 工程 + 34 安全 + 33 助手 + 5 Provider + 7 双引擎 + 17 集成 + 13 研究), 已完成 74 项, 预估 ~212 人天
lambdagentpaas 的理论层是同类最优。Phase 1-4 已完成核心工程基础设施(async/流式/重试/取消/隔离/安全),从"研究原型"升级为"可用框架"。Phase 5 聚焦 33 个内置工具,对标 Claude Code 全部关键能力。Phase 7 (集成层) 是商业化关键——不要求用户换框架,而是作为静态分析插件接入 LangChain/CrewAI/AutoGen 生态。MCP Server (I01-I03, 4 人天) 是最高杠杆任务:300 行代码让所有 AI IDE 获得 λA 类型检查、成本预测和并行安全验证能力。38+ 个真实 GitHub issue(含 Claude Code $342 单 session、10-20x 成本 28 天等)证明每一项静态检查都有对应的真实需求。