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

FileResponsibilityAction
src/domain/types.ts共享 type:TaskAction union / ApplyResult / RejectReason / ActionContext 拆 2 type / ContactLifecycleCreate
src/domain/policy-guard.ts9 个 check function + computeAllowedTypeCategories + computeNextDueAtCreate
src/domain/task-orchestrator.tsapplyTaskAction() + closeAllOpenForContact() (Drizzle 返回 SQL[])Create
src/domain/contact-writer.tsupsertIdentity() / setDNC() / touchActivity()Create
src/domain/timeline-writer.tswriteTimelineEvent() + TIMELINE_EVENT_TYPES + Zod payload schemasCreate
src/domain/task-orchestrator-sql.tsbuildTaskActionSQL() / buildCloseAllOpenSQL() (raw SQL 给 studio-api)Create
src/domain/idempotency.tsbuildProgressIdempotencyKey() 4 source patternCreate
src/domain/index.tsbarrel exportCreate
src/index.ts./domain export pathModify
package.json./domain 到 exports map + bump 1.1.0 → 1.2.0Modify

测试文件全 mirror 在 tests/domain/:

Test file覆盖
tests/domain/policy-guard.test.ts9 个 check 各自 happy + reject + edge case
tests/domain/task-orchestrator.test.ts6 action × happy + reject reasons + 并发场景
tests/domain/contact-writer.test.tsNAME_TRUST 6 level × winner / DNC sticky / forward-only
tests/domain/timeline-writer.test.tsZod schema enforcement per eventType + idempotency
tests/domain/task-orchestrator-sql.test.tsraw SQL output 语义 match Drizzle 版
tests/domain/idempotency.test.ts4 source key format 各自 + edge case
tests/integration/task-orchestrator.test.ts跑真实 Neon,verify SQL contract(human authority guard / partial unique / CTE)

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

Type用途来源
TaskAction discriminated union6 个 action 的 payload 类型normative-spec §3.1 line 157-163
CreateOpenPayload / CreateClosedPayload / ClosePayload / UpdatePayload / RecordProgressPayload / ReopenPayload每个 action 各自的 payloadspec §3.1
ApplyResult union(allow / reject / needs_review)Orchestrator 返回值spec §3.1 line 165-168
RejectReason string literal union(13 values)reject 原因spec §3.1 line 170-184
BaseActionContext共享 context 字段spec §3.1 line 189-194
CreateActionContext extends Base + contactPhonecreate_open / create_closed 用spec §3.1 line 196-199
TaskIdActionContext = BaseActionContext4 个 taskId action 用spec §3.1 line 201-202
ActionContext = union of twospec §3.1 line 204
ContactLifecycle(lifecycleStage / lifecycleState / leadStatus / doNotContact)computeAllowedTypeCategories 输入spec §3.6 line 450-455
ApplyResultSQL(raw SQL 版返回)buildTaskActionSQL() 返回spec §3.2 line 265-269

Step 2: Tests scenarios

原则:types.ts 是纯 type definition,没有 runtime logic。只 test 真 invariant —— type narrowing 在使用端工作正常(防写法 drift)。

#Scenario为什么 test
1TaskAction discriminated union — narrow by action field 后,payload 类型 narrowed 到对应 specific payload防有人改 union 丢失 discriminator
2ApplyResult narrow by status —— allow 状态有 statements,rejectreason + details防 union 写错
3CreateActionContext 要求 contactPhone non-optional,TaskIdActionContext 不要求防 S6 fix(spec §3.1 ActionContext 拆 2 type)被回退

不写其他 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):

#Input contact stateExpected output为什么
1doNotContact=true(其他字段任意)[]DNC 不允许任何 outreach create
2lifecycleState='terminal' 非 churned[]terminal lead 不允许 outreach
3lifecycleStage='lead', state='active'['lead_follow_up', 'booked_not_converted']AI 不允许 lead_outreach(lead-tracking 专属)
4lifecycleStage='lead', state='paused'['lead_follow_up', 'booked_not_converted']paused 不影响 allowed set(spec 没明说但 active/paused 同等)
5lifecycleStage='member', state='active'['cancellation_risk', 'retention', 'upgrade', 'renewal', 'referral']5 个 member 专属 category
6lifecycleStage='churned', state='active'['win_back']churned re-engage 只允许 win_back
7lifecycleStage='churned', state='terminal'[]terminal churned 不允许(re-engage 必须先被 AI 判成 active)
8lifecycleStage='unknown'(任何 state)[]unknown 不允许 outreach

