import type { DatabaseSync } from "node:sqlite"; import { insertExperienceCandidateInTransaction, listEligibleExperiences } from "../experience/storage.ts"; import { EXPERIENCE_AUTHORITIES, EXPERIENCE_KINDS, EXPERIENCE_SCOPE_KINDS, type ExperienceAuthority, type ExperienceHost, type ExperienceKind, type ExperienceScope } from "../experience/types.ts"; import { normalizeUserId } from "../storage/private-root.ts"; import { canonicalJson, checksumJson, sha256Hex } from "../storage/checksum.ts"; import { containsUnredactedSensitiveText, redactJson } from "../storage/redaction.ts"; import type { ValidatedObservationRecord } from "./observations.ts"; import { observationKey } from "./observations.ts"; import { consolidateProposalBatch, recordProposalReadCoverageInTransaction, recordZeroProposalReadCoverage, type ConsolidationResult, type Watermark } from "./commit.ts"; import type { ProposalBatch, ProposalSourceRef } from "./proposals.ts"; export type ModelProposalKind = ExperienceKind | "habit_candidate" | "correction_split"; export interface ModelOutputSourceRef extends ProposalSourceRef {} export interface HabitCandidateModelProposal { proposal_id: string; kind: "habit_candidate"; candidate_key: string; condition: string; behavior: string; polarity: 1 | -1; confidence_bp: number; source_refs: ModelOutputSourceRef[]; evidence_summary?: string; evidence_stage?: "collecting" | "reviewable"; ambiguous?: false; } export interface CorrectionSplitModelProposal { proposal_id: string; kind: "correction_split"; candidate_key: string; old_condition: string; old_behavior: string; new_condition: string; new_behavior: string; confidence_bp: number; source_refs: ModelOutputSourceRef[]; evidence_summary?: string; evidence_stage?: "collecting" | "reviewable"; ambiguous?: false; } export interface ExperienceCandidateModelProposal { proposal_id: string; kind: ExperienceKind; candidate_key: string; scope: ExperienceScope; authority: ExperienceAuthority; applicability: string; content: string; rationale?: string; exceptions: string[]; confidence_bp: number; source_refs: ModelOutputSourceRef[]; evidence_summary?: string; ambiguous?: false; } export type ValidatedModelProposal = ExperienceCandidateModelProposal | HabitCandidateModelProposal | CorrectionSplitModelProposal; export interface ValidatedModelOutputBatch { schema_version: 1; user_id: string; file_generation: string; batch_id: string; model: string; created_at: string; seq_start: number; seq_end: number; read_checksum: string; proposals: ValidatedModelProposal[]; checksum: string; } const MODEL_OUTPUT_KEYS = new Set(["schema_version", "user_id", "file_generation", "batch_id", "model", "created_at", "observations_read", "proposals"]); const OBSERVATIONS_READ_KEYS = new Set(["seq_start", "seq_end", "checksum"]); const HABIT_KEYS = new Set(["proposal_id", "kind", "candidate_key", "condition", "behavior", "polarity", "confidence_bp", "source_refs", "evidence_summary", "evidence_stage", "ambiguous"]); const CORRECTION_KEYS = new Set(["proposal_id", "kind", "candidate_key", "old_condition", "old_behavior", "new_condition", "new_behavior", "confidence_bp", "source_refs", "evidence_summary", "evidence_stage", "ambiguous"]); const EXPERIENCE_KEYS = new Set(["proposal_id", "kind", "candidate_key", "scope", "authority", "applicability", "content", "rationale", "exceptions", "confidence_bp", "source_refs", "evidence_summary", "ambiguous"]); const SCOPE_KEYS = new Set(["kind", "key"]); const EXPERIENCE_KIND_SET = new Set(EXPERIENCE_KINDS); const EXPERIENCE_SCOPE_SET = new Set(EXPERIENCE_SCOPE_KINDS); const EXPERIENCE_AUTHORITY_SET = new Set(EXPERIENCE_AUTHORITIES); const REF_KEYS = new Set(["file_generation", "seq", "checksum"]); function assertExactKeys(value: Record, allowed: Set, label: string): void { for (const key of Object.keys(value)) { if (!allowed.has(key)) throw new Error(`${label} has unsupported field: ${key}`); } } function assertSafeToken(value: unknown, label: string, max = 160): string { if (typeof value !== "string" || value.length < 1 || value.length > max || /[\x00-\x08\x0b\x0c\x0e-\x1f\x7f]/.test(value)) throw new Error(`Invalid ${label}`); return value; } function assertSafeText(value: unknown, label: string, max = 1000): string { const text = assertSafeToken(value, label, max); if (containsUnredactedSensitiveText(text)) throw new Error(`${label} contains unredacted sensitive text`); return text; } function assertGeneralizedHabitText(text: string, label: string): void { if (/\b(?:agent experience|pi-experiences|experience-consolidate)\b/i.test(text)) throw new Error(`${label} appears overfit to one project`); if (/\bv?\d+\.\d+\.\d+(?:[-+][A-Za-z0-9._-]+)?\b/.test(text)) throw new Error(`${label} appears overfit to one version`); if (/(^|[\s("'`])(?:~\/|\.\.?\/|\/[A-Za-z0-9._-])/.test(text)) throw new Error(`${label} appears overfit to one file path`); if (/\b[a-f0-9]{12,}\b/i.test(text)) throw new Error(`${label} appears overfit to one hash or screenshot`); } function assertGeneration(value: unknown): string { const generation = assertSafeToken(value, "file_generation", 80); if (!/^[A-Za-z0-9._-]+$/.test(generation)) throw new Error("Invalid file_generation"); return generation; } function assertChecksum(value: unknown, label = "checksum"): string { const checksum = assertSafeToken(value, label, 128); if (!/^[a-f0-9]{64}$/.test(checksum)) throw new Error(`Invalid ${label}`); return checksum; } function assertSeq(value: unknown, label: string): number { if (!Number.isInteger(value) || Number(value) < 1) throw new Error(`Invalid ${label}`); return Number(value); } function assertConfidence(value: unknown): number { if (!Number.isInteger(value) || Number(value) < 0 || Number(value) > 10000) throw new Error("Invalid confidence_bp"); return Number(value); } function validateSourceRef(value: unknown, expectedGeneration: string, seqStart: number, seqEnd: number): ModelOutputSourceRef { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Invalid model source ref"); const ref = value as Record; assertExactKeys(ref, REF_KEYS, "model source ref"); const fileGeneration = assertGeneration(ref.file_generation); if (fileGeneration !== expectedGeneration) throw new Error("Model source ref generation mismatch"); const seq = assertSeq(ref.seq, "model source seq"); if (seq < seqStart || seq > seqEnd) throw new Error("Model source ref outside read coverage"); return { file_generation: fileGeneration, seq, checksum: assertChecksum(ref.checksum, "source checksum") }; } function validateRefs(value: unknown, expectedGeneration: string, seqStart: number, seqEnd: number): ModelOutputSourceRef[] { if (!Array.isArray(value) || value.length < 1 || value.length > 20) throw new Error("Invalid model source_refs"); return value.map((ref) => validateSourceRef(ref, expectedGeneration, seqStart, seqEnd)); } export function validateModelOutputSourceRefs(output: ValidatedModelOutputBatch, observations: ValidatedObservationRecord[]): void { const byKey = new Map(observations.map((record) => [`${record.file_generation}:${record.seq}`, record])); for (const proposal of output.proposals) { for (const ref of proposal.source_refs) { const record = byKey.get(`${ref.file_generation}:${ref.seq}`); if (!record) throw new Error("Model source ref missing observation"); if (record.user_id !== output.user_id || record.checksum !== ref.checksum) throw new Error("Model source ref checksum mismatch"); } } } function validateProposal(value: unknown, seenIds: Set, generation: string, seqStart: number, seqEnd: number): ValidatedModelProposal { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Invalid model proposal"); const proposal = value as Record; if (proposal.ambiguous === true) throw new Error("Ambiguous model proposal"); if (proposal.ambiguous !== undefined && proposal.ambiguous !== false) throw new Error("Invalid ambiguous flag"); const kind = proposal.kind; if (typeof kind !== "string" || (!EXPERIENCE_KIND_SET.has(kind) && kind !== "habit_candidate" && kind !== "correction_split")) { throw new Error("Unsupported model proposal kind"); } assertExactKeys( proposal, EXPERIENCE_KIND_SET.has(kind) ? EXPERIENCE_KEYS : kind === "habit_candidate" ? HABIT_KEYS : CORRECTION_KEYS, "model proposal", ); const proposalId = assertSafeToken(proposal.proposal_id, "proposal_id"); if (seenIds.has(proposalId)) throw new Error("Duplicate model proposal_id"); seenIds.add(proposalId); const base = { proposal_id: proposalId, candidate_key: assertSafeToken(proposal.candidate_key, "candidate_key"), confidence_bp: assertConfidence(proposal.confidence_bp), source_refs: validateRefs(proposal.source_refs, generation, seqStart, seqEnd), ...(proposal.evidence_summary === undefined ? {} : { evidence_summary: assertSafeText(proposal.evidence_summary, "evidence_summary") }), ...(proposal.evidence_stage === undefined ? {} : { evidence_stage: proposal.evidence_stage === "collecting" || proposal.evidence_stage === "reviewable" ? proposal.evidence_stage : (() => { throw new Error("Invalid evidence_stage"); })() }), ...(proposal.ambiguous === undefined ? {} : { ambiguous: false as const }), }; if (EXPERIENCE_KIND_SET.has(kind)) { if (!proposal.scope || typeof proposal.scope !== "object" || Array.isArray(proposal.scope)) throw new Error("Invalid experience scope"); const scope = proposal.scope as Record; assertExactKeys(scope, SCOPE_KEYS, "experience scope"); if (typeof scope.kind !== "string" || !EXPERIENCE_SCOPE_SET.has(scope.kind)) throw new Error("Invalid experience scope kind"); const scopeKey = scope.key; if (scope.kind === "user" ? scopeKey !== undefined : typeof scopeKey !== "string" || scopeKey.length === 0 || scopeKey.length > 500) { throw new Error("Invalid experience scope key"); } if (typeof proposal.authority !== "string" || !EXPERIENCE_AUTHORITY_SET.has(proposal.authority)) throw new Error("Invalid experience authority"); const applicability = assertSafeText(proposal.applicability, "applicability", 8_000); const content = assertSafeText(proposal.content, "content", 8_000); const rationale = proposal.rationale === undefined ? undefined : assertSafeText(proposal.rationale, "rationale", 8_000); if (!Array.isArray(proposal.exceptions) || proposal.exceptions.length > 32) throw new Error("Invalid experience exceptions"); const exceptions = proposal.exceptions.map((exception, index) => assertSafeText(exception, `exceptions[${index}]`, 2_000)); if ((kind === "decision" || kind === "episode") && !rationale) throw new Error(`${kind} requires rationale`); if (kind === "habit") { assertGeneralizedHabitText(applicability, "applicability"); assertGeneralizedHabitText(content, "content"); } return { ...base, kind: kind as ExperienceKind, scope: scopeKey === undefined ? { kind: scope.kind } : { kind: scope.kind, key: scopeKey }, authority: proposal.authority as ExperienceAuthority, applicability, content, ...(rationale === undefined ? {} : { rationale }), exceptions, } as ExperienceCandidateModelProposal; } if (kind === "habit_candidate") { if (proposal.polarity !== 1 && proposal.polarity !== -1) throw new Error("Invalid model polarity"); const condition = assertSafeText(proposal.condition, "condition"); const behavior = assertSafeText(proposal.behavior, "behavior"); assertGeneralizedHabitText(condition, "condition"); assertGeneralizedHabitText(behavior, "behavior"); return { ...base, kind, condition, behavior, polarity: proposal.polarity }; } const oldCondition = assertSafeText(proposal.old_condition, "old_condition"); const oldBehavior = assertSafeText(proposal.old_behavior, "old_behavior"); const newCondition = assertSafeText(proposal.new_condition, "new_condition"); const newBehavior = assertSafeText(proposal.new_behavior, "new_behavior"); assertGeneralizedHabitText(oldCondition, "old_condition"); assertGeneralizedHabitText(oldBehavior, "old_behavior"); assertGeneralizedHabitText(newCondition, "new_condition"); assertGeneralizedHabitText(newBehavior, "new_behavior"); if (oldCondition === newCondition && oldBehavior === newBehavior) throw new Error("Invalid correction_split replacement"); return { ...base, kind, old_condition: oldCondition, old_behavior: oldBehavior, new_condition: newCondition, new_behavior: newBehavior }; } export function validateModelOutputBatch(value: unknown, expectedUserId?: string): ValidatedModelOutputBatch { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Invalid model output"); const batch = value as Record; assertExactKeys(batch, MODEL_OUTPUT_KEYS, "model output"); if (batch.schema_version !== 1) throw new Error("Unsupported model output schema_version"); const userId = normalizeUserId(assertSafeToken(batch.user_id, "user_id", 120)); if (expectedUserId !== undefined && userId !== normalizeUserId(expectedUserId)) throw new Error("Model output user_id mismatch"); const generation = assertGeneration(batch.file_generation); const createdAt = assertSafeToken(batch.created_at, "created_at", 80); if (Number.isNaN(Date.parse(createdAt))) throw new Error("Invalid model output created_at"); if (!batch.observations_read || typeof batch.observations_read !== "object" || Array.isArray(batch.observations_read)) throw new Error("Invalid observations_read"); const observationsRead = batch.observations_read as Record; assertExactKeys(observationsRead, OBSERVATIONS_READ_KEYS, "observations_read"); const seqStart = assertSeq(observationsRead.seq_start, "seq_start"); const seqEnd = assertSeq(observationsRead.seq_end, "seq_end"); if (seqEnd < seqStart) throw new Error("Invalid observations_read range"); const readChecksum = assertChecksum(observationsRead.checksum, "observations_read checksum"); if (!Array.isArray(batch.proposals) || batch.proposals.length > 200) throw new Error("Invalid model proposal list"); const seenIds = new Set(); const proposals = batch.proposals.map((proposal) => validateProposal(proposal, seenIds, generation, seqStart, seqEnd)); const normalized = { schema_version: 1 as const, user_id: userId, file_generation: generation, batch_id: assertSafeToken(batch.batch_id, "batch_id"), model: assertSafeToken(batch.model, "model", 120), created_at: createdAt, seq_start: seqStart, seq_end: seqEnd, read_checksum: readChecksum, proposals, }; return { ...normalized, checksum: checksumJson({ schema: "agent_experience_model_output_v1", batch: JSON.parse(canonicalJson(normalized)) }) }; } export function modelOutputToProposalBatch(batch: ValidatedModelOutputBatch): ProposalBatch { const proposals = batch.proposals.flatMap((proposal): ProposalBatch["proposals"] => { if (EXPERIENCE_KIND_SET.has(proposal.kind)) throw new Error("Typed experience proposal requires typed commit path"); if (proposal.kind === "habit_candidate") { return [{ proposal_id: proposal.proposal_id, kind: "habit_candidate", candidate_key: proposal.candidate_key, condition: proposal.condition, behavior: proposal.behavior, polarity: proposal.polarity, confidence_bp: proposal.confidence_bp, source_refs: proposal.source_refs, ...(proposal.evidence_summary === undefined ? {} : { evidence_summary: proposal.evidence_summary }), ...(proposal.evidence_stage === undefined ? {} : { evidence_stage: proposal.evidence_stage }), }]; } return [ { proposal_id: `${proposal.proposal_id}-old-negative`, kind: "habit_candidate", candidate_key: `${proposal.candidate_key}:old`, condition: proposal.old_condition, behavior: proposal.old_behavior, polarity: -1, confidence_bp: proposal.confidence_bp, source_refs: proposal.source_refs, ...(proposal.evidence_summary === undefined ? {} : { evidence_summary: proposal.evidence_summary }), correction_role: "old_negative", correction_group_id: proposal.proposal_id, ...(proposal.evidence_stage === undefined ? {} : { evidence_stage: proposal.evidence_stage }), }, { proposal_id: `${proposal.proposal_id}-new-positive`, kind: "habit_candidate", candidate_key: `${proposal.candidate_key}:new`, condition: proposal.new_condition, behavior: proposal.new_behavior, polarity: 1, confidence_bp: proposal.confidence_bp, source_refs: proposal.source_refs, ...(proposal.evidence_summary === undefined ? {} : { evidence_summary: proposal.evidence_summary }), correction_role: "replacement", correction_group_id: proposal.proposal_id, ...(proposal.evidence_stage === undefined ? {} : { evidence_stage: proposal.evidence_stage }), }, ]; }); return { schema_version: 1, user_id: batch.user_id, batch_id: batch.batch_id, created_at: batch.created_at, proposals }; } function stableId(prefix: string, value: unknown): string { return `${prefix}-${sha256Hex(canonicalJson(value)).slice(0, 40)}`; } function quarantineRowChecksum(row: { user_id: string; file_generation: string; seq_start: number; seq_end: number; reason: string; model: string; output_json: string; checksum: string; created_at: string }): string { return checksumJson({ table: "model_output_quarantine", row }); } function pendingReviewChecksum(row: { user_id: string; kind: string; status: string; payload_json: string }): string { return checksumJson({ table: "pending_review", row }); } export function insertPendingReview(db: any, input: { userId: string; kind: string; payload: unknown; createdAt: string }): { id: string; inserted: boolean; checksum: string } { const userId = normalizeUserId(input.userId); const payload = redactJson(input.payload ?? {}); const payloadJson = canonicalJson(payload); if (payloadJson.length > 24000) throw new Error("Pending review payload too large"); const checksum = pendingReviewChecksum({ user_id: userId, kind: input.kind, status: "open", payload_json: payloadJson }); const id = stableId("pending", { user_id: userId, kind: input.kind, checksum }); const existing = db.prepare("SELECT id, checksum FROM pending_review WHERE id = ?").get(id); if (existing) { if (existing.checksum !== checksum) throw new Error("Pending review stable id collision"); return { id, inserted: false, checksum }; } db.prepare("INSERT INTO pending_review (id, user_id, kind, status, payload_json, checksum, created_at, updated_at) VALUES (?, ?, ?, 'open', ?, ?, ?, ?)") .run(id, userId, input.kind, payloadJson, checksum, input.createdAt, input.createdAt); return { id, inserted: true, checksum }; } export function insertModelOutputQuarantine(db: any, input: { userId: string; fileGeneration: string; seqStart: number; seqEnd: number; reason: string; model: string; output: unknown; createdAt: string }): { id: string; inserted: boolean; checksum: string } { const userId = normalizeUserId(input.userId); if (!Number.isInteger(input.seqStart) || !Number.isInteger(input.seqEnd) || input.seqStart < 1 || input.seqEnd < input.seqStart) throw new Error("Invalid quarantine range"); const redacted = redactJson(input.output ?? {}); const outputJson = canonicalJson(redacted); if (outputJson.length > 24000) throw new Error("Quarantine output too large"); const checksum = checksumJson({ schema: "agent_experience_model_output_quarantine_v1", output: JSON.parse(outputJson) }); const id = stableId("quarantine", { user_id: userId, file_generation: input.fileGeneration, seq_start: input.seqStart, seq_end: input.seqEnd, reason: input.reason, checksum }); const rowChecksum = quarantineRowChecksum({ user_id: userId, file_generation: input.fileGeneration, seq_start: input.seqStart, seq_end: input.seqEnd, reason: input.reason, model: input.model, output_json: outputJson, checksum, created_at: input.createdAt }); const existing = db.prepare("SELECT id, checksum, row_checksum FROM model_output_quarantine WHERE id = ?").get(id); if (existing) { if (existing.checksum !== checksum || existing.row_checksum !== rowChecksum) throw new Error("Quarantine stable id collision"); return { id, inserted: false, checksum }; } db.prepare("INSERT INTO model_output_quarantine (id, user_id, file_generation, seq_start, seq_end, reason, model, output_json, checksum, created_at, row_checksum) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)") .run(id, userId, input.fileGeneration, input.seqStart, input.seqEnd, input.reason, input.model, outputJson, checksum, input.createdAt, rowChecksum); return { id, inserted: true, checksum }; } function normalizedIdentityText(value: string): string { return value.trim().replace(/\s+/g, " ").toLowerCase(); } function proposalIdentityForConflict(proposal: ValidatedModelProposal): string { if (EXPERIENCE_KIND_SET.has(proposal.kind)) { const typed = proposal as ExperienceCandidateModelProposal; return canonicalJson({ kind: typed.kind, scope: typed.scope, authority: typed.authority, applicability: normalizedIdentityText(typed.applicability), content: normalizedIdentityText(typed.content), }); } if (proposal.kind === "habit_candidate") return canonicalJson({ kind: proposal.kind, condition: normalizedIdentityText(proposal.condition), behavior: normalizedIdentityText(proposal.behavior), polarity: proposal.polarity }); return canonicalJson({ kind: proposal.kind, old_condition: normalizedIdentityText(proposal.old_condition), old_behavior: normalizedIdentityText(proposal.old_behavior), new_condition: normalizedIdentityText(proposal.new_condition), new_behavior: normalizedIdentityText(proposal.new_behavior) }); } function findCandidateKeyConflict(output: ValidatedModelOutputBatch): { candidate_key: string; identities: string[] } | null { const byKey = new Map>(); for (const proposal of output.proposals) { const set = byKey.get(proposal.candidate_key) || new Set(); set.add(proposalIdentityForConflict(proposal)); byKey.set(proposal.candidate_key, set); } for (const [candidate_key, identities] of byKey) { if (identities.size > 1) return { candidate_key, identities: [...identities].sort() }; } return null; } interface TypedExperienceConsolidationResult { user_id: string; file_generation: string; candidate_ids: string[]; evidence_ids: []; watermark_after: null; read_watermark_after: Watermark; audit_id: string; inserted: { candidates: number; evidence: 0; audit: 1; watermark: 0; read_watermark: 0 | 1 }; } function commitTypedExperienceProposals(input: { db: DatabaseSync; userId: string; output: ValidatedModelOutputBatch; observations: ValidatedObservationRecord[]; sourceLast: ValidatedObservationRecord; host: ExperienceHost; }): TypedExperienceConsolidationResult { const proposals = input.output.proposals as ExperienceCandidateModelProposal[]; const observationByKey = new Map(input.observations.map(record => [observationKey(record), record])); const current = listEligibleExperiences(input.db, { userId: input.userId, now: input.output.created_at }); const candidateIds: string[] = []; let insertedCandidates = 0; input.db.exec("BEGIN IMMEDIATE"); try { for (const proposal of proposals) { const sources = proposal.source_refs.map(ref => { const observation = observationByKey.get(`${ref.file_generation}:${ref.seq}`); if (!observation) throw new Error("Typed experience source observation is unavailable"); return observation; }); if (proposal.authority === "explicit_user" && sources.every(source => source.origin.source === "advisor_finding")) { throw new Error("Advisor findings cannot establish explicit-user authority"); } const conflictsWith = proposal.kind === "episode" ? [] : current .filter(record => record.kind === proposal.kind && record.scope.kind === proposal.scope.kind && record.scope.key === proposal.scope.key && normalizedIdentityText(record.applicability) === normalizedIdentityText(proposal.applicability) && normalizedIdentityText(record.content) !== normalizedIdentityText(proposal.content)) .map(record => record.id) .sort(); const id = stableId("experience", { userId: input.userId, candidateKey: proposal.candidate_key, kind: proposal.kind, scope: proposal.scope, applicability: normalizedIdentityText(proposal.applicability), content: normalizedIdentityText(proposal.content), }); const existed = !!input.db.prepare("SELECT 1 FROM experiences WHERE user_id = ? AND id = ?").get(input.userId, id); insertExperienceCandidateInTransaction(input.db, { id, userId: input.userId, kind: proposal.kind, scope: proposal.scope, authority: proposal.authority, applicability: proposal.applicability, content: proposal.content, ...(proposal.rationale === undefined ? {} : { rationale: proposal.rationale }), exceptions: proposal.exceptions, confidenceBp: proposal.confidence_bp, validFrom: input.output.created_at, lastConfirmedAt: input.output.created_at, supersedes: [], conflictsWith, provenance: sources.map(source => ({ source: source.origin.source === "advisor_finding" ? "advisor_finding" as const : "conversation" as const, host: input.host, evidenceId: `observation:${source.file_generation}:${source.seq}:${source.checksum}`, observedAt: source.created_at, })), }, { now: input.output.created_at }); candidateIds.push(id); if (!existed) insertedCandidates += 1; } const readCoverage = recordProposalReadCoverageInTransaction({ db: input.db, userId: input.userId, fileGeneration: input.output.file_generation, seqStart: input.output.seq_start, last: input.sourceLast, createdAt: input.output.created_at, }); const auditPayload = { user_id: input.userId, file_generation: input.output.file_generation, batch_checksum: input.output.checksum, candidate_ids: candidateIds, action: "commit_typed_experiences", }; const auditId = stableId("audit", auditPayload); const auditData = canonicalJson(auditPayload); const auditChecksum = checksumJson({ table: "consolidation_audit", id: auditId, data: auditPayload }); input.db.prepare(`INSERT OR IGNORE INTO consolidation_audit (id, user_id, file_generation, proposal_batch_checksum, action, data_json, checksum, created_at) VALUES (?, ?, ?, ?, 'commit_typed_experiences', ?, ?, ?)`).run( auditId, input.userId, input.output.file_generation, input.output.checksum, auditData, auditChecksum, input.output.created_at, ); input.db.exec("COMMIT"); return { user_id: input.userId, file_generation: input.output.file_generation, candidate_ids: candidateIds, evidence_ids: [], watermark_after: null, read_watermark_after: readCoverage.watermark_after, audit_id: auditId, inserted: { candidates: insertedCandidates, evidence: 0, audit: 1, watermark: 0, read_watermark: readCoverage.inserted.read_watermark, }, }; } catch (error) { try { input.db.exec("ROLLBACK"); } catch {} throw error; } } export async function processValidatedModelOutput(input: { db: DatabaseSync; userId: string; output: ValidatedModelOutputBatch; observations: ValidatedObservationRecord[]; host?: ExperienceHost; expectedRange?: { file_generation: string; seq_start: number; seq_end: number; read_checksum: string }; semantic?: Parameters[0]["semantic"]; }): Promise< | ConsolidationResult | TypedExperienceConsolidationResult | { user_id: string; file_generation: string; candidate_ids: []; evidence_ids: []; watermark_after: null; read_watermark_after?: Watermark; pending_review_id?: string; inserted: { read_watermark?: 0 | 1; pending_review?: 0 | 1 }; } > { const userId = normalizeUserId(input.userId); if (input.output.user_id !== userId) throw new Error("Model output user mismatch"); validateModelOutputSourceRefs(input.output, input.observations); if (input.expectedRange) { if (input.output.file_generation !== input.expectedRange.file_generation || input.output.seq_start !== input.expectedRange.seq_start || input.output.seq_end !== input.expectedRange.seq_end || input.output.read_checksum !== input.expectedRange.read_checksum) throw new Error("Model output expected range mismatch"); } const sourceLast = input.observations.find((record) => record.file_generation === input.output.file_generation && record.seq === input.output.seq_end); if (!sourceLast || sourceLast.checksum !== input.output.read_checksum) throw new Error("Model output read coverage mismatch"); const typedProposalCount = input.output.proposals.filter(proposal => EXPERIENCE_KIND_SET.has(proposal.kind)).length; if (typedProposalCount > 0 && typedProposalCount !== input.output.proposals.length) { throw new Error("Model output cannot mix typed and legacy proposals"); } const conflict = findCandidateKeyConflict(input.output); if (conflict) { let pending: { id: string; inserted: boolean } | undefined; let readCoverage: { watermark_after: Watermark; inserted: { read_watermark: 0 | 1 } } | undefined; input.db.exec("BEGIN IMMEDIATE"); try { pending = insertPendingReview(input.db, { userId, kind: "candidate_key_conflict", payload: { file_generation: input.output.file_generation, seq_start: input.output.seq_start, seq_end: input.output.seq_end, conflict }, createdAt: input.output.created_at }); readCoverage = recordProposalReadCoverageInTransaction({ db: input.db, userId, fileGeneration: input.output.file_generation, seqStart: input.output.seq_start, last: sourceLast, createdAt: input.output.created_at }); input.db.exec("COMMIT"); } catch (error) { try { input.db.exec("ROLLBACK"); } catch {} throw error; } return { user_id: userId, file_generation: input.output.file_generation, candidate_ids: [], evidence_ids: [], watermark_after: null, read_watermark_after: readCoverage.watermark_after, pending_review_id: pending.id, inserted: { pending_review: pending.inserted ? 1 : 0, read_watermark: readCoverage.inserted.read_watermark } }; } if (input.output.proposals.length === 0) { const zero = recordZeroProposalReadCoverage({ db: input.db, userId, fileGeneration: input.output.file_generation, seqStart: input.output.seq_start, last: sourceLast, createdAt: input.output.created_at }); return { user_id: userId, file_generation: input.output.file_generation, candidate_ids: [], evidence_ids: [], watermark_after: null, read_watermark_after: zero.watermark_after, inserted: zero.inserted }; } if (typedProposalCount > 0) { return commitTypedExperienceProposals({ db: input.db, userId, output: input.output, observations: input.observations, sourceLast, host: input.host ?? "pi", }); } return consolidateProposalBatch({ db: input.db, userId, proposalBatch: modelOutputToProposalBatch(input.output), observations: input.observations, readCoverage: { seq_start: input.output.seq_start, last: sourceLast }, semantic: input.semantic }); }