For AI agents: the complete documentation index is available at /llms.txt, the full documentation bundle is available at /llms-full.txt, and this page is available as Markdown at /product-design/v2/unified-pipeline/original/unified-pipeline-target-state-opus.md.

Unified Pipeline Target State (Opus)

Path: docs/product-design/v2/unified-pipeline/unified-pipeline-target-state-opus.md Date: 2026-05-31 Author: Opus session(基于 canonical brief + 代码调研) Status: Design proposal -- 待 Max review

前序文档:


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。下表逐一审视每个对象的定义清晰度和已知问题。

Business Object一句话定义定义清晰度具体问题
Contact一个电话号码在一个门店下的客户实体,是所有通话、消息、lead、task 的聚合挂载点混合contacts.actionNeeded/suggestedActions 与 tasks.actionNeeded/suggestedActions 语义重复。schema 注释明确说 tasks 已取代 contacts 级别的 action 字段(tasks.ts 行 77),但 contacts-analyzer 仍然同时写两处。contacts 级别的 actionNeeded 是 stale snapshot -- 员工关了 task 但 contacts.actionNeeded 仍为 true,直到下一次 contact analysis 才刷新。6 个 writer 写 contacts 表,name trust scoring(WEBHOOK=20 / AI_TRANSCRIPT=40 / RC_API=60 / LEAD=80 / STAFF=100)分散在 5 个 Lambda 代码库,无共享 enforcement
Task员工需要对某个 contact 执行的待办事项,由 AI 或代码创建,员工或 AI 关闭混合task 同时被当作"待办"和"业绩记录",但没有字段区分这两个语义。closeResult 18 个值混了 3 层含义:progress(no_answer/left_voicemail/callback_later)、disposition(attempted 占 55% 是兜底值)、business outcome(converted/booked/cancelled 等)。Neon 实测 682 条 closed task 全部 closeType='auto_closed'、closedByStaffName='system',manual_closed 为 0。closeResult 枚举在 schema(18 值)、prompt(13 值)、前端 UI(16 值,缺 booked/cancelled)三方不同步。tasks 表有 4 个 CHECK 约束,但 contacts 表 0 个 -- 防护等级不对称
Call一次 RingCentral 通话的完整记录,包含物理状态、业务结果和行动信号清晰call 有干净的 3 层分离设计:callState 物理层 / primaryOutcomeResult 业务层 / followUpNeeded 行动层。问题在跨表映射:callState 已记录 disposition(20,364 条 calls),但 task.closeResult 又重复记录同样的 disposition(仅 16 条,2.3%)-- 这是跨粒度错误,disposition 是 per-call 事实,不应出现在 per-objective 的 task 层
Lead从 email 解析出的潜在客户信息,是 contact 生命周期的起点之一清晰leads 表是单 writer(lead-tracking 独占),PK 是 composite dedup key。lead-tracking 创建的 lead_outreach task 全部硬编码(priority='high'、dueAt=5min SLA),不经 AI 判断 -- 正确设计(速度优先),但绕过了 Task Orchestrator 统一入口
Message一条 RingCentral SMS 或 Voicemail 消息清晰messages 表是单 writer(message-processor 独占)。message-processor 不触发 contacts-analyzer(2026-03-27 设计决策)。客户发 "STOP" 要等到次日 cron 才能被 AI 看到并标 DNC
Timelinecontact 维度的 append-only 审计事件流中等有 16 种 eventType 和 typed eventCategory,但 newValue/oldValue 是 JSONB convention,没有 per-eventType 的 payload schema。4 个 pipeline 各自定义 event_type 命名和 payload 结构。studio-api 关闭 task 时不写 contact_timeline(只写 task_events 审计表),两套审计系统并存

2.2 Writer 清单

contacts 表(7 个 Writer)

