Unified Pipeline Target State (Opus)
Path:
docs/product-design/v2/unified-pipeline/unified-pipeline-target-state-opus.mdDate: 2026-05-31 Author: Opus session(基于 canonical brief + 代码调研) Status: Design proposal -- 待 Max review前序文档:
- Unified Pipeline 设计 Brief -- canonical brief(本 deliverable 的任务书)
- Task Pipeline Deliverable (Codex) -- 已讨论通过的 task 设计
- Pipeline 全景分析 -- 代码级 trace
- Name Trust Scoring -- contacts 名字信任层级定义
- Task Enums -- task 枚举值清单
- Store-Level Isolation -- 唯一隔离键 =
store_id
1. Executive Summary
retaintive 当前有 5 条 pipeline(call analysis、contact analysis、SMS、lead、analytics)独立写入 3 张共享 Neon 表:contacts(7 条写入路径)、tasks(10 条写入路径)、contact_timeline(15 条写入路径)。每条 pipeline 各自实现 UPSERT/UPDATE 逻辑,trust score 判定、name COALESCE 策略、lifecycle 流转规则散落在不同 Lambda 里,没有统一的 mutation 层。结果是同一张表的写入行为因调用方不同而不一致 -- 例如 lead-processor 做 UPSERT 时不带 trust score,而 transcribe-processor 用 trust=60,contacts-analyzer 用 trust=40/80,studio-api 硬编码 trust=100。这不是"代码风格不统一",是业务语义在 5 个地方各说各话。
Phase 1 的目标是建三个 shared mutation module:Task Orchestrator(统一 task 的创建、关闭、reopen、状态流转)、Contact Writer(统一 contacts 的 UPSERT 策略、trust-based field resolution、lifecycle 流转)、Timeline Writer(统一 contact_timeline 的 event 写入和 payload schema 校验)。加一个 Policy Guard 做跨表约束(DNC cascade、lifecycle 合法性检查)。所有 pipeline 调用这些 module 而不是直接拼 SQL。
目标架构分两层:
- Layer 1(Shared Mutation) -- 上述四个 module,定义"怎么写表":trust 仲裁、状态机合法转移、DNC 级联、timeline payload schema
- Layer 2(Processing Capability) -- 各 pipeline 的业务逻辑,定义"写什么内容":AI 分析结果、SMS 信号解析、lead 分配规则。Layer 2 调用 Layer 1,不直接碰表
关于是否现在切换到 tool calling:不。当前 AI 调用模式是 one-time structured output(单次调用,Zod schema 约束返回),成功率 95-99%,满足 Phase 1 需求。Tool calling 适合交互式场景(AI 需要多轮查询再决策),是 Phase 2 的事。关键洞察是 shared mutation module 做好了,tool calling 的 backend 实现自然就有了 -- future tool 的 execute 函数就是调用同一套 Task Orchestrator / Contact Writer,不需要重写。先把 mutation 层统一,再考虑调用方式的演进(详见 Section 5)。
2. Current State Inventory
2.1 Business Object 健康检查
系统有 6 个核心 business object。下表逐一审视每个对象的定义清晰度和已知问题。
2.2 Writer 清单
contacts 表(7 个 Writer)
关键冲突:lastActivityAt 4 个 writer 用 4 种语义(2 个 NOW() 直接覆盖、1 个 COALESCE(GREATEST(...)) forward-only、1 个 startTime ?? NOW())。firstName/lastName 只有 lead-processor 用 COALESCE 防覆盖,其余直接覆盖。
tasks 表(10 个 Write Path)
关键冲突:3 个独立入口各自用不同的去重策略、closeResult 枚举子集、dueAt 计算逻辑。没有统一的 Task Orchestrator。
contact_timeline 表(15+ 个 Write Path)
所有 writer 通过 buildTimelineValues() 或 buildContactTimelineInsertSQL() 写入。INSERT-only(append-only audit log),ON CONFLICT DO NOTHING 用 idempotencyKey 去重。
关键问题:occurredAt 两种策略(源记录时间 vs NOW())逻辑正确但无 shared enforcement。studio-api 的 close 操作不写 contact_timeline(只写 task_events 审计表),两套审计系统并存。
2.3 Prompt 清单
系统有 6 个 prompt surface,分布在 2 条 AI pipeline 中。
核心分工:Pipeline 1(Call Analysis)的 4 个 prompt surface 只产信号(写 calls 表),不直接操作 task 或 contact lifecycle。Pipeline 2(Contact Analysis)是唯一的 AI -> 状态变更 入口。这是 "AI proposes, code executes" pattern 的现状实现 -- 但只覆盖了 contacts-analyzer 路径,lead-processor 和 studio-api 绕过了这个 pattern。这个分工在 target state 下保持不变(见 Section 3.4)。
2.4 Current Workflow 全景图(只画写入路径)
只画"谁写了什么共享表",不画 read。Pipeline 内部的处理步骤(triage/classify/coaching 等)见 §2.3 Prompt Inventory。
核心问题:
- contacts 有 6 个写入方,trust score 逻辑各写各的(60 / 40 / 无 / 100),没有共享 enforcement
- tasks 有 5 个写入方(lead-processor、contacts-analyzer、studio-api 4 种操作、DNC cascade),去重策略、closeResult、dueAt 计算全不一样
- 每个写入方都直接碰共享表 — 没有统一的 mutation 入口
- studio-api 的"没打通 → 关旧建新" 是把 progress 当 close 处理
3. Target State Design
3.1 Two-Layer Architecture Overview
统一 pipeline 架构把系统分成两层。Layer 1 是 shared mutation modules -- 所有写入 contacts、tasks、contact_timeline 的操作必须经过的唯一入口。Layer 2 是 processing capability modules -- 各条 pipeline 的业务处理逻辑,产出结构化提案后交给 Layer 1 执行写入。
核心关系:Layer 2 产出 proposal,Layer 1 执行 mutation。
这和现状的根本区别:现在 lead-processor 直接 INSERT tasks、contacts-analyzer 通过 writeAnalysisWithTasks() 写 tasks、studio-api 通过 /v2/tasks/close 写 tasks -- 三条路径各自实现去重、closeResult 校验、dueAt 计算、timeline 审计,行为不一致(见 Section 2.2)。Target state 下这三条路径都改为调用同一个 Task Orchestrator。
员工 UI 操作和 future AI agent 直接调用 Layer 1(不经过 Layer 2),因为它们本身不需要 AI 判断。
只画写入路径(跟 §2.4 对比看变化)。Pipeline 内部的处理步骤不展开。
Timeline Writer 被两类调用方调用:
- Pipeline 直接调 — 写 pipeline 自身的审计事件(
transcribe.completed、message.created、lead.created)。这些事件跟 task/contact mutation 无关,是"发生了什么"的记录 - Task Orchestrator / Contact Writer 内部调 — 写 mutation 引起的事件(
task.created、task.status_changed、contact.lifecycle_changed)。这些事件是 mutation 的副作用
跟 §2.4 现状图的对比:
3.2 Shared Mutation Module Contracts
3.2.1 Task Orchestrator
统一的 task 写入入口。所有 task 状态变更 -- 无论来源是 AI proposal、lead webhook、员工 UI、future tool calling -- 都经过这一个 contract。
Input:
Payload 类型:
Guarantees:
Callers:
3.2.2 Contact Writer
统一的 contacts UPSERT 入口。解决 6 个 writer 各自实现 name trust scoring 的问题(trust 层级定义见 name-trust.md)。
Input:
Guarantees:
Callers:
3.2.3 Timeline Writer
统一的 contact_timeline INSERT 入口。
Input:
Event Type Catalog(每种事件的 payload 有 Zod schema):
Guarantees:幂等(idempotencyKey unique constraint)、payload schema 校验(Zod discriminated union)、INSERT-only(append-only audit log)。不改 DB 结构(oldValue/newValue 仍是 JSONB),类型安全在 application 层保证。
3.2.4 Policy Guard
前置拦截层。在 Task Orchestrator / Contact Writer / Timeline Writer 执行写入之前检查操作是否合规。
Layer 2 module 不直接调用 Policy Guard -- 它们调用 Layer 1,由 Layer 1 内部触发 Policy Guard。
3.3 Processing Capability Module Contracts
3.3.1 Call Analysis Module
位置:callytics-infrastructure/lambda/ai-analysis-processor/
职责:分析单通电话,产出 call facts,写入 calls 表。不直接操作 tasks -- prompt 三处明令禁止(prompts.ts:67,197,218)。
调用 Layer 1:Contact Writer(upsertContact({source: 'AI_TRANSCRIPT', trust: 40/80, ...}))+ Timeline Writer(call.created)。follow_up_needed = yes 时向 dailyBatchQueue.fifo 发 SQS 消息触发 Contact Analysis Reprocess。
3.3.2 SMS Signal Module
位置:callytics-infrastructure/lambda/message-processor/
职责:接收 SMS/VM,持久化消息记录,代码层做轻量信号分类。不调用 AI。
3.3.3 Contact Analysis Reprocess
位置:callytics-infrastructure/lambda/contacts-analyzer/
职责:聚合跨通话/跨消息/跨 lead 信号,AI 产出 taskDecisions[] + 18 个 contact 分析字段。AI 输出的 taskDecisions[] 逐条转为 applyTaskAction() 调用,18 个 contact 分析字段通过 upsertContact() 写入。每个 task mutation 由 Task Orchestrator 内部自动写 timeline。
关键约束:task_close 触发的重分析可能再次建 task -- Task Orchestrator 的 duplicate guard 是打断无限循环的安全网。
3.3.4 Lead Downstream
位置:lead-tracking/ repo
职责:新 lead 到达后,确定性地创建 contact + task + timeline。不调用 AI。直接调用 Task Orchestrator(applyTaskAction({ action: 'create', ... })),不需要单独的 module。所有逻辑是确定性的,Task Orchestrator 已覆盖 task 部分。
3.4 Code vs AI Responsibility Framework
核心原则:AI proposes, code executes。 AI 只输出结构化 proposal(Zod 验证,success rate 95-99%),代码负责触发控制、输入限定、输出校验、状态转换、事务/幂等、审计、安全边界。
4. Phase 1 Implementation Design
Phase 1 按产品可见度和依赖关系排序:Task Orchestrator 先行(task deliverable 已 approved + dashboard/前端直接消费),Timeline Writer 同步跟进(task 事件必须写 timeline),Contact Writer 随后(identity/DNC 逻辑相对独立),Invocation Contracts 只做设计不实现。
4.1 Task Orchestrator(第一优先)
要建什么
applyTaskAction -- 一个 pure function module(不是 Lambda),住在 callytics-common/src/domain/task-orchestrator.ts。接收 action 描述,返回一组 Drizzle statements(INSERT/UPDATE + timeline INSERT),由调用方放进 db.batch() 原子执行。
各 caller 怎么改
lead-processor(persist-downstream.ts:190-216):现在直接 db.insert(tasks).values({...}).onConflictDoNothing() + buildTimelineValues() 拼 timeline。改后调用 applyTaskAction({ action: 'create', params: { typeCategory: 'lead_outreach', priority: 'high', sourceType: 'lead', dueAt: computeLeadSLA() } }),返回的 statements 塞进现有的 db.batch()。dueAt 计算(含 quiet hours)、幂等约束、timeline 写入全由 Orchestrator 统一处理。
contacts-analyzer(neon-repository.ts:540-860):现在遍历 taskDecisions[],per-decision 拼 SQL。改后 taskDecisions.map(d => applyTaskAction(d, context)) -> flatMap 所有 statements -> 跟 contacts UPDATE 一起塞进 db.batch()。create/close/update 的硬编码逻辑集中到 Orchestrator。
studio-api(close.ts:105-228):现在 raw SQL UPDATE + 独立 timeline + 独立 auto-follow-up。改后调用 applyTaskAction({ action: 'close', ... })。关键变化:no_answer / left_voicemail 不再走 close + recreate,改为 record_progress(见下方)。Gap 补齐:close 时补写 contacts.lastActivityAt + SQS 触发 contacts-analyzer(source: 'task_close')。
record_progress:新 action type
根问题:员工拨打电话没打通,当前 close.ts:199 逻辑是关闭旧 task + 建新 task。本质是把 progress 事件("这次没打通")编码为 close + create。
上游调研数据支持这个改动:
manual_closed= 0(682 条 closed task 全是 system auto-close)-- close.ts 的"员工手动关闭"路径从未走通- disposition 类 closeResult 仅占 2.3%(16/682)
- disposition 是 per-call 粒度事实,calls 表已有 20,364 条记录(tasks 表的 12 倍)-- 跨粒度错误
Contract:
对 closeResult 枚举的清理:4 个 progress 值(no_answer, left_voicemail, callback_later, attempted)从 TASK_CLOSE_RESULT 移出。attempted(占 55%,AI 兜底值)语义分离后不再需要 -- 无法判断 outcome 时不 close。剩余 14 个值全是真 outcome。详见 task-enums.md。
Progress events 走 contact_timeline,不新建表
不建 task_progress_events 表。contact_timeline 已有完整的 typed event 基础设施,新增 task.progress_recorded eventType 即可:
查询"这个 task 被尝试了几次" = SELECT COUNT(*) FROM contact_timeline WHERE entity_id = $taskId AND event_type = 'task.progress_recorded'。
4.2 Timeline Writer + Event Catalog(与 Task Orchestrator 同步)
当前有两套 builder:buildTimelineValues()(返回 Drizzle row object,Lambda 用)和 buildContactTimelineInsertSQL()(返回 raw SQL,studio-api 用)。两套 builder 已共享同一个 BuildTimelineValuesParams interface。
Timeline Writer 要做的:不替换这两个 builder(它们是不同 DB client 的适配),而是在上层封装一个 TimelineWriter,让调用方不再直接拼参数:
TypedTimelineEvent 是 discriminated union,每种 eventType 的 payload 有明确的 TypeScript type。Event Catalog 把每种事件的 payload 用 Zod schema 定义(schema 见 Section 3.2.3 的 Event Type Catalog 表)。
4.3 Contact Writer(下一步)
Phase 1 只统一 4 条 UPSERT/UPDATE 路径中的共同字段:
identity + name(trust score enforcement):ContactWriter.upsertIdentity(phone, storeId, { firstName, lastName, trustScore }) -- 内部实现 trust score 比较,所有 caller 走同一份 SQL。
lastActivityAt:ContactWriter.touchActivity(phone, storeId, occurredAt) -- 统一用 forward-only 语义(GREATEST(existing, newTime)),防止乱序事件回退时间。
DNC:ContactWriter.setDNC(phone, storeId, doNotContact, updatedBy) -- dnc-cascade.ts 调用 Contact Writer 而不是直接拼 SQL。
lifecycle / leadStatus guard:Phase 1 加 observe-only guard -- Contact Writer 在写入前 log 状态转换,积累数据。Phase 2 再改为 enforce。理由:lifecycle 状态机的合法转换规则当前只在 prompt 文字里定义,没有代码 enforcement,先 observe 再设计。
Contact Writer 不做的:不统一 AI 画像字段(customerSummary 等 18 个)-- 只有 contacts-analyzer 写,不存在多 writer 竞争。
4.4 Invocation Contracts(Phase 1 只设计)
定义 4 个 Processing Capability Module 的调用 contract,Phase 1 不实现 module 本身。
4.5 Reuse Scenario Validation
用 Phase 1 的 3 个 module 验证 4 个跨 pipeline 场景。
Scenario 1: Call retry(员工连打 3 次没人接)
- Call 1 (
no_answer): ai-analysis-processor 写 call facts -> SQS -> contacts-analyzer AI 建 task viaTask Orchestrator.create() - Call 2 (
no_answer): contacts-analyzer AI 看到 pending task 已存在 ->Task Orchestrator.update()更新 suggestedActions(uq_tasks_pending_contact_category保证不重复建) - 员工 UI 标记"没打通":
Task Orchestrator.record_progress({progressType: 'no_answer'})-> dueAt 推后 1 天,task 保持 pending - Call 3(接通, booked): contacts-analyzer AI 判断 outcome ->
Task Orchestrator.close({closeResult: 'booked'})
验证点:disposition 走 progress event,outcome 走 close。同一个 task 全程 pending 直到有真实 outcome。
Scenario 2: SMS "STOP"(DNC 级联关闭)
- Contact Writer.setDNC(true) -> UPDATE contacts
- DNC cascade: 每个 pending task 调
Task Orchestrator.close({closeResult: 'do_not_contact'})-- 统一的 close 逻辑 - Timeline Writer:
contact.dnc_changed+message.created
验证点:DNC cascade 通过 Task Orchestrator 关闭,保证 close 逻辑跟 AI auto-close / staff manual-close 完全一致。
Scenario 3: New store onboard(批量历史数据分析)
每个 contact 的 reprocess 走现有 contacts-analyzer -> Task Orchestrator 路径。reprocessRunId 通过 uq_tasks_analysis_run_category 保证不重复建 task。Contact Writer 的 trust score 保证不覆盖已有的高信任名字。
Scenario 4: Prompt update batch reanalysis
走现有 contacts-analyzer -> Task Orchestrator 路径。AI 可能关旧 task + 建新 task + 更新 priority,全部经过 Orchestrator 的幂等约束和 duplicate guard。不需要特殊逻辑。
5. One-time AI Call vs Tool Calling 兼容性
5.1 当前模式:one-time AI call + code orchestration
contacts-analyzer 是最完整的例子:1 次 AI 调用产出 ContactsAnalysis(含 taskDecisions[]),代码层遍历 decisions 执行 INSERT/UPDATE/timeline。AI 不直接操作数据库,不选择调用哪个函数。适合 batch/async 场景。
5.2 未来模式:tool calling(interactive)
Tool calling 让 AI 主动选择调用哪个 function:voice agent(实时决定查 task history)、manager Q&A(诊断 lead 为什么没有 task)、staff assistant(批量关闭过期 task)。
关键区别:one-time call 是 code 编排 AI output,tool calling 是 AI 编排 code execution。
5.3 Shared module 在两种模式下的工作方式
Contact Writer、Timeline Writer、DNC Cascade 同理 -- tool handler 是 thin wrapper,调用同一个 shared module。两种模式共存:同一个 store 可以同时有 Mode A(cron 跑 contact analysis)和 Mode B(manager 通过 chat 问诊断)。Shared module 保证两种模式对数据的修改逻辑一致。
5.4 决策标准
5.5 Tool calling 前置条件(Phase 1 不实现)
Migration path:Phase 1 建好 shared module 后,Phase 2 加 tool calling 只需要:定义 tool schema(复用 Zod types)-> 写 tool handler(thin wrapper)-> 加 RBAC + budget cap -> 注册 tools 到 AI runtime。不需要改 shared module 本身。
6. Open Questions for Max
架构决策
-
Task Orchestrator 的代码位置:放
callytics-common/src/domain/task-orchestrator.ts(所有 repo 共享)还是callytics-infrastructure/lambda/shared/domain/(只 Lambda 用)?取决于 studio-api 是否能直接依赖callytics-common(当前 studio-api 用 raw SQL 不用 Drizzle ORM)。如果放callytics-common,studio-api 需要一个 raw SQL adapter 层。 -
contacts.actionNeeded/suggestedActions何时清理:schema 注释说 tasks 已取代 contacts 级别的 action 字段,但 contacts-analyzer 仍然同时写两处。Phase 1 是否正式标记这两个字段为 deprecated?前端 dashboard 是否还在读 contacts 级别的 actionNeeded? -
neon-sync Lambda:确认为死代码后,是否在 Phase 1 直接删除?
record_progress 设计决策
-
attempted迁移策略:现有 378 条closeResult='attempted'的 closed task,是否需要数据迁移?还是只影响新数据,历史数据保持原样? -
前端 UI 改动范围:
record_progress需要 studio-api 新端点 + 前端把no_answer/left_voicemail/callback_later从"关闭 task"弹窗移到"记录进度"按钮。这个前端改动是否在 Phase 1 scope 内?还是先只改 backend contract,前端保持现有 close 流程(Orchestrator 内部把这 3 个值自动转为 record_progress)?
实施顺序
-
Phase 1 的 PR 拆分:建议分 4 个 PR -- (a) Task Orchestrator + tests (b) Timeline Writer + Event Catalog (c) contacts-analyzer 迁移到 Orchestrator (d) lead-processor + studio-api 迁移。是否同意?
-
closeResult 枚举三方不同步(schema 18 / prompt 13 / 前端 UI 16):Phase 1 是否顺带对齐这三方?还是只清理 progress 值,其他不同步留 #199 追踪?
Policy Guard scope
-
Rate limit 的具体阈值:Policy Guard 的
create_taskrate limit -- 具体 N 分钟 M 次的值?还是 Phase 1 先不加 rate limit,只加 DNC guard + storeId guard + duplicate guard? -
lifecycle guard observe-only 的持续时间:Contact Writer 的 lifecycle 状态转换 observe-only guard 运行多久后转 enforce?