> For AI agents: the complete documentation index is available at /llms.txt, the full documentation bundle is available at /llms-full.txt.

# 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

| File                                  | Responsibility                                                                                              | Action     |
| ------------------------------------- | ----------------------------------------------------------------------------------------------------------- | ---------- |
| `src/domain/types.ts`                 | 共享 type:`TaskAction` union / `ApplyResult` / `RejectReason` / `ActionContext` 拆 2 type / `ContactLifecycle` | **Create** |
| `src/domain/policy-guard.ts`          | 9 个 check function + `computeAllowedTypeCategories` + `computeNextDueAt`                                    | **Create** |
| `src/domain/task-orchestrator.ts`     | `applyTaskAction()` + `closeAllOpenForContact()` (Drizzle 返回 SQL\[])                                        | **Create** |
| `src/domain/contact-writer.ts`        | `upsertIdentity()` / `setDNC()` / `touchActivity()`                                                         | **Create** |
| `src/domain/timeline-writer.ts`       | `writeTimelineEvent()` + `TIMELINE_EVENT_TYPES` + Zod payload schemas                                       | **Create** |
| `src/domain/task-orchestrator-sql.ts` | `buildTaskActionSQL()` / `buildCloseAllOpenSQL()` (raw SQL 给 studio-api)                                    | **Create** |
| `src/domain/idempotency.ts`           | `buildProgressIdempotencyKey()` 4 source pattern                                                            | **Create** |
| `src/domain/index.ts`                 | barrel export                                                                                               | **Create** |
| `src/index.ts`                        | 加 `./domain` export path                                                                                    | Modify     |
| `package.json`                        | 加 `./domain` 到 exports map + bump 1.1.0 → 1.2.0                                                             | Modify     |

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

| Test file                                     | 覆盖                                                                         |
| --------------------------------------------- | -------------------------------------------------------------------------- |
| `tests/domain/policy-guard.test.ts`           | 9 个 check 各自 happy + reject + edge case                                    |
| `tests/domain/task-orchestrator.test.ts`      | 6 action × happy + reject reasons + 并发场景                                   |
| `tests/domain/contact-writer.test.ts`         | NAME\_TRUST 6 level × winner / DNC sticky / forward-only                   |
| `tests/domain/timeline-writer.test.ts`        | Zod schema enforcement per eventType + idempotency                         |
| `tests/domain/task-orchestrator-sql.test.ts`  | raw SQL output 语义 match Drizzle 版                                          |
| `tests/domain/idempotency.test.ts`            | 4 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 union                                                                                           | 6 个 action 的 payload 类型           | normative-spec §3.1 line 157-163 |
| `CreateOpenPayload` / `CreateClosedPayload` / `ClosePayload` / `UpdatePayload` / `RecordProgressPayload` / `ReopenPayload` | 每个 action 各自的 payload             | spec §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 + `contactPhone`                                                                        | create\_open / create\_closed 用   | spec §3.1 line 196-199           |
| `TaskIdActionContext` = BaseActionContext                                                                                  | 4 个 taskId action 用               | spec §3.1 line 201-202           |
| `ActionContext` = union of two                                                                                             | spec §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                                      |
| - | -------------------------------------------------------------------------------------------------------- | --------------------------------------------- |
| 1 | `TaskAction` discriminated union — narrow by `action` field 后,`payload` 类型 narrowed 到对应 specific payload | 防有人改 union 丢失 discriminator                   |
| 2 | `ApplyResult` narrow by `status` —— `allow` 状态有 `statements`,`reject` 有 `reason` + `details`             | 防 union 写错                                    |
| 3 | `CreateActionContext` 要求 `contactPhone` non-optional,`TaskIdActionContext` 不要求                           | 防 S6 fix(spec §3.1 ActionContext 拆 2 type)被回退 |

不写其他 test —— TS 编译时强制。Test 形态用 `// @ts-expect-error` 验证 negative narrowing 不编译。

### Step 3: Implementation

```typescript
// 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:

```bash
cd /Users/maxwsy/workspace/callytics-common && bun test tests/domain/types.test.ts
```

Expected: 3 tests PASS。

### Step 5: Commit

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

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

| # | Input                                                      | Expected                    | 为什么                              |
| - | ---------------------------------------------------------- | --------------------------- | -------------------------------- |
| 1 | `no_answer`, currentDueAt 任意, override=undefined           | `now + 60min`               | spec line 245                    |
| 2 | `left_voicemail`, override=undefined                       | `now + 24h`                 | spec line 246                    |
| 3 | `text_sent`, override=undefined                            | `now + 48h`                 | spec line 247                    |
| 4 | `callback_requested`, currentDueAt=`X`, override=undefined | `currentDueAt`(不变)          | spec line 248 — staff 必须手动设      |
| 5 | `follow_up_scheduled`, override=undefined                  | `currentDueAt`(不变)          | spec line 249                    |
| 6 | `customer_considering`, override=undefined                 | `now + 72h`                 | spec line 250                    |
| 7 | 任何 progressType, override=`2026-12-25`                     | `2026-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.storeId                             | `storeId='STORE_A'` → allow                                                             | `storeId=''` → reject `store_mismatch`                                                                                                   |
| **`checkDnc(snapshot, action)`** spec §3.6 line 438                             | ContactSnapshot.contact.doNotContact    | snapshot.contact=null(新 contact)→ allow / doNotContact=false → allow                    | doNotContact=true × action ∈ \{create\_open, create\_closed, update, record\_progress, reopen} → reject `dnc`;action='close' → **allow** |
| **`checkTaskBelongsToContact(snapshot, taskId)`**                               | ContactSnapshot.pendingTasks Set lookup | taskId ∈ snapshot.pendingTasks → allow                                                  | taskId ∉ snapshot.pendingTasks → reject `hallucinated_task_id`                                                                           |
| **`checkStateTransition(action, currentStatus)`**                               | task.status from snapshot               | close 对 open → allow;reopen 对 closed → allow                                            | close 对 closed → reject `task_not_open` / reopen 对 open → reject `task_not_closed` / update 对 closed → reject `task_not_open`            |
| **`checkCreateClosedHasNoOpenTask(snapshot, typeCategory)`** spec §3.6 line 442 | ContactSnapshot.pendingTasks Set        | 无同 typeCategory open task → allow                                                       | 已有 → reject `create_closed_with_open_task`                                                                                               |
| **`checkTypeCategoryAllowed(ctx, contact, typeCategory)`** spec §3.6 line 443   | ContactLifecycle                        | actor.type='staff'/'system' → skip allow / 'ai\_agent' + typeCategory ∈ allowed → allow | 'contact\_analysis' + typeCategory ∉ allowed → reject `invalid_type_category`                                                            |
| **`checkCloseNoteRequired(payload)`**                                           | payload                                 | `closeResult='other'` + closeNote 非空 → allow                                            | `closeResult='other'` + closeNote=''/undefined → reject `close_note_required`                                                            |
| **`checkConfidence(confidence)`**                                               | proposal confidence                     | undefined / ≥ 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):

| # | Scenario                                                                    | Expected                                                                    |
| - | --------------------------------------------------------------------------- | --------------------------------------------------------------------------- |
| 1 | 新 contact(尚未 INSERT)→ `loadSnapshot` 返回 `{contact: null, pendingTasks: []}` | 不 throw,后续 check 能 handle null contact                                      |
| 2 | 已存在 contact + 2 open task(过渡期 1 个 `'pending'` + 1 个 `'open'`)→ 全 cover      | snapshot.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):

| Source | Input                                                                                               | Expected 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`

```typescript
// 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`

```typescript
// 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`:

```typescript
// 真实 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 一致)

```typescript
// 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)

```typescript
// 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/decision | DB call/batch (20 decisions)      |
| ----------------------------- | ---------------- | --------------------------------- |
| Plan 02 旧设计(per-check SELECT) | 4-5              | **80-100**                        |
| **新设计(snapshot pattern)**     | 0(snapshot 已加载)  | **2**(1 contact + 1 pendingTasks) |

### Step 5: Implementation — `idempotency.ts`

```typescript
// 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:

```bash
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

```bash
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:**&#x6CBF;用现有 `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:

| Action            | Happy path scenarios                                                                                                                                                                                                  | Reject scenarios                                                                                                                                                                                                                                                                  |
| ----------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| `create_open`     | 1.staff manual create → allow + INSERT statement / 2.AI 提案 confidence=0.8 typeCategory ∈ allowed → allow / 3.lead-processor `sourceType='lead'` skip type check → allow                                               | 1.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_closed`   | 1.staff inbound 当场成交 + sourceCallId 唯一 → allow + INSERT closed                                                                                                                                                        | 1.同上 4 个;2.已有 open task → create\_closed\_with\_open\_task;3.同 sourceCallId 已建过 closed → duplicate;4.closeResult='other' 无 closeNote → close\_note\_required                                                                                                                      |
| `close`           | 1.staff close open task → allow / 2.AI close + aiRunStartedAt 新于 task.updated\_at → allow                                                                                                                             | 1.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) |
| `update`          | 1.staff 改 dueAt → allow / 2.AI 改 priority → allow                                                                                                                                                                     | 1.task closed → task\_not\_open;2.AI stale → stale\_proposal;3.hallucinated taskId → hallucinated\_task\_id;4.DNC contact → dnc(update 在 DNC allowed 列外)                                                                                                                          |
| `record_progress` | 1.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                                                                                                                                             |
| `reopen`          | 1.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 验)**:

| #  | Invariant                                                                   | Test setup                                                                                                                                |
| -- | --------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------- |
| C1 | `record_progress` CTE 防 attempt\_count 翻倍                                   | seed task w/ attempt=2 + seed 1 progress event;再调 `record_progress` with **same** idempotency\_key → expect attempt\_count 仍 = 2(不+1)     |
| C2 | `close` SQL `WHERE status='open' AND updated_at <= $aiRunStartedAt` 防 stale | seed task w/ updated\_at=NOW;模拟 AI aiRunStartedAt=NOW-1h;调 close → expect 0 rows updated → reject `stale_proposal`                        |
| C3 | `create_open` partial unique 防 duplicate                                    | seed 1 open task (phone=A, store=B, typeCategory=C);调 `create_open` 同 (A,B,C) → expect SQL 报 partial unique conflict → reject `duplicate` |
| C4 | `create_closed` partial unique 防同 sourceCallId 重复建                          | seed 1 closed task w/ source\_call\_id='CALL1';再调 create\_closed 同 sourceCallId → expect duplicate                                        |
| C5 | `closeAllOpenForContact()` 关多个 open task + 写多条 timeline 原子                  | seed 3 open tasks for contact;调 closeAllOpen → expect 3 UPDATE + 3 timeline INSERT 全成功                                                    |

**`closeAllOpenForContact()` happy + edge**:

| # | Scenario                                                                                                           | Expected                                                                        |
| - | ------------------------------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------- |
| 1 | contact 有 3 open task → bulk close 全部                                                                              | 返回 closedTaskIds 含 3 个 / 生成 3 个 timeline event(每 task 1 个 task.status\_changed) |
| 2 | contact 无 open task → no-op                                                                                        | closedTaskIds=\[] / 不生成 SQL statement                                           |
| 3 | DNC trigger 用例:`closeResult='do_not_contact'` + `closeNote='Auto-closed by DNC hard stop'` + `actor.type='system'` | 全部 closeType='auto\_closed'                                                     |

### Step 2: Implementation — `applyTaskAction` dispatch

```typescript
// 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)。

