Unified Pipeline Architecture 设计 Brief

这是什么: 给新 session 的设计任务。不是 session 总结,不是 bug 修复清单。 当前状态: Canonical unified pipeline brief。后续设计 session 优先读这份;unified-pipeline-reference-checklist.md 是 checklist/reference,unified-pipeline-research-appendix.md 是代码调研 appendix。 日期: 2026-05-31 前序工作: 2026-05-30 完成了 task pipeline 设计(team 会议讨论通过)。在设计 task pipeline 的过程中发现一个更大的问题:task 只是系统中的一条 pipeline,所有 pipeline(call/SMS/lead/contact analysis)共享同一批表但各自独立设计,职责划分和写入逻辑不统一。本 brief 是从 task pipeline 出发向外辐射的全系统设计。 前序材料:


你要做什么

审视 retaintive 现有系统的所有 pipeline 和 business object,检查哪些设计是对的、哪些有问题、哪些概念混在一起了。然后基于检查结果,设计一个统一的架构方向:怎么把散落的写入逻辑统一起来、代码和 AI 各自负责什么、现有模式和未来 tool calling 模式怎么衔接。

这不是从零设计一个新系统,是看我们现在有什么、哪些对哪些不对、然后怎么把它弄成统一的。

你的角色:一个精通 AI 技术、懂 SaaS 产品架构的 Principal Engineer / AI Architect。

阅读顺序

给团队 review 时,建议按"结论先行"读:

  1. Target State Architecture:先看最终理想状态长什么样。
  2. Phase 1 Implementation Design:再看第一阶段怎么落地。
  3. Current State Research:最后看代码事实和 gap evidence。

如果是新 session 要继续做 design work,再按这个顺序读:

  1. 本文:理解任务目标、问题边界、acceptance criteria。
  2. Current State Research:先看当前代码事实,不把 target docs 当实现。
  3. Phase 1 Implementation Design:看第一阶段怎么落地。
  4. Target State Architecture:看最理想的长期状态。

我们遇到的困境

在做 task pipeline 设计时发现了一个更大的问题。task 只是系统中的一条 pipeline,但我们还有 call analysis、contact analysis、SMS、lead 等多条 pipeline。这些 pipeline 各自独立设计,导致:

困境 1:每加一个新功能,改动成本指数增长。

比如现在要加 "SMS 触发 AI 分析" — 不是简单写一个新模块就行。你需要:在 message-processor 里加触发逻辑 → contacts-analyzer 的 prompt 要改 → contacts 表的写入逻辑要跟其他 pipeline 对齐 → task 创建要走 Orchestrator → timeline 要写新 event type。每一步都要跟其他 pipeline 的实现保持一致,但这些实现散落在不同 Lambda 里,每次都要翻一遍确认"别人是怎么写的"。

困境 2:代码和 AI 的职责边界是 case-by-case 画的,没有统一标准。

call analysis 的 3-stage prompt pipeline 里,triage 是 AI 做的;但 pre-triage(关键词匹配)是代码做的。contacts-analyzer 的 task 决策是 AI 提案 + 代码执行。lead-processor 的 task 创建是纯代码硬编码的。每条 pipeline 各自决定"这个该 AI 做还是代码做",没有一个统一的判断框架。以后再加新 pipeline 或改现有的,又要重新想一遍。

困境 3:不确定现有设计对不对。

task 和 contact 的某些概念可能混在一起了(比如 contact 的 actionNeeded 和 task 的 actionNeeded 是什么关系?)。call 的 disposition 和 task 的 progress 是不是同一件事?这些不确定需要做一次系统性检查——不是说设计错了,是不确定,需要验证。

困境 4:想加新能力(retry、onboard 扫描、batch 重分析)时,发现没有可复用的模块。

比如员工想 "retry 一次这通电话的 AI 分析",或者新店铺 onboard 想 "扫描过去两周的通话"。这些操作不应该是新 pipeline,应该是复用现有模块的不同调用方式。但现在模块之间耦合太紧,没有设计成可复用的。

困境 5:不知道未来方向应该怎么走。

现在是 one-time AI call 模式(把所有数据塞进 prompt 一次性调 AI),以后可能要用 tool calling 模式(AI 自己决定要查什么数据)。但不确定现在的架构怎么演进到那个方向,也不确定 tool calling 对我们有没有必要。

核心问题(一句话)

系统越来越复杂,每条 pipeline 都手动定义 + 硬编码设计,改动和调试成本指数增长。有没有更好的方式?哪些必须靠代码,哪些可以交给 AI / tool use / 其他工具?业界在怎么做?


Expected Output

