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+ 工程审计与实施基线;实施 pipeline 变更前必须重新核查 live code。

日期: 2026-06-01 状态: 合并 Opus + Codex 两份独立设计后的统一版本,待 Max review 环境前提: 本文以当时 test / new architecture 调研为主;production Task V2 行为需要单独 audit,不默认与 test 一致。 设计准则: 产品设计准则(Schema-first → API-first → State-machine-first → Prompt-last / Code as Guardrail, AI as Judgment) 历史前置设计: Task Pipeline Deliverable (Codex) — superseded Task contract,仅用于追溯 Task Orchestrator 的演进依据

原始材料(保留在本目录,不删除):


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 为例,改前改后的区别:

改前: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;为避免混淆,正文尽量写全名。

架构总览

只画写入路径。

当前写入全景图(现状)

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

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

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 PromptTask Decision Prompt(拆文本,不拆调用)。

Prompt / Decision Surface Inventory

#SurfaceType做什么写哪张表Phase 1 影响
Pre-triageCode stage纯代码。用关键词匹配过滤掉 ~40% 不值得分析的通话(自动语音、拨号音等),省掉后面 AI 调用的成本不写表不动
TriageAI promptAI 只看通话前 500 个字符,快速判断这通电话值不值得做完整分析。不值得的直接跳过后面所有 stage不写表不动
ClassifyAI promptAI 读完整 transcript,分析这通电话的分类(sales / cancellation / billing ...)、业务结果(成交 / 未成交 / 取消 ...)、是否需要跟进。产出 33 个字段写入 calls 表。其中 followUpNeeded 是给 Contact Analyzer 的 signalcalls不动
VerifyAI prompt(conditional / observe-only)当 Classify 的 subcategory 落在容易混淆的边界区(如 cancellation vs billing),AI 独立重新判断一次。当前 observe-only(只记 log 不改数据),积累准确率数据后再决定是否启用calls(observe-only)不动
CoachingAI prompt对真人通话(>30 秒),AI 生成给前台的辅导建议:哪句话说得好、哪句可以改进、具体怎么说更好。写入 calls 表的 practical_coaching 字段calls不动
Contact AnalyzerAI prompt系统里当前唯一触发 task/contact 业务变更 proposal 的 prompt。由 SQS 触发(每通电话分析完 + 每日定时 cron),聚合一个客户的所有通话、短信、lead、已有 task 的历史,AI 做跨互动判断:该建 task 吗?该关 task 吗?客户状态变了吗?当前输出 taskDecisions[](create / close / update)+ 18 个 contact 分析字段。followUpNeeded 是它读的 signal 之一(输入证据,不是触发条件,也不该直接当结论沿用)contacts + tasks + contact_timeline核心是改下游 + contractTaskDecision schema 加 record_progress / create_closed,执行逻辑从自己拼 SQL 改成调 Task Orchestrator / Contact Writer。拆 prompt 与此解耦、可独立做(见下方「Phase 1」与「拆 prompt 的方式」)
Task PlaybookAI 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)

┌───────────────────────────────────────────────────────────┐
│  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 改成调共享写入层。

┌───────────────────────────────────────────────────────────┐
│  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 Stateprompt 重构:拆 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 砍一大半。这条路同时也是触发层优化的根治方案(见「触发机制与优化路线」)。

拆文本后的 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 / reopenDrizzle 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 源记录
leadsLead pipeline 独占per-lead 源记录
contactsContact Writer客户聚合画像
tasksTask Orchestrator业务事项当前快照
task_progress_eventsTask Orchestratortask 进展 SoT(每次尝试的独立记录)
contact_timelineTimeline Writercontact 维度审计 / feed 投影

task_progress_eventscontact_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 taskcontacts退役(分 2 phase)(见下)
tasks.actionNeededtask 只要 status='open' 就代表"需要处理";这个 boolean 是 status 的重复表达tasks退役(见下)

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

两个 actionNeeded 都应退役(由 open task 派生)

tasks.actionNeededcontacts.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.tsc.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 / contacts-schema)。

AI 永远不负责的事

