core/loaders/record.loader.ts

core/loaders/record.loader.ts is a file in GovLab Patterns. 86 lines of code and 23 definitions.

import { createReadStream, readFileSync, readdirSync, statSync } from "node:fs";
import { DataLoadError } from "#core/classifiers/schema.classifier";
import type { LoadedData } from "#types/record.types";
import { NO_RECORDS } from "#configuration/strings/schema.strings";
import { createInterface } from "node:readline";
import { isRecord } from "#core/predicates/record.predicate";
import { join } from "node:path";
import { scanFloatFields } from "#core/parsers/field.parser";

const DEPTH_JSONL = 1;
const DEPTH_ARRAY = 2;
const DEPTH_OBJECT = 3;
const ARRAY_OPEN = "[";
const LINE_BREAK = "\n";
const SOURCE_SUFFIXES: readonly string[] = [".json", ".jsonl"];

interface JsonDocument {
    ok: boolean;
    value: unknown;
}

const parseJson = function parseJson(text: string): unknown {
    return JSON.parse(text);
};

const extractRecords = function extractRecords(parsed: unknown): unknown[] {
    if (Array.isArray(parsed)) {
        return parsed;
    }
    const found = isRecord(parsed) ? Object.values(parsed).find((candidate) => Array.isArray(candidate)) : undefined;
    if (Array.isArray(found)) {
        return found;
    }
    throw new DataLoadError(NO_RECORDS);
};

const parseJsonl = function parseJsonl(text: string): unknown[] {
    return text
        .split(LINE_BREAK)
        .filter((line) => line.trim().length > 0)
        .map(parseJson);
};

const wholeDocument = function wholeDocument(text: string): JsonDocument {
    try {
        return { ok: true, value: parseJson(text) };
    } catch (error) {
        if (!(error instanceof SyntaxError)) {
            throw error;
        }
        return { ok: false, value: null };
    }
};

const loadFile = function loadFile(path: string): LoadedData {
    const text = readFileSync(path, "utf8");
    if (text.trimStart().startsWith(ARRAY_OPEN)) {
        return { floatFields: scanFloatFields(text, DEPTH_ARRAY), records: extractRecords(parseJson(text)) };
    }
    const document = wholeDocument(text);
    if (document.ok) {
        return { floatFields: scanFloatFields(text, DEPTH_OBJECT), records: extractRecords(document.value) };
    }
    return { floatFields: scanFloatFields(text, DEPTH_JSONL), records: parseJsonl(text) };
};

const isSource = function isSource(name: string): boolean {
    return SOURCE_SUFFIXES.some((suffix) => name.endsWith(suffix));
};

export const discoverSources = function discoverSources(dir: string): string[] {
    return readdirSync(dir)
        .filter(isSource)
        .sort((a, b) => a.localeCompare(b))
        .map((name) => join(dir, name));
};

const mergeSources = function mergeSources(paths: readonly string[]): LoadedData {
    const loaded = paths.map(loadFile);
    return {
        floatFields: new Set(loaded.flatMap((entry) => [...entry.floatFields])),
        records: loaded.flatMap((entry) => entry.records),
    };
};

export const loadRecords = function loadRecords(path: string): LoadedData {
    return statSync(path).isDirectory() ? mergeSources(discoverSources(path)) : loadFile(path);
};

export const streamJsonlFile = async function* streamJsonlFile(path: string): AsyncGenerator {
    const reader = createInterface({ crlfDelay: Number.POSITIVE_INFINITY, input: createReadStream(path, "utf8") });
    for await (const line of reader) {
        const trimmed = line.trim();
        if (trimmed.length > 0) {
            yield parseJson(trimmed);
        }
    }
};