/** * @copyright Sister Software * @license AGPL-3.0 * @author Teffen Ellis, et al. * * Adapter runner — drives a `CorpusAdapter` to completion and writes `canonical.jsonl` plus a * `MANIFEST.json` per adapter. It does not align, tokenize, synthesize, or write Parquet. */ import { type WriteStream } from "@mailwoman/core/fs/streams"; import { type PathBuilderLike } from "path-ts"; import type { AdapterRegistry } from "#adapters/registry"; import type { AdapterOptions, CorpusAdapter } from "#types"; /** * Snapshot of the runner's state, emitted on every progress tick. */ export interface RunnerProgress { adapterID: string; /** * Total rows the adapter has yielded, before dedup. */ yielded: number; written: number; bytes: number; elapsed_ms: number; } /** * Base-2 logarithm of the fingerprint table's slot count. * * 2^27 slots hold 93,952,409 keys before the load limit. * The largest measured adapter has 57,570,829 distinct keys. * The table occupies 2.0 GiB outside the V8 heap. */ export declare const DEFAULT_DEDUP_SLOTS_LOG2 = 27; /** * How many distinct dedup keys a `Set`-backed run holds before it stops adding new ones. * * A V8 `Set` refuses a 16,777,216th entry outright, so this cap was a real bound rather than a preference. * It applies only to a run that opts out of the fingerprint table * through {@linkcode RunAdapterOptions.dedupMaxSize}. */ export declare const DEFAULT_DEDUP_MAX_SIZE = 10000000; /** * Where a run held its dedup keys. */ export declare const DedupStore: { /** * 128-bit fingerprints in a flat `Uint32Array`, bounded by the array the run sized. */ readonly FingerprintTable: "fingerprint-table"; /** * A V8 `Set` of key strings, capped because V8 refuses a 16,777,216th entry. */ readonly CappedSet: "capped-set"; }; export type DedupStore = (typeof DedupStore)[keyof typeof DedupStore]; export interface RunAdapterOptions { adapter: CorpusAdapter; adapterOptions: AdapterOptions; /** * Hold dedup keys in a V8 `Set` capped at this many, rather than in a fingerprint table. * * The cap is what the fingerprint table replaces. * A run that sets this reproduces the pre-2026-09-28 behavior. * * Past the cap, the runner still drops a duplicate of a key it holds * and writes a duplicate of a key first seen after the cap. * * Over `v0.7.0-de-holdout` that wrote 1,490,992 duplicate rows across three adapters. * A test sets it low to reach exhaustion in a few rows. * * After reaching the cap, the runner still drops duplicates for keys already in the set. * It writes duplicates for keys first seen after the cap. */ dedupMaxSize?: number; /** * Base-2 logarithm of the fingerprint table's slot count. * * Defaults to {@linkcode DEFAULT_DEDUP_SLOTS_LOG2}. * A test sets it low to keep the allocation small. */ dedupSlotsLog2?: number; /** * Root output directory. * The runner creates `//` under it. */ outputDir: PathBuilderLike; /** * Corpus version stamped onto every row, locked together with the tokenizer version. */ corpusVersion: string; /** * The `source` id stamped on every emitted row, when it differs from the adapter's own id. * * A source id is a wire identifier keyed by `source_weights`, so re-using a name a built * corpus already contains makes this run's rows indistinguishable from that corpus's. * When absent, rows use `adapter.id`. */ sourceName?: string; /** * Invoked every `progressEvery` rows yielded and once at the end. * A thrown error aborts the run. */ onProgress?: (snapshot: RunnerProgress) => void; /** * Yielded-row interval at which `onProgress` fires. * * Defaults to 1000 rows per callback. * The runner always emits a terminal update. */ progressEvery?: number; } /** * Return value of `runAdapter`, the same shape as the manifest written to disk. */ export interface AdapterRunManifest { adapter_id: string; corpus_version: string; default_license: string; description: string; yielded: number; written: number; deduped: number; /** * Input records the adapter refused or trimmed before yielding, keyed by reason. * * See {@link AdapterOptions.dropped} for the key shapes. * An empty object means the adapter does not count its drops, so the number it dropped is unmeasured. */ dropped: Record; /** * The `yielded` count at which the dedup set stopped growing, or `null` where it never did. * * A run reporting a number here deduplicated its rows completely up to that point. * After that point, it still drops duplicates for keys already held. * * It writes duplicates for keys first seen after the cap. * `deduped` by itself cannot identify which, so a consumer comparing two builds' * duplicate counts needs this beside it. */ dedup_exhausted_at_yielded: number | null; /** * How the run held its dedup keys: `fingerprint-table` or `capped-set`. * * The two reach different row counts on the same input, so a consumer comparing * two builds needs to know which each used. * `capped-set` writes a duplicate of any key first seen after its cap. */ dedup_store: DedupStore; /** * Distinct dedup keys the run held at the end. * * Under `fingerprint-table` this is a fingerprint count, so two keys sharing all 128 bits count once. */ dedup_keys: number; bytes: number; sha256: string; jsonl_path: string; started_at: string; ended_at: string; elapsed_ms: number; } /** * Drive a single adapter to completion, writing its jsonl and manifest under `outputDir//`. */ export declare function runAdapter(opts: RunAdapterOptions): Promise; /** * Drive every adapter in a registry sequentially, stopping on the first failure. */ export declare function runAllAdapters(registry: AdapterRegistry, common: Omit & { adapterOptionsFor?: (a: CorpusAdapter) => AdapterOptions; }): Promise; /** * Resolve once everything written to `stream` so far has reached the file. * * A zero-length write acts as a barrier. * Its callback runs after earlier queued writes. * * A caller can then record a byte offset for the bytes in the file. * `drain` fires only when the buffer was full. * * `bytesWritten` excludes queued data. * Neither value answers when all earlier writes have reached the file. */ export declare function flushStream(stream: WriteStream): Promise; /** * Promise-ify a single `drain` or `close` emission, detaching the loser listener * so a long-lived stream does not accumulate an orphan handler per wait. */ export declare function once(emitter: WriteStream, event: "drain" | "close"): Promise; //# sourceMappingURL=runner.d.ts.map