WriterLambda/ServiceOperationTrust ScoreDedup 策略关键行为差异
transcribe-processorcallytics-infrastructureUPSERT60 (RC_API)ON CONFLICT (phone, store_id) DO UPDATElastActivityAt = NOW()(直接覆盖,不防回退);firstName/lastName 直接覆盖
ai-analysis-processorcallytics-infrastructureUPSERT40/80ON CONFLICT (phone, store_id) DO UPDATElastActivityAt = input.timeline.startTime ?? NOW()(直接覆盖);firstName/lastName 直接覆盖
message-processorcallytics-infrastructureUPSERTN/A(不写 name)ON CONFLICT (phone, store_id) DO UPDATElastActivityAt = COALESCE(GREATEST(existing, eventTime), eventTime)(唯一防回退的 writer)
lead-processorlead-trackingUPSERTN/AON CONFLICT (phone, store_id) DO UPDATElastActivityAt forward-only;firstName = COALESCE(existing, new)(只补写不覆盖);仅 INSERT 时设 lifecycle
contacts-analyzercallytics-infrastructureUPDATE only40 (AI)无 UPSERT -- contact 必须已存在写 18 个 AI 字段 + lastContactAnalysisAt;DNC sticky 逻辑(不可 true->false)
studio-apistudio-website-monorepoUPDATE100 (STAFF)无 UPSERT只写 firstName/lastName(员工手动修正);raw SQL 非 Drizzle
neon-synccallytics-infrastructureLEGACYN/A旧 schema(无 storeId)死代码 -- schema 与现状不兼容,部署即崩

关键冲突:lastActivityAt 4 个 writer 用 4 种语义(2 个 NOW() 直接覆盖、1 个 COALESCE(GREATEST(...)) forward-only、1 个 startTime ?? NOW())。firstName/lastName 只有 lead-processor 用 COALESCE 防覆盖,其余直接覆盖。

tasks 表(10 个 Write Path)

WriterLambda/ServiceOperation关键行为差异
lead-processor CREATElead-trackingINSERT全硬编码:typeCategory='lead_outreach'/priority='high'/dueAt=5min SLA/sourceType='lead';不经 AI
contacts-analyzer CREATEcallytics-infrastructureINSERTAI 决定 typeCategory/priority;代码硬编码 status='pending'/sourceType='contact_analysis'
contacts-analyzer CLOSEcallytics-infrastructureUPDATEAI 选 closeResult(11 值枚举);代码硬编码 closeType='auto_closed'/closedByStaffName='system'
contacts-analyzer UPDATEcallytics-infrastructureUPDATEAI 更新 priority/suggestedActions/dueAt;priority 变更触发 dueAt 重算
DNC cascade CLOSEcallytics-commonUPDATEAI 判断 DNC=true 时自动关闭所有 pending tasks;closeResult='do_not_contact'
studio-api CLOSEstudio-website-monorepoUPDATE员工 UI 选 closeResult;closeType='manual_closed';raw SQL;不写 contact_timeline
studio-api auto-follow-upstudio-website-monorepoINSERT员工关 task 时 no_answer/left_voicemail 自动建新 task;双重死代码(manual_closed=0 + disposition 类 closeResult 仅 2.3%)
studio-api REOPENstudio-website-monorepoUPDATE员工重开已关闭 task
studio-api NOTEstudio-website-monorepoUPDATE只改 note 字段
studio-api POSTPONEstudio-website-monorepoUPDATE只改 due_at;端点未实施

关键冲突:3 个独立入口各自用不同的去重策略、closeResult 枚举子集、dueAt 计算逻辑。没有统一的 Task Orchestrator。

contact_timeline 表(15+ 个 Write Path)

所有 writer 通过 buildTimelineValues() 或 buildContactTimelineInsertSQL() 写入。INSERT-only(append-only audit log),ON CONFLICT DO NOTHING 用 idempotencyKey 去重。

WritereventTypeoccurredAt 策略
lead-processorlead.created源记录时间(lead.receivedAt)
transcribe-processorcall.status_changedRC event time
ai-analysis-processorcall.created源记录时间(call.startTime)
message-processormessage.created源记录时间(message.creationTime)
contacts-analyzertask.created / task.status_changed / task.updated / contact.lifecycle_changed / contact.temperature_changedNOW()(transaction 时间)
studio-apitask.status_changed / task.due_at_changedNOW()

关键问题:occurredAt 两种策略(源记录时间 vs NOW())逻辑正确但无 shared enforcement。studio-api 的 close 操作不写 contact_timeline(只写 task_events 审计表),两套审计系统并存。

2.3 Prompt 清单

系统有 6 个 prompt surface,分布在 2 条 AI pipeline 中。

