ENGINEERING_GAP_ANALYSIS.md 87 KB

LambdagentPaaS 工程化改善分析

对标 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 中等

二、P0 — 致命短板 (不解决无法生产使用)

P0-1: 全链路异步化 + 流式输出

现状问题

executor.pyreduce() 是同步函数,LLM 调用阻塞整个线程:

# executor.py - 当前实现
def _reduce_lam(self, lam, input_val, ctx):
    result = self.llm_adapter.call(...)  # 阻塞 30-60 秒
    return result                         # 必须等完整响应

后果

  • 10 步 ReAct loop × 30 秒/步 = 5 分钟用户看不到任何输出
  • PaaS 服务中一个请求阻塞一个 Worker 线程
  • 无法实现"边生成边执行工具"

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 — 不做流式,用户体验和服务吞吐都不可接受


P0-2: 超时 + 重试 + 断路器

现状问题

几乎所有 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)  # 不可配置

后果

  • 网络抖动 → 任务彻底失败
  • MCP 服务器慢 → 整个 Agent 卡死
  • 无法区分"暂时故障"和"永久故障"

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 — 没有重试和超时,任何网络问题都导致任务失败


P0-3: 取消机制 (AbortController 等价物)

现状问题

一旦 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 — 长任务无法取消 = 不可控


P0-4: 智能体文件隔离

现状问题

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 → 文件冲突/覆盖/损坏!

后果

  • Agent A 写入 main.py 的改动被 Agent B 覆盖 → 静默数据丢失
  • 多 Agent 并行 pip install 不同版本 → 依赖冲突
  • Agent 执行 rm -rf temp/ 影响其他 Agent 的临时文件
  • 无法回滚单个 Agent 的文件修改(全在一个 git 状态中)
  • PaaS 多租户场景下,不同 tenant 的 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)

隔离效果:

  • 每个 Agent 有独立目录: .claude/worktrees/<agent-slug>/
  • 每个 Agent 有独立 Git 分支: 可独立 commit,不影响主分支
  • 共享大目录: node_modules 等通过 symlink 共享,不浪费磁盘
  • 变更检测: Agent 完成后可检查是否有文件变更
  • 自动清理: Agent 结束后 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 写文件场景下,无隔离 = 数据损坏


三、P1 — 严重短板 (影响可靠性和扩展性)

P1-1: 并发安全修复

问题 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


P1-2: 上下文窗口管理

现状问题

# 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 或超出上下文


P1-3: 工具输入校验

现状问题

# 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


P1-4: GroupChat 防死锁 + 状态爆炸

现状问题

# 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


P1-5: Checkpoint 支持中间状态恢复

现状问题

# 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+ 步)没有中间恢复 = 浪费大量计算


四、P2 — 建议改善 (提升竞争力)

P2-1: 结构化可观测性

现状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

P2-2: Token 预算管理

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

P2-3: 工具并发安全声明

借鉴 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 顺序执行

P2-4: 速率限制器

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

P2-5: MCP 连接池 + 断路器

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 连接池 + 断路器

六、TODO LIST

完整任务清单,按执行顺序排列。每项标注预估工作量和前置依赖。

Phase 1: P0 生存基础

  • [x] T01: 实现 lambdagent/async_core.py — 异步 Term 基类 2d

    • 新增 AsyncTerm 基类,async def apply() 签名
    • 所有 Term 子类迁移为 async (Lam, Compose, Loop, If, Route, Par, Guard, Memory)
    • 保留同步 Term.apply() 作为兼容层
  • [x] T02: 实现 lambdagent/agentruntime/llm_adapter.py 流式接口 [3d]

    • 新增 async def stream() 方法 → AsyncGenerator[str, None]
    • Anthropic: client.messages.stream() + text_stream
    • OpenAI: client.chat.completions.create(stream=True)
    • DashScope: SSE 流式 HTTP
    • Ollama: stream=True 参数
    • 保留同步 call() 作为兼容
  • [x] T03: 重构 lambdagent/agentruntime/executor.py 为 async [3d]

    • reduce()async def reduce()
    • 所有 _reduce_* 方法异步化
    • 依赖: T01, T02
  • [x] T04: 重构 lambdagent/agentruntime/react_engine.py 为 async + streaming [3d]

    • 7-phase loop 异步化
    • 工具调用超时: asyncio.wait_for(tool(input), timeout=tool_timeout)
    • 依赖: T03
  • [x] T05: 实现 lambdagent/retry.py — 统一重试/超时/断路器 [2d]

    • RetryPolicy dataclass (max_attempts, base_delay, jitter, retryable_errors)
    • async def with_retry(fn, policy) — 指数退避 + jitter
    • CircuitBreaker class (closed/open/half_open 三状态)
    • 集成到 LLMAdapter 和 MCPClient
  • [x] T06: 为所有 I/O 调用添加可配置超时 [1d]

    • LLMAdapter.call(): timeout 参数 (默认 120s)
    • MCPClient.call_tool(): timeout 参数 (默认 30s)
    • Lam._call_llm(): 移除硬编码 120s,使用 RuntimeConfig.timeout
    • shell_tool(): subprocess.run timeout (默认 30s)
    • 依赖: T05
  • [x] T07: 实现 lambdagent/cancellation.py — CancellationToken [2d]

    • CancellationToken(parent=None): 层级化取消
    • cancel(reason): 递归取消子 token (WeakRef 防泄漏)
    • check(): 检查并抛出 CancelledError
    • child(): 创建子 token
    • 集成到 Executor.reduce(), Loop, GroupChat, AsyncPar
  • [x] T08: 实现 lambdagent/isolation.py — 智能体文件隔离 [4d]

    • IsolationLevel enum: NONE / DIRECTORY / WORKTREE
    • IsolatedWorkspace dataclass: workspace_path, branch, has_changes(), get_diff(), commit(), merge_back()
    • WorkspaceManager: create(), cleanup(), cleanup_all()
    • Git Worktree 创建 + symlink 大目录 + 路径穿越校验
    • 临时目录隔离 (无 Git 场景降级)
    • 依赖: 无
  • [x] T09: WorkspaceManager 集成到 Executor [2d]

    • Par/AsyncPar/GroupChat 自动为每个 Agent 创建 workspace
    • shell_tool 自动使用 workspace CWD
    • Context 注入 __workspace__ binding
    • 依赖: T08, T03
  • [x] T10: 变更合并策略 (ChangeConsolidator) [2d]

    • MergeStrategy: AUTO_MERGE / COORDINATOR / MANUAL
    • 自动合并无冲突变更
    • 冲突检测 + 报告
    • 集成到 GroupChat/Par 的 finally 清理流程
    • 依赖: T08, T09
  • [x] T11: from_config 编译器支持 isolation 字段 [1d]

    • YAML schema 新增 isolation.level, isolation.symlink, isolation.merge_strategy
    • 编译器将 isolation 配置传递给 Executor
    • ���赖: T08
  • [x] T12: 集成测试 — Phase 1 全部功能 [2d]

    • 异步 executor 端到端测试
    • 流式输出验证
    • 超时/重试场景测试
    • 取消传播测试
    • 文件隔离: 多 Agent 同时写文件不冲突
    • Worktree 创建/清理/合并测试

