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

# Unified Pipeline 统一设计

> **Historical architecture snapshot（2026-07-23 校准）**：本文保留 shared mutation module、Policy Guard、Timeline Writer、0 / 1 / N decisions 和 `create_closed` 等 Task V2 架构依据。`task_progress_events`、global `closeResult` 扩展和 Lead relay 等旧 proposal 不能直接实施；其中 `create_closed` 是当前必须保持的一等产品语义。当前路线见 [Task V2+ 工程审计与实施基线](/product-design/v2/tasks-feature/task-v2-plus-engineering-audit.md)；实施 pipeline 变更前必须重新核查 live code。
>
> **日期**: 2026-06-01
> **状态**: 合并 Opus + Codex 两份独立设计后的统一版本，待 Max review
> **环境前提**: 本文以当时 test / new architecture 调研为主；production Task V2 行为需要单独 audit，不默认与 test 一致。
> **设计准则**: [产品设计准则](/product-design/index.md)（Schema-first → API-first → State-machine-first → Prompt-last / Code as Guardrail, AI as Judgment）
> **历史前置设计**: [Task Pipeline Deliverable (Codex)](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) — superseded Task contract，仅用于追溯 Task Orchestrator 的演进依据
>
> **原始材料**（保留在本目录，不删除）:
>
> - [Target State (Opus)](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md) — current state inventory + target state + Phase 1
> - [Target State (Codex)](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md) — 理想终态架构
> - [Phase 1 (Codex)](/product-design/v2/unified-pipeline/original/unified-pipeline-phase-1-codex.md) — Phase 1 落地方案
> - [Research 材料](research/) — canonical brief、现状调研、checklist、appendix

***

## 1. Historical proposed Target State

所有写入方通过同一套共享函数改变系统状态，不直接碰共享表。

### 什么是 module（共享函数）

现在系统的现状：很多写入逻辑其实已经收口了 —— name trust scoring 已集中到 `NAME_TRUST` 共享常量（`name-trust.ts`），dueAt 计算已抽出 `computeDueAt`（3 个 Lambda 共用），timeline 写入已统一走 `buildTimelineValues` / `buildContactTimelineInsertSQL`，DNC cascade 已共享（`dnc-cascade.ts`）。剩下两个真实缺口：(1) task mutation（create/close/update）还没有统一的 Task Orchestrator，状态机语义散在 contacts-analyzer / lead-processor / studio-api 的各自实现里；(2) `task_progress_events` 表不存在，导致 `no_answer` / `left_voicemail` 这类进展被错塞进 closeResult。本设计收口这两个缺口，不是从零统一（大部分共享 helper 已就位）。

Module 就是把这些散落在多个地方的写入逻辑，抽到一个共享的 TypeScript 文件里（住在 `callytics-common/src/domain/`）。不是新 Lambda，不是新 SQS，不是新进程。以 task mutation 为例，改前改后的区别：

```text
改前：task mutation 散在三处各自实现
      （去重和 timeline 已共享，但 create/close/update 的状态机逻辑各写各的）
      lead-processor / contacts-analyzer / studio-api → 各自的 task 写入

改后：三条路都调 applyTaskAction()
      task 状态机 + create_closed / record_progress 等语义统一
      lead-processor / contacts-analyzer / studio-api → applyTaskAction() → tasks 表
```

每个 caller 传不同的 action + payload（create / create\_closed / close / record\_progress ...），统一的是规则（去重、状态机、timeline 审计），不统一的是内容。

**Terminology note**: 本文里的 `Task Orchestrator` 只指 task domain 的共享函数 / module，负责 `create` / `create_closed` / `close` / `update` / `record_progress` / `reopen` 这些 task mutation。它不是 workflow engine，不是新的 Lambda，也不是整个 unified pipeline 的总调度器。下文如果说 “Orchestrator”，都应该理解为 `Task Orchestrator`；为避免混淆，正文尽量写全名。

### 架构总览

> 只画写入路径。

```mermaid
%%{init: {'theme': 'base', 'themeVariables': {'primaryColor': '#e8eaf6', 'lineColor': '#666'}}}%%
flowchart LR
    subgraph entrypoints["写入方"]
        CallPipeline["通话分析"]
        ContactAnalysis["联系人分析"]
        SmsPipeline["短信处理"]
        LeadPipeline["Lead 处理"]
        StaffUI["员工操作"]
        FutureAgent["Future AI Agent"]
    end

    subgraph judgment["判断层"]
        CodeRules["代码规则<br/>exact STOP / 去重 / 过滤"]
        AIProposal["AI 结构化提案<br/>taskDecisions / contact 分析"]
    end

    subgraph mutation["共享写入层"]
        Guard["Policy Guard<br/>allow / reject / needs_review"]
        TaskOrc["Task Orchestrator<br/>create / create_closed / close /<br/>update / record_progress / reopen"]
        ContWriter["Contact Writer<br/>identity / DNC / trust"]
        TlWriter["Timeline Writer<br/>event catalog + Zod schema"]
    end

    subgraph source_records["源记录"]
        Calls[("calls")]
        Messages[("messages")]
        Leads[("leads")]
    end

    subgraph shared_tables["共享表"]
        Tasks[("tasks")]
        Progress[("task_progress_events")]
        Contacts[("contacts")]
        Timeline[("contact_timeline")]
    end

    CallPipeline --> Calls
    SmsPipeline --> Messages
    LeadPipeline --> Leads

    CallPipeline --> CodeRules
    CallPipeline -.->|"followUpNeeded<br/>(reference, not trigger)"| ContactAnalysis
    ContactAnalysis --> AIProposal
    SmsPipeline --> CodeRules
    LeadPipeline --> CodeRules
    StaffUI --> Guard
    FutureAgent --> Guard

    CodeRules --> Guard
    AIProposal --> Guard

    Guard --> TaskOrc
    Guard --> ContWriter

    TaskOrc --> Tasks
    TaskOrc --> Progress
    TaskOrc -->|"task 事件"| TlWriter
    ContWriter --> Contacts
    ContWriter -->|"contact 事件"| TlWriter

    CallPipeline -->|"source event"| TlWriter
    SmsPipeline -->|"source event"| TlWriter
    LeadPipeline -->|"source event"| TlWriter

    TlWriter --> Timeline

    style mutation fill:#dbeafe,stroke:#2563eb
    style Guard fill:#fef3c7,stroke:#d97706
```

### 当前写入全景图（现状）

跟上面的 target state 对比：现在每个写入方**物理上直接执行自己的 SQL** 碰共享表，没有统一入口函数。

> 注意区分"物理路径"和"逻辑共享"：下图画的是物理写入路径（谁的 SQL 写哪张表）。逻辑层面 trust scoring（`NAME_TRUST`）、timeline（`buildTimelineValues`）、dueAt（`computeDueAt`）、DNC cascade 其实**已经抽成共享函数**——但每个 caller 仍各自调用并自己执行 SQL（这就是「弱版本」：共享函数、各 caller 自己写）。Target State 的统一不是改成"物理上只有一个地方能写"，而是让 task mutation 也经过共享的 Task Orchestrator。

```mermaid
%%{init: {'theme': 'base', 'themeVariables': {'primaryColor': '#e8eaf6', 'lineColor': '#666'}}}%%
flowchart LR
    subgraph writers["写入方"]
        P1["通话分析<br/>transcribe + ai-analysis"]
        P2["联系人分析<br/>contacts-analyzer"]
        P3["短信处理<br/>message-processor"]
        P4["Lead 处理<br/>lead-processor"]
        P5["员工操作<br/>studio-api"]
        DNC["DNC cascade<br/>共享代码"]
    end

    subgraph shared["共享表"]
        C[("contacts")]
        T[("tasks")]
        TL[("contact_timeline")]
    end

    P1 -->|"UPSERT trust=60"| C
    P1 -->|"INSERT"| TL
    P2 -->|"UPDATE 18 AI 字段 trust=40"| C
    P2 -->|"CREATE/CLOSE/UPDATE"| T
    P2 -->|"INSERT"| TL
    P3 -->|"UPSERT trust=60"| C
    P3 -->|"INSERT"| TL
    P4 -->|"UPSERT 无 trust"| C
    P4 -->|"INSERT lead_outreach"| T
    P4 -->|"INSERT"| TL
    P5 -->|"UPDATE trust=100"| C
    P5 -->|"4 种直接写"| T
    P5 -->|"reopen/postpone/note INSERT<br/>close 缺失"| TL
    DNC -->|"DNC=true"| C
    DNC -->|"关闭 open tasks"| T
    DNC -->|"INSERT"| TL
```

