Flow 3: Contact Analysis

Current legacy write-path snapshot(verified 2026-05-04; boundary updated 2026-07-15):本文用于追溯旧 contacts-analyzertaskDecisions[] / typeCategory / closeResult 写法,不是 Target Task contract,且 Current 细节必须重新以 live code 核查。Target 见 Task System Design V3

contacts-analyzer Lambda 在客户 contact 累积新通话/消息/lead 后,跑 1 次 AI 调用得到客户状态 + per-task 决策(close/create/update),原子写入 contacts + tasks + contact_timeline。本文给写入字段映射 + 安全网 + 待办,供改 Lambda 行为或追溯字段所有权时参考。

触发: 6 个来源(见 总览) · 错误处理: 阻塞(SQS batchItemFailures,最多 3 次后进 DLQ) · 代码 SoT: callytics-infrastructure repo 的 lambda/contacts-analyzer/

::: details 实施溯源(2026-05-04 verified against source code) 历史设计依据: .claude/specs/2026-03-28-flow3-contact-analysis-design.md(旧 AI-driven task decision 设计)

关键事实:

  • storeId 已写入 contacts + tasks + contact_timeline
  • account_id rename 完成(common migration 0016);下方设计中残留 siteId 字眼按 accountId 理解
  • franchiseId 从 contact 行读取,不硬编码
  • writeAnalysisWithTasks() 是唯一写入路径;旧 writeAnalysis() + TASK_ENABLED_SITES feature flag 已移除
  • Prompt 默认传 pending + closed tasks 给 AI

:::

一、数据流概览