AI 只输出结构化提案(Zod 验证),以下由代码负责:

  • final task state(状态转换)
  • idempotency(幂等)
  • permission / DNC hard stop
  • store isolation
  • transaction boundary
  • audit trail

Code vs AI 职责

决策负责层
是否触发 pipeline / exact STOP DNC代码
自然语言 DNCAI 提案 + 代码验证(DNC sticky,AI 不可反转)
Task typeCategoryAI 在代码提供的 allowed set 内选
Progress vs Close代码状态机(no_answer/left_voicemail 是 progress,不是 close)
Name trust arbitration代码(CASE WHEN $trust > COALESCE(...) SQL)
closeResultAI 在 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 + 建新 taskrecord_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 modulesTool 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_taskTaskOrchestrator.create()
close_taskTaskOrchestrator.close()
record_task_progressTaskOrchestrator.recordProgress()
update_contactContactWriter.update()
analyze_callCallAnalysisModule.analyzeCall()
reanalyze_contactContactAnalysisModule.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) 当时被作为 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_eventsNeon PostgreSQL。幂等:idempotencyKey unique + (taskId, callId, progressType) partial unique
Task Orchestratorcallytics-common/src/domain/task-orchestrator.tsapplyTaskAction() — 6 种 action:create_open / create_closed / close / update / record_progress / reopen。返回 Drizzle statements,调用方放进 db.batch() 原子执行
POST /v2/tasks/:taskId/progressstudio-api 新端点

applyTaskAction() 并发/幂等 contract(每 action SQL 层显式)