PromptPipeline是否触发状态变更核心行为
Pre-TriageCall Analysis间接 -- 命中 keyword 规则时 short-circuit无 AI,纯代码规则(pre-triage.ts keyword matching),拦截约 40% 的 triage 成本
TriageCall Analysis间接 -- isWorthAnalyzing=false 跳过后续 stageAI 闸门,只传前 500 字符
ClassifyCall Analysis间接 -- followUpNeeded=yes 触发 SQS fan-out完整 transcript AI 分析,产出 33 个 per-call 字段,Zod schema 校验
VerifyCall Analysis否 -- 仅修正 subcategory条件触发(仅 CONFUSION_PAIRS),Classification 后校验层
CoachingCall Analysis否 -- 只写 calls 表 coaching 字段条件触发(human_conversation && duration>30s)
Contact AnalyzerContact Analysis是 -- 直接触发多表状态变更AI 产出 taskDecisions[] + 18 个 contact 字段。代码层 8 重安全网:Zod validation -> unique index x2 -> CHECK constraint -> WHERE guard -> pre-batch code dedup -> ON CONFLICT DO NOTHING -> FIFO SQS 串行

核心分工: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。

核心问题:

  1. contacts 有 6 个写入方,trust score 逻辑各写各的(60 / 40 / 无 / 100),没有共享 enforcement
  2. tasks 有 5 个写入方(lead-processor、contacts-analyzer、studio-api 4 种操作、DNC cascade),去重策略、closeResult、dueAt 计算全不一样
  3. 每个写入方都直接碰共享表 — 没有统一的 mutation 入口
  4. 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 被两类调用方调用:

  1. Pipeline 直接调 — 写 pipeline 自身的审计事件(transcribe.completed、message.created、lead.created)。这些事件跟 task/contact mutation 无关,是"发生了什么"的记录
  2. Task Orchestrator / Contact Writer 内部调 — 写 mutation 引起的事件(task.created、task.status_changed、contact.lifecycle_changed)。这些事件是 mutation 的副作用

跟 §2.4 现状图的对比:

现状目标
6 个写入方各自直接写 contacts / tasks / timeline所有写入方先产出结构化提案,交给共享写入层
trust score 逻辑散落 5 个 LambdaContact Writer 统一处理 trust 仲裁
task 去重/close/dueAt 各写各的Task Orchestrator 统一处理状态机
timeline event_type 各自定义Timeline Writer 统一 event catalog
没打通 → 关旧 task + 建新 taskrecord_progress 新 action,task 保持 open
DNC cascade 是唯一的共享逻辑Policy Guard 集中所有规则(DNC、storeId、权限)
没有 task_progress_events 表新增 task_progress_events,记录每次尝试

3.2 Shared Mutation Module Contracts

3.2.1 Task Orchestrator

统一的 task 写入入口。所有 task 状态变更 -- 无论来源是 AI proposal、lead webhook、员工 UI、future tool calling -- 都经过这一个 contract。

Input:

applyTaskAction({
  action: 'create' | 'update' | 'close' | 'record_progress',
  identity: {
    contactPhone: string,
    storeId: string,       // 唯一隔离键, 见 store-level-isolation.md
    accountId: string,
    franchiseId: string,
  },
  actor: {
    type: 'contact_analysis' | 'lead_webhook' | 'staff' | 'system' | 'ai_agent',
    id?: string,           // staffId / analysisRunId / agentId
    name?: string,         // 审计用
  },
  payload: CreatePayload | UpdatePayload | ClosePayload | ProgressPayload,
  idempotencyKey: string,
})

Payload 类型:

actionpayload 字段来源
createtypeCategory, priority, suggestedActions[], sourceType, evidence?, closeResult?(仅 create_closed)AI proposal 或 lead-processor 硬编码
updatetaskId, priority?, suggestedActions?, dueAt?AI proposal 或员工调整
closetaskId, closeResult(仅 business outcome 枚举), closeNote?, closeTypeAI proposal 或员工操作
record_progresstaskId, progressType(no_answer / left_voicemail / text_sent / callback_requested / follow_up_scheduled), channel, callId?, messageId?, note?, nextDueAtHint?AI proposal 或员工操作

Guarantees:

保证实现方式
幂等idempotencyKey -> DB unique constraint;record_progress 额外用 (taskId, callId, progressType) partial unique
原子性单个 db.batch() 写 tasks + contact_timeline,全成功或全回滚
审计每个 mutation 自动生成 contact_timeline event,带 actor 信息
租户隔离identity.storeId 必须非空;close/update 的 WHERE 加 AND store_id = $storeId guard
状态机create 只能产 pending(或 create_closed);close 要求当前 status=pending;record_progress 要求当前 status=pending

