> For AI agents: the complete documentation index is available at /llms.txt, the full documentation bundle is available at /llms-full.txt.

# Unified Pipeline Architecture 设计 Brief

> **这是什么**: 给新 session 的设计任务。不是 session 总结，不是 bug 修复清单。
> **当前状态**: Canonical unified pipeline brief。后续设计 session 优先读这份；`unified-pipeline-reference-checklist.md` 是 checklist/reference，`unified-pipeline-research-appendix.md` 是代码调研 appendix。
> **日期**: 2026-05-31
> **前序工作**: 2026-05-30 完成了 task pipeline 设计（team 会议讨论通过）。在设计 task pipeline 的过程中发现一个更大的问题：task 只是系统中的一条 pipeline，所有 pipeline（call/SMS/lead/contact analysis）共享同一批表但各自独立设计，职责划分和写入逻辑不统一。本 brief 是从 task pipeline 出发向外辐射的全系统设计。
> **前序材料**:
>
> - [Task Pipeline Deliverable (Codex)](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) — 已讨论通过的 task 架构方案
> - [Task Pipeline Architecture Addendum](/product-design/v2/tasks-feature/design/research/task-pipeline-architecture-addendum.md) — Task Orchestrator 如何接入 unified architecture
> - [Task Pipeline 设计 Brief](/product-design/v2/tasks-feature/design/research/task-pipeline-design-brief.md) — task 设计的方法论和思考框架（本 brief 的格式参考）
> - [Unified Pipeline 现状调研 (Codex)](/product-design/v2/unified-pipeline/research/unified-pipeline-current-state-codex.md) — 当前实现事实调研，区分 target docs / test current code / prod legacy
> - [Unified Pipeline Phase 1 (Codex)](/product-design/v2/unified-pipeline/original/unified-pipeline-phase-1-codex.md) — 第一阶段落地方案
> - [Unified Pipeline Target State (Codex)](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md) — 最理想的目标架构
> - [Unified Pipeline Research Appendix](/product-design/v2/unified-pipeline/research/unified-pipeline-research-appendix.md) — 2026-05-31 代码调研结果，用作 facts / appendix
> - [Pipeline 全景分析](/system-design/pipeline-architecture-analysis.md) — 5 条 pipeline 的代码级 trace + cross-pipeline 依赖矩阵
> - [产品设计准则](/product-design/v2.md) — Schema-first → API-first → State-machine-first → Prompt-last

***

## 你要做什么

**审视 retaintive 现有系统的所有 pipeline 和 business object，检查哪些设计是对的、哪些有问题、哪些概念混在一起了。然后基于检查结果，设计一个统一的架构方向：怎么把散落的写入逻辑统一起来、代码和 AI 各自负责什么、现有模式和未来 tool calling 模式怎么衔接。**

这不是从零设计一个新系统，是看我们现在有什么、哪些对哪些不对、然后怎么把它弄成统一的。

你的角色：一个精通 AI 技术、懂 SaaS 产品架构的 Principal Engineer / AI Architect。

### 阅读顺序

给团队 review 时，建议按"结论先行"读：

1. [Target State Architecture](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md)：先看最终理想状态长什么样。
2. [Phase 1 Implementation Design](/product-design/v2/unified-pipeline/original/unified-pipeline-phase-1-codex.md)：再看第一阶段怎么落地。
3. [Current State Research](/product-design/v2/unified-pipeline/research/unified-pipeline-current-state-codex.md)：最后看代码事实和 gap evidence。

如果是新 session 要继续做 design work，再按这个顺序读：

1. 本文：理解任务目标、问题边界、acceptance criteria。
2. [Current State Research](/product-design/v2/unified-pipeline/research/unified-pipeline-current-state-codex.md)：先看当前代码事实，不把 target docs 当实现。
3. [Phase 1 Implementation Design](/product-design/v2/unified-pipeline/original/unified-pipeline-phase-1-codex.md)：看第一阶段怎么落地。
4. [Target State Architecture](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md)：看最理想的长期状态。

***

## 我们遇到的困境

在做 task pipeline 设计时发现了一个更大的问题。task 只是系统中的一条 pipeline，但我们还有 call analysis、contact analysis、SMS、lead 等多条 pipeline。这些 pipeline 各自独立设计，导致：

**困境 1：每加一个新功能，改动成本指数增长。**

比如现在要加 "SMS 触发 AI 分析" — 不是简单写一个新模块就行。你需要：在 message-processor 里加触发逻辑 → contacts-analyzer 的 prompt 要改 → contacts 表的写入逻辑要跟其他 pipeline 对齐 → task 创建要走 Orchestrator → timeline 要写新 event type。每一步都要跟其他 pipeline 的实现保持一致，但这些实现散落在不同 Lambda 里，每次都要翻一遍确认"别人是怎么写的"。

**困境 2：代码和 AI 的职责边界是 case-by-case 画的，没有统一标准。**

call analysis 的 3-stage prompt pipeline 里，triage 是 AI 做的；但 pre-triage（关键词匹配）是代码做的。contacts-analyzer 的 task 决策是 AI 提案 + 代码执行。lead-processor 的 task 创建是纯代码硬编码的。每条 pipeline 各自决定"这个该 AI 做还是代码做"，没有一个统一的判断框架。以后再加新 pipeline 或改现有的，又要重新想一遍。

**困境 3：不确定现有设计对不对。**

task 和 contact 的某些概念可能混在一起了（比如 contact 的 `actionNeeded` 和 task 的 `actionNeeded` 是什么关系？）。call 的 disposition 和 task 的 progress 是不是同一件事？这些不确定需要做一次系统性检查——不是说设计错了，是不确定，需要验证。

**困境 4：想加新能力（retry、onboard 扫描、batch 重分析）时，发现没有可复用的模块。**

比如员工想 "retry 一次这通电话的 AI 分析"，或者新店铺 onboard 想 "扫描过去两周的通话"。这些操作不应该是新 pipeline，应该是复用现有模块的不同调用方式。但现在模块之间耦合太紧，没有设计成可复用的。

**困境 5：不知道未来方向应该怎么走。**