```typescript
// 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)

```typescript
// 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`

```typescript
// 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

```typescript
// 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

```bash
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

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

| # | Scenario                                                                          | Expected                                     |
| - | --------------------------------------------------------------------------------- | -------------------------------------------- |
| 1 | 新 contact(不存在)→ INSERT 完整 row,first\_name\_trust\_score = 输入 trustScore           | row 拿回 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 行为)                        |
| 5 | franchiseId / accountId NULL → UPSERT 失败(NOT NULL violation in DB)                | caller 检查必传                                  |
| 6 | UPSERT 时 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`):

| # | Scenario                                                                                                                                                        | Expected                                        | Why                                                              |
| - | --------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------- | ---------------------------------------------------------------- |
| 1 | contact 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 永远 sticky                   | **silent bug** —— blind UPDATE SET 会丢 staff 审计 trail。这条 test 真守住 |
| 3 | updatedBy='staff' / 'ai' / 'system' 三种 false→true 都接受                                                                                                           | 三种都正确 stamp                                     | enum 完整                                                          |
| 4 | **没有 unsetDNC API**(sticky 不可反转)                                                                                                                                | TypeScript 编译时只有 setDNC,无 false 写入路径            | type-level guard                                                 |

**`touchActivity()` —— forward-only**:

| # | Scenario                                                                         | Expected                             |
| - | -------------------------------------------------------------------------------- | ------------------------------------ |
| 1 | contact 不存在 → 写 lastActivityAt = NOW()                                           | INSERT 含 lastActivityAt              |
| 2 | contact lastActivityAt=2026-01-01,activityAt=2026-06-01 → 更新到 2026-06-01         | row\.last\_activity\_at = 2026-06-01 |
| 3 | contact lastActivityAt=2026-06-01,activityAt=2026-01-01(更旧!)→ **不动**,GREATEST 保护 | row\.last\_activity\_at 仍 2026-06-01 |

### Step 2: Implementation

```typescript
// 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`**:

```typescript
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

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

| eventType                     | Required payload fields                                | Test 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**:

| # | Scenario                                                    | Expected                                      |
| - | ----------------------------------------------------------- | --------------------------------------------- |
| 1 | 合法 payload + eventType match → 返回 INSERT SQL                | SQL contains correct eventType + payload JSON |
| 2 | payload 不 match eventType Zod → throw at runtime            | error has eventType prefix                    |
| 3 | 同 idempotencyKey 重复调 → 返回 SQL with `ON CONFLICT DO NOTHING` | impl 已加 ON CONFLICT                           |

### Step 2: Implementation

```typescript
// 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**:

```typescript
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

```bash
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 验证**:

| # | Scenario                                                                 | Expected                                                                                    |
| - | ------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------- |
| 1 | `close` action → output `statements[]` 含 UPDATE SQL string + params      | string 含 `UPDATE tasks SET status='closed'`,params 数组 length 正确                             |
| 2 | `record_progress` action → output 含 CTE SQL                              | string 含 `WITH inserted AS (INSERT ...) UPDATE tasks ...`                                   |
| 3 | reject case → output `{status: 'reject', reason, details}` 跟 Drizzle 版相同 | reason 用相同 string literal                                                                   |
| 4 | allow case parity                                                        | Drizzle 版和 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

```typescript
// 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

```bash
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

```typescript
// 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`

```json
{
  "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

```bash
cd /Users/maxwsy/workspace/callytics-common && bun test && bun run build
```

Expected: all PASS + dist/ contains domain/ subdir。

### Step 4: Commit + Push + Open PR

```bash
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。