### Prompt Pipeline

先把 prompt 数量说清楚，避免把 code stage、current prompt、future read-only prompt 混在一起：

- **Current live code / Phase 1 write path**: 6 个 active decision surface = 1 个纯代码 stage + 5 个 AI prompt。5 个 AI prompt 是 `Triage` / `Classify` / conditional `Verify` / `Coaching` / `Contact Analyzer`
- **Task Playbook**: planned read-only prompt，来自 Task Pipeline Deliverable。它读 task + contact 生成员工话术，不产出 DB mutation，不是 Phase 1 写入路径的前置条件
- **Target State**: `Contact Analyzer` 拆成 `Contact Profile Prompt` + `Task Decision Prompt`（**拆 prompt 文本/组装，仍是一次 LLM 调用、一个 JSON 输出 —— 见下方「拆 prompt 的方式」**）。如果加上 `Task Playbook`，target state 是 1 个纯代码 stage + 7 个 AI prompt surface

只有 `Contact Analyzer` 当前会产出 task/contact business mutation proposal。Phase 1 不拆它，只改下游写入方式；Target State 再把它拆成 `Contact Profile Prompt` 和 `Task Decision Prompt`（拆文本，不拆调用）。

#### Prompt / Decision Surface Inventory

| # | Surface              | Type                                  | 做什么                                                                                                                                                                                                                                                                          | 写哪张表                                      | Phase 1 影响                                                                                                                                                                            |
| - | -------------------- | ------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| ① | Pre-triage           | Code stage                            | 纯代码。用关键词匹配过滤掉 \~40% 不值得分析的通话（自动语音、拨号音等），省掉后面 AI 调用的成本                                                                                                                                                                                                                        | 不写表                                       | 不动                                                                                                                                                                                    |
| ② | Triage               | AI prompt                             | AI 只看通话前 500 个字符，快速判断这通电话值不值得做完整分析。不值得的直接跳过后面所有 stage                                                                                                                                                                                                                        | 不写表                                       | 不动                                                                                                                                                                                    |
| ③ | Classify             | AI prompt                             | AI 读完整 transcript，分析这通电话的分类（sales / cancellation / billing ...）、业务结果（成交 / 未成交 / 取消 ...）、是否需要跟进。产出 33 个字段写入 `calls` 表。其中 `followUpNeeded` 是给 Contact Analyzer 的 signal                                                                                                        | `calls`                                   | 不动                                                                                                                                                                                    |
| ④ | Verify               | AI prompt（conditional / observe-only） | 当 Classify 的 subcategory 落在容易混淆的边界区（如 cancellation vs billing），AI 独立重新判断一次。当前 observe-only（只记 log 不改数据），积累准确率数据后再决定是否启用                                                                                                                                                      | `calls`（observe-only）                     | 不动                                                                                                                                                                                    |
| ⑤ | Coaching             | AI prompt                             | 对真人通话（>30 秒），AI 生成给前台的辅导建议：哪句话说得好、哪句可以改进、具体怎么说更好。写入 `calls` 表的 `practical_coaching` 字段                                                                                                                                                                                       | `calls`                                   | 不动                                                                                                                                                                                    |
| ⑥ | **Contact Analyzer** | AI prompt                             | 系统里当前唯一触发 task/contact 业务变更 proposal 的 prompt。由 SQS 触发（每通电话分析完 + 每日定时 cron），聚合一个客户的所有通话、短信、lead、已有 task 的历史，AI 做跨互动判断：该建 task 吗？该关 task 吗？客户状态变了吗？当前输出 `taskDecisions[]`（create / close / update）+ 18 个 contact 分析字段。`followUpNeeded` 是它读的 signal 之一（输入证据，不是触发条件，也不该直接当结论沿用） | `contacts` + `tasks` + `contact_timeline` | **核心是改下游 + contract**：`TaskDecision` schema 加 `record_progress` / `create_closed`，执行逻辑从自己拼 SQL 改成调 Task Orchestrator / Contact Writer。拆 prompt 与此解耦、可独立做（见下方「Phase 1」与「拆 prompt 的方式」） |
| ⑦ | Task Playbook        | AI prompt（planned / read-only）        | 给前台用的执行指南。读一个 task 的上下文（类型、优先级、客户历史、上次通话内容），生成具体话术：怎么开场、能提供什么选项、哪些话不能说、什么时候可以关这个 task。不写任何表                                                                                                                                                                                  | 不写表（纯读取）                                  | 不阻塞 Phase 1 写入层；可以跟 UI 一起单独落地                                                                                                                                                         |

Task 相关的 AI 不应该混成一个概念：

- **Task Decision Prompt**: 决定是否 create / create\_closed / close / update / record\_progress。它产出 mutation proposal，必须走 Policy Guard + Task Orchestrator。Phase 1 先继续嵌在 Contact Analyzer 里；Target State 再拆成独立 prompt stage
- **Task Playbook Prompt**: task 已经存在以后，生成 staff-facing 话术和执行建议。它只读 task/contact/call context，不写 DB，不应该绕过 Task Orchestrator

#### 现在（Before）

```text
┌───────────────────────────────────────────────────────────┐
│  Call Analysis (ai-analysis-processor)                     │
│                                                           │
│  ① Pre-triage → ② Triage → ③ Classify → ④ Verify        │
│  → ⑤ Coaching                                            │
│                                                           │
│  只写 calls 表（33 字段 + followUpNeeded + coaching）      │
└──────────────────────────┬────────────────────────────────┘
                           │ SQS 触发（每通电话都触发）
                           ▼
┌───────────────────────────────────────────────────────────┐
│  ⑥ Contact Analyzer（一个 1000 行的大 prompt）             │
│                                                           │
│  读：calls + messages + leads + open tasks                 │
│                                                           │
│  一次 AI 调用同时输出：                                    │
│    ├── 18 个 contact 分析字段（画像）                      │
│    └── taskDecisions[]（create / close / update）          │
│                                                           │
│  代码自己拼 SQL 写库（部分逻辑已共享，部分还散）：          │
│    contacts 表 ← trust scoring 已用共享 NAME_TRUST 常量    │
│    tasks 表    ← 去重靠 DB 约束已一致，但 create/close/    │
│                  update 状态机语义仍各写各的（无 Orchestrator）│
│    timeline    ← 已统一走 buildTimelineValues             │
└───────────────────────────────────────────────────────────┘

问题：
• task 判断和 contact 画像耦合在一个大 prompt 里（prompt 太长难维护）
• task mutation 缺统一 Task Orchestrator —— create/close/update 状态机散在 contacts-analyzer / lead-processor / studio-api（注：trust / timeline / dueAt / DNC cascade 已共享，不在此列）
• TaskDecision schema 缺 create_closed / record_progress → 当场成交和多次跟进进度无法被正确记录
• 没打通 → 关旧 task + 建新 task → 打 5 次 = 5 个 task
• `task_progress_events` 表不存在 → no_answer / left_voicemail 被错塞进 closeResult
```

#### Phase 1：核心是写入层；拆 prompt 并行外包

Phase 1 的核心改动是**写入层**（Task Orchestrator / Contact Writer / Timeline Writer / Policy Guard + `task_progress_events`）。

拆 prompt（画像段 / task 决策段）**已经并行外包给独立 engineer 做**，跟写入层改造解耦 —— 拆完直接拿来接入即可。两者唯一的接触点是 `TaskDecision` schema 要加 `record_progress` / `create_closed` 等 6 个 action 命名（见下方「Phase 1 实施步骤」），这个 schema 改动是 Phase 1 自己的 deliverable，**不管 prompt 是否已拆都要做**。

下图展示写入层改造（prompt 是否已拆不影响这张图的写入路径）：Contact Analyzer 一次 AI 调用输出 contact 字段 + `taskDecisions[]`，下游从自己拼 SQL 改成调共享写入层。

```text
┌───────────────────────────────────────────────────────────┐
│  Call Analysis (ai-analysis-processor)                     │
│                                                           │
│  ① Pre-triage → ② Triage → ③ Classify → ④ Verify        │
│  → ⑤ Coaching                                            │
│                                                           │
│  只写 calls 表（33 字段 + followUpNeeded + coaching）      │
└──────────────────────────┬────────────────────────────────┘
                           │ SQS 触发（每通电话都触发）
                           ▼
┌───────────────────────────────────────────────────────────┐
│  ⑥ Contact Analyzer（一次 AI 调用，拆不拆 prompt 都一样）  │
│                                                           │
│  读：calls + messages + leads + open tasks                 │
│                                                           │
│  一次 AI 调用输出：                                        │
│    ├── 18 个 contact 分析字段（画像）                      │
│    └── taskDecisions[]                                    │
│         ★ 新增 record_progress + create_closed             │
└──────────────────────────┬────────────────────────────────┘
                           │ ★ Phase 1 改动：不再自己拼 SQL
                           ▼
┌───────────────────────────────────────────────────────────┐
│  ╔═══════════════════════════════════════════════════════╗ │
│  ║  共享写入层（Phase 1 新建）                           ║ │
│  ║                                                       ║ │
│  ║  Policy Guard                                         ║ │
│  ║       │                                               ║ │
│  ║       ├──→ Contact Writer ──→ contacts 表             ║ │
│  ║       │                                               ║ │
│  ║       └──→ Task Orchestrator ──→ tasks 表             ║ │
│  ║            ──→ task_progress_events 表                 ║ │
│  ║                                                       ║ │
│  ║  Timeline Writer ──→ contact_timeline 表              ║ │
│  ╚═══════════════════════════════════════════════════════╝ │
└──────────────────────────┬────────────────────────────────┘
                           │
                           ▼
┌───────────────────────────────────────────────────────────┐
│  ⑦ Task Playbook（只读，不写表）                          │
│  读 task + contact → 生成员工话术                          │
└───────────────────────────────────────────────────────────┘
```

#### 历史方案为何不直接一步到当时的 Target State

这里分成 Current / Phase 1 / Target State，不是因为 Target State 不对，而是因为这两类改动风险不一样：

| 层级           | 主要改什么                                                                                            | 是否改变 AI 判断方式                                                                                                                                                | 主要解决什么                                             |
| ------------ | ------------------------------------------------------------------------------------------------ | ----------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------- |
| Current      | 记录现在代码真实长什么样                                                                                     | 不适用                                                                                                                                                         | 给后面的设计一个现实基线                                       |
| Phase 1      | 写入层：Task Orchestrator / Contact Writer / Timeline Writer / Policy Guard / `task_progress_events` | **基本不变**。schema 加 6 个 action（`create_open` / `create_closed` / `close` / `update` / `record_progress` / `reopen`）。拆 prompt 由独立 engineer 并行交付，Phase 1 拿来直接接入 | 先把 DB mutation、幂等、状态机、timeline 审计统一起来              |
| Target State | prompt 重构：拆 Contact Analyzer 已外包给独立 engineer（拆文本/组装，仍一次 LLM 调用），再接 Task Playbook Prompt          | **判断方式不变**（仍一次调用、一个 JSON 输出，不增加 latency）。改的是 prompt 文本组织，可分别维护和 eval                                                                                        | 让 prompt 职责更清楚，降低单 prompt 认知负担，以后单独迭代 task 判断和员工话术 |

所以 Phase 1 和 Target State 的区别是：

- **Phase 1 是 mutation architecture refactor**：先把"谁能写 tasks / contacts / timeline、怎么写、怎么审计"定下来。它解决的是 correctness 和 shared infrastructure
- **Target State 是 prompt architecture refactor**：把"哪个 prompt 负责画像、哪个 prompt 负责 task decision、哪个 prompt 负责 playbook"拆清楚。**已经并行外包给独立 engineer**，Phase 1 不阻塞这条线

测试环境没有 prod migration 压力，所以 Phase 1 可以做得更快；拆好的 prompt 一交付就可以接入。两条线的接触点只有 `TaskDecision` 的 6 个 action 命名（见「Phase 1 实施步骤」），双方对齐这一点即可独立推进。

#### Historical Target：Contact Analyzer 拆成画像 + task 决策两段

现在 Contact Analyzer 是一个 1000+ 行的大 prompt，同时做客户画像分析和 task 判断。读真实 prompt 后确认：画像段（SECTION 1-3）和 task 决策段（SECTION 4）已经是天然两块，中间只有一条单向依赖 —— task 决策需要先知道 `lifecycleStage` 才能选 `typeCategory`。所以"拆开"是顺势而为，不是发明新结构。

##### 拆 prompt 的方式：拆文本/组装，**不拆 LLM 调用**（经 Codex + Gemini cross-validate）

拆有两种做法，必须区分清楚：

| 做法           | 怎么拆                                                                    | LLM 调用 | input token                  | 成本                             | 取舍                                                                                                  |
| ------------ | ---------------------------------------------------------------------- | ------ | ---------------------------- | ------------------------------ | --------------------------------------------------------------------------------------------------- |
| **拆文本（采用）**  | 画像段、task 段拆成两个 prompt 模板 / 组装函数，**拼成一个 system prompt，一次调用，一个 JSON 输出** | 1 次    | 不变                           | ≈ 零                            | ✅ 解决"prompt 太长难维护"，可分别 eval，无新 failure mode                                                         |
| 拆调用（**不采用**） | 先调一次 LLM 出画像 → 画像结果喂进第二次调用做 task 决策                                    | 2 次    | ≈ 翻倍（两次都带 calls/messages 历史） | latency +1 round-trip、token 翻倍 | ❌ 引入"画像成功 / task 失败"半成功态；且画像段的 `actionNeeded` 与 task 段的 `taskDecisions[]` 一致性约束会断裂，需额外 reconcile 代码 |

**采用拆文本**。理由：当前痛点是 prompt 长度和可维护性，拆文本零成本解决；画像→task 是单向依赖，拆调用收益有限却引入新成本。三方（读真实 prompt 的分析 + Codex + Gemini）一致不建议把"拆成两次 LLM 调用"作为目标态。

##### 未来演进：画像持久化 + 增量（不经过"两次完整调用"）

如果将来积累到"task 决策质量被画像 prompt 拖累"的实测证据，下一步**不是**升级成两次完整调用，而是：把画像结果**持久化**（其实 `contacts.lifecycleStage` 等字段已经在存），task 决策 prompt 只吃「已存画像 + 最近新增的 call/SMS delta + 当前 open tasks」，不再重读全部历史 —— input token 砍一大半。这条路同时也是触发层优化的根治方案（见「触发机制与优化路线」）。