现在是 one-time AI call 模式（把所有数据塞进 prompt 一次性调 AI），以后可能要用 tool calling 模式（AI 自己决定要查什么数据）。但不确定现在的架构怎么演进到那个方向，也不确定 tool calling 对我们有没有必要。

### 核心问题（一句话）

**系统越来越复杂，每条 pipeline 都手动定义 + 硬编码设计，改动和调试成本指数增长。有没有更好的方式？哪些必须靠代码，哪些可以交给 AI / tool use / 其他工具？业界在怎么做？**

***

## Expected Output

新 session 完成后应该交付这些：

1. **Schema / DB Contract Check** — 以 `callytics-common/src/db/schema/` 为 source of truth，检查 `contacts` / `tasks` / `calls` / `messages` / `leads` / `contact_timeline` 当前 schema 是否支持目标职责。
2. **Current Design / Prompt Inventory** — 列出现有 pipeline、prompt surface、structured output、code-only rules、writer 入口。重点不是数 prompt，而是判断每个 prompt owns 哪类 judgment。
3. **Business object 健康检查** — task、contact、call、lead、message、timeline 各自的职责边界是否清晰，有没有概念混淆。先检查，再设计。
4. **Current vs Target Workflow 图** — 至少两张图：现状（code/prompt/writer 混在一起）vs 目标（pipeline 负责触发和 judgment，shared modules 负责 mutation）。每个节点标清 `[Code]` / `[Prompt]` / `[Hybrid]` / `[Shared Code]`。
5. **两层架构设计** — Layer 1: Shared Mutation Modules（Task Orchestrator + Contact Writer + Timeline Writer + Policy Guard）；Layer 2: Processing Capability / Invocation Modules（Call Analysis、SMS Signal、Contact Analysis Reprocess、Lead Downstream）。
6. **Code vs AI 职责分界框架** — 一个适用于所有 pipeline 的通用标准（不是每条 pipeline 单独决定）。
7. **Schema / API / AI Agent 反推检查** — 从 Human UI、Backend API、AI structured output、future tool API、voice agent、audit/replay 的视角反推数据模型是否对。
8. **现有模式 vs tool calling 模式的设计兼容** — 确保方案 A 做好后，方案 B 可以无缝切换（见下方对比表）。
9. **Phase 1 Implementation Design** — 明确第一阶段应该做什么，不只是愿景。至少覆盖 Task Orchestrator contract、record\_progress API、task\_progress\_events、timeline projection、contacts-analyzer taskDecisions schema，以及 Call/SMS/Lead invocation contract 的边界。

> 不是写代码。是设计 + Max review 后才执行。
>
> **呈现要求**: 产出必须让不在这个 session 里的人也能 10 分钟看懂核心结论。每个交付物先用一句话说结论，再展开细节。用表格和图而不是长段落。技术细节不省略但要有层次（先 what + why，想看 how 再往下读）。

### Design Acceptance Criteria

设计合格的标准（不满足就不算完成）：

- 必须回答：**现在是否应该改成 tool calling 模式**。不能含糊
- 必须列出：**每个 business object（task、contact、call、lead、message、timeline）的写入入口** — 谁写、通过什么模块写、有没有统一
- 必须覆盖：至少 3 个**复用场景**的设计（如：call retry、新用户 onboard 扫描、SMS 触发 AI 分析）
- 必须覆盖：**两层架构**。不要只设计 shared writers，也要设计 Call/SMS/Lead/reprocess 这些 reusable processing capability 怎么被调用
- 必须明确：**Call Analysis Module / SMS Signal Module 不是新的 writer**，它们是 invocation/capability contract，写入仍通过 shared mutation modules
- 必须产出：**Schema/API check** 和 **current vs target workflow diagrams**
- 必须区分：**Phase 1 必做**（统一抽象）和 **future-compatible 设计**（tool calling、voice agent）

### Non-goals（这次设计明确不做的）

- ❌ 不实现 tool calling / agentic loop — 只确保架构兼容
- ❌ 不设计 voice agent 产品 — 只确保 tool API 接口预留
- ❌ 不重写现有 pipeline — 设计的是"共享模块"怎么插入现有 pipeline
- ❌ 不做 production rollout 方案
- ❌ 不改代码、不跑 DB 操作、不 push、不 PR

### 必须显式回答的冲突问题

Brief 里有几组方向天然有张力。**不要默认调和，每个都逐一给出明确结论 + 理由：**

1. **"每条 pipeline 各自优化" vs "统一成共享模块"** — 统一的边界在哪？所有 contacts 写入都走 Contact Writer，还是允许简单 writer（如 message-processor 只更新 lastActivityAt）直接写？
2. **"one-time AI call（成本低、可预测）" vs "tool calling（灵活、可扩展）"** — 什么场景用 A 什么场景用 B？是非此即彼还是共存？
3. **"prompt 做更多判断（灵活）" vs "代码做更多控制（确定性）"** — 边界画在哪？contacts-analyzer 的 1040 行 prompt 是应该继续扩大还是应该拆分？
4. **"现在就做好统一抽象" vs "先做 task pipeline 再说"** — 统一抽象是不是 premature？还是现在不做以后更难？
5. **"contact 的 `actionNeeded` / `suggestedActions`" vs "task 的 `actionNeeded` / `suggestedActions`"** — 是冗余？还是不同语义？应该保留还是合并？
6. **"SMS 不触发 AI（省钱）" vs "有意义 SMS 需要即时跟进（及时性）"** — 怎么定义"有意义"？是代码规则过滤还是 AI 判断？
7. **"Call/SMS 是不是也要 module" vs "只有 shared tables 才需要 module"** — Call/SMS 不需要 writer module，但需要 processing capability / invocation contract。不要把这两类 module 混在一起。

***

## 背景：两种模式的对比

这是这次设计的核心框架。所有设计决策都在回答：**这个操作在哪种模式下怎么实现。**

### 方案 A: One-time AI call + 统一抽象（现在做）

```
代码提前聚合所有数据 → 一次性喂给 AI → AI 返回结构化 JSON → 代码验证执行

改进点（这次设计的目标）：
1. Task Orchestrator — 统一 task 写入
2. Contact Writer  — 统一 contacts 写入（trust score 集中）
3. Timeline Writer — 统一 timeline 写入（event catalog 集中）
4. Schema-driven prompt — JSON Schema 定义 AI 能做什么
```