新 session 完成后应该交付这些:

  1. Schema / DB Contract Check — 以 callytics-common/src/db/schema/ 为 source of truth,检查 contacts / tasks / calls / messages / leads / contact_timeline 当前 schema 是否支持目标职责。
  2. Current Design / Prompt Inventory — 列出现有 pipeline、prompt surface、structured output、code-only rules、writer 入口。重点不是数 prompt,而是判断每个 prompt owns 哪类 judgment。
  3. Business object 健康检查 — task、contact、call、lead、message、timeline 各自的职责边界是否清晰,有没有概念混淆。先检查,再设计。
  4. Current vs Target Workflow 图 — 至少两张图:现状(code/prompt/writer 混在一起)vs 目标(pipeline 负责触发和 judgment,shared modules 负责 mutation)。每个节点标清 [Code] / [Prompt] / [Hybrid] / [Shared Code]
  5. 两层架构设计 — Layer 1: Shared Mutation Modules(Task Orchestrator + Contact Writer + Timeline Writer + Policy Guard);Layer 2: Processing Capability / Invocation Modules(Call Analysis、SMS Signal、Contact Analysis Reprocess、Lead Downstream)。
  6. Code vs AI 职责分界框架 — 一个适用于所有 pipeline 的通用标准(不是每条 pipeline 单独决定)。
  7. Schema / API / AI Agent 反推检查 — 从 Human UI、Backend API、AI structured output、future tool API、voice agent、audit/replay 的视角反推数据模型是否对。
  8. 现有模式 vs tool calling 模式的设计兼容 — 确保方案 A 做好后,方案 B 可以无缝切换(见下方对比表)。
  9. Phase 1 Implementation Design — 明确第一阶段应该做什么,不只是愿景。至少覆盖 Task Orchestrator contract、record_progress API、task_progress_events、timeline projection、contacts-analyzer taskDecisions schema,以及 Call/SMS/Lead invocation contract 的边界。

不是写代码。是设计 + Max review 后才执行。

呈现要求: 产出必须让不在这个 session 里的人也能 10 分钟看懂核心结论。每个交付物先用一句话说结论,再展开细节。用表格和图而不是长段落。技术细节不省略但要有层次(先 what + why,想看 how 再往下读)。

Design Acceptance Criteria

设计合格的标准(不满足就不算完成):

  • 必须回答:现在是否应该改成 tool calling 模式。不能含糊
  • 必须列出:每个 business object(task、contact、call、lead、message、timeline)的写入入口 — 谁写、通过什么模块写、有没有统一
  • 必须覆盖:至少 3 个复用场景的设计(如:call retry、新用户 onboard 扫描、SMS 触发 AI 分析)
  • 必须覆盖:两层架构。不要只设计 shared writers,也要设计 Call/SMS/Lead/reprocess 这些 reusable processing capability 怎么被调用
  • 必须明确:Call Analysis Module / SMS Signal Module 不是新的 writer,它们是 invocation/capability contract,写入仍通过 shared mutation modules
  • 必须产出:Schema/API checkcurrent vs target workflow diagrams
  • 必须区分:Phase 1 必做(统一抽象)和 future-compatible 设计(tool calling、voice agent)

Non-goals(这次设计明确不做的)

  • ❌ 不实现 tool calling / agentic loop — 只确保架构兼容
  • ❌ 不设计 voice agent 产品 — 只确保 tool API 接口预留
  • ❌ 不重写现有 pipeline — 设计的是"共享模块"怎么插入现有 pipeline
  • ❌ 不做 production rollout 方案
  • ❌ 不改代码、不跑 DB 操作、不 push、不 PR

必须显式回答的冲突问题

Brief 里有几组方向天然有张力。不要默认调和,每个都逐一给出明确结论 + 理由:

  1. "每条 pipeline 各自优化" vs "统一成共享模块" — 统一的边界在哪?所有 contacts 写入都走 Contact Writer,还是允许简单 writer(如 message-processor 只更新 lastActivityAt)直接写?
  2. "one-time AI call(成本低、可预测)" vs "tool calling(灵活、可扩展)" — 什么场景用 A 什么场景用 B?是非此即彼还是共存?
  3. "prompt 做更多判断(灵活)" vs "代码做更多控制(确定性)" — 边界画在哪?contacts-analyzer 的 1040 行 prompt 是应该继续扩大还是应该拆分?
  4. "现在就做好统一抽象" vs "先做 task pipeline 再说" — 统一抽象是不是 premature?还是现在不做以后更难?
  5. "contact 的 actionNeeded / suggestedActions" vs "task 的 actionNeeded / suggestedActions" — 是冗余?还是不同语义?应该保留还是合并?
  6. "SMS 不触发 AI(省钱)" vs "有意义 SMS 需要即时跟进(及时性)" — 怎么定义"有意义"?是代码规则过滤还是 AI 判断?
  7. "Call/SMS 是不是也要 module" vs "只有 shared tables 才需要 module" — Call/SMS 不需要 writer module,但需要 processing capability / invocation contract。不要把这两类 module 混在一起。

背景:两种模式的对比

这是这次设计的核心框架。所有设计决策都在回答:这个操作在哪种模式下怎么实现。

方案 A: One-time AI call + 统一抽象(现在做)

代码提前聚合所有数据 → 一次性喂给 AI → AI 返回结构化 JSON → 代码验证执行

改进点(这次设计的目标):
1. Task Orchestrator — 统一 task 写入
2. Contact Writer  — 统一 contacts 写入(trust score 集中)
3. Timeline Writer — 统一 timeline 写入(event catalog 集中)
4. Schema-driven prompt — JSON Schema 定义 AI 能做什么

方案 B: Tool calling(以后做,基于 A)

定义 tools:
  - read_contact(phone, storeId)
  - get_recent_calls(phone, limit)
  - create_task(type, category, priority, ...)
  - close_task(taskId, closeResult)
  - update_contact(phone, storeId, fields)

