configuration/principle/data/pipeline.data.json

configuration/principle/data/pipeline.data.json is a file in GovLab Context. 381 lines of code and 0 definitions.

{
    "category": "Streaming / Pipeline / Dataflow Processing",
    "check": {
        "population": "every stage, stream source, accumulator and data edge of a pipeline",
        "freshness": "a verdict stands until a stage, its contract or the input volume changes",
        "refusal": "the pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency",
        "observation": "stage contracts read from source, plus lag, throughput and memory measured under load",
        "evidence": "none: the catalog states this check as a class, so a watched run belongs to each system that adopts it",
        "authority": "the stage contract, which each stage's behavior under load is compared against"
    },
    "records": [
        {
            "id": "streaming-architecture",
            "name": "Streaming Architecture",
            "definition": "A convention of processing data record by record as it arrives, with backpressure, instead of collecting it first.",
            "type": "style",
            "scope": [
                "data processing",
                "integration"
            ],
            "requires": [
                "Event Stream",
                "Backpressure"
            ],
            "reinforces": ["Single-Pass Processing"],
            "enables": ["Continuous Processing"],
            "conflicts_with": [],
            "tensions_with": [
                "Ordering/State",
                "Batch-Only Processing"
            ],
            "violated_by": ["lexicon:full-materialization"],
            "detected_by": ["unbounded collection over stream source"],
            "measured_by": [
                "lag",
                "throughput",
                "memory usage"
            ],
            "refactored_by": ["architecture:backpressure"],
            "enforced_by": ["load/memory tests"],
            "severity": "contextual",
            "exemplar": {
                "before": "const foos = await source.readAll();\nconst results = foos.map(transformFoo);\nawait sink.writeAll(results);",
                "after": "for await (const foo of source.stream()) {\n  await sink.write(transformFoo(foo));\n}",
                "lang": "ts"
            }
        },
        {
            "id": "single-pass-processing",
            "name": "Single-Pass Processing",
            "definition": "A design rule that large input is read once, with every result it feeds computed in that one pass.",
            "type": "principle",
            "scope": [
                "algorithm",
                "stream",
                "parser"
            ],
            "requires": ["Forward-Only State Model"],
            "reinforces": ["Memory Efficiency"],
            "enables": ["Large Input Handling"],
            "conflicts_with": [],
            "tensions_with": [
                "Global Optimization",
                "Multi-Pass Full Materialization"
            ],
            "violated_by": ["lexicon:repeated-full-scan"],
            "detected_by": ["multiple loops/materializations over same large input"],
            "measured_by": [
                "pass count",
                "memory use"
            ],
            "refactored_by": [
                "lexicon:fuse-passes",
                "architecture:iterator-pattern"
            ],
            "enforced_by": ["performance review"],
            "severity": "contextual",
            "exemplar": {
                "before": "const names = foos.map(foo => foo.name);\nconst active = foos.filter(foo => foo.active);\nconst total = foos.reduce((sum, foo) => sum + foo.count, 0);",
                "after": "const names: string[] = [];\nconst active: Foo[] = [];\nlet total = 0;\nfor (const foo of foos) {\n  names.push(foo.name);\n  if (foo.active) active.push(foo);\n  total += foo.count;\n}",
                "lang": "ts"
            }
        },
        {
            "id": "pipeline-architecture",
            "name": "Pipeline Architecture",
            "definition": "A design pattern that splits processing into ordered stages, each with a declared input and output contract.",
            "type": "pattern",
            "scope": [
                "processing",
                "dataflow",
                "build"
            ],
            "requires": ["Stage Contracts"],
            "reinforces": [
                "Composability",
                "Streaming"
            ],
            "enables": ["Stepwise Transformation"],
            "conflicts_with": ["Monolithic Processing Function"],
            "tensions_with": ["Error Propagation/Debugging"],
            "violated_by": ["lexicon:monolithic-processing-function"],
            "detected_by": ["long procedural transformation chain"],
            "measured_by": [
                "stage cohesion",
                "stage contract coverage"
            ],
            "refactored_by": ["lexicon:define-contract"],
            "enforced_by": ["pipeline tests"],
            "severity": "recommended",
            "exemplar": {
                "before": "function processFoo(raw: string) {\n  const parsed = JSON.parse(raw);\n  const validated = validateFoo(parsed);\n  const normalized = normalizeFoo(validated);\n  return saveFoo(normalized);\n}",
                "after": "const fooPipeline = pipeline(\n  parseJson,\n  validateWith(FooSchema),\n  normalizeFoo,\n  saveFoo,\n);\nfooPipeline.run(raw);",
                "lang": "ts"
            }
        },
        {
            "id": "lazy-evaluation",
            "name": "Lazy Evaluation",
            "definition": "An approach in which a value is computed only when something consumes it.",
            "type": "approach",
            "scope": [
                "computation",
                "collection",
                "stream"
            ],
            "requires": ["Deferred Execution Semantics"],
            "reinforces": ["Memory Efficiency"],
            "enables": ["Avoiding Unneeded Work"],
            "conflicts_with": [],
            "tensions_with": [
                "Debuggability/Resource Lifetime",
                "Eager Full Materialization"
            ],
            "violated_by": ["lexicon:full-materialization"],
            "detected_by": ["eager loading of large unused data"],
            "measured_by": [
                "avoided work",
                "memory reduction"
            ],
            "refactored_by": ["architecture:iterator-pattern"],
            "enforced_by": ["performance tests"],
            "severity": "contextual",
            "exemplar": {
                "before": "const normalized = millionFoos.map(normalizeFoo);\nconst active = normalized.filter(foo => foo.active);\nconst firstTen = active.slice(0, 10);",
                "after": "const firstTen = sequence(millionFoos)\n  .map(normalizeFoo)\n  .filter(foo => foo.active)\n  .take(10)\n  .toArray();",
                "lang": "ts"
            }
        },
        {
            "id": "sequential-access",
            "name": "Sequential Access",
            "definition": "A design pattern that reads a source in its stored order through a cursor or scan, instead of looking records up one by one.",
            "type": "pattern",
            "scope": [
                "file",
                "stream",
                "iterator"
            ],
            "requires": ["Ordered Read Model"],
            "reinforces": ["Memory Efficiency"],
            "enables": ["Large Data Processing"],
            "conflicts_with": [],
            "tensions_with": [
                "Lookup Performance",
                "Random Access Requirement"
            ],
            "violated_by": ["lexicon:random-access-over-a-stream"],
            "detected_by": ["seek/index assumptions on sequential source"],
            "measured_by": ["access pattern cost"],
            "refactored_by": [],
            "enforced_by": ["performance tests"],
            "severity": "contextual",
            "exemplar": {
                "before": "for (const id of fooIds) await fooStore.randomRead(id);",
                "after": "for await (const foo of fooStore.scan({ orderBy: \"id\" })) {\n  processFoo(foo);\n}",
                "lang": "ts"
            }
        },
        {
            "id": "forward-only-processing",
            "name": "Forward-Only Processing",
            "definition": "A rule or precondition that a stream is read once in order, keeping only rolling state, so no decision ever needs to revisit earlier input.",
            "aliases": ["No Backtracking Requirement"],
            "type": "constraint",
            "scope": [
                "stream",
                "parser",
                "iterator"
            ],
            "requires": [],
            "reinforces": ["Single-Pass Processing"],
            "enables": ["Streaming Parsers"],
            "conflicts_with": [],
            "tensions_with": [
                "Complex Grammar/Global State",
                "Backtracking Algorithm"
            ],
            "violated_by": ["lexicon:look-back-buffering"],
            "detected_by": ["buffering full stream to look back"],
            "measured_by": [
                "buffer size",
                "pass count"
            ],
            "refactored_by": [
                "architecture:windowing",
                "lexicon:tolerant-reader"
            ],
            "enforced_by": ["memory tests"],
            "severity": "contextual",
            "exemplar": {
                "before": "const cursor = fooStream.cursor();\ncursor.next();\ncursor.previous();\ncursor.seek(0);",
                "after": "for await (const foo of fooStream) {\n  await processFoo(foo);\n}",
                "lang": "ts"
            }
        },
        {
            "id": "dataflow-architecture",
            "name": "Dataflow Architecture",
            "definition": "A convention of arranging processing as a graph of stages connected by explicit data edges, where a stage runs when its inputs arrive.",
            "type": "style",
            "scope": [
                "processing",
                "workflow",
                "stream"
            ],
            "requires": [
                "Data Dependencies",
                "Stages"
            ],
            "reinforces": ["Pipeline Architecture"],
            "enables": ["Parallel/Stream Processing"],
            "conflicts_with": ["Control-Flow-Centric Monolith"],
            "tensions_with": ["State Coordination"],
            "violated_by": ["lexicon:hidden-dependency"],
            "detected_by": ["implicit shared state in pipeline"],
            "measured_by": ["data dependency clarity"],
            "refactored_by": ["architecture:pipeline-architecture"],
            "enforced_by": ["pipeline contracts"],
            "severity": "contextual",
            "exemplar": {
                "before": "controller.runFoo();\ncontroller.runBar();\ncontroller.runBaz();",
                "after": "const graph = dataflow()\n  .source(\"foo\", fooSource)\n  .map(\"bar\", \"foo\", toBar)\n  .map(\"baz\", \"bar\", toBaz)\n  .sink(\"output\", \"baz\", bazSink);\nawait graph.run();",
                "lang": "ts"
            }
        },
        {
            "id": "stateless-processing",
            "name": "Stateless Processing",
            "definition": "A design rule that a processor derives each result only from its input and explicitly passed context.",
            "type": "principle",
            "scope": [
                "function",
                "stream processor",
                "service"
            ],
            "requires": [
                "Explicit Inputs",
                "No Hidden State"
            ],
            "reinforces": [
                "Scalability",
                "Testability"
            ],
            "enables": ["Parallel Processing"],
            "conflicts_with": ["Stateful Hidden Accumulation"],
            "tensions_with": ["Stateful Business Rules"],
            "violated_by": ["lexicon:stateful-hidden-accumulation"],
            "detected_by": ["mutable state across records/requests"],
            "measured_by": ["stateful operator count"],
            "refactored_by": [
                "lexicon:externalize-session-state",
                "lexicon:pass-context-explicitly"
            ],
            "enforced_by": [
                "code review",
                "tests"
            ],
            "severity": "recommended",
            "exemplar": {
                "before": "class FooProcessor {\n  private previous?: Foo;\n  process(foo: Foo) {\n    const result = merge(this.previous, foo);\n    this.previous = foo;\n    return result;\n  }\n}",
                "after": "function processFoo(foo: Foo, context: Readonly<FooContext>): FooResult {\n  return deriveFooResult(foo, context);\n}",
                "lang": "ts"
            }
        },
        {
            "id": "windowing",
            "name": "Windowing",
            "definition": "A mechanism that groups an unbounded stream into bounded windows of event time and aggregates each window separately.",
            "type": "mechanism",
            "scope": [
                "data processing",
                "streaming",
                "aggregation"
            ],
            "requires": ["Event Time"],
            "reinforces": [
                "Streaming Architecture",
                "Bounded State"
            ],
            "enables": ["Bounded Aggregation over Unbounded Streams"],
            "conflicts_with": ["Unbounded Accumulation"],
            "tensions_with": ["Late-Data Handling"],
            "violated_by": ["lexicon:unbounded-accumulation"],
            "detected_by": ["unbounded accumulator over a stream"],
            "measured_by": ["aggregation state growth rate"],
            "refactored_by": [],
            "enforced_by": ["streaming design review"],
            "severity": "contextual",
            "exemplar": {
                "before": "const total = allFooEvents.reduce((sum, event) => sum + event.value, 0);",
                "after": "for await (const window of fooStream.tumbling({ seconds: 60 })) {\n  emit(window.start, window.events.reduce((sum, event) => sum + event.value, 0));\n}",
                "lang": "ts"
            }
        },
        {
            "id": "fan-out-fan-in",
            "name": "Fan-out/Fan-in",
            "aliases": ["Upward Reduction"],
            "definition": "A design pattern that splits work into independent parts processed in parallel, then merges their results.",
            "type": "pattern",
            "scope": [
                "data processing",
                "parallelism",
                "pipeline"
            ],
            "requires": ["Independent Work Units"],
            "reinforces": [
                "Parallelism",
                "Throughput"
            ],
            "enables": [
                "Parallel Branch Processing",
                "Result Aggregation"
            ],
            "conflicts_with": [],
            "tensions_with": [
                "Coordination Overhead",
                "Serial Item Processing"
            ],
            "violated_by": ["lexicon:sequential-bottleneck"],
            "detected_by": ["serial loop over parallelizable work"],
            "measured_by": ["parallelism utilization"],
            "refactored_by": [],
            "enforced_by": ["pipeline design review"],
            "severity": "contextual",
            "exemplar": {
                "before": "const report = await buildFullFooReport(foos);",
                "after": "const partials = await fanOut(partition(foos), buildPartialFooReport);\nconst report = fanIn(partials, mergeFooReports);",
                "lang": "ts"
            }
        },
        {
            "id": "batch-vs-stream",
            "name": "Batch-vs-Stream",
            "definition": "An approach in which the processing model, periodic batch or continuous stream, is chosen by how fresh the data must be.",
            "type": "approach",
            "scope": [
                "data processing",
                "latency",
                "architecture"
            ],
            "requires": ["Latency Requirement Clarity"],
            "reinforces": ["Fitness for Purpose"],
            "enables": ["Latency-Appropriate Processing Model"],
            "conflicts_with": ["One-Size-Fits-All Processing"],
            "tensions_with": ["Operational Duplication"],
            "violated_by": ["lexicon:one-size-fits-all-processing"],
            "detected_by": ["batch cadence mismatched to freshness requirements"],
            "measured_by": ["data-freshness lag vs requirement"],
            "refactored_by": [],
            "enforced_by": ["data architecture review"],
            "severity": "contextual",
            "exemplar": {
                "before": "schedule.daily(() => reprocessAllFoos());",
                "after": "fooStream.subscribe(foo => processFoo(foo));",
                "lang": "ts"
            }
        }
    ]
}