### 方案 B: Tool calling（以后做，基于 A）

```
定义 tools:
  - read_contact(phone, storeId)
  - get_recent_calls(phone, limit)
  - create_task(type, category, priority, ...)
  - close_task(taskId, closeResult)
  - update_contact(phone, storeId, fields)

AI 自己决定要查什么 → 调 tool → 拿结果 → 继续推理 → 再调 tool → 完成

每个 tool 的实现 = 调方案 A 的统一模块
```

### 对比

|    | 方案 A（现在）               | 方案 B（以后）                  |
| -- | ---------------------- | ------------------------- |
| 信息 | 代码提前聚合所有数据塞进 prompt    | AI 自己决定要查什么数据             |
| 逻辑 | prompt 里定义"遇到 X 做 Y"   | AI 根据情况灵活选择工具             |
| 扩展 | 加新功能 = 改 prompt + 改代码  | 加新功能 = 注册一个新 tool         |
| 成本 | 一次 AI 调用（prompt 大但调用少） | 多轮 AI 调用（每轮 prompt 小但调用多） |
| 适用 | batch 处理（每天处理数百联系人）    | 交互式（用户问一个问题、AI 诊断）        |
| 延迟 | 可预测（一次调用）              | 不可预测（多轮往返）                |

### 关键设计原则

**A 做好了，B 自然就有了。** Task Orchestrator 就是 `create_task` tool 的实现；Contact Writer 就是 `update_contact` tool 的实现。不需要做两次设计。

**A 和 B 可以共存。** batch 场景继续用 A（cost-efficient），交互场景用 B（flexible）。不是非此即彼。随时可以切，也可以两种模式同时跑。

***

## 怎么做

### Step 1: 理解现有系统的 business objects

**不要直接看代码实现，先从业务语义出发。** 搞清楚系统里有几个核心 business object，每个的 lifecycle 是什么，它们之间的关系是什么。

**需要回答的问题**：

| Object             | 该问的                                                                                                   | 参考                                                                                                              |
| ------------------ | ----------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------- |
| **Contact**        | 一个 contact 从出现到最终状态，经历哪些阶段？lifecycleStage（lead/member/churned）的转换规则是什么？谁有权改？                          | `callytics-common/src/db/schema/contacts.ts`                                                                    |
| **Task**           | task 的 lifecycle（open → closed）已经在 deliverable 里设计好了，这里检查跟 contact lifecycle 有没有冲突                    | [Task Pipeline Deliverable (Codex)](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) |
| **Call**           | 一通电话从 RC webhook 到最终存储，数据会被改几次、被谁改？                                                                   | `callytics-common/src/db/schema/calls.ts`                                                                       |
| **Lead**           | lead 从邮件进来到被联系，lifecycle 是什么？lead 和 contact 的关系是什么？（同一个人同一个电话号码？）                                     | `callytics-common/src/db/schema/leads.ts`                                                                       |
| **Message**        | SMS 消息的 lifecycle。message 和 contact 的关系。SMS 内容是否应该触发 AI 分析？                                           | `callytics-common/src/db/schema/messages.ts`                                                                    |
| **Timeline Event** | contact\_timeline 是 source of truth 还是 projection？它跟 task\_progress\_events（deliverable 提议的新表）的关系是什么？ | `callytics-common/src/db/schema/contact-timeline.ts`                                                            |

**特别要检查的**：

- task 和 contact 的概念有没有混在一起（比如：contact 的 `actionNeeded` 字段和 task 的 `actionNeeded` 字段是什么关系？）
- call 的 disposition（`callState`）和 task 的 progress 是什么关系？是两个独立的概念还是同一件事？
- lead 和 contact 在什么节点合并？lead-processor 写 contacts 表时，如果 contact 已存在怎么处理？

### Step 2: 设计统一抽象层（两层）

这两层解决不同的问题，但必须一起设计。

**Layer 1: Shared Mutation Modules — "写的统一"**

多条 pipeline 往同一张表写数据，写入逻辑必须统一。

**Task Orchestrator**（已有初步设计，需细化）

- 输入：action（create/update/close）+ 参数
- 输出：写入 tasks + contact\_timeline（原子事务）
- 保证：幂等（per-category dedup）、权限（storeId 隔离）、audit（timeline event）
- 调用方：lead-processor、contacts-analyzer、studio-api、未来的 AI agent

**Contact Writer**（新设计）

- 输入：phone + storeId + 要更新的字段 + 来源（trust level）
- 输出：UPSERT contacts + contact\_timeline
- 保证：trust score 逻辑集中（WEBHOOK=20 \< AI\_TRANSCRIPT=40 \< RC\_API=60 \< LEAD=80 \< STAFF=100）、只升不降、GREATEST(lastActivityAt)
- 调用方：所有写 contacts 的 pipeline

**Timeline Writer**（新设计）

- 输入：event\_type + event\_category + entity + actor + payload
- 输出：INSERT contact\_timeline
- 保证：idempotencyKey 去重、event catalog validation（只允许预定义的 event\_type）
- 调用方：所有写 timeline 的 pipeline

**Policy Guard**（横切层，不是单独 writer）

- 输入：actor + storeId + requested action + current state + idempotency key
- 输出：allow / reject / needs\_review + reason
- 保证：tenant isolation、DNC、权限、幂等、状态机转换、schema validation、AI confidence gate
- 调用方：Task Orchestrator / Contact Writer / Timeline Writer，以及 future tool API
- 设计原则：AI 可以提出 action，但 Policy Guard 决定这个 action 能不能执行

**Layer 2: Processing Capability Modules — "调的统一"**

同一套处理能力要被多种入口复用。这些不是"多 writer 争抢同一张表"的问题，而是"同一个能力怎么被 normal flow / retry / onboard / reconciliation 调用"的问题。

**Call Analysis Module**（从现有代码提取 invocation contract）

- 核心问题：`analyze_call(callId, mode)` 怎么被调用？
- mode 至少包括：normal（正常 pipeline 流转）/ retry（员工点"重新分析"）/ backfill（新用户 onboard 扫描历史）/ reconciliation（修复遗漏）
- 需要回答：retry 是否清 `aiAnalysisCompletedAt`？重新分析后是否重新触发 contacts-analyzer？如何避免重复 AI cost 和并发 double run？onboard 扫描是重走 transcribe 还是只对已有 transcripts 跑 AI？
- 现状：入口主要是 S3 transcription event → SQS → ai-analysis-processor，已有 `tryClaimAnalysis()` 防并发。但还不是一个可被多入口调用的抽象模块

**Contact Analysis Reprocess Module**（现有 `contact-analyzer-reprocess` Lambda）

- 核心问题：已有 fan-out 到 contacts-analyzer queue 的能力，但只是 contact-level reprocess，不是 call-level retry
- 需要回答：reprocess 和正常 cron 走同一个路径吗？reprocess 后写入用同一套 Contact Writer / Task Orchestrator 吗？

**SMS Signal Module**（部分已有雏形，需要设计完整 contract）

- 核心问题：一条 SMS 到了以后，系统怎么判断它是不是有业务意义，以及要不要触发后续 workflow
- 现状雏形：message-processor 已有 exact "STOP" 检测（`stop-keyword.ts`）→ 设 DNC + 关 open tasks（`neon-repository.ts:299`）
- 需要设计的完整 flow：

```
inbound SMS
  → store message（message-processor 现有）
  → SMS Signal Module
      → exact STOP: code DNC（现有）
      → trivial (thanks/ok/emoji): no-op
      → meaningful: trigger contacts-analyzer
      → natural-language DNC: AI judgment + code DNC guard
```

- **关键设计原则**：SMS AI analysis 不应该成为一个新的 task writer。它应该只产出 signal（meaningful / DNC / booking intent / cancellation intent / complaint），然后交给 Contact Analyzer + Task Orchestrator + Contact Writer + Timeline Writer 去做真正的状态变化
- 需要回答：code filter 做多少（exact keyword match）vs AI 做多少（natural language intent）？debounce？

**Lead Downstream Invocation**（现有 lead-processor，检查 contract）

- 核心问题：lead 进来后触发的下游链（contacts UPSERT + task create + timeline）是否走 Layer 1 的共享模块
- 现状：lead-processor 直接 `db.batch()` 原子写 3 张表，没走 Task Orchestrator
- 需要回答：lead-processor 应该改成调 Task Orchestrator 创建 lead\_outreach task 吗？还是保持现有的 hardcoded 逻辑（因为它是纯确定性的，不经 AI）？

**两层的关系**：Layer 2 的 module 跑完之后，结果通过 Layer 1 的 module 写入。比如 Call Analysis Module 分析完一通电话 → 调 Contact Writer 更新 contacts → 调 Timeline Writer 写 timeline。SMS Signal Module 判断"有意义" → 触发 Contact Analysis → 结果再通过 Task Orchestrator 写 tasks。**如果只设计了 Layer 1 没设计 Layer 2，以后做 retry 时会发现 retry 绕过了统一写入。**

**Layer 2 的 Invocation Contract 清单**（design 必须覆盖，实现可以分阶段）：

| Contract                   | 解决什么                                                                                  |
| -------------------------- | ------------------------------------------------------------------------------------- |
| Call Analysis Invocation   | `analyze_call(callId, mode)` — normal / retry / backfill / reconciliation             |
| SMS Signal Invocation      | `evaluate_sms(messageId)` — STOP / trivial / meaningful / NL-DNC                      |
| Contact Analysis Reprocess | `reanalyze_contact(phone, storeId, reason)` — cron / per-call / on-demand / reprocess |
| Lead Downstream Invocation | `persist_lead_downstream(leadRow)` — 是否走 Task Orchestrator                            |

**Phase 1 执行策略**：第一刀先做 Layer 1（Task Orchestrator + Contact Writer + Timeline Writer）。但 **design 必须把 Layer 2 的 invocation contract 全部写进去**，否则后面做 retry/onboard/SMS AI trigger 时又会发现缺抽象。执行可以分阶段，设计不能只看一个角落。

### Step 3: 检查每条 pipeline 在统一架构下的职责

改造后每条 pipeline 的职责应该是：

```
Pipeline 专有职责:
  触发检测 → 数据获取 → AI 调用（如果需要）→ 解析 AI 输出

共享模块调用:
  → Contact Writer（更新 contacts）
  → Task Orchestrator（创建/更新/关闭 tasks）
  → Timeline Writer（记录审计事件）
```

需要逐条 pipeline 画出来，标清楚哪些是"专有"哪些是"共享"。

### Step 4: 定义 Code vs AI 的统一分界框架

建立一个适用于所有 pipeline 的通用判断标准：

| 判断维度         | Code 做 | AI 做  | Hybrid          |
| ------------ | ------ | ----- | --------------- |
| 需要确定性结果？     | ✅      |       |                 |
| 需要理解自然语言？    |        | ✅     |                 |
| 涉及状态转换？      | ✅（执行）  | ✅（提案） | AI 提案 + Code 执行 |
| 涉及金钱/权限/DNC？ | ✅      |       |                 |
| 需要分类/评分？     |        | ✅     | AI 分类 + Code 校验 |
| 需要生成文本？      |        | ✅     |                 |
| 有已知规则可以写死？   | ✅      |       |                 |

把这个框架应用到每条 pipeline 的每个决策点上，产出完整的职责矩阵。

### Step 5: 设计可复用操作

列出系统里应该被设计成"可以被任何入口调用"的操作。每个操作定义：谁能调、输入什么、做什么、返回什么。

初步清单（需要验证和补充）：

| 操作             | 触发方                                           | 做什么                                    |
| -------------- | --------------------------------------------- | -------------------------------------- |
| 分析单个联系人        | cron / per-call / on-demand / retry / onboard | 聚合数据 → AI 分析 → 写 contacts + tasks      |
| 重新处理一通电话       | 员工点 retry / reconciliation                    | 重新跑 transcribe → AI analysis → 写 calls |
| 新用户 onboard 扫描 | 新店铺接入                                         | 扫描过去 N 天的 calls/messages → batch 分析    |
| SMS 触发 AI 分析   | SMS 到达且内容有意义                                  | 触发 contacts-analyzer 重新分析该联系人          |
| 手动创建/关闭 task   | 员工 UI                                         | 调 Task Orchestrator                    |
| 批量重分析          | 管理员 / prompt 更新后                              | 对一批联系人重跑 contacts-analyzer             |

### Step 6: 验证 tool calling 兼容性

对 Step 5 的每个操作，验证：如果以后改成 tool calling 模式，这个操作能不能直接变成一个 tool？

```
方案 A（现在）:
  contacts-analyzer 代码里调 TaskOrchestrator.create(...)

方案 B（以后）:
  AI 调用 create_task tool → tool 实现里调 TaskOrchestrator.create(...)
```

如果答案是"能"，说明设计是 future-compatible 的。如果答案是"不能"，需要调整 Step 2 的模块设计。

### Step 7: 从 API / AI Agent 视角反推数据模型

这一步不是要现在做 agent，而是用 future callers 反推 Phase 1 的接口有没有设计对。一个 shared module 如果只能被当前 Lambda 调用，不能被 API / CLI / AI agent 复用，那它还不是稳定抽象。

需要逐个检查：

| 调用方                  | 需要看到什么                                                           | 反推问题                                            |
| -------------------- | ---------------------------------------------------------------- | ----------------------------------------------- |
| Human UI             | task 状态、progress、close result、timeline                           | UI 是否需要解析 JSONB 才知道发生了什么？                       |
| Backend API          | 结构化 task/contact/call/message 对象                                 | API contract 是否能表达 retry/backfill/reprocess？    |
| AI structured output | action proposal + typed params                                   | prompt output schema 是否等同于 module input schema？ |
| Future tool API      | `create_task` / `close_task` / `update_contact` / `analyze_call` | tool 是否只是 shared module 的薄包装？                   |
| Voice agent          | 低延迟、幂等、可拒绝、可审计                                                   | 每个 tool call 是否有 policy gate 和 idempotency key？ |
| Audit / replay       | 谁提议、谁执行、为什么拒绝                                                    | 是否能从 events 重放一次决策链？                            |

设计结论要明确区分：

- **Domain API**：给产品和后端用，比如 `TaskOrchestrator.closeTask()`
- **AI Action Schema**：给 prompt structured output 用，比如 `taskDecisions[].action = "close"`
- **Tool API**：给 future agent 用，比如 `close_task(taskId, result, evidence)`

三者可以共享同一套 schema / validation，但不一定暴露完全一样的字段。

***

## 思考方式

**"Code as Guardrail, AI as Judgment"** — 业界在用的 production pattern，多个名字：

- "AI Proposes, Deterministic Layer Disposes"
- "Pre-execution Approval Gate"
- "Graduated Autonomy"

核心：

| 职责     | 谁做                       | 为什么                             |
| ------ | ------------------------ | ------------------------------- |
| 状态机转换  | 代码                       | 必须可审计、可逆、exactly-once           |
| 权限检查   | 代码（Policy rules）         | 不能被 prompt injection 影响         |
| 幂等     | 代码（DB constraints、dedup） | LLM 不能保证确定性行为                   |
| 意图分类   | AI                       | 需要语义理解                          |
| 内容生成   | AI                       | 需要自然语言能力                        |
| 动作参数提取 | AI 提取 + 代码验证             | Hybrid：AI 提取，Zod/JSON Schema 校验 |

生产数据：使用 structured output（JSON Schema 约束 AI 输出格式）的 agent 动作成功率 95-99%，而解析非结构化文本的只有 70-85%。

**你的 task pipeline deliverable 里的 "AI proposes, code executes" 完全是这个 pattern。** 你已经在 contacts-analyzer 上做了（taskDecisions\[] = AI 提案 + 代码执行）。这次的任务是检查其他 pipeline 是不是也该统一到这个模式。

**"Schema-first" 对应到这里就是**：用 JSON Schema / Zod 定义 AI 能做什么操作、每个操作的参数是什么、哪些值是允许的。规定一堆 JSON 然后按照 JSON 的结构生成 prompt 和数据库结构。

***

## 前序调研发现的关键事实（设计时要知道的）

这些不是"要修的 bug 清单"，而是做 system design 时的**输入事实**——你需要知道系统实际上怎么运作的，才能设计对。**不要逐个解决这些小问题——做好顶层设计，这些问题会同时被解决。以面破题，不以点破题。**

**现有 AI signal chain 怎么触发 task：**

- Per-call AI stage 分析每通电话，产出 `fu=yes/no` (follow-up needed) + `out=success/attempted/...` (outcome)
- Contact analysis stage 消费 per-call 输出，做跨通话的 contact 级别分析，产出 `taskDecisions[]` (create/close/update)
- 这个 `fu=` 信号是 AI 对 AI 的传递——per-call AI 告诉 contact-level AI "要不要跟进"
- 注意：这里的 per-call / contact-level 是现有 AI 处理链，不要和上文的 Layer 1 Shared Mutation / Layer 2 Processing Capability 混用

**task 创建只有 3 条路径：**

1. lead-tracking 直接建（lead\_outreach，不经 AI）
2. contact analysis AI 建（taskDecisions\[].create）
3. 员工手动建（studio-api）

- 所有路径都只产 status='pending'。没有"一步建成 closed"的路径

**task 关闭有 2 条路径：**

1. AI auto-close（contact analysis 的 taskDecisions\[].close，必须引用已有 taskId）
2. 员工 manual-close（studio-api close.ts）

**contacts 表有多个 writer：**

- 每个 pipeline 各自实现 UPSERT 逻辑和 trust score
- name trust score（WEBHOOK=20 \< AI\_TRANSCRIPT=40 \< RC\_API=60 \< LEAD=80 \< STAFF=100）散落在不同 Lambda

**contact\_timeline 是统一审计线：**

- 16 种事件类型，覆盖 task 全生命周期
- 多态 actor（system / call\_analysis / contact\_analysis / staff / lead\_webhook）
- AI forensic 字段（aiRunId / aiModelUsed / aiConfidence）

**disposition 数据已经在 calls 表里：**

- `calls.callState` = human\_conversation / voicemail / no\_answer / busy\_signal / system\_error
- `calls.followUpNeeded` + `calls.followUpReasons` = AI 判断的 follow-up 信号
- 这些都是 per-call 的，不需要在其他表重复存

***

## 前序调研成果 — 分三类: Facts / Hypotheses / Calibration

前序 session 做了大量调研和代码分析。**这些是校准用的，不是替代你自己思考的答案。** 按可信度分成三类：

### Facts（已用代码/数据验证，可直接信）

| 事实                                                                          | 验证方式                           |
| --------------------------------------------------------------------------- | ------------------------------ |
| contacts 表有多条 pipeline 同时写入，trust score 逻辑各写各的                              | 读 5 个 Lambda 的 neon-repository |
| tasks 表有 3 个独立写入入口（lead-processor / contacts-analyzer / studio-api），去重策略不一样 | 读代码                            |
| contact\_timeline 有多个 writer，event\_type 各自定义                               | 读代码                            |
| storeId null 时各 pipeline 行为不一致（有的跳过、有的继续、有的 DDB fallback）                   | 读代码                            |
| contacts-analyzer prompt 已 1040 行，每次改一个地方可能影响一大片                            | 读 prompt-builder.ts            |
| call analysis 有 pre-triage 代码层（关键词匹配，节省 \~40% AI 调用）                        | 读 pre-triage.ts                |
| Task Pipeline Deliverable 的 "AI proposes, code executes" 已设计并讨论通过           | team 会议 2026-05-30             |

### Hypotheses（前序分析的推断，需要验证后才能采纳）

| 假设                                                    | 来源            | 需要验证什么                                    |
| ----------------------------------------------------- | ------------- | ----------------------------------------- |
| 所有 contacts 写入应该统一成 Contact Writer module             | pipeline 全景分析 | 会不会 over-engineering？简单 writer 是否值得走统一模块？ |
| 所有 timeline 写入应该统一成 Timeline Writer module            | pipeline 全景分析 | event catalog 是否需要预定义所有 event\_type？      |
| contact 的 `actionNeeded` 和 task 的 `actionNeeded` 概念混淆 | 代码对比          | 具体怎么混的？是否需要去掉一个？                          |
| 方案 A 做好了方案 B 自然就有                                     | 架构推理          | 是否真的无缝？有没有 A 的设计会卡住 B？                    |

### Calibration（业界校准 — 2026-05-31 web search）

> **web search 是校准用的，不替代系统设计结论。最终判断必须回到 retaintive 自己的架构和成本。**

| 校准点                                 | 业界怎么做                                                                                              | 来源                              |
| ----------------------------------- | -------------------------------------------------------------------------------------------------- | ------------------------------- |
| 大厂 AI SaaS 架构                       | 每个用例一个专门 agent + 共享基础设施，不是一个万能 agent（Salesforce Agentforce 12,000+ 客户、HubSpot Breeze、Intercom Fin） | Salesforce · HubSpot · Intercom |
| LangGraph / Temporal 生产落地           | 真的有人用（LinkedIn/Uber/Replit），但 prototype→production 落差巨大，40% 项目被取消                                  | Gartner · LangChain blog        |
| "Code as Guardrail, AI as Judgment" | 成熟 pattern，structured output 成功率 95-99% vs 非结构化 70-85%                                             | 多源                              |
| 2-3 人小团队                            | Level 2-3 agent realistic（单 agent + tool calling），multi-agent overkill                             | 多源                              |
| Voice agent 后端需求                    | HTTP endpoint 接收 tool call POST + 快速响应（\<2s） + 幂等。Vapi/Retell/Bland 三家核心模式一样                       | Vapi · Retell · Bland docs      |
| Tool calling 生产模式                   | OpenRouter 标准化 tool calling 接口，poly-model routing（Opus 规划 + Sonnet 执行）                             | OpenRouter docs                 |
| 小团队推荐 pattern                       | Router Agent：便宜模型（Haiku，1/10 成本）做意图分类 → 转给专门逻辑                                                     | 多源                              |

***

## 已知的系统现状

不展开，只列要点。详细 trace 见 [Pipeline 全景分析](/system-design/pipeline-architecture-analysis.md)。

### 5 条活跃 pipeline

| Pipeline         | 触发                                 | AI？                                | 写哪些**共享**表                         |
| ---------------- | ---------------------------------- | ---------------------------------- | ---------------------------------- |
| Call Analysis    | RC Webhook → SQS                   | 3 prompt（triage/classify/coaching） | contacts, contact\_timeline        |
| Contact Analysis | dailyBatchQueue + cron + on-demand | 1 prompt（contact-analyzer）         | contacts, tasks, contact\_timeline |
| SMS              | RC Webhook → SQS                   | Code-only exact STOP；暂无 SMS AI     | contacts, contact\_timeline        |
| Lead             | IMAP → EventBridge                 | 无                                  | contacts, tasks, contact\_timeline |
| Analytics        | EventBridge cron                   | 读 call-analysis                    | 不写共享表                              |

> 每条 pipeline 也写自己的**专有表**（calls、messages、leads、DDB call-analysis），那些是 pipeline 内部的事。这次关注的是**多条 pipeline 同时写的表**怎么统一。

### 多条 pipeline 同时写的表（需要统一的部分）

| 表                     | 哪些 pipeline 写                                              | 问题                                      |
| --------------------- | ---------------------------------------------------------- | --------------------------------------- |
| **contacts**          | Call Analysis + Contact Analysis + SMS + Lead + studio-api | name trust score 逻辑散落在每个 pipeline 里各写各的 |
| **contact\_timeline** | Call Analysis + Contact Analysis + SMS + Lead + studio-api | event\_type 各自定义，没有统一的 event catalog    |
| **tasks**             | Contact Analysis + Lead + studio-api                       | 去重策略、closeResult、dueAt 计算全不一样           |

***

## Max 在对话中表达的方向

**这些是他的想法和倾向，不是 finalized spec：**

- 从 task pipeline 设计出发，发现所有 pipeline 需要统一设计
- "规定一堆 JSON 然后按照 JSON 的结构生成 prompt 和数据库结构比较靠谱，别的感觉想的很好都没落地过"
- 认同 "Code as Guardrail, AI as Judgment" pattern；该原则已写入产品设计准则
- 想有一个**最终构思（终态方向）**，然后看现有的怎么靠拢
- 先做一个版本 buying more time，设计要符合 long-term scalability
- 现在没有 voice agent，但以后可能有，架构要兼容
- **先把基础的东西（task、lead、contact、call）弄清楚**，再拓展高级功能
- 关心可复用性：call retry、新用户 onboard 扫描、批量重处理 这些操作能不能设计成模块
- 不确定现有设计对不对，需要做系统性检查（不是说设计错了，是需要验证）
- ai-integration-roadmap.md 是很久以前写的，可能 outdated，只作过去参考

