Lambda 写入矩阵

每个 Lambda 对每张 Neon 表的字段级写入职责清单,给开发者改 Lambda 写入逻辑前对照。字段所有权是硬约束 — 违反会导致 last-write-wins 竞态或数据不一致。

::: details 实施溯源 + 落地状态(2026-05-04) 设计文档:

  • Pipeline 设计: .claude/specs/2026-03-25-task-pipeline-write-flow.md
  • Flow 3(AI-driven task decisions): .claude/specs/2026-03-28-flow3-contact-analysis-design.md

Schema 状态(callytics-common 0.34.0 + migration 0014-0023):

  • 主业务表 contacts / calls / messages / leads / contact_timeline / tasks / staff,辅助表 rc-stores / store-config / user-stores 均具备 store_id
  • task_events 仍只有 account_id,无 store_id — audit 表暂保留 account-level isolation,store-level 终态后续补
  • contacts 复合主键 = (phone, store_id),store_id NOT NULL(migration 0014)
  • staff.store_id NOT NULL(migration 0015)
  • site_idaccount_id rename(migration 0016)
  • 其余表 store_id 仍 nullable(calls / leads / messages / contact_timeline / tasks),writer 全写但保留 NULL 兼容历史行,各表均有 idx_<table>_store_id WHERE store_id IS NOT NULL partial index

Lambda writer 当前 store_id 写入状态(2026-05-04):

  • transcribe-processorcalls.store_id + contact_timeline.store_id(call.status_changed,Phase 1 起 dual writer)
  • ai-analysis-processorcontacts/calls/contact_timeline.store_id
  • contacts-analyzercontacts/tasks/contact_timeline.store_id
  • message-processormessages/contacts/contact_timeline.store_id
  • lead-tracking(独立 repo)→ leads/contacts/tasks/contact_timeline.store_id