AI 自己决定要查什么 → 调 tool → 拿结果 → 继续推理 → 再调 tool → 完成

每个 tool 的实现 = 调方案 A 的统一模块

对比

方案 A(现在)方案 B(以后)
信息代码提前聚合所有数据塞进 promptAI 自己决定要查什么数据
逻辑prompt 里定义"遇到 X 做 Y"AI 根据情况灵活选择工具
扩展加新功能 = 改 prompt + 改代码加新功能 = 注册一个新 tool
成本一次 AI 调用(prompt 大但调用少)多轮 AI 调用(每轮 prompt 小但调用多)
适用batch 处理(每天处理数百联系人)交互式(用户问一个问题、AI 诊断)
延迟可预测(一次调用)不可预测(多轮往返)

关键设计原则

A 做好了,B 自然就有了。 Task Orchestrator 就是 create_task tool 的实现;Contact Writer 就是 update_contact tool 的实现。不需要做两次设计。

A 和 B 可以共存。 batch 场景继续用 A(cost-efficient),交互场景用 B(flexible)。不是非此即彼。随时可以切,也可以两种模式同时跑。


怎么做

Step 1: 理解现有系统的 business objects

不要直接看代码实现,先从业务语义出发。 搞清楚系统里有几个核心 business object,每个的 lifecycle 是什么,它们之间的关系是什么。

需要回答的问题

Object该问的参考
Contact一个 contact 从出现到最终状态,经历哪些阶段?lifecycleStage(lead/member/churned)的转换规则是什么?谁有权改?callytics-common/src/db/schema/contacts.ts
Tasktask 的 lifecycle(open → closed)已经在 deliverable 里设计好了,这里检查跟 contact lifecycle 有没有冲突Task Pipeline Deliverable (Codex)
Call一通电话从 RC webhook 到最终存储,数据会被改几次、被谁改?callytics-common/src/db/schema/calls.ts
Leadlead 从邮件进来到被联系,lifecycle 是什么?lead 和 contact 的关系是什么?(同一个人同一个电话号码?)callytics-common/src/db/schema/leads.ts
MessageSMS 消息的 lifecycle。message 和 contact 的关系。SMS 内容是否应该触发 AI 分析?callytics-common/src/db/schema/messages.ts
Timeline Eventcontact_timeline 是 source of truth 还是 projection?它跟 task_progress_events(deliverable 提议的新表)的关系是什么?callytics-common/src/db/schema/contact-timeline.ts

特别要检查的

  • task 和 contact 的概念有没有混在一起(比如:contact 的 actionNeeded 字段和 task 的 actionNeeded 字段是什么关系?)
  • call 的 disposition(callState)和 task 的 progress 是什么关系?是两个独立的概念还是同一件事?
  • lead 和 contact 在什么节点合并?lead-processor 写 contacts 表时,如果 contact 已存在怎么处理?

Step 2: 设计统一抽象层(两层)

这两层解决不同的问题,但必须一起设计。

Layer 1: Shared Mutation Modules — "写的统一"

多条 pipeline 往同一张表写数据,写入逻辑必须统一。

Task Orchestrator(已有初步设计,需细化)

  • 输入:action(create/update/close)+ 参数
  • 输出:写入 tasks + contact_timeline(原子事务)
  • 保证:幂等(per-category dedup)、权限(storeId 隔离)、audit(timeline event)
  • 调用方:lead-processor、contacts-analyzer、studio-api、未来的 AI agent

Contact Writer(新设计)

  • 输入:phone + storeId + 要更新的字段 + 来源(trust level)
  • 输出:UPSERT contacts + contact_timeline
  • 保证:trust score 逻辑集中(WEBHOOK=20 < AI_TRANSCRIPT=40 < RC_API=60 < LEAD=80 < STAFF=100)、只升不降、GREATEST(lastActivityAt)
  • 调用方:所有写 contacts 的 pipeline

Timeline Writer(新设计)

  • 输入:event_type + event_category + entity + actor + payload
  • 输出:INSERT contact_timeline
  • 保证:idempotencyKey 去重、event catalog validation(只允许预定义的 event_type)
  • 调用方:所有写 timeline 的 pipeline

Policy Guard(横切层,不是单独 writer)

  • 输入:actor + storeId + requested action + current state + idempotency key
  • 输出:allow / reject / needs_review + reason
  • 保证:tenant isolation、DNC、权限、幂等、状态机转换、schema validation、AI confidence gate
  • 调用方:Task Orchestrator / Contact Writer / Timeline Writer,以及 future tool API
  • 设计原则:AI 可以提出 action,但 Policy Guard 决定这个 action 能不能执行

Layer 2: Processing Capability Modules — "调的统一"

同一套处理能力要被多种入口复用。这些不是"多 writer 争抢同一张表"的问题,而是"同一个能力怎么被 normal flow / retry / onboard / reconciliation 调用"的问题。

Call Analysis Module(从现有代码提取 invocation contract)

  • 核心问题:analyze_call(callId, mode) 怎么被调用?
  • mode 至少包括:normal(正常 pipeline 流转)/ retry(员工点"重新分析")/ backfill(新用户 onboard 扫描历史)/ reconciliation(修复遗漏)
  • 需要回答:retry 是否清 aiAnalysisCompletedAt?重新分析后是否重新触发 contacts-analyzer?如何避免重复 AI cost 和并发 double run?onboard 扫描是重走 transcribe 还是只对已有 transcripts 跑 AI?
  • 现状:入口主要是 S3 transcription event → SQS → ai-analysis-processor,已有 tryClaimAnalysis() 防并发。但还不是一个可被多入口调用的抽象模块

