SPEC.md 50 KB

AgentPaaS 需求规格文档

将 lambdagent 计算内核包装为生产级 Agent 平台即服务

版本: 1.0 | 前置: RUNTIME_SPEC.md, FROM_CONFIG_SPEC.md | 内核: lambdagent


0. 为什么需要 AgentPaaS

lambdagent 解决了 "如何定义和执行一个 agent" 的问题——它是进程内的 β-规约引擎。但从单机引擎到可商用的 Agent 平台,还缺少:

lambdagent 能做的               lambdagent 做不了的
─────────────────               ──────────────────
YAML → Term 编译                 多租户隔离
β-规约执行                       HTTP API 服务化
ReAct 循环                       认证 / 鉴权 / 限流
MCP / A2A 协议                   异步任务队列
Trace 记录                       分布式可观测性
Checkpoint 保存                  计费 / 配额
                                 Agent 版本管理
                                 灰度发布 / 流量治理
                                 密钥管理

AgentPaaS 补全右列,左列保持不动。

┌─────────────────────────────────────────────────────┐
│                     AgentPaaS                       │
│                                                     │
│  ┌─────────┐ ┌──────────┐ ┌───────┐ ┌───────────┐  │
│  │ Gateway  │ │ Registry │ │Engine │ │ Dashboard  │  │
│  │ 认证/限流 │ │ 注册/版本 │ │调度/沙箱│ │ Web UI   │  │
│  └────┬────┘ └────┬─────┘ └───┬───┘ └───────────┘  │
│       │           │           │                     │
│  ┌────▼───────────▼───────────▼──────────────────┐  │
│  │            Tenant / Billing / Secrets          │  │
│  │         多租户 · 计费 · 密钥 · 可观测性          │  │
│  └───────────────────┬───────────────────────────┘  │
│                      │                              │
├──────────────────────┼──────────────────────────────┤
│                      ▼                              │
│  ┌───────────────────────────────────────────────┐  │
│  │              lambdagent (内核)                  │  │
│  │  from_config · Runtime · Executor · Context    │  │
│  │  Term · ReAct · MCP · A2A · Checkpoint         │  │
│  └───────────────────────────────────────────────┘  │
└─────────────────────────────────────────────────────┘

1. 设计原则

编号 原则 说明
P1 内核零侵入 AgentPaaS 只调用 lambdagent 公共 API,不 fork、不 monkey-patch
P2 API First 所有功能通过 REST/gRPC 暴露;Dashboard、CLI 都是 API 消费者
P3 渐进式复杂度 单机 SQLite + 内存队列可跑 MVP;生产上 PostgreSQL + Redis + K8s
P4 形式化可追溯 保留 lambdagent β-规约 Trace,上层叠加 OpenTelemetry 分布式 Trace
P5 开放协议 原生支持 MCP、A2A;不锁定用户到私有协议
P6 安全默认 密钥加密存储、最小权限、所有变更可审计

2. 核心等式

AgentPaaS.serve(request) =
    let tenant  = Gateway.authenticate(request)                in  ← 认证
    let quota   = Tenant.check_quota(tenant)                   in  ← 配额
    let yaml    = Registry.load(request.agent_id, version)     in  ← 加载
    let term    = lambdagent.from_config(yaml)                 in  ← 编译 (内核)
    let ctx     = Context.new(secrets=Vault.get(tenant))       in  ← 上下文
    let result  = Engine.execute(term, request.input, ctx)     in  ← 执行 (内核)
    let _       = Billing.record(tenant, ctx.trace)            in  ← 计费
    let _       = Observability.export(ctx.trace)              in  ← 观测
    (result, ctx.trace)

Engine.execute(term, input, ctx) =
    match execution_mode with
    | Sync    → lambdagent.Runtime().reduce(term, input, ctx)
    | Async   → JobQueue.submit(term, input, ctx) → job_id
    | Stream  → SSEStream(lambdagent.Runtime().reduce_iter(term, input, ctx))

3. API 规格

3.1 总览

基础路径: /api/v1
认证方式: Bearer Token (API Key 或 JWT)
内容类型: application/json
错误格式: { "error": { "code": string, "message": string, "details": object } }

3.2 Agent 管理

3.2.1 创建 Agent

POST /api/v1/agents
Authorization: Bearer {api_key}
Content-Type: application/json

Request:
{
    "name": "customer-support",
    "description": "客户支持 agent",
    "config": { ... },              // lambdagent YAML 等价的 JSON
    "tags": ["support", "chat"],
    "environment": "production"
}

Response: 201 Created
{
    "agent_id": "ag_3f8a1b2c",
    "version": 1,
    "created_at": "2026-03-26T10:00:00Z",
    "endpoint": "/api/v1/agents/ag_3f8a1b2c"
}

3.2.2 查询 Agent

GET /api/v1/agents/{agent_id}
GET /api/v1/agents/{agent_id}/versions            // 版本列表
GET /api/v1/agents/{agent_id}/versions/{version}  // 特定版本
GET /api/v1/agents?tag=support&model=claude       // 搜索

3.2.3 更新 Agent

PUT /api/v1/agents/{agent_id}

自动递增版本号。旧版本保留,可回滚:
POST /api/v1/agents/{agent_id}/rollback
{ "target_version": 2 }

3.2.4 删除 Agent

DELETE /api/v1/agents/{agent_id}

软删除。30 天后物理清除。

3.3 Agent 执行

3.3.1 同步执行

POST /api/v1/agents/{agent_id}/run
POST /api/v1/agents/{agent_id}@v2/run             // 指定版本

Request:
{
    "input": "帮我查询订单 #12345 的状态",
    "parameters": {                                // 运行时参数覆盖
        "model.temperature": 0.3,
        "react.maxSteps": 5
    },
    "context": {                                   // 额外上下文
        "user_id": "u_abc",
        "session_id": "s_xyz"
    }
}