computeNextDueAt(progressType, currentDueAt, override) — 6 progressType × override 有/无(spec §3.3):

#InputExpected为什么
1no_answer, currentDueAt 任意, override=undefinednow + 60minspec line 245
2left_voicemail, override=undefinednow + 24hspec line 246
3text_sent, override=undefinednow + 48hspec line 247
4callback_requested, currentDueAt=X, override=undefinedcurrentDueAt(不变)spec line 248 — staff 必须手动设
5follow_up_scheduled, override=undefinedcurrentDueAt(不变)spec line 249
6customer_considering, override=undefinednow + 72hspec line 250
7任何 progressType, override=2026-12-252026-12-25(override wins)spec line 257 caller override 总赢

Policy Guard 8 checks(snapshot-based,无 DB call)—— 每个 check pure function over ContactSnapshot:

Check输入Happy(allow)Reject scenarios
checkStoreId(ctx)ctx.storeIdstoreId='STORE_A' → allowstoreId='' → reject store_mismatch
checkDnc(snapshot, action) spec §3.6 line 438ContactSnapshot.contact.doNotContactsnapshot.contact=null(新 contact)→ allow / doNotContact=false → allowdoNotContact=true × action ∈ {create_open, create_closed, update, record_progress, reopen} → reject dnc;action='close' → allow
checkTaskBelongsToContact(snapshot, taskId)ContactSnapshot.pendingTasks Set lookuptaskId ∈ snapshot.pendingTasks → allowtaskId ∉ snapshot.pendingTasks → reject hallucinated_task_id
checkStateTransition(action, currentStatus)task.status from snapshotclose 对 open → allow;reopen 对 closed → allowclose 对 closed → reject task_not_open / reopen 对 open → reject task_not_closed / update 对 closed → reject task_not_open
checkCreateClosedHasNoOpenTask(snapshot, typeCategory) spec §3.6 line 442ContactSnapshot.pendingTasks Set无同 typeCategory open task → allow已有 → reject create_closed_with_open_task
checkTypeCategoryAllowed(ctx, contact, typeCategory) spec §3.6 line 443ContactLifecycleactor.type='staff'/'system' → skip allow / 'ai_agent' + typeCategory ∈ allowed → allow'contact_analysis' + typeCategory ∉ allowed → reject invalid_type_category
checkCloseNoteRequired(payload)payloadcloseResult='other' + closeNote 非空 → allowcloseResult='other' + closeNote=''/undefined → reject close_note_required
checkConfidence(confidence)proposal confidenceundefined / ≥ MIN_AI_CONFIDENCE → allow< threshold → reject low_confidence

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):

#ScenarioExpected
1新 contact(尚未 INSERT)→ loadSnapshot 返回 {contact: null, pendingTasks: []}不 throw,后续 check 能 handle null contact
2已存在 contact + 2 open task(过渡期 1 个 'pending' + 1 个 'open')→ 全 coversnapshot.pendingTasks.length === 2,status IN ('pending','open') filter 工作
3跨 store(另一 storeId 的 task 不漏)→ tenant isolation不漏

buildProgressIdempotencyKey() — 4 source pattern(spec §3.2 table line 296-301):

SourceInputExpected key format
phone{taskId:'T1', callId:'C1', progressType:'no_answer'}progress:T1:call:C1:no_answer
sms{taskId:'T1', messageId:'M1', progressType:'text_sent'}progress:T1:message:M1:text_sent
manual{taskId:'T1', actorId:'staff1', occurredAt:Date('2026-01-01'), progressType:'callback_requested'}progress:T1:manual:staff1:2026-01-01T00:00:00.000Z:callback_requested
system{taskId:'T1', runId:'R1', progressType:'follow_up_scheduled'}progress:T1:system:R1:follow_up_scheduled

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 对比:

设计DB call/decisionDB call/batch (20 decisions)
Plan 02 旧设计(per-check SELECT)4-580-100
新设计(snapshot pattern)0(snapshot 已加载)2(1 contact + 1 pendingTasks)

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
"

Task 3: task-orchestrator.ts — applyTaskAction() + closeAllOpenForContact() (Drizzle)

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:

ActionHappy path scenariosReject scenarios
create_open1.staff manual create → allow + INSERT statement / 2.AI 提案 confidence=0.8 typeCategory ∈ allowed → allow / 3.lead-processor sourceType='lead' skip type check → allow1.storeId 缺 → store_mismatch / 2.DNC contact → dnc / 3.AI confidence=0.5 → low_confidence / 4.AI 选 invalid typeCategory → invalid_type_category / 5.同 (phone,store,typeCategory) 已有 open task → duplicate (partial unique 触发,SQL 层)
create_closed1.staff inbound 当场成交 + sourceCallId 唯一 → allow + INSERT closed1.同上 4 个;2.已有 open task → create_closed_with_open_task;3.同 sourceCallId 已建过 closed → duplicate;4.closeResult='other' 无 closeNote → close_note_required
close1.staff close open task → allow / 2.AI close + aiRunStartedAt 新于 task.updated_at → allow1.taskId 不属于 contact → hallucinated_task_id;2.task 已 closed → task_not_open;3.AI 提案 stale(task 在 aiRunStartedAt 后被 staff 改了)→ SQL UPDATE 0 rows → stale_proposal;4.closeResult='other' 无 closeNote → close_note_required;5.DNC contact 仍 allow(close 是 DNC allowed action)
update1.staff 改 dueAt → allow / 2.AI 改 priority → allow1.task closed → task_not_open;2.AI stale → stale_proposal;3.hallucinated taskId → hallucinated_task_id;4.DNC contact → dnc(update 在 DNC allowed 列外)
record_progress1.staff 记 no_answer → INSERT progress event + attempt_count+=1 / 2.AI 记 callback_requested → 同上 + nextDueAt 不变 / 3.同 callId/progressType 已记过 → ON CONFLICT idempotency_key DO NOTHING + attempt_count 不变(CTE 保护)1.task closed → task_not_open;2.DNC contact → dnc(record_progress 在 DNC allowed 列外);3.hallucinated taskId → hallucinated_task_id
reopen1.staff reopen closed → allow + UPDATE status='open' + close 字段清空1.task open → task_not_closed;2.DNC contact → dnc(reopen 在 DNC allowed 列外);3.hallucinated taskId → hallucinated_task_id

Critical SQL invariants(integration test 验):

#InvariantTest setup
C1record_progress CTE 防 attempt_count 翻倍seed task w/ attempt=2 + seed 1 progress event;再调 record_progress with same idempotency_key → expect attempt_count 仍 = 2(不+1)
C2close SQL WHERE status='open' AND updated_at <= $aiRunStartedAt 防 staleseed task w/ updated_at=NOW;模拟 AI aiRunStartedAt=NOW-1h;调 close → expect 0 rows updated → reject stale_proposal
C3create_open partial unique 防 duplicateseed 1 open task (phone=A, store=B, typeCategory=C);调 create_open 同 (A,B,C) → expect SQL 报 partial unique conflict → reject duplicate
C4create_closed partial unique 防同 sourceCallId 重复建seed 1 closed task w/ source_call_id='CALL1';再调 create_closed 同 sourceCallId → expect duplicate
C5closeAllOpenForContact() 关多个 open task + 写多条 timeline 原子seed 3 open tasks for contact;调 closeAllOpen → expect 3 UPDATE + 3 timeline INSERT 全成功

closeAllOpenForContact() happy + edge:

#ScenarioExpected
1contact 有 3 open task → bulk close 全部返回 closedTaskIds 含 3 个 / 生成 3 个 timeline event(每 task 1 个 task.status_changed)
2contact 无 open task → no-opclosedTaskIds=[] / 不生成 SQL statement
3DNC trigger 用例:closeResult='do_not_contact' + closeNote='Auto-closed by DNC hard stop' + actor.type='system'全部 closeType='auto_closed'

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 略

Step 5: Implementation — closeAllOpenForContact

// 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
"

Task 4: contact-writer.ts — upsertIdentity / setDNC / touchActivity

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):