1 次 AI 调用 + 1 次 transaction 写入。AI 输出客户状态字段 + per-task 决策(close/create/update),系统执行写入。

  6 触发源 → SQS {phone, storeId, source?}(FIFO,MessageGroupId = `phone#storeId`;PR #676 SQS contract 已从 `{phone, franchiseId, accountId}` 迁到 `{phone, storeId}`)

    ┌───────────▼────────────────────────────────────────┐
    │  读取 5 张表                                        │
    │    contacts → summary + lifecycle + lastAnalysisAt  │
    │    calls/messages/leads → 增量数据                   │
    │    tasks → pending + closed (全量)                   │
    │                                                     │
    │  增量窗口:                                           │
    │    lastContactAnalysisAt = NULL → 30d 首次全量       │
    │    有值 → since lastContactAnalysisAt               │
    └───────────┬────────────────────────────────────────┘

    ┌───────────▼────────────────────────────────────────┐
    │  构建 Prompt → AI 调用 (1 次)                       │
    │  → 客户状态字段 + taskDecisions[] (Zod 验证)        │
    └───────────┬────────────────────────────────────────┘

    ╔═══════════▼═══════════════════════════════════════╗
    ║  db.batch() — 1 次 HTTP,全成功 or 全回滚          ║
    ║                                                   ║
    ║  ③ UPDATE contacts       ← AI 字段 + system 字段  ║
    ║  ④ AI-driven task decisions:                      ║
    ║     遍历 taskDecisions[],per-task 执行:           ║
    ║     - close  → UPDATE task (status+closeResult)   ║
    ║     - create → INSERT 新 task(可多个)             ║
    ║     - update → UPDATE task (priority/actions)     ║
    ║  ⑤ INSERT timeline       ← ④ 执行结果的审计记录   ║
    ╚═══════════════════════════════════════════════════╝

    ⑥ CloudWatch 日志

二、输入:什么数据喂给 AI

读 5 张表(contacts / calls / messages / leads / tasks),拼成 prompt 给 AI。增量窗口由 contacts 行的 lastContactAnalysisAt 决定。

2.1 contacts 表

字段为什么读
customerSummarycompressed memory — 上次分析的精华,AI 用它作为历史上下文
lastContactAnalysisAt增量读取的时间边界(NULL = 首次分析)
leadStatus当前漏斗阶段
lifecycleStagelead / member / churned
firstName, lastName客户身份

2.2 增量数据:calls + messages + leads

if (lastContactAnalysisAt === null) {
  // 首次分析:读 30d 全量,建立基线 customerSummary
  时间边界 = NOW() - 30d
} else {
  // 后续:只读上次分析后的新数据
  时间边界 = lastContactAnalysisAt
}
查询条件读取字段
callscustomerPhone=? AND franchiseId=? AND startTime > 时间边界(PR #676 起 reader 不再带 accountId)startTime, direction, duration, executiveSummary, primaryCategory, callState
messagesfromPhoneNumber=? AND franchiseId=? AND creationTime > 时间边界creationTime, direction, type, subject
leadsphone=? AND franchiseId=? AND receivedAt > 时间边界(infra#518 已修复)firstName, lastName, receivedAt, leadType

2.3 tasks 表

读取该 contact 的所有 tasks(per-contact 通常 1-5 条,数据量极小):

状态读取字段AI 用来做什么
pendingtypeCategory, priority, dueAt, suggestedActions知道已有什么 task → 不重复建
closedtypeCategory, closeResult, closeNote, closedAtcloseResult 最关键 — 员工说 "booked" → AI 建 booked_not_converted 跟进

2.4 Prompt 结构

System Prompt(静态,可缓存)
  └─ 角色定义 + 输出格式要求

User Message(动态,per-contact)
  ├─ CUSTOMER PROFILE: phone, siteId, lifecycleStage, leadStatus
  ├─ PREVIOUS SUMMARY: customerSummary
  ├─ PENDING TASKS:
  │     - [LEAD_OUTREACH] high priority, due in 2h
  │       Suggested: "Call within SLA window"
  ├─ RECENTLY CLOSED TASKS:
  │     - [LEAD_FOLLOW_UP] closed Mar 25 by Sarah
  │       Result: booked | Note: "Customer said Thursday 3pm for trial"
  ├─ NEW CALLS (since last analysis):
  │     - Mar 26 10:00 (inbound): "Customer called back, ready to sign up"
  ├─ NEW MESSAGES (since last analysis):
  │     - Mar 25 16:00 (inbound): "I'll come Thursday for the trial"
  └─ NEW LEADS (since last analysis):
        (none)

当前 prompt 默认传入 tasks: { pending, closed },渲染为 PENDING TASKS / RECENTLY CLOSED TASKS section。TASK_ENABLED_SITES feature flag 已移除。

2.5 跳过条件

  • 如果 calls + messages + leads 全部为空 无 customerSummary → skip(无数据可分析)

2.6 customerSummary 准确性 — Prompt 要求

customerSummary 是 compressed memory:AI 每次读旧 summary + 新数据 → 重写(不追加)。准确度靠 prompt,system prompt 必须包含 3 条:

  1. Previous summary 是 context,不是绝对真理 — 新数据可以 override
  2. 保留关键细节 — dates / commitments / objections / complaints
  3. 区分确认 vs 可能 — "customer confirmed Thursday 3pm" vs "customer mentioned possibly coming"

为什么不保存历史 summary:contact_timeline 已记录精确事件历史,customerSummary 只需要"当前快照"。


三、AI 输出字段路由表 ★

AI 输出落到 contacts / tasks / contact_timeline 三张表;系统补充 timestamp / id / store 隔离字段。

3.1 AI 输出 → 3 张表

字段集由 ContactsAnalysisSchema Zod schema 定义。

客户状态字段(写入 contacts 表):

#AI 字段类型→ contacts→ timeline决定/触发Zod
1customerSummarytext
2leadStatustext enum (12)
3leadStatusReasontext
4lifecycleStagetext enum (4)✅ 6c new变了 → 写 6c 事件
5lifecycleStatetext enum (3)✅ 6c new变了 → 写 6c 事件
6purchaseIntenttext enum (3)
7purchaseIntentReasontext
8goalsjsonb
9leadObjectionsjsonb
10leadRejectionReasonsjsonb
11actionNeededboolean信号字段,不直接触发 task
12actionNeededReasontext
13suggestedActionsjsonbcontacts 级别建议(legacy 兼容)
14typeCategorytext enum (9)legacy 兼容,task 操作由 taskDecisions 驱动
15prioritytext enum (3)legacy 兼容
16doNotContactboolean
17hasOpenComplaintboolean

taskDecisions[](AI-driven task 操作,写入 tasks + timeline):

action→ tasks→ timeline说明
close✅ UPDATE status→closed + closeResult + closeNote✅ task.status_changedAI 从 11 值 closeResult 枚举选,reason → closeNote
create✅ INSERT 新 task(可多个,不同 typeCategory)✅ task.createdonConflictDoNothing dedup
update✅ UPDATE priority / suggestedActions / dueAt✅ task.updated同件事,情况变了

taskDecisions 驱动所有 task 操作:actionNeeded 留在 contacts 当信号字段,但 task 的 close / create / update 完全由 taskDecisions[] 决定,不再由 actionNeeded=true 触发。

AI 能关任何 sourceType 的 task:不限 sourceType='contact_analysis',也能关 sourceType='lead' 的 lead_outreach task。

Legacy vocabularytypeCategory / closeResult 与 Follow-up relay 的历史定义见 tasks-field-design.mdtask-lifecycle.md;两者均已 superseded。Target 只以 Task System Design V3 为准。

3.2 System 补充字段(非 AI 输出)

字段来源→ contacts→ tasks→ timeline
lastContactAnalysisAtNOW()
updatedAtNOW()(显式,$onUpdate 在 tx 内不生效)
taskIdcrypto.randomUUID()(create action)✅ entityId
contactAnalysisRunIdsqsRecord.messageId(SQS envelope 天然 UUID,直接使用)
taskTypedecision.typeCategory 派生:lead_outreach→自身,其他→follow_up✅ newValue
sourceType'contact_analysis' 硬编码
dueAtcomputeDueAt(priority)(high=4h / med=24h / low=72h)
status='closed'close action✅ 6a newValue
closeType='auto_closed'close action✅ 6a newValue
closedByStaffName='system'close action
closeResultAI 从 11 值枚举选(close action 的 decision.closeResult)✅ 6a newValue
closeNoteAI 写原因(close action 的 decision.reason)✅ 6a newValue
storeIdresolvePhoneIdentity(calls[0].from, calls[0].to) → storeId(fallback messages[0])
storePhone同上,resolvePhoneIdentity 返回的 storePhone— contacts 无 store_phone
predecessorTaskIdclose→create 同 typeCategory 时记录前任 taskId✅ newValue
leadTemperaturecomputeTemperature()⚠️ 待移除 #32

store_id 覆盖现状:8 张 Neon 表(contacts / calls / messages / leads / contact_timeline / staff / tasks / contact_analysis_runs)均具备 store_id 列。contacts PK = (phone, store_id);staff.store_id NOT NULL;其余 6 表 nullable,writer 全写,老数据靠 backfill 脚本补齐。

相关文档:store-level-isolation.md(终态 SoT)· infra#628(实施审计)。

3.3 Timeline 事件路由(Step ⑤)

Timeline 记录 ④ 中 taskDecisions 的执行结果(audit records),不是 AI 决策本身。

事件触发条件oldValuenewValue
6a task.status_changedtaskDecisions[].action='close'(AI per-task 决策){status: 'pending'}{status: 'closed', closeType: 'auto_closed', closeResult: '<AI选择>', closeNote: '<AI原因>'}
6b task.createdtaskDecisions[].action='create'(可多个,不同 typeCategory){taskType, typeCategory, priority}
6b+ task.updatedtaskDecisions[].action='update'(优先级/建议更新){priority: old, suggestedActions: old}{priority: new, suggestedActions: new, reason: '<AI原因>'}
6c contact.lifecycle_changedAI 输出 ≠ DB 旧值(AI vs DB diff){stage, state} DB READ{stage, state} AI
6d contact.temperature_changed规则计算 ≠ DB 旧值(⚠️ 待移除){temp} DB READ{temp, reason} RULE

四、执行流程

唯一路径是 writeAnalysisWithTasks(),所有触发走同一原子写入。底层用 db.batch() 不是 db.transaction() —— neon-http driver 的 db.transaction() 会直接 THROW;db.batch() 把所有 statement 打包成 1 次 HTTP,在 DB 内 BEGIN/COMMIT,全成功或全回滚。

4.1 AI 之后、batch 之前(后处理)

步骤动作
派生 contactAnalysisRunId直接用 SQS envelope 的 messageId,天然 UUID
派生 leadTemperature规则计算 ⚠️ 待移除 #32
Pre-batch guard比对 AI 的 create 决策和现有 pending tasks,已有 typeCategory 不生成 INSERT statement(避免无意义的 onConflictDoNothing)

4.2 遍历 taskDecisions[] 生成 statements

action生成的 statement关键字段
closeUPDATE tasks + INSERT timeline(task.status_changed)closeResult = AI 从 11 值枚举选;closeNote = AI 原因;closedByStaffName='system';WHERE 带 status='pending' guard 防关已关
createINSERT tasks + INSERT timeline(task.created)taskId = 新 UUID;taskTypetypeCategory 派生(lead_outreach→自身,其他→follow_up);onConflictDoNothing 兜底 unique index
updateUPDATE tasks + INSERT timeline(task.updated)只 set 提供的字段(priority / suggestedActions / dueAt);priority 变了会重算 dueAt;WHERE 同样 pending guard

4.3 batch 内 statement 顺序

  1. UPDATE contacts(AI 字段 + lastContactAnalysisAt=NOW() + updatedAt=NOW())
  2. 所有 close/create/update statements + 对应 timeline INSERT
  3. lifecycle 变了 → INSERT timeline(contact.lifecycle_changed),没变就不加这条 statement

batch 后 ⑥ 输出 CloudWatch 日志:tasksClosedCount / tasksCreatedCount / tasksUpdatedCount / timelineEventsCount / contactAnalysisRunId / transactionDurationMs / source / isFirstAnalysis

Statement 数量:最少 1 条(只有 ③,taskDecisions 为空);典型 3-5 条;最多 ~15 条(多个 close/create/update + 对应 timeline + lifecycle)。

⚠️ batch 内 3 条硬约束:updatedAt 必须显式 sql\NOW()`($onUpdate 在 batch 内不生效);ON CONFLICT DO NOTHING不导致 batch 回滚;不能在 batch 内用Promise.all`。


五、安全网

9 层保护防 AI 误决策 / 并发 / 重试 / 数据损坏。任一保护被绕过都需要 incident review。

保护层机制防什么
Zod validationAI 输出 → ContactsAnalysisSchema.parse() + TaskDecision discriminated union非法类型/枚举值/action 类型
uq_tasks_analysis_run_category(contactAnalysisRunId, typeCategory) WHERE NOT NULLSQS 重试重复
uq_tasks_pending_contact_category(contactPhone, franchiseId, accountId, typeCategory) WHERE pending ⚠️ 仍是老 3-key,store-level isolation 后未升级到 (phone, storeId, typeCategory) — 见 infra#776并发重复(多 store 共享 phone 场景失效)
chk_tasks_closed_integrity(pending → no closeType) OR (closed → has closeType)状态不一致
Close/update WHERE guardeq(tasks.status, 'pending') 在每个 close/update 的 WHERE 中防止关/改已关的 task
AI 引用不存在的 taskIdUPDATE 影响 0 行,batch 继续安全(无数据损坏),log warning 监控
Create 预检查batch 前比对 AI 的 create 和现有 pending,已存在的 typeCategory 不生成 statements避免生成无意义的 onConflictDoNothing
Transaction rollback撞任一 index → 整个 tx 回滚 → SQS ack 成功幂等保证
FIFO SQS同一 contact 串行(MessageGroupId = phone#storeId,PR #676 起;cron + per-call 一致)防同 contact 并发分析

六、待加入 / 待移除

项目状态
Prompt 加 pending + closed tasks 数据✅ 已落地,渲染为 PENDING TASKS / RECENTLY CLOSED TASKS section
Zod schema 加 taskDecisions + doNotContact + hasOpenComplaint✅ 已落地(TaskDecision discriminated union)
writeAnalysisWithTasks() 新路径(AI-driven task decisions)✅ 已落地(唯一写入路径)
FIFO SQS queue(dailyBatchQueue Standard → FIFO)✅ 已落地,同 contact 串行,不同 contact 并行
SMS 正文内容传给 AI未来 — 当前只传 direction/type/subject
leadTemperature / leadTemperatureReason⚠️ 待移除 #32 — 从 timeline 到达时间实时计算
contact.temperature_changed timeline⚠️ 随 temperature 一起移除