franchise_id 来源(读源代码核对):

  • message-processor:parseFranchiseFromClientId(_client_id, _account_id) 解析,不是硬编码
  • ai-analysis-processor:从 SQS 上游 metadata 传入
  • contacts-analyzer:从 contact 行的 franchiseId 读取
  • lead-tracking:StoresV2 leadEmail 映射 → parseBrandFromFranchise()
  • 当前所有 writer 都从外部数据派生 franchise_id,不硬编码(历史 P0 #632 fix:阻止"_client_id 整串含 site 后缀直接写入"的 bug)。

:::

全局数据流

图例:实线 = 写入 · 虚线 = 读取 / DEFERRED · 🆕 = Task 系统新增 · * = DEFERRED(V1 不实现)


一、写入模式总览

1.1 所有 Lambda 统一用 Neon HTTP driver + Drizzle ORM

连接方式:@neondatabase/serverless HTTP driver → drizzle-orm/neon-http。每条 SQL 是一个 HTTP POST。

1.2 多表写入:db.batch() vs 分开 await

写法HTTP 请求数原子性用在哪
await db.batch([stmt1, stmt2, stmt3])1✅ 全成功 or 全回滚写入之间有因果依赖
分开 await db.insert(...)N❌ 各自独立有明确"主记录",其他是补充
db.transaction(async tx => {...})neon-http driver 不支持,会 THROW

db.batch() 底层调用 Neon 的 sql.transaction() — 所有 SQL 打包成 1 个 HTTP 请求,DB 内执行 BEGIN/COMMIT/ROLLBACK。

1.3 五个 Flow 写入模式

所有 Flow 的 Neon 写入统一用 db.batch() 打包(2026-03-28 Neon-primary 决策)。DDB 写入是过渡期冗余,独立标记不参与 batch;SQS 发送在 batch 外,batch 成功后才发,非阻塞。

FlowLambdaNeon 写入模式原因
1lead-tracking3-Lambda 拆分:poller UPSERT leads(fail-open,失败 → neon-write-dlq) → NeonRetryProcessor retry + LeadCreated → EventBridge/SQS → lead-processor db.batch([contacts, tasks, timeline])leads 写入失败走 DLQ 异步重试,email still SEEN;downstream 失败走 SQS retry,避免互相阻塞
2Per-Call Analysisdb.batch([calls, contacts, timeline])calls+contacts+timeline 原子;DDB calls 过渡期独立写
3Contact Analysisdb.batch([contacts, ...taskDecisions, timeline])关旧 task + 建新 task 必须原子(否则"幽灵分析")
4astudio-api Closeraw SQL × 2(UPDATE tasks → INSERT task_events)当前未原子化,前者成功后者失败 → task 已关但无 audit;后续可改用 db.batch()
4bstudio-api Postpone❌ 未实施暂无端点;终态用 db.batch()
5message-processorper-entry db.batch([messages, contacts, timeline])每条 SQS record 一个 batch,失败 → 该 entry 进 batchItemFailures 重试

1.4 踩坑注意事项

::: warning db.batch() 生产注意(社区踩坑总结)

  1. db.batch() 内不能 Promise.all — 并行语句会丢 transaction scope(Drizzle #2200
  2. $onUpdate(() => new Date()) 在 batch 内不生效 — Drizzle 的 $onUpdate 是 query builder hook,不是 DB trigger。batch 内所有 UPDATE 必须显式写 updatedAt: sql\NOW()``
  3. ON CONFLICT DO NOTHING 不导致 batch 回滚 — PostgreSQL 的 conflict 处理是静默跳过,不抛错。dedup index 在 batch 内安全工作
  4. Pin driver 版本@neondatabase/serverlessdrizzle-orm 之间有版本兼容性问题,社区报告升级后运行时报错(Drizzle #5208
  5. SQS 提供 Lambda 层重试db.batch() 解决 SQL 原子性,SQS batchItemFailures 解决 Lambda crash 重试。两者配合 = 幂等 + 原子

:::

二、总览矩阵

Lambda触发输入(原始数据从哪来)AI 分析前从 Neon 读什么写入哪些表详情
lead-trackingIMAP email(EventBridge 5min 轮询)leads ✅ · contacts 🆕 · tasks 🆕 · timeline 🆕详情
Per-Call AnalysisSQS(transcript,来自 transcribe-processor)contacts(customerSummary)+ calls(24h) + messages(24h)calls ✅ · contacts ✅ · timeline 🆕详情
Contact AnalysisSQS {phone, franchiseId, siteId}contacts(summary)+ 增量 calls/msgs/leads + taskscontacts ✅ · tasks 🆕 · timeline 🆕详情
message-processorSQS(RingCentral SMS/VM webhook)messages ✅ · contacts 🆕 · timeline 🆕详情
studio-apiAPI request(员工 UI 操作)tasks(校验 status=pending)*tasks* · contacts* · timeline*详情

✅ = 现有 🆕 = V1 新增 * = DEFERRED — = 不读 Neon(数据全在触发输入里)

Per-Call Analysis 读取:当前代码从 DDB 读(customer-history.ts),计划迁移到 Neon。如果 customerSummary 存在就读,不存在就跳过(graceful degradation)。

Contact Analysis 读取:customerSummary 是 compressed memory。lastContactAnalysisAt 为 NULL(首次)→ 读 30d 全量建立基线。有值 → 读 lastContactAnalysisAt 之后的增量。详见 contact-analysis-writes.md


三、Contact Analysis 触发源

6 个触发源,统一发 SQS 消息到同一 schema(Zod-validated by SqsMessageBodySchema):

{
  phone: string,        // 必填
  storeId: string,      // 必填 — contacts PK 组件(post-migration 0014)
  source?: 'cron' | 'on_demand' | 'per_call_analysis' | 'task_close' | 'showed' | 'trialed',
  // optional observability trace 字段 — 老 producer 不传也能解析
  callId?: string,
  telephonySessionId?: string,
  analysisCompletedAt?: number,
  queuedAt?: number,
}

franchiseId / accountId 不再传 — Consumer(contacts-analyzer)从 contact 行通过 PK lookup 自行读取。上游 producer 只需要 phone + storeId

技术约束:SQS ESM filtering 只支持 body key(不支持 MessageAttributes)。

FIFO Queue 已上线 — dailyBatchQueue 从 Standard 切换为 FIFO 并迁移到 SqsStack。触发源发送消息必须提供 MessageGroupIdMessageDeduplicationId,例如 ai-analysis-processor 用 ${phone}#${storeId} 作 group key、per_call_analysis-${phone}-${callId} 作 dedup id。

#触发源SQS Queuesource状态
Cron 每天 06:00 UTCdailyBatchQueue (FIFO)'cron'✅ ENABLED in CDK
On-demand(前端 Refresh AI)dailyBatchQueue(V1)/ expressQueue*'on_demand'✅ 后端现有,前端 DEFERRED
Per-Call Analysis(每通电话)dailyBatchQueue(V1)/ expressQueue*'per_call_analysis'✅ 已实施(ai-analysis-processor non-blocking SQS send)
Task close(员工关闭 Task)dailyBatchQueue (FIFO) / expressQueue*'task_close'⚠️ 部分实施 — studio-api 已能关 task + 写 audit,timeline / contacts.lastActivityAt / SQS 触发暂未补齐(详见 studio-api-writes.md
T5showed(4h 黄金窗口)expressQueue*'showed'🆕 DEFERRED
T6trialed(2-4h 黄金窗口)expressQueue*'trialed'🆕 DEFERRED

* expressQueue 尚未创建(Phase D5)。V1 期间 ❷❸❹ 用 dailyBatchQueue。 ❷ On-demand 后端 API 已存在,但前端 Refresh AI 按钮随员工 UI 一起 DEFERRED。

字段来源

触发源phone / storeId 从哪来
❶ Croncontacts 表遍历(dispatch fan-out 时发 per-contact 消息)
❷ On-demand前端 API request payload
❸ Per-Call Analysisai-analysis-processor 解析 storeId 后发送(contactPhone + storeId 都解析才触发)
❹ Task closetasks 表的 contact_phone + store_id
T5/T6Contact Analysis 自身检测到 leadStatus 变化后发送

V1 不做 Cooldown

不做应用层 cooldown,重复分析由 DB index 保护(uq_tasks_pending_contact_category 防重复 task,uq_tasks_analysis_run_category 防 SQS 重试)。成本 ~200 calls/day × $0.005 ≈ $1/day,可接受。lastContactAnalysisAt 每次分析仍写 NOW(),未来加 cooldown 时可直接读取。

T5/T6 推荐实现(DEFERRED)

Contact Analysis 自检测 leadStatus 变化 → EventBridge Scheduler 延迟触发(showed=4h, trialed=2-4h) — SQS DelaySeconds 最大 15min 不够用。Loop guard:source 已是 showed/trialed 时不再检测。


四、contacts 表写入详情

contacts 是中心化客户档案,4 个 Lambda 各写不同字段。

字段所有权

Per-Call Analysis 只写"活动标记"(什么时候发生了什么),Contact Analysis 只写"AI 判断"(客户状态是什么) — 两者字段互不重叠

分类字段所有者
身份firstName, lastNamePer-Call Analysis(初始值由 lead-tracking 设)
活动时间lastActivityAtPer-Call Analysis + lead-tracking + message-processor + studio-api
业务数据hasCardOnFilePer-Call Analysis(credit_card_captured AI 检测 → contacts 直写)
运营notes员工(studio-api,DEFERRED)
初始状态lifecycleStage='lead', leadStatus='new'lead-tracking(仅 INSERT)
AI 画像customerSummary, leadStatus, leadStatusReason, lifecycleStage, lifecycleState, purchaseIntent, purchaseIntentReason, goals, leadObjections, leadRejectionReasonsContact Analysis
AI 行动actionNeeded, actionNeededReason, suggestedActionsContact Analysis
AI 风险doNotContact, doNotContactUpdatedBy, hasOpenComplaintContact Analysis
AI 追踪lastContactAnalysisAtContact Analysis
系统createdAt, updatedAtDB / 所有 Lambda

三个 UPSERT contacts 对比

lead-tracking / Per-Call Analysis / message-processor 都 UPSERT contacts,但写的字段不同 — 数据源决定能写什么。

INSERT 时写入:

字段lead-trackingPer-Call Analysismessage-processor差异原因
phone, franchiseId, siteIdPK,三者都需要
lifecycleStage='lead'只有 lead-tracking 知道这是 lead
leadStatus='new'同上
firstName✅(email 解析)✅(call AI 提取)SMS webhook 不含名字(#521 未来可从 RC API 补写)
lastName同上
lastActivityAtNOW()NOW()NOW()三者一致 — 任何互动都是活动
firstAttemptedAt 待移除通话才算"尝试联系"
firstConnectedAt 待移除接通才算"首次接通"

ON CONFLICT 时更新:

字段lead-trackingPer-Call Analysismessage-processor差异原因
lastActivityAtNOW()(覆盖)NOW()(覆盖)NOW()(覆盖)三者一致
firstNameCOALESCE(只补写)覆盖(call AI 更准)lead email 名字不如通话准,用 COALESCE 保留已有值
lastNameCOALESCE覆盖同上
lifecycleStage不覆盖已有 contact 由 Contact Analysis 管理
leadStatus不覆盖同上
updatedAtNOW()NOW()NOW()三者一致

详见各 per-Lambda 文件的 Step ② section。


五、contact_timeline 表写入详情

contact_timeline 记录所有客户互动和状态变更事件,5 个 Lambda 共写 9 种事件类型。每条都包含公共必填字段(contactPhone / franchiseId / siteId / eventType / eventCategory / occurredAt / actorType)。

Lambda(文件)StepeventTypeentityTypeentityIdoldValuenewValueoccurred_atactorType
lead-tracking(Step ④lead.createdleadlead.idlead.receivedAtlead_webhook
Per-Call Analysis(Step ④call.createdcallcall.telephonySessionId{callDirection, staffName, duration}call.startTimecall_analysis
message-processor(Step ③message.createdmessageString(message.id){direction, messageType}message.creationTimesystem
Contact Analysis(Step ⑥task.createdtaskString(taskId){taskType, typeCategory, priority}NOW()(tx)contact_analysis
Contact Analysis(Step ⑥task.status_changedtaskString(closed_task_id){status: pending}{status: closed, closeType: auto_closed, closeResult: AI(11值)}NOW()(tx)contact_analysis
Contact Analysis(Step ⑥task.updatedtaskString(taskId){priority: old, suggestedActions: old}{priority: new, suggestedActions: new, reason}NOW()(tx)contact_analysis
Contact Analysis(Step ⑥contact.lifecycle_changedcontact{lifecycleStage: old, lifecycleState: old}{lifecycleStage: new, lifecycleState: new}NOW()(tx)contact_analysis
Contact Analysis(Step ⑥contact.temperature_changedcontact{leadTemperature: old}{leadTemperature: new, reason}NOW()(tx)contact_analysis
studio-api(4a Step ③)*task.status_changedtaskString(task.taskId){status: pending}{status: closed, closeType: manual_closed, closeResult}NOW()staff
studio-api(4b Step ②)*task.due_at_changedtaskString(task.taskId){dueAt: old}{dueAt: new, reason}NOW()staff

* = DEFERRED

occurred_at 规则:非 transaction 写入(lead / call / SMS)用源记录时间防排序漂移;transaction 内(Contact Analysis Step ⑥)用 NOW(),所有语句共享同一时间。

错误处理:5 个 Lambda 的 timeline INSERT 都在 db.batch() 内,batch 阻塞(失败 → 整个 batch 回滚 → SQS 重试)。lead-tracking 的 prod Neon 写入阻塞,test Neon fire-and-forget;其他 Lambda 的 SQS send / DDB 旁路写入是非阻塞。studio-api Close 当前是 raw SQL × 2(UPDATE tasks + INSERT task_events)未原子化;不写 contact_timeline。


六、字段清理计划

有了 contact_timeline 后,contacts 表上部分字段可从 timeline 派生 → 计划移除。temperature 相关字段和 timeline 事件也将一并移除。

已移除的 contacts 字段

字段替代方案PR
leadStatusChangelogtimeline contact.lifecycle_changed 事件callytics-common#33

待移除的 contacts 字段(需 engineer review,callytics-common#32)

Phase 1:timeline 有数据后可移除

移除字段替代方案
acquiredAtMIN(occurred_at) FROM contact_timeline WHERE event_type='lead.created'
firstAttemptedAtMIN(occurred_at) FROM contact_timeline WHERE event_type='call.status_changed' AND new_value->>'status'='Setup' AND new_value->>'direction'='Outbound'
Phase 1 起(infra#592 / #681 merged 2026-04-29):用 call.status_changed Setup 事件 — 拨号瞬间记录,不等 AI 分析(原 call.created 是 Disconnected 后回填,失去实时性)
firstConnectedAtMIN(occurred_at) FROM contact_timeline WHERE event_type='call.status_changed' AND new_value->>'status'='Answered'
Phase 1 起:RC Answered = 路由层接通(不一定是 human conversation,见 SoT § 5)。如果只算"真说话",仍需读 calls.callState='human_conversation'(AI 分析填)
lastComplaintAttimeline complaint 事件
leadTemperature / leadTemperatureReason从 timeline 到达时间实时计算(NOW() - occurredAt),无需存储或中间函数

Phase 2:Task Pipeline 完成 + 前端迁移后可移除

移除字段替代方案为什么等
actionNeededEXISTS(SELECT 1 FROM tasks WHERE status='pending')有了 tasks 表后这是 stale snapshot — 员工关了 task 但 contacts.actionNeeded 仍为 true
actionNeededReasontasks 表同名字段跟随 actionNeeded
suggestedActionstasks 表同名字段(更新鲜)同上。需删 idx_contacts_action_needed

待移除的 timeline 事件

contact.temperature_changed — 随 temperature 字段一起移除。

为什么 timeline 比 contacts 列更好

重复进入场景(客户 3 月来 → 冷却 → 9 月再来)时,contacts 列的 acquiredAt 永远冻结在 3 月。Timeline 可查当前周期的 lead.created,准确算 Speed to Lead。

相关决策变更

  • typeCategory 门控:temperature≠cold → 改为 AI prompt 指令("lastActivityAt > 5d 不建 lead_follow_up")
  • Stale override:>30d → neglected → 改为"30 天无活动 + 无 pending task"查询

当前状态:新代码(lead-tracking / message-processor)不写这些字段。Contact Analysis 现有代码仍写,等 timeline 有数据 + 验证后停止。


七、已知问题与代码隐患

Expert Panel 代码验证发现 10 个待修复问题,按严重度排序。

#问题严重度当前状态GH Issue
1getRecentLeadsfranchiseId 过滤P0✅ 已修复 — 查询已加 franchiseId 过滤#518
2Per-Call Analysis 当前不触发 Contact AnalysisP0✅ 已修复 — ai-analysis-processor 在 batch 写完后 non-blocking 发 SQS#519
3lastContactAnalysisAt 从未写入P1✅ 已修复 — Contact Analysis 每次写 NOW()#519
4ContactData interface 缺字段P1✅ 已修复#519
5Cooldown 未实现时 duplicate taskV1 设计不做 Cooldown — DB index uq_tasks_pending_contact_category 保护数据完整性
6冗余 contacts 字段清理(9 个字段)P2待 timeline 有数据后移除 6 个 + Phase B 后移除 actionNeeded 三件套callytics-common#32
7writeAnalysis() 不在 transactionP1✅ 已修复 — 唯一路径是 writeAnalysisWithTasks()db.batch() 原子写入
8Reserved concurrency=2P2多触发源时升到 5
9Per-Call Analysis 从 DDB 读历史 → 迁移到 NeonP2customer-history.ts 改为读 Neon contacts + calls(24h) + messages(24h)
10Contact Analysis 30 天全量读 → 增量优化P2customerSummary + 增量(NULL → 30d 首次,有值 → since lastContactAnalysisAt)

八、与其他文档的关系

文档关系
Task Pipeline 写入流程(.claude/specs/2026-03-25-task-pipeline-write-flow.md)Pipeline 级流程设计(本目录是字段级补充)
Expert Panel 审查报告(.claude/specs/2026-03-25-task-pipeline-expert-review.md)技术决策背景和争议记录
Task 字段设计tasks 表业务字段定义
Contacts 数据来源与字段设计contacts 表字段设计来源
Backend 数据管道数据管道总体架构