{"_id":"@a0n/aeon-pipelines","_rev":"2-7d53f7e683e40bbcf9c3c052f18c79ba","name":"@a0n/aeon-pipelines","dist-tags":{"latest":"1.0.0"},"versions":{"1.0.0":{"name":"@a0n/aeon-pipelines","version":"1.0.0","_id":"@a0n/aeon-pipelines@1.0.0","maintainers":[{"name":"buley","email":"buley@outlook.com"}],"dist":{"shasum":"fc0a4b3575e8369f9b24870f9f015e3313bfb36d","tarball":"https://registry.npmjs.org/@a0n/aeon-pipelines/-/aeon-pipelines-1.0.0.tgz","fileCount":65,"integrity":"sha512-njSIZg1QHZrEJWt/y0uTvlbmgFJiH5jhSDD9el4pHc2MWGth0IDVqk79chiY+qMRW6l+C4iK6pKcd8aV+WppiA==","signatures":[{"sig":"MEYCIQDJ0eXsi3bVsrL646SlnrOl2BSxc6Zg0d/p4A1vHr6IggIhAMvtoRq4sYfghdUcSIDtozx1W4eih/V0CsybOixC/O6W","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":267802},"main":"src/index.ts","type":"module","_from":"file:a0n-aeon-pipelines-1.0.0.tgz","types":"src/index.ts","exports":{".":"./src/index.ts","./bridge":"./src/flow-bridge.ts","./pipeline":"./src/pipeline.ts","./strategies":"./src/strategies/index.ts"},"scripts":{"test":"vitest run","typecheck":"tsc --noEmit","test:watch":"vitest"},"_npmUser":{"name":"buley","email":"buley@outlook.com"},"_resolved":"/private/var/folders/kf/tkq8cmd54fl9c1wh414dj3240000gn/T/23e9c99ba21993378192ae06af51eb34/a0n-aeon-pipelines-1.0.0.tgz","_integrity":"sha512-njSIZg1QHZrEJWt/y0uTvlbmgFJiH5jhSDD9el4pHc2MWGth0IDVqk79chiY+qMRW6l+C4iK6pKcd8aV+WppiA==","_npmVersion":"11.11.0","description":"Transport-agnostic fork/race/fold pipeline engine with fluidic routing and quantum-inspired state management","directories":{},"_nodeVersion":"25.8.0","_hasShrinkwrap":false,"devDependencies":{"vitest":"^2.1.0","typescript":"^5.8.0"},"peerDependencies":{"@a0n/aeon":"5.0.1","@a0n/gnosis":">=1.0.0"},"peerDependenciesMeta":{"@a0n/aeon":{"optional":true},"@a0n/gnosis":{"optional":true}},"_npmOperationalInternal":{"tmp":"tmp/aeon-pipelines_1.0.0_1773802296591_0.41488796522431604","host":"s3://npm-registry-packages-npm-production"},"deprecated":"This package is deprecated and no longer supported. Do not use."}},"time":{"created":"2026-03-18T02:51:36.438Z","modified":"2026-03-27T20:49:04.008Z","1.0.0":"2026-03-18T02:51:36.735Z"},"description":"Transport-agnostic fork/race/fold pipeline engine with fluidic routing and quantum-inspired state management","maintainers":[{"name":"buley","email":"buley@outlook.com"}],"readme":"# @a0n/aeon-pipelines\n\nA transport-agnostic execution engine for fork/race/fold/vent orchestration. Primitives, composition, fold strategies, scheduling, frame encoding, and topology diagnostics in one package.\n\n## Install\n\n```bash\nbun add @a0n/aeon-pipelines\n```\n\n```bash\nnpm install @a0n/aeon-pipelines   # or yarn / pnpm\n```\n\n## Quick Start\n\n### Race (speed)\n\n```ts\nimport { Pipeline } from '@a0n/aeon-pipelines';\n\n// First to resolve wins, losers are vented\nconst { result } = await Pipeline.from([\n  () => fetch('/api/primary').then((r) => r.json()),\n  () => fetch('/api/fallback').then((r) => r.json()),\n]).race();\n```\n\n### Race (value)\n\n```ts\n// All complete, judge picks the winner -- race on quality, not speed\nconst { result } = await Pipeline.from([\n  () => compress(data, 'gzip'),\n  () => compress(data, 'brotli'),\n  () => compress(data, 'deflate'),\n]).race((results) => {\n  // Pick the smallest output\n  let best = 0;\n  for (let i = 1; i < results.length; i++) {\n    if (results[i].length < results[best].length) best = i;\n  }\n  return best;\n});\n```\n\n### Fold\n\n```ts\n// Wait for all, merge results via strategy\nconst total = await Pipeline.from([\n  () => countFromDB(),\n  () => countFromCache(),\n  () => countFromAPI(),\n]).fold({\n  type: 'merge-all',\n  merge: (results) => Array.from(results.values()).reduce((a, b) => a + b, 0),\n});\n```\n\n### Wallington Rotation + Worthington Whip\n\n```ts\nimport { wallingtonRotation, worthingtonWhip } from '@a0n/aeon-pipelines';\n\n// Chunk rotation: process chunk N+1 while sending chunk N\nconst rotated = await wallingtonRotation(\n  [1, 2, 3, 4, 5, 6, 7, 8],\n  [\n    (chunk) => chunk.map((v) => v + 1),\n    (chunk) => chunk.map((v) => v * 2),\n  ],\n);\n\n// Shard-level: fork shards, rotate per shard, fold collapse\nconst whipped = await worthingtonWhip([1, 2, 3, 4, 5, 6, 7, 8], {\n  shardCount: 2,\n  stages: [(chunk) => chunk.map((v) => v + 10)],\n});\n```\n\n### SequencePipeline (linear-feeling API)\n\n```ts\nconst output = await Pipeline.sequence([1, 2, 3, 4, 5, 6, 7, 8])\n  .through([\n    (chunk) => chunk.map((v) => v + 1),\n    (chunk) => chunk.map((v) => v * 2),\n  ])\n  .rotate({ chunks: 'auto' })\n  .shard(2)\n  .collapse()\n  .run();\n```\n\n### Metrics + Topology\n\n```ts\nconst pipeline = new Pipeline({ capacity: 128, intrinsicBeta1: 16 });\nconst branch = pipeline.fork([() => Promise.resolve(1), () => Promise.resolve(2)]);\nawait branch.fold({ type: 'winner-take-all' });\n\nconst metrics = pipeline.metrics();  // Reynolds number, Betti number, active streams\nconst bule = pipeline.bule();        // Topological deficit diagnostic\n```\n\n## API\n\n### Core Engine\n\n| Export | Description |\n|--------|-------------|\n| `Pipeline` | DAG orchestrator with metrics, backpressure, multiplexing |\n| `Superposition` | Chainable builder: `.fork().race()`, `.vent(predicate)`, `.tunnel()` |\n| `SequencePipeline` | Linear API: `.sequence().through().rotate().shard().run()` |\n\n### Primitives\n\n| Export | Description |\n|--------|-------------|\n| `fork(workFns)` | Create N parallel streams. Beta-1 increases by N-1. |\n| `race(streams, allStreams?, judge?)` | First to complete wins (or judge picks winner from all results) |\n| `fold(streams, strategy)` | Wait for all, merge via strategy |\n| `ventStream(stream)` | Prune branch, propagate down never across |\n| `wallingtonRotation(items, stages)` | Chunk-level pipelined rotation |\n| `worthingtonWhip(items, options)` | Shard-level fork + rotate + fold |\n| `laminar(data, config)` | Codec-racing pipeline (gzip/brotli/deflate per chunk) |\n\n### Fold Strategies\n\n| Strategy | Behavior |\n|----------|----------|\n| `winner-take-all` | Pick first (or custom selector) |\n| `quorum` | N/M agreement threshold |\n| `merge-all` | Collect and merge all results |\n| `consensus` | Constructive (agree) or destructive (conflict detection) |\n| `weighted` | Weighted sum with per-stream weights |\n| `custom` | Your own fold function |\n\n### Race Judge\n\nPass a `RaceJudge<T>` lambda to `race()` to race on value instead of speed:\n\n```ts\ntype RaceJudge<T> = (results: readonly T[]) => number;  // return index of winner\n```\n\nWhen omitted, race uses first-to-finish semantics (backward compatible).\n\n### Execution Paths\n\nThree paths, fastest first:\n\n1. **Frame-native** -- direct `Promise.race`/`Promise.allSettled` on raw work functions, zero Stream overhead\n2. **GG topology** -- compiled to Gnosis runtime (when available)\n3. **TS stream machine** -- full state machine with AbortController per stream\n\n### Quantum Modalities\n\n| Export | Description |\n|--------|-------------|\n| `tunnel(predicate)` | Early exit when confidence threshold met |\n| `interfere(mode, compare)` | Constructive (consensus) or destructive (conflict) observation |\n| `measure(streams)` | Pure observation without execution |\n| `entangle(streams)` | Coupled stream execution |\n| `search(config)` | Grover-style oracle search over candidates |\n\n### Diagnostics\n\n| Export | Description |\n|--------|-------------|\n| `ReynoldsTracker` | Flow regime tracking (laminar/transitional/turbulent) |\n| `TopologyAnalyzer` | Compute Betti numbers, fork-join pairs, DAG validation |\n| `TopologySampler` | Statistical sampling of live topology |\n| `computeBuleMeasure` | Topological deficit: intrinsic vs actual beta-1 |\n\n### Transport\n\n| Export | Description |\n|--------|-------------|\n| `encodeFlowFrame` / `decodeFlowFrame` | 10-byte self-describing frame headers |\n| `encodeBatch` / `decodeBatch` | Batch frame encoding |\n| `createFrame` | Work unit with streamId, sequence, flags |\n| `FrameReassembler` | Out-of-order reassembly with loss detection |\n\nWire format flags: `FORK=0x01`, `RACE=0x02`, `FOLD=0x04`, `VENT=0x08`, `FIN=0x10`, `LAMINAR=0x40`\n\n## Development\n\n```bash\ncd open-source/aeon-pipelines\nnpx vitest run          # 173 tests across 17 files\nnpx vitest              # watch mode\n```\n\n","readmeFilename":"README.md"}