```text
拆文本后的 Contact Analyzer（仍一次 LLM 调用）：

┌───────────────────────────────────────────────────────────┐
│  Call Analysis（不变）→ 写 calls 表 + 产出 per-call facts    │
└──────────────────────────┬────────────────────────────────┘
                           │ 触发 = (contactPhone && storeId)
                           ▼
┌───────────────────────────────────────────────────────────┐
│  Context Builder：聚合 calls + messages + leads             │
│  + contacts + open tasks（一次性读，下面两段共享）           │
└──────────────────────────┬────────────────────────────────┘
                           ▼
┌───────────────────────────────────────────────────────────┐
│  一次 LLM 调用，system prompt = 模板1 + 模板2 拼接           │
│  ┌─ 模板1：画像段 ────────┐  ┌─ 模板2：task 决策段 ───────┐ │
│  │ lifecycle/leadStatus  │─▶│ 读画像结果选 typeCategory   │ │
│  │ /DNC/summary/risk/    │依赖│ → taskDecisions[]          │ │
│  │  goals/customerSummary│  │ (create/close/update/       │ │
│  │  (不可派生的语义判断)  │  │  record_progress/create_closed)│
│  └───────────────────────┘  └────────────────────────────┘ │
│  两个模板独立维护/eval，但运行时一次调用、一个 JSON 输出     │
└──────────────────────────┬────────────────────────────────┘
                           ▼
┌───────────────────────────────────────────────────────────┐
│  MutationPlan（数据结构，不是执行方式）                     │
│  { identity, actor, source, idempotencyKey,                 │
│    proposedProfile, contactActions[], taskActions[],        │
│    timelineEvents[] }   ← contactActions[]/taskActions[] 可空 │
└──────────────────────────┬────────────────────────────────┘
                           ▼
┌───────────────────────────────────────────────────────────┐
│  ╔═══════════════════════════════════════════════════════╗ │
│  ║  Policy Guard（审整个 MutationPlan）                   ║ │
│  ║  执行顺序（代码保证）：                                 ║ │
│  ║    1. Contact Writer ──→ contacts 表                  ║ │
│  ║    2. Task Orchestrator ──→ tasks + progress 表       ║ │
│  ║    3. Timeline Writer ──→ contact_timeline 表         ║ │
│  ║  三步放在一个 db.batch() 里原子提交                    ║ │
│  ╚═══════════════════════════════════════════════════════╝ │
└──────────────────────────┬────────────────────────────────┘
                           ▼
┌───────────────────────────────────────────────────────────┐
│  Task Playbook Prompt（只读，不写表）→ 生成员工话术         │
└───────────────────────────────────────────────────────────┘
```

#### 关键设计决策

- **Call Analysis 不升级**。继续只产 per-call facts/signals，不读 SMS/leads/tasks，不做 task decision。Task 判断需要跨通话的历史上下文（open tasks、多次 no\_answer 的累积），per-call prompt 没有这些信息
- **拆文本不拆调用**。画像段和 task 决策段拆成两个模板分别维护，但拼成一个 system prompt 一次调用。不引入第二次 LLM 调用（避免 latency 翻倍 + 一致性约束断裂），详见上方「拆 prompt 的方式」
- **先画像，后任务**。task 决策段需要知道客户当前状态（lifecycle、leadStatus、DNC）才能决定建什么 task。这个依赖在一次调用内由 prompt 内部顺序保证（模板1 在模板2 之前）
- **MutationPlan 是信封，不是执行方式**。contactActions\[] 可以为空。执行顺序由代码保证：contact → task → timeline
- **`followUpNeeded` 不是 trigger，是参考信号**。contacts-analyzer 的触发条件是 `contactPhone && storeId`（每通电话都触发）；`followUpNeeded` 是上游 per-call AI 写进 `calls` 表的字段，Contact Analyzer 读它当 prompt context（渲染成 `fu=yes(...)` 一行），不参与触发判断也不进任何代码分支。三个信号字段不是一条链：`followUpNeeded`（calls 表，per-call 输出，被参考）/ `taskDecisions[]`（不落库的中间变量）/ `actionNeeded`（落 contacts 表的真状态）各处不同层 —— 详见「信号字段定位」

### 4 个共享 module

| Module            | 职责                                                                                                                                                    | 返回值                                                         |
| ----------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------- |
| Policy Guard      | 写入前检查：DNC / storeId / 状态机 / 幂等 / 去重 / **AI hallucination guard**（AI close/update 必须引用真实存在的 open task）/ **human authority guard**（员工的改动不被过期 AI 提案静默覆盖） | `allow` / `reject` / `needs_review`（Phase 1 不走该分支，见 Step 4） |
| Task Orchestrator | 统一 task 写入：create / create\_closed / close / update / record\_progress / reopen                                                                       | Drizzle statements                                          |
| Contact Writer    | 统一 contacts 写入：identity（trust scoring）/ DNC（sticky）/ lastActivityAt（forward-only）                                                                     | Drizzle statements                                          |
| Timeline Writer   | 统一 contact\_timeline 写入：event catalog + Zod schema / 幂等                                                                                               | INSERT statement                                            |

Policy Guard 的 `needs_review`：低置信度的 AI 提案不自动落库，进入人工确认队列。这是 AI autonomy 控制的核心机制。

Timeline Writer 接收两种写入：**source event**（"一通电话打完了" — 不过 Policy Guard，只记录事实）和 **mutation audit**（"AI 建了一个 task" — Task Orchestrator / Contact Writer 内部自动写的，已过 Guard）。

### 7 张表

| 表                      | 谁写                | 角色                     |
| ---------------------- | ----------------- | ---------------------- |
| `calls`                | 通话分析 pipeline 独占  | per-call 源记录 + AI 分析结果 |
| `messages`             | 短信处理 pipeline 独占  | per-message 源记录        |
| `leads`                | Lead pipeline 独占  | per-lead 源记录           |
| `contacts`             | Contact Writer    | 客户聚合画像                 |
| `tasks`                | Task Orchestrator | 业务事项当前快照               |
| `task_progress_events` | Task Orchestrator | task 进展 SoT（每次尝试的独立记录） |
| `contact_timeline`     | Timeline Writer   | contact 维度审计 / feed 投影 |

`task_progress_events` 和 `contact_timeline` 并存：progress event 会同步投影到 timeline，但 SoT 在 `task_progress_events`（按 task 维度查询用独立表，不用从 timeline JSONB 里 parse）。

### 信号字段定位（followUpNeeded / actionNeeded / taskDecisions）

这三个字段常被当成"一条 AI 信号链"，但读真实代码后它们其实在不同层，归属不同 prompt，落不落库也不同。拆 prompt 前先把它们定位清楚：

| 字段                                              | 真身                                                           | 落哪                       | 处理                                                                    |
| ----------------------------------------------- | ------------------------------------------------------------ | ------------------------ | --------------------------------------------------------------------- |
| `followUpNeeded`（+ `followUpReasons`）           | **上游 per-call AI 的事实输出**（这通电话要不要跟进），不可派生 + 有独立 dashboard 消费者 | `calls` 表（per-call AI 写） | **保留**。与拆 Contact Analyzer 正交 —— 它在上游，是 Contact Analyzer 的输入证据之一，不是结论 |
| `taskDecisions[]`                               | **中间变量**（AI 的"指令"，代码立刻翻译成 tasks 行的瞬时载体）                      | **不落库**（从来不是列）           | 已是最佳形态。task 决策段输出                                                     |
| `contacts.actionNeeded`（+ `actionNeededReason`） | "这个客户要不要跟进" —— 但在 task 4 层 model 下，它 = 是否存在 open task        | `contacts` 表             | **退役（分 2 phase）**（见下）                                                 |
| `tasks.actionNeeded`                            | task 只要 `status='open'` 就代表"需要处理"；这个 boolean 是 status 的重复表达  | `tasks` 表                | **退役**（见下）                                                            |

关键纠正：`followUpNeeded` 不是 Contact Analyzer 的触发器（触发是 `contactPhone && storeId`），它是被 Contact Analyzer **读去当 prompt 参考**的信号，读完不二次落库，**且不该被直接当结论沿用** —— task 决策要基于跨通话上下文重新判断（per-call 的 followUpNeeded 是单通粒度，task 是跨通话粒度）。

#### 两个 `actionNeeded` 都应退役（由 open task 派生）

`tasks.actionNeeded` 和 `contacts.actionNeeded` 是同一个冗余的两半，根因都是"旧 task 系统不够强、没法可靠地用 open task 表示'要不要跟进'"。task 4 层 model 把 task 做成可靠的 work object（`status='open'` = 待办）后，两个都能由 open task 派生：

- **`tasks.actionNeeded`**（Phase 1 直接删）：
  - 派生公式：`actionNeeded === (status === 'open')`
  - 自然语言：task 只要 open 就代表"需要处理"，actionNeeded 是 status 的重复表达。从历史代码看，所有 caller（lead-processor / contacts-analyzer / studio-api）create 时都硬编码写 `true`，从来没人写 `false` —— 这个 boolean 没有信息量
  - 实施：Phase 1 直接 DROP COLUMN
- **`contacts.actionNeeded`**（Phase 1 改读路径，Phase 2 DROP COLUMN）：
  - 派生公式：`actionNeeded === EXISTS(SELECT 1 FROM tasks WHERE contact_phone=? AND store_id=? AND status='open')`
  - 自然语言：「该 contact 要不要跟进」等价于「该 contact 名下有没有 open task」，task 是可靠的 SoT
  - **实施分 2 步**（因为 studio-api 真在读 `c.action_needed`：contacts list query / leads KPI summary）：
    - **Phase 1**：改 `routes/v3/contacts.ts` + `routes/v3/leads.ts` 把 `c.action_needed` 替换成 `EXISTS(...)` 子查询；contact\_analyzer 暂时仍写该字段（保险，observe-only）；**不动 schema**
    - **Phase 2**：observe 一段时间确认 EXISTS 口径没漂移后，DROP COLUMN + 停 contact\_analyzer 写入 + 删 `idx_contacts_action_needed` index
  - Leads KPI 口径同步改成"有 open task 的 contact 数"

这呼应设计准则 **Code as Guardrail, AI as Judgment**：「要不要跟进」是 open task 的确定性派生，不该有第二个独立真值源。字段级冗余说明见 schema 文档（[tasks-schema](/product-design/v2/tasks-feature/tasks-schema.md) / [contacts-schema](/product-design/v2/contacts-feature/contacts-schema.md)）。

### AI 永远不负责的事

AI 只输出结构化提案（Zod 验证），以下由代码负责：

- final task state（状态转换）
- idempotency（幂等）
- permission / DNC hard stop
- store isolation
- transaction boundary
- audit trail

### Code vs AI 职责

| 决策                             | 负责层                                                   |
| ------------------------------ | ----------------------------------------------------- |
| 是否触发 pipeline / exact STOP DNC | 代码                                                    |
| 自然语言 DNC                       | AI 提案 + 代码验证（DNC sticky，AI 不可反转）                      |
| Task typeCategory              | AI 在代码提供的 allowed set 内选                              |
| Progress vs Close              | 代码状态机（no\_answer/left\_voicemail 是 progress，不是 close） |
| Name trust arbitration         | 代码（`CASE WHEN $trust > COALESCE(...)` SQL）            |
| closeResult                    | AI 在 allowed enum 内选，代码校验                             |
| 幂等 / 并发 / 事务 / storeId 隔离      | 代码 + DB                                               |
| suggestedActions 内容            | AI                                                    |

### 现状 vs 目标

> 现状栏区分"已收口"和"仍缺"——多数共享 helper 已就位，真正的缺口集中在 task mutation 状态机。

| 现状                                                                                                                                        | 目标                                                                                                           |
| ----------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------ |
| trust score 已集中到 `NAME_TRUST` 共享常量（已收口）                                                                                                   | 由 Contact Writer 封装调用                                                                                        |
| timeline 写入已统一走 `buildTimelineValues` / `buildContactTimelineInsertSQL`（已收口）                                                              | 由 Timeline Writer 封装 + event catalog                                                                         |
| dueAt 已抽出 `computeDueAt`（已收口；studio-api 是唯一例外，hardcode INTERVAL）                                                                          | Task Orchestrator 内统一调用                                                                                      |
| DNC cascade 已共享 `dnc-cascade.ts`（已收口）                                                                                                     | 并入 Policy Guard                                                                                              |
| **task create/close/update 状态机仍各写各的（缺 Task Orchestrator）**                                                                                | Task Orchestrator 统一状态机                                                                                      |
| **没打通 → 关旧 task + 建新 task**                                                                                                               | `record_progress`：task 保持 open，追加进展记录                                                                        |
| **没有 `task_progress_events` 表**                                                                                                           | 新增                                                                                                           |
| **可派生字段被 AI 单独输出 / 单独存（`tasks.actionNeeded` / `taskType` 等）**                                                                             | AI 只输出不可派生判断，可派生值由代码派生（见 schema 文档冗余说明）。`contacts.lifecycleState` 不在此列 —— 撤销退役决定,churned re-engage 需 AI 语义判断 |
| **`tasks.status` 用 `'pending' / 'closed'`（命名暗示"还没开始"）**                                                                                   | 改名为 `'open' / 'closed'`（schema + code + UI 同步）—— 测试环境直接改，prod 无迁移压力                                          |
| **`tasks.store_id` 仍 nullable**（partial unique 用 `WHERE store_id IS NOT NULL` 跳过历史 NULL 老 task）                                           | Phase 1 backfill 老行 + `SET NOT NULL` + 删 partial WHERE。Tenant isolation DB-level 强制                          |
| studio-api close **已写** contact\_timeline（V1.8 后用 `buildContactTimelineInsertSQL` + `logTimelineEvent`，event\_type=`task.status_changed`） | 保持，由 Task Orchestrator + Timeline Writer 内部自动写                                                               |

### One-time AI Call vs Tool Calling

两种模式共享同一套 `Task Orchestrator` / `Contact Writer` / `Timeline Writer` / `Policy Guard`。

| 维度    | Mode A: one-time call（现在）                               | Mode B: tool calling（未来）           |
| ----- | ------------------------------------------------------- | ---------------------------------- |
| AI 角色 | 数据产出者（代码编排 AI output）                                   | 编排者（AI 决定调哪个 tool）                 |
| 触发    | Event-driven（SQS / cron）                                | User/agent-initiated（API / voice）  |
| 延迟    | 分钟级                                                     | 秒级                                 |
| 调用方式  | 代码遍历 AI output → 调 `Task Orchestrator` / writer modules | Tool handler → 调同一套 shared modules |

Phase 1 建好 shared modules 后，tool calling 的 backend 调用路径已经有了（tool handler 调 shared modules）。但真正开放 tool calling 还需要额外前置条件：RBAC（谁能调哪些 tool）、budget cap（防止 AI 无限调用）、approval threshold（低置信度进 needs\_review）、audit actor（记录 actorType='ai\_agent'）、idempotency（防重复调用）。这些在 Phase 1 → Target State 的过渡期间完成。

| Tool                   | 调用的 module                                 |
| ---------------------- | ------------------------------------------ |
| `create_task`          | `TaskOrchestrator.create()`                |
| `close_task`           | `TaskOrchestrator.close()`                 |
| `record_task_progress` | `TaskOrchestrator.recordProgress()`        |
| `update_contact`       | `ContactWriter.update()`                   |
| `analyze_call`         | `CallAnalysisModule.analyzeCall()`         |
| `reanalyze_contact`    | `ContactAnalysisModule.reanalyzeContact()` |

### 加新 workflow 应该是什么体验

Target state 下加一个新 workflow 只需要 5 步：

1. 定义事件 / 触发条件
2. 决定判断层是代码、AI 还是混合
3. 复用或新增 action schema
4. 调共享 module
5. 加 test + event catalog payload

不需要：设计新的 DB writer / 发明新的 timeline payload / 重新决定 DNC 行为 / 重新决定 task 去重 / 重新决定 store 隔离 / 教 prompt 怎么操作数据库。

***

## 2. Phase 1 — 实施步骤

Phase 1 不改 pipeline runtime（SQS/Lambda 拓扑不变，不需要新的 AWS 资源），只把共享表写入收口到共享函数。所有 module 是 in-process TypeScript 函数。

### Step 1: 历史方案中的 Task Pipeline Deliverable（已 superseded）

[Task Pipeline Deliverable (Codex)](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) 当时被作为 Phase 1 基础；以下内容是历史方案，不得绕过 V3 Task contract 直接落地。它定义了：

