Plan 02 — callytics-common Modules Implementation Plan
For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (- [ ]) syntax for tracking.
Goal: Implement the 4 shared mutation modules(Task Orchestrator + Policy Guard + Contact Writer + Timeline Writer)+ their raw-SQL adapter for studio-api,as @retaintive/common/domain/*. After this plan ships, callers in Plans 03-06 swap their inline SQL to one-line module calls.
Architecture: 5 个 src/domain/*.ts 文件,每个一种职责 — task-orchestrator.ts(state machine + 6 action SQL builders),policy-guard.ts(9 个 check + computeAllowedTypeCategories),contact-writer.ts(identity/DNC/activity 3 SQL builders),timeline-writer.ts(writeTimelineEvent + Zod event payload validation),task-orchestrator-sql.ts(raw SQL helper 给 studio-api)。沿用现有代码风格:client.batch() / sql.transaction() 做事务边界,raw path 按 buildContactTimelineInsertSQL() 的方式显式返回 { sql, params },不引入新的 builder framework。
Tech Stack: TypeScript 5 / Drizzle ORM / Zod 3 / bun test / vitest integration / PostgreSQL 16(Neon)
Spec source: docs/product-design/v2/unified-pipeline/implementation-plan/normative-spec.md §3 全节 + §5.0 invariants
Dependency: Plan 01 merged + @retaintive/common 1.1.0 published(needs new taskProgressEvents table + tasks.executor_type / attempt_count / source_call_id columns)
File Structure
测试文件全 mirror 在 tests/domain/:
Task 1: types.ts — Shared interfaces (foundation)
Files:
- Create:
src/domain/types.ts
- Test:
tests/domain/types.test.ts(只 verify 类型 narrowing,可选)
Why this task first: Task 2-5 都 import 这里的 type。先 ship 让其他 task 可独立 review。
Step 1: Map what types this file owns
Step 2: Tests scenarios
原则:types.ts 是纯 type definition,没有 runtime logic。只 test 真 invariant —— type narrowing 在使用端工作正常(防写法 drift)。
不写其他 test —— TS 编译时强制。Test 形态用 // @ts-expect-error 验证 negative narrowing 不编译。
Step 3: Implementation
// src/domain/types.ts
import type { SQL } from 'drizzle-orm';
import type {
TaskTypeCategory,
TaskPriority,
TaskCloseType,
TaskCloseResult,
SuggestedAction,
} from '../db/schema/tasks';
import type {
TaskProgressType,
TaskChannel,
TaskActorType,
} from '../db/schema/task-progress-events';
import type { DrizzleClient } from '../db/client';
// contacts.lead_status is currently text in live schema, not an exported enum.
// Keep this as a domain alias until contacts schema owns a typed LEAD_STATUS const.
export type LeadStatus = string;
// ─── Action payloads ────────────────────────────────────────────
export interface CreateOpenPayload {
typeCategory: TaskTypeCategory;
suggestedActions: SuggestedAction[];
/** Optional override of dueAt; default = computeDueAt(priority) */
dueAt?: Date;
/** Phase 1: lead-processor passes 'lead', AI passes 'contact_analysis' */
sourceType: 'lead' | 'contact_analysis' | 'manual';
sourceLeadId?: string;
contactAnalysisRunId?: string;
/** AI 提案的 confidence (0-1);staff 操作不传 */
confidence?: number;
}
export interface CreateClosedPayload {
typeCategory: TaskTypeCategory;
suggestedActions: SuggestedAction[];
sourceType: 'lead' | 'contact_analysis' | 'manual';
/** REQUIRED: 触发 "当场办成" 的 callId, 作为 partial unique 防重 */
sourceCallId: string;
closeResult: TaskCloseResult;
closeNote?: string; // required when closeResult='other'
confidence?: number;
}
export interface ClosePayload {
taskId: string;
closeResult: TaskCloseResult;
closeNote?: string; // required when closeResult='other'
confidence?: number;
}
export interface UpdatePayload {
taskId: string;
/** Optional fields to patch (only `priority` / `dueAt` / `suggestedActions` for Phase 1) */
priority?: TaskPriority;
dueAt?: Date;
suggestedActions?: SuggestedAction[];
confidence?: number;
}
export interface RecordProgressPayload {
taskId: string;
progressType: TaskProgressType;
channel: TaskChannel;
callId?: string;
messageId?: string;
note?: string;
/** Optional: staff manually sets nextDueAt;否则 Orchestrator 用 computeNextDueAt() */
nextDueAtOverride?: Date;
}
export interface ReopenPayload {
taskId: string;
/** Optional reason recorded in timeline */
reason?: string;
}
export type TaskAction =
| { action: 'create_open'; payload: CreateOpenPayload }
| { action: 'create_closed'; payload: CreateClosedPayload }
| { action: 'close'; payload: ClosePayload }
| { action: 'update'; payload: UpdatePayload }
| { action: 'record_progress'; payload: RecordProgressPayload }
| { action: 'reopen'; payload: ReopenPayload };
// ─── Result types ───────────────────────────────────────────────
export type RejectReason =
| 'dnc'
| 'store_mismatch'
| 'low_confidence'
| 'stale_proposal'
| 'duplicate'
| 'task_not_open'
| 'task_not_closed'
| 'create_closed_with_open_task'
| 'invalid_state_transition'
| 'hallucinated_task_id'
| 'invalid_type_category'
| 'invalid_progress_type'
| 'invalid_close_result'
| 'close_note_required';
export type ApplyResult =
| {
status: 'allow';
statements: SQL[];
/**
* Statements whose RETURNING row count is semantically meaningful.
* Caller/executor MUST map 0 rows to the given reject reason.
* Examples:
* - create_open CTE inserted 0 rows → duplicate
* - close/update CTE updated 0 rows → stale_proposal / task_not_open
* - reopen CTE updated 0 rows → task_not_closed
*/
resultChecks?: Array<{ statementIndex: number; zeroRowsReason: RejectReason }>;
/** Stable ids generated before SQL rendering, so timeline payloads match task rows. */
generatedIds?: { taskId?: string; idempotencyKey?: string };
}
| { status: 'reject'; reason: RejectReason; details: string }
| { status: 'needs_review'; details: string }; // Phase 1 不返回,Phase 2 预留
export type ApplyResultSQL =
| {
status: 'allow';
statements: { sql: string; params: unknown[] }[];
resultChecks?: Array<{ statementIndex: number; zeroRowsReason: RejectReason }>;
generatedIds?: { taskId?: string; idempotencyKey?: string };
}
| { status: 'reject'; reason: RejectReason; details: string }
| { status: 'needs_review'; details: string };
// ─── Context types ──────────────────────────────────────────────
export interface BaseActionContext {
storeId: string;
actor: {
type: 'staff' | 'system' | 'ai_agent' | 'contact_analysis';
id?: string;
name?: string;
};
/** AI 提案才传,human authority guard 用 */
aiRunStartedAt?: Date;
db: DrizzleClient;
}
/** create_open / create_closed —— caller 必传 contactPhone(新建 task)*/
export interface CreateActionContext extends BaseActionContext {
contactPhone: string;
/** Identity for upserting contact in same transaction */
franchiseId: string;
accountId: string;
}
/** close / update / record_progress / reopen —— caller 只传 taskId */
export type TaskIdActionContext = BaseActionContext;
export type ActionContext = CreateActionContext | TaskIdActionContext;
// ─── Domain types ───────────────────────────────────────────────
export interface ContactLifecycle {
lifecycleStage: 'lead' | 'member' | 'churned' | 'unknown';
lifecycleState: 'active' | 'paused' | 'terminal';
leadStatus?: LeadStatus;
doNotContact: boolean;
}
Step 4: Test write + run
Implementer 按 case map 1-3 写 type narrowing tests(用 // @ts-expect-error pattern)。Run:
cd /Users/maxwsy/workspace/callytics-common && bun test tests/domain/types.test.ts
Expected: 3 tests PASS。
Step 5: Commit
git add src/domain/types.ts tests/domain/types.test.ts
git commit -m "feat(domain): types — TaskAction union + ActionContext split + ApplyResult
Phase 1 Plan 02 Task 1.
Spec: normative-spec.md §3.1
"
Task 2: policy-guard.ts — 9 checks + computeAllowedTypeCategories + computeNextDueAt
Files:
- Create:
src/domain/policy-guard.ts
- Create:
src/domain/idempotency.ts
- Test:
tests/domain/policy-guard.test.ts
- Test:
tests/domain/idempotency.test.ts
Step 1: Map test scenarios
computeAllowedTypeCategories(contact) — 7 case 全枚举(spec §3.6 line 489-497):
computeNextDueAt(progressType, currentDueAt, override) — 6 progressType × override 有/无(spec §3.3):
Policy Guard 8 checks(snapshot-based,无 DB call)—— 每个 check pure function over ContactSnapshot:
TOCTOU notes(为什么 snapshot 而不是 per-check SELECT)—— pattern source contacts-analyzer:413, 543-557:
- production 用
pendingTasks snapshot + in-memory knownPendingIds / effectivePendingCategories Set 做 check
- 真 enforcement 在 SQL 层:partial unique index /
WHERE status='pending' clause / sticky CASE WHEN
- module 层不 SELECT 比较 = 不引入 TOCTOU race + 不引入 N+1 round-trip
Snapshot loader test(net new):
buildProgressIdempotencyKey() — 4 source pattern(spec §3.2 table line 296-301):
Edge case:
- 同一组 input 调两次 → 输出完全相同(deterministic,保证幂等)
- callId + messageId 同时给 → throw(spec 未明说但语义冲突应 fail-fast)
- 都不给 → throw(必须有至少一种 evidence)
Step 2: Implementation — computeAllowedTypeCategories
// src/domain/policy-guard.ts
import type { TaskTypeCategory } from '../db/schema/tasks';
import type { ContactLifecycle, RejectReason, TaskAction, ActionContext } from './types';
const MEMBER_CATEGORIES: TaskTypeCategory[] = [
'cancellation_risk', 'retention', 'upgrade', 'renewal', 'referral',
];
const LEAD_AI_CATEGORIES: TaskTypeCategory[] = [
'lead_follow_up', 'booked_not_converted',
];
const CHURNED_CATEGORIES: TaskTypeCategory[] = ['win_back'];
export function computeAllowedTypeCategories(
contact: ContactLifecycle,
): TaskTypeCategory[] {
if (contact.doNotContact) return [];
if (contact.lifecycleState === 'terminal') return [];
if (contact.lifecycleStage === 'lead') return LEAD_AI_CATEGORIES;
if (contact.lifecycleStage === 'member') return MEMBER_CATEGORIES;
if (contact.lifecycleStage === 'churned') return CHURNED_CATEGORIES;
return []; // unknown
}
Step 3: Implementation — computeNextDueAt
// src/domain/policy-guard.ts (continued)
import type { TaskProgressType } from '../db/schema/task-progress-events';
const PROGRESS_NEXT_DUE_INTERVAL_MINUTES: Record<TaskProgressType, number | null> = {
no_answer: 60,
left_voicemail: 60 * 24,
text_sent: 60 * 24 * 2,
callback_requested: null,
follow_up_scheduled: null,
customer_considering: 60 * 24 * 3,
};
export function computeNextDueAt(
progressType: TaskProgressType,
currentDueAt: Date | null,
callerOverride?: Date,
): Date | null {
if (callerOverride) return callerOverride;
const interval = PROGRESS_NEXT_DUE_INTERVAL_MINUTES[progressType];
if (interval === null) return currentDueAt;
return new Date(Date.now() + interval * 60_000);
}
Step 4: Implementation — 8 checks(snapshot-based pure functions)
Pattern source(do NOT reinvent) —— contacts-analyzer/src/infrastructure/neon-repository.ts:413, 543-557:
// 真实 codebase 现有 pattern:caller 顶部读 1 次,后续 check 用内存 Set
const pendingTasks = await getTasksForContact(client, phone, storeId); // 1 read
const knownPendingIds = new Set(pendingTasks.map(t => t.taskId));
const closingCategories = new Set(taskDecisions
.filter(d => d.action === 'close' && knownPendingIds.has(d.taskId))
.map(d => d.typeCategory));
const effectivePendingCategories = new Set(pendingTasks
.filter(t => !closingCategories.has(t.typeCategory))
.map(t => t.typeCategory));
// 后续 8 个 check 全在内存 Set 里查 —— 不再回 DB
Phase 1 Policy Guard 直接沿用这个 pattern。Module 层不做 SELECT contacts WHERE ... / SELECT tasks WHERE ... per-check —— 那会引入:
- 8N round-trip(N = 单次 caller 处理的 decision 数;contacts-analyzer 一次能跑 20 decision)
- TOCTOU race:per-check SELECT 到最终
client.batch() 之间状态可能变。production 用 snapshot Set + SQL 层 enforcement(partial unique / WHERE status='pending' / sticky CASE)天然避免
Snapshot loader(Orchestrator 调一次,跟现有 caller pattern 一致)
// src/domain/policy-guard.ts
import { tasks, contacts } from '../db/schema';
import { eq, and, sql } from 'drizzle-orm';
import type { DrizzleClient } from '../db/types';
import type { TaskAction, ActionContext } from './types';
const MIN_AI_CONFIDENCE = Number(process.env['MIN_AI_CONFIDENCE'] ?? 0.7);
export type CheckResult =
| { ok: true }
| { ok: false; reason: RejectReason; details: string };
/** Caller 顶部读一次的 snapshot —— 跟 contacts-analyzer:413 同形态 */
export interface ContactSnapshot {
contact: {
doNotContact: boolean;
lifecycleStage: ContactLifecycle['lifecycleStage'];
lifecycleState: ContactLifecycle['lifecycleState'];
leadStatus?: ContactLifecycle['leadStatus'];
} | null; // null 表示新 contact(尚未 INSERT)
pendingTasks: Array<{ taskId: string; typeCategory: TaskTypeCategory }>;
}
/** Orchestrator 顶部调一次,后续 check 全用 snapshot */
export async function loadSnapshot(
client: DrizzleClient,
contactPhone: string,
storeId: string,
): Promise<ContactSnapshot> {
const [contactRow] = await client
.select({
doNotContact: contacts.doNotContact,
lifecycleStage: contacts.lifecycleStage,
lifecycleState: contacts.lifecycleState,
leadStatus: contacts.leadStatus,
})
.from(contacts)
.where(and(eq(contacts.phone, contactPhone), eq(contacts.storeId, storeId)))
.limit(1);
const pendingTasks = await client
.select({ taskId: tasks.taskId, typeCategory: tasks.typeCategory })
.from(tasks)
.where(and(
eq(tasks.contactPhone, contactPhone),
eq(tasks.storeId, storeId),
sql`${tasks.status} IN ('pending', 'open')`, // transitional 3-value
));
return { contact: contactRow ?? null, pendingTasks };
}
8 checks(全 pure function,无 DB call)
// 1. storeId 非空 — pure
export function checkStoreId(ctx: ActionContext): CheckResult {
if (!ctx.storeId || ctx.storeId.length === 0) {
return { ok: false, reason: 'store_mismatch', details: 'storeId is required' };
}
return { ok: true };
}
// 2. DNC hard stop — snapshot lookup,无 DB call
const DNC_ALLOWED_ACTIONS = new Set(['close']);
export function checkDnc(
snapshot: ContactSnapshot,
actionName: TaskAction['action'],
): CheckResult {
// 新 contact(snapshot.contact === null)默认 doNotContact=false
if (!snapshot.contact || !snapshot.contact.doNotContact) return { ok: true };
if (DNC_ALLOWED_ACTIONS.has(actionName)) return { ok: true };
return { ok: false, reason: 'dnc', details: `contact is DNC; only close is allowed` };
}
// 3. AI hallucination guard — Set lookup,无 DB call
export function checkTaskBelongsToContact(
snapshot: ContactSnapshot,
taskId: string,
): CheckResult {
const knownIds = new Set(snapshot.pendingTasks.map(t => t.taskId));
if (!knownIds.has(taskId)) {
return {
ok: false,
reason: 'hallucinated_task_id',
details: `taskId ${taskId} not in this contact's open tasks`,
};
}
return { ok: true };
}
// 4. State transition — pure(不需要 snapshot,因为 caller 已 narrow 到 specific decision)
export function checkStateTransition(
action: TaskAction['action'],
taskCurrentStatus: 'pending' | 'open' | 'closed', // transitional 3-value
): CheckResult {
const isOpen = taskCurrentStatus === 'pending' || taskCurrentStatus === 'open';
if (action === 'close' && !isOpen) {
return { ok: false, reason: 'task_not_open', details: 'close requires status=open' };
}
if (action === 'reopen' && taskCurrentStatus !== 'closed') {
return { ok: false, reason: 'task_not_closed', details: 'reopen requires status=closed' };
}
if ((action === 'update' || action === 'record_progress') && !isOpen) {
return { ok: false, reason: 'task_not_open', details: `${action} requires status=open` };
}
return { ok: true };
}
// 5. create_closed dedup — snapshot Set lookup,无 DB call
// NOTE: SQL 层 partial unique index `uq_tasks_create_closed_evidence` 是
// 最终安全网(Plan 01 Task 2 Step 8)—— 即使这条 check 漏了,INSERT 仍会
// 因 unique 冲突 fail,caller `RETURNING *` 0 rows → reject 'duplicate'。
// 这里只是 fail-fast in-memory check,跟 lead-tracking:201-204 同形态。
export function checkCreateClosedHasNoOpenTask(
snapshot: ContactSnapshot,
typeCategory: TaskTypeCategory,
): CheckResult {
// 注意:正在被这次 batch close 的 typeCategory 不算"已有 open"
// —— caller 应在 closingCategories 计算后传 effectivePendingCategories
const hasOpen = snapshot.pendingTasks.some(t => t.typeCategory === typeCategory);
if (hasOpen) {
return {
ok: false,
reason: 'create_closed_with_open_task',
details: `open task exists for typeCategory=${typeCategory}; use close instead`,
};
}
return { ok: true };
}
// 6. typeCategory allowed set — pure
const SKIP_TYPE_CHECK_ACTORS = new Set(['staff', 'system']);
export function checkTypeCategoryAllowed(
ctx: ActionContext,
contact: ContactLifecycle,
typeCategory: TaskTypeCategory,
): CheckResult {
if (SKIP_TYPE_CHECK_ACTORS.has(ctx.actor.type)) return { ok: true };
const allowed = computeAllowedTypeCategories(contact);
if (!allowed.includes(typeCategory)) {
return {
ok: false,
reason: 'invalid_type_category',
details: `typeCategory=${typeCategory} not in allowed set [${allowed.join(',')}]`,
};
}
return { ok: true };
}
// 7. closeNote required when closeResult='other' — pure
export function checkCloseNoteRequired(payload: {
closeResult?: string;
closeNote?: string;
}): CheckResult {
if (payload.closeResult === 'other' && !payload.closeNote?.trim()) {
return {
ok: false, reason: 'close_note_required',
details: 'closeNote is required when closeResult=other',
};
}
return { ok: true };
}
// 8. Confidence threshold — pure
export function checkConfidence(confidence?: number): CheckResult {
if (confidence === undefined) return { ok: true };
if (confidence < MIN_AI_CONFIDENCE) {
return {
ok: false, reason: 'low_confidence',
details: `confidence=${confidence} < MIN_AI_CONFIDENCE=${MIN_AI_CONFIDENCE}`,
};
}
return { ok: true };
}
// 9. Human authority guard 是 SQL 层做的(WHERE updated_at <= aiRunStartedAt RETURNING)
// 不在 module 层 SELECT 比较 —— 跟 contacts-analyzer pattern 一致,SQL 层
// enforcement 是真安全网。详见 task-orchestrator.ts SQL builder。
Round-trip 对比:
Step 5: Implementation — idempotency.ts
// src/domain/idempotency.ts
import type { TaskProgressType } from '../db/schema/task-progress-events';
export interface PhoneProgressEvidence {
source: 'phone';
callId: string;
}
export interface SmsProgressEvidence {
source: 'sms';
messageId: string;
}
export interface ManualProgressEvidence {
source: 'manual';
actorId: string;
occurredAt: Date;
}
export interface SystemProgressEvidence {
source: 'system';
runId: string;
}
export type ProgressEvidence =
| PhoneProgressEvidence
| SmsProgressEvidence
| ManualProgressEvidence
| SystemProgressEvidence;
export function buildProgressIdempotencyKey(
taskId: string,
evidence: ProgressEvidence,
progressType: TaskProgressType,
): string {
switch (evidence.source) {
case 'phone':
return `progress:${taskId}:call:${evidence.callId}:${progressType}`;
case 'sms':
return `progress:${taskId}:message:${evidence.messageId}:${progressType}`;
case 'manual':
return `progress:${taskId}:manual:${evidence.actorId}:${evidence.occurredAt.toISOString()}:${progressType}`;
case 'system':
return `progress:${taskId}:system:${evidence.runId}:${progressType}`;
}
}
Step 6: Test write + run
Implementer 按 case map 写:
- 8
computeAllowedTypeCategories happy/edge case
- 7
computeNextDueAt scenarios
- 8 checks × happy+reject scenarios = ~20 test
- 4
buildProgressIdempotencyKey source + 2 edge case
Run:
cd /Users/maxwsy/workspace/callytics-common && bun test tests/domain/policy-guard.test.ts tests/domain/idempotency.test.ts
Expected: all PASS。
Step 7: Commit
git add src/domain/policy-guard.ts src/domain/idempotency.ts tests/domain/policy-guard.test.ts tests/domain/idempotency.test.ts
git commit -m "feat(domain): policy-guard snapshot-based + computeAllowedTypeCategories + computeNextDueAt + idempotency
Phase 1 Plan 02 Task 2.
Pattern source: 沿用 contacts-analyzer:413, 543-557 现有 snapshot+Set lookup pattern。
Module 层 8 checks 全 pure function over ContactSnapshot — 不 SELECT,不引入
N+1 round-trip,不引入 TOCTOU race。真 enforcement 在 SQL 层(partial unique
index / WHERE clause / sticky CASE WHEN)。
- loadSnapshot(): caller 顶部调 1 次,读 contact row + pendingTasks
- 8 Policy Guard checks (storeId / DNC / hallucination / state /
create_closed-no-open / typeCategory / closeNote / confidence)
- computeAllowedTypeCategories: 7 case 表 cover lifecycle x DNC x lead/member/churned/unknown
- computeNextDueAt: per-progressType fixed interval + caller override
- buildProgressIdempotencyKey: 4 source pattern (phone/sms/manual/system)
- Human authority guard 由 SQL conditional WHERE 做,不在 module 层 SELECT
Spec: normative-spec.md §3.6 + §3.3 + §3.2
"
Files:
- Create:
src/domain/task-orchestrator.ts
- Test:
tests/domain/task-orchestrator.test.ts(unit)
- Test:
tests/integration/task-orchestrator.test.ts(连 Neon,验 SQL contract 真生效)
Atomicity rule:沿用现有 client.batch() pattern 做事务边界。只有当 audit/progress write 必须依赖前一个 mutation 是否实际 INSERT/UPDATE RETURNING 成功时,才把该 mutation + audit 写成单条 CTE SQL unit。不要引入新的 transaction framework;CTE 只是为了表达 conditional dependent write。
Step 1: Map test scenarios
applyTaskAction() 6 action × happy + reject —— spec §3.1 + Plan 06 normative §5 invariants:
Critical SQL invariants(integration test 验):
closeAllOpenForContact() happy + edge:
Step 2: Implementation — applyTaskAction dispatch
// src/domain/task-orchestrator.ts
import { eq, and, sql, type SQL } from 'drizzle-orm';
import { tasks, contacts } from '../db/schema';
import { taskProgressEvents } from '../db/schema/task-progress-events';
import { contactTimeline } from '../db/schema/contact-timeline';
import { computeDueAt } from '../utils/due-at';
import { NAME_TRUST } from '../db/schema/name-trust';
import {
checkStoreId,
checkDnc,
checkTaskBelongsToContact,
checkStateTransition,
checkCreateClosedHasNoOpenTask,
checkTypeCategoryAllowed,
checkCloseNoteRequired,
checkConfidence,
computeNextDueAt,
} from './policy-guard';
import { buildProgressIdempotencyKey } from './idempotency';
import type {
TaskAction,
ApplyResult,
ActionContext,
CreateActionContext,
ContactLifecycle,
} from './types';
export async function applyTaskAction(
action: TaskAction,
ctx: ActionContext,
): Promise<ApplyResult> {
// Global checks (all actions)
const storeCheck = checkStoreId(ctx);
if (!storeCheck.ok) return reject(storeCheck.reason, storeCheck.details);
switch (action.action) {
case 'create_open':
return handleCreateOpen(action.payload, ctx as CreateActionContext);
case 'create_closed':
return handleCreateClosed(action.payload, ctx as CreateActionContext);
case 'close':
return handleClose(action.payload, ctx);
case 'update':
return handleUpdate(action.payload, ctx);
case 'record_progress':
return handleRecordProgress(action.payload, ctx);
case 'reopen':
return handleReopen(action.payload, ctx);
}
}
function reject(reason: ApplyResult & { status: 'reject' }['reason'], details: string): ApplyResult {
return { status: 'reject', reason, details };
}
Step 3: Implementation — applyTaskAction orchestrator(snapshot pattern)
关键修正(2026-06-02 audit fix):整个 orchestrator dispatch path 不再 per-decision SELECT。
顶部用 Task 2 的 loadSnapshot() 一次性加载 contact + pendingTasks,后续 8 checks 全
pure function over snapshot,跟 contacts-analyzer:413, 543-557 verbatim 同形态(snapshot Set lookup,无 DB round-trip)。
// src/domain/task-orchestrator.ts
import type { CreateOpenPayload, CreateClosedPayload } from './types';
import {
loadSnapshot,
checkStoreId, checkDnc, checkTaskBelongsToContact, checkStateTransition,
checkCreateClosedHasNoOpenTask, checkTypeCategoryAllowed, checkConfidence,
checkCloseNoteRequired,
type ContactSnapshot,
} from './policy-guard';
/**
* 顶层 batch entry — caller 传入 N 个 decisions(可能同 contact 多个 action),
* orchestrator 在 batch 开头读 1 次 snapshot(contact row + open tasks),
* 后续每个 action handler 复用同一 snapshot。
*
* 跟 contacts-analyzer:413, 543-557 verbatim 同形态:
* const knownPendingIds = new Set(pendingTasks.map(t => t.taskId));
* const closingCategories = new Set(decisions.filter(d => d.action === 'close')...);
*
* Round-trip 对比(20 decision batch):
* 旧设计: 1 + 20*4 ≈ 80 SELECT(contact + per-action 3 checks)
* 新设计: 2 SELECT(loadSnapshot 内的 contact + pendingTasks)
*/
export async function applyTaskActions(
decisions: TaskAction[],
ctx: ActionContext,
): Promise<ApplyResult[]> {
// 1. storeId 非空 gate(pure check,无 snapshot 依赖)
const storeCheck = checkStoreId(ctx);
if (!storeCheck.ok) return decisions.map(() => reject(storeCheck.reason, storeCheck.details));
// 2. 顶部读一次 snapshot — 不 per-decision SELECT
const snapshot = await loadSnapshot(ctx.db, ctx.contactPhone, ctx.storeId);
// 3. dispatch — 全部 handler 接 snapshot 而非 ctx + phone
return Promise.all(
decisions.map(d => dispatchAction(d, ctx, snapshot)),
);
}
function dispatchAction(
decision: TaskAction,
ctx: ActionContext,
snapshot: ContactSnapshot,
): ApplyResult {
switch (decision.action) {
case 'create_open': return handleCreateOpen(decision.payload, ctx, snapshot);
case 'create_closed': return handleCreateClosed(decision.payload, ctx, snapshot);
case 'close': return handleClose(decision.payload, ctx, snapshot);
case 'update': return handleUpdate(decision.payload, ctx, snapshot);
case 'record_progress': return handleRecordProgress(decision.payload, ctx, snapshot);
case 'reopen': return handleReopen(decision.payload, ctx, snapshot);
}
}
function handleCreateOpen(
payload: CreateOpenPayload,
ctx: CreateActionContext,
snapshot: ContactSnapshot,
): ApplyResult {
// confidence check(pure)
const confCheck = checkConfidence(payload.confidence);
if (!confCheck.ok) return reject(confCheck.reason, confCheck.details);
// DNC + typeCategory checks 全用 snapshot — 无 DB call
const dncCheck = checkDnc(snapshot, 'create_open');
if (!dncCheck.ok) return reject(dncCheck.reason, dncCheck.details);
// snapshot.contact === null 表示新 contact,checkTypeCategoryAllowed 内部 handle
const typeCheck = checkTypeCategoryAllowed(ctx, snapshot, payload.typeCategory);
if (!typeCheck.ok) return reject(typeCheck.reason, typeCheck.details);
// Derive priority from max(suggestedActions[].priority) — spec §3.1 S1
const priority = derivePriority(payload.suggestedActions);
const dueAt = payload.dueAt ?? computeDueAt(priority);
/*
* Generate taskId before SQL rendering. Timeline payload must reference the
* same taskId that is inserted into tasks. Do NOT rely on DB default UUID +
* caller-side chaining; ApplyResult is a static mutation plan.
*/
const taskId = crypto.randomUUID();
/*
* One SQL unit owns both task insert and timeline audit. The timeline SELECT
* reads from inserted, so duplicate/no-op insert writes zero timeline rows.
* RETURNING task_id lets the executor map 0 rows -> duplicate.
*
* Drizzle query builder cannot express this conditional timeline dependency
* cleanly; use sql`` CTE for create_open/create_closed.
*/
const createWithTimelineStmt = buildCreateOpenWithTimelineSQL({
taskId,
ctx,
payload,
priority,
dueAt,
});
return {
status: 'allow',
statements: [createWithTimelineStmt],
resultChecks: [{ statementIndex: 0, zeroRowsReason: 'duplicate' }],
generatedIds: { taskId },
};
}
function derivePriority(
actions: { priority: 'high' | 'medium' | 'low' }[],
): 'high' | 'medium' | 'low' {
const order = { high: 3, medium: 2, low: 1 } as const;
let max: 'high' | 'medium' | 'low' = 'low';
for (const a of actions) {
if (order[a.priority] > order[max]) max = a.priority;
}
return max;
}
Step 4: Implementation — handleClose, handleRecordProgress (其余 actions 类似 pattern)
// src/domain/task-orchestrator.ts (continued)
function handleClose(
payload: { taskId: string; closeResult: any; closeNote?: string; confidence?: number },
ctx: ActionContext,
snapshot: ContactSnapshot,
): ApplyResult {
// confidence + closeNote(pure checks)
const confCheck = checkConfidence(payload.confidence);
if (!confCheck.ok) return reject(confCheck.reason, confCheck.details);
const noteCheck = checkCloseNoteRequired(payload);
if (!noteCheck.ok) return reject(noteCheck.reason, noteCheck.details);
// taskId 归属 + 当前状态 — 全 snapshot lookup,无 DB call
const belongsCheck = checkTaskBelongsToContact(snapshot, payload.taskId);
if (!belongsCheck.ok) return reject(belongsCheck.reason, belongsCheck.details);
// close 在 DNC 下允许(spec §3.6 line 438)— 无需 checkDnc
// state transition check — snapshot.pendingTasks 已 narrow 到 open tasks,
// 能找到说明 status='open'/'pending'。Stale proposal 由 SQL CTE WHERE updated_at <= aiRunStartedAt 守。
// Build conditional UPDATE: WHERE status IN ('pending','open') AND updated_at <= $aiRunStartedAt
const whereClauses = [
eq(tasks.taskId, payload.taskId),
sql`${tasks.status} IN ('pending', 'open')`,
];
if (ctx.aiRunStartedAt) {
whereClauses.push(sql`${tasks.updatedAt} <= ${ctx.aiRunStartedAt}`);
}
/*
* One SQL unit owns UPDATE + timeline. The timeline INSERT reads from updated,
* so stale/closed/no-op updates do not create audit rows.
*/
const closeStmt = buildCloseWithTimelineSQL({
ctx,
payload,
whereClauses,
});
return {
status: 'allow',
statements: [closeStmt],
resultChecks: [{
statementIndex: 0,
zeroRowsReason: ctx.aiRunStartedAt ? 'stale_proposal' : 'task_not_open',
}],
};
}
function handleRecordProgress(
payload: {
taskId: string;
progressType: any;
channel: any;
callId?: string;
messageId?: string;
note?: string;
nextDueAtOverride?: Date;
},
ctx: ActionContext,
snapshot: ContactSnapshot,
): ApplyResult {
// taskId 归属 + DNC + state — 全 snapshot pure checks,无 DB call
const belongsCheck = checkTaskBelongsToContact(snapshot, payload.taskId);
if (!belongsCheck.ok) return reject(belongsCheck.reason, belongsCheck.details);
// DNC check — record_progress is NOT allowed under DNC (spec §3.6 line 438)
const dncCheck = checkDnc(snapshot, 'record_progress');
if (!dncCheck.ok) return reject(dncCheck.reason, dncCheck.details);
// Build idempotency_key
const evidence = payload.callId
? { source: 'phone' as const, callId: payload.callId }
: payload.messageId
? { source: 'sms' as const, messageId: payload.messageId }
: ctx.actor.id
? { source: 'manual' as const, actorId: ctx.actor.id, occurredAt: new Date() }
: { source: 'system' as const, runId: 'phase1' };
const idempotencyKey = buildProgressIdempotencyKey(payload.taskId, evidence, payload.progressType);
// dueAt 查从 snapshot.pendingTasks 拿不到 — 需 caller 在 payload 传或额外 loader。
// 简化:Orchestrator 接 payload.currentDueAt(caller 已知),或 nextDueAtOverride 必填。
const currentDueAt = (payload as any).currentDueAt ?? null;
const nextDueAt = computeNextDueAt(payload.progressType, currentDueAt, payload.nextDueAtOverride);
/*
* CTE: only increment attempt_count and write timeline if INSERT actually
* inserted a new progress event (B8). Duplicate idempotency_key returns 0 rows.
*/
const progressStmt = buildRecordProgressWithTimelineSQL({
ctx,
payload,
nextDueAt,
idempotencyKey,
});
return {
status: 'allow',
statements: [progressStmt],
resultChecks: [{ statementIndex: 0, zeroRowsReason: 'duplicate' }],
generatedIds: { idempotencyKey },
};
}
// handleCreateClosed / handleUpdate / handleReopen — same pattern,impl 略
// src/domain/task-orchestrator.ts (continued)
import type { TaskCloseResult } from '../db/schema/tasks';
export async function closeAllOpenForContact(
params: {
contactPhone: string;
storeId: string;
closeResult: TaskCloseResult;
closeNote: string;
actor: ActionContext['actor'];
},
ctx: { db: DrizzleClient },
): Promise<
| { status: 'allow'; statements: SQL[]; closedTaskIds: string[] }
| { status: 'reject'; reason: any; details: string }
> {
// Optional pre-read for logging / idempotency metadata only.
// Correctness must come from UPDATE ... RETURNING inside the SQL unit below.
const openTasks = await ctx.db
.select({ taskId: tasks.taskId })
.from(tasks)
.where(and(
eq(tasks.contactPhone, params.contactPhone),
eq(tasks.storeId, params.storeId),
sql`${tasks.status} IN ('pending', 'open')`,
));
if (openTasks.length === 0) {
return { status: 'allow', statements: [], closedTaskIds: [] };
}
const candidateTaskIds = openTasks.map(t => t.taskId);
/*
* Single SQL unit:
* WITH updated AS (UPDATE tasks ... RETURNING task_id, ...)
* INSERT INTO contact_timeline (...) SELECT ... FROM updated
*
* This avoids the SELECT-before-UPDATE race where a task is preselected but
* no longer open by the time UPDATE runs. Timeline rows are sourced only from
* rows actually updated.
*/
const bulkCloseStmt = buildCloseAllOpenWithTimelineSQL(params);
return {
status: 'allow',
statements: [bulkCloseStmt],
closedTaskIds: candidateTaskIds,
};
}
Step 6: Implementation — conditional mutation + timeline SQL helpers
// src/domain/task-orchestrator.ts (continued)
// (impl 内部 helpers,跟 timeline-writer.ts 协作 — 见 Task 5)
declare function buildCreateOpenWithTimelineSQL(params: {
taskId: string;
ctx: CreateActionContext;
payload: CreateOpenPayload;
priority: TaskPriority;
dueAt: Date;
}): SQL;
declare function buildCloseWithTimelineSQL(params: {
ctx: ActionContext;
payload: ClosePayload;
whereClauses: SQL[];
}): SQL;
declare function buildRecordProgressWithTimelineSQL(params: {
ctx: ActionContext;
payload: RecordProgressPayload;
task: { contactPhone: string; status: string; dueAt: Date | null };
nextDueAt: Date | null;
idempotencyKey: string;
}): SQL;
declare function buildCloseAllOpenWithTimelineSQL(params: {
contactPhone: string;
storeId: string;
closeResult: TaskCloseResult;
closeNote: string;
actor: ActionContext['actor'];
}): SQL;
实际 impl 在 Task 5(timeline-writer.ts)。Orchestrator 这里只 declare,Task 5 ship 后 import。
Step 7: Test write + run
Implementer 按 case map:
- 6 action × happy + reject = ~31 test(含 update under DNC reject)
- C1-C5 SQL invariants = 5 integration test(连 test Neon)
- closeAllOpenForContact 3 scenarios
cd /Users/maxwsy/workspace/callytics-common && bun test tests/domain/task-orchestrator.test.ts
# 然后跑 integration:
bun test tests/integration/task-orchestrator.test.ts
Expected: all PASS。
Step 8: Commit
git add src/domain/task-orchestrator.ts tests/domain/task-orchestrator.test.ts tests/integration/task-orchestrator.test.ts
git commit -m "feat(domain): task-orchestrator — applyTaskAction 6 actions + closeAllOpenForContact
Phase 1 Plan 02 Task 3.
- applyTaskAction: 6 action dispatch + per-action Policy Guard checks
- record_progress: CTE 防 attempt_count 翻倍(B8 fix)
- close/update/reopen: human authority guard via SQL conditional WHERE
- create_closed: source_call_id evidence + create_closed_with_open_task check
- tasks.priority 派生自 max(suggestedActions[].priority) (S1)
- closeAllOpenForContact: bulk close + N timeline INSERT (B7 helper)
SQL invariants verified via integration test against Neon test env.
Spec: normative-spec.md §3.1 + §5.0 invariants
"
Files:
- Create:
src/domain/contact-writer.ts
- Test:
tests/domain/contact-writer.test.ts(unit)
- Test:
tests/integration/contact-writer.test.ts(NAME_TRUST winner 验证连 Neon)
Step 1: Map test scenarios
upsertIdentity() —— NAME_TRUST winner(spec §3.4 + name-trust.md):
setDNC() —— sticky CASE WHEN audit trail(spec §3.4 line 377-381;pattern source message-processor:332-335):
touchActivity() —— forward-only:
Step 2: Implementation
// src/domain/contact-writer.ts
import { sql, type SQL } from 'drizzle-orm';
import { contacts } from '../db/schema/contacts';
import { NAME_TRUST } from '../db/schema/name-trust';
import type { NameTrustSource } from '../db/schema/name-trust';
export interface UpsertIdentityParams {
phone: string;
storeId: string;
franchiseId: string; // NOT NULL in DB
accountId: string; // NOT NULL in DB
firstName?: string;
lastName?: string;
trustScore: number; // 用 NAME_TRUST 常量
activityAt: Date;
}
// ─── upsertIdentity — composes 5-caller-proven helpers ─────────
//
// Pattern source(do NOT reinvent):
// - buildIdentityFields → callytics-common/src/db/helpers/identity-fields.ts:69
// (4 ID fields with NOT NULL storeId narrow — 5 caller verbatim usage:
// transcribe-processor:442, ai-analysis-processor:376, message-processor:301,
// lead-tracking:139, contacts-analyzer:524)
// - lastActivityAtForward → callytics-common/src/db/helpers/last-activity.ts:29
// (forward-only GREATEST SQL fragment)
// - NAME_TRUST CASE WHEN → currently inlined verbatim at
// message-processor:322-324 + contacts-analyzer:524-526.
// **Phase 1 这里 extract 成 helper `nameTrustWinnerSet()`**,顺手收口
// 5 caller 的 copy-paste(一次性 cleanup,不是新发明)
//
// Why compose: 5 个 caller 6 个月稳定运行,helper 已 prove。重写 monolithic SQL
// = 漏 sticky / 漏 forward-only / 不一致命名风险。
export function nameTrustWinnerSet(params: {
trustScore: number;
firstName?: string;
lastName?: string;
eventOccurredAt: Date;
}) {
// Extracted verbatim from message-processor:322-324 + contacts-analyzer:524-526.
// Returns Drizzle .set() shape object — caller spreads into .onConflictDoUpdate({set:...}).
return {
...(params.firstName != null && {
firstName: sql`CASE WHEN ${params.trustScore} >= COALESCE(${contacts.firstNameTrustScore}, 0) THEN ${params.firstName} ELSE ${contacts.firstName} END`,
firstNameTrustScore: sql`CASE WHEN ${params.trustScore} >= COALESCE(${contacts.firstNameTrustScore}, 0) THEN ${params.trustScore} ELSE ${contacts.firstNameTrustScore} END`,
firstNameUpdatedAt: sql`CASE WHEN ${params.trustScore} >= COALESCE(${contacts.firstNameTrustScore}, 0) THEN ${params.eventOccurredAt} ELSE ${contacts.firstNameUpdatedAt} END`,
}),
...(params.lastName != null && {
lastName: sql`CASE WHEN ${params.trustScore} >= COALESCE(${contacts.firstNameTrustScore}, 0) THEN ${params.lastName} ELSE ${contacts.lastName} END`,
}),
};
}
export function upsertIdentity(
client: DrizzleClient,
params: UpsertIdentityParams,
): SQL {
// Compose existing helpers — does NOT emit monolithic raw SQL.
// NOTE: storeId NOT NULL guard enforced by buildIdentityFields type.
return client.insert(contacts).values({
...buildIdentityFields({
phone: params.phone,
storeId: params.storeId,
franchiseId: params.franchiseId,
accountId: params.accountId,
}),
firstName: params.firstName ?? null,
lastName: params.lastName ?? null,
firstNameTrustScore: params.trustScore,
firstNameUpdatedAt: params.activityAt,
lastActivityAt: lastActivityAtForward(params.activityAt),
}).onConflictDoUpdate({
target: [contacts.phone, contacts.storeId],
set: {
...nameTrustWinnerSet({
trustScore: params.trustScore,
firstName: params.firstName,
lastName: params.lastName,
eventOccurredAt: params.activityAt,
}),
lastActivityAt: lastActivityAtForward(params.activityAt),
updatedAt: sql`NOW()`,
},
}) as unknown as SQL;
}
// ─── setDNC — sticky pattern verbatim from message-processor:332-335 ────
//
// Pattern source(do NOT reinvent):
// - message-processor/src/infrastructure/neon-repository.ts:332-335
// - contacts-analyzer/src/infrastructure/neon-repository.ts:510-515
// Both use the same sticky CASE WHEN to PRESERVE the original updatedBy
// when DNC was already set. A blind UPDATE SET do_not_contact_updated_by
// = $updatedBy LOSES the original audit trail (TCPA regression).
export interface SetDncParams {
phone: string;
storeId: string;
/** Only stamped on the false→true transition. If DNC was already true,
* the existing updatedBy is preserved (sticky audit trail). */
updatedBy: 'staff' | 'ai' | 'system';
}
export function setDNC(client: DrizzleClient, params: SetDncParams): SQL {
return client.update(contacts)
.set({
...stickyDncFragments({ newDnc: true, updatedBy: params.updatedBy }),
updatedAt: sql`NOW()`,
})
.where(and(
eq(contacts.phone, params.phone),
eq(contacts.storeId, params.storeId),
)) as unknown as SQL;
}
/**
* Sticky DNC SQL fragments — 返回 row write 内嵌入 .set({...}) 的 spread fragment。
*
* 用于 caller 已有一个大 UPDATE statement(如 contacts-analyzer 写 18 个 analysis
* field)需要复用 sticky audit pattern 的场景 — 避免拆 statement、避免 inline
* verbatim CASE WHEN(reinvent)。
*
* Pattern source verbatim: message-processor:332-335。
*/
export function stickyDncFragments(params: {
/** false 时表示不动 DNC(用于 newDnc 取自 nullable AI output);true 时 sticky 翻转 */
newDnc: boolean;
/** 仅在 false → true transition 时 stamp。若 DNC 已 true,保留原 updatedBy */
updatedBy: 'staff' | 'ai' | 'system';
}): {
doNotContact: SQL;
doNotContactUpdatedBy: SQL;
} {
return {
// 永不 false → true 反转 → 永不解 DNC
doNotContact: sql`CASE
WHEN ${contacts.doNotContact} = true THEN true
ELSE ${params.newDnc}
END`,
// 第一次 set 的 updatedBy 永远 sticky(TCPA audit invariant)
doNotContactUpdatedBy: sql`CASE
WHEN ${contacts.doNotContact} = true THEN ${contacts.doNotContactUpdatedBy}
WHEN ${params.newDnc} = true THEN ${params.updatedBy}
ELSE ${contacts.doNotContactUpdatedBy}
END`,
};
}
// ─── touchActivity — composes lastActivityAtForward ─────────────────
//
// Pattern source(do NOT reinvent):
// - lastActivityAtForward → callytics-common/src/db/helpers/last-activity.ts:29
// This helper produces the GREATEST(...) fragment;touchActivity wraps it
// in a standalone UPDATE for callers that don't have a surrounding UPSERT
// (e.g. existing-only contact, no identity write needed).
export interface TouchActivityParams {
phone: string;
storeId: string;
/** Event-derived timestamp(msg.creationTime / call.startTime),
* NOT new Date(). See lastActivityAtForward provenance contract. */
activityAt: Date;
}
export function touchActivity(
client: DrizzleClient,
params: TouchActivityParams,
): SQL {
return client.update(contacts)
.set({
lastActivityAt: lastActivityAtForward(params.activityAt),
updatedAt: sql`NOW()`,
})
.where(and(
eq(contacts.phone, params.phone),
eq(contacts.storeId, params.storeId),
)) as unknown as SQL;
}
Imports needed at top of contact-writer.ts:
import { sql, and, eq } from 'drizzle-orm';
import { contacts } from '../db/schema';
import { buildIdentityFields } from '../db/helpers/identity-fields';
import { lastActivityAtForward } from '../db/helpers/last-activity';
import type { DrizzleClient } from '../db/types';
Where nameTrustWinnerSet() lives:推荐放 callytics-common/src/db/helpers/name-trust-winner.ts(跟其他 helper 同位置),不要藏在 domain/contact-writer.ts 里 —— 让 5 个现有 caller(transcribe / ai-analysis / message / contacts-analyzer / lead-tracking)有机会以后 PR 一起切换到这个 helper。Phase 1 不强迫 5 caller migrate,但 helper 位置预留。
Step 3: Test write + run
Implementer 按 case map:
- 6 upsertIdentity scenarios
- 3 setDNC scenarios
- 3 touchActivity scenarios
Integration test 跑真实 Neon verify SQL 行为(NAME_TRUST winner / forward-only)。
Step 4: Commit
git add src/domain/contact-writer.ts tests/domain/contact-writer.test.ts tests/integration/contact-writer.test.ts
git commit -m "feat(domain): contact-writer — compose existing helpers,sticky DNC,forward-only activity
Phase 1 Plan 02 Task 4.
Pattern source(do NOT reinvent):
- buildIdentityFields → callytics-common/src/db/helpers/identity-fields.ts:69
- lastActivityAtForward → callytics-common/src/db/helpers/last-activity.ts:29
- NAME_TRUST sticky CASE WHEN → message-processor:322-324 verbatim
(extract 成 nameTrustWinnerSet helper,5 caller 以后可一起 cleanup)
- sticky DNC CASE WHEN → message-processor:332-335 verbatim
(TCPA audit trail — DNC=true 时保留原 updatedBy,防覆盖)
- upsertIdentity: compose buildIdentityFields + lastActivityAtForward + nameTrustWinnerSet
- setDNC: sticky CASE WHEN(do_not_contact_updated_by 一旦 set 不被覆盖)
- touchActivity: compose lastActivityAtForward(standalone UPDATE wrapper)
Spec: normative-spec.md §3.4
"
Task 5: timeline-writer.ts — TIMELINE_EVENT_TYPES + Zod payload + writeTimelineEvent
Files:
- Create:
src/domain/timeline-writer.ts
- Test:
tests/domain/timeline-writer.test.ts(Zod validation per eventType)
Step 1: Map test scenarios
Per-eventType Zod payload schema validation(spec §3.5 line 404-414):
writeTimelineEvent() invariants:
Step 2: Implementation
// src/domain/timeline-writer.ts
import { z } from 'zod';
import { sql, type SQL } from 'drizzle-orm';
import { contactTimeline } from '../db/schema/contact-timeline';
import { buildTimelineValues } from '../db/helpers/timeline-values';
export const TIMELINE_EVENT_TYPES = [
'task.created', 'task.status_changed', 'task.updated',
'task.progress_recorded', 'task.note_updated',
'contact.lifecycle_changed', 'contact.lead_status_changed', 'contact.dnc_changed',
'contact.complaint_opened', 'contact.complaint_resolved',
'contact_analysis.completed',
'transcribe.completed', 'call_analysis.completed',
'call.status_changed', 'message.created',
'lead.created',
] as const;
export type TimelineEventType = (typeof TIMELINE_EVENT_TYPES)[number];
// ─── Per-eventType payload schemas ──────────────────────────────
const taskCreatedSchema = z.object({
taskId: z.string().uuid(),
typeCategory: z.string(),
sourceType: z.enum(['lead', 'contact_analysis', 'manual']),
});
const taskStatusChangedSchema = z.object({
taskId: z.string().uuid(),
oldStatus: z.enum(['pending', 'open', 'closed']),
newStatus: z.enum(['open', 'closed']),
closeResult: z.string().optional(),
});
const taskProgressRecordedSchema = z.object({
taskId: z.string().uuid(),
progressType: z.enum([
'no_answer','left_voicemail','text_sent',
'callback_requested','follow_up_scheduled','customer_considering',
]),
channel: z.enum(['phone','sms','voicemail','email']),
callId: z.string().optional(),
messageId: z.string().optional(),
});
// ... 类似 13 个其余 schema
const PAYLOAD_SCHEMAS: Record<TimelineEventType, z.ZodSchema> = {
'task.created': taskCreatedSchema,
'task.status_changed': taskStatusChangedSchema,
'task.progress_recorded': taskProgressRecordedSchema,
// ... 13 more
'task.updated': z.object({ taskId: z.string().uuid(), changedFields: z.array(z.string()) }),
'task.note_updated': z.object({ taskId: z.string().uuid(), note: z.string() }),
'contact.lifecycle_changed': z.object({ phone: z.string(), oldStage: z.string(), newStage: z.string() }),
'contact.lead_status_changed': z.object({ phone: z.string(), oldStatus: z.string(), newStatus: z.string() }),
'contact.dnc_changed': z.object({ phone: z.string(), dncValue: z.boolean(), updatedBy: z.string() }),
'contact.complaint_opened': z.object({ phone: z.string(), complaintSummary: z.string() }),
'contact.complaint_resolved': z.object({ phone: z.string() }),
'contact_analysis.completed': z.object({ phone: z.string(), runId: z.string() }),
'transcribe.completed': z.object({ callId: z.string() }),
'call_analysis.completed': z.object({
callId: z.string(),
primaryOutcomeResult: z.string().optional(),
followUpNeeded: z.boolean().optional(),
}),
'call.status_changed': z.object({ callId: z.string(), oldStatus: z.string(), newStatus: z.string() }),
'message.created': z.object({ messageId: z.string(), direction: z.enum(['inbound','outbound']), messageType: z.string() }),
'lead.created': z.object({ leadId: z.string(), leadType: z.string(), source: z.string() }),
};
export interface WriteTimelineEventParams {
eventType: TimelineEventType;
contactPhone: string;
storeId: string;
franchiseId: string;
accountId: string;
actor: {
type: string;
subjectId?: string;
name?: string;
};
payload: unknown;
idempotencyKey: string;
occurredAt?: Date;
entityType?: string;
entityId?: string;
}
/**
* writeTimelineEvent —— wraps existing helpers, does NOT replace them.
*
* Pattern source(do NOT reinvent):
* - buildTimelineValues → callytics-common/src/db/helpers/timeline-values.ts:92
* - buildContactTimelineInsertSQL → callytics-common/src/db/helpers/contact-timeline-insert-sql.ts:35
*
* What this module adds(net new — justified):
* - Per-eventType Zod payload validation(防 JSONB content drift)
* - deriveCategory(eventType prefix → eventCategory mapping)
*
* What this module deliberately does NOT do:
* - Issue a custom raw `INSERT ... ON CONFLICT (idempotency_key) DO NOTHING`
* —— contact-timeline-insert-sql.ts:18-29 自己 verbatim 注释 "DO NOT write
* ON CONFLICT (idempotency_key)" 因为 `idx_timeline_idempotency` 是 PARTIAL
* unique index(`WHERE idempotency_key IS NOT NULL`),Postgres 要求 conflict
* target 跟 partial predicate 完全 match,否则报错 "no unique or exclusion
* constraint matching ON CONFLICT specification"。
* 正确做法:用 bare `ON CONFLICT DO NOTHING`(no target)—— `onConflictDoNothing()`
* 不传 target,Drizzle 渲染成 bare form。
* - Build values from scratch —— buildTimelineValues 已 cover 全字段映射。
*/
export function writeTimelineEvent(
client: DrizzleClient,
params: WriteTimelineEventParams,
): SQL {
// 1. Net new: Zod payload validation gate
const schema = PAYLOAD_SCHEMAS[params.eventType];
const validated = schema.parse(params.payload); // throws on invalid shape
// 2. Compose existing helper(do NOT rebuild values shape)
const values = buildTimelineValues({
contactPhone: params.contactPhone,
franchiseId: params.franchiseId,
accountId: params.accountId,
storeId: params.storeId,
eventType: params.eventType,
eventCategory: deriveCategory(params.eventType),
entityType: params.entityType,
entityId: params.entityId,
newValue: validated as Record<string, unknown>,
occurredAt: params.occurredAt ?? new Date(),
actorType: params.actor.type as any,
actorName: params.actor.name,
actorSubjectId: params.actor.subjectId,
idempotencyKey: params.idempotencyKey,
});
// 3. Bare ON CONFLICT DO NOTHING(no target — matches partial index caveat)
// Pattern: 5 个 Lambda 现有 `.onConflictDoNothing()` 用法 verbatim
return client.insert(contactTimeline)
.values(values)
.onConflictDoNothing() as unknown as SQL;
}
/** Raw SQL variant for studio-api(neon-http path).Compose existing helper. */
export function buildTimelineEventSQL(
params: WriteTimelineEventParams,
): { text: string; params: unknown[] } {
const schema = PAYLOAD_SCHEMAS[params.eventType];
const validated = schema.parse(params.payload);
// buildContactTimelineInsertSQL 已经 emit `INSERT ... ON CONFLICT DO NOTHING`
// (no target),直接 wrap。
return buildContactTimelineInsertSQL({
contactPhone: params.contactPhone,
franchiseId: params.franchiseId,
accountId: params.accountId,
storeId: params.storeId,
eventType: params.eventType,
eventCategory: deriveCategory(params.eventType),
entityType: params.entityType,
entityId: params.entityId,
newValue: validated as Record<string, unknown>,
occurredAt: params.occurredAt ?? new Date(),
actorType: params.actor.type as any,
actorName: params.actor.name,
actorSubjectId: params.actor.subjectId,
idempotencyKey: params.idempotencyKey,
});
}
function deriveCategory(et: TimelineEventType): 'task' | 'contact' | 'call' | 'message' | 'lead' {
if (et.startsWith('task.')) return 'task';
if (et.startsWith('contact')) return 'contact';
if (et.startsWith('call')) return 'call';
if (et.startsWith('message')) return 'message';
if (et.startsWith('lead')) return 'lead';
return 'contact';
}
Imports needed at top:
import { z } from 'zod';
import type { SQL } from 'drizzle-orm';
import { contactTimeline } from '../db/schema/contact-timeline';
import { buildTimelineValues } from '../db/helpers/timeline-values';
import { buildContactTimelineInsertSQL } from '../db/helpers/contact-timeline-insert-sql';
import type { DrizzleClient } from '../db/types';
Step 3: Test write + run
Implementer 按 case map 写 16 个 happy test(每 eventType 一个)+ 至少 3 reject test(随便选 3 个 eventType 故意传错 payload)+ 1 invariant test(同 idempotencyKey 不重复 INSERT)。
Step 4: Commit
git add src/domain/timeline-writer.ts tests/domain/timeline-writer.test.ts
git commit -m "feat(domain): timeline-writer — wrap buildTimelineValues + Zod payload gate
Phase 1 Plan 02 Task 5.
Pattern source(do NOT reinvent):
- buildTimelineValues → callytics-common/src/db/helpers/timeline-values.ts:92
(Drizzle path,5 caller verbatim usage)
- buildContactTimelineInsertSQL → callytics-common/src/db/helpers/
contact-timeline-insert-sql.ts:35 (raw SQL path,studio-api 用)
Net new(justified):
- 16 eventType TimelineEventType union(加 task.progress_recorded)
- PAYLOAD_SCHEMAS Zod per-eventType validation gate(防 JSONB content drift)
- deriveCategory(eventType prefix → eventCategory mapping)
Bug avoided:
- NOT写 ON CONFLICT (idempotency_key) DO NOTHING — partial unique index
WHERE clause makes target inference fail with PG error. 用 bare
ON CONFLICT DO NOTHING(no target),via .onConflictDoNothing()
(Drizzle 默认无 target)/ buildContactTimelineInsertSQL 已 emit bare form。
See contact-timeline-insert-sql.ts:18-29 verbatim caveat。
Spec: normative-spec.md §3.5
"
Task 6: task-orchestrator-sql.ts — raw SQL adapter for studio-api
Files:
- Create:
src/domain/task-orchestrator-sql.ts
- Test:
tests/domain/task-orchestrator-sql.test.ts
Step 1: Map test scenarios
buildTaskActionSQL() —— output 验证:
buildCloseAllOpenSQL() — 与 closeAllOpenForContact 行为 match,只 output 差。
Key invariant:对同一 input,raw SQL 版和 Drizzle 版的 reject reason 必须 exactly match(防 caller 收到不一致 reason)。
Step 2: Implementation
// src/domain/task-orchestrator-sql.ts
import type { TaskAction, ApplyResultSQL, ActionContext } from './types';
// Note:
// - Do NOT convert Drizzle SQL objects with `.toQuery()` / `as any`.
// That is not a stable public contract and can silently produce empty SQL.
// - Keep the current helper style: explicit raw SQL builders returning
// `{ sql, params }`, same pattern as buildContactTimelineInsertSQL().
// - Share Policy Guard / idempotency / due-date helpers; duplicate only the
// final SQL rendering where Drizzle and neon-http need different shapes.
export async function buildTaskActionSQL(
action: TaskAction,
ctx: ActionContext,
): Promise<ApplyResultSQL> {
switch (action.action) {
case 'close':
return buildCloseTaskSQL(action.payload, ctx);
case 'record_progress':
return buildRecordProgressSQL(action.payload, ctx);
// create_open / create_closed / update / reopen follow the same explicit
// raw helper pattern. Each helper returns complete transactional units
// plus resultChecks; callers only wrap them in sql.transaction().
default:
return buildOtherTaskActionSQL(action, ctx);
}
}
export async function buildCloseAllOpenSQL(
params: Parameters<typeof import('./task-orchestrator').closeAllOpenForContact>[0],
ctx: { sql: unknown },
): Promise<{ status: 'allow'; statements: { sql: string; params: unknown[] }[]; closedTaskIds: string[] }
| { status: 'reject'; reason: any; details: string }> {
return buildCloseAllOpenForContactSQL(params, ctx);
}
Step 3: Test + commit
git add src/domain/task-orchestrator-sql.ts tests/domain/task-orchestrator-sql.test.ts
git commit -m "feat(domain): task-orchestrator-sql — raw SQL adapter for studio-api
Phase 1 Plan 02 Task 6.
- buildTaskActionSQL / buildCloseAllOpenSQL 共享 Policy Guard / idempotency / due-date helper
- raw SQL helper 显式返回 `{ sql, params }`,沿用 `buildContactTimelineInsertSQL()` 风格,不通过 Drizzle SQL.toQuery() 转换
- caller (studio-api) 用 sql.transaction(statements.map(...)) 包裹
Spec: normative-spec.md §3.2
"
Task 7: Export + version bump + publish PR
Files:
- Create:
src/domain/index.ts
- Modify:
src/index.ts and package.json
Step 1: barrel export
// src/domain/index.ts
export * from './types';
export * from './policy-guard';
export * from './task-orchestrator';
export * from './task-orchestrator-sql';
export * from './contact-writer';
export * from './timeline-writer';
export * from './idempotency';
Step 2: Modify package.json exports
{
"exports": {
".": { /* existing */ },
"./domain": {
"types": "./dist/domain/index.d.ts",
"import": "./dist/domain/index.js"
},
// ... existing other paths
},
"version": "1.2.0" // bumped from 1.1.0 (Plan 01)
}
Step 3: Run full test suite + build
cd /Users/maxwsy/workspace/callytics-common && bun test && bun run build
Expected: all PASS + dist/ contains domain/ subdir。
Step 4: Commit + Push + Open PR
git add src/domain/index.ts src/index.ts package.json
git commit -m "chore: export @retaintive/common/domain + bump to 1.2.0"
git push -u origin <feature-branch>
gh pr create --title "feat(domain): Phase 1 Plan 02 — 4 shared modules" --body-file <(cat << 'EOF'
## Plan 02: 4 共享 module
Implements:
- task-orchestrator.ts (applyTaskAction 6 actions + closeAllOpenForContact)
- policy-guard.ts (9 checks + computeAllowedTypeCategories + computeNextDueAt)
- contact-writer.ts (upsertIdentity / setDNC / touchActivity)
- timeline-writer.ts (16 eventType + Zod payload)
- task-orchestrator-sql.ts (raw SQL adapter for studio-api)
- idempotency.ts (4 source pattern)
Test coverage:
- Unit: 8+9+30+12+19 = ~78 tests
- Integration (Neon): 5 SQL invariants (CTE attempt_count / human authority / partial unique x2 / closeAllOpen atomic)
Spec: normative-spec.md §3 全节 + §5.0 invariants
EOF
)
Self-Review Checklist
Spec coverage
- ✅ §3.1 TaskAction union → Task 1 types.ts
- ✅ §3.1 6 action dispatch → Task 3 task-orchestrator.ts
- ✅ §3.1 ActionContext 拆 2 type (S6) → Task 1
- ✅ §3.1 RejectReason 13 values → Task 1
- ✅ §3.1 priority 派生 (S1) → Task 3 handleCreateOpen
derivePriority
- ✅ §3.1 closeAllOpenForContact (B7) → Task 3
- ✅ §3.2 buildTaskActionSQL (raw SQL) → Task 6
- ✅ §3.2 idempotency_key 4 source (B8 CTE) → Task 2 idempotency.ts + Task 3 handleRecordProgress
- ✅ §3.3 computeNextDueAt → Task 2
- ✅ §3.4 ContactWriter 3 fragments (B5 franchiseId/accountId) → Task 4
- ✅ §3.5 TIMELINE_EVENT_TYPES 16 + Zod → Task 5
- ✅ §3.6 9 checks + computeAllowedTypeCategories 7 case (S2) → Task 2
- ✅ §3.6 create_closed-has-no-open guard (S3) → Task 2
- ✅ §3.6 DNC reject reopen (I6) → Task 2 checkDnc
- ✅ §5.0 invariants C1-C5 → Task 3 integration tests
Placeholder scan
- ✅ 每 task 有具体 file path + line / new file 标 Create
- ✅ Test scenarios 全表格列清(无 "write tests for the above")
- ✅ Implementation 关键 SQL pattern 完整给(CTE / CASE WHEN / partial unique 触发)
- ⚠️ conditional mutation + timeline helpers 在 Task 3
declare 引用,实际 impl 与 Task 5 TimelineWriter 协作。implementer 必须按 task 顺序做 — 跳序会撞 unresolved import
Type consistency
- ✅
TaskAction / ApplyResult / RejectReason 在 types.ts 定义,Task 2-6 全 import 自此
- ✅
RejectReason 字符串字面在 policy-guard / task-orchestrator / task-orchestrator-sql 三处都 match types.ts
- ✅
computeNextDueAt signature 在 Task 2 定义,Task 3 handleRecordProgress 用 — 参数顺序 match
- ✅
TIMELINE_EVENT_TYPES 在 Task 5 定义,Task 3 timeline helpers 用对应 event names(task.created / task.status_changed / task.progress_recorded)
Outstanding
无。Plan 02 完成 = @retaintive/common@1.2.0 含 ./domain export,Plans 03-06 caller 升级 dependency 即可使用。
Execution Handoff
Option 1: Subagent-Driven — 每 Task 一 fresh subagent,review 之间 iterate
Option 2: Inline — superpowers:executing-plans 跑全 7 task,checkpoint 之间停
Plan 02 是最大 module,建议 subagent — 每 task 完成后 user review unit test 覆盖度 + SQL contract 是否 sane,再进下一 task。