Pipeline 全景分析
Historical code audit / Non-normative direction(2026-07-15):本文是 2026-05-31 的 pipeline snapshot,保留问题发现与 shared mutation 方向;其中旧
Task Pipeline Deliverable、task_progress_events、globalcloseResult、create_closed和独立lead_outreachTask 不再是 Target。Task 设计见 Task System Design V3;任何 Current 实现结论都必须重新以 live code 核查。读者:产品/工程团队成员、AI Agent — 理解系统所有数据 pipeline 的完整流转、交叉依赖、code/prompt 职责边界,以及从 Task Pipeline 设计出发的统一架构方向。
数据来源:2026-05-31 从
callytics-infrastructure、callytics-common、lead-tracking代码直接验证。以代码为准,文档描述如与代码不符已标注。与其他文档的关系:
- 系统架构总览 — repo 职责、runtime 部署
- Backend 数据管道 — 7 条数据线路的入口/出口/存储描述
- Backend 架构 Pattern — Layered vs Hexagonal 选型
- Task System Design V3 — 现行 Task / AI runtime contract
- Task Pipeline 设计 Deliverable — historical Task Orchestrator 设计依据,不是现行 contract
一、Pipeline 全景图
系统当前有 5 条活跃数据 pipeline + 2 条辅助 pipeline。每条 pipeline 的触发、处理、写入、下游消费路径如下。
二、逐条 Pipeline 详解
Call Analysis(通话分析)
触发:RC Webhook → ringcentralSubscriptionService → SQS transcribe-queue(5 分钟延迟等录音就绪)
阶段链:
3-Stage AI Prompt Pipeline(ai-analysis-processor/src/core/stages/):
写入表(代码验证):
下游 fan-out:当 follow_up_needed=yes 时,向 dailyBatchQueue.fifo 发 SQS 消息,触发 Contact Analysis 对该联系人重新做跨通话分析。
架构 pattern:Hexagonal — handler.ts 是 composition root,core/stages/ 含 pipeline use case,core/protocols.ts 定义 ports,infrastructure/ 含 AI/Neon/DDB/S3/SQS adapters。这是系统中最成熟的架构实现。
Contact Analysis(联系人跨通话分析 + Task 生命周期)
触发方式(3 种):
处理流程(handler.ts 第 154-331 行):
AI Prompt 输出(core/prompt-builder.ts):
Task 生命周期管理(代码路径:infrastructure/neon-repository.ts:writeAnalysisWithTasks()):
写入表:
架构 pattern:Semi-hexagonal — 有 infrastructure/protocols.ts 定义 ports(AIClient、ContactsReader、ContactsWriter),有 infrastructure/neon-repository.ts 和 infrastructure/ai-client.ts 作 adapters。但 core/ 只有 models.ts + prompt-builder.ts,没有显式的 pipeline/use-case 层。业务逻辑主要在 handler.ts:analyzeContact() 和 neon-repository.ts:writeAnalysisWithTasks() 里。
SMS Processing(短信处理)
触发:RC Webhook → ringcentralSubscriptionService → SQS message-processing-queue
处理(message-processor/src/core/message-processing.ts):
- 从 RC API fetch message 详情
- S3 保存 MMS 附件
- 原子
db.batch()写入 Neon
写入表:
关键特征:
- 不调用 AI — 纯数据持久化管道
- 不触发下游 pipeline — 不向 dailyBatchQueue 发消息
- SMS 内容要等到 Contact Analysis 的每日 cron 或 per-call 触发时才会被 AI 分析
架构 pattern:Hexagonal(轻量) — 有 core/ + infrastructure/ 分层,但没有 AI 相关 ports。
Lead Processing(线索处理)
触发:lead-tracking Lambda(IMAP 轮询邮箱)→ EventBridge LeadCreated → SQS → lead-processor
处理(lead-processor/src/core/persist-downstream.ts):
单个原子 db.batch() 写入 3 个实体:
关键特征:
- 不调用 AI — Task 创建是硬编码的确定性规则(5 分钟 SLA、high priority)
- 直接写 tasks 表 — 不经过任何 Task Orchestrator 或 contacts-analyzer
suggestedActions是模板 — 不是 AI 生成的,是代码里写死的默认建议onConflictDoNothing— 通过uq_tasks_source_lead唯一索引防重复
架构 pattern:Thin handler + core function — handler.ts 做 SQS 解析和 Neon client 管理,core/persist-downstream.ts 是纯业务逻辑(no AWS SDK import)。
Analytics / Periodic(定时报告和监控)
三、跨 Pipeline 依赖矩阵
3.1 共享 Neon 表 — 多 Writer 冲突风险
3.2 SQS 队列 — 生产/消费关系
3.2.1 Onboarding historical backfill 的正确位置
新店 onboarding 需要“把过去 14 天数据导入系统”时,不应该新写一条绕过现有 pipeline 的 ETL。
Calls 已经有可复用的机制:reconciliation-worker 支持 reconciliation_window.startTime/endTime,检测 missing calls 后把 synthetic webhook 投进 transcribe-queue。Onboarding 只需要一个更清晰的 operator trigger。
Messages 目前还没有同级的 historical backfill job。studio-api 能按 storeId/dateFrom/dateTo 读取 RingCentral message-store,message-processor 能持久化 message-store webhook,但中间缺“按日期拉历史消息并复用 message-processor 写入”的 job(跟踪:retaintive/callytics-infrastructure#1156)。
3.3 Pipeline 交叉点
核心交叉点:
- P1 → P2 的 per-call fan-out — ai-analysis-processor 通过
follow_up_needed=yes触发 contacts-analyzer 重新分析 - P2 读取所有其他 pipeline 的产出 — contacts-analyzer 聚合 calls(P1 产出)+ messages(P3 产出)+ leads(P4 产出)做跨维度 AI 分析
- contacts 表是 4 个 pipeline 的共享写入点 — 是系统中 writer 最多的表
- tasks 表有 2 个独立写入入口 — lead-processor 直接创建 lead_outreach,contacts-analyzer 管理 follow_up 生命周期
四、Code vs Prompt 职责边界(现状)
4.1 各 Pipeline 的职责分配
4.2 Task Pipeline Deliverable 提出的理想模式
Deliverable 定义了一个 8 步 responsibility chain(第 2 节):
核心原则:"AI proposes, code executes." AI 只输出 taskDecisions[] 提案,所有状态变更由 Orchestrator 代码验证并执行。
4.3 现状 vs 理想的差距
五、Schema 健康检查
5.1 Task 枚举对齐状态
5.2 关键字段覆盖
六、Gap Analysis — 从 Task Pipeline 推广到全系统
Gap 0(P0):NeonRetryProcessor DLQ retry 不发 LeadCreated 事件 ✅ 已修复(2026-06-08,lead-tracking#201)
原现状:lead-tracking 的 Neon 写入失败时进 DLQ,NeonRetryProcessor 负责 retry。但 retry 成功后不发 EventBridge LeadCreated 事件。
原后果:DLQ retry 成功的 leads 在 Neon leads 表里存在,但永远没有下游 contacts/tasks/timeline。
修复(lead-tracking#201 PR):
- 抽
publishLeadCreatedEvent到共享 modulesrc/lead-created-event.ts,poller 和 NeonRetryProcessor 共用同一份事件契约 - NeonRetryProcessor
persistLeadPipeline成功后调用publishLeadCreatedEvent(row, tenantId),使用 persistLeadPipeline 返回的 tenantId 避免下游再查 control-plane - 显式检查
PutEventsCommand返回的FailedEntryCount > 0(AWS SDK 在 200 OK + per-entry failure 时不抛,会静默丢失事件) - 给 NeonRetryProcessor Lambda 加
events:PutEventsIAM grant(之前只有 poller 有) - 区分 stage-tagged log:Neon 失败 vs EventBridge publish 失败,操作人员能清晰判断哪一步出问题
- 详细写入流程见 lead-tracking 写入流程
Gap 1:Task 写入缺乏统一 Orchestrator
现状(3 个独立入口):
- lead-processor — 直接 INSERT tasks(lead_outreach),hardcoded 5-min SLA、hardcoded suggestedActions,去重靠
uq_tasks_source_leadunique index - contacts-analyzer —
writeAnalysisWithTasks()处理 AI 产出的 CREATE/UPDATE/CLOSE,去重靠effectivePendingCategoriescode check + DB unique constraint - studio-api — HTTP API 直接写(manual CREATE/CLOSE/REOPEN),无 AI 判断,无共享 state machine
行为不一致:
建议方向:Deliverable 的 Task Orchestrator 应该成为 唯一的 task mutation 入口。lead-processor 调 Orchestrator 的 createTask(type='lead_outreach', source='lead', ...) 而不是直接 INSERT。lead-processor 的 deterministic create 可以作为 Orchestrator 的 bypass-AI 模式。
Gap 2:SMS Pipeline 没有 AI 分析
现状:message-processor 只做数据持久化,不做任何语义分析。SMS 内容要等 Contact Analysis 的 cron(最多 24 小时后)才会被 AI 看到。
问题:
- 客户发 "STOP" → 应立即标 DNC + 关闭所有 pending tasks → 现在要等到次日凌晨
- 客户 SMS 回复表达高意向 → 应立即触发 contacts-analyzer → 现在没有 per-message fan-out
建议方向:在 message-processor 末尾加一步 Code Filter(不是 AI):
- 检测 "STOP"/"UNSUBSCRIBE" → 立即写 DNC + 触发 task close
- 检测有意义回复(非自动回复、非单字) → 向 dailyBatchQueue 发 SQS 触发 Contact Analysis
Gap 3:contact_timeline 缺乏统一写入标准
现状:4 个 pipeline 各自写 contact_timeline,event_type 命名、actor_type 取值、newValue 结构都是各 pipeline 自行定义。
问题:没有中央的 "timeline event catalog",新增 event type 时容易命名不一致或遗漏字段。
建议方向:在 callytics-common 里定义 timeline event schema(event_type 枚举 + 每种 event 的 required fields),各 pipeline 的 timeline 写入通过 buildTimelineValues() helper 强制走 schema validation。
Gap 4:contacts 表 name trust scoring 分散,storeId null 策略不统一
Name Trust 现状:contacts 表的 firstName/lastName 由 6 个 writer 写入,各自带不同的信任分数:
这个 trust scoring 逻辑(SQL CASE WHEN + GREATEST)分散在 5 个 Lambda 代码库中,没有共享的 enforcement layer。如果某个 pipeline 的 UPSERT 忘了带 trust guard,低信任来源可以覆盖高信任来源的姓名。
storeId null 现状:4 个 pipeline 对 null storeId 的行为不一致:
建议:写一份 store-id-null-policy.md 明确每个 pipeline 的策略,并在 storeid-coverage-monitor 增加 per-pipeline breakdown 维度。
Gap 5:Prompt 输入不完整
现状:contacts-analyzer prompt 的 PENDING TASKS 输入只包含 taskId、typeCategory、priority、dueAt、first suggested action。不包含:
- 员工手动调整的 dueAt 变更历史(
dueAtChangelog) - task progress events(Deliverable 提议的新表,当前不存在)
- 来自其他 pipeline 的最新 SMS 内容(如果距离 cron 运行时刚收到)
问题:AI 做 task 决策时缺少关键上下文,导致可能重复创建已经在跟进的 task,或覆盖员工手动调整的 dueAt。
Gap 6:Per-Call AI(P1)和 Contact-Level AI(P2)的职责边界需要重新审视
现状:
- P1 Classification 输出
follow_up.needed+ reason 码 → 是 P2 的触发信号 - P1 Classification 输出
outcome.result(booked/cancelled/pending_follow_up)→ 是 P2 做 task 决策的证据之一 - P1 不直接操作 task → 这是设计意图("单通信息不足以判断 task")
问题:P1 的 follow_up.needed=yes 触发 P2 重新分析,但 P2 的 AI 还需要自己重新读一遍那通电话的内容来做判断。这意味着同一通电话被 AI 分析了 两次(P1 一次 + P2 一次),且 P2 的分析范围更大(全量 calls + messages + leads)。
这不一定是问题 — P1 和 P2 的分析粒度不同(单通 vs 跨通话),两次分析的目的不同。但 token 成本可以优化:P2 可以直接消费 P1 的结构化输出(category/outcome/follow_up)作为预处理过的事实,而不是从 transcript 重新推断。
七、架构一致性评估
建议:不需要所有 pipeline 都达到 P1 的 Hexagonal 深度。判断标准是"业务规则是否值得从 runtime/vendor 细节里剥离出来"(引自 backend-patterns.md §6)。P3 和 P4 业务规则简单,thin handler 够用。P2 的 task mutation 逻辑值得抽成独立的 Orchestrator core。
八、Prompt 设计对 Pipeline 架构的影响
当你为每个 pipeline 设计 prompt 时,需要清楚 prompt 能读到什么数据:
关键洞察:prompt 的输入质量直接取决于上游 pipeline 写入的数据质量。如果 SMS pipeline 不做 AI 分析,那 P2 prompt 里的 RECENT MESSAGES 就只有原始文本,没有预处理的 intent/sentiment。如果 lead-processor 不记录 progress events,P2 prompt 就不知道员工已经打了几次电话给这个 lead。
九、实际运行架构全景(2026-06 code-verified)
从代码直接验证的端到端架构 — 从外部源到 Neon 写入,含 shared mutation layer 和尚未接入的 Control Plane / studio-api。
十、建议的统一架构方向
基于 Task Pipeline Deliverable 的 "Task Orchestrator" 模式,推广到整个系统的统一架构:
这不是要重写所有 pipeline — 是一个渐进式的统一方向,按优先级排序: