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"
}
}
]
}