***

## 材料清单

> **只有代码是 source of truth。其余都是不同人在不同时间的想法，不一定对。**

### A. 代码 (✅ Source of Truth)

| 文件                                                                              | 看什么                                                                     |
| ------------------------------------------------------------------------------- | ----------------------------------------------------------------------- |
| `callytics-common/src/db/schema/` 全部文件                                          | 所有 business object 的 schema 定义、枚举、约束                                    |
| `callytics-infrastructure/lambda/*/src/`                                        | 每条 pipeline 的实际实现 — trigger、读写、AI 调用                                    |
| `callytics-infrastructure/lambda/contacts-analyzer/src/core/prompt-builder.ts`  | contacts-analyzer prompt（1040 行）— SECTION 4 是 task 决策                   |
| `callytics-infrastructure/lambda/ai-analysis-processor/src/core/stages/`        | per-call AI 3-stage pipeline（pre-triage → triage → classify → coaching） |
| `callytics-infrastructure/lambda/lead-processor/src/core/persist-downstream.ts` | lead 下游写入（contacts + tasks + timeline 原子 batch）                         |
| `callytics-infrastructure/lambda/message-processor/src/core/`                   | SMS 处理逻辑                                                                |
| `studio-website-monorepo/apps/api/src/routes/v3/`                               | studio-api 的 task/contact 读写 endpoint                                   |

### B. 设计文档（已确认的方向）

| 文件                                                                                                              | 看什么                                                           |
| --------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------- |
| [Task Pipeline Deliverable (Codex)](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) | Task Orchestrator 设计、4-Layer Object Model、Code vs Prompt 职责矩阵 |
| [Task Pipeline 设计 Brief](/product-design/v2/tasks-feature/design/research/task-pipeline-design-brief.md)        | 设计方法论和 Max 的方向性想法                                             |
| [产品设计准则](/product-design/v2.md)                                                                                 | Schema-first → API-first → State-machine-first → Prompt-last  |
| [Pipeline 全景分析](/system-design/pipeline-architecture-analysis.md)                                               | 5 条 pipeline 的代码级 trace + gap analysis                        |

### C. 参考材料（⚠️ 不一定对，当参考不当结论）

| 文件                                                             | 来源    | 说明                                                              |
| -------------------------------------------------------------- | ----- | --------------------------------------------------------------- |
| [AI 集成路线图](../../ai/product/ai-integration-roadmap.md)         | 早期设计  | agentic loop + 6 gym tools + 5 级金字塔。**Max 说可能 outdated，只作过去参考** |
| [Voice Agent 可行性](../../ai/product/voice-agent-feasibility.md) | 调研报告  | 接 Vapi/Retell 的技术可行性。voice agent 不是这次的目标，但架构要兼容                 |
| [Backend 架构 Pattern](/system-design/backend-patterns.md)       | 架构文档  | Layered vs Hexagonal 选型标准                                       |
| [Task 枚举清单](tasks-feature/task-enums.md)                       | V2 文档 | closeResult 18 值、typeCategory 9 值 — 代码 SoT                      |

### D. Neon Test 环境 (✅ 可查真实数据)

- Project: `callytics-test` (restless-boat-70724564, us-west-2)
- 用 Neon MCP 跑 SELECT 验证假设，不写入

### E. AWS (✅ 已登录)

- Account: 237206024479, AdminAccess role
- 可看 Lambda / CloudWatch，用于验证 prompt 实际输出和 pipeline 行为

***

## 补充维度（设计时别漏）

### 多门店 / store 级隔离

所有共享模块必须以 `store_id` 作为隔离键。Task Orchestrator 的幂等约束 `uq_tasks_pending_contact_category` 按 `(contactPhone, storeId, typeCategory)` 做。Contact Writer 的 UPSERT 以 `(phone, storeId)` 为 PK。同一个电话号码在不同门店是独立的 contact。

### 前台一天的实际操作流程（场景 robust 的基础）

设计必须对得上前台真实的工作方式，不能只看代码：

1. **早上开工**：打开系统看待办清单——"今天有哪些客人要跟进"
2. **白天逐个处理**：挑一件 → 打电话/发短信 → 当场办成 / 没当场办成回头来 / 没打通再试 / DNC
3. **当场来的**：客户直接打进来当场预约 — 没经过待办但算业绩
4. **一天结束**：店长看战绩（完成几件、几个签约/挽留）

### 未来 AI agent 的操作流程（架构要兼容但不是这次的目标）

- AI agent 自动外呼 → 自动判断对话结果 → 调 Task Orchestrator 关闭/更新 task
- 需要 `executorType`（human / ai\_agent / system）区分"谁做的"
- 低置信度的 AI 决策需要 human review
- 共享模块的设计要让 AI agent 和人类走**同一个接口**

### SMS 处理方式

**现状**：message-processor 存储 `messages`，并对 exact inbound STOP 做 code-only DNC + task close cascade；它还不触发 AI / contacts-analyzer 来处理 meaningful SMS。

**设计时要思考**：

- "STOP" 这种关键词是否应该在 code filter 层立即处理（不等 AI）？
- 有意义的 SMS 回复是否应该触发 contacts-analyzer？
- 怎么区分"有意义"vs"无价值"（thanks / ok / emoji / 自动回复）？
- SMS 还没有完全想清楚，先 focus 电话。但 pipeline 设计要给 SMS 留位置

### 可复用操作的场景

以下场景直接验证共享模块的设计是否合理：

| 场景                        | 触发方               | 需要调哪些模块                                                                   |
| ------------------------- | ----------------- | ------------------------------------------------------------------------- |
| 员工说"我想 retry 这通电话的 AI 分析" | UI 按钮             | 重跑 AI analysis → Contact Writer → Task Orchestrator                       |
| 新用户 onboard，想扫描过去两周的通话    | onboard 流程        | batch 调 transcribe + AI analysis + contacts-analyzer                      |
| prompt 更新了，想对所有联系人重分析     | 管理员               | batch 调 contacts-analyzer                                                 |
| SMS 收到 "STOP"             | message-processor | Code filter → Contact Writer（标 DNC）→ Task Orchestrator（关所有 pending tasks） |