Callers:

调用方action说明
contacts-analyzercreate, update, close, record_progressAI taskDecisions[] 逐条转为调用
lead-processorcreate确定性创建 lead_outreach,payload 硬编码
studio-apiclose, record_progress, update员工 UI 操作,actor.type = 'staff'
DNC cascadeclose批量 close 所有 pending tasks,closeResult = 'do_not_contact'
future AI agentcreate, close, record_progresstool calling API 调用同一 contract(见 Section 5.3)

3.2.2 Contact Writer

统一的 contacts UPSERT 入口。解决 6 个 writer 各自实现 name trust scoring 的问题(trust 层级定义见 name-trust.md)。

Input:

upsertContact({
  identity: { contactPhone, storeId, accountId, franchiseId },
  source: 'RC_API' | 'AI_TRANSCRIPT' | 'WEBHOOK' | 'LEAD' | 'STAFF' | 'CONTACT_ANALYSIS',
  trust: 20 | 40 | 60 | 80 | 100,
  patch: {
    firstName?: string,
    lastName?: string,
    lastActivityAt?: Date,
    customerSummary?: string,       // 仅 CONTACT_ANALYSIS
    lifecycleStage?: string,        // 仅 CONTACT_ANALYSIS
    leadStatus?: string,            // 仅 CONTACT_ANALYSIS
    doNotContact?: boolean,
    // ... 其余 AI 字段
  },
  actor: ActorInfo,
  observedAt: Date,   // 事件发生时间, 不是写入时间
})

Guarantees:

保证实现方式
Name trust 一致性CASE WHEN $trust > COALESCE(contacts.name_trust, 0) THEN $firstName ELSE contacts.first_name END,所有 caller 走同一份 SQL
DNC stickydoNotContact = true 一旦设置,仅 STAFF source + trust=100 可以 unstick
租户隔离PK = (phone, store_id),物理上不可能跨 store 写入

Callers:

调用方sourcetrustpatch 范围
transcribe-processorRC_API60firstName, lastName, lastActivityAt
ai-analysis-processorAI_TRANSCRIPT40/80firstName, lastName, lastActivityAt
message-processorRC_API60firstName, lastActivityAt, doNotContact
lead-processorLEAD80firstName, lastName, lifecycleStage, leadStatus
contacts-analyzerCONTACT_ANALYSIS4018 个 AI 分析字段
studio-apiSTAFF100firstName, lastName

3.2.3 Timeline Writer

统一的 contact_timeline INSERT 入口。

Input:

writeTimelineEvent({
  eventType: TimelineEventType,   // 严格枚举, 见 Event Catalog
  eventCategory: 'activity' | 'progress' | 'lifecycle' | 'task' | 'system',
  entity: { type: 'call' | 'message' | 'task' | 'lead' | 'contact', id: string },
  contact: { phone, storeId, accountId, franchiseId },
  actor: ActorInfo,
  payload: Record<string, unknown>,  // typed per eventType via Zod discriminated union
  idempotencyKey: string,
  occurredAt: Date,
})

Event Type Catalog(每种事件的 payload 有 Zod schema):

eventTypeeventCategory新增?payload 必需字段
call.createdactivity否callId, direction, callState
call.status_changedactivity否callId, oldState, newState
message.receivedactivity否messageId, direction, content_preview
message.sentactivity否messageId, direction
lead.createdactivity否leadId, source
task.createdtask否taskId, typeCategory, priority
task.status_changedtask否taskId, oldStatus, newStatus, closeResult?
task.updatedtask否taskId, changedFields
task.progress_recordedprogress新增taskId, progressType, channel, attemptCount
contact.lifecycle_changedlifecycle否oldStage, newStage
contact.dnc_changedlifecycle否doNotContact, reason
contact.analysis_completedsystem否analysisRunId, model

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 执行写入之前检查操作是否合规。

规则检查内容拒绝行为
DNC guardcontacts.doNotContact = true 时禁止 create_task(outreach 类){allowed: false, reason: 'DNC_BLOCKED'}
storeId requiredidentity.storeId 必须非空{allowed: false, reason: 'MISSING_STORE_ID'}
storeId cross-checktask 的 store_id 必须等于 identity.storeId{allowed: false, reason: 'STORE_MISMATCH'}
permission checkactor.type = 'staff' 时校验 store 访问权限{allowed: false, reason: 'PERMISSION_DENIED'}
duplicate guardcreate_task 时同 typeCategory 已有 pending task{allowed: false, reason: 'DUPLICATE_PENDING'}

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。