Contact Analysis Reprocess Module(现有 contact-analyzer-reprocess Lambda)

  • 核心问题:已有 fan-out 到 contacts-analyzer queue 的能力,但只是 contact-level reprocess,不是 call-level retry
  • 需要回答:reprocess 和正常 cron 走同一个路径吗?reprocess 后写入用同一套 Contact Writer / Task Orchestrator 吗?

SMS Signal Module(部分已有雏形,需要设计完整 contract)

  • 核心问题:一条 SMS 到了以后,系统怎么判断它是不是有业务意义,以及要不要触发后续 workflow
  • 现状雏形:message-processor 已有 exact "STOP" 检测(stop-keyword.ts)→ 设 DNC + 关 open tasks(neon-repository.ts:299
  • 需要设计的完整 flow:
inbound SMS
  → store message(message-processor 现有)
  → SMS Signal Module
      → exact STOP: code DNC(现有)
      → trivial (thanks/ok/emoji): no-op
      → meaningful: trigger contacts-analyzer
      → natural-language DNC: AI judgment + code DNC guard
  • 关键设计原则:SMS AI analysis 不应该成为一个新的 task writer。它应该只产出 signal(meaningful / DNC / booking intent / cancellation intent / complaint),然后交给 Contact Analyzer + Task Orchestrator + Contact Writer + Timeline Writer 去做真正的状态变化
  • 需要回答:code filter 做多少(exact keyword match)vs AI 做多少(natural language intent)?debounce?

Lead Downstream Invocation(现有 lead-processor,检查 contract)

  • 核心问题:lead 进来后触发的下游链(contacts UPSERT + task create + timeline)是否走 Layer 1 的共享模块
  • 现状:lead-processor 直接 db.batch() 原子写 3 张表,没走 Task Orchestrator
  • 需要回答:lead-processor 应该改成调 Task Orchestrator 创建 lead_outreach task 吗?还是保持现有的 hardcoded 逻辑(因为它是纯确定性的,不经 AI)?

两层的关系:Layer 2 的 module 跑完之后,结果通过 Layer 1 的 module 写入。比如 Call Analysis Module 分析完一通电话 → 调 Contact Writer 更新 contacts → 调 Timeline Writer 写 timeline。SMS Signal Module 判断"有意义" → 触发 Contact Analysis → 结果再通过 Task Orchestrator 写 tasks。如果只设计了 Layer 1 没设计 Layer 2,以后做 retry 时会发现 retry 绕过了统一写入。

Layer 2 的 Invocation Contract 清单(design 必须覆盖,实现可以分阶段):

Contract解决什么
Call Analysis Invocationanalyze_call(callId, mode) — normal / retry / backfill / reconciliation
SMS Signal Invocationevaluate_sms(messageId) — STOP / trivial / meaningful / NL-DNC
Contact Analysis Reprocessreanalyze_contact(phone, storeId, reason) — cron / per-call / on-demand / reprocess
Lead Downstream Invocationpersist_lead_downstream(leadRow) — 是否走 Task Orchestrator

Phase 1 执行策略:第一刀先做 Layer 1(Task Orchestrator + Contact Writer + Timeline Writer)。但 design 必须把 Layer 2 的 invocation contract 全部写进去,否则后面做 retry/onboard/SMS AI trigger 时又会发现缺抽象。执行可以分阶段,设计不能只看一个角落。

Step 3: 检查每条 pipeline 在统一架构下的职责

改造后每条 pipeline 的职责应该是:

Pipeline 专有职责:
  触发检测 → 数据获取 → AI 调用(如果需要)→ 解析 AI 输出

共享模块调用:
  → Contact Writer(更新 contacts)
  → Task Orchestrator(创建/更新/关闭 tasks)
  → Timeline Writer(记录审计事件)

需要逐条 pipeline 画出来,标清楚哪些是"专有"哪些是"共享"。

Step 4: 定义 Code vs AI 的统一分界框架

建立一个适用于所有 pipeline 的通用判断标准:

判断维度Code 做AI 做Hybrid
需要确定性结果?
需要理解自然语言?
涉及状态转换?✅(执行)✅(提案)AI 提案 + Code 执行
涉及金钱/权限/DNC?
需要分类/评分?AI 分类 + Code 校验
需要生成文本?
有已知规则可以写死?

把这个框架应用到每条 pipeline 的每个决策点上,产出完整的职责矩阵。

Step 5: 设计可复用操作

列出系统里应该被设计成"可以被任何入口调用"的操作。每个操作定义:谁能调、输入什么、做什么、返回什么。

初步清单(需要验证和补充):

操作触发方做什么
分析单个联系人cron / per-call / on-demand / retry / onboard聚合数据 → AI 分析 → 写 contacts + tasks
重新处理一通电话员工点 retry / reconciliation重新跑 transcribe → AI analysis → 写 calls
新用户 onboard 扫描新店铺接入扫描过去 N 天的 calls/messages → batch 分析
SMS 触发 AI 分析SMS 到达且内容有意义触发 contacts-analyzer 重新分析该联系人
手动创建/关闭 task员工 UI调 Task Orchestrator
批量重分析管理员 / prompt 更新后对一批联系人重跑 contacts-analyzer

Step 6: 验证 tool calling 兼容性

对 Step 5 的每个操作,验证:如果以后改成 tool calling 模式,这个操作能不能直接变成一个 tool?

方案 A(现在):
  contacts-analyzer 代码里调 TaskOrchestrator.create(...)

方案 B(以后):
  AI 调用 create_task tool → tool 实现里调 TaskOrchestrator.create(...)

如果答案是"能",说明设计是 future-compatible 的。如果答案是"不能",需要调整 Step 2 的模块设计。

Step 7: 从 API / AI Agent 视角反推数据模型

这一步不是要现在做 agent,而是用 future callers 反推 Phase 1 的接口有没有设计对。一个 shared module 如果只能被当前 Lambda 调用,不能被 API / CLI / AI agent 复用,那它还不是稳定抽象。

需要逐个检查:

调用方需要看到什么反推问题
Human UItask 状态、progress、close result、timelineUI 是否需要解析 JSONB 才知道发生了什么?
Backend API结构化 task/contact/call/message 对象API contract 是否能表达 retry/backfill/reprocess?
AI structured outputaction proposal + typed paramsprompt output schema 是否等同于 module input schema?
Future tool APIcreate_task / close_task / update_contact / analyze_calltool 是否只是 shared module 的薄包装?
Voice agent低延迟、幂等、可拒绝、可审计每个 tool call 是否有 policy gate 和 idempotency key?
Audit / replay谁提议、谁执行、为什么拒绝是否能从 events 重放一次决策链?

设计结论要明确区分:

  • Domain API:给产品和后端用,比如 TaskOrchestrator.closeTask()
  • AI Action Schema:给 prompt structured output 用,比如 taskDecisions[].action = "close"
  • Tool API:给 future agent 用,比如 close_task(taskId, result, evidence)

三者可以共享同一套 schema / validation,但不一定暴露完全一样的字段。


思考方式

"Code as Guardrail, AI as Judgment" — 业界在用的 production pattern,多个名字:

  • "AI Proposes, Deterministic Layer Disposes"
  • "Pre-execution Approval Gate"
  • "Graduated Autonomy"

核心:

职责谁做为什么
状态机转换代码必须可审计、可逆、exactly-once
权限检查代码(Policy rules)不能被 prompt injection 影响
幂等代码(DB constraints、dedup)LLM 不能保证确定性行为
意图分类AI需要语义理解
内容生成AI需要自然语言能力
动作参数提取AI 提取 + 代码验证Hybrid:AI 提取,Zod/JSON Schema 校验

生产数据:使用 structured output(JSON Schema 约束 AI 输出格式)的 agent 动作成功率 95-99%,而解析非结构化文本的只有 70-85%。

你的 task pipeline deliverable 里的 "AI proposes, code executes" 完全是这个 pattern。 你已经在 contacts-analyzer 上做了(taskDecisions[] = AI 提案 + 代码执行)。这次的任务是检查其他 pipeline 是不是也该统一到这个模式。

"Schema-first" 对应到这里就是:用 JSON Schema / Zod 定义 AI 能做什么操作、每个操作的参数是什么、哪些值是允许的。规定一堆 JSON 然后按照 JSON 的结构生成 prompt 和数据库结构。


前序调研发现的关键事实(设计时要知道的)

这些不是"要修的 bug 清单",而是做 system design 时的输入事实——你需要知道系统实际上怎么运作的,才能设计对。不要逐个解决这些小问题——做好顶层设计,这些问题会同时被解决。以面破题,不以点破题。

现有 AI signal chain 怎么触发 task:

  • Per-call AI stage 分析每通电话,产出 fu=yes/no (follow-up needed) + out=success/attempted/... (outcome)
  • Contact analysis stage 消费 per-call 输出,做跨通话的 contact 级别分析,产出 taskDecisions[] (create/close/update)
  • 这个 fu= 信号是 AI 对 AI 的传递——per-call AI 告诉 contact-level AI "要不要跟进"
  • 注意:这里的 per-call / contact-level 是现有 AI 处理链,不要和上文的 Layer 1 Shared Mutation / Layer 2 Processing Capability 混用

task 创建只有 3 条路径:

  1. lead-tracking 直接建(lead_outreach,不经 AI)
  2. contact analysis AI 建(taskDecisions[].create)
  3. 员工手动建(studio-api)
  • 所有路径都只产 status='pending'。没有"一步建成 closed"的路径

task 关闭有 2 条路径:

  1. AI auto-close(contact analysis 的 taskDecisions[].close,必须引用已有 taskId)
  2. 员工 manual-close(studio-api close.ts)

contacts 表有多个 writer:

  • 每个 pipeline 各自实现 UPSERT 逻辑和 trust score
  • name trust score(WEBHOOK=20 < AI_TRANSCRIPT=40 < RC_API=60 < LEAD=80 < STAFF=100)散落在不同 Lambda

contact_timeline 是统一审计线:

  • 16 种事件类型,覆盖 task 全生命周期
  • 多态 actor(system / call_analysis / contact_analysis / staff / lead_webhook)
  • AI forensic 字段(aiRunId / aiModelUsed / aiConfidence)

disposition 数据已经在 calls 表里:

  • calls.callState = human_conversation / voicemail / no_answer / busy_signal / system_error
  • calls.followUpNeeded + calls.followUpReasons = AI 判断的 follow-up 信号
  • 这些都是 per-call 的,不需要在其他表重复存

前序调研成果 — 分三类: Facts / Hypotheses / Calibration

前序 session 做了大量调研和代码分析。这些是校准用的,不是替代你自己思考的答案。 按可信度分成三类:

Facts(已用代码/数据验证,可直接信)

事实验证方式
contacts 表有多条 pipeline 同时写入,trust score 逻辑各写各的读 5 个 Lambda 的 neon-repository
tasks 表有 3 个独立写入入口(lead-processor / contacts-analyzer / studio-api),去重策略不一样读代码
contact_timeline 有多个 writer,event_type 各自定义读代码
storeId null 时各 pipeline 行为不一致(有的跳过、有的继续、有的 DDB fallback)读代码
contacts-analyzer prompt 已 1040 行,每次改一个地方可能影响一大片读 prompt-builder.ts
call analysis 有 pre-triage 代码层(关键词匹配,节省 ~40% AI 调用)读 pre-triage.ts
Task Pipeline Deliverable 的 "AI proposes, code executes" 已设计并讨论通过team 会议 2026-05-30

Hypotheses(前序分析的推断,需要验证后才能采纳)

假设来源需要验证什么
所有 contacts 写入应该统一成 Contact Writer modulepipeline 全景分析会不会 over-engineering?简单 writer 是否值得走统一模块?
所有 timeline 写入应该统一成 Timeline Writer modulepipeline 全景分析event catalog 是否需要预定义所有 event_type?
contact 的 actionNeeded 和 task 的 actionNeeded 概念混淆代码对比具体怎么混的?是否需要去掉一个?
方案 A 做好了方案 B 自然就有架构推理是否真的无缝?有没有 A 的设计会卡住 B?

web search 是校准用的,不替代系统设计结论。最终判断必须回到 retaintive 自己的架构和成本。

校准点业界怎么做来源
大厂 AI SaaS 架构每个用例一个专门 agent + 共享基础设施,不是一个万能 agent(Salesforce Agentforce 12,000+ 客户、HubSpot Breeze、Intercom Fin)Salesforce · HubSpot · Intercom
LangGraph / Temporal 生产落地真的有人用(LinkedIn/Uber/Replit),但 prototype→production 落差巨大,40% 项目被取消Gartner · LangChain blog
"Code as Guardrail, AI as Judgment"成熟 pattern,structured output 成功率 95-99% vs 非结构化 70-85%多源
2-3 人小团队Level 2-3 agent realistic(单 agent + tool calling),multi-agent overkill多源
Voice agent 后端需求HTTP endpoint 接收 tool call POST + 快速响应(<2s) + 幂等。Vapi/Retell/Bland 三家核心模式一样Vapi · Retell · Bland docs
Tool calling 生产模式OpenRouter 标准化 tool calling 接口,poly-model routing(Opus 规划 + Sonnet 执行)OpenRouter docs
小团队推荐 patternRouter Agent:便宜模型(Haiku,1/10 成本)做意图分类 → 转给专门逻辑多源

已知的系统现状

不展开,只列要点。详细 trace 见 Pipeline 全景分析

5 条活跃 pipeline

Pipeline触发AI?写哪些共享
Call AnalysisRC Webhook → SQS3 prompt(triage/classify/coaching)contacts, contact_timeline
Contact AnalysisdailyBatchQueue + cron + on-demand1 prompt(contact-analyzer)contacts, tasks, contact_timeline
SMSRC Webhook → SQSCode-only exact STOP;暂无 SMS AIcontacts, contact_timeline
LeadIMAP → EventBridgecontacts, tasks, contact_timeline
AnalyticsEventBridge cron读 call-analysis不写共享表

每条 pipeline 也写自己的专有表(calls、messages、leads、DDB call-analysis),那些是 pipeline 内部的事。这次关注的是多条 pipeline 同时写的表怎么统一。

多条 pipeline 同时写的表(需要统一的部分)

哪些 pipeline 写问题
contactsCall Analysis + Contact Analysis + SMS + Lead + studio-apiname trust score 逻辑散落在每个 pipeline 里各写各的
contact_timelineCall Analysis + Contact Analysis + SMS + Lead + studio-apievent_type 各自定义,没有统一的 event catalog
tasksContact Analysis + Lead + studio-api去重策略、closeResult、dueAt 计算全不一样

Max 在对话中表达的方向

这些是他的想法和倾向,不是 finalized spec:

  • 从 task pipeline 设计出发,发现所有 pipeline 需要统一设计
  • "规定一堆 JSON 然后按照 JSON 的结构生成 prompt 和数据库结构比较靠谱,别的感觉想的很好都没落地过"
  • 认同 "Code as Guardrail, AI as Judgment" pattern;该原则已写入产品设计准则
  • 想有一个最终构思(终态方向),然后看现有的怎么靠拢
  • 先做一个版本 buying more time,设计要符合 long-term scalability
  • 现在没有 voice agent,但以后可能有,架构要兼容
  • 先把基础的东西(task、lead、contact、call)弄清楚,再拓展高级功能
  • 关心可复用性:call retry、新用户 onboard 扫描、批量重处理 这些操作能不能设计成模块
  • 不确定现有设计对不对,需要做系统性检查(不是说设计错了,是需要验证)
  • ai-integration-roadmap.md 是很久以前写的,可能 outdated,只作过去参考

材料清单

只有代码是 source of truth。其余都是不同人在不同时间的想法,不一定对。

A. 代码 (✅ Source of Truth)

文件看什么
callytics-common/src/db/schema/ 全部文件所有 business object 的 schema 定义、枚举、约束
callytics-infrastructure/lambda/*/src/每条 pipeline 的实际实现 — trigger、读写、AI 调用
callytics-infrastructure/lambda/contacts-analyzer/src/core/prompt-builder.tscontacts-analyzer prompt(1040 行)— SECTION 4 是 task 决策
callytics-infrastructure/lambda/ai-analysis-processor/src/core/stages/per-call AI 3-stage pipeline(pre-triage → triage → classify → coaching)
callytics-infrastructure/lambda/lead-processor/src/core/persist-downstream.tslead 下游写入(contacts + tasks + timeline 原子 batch)
callytics-infrastructure/lambda/message-processor/src/core/SMS 处理逻辑
studio-website-monorepo/apps/api/src/routes/v3/studio-api 的 task/contact 读写 endpoint

B. 设计文档(已确认的方向)

文件看什么
Task Pipeline Deliverable (Codex)Task Orchestrator 设计、4-Layer Object Model、Code vs Prompt 职责矩阵
Task Pipeline 设计 Brief设计方法论和 Max 的方向性想法
产品设计准则Schema-first → API-first → State-machine-first → Prompt-last
Pipeline 全景分析5 条 pipeline 的代码级 trace + gap analysis

C. 参考材料(⚠️ 不一定对,当参考不当结论)

文件来源说明
AI 集成路线图早期设计agentic loop + 6 gym tools + 5 级金字塔。Max 说可能 outdated,只作过去参考
Voice Agent 可行性调研报告接 Vapi/Retell 的技术可行性。voice agent 不是这次的目标,但架构要兼容
Backend 架构 Pattern架构文档Layered vs Hexagonal 选型标准
Task 枚举清单V2 文档closeResult 18 值、typeCategory 9 值 — 代码 SoT

D. Neon Test 环境 (✅ 可查真实数据)

  • Project: callytics-test (restless-boat-70724564, us-west-2)
  • 用 Neon MCP 跑 SELECT 验证假设,不写入

E. AWS (✅ 已登录)

  • Account: 237206024479, AdminAccess role
  • 可看 Lambda / CloudWatch,用于验证 prompt 实际输出和 pipeline 行为

补充维度(设计时别漏)

多门店 / store 级隔离

所有共享模块必须以 store_id 作为隔离键。Task Orchestrator 的幂等约束 uq_tasks_pending_contact_category(contactPhone, storeId, typeCategory) 做。Contact Writer 的 UPSERT 以 (phone, storeId) 为 PK。同一个电话号码在不同门店是独立的 contact。

前台一天的实际操作流程(场景 robust 的基础)

设计必须对得上前台真实的工作方式,不能只看代码:

  1. 早上开工:打开系统看待办清单——"今天有哪些客人要跟进"
  2. 白天逐个处理:挑一件 → 打电话/发短信 → 当场办成 / 没当场办成回头来 / 没打通再试 / DNC
  3. 当场来的:客户直接打进来当场预约 — 没经过待办但算业绩
  4. 一天结束:店长看战绩(完成几件、几个签约/挽留)

未来 AI agent 的操作流程(架构要兼容但不是这次的目标)

  • AI agent 自动外呼 → 自动判断对话结果 → 调 Task Orchestrator 关闭/更新 task
  • 需要 executorType(human / ai_agent / system)区分"谁做的"
  • 低置信度的 AI 决策需要 human review
  • 共享模块的设计要让 AI agent 和人类走同一个接口

SMS 处理方式

现状:message-processor 存储 messages,并对 exact inbound STOP 做 code-only DNC + task close cascade;它还不触发 AI / contacts-analyzer 来处理 meaningful SMS。

设计时要思考

  • "STOP" 这种关键词是否应该在 code filter 层立即处理(不等 AI)?
  • 有意义的 SMS 回复是否应该触发 contacts-analyzer?
  • 怎么区分"有意义"vs"无价值"(thanks / ok / emoji / 自动回复)?
  • SMS 还没有完全想清楚,先 focus 电话。但 pipeline 设计要给 SMS 留位置

可复用操作的场景

以下场景直接验证共享模块的设计是否合理:

场景触发方需要调哪些模块
员工说"我想 retry 这通电话的 AI 分析"UI 按钮重跑 AI analysis → Contact Writer → Task Orchestrator
新用户 onboard,想扫描过去两周的通话onboard 流程batch 调 transcribe + AI analysis + contacts-analyzer
prompt 更新了,想对所有联系人重分析管理员batch 调 contacts-analyzer
SMS 收到 "STOP"message-processorCode filter → Contact Writer(标 DNC)→ Task Orchestrator(关所有 pending tasks)

如果这些场景都能通过"调共享模块"完成,而不是写新的专门 pipeline,说明设计是对的。

API / CLI / AI agent 可读性

设计时要考虑未来 API 的形状。以后可能有第三方系统、CLI、客户的 AI agent 通过 OAuth/API 读写 task 和 contact。

设计时要显式回答:

  • 外部调用方如何区分 task 的生命周期状态 / 执行进展 / 最终业务结果?
  • API 返回应该是结构化对象,不是要求调用方解析 timeline JSONB 或 prompt 文本
  • 共享模块的接口设计要能直接变成 tool calling 的 tool

Deep Gaps / Open Design Questions

这些不是立即要写代码的 todo,而是 design review 必须讨论清楚的风险点。很多后续 bug 会从这些边界没定义开始。

Gap为什么重要设计时要回答
Processing intent modelnormal / retry / backfill / reconciliation / on-demand / task_close 如果只是散落的 string,后面会变成不可控分支是否需要统一 processingIntent enum?每个 intent 能不能触发 AI、能不能写 shared state、能不能覆盖旧结果?
Retry semanticsretry 一通电话可能影响 callscontactstasks、timeline 和 AI costretry 是追加 attempt 还是覆盖字段?是否清 aiAnalysisCompletedAt / s3AnalysisPath?是否重新触发 contacts-analyzer?
Source of truth vs projectiontask_progress_eventscontact_timeline、contact action summary 可能重复表达同一件事哪张表是 source of truth?哪些只是 projection?projection 失败能不能重建?
storeId unresolved policy现在不同 pipeline 对 storeId = null 行为不一致是 skip、quarantine、fallback、retry,还是允许某些 source 先落库后补 store?
Idempotency taxonomy每个 Lambda 自己做 dedup 会导致跨入口重复写idempotency key 按 event、message、call、contact-run、task-action 还是 external request 生成?
Human override vs AI race员工 close task 的同时 AI auto-close / update task,可能覆盖人的判断staff action 是否永远优先?AI proposal 发现 task 已被人工修改时是 reject、append note 还是 reopen?
Payload catalog / schema versioningcontact_timeline.newValue、prompt output、future tool schema 都会演进每个 event payload 是否要 version?旧 prompt output 如何兼容新 module input?
Permission / RBAC同一个 shared module 未来会被 staff、system、AI agent、external API 调用actor 能力是否分级?AI agent 是否只能 propose,不能直接 execute 高风险 action?
Cost budget / debounceSMS meaningful reply、tool calling、batch reprocess 都可能放大 AI 调用量是否需要 per-store budget、cooldown、debounce window、batch cap?
Rollout / migration现有 enum 和历史数据已经存在,不能只设计新世界Phase 1 怎么兼容旧 contact_timeline.event_type、tasks status、contacts action summary?
Testing strategy共享模块一旦错,会影响所有 pipeline是否需要 contract tests、golden prompt output tests、replay tests、idempotency/race tests?
Analytics definitionstask 完成、call 成功、contact converted、business win 可能不是同一层概念dashboard 应该按 attempts、objectives completed、closed won、retained member 里的哪个口径算?

需要特别小心的一点:不要把所有问题都塞进 Task Orchestrator。Task Orchestrator 只管 task lifecycle 和 task progress;Contact Writer 管 contact aggregate state;Timeline Writer 管 audit/projection;Call/SMS Module 管 processing invocation。职责一旦混了,Phase 1 会重新变成一个大而全的 pipeline。


设计时要考虑到(high-level,不需要写具体实现)

  • AI 判断失败时的 fallback — prompt 返回无法解析 / Lambda timeout / Zod 校验不过,代码怎么兜底?
  • contacts-analyzer 有 6 个触发路径 — cron / on-demand / P1 followUp / Task close / T5 showed / T6 trialed。设计要覆盖所有入口
  • 并发 / 幂等 / race condition — daily batch 和 per-call 可能同时分析同一个 contact;lead-tracking 建 task 的同时 AI 也想建;员工 close 的同时 AI auto-close
  • 成本控制 — 方案 A 的成本是可控的;方案 B 多轮调用的成本需要 budget cap
  • 业界做法是参考不是圣经 — 2026 年的技术限制和 2024 年不同(LLM 成本降 10-50x)。不要因为"Salesforce 这么做"就照搬
  • 产品设计准则更新 — "Code as Guardrail, AI as Judgment" 这个 pattern 已写进 product-design/index.md,后续 AI workflow 默认按这个原则设计
  • 设计边界:如果需要改,什么都可以改 — 所有 pipeline / schema / prompt / 前端都在 scope。这是 test 环境

约束

  • 先理解全局再看局部 — 先把所有 business object 的关系搞清楚
  • 先设计再写代码 — 画图、确认逻辑、Max review 之后才动手
  • Code as source of truth — 每个结论先读代码验证,不信文档 narrative
  • 不要以点破题,以面破题 — 不要逐个修小问题,设计好统一抽象后小问题自然消解
  • 方案 A 优先 — 先做好统一抽象(one-time AI call),tool calling 留给交互式场景
  • 检查而不是重写 — 检查现有设计对不对,有问题的修正,没问题的保留
  • 从正确性出发 — 不要因为现在是这样就继续走,test 环境什么都可以改
  • 保持 high-level — 这是 system design,不是 code review
  • 不自动 push/PR