# 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 认证机制 ```python # 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) ```json // agents.traffic_rules 字段 { "rules": [ { "version": 4, "weight": 10, "match": { "header": "X-Canary: true" } }, { "version": 3, "weight": 90, "match": "default" } ] } ``` 路由逻辑: ```python 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(调度器) ```python 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 ```python 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 重试与断路器 ```python 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 ```python 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 配额管理 ```python 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 告警规则 ```yaml 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 注入方式 ```python # 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 成本计算 ```python # 单价表 (可配置, 按模型) 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. 配置规格 ```python # 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 ```yaml 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** |