如果这些场景都能通过"调共享模块"完成，而不是写新的专门 pipeline，说明设计是对的。

### API / CLI / AI agent 可读性

设计时要考虑未来 API 的形状。以后可能有第三方系统、CLI、客户的 AI agent 通过 OAuth/API 读写 task 和 contact。

设计时要显式回答：

- 外部调用方如何区分 task 的生命周期状态 / 执行进展 / 最终业务结果？
- API 返回应该是结构化对象，不是要求调用方解析 timeline JSONB 或 prompt 文本
- 共享模块的接口设计要能直接变成 tool calling 的 tool

***

## Deep Gaps / Open Design Questions

这些不是立即要写代码的 todo，而是 design review 必须讨论清楚的风险点。很多后续 bug 会从这些边界没定义开始。

| Gap                                 | 为什么重要                                                                                                     | 设计时要回答                                                                                            |
| ----------------------------------- | --------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------- |
| Processing intent model             | `normal` / `retry` / `backfill` / `reconciliation` / `on-demand` / `task_close` 如果只是散落的 string，后面会变成不可控分支 | 是否需要统一 `processingIntent` enum？每个 intent 能不能触发 AI、能不能写 shared state、能不能覆盖旧结果？                     |
| Retry semantics                     | retry 一通电话可能影响 `calls`、`contacts`、`tasks`、timeline 和 AI cost                                              | retry 是追加 attempt 还是覆盖字段？是否清 `aiAnalysisCompletedAt` / `s3AnalysisPath`？是否重新触发 contacts-analyzer？ |
| Source of truth vs projection       | `task_progress_events`、`contact_timeline`、contact action summary 可能重复表达同一件事                               | 哪张表是 source of truth？哪些只是 projection？projection 失败能不能重建？                                          |
| `storeId` unresolved policy         | 现在不同 pipeline 对 `storeId = null` 行为不一致                                                                    | 是 skip、quarantine、fallback、retry，还是允许某些 source 先落库后补 store？                                       |
| Idempotency taxonomy                | 每个 Lambda 自己做 dedup 会导致跨入口重复写                                                                             | idempotency key 按 event、message、call、contact-run、task-action 还是 external request 生成？              |
| Human override vs AI race           | 员工 close task 的同时 AI auto-close / update task，可能覆盖人的判断                                                    | staff action 是否永远优先？AI proposal 发现 task 已被人工修改时是 reject、append note 还是 reopen？                    |
| Payload catalog / schema versioning | `contact_timeline.newValue`、prompt output、future tool schema 都会演进                                         | 每个 event payload 是否要 version？旧 prompt output 如何兼容新 module input？                                  |
| Permission / RBAC                   | 同一个 shared module 未来会被 staff、system、AI agent、external API 调用                                              | actor 能力是否分级？AI agent 是否只能 propose，不能直接 execute 高风险 action？                                       |
| Cost budget / debounce              | SMS meaningful reply、tool calling、batch reprocess 都可能放大 AI 调用量                                            | 是否需要 per-store budget、cooldown、debounce window、batch cap？                                         |
| Rollout / migration                 | 现有 enum 和历史数据已经存在，不能只设计新世界                                                                                | Phase 1 怎么兼容旧 `contact_timeline.event_type`、tasks status、contacts action summary？                 |
| Testing strategy                    | 共享模块一旦错，会影响所有 pipeline                                                                                    | 是否需要 contract tests、golden prompt output tests、replay tests、idempotency/race tests？               |
| Analytics definitions               | task 完成、call 成功、contact converted、business win 可能不是同一层概念                                                  | dashboard 应该按 attempts、objectives completed、closed won、retained member 里的哪个口径算？                   |

**需要特别小心的一点**：不要把所有问题都塞进 Task Orchestrator。Task Orchestrator 只管 task lifecycle 和 task progress；Contact Writer 管 contact aggregate state；Timeline Writer 管 audit/projection；Call/SMS Module 管 processing invocation。职责一旦混了，Phase 1 会重新变成一个大而全的 pipeline。

***

## 设计时要考虑到（high-level，不需要写具体实现）

- **AI 判断失败时的 fallback** — prompt 返回无法解析 / Lambda timeout / Zod 校验不过，代码怎么兜底？
- **contacts-analyzer 有 6 个触发路径** — cron / on-demand / P1 followUp / Task close / T5 showed / T6 trialed。设计要覆盖所有入口
- **并发 / 幂等 / race condition** — daily batch 和 per-call 可能同时分析同一个 contact；lead-tracking 建 task 的同时 AI 也想建；员工 close 的同时 AI auto-close
- **成本控制** — 方案 A 的成本是可控的；方案 B 多轮调用的成本需要 budget cap
- **业界做法是参考不是圣经** — 2026 年的技术限制和 2024 年不同（LLM 成本降 10-50x）。不要因为"Salesforce 这么做"就照搬
- **产品设计准则更新** — "Code as Guardrail, AI as Judgment" 这个 pattern 已写进 `product-design/index.md`，后续 AI workflow 默认按这个原则设计
- **设计边界：如果需要改，什么都可以改** — 所有 pipeline / schema / prompt / 前端都在 scope。这是 test 环境

***

## 约束

- **先理解全局再看局部** — 先把所有 business object 的关系搞清楚
- **先设计再写代码** — 画图、确认逻辑、Max review 之后才动手
- **Code as source of truth** — 每个结论先读代码验证，不信文档 narrative
- **不要以点破题，以面破题** — 不要逐个修小问题，设计好统一抽象后小问题自然消解
- **方案 A 优先** — 先做好统一抽象（one-time AI call），tool calling 留给交互式场景
- **检查而不是重写** — 检查现有设计对不对，有问题的修正，没问题的保留
- **从正确性出发** — 不要因为现在是这样就继续走，test 环境什么都可以改
- **保持 high-level** — 这是 system design，不是 code review
- **不自动 push/PR**
