import assert from "node:assert/strict"; import { readFileSync } from "node:fs"; import test from "node:test"; import { RECTIFICATION_TELEMETRY_BIRTH_TIME_SOURCES, RECTIFICATION_TELEMETRY_FIELDS, RECTIFICATION_TELEMETRY_LIVE_WINDOW_MS, RECTIFICATION_TELEMETRY_STOP_REASONS, buildRectificationTelemetryRow, rectificationTelemetryEligible, rectificationTelemetryRpcArgs, resetRectificationTelemetryForTests, scheduleRectificationTelemetry, telemetryBirthTimeSource, telemetryQuestionKind, telemetryStopReason, type RectificationTelemetryInput, } from "../src/lib/rectification-agentic/v9/telemetry.ts"; import { budgetExhausted } from "../src/lib/rectification-agentic/core/rectification-decision.ts"; import { DEFAULT_MAX_DISCRIMINATION_ROUNDS } from "../src/lib/rectification-agentic/core/types.ts"; const USER_ID = "11111111-1111-4111-8111-111111111111"; const CASE_ID = "22222222-2222-4222-8222-222222222222"; const NOW = Date.parse("2026-09-26T08:00:00.000Z"); // Fictional, deliberately identifying-looking content: none of it may reach the row. const FICTIONAL_NAME = "测试虚构人甲"; const FICTIONAL_BIRTH_DATE = "1990-05-12"; const FICTIONAL_PLACE = "虚构市示例区"; const FICTIONAL_QUOTE = "2015 年 9 月我在虚构市示例区换了工作"; function focus( questionId: string, status: string, minute: number, schema: Record = {}, ) { const at = new Date(NOW - (60 - minute) * 60_000).toISOString(); return { questionId, expectedAnswerSchema: schema, status, askedAt: at, resolvedAt: status === "active" ? null : at }; } function input(overrides: Partial = {}): RectificationTelemetryInput { return { dossier: { case: { candidateRange: { start_time: "04:30", end_time: "05:30" }, stage: "minute", birthTimeSource: "approximate", skillVersion: "10.0.30", lastActivityAt: new Date(NOW - 30_000).toISOString(), }, turns: [ { createdAt: new Date(NOW - 45 * 60_000).toISOString() }, { createdAt: new Date(NOW - 10 * 60_000).toISOString() }, ], evidence: [ { status: "confirmed", summary: FICTIONAL_QUOTE, domain: "career" }, { status: "draft", summary: `${FICTIONAL_NAME} ${FICTIONAL_BIRTH_DATE}`, domain: "family" }, { status: "superseded", summary: FICTIONAL_PLACE, domain: "relocation" }, { status: "rejected", summary: "不算", domain: "other" }, ] as unknown as RectificationTelemetryInput["dossier"]["evidence"], latestResult: { algorithmVersion: "rectification-v5-matrix-scoring-9", policyVersion: "rectification-policy-v1", decisionReceipt: null, createdAt: new Date(NOW - 5 * 60_000).toISOString(), }, }, focuses: [ focus("collect:targeted:career", "resolved", 1), focus("collect:targeted:relationship", "declined", 2), focus("collect:guided:window:2015-09", "resolved", 3), focus("probe:career.2015.change", "resolved", 4, { probe_year: 2015, probe_id: "p1" }), focus("probe:varga.d9:style", "resolved", 5, { choice_kind: "varga_style" }), focus("collect:invite:more", "resolved", 6), focus("block:morning", "resolved", 7, { choice_kind: "block_choice" }), focus("collect:targeted:finance", "superseded", 8), ], decision: { sessionOutcome: "completed_with_range", stopReason: "probe_pool_exhausted", precisionGateMet: false, credibleRange: ["04:50", "05:06"], separation: { ranked: [{}, {}, {}, {}] }, }, rangeDelivery: { range: ["04:52", "05:04"], columns: [{ probability_percent: 41 }, { probability_percent: 34 }, { probability_percent: 25 }], }, now: NOW, ...overrides, }; } test("telemetry row keys are exactly the allowlist (new fields must change this test)", () => { const row = buildRectificationTelemetryRow(input()); assert.deepEqual(Object.keys(row).sort(), [...RECTIFICATION_TELEMETRY_FIELDS].sort()); assert.deepEqual([...RECTIFICATION_TELEMETRY_FIELDS].sort(), [ "algorithm_version", "birth_time_source", "candidate_count", "duration_seconds", "experiences_added", "policy_version", "precision_gate_met", "questions_dated_probe", "questions_guided", "questions_open", "questions_personality", "questions_targeted", "questions_total", "range_width_minutes", "skill_version", "stop_reason", "top_two_gap_points", "window_radius_minutes", ]); }); test("telemetry insert payload is the allowlist plus the two ids used only for ownership and dedupe", () => { const args = rectificationTelemetryRpcArgs({ userId: USER_ID, caseId: CASE_ID }, buildRectificationTelemetryRow(input())); assert.deepEqual( Object.keys(args).sort(), ["p_case_id", "p_user_id", ...RECTIFICATION_TELEMETRY_FIELDS.map((field) => `p_${field}`)].sort(), ); }); test("telemetry row holds only numbers, booleans, closed enums and version ids", () => { const row = buildRectificationTelemetryRow(input()); const serialized = JSON.stringify(row); for (const forbidden of [FICTIONAL_NAME, FICTIONAL_BIRTH_DATE, FICTIONAL_PLACE, FICTIONAL_QUOTE, USER_ID, CASE_ID, "04:52", "04:30"]) { assert.equal(serialized.includes(forbidden), false, `row must not contain ${forbidden}`); } const textFields = new Set(["birth_time_source", "stop_reason", "algorithm_version", "policy_version", "skill_version"]); for (const [key, value] of Object.entries(row)) { if (value === null || typeof value === "boolean") continue; if (typeof value === "number") { assert.ok(Number.isInteger(value) && value >= 0, `${key} must be a non-negative integer`); continue; } assert.ok(textFields.has(key), `${key} must not be text`); assert.match(String(value), /^[A-Za-z0-9][A-Za-z0-9._:+-]{0,63}$/); } assert.ok((RECTIFICATION_TELEMETRY_BIRTH_TIME_SOURCES as readonly string[]).includes(row.birth_time_source)); assert.ok((RECTIFICATION_TELEMETRY_STOP_REASONS as readonly string[]).includes(row.stop_reason)); }); test("telemetry row values reflect the delivered card and the new fewer-probes flow", () => { const row = buildRectificationTelemetryRow(input()); assert.equal(row.window_radius_minutes, 30); assert.equal(row.birth_time_source, "approximate"); // superseded focuses were replaced before the user answered them assert.equal(row.questions_total, 7); assert.equal(row.questions_targeted, 2); assert.equal(row.questions_guided, 1); assert.equal(row.questions_dated_probe, 1); assert.equal(row.questions_personality, 1); assert.equal(row.questions_open, 1); assert.equal(row.experiences_added, 2); // the card range wins over the decision range assert.equal(row.range_width_minutes, 12); assert.equal(row.candidate_count, 4); // 41 − 34 ≥ 5: this card showed percentages (RANGE_DELIVERY_PERCENT_MIN_GAP) assert.equal(row.top_two_gap_points, 7); assert.equal(row.stop_reason, "pool_exhausted"); assert.equal(row.precision_gate_met, false); assert.equal(row.duration_seconds, 45 * 60); assert.equal(row.algorithm_version, "rectification-v5-matrix-scoring-9"); assert.equal(row.policy_version, "rectification-policy-v1"); assert.equal(row.skill_version, "10.0.30"); }); test("telemetry drops versions that are not plain version ids and clamps odd values", () => { const base = input(); const row = buildRectificationTelemetryRow({ ...base, dossier: { ...base.dossier, case: { ...base.dossier.case, skillVersion: "10.0 我的版本", candidateRange: null }, turns: [], latestResult: { algorithmVersion: "x".repeat(80), policyVersion: " ", decisionReceipt: null }, }, rangeDelivery: { range: null, columns: [{ probability_percent: 100 }] }, decision: { ...base.decision, credibleRange: null, precisionGateMet: null }, }); assert.equal(row.skill_version, null); assert.equal(row.algorithm_version, null); assert.equal(row.policy_version, null); assert.equal(row.window_radius_minutes, null); assert.equal(row.range_width_minutes, null); assert.equal(row.top_two_gap_points, null); assert.equal(row.precision_gate_met, null); assert.equal(row.duration_seconds, 0); }); test("question kinds come from server-built focus ids and schemas", () => { assert.equal(telemetryQuestionKind({ questionId: "collect:targeted:career", expectedAnswerSchema: {} }), "targeted"); assert.equal(telemetryQuestionKind({ questionId: "collect:guided:window:2015-09", expectedAnswerSchema: {} }), "guided"); assert.equal(telemetryQuestionKind({ questionId: "probe:x", expectedAnswerSchema: { probe_year: 2019 } }), "dated_probe"); assert.equal(telemetryQuestionKind({ questionId: "d9_relationship:relationship_style", expectedAnswerSchema: { choice_kind: "varga_style" } }), "personality"); assert.equal(telemetryQuestionKind({ questionId: "x", expectedAnswerSchema: { tie_break_round: true } }), "personality"); assert.equal(telemetryQuestionKind({ questionId: "probe:nakshatra.trait", expectedAnswerSchema: {} }), "personality"); assert.equal(telemetryQuestionKind({ questionId: "collect:relationship:collect_method_evidence", expectedAnswerSchema: {} }), "open"); assert.equal(telemetryQuestionKind({ questionId: "collect:other:more", expectedAnswerSchema: {} }), "open"); assert.equal(telemetryQuestionKind({ questionId: "window_widen:widen_window", expectedAnswerSchema: { choice_kind: "widen_window" } }), "other"); }); test("birth time source collapses to the four categories", () => { assert.equal(telemetryBirthTimeSource({ birthTimeSource: "hospital_record" }), "hospital_record"); assert.equal(telemetryBirthTimeSource({ birthTimeSource: "period_only" }), "period_only"); assert.equal(telemetryBirthTimeSource({ birthTimeSource: "unknown" }), "unknown"); assert.equal(telemetryBirthTimeSource({ birthTimeSource: "approximate", stage: "block_scan" }), "unknown"); assert.equal(telemetryBirthTimeSource({ birthTimeSource: "family_exact" }), "approximate"); assert.equal(telemetryBirthTimeSource({ birthTimeSource: null }), "approximate"); }); test("stop reason: user stop, converged, round cap, user ran out, pool exhausted", () => { const base = { sessionOutcome: "completed_with_range", precisionGateMet: false, stopReason: "probe_pool_exhausted", roundCapReached: false, lastCollectDeclined: false, }; assert.equal(telemetryStopReason({ ...base, sessionOutcome: "provisional_range_user_stopped", precisionGateMet: true }), "user_stopped"); assert.equal(telemetryStopReason({ ...base, precisionGateMet: true, roundCapReached: true }), "converged"); assert.equal(telemetryStopReason({ ...base, roundCapReached: true, lastCollectDeclined: true }), "round_cap"); assert.equal(telemetryStopReason({ ...base, stopReason: "user_uncertainty_too_high" }), "user_no_more"); assert.equal(telemetryStopReason({ ...base, lastCollectDeclined: true }), "user_no_more"); assert.equal(telemetryStopReason(base), "pool_exhausted"); assert.equal(telemetryStopReason({ ...base, stopReason: "tied_first" }), "pool_exhausted"); // round_cap uses the decision layer's own fuse, not a second copy of it assert.equal(budgetExhausted({ inferenceRounds: DEFAULT_MAX_DISCRIMINATION_ROUNDS }), true); assert.equal(budgetExhausted({ inferenceRounds: DEFAULT_MAX_DISCRIMINATION_ROUNDS - 1 }), false); }); test("a declined last collect question marks the row user_no_more", () => { const base = input(); const row = buildRectificationTelemetryRow({ ...base, focuses: [...base.focuses, focus("collect:targeted:health", "declined", 30)], }); assert.equal(row.stop_reason, "user_no_more"); }); test("only a live range-card delivery is eligible", () => { assert.equal(rectificationTelemetryEligible(input()), true); assert.equal(rectificationTelemetryEligible(input({ decision: { ...input().decision, sessionOutcome: "collect_evidence" }, })), false); assert.equal(rectificationTelemetryEligible(input({ rangeDelivery: null })), false); assert.equal(rectificationTelemetryEligible(input({ rangeDelivery: { range: ["04:52", "05:04"], columns: [] } })), false); const stale = input(); assert.equal(rectificationTelemetryEligible({ ...stale, dossier: { ...stale.dossier, case: { ...stale.dossier.case, lastActivityAt: new Date(NOW - RECTIFICATION_TELEMETRY_LIVE_WINDOW_MS - 1).toISOString() }, }, }), false); assert.equal(rectificationTelemetryEligible(input({ decision: { ...input().decision, sessionOutcome: "provisional_range_user_stopped" }, })), true); }); function recordingClient(result: { data: unknown; error: { message?: string; code?: string } | null } | Error) { const calls: Array<{ fn: string; args: Record }> = []; return { calls, client: { rpc(fn: string, args: Record) { calls.push({ fn, args }); if (result instanceof Error) throw result; return Promise.resolve(result); }, }, }; } function captureWarnings() { const lines: string[] = []; const original = console.warn; console.warn = (...parts: unknown[]) => { lines.push(parts.map(String).join(" ")); }; return { lines, restore: () => { console.warn = original; } }; } test("schedule writes once per Case, after the response, with the allowlisted payload", async () => { resetRectificationTelemetryForTests(); const recorder = recordingClient({ data: true, error: null }); const pending = scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => input(), }); assert.equal(recorder.calls.length, 0, "the write must not run on the request's own tick"); await pending; assert.equal(recorder.calls.length, 1); assert.equal(recorder.calls[0]!.fn, "record_rectification_telemetry"); assert.deepEqual( Object.keys(recorder.calls[0]!.args).sort(), ["p_case_id", "p_user_id", ...RECTIFICATION_TELEMETRY_FIELDS.map((field) => `p_${field}`)].sort(), ); await scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => input(), }); assert.equal(recorder.calls.length, 1, "the same process does not re-send a Case"); }); test("schedule skips read-only, non-delivery and unavailable builds without writing", async () => { resetRectificationTelemetryForTests(); const recorder = recordingClient({ data: true, error: null }); await scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, readOnly: true, build: () => input() }); await scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => null }); await scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => input({ decision: { ...input().decision, sessionOutcome: "collect_evidence" } }), }); assert.equal(recorder.calls.length, 0); // a Case that was not eligible yet is still recorded once it delivers await scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => input() }); assert.equal(recorder.calls.length, 1); }); test("write failures never throw and log only a reason code, never the payload or ids", async () => { for (const failure of [ { data: null, error: { message: "rectification_telemetry_case_not_found", code: "P0002" } }, { data: null, error: { message: `insert failed for ${CASE_ID} ${FICTIONAL_NAME}` } }, new Error(`socket closed ${USER_ID}`), ]) { resetRectificationTelemetryForTests(); const recorder = recordingClient(failure); const warnings = captureWarnings(); try { await assert.doesNotReject(scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => input(), })); } finally { warnings.restore(); } assert.equal(warnings.lines.length, 1); const line = warnings.lines[0]!; assert.match(line, /^\[rectification-telemetry\] write skipped reason=[A-Za-z0-9_]+$/); for (const forbidden of [CASE_ID, USER_ID, FICTIONAL_NAME, "04:52"]) { assert.equal(line.includes(forbidden), false); } } }); test("a throwing build is swallowed on the request path", async () => { resetRectificationTelemetryForTests(); const recorder = recordingClient({ data: true, error: null }); const warnings = captureWarnings(); let pending: Promise | null = null; try { assert.doesNotThrow(() => { pending = scheduleRectificationTelemetry({ accounting: recorder.client, userId: USER_ID, caseId: CASE_ID, build: () => { throw new TypeError(`bad ${FICTIONAL_BIRTH_DATE}`); }, }); }); await pending; } finally { warnings.restore(); } assert.equal(recorder.calls.length, 0); assert.deepEqual(warnings.lines, ["[rectification-telemetry] write skipped reason=TypeError"]); }); const migration = readFileSync( new URL("../supabase/migrations/20260926010000_rectification_telemetry.sql", import.meta.url), "utf8", ); function tableColumns(sql: string, table: string): string[] { const start = sql.indexOf(`create table if not exists public.${table} (`); assert.ok(start >= 0, `${table} missing`); const body = sql.slice(sql.indexOf("(", start) + 1, sql.indexOf("\n);", start)); return body .split("\n") .map((line) => line.trim()) .filter((line) => /^[a-z_]+ (?:uuid|date|integer|text|boolean)\b/.test(line)) .map((line) => line.split(" ")[0]!); } test("the database row has exactly the allowlisted columns and no identity column", () => { assert.deepEqual( tableColumns(migration, "rectification_telemetry"), ["id", "recorded_week", ...RECTIFICATION_TELEMETRY_FIELDS], ); assert.deepEqual(tableColumns(migration, "rectification_telemetry_reported_cases"), ["case_id"]); const telemetryTable = migration.slice( migration.indexOf("create table if not exists public.rectification_telemetry ("), migration.indexOf("create index if not exists rectification_telemetry_week_idx"), ); assert.doesNotMatch(telemetryTable, /user_id|case_id|session_id|birth_date|birth_place|email|name text|created_at|timestamptz/); }); test("the write function takes exactly the allowlisted parameters", () => { const start = migration.indexOf("create or replace function public.record_rectification_telemetry("); const signature = migration.slice(start, migration.indexOf(")\nreturns boolean", start)); const params = [...signature.matchAll(/\n (p_[a-z_]+) /g)].map((match) => match[1]); assert.deepEqual(params, ["p_user_id", "p_case_id", ...RECTIFICATION_TELEMETRY_FIELDS.map((field) => `p_${field}`)]); }); test("migration is additive, RLS-guarded, server-write / admin-aggregate-read, 180-day retention", () => { assert.doesNotMatch(migration, /alter table public\.(?!rectification_telemetry)/); assert.doesNotMatch(migration, /drop (?:table|column|function)/i); assert.match(migration, /alter table public\.rectification_telemetry enable row level security/); assert.match(migration, /alter table public\.rectification_telemetry_reported_cases enable row level security/); assert.match(migration, /references public\.agentic_rectification_cases\(id\) on delete cascade/); assert.doesNotMatch(migration, /grant (?:select|insert|update|delete|all)[^;]*on table public\.rectification_telemetry/); assert.match(migration, /grant execute on function public\.record_rectification_telemetry\([\s\S]*?\) to service_role;/); assert.match(migration, /grant execute on function public\.rectification_telemetry_summary\(integer\) to admin_runtime;/); assert.doesNotMatch(migration, /rectification_telemetry_summary\(integer\) to (?:service_role|authenticated|anon|app_runtime)/); assert.equal((migration.match(/recorded_week < \(pg_catalog\.timezone\('UTC', pg_catalog\.now\(\)\)\)::date - 180/g) ?? []).length, 2); assert.match(migration, /b\.today - 180/); assert.match(migration, /least\(coalesce\(p_weeks, 12\), 26\)/); }); test("the case GET route schedules telemetry without awaiting it", () => { const route = readFileSync(new URL("../src/app/api/rectification/cases/[caseId]/route.ts", import.meta.url), "utf8"); assert.match(route, /onProjected: \(\{ decision, rangeDelivery, readOnly \}\) => \{\s*void scheduleRectificationTelemetry\(/); assert.doesNotMatch(route, /await scheduleRectificationTelemetry/); assert.match(route, /return NextResponse\.json\(await dossierResponseWithIdentity/); });