Phase 2: P1 可靠运行

  • [x] T13: Par/AsyncPar 真并行 [2d]

    • Par: asyncio.gather() 替代顺序 generator comprehension
    • AsyncPar: 改用 asyncio (而非 ThreadPoolExecutor),解决 ctx 线程安全问题
    • 每个并行 Agent 获得 ctx.fork() 独立上下文
    • 依���: T01, T08
  • [x] T14: Context 上下文压缩 (ContextManager) [3d]

    • should_compact(state): token 估算超过 80% 阈值
    • compact(state, llm): 旧步骤摘要 + 保留最近 N 步
    • 集成到 ReActEngine 的 UPDATE 阶段
    • 依赖: T02, T04
  • [x] T15: Pydantic 工具输入校验 [2d]

    • ToolSchema(BaseModel) 基类
    • ValidatedTool(name, fn, schema): apply 时自动校验
    • 内置 shell_tool, terminate 等添加 schema
    • 校验失败返回错误提示,不崩溃
    • 依赖: 无
  • [x] T16: GroupChat 防状态爆炸 [2d]

    • _build_speaker_input(): 智能上下文窗口
    • 保留 topic + 自己的发言 + 最近 N 轮
    • 旧轮摘要压缩
    • 依赖: T14
  • [x] T17: 中间状态 Checkpoint [3d]

    • ExecutionCheckpoint: execution_stack + loop_state + groupchat_state
    • StackFrame: term_type, step_index, local_state
    • Loop/GroupChat 每 N 步自动保存
    • 恢复时从 checkpoint.step_index 继续,不从头
    • 依赖: T03, T04

Phase 3: P2 竞争力提升 + 安全加固

  • [x] T18: OpenTelemetry 集成 [2d]

    • InstrumentedExecutor: 每个 reduce 步骤生成 span
    • Attributes: term.name, input.length, tokens.used, duration_ms
    • 依赖: T03
  • [x] T19: Token 预算管理 (TokenBudget) [1d]

    • estimate_cost(), can_afford(), record(), remaining()
    • 集成到 Executor: 预算耗尽时优雅停止
    • 依赖: T02
  • [x] T20: 工具并发安全声明 [1d]

    • Tool(name, fn, concurrent_safe=False) 参数
    • Executor Route 按安全标记分组: safe 并行, unsafe 顺序
    • 依赖: T13
  • [x] T21: 速率限制器 (RateLimiter) [1d]

    • 令牌桶算法
    • 集成到 LLMAdapter: 每次调用前 acquire()
    • 依赖: T02
  • [x] T22: MCP 连接池 + 断路器 [2d]

    • ResilientMCPClient: aiohttp 连接池 + CircuitBreaker
    • 工具发现缓存 (避免重复 list_tools)
    • 依赖: T05

Phase 4: 安全加固 (来源: lambdagent/SECURITYSPEC.md + agentpaas/SECURITYSPEC.md)

以下安全项来源于两份 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

    • 所有 Tool 调用可配置 timeout (默认 30s),进程级而非仅 asyncio
    • 依赖: T06
  • [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)
    • Agent config 字段大小限制
    • SQL LIKE 查询转义通配符 %, _
    • 依赖: 无
  • [x] S09: 移除加密 fallback 0.5d

    • 当前加密失败 fallback 到 base64 (等于无加密),生产模式强制 cryptography
    • 依赖: 无
  • [x] S10: HTTPS 强制中间件 0.5d

    • 生产模式下强制 HTTPS,拒绝 HTTP 请求
    • 依赖: 无
  • [x] S11: ToolGateway 审计日志接入 PaaS 1d

    • gateway.audit.entries 暴露到 /api/v1/traces/{run_id} 端点
    • 依赖: 无
  • [x] S12: HIGH-risk 确认回调 webhook 2d

    • highRiskConfirmation: true 时通过 API/webhook 通知 tenant,暂停等待确认
    • 依赖: S06

