import assert from "node:assert/strict"; import { randomUUID } from "node:crypto"; import { spawnSync } from "node:child_process"; import { fileURLToPath } from "node:url"; import { closeLocalPostgresDataPool, createLocalPostgresDataClient } from "../src/lib/db/local-postgres-client-core.ts"; import { resolveActiveSkillPackage } from "../src/lib/skill-package-registry.ts"; import { publicNextAction } from "../src/lib/rectification-agentic/core/rectification-decision.ts"; import { rawMinuteWeights, segmentsForMinutes, summarizeSegmentWeights, targetChartsForDomain } from "../src/lib/rectification-agentic/core/segment-summary.ts"; import { applyRectificationChoice, inferenceForPersistedAnswer, persistNextInterviewAfterChoice } from "../src/lib/rectification-agentic/v9/answer-choice.ts"; import { CHOICE_ACTION, outcomeIdForOption } from "../src/lib/rectification-agentic/v9/choice-action.ts"; import type { ChoiceKey } from "../src/lib/rectification-agentic/v9/choice-card.ts"; import { decideAfterInferenceChange, rectificationFollowupCatalog } from "../src/lib/rectification-agentic/v9/decision-from-dossier.ts"; import type { EvidenceKind } from "../src/lib/rectification-agentic/v9/evidence-model.ts"; import { resolveRectificationProductDomain } from "../src/lib/rectification-agentic/v9/product-domain.ts"; import { scoreAndPersistCurrentEvidence } from "../src/lib/rectification-agentic/v9/score-persist.ts"; import { confirmV9Evidence, loadV9CaseDossier, proposeV9Evidence } from "../src/lib/rectification-agentic/v9/tool-service.ts"; import { startPostgresFixture } from "../tests/helpers/postgres-fixture.ts"; import { segmentReplayOracleAnswer } from "./rectification-segment-oracle.ts"; // Fictional owned Case seeds only. Production turn/evidence/choice calls are // never replaced by fixture inserts, including when the original ACL refuses. type ReplayRow = { case_key: string; radius: number; oracle_truth_minute: string; snapshot: Record & { birth_date: string; reported_birth_time: string }; range: { start_time: string; end_time: string }; events: { event_kind: EvidenceKind; domain: string; date_start: string; date_end: string; precision: string; summary: string }[]; }; const chunks: Buffer[] = []; for await (const chunk of process.stdin) chunks.push(Buffer.from(chunk)); const input = JSON.parse(Buffer.concat(chunks).toString("utf8")) as { rows: ReplayRow[] }; const records: Record[] = []; const previousOrder = process.env.RECTIFICATION_SEGMENT_ORDER; const quote = (value: unknown) => `'${String(value).replaceAll("'", "''")}'`; const fixture = startPostgresFixture(); const url = fixture.connectionUrl("service_runtime", "service-runtime-test-password"); const authenticatedUrl = fixture.connectionUrl("app_runtime", "app-runtime-test-password"); let fatal: string | null = null; try { const migrated = spawnSync(process.execPath, [fileURLToPath(new URL("./db-migrate.mjs", import.meta.url))], { encoding: "utf8", env: { ...process.env, SCHEMA_DATABASE_URL: fixture.connectionUrl("schema_owner", "schema-owner-test-password") }, }); assert.equal(migrated.status, 0, migrated.stderr); const service = createLocalPostgresDataClient(url, null, "service_role"); const skill = resolveActiveSkillPackage("jyotish-birth-time-rectification"); for (const row of input.rows) { for (const enabled of [false, true]) { process.env.RECTIFICATION_SEGMENT_ORDER = enabled ? "on" : "off"; const userId = randomUUID(), caseId = randomUUID(), sessionId = randomUUID(), sourceSessionId = randomUUID(); const trace: Record = { case_key: row.case_key, radius: row.radius, ordering: enabled ? "on-experiment" : "off-default", opening: "explicit-fictional-case-seed-not-opening-success", asked: [], gates: [], stop: null, stage: "seed", metrics: null }; const ledger = () => JSON.parse(fixture.psql(`select json_build_object( 'turns',(select count(*) from public.agentic_rectification_turns where case_id=${quote(caseId)}), 'evidence',(select count(*) from public.agentic_rectification_evidence where case_id=${quote(caseId)}), 'results',(select count(*) from public.agentic_rectification_results where case_id=${quote(caseId)}), 'actions',(select count(*) from public.agentic_rectification_choice_actions where case_id=${quote(caseId)}), 'transitions',(select count(*) from public.agentic_rectification_inference_transitions where case_id=${quote(caseId)}))`)); try { fixture.psqlAs("identity_runtime", "identity-runtime-test-password", `insert into identity.users(id,name,email,email_verified) values(${quote(userId)},'Fictional Public Replay Owner',${quote(`${userId}@example.invalid`)},true)`); const snapshot = row.snapshot; fixture.psql(`update public.profiles set birth_date=${quote(snapshot.birth_date)},reported_birth_time=${quote(snapshot.reported_birth_time)}, birth_time_source='approximate',birth_time_status='reported',uncertainty_before_minutes=${row.radius},uncertainty_after_minutes=${row.radius}, latitude=${Number(snapshot.latitude)},longitude=${Number(snapshot.longitude)},timezone_offset=${Number(snapshot.timezone_offset)}, timezone_id=${snapshot.timezone_id ? quote(snapshot.timezone_id) : "null"} where id=${quote(userId)}; insert into public.chat_sessions(id,user_id,title,theme,session_type,messages) values (${quote(sourceSessionId)},${quote(userId)},'Fictional replay source','general','consultation','[]'), (${quote(sessionId)},${quote(userId)},'Fictional persisted replay','general','birth_time_rectification','[]')`); // Match POST /cases/open: owned source reads use the authenticated // session client; the accounting client is reserved for Case RPCs. const authenticated = createLocalPostgresDataClient(authenticatedUrl, { id: userId, email: null }, "authenticated"); const domain = await resolveRectificationProductDomain( authenticated as unknown as Parameters[0], userId, { intent: "homepage", sourceSessionId, requestId: randomUUID() }, ); fixture.psql(`insert into public.agentic_rectification_cases(id,user_id,session_id,status,skill_name,skill_version,skill_sha256,skill_source_commit, baseline_profile_fingerprint,baseline_birth_snapshot,candidate_range,rectification_domain) values (${quote(caseId)},${quote(userId)},${quote(sessionId)},'collecting_evidence',${quote(skill.name)},${quote(skill.version)},${quote(skill.sha256)}, ${skill.sourceCommit ? quote(skill.sourceCommit) : "null"},${quote("a".repeat(64))},${quote(JSON.stringify(snapshot))},${quote(JSON.stringify(row.range))},${quote(domain)})`); const seeded = await loadV9CaseDossier(service, userId, caseId); assert.equal(seeded.case.rectificationDomain, domain); trace.domain = seeded.case.rectificationDomain; trace.targets = targetChartsForDomain(seeded.case.rectificationDomain); trace.stage = "append_turn_v10"; trace.ledger_before = ledger(); // Same V10 request-idempotent overload as agent-run-prepare.ts; the V9 overload is revoked from service_role. const appended = await service.rpc("append_agentic_rectification_turn", { p_user_id: userId, p_case_id: caseId, p_user_message: row.events.map(event => event.summary).join("; "), p_assistant_message: null, p_model_name: "public-benchmark-oracle-not-provider", p_model_version: null, p_status: "pending", p_request_id: randomUUID() }); assert.equal(appended.error, null, JSON.stringify(appended.error)); const turn = { turnId: (appended.data as { turn_id: string }).turn_id }; trace.stage = "evidence"; for (const event of row.events) { const proposed = await proposeV9Evidence(service, userId, caseId, { sourceTurnId: turn.turnId, quote: event.summary, subject: event.domain === "family" ? "family" : "self", eventKind: event.event_kind, domain: event.domain, occurredFrom: event.date_start, occurredTo: event.date_end, datePrecision: event.precision, summary: event.summary, }); assert.equal(proposed.outcome, "accepted", proposed.errorCode ?? "evidence rejected"); assert.ok(proposed.evidenceId); await confirmV9Evidence(service, userId, caseId, proposed.evidenceId); } trace.stage = "scoreAndPersistCurrentEvidence"; await scoreAndPersistCurrentEvidence({ accounting: service, userId, caseId }); for (let round = 0; round < 6; round++) { let dossier = await loadV9CaseDossier(service, userId, caseId); const state = inferenceForPersistedAnswer(dossier.latestResult?.decisionReceipt); const decision = decideAfterInferenceChange({ dossier, state, userStopped: false, birthDate: snapshot.birth_date }); const catalog = rectificationFollowupCatalog(dossier.latestResult, dossier.evidence); (trace.gates as unknown[]).push({ round, session_outcome: decision.sessionOutcome, stop_reason: decision.stopReason ?? null, next_action: publicNextAction(decision), revision: state?.revision ?? null, eligible_event_keys: catalog.eventProbes.map(probe => probe.semantic_key).sort(), contrast_keys: catalog.contrastPacket?.probes.map(probe => probe.semanticKey).sort() ?? [] }); trace.stage = "persistNextInterviewAfterChoice"; const interview = await persistNextInterviewAfterChoice({ accounting: service, userId, caseId, dossier, decisionState: state, decision, nextAction: publicNextAction(decision), birthDate: snapshot.birth_date, askedTurnId: turn.turnId }); if (!interview.choiceReady || interview.terminalNote) { trace.stop = interview.terminalNote ? "production-terminal" : "production-no-renderable-choice"; break; } dossier = await loadV9CaseDossier(service, userId, caseId); const focus = dossier.conversationSummary.activeFocus; assert.ok(focus); const ownedState = inferenceForPersistedAnswer(dossier.latestResult?.decisionReceipt); const schema = focus.expectedAnswerSchema; const probe = ownedState?.probes.find(probe => probe.id === schema.probe_id); if (!probe) { trace.stop = "oracle-missing-for-production-focus"; (trace.gates as unknown[]).push({ intent: focus.intent, source: interview.followup?.source ?? null, choice_kind: schema.choice_kind ?? null, scoring: schema.scoring !== false }); break; } const answerClass = segmentReplayOracleAnswer(probe, row.oracle_truth_minute, ownedState?.transitions); const choice = schema.choice as { options?: { key: ChoiceKey }[] } | undefined; const option = choice?.options?.find(option => outcomeIdForOption(option.key, schema) === answerClass); if (!option) { trace.stop = "oracle-no-renderable-class-option"; break; } const before = ledger(); trace.stage = "applyRectificationChoice"; const applied = await applyRectificationChoice(service, { userId, caseId, sessionId, actionId: randomUUID(), action: CHOICE_ACTION, focusId: focus.id, optionId: option.key, expectedRevision: ownedState!.revision, deferFollowup: true }); assert.equal(applied.applied, true); const reloaded = await loadV9CaseDossier(service, userId, caseId); const afterState = inferenceForPersistedAnswer(reloaded.latestResult?.decisionReceipt); (trace.asked as unknown[]).push({ round, semantic_key: probe.semantic_key, optionId: option.key, answer_class: outcomeIdForOption(option.key, schema), source: probe.source, choice_kind: schema.choice_kind, scoring: schema.scoring !== false, ledger_before: before, ledger_after: ledger(), revision_before: ownedState!.revision, revision_after: afterState?.revision ?? null }); assert.equal(afterState?.revision, ownedState!.revision + 1); if (round === 5) trace.stop = "six-question-budget"; } const final = await loadV9CaseDossier(service, userId, caseId); const state = inferenceForPersistedAnswer(final.latestResult?.decisionReceipt); if (state?.segment_minutes?.length) { const offsets = Object.fromEntries(state.segment_minutes.map(minute => [minute.time, minute.offset])); const truthOffset = offsets[row.oracle_truth_minute]; assert.notEqual(truthOffset, undefined); const weights = rawMinuteWeights(state.candidates.map(candidate => ({ time: candidate.time, score: candidate.raw_posterior_score ?? 0, cluster_times: candidate.cluster_times, eliminated: candidate.raw_eliminated }))); trace.metrics = Object.fromEntries((["D9", "D10"] as const).map(chart => [chart, summarizeSegmentWeights(weights, segmentsForMinutes(state.segment_minutes!, [chart]), offsets, truthOffset)])); } trace.stage = "completed"; trace.ledger_after = ledger(); } catch (error) { trace.stop = "blocked"; trace.error = error instanceof Error ? error.message : "unknown_error"; trace.ledger_after = ledger(); } records.push(trace); } } } catch (error) { fatal = error instanceof Error ? error.message : "runner_failure"; } finally { if (previousOrder === undefined) delete process.env.RECTIFICATION_SEGMENT_ORDER; else process.env.RECTIFICATION_SEGMENT_ORDER = previousOrder; await closeLocalPostgresDataPool(url); await closeLocalPostgresDataPool(authenticatedUrl); fixture.stop(); } const aggregate = Object.fromEntries([...new Set(input.rows.map(row => row.radius))].flatMap(radius => ["off-default", "on-experiment"].map(ordering => { const rows = records.filter(row => row.radius === radius && row.ordering === ordering); const measured = rows.filter(row => row.metrics !== null); const metrics = Object.fromEntries((["D9", "D10"] as const).map(chart => { const values = measured.map(row => (row.metrics as Record)[chart]!); return [chart, { measured: values.length, truth_kept: values.filter(value => value.truth_retained).length, top_hit: values.filter(value => value.top_is_truth).length }]; })); return [`${radius}:${ordering}`, { attempted: rows.length, measured: measured.length, blocked: rows.filter(row => row.stop === "blocked").length, asked: rows.reduce((total, row) => total + (row.asked as unknown[]).length, 0), completed_six: rows.filter(row => (row.asked as unknown[]).length === 6).length, metrics, redline_status: measured.length === 77 ? "requires-independent-cell-comparison" : "not-evaluable" }]; }))); const blocked = fatal !== null || records.some(row => row.stop === "blocked" || row.metrics === null); process.stdout.write(JSON.stringify({ scope: "explicit-seeded persisted replay; off-default vs on-experiment; no opening success or provider claim", fatal, records, aggregate, blocked, distinct_cases: new Set(input.rows.map(row => row.case_key)).size, requested_traces: input.rows.length * 2, completed_six: records.filter(row => (row.asked as unknown[]).length === 6).length }) + "\n"); process.exitCode = blocked ? 1 : 0;