Flow 1: Lead Tracking
Repo: lead-tracking(独立 repo + CDK stack) · 触发: EventBridge 每 5 分钟 IMAP 轮询 · 架构: 3 Lambda(poller leads-only + neon-retry-processor DLQ 重试 + lead-processor contacts/tasks/timeline)· Neon 写入: poller prod/test 双写,retry path 只重试 prod
::: details 字段命名约定
accountId = RingCentral OAuth providerAccountId,对应 Neon 列 account_id(callytics-common 0.21.0 起从 siteId rename)。下文遗留的 siteId 字眼按 accountId 理解。storeId 是 store-level isolation 的唯一隔离键(UUID),accountId 仅作审计。
:::
数据流概览
不调用 AI,所有值来自 email 解析 [EMAIL] 或系统硬编码 [SYSTEM]。
3 Lambda 拆分理由:
- Poller 和 lead-processor 拆分:leads 写入失败(circuit breaker + DLQ)和 downstream 写入失败(SQS retry)解耦,避免互相阻塞。
- NeonRetryProcessor 拆分:poller 走 fail-open 不阻塞 IMAP SEEN 标记;真正的 Neon 重试由独立 Lambda 异步消费 DLQ,降低 poller 延迟和复杂度。
- retry path 和 happy path 必须发同一份 LeadCreated 事件(issue #201),否则 DLQ 救回的 leads 永远没有下游 contacts/tasks/timeline。
错误处理
Poller(src/poller.ts)
NeonRetryProcessor(src/neon-retry-processor.ts)— 本节由 issue #201 引入
下游 lead-processor(callytics-infrastructure)
Duplicate LeadCreated 风险与下游 idempotency
NeonRetryProcessor 写 Neon 成功后再发 LeadCreated。若 EventBridge publish 失败 → 整条 SQS 消息 redrive → 第二次执行 persistLeadPipeline + 重发 LeadCreated → 下游可能收到 2 次同一 lead 的 LeadCreated。
下游 lead-processor 靠 3 个 idempotency 抗住:
contactsON CONFLICT (phone, storeId) DO UPDATE — UPSERT 幂等tasksON CONFLICT (sourceLeadId) DO NOTHING(uq_tasks_source_lead)— 同 lead 不创建第 2 个 taskcontact_timelineON CONFLICT (idempotencyKey) DO NOTHING —lead.created:${row.id}作为 key
persistLeadPipeline 自己也是 idempotent(leads ON CONFLICT DO UPDATE),双写主要影响 syncedAt bump,无业务副作用。
读取数据
Step ①:UPSERT leads
Poller Lambda 执行 Drizzle upsert,PK 是 composite dedup key({email}#{phone}#{date} 或 {email}#{receivedAt})。
INSERT 时写入
ON CONFLICT (id) DO UPDATE
Step ②:发布 EventBridge LeadCreated 事件
leads UPSERT 成功后由 poller 发布。Event payload: {Source: 'lead-tracking', DetailType: 'LeadCreated', Detail: <LeadRow JSON>}。
前置条件(全部非 null 才发布):
缺任一字段 → 跳过 EventBridge 发布,只写 leads 表,log [EVENTBRIDGE] Skipping — missing required fields。
CDK 路由:LeadCreatedRule-${region} Rule → SQS → Lead Processor Lambda。
Step ③:UPSERT contacts(轻量)
确保 contact 记录存在 + 设置初始状态。Lead Processor 在 db.batch() 内与 ④⑤ 原子执行。
INSERT 时写入
ON CONFLICT (phone, storeId) DO UPDATE
ON CONFLICT 只补写身份和 activity 字段,不写 leadStatus 和 lifecycleStage(这两个已有 contact 的字段由 Contact Analysis 管理)。
不写的字段
Step ④:INSERT tasks
Lead Processor 在 db.batch() 内与 ③⑤ 原子执行,创建 1 条 high-priority lead_outreach task。
去重保护(ON CONFLICT DO NOTHING)
撞任一 index 安全跳过,不报错。
Step ⑤:INSERT contact_timeline
Lead Processor 在 db.batch() 内与 ③④ 原子执行。所有 timeline writer 通过共享 buildTimelineValues() helper(@retaintive/common/db)写入,统一字段约定。
ON CONFLICT DO NOTHING 通过 idempotencyKey 去重,SQS 重试幂等。
CDK 基础设施
资源定义在 lib/leadtracking-stack.ts(lead-tracking repo)和 lib/stacks/sqs-stack.ts / eventbridge-stack.ts(callytics-infrastructure repo)。
lead-tracking repo(us-east-1)
callytics-infrastructure repo(us-west-2)
⚠️ 已知缺失:DLQ CloudWatch alarm
neon-write-dlq 和 neon-write-final-dlq 当前没有 CloudWatch Alarm → 消息进入终态 DLQ 时静默,只能从 SQS console 看到。
对比 callytics-infrastructure lib/stacks/monitoring-stack.ts 已有完整 alarm pattern(LeadProcessorDLQAlarm / TranscribeDLQAlarm 等,threshold=1, 经 alarmTopic SNS → alarm-dispatcher Lambda → Lark 卡片)。lead-tracking 需要单独 issue 跟踪补 alarm。
Test/Prod 双写模式
Poller 仍同时写 prod 和 test Neon;NeonRetryProcessor 只消费 prod write DLQ 并重试 prod Neon。下游 lead-processor 在 callytics-infrastructure 里按当前 stack 的 environment 写对应 Neon,不是 prod Lambda 再 fire-and-forget 写 test。
SSM 参数: