Unified Pipeline 统一设计
Historical architecture snapshot(2026-07-23 校准):本文保留 shared mutation module、Policy Guard、Timeline Writer、0 / 1 / N decisions 和
create_closed等 Task V2 架构依据。task_progress_events、globalcloseResult扩展和 Lead relay 等旧 proposal 不能直接实施;其中create_closed是当前必须保持的一等产品语义。当前路线见 Task V2+ 工程审计与实施基线;实施 pipeline 变更前必须重新核查 live code。日期: 2026-06-01 状态: 合并 Opus + Codex 两份独立设计后的统一版本,待 Max review 环境前提: 本文以当时 test / new architecture 调研为主;production Task V2 行为需要单独 audit,不默认与 test 一致。 设计准则: 产品设计准则(Schema-first → API-first → State-machine-first → Prompt-last / Code as Guardrail, AI as Judgment) 历史前置设计: Task Pipeline Deliverable (Codex) — superseded Task contract,仅用于追溯 Task Orchestrator 的演进依据
原始材料(保留在本目录,不删除):
- Target State (Opus) — current state inventory + target state + Phase 1
- Target State (Codex) — 理想终态架构
- Phase 1 (Codex) — Phase 1 落地方案
- Research 材料 — canonical brief、现状调研、checklist、appendix
1. Historical proposed Target State
所有写入方通过同一套共享函数改变系统状态,不直接碰共享表。
什么是 module(共享函数)
现在系统的现状:很多写入逻辑其实已经收口了 —— name trust scoring 已集中到 NAME_TRUST 共享常量(name-trust.ts),dueAt 计算已抽出 computeDueAt(3 个 Lambda 共用),timeline 写入已统一走 buildTimelineValues / buildContactTimelineInsertSQL,DNC cascade 已共享(dnc-cascade.ts)。剩下两个真实缺口:(1) task mutation(create/close/update)还没有统一的 Task Orchestrator,状态机语义散在 contacts-analyzer / lead-processor / studio-api 的各自实现里;(2) task_progress_events 表不存在,导致 no_answer / left_voicemail 这类进展被错塞进 closeResult。本设计收口这两个缺口,不是从零统一(大部分共享 helper 已就位)。
Module 就是把这些散落在多个地方的写入逻辑,抽到一个共享的 TypeScript 文件里(住在 callytics-common/src/domain/)。不是新 Lambda,不是新 SQS,不是新进程。以 task mutation 为例,改前改后的区别:
每个 caller 传不同的 action + payload(create / create_closed / close / record_progress ...),统一的是规则(去重、状态机、timeline 审计),不统一的是内容。
Terminology note: 本文里的 Task Orchestrator 只指 task domain 的共享函数 / module,负责 create / create_closed / close / update / record_progress / reopen 这些 task mutation。它不是 workflow engine,不是新的 Lambda,也不是整个 unified pipeline 的总调度器。下文如果说 “Orchestrator”,都应该理解为 Task Orchestrator;为避免混淆,正文尽量写全名。
架构总览
只画写入路径。
当前写入全景图(现状)
跟上面的 target state 对比:现在每个写入方物理上直接执行自己的 SQL 碰共享表,没有统一入口函数。
注意区分"物理路径"和"逻辑共享":下图画的是物理写入路径(谁的 SQL 写哪张表)。逻辑层面 trust scoring(
NAME_TRUST)、timeline(buildTimelineValues)、dueAt(computeDueAt)、DNC cascade 其实已经抽成共享函数——但每个 caller 仍各自调用并自己执行 SQL(这就是「弱版本」:共享函数、各 caller 自己写)。Target State 的统一不是改成"物理上只有一个地方能写",而是让 task mutation 也经过共享的 Task Orchestrator。
Prompt Pipeline
先把 prompt 数量说清楚,避免把 code stage、current prompt、future read-only prompt 混在一起:
- Current live code / Phase 1 write path: 6 个 active decision surface = 1 个纯代码 stage + 5 个 AI prompt。5 个 AI prompt 是
Triage/Classify/ conditionalVerify/Coaching/Contact Analyzer - Task Playbook: planned read-only prompt,来自 Task Pipeline Deliverable。它读 task + contact 生成员工话术,不产出 DB mutation,不是 Phase 1 写入路径的前置条件
- Target State:
Contact Analyzer拆成Contact Profile Prompt+Task Decision Prompt(拆 prompt 文本/组装,仍是一次 LLM 调用、一个 JSON 输出 —— 见下方「拆 prompt 的方式」)。如果加上Task Playbook,target state 是 1 个纯代码 stage + 7 个 AI prompt surface
只有 Contact Analyzer 当前会产出 task/contact business mutation proposal。Phase 1 不拆它,只改下游写入方式;Target State 再把它拆成 Contact Profile Prompt 和 Task Decision Prompt(拆文本,不拆调用)。
Prompt / Decision Surface Inventory
Task 相关的 AI 不应该混成一个概念:
- Task Decision Prompt: 决定是否 create / create_closed / close / update / record_progress。它产出 mutation proposal,必须走 Policy Guard + Task Orchestrator。Phase 1 先继续嵌在 Contact Analyzer 里;Target State 再拆成独立 prompt stage
- Task Playbook Prompt: task 已经存在以后,生成 staff-facing 话术和执行建议。它只读 task/contact/call context,不写 DB,不应该绕过 Task Orchestrator
现在(Before)
Phase 1:核心是写入层;拆 prompt 并行外包
Phase 1 的核心改动是写入层(Task Orchestrator / Contact Writer / Timeline Writer / Policy Guard + task_progress_events)。
拆 prompt(画像段 / task 决策段)已经并行外包给独立 engineer 做,跟写入层改造解耦 —— 拆完直接拿来接入即可。两者唯一的接触点是 TaskDecision schema 要加 record_progress / create_closed 等 6 个 action 命名(见下方「Phase 1 实施步骤」),这个 schema 改动是 Phase 1 自己的 deliverable,不管 prompt 是否已拆都要做。
下图展示写入层改造(prompt 是否已拆不影响这张图的写入路径):Contact Analyzer 一次 AI 调用输出 contact 字段 + taskDecisions[],下游从自己拼 SQL 改成调共享写入层。
历史方案为何不直接一步到当时的 Target State
这里分成 Current / Phase 1 / Target State,不是因为 Target State 不对,而是因为这两类改动风险不一样:
所以 Phase 1 和 Target State 的区别是:
- Phase 1 是 mutation architecture refactor:先把"谁能写 tasks / contacts / timeline、怎么写、怎么审计"定下来。它解决的是 correctness 和 shared infrastructure
- Target State 是 prompt architecture refactor:把"哪个 prompt 负责画像、哪个 prompt 负责 task decision、哪个 prompt 负责 playbook"拆清楚。已经并行外包给独立 engineer,Phase 1 不阻塞这条线
测试环境没有 prod migration 压力,所以 Phase 1 可以做得更快;拆好的 prompt 一交付就可以接入。两条线的接触点只有 TaskDecision 的 6 个 action 命名(见「Phase 1 实施步骤」),双方对齐这一点即可独立推进。
Historical Target:Contact Analyzer 拆成画像 + task 决策两段
现在 Contact Analyzer 是一个 1000+ 行的大 prompt,同时做客户画像分析和 task 判断。读真实 prompt 后确认:画像段(SECTION 1-3)和 task 决策段(SECTION 4)已经是天然两块,中间只有一条单向依赖 —— task 决策需要先知道 lifecycleStage 才能选 typeCategory。所以"拆开"是顺势而为,不是发明新结构。
拆 prompt 的方式:拆文本/组装,不拆 LLM 调用(经 Codex + Gemini cross-validate)
拆有两种做法,必须区分清楚:
采用拆文本。理由:当前痛点是 prompt 长度和可维护性,拆文本零成本解决;画像→task 是单向依赖,拆调用收益有限却引入新成本。三方(读真实 prompt 的分析 + Codex + Gemini)一致不建议把"拆成两次 LLM 调用"作为目标态。
未来演进:画像持久化 + 增量(不经过"两次完整调用")
如果将来积累到"task 决策质量被画像 prompt 拖累"的实测证据,下一步不是升级成两次完整调用,而是:把画像结果持久化(其实 contacts.lifecycleStage 等字段已经在存),task 决策 prompt 只吃「已存画像 + 最近新增的 call/SMS delta + 当前 open tasks」,不再重读全部历史 —— input token 砍一大半。这条路同时也是触发层优化的根治方案(见「触发机制与优化路线」)。
关键设计决策
- Call Analysis 不升级。继续只产 per-call facts/signals,不读 SMS/leads/tasks,不做 task decision。Task 判断需要跨通话的历史上下文(open tasks、多次 no_answer 的累积),per-call prompt 没有这些信息
- 拆文本不拆调用。画像段和 task 决策段拆成两个模板分别维护,但拼成一个 system prompt 一次调用。不引入第二次 LLM 调用(避免 latency 翻倍 + 一致性约束断裂),详见上方「拆 prompt 的方式」
- 先画像,后任务。task 决策段需要知道客户当前状态(lifecycle、leadStatus、DNC)才能决定建什么 task。这个依赖在一次调用内由 prompt 内部顺序保证(模板1 在模板2 之前)
- MutationPlan 是信封,不是执行方式。contactActions[] 可以为空。执行顺序由代码保证:contact → task → timeline
followUpNeeded不是 trigger,是参考信号。contacts-analyzer 的触发条件是contactPhone && storeId(每通电话都触发);followUpNeeded是上游 per-call AI 写进calls表的字段,Contact Analyzer 读它当 prompt context(渲染成fu=yes(...)一行),不参与触发判断也不进任何代码分支。三个信号字段不是一条链:followUpNeeded(calls 表,per-call 输出,被参考)/taskDecisions[](不落库的中间变量)/actionNeeded(落 contacts 表的真状态)各处不同层 —— 详见「信号字段定位」
4 个共享 module
Policy Guard 的 needs_review:低置信度的 AI 提案不自动落库,进入人工确认队列。这是 AI autonomy 控制的核心机制。
Timeline Writer 接收两种写入:source event("一通电话打完了" — 不过 Policy Guard,只记录事实)和 mutation audit("AI 建了一个 task" — Task Orchestrator / Contact Writer 内部自动写的,已过 Guard)。
7 张表
task_progress_events 和 contact_timeline 并存:progress event 会同步投影到 timeline,但 SoT 在 task_progress_events(按 task 维度查询用独立表,不用从 timeline JSONB 里 parse)。
信号字段定位(followUpNeeded / actionNeeded / taskDecisions)
这三个字段常被当成"一条 AI 信号链",但读真实代码后它们其实在不同层,归属不同 prompt,落不落库也不同。拆 prompt 前先把它们定位清楚:
关键纠正:followUpNeeded 不是 Contact Analyzer 的触发器(触发是 contactPhone && storeId),它是被 Contact Analyzer 读去当 prompt 参考的信号,读完不二次落库,且不该被直接当结论沿用 —— task 决策要基于跨通话上下文重新判断(per-call 的 followUpNeeded 是单通粒度,task 是跨通话粒度)。
两个 actionNeeded 都应退役(由 open task 派生)
tasks.actionNeeded 和 contacts.actionNeeded 是同一个冗余的两半,根因都是"旧 task 系统不够强、没法可靠地用 open task 表示'要不要跟进'"。task 4 层 model 把 task 做成可靠的 work object(status='open' = 待办)后,两个都能由 open task 派生:
tasks.actionNeeded(Phase 1 直接删):- 派生公式:
actionNeeded === (status === 'open') - 自然语言:task 只要 open 就代表"需要处理",actionNeeded 是 status 的重复表达。从历史代码看,所有 caller(lead-processor / contacts-analyzer / studio-api)create 时都硬编码写
true,从来没人写false—— 这个 boolean 没有信息量 - 实施:Phase 1 直接 DROP COLUMN
- 派生公式:
contacts.actionNeeded(Phase 1 改读路径,Phase 2 DROP COLUMN):- 派生公式:
actionNeeded === EXISTS(SELECT 1 FROM tasks WHERE contact_phone=? AND store_id=? AND status='open') - 自然语言:「该 contact 要不要跟进」等价于「该 contact 名下有没有 open task」,task 是可靠的 SoT
- 实施分 2 步(因为 studio-api 真在读
c.action_needed:contacts list query / leads KPI summary):- Phase 1:改
routes/v3/contacts.ts+routes/v3/leads.ts把c.action_needed替换成EXISTS(...)子查询;contact_analyzer 暂时仍写该字段(保险,observe-only);不动 schema - Phase 2:observe 一段时间确认 EXISTS 口径没漂移后,DROP COLUMN + 停 contact_analyzer 写入 + 删
idx_contacts_action_neededindex
- Phase 1:改
- Leads KPI 口径同步改成"有 open task 的 contact 数"
- 派生公式:
这呼应设计准则 Code as Guardrail, AI as Judgment:「要不要跟进」是 open task 的确定性派生,不该有第二个独立真值源。字段级冗余说明见 schema 文档(tasks-schema / contacts-schema)。
AI 永远不负责的事
AI 只输出结构化提案(Zod 验证),以下由代码负责:
- final task state(状态转换)
- idempotency(幂等)
- permission / DNC hard stop
- store isolation
- transaction boundary
- audit trail
Code vs AI 职责
现状 vs 目标
现状栏区分"已收口"和"仍缺"——多数共享 helper 已就位,真正的缺口集中在 task mutation 状态机。
One-time AI Call vs Tool Calling
两种模式共享同一套 Task Orchestrator / Contact Writer / Timeline Writer / Policy Guard。
Phase 1 建好 shared modules 后,tool calling 的 backend 调用路径已经有了(tool handler 调 shared modules)。但真正开放 tool calling 还需要额外前置条件:RBAC(谁能调哪些 tool)、budget cap(防止 AI 无限调用)、approval threshold(低置信度进 needs_review)、audit actor(记录 actorType='ai_agent')、idempotency(防重复调用)。这些在 Phase 1 → Target State 的过渡期间完成。
加新 workflow 应该是什么体验
Target state 下加一个新 workflow 只需要 5 步:
- 定义事件 / 触发条件
- 决定判断层是代码、AI 还是混合
- 复用或新增 action schema
- 调共享 module
- 加 test + event catalog payload
不需要:设计新的 DB writer / 发明新的 timeline payload / 重新决定 DNC 行为 / 重新决定 task 去重 / 重新决定 store 隔离 / 教 prompt 怎么操作数据库。
2. Phase 1 — 实施步骤
Phase 1 不改 pipeline runtime(SQS/Lambda 拓扑不变,不需要新的 AWS 资源),只把共享表写入收口到共享函数。所有 module 是 in-process TypeScript 函数。
Step 1: 历史方案中的 Task Pipeline Deliverable(已 superseded)
Task Pipeline Deliverable (Codex) 当时被作为 Phase 1 基础;以下内容是历史方案,不得绕过 V3 Task contract 直接落地。它定义了:
- Task 的 4 层 Object Model(Activity / Progress / Lifecycle / Outcome)
record_progress语义:没打通不关 task,追加进展记录(解决"打 5 次没接 = 5 个 task"的问题)create_closed语义:当场完成的业务目标(如 inbound 当场成交)直接建一个 closed task,让 dashboard 能看到closeResult枚举清理:progress 值(no_answer / left_voicemail)移出 closeResult,只留 business outcome- API contract(v3 endpoints)
- Prompt contract(输入 allowed set、输出 proposed mutations)
- UI semantics(Progress Update vs Close Task 两种操作)
- Code vs Prompt 职责矩阵
Step 2: 新建 Task Orchestrator + task_progress_events 表 + status 改名
applyTaskAction() 并发/幂等 contract(每 action SQL 层显式)
task_progress_events.idempotency_key 生成规则(caller 端)
这是 2026-06 historical proposal 的 standalone table 去重设计,不是当前实施方向。V3 proposal 中的 Timeline 使用分层的 commandId、eventIdempotencyKey 与 sourceInteractionKey,见 Task V3 Target Design Proposal §9.4。
Step 3: 新建 Timeline Writer + Event Catalog
Task 事件必须写 timeline,所以 Timeline Writer 要跟 Task Orchestrator 同步就位。
Step 4: 新建 Contact Writer + Policy Guard + tasks.store_id NOT NULL
Contact Writer Phase 1 不统一 AI 画像字段(customerSummary 等 18 个),因为只有 contacts-analyzer 写,不存在多 writer 竞争。
Policy Guard 返回值
Step 5: 4 个 caller 逐个迁移
现有 API 端点(close / reopen / postpone)外部行为不变,内部改成调 Task Orchestrator。
Prompt schema 改动:contacts-analyzer TaskDecision 6 个 action 命名为 create_open / create_closed / close / update / record_progress / reopen(与 Task Orchestrator applyTaskAction() 对齐)。由独立 engineer 并行交付,Phase 1 拿来直接接入。
Phase 1 不做
不切换 tool calling runtime、不引入 workflow engine、不合并 Lambda、不做 meaningful SMS 分类、不做 production Task V2 数据迁移。
3. Historical Phase 1 → 当时的 Target State
Phase 1 完成后基础架构就位。Target State 加更多 caller 和扩展。
Tool Calling 前置条件
SMS meaningful(Phase 2)
当前 message-processor 只做:存消息 + exact STOP DNC cascade。Phase 2 加代码分类器(trivial / meaningful / natural-language DNC),meaningful SMS 即时触发 Contact Analysis Reprocess。
触发机制与优化路线
Contact Analyzer 有 4 个触发源(per_call_analysis / cron / on_demand / reprocess),共享同一条 FIFO 队列 + 同一个 worker。现状能跑(DeepSeek 成本低 + daily cron 06:00 兜底 + 即时性对前台有价值),但触发机制有重复分析,量大后需优化。详细 trace + 修复方案见 GitHub issue #365。
根因:cooldown gate(contacts-analyzer/handler.ts)只对 source === 'reprocess' 生效,per_call / cron 不走它;而它检查的 lastContactAnalysisAt 字段每次分析都在写(基础设施现成,只是没接上)。
优化路线(非紧急,量大再做):
- 触发层(冗余 1/2):per_call 接入 cooldown / 短 debounce;cron 查询跳过"自上次分析后无新活动"的 contact。低风险,
lastContactAnalysisAt已就绪。 - 计算层(冗余 5):即「拆 prompt 的方式」里说的画像持久化 + 增量分析 —— task 决策只吃「已存画像 + 新增 delta」,不重读全史。优先级低于触发层(先削触发次数 ROI 更高)。
这两层都符合"先做一个能跑的版本(拆 prompt 文本 + Phase 1 写入层),scalability 优化随量增长再做"的策略。
4. Appendix
Processing Capability Module contracts
Target state 下 4 个 Processing Capability Module 是可复用的调用接口(Layer 2),它们产出结构化提案交给共享写入层(Layer 1)执行。
分歧决策
详细参考
Open Questions
Implementation Handoff Checklist
这部分不改变 design 结论,只是在开始 implementation 前把需要落到 PR / migration / test plan 的 delta 列出来,避免设计文档变成工程细节正文。