- Task 的 4 层 Object Model（Activity / Progress / Lifecycle / Outcome）
- `record_progress` 语义：没打通不关 task，追加进展记录（解决"打 5 次没接 = 5 个 task"的问题）
- `create_closed` 语义：当场完成的业务目标（如 inbound 当场成交）直接建一个 closed task，让 dashboard 能看到
- `closeResult` 枚举清理：progress 值（no\_answer / left\_voicemail）移出 closeResult，只留 business outcome
- API contract（v3 endpoints）
- Prompt contract（输入 allowed set、输出 proposed mutations）
- UI semantics（Progress Update vs Close Task 两种操作）
- Code vs Prompt 职责矩阵

### Step 2: 新建 Task Orchestrator + `task_progress_events` 表 + `status` 改名

| 新建项 / 改动                          | 位置                                                                                                                                                                                                                                                                        |
| --------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| **`tasks.status` 改名**             | `callytics-common/src/db/schema/tasks.ts:16` `TASK_STATUS = ['open', 'closed']`（原 `'pending' / 'closed'`）+ DEFAULT `'open'`。测试环境 `UPDATE tasks SET status='open' WHERE status='pending'` 一次性迁移。全 repo grep `'pending'` sweep（4 个 Lambda + studio-api + CHECK constraints） |
| `task_progress_events` 表          | Neon PostgreSQL。幂等：`idempotencyKey` unique + `(taskId, callId, progressType)` partial unique                                                                                                                                                                              |
| Task Orchestrator                 | `callytics-common/src/domain/task-orchestrator.ts`。`applyTaskAction()` — 6 种 action：`create_open` / `create_closed` / `close` / `update` / `record_progress` / `reopen`。返回 Drizzle statements，调用方放进 `db.batch()` 原子执行                                                     |
| `POST /v2/tasks/:taskId/progress` | studio-api 新端点                                                                                                                                                                                                                                                            |

#### `applyTaskAction()` 并发/幂等 contract（每 action SQL 层显式）

| Action            | DB 幂等 / 并发约束                                                                                                                                                                                                                      | 备注                                                                                                                                                                   |
| ----------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `create_open`     | `uq_tasks_pending_contact_category` partial unique on `(contactPhone, storeId, typeCategory) WHERE status='open'`（现有 index 的 `status='pending'` 改成 `'open'`）                                                                      | 防同 contact 同 typeCategory 重复 open task                                                                                                                               |
| `create_closed`   | INSERT 行必带 `tasks.source_call_id` (`create_closed` 的 evidence reference，对应触发"当场办成"的 callId)；**新增 partial unique on `(contactPhone, storeId, typeCategory, source_call_id) WHERE status='closed' AND source_call_id IS NOT NULL`** | 防同一通电话当场办成的 task 被重复建 closed                                                                                                                                         |
| `close`           | `UPDATE ... SET status='closed' WHERE task_id=$1 AND status='open' AND updated_at <= $aiRunStartedAt RETURNING *` — 0 rows = reject                                                                                               | human authority guard：staff 已动过则 AI 提案 reject                                                                                                                        |
| `update`          | `UPDATE ... SET ... WHERE task_id=$1 AND status='open' AND updated_at <= $aiRunStartedAt RETURNING *`                                                                                                                             | 同 close                                                                                                                                                              |
| `record_progress` | `INSERT INTO task_progress_events ... ON CONFLICT (idempotency_key) DO NOTHING` + `UPDATE tasks SET attempt_count = attempt_count + 1, due_at = $newDueAt WHERE task_id=$1`                                                       | **用 `idempotency_key` 而不是 `(task_id, call_id, progress_type)` ON CONFLICT — 因为 PostgreSQL `NULL ≠ NULL`，SMS progress 无 callId 时复合 key 不去重。`idempotency_key` 生成规则见下** |
| `reopen`          | `UPDATE ... SET status='open' WHERE task_id=$1 AND status='closed' RETURNING *`                                                                                                                                                   | 0 rows = task 不是 closed 状态，reject                                                                                                                                    |

#### `task_progress_events.idempotency_key` 生成规则（caller 端）

| Progress 来源                                | idempotency\_key 格式                                                            |
| ------------------------------------------ | ------------------------------------------------------------------------------ |
| 电话 progress（callId 非空）                     | `progress:{taskId}:call:{callId}:{progressType}`                               |
| SMS progress（messageId 非空）                 | `progress:{taskId}:message:{messageId}:{progressType}`                         |
| 手动 progress（staff 触发，无 callId / messageId） | `progress:{taskId}:manual:{actorId}:{occurredAt.toISOString()}:{progressType}` |
| 系统兜底（cron / reconciliation）                | `progress:{taskId}:system:{runId}:{progressType}`                              |