P2 — 中优先级安全 (纵深防御)

  • [ ] S13: Checkpoint 加密 1d

    • Checkpoint JSON 使用 AES-256-GCM 加密,文件权限 0600
    • 依赖: 无
  • [ ] S14: Memory 按 Agent/Tenant 隔离 2d

    • 添加 scope_key = f"{tenant_id}:{agent_id}" 前缀,Redis/SQLite 表隔离
    • 依赖: T08
  • [ ] S15: Channel 访问控制 1d

    • Agent 只能访问显式声明的 Channel,创建时绑定 agent_id 列表
    • 依赖: 无
  • [x] S16: Prompt 注入防御 1d

    • 默认 system prompt 加入注入抵抗指令,工具输出标记为 untrusted data
    • 依赖: 无
  • [ ] S17: Token 用量硬限制 1d

    • 每次执行上下文 token 总量硬上限,超限优雅终止
    • 依赖: T19
  • [ ] S18: ToolGateway 路径 ACL 2d

    • 执行 SandboxPolicy 的 allowed_read_paths / allowed_write_paths
    • 依赖: 无
  • [ ] S19: ToolGateway 网络 ACL 1d

    • SandboxPolicy.network=False 时阻止工具出站网络
    • 依赖: 无
  • [x] S20: API Key 哈希升级 PBKDF2 1d

    • 从裸 SHA-256 升级为 PBKDF2 + salt
    • 依赖: 无
  • [ ] S21: 安全审计日志 2d

    • 记录: key 创建/吊销, tenant 创建, agent 部署, secret 访问, 认证失败
    • 依赖: 无
  • [ ] S22: WeChat 凭证加密 0.5d

    • ~/.agentpaas/wechat/credentials.json 用 master key 加密
    • 依赖: 无
  • [x] S23: 安全响应头 0.5d

    • HSTS, X-Frame-Options, Content-Security-Policy, X-Content-Type-Options
    • 依赖: 无
  • [x] S24: 请求 ID 追踪 0.5d

    • 每请求分配唯一 ID,贯穿日志链路
    • 依赖: 无
  • [ ] S25: Per-Tenant Gateway Policy 覆盖 1d

    • Admin 可强制全 tenant 使用 strict() 策略
    • 依赖: S11
  • [ ] S26: ToolGateway 审计接入 OpenTelemetry 1d

    • 审计日志作为 OTel span 发出
    • 依赖: T18

P3 — 低优先级安全 (持续改进)

  • [ ] S27: L2 容器沙箱 5d

    • Docker/cgroup 完全隔离不信任 Agent
    • 依赖: T08
  • [ ] S28: YAML 解析器 Fuzz 测试 2d

    • from_config() 模糊测试
    • 依赖: 无
  • [ ] S29: 多 Agent 安全审计 2d

    • 审计 Channel/GroupChat/AsyncPar 安全边界 (共享 ctx 线程安全, Channel 泄漏, Handoff 劫持)
    • 依赖: T13
  • [ ] S30: macOS sandbox-exec(1) 内核级隔离 3d

    • 替代 Python monkey-patch fallback
    • 依赖: 无
  • [ ] S31: 依赖漏洞扫描 0.5d

    • 集成 pip-audit / safety 到 CI
    • 依赖: 无
  • [ ] S32: Per-Agent 工具速率限制 1d

    • 每分钟每 Agent 工具调用上限
    • 依赖: S06
  • [ ] S33: API Key 格式升级 0.5d

    • secrets.token_hex(16)secrets.token_urlsafe(32)
    • 依赖: 无
  • [ ] S34: WeChat token 轮换策略 1d

    • 自动刷新 WeChat bot token
    • 依赖: 无

Phase 5: 个人助手 & 编程助手工具集

对标 Claude Code 60+ 内置工具,构建完整的编程助手 / 个人电脑助手能力。 每个工具是一个 Lambda 项 λx. tool(x),通过 ValidatedTool + ToolGateway 接入。

Tier 1 — 基础工具层 (MVP)

  • [x] A01: 结构化文件读取工具 ReadFile [1.5d]

    • 行号显示、offset/limit 分页读取、编码自动检测
    • 大文件安全(>2000 行自动截断)
    • PDF 读取(基于 PyPDF2pdfplumber,按页范围)
    • 图片识别(返回路径 + 元信息,配合多模态 LLM 描述)
    • Jupyter Notebook (.ipynb) 解析(cell 内容 + 输出)
    • Schema: ReadFileSchema(file_path, offset?, limit?, pages?)
    • 依赖: 无
  • [x] A02: 精确文件编辑工具 EditFile [2d]

    • 基于 old_string → new_string 的精确替换
    • 唯一性校验: old_string 在文件中必须唯一,否则报错
    • replace_all 模式: 替换所有匹配
    • 编辑前自动备份 (.bak)
    • 保留原始缩进和行尾
    • Schema: EditFileSchema(file_path, old_string, new_string, replace_all?)
    • 依赖: A01 (需先读取才能编辑)
  • [x] A03: 文件写入工具 WriteFile [0.5d]

    • 创建新文件或覆盖已有文件
    • 自动创建中间目录 (os.makedirs)
    • 安全检查: 拒绝写入 .env、credentials 等敏感文件名
    • Schema: WriteFileSchema(file_path, content)
    • 依赖: 无
  • [x] A04: 文件搜索工具 ListFiles (Glob) [1d]

    • 基于 glob 模式搜索 (**/*.py, src/**/*.ts)
    • 按修改时间排序
    • 自动忽略 .git, node_modules, pycache, .venv
    • 限制返回数量 (默认 100)
    • Schema: ListFilesSchema(pattern, path?, max_results?)
    • 依赖: 无
  • [x] A05: 内容搜索工具 SearchContent (Grep) [1.5d]

    • 正则表达式搜索文件内容
    • 上下文行 (-A, -B, -C)
    • 文件类型过滤 (--type py, js, etc.)
    • 输出模式: content / files_only / count
    • 多行模式 (跨行匹配)
    • 底层优先使用 ripgrep (rg),fallback 到纯 Python re
    • Schema: SearchContentSchema(pattern, path?, glob?, context?, output_mode?)
    • 依赖: 无
  • [x] A06: 增强 Shell 执行 Bash [2d]

    • 工作目录跨命令持久化 (session CWD)
    • 后台运行模式 (run_in_background=True)
    • 环境变量继承自用户 shell profile
    • 智能输出截断 (保留头尾,中间省略)
    • 交互命令检测与拒绝 (vim, less, ssh 等)
    • 可配置超时 (默认 120s,最大 600s)
    • Schema: BashSchema(command, timeout?, run_in_background?, working_dir?)
    • 依赖: 无 (增强现有 ShellTool)
  • [x] A07: Git 工作流工具集 [2d]

    • GitStatus: 查看工作区状态 + 未追踪文件
    • GitDiff: 查看 staged/unstaged 变更
    • GitLog: 查看提交历史 (--oneline, -n)
    • GitCommit: 暂存 + 提交 (分析变更自动生成 message)
    • GitBranch: 列出/创建/切换分支
    • 所有操作通过 ToolGateway 权限管控 (push/force 为 HIGH 风险)
    • Schema: 每个子工具独立 Schema
    • 依赖: 无 (封装 subprocess git 命令)
  • [x] A08: 内置工具注册表 + from_config 集成 [1d]

    • lambdagent/builtin_tools/__init__.py: 统一导出所有内置工具
    • BUILTIN_TOOLS 注册表 (Dict[str, ValidatedTool])
    • from_config 编译器: mcp.localTools 支持内置工具名自动解析
    • YAML 配置: localTools: [ReadFile, EditFile, WriteFile, Bash, terminate]
    • 依赖: A01-A07
  • [x] A09: 流式终端 UI (基础版) [3d]

    • 基于 rich 库的终端渲染
    • Token 逐字流式输出 (消费 AsyncReActEngine 的 StreamEvent)
    • 工具调用展示 (折叠模式: 显示工具名 + 耗时,展开显示输入/输出)
    • Spinner 进度条 (思考中 / 执行工具中)
    • 彩色 β-reduction 追踪
    • 依赖: T01 (async), T02 (streaming)
  • [x] A10: 集成测试 — MVP 全部工具 [2d]

    • ReadFile: 普通文件/大文件/编码/不存在的文件
    • EditFile: 精确替换/唯一性失败/replace_all
    • WriteFile: 新建/覆盖/敏感文件拒绝
    • ListFiles: glob 匹配/排序/忽略规则
    • SearchContent: 正则/上下文/类型过滤
    • Bash: 普通命令/超时/后台/危险命令拦截
    • Git: status/diff/commit/branch
    • 端到端: ReAct Agent 使用内置工具完成编程任务
    • 依赖: A01-A08

