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/suggestedActionstasks.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 NOTHINGidempotencyKey 去重。

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 -- 所有写入 contactstaskscontact_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.completedmessage.createdlead.created)。这些事件跟 task/contact mutation 无关,是"发生了什么"的记录
  2. Task Orchestrator / Contact Writer 内部调 — 写 mutation 引起的事件(task.createdtask.status_changedcontact.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, progressTypeno_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.createdactivitycallId, direction, callState
call.status_changedactivitycallId, oldState, newState
message.receivedactivitymessageId, direction, content_preview
message.sentactivitymessageId, direction
lead.createdactivityleadId, source
task.createdtasktaskId, typeCategory, priority
task.status_changedtasktaskId, oldStatus, newStatus, closeResult?
task.updatedtasktaskId, changedFields
task.progress_recordedprogress新增taskId, progressType, channel, attemptCount
contact.lifecycle_changedlifecycleoldStage, newStage
contact.dnc_changedlifecycledoNotContact, reason
contact.analysis_completedsystemanalysisRunId, 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 NOTHINGdb.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-processorpersist-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-analyzerneon-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-apiclose.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。

lastActivityAtContactWriter.touchActivity(phone, storeId, occurredAt) -- 统一用 forward-only 语义(GREATEST(existing, newTime)),防止乱序事件回退时间。

DNCContactWriter.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?