ActionDB 幂等 / 并发约束备注
create_openuq_tasks_pending_contact_category partial unique on (contactPhone, storeId, typeCategory) WHERE status='open'(现有 index 的 status='pending' 改成 'open'防同 contact 同 typeCategory 重复 open task
create_closedINSERT 行必带 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
closeUPDATE ... SET status='closed' WHERE task_id=$1 AND status='open' AND updated_at <= $aiRunStartedAt RETURNING * — 0 rows = rejecthuman authority guard:staff 已动过则 AI 提案 reject
updateUPDATE ... SET ... WHERE task_id=$1 AND status='open' AND updated_at <= $aiRunStartedAt RETURNING *同 close
record_progressINSERT INTO task_progress_events ... ON CONFLICT (idempotency_key) DO NOTHING + UPDATE tasks SET attempt_count = attempt_count + 1, due_at = $newDueAt WHERE task_id=$1idempotency_key 而不是 (task_id, call_id, progress_type) ON CONFLICT — 因为 PostgreSQL NULL ≠ NULL,SMS progress 无 callId 时复合 key 不去重。idempotency_key 生成规则见下
reopenUPDATE ... 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 使用分层的 commandIdeventIdempotencyKeysourceInteractionKey,见 Task V3 Target Design Proposal §9.4

Step 3: 新建 Timeline Writer + Event Catalog

新建项位置
Timeline Writercallytics-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 Writercallytics-common/src/domain/contact-writer.ts。Phase 1 最小版:upsertIdentity()(trust scoring)/ touchActivity()(forward-only lastActivityAt)/ setDNC()。lifecycle guard 先 observe-only 积累数据,Phase 2 再 enforce
Policy Guardcallytics-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 行为
allowOrchestrator 执行 SQL
rejectlog + skip。低置信度提案也走这条(不引入 needs_review queue/table)。原因:阈值与人工复核流程尚未 calibrate,Phase 1 先用 CloudWatch metric 收数据(mutation_rejected_total{reason=...}),Phase 2 再决定要不要建 needs_review 队列、谁审、SLA
needs_reviewtype 保留为 contract 第三态,Phase 1 不走这条分支。Phase 2 接 needs_review 实现时不用改 type,只改 dispatch 逻辑

Step 5: 4 个 caller 逐个迁移

顺序Caller现在Phase 1
5astudio-api closeraw SQL + 关旧建新(已写 timeline)Task Orchestratorno_answer/left_voicemail 改走 record_progress。timeline 写入由 Orchestrator + Timeline Writer 内部完成(不再 caller 自己拼 buildContactTimelineInsertSQL
5bcontacts-analyzer遍历 taskDecisions[] 自己拼 SQLtaskDecisions.map(d => applyTaskAction(d))
5clead-processor直接 db.insert(tasks)applyTaskAction({ action: 'create_open' })
5dmessage-processor STOP直接拼 DNC cascade SQLContactWriter.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 isolationcontacts PK 已含 store_id;tasks.store_id 需改 NOT NULL
RBACstudio-api 有 store-level 校验,需泛化
Audit trailTimeline 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

#冗余场景真冗余?现有保护
1per_call 已分析 → 当晚 cron 又分析一次是(故意兜底,实现可优化)❌ cron 查询不看 last_contact_analysis_at
2一天打 N 通电话 → per_call 触发 N 次是,最浪费❌ 完全无 cooldown(dedup id 含 callId,FIFO 去重失效)
3cron 同天对同 contact 重复发理论/轻微✅ FIFO dedup(5min 窗口)
4reprocess 撞上 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 Analysisanalyze_call(callId, mode)transcript / per-call AI / call 分类task 创建
SMS Signalevaluate_sms(messageId, mode)exact STOP / trivial filter / meaningful signal直接写 task
Contact Analysisreanalyze_contact(phone, storeId, reason)聚合跨互动上下文 + AI 判断直接 DB mutation
Lead Downstreamprocess_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
Writer 清单(contacts 7 / tasks 10 / timeline 15+)Opus §2.2
Contact Writer 完整 contractOpus §3.2.2
Timeline Writer Event Catalog(12 种 eventType)Opus §3.2.3
Code vs AI 完整职责矩阵Opus §3.4 + Task Pipeline Deliverable §7
Reuse Scenario Validation(4 个场景)Opus §4.5
Business Object 语义定义Codex §Ideal Business Object Semantics
Target AI Pattern + AI never owns 清单Codex §Target AI Pattern
加新 workflow 的体验定义Codex §What Adding a New Workflow Should Feel Like

Open Questions

问题背景 / 决策状态
Task Orchestrator 放 callytics-common 还是 callytics-infrastructurestudio-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 + 删 indexactionNeededReason / 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

AreaDeltaNotes
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_idnullable → NOT NULL测试环境 backfill NULL 老行(反查 PhoneStoreAssignments)或直接 DELETE;ALTER COLUMN SET NOT NULL;删 idx_tasks_store_idWHERE 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.actionNeededDROP COLUMN派生公式 (status === 'open'),无信息量。删字段 + 删 idx_tasks_action_needed(如有)
tasks.taskTypeDROP COLUMNtypeCategory 重复,已有 chk_tasks_task_type_category_consistency CHECK 锁定
contact_timeline新增 / 规范 event catalog至少覆盖 task created / updated / closed / progress recorded;payload 需要 Zod schema
tasks.closeResult18 → 15 values移出 4 个 progress 值(no_answer / left_voicemail / callback_later / attempted)到 task_progress_events.progressType;新增 unable_to_reach
tasks.closeType2 → 3 values新增 create_closed
contacts.actionNeededPhase 1 改读路径,Phase 2 DROP COLUMNPhase 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
contacts AI 画像字段Phase 1 不迁移全部 18 个只收口 identity / DNC / trust / lastActivityAt;AI profile fields 留后续

API / Contract Delta

AreaDeltaNotes
Studio API新增 POST /v2/tasks/:taskId/progressno_answer / left_voicemail 走 progress,不再 close + recreate
Existing task routesclose / reopen / postpone 内部改调 TaskOrchestrator外部 API 尽量兼容,内部统一 policy / idempotency / timeline
Contact Analyzer promptTaskDecisionrecord_progress / create_closed先兼容旧 create/update/close,再改 prompt schema
Shared contractsTaskAction / ContactAction / TimelineEventPayload / ProcessingIntent放在 shared domain package,避免归属于某个 Lambda
Tool callingPhase 1 不暴露 tool runtime只保证 shared modules 未来可被 thin tool handler 调用

Test / Rollout Delta

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