Tier 2 — 智能层 (编程助手 V1)

  • [x] A11: 代码搜索引擎 CodeSearch [3d]

    • 语义代码搜索: 类定义 (class Foo)、函数定义 (def/func/function)
    • 基于 ripgrep + 语言感知正则 (不依赖 AST)
    • 支持 Python, TypeScript/JS, Go, Java, Rust 语法模式
    • 符号引用查找 (grep 变体)
    • Schema: CodeSearchSchema(query, language?, type?, path?)
    • 依赖: A05
  • [x] A12: 项目结构概览 ProjectMap [2d]

    • 递归生成目录树 (忽略 .git/node_modules 等)
    • 关键文件自动识别 (README, package.json, pyproject.toml, Makefile 等)
    • 文件统计 (按语言分类,行数统计)
    • 可选: LLM 摘要 (对每个关键文件生成一行描述)
    • Schema: ProjectMapSchema(path?, depth?, include_summary?)
    • 依赖: A04
  • [x] A13: 测试运行器 RunTests [2d]

    • 自动检测测试框架 (pytest, jest, go test, cargo test)
    • 执行测试 + 结构化输出 (passed/failed/error 计数)
    • 失败用例详情提取 (文件、行号、错误消息)
    • 单文件/目录/全项目 三种粒度
    • Schema: RunTestsSchema(target?, framework?, verbose?)
    • 依赖: A06
  • [x] A14: 任务管理器 TaskManager [1d]

    • TaskCreate(subject, description) → task_id
    • TaskUpdate(task_id, status): pending → in_progress → completed
    • TaskList(): 列出所有任务
    • 内存数据结构 + JSON 持久化到 .lambdagent/tasks.json
    • 依赖: 无
  • [x] A15: 权限审批 CLI 交互 [1.5d]

    • ToolGateway confirm_callback 实现: 终端 input("Allow? [y/n/always]")
    • "always" 模式: 本次会话记忆已授权的工具模式
    • 显示工具名、风险等级、输入预览
    • 超时自动拒绝 (30s)
    • 依赖: 无 (接入现有 ToolGateway)
  • [x] A16: 项目级配置 .lambdagent.md [0.5d]

    • Agent 启动时自动加载项目根目录的 .lambdagent.md
    • 内容拼接到 systemPrompt 末尾
    • 支持 from_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) → 自定义恢复
    • 在 AsyncExecutor.reduce() 各分支注入 hook 调用点
    • YAML 配置: hooks.pre_tool_call: "python script.py"
    • 依赖: T03 (AsyncExecutor)
  • [x] A18: Notebook 编辑工具 NotebookEdit [3d]

    • 读取 .ipynb: 解析 JSON, 提取所有 cell (code + markdown + output)
    • Cell 级操作: 插入/删除/修改/移动 cell
    • 执行 cell: 通过 jupyter_client 连接 kernel (可选)
    • Schema: NotebookEditSchema(path, cell_index, action, content?)
    • 依赖: A01
  • [x] A19: Web 搜索工具 WebSearch [1d]

    • 接入搜索 API (Brave Search / SerpAPI / DuckDuckGo)
    • 返回 title + url + snippet 结构化结果
    • 可配置 API key 和搜索引擎
    • Schema: WebSearchSchema(query, max_results?)
    • 依赖: 无
  • [x] A20: Web 页面获取 WebFetch [1.5d]

    • HTTP GET 获取网页内容
    • HTML → Markdown 转换 (使用 markdownifyhtml2text)
    • 自动截断 (默认 max 5000 字符)
    • 支持 CSS 选择器提取特定内容
    • Schema: WebFetchSchema(url, selector?, max_length?)
    • 依赖: 无
  • [x] A21: 模型自动 Fallback [1.5d]

    • LLMAdapter 增加 fallback 链: [claude-sonnet, gpt-4o, qwen-max]
    • 主模型失败时自动切换 (429/500/超时)
    • 配合 CircuitBreaker: 主模型 circuit open → 自动用 fallback
    • YAML 配置: model.fallback: [openai/gpt-4o, dashscope/qwen-max]
    • 依赖: T05 (CircuitBreaker)
  • [x] A22: 流式终端 UI (完整版) [5d]

    • 多 Agent 面板 (并行 Agent 分栏显示)
    • 工具调用实时预览 (输入/输出折叠)
    • Token 计数实时显示
    • 成本估算实时更新
    • 历史会话浏览
    • 键盘快捷键 (Ctrl+C 取消, Ctrl+Z 撤销上一步)
    • 依赖: A09
  • [x] A23: 集成测试 — V2 全部功能 [2d]

    • CodeSearch + ProjectMap 端到端
    • TestRunner 多框架测试
    • Hook 系统触发验证
    • Notebook 读写
    • Web 搜索 + 获取
    • Fallback 模型切换
    • 完整编程场景: "读取项目 → 理解代码 → 修改 → 测试 → 提交"
    • 依赖: A11-A21