#ScenarioExpected
1新 contact(不存在)→ INSERT 完整 row,first_name_trust_score = 输入 trustScorerow 拿回 trustScore match
2已存在 contact,trustScore=80(LEAD)旧 / 现在写 trustScore=100(STAFF)first_name 更新为新值,trust_score 更新为 100
3已存在 contact,trustScore=80(LEAD)旧 / 现在写 trustScore=40(AI_TRANSCRIPT)first_name 不变,trust_score 保留 80
4已存在 contact,trustScore=60 / 现在写 trustScore=60(同分)用新值(同分新覆盖旧 — spec 行为)
5franchiseId / accountId NULL → UPSERT 失败(NOT NULL violation in DB)caller 检查必传
6UPSERT 时 storeId 不变 但 franchiseId 变了 → 老 row 的 franchiseId 不更新(只走 INSERT path 时设)验 ON CONFLICT DO UPDATE 不动 franchise/account

setDNC() —— sticky CASE WHEN audit trail(spec §3.4 line 377-381;pattern source message-processor:332-335):

#ScenarioExpectedWhy
1contact doNotContact=false → setDNC(updatedBy='staff')→ DNC=true,doNotContactUpdatedBy='staff'false→true transition stamps caller's updatedBy正常
2关键 — TCPA audit invariant —— contact doNotContact=true(updatedBy='staff')→ setDNC(updatedBy='ai')→ DNC=true,doNotContactUpdatedBy 仍 'staff'(不被 'ai' 覆盖)第一次 set 的 updatedBy 永远 stickysilent bug —— blind UPDATE SET 会丢 staff 审计 trail。这条 test 真守住
3updatedBy='staff' / 'ai' / 'system' 三种 false→true 都接受三种都正确 stampenum 完整
4没有 unsetDNC API(sticky 不可反转)TypeScript 编译时只有 setDNC,无 false 写入路径type-level guard

touchActivity() —— forward-only:

#ScenarioExpected
1contact 不存在 → 写 lastActivityAt = NOW()INSERT 含 lastActivityAt
2contact lastActivityAt=2026-01-01,activityAt=2026-06-01 → 更新到 2026-06-01row.last_activity_at = 2026-06-01
3contact lastActivityAt=2026-06-01,activityAt=2026-01-01(更旧!)→ 不动,GREATEST 保护row.last_activity_at 仍 2026-06-01

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):

eventTypeRequired payload fieldsTest happy + reject
task.created{taskId, typeCategory, sourceType}happy:full payload pass / reject: missing taskId → Zod error
task.status_changed{taskId, oldStatus, newStatus, closeResult?}happy / reject:newStatus 不在 enum
task.updated{taskId, changedFields:[]}happy / reject:changedFields not array
task.progress_recorded{taskId, progressType, channel, callId?, messageId?}happy / reject:progressType invalid
task.note_updated{taskId, note}happy
contact.lifecycle_changed{phone, oldStage, newStage}happy
contact.lead_status_changed{phone, oldStatus, newStatus}happy
contact.dnc_changed{phone, dncValue, updatedBy}happy / reject:dncValue not boolean
contact.complaint_opened{phone, complaintSummary}happy
contact.complaint_resolved{phone}happy
contact_analysis.completed{phone, runId}happy
transcribe.completed{callId}happy
call_analysis.completed{callId, primaryOutcomeResult, followUpNeeded}happy
call.status_changed{callId, oldStatus, newStatus}happy
message.created{messageId, direction, messageType}happy
lead.created{leadId, leadType, source}happy

writeTimelineEvent() invariants:

#ScenarioExpected
1合法 payload + eventType match → 返回 INSERT SQLSQL contains correct eventType + payload JSON
2payload 不 match eventType Zod → throw at runtimeerror has eventType prefix
3同 idempotencyKey 重复调 → 返回 SQL with ON CONFLICT DO NOTHINGimpl 已加 ON CONFLICT

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 验证:

#ScenarioExpected
1close action → output statements[] 含 UPDATE SQL string + paramsstring 含 UPDATE tasks SET status='closed',params 数组 length 正确
2record_progress action → output 含 CTE SQLstring 含 WITH inserted AS (INSERT ...) UPDATE tasks ...
3reject case → output {status: 'reject', reason, details} 跟 Drizzle 版相同reason 用相同 string literal
4allow case parityDrizzle 版和 raw SQL 版对同一 action 产生相同 resultChecks / generatedIds 语义(不要求 SQL string 完全相同)

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: Inlinesuperpowers:executing-plans 跑全 7 task,checkpoint 之间停

Plan 02 是最大 module,建议 subagent — 每 task 完成后 user review unit test 覆盖度 + SQL contract 是否 sane,再进下一 task。