# core/pipelines/snapshot.pipeline.ts

> 77 lines of code and 18 definitions.

Tree: GovLab Patterns
Language: typescript
Layer: processing
Canonical: https://banes-lab.com/anatomy/patterns#file-patterns-core-pipelines-snapshot-pipeline-ts
Source text: https://banes-lab.com/source/patterns/core/pipelines/snapshot.pipeline.ts.txt

Listed in [core/pipelines](https://banes-lab.com/api/source/patterns/core/pipelines.md), after [core/pipelines/record.pipeline.ts](https://banes-lab.com/source/patterns/core/pipelines/record.pipeline.ts.md).

## Definitions

- `snapshotOf` (lexical_declaration, line 26)
- `emitWindow` (lexical_declaration, line 58)
- `windowedReports` (lexical_declaration, line 42, exported)
- `contextOf` (lexical_declaration, line 21)
- `stepWindow` (lexical_declaration, line 31)
- `streamReports` (lexical_declaration, line 68, exported)
- `StreamContext` (interface_declaration, line 11)
- `StreamState` (interface_declaration, line 16)
- `schema` (lexical_declaration, line 22)
- `{ nodes, findings }` (lexical_declaration, line 27)
- `context` (lexical_declaration, line 47, exported)
- `size` (lexical_declaration, line 48, exported)
- `processed` (lexical_declaration, line 49, exported)
- `start` (lexical_declaration, line 50, exported)
- `chunk` (lexical_declaration, line 51, exported)
- `step` (lexical_declaration, line 63)
- `window` (lexical_declaration, line 73, exported)
- `state` (lexical_declaration, line 74, exported)

## Contained in

- [core/pipelines](https://banes-lab.com/anatomy/patterns/folder-patterns-core-pipelines.md)

## Uses

- [core/converters/report.converter.ts](https://banes-lab.com/source/patterns/core/converters/report.converter.ts.md)
- [core/coordinators/field.coordinator.ts](https://banes-lab.com/source/patterns/core/coordinators/field.coordinator.ts.md)
- [core/factories/field.factory.ts](https://banes-lab.com/source/patterns/core/factories/field.factory.ts.md)
- [core/factories/graph.factory.ts](https://banes-lab.com/source/patterns/core/factories/graph.factory.ts.md)
- [core/resolvers/representation.resolver.ts](https://banes-lab.com/source/patterns/core/resolvers/representation.resolver.ts.md)

## Used by

- [runtime/entrypoints/pattern.entrypoint.ts](https://banes-lab.com/source/patterns/runtime/entrypoints/pattern.entrypoint.ts.md)

## Linked from

- [core/converters](https://banes-lab.com/anatomy/patterns/folder-patterns-core-converters.md)
- [core/coordinators](https://banes-lab.com/anatomy/patterns/folder-patterns-core-coordinators.md)
- [core/factories](https://banes-lab.com/anatomy/patterns/folder-patterns-core-factories.md)
- [core/pipelines](https://banes-lab.com/anatomy/patterns/folder-patterns-core-pipelines.md)
- [core/resolvers](https://banes-lab.com/anatomy/patterns/folder-patterns-core-resolvers.md)
- [runtime/entrypoints](https://banes-lab.com/anatomy/patterns/folder-patterns-runtime-entrypoints.md)

## Source

```typescript
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);
    }
};
```
