/** * @copyright Sister Software * @license AGPL-3.0 * @author Teffen Ellis, et al. * * `mailwoman corpus audit` — per-source parquet-file count vs source_weight diagnostic. * * Reads a corpus dir's manifest.json (or scans the parquet files directly), counts files per * source, optionally loads a training config to pair the counts with the configured * source_weights. It reports the estimated sampled-row distribution at training time. * * Would have caught v0.3.0's "NAD = 411/674 train files × 2.0 weight = ~75% of sampled mix" * finding before the v0.3.0 retrospective surfaced it. * * Emits warnings to stderr and the audit table to stdout. never throws on an empty corpus. */ import { pathExists, readLocalJSONFile, readLocalTextFile } from "@mailwoman/core/fs/readers" import { basename, type PathBuilderLike, resolvePathBuilder } from "path-ts" import { Globerator } from "spliterator/node/fs" import { readConfigView } from "#source-register/effective-manifest" /** * Share of the sampled mix one source may hold before the mix is flagged as dominated by it. */ const DOMINANT_SOURCE_SHARE = 0.4 /** * How far the leading source may outweigh the runner-up before the mix is flagged. * * Catches the case where no single source clears {@link DOMINANT_SOURCE_SHARE} * but the distribution is still lopsided. */ const MAX_TOP_TO_RUNNER_UP_RATIO = 1.5 /** * Options for {@linkcode audit}. */ export interface AuditOpts { corpusDir: PathBuilderLike configPath?: PathBuilderLike /** * Sample at most N parquet files per split when counting sources. * * Default 100 for speed. * Bump to read the full set on a slow run. * * The first row of each file determines its source — corpus-v0.2.0+ files are 100% * source-segregated, so a one-row read is authoritative. */ sampleFileCount?: number } interface FileCountStats { /** * Parquet files per source per split */ bySplit: Record> /** * Total parquet files counted (may be less than file count if `sampleFileCount` caps reads) */ totalCounted: number /** * Total parquet files on disk — equals totalCounted unless capped */ totalFiles: number } interface ParsedConfig { sourceWeights: Record } /** * Read a training config's `source_weights`, or `null` when the file is absent. * * The scan lives in `readConfigView`, which the effective-manifest reader also calls, * so one reading of the hand-written config subset serves both. * This wrapper keeps the absent-file case a `null` rather than a throw, because * `corpus audit` reports a corpus without a config rather than refusing to audit it. */ async function parseConfig(configPath: PathBuilderLike): Promise { if (!(await pathExists(configPath))) return null return { sourceWeights: readConfigView(await readLocalTextFile(configPath)).sourceWeights } } /** * Scan a corpus directory's parquet files (typically under /train, /val, /test) * and count files per source per split. */ async function scanParquetFiles(corpusDir: PathBuilderLike, sampleCount: number): Promise { const stats: FileCountStats = { bySplit: {}, totalCounted: 0, totalFiles: 0 } for (const split of ["train", "val", "test"]) { const splitDir = resolvePathBuilder(corpusDir, split) if (!(await pathExists(splitDir))) continue const files = await Globerator.files("parquet", { cwd: splitDir, absolute: false, recursive: false }).toSorted() stats.totalFiles += files.length const sampleEvery = Math.max(1, Math.floor(files.length / sampleCount)) const sampled = files.filter((_, i) => i % sampleEvery === 0).slice(0, sampleCount) const splitMap: Record = {} // We can't read parquet without a dep, so we infer source from filenames where possible. // The corpus build typically writes deterministically by source — // fall back to "" when filename gives no hint. // For accurate per-source counts on real corpora, the manifest.json route below is preferred. for (const f of sampled) { const inferred = inferSourceFromFilename(f) splitMap[inferred] = (splitMap[inferred] ?? 0) + 1 } // Scale to estimated full-file counts. const scale = files.length / Math.max(sampled.length, 1) for (const k of Object.keys(splitMap)) { splitMap[k] = Math.round(splitMap[k]! * scale) } stats.bySplit[split] = splitMap stats.totalCounted += files.length } return stats } function inferSourceFromFilename(filename: string): string { // Many corpus builds write part--.parquet or part-.parquet. // The latter (current build at corpus-v0.3.0) gives no source signal in the filename. // See manifestScan() for the authoritative path. // Return "" so the caller flags this case. const m = basename(filename).match(/part-([\w-]+)-\d+\.parquet$/) if (m && m[1] !== undefined) return m[1] return "" } /** * Known source name prefixes. * * Corpus-v0.3.0 uses these as `source_id` prefixes. * The parser matches the longest prefix that fits a given `first_source_id` * recovers the canonical source name. * * Order matters: longer prefixes must be tried first so `usgov-nad-...` matches * `usgov-nad` rather than `usgov`. * Sorted descending by length at use site. */ const KNOWN_SOURCE_PREFIXES: ReadonlyArray = [ "wof-admin", "wof-postalcode", "ban", "tiger", "usgov-nad", "usgov-nppes", "usgov-hrsa-fqhc", "usgov-imls-pls", "state-ia-contractors", "state-tx-notaries", "state-ny-notaries", "openaddresses", // Synthetic adversarial sources (corpus-v0.4.0+, Thread B). "deepseek-kryptonite", "deepseek-translit-cyrl", "deepseek-translit-jpan", "deepseek-translit-hans", "deepseek-translit-hang", "deepseek-translit-armn", ] /** * Extract the source-name prefix from a `first_source_id` value. */ function sourceFromID(sourceID: string, knownPrefixes: readonly string[]): string { // Sort longest-first so usgov-nad beats usgov, wof-admin beats wof. const sorted = [...knownPrefixes].toSorted((a, b) => b.length - a.length) for (const prefix of sorted) { if (sourceID.startsWith(prefix + "-") || sourceID === prefix) return prefix } return "" } /** * Prefer reading manifest.json when present — uses each file's `first_source_id` + * prefix matching to recover the source name. * * Falls back to scanParquetFiles when manifest is absent. * * Note: corpus-v0.3.0 files can mix sources (see `last_source_id` differing from `first_source_id`). * The first-row source is an approximation. * * A full read of the parquet source column would be authoritative but requires a parquet dependency. * * For audit purposes the first-row approximation is accurate within ~5% for the * corpus-v0.3.0 shape (most files are >95% one source). */ async function manifestScan( corpusDir: PathBuilderLike, knownPrefixes: readonly string[] ): Promise { const manifestPath = resolvePathBuilder(corpusDir, "MANIFEST.json") if (!(await pathExists(manifestPath))) return null const manifest = await readLocalJSONFile<{ slices?: Array<{ split: string; source?: string | null; first_source_id?: string | null }> }>(manifestPath) if (!Array.isArray(manifest.slices)) return null const bySplit: Record> = {} for (const file of manifest.slices) { const split = file.split const src = file.source ?? sourceFromID(file.first_source_id ?? "", knownPrefixes) bySplit[split] ??= {} bySplit[split][src] = (bySplit[split][src] ?? 0) + 1 } const total = Object.values(bySplit).reduce((sum, m) => sum + Object.values(m).reduce((a, b) => a + b, 0), 0) return { bySplit, totalCounted: total, totalFiles: total } } interface AuditRow { source: string files: number filePct: number weight: number | "—" effectiveSamplePct: number | "—" overweightFactor?: number } function buildAuditRows(stats: Record, weights: Record): AuditRow[] { const totalCounted = Object.values(stats).reduce((a, b) => a + b, 0) const allSources = new Set([...Object.keys(stats), ...Object.keys(weights)]) const rows: AuditRow[] = [] // Compute effective sample weight: file_count × source_weight. // Sources with no weight get the "—" marker (loader skips them). const sampleWeights: Array<[string, number]> = [] for (const src of allSources) { const files = stats[src] ?? 0 const weight = weights[src] const effective = weight !== undefined ? files * weight : 0 sampleWeights.push([src, effective]) } const totalSampleWeight = sampleWeights.reduce((a, [, w]) => a + w, 0) for (const src of allSources) { const files = stats[src] ?? 0 const weight = weights[src] ?? "—" const effective = typeof weight === "number" ? (files * weight) / Math.max(totalSampleWeight, 1) : "—" rows.push({ source: src, files, filePct: totalCounted > 0 ? files / totalCounted : 0, weight, effectiveSamplePct: typeof effective === "number" ? effective : "—", }) } // Flag the dominator: empirically calibrated against the v0.3.0 → v0.4.0 retrospective. // v0.3.0 had usgov-nad at 52% effective sample (1.9× ban); the resulting label-space // dilution was responsible for the coarse-F1 regression. // So flag a source as "concentration warning" when it's above 40% effective sample // or more than 1.5× the next-highest. const numeric = rows.filter((r) => typeof r.effectiveSamplePct === "number") as Array< AuditRow & { effectiveSamplePct: number } > numeric.sort((a, b) => b.effectiveSamplePct - a.effectiveSamplePct) if (numeric.length) { const top = numeric[0]! const next = numeric[1]?.effectiveSamplePct ?? 0 if ( top.effectiveSamplePct > DOMINANT_SOURCE_SHARE || (next > 0 && top.effectiveSamplePct / next > MAX_TOP_TO_RUNNER_UP_RATIO) ) { top.overweightFactor = next > 0 ? top.effectiveSamplePct / next : Infinity } } rows.sort((a, b) => b.files - a.files) return rows } function formatPct(v: number | "—"): string { if (v === "—") return "—" return `${(v * 100).toFixed(1)}%` } function printReport( corpusDir: PathBuilderLike, configPath: PathBuilderLike | undefined, stats: FileCountStats, rows: AuditRow[] ): void { console.log(`\nCorpus audit — ${corpusDir}`) if (configPath) { console.log(`Config: ${configPath}`) } console.log( `Total files: ${stats.totalCounted}${stats.totalFiles !== stats.totalCounted ? ` (${stats.totalFiles} files on disk)` : ""}` ) console.log("") const trainStats = stats.bySplit["train"] if (trainStats) { const total = Object.values(trainStats).reduce((a, b) => a + b, 0) console.log(`Train split: ${total} files`) console.log("") const headers = ["source", "files", "file %", "weight", "eff. sample %"] const widths = [22, 8, 10, 8, 14] const fmtRow = (cells: string[]) => cells.map((c, i) => c.padEnd(widths[i]!)).join(" ") console.log(fmtRow(headers)) console.log(fmtRow(widths.map((w) => "─".repeat(w)))) for (const row of rows) { console.log( fmtRow([ row.source, String(row.files), formatPct(row.filePct), typeof row.weight === "number" ? row.weight.toFixed(2) : "—", formatPct(row.effectiveSamplePct), ]) ) } console.log("") const dominator = rows.find((r) => r.overweightFactor !== undefined) if (dominator) { const factor = dominator.overweightFactor const factorStr = factor === Infinity ? "∞" : factor?.toFixed(1) console.error( `⚠ Concentration: ${dominator.source} would sample ${formatPct(dominator.effectiveSamplePct)} ` + `of training rows (${factorStr}× the next-highest). ` + `Past lesson: v0.3.0's NAD at ~52% caused the 21-label coarse regression. ` + `Consider lowering this source's weight or boosting others.` ) } else { console.log("✓ No single-source concentration (top source < 40% effective sample AND < 1.5× next).") } const missingWeights = rows.filter((r) => r.weight === "—" && r.files > 0) if (missingWeights.length && configPath) { console.error( `⚠ Sources present in corpus but absent from config.source_weights ` + `(loader will skip them): ${missingWeights.map((r) => r.source).join(", ")}` ) } const orphanWeights = rows.filter((r) => typeof r.weight === "number" && r.files === 0) if (orphanWeights.length) { console.error( `⚠ Sources weighted in config but no parquet files found in corpus ` + `(no-op weights): ${orphanWeights.map((r) => r.source).join(", ")}` ) } } } export async function audit(opts: AuditOpts): Promise { const config = opts.configPath ? await parseConfig(opts.configPath) : null // Compose the known-prefix list from both the hardcoded set and any extra names in the // config (forward-compat for future adapters added before this file is updated). const prefixes = [...new Set([...KNOWN_SOURCE_PREFIXES, ...Object.keys(config?.sourceWeights ?? {})])] const stats = (await manifestScan(opts.corpusDir, prefixes)) ?? (await scanParquetFiles(opts.corpusDir, opts.sampleFileCount ?? 100)) const trainStats = stats.bySplit["train"] ?? {} const rows = buildAuditRows(trainStats, config?.sourceWeights ?? {}) printReport(opts.corpusDir, opts.configPath, stats, rows) }