diff --git a/data-pipeline/package.json b/data-pipeline/package.json index 515057a..84d59d0 100644 --- a/data-pipeline/package.json +++ b/data-pipeline/package.json @@ -6,7 +6,8 @@ "scripts": { "test": "vitest run", "test:watch": "vitest", - "pipeline:run": "tsx --env-file .env pipeline.ts" + "pipeline:run": "tsx --env-file .env pipeline.ts", + "pipeline:replay": "tsx replay.ts" }, "dependencies": { "@lila/shared": "workspace:*", diff --git a/data-pipeline/pipeline.ts b/data-pipeline/pipeline.ts index baf40df..c8e25c2 100644 --- a/data-pipeline/pipeline.ts +++ b/data-pipeline/pipeline.ts @@ -46,6 +46,7 @@ type ListStats = { staged: number; skipped: number; rejected: number; + normalized: number; failedBatches: number; }; @@ -217,6 +218,7 @@ const main = async (): Promise => { staged: 0, skipped: 0, rejected: 0, + normalized: 0, failedBatches: 0, }; allStats.push(stats); @@ -276,6 +278,12 @@ const main = async (): Promise => { stageEntry(db, result.entry); covered.add(result.entry.headword); stats.staged++; + if (result.normalizations.length > 0) { + stats.normalized++; + for (const note of result.normalizations) { + console.log(` normalized "${result.entry.headword}" ${note}`); + } + } } else if (result.status === "empty") { covered.add(result.headword); stats.skipped++; @@ -311,7 +319,7 @@ const main = async (): Promise => { 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`, + `${stats.list.sourceLanguage}/${stats.list.pos}: ${stats.staged} staged (${stats.normalized} normalized), ${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); diff --git a/data-pipeline/replay.ts b/data-pipeline/replay.ts new file mode 100644 index 0000000..af74f65 --- /dev/null +++ b/data-pipeline/replay.ts @@ -0,0 +1,226 @@ +/** + * Re-validates saved Gemini responses from `responses/` without calling the API. + * + * Every batch response is written to disk before it is parsed, so a change to + * the validation rules can be applied retroactively to everything already + * generated. Use this after editing `validate.ts` to see what the change would + * do, and to recover entries that the old rules rejected. + * + * Reports by default; pass --write to stage recovered entries into staging.db. + * Staging is idempotent, so re-running is safe. + */ +import { readdirSync, readFileSync } 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 { parseEntries } from "./gemini.js"; +import { openStaging, stageEntry, countStagedRows } from "./staging.js"; +import { validateEntry, type ValidationContext } from "./validate.js"; + +const ROOT = import.meta.dirname; +const RESPONSES_DIR = path.join(ROOT, "responses"); +const STAGING_DB_PATH = path.join(ROOT, "db", "staging.db"); +const STAGING_SCHEMA_PATH = path.join(ROOT, "db", "schema.sql"); + +type SavedResponse = { + sourceLanguage: SupportedLanguageCode; + pos: SupportedPos; + targetLanguages: SupportedLanguageCode[]; + words: string[]; + rawText: string; +}; + +type Tally = { + files: number; + unreadable: number; + entries: number; + valid: number; + normalized: number; + empty: number; + invalid: number; + staged: number; + alreadyStaged: number; +}; + +const isRecord = (value: unknown): value is Record => + typeof value === "object" && value !== null && !Array.isArray(value); + +const isStringArray = (value: unknown): value is string[] => + Array.isArray(value) && value.every((item) => typeof item === "string"); + +/** Saved files are trusted-but-verified: they are ours, but they are on disk. */ +const parseSavedResponse = (raw: unknown): SavedResponse | null => { + if (!isRecord(raw)) return null; + const { sourceLanguage, pos, targetLanguages, words, rawText } = raw; + if ( + typeof sourceLanguage !== "string" || + !(SUPPORTED_LANGUAGE_CODES as readonly string[]).includes(sourceLanguage) || + typeof pos !== "string" || + !(SUPPORTED_POS as readonly string[]).includes(pos) || + !isStringArray(targetLanguages) || + !isStringArray(words) || + typeof rawText !== "string" + ) { + return null; + } + return { + sourceLanguage: sourceLanguage as SupportedLanguageCode, + pos: pos as SupportedPos, + targetLanguages: targetLanguages as SupportedLanguageCode[], + words, + rawText, + }; +}; + +const main = (): void => { + const { values } = parseArgs({ + args: + process.argv[2] === "--" ? process.argv.slice(3) : process.argv.slice(2), + options: { + write: { type: "boolean", default: false }, + langs: { type: "string" }, + verbose: { type: "boolean", default: false }, + }, + }); + + const langFilter = + values.langs === undefined + ? null + : new Set(values.langs.split(",").map((code) => code.trim())); + + let files: string[]; + try { + files = readdirSync(RESPONSES_DIR) + .filter((name) => name.endsWith(".json")) + .sort(); + } catch { + console.error(`no responses directory at ${RESPONSES_DIR}`); + process.exitCode = 1; + return; + } + + if (files.length === 0) { + console.log("no saved responses to replay"); + return; + } + + const db = values.write + ? openStaging(STAGING_DB_PATH, STAGING_SCHEMA_PATH) + : null; + + const tally: Tally = { + files: 0, + unreadable: 0, + entries: 0, + valid: 0, + normalized: 0, + empty: 0, + invalid: 0, + staged: 0, + alreadyStaged: 0, + }; + const recovered: string[] = []; + const stillInvalid = new Map(); + + for (const name of files) { + let saved: SavedResponse | null; + try { + saved = parseSavedResponse( + JSON.parse(readFileSync(path.join(RESPONSES_DIR, name), "utf-8")), + ); + } catch { + saved = null; + } + if (saved === null) { + tally.unreadable++; + console.warn(` skipping unreadable response file: ${name}`); + continue; + } + if (langFilter !== null && !langFilter.has(saved.sourceLanguage)) continue; + tally.files++; + + let entries: unknown[]; + try { + entries = parseEntries(saved.rawText); + } catch (error) { + tally.unreadable++; + console.warn( + ` skipping ${name}: ${error instanceof Error ? error.message : String(error)}`, + ); + continue; + } + + const ctx: ValidationContext = { + sourceLanguage: saved.sourceLanguage, + pos: saved.pos, + targetLanguages: saved.targetLanguages, + inputWords: new Set(saved.words), + }; + + for (const entry of entries) { + tally.entries++; + const result = validateEntry(entry, ctx); + if (result.status === "valid") { + tally.valid++; + if (result.normalizations.length > 0) { + tally.normalized++; + recovered.push( + `${saved.sourceLanguage}/${result.entry.headword}: ${result.normalizations.join("; ")}`, + ); + } + if (db !== null) { + const outcome = stageEntry(db, result.entry); + if (outcome === "staged") tally.staged++; + else tally.alreadyStaged++; + } + } else if (result.status === "empty") { + tally.empty++; + } else { + tally.invalid++; + const key = result.errors.join(" | "); + stillInvalid.set(key, [...(stillInvalid.get(key) ?? []), name]); + } + } + } + + if (values.verbose && recovered.length > 0) { + console.log("\n— normalized entries —"); + for (const line of recovered) console.log(` ${line}`); + } + + if (stillInvalid.size > 0) { + console.log("\n— still invalid —"); + for (const [errors, sources] of stillInvalid) { + console.log(` (${sources.length}×) ${errors}`); + } + } + + console.log("\n— replay summary —"); + console.log(` response files replayed: ${tally.files}`); + if (tally.unreadable > 0) + console.log(` unreadable files: ${tally.unreadable}`); + console.log(` entries seen: ${tally.entries}`); + console.log(` valid: ${tally.valid}`); + console.log(` of which normalized: ${tally.normalized}`); + console.log(` empty (skipped): ${tally.empty}`); + console.log(` invalid: ${tally.invalid}`); + + if (db !== null) { + console.log(` newly staged: ${tally.staged}`); + console.log(` already staged: ${tally.alreadyStaged}`); + const totals = countStagedRows(db); + console.log( + ` staging.db totals: ${totals.words} words, ${totals.senses} senses, ${totals.translations} translations`, + ); + db.close(); + } else { + console.log("\n (report only — pass --write to stage recovered entries)"); + } +}; + +main(); diff --git a/data-pipeline/tests/validate.test.ts b/data-pipeline/tests/validate.test.ts index 53512d8..c5cf88d 100644 --- a/data-pipeline/tests/validate.test.ts +++ b/data-pipeline/tests/validate.test.ts @@ -1,5 +1,9 @@ import { describe, it, expect } from "vitest"; -import { validateEntry, type ValidationContext } from "../validate.js"; +import { + validateEntry, + type ValidationContext, + type ValidationResult, +} from "../validate.js"; const ctx: ValidationContext = { sourceLanguage: "es", @@ -253,21 +257,6 @@ describe("validateEntry", () => { ); }); - it("rejects a translation difficulty below the sense difficulty", () => { - const easyTranslations = translations(); - expect( - errorsOf( - entry({ - senses: [ - sense({ difficulty: "medium", translations: easyTranslations }), - ], - }), - ), - ).toContainEqual( - expect.stringContaining('is lower than sense difficulty "medium"'), - ); - }); - it("allows a translation difficulty above the sense difficulty", () => { const harder = translations(); harder[2] = { @@ -282,4 +271,97 @@ describe("validateEntry", () => { ); expect(result.status).toBe("valid"); }); + + describe("sense difficulty floor", () => { + const senseDifficultyOf = (result: ValidationResult): string => + result.status === "valid" + ? (result.entry.senses[0]?.difficulty ?? "") + : ""; + + it("floors a sense to its easiest translation instead of rejecting", () => { + // The real "Ellbogen" case: sense tagged medium, but some translations + // are correctly easy. The entry is good; only the derived label is wrong. + const mixed = translations(); + mixed[1] = { + target_language: "it", + word: "gomito", + gender: "masculine", + difficulty: "medium", + }; + const result = validateEntry( + entry({ + senses: [sense({ difficulty: "medium", translations: mixed })], + }), + ctx, + ); + expect(result.status).toBe("valid"); + expect(senseDifficultyOf(result)).toBe("easy"); + }); + + it("reports what it changed", () => { + const result = validateEntry( + entry({ senses: [sense({ difficulty: "hard" })] }), + ctx, + ); + expect(result.status).toBe("valid"); + if (result.status !== "valid") return; + expect(result.normalizations).toHaveLength(1); + expect(result.normalizations[0]).toContain('"hard" → "easy"'); + }); + + it("leaves an already-consistent sense untouched", () => { + const result = validateEntry(entry(), ctx); + expect(result.status).toBe("valid"); + if (result.status !== "valid") return; + expect(result.normalizations).toEqual([]); + expect(senseDifficultyOf(result)).toBe("easy"); + }); + + it("never raises a sense above its own label", () => { + // All translations medium, sense easy: the concept stays easy, because + // sense difficulty gates the meaning, not the word (design-doc §4). + const allMedium = translations().map((t) => ({ + ...t, + difficulty: "medium", + })); + const result = validateEntry( + entry({ senses: [sense({ translations: allMedium })] }), + ctx, + ); + expect(result.status).toBe("valid"); + if (result.status !== "valid") return; + expect(result.normalizations).toEqual([]); + expect(senseDifficultyOf(result)).toBe("easy"); + }); + + it("floors each sense independently", () => { + const easySense = sense({ sense_index: 0, difficulty: "medium" }); + const hardSense = sense({ + sense_index: 1, + difficulty: "hard", + translations: translations().map((t) => ({ + ...t, + difficulty: "medium", + })), + }); + const result = validateEntry( + entry({ senses: [easySense, hardSense] }), + ctx, + ); + expect(result.status).toBe("valid"); + if (result.status !== "valid") return; + expect(result.entry.senses.map((s) => s.difficulty)).toEqual([ + "easy", + "medium", + ]); + expect(result.normalizations).toHaveLength(2); + }); + + it("does not mutate the input entry", () => { + const input = entry({ senses: [sense({ difficulty: "medium" })] }); + validateEntry(input, ctx); + const senses = input["senses"] as Record[]; + expect(senses[0]?.["difficulty"]).toBe("medium"); + }); + }); }); diff --git a/data-pipeline/validate.ts b/data-pipeline/validate.ts index 75d5f04..8f84a94 100644 --- a/data-pipeline/validate.ts +++ b/data-pipeline/validate.ts @@ -39,9 +39,14 @@ export type ValidationContext = { /** * "empty" is the contract's way of saying "not a valid word of this POS" * (senses: []) — the word is skipped, not rejected. + * + * A "valid" result carries the *normalized* entry, which may differ from the + * input — see `applySenseDifficultyFloor`. `normalizations` is a human-readable + * record of any repair applied, so the pipeline can report how often the model + * needed correcting instead of silently papering over it. */ export type ValidationResult = - | { status: "valid"; entry: GeminiWordEntry } + | { status: "valid"; entry: GeminiWordEntry; normalizations: string[] } | { status: "empty"; headword: string } | { status: "invalid"; errors: string[] }; @@ -82,7 +87,6 @@ const checkStringArray = ( const checkTranslation = ( raw: unknown, label: string, - senseDifficulty: DifficultyLevel | null, ctx: ValidationContext, errors: string[], ): void => { @@ -99,18 +103,12 @@ const checkTranslation = ( if (!isNonEmptyString(raw["word"])) { errors.push(`${label}: word must be a non-empty string`); } - const difficulty = raw["difficulty"]; - if (!isDifficulty(difficulty)) { + // A translation ranked below its sense is not an error — the sense is the + // derived value and gets floored to match. See applySenseDifficultyFloor. + if (!isDifficulty(raw["difficulty"])) { errors.push( `${label}: difficulty must be one of ${DIFFICULTY_LEVELS.join(", ")}`, ); - } else if ( - senseDifficulty !== null && - difficultyRank(difficulty) < difficultyRank(senseDifficulty) - ) { - errors.push( - `${label}: difficulty "${difficulty}" is lower than sense difficulty "${senseDifficulty}"`, - ); } if ( typeof target === "string" && @@ -157,13 +155,7 @@ const checkSense = ( return; } translations.forEach((translation, i) => { - checkTranslation( - translation, - `${label}.translations[${i}]`, - senseDifficulty, - ctx, - errors, - ); + checkTranslation(translation, `${label}.translations[${i}]`, ctx, errors); }); const wordsPerTarget = new Map(); @@ -196,6 +188,45 @@ const checkSense = ( } }; +/** + * Lowers a sense's difficulty to that of its easiest translation when the model + * tagged the sense higher. + * + * Why this is a repair and not a rejection: design-doc §5.1 filters sense + * difficulty as a *ceiling* and translation difficulty as an *exact* target, so + * a translation ranked below its own sense can never be served — the sense is + * gated out at exactly the level where that translation would be the answer. + * The row is dead data. The prompt already defines sense difficulty as the + * easiest translation difficulty in the sense, which makes it a derived value + * rather than an independent judgement, so it is recomputed here instead of + * discarding an otherwise-good entry. + * + * The floor only ever lowers. Raising a sense to match its translations would + * gate a concept out of levels it belongs in, and would collapse the + * concept-vs-word distinction that the two difficulty columns exist to express + * (design-doc §4). + */ +const applySenseDifficultyFloor = ( + entry: GeminiWordEntry, + normalizations: string[], +): GeminiWordEntry => ({ + ...entry, + senses: entry.senses.map((sense) => { + const easiest = sense.translations.reduce( + (lowest, translation) => + difficultyRank(translation.difficulty) < difficultyRank(lowest) + ? translation.difficulty + : lowest, + sense.difficulty, + ); + if (easiest === sense.difficulty) return sense; + normalizations.push( + `senses[${sense.sense_index}]: difficulty "${sense.difficulty}" → "${easiest}" (floored to easiest translation)`, + ); + return { ...sense, difficulty: easiest }; + }), +}); + export const validateEntry = ( raw: unknown, ctx: ValidationContext, @@ -235,5 +266,10 @@ export const validateEntry = ( if (errors.length > 0) { return { status: "invalid", errors }; } - return { status: "valid", entry: raw as GeminiWordEntry }; + const normalizations: string[] = []; + return { + status: "valid", + entry: applySenseDifficultyFloor(raw as GeminiWordEntry, normalizations), + normalizations, + }; }; diff --git a/documentation/DATA_PIPELINE.md b/documentation/DATA_PIPELINE.md index 02b07eb..2c55267 100644 --- a/documentation/DATA_PIPELINE.md +++ b/documentation/DATA_PIPELINE.md @@ -70,6 +70,7 @@ co-located unit tests in `data-pipeline/tests/`. | `data-pipeline/validate.ts` | ✅ Per-entry validation (design-doc §6.4) | | `data-pipeline/staging.ts` | ✅ SQLite writes, one transaction per word | | `data-pipeline/pipeline.ts` | ✅ Orchestrator with CLI flags, resumable | +| `data-pipeline/replay.ts` | ✅ Re-validates saved responses without API calls | | `data-pipeline/db/schema.sql` | ✅ SQLite staging schema | | `data-pipeline/db/staging.db` | ✅ Populated — gitignored | | `data-pipeline/responses/` | ✅ Raw Gemini responses, one JSON per batch — gitignored | @@ -91,8 +92,11 @@ gemini.ts buildEntriesResponseSchema(): OpenAPI schema with enums narro generateContent(): 5 attempts, retries 429/500/503, honours the API's own retryDelay; rejects non-STOP finishReason validate.ts validateEntry() → "valid" | "empty" | "invalid" + applySenseDifficultyFloor(): lowers a sense to its easiest + translation, reported via the result's `normalizations` staging.ts openStaging(), getStagedHeadwords(), stageEntry(), countStagedRows() pipeline.ts orchestration, batching, rate-limit delay, rejection logging +replay.ts re-validates responses/ with the current rules, no API calls ``` Two properties worth knowing: @@ -155,12 +159,14 @@ actually in the input batch. Validation is the safety net, not the prompt — every entry is checked before it reaches SQLite, and rejects go to a log for review rather than silently disappearing. Target reject rate is under 10%. -> **Known systemic rejection cause (open).** Effectively all current rejections are -> `translation difficulty lower than sense difficulty`. The model tends to tag a sense -> `medium` while correctly tagging some of its translations `easy`. Since the prompt already -> defines sense difficulty as the easiest translation difficulty in the sense, this value is -> derivable and the entry is otherwise good — normalizing it in `validate.ts` instead of -> rejecting would recover these words. Not yet implemented. +> **Systemic rejection cause — resolved.** Effectively every rejection used to be +> `translation difficulty lower than sense difficulty`: the model tags a sense `medium` +> while correctly tagging some of its translations `easy`. The prompt defines sense +> difficulty as the easiest translation difficulty in the sense, which makes it a derived +> value, so `validate.ts` now floors it (`applySenseDifficultyFloor`) instead of rejecting +> the entry. The floor only ever lowers — raising a sense would gate a concept out of levels +> it belongs in and collapse the concept-vs-word split of design-doc §4. Replaying every +> saved response with the new rule took the reject rate from 20 entries to **zero**. --- @@ -191,6 +197,40 @@ pnpm --filter @lila/pipeline pipeline:run -- --langs de --max-batches 1 --dry-ru Because runs are resumable, interrupting with Ctrl-C is safe — already-staged words are skipped on the next run. +### Replaying saved responses + +`replay.ts` re-validates everything in `responses/` **without calling the API**, which is how +a validation-rule change is applied to data that has already been generated. + +```bash +pnpm --filter @lila/pipeline pipeline:replay # report only +pnpm --filter @lila/pipeline pipeline:replay -- --write # stage recovered entries +pnpm --filter @lila/pipeline pipeline:replay -- --verbose --langs de +``` + +It reports valid / normalized / empty / invalid counts and groups anything still invalid by +error. Staging is idempotent, so `--write` is safe to repeat. One limitation: it can only +_insert_. A word already in `staging.db` is left untouched, so a rule change cannot repair +rows that were staged under the old rules — only recover ones that were rejected. + +### API quota + +⚠️ **The free tier allows far fewer requests than the batch count needs.** Observed +2026-08-20: `gemini-3.6-flash` returned `HTTP 429 … limit: 20` on metric +`generate_content_free_tier_requests` after **17 successful requests in one day**, and did +not recover across ~55 minutes of retrying — a daily window, not a per-minute one. The +`retryDelay: ~59s` in the 429 payload is generic backoff advice and misleads the retry logic +into grinding. + +Consequences to plan around: + +- What is rationed is **requests**, not words, so `--batch-size` is the cheap lever. The + remaining ~7,300 words are ~365 requests at batch size 20, but only ~74 at batch size 100. +- `gemini.ts` currently treats 429 as retryable and burns all 5 `MAX_ATTEMPTS` against a + quota that will not clear for hours. It should distinguish rate-limiting from daily + exhaustion and abort the run. +- Replays cost nothing, which is why raw responses are persisted. + The pipeline reads `.env` from the repo root: `GEMINI_API_KEY`, plus `PIPELINE_POSTGRES_USER` / `PIPELINE_POSTGRES_PASSWORD` / `PIPELINE_POSTGRES_DB` / `PIPELINE_DATABASE_URL`. The pipeline database is deliberately separate from the app database (`:5432`) so pipeline work can never damage dev data. --- diff --git a/documentation/pipeline/roadmap.md b/documentation/pipeline/roadmap.md index 43e4e38..52b4a4d 100644 --- a/documentation/pipeline/roadmap.md +++ b/documentation/pipeline/roadmap.md @@ -176,8 +176,9 @@ Re-runs skip words that are already staged. `--dry-run`. - [ ] Run the pipeline for all 5 languages (nouns only) — **in progress**, German first -- [ ] Review the rejection log, fix prompt issues, re-run failed batches - — reviewed; one systemic cause found (see below), fix not yet applied +- [x] Review the rejection log, fix prompt issues, re-run failed batches — + one systemic cause found and fixed in `validate.ts`; recovered from + `responses/` via `replay.ts`, no re-run needed. Reject rate now zero. - [ ] Spot-check 50 random entries for correctness **Dependencies:** Phase 2 complete. @@ -199,15 +200,25 @@ resumed correctly from previously staged words), raw responses are on disk, and the reject rate is tracking well under 10%. The two open items are the systemic rejection cause and the hard-tier shortfall below. -**Open issue — systemic rejection cause.** Effectively every rejection is +**Resolved — systemic rejection cause.** Every rejection was `translation difficulty lower than sense difficulty`: the model tags a sense `medium` while correctly tagging some of its translations `easy`. Example: `Ellbogen` with sense `medium` but `elbow` (en) and `codo` (es) as `easy` — the translations are right and the sense label is wrong, yet the whole entry -is discarded. The prompt already defines sense difficulty as the easiest -translation difficulty in the sense, so the value is derivable. Normalizing it -in `validate.ts` rather than rejecting would recover these entries and can be -replayed against `responses/` without new API calls. +was discarded. Since the prompt defines sense difficulty as the easiest +translation difficulty in the sense, the value is derived, so `validate.ts` +now floors it (`applySenseDifficultyFloor`) rather than rejecting. The floor +only lowers; raising a sense would gate a concept out of levels it belongs in. +Replaying all 38 saved responses took the reject rate from 20 entries to zero, +and `staging.db` now holds 746 words / 774 senses / 3,436 translations with no +sense ranked above its easiest translation. + +**Open issue — API quota is the binding constraint.** The free tier allows +~20 requests/day for `gemini-3.6-flash` (observed 2026-08-20), not the ~1,000 +previously assumed. At `--batch-size 20` the remaining ~7,300 words are ~365 +requests, i.e. roughly 18 days of waiting. Requests are what is rationed, not +words, so raising `--batch-size` is the cheap lever (~74 requests at 100). +Paid billing is the alternative. Decide before scheduling the rest of the run. **Open issue — the `hard` tier is nearly empty.** Generated difficulty skews heavily easy; `hard` translations are well under 1% of staged rows. Because