Response: 200 OK
{
    "run_id": "run_7d9e4f",
    "status": "completed",
    "output": "订单 #12345 当前状态为已发货...",
    "usage": {
        "input_tokens": 1250,
        "output_tokens": 380,
        "total_tokens": 1630,
        "steps": 3,
        "duration_ms": 4520
    },
    "workspace_path": "/var/agentpaas/agents/.../workspace/run_20260606_143000"
}

v1.0 实现说明: usage.estimated_cost_usd 和顶层 trace_id 字段曾 在设计草案里出现, 但 agentpaas/src/agentpaas/api/v1/agents.py 当前 的响应构造没有填充它们。前者依赖 lambdagent.cost_grade 的接入, 后者需要 OpenTelemetry trace 注入到 response 包装层——都列入 v1.1 backlog (审计 docs/AUDIT_2026-06-05.md #75)。

3.3.2 异步执行

POST /api/v1/agents/{agent_id}/jobs

Response: 202 Accepted
{
    "job_id": "job_8e5f2a",
    "status": "queued",
    "poll_url": "/api/v1/jobs/job_8e5f2a"
}

GET /api/v1/jobs/{job_id}

Response:
{
    "job_id": "job_8e5f2a",
    "status": "running" | "completed" | "failed" | "cancelled",
    "progress": { "current_step": 3, "max_steps": 10 },
    "output": "...",                               // status=completed 时存在
    "error": { ... },                              // status=failed 时存在
    "created_at": "...",
    "started_at": "...",
    "completed_at": "..."
}

POST /api/v1/jobs/{job_id}/cancel                  // 取消任务

3.3.3 流式执行 (SSE)

POST /api/v1/agents/{agent_id}/run/stream
Accept: text/event-stream

SSE 事件流:
event: step
data: {"step": 1, "type": "think", "content": "我需要查询订单系统..."}

event: step
data: {"step": 2, "type": "tool_call", "tool": "order_api", "input": {...}}

event: step
data: {"step": 2, "type": "tool_result", "output": {...}}

event: step
data: {"step": 3, "type": "answer", "content": "订单 #12345 已发货..."}

event: done
data: {"run_id": "run_7d9e4f", "usage": {...}}

3.3.4 WebSocket 交互 (planned — 未在 v1.0 实现)

v1.0 状态: 此端点为设计文档, 代码侧暂未实现。v1.0 的实时交互 通过 §3.3.3 的 SSE 流式端点完成, SSE 已经覆盖了大多数 WS 想解决的 场景。在确定真实需求 (双向 / 多轮持续会话) 之前, WS 实现暂不立项。 审计参见 docs/AUDIT_2026-06-05.md #74

WS /api/v1/agents/{agent_id}/ws    (planned)

→ Client:  {"type": "message", "input": "查询订单 #12345"}
← Server:  {"type": "thinking", "content": "正在分析..."}
← Server:  {"type": "tool_use", "tool": "order_api", ...}
← Server:  {"type": "response", "content": "订单已发货..."}
→ Client:  {"type": "message", "input": "帮我修改地址"}
← Server:  ...
→ Client:  {"type": "close"}

3.4 认证与租户

3.4.1 API Key 管理

POST   /api/v1/auth/keys                          // 创建 API Key
GET    /api/v1/auth/keys                          // 列出 API Key
DELETE /api/v1/auth/keys/{key_id}                 // 吊销 API Key

创建时指定 scope:
{
    "name": "production-key",
    "scopes": ["agents:read", "agents:execute"],   // 细粒度权限
    "rate_limit": 100,                             // 每分钟请求上限
    "expires_at": "2027-01-01T00:00:00Z"          // 可选过期时间
}

3.4.2 租户管理 (Admin)

POST   /api/v1/admin/tenants                      // 创建租户
GET    /api/v1/admin/tenants/{tenant_id}          // 查询租户
PUT    /api/v1/admin/tenants/{tenant_id}/quota    // 设置配额
GET    /api/v1/admin/tenants/{tenant_id}/usage    // 用量统计

3.5 可观测性

GET /api/v1/traces/{trace_id}                     // 查询执行轨迹
GET /api/v1/agents/{agent_id}/metrics             // Agent 级指标
GET /api/v1/metrics/overview                      // 全局指标概览

4. 数据模型

4.1 ER 关系

Tenant ─1:N─ User ─1:N─ APIKey
  │
  ├─1:N─ Agent ─1:N─ AgentVersion
  │                       │
  │                       └─ config (YAML JSON)
  │                       └─ compiled_term (缓存)
  │
  ├─1:N─ Run ─1:N─ TraceSpan
  │        │
  │        └─ job_id (nullable, 异步任务时)
  │
  ├─1:N─ Job
  │
  └─1:N─ UsageRecord

4.2 核心表结构

tenants

字段 类型 说明
id UUID 主键
name VARCHAR(128) 租户名称
plan ENUM free / pro / enterprise
quota_tokens_monthly BIGINT 月 token 配额
quota_concurrency INT 最大并发执行数
quota_agents INT 最大 agent 数
status ENUM active / suspended / deleted
created_at TIMESTAMPTZ 创建时间

agents

字段 类型 说明
id VARCHAR(16) 主键, ag_ 前缀
tenant_id UUID FK → tenants
name VARCHAR(128) Agent 名称
description TEXT 描述
current_version INT 当前生效版本号
tags JSONB 标签数组
environment VARCHAR(32) dev / staging / production
traffic_rules JSONB 灰度规则 (nullable)
status ENUM active / paused / deleted
created_at TIMESTAMPTZ 创建时间
updated_at TIMESTAMPTZ 最后更新

agent_versions

字段 类型 说明
agent_id VARCHAR(16) FK → agents
version INT 版本号
config JSONB 完整 YAML 配置 (JSON 序列化)
config_hash VARCHAR(64) SHA-256, 用于去重
changelog TEXT 变更说明
created_by UUID FK → users
created_at TIMESTAMPTZ 创建时间
PK (agent_id, version) 联合主键

runs

字段 类型 说明
id VARCHAR(16) 主键, run_ 前缀
agent_id VARCHAR(16) FK → agents
agent_version INT 执行时的版本
tenant_id UUID FK → tenants
job_id VARCHAR(16) FK → jobs (nullable)
input JSONB 输入
output TEXT 输出
status ENUM running / completed / failed / cancelled
input_tokens INT 输入 token 数
output_tokens INT 输出 token 数
steps INT β-规约步数
duration_ms INT 执行耗时
error JSONB 错误详情 (nullable)
trace_id VARCHAR(32) OpenTelemetry trace ID
created_at TIMESTAMPTZ 创建时间
completed_at TIMESTAMPTZ 完成时间

jobs

字段 类型 说明
id VARCHAR(16) 主键, job_ 前缀
tenant_id UUID FK → tenants
agent_id VARCHAR(16) FK → agents
priority INT 优先级 (0=最高)
status ENUM queued / running / completed / failed / cancelled
input JSONB 执行输入
worker_id VARCHAR(64) 执行该 job 的 worker (nullable)
attempts INT 已尝试次数
max_attempts INT 最大重试次数, 默认 3
scheduled_at TIMESTAMPTZ 计划执行时间 (nullable)
started_at TIMESTAMPTZ 开始执行时间
completed_at TIMESTAMPTZ 完成时间
created_at TIMESTAMPTZ 创建时间

api_keys

字段 类型 说明
id UUID 主键
tenant_id UUID FK → tenants
user_id UUID FK → users
key_hash VARCHAR(64) SHA-256(key), 不存明文
key_prefix VARCHAR(8) 前 8 位, 用于识别 ap_xxxx...
name VARCHAR(128) Key 名称
scopes JSONB 权限列表
rate_limit INT 每分钟请求上限
expires_at TIMESTAMPTZ 过期时间 (nullable)
last_used_at TIMESTAMPTZ 最后使用时间
status ENUM active / revoked
created_at TIMESTAMPTZ 创建时间

5. Gateway 层

5.1 请求处理流水线

Request
  │
  ▼
┌────────────────┐
│ TLS 终止       │
└───────┬────────┘
        ▼
┌────────────────┐     ┌─────────────────────────────────┐
│ 认证中间件      │────▶│ API Key → SHA-256 → 查 api_keys │
│ (AuthMiddleware)│     │ JWT → 验签 → 解析 claims        │
└───────┬────────┘     └─────────────────────────────────┘
        ▼
┌────────────────┐     ┌─────────────────────────────────┐
│ 租户上下文注入   │────▶│ 从 token 提取 tenant_id         │
│ (TenantCtx)    │     │ 注入 request.state.tenant       │
└───────┬────────┘     └─────────────────────────────────┘
        ▼
┌────────────────┐     ┌─────────────────────────────────┐
│ 配额检查        │────▶│ 检查月 token 余量               │
│ (QuotaCheck)   │     │ 检查并发执行数                   │
│                │     │ 超限 → 429 Too Many Requests     │
└───────┬────────┘     └─────────────────────────────────┘
        ▼
┌────────────────┐     ┌─────────────────────────────────┐
│ 限流            │────▶│ 令牌桶算法                      │
│ (RateLimiter)  │     │ 键 = tenant_id:api_key_id       │
│                │     │ Redis INCR + EXPIRE              │
└───────┬────────┘     └─────────────────────────────────┘
        ▼
┌────────────────┐     ┌─────────────────────────────────┐
│ 请求校验        │────▶│ Pydantic 模型校验 request body  │
│ (Validation)   │     │ 非法输入 → 422                   │
└───────┬────────┘     └─────────────────────────────────┘
        ▼
   Route Handler

5.2 认证机制

# API Key 认证 (主要方式)
# 格式: ap_xxxxxxxxxxxxxxxx (32 字符随机串, 前缀 ap_)
# 存储: 只存 SHA-256 哈希, 创建时返回一次明文

class APIKeyAuth:
    async def authenticate(self, key: str) -> Tenant:
        key_hash = sha256(key)
        record = await db.api_keys.find(key_hash=key_hash, status="active")
        if not record:
            raise AuthError("Invalid API key", code=401)
        if record.expires_at and record.expires_at < now():
            raise AuthError("API key expired", code=401)
        await db.api_keys.update(record.id, last_used_at=now())
        return await db.tenants.get(record.tenant_id)

5.3 限流策略

三级限流:

Level 1 — 全局:        10000 req/min (防 DDoS)
Level 2 — 租户:        按 plan 配置 (free=60, pro=600, enterprise=自定义)
Level 3 — 单 API Key:  创建时指定, 默认 = 租户限额

算法: 滑动窗口 (Redis ZSET)
键:   rate:{tenant_id}:{window_start}
超限响应:
    HTTP 429
    Headers:
        X-RateLimit-Limit: 600
        X-RateLimit-Remaining: 0
        X-RateLimit-Reset: 1711440060
        Retry-After: 32

6. Registry(注册中心)

6.1 职责

Registry 管理 Agent 的完整生命周期:定义 → 版本 → 发布 → 发现 → 退役。

6.2 版本管理

每次 PUT /agents/{id} 自动递增版本:

v1 ──── v2 ──── v3 (current) ──── v4
                  ▲
                  └── rollback 目标

规则:
    - 版本号单调递增, 不可跳跃
    - config 不变时 (config_hash 相同) 不创建新版本
    - current_version 指向生效版本
    - 旧版本永久保留 (除非租户主动清理)
    - rollback = 将 current_version 指向目标版本, 不创建新版本

6.3 流量治理 (Canary)

// agents.traffic_rules 字段
{
    "rules": [
        {
            "version": 4,
            "weight": 10,
            "match": { "header": "X-Canary: true" }
        },
        {
            "version": 3,
            "weight": 90,
            "match": "default"
        }
    ]
}

路由逻辑:

def resolve_version(agent: Agent, request: Request) -> int:
    if not agent.traffic_rules:
        return agent.current_version

    for rule in agent.traffic_rules["rules"]:
        if matches_header(rule.get("match"), request):
            return rule["version"]

    # 按 weight 加权随机
    return weighted_random(agent.traffic_rules["rules"])

6.4 A2A 服务发现

集成 lambdagent 的 A2A 协议:

GET /.well-known/agent.json

Response:
{
    "name": "customer-support",
    "description": "客户支持 agent",
    "url": "https://platform.example.com/api/v1/agents/ag_3f8a1b2c",
    "capabilities": ["order-query", "address-update"],
    "protocol": "a2a/1.0",
    "authentication": {
        "type": "bearer",
        "token_url": "https://platform.example.com/api/v1/auth/token"
    }
}

注册中心自动生成 Agent Card, 外部系统可通过 A2A 协议发现和调用。

7. Execution Engine(执行引擎)

7.1 架构

┌─────────────────────────────────────────────┐
│              Execution Engine                │
│                                             │
│  ┌───────────┐    ┌──────────────────────┐  │
│  │ Dispatcher │───▶│     Job Queue        │  │
│  │ (路由请求)  │    │  (Redis / PostgreSQL) │  │
│  └───────────┘    └──────────┬───────────┘  │
│                              │              │
│                    ┌─────────▼──────────┐   │
│                    │    Worker Pool     │   │
│                    │                    │   │
│                    │  ┌──────────────┐  │   │
│                    │  │   Sandbox    │  │   │
│                    │  │ ┌──────────┐ │  │   │
│                    │  │ │lambdagent│ │  │   │
│                    │  │ │ Runtime  │ │  │   │
│                    │  │ └──────────┘ │  │   │
│                    │  └──────────────┘  │   │
│                    │                    │   │
│                    │  ┌──────────────┐  │   │
│                    │  │   Sandbox    │  │   │
│                    │  │     ...      │  │   │
│                    │  └──────────────┘  │   │
│                    └────────────────────┘   │
└─────────────────────────────────────────────┘

7.2 Dispatcher(调度器)

class Dispatcher:
    """
    将执行请求路由到合适的执行模式。
    """

    async def dispatch(self, agent: Agent, input: Any, mode: str, ctx: TenantContext):
        # 1. 编译 (或从缓存取)
        term = self.compile_cache.get_or_compile(agent)

        # 2. 构建 lambdagent Context
        la_ctx = self._build_context(ctx)

        match mode:
            case "sync":
                return await self.worker_pool.execute(term, input, la_ctx,
                    timeout=agent.config.get("timeout", 120))

            case "async":
                job = Job(agent_id=agent.id, input=input, priority=ctx.priority)
                await self.job_queue.enqueue(job)
                return job

            case "stream":
                return self.worker_pool.execute_stream(term, input, la_ctx)

7.3 Worker 与 Sandbox

class Sandbox:
    """
    隔离的 agent 执行环境。

    隔离级别 (可配):
        Level 0 — 同进程      (开发模式, 最快, 无隔离)
        Level 1 — 子进程      (默认, 进程级隔离, 内存限制)
        Level 2 — 容器        (生产, Docker 容器, 完全隔离)
    """

    def __init__(self, level: int, limits: ResourceLimits):
        self.level = level
        self.limits = limits   # cpu_seconds, memory_mb, network

    async def execute(self, term, input, ctx, timeout: int):
        match self.level:
            case 0:
                return await self._exec_inprocess(term, input, ctx, timeout)
            case 1:
                return await self._exec_subprocess(term, input, ctx, timeout)
            case 2:
                return await self._exec_container(term, input, ctx, timeout)

class ResourceLimits:
    cpu_seconds: int = 300          # 最大 CPU 时间
    memory_mb: int = 512            # 最大内存
    max_steps: int = 100            # 最大 β-规约步数
    max_tokens: int = 100_000       # 最大 token 总消耗
    network_enabled: bool = True    # 是否允许网络访问

7.4 Job Queue

任务状态机:

                ┌─── cancel ───┐
                ▼              │
queued ──▶ running ──▶ completed
  │          │
  │          ├──▶ failed ──▶ queued (retry, if attempts < max)
  │          │
  └── cancel ▶ cancelled

优先级: 0 (最高) ~ 9 (最低)
    - 同步请求转异步 (超时降级): priority=0
    - 用户主动提交异步任务: priority=5
    - 批量任务: priority=9

队列实现:
    开发: asyncio.PriorityQueue (内存)
    生产: Redis Sorted Set (ZADD score=priority*1e12+timestamp)

7.5 重试与断路器

class RetryPolicy:
    max_attempts: int = 3
    base_delay: float = 1.0         # 秒
    max_delay: float = 60.0
    exponential_base: float = 2.0
    retryable_errors: list = [
        "RateLimitError",           # LLM API 限流
        "TimeoutError",             # LLM API 超时
        "ServiceUnavailable",       # 503
    ]

class CircuitBreaker:
    """
    熔断器: 防止 LLM API 故障引发雪崩。

    状态:
        CLOSED  → 正常放行
        OPEN    → 快速失败 (不调 LLM)
        HALF    → 试探性放行少量请求

    触发条件:
        - 连续 5 次失败 → OPEN
        - OPEN 持续 30s → HALF
        - HALF 中 1 次成功 → CLOSED
        - HALF 中 1 次失败 → OPEN
    """
    failure_threshold: int = 5
    recovery_timeout: float = 30.0

7.6 Compile Cache

class CompileCache:
    """
    缓存 YAML → Term 编译结果。

    键: config_hash (agent_versions.config_hash)
    值: 编译后的 Term 对象
    策略: LRU, 最大 1000 个 Term

    编译是确定性的 (FROM_CONFIG_SPEC P4), 所以 hash 相同 → Term 相同。
    """

    def get_or_compile(self, agent: Agent) -> Term:
        version = agent.current_version
        config_hash = agent.versions[version].config_hash

        if config_hash in self.cache:
            return self.cache[config_hash]

        term = lambdagent.from_config(agent.versions[version].config)
        self.cache[config_hash] = term
        return term

8. Multi-Tenancy(多租户)

8.1 隔离模型

                    ┌──────────────────────┐
     Tenant A       │       Tenant B       │
    ┌──────────┐    │    ┌──────────┐      │
    │ Agent 1  │    │    │ Agent 3  │      │
    │ Agent 2  │    │    │ Agent 4  │      │
    ├──────────┤    │    ├──────────┤      │
    │ Secrets  │    │    │ Secrets  │      │
    ├──────────┤    │    ├──────────┤      │
    │ Runs     │    │    │ Runs     │      │
    ├──────────┤    │    ├──────────┤      │
    │ Usage    │    │    │ Usage    │      │
    └──────────┘    │    └──────────┘      │
                    └──────────────────────┘

隔离维度:
    数据隔离:   所有表按 tenant_id 分区, 查询自动注入 WHERE tenant_id = ?
    执行隔离:   不同租户的 agent 在不同 Sandbox 中执行
    配额隔离:   每个租户独立的 token/并发/存储 配额
    密钥隔离:   每个租户独立的加密密钥空间

8.2 RBAC

角色定义:

Admin:
    - tenants:*           (管理租户)
    - agents:*            (所有 agent 操作)
    - keys:*              (管理 API Key)
    - billing:read        (查看账单)
    - secrets:*           (管理密钥)

Developer:
    - agents:read,write,execute   (创建/修改/执行 agent)
    - keys:read,write             (管理自己的 API Key)
    - runs:read                   (查看执行记录)
    - secrets:read                (读取密钥, 不可查看明文)

Viewer:
    - agents:read          (查看 agent 定义)
    - runs:read            (查看执行记录)
    - billing:read         (查看账单)

8.3 配额管理

class QuotaManager:
    """
    配额检查与扣减。

    检查时机:
        1. 请求入口 (Gateway 层) — 预检, 快速拒绝明显超限请求
        2. 执行前 (Engine 层) — 精确检查, 原子扣减
        3. 执行后 (Billing 层) — 记录实际消耗, 补偿预扣差额
    """

    async def check_and_reserve(self, tenant_id: UUID, estimated_tokens: int):
        quota = await self.get_quota(tenant_id)
        used = await self.get_monthly_usage(tenant_id)
        remaining = quota.tokens_monthly - used.tokens

        if estimated_tokens > remaining:
            raise QuotaExceededError(
                limit=quota.tokens_monthly,
                used=used.tokens,
                requested=estimated_tokens
            )

        # 预扣 (乐观锁)
        await self.reserve(tenant_id, estimated_tokens)

    async def settle(self, tenant_id: UUID, reserved: int, actual: int):
        # 结算: 退还预扣多余部分
        diff = reserved - actual
        if diff > 0:
            await self.release(tenant_id, diff)

9. Observability(可观测性)

9.1 三支柱

Trace (追踪)
    │
    │  一次 agent 执行的完整调用链
    │
    │  AgentPaaS span (gateway → dispatch → execute)
    │      └── lambdagent span (β-规约步骤)
    │              ├── LLM call span
    │              ├── Tool call span
    │              └── Memory read/write span
    │
    │  实现: OpenTelemetry SDK → Jaeger / Tempo
    │
Metrics (指标)
    │
    │  聚合数值, 用于告警和仪表盘
    │
    │  agentpaas_requests_total{agent_id, status, method}
    │  agentpaas_request_duration_seconds{agent_id, quantile}
    │  agentpaas_tokens_total{agent_id, tenant_id, direction}
    │  agentpaas_active_jobs{tenant_id}
    │  agentpaas_queue_depth{priority}
    │  agentpaas_circuit_breaker_state{provider}
    │  agentpaas_cache_hit_ratio{cache}
    │
    │  实现: Prometheus client → Prometheus → Grafana
    │
Logs (日志)
    │
    │  结构化事件, 用于排查
    │
    │  {
    │      "timestamp": "...",
    │      "level": "info",
    │      "service": "agentpaas",
    │      "trace_id": "tr_a1b2c3",
    │      "tenant_id": "...",
    │      "agent_id": "ag_3f8a1b2c",
    │      "event": "agent.run.completed",
    │      "duration_ms": 4520,
    │      "tokens": 1630
    │  }
    │
    │  实现: structlog → stdout (JSON) → Loki / ELK

9.2 Trace 桥接

AgentPaaS 的分布式 Trace 与 lambdagent 的 Context.trace 桥接:

lambdagent Context.trace:
    [TraceEntry(step=1, type="think", input=..., output=..., duration_ms=...),
     TraceEntry(step=2, type="tool_call", tool="order_api", ...),
     TraceEntry(step=3, type="answer", ...)]

桥接逻辑:
    for entry in ctx.trace:
        with tracer.start_span(entry.type) as span:
            span.set_attribute("step", entry.step)
            span.set_attribute("duration_ms", entry.duration_ms)
            span.set_attribute("tokens", entry.tokens)
            if entry.type == "tool_call":
                span.set_attribute("tool", entry.tool)

9.3 告警规则

alerts:
  - name: high_error_rate
    condition: rate(agentpaas_requests_total{status="error"}[5m]) > 0.05
    severity: critical
    message: "Agent {agent_id} 错误率超过 5%"

  - name: token_budget_warning
    condition: agentpaas_tokens_monthly_used / agentpaas_tokens_monthly_quota > 0.8
    severity: warning
    message: "租户 {tenant_id} 月 token 用量已达 80%"

  - name: high_latency
    condition: histogram_quantile(0.95, agentpaas_request_duration_seconds) > 30
    severity: warning
    message: "Agent {agent_id} P95 延迟超过 30s"

  - name: queue_backlog
    condition: agentpaas_queue_depth > 100
    severity: warning
    message: "任务队列积压超过 100"

10. Secret Management(密钥管理)

10.1 存储模型

secrets 表:

| 字段 | 类型 | 说明 |
|------|------|------|
| id | UUID | 主键 |
| tenant_id | UUID | FK → tenants |
| name | VARCHAR(128) | 密钥名 (如 "openai_api_key") |
| encrypted_value | BYTEA | AES-256-GCM 加密后的值 |
| environment | VARCHAR(32) | dev / staging / production |
| created_by | UUID | FK → users |
| updated_at | TIMESTAMPTZ | 最后更新时间 |

加密方案:
    主密钥 (Master Key): 环境变量 AGENTPAAS_MASTER_KEY 或 KMS
    加密: AES-256-GCM(master_key, plaintext) → (ciphertext, nonce, tag)
    存储: nonce + tag + ciphertext 合并存入 encrypted_value
    解密: 仅在 agent 执行时在 Worker 进程内解密, 不经过 API 层

10.2 注入方式

# Agent 执行时, 密钥注入 lambdagent Context:

async def inject_secrets(ctx: Context, tenant_id: UUID, env: str):
    secrets = await vault.get_all(tenant_id, environment=env)
    for secret in secrets:
        # 密钥以环境变量形式注入, lambdagent 通过 os.environ 或 ctx 读取
        ctx.set(f"secret.{secret.name}", vault.decrypt(secret.encrypted_value))

# lambdagent agent 配置中引用:
# model:
#   provider: openai
#   apiKey: ${secret.openai_api_key}    ← AgentPaaS 注入

11. Billing & Metering(计费与计量)

11.1 计量点

每次 agent 执行记录:

UsageRecord:
    tenant_id       — 归属租户
    agent_id        — 归属 agent
    run_id          — 关联执行
    model           — 使用的模型 (如 "claude-sonnet-4-20250514")
    input_tokens    — 输入 token 数
    output_tokens   — 输出 token 数
    tool_calls      — 工具调用次数
    duration_ms     — 执行耗时
    timestamp       — 发生时间

来源: lambdagent Context.trace 中的 TraceEntry
精度: 使用 LLM 返回的精确 token 计数 (非估算)

11.2 成本计算

# 单价表 (可配置, 按模型)
PRICING = {
    "claude-sonnet-4-20250514": {"input": 3.0, "output": 15.0},   # $/M tokens
    "gpt-4o":                   {"input": 2.5, "output": 10.0},
    "qwen-max":                 {"input": 1.0, "output": 2.0},
}

def calculate_cost(record: UsageRecord) -> Decimal:
    prices = PRICING[record.model]
    input_cost  = Decimal(record.input_tokens)  / 1_000_000 * Decimal(str(prices["input"]))
    output_cost = Decimal(record.output_tokens) / 1_000_000 * Decimal(str(prices["output"]))
    return input_cost + output_cost

11.3 用量报告 API

GET /api/v1/billing/usage?period=2026-03&group_by=agent

Response:
{
    "period": "2026-03",
    "tenant_id": "...",
    "total_cost_usd": 142.50,
    "total_tokens": 12_500_000,
    "breakdown": [
        {
            "agent_id": "ag_3f8a1b2c",
            "agent_name": "customer-support",
            "runs": 2340,
            "input_tokens": 5_200_000,
            "output_tokens": 1_800_000,
            "cost_usd": 89.40
        },
        ...
    ]
}

12. 项目结构

agentpaas/
├── api/                          # HTTP 接口层
│   ├── __init__.py
│   ├── app.py                    # FastAPI 应用入口
│   ├── deps.py                   # 依赖注入 (DB session, current_tenant, ...)
│   ├── v1/
│   │   ├── __init__.py
│   │   ├── agents.py             # Agent CRUD + 执行端点
│   │   ├── jobs.py               # 异步任务管理端点
│   │   ├── auth.py               # 认证端点 (API Key 管理)
│   │   ├── billing.py            # 计费查询端点
│   │   ├── admin.py              # 租户管理端点 (Admin only)
│   │   └── discovery.py          # A2A 发现端点
│   └── middleware/
│       ├── __init__.py
│       ├── auth.py               # 认证中间件 (API Key / JWT)
│       ├── rate_limit.py         # 限流中间件
│       ├── tenant.py             # 租户上下文注入
│       └── request_id.py         # 请求 ID 生成
│
├── engine/                       # 执行引擎层
│   ├── __init__.py
│   ├── dispatcher.py             # 请求调度 (sync/async/stream)
│   ├── sandbox.py                # 沙箱隔离 (process/container)
│   ├── worker.py                 # Worker 进程 (消费 Job Queue)
│   ├── queue.py                  # Job Queue 抽象 (memory/redis)
│   ├── stream.py                 # SSE / WebSocket 流式适配
│   ├── retry.py                  # 重试策略 + 断路器
│   └── cache.py                  # Term 编译缓存
│
├── registry/                     # Agent 注册中心
│   ├── __init__.py
│   ├── store.py                  # Agent CRUD 持久化
│   ├── version.py                # 版本管理 + 回滚
│   ├── traffic.py                # 流量治理 (Canary)
│   └── discovery.py              # A2A 服务发现
│
├── tenant/                       # 多租户
│   ├── __init__.py
│   ├── models.py                 # Tenant / User / Role ORM
│   ├── isolation.py              # 数据隔离 (自动注入 tenant_id)
│   ├── quota.py                  # 配额检查 + 预扣 + 结算
│   └── rbac.py                   # 角色权限检查
│
├── observability/                # 可观测性
│   ├── __init__.py
│   ├── tracing.py                # OpenTelemetry 集成 + lambdagent Trace 桥接
│   ├── metrics.py                # Prometheus 指标定义 + 导出
│   └── logging.py                # structlog 配置
│
├── secrets/                      # 密钥管理
│   ├── __init__.py
│   ├── vault.py                  # AES-256-GCM 加解密
│   └── inject.py                 # 执行时密钥注入
│
├── billing/                      # 计费
│   ├── __init__.py
│   ├── metering.py               # 用量记录
│   ├── pricing.py                # 成本计算
│   └── report.py                 # 用量报告生成
│
├── db/                           # 数据库
│   ├── __init__.py
│   ├── models.py                 # SQLAlchemy ORM 模型
│   ├── session.py                # 数据库连接管理
│   └── migrations/               # Alembic 迁移
│       └── ...
│
├── cli/                          # 管理 CLI
│   └── main.py                   # typer CLI (create-tenant, import-agent, ...)
│
├── config.py                     # 全局配置 (Pydantic Settings)
├── __init__.py
├── pyproject.toml
├── Dockerfile
├── docker-compose.yml            # 本地开发 (PostgreSQL + Redis)
├── SPEC.md                       # ← 本文档
├── IDEA.md
└── README.md

13. 配置规格

# config.py — Pydantic Settings, 支持环境变量 + .env 文件

class AgentPaaSConfig(BaseSettings):
    # 服务
    host: str = "0.0.0.0"
    port: int = 8000
    workers: int = 4
    debug: bool = False

    # 数据库
    database_url: str = "postgresql+asyncpg://localhost/agentpaas"
    # 开发模式可用: "sqlite+aiosqlite:///./agentpaas.db"

    # Redis
    redis_url: str = "redis://localhost:6379/0"
    # 开发模式可用: None (降级为内存实现)

    # lambdagent
    lambdagent_sandbox_level: int = 1           # 0=inprocess, 1=subprocess, 2=container
    lambdagent_default_timeout: int = 120       # 秒
    lambdagent_max_steps: int = 100
    lambdagent_max_tokens: int = 100_000

    # 认证
    master_key: str                             # 密钥加密主密钥 (必填)
    jwt_secret: str = ""                        # JWT 签名密钥 (可选)
    jwt_algorithm: str = "HS256"

    # 限流
    global_rate_limit: int = 10_000             # 全局 req/min
    default_tenant_rate_limit: int = 60         # 默认租户 req/min

    # 可观测性
    otel_endpoint: str = ""                     # OpenTelemetry Collector
    otel_service_name: str = "agentpaas"
    prometheus_enabled: bool = True
    log_level: str = "INFO"
    log_format: str = "json"                    # json / console

    class Config:
        env_prefix = "AGENTPAAS_"
        env_file = ".env"

14. 部署架构

14.1 单机开发模式

                    ┌─────────────┐
                    │   FastAPI   │
                    │  (uvicorn)  │
                    ├─────────────┤
                    │   SQLite    │
                    │   Memory Q  │
                    │   L0 Sandbox│
                    └─────────────┘

启动: agentpaas serve --dev
依赖: Python 3.10+, lambdagent
无需: PostgreSQL, Redis, Docker

14.2 生产 Docker Compose

services:
  api:
    build: .
    ports: ["8000:8000"]
    environment:
      AGENTPAAS_DATABASE_URL: postgresql+asyncpg://postgres:pass@db/agentpaas
      AGENTPAAS_REDIS_URL: redis://redis:6379/0
      AGENTPAAS_MASTER_KEY: ${MASTER_KEY}
    depends_on: [db, redis]

  worker:
    build: .
    command: agentpaas worker --concurrency 4
    environment: *api-env
    depends_on: [db, redis]

  db:
    image: postgres:16
    volumes: ["pgdata:/var/lib/postgresql/data"]

  redis:
    image: redis:7-alpine

  prometheus:
    image: prom/prometheus
    volumes: ["./prometheus.yml:/etc/prometheus/prometheus.yml"]

  jaeger:
    image: jaegertracing/all-in-one
    ports: ["16686:16686"]

14.3 Kubernetes 生产部署

┌─────────────────────────────────────────────────────────┐
│                    Kubernetes Cluster                    │
│                                                         │
│  ┌──────────┐   ┌──────────┐   ┌──────────────────┐    │
│  │ Ingress  │──▶│ API Pods │──▶│  Worker Pods      │    │
│  │ (nginx)  │   │ (HPA)    │   │  (HPA by queue)   │    │
│  └──────────┘   └──────────┘   └──────────────────┘    │
│                       │                │                │
│                  ┌────▼────┐     ┌─────▼──────┐        │
│                  │PostgreSQL│     │   Redis     │        │
│                  │(StatefulSet)   │(StatefulSet)│        │
│                  └─────────┘     └────────────┘        │
│                                                         │
│  HPA 策略:                                              │
│    API Pods:    CPU > 60% → 扩容, 最小 2, 最大 20       │
│    Worker Pods: queue_depth > 10 → 扩容, 最小 1, 最大 50│
└─────────────────────────────────────────────────────────┘

15. 错误码规范

错误响应格式:
{
    "error": {
        "code": "QUOTA_EXCEEDED",
        "message": "Monthly token quota exceeded",
        "details": {
            "limit": 1000000,
            "used": 998500,
            "requested": 5000
        }
    }
}

HTTP 状态码映射:

| 状态码 | 错误码 | 说明 |
|--------|--------|------|
| 400 | INVALID_REQUEST | 请求格式错误 |
| 400 | INVALID_CONFIG | Agent YAML 配置无效 |
| 401 | UNAUTHORIZED | API Key 缺失或无效 |
| 401 | KEY_EXPIRED | API Key 已过期 |
| 403 | FORBIDDEN | 无权限执行此操作 |
| 404 | AGENT_NOT_FOUND | Agent 不存在 |
| 404 | JOB_NOT_FOUND | Job 不存在 |
| 404 | VERSION_NOT_FOUND | Agent 版本不存在 |
| 409 | AGENT_ALREADY_EXISTS | Agent 名称冲突 |
| 422 | VALIDATION_ERROR | 请求参数校验失败 |
| 429 | RATE_LIMITED | 请求频率超限 |
| 429 | QUOTA_EXCEEDED | Token 配额超限 |
| 429 | CONCURRENCY_EXCEEDED | 并发执行数超限 |
| 500 | INTERNAL_ERROR | 内部错误 |
| 502 | LLM_ERROR | LLM 提供商返回错误 |
| 503 | CIRCUIT_OPEN | 断路器打开, LLM 暂不可用 |
| 504 | EXECUTION_TIMEOUT | Agent 执行超时 |

16. 与 lambdagent 的接口边界

AgentPaaS 使用的 lambdagent 公共 API (且仅限这些):

编译层:
    lambdagent.from_config(yaml: dict) → Term
    lambdagent.from_config(path: str) → Term

执行层:
    lambdagent.Runtime(config: dict) → Runtime
    Runtime.reduce(term: Term, input: Any, ctx: Context) → Any

上下文:
    lambdagent.Context() → Context
    Context.set(key: str, value: Any)
    Context.trace → list[TraceEntry]
    Context.total_tokens → int

追踪:
    TraceEntry.step: int
    TraceEntry.type: str
    TraceEntry.input: Any
    TraceEntry.output: Any
    TraceEntry.duration_ms: int
    TraceEntry.tokens: int
    TraceEntry.model: str
    TraceEntry.tool: str | None

Checkpoint:
    lambdagent.Checkpoint.save(ctx: Context, path: str)
    lambdagent.Checkpoint.load(path: str) → Context

A2A:
    lambdagent.A2AServer(agent: Term, metadata: dict)

如果 lambdagent 需要扩展接口 (如 reduce_iter 用于流式),
应在 lambdagent 侧新增, 而非 AgentPaaS 侧 hack。

17. MVP 路线图

Phase 1: 骨架可跑(2 周)

目标: 能通过 HTTP API 创建 agent 并同步执行

交付:
    [x] FastAPI 项目骨架 + config.py
    [x] SQLite 数据库 + Agent/Run ORM
    [x] POST /agents (创建)
    [x] GET /agents/{id} (查询)
    [x] POST /agents/{id}/run (同步执行, L0 sandbox)
    [x] API Key 认证 (简单版)
    [x] 基础 structlog 日志
    [x] docker-compose (PostgreSQL + 应用)

Phase 2: 异步与隔离(4 周)

目标: 支持异步执行、流式输出、基本多租户

交付:
    [ ] Redis Job Queue + Worker 进程
    [ ] POST /agents/{id}/jobs (异步执行)
    [ ] POST /agents/{id}/stream (SSE)
    [ ] L1 Sandbox (子进程隔离)
    [ ] 多租户数据隔离
    [ ] 配额检查 + 预扣
    [ ] 令牌桶限流
    [ ] Agent 版本管理 + 回滚
    [ ] OpenTelemetry Trace
    [ ] Prometheus Metrics

Phase 3: 平台化(8 周)

目标: 生产可用的完整平台

交付:
    [ ] L2 Sandbox (容器隔离)
    [ ] RBAC 权限体系
    [ ] 密钥管理 (AES-256-GCM)
    [ ] 计费系统 + 用量报告
    [ ] Canary 灰度发布
    [ ] A2A 服务发现
    [ ] 告警规则引擎
    [ ] CLI 管理工具
    [ ] K8s Helm Chart
    [ ] Web Dashboard (基础版)

18. Implementation Status (as of v0.1.0)

Phase 1: Skeleton ✅ COMPLETE

Feature File(s) Status
FastAPI skeleton + config api/app.py, config.py
SQLite DB + 7 tables ORM db/models.py, db/session.py
Agent CRUD API api/v1/agents.py
Sync execution /run api/v1/agents.py
API Key auth (SHA-256) api/middleware/auth.py
Tenant management api/v1/admin.py
Structured logging observability/logging.py
Docker + docker-compose Dockerfile, docker-compose.yml

Phase 2: Async + Isolation ✅ COMPLETE

Feature File(s) Status
Job Queue (memory + Redis) engine/queue.py
Worker process engine/worker.py
Dispatcher (sync/async/stream) engine/dispatcher.py
SSE stream + WebSocket engine/stream.py
L0/L1 Sandbox engine/sandbox.py
Retry + Circuit Breaker engine/retry.py
Compile Cache (LRU) engine/cache.py
Agent version management registry/version.py
Canary traffic routing registry/traffic.py
A2A agent discovery registry/discovery.py
Multi-tenant data isolation tenant/isolation.py
Token quota + concurrency tenant/quota.py
RBAC (admin/dev/viewer) tenant/rbac.py
OpenTelemetry trace bridge observability/tracing.py
Prometheus-style metrics observability/metrics.py
AES-256-GCM secret vault secrets/vault.py
Secret injection secrets/inject.py
Token metering billing/metering.py
Per-model pricing (6 models) billing/pricing.py
Monthly usage reports billing/report.py

Phase 3: Platform (In Progress)

Feature Status Notes
L2 Container sandbox Requires Docker-in-Docker
Web Dashboard React + React Flow planned
K8s Helm Chart
Full RBAC enforcement 🔶 Model defined, middleware pending
Webhook notifications

CLI Provider Management ✅

Feature Status
provider list (7 built-in)
provider add (custom OpenAI-compatible)
provider test (connection test)
provider models (query available models)
provider remove

Built-in providers: Anthropic, OpenAI, DashScope, DeepSeek, Zhipu AI, Moonshot, Ollama

Test Coverage

Module Tests Status
lambdagent.core 7
lambdagent.primitives 10
lambdagent.extensions 10
lambdagent.dataset 2
lambdagent.fromconfig 9
lambdagent.agentruntime 16
agentpaas (engine/tenant/billing) 16
Total 70 All pass (0.24s)

Code Statistics

Component Files Lines
lambdagent (kernel) 48 ~12,800
agentpaas (platform) 50 ~2,000
tests 8 ~640
Total 106 ~15,400