core/pipelines/snapshot.pipeline.ts

core/pipelines/snapshot.pipeline.ts is a file in GovLab Patterns. 77 lines of code and 18 definitions.

import type { AnalyzeOptions, AnalyzeReport, WindowSnapshot } from "#types/record.types";
import { assemble, feed } from "#core/coordinators/field.coordinator";
import type { FieldAnalyzer } from "#core/analyzers/field.analyzer";
import type { FieldSchema } from "#types/schema.types";
import { buildAnalyzers } from "#core/factories/field.factory";
import { detectSchema } from "#core/analyzers/schema.analyzer";
import { graphOf } from "#core/factories/graph.factory";
import { inferMapping } from "#core/resolvers/representation.resolver";
import { toReport } from "#core/converters/report.converter";

interface StreamContext {
    analyzers: Map<string, FieldAnalyzer>;
    schema: FieldSchema[];
}

interface StreamState {
    context: StreamContext | null;
    processed: number;
}

const contextOf = function contextOf(window: readonly unknown[], options: AnalyzeOptions): StreamContext {
    const schema = detectSchema(window, options.floatFields);
    return { analyzers: buildAnalyzers(options.mapping ?? inferMapping(schema)), schema };
};

const snapshotOf = function snapshotOf(context: StreamContext, processed: number): AnalyzeReport {
    const { nodes, findings } = assemble(context.analyzers);
    return toReport(processed, { findings, graph: graphOf(nodes), schema: context.schema });
};

const stepWindow = function stepWindow(
    window: readonly unknown[],
    state: StreamState,
    options: AnalyzeOptions,
): { state: StreamState; snapshot: WindowSnapshot } {
    const context = state.context ?? contextOf(window, options);
    feed(context.analyzers, window);
    const processed = state.processed + window.length;
    return { snapshot: { count: processed, report: snapshotOf(context, processed) }, state: { context, processed } };
};

export const windowedReports = function* windowedReports(
    records: readonly unknown[],
    windowSize: number,
    options: AnalyzeOptions = {},
): Generator<WindowSnapshot> {
    const context = contextOf(records, options);
    const size = windowSize > 0 ? windowSize : records.length;
    let processed = 0;
    for (let start = 0; start < records.length; start += size) {
        const chunk = records.slice(start, start + size);
        feed(context.analyzers, chunk);
        processed += chunk.length;
        yield { count: processed, report: snapshotOf(context, processed) };
    }
};

const emitWindow = function* emitWindow(
    window: readonly unknown[],
    state: StreamState,
    options: AnalyzeOptions,
): Generator<WindowSnapshot, StreamState> {
    const step = stepWindow(window, state, options);
    yield step.snapshot;
    return step.state;
};

export const streamReports = async function* streamReports(
    source: AsyncIterable<unknown> | Iterable<unknown>,
    windowSize: number,
    options: AnalyzeOptions = {},
): AsyncGenerator<WindowSnapshot> {
    const window: unknown[] = [];
    let state: StreamState = { context: null, processed: 0 };
    for await (const record of source) {
        window.push(record);
        if (window.length >= windowSize) {
            state = yield* emitWindow(window, state, options);
            window.length = 0;
        }
    }
    if (window.length > 0) {
        yield* emitWindow(window, state, options);
    }
};