Tier 4 — 对标 Claude Code 关键差异能力

来源: Claude Code 33 个内置工具分析 (CC工具.md)。 以下是 A01-A23 未覆盖但对编程助手体验至关重要的能力。

  • [ ] A24: LSP 代码智能工具 [5d] ★★★ 关键

    • 集成 Language Server Protocol,提供真正的代码智能(非 grep 近似)
    • 核心能力:
    • goto_definition(file, line, col) → 跳转到符号定义
    • find_references(file, line, col) → 查找所有引用
    • hover(file, line, col) → 获取类型/文档信息
    • diagnostics(file) → 获取类型错误和警告
    • workspace_symbols(query) → 工作区符号搜索
    • 实现方案: 启动 LSP 子进程 (pylsp/tsserver/gopls),通过 JSON-RPC stdio 通信
    • 文件编辑后自动触发 diagnostics,报告引入的类型错误
    • Schema: LSPSchema(action, file_path, line?, column?, query?)
    • 依赖: 无 (需安装对应语言的 LSP server)
  • [ ] A25: 用户交互工具 AskUser [1.5d] ★★ 重要

    • 向用户提问以消歧或收集需求(而非盲目猜测)
    • 支持模式:
    • 开放式问题: AskUser("你希望用哪种方案?")
    • 单选: AskUser("选择框架", options=["React", "Vue", "Svelte"])
    • 多选: AskUser("选择功能", options=[...], multi=True)
    • 确认: AskUser("确认删除 main.py?", confirm=True)
    • CLI 实现: input() 交互 + rich 渲染选项列表
    • Agent 调用此工具时暂停执行,等待用户回复
    • Schema: AskUserSchema(question, options?, multi?, confirm?)
    • 依赖: A09 (终端 UI)
  • [ ] A26: 规划模式 PlanMode [2d] ★★ 重要

    • 两阶段工作流: 先探索 + 规划 → 用户批准 → 再执行
    • EnterPlanMode: 切换到只读模式(禁止 Edit/Write/Bash),只允许 Read/Glob/Grep/LSP
    • 计划写入 .lambdagent/plan.md 文件
    • ExitPlanMode: 提交计划,展示给用户审批
    • 用户批准后恢复全部工具权限,按计划执行
    • 在 ToolGateway 层实现: plan mode 时 GatewayPolicy.allowed_tools 仅保留只读工具
    • 依赖: A15 (权限审批)
  • [ ] A27: Worktree 工具 (用户级) [1d]

    • 将 IsolatedWorkspace 能力暴露为用户可调用的工具
    • EnterWorktree(slug) → 创建 Git Worktree,切换 CWD
    • ExitWorktree(action="keep"|"remove") → 退出,选择保留或清理
    • 自动检测 Worktree 中的变更,退出时提示 commit/discard
    • 底层复用 WorkspaceManager.create() / cleanup()
    • 依赖: T08 (isolation.py, 已完成)
  • [ ] A28: 工具延迟加载 ToolSearch [2d] ★★ 重要

    • 当 MCP 工具数量 >30 时,不全部注入 system prompt
    • 初始只暴露核心内置工具 + ToolSearch 元工具
    • ToolSearch(query, max_results=5) → 返回匹配的工具定义 (name + schema + description)
    • 模型在下一轮调用发现的实际工具
    • 实现: 在 _compile_tools() 中,工具数超阈值时将 MCP 工具放入 deferred 池
    • 减少 system prompt token 消耗 (大 MCP 服务器可有 50+ 工具)
    • Schema: ToolSearchSchema(query, max_results?)
    • 依赖: A08 (工具注册表)
  • [ ] A29: 定时调度 Cron [2d]

    • 会话内定时任务调度 (类似 Claude Code 的 CronCreate/Delete/List)
    • CronCreate(schedule, prompt, recurring=True) → 创建定时 prompt
    • CronDelete(cron_id) → 取消
    • CronList() → 列出所有
    • schedule: 标准 5 字段 cron 表达式 或 interval 格式 ("5m", "1h")
    • 非 recurring 任务执行一次后自动删除
    • recurring 任务最长 7 天有效 (防止遗忘)
    • 实现: asyncio.create_task + 内存调度器
    • 依赖: T01 (async)
  • [ ] A30: Agent Team 工具封装 [2d]

    • 将 GroupChat/AsyncPar/Handoff 能力暴露为运行时可调用的工具
    • AgentSpawn(config, input, background=False) → 启动子 Agent
    • background=True: 后台运行,返回 agent_id
    • isolation="worktree": 自动创建隔离工作区
    • AgentMessage(agent_id, message) → 向运行中的子 Agent 发消息
    • AgentStatus(agent_id) → 查看子 Agent 状态
    • 子 Agent 拥有独立 Context (避免污染父 Agent 上下文)
    • 底层复用 AsyncPar + IsolatedWorkspace + CancellationToken
    • Schema: AgentSpawnSchema(config, input, background?, isolation?)
    • 依赖: T01, T08, A27
  • [ ] A31: Skill 工具封装 [1d]

    • 将 SkillPack/SkillRegistry 暴露为运行时工具
    • SkillRun(skill_name, args) → 执行已注册的 Skill
    • Skill 定义: Markdown 文件中的 prompt 模板 + 工具列表
    • 自动发现 .lambdagent/skills/*.md 目录下的 Skill 文件
    • 底层复用 lambdagent.skills.SkillRegistry
    • 依赖: 无 (SkillPack 已存在)
  • [ ] A32: MCP 资源工具 [1d]

    • ListMcpResources(server?) → 列出 MCP Server 暴露的资源
    • ReadMcpResource(uri) → 读取特定 MCP 资源内容
    • 底层复用 MCPServer.read_resource(uri)
    • Schema: ListMcpResourcesSchema(server?), ReadMcpResourceSchema(uri)
    • 依赖: 无 (MCPServer 已有 read_resource)
  • [ ] A33: 集成测试 — 对标 CC 全部关键能力 [2d]

    • LSP: goto_definition + find_references 在 Python 项目中测试
    • AskUser: 模拟用户交互
    • PlanMode: 进入 → 只读验证 → 退出 → 全权限恢复
    • ToolSearch: 延迟加载 + 按关键词发现
    • AgentSpawn: 父子 Agent 隔离 + 取消传播
    • 端到端场景: "分析项目 → 规划方案 → 用户批准 → 多 Agent 并行实现 → 测试 → 提交"
    • 依赖: A24-A32

Tier 5 — 知识库 + OCR + 文档生成

基于已有 RAG 模块(rag.py),补充文档处理全链路能力。

  • [x] A34: 文档分块器 ChunkSplitter [1d]

    • 将长文档自动分割为适合检索的段落
    • 策略: 按段落 / 按 token 数 / 按标题层级 / 固定大小 + 重叠
    • split(text, strategy="paragraph", chunk_size=500, overlap=50) → List[str]
    • 支持 Markdown 标题感知分割(按 # / ## / ### 切分)
    • 集成到 create_rag(): create_rag(documents, chunk=True, chunk_size=500)
    • 依赖: 无
  • [x] A35: OCR 工具 OCRTool [1d]

    • 图片/扫描 PDF → 文字提取
    • 后端优先级: PaddleOCR > Tesseract > 云端 API
    • OCRTool(backend="auto"): 自动检测可用后端
    • ocr({"file_path": "scan.png"}) → 提取的文字
    • ocr({"file_path": "scan.pdf", "pages": "1-3"}) → PDF 逐页 OCR
    • 集成到 ReadFile: 扫描型 PDF (无文本层) 自动降级到 OCR
    • Schema: OCRSchema(file_path, pages?, language?)
    • 注册到 BUILTIN_TOOLS
    • 依赖: A01 (ReadFile)
  • [x] A36: 文档生成工具 DocGen [1.5d]

    • 从 Markdown 生成 PDF / Word / HTML
    • doc_gen({"source": "report.md", "format": "pdf", "output": "report.pdf"})
    • PDF 后端: weasyprint (推荐) > fpdf2 > Bash(pandoc)
    • Word 后端: python-docx
    • HTML: 直接转换(Markdown → HTML)
    • 支持模板: 内置论文/报告/简历模板
    • Schema: DocGenSchema(source, format, output, template?)
    • 注册到 BUILTIN_TOOLS
    • 依赖: 无
  • [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 序列化
    • 支持增量更新(文件变更时重新索引变更部分)
    • 注册到 BUILTIN_TOOLS
    • 依赖: A34 (分块器), rag.py (RAGTool)
  • [x] A38: 集成测试 — 知识库全链路 [1d]

    • OCR: 图片 → 文字提取 → 验证
    • 分块: 长文档 → 分块 → 检索 → 验证 top_k 相关性
    • 文档生成: Markdown → PDF/HTML 验证
    • E2E: "导入一个 PDF 文档到知识库 → OCR 处理 → 提问检索 → 生成报告"
    • 依赖: A34-A37

Phase 6: PaaS 平台服务 — JARVIS 级能力

将 JARVIS 的核心能力抽象为 PaaS 服务,任何 Agent 通过 YAML 配置即可消费。 Agent 不需要改代码,只需 services.memory: enabled

  • [x] P01: Memory Service — 持久化记忆服务 [3d] ★★★★★

    • PaaS API: POST /api/v1/memory/{agent_id}/store, GET .../recall, DELETE .../forget
    • 三层记忆: Working (会话) + Episodic (近期摘要) + Semantic (永久知识)
    • 存储: SQLite (dev) / Redis (prod)
    • 每次会话结束自动提取关键信息存入 Semantic
    • 每次会话开始自动召回相关记忆注入 system prompt
    • Agent 侧工具: MemoryStore, MemoryRecall, MemoryProfile
    • 依赖: 已有 MemoryBackend + RAG
  • [x] P02: Scheduler Service — 持久化调度服务 [2d] ★★★★

    • PaaS API: POST /api/v1/schedules, GET, DELETE, PATCH
    • 任务类型: cron 表达式 / interval / once / event-triggered
    • 每个任务 = 一次 Agent 执行 (agent_id + input + notify_channel)
    • 持久化: ~/.lambdagent/schedules.json 或 DB
    • 支持暂停/恢复/查看执行历史
    • Agent 侧工具: ScheduleCreate, ScheduleList, ScheduleDelete
    • 依赖: 无
  • [x] P03: Event Bus — 系统事件总线 [3d] ★★★★

    • PaaS API: POST /api/v1/events/subscribe, POST .../publish
    • 系统事件: file.created/modified, disk.low, battery.low, wifi.changed, app.launched
    • Git 事件: push, conflict, merge
    • Agent 事件: task.completed, task.failed
    • Agent 可订阅事件并绑定自动动作
    • 实现: watchdog (文件) + 平台轮询 (系统)
    • 每个 Agent YAML 声明订阅: events: [system.file.created:~/Desktop/*]
    • 依赖: Hook Layer 3
  • [x] P04: Notification Service — 多渠道通知 [2d] ★★★★

    • PaaS API: POST /api/v1/notify
    • 渠道: terminal (REPL banner) / macos (通知中心) / webhook / email / wechat
    • 优先级: low / normal / high / urgent
    • Agent 侧工具: Notify(message, channel, priority)
    • 已有基础: app(notify) macOS 工具、wechat.py
    • 依赖: 无
  • [x] P05: Profile Service — 用户画像自动构建 [1.5d] ★★★

    • PaaS API: GET /api/v1/profile, PUT
    • 自动收集: 常用命令、编程语言偏好、工作时间段、交互风格
    • 每个 Agent 启动时自动注入 Profile 到 context
    • 存储: ~/.lambdagent/user_profile.json 或 DB
    • 依赖: P01 (Memory Service)
  • [x] P06: Persona Engine — 可配置人格 [2d] ★★★

    • YAML 配置: persona.name, persona.style, persona.proactivity, persona.wake_word
    • 内置模板: JARVIS (正式) / lambda🐑 (活泼) / Friday (干练) / 自定义
    • 编译时将 persona 转换为 system prompt 组件
    • 支持情境切换(工作模式 vs 休闲模式)
    • 依赖: P05 (Profile)
  • [x] P07: Multimodal Gateway — 多模态网关 [5d] ★★★

    • PaaS API: POST /api/v1/multimodal/stt, POST .../tts, POST .../vision
    • STT: whisper (本地) / 云端 ASR
    • TTS: edge-tts / macOS say / 云端
    • Vision: Claude Vision / GPT-4V / 本地 VLLM
    • Agent 侧工具: Listen, Speak, See
    • 依赖: 无 (独立服务)
  • [x] P08: Learning Service — 自适应学习 [2d] ★★

    • PaaS API: POST /api/v1/learning/feedback, GET .../strategies
    • 记录工具调用成功/失败模式
    • 自动提取策略规则(如 "WriteFile 大文件用 Bash+python3")
    • 下次遇到类似场景自动应用学到的策略
    • 依赖: P01 (Memory), Hook Layer 3

TODO 统计

类别 数量 已完成 预估总工作量
工程改善 (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 人天

TODO 依赖图

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

建议并行路线

  • 路线 A (核心引擎): T01 → T02 → T03 → T04 → T05 → T06 → T07 ✅ 已完成
  • 路线 B (文件隔离): T08 → T09 → T10 → T11 ✅ 已完成
  • 路线 C (独立优化): T15 ✅ 已完成
  • 路线 D (安全加固): S01-S12 大部分已完成 ✅
  • 路线 E (助手 MVP): A01-A07 → A08 → A09 → A10 (可立即启动)
  • 路线 F (助手 V1): A11-A16 (依赖路线 E)
  • 路线 G (助手 V2): A17-A23 (依赖路线 F)
  • 路线 H (对标 CC): A24 (LSP, 独立) + A25-A26 (交互/规划) + A28-A32 (可与 G 并行)
  • 路线 I (Provider 统一): L01 → L02 → L03 → L04 → L05 (消除幻觉,统一多模型)
  • 路线 K (双引擎切换 — 可立即启动,无前置依赖):
    • E01 → E02/E03 (并行) → E04 → E05 → E06 → E07
    • E01 (Engine 接口) 无依赖,可今天开始
    • E02 (RecursiveEngine) 和 E03 (CEKEngine) 可并行开发
    • E07 (一致性测试) 是最终验收门
  • 路线 J (集成层 — 可立即启动,无前置依赖):
    • J-1 (MCP Server): I01 → I02 → I03 [4d, 最高优先级]
    • J-2 (GitHub Action): I04 → I05 [3d, 与 J-1 并行]
    • J-3 (提取器): I06 → I07/I08/I09 (并行) [6.5d]
    • J-4 (运行时守卫): I10 → I11/I12/I13 (并行) → I14 [9d, 依赖 J-3]
    • J-5 (REST API): I15 [1d, 独立]
    • J-6 (IDE 扩展): I16 → I17 [4d, 独立]

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 发布)

Phase 6: LLM Provider 统一 + 会话持久化 (消除幻觉)

  • [ ] L01: 实现 lambdagent/providers/base.py — LLMProvider 接口 [1d]

    • 统一接口: chat(messages: list[dict]) -> str
    • 所有 provider 实现同一个 messages-in / text-out 契约
    • 与现有 Lam._call_llm() 解耦
  • [ ] L02: 实现 4 个 Provider [2d]

    • AnthropicProvider: client.messages.create(messages=...)
    • OpenAICompatProvider: 覆盖 OpenAI / Ollama / DashScope / DeepSeek / Moonshot
    • ClaudeCodeProvider: 优先 --resume,fallback 到 messages 拼接
    • 依赖: L01
  • [ ] L03: 实现 lambdagent/conversation.py — ConversationLam [2d]

    • 维护 messages: list[dict] 对话历史
    • 上下文窗口管理: 滑窗 + 摘要,按 maxHistoryTokens 控制
    • 适配各 provider 的上下文窗口差异 (32K~200K)
    • reset_session() 清空历史
    • 依赖: L01, L02
  • [ ] L04: from_config 编译器集成 [1d]

    • _compile_lam(): 根据 provider 创建对应 Provider + ConversationLam
    • YAML 新增: model.conversation: true (默认开启), model.maxHistoryTokens: 80000
    • react_step 简化: 首步传完整输入,后续只传最新工具结果
    • 移除 _compress_state 依赖(ConversationLam 内部管理历史)
    • 依赖: L03
  • [ ] L05: 端到端验证 + PaaS 集成 [1d]

    • agentpaas chat 统一走 from_config + ConversationLam
    • 测试: Claude Code / Ollama / DashScope 三种 provider
    • 验证: 多步任务零幻觉, 工具调用正确, commit+push 成功
    • 更新 docs/hallucination-root-cause-and-fix.md
    • 依赖: L04

Phase 6.5: 双引擎切换 — Recursive / CEK 可配置执行引擎

目标: 将 Python 调用栈执行(Recursive)和 Agent CEK Machine 统一在同一个 Engine 抽象下,通过配置切换。简单 agent 用 Recursive(低开销),复杂/高成本 agent 用 CEK(可暂停/恢复/成本熔断/精确追踪)。

原则: Recursive 是默认引擎(向后兼容),CEK 是可选升级。双引擎对 125 个测试用例必须产生完全一致的结果。

  • [ ] E01: 实现 lambdagent/agentruntime/engine.py — Engine 抽象接口 [0.5d]

    • EngineMode 枚举: RECURSIVE | CEK
    • Engine ABC: execute(term, input, **opts) -> EngineResult, execute_async(...)
    • EngineResult 统一输出: value, trace: List[TraceRecord], cost: CostVector, steps: int
    • CEK 独有字段: final_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 格式
    • 从 trace 汇总 CostVector(tokens, latency, money)
    • final_statetransitions 输出为 None
    • 依赖: E01
  • [ ] E03: 实现 CEKEngine — 封装现有 CEKMachine [2d]

    • 封装 CEKMachine.run(term, input)Engine.execute() 接口
    • 实现 execute_async(): Yield 点映射为 await(LLM/Tool 异步调用)
    • 逐步成本检查: 每次 Yield 后比较 state.cost.moneycost_budget
    • 逐步循环检测: 连续 N 步 state hash 相同 → InfiniteLoopDetected
    • Transition 列表转换为统一 TraceRecord 格式
    • K 栈可读描述: _format_kont(kont)"compK(g, loopK(b,c,7, halt))"
    • 依赖: E01
  • [ ] E04: 统一 TraceRecord 格式与转换 [1d]

    • TraceRecord dataclass: step, term_name, action, input_summary, output_summary, duration_ms, cost, cumulative_cost
    • CEK 独有字段: continuation: Optional[str](K 栈描述), yield_type: Optional[str]("llm"|"tool"|None)
    • RecursiveEngine._convert_trace(): ctx.trace → List[TraceRecord]
    • CEKEngine._convert_transitions(): List[Transition] → List[TraceRecord](跳过 τ 静默转移或可选保留)
    • 依赖: E02, E03
  • [ ] 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 字段选择 RecursiveEngineCEKEngine
      • 向后兼容: 无 runtime.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

      Phase 7: 集成层 — lambdagent 作为外部框架的静态分析插件

      核心策略: 不要求用户换框架。用户继续用 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

      Phase 7a: MCP Server — AI IDE 集成 (最高优先级)

      • I01: 实现 lambdagent/mcp_server.py — MCP Server 入口 [2d]
      • 暴露 4 个 MCP 工具: lint_agent_config, estimate_agent_cost, check_agent_types, check_parallel_safety
      • lint_agent_config: 调用已有 lint_config(), 支持 LangChain/CrewAI/AutoGen/Dify/通用 YAML
      • estimate_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

      Phase 7b: GitHub Action — CI/CD 守门

      • I04: 实现 lambdagent/agent-lint-action — GitHub Action [2d]
      • Docker Action 封装 lambdagent lint --format json
      • 输入参数: paths, fail-on (error|warn|info), frameworks (auto|指定), cost-threshold, type-check, parallel-check
      • 输出: PR Comment (Markdown 表格) + Check Annotation (行级标注)
      • 支持 monorepo: 自动发现 **/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

      Phase 7c: 框架配置提取器

      • I06: 实现 lambdagent/extractors/base.py — 提取器接口 [0.5d]
      • FrameworkExtractor ABC: detect(obj) -> bool, extract(obj) -> dict
      • extract_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 链 (RunnableSequencetype: chain)

      • 依赖: I06

      • [ ] I08: 实现 lambdagent/extractors/crewai_extractor.py [2d]

      • Crew 提取: agents (role/goal/backstory → systemPrompt), tasks, process type

      • 映射: sequentialtype: chain, hierarchicaltype: 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

      Phase 7d: Python 中间件 — 运行时守卫

      • I10: 实现 lambdagent_guard/core.py — 守卫核心 [2d]
      • GuardedExecutor: 编译时检查 (type_check + cost_budget + store_independence) + 运行时 hook
      • 运行时 hook: 每步成本累积 (CEK Yield 等价), 循环检测 (state hash 重复 3 次), 成本熔断
      • CostBudgetExceeded, 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_stepagent.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

      Phase 7e: REST API 分析端点

      • I15: 新增 /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)
      • 依赖: 无 (底层函数已存在)

      Phase 7f: IDE 扩展

      • I16: 实现 VS Code 扩展 lambdagent.agent-lint [3d]
      • Language Server Protocol (LSP) 集成: YAML 文件保存时自动 lint
      • Inline diagnostics: 行级错误/警告标注
      • Status bar: λA: 2 errors | Est. cost: $6.00/run | Success: 0.16%
      • Quick Fix: 针对每条 lint 规则的自动修复建议
      • 底层调用 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 天等)证明每一项静态检查都有对应的真实需求。