Imported from uhop/stream-chain (
AGENTS.md). Install upstream withnpx skills add uhop/stream-chain. Copyright stays with the author.
AGENTS.md — stream-chain
stream-chaincreates a chain of streams out of regular functions, asynchronous functions, generators, Node streams, and Web streams, with proper per-item backpressure. The default chain returns a NodeDuplex; subpath variants run natively on Web Streams (stream-chain/web) or pure async iterables (stream-chain/core). Zero runtime dependencies.
For project structure, module dependencies, and the architecture overview see ARCHITECTURE.md. For detailed usage docs and API references see the wiki. For migrating from 3.x see Migration-V3-to-V4.
Setup
This project uses a git submodule for the wiki:
git clone --recursive https://github.com/uhop/stream-chain.git
cd stream-chain
npm install
Commands
- Install:
npm install - Test:
npm test(runstape6 --flags FO) - Test (Bun):
npm run test:bun - Test (Deno):
npm run test:deno - Test (sequential):
npm run test:seq(alsotest:seq:bun,test:seq:deno) - Test (browser):
npm run test:browser— drives headless Chromium viatape-six-puppeteer; auto-startstape6-serveron port55555(env-overridable, avoids the default3000collision). Browser-safe test set is selected bytape6.tests(tests/core/+tests/web/);tape6.cli(tests/node/) is skipped in browser context. Nothing downloads at install time:.npmrcsetsignore-scripts=true, which blocks thetape-six-puppeteer/puppeteerpostinstalls on every npm version, and.puppeteerrc.cjssetsskipDownload. Do not swap it forpackage.json'sallowScripts— that field needs npm 11.16+, so it is silently ignored by the npm 10.9 that Node 22 bundles, and the CI matrix still tests Node 22. The browser resolves from puppeteer's cache (one-timenpm run browser:install) orPUPPETEER_EXECUTABLE_PATH. - Test (single file):
node tests/<bucket>/test-<name>.js(bucket iscore,web, ornode) - TypeScript check:
npm run ts-check - JavaScript type check (dual tsconfig):
npm run js-check - TypeScript tests:
npm run ts-test(alsots-test:bun,ts-test:deno) - Bench:
npm run bench -- bench/<name>.js - Lint:
npm run lint(Prettier check) - Lint fix:
npm run lint:fix(Prettier write)
Project structure
stream-chain/
├── package.json # Package config; "tape6" section configures test discovery
├── src/ # Source code
│ ├── index.js # /node entry: chain() factory + asStream + asWebStream + gen + re-exports
│ ├── index.d.ts # TypeScript definitions for the /node public API
│ ├── defs.js # Special values (none, stop, many, finalValue, flushable, etc.) + Web/Node stream type guards
│ ├── defs.d.ts # TypeScript definitions for defs
│ ├── exec.js # Shared sync-when-possible value-or-promise executor (engine behind gen/fun/asStream/asWebStream)
│ ├── exec.d.ts
│ ├── gen.js # Push→pull async-generator bridge over exec
│ ├── gen.d.ts
│ ├── fun.js # Creates function pipeline from functions (sync-first; collects via exec.next; exported via /web and /core)
│ ├── fun.d.ts
│ ├── dataSource.js # Coerces a function or iterable to an iterator-producing function (substrate-agnostic)
│ ├── dataSource.d.ts
│ ├── asStream.js # Wraps a function as a Node Duplex with per-item backpressure
│ ├── asStream.d.ts
│ ├── asWebStream.js # Wraps a function as a Web Streams {readable, writable} duplex pair
│ ├── asWebStream.d.ts
│ ├── typed-streams.js # TypeScript helpers: TypedReadable, TypedWritable, TypedDuplex, TypedTransform
│ ├── typed-streams.d.ts
│ ├── node/ # Subpath: stream-chain/node — canonical Node Streams chain (re-export of root)
│ │ ├── index.js
│ │ └── index.d.ts
│ ├── web/ # Subpath: stream-chain/web — native Web Streams chain (no node:stream)
│ │ ├── index.js # chain() over {readable, writable} duplex pairs
│ │ └── index.d.ts
│ ├── core/ # Subpath: stream-chain/core — substrate-free async-iterable chain
│ │ ├── index.js # chain() returning a callable async-generator factory
│ │ └── index.d.ts
│ ├── jsonl/ # JSONL (line-separated JSON) support
│ │ ├── parser.js # JSONL parser function (returns gen() pipeline)
│ │ ├── parser.d.ts
│ │ ├── parserStream.js # JSONL parser as a Node Duplex
│ │ ├── parserStream.d.ts
│ │ ├── parserWebStream.js # JSONL parser as a Web Streams duplex pair
│ │ ├── parserWebStream.d.ts
│ │ ├── stringerStream.js # JSONL stringer as a Node Transform
│ │ ├── stringerStream.d.ts
│ │ ├── stringerWebStream.js # JSONL stringer as a Web Streams TransformStream
│ │ └── stringerWebStream.d.ts
│ └── utils/ # Utility functions
│ ├── take.js # Take N items from stream
│ ├── takeWhile.js # Take items while condition is true
│ ├── takeWithSkip.js # Skip then take
│ ├── skip.js # Skip N items
│ ├── skipWhile.js # Skip items while condition is true
│ ├── fold.js # Reduce/fold stream to single value
│ ├── reduce.js # Alias for fold
│ ├── scan.js # Running accumulator (like fold but emits each step)
│ ├── batch.js # Group items into fixed-size arrays
│ ├── readableFrom.js # Convert iterable to Node Readable stream
│ ├── readableWebStreamFrom.js # Convert iterable to Web Streams ReadableStream
│ ├── reduceStream.js # Reduce as a Node Writable stream (.accumulator)
│ ├── reduceWebStream.js # Reduce as a Web WritableStream ({writable, result, accumulator})
│ ├── fixUtf8Stream.js # Fix multi-byte UTF-8 splits across chunks
│ ├── lines.js # Split byte stream into lines
│ ├── streamPuller.js # Wrap Node Readable as a non-destructive async iterator
│ ├── webStreamPuller.js # Wrap Web ReadableStream as a non-destructive async iterator
│ └── *.d.ts # TypeScript definitions for each utility
├── tests/ # Test files organized by environment (see "Tests" below)
│ ├── core/ # Substrate-agnostic — runs in browser AND CLI (uses /web chain internally)
│ ├── web/ # Web Streams — runs in browser AND CLI
│ ├── node/ # Node Streams / node:* APIs — runs only in CLI
│ ├── helpers.js # Node-stream test helpers (Readable/Writable factories) — re-exports web-helpers
│ ├── web-helpers.js # Pure + Web Streams helpers (delay, webStreamToArray, writeAndCollect, runChain)
│ ├── data/ # Test fixtures (referenced by tests/node/test-jsonl-*.js)
│ └── manual/ # Manual test scripts (not part of the automated suite)
├── bench/ # Benchmarks
├── wiki/ # GitHub wiki documentation (git submodule)
└── .github/ # CI workflows, Dependabot config
Code style
- ESM throughout (
"type": "module"in package.json). Source usesimportsyntax. - No transpilation — code runs directly.
- Prettier for formatting (see
.prettierrc): 100 char width, single quotes, no bracket spacing, no trailing commas, arrow parens "avoid". - 2-space indentation.
- Semicolons are enforced by Prettier (default
semi: true). - All public modules declare both
export default Xandexport {X}for the same value (default = ESM DX, named = CJS destructure + cleaner re-exports). See fleet slice 17. // @ts-self-types="./<file>.d.ts"directive at the top of every.js; JSDoc lives in the paired.d.ts, not in.js.- The package is
stream-chain. Internal symbols useSymbol.for('object-stream.*').
Critical rules
- Zero runtime dependencies. Never add packages to
dependencies. OnlydevDependenciesare allowed. - Do not modify or delete test expectations without understanding why they changed.
- Comments are short why-markers only — a non-trivial decision or constraint, or an algorithm reference; never narrate what the code does. See fleet slice 20.
- Keep
src/index.jsandsrc/index.d.tsin sync. All public API is exported fromindex.jsand typed inindex.d.ts. - Keep
.jsand.d.tsfiles in sync for all modules undersrc/. - Object mode by default.
chain()(the /node variant) defaults to{writableObjectMode: true, readableObjectMode: true}. - Per-item backpressure must be preserved.
asStreamandasWebStreamdrive the shared executor (exec.next/exec.flush), whosepushreturn is honored: when an enqueue backpressures it returns a Promise and the executor suspends at that push, resuming on drain. Keeps the queue at hwm+1 under unboundedmany()/generator expansion, with O(1) live allocation (one resume closure per actual suspension, not per element). Do not change the executor to ignore thepushreturn or to eagerly chain per element. - Source generators must be released on abnormal termination. When iteration through a source generator ends abnormally — a downstream stage throwing, a consumer cancelling (for-await
break→CANCEL), orstop—exec'snextGendriver callsit.return()(awaited for an async generator) before re-throwing the original error, so the source'sfinally {}runs. Without it a resource-owning source (e.g.asyncBlockReader'sFileHandle) leaks — cleanup deferred to GC, which on Node raisesERR_INVALID_STATE. Matchesfor await…of'sAsyncIteratorClose(the original error wins over any cleanup error). Don't drop the abort wrapper or let it mask the original error. Only the purepipe/drain(gen()) path needs this;asStream/asWebStreamclean up via the streamdestroylifecycle. - Special values in generators. From a generator, don't yield
defs.noneordefs.many(...)— express them natively (skip withcontinue, emit multiple via separateyields oryield* anIterable). Butdefs.stopanddefs.finalValue(...)ARE supported from generators (a plain generator can't express them):stopterminates the whole pipeline,finalValue(x)emitsxand skips the segment's remaining functions. After issuing either,return. In regular functions all four markers are return values. See wiki/defs.md. chain.asStream/chain.genare override hooks — internal references go through the static-property indirection so users can monkey-patch. Don't refactor to direct imports.
Architecture
chain(fns, options)is the main entry point (default = /node). Returns a NodeDuplexwith.streams,.input,.outputproperties.stream-chain/webexposes a parallelchain()that returns{readable, writable, streams, input, output}— a native Web Streams duplex pair.stream-chain/coreexposes a callable async-iterable factory — no Node streams, no Web Streams. Browser-safe and substrate-free. Input handling:null/undefined→ empty; strings and other non-iterables (numbers, booleans, plain objects, …) → passed through as a single value; arrays / generators / async iterables / Maps / Sets → iterated.- Functions in a chain are grouped together via
gen()for efficiency (unlessnoGrouping: true). exec(...fns)(src/exec.js) is the shared sync-when-possible, value-or-promise executor — the single engine behindgen,fun,asStream, andasWebStream(it replaced the old per-wrapperasync applyFns). It threads a value through the function-list, emits terminal values via apushcallback, and stays synchronous until the first real promise (async stage, thenable, or backpressuring push) appears. Internal — not a public export.gen(...fns)creates an async generator pipeline — a push→pull bridge overexec.next. Handles all special return values from regular functions:none,stop,many(),finalValue(), flushable.fun(...fns)creates a function pipeline (sync when possible). Collects all outputs into aManyper input, so memory scales with output size — not safe for unbounded pipelines. Intentionally NOT on the defaultstream-chain//nodeexport; requires an explicit import fromstream-chain/fun.js(also re-exported via/weband/core). The friction is deliberate.asStream(fn[, options])wraps a function as a NodeDuplexwith per-item backpressure.asWebStream(fn[, options])wraps a function as a Web Streams{readable, writable}pair with per-item backpressure.- Special return values are defined in
defs.js:none(skip),stop(terminate),many(values)(emit multiple),finalValue(value)(skip rest of chain),flushable(fn)(called at stream end). - Web Streams type guards (
isReadableWebStream,isWritableWebStream,isDuplexWebStream) live indefs.jsand are re-exported fromindex.jsandweb/index.js. - The
/nodechain adapts Web Stream objects to Node streams viaReadable.fromWeb()/Writable.fromWeb()/Duplex.fromWeb()with{objectMode: true}. The/webchain handles them natively. - JSONL support is in
src/jsonl/— parser and stringer for line-separated JSON. Parser emits{key, value}per line; empty lines are dropped. Error handling:ignoreErrors: truedrops failed lines but the counter still bumps (gappy keys; back-compat);errorIndicator(presence-checked option —errorIndicator: undefinedis meaningful) substitutes a value or calls a function(error, input, reviver) => unknownwhoseundefinedreturn drops without bumping the counter. Stream wrappers (parserStream,parserWebStream) forward both. Raw export:jsonlParser(per-line factory, nofixUtf8Stream/linesfront). Function-pipeline stringer lives atsrc/jsonl/stringer.js(flushable);stringerStream/stringerWebStreamkeep their Transform / TransformStream shapes. Factory-bundled entries atsrc/{node,web}/jsonl/{parser,stringer}.jscarry.asStream/.asWebStreammethods (Web entries omit.asStreamto stay browser-safe); thesrc/{node,web}/jsonl/index.jsbarrels export{jsonlParser, jsonlStringer}(resolvable asstream-chain/node/jsonl/stream-chain/web/jsonlvia packageexports). They exist so stream-json's deprecated JSONL users can migrate imports to stream-chain with unchanged call sites. - File-edge JSONL components in
src/jsonl/file/(Node-only):parseFile(options)returnsgen(asyncBlockReader, parser);stringerToFile(path, options)returnsgen(stringer, asyncBlockWriter). Drive withpipe(...)+drain(...)(fromsrc/utils/) so the writer's flushable closes the file. Block-I/O primitivesasyncBlockReader/asyncBlockWriterand the substrate-free helperspipe/drainlive insrc/utils/. Round-trip is ~40% faster than the equivalentfs streams + parserStream + stringerStreampipeline; pure parse-and-count via for-await is slower (per-token gen-bridge cost) — seebench/jsonl-file.js. - Utility functions in
src/utils/provide common stream operations: slicing (take,skip), folding (fold,scan), batching, line splitting, UTF-8 fixing, and async-iterator wrappers (makeStreamPuller,makeWebStreamPuller).
Writing tests
import test from 'tape-six';
import chain from 'stream-chain';
import {Readable} from 'node:stream';
test('example', async t => {
const output = [];
const pipeline = chain([x => x * x]);
const source = new Readable({objectMode: true, read() {}});
source.pipe(pipeline);
pipeline.on('data', chunk => output.push(chunk));
pipeline.on('end', () => {
t.deepEqual(output, [1, 4, 9]);
});
source.push(1);
source.push(2);
source.push(3);
source.push(null);
});
- Test files use
tape-six:.jsfor runtime tests,.tsfor TypeScript typing tests,.cjsfor CommonJS tests. - Test file naming convention:
test-*.jsandtest-*.ts. - Tests are configured in
package.jsonunder the"tape6"section. Three buckets (per the user's environment-by-directory convention):tests/core/— substrate-agnostic. Use therunChain(transducers, input) → Promise<output>helper fromtests/web-helpers.js, which internally drives a/webchain. Runs in browser AND CLI (Web Streams are universal in Node 22+/Deno/Bun).tests/web/— Web Streams substrate (asWebStream,/webchain,webStreamPuller). Runs in browser AND CLI.tests/node/— Node Streams substrate (asStream, JSONL vianode:fs+node:zlib,streamPuller, etc.). Runs only in CLI. Anything that importsnode:*or transitively pullstests/helpers.js'sReadable/Writablefactories belongs here.tape6.tests=tests/core+tests/web(both buckets — browser-runnable).tape6.cli=tests/node(added only in non-browser context pertape-six'sresolveTestsrules — seenode_modules/tape-six/TESTING.md§"Configuring test discovery").
- Test files should be directly executable:
node tests/<bucket>/test-foo.js.
Key conventions
- Do not add dependencies unless absolutely necessary — the library is intentionally zero-dependency.
- All public API is exported from
src/index.jsand typed insrc/index.d.ts. Keep them in sync. - Wiki documentation lives in the
wiki/submodule. - Symbols use the
object-streamnamespace:Symbol.for('object-stream.none'), etc. - The library is ESM. CJS consumers use destructure:
const {chain} = require('stream-chain'). The bare-callableconst chain = require('stream-chain')form from 3.x is gone. - Supported Node majors: 22, 24, 26 (latest minor of each).