Signal type调用 Layer 1操作
STOP(exact keyword)Contact Writer + Task OrchestratorsetDNC(true) + 批量 close pending tasks
trivial(ok/thanks/emoji)Contact Writer + Timeline WritertouchActivity + timeline event
meaningful(购买意向/取消/投诉)Contact Writer + Timeline Writer + SQS同 trivial + 触发 Contact Analysis Reprocess

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

决策类型负责层retaintive 具体例子
是否触发 pipelineCodePre-Triage 代码闸门拦截约 40% 的 AI 调用;SMS STOP/trivial/meaningful 由代码分类
Task typeCategory 选择Hybrid: AI 在 code-provided allowed set 内选Code 按 lifecycleStage 算 allowed set,AI 在范围内选
Task priority 判断Hybrid: AI 建议 + code 算 dueAtAI 输出 high/medium/low;computeDueAt(priority) 代码算具体时间
Progress vs Close 判断Code state machineno_answer/left_voicemail 是 progress 不是 close -- 产品语义定义(见 Section 4.1 record_progress)
closeResult 选择Hybrid: AI 在 allowed enum 内选AI 从 allowed set 选,prompt 被禁止输出 progress-like 值
DNC 判断Hybrid: code exact match + AI 自然语言message-processor 代码检测 exact STOP;contacts-analyzer AI 检测自然语言 DNC
Name trust arbitrationCodeContact Writer 的 CASE WHEN $trust > COALESCE(...) SQL
幂等 / 并发 / 事务Code + DBON CONFLICT DO NOTHING、db.batch()、FIFO SQS 串行化
storeId 隔离CodeWHERE 子句的 AND store_id = $storeId guard
suggestedActions 内容AIAI 生成个性化话术;lead-processor 的模板是唯一 fallback

核心原则: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() 原子执行。

type TaskAction =
  | { action: 'create'; params: TaskCreateParams }
  | { action: 'close'; params: TaskCloseParams }
  | { action: 'update'; params: TaskUpdateParams }
  | { action: 'record_progress'; params: TaskProgressParams }

interface TaskActionResult {
  statements: DrizzleStatement[];
  timelineEntries: TimelineEntry[];
}

function applyTaskAction(action: TaskAction, context: TaskContext): TaskActionResult;

各 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:

interface TaskProgressParams {
  taskId: string;
  progressType: 'no_answer' | 'left_voicemail' | 'callback_later';
  note?: string;
  // Orchestrator 内部行为:
  //   1. 不改 task.status (保持 pending)
  //   2. 根据 progressType 重新计算 dueAt (no_answer +1d, left_voicemail +2d, callback_later +3d)
  //   3. 写 timeline event (eventType='task.progress_recorded')
}

对 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 即可:

{
  eventType: 'task.progress_recorded',
  entityType: 'task',
  entityId: taskId,
  newValue: {
    progressType: 'no_answer',
    newDueAt: '2026-06-02T10:00:00Z',
    attemptNumber: 3,
  },
  actorType: 'staff',
}

查询"这个 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,让调用方不再直接拼参数:

class TimelineWriter {
  buildEntry(event: TypedTimelineEvent, context: TimelineContext): DrizzleRow;
  buildSQL(event: TypedTimelineEvent, context: TimelineContext): ContactTimelineInsertSQL;
}

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 本身。

ModuleContract 入口Phase 1 改动
Call Analysisanalyze_call(callId, mode)不拆,但 contract 定义让 future tool calling 有明确接口
SMS Signalevaluate_sms(messageId): SMSSignal不拆,但 signal 分类(STOP/trivial/meaningful/NL_DNC)明确化
Lead Downstream直接调用 Task Orchestrator不需要单独 module
Reprocessreprocess(runId, scope)现有 source='reprocess' 机制已支持

4.5 Reuse Scenario Validation

用 Phase 1 的 3 个 module 验证 4 个跨 pipeline 场景。

