lila/data-pipeline/pipeline.ts
lila 37c978e230 implementing phase 3 pipeline: gemini structured output, validation, sqlite staging
Prompt is now a template (fixes the hardcoded en/es leftovers in rules
2, 3, 15, 16, 26, 31). pipeline.ts replaces the pseudocode: wordlist
normalization, skip-already-staged idempotency, batches of 20 against
gemini-3.6-flash with responseSchema, raw responses persisted per batch,
per-entry validation with rejection log, one transaction per word into
db/staging.db. Flags: --langs --pos --max-batches --delay-ms --dry-run.

Smoke run: 40/40 words staged (de+es, one batch each), 0 rejections.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-09 19:04:04 +02:00

327 lines
9.5 KiB
TypeScript

import { appendFileSync, mkdirSync, writeFileSync } from "node:fs";
import path from "node:path";
import { parseArgs } from "node:util";
import {
SUPPORTED_LANGUAGE_CODES,
SUPPORTED_POS,
type SupportedLanguageCode,
type SupportedPos,
} from "@lila/shared";
import {
DEFAULT_MODEL,
buildEntriesResponseSchema,
generateContent,
parseEntries,
} from "./gemini.js";
import { loadPromptTemplate, renderPrompt } from "./promptTemplate.js";
import { discoverSourceLists, type SourceList } from "./sourceLists.js";
import {
countStagedRows,
getStagedHeadwords,
openStaging,
stageEntry,
} from "./staging.js";
import { validateEntry, type ValidationContext } from "./validate.js";
const ROOT = import.meta.dirname;
const SOURCE_DATA_DIR = path.join(ROOT, "source-data");
const STAGING_DB_PATH = path.join(ROOT, "db", "staging.db");
const STAGING_SCHEMA_PATH = path.join(ROOT, "db", "schema.sql");
const RESPONSES_DIR = path.join(ROOT, "responses");
const REJECTIONS_DIR = path.join(ROOT, "rejections");
type CliOptions = {
langs: SupportedLanguageCode[] | null;
pos: SupportedPos;
maxBatches: number | null;
batchSize: number;
delayMs: number;
dryRun: boolean;
};
type ListStats = {
list: SourceList;
pending: number;
batchesRun: number;
staged: number;
skipped: number;
rejected: number;
failedBatches: number;
};
const parseCli = (): CliOptions => {
// pnpm forwards the "--" separator itself (pnpm pipeline:run -- --langs …);
// drop it so the flags after it are parsed as flags, not positionals.
const args = process.argv.slice(2);
if (args[0] === "--") args.shift();
const { values } = parseArgs({
args,
options: {
langs: { type: "string" },
pos: { type: "string", default: "noun" },
"max-batches": { type: "string" },
"batch-size": { type: "string", default: "20" },
"delay-ms": { type: "string", default: "6000" },
"dry-run": { type: "boolean", default: false },
},
});
const pos = values.pos as SupportedPos;
if (!(SUPPORTED_POS as readonly string[]).includes(pos)) {
throw new Error(`--pos must be one of ${SUPPORTED_POS.join(", ")}`);
}
let langs: SupportedLanguageCode[] | null = null;
if (values.langs !== undefined) {
langs = values.langs.split(",").map((code) => {
const trimmed = code.trim() as SupportedLanguageCode;
if (!(SUPPORTED_LANGUAGE_CODES as readonly string[]).includes(trimmed)) {
throw new Error(`--langs: unknown language code "${trimmed}"`);
}
return trimmed;
});
}
return {
langs,
pos,
maxBatches:
values["max-batches"] !== undefined
? Number(values["max-batches"])
: null,
batchSize: Number(values["batch-size"]),
delayMs: Number(values["delay-ms"]),
dryRun: values["dry-run"],
};
};
const chunk = <T>(items: readonly T[], size: number): T[][] => {
const chunks: T[][] = [];
for (let i = 0; i < items.length; i += size) {
chunks.push(items.slice(i, i + size));
}
return chunks;
};
const sleep = (ms: number): Promise<void> =>
new Promise((resolve) => setTimeout(resolve, ms));
const rejectionFile = (list: SourceList): string =>
path.join(REJECTIONS_DIR, `${list.sourceLanguage}-${list.pos}.jsonl`);
const logRejection = (
list: SourceList,
headword: string | null,
errors: string[],
entry: unknown,
): void => {
mkdirSync(REJECTIONS_DIR, { recursive: true });
const line = JSON.stringify({
at: new Date().toISOString(),
sourceLanguage: list.sourceLanguage,
pos: list.pos,
headword,
errors,
entry,
});
appendFileSync(rejectionFile(list), `${line}\n`);
};
const saveRawResponse = (
list: SourceList,
batchIndex: number,
model: string,
words: readonly string[],
targetLanguages: readonly SupportedLanguageCode[],
rawText: string,
): void => {
mkdirSync(RESPONSES_DIR, { recursive: true });
const stamp = new Date().toISOString().replaceAll(":", "-");
const file = path.join(
RESPONSES_DIR,
`${list.sourceLanguage}-${list.pos}-${stamp}-batch${batchIndex}.json`,
);
writeFileSync(
file,
JSON.stringify(
{
model,
sourceLanguage: list.sourceLanguage,
pos: list.pos,
targetLanguages,
words,
receivedAt: new Date().toISOString(),
rawText,
},
null,
2,
),
);
};
const headwordOf = (entry: unknown): string | null => {
if (typeof entry === "object" && entry !== null && !Array.isArray(entry)) {
const headword = (entry as Record<string, unknown>)["headword"];
if (typeof headword === "string") return headword;
}
return null;
};
const main = async (): Promise<void> => {
const options = parseCli();
const model = process.env["GEMINI_MODEL"] ?? DEFAULT_MODEL;
const apiKey = process.env["GEMINI_API_KEY"];
if (apiKey === undefined && !options.dryRun) {
throw new Error("GEMINI_API_KEY is not set (data-pipeline/.env)");
}
const allLists = discoverSourceLists(SOURCE_DATA_DIR);
console.log(`found ${allLists.length} source lists:\n`);
for (const list of allLists) {
console.log(
` ${list.sourceLanguage}: ${list.pos} (${list.words.length} unique words)`,
);
}
const lists = allLists.filter(
(list) =>
list.pos === options.pos &&
(options.langs === null || options.langs.includes(list.sourceLanguage)),
);
if (lists.length === 0) {
console.log("\nnothing matches the requested --langs/--pos, exiting");
return;
}
const template = loadPromptTemplate();
const db = openStaging(STAGING_DB_PATH, STAGING_SCHEMA_PATH);
const allStats: ListStats[] = [];
let firstApiCall = true;
for (const list of lists) {
const staged = getStagedHeadwords(db, list.sourceLanguage, list.pos);
const pending = list.words.filter((word) => !staged.has(word));
const targetLanguages = SUPPORTED_LANGUAGE_CODES.filter(
(code) => code !== list.sourceLanguage,
);
const batches = chunk(pending, options.batchSize).slice(
0,
options.maxBatches ?? Number.POSITIVE_INFINITY,
);
console.log(
`\n${list.sourceLanguage}/${list.pos}: ${list.words.length} unique, ${staged.size} already staged, ${pending.length} pending → running ${batches.length} batch(es)`,
);
const stats: ListStats = {
list,
pending: pending.length,
batchesRun: 0,
staged: 0,
skipped: 0,
rejected: 0,
failedBatches: 0,
};
allStats.push(stats);
for (const [batchIndex, words] of batches.entries()) {
const prompt = renderPrompt(template, {
sourceLanguage: list.sourceLanguage,
pos: list.pos,
targetLanguages,
words,
});
if (options.dryRun) {
console.log(
` [dry-run] batch ${batchIndex + 1}/${batches.length}: ${words.join(", ")}`,
);
stats.batchesRun++;
continue;
}
if (!firstApiCall) {
await sleep(options.delayMs);
}
firstApiCall = false;
try {
const rawText = await generateContent(
apiKey as string,
model,
prompt,
buildEntriesResponseSchema(
list.sourceLanguage,
list.pos,
targetLanguages,
),
);
saveRawResponse(
list,
batchIndex + 1,
model,
words,
targetLanguages,
rawText,
);
const entries = parseEntries(rawText);
const ctx: ValidationContext = {
sourceLanguage: list.sourceLanguage,
pos: list.pos,
targetLanguages,
inputWords: new Set(words),
};
const covered = new Set<string>();
for (const entry of entries) {
const result = validateEntry(entry, ctx);
if (result.status === "valid") {
stageEntry(db, result.entry);
covered.add(result.entry.headword);
stats.staged++;
} else if (result.status === "empty") {
covered.add(result.headword);
stats.skipped++;
console.log(
` skipped "${result.headword}" (no valid ${list.pos} senses)`,
);
} else {
const headword = headwordOf(entry);
if (headword !== null) covered.add(headword);
logRejection(list, headword, result.errors, entry);
stats.rejected++;
}
}
for (const word of words) {
if (!covered.has(word)) {
logRejection(list, word, ["missing from Gemini response"], null);
stats.rejected++;
}
}
stats.batchesRun++;
console.log(
` batch ${batchIndex + 1}/${batches.length} done — ${stats.staged} staged, ${stats.rejected} rejected, ${stats.skipped} skipped`,
);
} catch (error) {
stats.failedBatches++;
console.error(
` batch ${batchIndex + 1}/${batches.length} FAILED: ${error instanceof Error ? error.message : String(error)}`,
);
}
}
}
console.log("\n— summary —");
for (const stats of allStats) {
console.log(
`${stats.list.sourceLanguage}/${stats.list.pos}: ${stats.staged} staged, ${stats.skipped} skipped, ${stats.rejected} rejected, ${stats.failedBatches} failed batch(es), ${stats.pending - stats.staged - stats.skipped - stats.rejected} still pending`,
);
}
const totals = countStagedRows(db);
console.log(
`staging.db totals: ${totals.words} words, ${totals.senses} senses, ${totals.translations} translations`,
);
db.close();
};
main().catch((error: unknown) => {
console.error(error instanceof Error ? error.message : error);
process.exitCode = 1;
});