这是 2026-06 historical proposal 的 standalone table 去重设计，不是当前实施方向。V3 proposal 中的 Timeline 使用分层的 `commandId`、`eventIdempotencyKey` 与 `sourceInteractionKey`，见 [Task V3 Target Design Proposal §9.4](/product-design/v3/tasks-feature/task-domain-lifecycle.md#94-timeline--activity)。

### Step 3: 新建 Timeline Writer + Event Catalog

| 新建项             | 位置                                                                                                       |
| --------------- | -------------------------------------------------------------------------------------------------------- |
| Timeline Writer | `callytics-common/src/domain/timeline-writer.ts`。封装现有 `buildTimelineValues()` + Event Catalog Zod schema |

Task 事件必须写 timeline，所以 Timeline Writer 要跟 Task Orchestrator 同步就位。

### Step 4: 新建 Contact Writer + Policy Guard + `tasks.store_id` NOT NULL

| 新建项 / 改动                                | 位置                                                                                                                                                                                                                                                                             |
| --------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| Contact Writer                          | `callytics-common/src/domain/contact-writer.ts`。Phase 1 最小版：`upsertIdentity()`（trust scoring）/ `touchActivity()`（forward-only lastActivityAt）/ `setDNC()`。lifecycle guard 先 observe-only 积累数据，Phase 2 再 enforce                                                                |
| Policy Guard                            | `callytics-common/src/domain/policy-guard.ts`。DNC guard + storeId guard + duplicate guard + state transition + idempotency + AI hallucination guard（close/update 引用的 taskId 必须是当前 contact 真实存在的 open task）+ human authority guard（conditional update：员工已改动的 task 不被过期 AI 提案覆盖） |
| **`tasks.store_id` NOT NULL migration** | 测试环境：先 `UPDATE` backfill NULL 老行（按 `contact_phone + franchise_id + account_id` 反查 `PhoneStoreAssignments`）或直接 `DELETE`（测试数据无价值）→ `ALTER COLUMN store_id SET NOT NULL` → 删 `idx_tasks_store_id` 上的 `WHERE store_id IS NOT NULL` partial 条件                                      |

Contact Writer Phase 1 不统一 AI 画像字段（`customerSummary` 等 18 个），因为只有 contacts-analyzer 写，不存在多 writer 竞争。

#### Policy Guard 返回值

| 返回值            | Phase 1 行为                                                                                                                                                                                    |
| -------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `allow`        | Orchestrator 执行 SQL                                                                                                                                                                           |
| `reject`       | log + skip。低置信度提案也走这条（不引入 needs\_review queue/table）。原因：阈值与人工复核流程尚未 calibrate，Phase 1 先用 CloudWatch metric 收数据（`mutation_rejected_total{reason=...}`），Phase 2 再决定要不要建 needs\_review 队列、谁审、SLA |
| `needs_review` | type 保留为 contract 第三态，Phase 1 不走这条分支。Phase 2 接 needs\_review 实现时不用改 type，只改 dispatch 逻辑                                                                                                       |

### Step 5: 4 个 caller 逐个迁移

| 顺序 | Caller                 | 现在                           | Phase 1                                                                                                                                                                 |
| -- | ---------------------- | ---------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| 5a | studio-api close       | raw SQL + 关旧建新（已写 timeline）  | 调 `Task Orchestrator`。`no_answer`/`left_voicemail` 改走 `record_progress`。timeline 写入由 Orchestrator + Timeline Writer 内部完成（不再 caller 自己拼 `buildContactTimelineInsertSQL`） |
| 5b | contacts-analyzer      | 遍历 `taskDecisions[]` 自己拼 SQL | `taskDecisions.map(d => applyTaskAction(d))`                                                                                                                            |
| 5c | lead-processor         | 直接 `db.insert(tasks)`        | 调 `applyTaskAction({ action: 'create_open' })`                                                                                                                          |
| 5d | message-processor STOP | 直接拼 DNC cascade SQL          | 调 `ContactWriter.setDNC()` + `TaskOrchestrator.closeAllOpen()`                                                                                                          |

现有 API 端点（close / reopen / postpone）外部行为不变，内部改成调 `Task Orchestrator`。

Prompt schema 改动：contacts-analyzer `TaskDecision` 6 个 action 命名为 `create_open` / `create_closed` / `close` / `update` / `record_progress` / `reopen`（与 Task Orchestrator `applyTaskAction()` 对齐）。**由独立 engineer 并行交付，Phase 1 拿来直接接入**。

### Phase 1 不做

不切换 tool calling runtime、不引入 workflow engine、不合并 Lambda、不做 meaningful SMS 分类、不做 production Task V2 数据迁移。

***

## 3. Historical Phase 1 → 当时的 Target State

Phase 1 完成后基础架构就位。Target State 加更多 caller 和扩展。

| Phase 1 已有                                | Target State 加的                         |
| ----------------------------------------- | --------------------------------------- |
| Task Orchestrator（6 种 action）             | Tool calling thin wrapper               |
| Contact Writer（identity / DNC / trust）    | 扩展到全部 18 个 AI 字段                        |
| Timeline Writer                           | 更多 event type                           |
| Policy Guard（DNC / storeId / dedup / 状态机） | RBAC、rate limit、AI autonomy 级别          |
| 4 个 caller 迁移完                            | Voice agent、SMS meaningful、API platform |

### Tool Calling 前置条件

| 前置条件                    | Phase 1 后状态                                          |
| ----------------------- | ---------------------------------------------------- |
| Shared mutation modules | 已就位                                                  |
| Tenant isolation        | contacts PK 已含 store\_id；tasks.store\_id 需改 NOT NULL |
| RBAC                    | studio-api 有 store-level 校验，需泛化                      |
| Audit trail             | Timeline Writer 需新增 `actorType = 'ai_agent'`         |
| Budget cap              | 需新建                                                  |

### SMS meaningful（Phase 2）

当前 message-processor 只做：存消息 + exact STOP DNC cascade。Phase 2 加代码分类器（trivial / meaningful / natural-language DNC），meaningful SMS 即时触发 Contact Analysis Reprocess。

### 触发机制与优化路线

Contact Analyzer 有 4 个触发源（`per_call_analysis` / `cron` / `on_demand` / `reprocess`），共享同一条 FIFO 队列 + 同一个 worker。现状能跑（DeepSeek 成本低 + daily cron 06:00 兜底 + 即时性对前台有价值），但触发机制有重复分析，**量大后需优化**。详细 trace + 修复方案见 **GitHub issue [#365](https://github.com/retaintive/docs/issues/365)**。

| # | 冗余场景                          | 真冗余？          | 现有保护                                        |
| - | ----------------------------- | ------------- | ------------------------------------------- |
| 1 | per\_call 已分析 → 当晚 cron 又分析一次 | 是（故意兜底，实现可优化） | ❌ cron 查询不看 `last_contact_analysis_at`      |
| 2 | 一天打 N 通电话 → per\_call 触发 N 次  | **是，最浪费**     | ❌ 完全无 cooldown（dedup id 含 callId，FIFO 去重失效） |
| 3 | cron 同天对同 contact 重复发         | 理论/轻微         | ✅ FIFO dedup（5min 窗口）                       |
| 4 | reprocess 撞上 cron             | 否             | ✅ cooldown gate 24h                         |
| 5 | 每次触发重读全部历史 + 跑大 prompt        | 是（成本乘数）       | ❌ 无缓存/增量                                    |

根因：cooldown gate（`contacts-analyzer/handler.ts`）只对 `source === 'reprocess'` 生效，per\_call / cron 不走它；而它检查的 `lastContactAnalysisAt` 字段每次分析都在写（基础设施现成，只是没接上）。

**优化路线（非紧急，量大再做）：**

1. **触发层**（冗余 1/2）：per\_call 接入 cooldown / 短 debounce；cron 查询跳过"自上次分析后无新活动"的 contact。低风险，`lastContactAnalysisAt` 已就绪。
2. **计算层**（冗余 5）：即「拆 prompt 的方式」里说的画像持久化 + 增量分析 —— task 决策只吃「已存画像 + 新增 delta」，不重读全史。优先级低于触发层（先削触发次数 ROI 更高）。

这两层都符合"先做一个能跑的版本（拆 prompt 文本 + Phase 1 写入层），scalability 优化随量增长再做"的策略。

***

## 4. Appendix

### Processing Capability Module contracts

Target state 下 4 个 Processing Capability Module 是可复用的调用接口（Layer 2），它们产出结构化提案交给共享写入层（Layer 1）执行。

| Module           | 调用接口                                        | 职责                                              | 不负责                                                            |
| ---------------- | ------------------------------------------- | ----------------------------------------------- | -------------------------------------------------------------- |
| Call Analysis    | `analyze_call(callId, mode)`                | transcript / per-call AI / call 分类              | task 创建                                                        |
| SMS Signal       | `evaluate_sms(messageId, mode)`             | exact STOP / trivial filter / meaningful signal | 直接写 task                                                       |
| Contact Analysis | `reanalyze_contact(phone, storeId, reason)` | 聚合跨互动上下文 + AI 判断                                | 直接 DB mutation                                                 |
| Lead Downstream  | `process_lead_downstream(leadId)`           | 确定性 lead → contact/task                         | 不需要单独 heavy module（直接调 `Contact Writer` + `Task Orchestrator`） |

### 分歧决策

| 分歧点                                                 | 决策                               | 理由                                                                                                                                       |
| --------------------------------------------------- | -------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------- |
| `task_progress_events` 独立表 vs 复用 `contact_timeline` | **独立表**                          | progress 是 task 维度高频查询，timeline JSONB parse 脏且慢。Opus 倾向复用 timeline；Codex Phase 1 文档定义独立表；Task Pipeline Deliverable 明确支持独立表。Final 决策采用独立表 |
| `reopen` action                                     | **保留**                           | studio-api 已有 reopen 路由                                                                                                                  |
| Actor 字段                                            | **先 3 字段（type / id / name），可扩展** | Phase 1 够用                                                                                                                               |
| Source Event Projection                             | **概念保留，不改代码**                    | 只是命名区分                                                                                                                                   |

### 详细参考

| 内容                                              | 位置                                                                                                                                                                                                           |
| ----------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| Business Object 健康检查                            | [Opus §2.1](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md)                                                                                                              |
| Writer 清单（contacts 7 / tasks 10 / timeline 15+） | [Opus §2.2](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md)                                                                                                              |
| Contact Writer 完整 contract                      | [Opus §3.2.2](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md)                                                                                                            |
| Timeline Writer Event Catalog（12 种 eventType）   | [Opus §3.2.3](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md)                                                                                                            |
| Code vs AI 完整职责矩阵                               | [Opus §3.4](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md) + [Task Pipeline Deliverable §7](/product-design/v2/tasks-feature/design/task-pipeline-deliverable-codex.md) |
| Reuse Scenario Validation（4 个场景）                | [Opus §4.5](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md)                                                                                                              |
| Business Object 语义定义                            | [Codex §Ideal Business Object Semantics](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md)                                                                                |
| Target AI Pattern + AI never owns 清单            | [Codex §Target AI Pattern](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md)                                                                                              |
| 加新 workflow 的体验定义                               | [Codex §What Adding a New Workflow Should Feel Like](/product-design/v2/unified-pipeline/original/unified-pipeline-target-state-codex.md)                                                                    |

### Open Questions

| 问题                                                                             | 背景 / 决策状态                                                                                                                                                                                                                               |
| ------------------------------------------------------------------------------ | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| Task Orchestrator 放 `callytics-common` 还是 `callytics-infrastructure`？          | studio-api 当前用 raw SQL 不用 Drizzle                                                                                                                                                                                                       |
| ~~`contacts.actionNeeded` 何时 deprecate？~~ **已决**                               | 由 `EXISTS(open task)` 派生。**Phase 1 改读路径（`routes/v3/contacts.ts` + `routes/v3/leads.ts` 用 EXISTS 子查询），字段保留 observe-only；Phase 2 DROP COLUMN + 停 writer + 删 index**。`actionNeededReason` / `suggestedActions` 归属仍待定（留 contacts 还是并入 task） |
| ~~open task 创建 action 叫 `create` 还是 `create_open`？~~ **已决**                    | 统一用 `create_open`（与 `create_closed` 对称）。prompt engineer 跟 schema 对齐                                                                                                                                                                     |
| ~~`tasks.status` 命名 `pending` vs `open`？~~ **已决**                              | 改为 `open` / `closed`（schema + code + UI 同步）。测试环境直接迁移                                                                                                                                                                                    |
| ~~Policy Guard `needs_review` queue / table / ownership？~~ **已决（Phase 2 再设计）** | Phase 1 低置信度直接 `reject` + log + CloudWatch metric 收数据。`needs_review` 作为 return type 第三态保留，但 Phase 1 不走该分支。数据积累后 Phase 2 再决定要不要建 queue / 谁审 / SLA                                                                                        |
| `neon-sync` Lambda 是否 Phase 1 删除？                                              | 确认死代码                                                                                                                                                                                                                                   |
| 前端 UI 是否 Phase 1 scope？                                                        | `record_progress` 需要前端改                                                                                                                                                                                                                 |
| `closeResult` 枚举三方不同步？                                                         | schema 18（含 booked/cancelled）/ Task Playbook prompt 仍含 stale `attempted` + `backend-gated booked/cancelled` / 前端 16                                                                                                                     |
| 现有 378 条 `closeResult='attempted'` 是否数据迁移？                                     | 还是只影响新数据                                                                                                                                                                                                                                |
| `processing_runs` 是否 first-class table？                                        | retry / backfill 需要 run ledger                                                                                                                                                                                                          |
| Timeline payload 是否 versioned？                                                 | event schema 会演进                                                                                                                                                                                                                        |
| Phase 1 PR 怎么拆？                                                                | 建议按 Step 2/3/4/5 各一个 PR + 1 个独立 PR 做 `status` 命名 sweep                                                                                                                                                                                  |

### Implementation Handoff Checklist

这部分不改变 design 结论，只是在开始 implementation 前把需要落到 PR / migration / test plan 的 delta 列出来，避免设计文档变成工程细节正文。

#### Schema Delta

| Area                      | Delta                                        | Notes                                                                                                                                                                                                                                     |
| ------------------------- | -------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `task_progress_events`    | 新建表                                          | task progress / attempt 的 source of truth；`idempotency_key` NOT NULL UNIQUE 是唯一去重键（不要用 `(taskId, callId, progressType)`，PostgreSQL `NULL ≠ NULL` 让 SMS progress 无 callId 时不去重）。生成规则见 Step 2                                               |
| `tasks.status`            | `'pending' / 'closed'` → `'open' / 'closed'` | schema enum + DEFAULT 改名；`UPDATE tasks SET status='open' WHERE status='pending'`；全 repo grep `'pending'` sweep；CHECK constraints `chk_tasks_closed_integrity` 等更新；partial unique index 的 `WHERE status='pending'` → `WHERE status='open'` |
| `tasks.store_id`          | nullable → NOT NULL                          | 测试环境 backfill NULL 老行（反查 `PhoneStoreAssignments`）或直接 DELETE；`ALTER COLUMN SET NOT NULL`；删 `idx_tasks_store_id` 的 `WHERE store_id IS NOT NULL` 条件                                                                                          |
| `tasks.source_call_id`    | 新增字段（text NULL）                              | `create_closed` task 的 evidence reference。新增 partial unique `(contactPhone, storeId, typeCategory, source_call_id) WHERE status='closed' AND source_call_id IS NOT NULL` 防同一通电话重复建 closed task                                            |
| `tasks.executorType`      | 新增字段（text NULL）                              | `human` \| `ai_agent` \| `system`，metrics 归因                                                                                                                                                                                              |
| `tasks.attemptCount`      | 新增字段（integer NOT NULL DEFAULT 0）             | 列表 snapshot，SoT 在 `task_progress_events`                                                                                                                                                                                                  |
| `tasks.actionNeeded`      | DROP COLUMN                                  | 派生公式 `(status === 'open')`，无信息量。删字段 + 删 `idx_tasks_action_needed`（如有）                                                                                                                                                                     |
| `tasks.taskType`          | DROP COLUMN                                  | 与 `typeCategory` 重复，已有 `chk_tasks_task_type_category_consistency` CHECK 锁定                                                                                                                                                                |
| `contact_timeline`        | 新增 / 规范 event catalog                        | 至少覆盖 task created / updated / closed / progress recorded；payload 需要 Zod schema                                                                                                                                                            |
| `tasks.closeResult`       | 18 → 15 values                               | 移出 4 个 progress 值（`no_answer` / `left_voicemail` / `callback_later` / `attempted`）到 `task_progress_events.progressType`；新增 `unable_to_reach`                                                                                              |
| `tasks.closeType`         | 2 → 3 values                                 | 新增 `create_closed`                                                                                                                                                                                                                        |
| `contacts.actionNeeded`   | Phase 1 改读路径，Phase 2 DROP COLUMN             | Phase 1：`routes/v3/contacts.ts` + `routes/v3/leads.ts` 用 EXISTS 子查询替代 `c.action_needed`；字段保留 observe-only。Phase 2：停 contact\_analyzer 写入 + DROP COLUMN + 删 `idx_contacts_action_needed`                                                   |
| `contacts.lifecycleState` | **保留**（撤销退役）                                 | churned re-engage 等 case 不可纯派生（需 AI 语义判断 transcripts）;`lead_declined` close result 设计依赖该字段。详见 [contacts-schema §3.2](/product-design/v2/contacts-feature/contacts-schema.md)                                                              |
| `contacts` AI 画像字段        | Phase 1 不迁移全部 18 个                           | 只收口 identity / DNC / trust / lastActivityAt；AI profile fields 留后续                                                                                                                                                                         |

#### API / Contract Delta

| Area                    | Delta                                                                        | Notes                                                         |
| ----------------------- | ---------------------------------------------------------------------------- | ------------------------------------------------------------- |
| Studio API              | 新增 `POST /v2/tasks/:taskId/progress`                                         | `no_answer` / `left_voicemail` 走 progress，不再 close + recreate |
| Existing task routes    | close / reopen / postpone 内部改调 `TaskOrchestrator`                            | 外部 API 尽量兼容，内部统一 policy / idempotency / timeline              |
| Contact Analyzer prompt | `TaskDecision` 加 `record_progress` / `create_closed`                         | 先兼容旧 `create/update/close`，再改 prompt schema                   |
| Shared contracts        | `TaskAction` / `ContactAction` / `TimelineEventPayload` / `ProcessingIntent` | 放在 shared domain package，避免归属于某个 Lambda                       |
| Tool calling            | Phase 1 不暴露 tool runtime                                                     | 只保证 shared modules 未来可被 thin tool handler 调用                  |

#### Test / Rollout Delta

| Area                   | Required check                                                                        |
| ---------------------- | ------------------------------------------------------------------------------------- |
| Unit tests             | `create/create_closed/update/close/record_progress/reopen` 每个 action path             |
| Policy tests           | DNC hard stop、store isolation、state transition、AI hallucination guard、staff override  |
| Idempotency tests      | 重复 action 不重复建 task / progress / timeline                                             |
| Caller migration tests | studio-api、contacts-analyzer、lead-processor、message-processor STOP 各自走 shared modules |
| Replay / compatibility | 旧 prompt output、旧 API close route、partial failure retry                               |
| Prod safety            | production Task V2 单独 audit；Phase 1 不做 destructive prod migration                     |
