diff --git a/scripts/import-benchmark-data.test.ts b/scripts/import-benchmark-data.test.ts new file mode 100644 index 00000000..2d86d7f9 --- /dev/null +++ b/scripts/import-benchmark-data.test.ts @@ -0,0 +1,298 @@ +import { describe, it, expect } from "vitest"; +import { chmod, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + commitWrites, + planImport, + type Corpus, + type WriteFs, +} from "./import-benchmark-data.ts"; +import { Report } from "./import-pathway-data.ts"; +import type { PathwayMetadataV2 } from "../src/types/pathwayMetadata.v2.d.ts"; +// A real publication block: publishers are a closed list in the schema. +import ieaSteps from "../src/data/iea/IEA-STEPS-2024_timeseries.json" with { type: "json" }; + +const V2 = "http://pathways.rmi.org/schema/pathwayTimeseries.v2.json"; + +const row = (over: Record = {}) => ({ + year: 2030, + geography: "Southeast Asia", + sector: "power", + sectorSegment: ["Power generation"], + technology: null, + metric: "capacity", + value: 1, + unit: "GW", + ...over, +}); + +const file = (over: Record = {}, rows = [row()]) => ({ + $schema: V2, + id: "IEA-X_timeseries", + pathwayId: ["IEA-X"], + name: "IEA X Timeseries Data", + description: "A test series.", + publication: ieaSteps.publication, + pathwayName: "X", + emissionsScope: "CO2", + data: rows, + ...over, +}); + +const corpus = (existing = false): Corpus => ({ + metadataGeographyById: new Map([ + ["IEA-X", { global: true, regions: { "Southeast Asia": ["TH", "VN"] } }], + ]), + metadataPathById: new Map([["IEA-X", "src/data/iea/IEA-X.json"]]), + occupiedPaths: new Set([ + "src/data/iea/IEA-X.json", + ...(existing + ? [ + "src/data/iea/X-legacy-name_timeseries.json", + "src/data/iea/OTHER.json", + ] + : []), + ]), + timeseriesById: new Map( + existing + ? [ + [ + "IEA-X_timeseries", + { path: "src/data/iea/X-legacy-name_timeseries.json", doc: {} }, + ], + ["OTHER_timeseries", { path: "src/data/iea/OTHER.json", doc: {} }], + ] + : [], + ), +}); + +const plan = (inputs: unknown[], c = corpus()) => { + const report = new Report(); + const writes = planImport( + inputs.map((data, i) => ({ name: `in/${i}.json`, data })), + c, + report, + ); + return { writes, report }; +}; + +describe("planImport", () => { + it("places a new file next to its pathway's metadata", () => { + const { writes, report } = plan([file()]); + expect(report.errors).toEqual([]); + expect(writes.map((w) => [w.path, w.replaces])).toEqual([ + ["src/data/iea/IEA-X_timeseries.json", false], + ]); + }); + + it("overwrites an existing file in place, keeping its name", () => { + const { writes, report } = plan([file()], corpus(true)); + expect(writes.map((w) => w.path)).toEqual([ + "src/data/iea/X-legacy-name_timeseries.json", + ]); + // ...and says which existing files the import leaves alone. + expect(report.lines.join("\n")).toContain("left as they are: OTHER.json"); + }); + + it("never lets a new id land on an existing file with another id", () => { + // ACE's real shape: the file name drops the id's prefix, so an input whose + // id happens to equal that name is new by id but not by path. + const { writes, report } = plan( + [file({ id: "X-legacy-name_timeseries" })], + corpus(true), + ); + expect(report.errors.join("\n")).toMatch( + /X-legacy-name_timeseries\.json is already taken/, + ); + expect(writes).toEqual([]); + }); + + it("treats a path differing only by case as taken", () => { + // On default macOS and Windows file systems these are one file. + const { writes, report } = plan([file({ id: "iea-x" })]); + expect(report.errors.join("\n")).toMatch( + /src\/data\/iea\/iea-x\.json is already taken/, + ); + expect(writes).toEqual([]); + }); + + it("reports a taken destination in the same run as the file's other errors", () => { + const { report } = plan( + [ + file({ id: "X-legacy-name_timeseries" }, [ + row({ geography: "Atlantis" }), + ]), + ], + corpus(true), + ); + const errors = report.errors.join("\n"); + expect(errors).toMatch(/"Atlantis" is not a geography/); + expect(errors).toMatch(/is already taken/); + }); + + it.each([ + [ + "an undeclared geography", + [row({ geography: "South East Asia" })], + /not a geography pathway IEA-X declares/, + ], + [ + "a segment of another sector", + [row({ sectorSegment: ["Ironmaking"] })], + /not a segment of Power/, + ], + [ + "a sentinel segment", + [row({ sectorSegment: ["Unspecified"] })], + /is not allowed/, + ], + [ + "a v1-shaped row with no sectorSegment", + [{ ...row(), sectorSegment: undefined }], + /sectorSegment/, + ], + ])("blocks the import on %s", (_, rows, message) => { + const { writes, report } = plan([file({}, rows)]); + expect(report.errors.join("\n")).toMatch(message); + expect(writes).toEqual([]); + }); + + it("blocks a file whose pathway has no metadata", () => { + const { report } = plan([file({ pathwayId: ["GONE"] })]); + expect(report.errors.join("\n")).toMatch( + /"GONE" is not the id of any pathway/, + ); + }); + + it("blocks a file still on the v1 schema", () => { + const { report } = plan([ + file({ + $schema: "http://pathways.rmi.org/schema/pathwayTimeseries.v1.json", + }), + ]); + expect(report.errors.length).toBeGreaterThan(0); + }); + + it("blocks two input files with the same id, and still checks the second", () => { + // One run should list every problem, so the duplicate's own errors are + // reported alongside the duplicate id. + const { writes, report } = plan([ + file(), + file({}, [row({ geography: "Atlantis" })]), + ]); + const errors = report.errors.join("\n"); + expect(errors).toMatch(/also used by in\/0\.json/); + expect(errors).toMatch(/"Atlantis" is not a geography/); + expect(writes.map((w) => w.path)).toEqual([ + "src/data/iea/IEA-X_timeseries.json", + ]); + }); + + it.each(["../escape", "iea/nested", ".hidden"])( + "blocks an id that is not a safe file name: %s", + (id) => { + const { writes, report } = plan([file({ id })]); + expect(report.errors.join("\n")).toMatch(/cannot be used as a file name/); + expect(writes).toEqual([]); + }, + ); +}); + +describe("commitWrites", () => { + /** An in-memory file system whose writes fail for chosen paths. */ + const memoryFs = (files: Record, failOn: string[] = []) => { + const failing = new Set(failOn); // each fails once; the restore succeeds + const io: WriteFs = { + readFile: async (path) => files[path] ?? null, + writeFile: async (path, text) => { + if (failing.delete(path)) { + // Like ENOSPC: the file is truncated before the write gives up. + files[path] = text.slice(0, 3); + throw new Error(`disk full at ${path}`); + } + files[path] = text; + }, + removeFile: async (path) => { + delete files[path]; + }, + }; + return { files, io }; + }; + + it("writes every staged file", async () => { + const { files, io } = memoryFs({ a: "old a" }); + await commitWrites( + [ + { path: "a", text: "new a" }, + { path: "b", text: "new b" }, + ], + io, + ); + expect(files).toEqual({ a: "new a", b: "new b" }); + }); + + it("restores earlier files and removes new ones when a later write fails", async () => { + const { files, io } = memoryFs({ a: "old a", c: "old c" }, ["c"]); + await expect( + commitWrites( + [ + { path: "a", text: "new a" }, + { path: "b", text: "new b" }, + { path: "c", text: "new c" }, + ], + io, + ), + ).rejects.toThrow(/the 3 file\(s\) touched were restored/); + // c included: its failed write had already truncated it. + expect(files).toEqual({ a: "old a", c: "old c" }); + }); + + it("restores the file whose own write failed, even when it is the first", async () => { + const { files, io } = memoryFs({ a: "old a" }, ["a"]); + await expect( + commitWrites([{ path: "a", text: "new a" }], io), + ).rejects.toThrow(); + expect(files).toEqual({ a: "old a" }); + }); + + it("writes nothing when an existing file cannot be backed up", async () => { + const { files, io } = memoryFs({ a: "old a" }); + io.readFile = async (path) => { + if (path === "b") throw new Error("permission denied"); + return files[path] ?? null; + }; + await expect( + commitWrites( + [ + { path: "a", text: "new a" }, + { path: "b", text: "new b" }, + ], + io, + ), + ).rejects.toThrow(/permission denied/); + expect(files).toEqual({ a: "old a" }); + }); + + it("does not delete an existing file it could not back up", async () => { + // A write-only file exists but cannot be read (EACCES). Counting that as + // "new" would make a rollback delete it; it must abort before writing. + const dir = await mkdtemp(join(tmpdir(), "import-benchmark-")); + const locked = join(dir, "locked.json"); + try { + await writeFile(locked, "old"); + await chmod(locked, 0o222); + await expect( + commitWrites([ + { path: locked, text: "new" }, + // Missing on read (so "new"), and its write fails: no such folder. + { path: join(dir, "no-such-folder", "x.json"), text: "{}" }, + ]), + ).rejects.toThrow(/EACCES/); + await chmod(locked, 0o644); + expect(await readFile(locked, "utf8")).toBe("old"); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); +}); diff --git a/scripts/import-benchmark-data.ts b/scripts/import-benchmark-data.ts new file mode 100644 index 00000000..5566b0e4 --- /dev/null +++ b/scripts/import-benchmark-data.ts @@ -0,0 +1,367 @@ +/** + * Import benchmark timeseries prepared by RMI/tpr_benchmark_data_preparation + * into src/data (#915, #902). + * + * Usage: + * npx ts-node --esm scripts/import-benchmark-data.ts --in [--dry-run] + * + * `` holds the prep repo's output: one JSON file per timeseries dataset, + * already in `pathwayTimeseries.v2.json` format. This script does not convert + * anything. It is the gate between the two repos: every file must pass the v2 + * schema and the cross-file checks in src/utils/validateTimeseries.ts (geography + * declared by the pathway's metadata, segments belonging to the row's sector) + * against the metadata already in src/data. + * + * Where each file goes: + * - an existing src/data timeseries with the same `id` is overwritten in + * place, keeping its file name (some, like ACE's ATS-2024_timeseries.json, + * do not follow the id); + * - otherwise a new `.json` is written next to the metadata of the + * file's first pathway (ids already end in `_timeseries`). + * Existing timeseries absent from the input are reported, never deleted. + * + * Like import-pathway-data.ts, any error blocks the whole import: nothing is + * written, every error is listed, and the run exits 1. A half-imported data + * set is worse than none. With --dry-run it reports and writes nothing either + * way, which is how the prep repo can check its output against TPR's contract. + */ +import { promises as fs } from "node:fs"; +import { basename, dirname, join, normalize } from "node:path"; +import { fileURLToPath } from "node:url"; +import * as prettier from "prettier"; +import { Report } from "./import-pathway-data.ts"; +import { validateFilesBySchema } from "../src/utils/validateData.ts"; +import { + PATHWAY_TIMESERIES_V2_ID, + validateTimeseries, + type MetadataGeographyById, +} from "../src/utils/validateTimeseries.ts"; +import { PATHWAY_METADATA_V2_ID } from "../src/utils/validateScopes.ts"; +import pathwayMetadataV1Schema from "../src/schema/pathwayMetadata.v1.json" with { type: "json" }; +import pathwayMetadataV2Schema from "../src/schema/pathwayMetadata.v2.json" with { type: "json" }; +import pathwayTimeseriesV2Schema from "../src/schema/pathwayTimeseries.v2.json" with { type: "json" }; +import { commonSchemas } from "../src/schema/common/index.ts"; +import type { PathwayMetadataV2 } from "../src/types/pathwayMetadata.v2.d.ts"; +import type { PathwayTimeseriesV2 } from "../src/types/pathwayTimeseries.v2.d.ts"; + +const DATA_DIR = "src/data"; +const METADATA_IDS = new Set([ + String(pathwayMetadataV1Schema.$id), + PATHWAY_METADATA_V2_ID, +]); +const TIMESERIES_ID_PATTERN = /pathwayTimeseries\.v\d+\.json$/; +const SAFE_FILE_ID = /^[A-Za-z0-9][A-Za-z0-9._-]*$/; + +type Json = Record; + +/** What src/data already holds, as far as an import needs to know. */ +export interface Corpus { + /** Each pathway's declared geography, by pathway id. */ + metadataGeographyById: MetadataGeographyById; + /** Where each pathway's metadata file lives, by pathway id. */ + metadataPathById: ReadonlyMap; + /** Existing timeseries files, by dataset `id`. */ + timeseriesById: ReadonlyMap; + /** + * Every file path under src/data. A new file must not land on one: file + * names do not always follow ids (ACE's `ATS-2024_timeseries.json` holds + * `ACE-ATS-2024_timeseries`), so a free path cannot be inferred from ids. + */ + occupiedPaths: ReadonlySet; +} + +export interface InputFile { + name: string; + data: unknown; +} + +export interface PlannedWrite { + path: string; + doc: PathwayTimeseriesV2; + replaces: boolean; +} + +const distinct = (values: Iterable) => [...new Set(values)].sort(); + +/** A path as a case-insensitive file system sees it. */ +const pathKey = (path: string) => normalize(path).toLowerCase(); + +function describe(doc: { data?: unknown }): string { + const rows = Array.isArray(doc.data) ? (doc.data as Json[]) : []; + const geographies = distinct(rows.map((r) => String(r.geography))); + const segments = distinct( + rows.flatMap((r) => + Array.isArray(r.sectorSegment) ? (r.sectorSegment as string[]) : [], + ), + ); + return ( + `${rows.length} rows; geography ${geographies.join(", ") || "-"}` + + (segments.length ? `; segments ${segments.join(", ")}` : "") + ); +} + +/** + * Validate the input files and decide where each goes. Pure apart from the + * report: errors are recorded there, and the caller writes only when there + * are none. + */ +export function planImport( + inputs: readonly InputFile[], + corpus: Corpus, + report: Report, +): PlannedWrite[] { + const { valid, invalid } = validateFilesBySchema( + [...inputs], + [ + pathwayMetadataV1Schema, + pathwayMetadataV2Schema, + pathwayTimeseriesV2Schema, + ...commonSchemas, + ], + ); + for (const problem of invalid) + for (const e of problem.errors) report.error(`${problem.name}: ${e}`); + + const planned: PlannedWrite[] = []; + // Paths compared ignoring case: on the default macOS and Windows file + // systems `iea-x.json` and `IEA-X.json` are one file. + const occupiedKeys = new Set([...corpus.occupiedPaths].map(pathKey)); + const plannedKeys = new Set(); + const seenIds = new Map(); + for (const record of valid) { + if (record.schemaId !== PATHWAY_TIMESERIES_V2_ID) { + report.error( + `${record.name}: $schema is ${record.schemaId}, expected ${PATHWAY_TIMESERIES_V2_ID}`, + ); + continue; + } + const doc = record.data as PathwayTimeseriesV2; + // Every check runs before the file is skipped, so one run lists every + // problem rather than the first one found. + const problems = validateTimeseries(doc, corpus.metadataGeographyById); + const earlier = seenIds.get(doc.id); + if (earlier) problems.push(`id "${doc.id}" is also used by ${earlier}`); + else seenIds.set(doc.id, record.name); + // A new file is named after its id, so the id must not be able to point + // anywhere else (no "/" or a leading "."). + if (!SAFE_FILE_ID.test(doc.id)) + problems.push( + `id "${doc.id}" cannot be used as a file name: use letters, digits, ".", "_" and "-", starting with a letter or digit`, + ); + // Where the file would go, worked out even when the checks above failed, + // so a placement problem is reported in the same run rather than the next. + const existing = corpus.timeseriesById.get(doc.id); + let path = existing?.path; + if (!existing) { + const metadataPath = corpus.metadataPathById.get(doc.pathwayId[0]); + if (metadataPath === undefined) { + // validateTimeseries reports unknown pathway ids; this only guards + // against that ever changing. + if (problems.length === 0) + problems.push(`no metadata file for ${doc.pathwayId[0]}`); + } else if (SAFE_FILE_ID.test(doc.id)) { + path = join(dirname(metadataPath), `${doc.id}.json`); + const key = pathKey(path); + if (occupiedKeys.has(key) || plannedKeys.has(key)) + problems.push( + `id "${doc.id}" is new, but ${path} is already taken by another file ` + + "(names compared ignoring case); an update must use the existing file's id", + ); + plannedKeys.add(key); + } + } + + for (const e of problems) report.error(`${record.name}: ${e}`); + if (problems.length > 0 || path === undefined) continue; + planned.push({ path, doc, replaces: existing !== undefined }); + report.note( + existing + ? `UPD ${path}: ${describe(existing.doc)} -> ${describe(doc)}` + : `NEW ${path}: ${describe(doc)}`, + ); + } + + const imported = new Set(seenIds.keys()); + const untouched = [...corpus.timeseriesById] + .filter(([id]) => !imported.has(id)) + .map(([, { path }]) => basename(path)); + if (untouched.length) + report.note( + `Not in this import, left as they are: ${untouched.sort().join(", ")}`, + ); + return planned; +} + +/** The file operations {@link commitWrites} needs; injectable for tests. */ +export interface WriteFs { + readFile(path: string): Promise; + writeFile(path: string, text: string): Promise; + removeFile(path: string): Promise; +} + +const nodeFs: WriteFs = { + readFile: async (path) => { + try { + return await fs.readFile(path, "utf8"); + } catch (error) { + // Only a missing file is "new". Any other failure (an unreadable but + // writable file, say) must stop the import: without a backup, a + // rollback would delete the file instead of restoring it. + if ((error as NodeJS.ErrnoException).code === "ENOENT") return null; + throw error; + } + }, + writeFile: (path, text) => fs.writeFile(path, text), + removeFile: (path) => fs.rm(path, { force: true }), +}; + +/** + * Write every staged file, or none: if any write fails, every file touched so + * far — including the one whose write failed, which may already be truncated + * — is restored to what it held before (or removed, if new), and the error is + * rethrown. This is what makes the import's all-or-nothing promise hold past + * validation, through to the disk. A restore that itself fails does not stop + * the others; the error names the files left as they are. + */ +export async function commitWrites( + staged: readonly { path: string; text: string }[], + io: WriteFs = nodeFs, +): Promise { + // Read every backup before writing anything; a read failure aborts here, + // with nothing touched yet. + const originals = await Promise.all( + staged.map(async ({ path }) => ({ path, text: await io.readFile(path) })), + ); + const touched: number[] = []; + try { + for (const [i, { path, text }] of staged.entries()) { + touched.push(i); // before the write: a failed write may still have changed the file + await io.writeFile(path, text); + } + } catch (error) { + const unrestored: string[] = []; + for (const i of touched.reverse()) { + const { path, text } = originals[i]; + try { + if (text === null) await io.removeFile(path); + else await io.writeFile(path, text); + } catch { + unrestored.push(path); + } + } + const outcome = unrestored.length + ? `could NOT restore ${unrestored.join(", ")}; check them with git` + : `the ${touched.length} file(s) touched were restored`; + throw new Error(`writing failed, and ${outcome}: ${String(error)}`); + } +} + +async function jsonFilesUnder(dir: string): Promise { + const out: string[] = []; + for (const d of await fs.readdir(dir, { withFileTypes: true })) { + const full = join(dir, d.name); + if (d.isDirectory()) out.push(...(await jsonFilesUnder(full))); + else if (d.name.endsWith(".json")) out.push(full); + } + return out; +} + +/** Read src/data into the shape {@link planImport} needs. */ +export async function readCorpus(dir = DATA_DIR): Promise { + const metadataGeographyById = new Map< + string, + PathwayMetadataV2["geography"] + >(); + const metadataPathById = new Map(); + const timeseriesById = new Map(); + const occupiedPaths = new Set(); + for (const path of await jsonFilesUnder(dir)) { + occupiedPaths.add(path); + let doc: Json; + try { + doc = JSON.parse(await fs.readFile(path, "utf8")) as Json; + } catch { + continue; // not a data document; schema:check reports it + } + const schema = String(doc.$schema ?? ""); + if (METADATA_IDS.has(schema)) { + metadataGeographyById.set( + String(doc.id), + doc.geography as PathwayMetadataV2["geography"], + ); + metadataPathById.set(String(doc.id), path); + } else if (TIMESERIES_ID_PATTERN.test(schema)) { + timeseriesById.set(String(doc.id), { path, doc }); + } + } + return { + metadataGeographyById, + metadataPathById, + timeseriesById, + occupiedPaths, + }; +} + +async function main() { + const args = process.argv.slice(2); + const dryRun = args.includes("--dry-run"); + const inAt = args.indexOf("--in"); + const inDir = inAt >= 0 ? args[inAt + 1] : undefined; + if (!inDir) { + console.error("Usage: import-benchmark-data.ts --in [--dry-run]"); + process.exit(2); + } + + const report = new Report(); + const inputs: InputFile[] = []; + for (const path of await jsonFilesUnder(inDir)) { + try { + inputs.push({ + name: path, + data: JSON.parse(await fs.readFile(path, "utf8")), + }); + } catch (e) { + report.error(`${path}: not valid JSON (${String(e)})`); + } + } + report.section(`Import of ${inputs.length} file(s) from ${inDir}`); + const planned = planImport(inputs, await readCorpus(), report); + + console.info(report.lines.join("\n")); + if (report.errors.length > 0) { + console.error( + `\n## ${report.errors.length} error(s) -- nothing written\n` + + "Fix these in the prepared data and re-run the import:\n" + + report.errors.map((e) => ` - ${e}`).join("\n"), + ); + process.exitCode = 1; + return; + } + if (dryRun) { + console.info(`\nDry run: ${planned.length} file(s) would be written.`); + return; + } + // Format everything before touching any file, so a formatting failure + // cannot leave a half-written import behind. + const staged = await Promise.all( + planned.map(async ({ path, doc }) => { + const options = (await prettier.resolveConfig(path)) ?? {}; + const text = await prettier.format(JSON.stringify(doc, null, 2), { + ...options, + parser: "json", + }); + return { path, text }; + }), + ); + await commitWrites(staged); + console.info( + `\nWrote ${planned.length} file(s). Run \`npm run build:timeseries\` and \`npm run schema:check\` next.`, + ); +} + +if (process.argv[1] && fileURLToPath(import.meta.url) === process.argv[1]) { + main().catch((e: unknown) => { + console.error(String(e instanceof Error ? e.stack : e)); + process.exit(1); + }); +} diff --git a/src/data/README.md b/src/data/README.md index ecb82cce..f397d4b9 100644 --- a/src/data/README.md +++ b/src/data/README.md @@ -191,6 +191,14 @@ A `*_timeseries.json` file holds the benchmark data series for one or more pathw `npm run schema:check` enforces both rules against the metadata files, since JSON Schema alone cannot look into another file. +**Timeseries files come from [RMI/tpr_benchmark_data_preparation](https://github.com/RMI/tpr_benchmark_data_preparation)**, which writes one v2 JSON file per dataset. Bring them in with: + +```bash +npx ts-node --esm scripts/import-benchmark-data.ts --in --dry-run +``` + +Drop `--dry-run` to write. The importer converts nothing: every file must already pass the v2 schema and both rules above, and any error blocks the whole import, listing every problem. A file whose `id` matches an existing timeseries overwrites it in place; a new one is written next to its first pathway's metadata, and never onto a path another file already holds. Writes are all-or-nothing too: if one fails, the files already written are restored. Existing files missing from the input are reported and left alone. Then run `npm run build:timeseries` and `npm run schema:check`. The prep repo can run the same command with `--dry-run` to check its output against this repo before handing it over. + Files still on `pathwayTimeseries.v1.json` can be moved with `scripts/codemod-timeseries-v1-to-v2.ts` (`--dry-run` first). Its `RENAMES` table lists, per pathway, which old label became which declared one; it stops on any label it cannot map rather than guess. ## Migrating an existing v1 file