/** * #866 (T10): record-class partition migrators (build only; execution is #861). * * Line-closed placement by each row's typed kind (never by source directory alone): * - ticket-provenance → /records.jsonl (live topology; human view cancelled #900) * - submission-ledger kinds → owning run session/submission-ledger/ * - attempt-history → owning run session/attempt-history/ * - unknown ownership → unbound/ (never silent discard) * - malformed → unbound with exact bytes + malformedRows entry * - misplaced: any foreign jsonl whose kind is one of the three classes */ import { createHash } from "node:crypto"; import { appendFile, mkdir, readFile, stat, writeFile } from "node:fs/promises"; import { dirname, join, relative, sep } from "node:path"; import { bookHistoricalRoots, findBookRunDirectory, findPlacedMigratingRun, findUniquePrincipalPlacedRun, isMigrationEnoent, isTicketNumberString, listMigrationBookKeys, listMigrationDirents, runCoordsFromSessionParent, runIdFromSubject, runRefFromBoundPath, ticketNumberFromSubject, } from "./book-topology-migration-placement.ts"; import { reconcileMigrationPartition, type BookTopologyMigrationContext, type BookTopologyPartitionMigrator, type MigrationItemOutcome, } from "./book-topology-migration.ts"; import { S4_SUBMISSION_LEDGER_KINDS } from "./sitian-appender.ts"; import { projectTicketProvenanceHeader } from "./ticket-provenance-contracts.ts"; const TICKET_PROVENANCE = "ticket-provenance"; const SUBMISSION_LEDGER = "submission-ledger"; const ATTEMPT_HISTORY = "attempt-history"; const MISPLACED = "misplaced-record-class"; type RecordClass = typeof TICKET_PROVENANCE | typeof SUBMISSION_LEDGER | typeof ATTEMPT_HISTORY; type JsonLine = | { readonly ok: true; readonly value: Record; readonly raw: string } | { readonly ok: false; readonly raw: string }; function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function parseJsonlLine(raw: string): JsonLine { try { const value: unknown = JSON.parse(raw); if (!isRecord(value)) return { ok: false, raw }; return { ok: true, value, raw }; } catch { return { ok: false, raw }; } } function sourceRelative(backupBooksDirectory: string, absolutePath: string): string { return relative(backupBooksDirectory, absolutePath).split(sep).join("/"); } function lineSource(backupBooksDirectory: string, filePath: string, index: number): string { return `${sourceRelative(backupBooksDirectory, filePath)}#${index + 1}`; } function stableKey(material: string): string { return createHash("sha256").update(material).digest("hex").slice(0, 32); } async function readJsonlLines(filePath: string): Promise { let text: string; try { text = await readFile(filePath, "utf8"); } catch (error) { if (isMigrationEnoent(error)) return []; throw error; } if (text.length === 0) return []; const lines = text.split("\n"); if (text.endsWith("\n")) lines.pop(); return lines; } async function ensureDir(path: string): Promise { await mkdir(path, { recursive: true }); } /** One provisional line written to a ticket dest before a bare volume may supersede it. */ type DestLinePlacement = { readonly raw: string; /** Index into the partition outcomes array — flipped to unbound if bare replaces dest. */ readonly outcomeIndex: number; }; /** Per-destination identity cache — one read per file, O(1) subsequent checks. */ class RecordWriteCache { private readonly identities = new Map>(); /** Dest paths written as whole #901 bare volumes — legacy must not append under them. */ private readonly bareVolumeFiles = new Set(); /** Lines currently on a dest that are still claimed `placed` in outcomes. */ private readonly destPlacements = new Map(); private async load(recordFile: string): Promise> { const cached = this.identities.get(recordFile); if (cached !== undefined) return cached; const set = new Set(); for (const line of await readJsonlLines(recordFile)) { if (line.trim() === "") continue; const parsed = parseJsonlLine(line); if (parsed.ok && typeof parsed.value.identity === "string") set.add(parsed.value.identity); } this.identities.set(recordFile, set); return set; } markBareVolume(recordFile: string): void { this.bareVolumeFiles.add(recordFile); this.identities.set(recordFile, new Set()); } isBareVolume(recordFile: string): boolean { return this.bareVolumeFiles.has(recordFile); } /** * After a bare whole-file write, register its lines so a later bare for the * same ticket can rehome them and flip dispositions (multi-bare honesty). */ setDestPlacements(recordFile: string, placements: readonly DestLinePlacement[]): void { this.destPlacements.set(recordFile, [...placements]); } /** * Install a bare volume as the sole ticket file. Any lines already placed on * this dest (legacy appends or a prior bare) are moved to unbound and their * outcomes flipped — report stays honest. */ async replaceWholeBare( recordFile: string, body: string, outcomes: MigrationItemOutcome[], unboundFileForRaw: (raw: string) => string, ): Promise { const prior = this.destPlacements.get(recordFile) ?? []; for (const placement of prior) { await this.append(unboundFileForRaw(placement.raw), placement.raw); const previous = outcomes[placement.outcomeIndex]; if (previous !== undefined) { outcomes[placement.outcomeIndex] = { disposition: "unbound", source: previous.source, }; } } this.destPlacements.delete(recordFile); await ensureDir(dirname(recordFile)); const text = body.endsWith("\n") ? body : `${body}\n`; await writeFile(recordFile, text, "utf8"); this.markBareVolume(recordFile); } async append( recordFile: string, raw: string, options?: { readonly outcomeIndex?: number }, ): Promise { await ensureDir(dirname(recordFile)); const line = raw.endsWith("\n") ? raw : `${raw}\n`; const parsed = parseJsonlLine(raw); if (parsed.ok && typeof parsed.value.identity === "string") { const existing = await this.load(recordFile); if (!existing.has(parsed.value.identity)) { existing.add(parsed.value.identity); await appendFile(recordFile, line, "utf8"); } // Identity hit skips a second physical copy, but the source row still claims // this dest — register outcome below so bare replace flips every claim. } else { if (!this.identities.has(recordFile)) { // Ensure subsequent identity loads see a live map even when first rows lack identity. this.identities.set(recordFile, await this.load(recordFile)); } await appendFile(recordFile, line, "utf8"); } // Register outcome claims on dest even when identity skipped a write. if (options?.outcomeIndex !== undefined && !this.bareVolumeFiles.has(recordFile)) { const list = this.destPlacements.get(recordFile) ?? []; list.push({ raw, outcomeIndex: options.outcomeIndex }); this.destPlacements.set(recordFile, list); } } } function roleFromPayload(payload: unknown): string | undefined { if (isRecord(payload) && typeof payload.role === "string" && payload.role.length > 0) { return payload.role; } return undefined; } function runIdFromPayload(payload: unknown): string | undefined { return isRecord(payload) && typeof payload.runId === "string" && payload.runId.length > 0 ? payload.runId : undefined; } function recordClassOfKind(kind: unknown): RecordClass | undefined { if (kind === TICKET_PROVENANCE) return TICKET_PROVENANCE; if (kind === ATTEMPT_HISTORY) return ATTEMPT_HISTORY; if (typeof kind === "string" && S4_SUBMISSION_LEDGER_KINDS.has(kind)) return SUBMISSION_LEDGER; return undefined; } async function listFilesRecursive(root: string, predicate: (name: string) => boolean): Promise { const out: string[] = []; async function walk(directory: string): Promise { for (const entry of await listMigrationDirents(directory)) { const path = join(directory, entry.name); if (entry.isDirectory()) await walk(path); else if (entry.isFile() && predicate(entry.name)) out.push(path); } } await walk(root); return out; } /** * Home partition sources under one record-class root: * - direct `/records.jsonl` (legacy flat book-root volume) * - hashed sub-volumes `//records.jsonl` * Both enter the same line-closed placement path; ENOENT yields empty reads. */ async function listVolumeRecordFiles(partitionDir: string): Promise { const files: string[] = [join(partitionDir, "records.jsonl")]; for (const entry of await listMigrationDirents(partitionDir)) { if (!entry.isDirectory()) continue; files.push(join(partitionDir, entry.name, "records.jsonl")); } return files; } /** * Paths already closed by the three home migrators — including the partition-root * direct records.jsonl. Criterion (not a partition name list): root record-class * trees and nested ticket-provenance. Run trees are NOT blanket-skipped — wrong-kind * rows inside a run must still rehome by content (#866 / #852 §落错). */ function isHomeRecordClassRelativePath(relPosix: string): boolean { const parts = relPosix.split("/"); if (parts[0] === TICKET_PROVENANCE || parts[0] === SUBMISSION_LEDGER || parts[0] === ATTEMPT_HISTORY) { return true; } if (parts.length >= 2 && isTicketNumberString(parts[0]!)) { // Partial-nest era + live ticket-root volume (records.jsonl at ticket dir). if (parts[1] === TICKET_PROVENANCE) return true; if (parts.length === 2 && parts[1] === "records.jsonl") return true; } return false; } /** Canonical nested path for a run-owned record class; ticket-provenance is never run-nested. */ function isCanonicalRunNestedRecordFile(relPosix: string, recordClass: RecordClass): boolean { if (recordClass === TICKET_PROVENANCE) return false; const needle = `/session/${recordClass}/records.jsonl`; return relPosix.endsWith(needle) || relPosix.endsWith(needle.slice(1)); } function runCoordsFromRelativePath(relPosix: string): { readonly runId: string; readonly role: string } | undefined { const parts = relPosix.split("/"); const runsIndex = parts.indexOf("runs"); if (runsIndex < 0 || runsIndex + 1 >= parts.length) return undefined; const leaf = parts[runsIndex + 1]!; const at = leaf.indexOf("@"); if (at <= 0 || at === leaf.length - 1) return undefined; return { runId: leaf.slice(0, at), role: leaf.slice(at + 1) }; } /** Typed owning run for a record-class row (subject / payload / sessionParent). */ function owningRunIdFromRecord(value: Record): string | undefined { return ( runIdFromSubject(value.subject) ?? runIdFromPayload(value.payload) ?? runCoordsFromSessionParent(value.sessionParent)?.runId ); } /** * Row already sits at the correct run's canonical nest — #865 carries the file. * Path-class alone is not enough: a submission-ledger row under run A that names * run B must still rehome (#866 run-principal). */ function isAlreadyHomeUnderOwningRun( withinBook: string, recordClass: RecordClass, value: Record, ): boolean { if (!isCanonicalRunNestedRecordFile(withinBook, recordClass)) return false; const pathCoords = runCoordsFromRelativePath(withinBook); if (pathCoords === undefined) return false; const owningRunId = owningRunIdFromRecord(value); return owningRunId !== undefined && owningRunId === pathCoords.runId; } function relativePathWithinRun(relPosix: string): string | undefined { const parts = relPosix.split("/"); const runsIndex = parts.indexOf("runs"); if (runsIndex < 0 || runsIndex + 2 >= parts.length) return undefined; return parts.slice(runsIndex + 2).join("/"); } function unboundCategoryFile( booksDirectory: string, bookKey: string, category: string, key: string, ): string { return join(booksDirectory, bookKey, "unbound", category, key, "records.jsonl"); } async function resolveRunDestination( context: BookTopologyMigrationContext, bookKey: string, runId: string, hints: { readonly role?: string; readonly sessionParent?: unknown }, ): Promise< | { readonly kind: "run"; readonly runDirectory: string; readonly disposition: "placed" | "unbound" } | { readonly kind: "unbound-key"; readonly key: string } > { const hasSessionParent = typeof hints.sessionParent === "string" && hints.sessionParent.length > 0; const bound = hasSessionParent ? runRefFromBoundPath( hints.sessionParent, bookHistoricalRoots(context.booksDirectory, context.backupBooksDirectory, bookKey), ) : undefined; // Path present but not bound to this book: never steal role or unique-run-guess. if (hasSessionParent && bound === undefined) { return { kind: "unbound-key", key: stableKey(runId) }; } if (bound !== undefined) { const existingDest = await findPlacedMigratingRun( context.booksDirectory, bookKey, bound.leaf, bound.sourceRelative, ); if (existingDest !== undefined) { return { kind: "run", runDirectory: existingDest.runDirectory, disposition: existingDest.disposition, }; } return { kind: "unbound-key", key: stableKey(bound.leaf) }; } if (hints.role !== undefined && hints.role.length > 0) { const existingDest = await findPlacedMigratingRun( context.booksDirectory, bookKey, `${runId}@${hints.role}`, ); if (existingDest !== undefined) { return { kind: "run", runDirectory: existingDest.runDirectory, disposition: existingDest.disposition, }; } return { kind: "unbound-key", key: stableKey(`${runId}@${hints.role}`) }; } // No path, no role: follow only a uniquely matching principal in this book. const uniquePrincipal = await findUniquePrincipalPlacedRun( context.booksDirectory, bookKey, runId, ); if (uniquePrincipal !== undefined) { return { kind: "run", runDirectory: uniquePrincipal.runDirectory, disposition: uniquePrincipal.disposition, }; } return { kind: "unbound-key", key: stableKey(runId) }; } async function placeTicketProvenanceLine( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, source: string, raw: string, value: Record | undefined, /** Outcomes length at call time = index of the outcome the caller is about to push. */ outcomeIndex: number, ): Promise { if (value === undefined) { await writes.append(unboundCategoryFile(context.booksDirectory, bookKey, TICKET_PROVENANCE, stableKey(raw)), raw); return { disposition: "unbound", source, malformed: true, malformedRaw: raw }; } const ticketNumber = ticketNumberFromSubject(value.subject); if (ticketNumber === undefined) { await writes.append(unboundCategoryFile(context.booksDirectory, bookKey, TICKET_PROVENANCE, stableKey(raw)), raw); return { disposition: "unbound", source }; } // Live topology (docs/dossier-topology.md): ticket root holds records.jsonl. const dest = join( context.booksDirectory, bookKey, String(ticketNumber), "records.jsonl", ); // Never append legacy SitianRecord under a bare #901 volume (header must stay first). if (writes.isBareVolume(dest)) { await writes.append( unboundCategoryFile(context.booksDirectory, bookKey, TICKET_PROVENANCE, stableKey(raw)), raw, ); return { disposition: "unbound", source }; } const existing = await readJsonlLines(dest); if (tryBareTicketProvenanceVolume(existing) !== undefined) { writes.markBareVolume(dest); await writes.append( unboundCategoryFile(context.booksDirectory, bookKey, TICKET_PROVENANCE, stableKey(raw)), raw, ); return { disposition: "unbound", source }; } await writes.append(dest, raw, { outcomeIndex }); return { disposition: "placed", source }; } async function placeRunOwnedLine( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, category: typeof SUBMISSION_LEDGER | typeof ATTEMPT_HISTORY, source: string, raw: string, value: Record | undefined, ): Promise { if (value === undefined) { await writes.append(unboundCategoryFile(context.booksDirectory, bookKey, category, stableKey(raw)), raw); return { disposition: "unbound", source, malformed: true, malformedRaw: raw }; } const fromParent = runCoordsFromSessionParent(value.sessionParent); const runId = runIdFromSubject(value.subject) ?? runIdFromPayload(value.payload) ?? fromParent?.runId; if (runId === undefined) { await writes.append(unboundCategoryFile(context.booksDirectory, bookKey, category, stableKey(raw)), raw); return { disposition: "unbound", source }; } const role = roleFromPayload(value.payload); const target = await resolveRunDestination(context, bookKey, runId, { ...(role === undefined ? {} : { role }), ...(value.sessionParent === undefined ? {} : { sessionParent: value.sessionParent }), }); if (target.kind === "unbound-key") { await writes.append(unboundCategoryFile(context.booksDirectory, bookKey, category, target.key), raw); return { disposition: "unbound", source }; } await writes.append(join(target.runDirectory, "session", category, "records.jsonl"), raw); return { disposition: target.disposition, source }; } /** Place one parsed row by its typed kind; unknown kinds fall back to the home category unbound. */ async function placeRecordClassLine( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, source: string, raw: string, value: Record, fallbackCategory: RecordClass, outcomeIndex: number, ): Promise { const recordClass = recordClassOfKind(value.kind); if (recordClass === TICKET_PROVENANCE) { return placeTicketProvenanceLine(context, bookKey, writes, source, raw, value, outcomeIndex); } if (recordClass === SUBMISSION_LEDGER || recordClass === ATTEMPT_HISTORY) { return placeRunOwnedLine(context, bookKey, writes, recordClass, source, raw, value); } // Unknown kind inside a record-class home volume: preserve under that home's unbound. await writes.append( unboundCategoryFile(context.booksDirectory, bookKey, fallbackCategory, stableKey(raw)), raw, ); return { disposition: "unbound", source }; } /** * #901 bare diary volume: first non-empty line is a header without SitianRecord * `kind`. Recognise the whole file as one ticket volume — do not split rows into unbound. */ function tryBareTicketProvenanceVolume( lines: readonly string[], ): { readonly ticket: number; readonly body: readonly string[] } | undefined { const body: string[] = []; for (const raw of lines) { if (raw.trim() === "") continue; body.push(raw); } if (body.length === 0) return undefined; const first = parseJsonlLine(body[0]!); if (!first.ok) return undefined; // Legacy SitianRecord rows carry a typed kind — leave them to line placement. if (recordClassOfKind(first.value.kind) !== undefined) return undefined; const header = projectTicketProvenanceHeader(first.value); if (header === undefined) return undefined; return { ticket: header.ticket, body }; } async function placeBareTicketProvenanceVolume( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, filePath: string, bare: { readonly ticket: number; readonly body: readonly string[] }, outcomes: MigrationItemOutcome[], ): Promise { const dest = join( context.booksDirectory, bookKey, String(bare.ticket), "records.jsonl", ); // Whole-file bare install. Prior dest lines (legacy or earlier bare) rehome to // unbound with flipped dispositions — header stays first; no silent erase. const body = bare.body.map((raw) => (raw.endsWith("\n") ? raw : `${raw}\n`)).join(""); await writes.replaceWholeBare(dest, body, outcomes, (raw) => unboundCategoryFile(context.booksDirectory, bookKey, TICKET_PROVENANCE, stableKey(raw)), ); const barePlacements: DestLinePlacement[] = []; for (let index = 0; index < bare.body.length; index += 1) { const source = lineSource(context.backupBooksDirectory, filePath, index); const outcomeIndex = outcomes.length; outcomes.push({ disposition: "placed", source }); barePlacements.push({ raw: bare.body[index]!, outcomeIndex }); } // Register so a later bare for the same ticket can supersede honestly. writes.setDestPlacements(dest, barePlacements); } async function migrateJsonlFileByKind( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, filePath: string, fallbackCategory: RecordClass, outcomes: MigrationItemOutcome[], ): Promise { const lines = await readJsonlLines(filePath); if (fallbackCategory === TICKET_PROVENANCE) { const bare = tryBareTicketProvenanceVolume(lines); if (bare !== undefined) { await placeBareTicketProvenanceVolume(context, bookKey, writes, filePath, bare, outcomes); return; } } for (let index = 0; index < lines.length; index += 1) { const raw = lines[index]!; if (raw.trim() === "") continue; const source = lineSource(context.backupBooksDirectory, filePath, index); const parsed = parseJsonlLine(raw); if (!parsed.ok) { // Malformed: preserve under the home partition's unbound bucket. if (fallbackCategory === TICKET_PROVENANCE) { outcomes.push( await placeTicketProvenanceLine( context, bookKey, writes, source, raw, undefined, outcomes.length, ), ); } else { outcomes.push(await placeRunOwnedLine(context, bookKey, writes, fallbackCategory, source, raw, undefined)); } continue; } outcomes.push( await placeRecordClassLine( context, bookKey, writes, source, raw, parsed.value, fallbackCategory, outcomes.length, ), ); } } const OFFERED_IDENTITIES = "offered-identities.jsonl" as const; /** * Merge offered-identities from one source volume into the destination ticket * dir by identity key. Never last-source-wins overwrite of the whole file. */ async function mergeOfferedIdentities( destDir: string, volumeDir: string, writes: RecordWriteCache, ): Promise { const from = join(volumeDir, OFFERED_IDENTITIES); let lines: readonly string[]; try { await stat(from); lines = await readJsonlLines(from); } catch (error) { if (isMigrationEnoent(error)) return; throw error; } const destFile = join(destDir, OFFERED_IDENTITIES); for (const raw of lines) { if (raw.trim() === "") continue; await writes.append(destFile, raw); } } /** Collect typed ticket numbers present in a source volume's records.jsonl. */ async function ticketNumbersInVolume(volumeDir: string): Promise { const lines = await readJsonlLines(join(volumeDir, "records.jsonl")); const bare = tryBareTicketProvenanceVolume(lines); if (bare !== undefined) return [bare.ticket]; const out: number[] = []; const seen = new Set(); for (const raw of lines) { if (raw.trim() === "") continue; const parsed = parseJsonlLine(raw); if (!parsed.ok) continue; if (recordClassOfKind(parsed.value.kind) !== TICKET_PROVENANCE) continue; const ticketNumber = ticketNumberFromSubject(parsed.value.subject); if (ticketNumber === undefined || seen.has(ticketNumber)) continue; seen.add(ticketNumber); out.push(ticketNumber); } return out; } async function mergeCompanionsFromVolume( context: BookTopologyMigrationContext, bookKey: string, volumeDir: string, writes: RecordWriteCache, ): Promise { for (const ticketNumber of await ticketNumbersInVolume(volumeDir)) { const destDir = join(context.booksDirectory, bookKey, String(ticketNumber)); await ensureDir(destDir); await mergeOfferedIdentities(destDir, volumeDir, writes); } } async function migrateHomePartitionBook( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, category: RecordClass, outcomes: MigrationItemOutcome[], ): Promise { const backupBook = join(context.backupBooksDirectory, bookKey); // listVolumeRecordFiles covers partition-root records.jsonl + hashed sub-volumes. const rootFiles = await listVolumeRecordFiles(join(backupBook, category)); for (const filePath of rootFiles) { await migrateJsonlFileByKind(context, bookKey, writes, filePath, category, outcomes); if (category === TICKET_PROVENANCE) { await mergeCompanionsFromVolume(context, bookKey, dirname(filePath), writes); } } if (category !== TICKET_PROVENANCE) return; // One ticket-dir walk: partial-nest /ticket-provenance/ and live // / root volumes share existence check, line placement, companions. for (const entry of await listMigrationDirents(backupBook)) { if (!entry.isDirectory() || !isTicketNumberString(entry.name)) continue; const ticketDir = join(backupBook, entry.name); for (const volumeDir of [ join(ticketDir, TICKET_PROVENANCE), ticketDir, ] as const) { const recordsFile = join(volumeDir, "records.jsonl"); try { await stat(recordsFile); } catch (error) { if (isMigrationEnoent(error)) continue; throw error; } await migrateJsonlFileByKind(context, bookKey, writes, recordsFile, TICKET_PROVENANCE, outcomes); await mergeCompanionsFromVolume(context, bookKey, volumeDir, writes); } } } function normalizeJsonlRaw(raw: string): string { return raw.endsWith("\n") ? raw : `${raw}\n`; } /** Rewrite dest file without the given exact raw lines (drop wrong-nest copies after rehome). */ async function scrubRawLinesFromFile( recordFile: string, rawLines: ReadonlySet, ): Promise { if (rawLines.size === 0) return; const lines = await readJsonlLines(recordFile); if (lines.length === 0) return; const kept: string[] = []; let changed = false; for (const raw of lines) { if (raw.trim() === "") continue; const normalized = normalizeJsonlRaw(raw); if (rawLines.has(normalized)) { changed = true; continue; } kept.push(normalized); } if (!changed) return; await ensureDir(dirname(recordFile)); await writeFile(recordFile, kept.join(""), "utf8"); } /** * Foreign jsonl rows whose typed kind is one of the three record classes. * Range = whole book minus home partitions (root records.jsonl, hashed volumes, * nested ticket-provenance) — those close under the three home migrators so T8 * does not double-count. Inside runs/: skip only rows already under their typed * owning run at the canonical nest (those travel with #865); path-class alone * never excuses a cross-run wrong principal. Every other typed row rehomes by * content; dest copies at the source nest are scrubbed so a recursive run cp * cannot keep a lying duplicate. */ async function migrateMisplacedBook( context: BookTopologyMigrationContext, bookKey: string, writes: RecordWriteCache, outcomes: MigrationItemOutcome[], ): Promise { const backupBook = join(context.backupBooksDirectory, bookKey); const files = await listFilesRecursive(backupBook, (name) => name.endsWith(".jsonl")); // backupRelPath → exact raw lines rehomed out of a source-run nest (wrong path or wrong principal). // Match by line bytes, not identity — rows without identity must still leave the lying source copy. const scrubPlans = new Map>(); for (const filePath of files) { const rel = sourceRelative(context.backupBooksDirectory, filePath); const withinBook = rel.split("/").slice(1).join("/"); // Home partition roots (incl. direct records.jsonl) are owned by home migrators. if (isHomeRecordClassRelativePath(withinBook)) continue; const lines = await readJsonlLines(filePath); for (let index = 0; index < lines.length; index += 1) { const raw = lines[index]!; if (raw.trim() === "") continue; const parsed = parseJsonlLine(raw); if (!parsed.ok) continue; const recordClass = recordClassOfKind(parsed.value.kind); if (recordClass === undefined) continue; // Correct principal + canonical nest — #865 carries the file as-is. if (isAlreadyHomeUnderOwningRun(withinBook, recordClass, parsed.value)) continue; const source = lineSource(context.backupBooksDirectory, filePath, index); if (recordClass === TICKET_PROVENANCE) { outcomes.push( await placeTicketProvenanceLine( context, bookKey, writes, source, raw, parsed.value, outcomes.length, ), ); } else { outcomes.push( await placeRunOwnedLine(context, bookKey, writes, recordClass, source, raw, parsed.value), ); } if (withinBook.includes("/runs/") || withinBook.startsWith("runs/")) { let set = scrubPlans.get(withinBook); if (set === undefined) { set = new Set(); scrubPlans.set(withinBook, set); } set.add(normalizeJsonlRaw(raw)); } } } // Drop rehomed lines from the source nest in dest (wrong path or wrong principal). for (const [withinBook, rawLines] of scrubPlans) { const coords = runCoordsFromRelativePath(withinBook); const withinRun = relativePathWithinRun(withinBook); if (coords === undefined || withinRun === undefined) continue; // Scrub plan already carries runId+role from the source leaf — keep the full identity. const destRun = await findBookRunDirectory( join(context.booksDirectory, bookKey), coords.runId, coords.role, ); if (destRun === undefined) continue; await scrubRawLinesFromFile(join(destRun.runDirectory, withinRun), rawLines); } } async function forEachBook( context: BookTopologyMigrationContext, body: (bookKey: string, writes: RecordWriteCache) => Promise, ): Promise { for (const bookKey of await listMigrationBookKeys(context.backupBooksDirectory)) { await body(bookKey, new RecordWriteCache()); } } export const ticketProvenancePartitionMigrator: BookTopologyPartitionMigrator = { partition: TICKET_PROVENANCE, async migrate(context) { const outcomes: MigrationItemOutcome[] = []; await forEachBook(context, (bookKey, writes) => migrateHomePartitionBook(context, bookKey, writes, TICKET_PROVENANCE, outcomes)); return reconcileMigrationPartition(TICKET_PROVENANCE, "lines", outcomes); }, }; export const submissionLedgerPartitionMigrator: BookTopologyPartitionMigrator = { partition: SUBMISSION_LEDGER, async migrate(context) { const outcomes: MigrationItemOutcome[] = []; await forEachBook(context, (bookKey, writes) => migrateHomePartitionBook(context, bookKey, writes, SUBMISSION_LEDGER, outcomes)); return reconcileMigrationPartition(SUBMISSION_LEDGER, "lines", outcomes); }, }; export const attemptHistoryPartitionMigrator: BookTopologyPartitionMigrator = { partition: ATTEMPT_HISTORY, async migrate(context) { const outcomes: MigrationItemOutcome[] = []; await forEachBook(context, (bookKey, writes) => migrateHomePartitionBook(context, bookKey, writes, ATTEMPT_HISTORY, outcomes)); return reconcileMigrationPartition(ATTEMPT_HISTORY, "lines", outcomes); }, }; export const misplacedRecordClassPartitionMigrator: BookTopologyPartitionMigrator = { partition: MISPLACED, async migrate(context) { const outcomes: MigrationItemOutcome[] = []; await forEachBook(context, (bookKey, writes) => migrateMisplacedBook(context, bookKey, writes, outcomes)); return reconcileMigrationPartition(MISPLACED, "lines", outcomes); }, }; export const BOOK_TOPOLOGY_RECORD_CLASS_MIGRATORS: readonly BookTopologyPartitionMigrator[] = [ ticketProvenancePartitionMigrator, submissionLedgerPartitionMigrator, attemptHistoryPartitionMigrator, misplacedRecordClassPartitionMigrator, ];