Scenario 1: Call retry(员工连打 3 次没人接)

  1. Call 1 (no_answer): ai-analysis-processor 写 call facts -> SQS -> contacts-analyzer AI 建 task via Task Orchestrator.create()
  2. Call 2 (no_answer): contacts-analyzer AI 看到 pending task 已存在 -> Task Orchestrator.update() 更新 suggestedActions(uq_tasks_pending_contact_category 保证不重复建)
  3. 员工 UI 标记"没打通": Task Orchestrator.record_progress({progressType: 'no_answer'}) -> dueAt 推后 1 天,task 保持 pending
  4. 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 级联关闭)

  1. Contact Writer.setDNC(true) -> UPDATE contacts
  2. DNC cascade: 每个 pending task 调 Task Orchestrator.close({closeResult: 'do_not_contact'}) -- 统一的 close 逻辑
  3. 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

SQS trigger -> Lambda -> 读取数据 -> 拼 prompt -> 一次 AI 调用 -> Zod validate -> code 执行 mutation

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 在两种模式下的工作方式

Mode A (one-time, batch/async):
  contacts-analyzer code -> AI 产出 taskDecisions[] -> code 遍历 ->
    TaskOrchestrator.applyTaskAction() -> statements 塞进 db.batch()

Mode B (tool calling, interactive):
  AI 调用 create_task tool -> tool handler 解析参数 ->
    TaskOrchestrator.applyTaskAction() -> tool handler 执行 db.batch() -> 返回结果给 AI

共享层: applyTaskAction() -- 同一个函数, 同样的校验/幂等/timeline

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 决策标准

维度Mode A: one-time callMode B: tool calling
触发方式Event-driven(SQS / cron / webhook)User/agent-initiated(API / voice)
AI 角色数据产出者(leaf node)编排者(决定调用什么)
延迟要求分钟级可接受需要秒级
并发模式高并发(每天 200+ calls)低并发(单用户交互)
现有场景per-call analysis, contact analysis, lead, SMS未来:voice agent, manager Q&A, staff assistant

5.5 Tool calling 前置条件(Phase 1 不实现)

前置条件当前状态
Tenant isolationcontacts PK 已含 store_id;但 tasks.store_id 仍 nullable
Permission / RBACstudio-api 有 store-level 校验,需要泛化成 tool-level RBAC
Audit trailTimeline Writer 需新增 actorSourceType='ai_tool'
Budget cap不存在,需新建
IdempotencyTask Orchestrator 已考虑,每个 action 的 timeline entry 都带 idempotencyKey

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

架构决策

  1. 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 层。

  2. contacts.actionNeeded / suggestedActions 何时清理:schema 注释说 tasks 已取代 contacts 级别的 action 字段,但 contacts-analyzer 仍然同时写两处。Phase 1 是否正式标记这两个字段为 deprecated?前端 dashboard 是否还在读 contacts 级别的 actionNeeded?

  3. neon-sync Lambda:确认为死代码后,是否在 Phase 1 直接删除?

record_progress 设计决策

  1. attempted 迁移策略:现有 378 条 closeResult='attempted' 的 closed task,是否需要数据迁移?还是只影响新数据,历史数据保持原样?

  2. 前端 UI 改动范围:record_progress 需要 studio-api 新端点 + 前端把 no_answer/left_voicemail/callback_later 从"关闭 task"弹窗移到"记录进度"按钮。这个前端改动是否在 Phase 1 scope 内?还是先只改 backend contract,前端保持现有 close 流程(Orchestrator 内部把这 3 个值自动转为 record_progress)?

实施顺序

  1. Phase 1 的 PR 拆分:建议分 4 个 PR -- (a) Task Orchestrator + tests (b) Timeline Writer + Event Catalog (c) contacts-analyzer 迁移到 Orchestrator (d) lead-processor + studio-api 迁移。是否同意?

  2. closeResult 枚举三方不同步(schema 18 / prompt 13 / 前端 UI 16):Phase 1 是否顺带对齐这三方?还是只清理 progress 值,其他不同步留 #199 追踪?

Policy Guard scope

  1. Rate limit 的具体阈值:Policy Guard 的 create_task rate limit -- 具体 N 分钟 M 次的值?还是 Phase 1 先不加 rate limit,只加 DNC guard + storeId guard + duplicate guard?

  2. lifecycle guard observe-only 的持续时间:Contact Writer 的 lifecycle 状态转换 observe-only guard 运行多久后转 enforce?