
Reactive flows built on async iterators.
streamix is a reactive flows library built on async iterators. It gives you synchronous state reads, composable derived values, lifecycle-aware scopes, and flows that work naturally with modern TypeScript.
npm install @epikodelabs/streamixCore concepts
Atoms
Atoms are reactive values: readable, writable, composable with derived, and consumable as async iterables when you need pipelines.
import { atom, derived, iterate, map, pipe, take } from '@epikodelabs/streamix';
const count = atom(0); // always has a value
const label = atom<string>(); // value arrives later
const summary = derived(() => `count is ${count.value}`);
count.next(5);
console.log(summary.value); // "count is 5"
const doubled = pipe(count, map(n => n * 2), take(1));
for await (const value of doubled) console.log(value);Operators
Operators transform flows. Sync and async callbacks are both supported.
import { filter, map, pipe, range, take } from '@epikodelabs/streamix';
const valid = pipe(
range(1, 100),
map(x => enrich(x)),
filter(x => x.valid),
take(10)
);
for await (const item of valid) {
console.log(item);
}Pipe signature:
pipe()is a standalone function — the source comes first, followed by any number of operators:pipe(source, op1, op2, ...). The source can be an atom, a flow, or any async iterable. This replaces the v2 method-chaining style (source.pipe(op1, op2)); the pipeline reads left-to-right and returns a new atom. Up to 16 operators keep full type inference — beyond that, the result falls back toAtom<any>. See the migration guide for the full v2 → v3 mapping.
Full catalog: audit, buffer, bufferCount, bufferUntil, bufferWhile, catchError, concatMap, debounce, defaultIfEmpty, delay, delayUntil, delayWhile, distinctUntilChanged, distinctUntilKeyChanged, endWith, exhaustMap, expand, filter, finalize, first, fork, groupBy, ignoreElements, last, map, mergeMap, observeOn, partition, reduce, sample, scan, select, share, shareReplay, skip, skipUntil, skipWhile, slidingPair, startWith, switchMap, take, takeUntil, takeWhile, tap, throttle, throwError, toArray, withLatestFrom.
Flow Factories
| Factory | Description |
|---|---|
addListener(target, event) | DOM / Node events |
combineLatest(...sources) | Latest value from each source, combined |
concat(...sources) | Sources run sequentially |
defer(factory) | Fresh flow per subscriber |
EMPTY | Completes immediately (empty() creates a fresh one) |
forkJoin(...sources) | Emits once when all complete |
from(source) | Arrays, iterables, generators, promises |
interval(ms) | Counter every ms milliseconds |
merge(...sources) | Interleaved concurrent emissions |
of(value) | Single value, then complete |
race(...sources) | First source to emit wins |
range(start, count) | Sequential integers |
retry(factory, maxRetries?, delay?) | Retry the source factory on error |
timer(delay, period?) | Delayed, optionally repeating |
zip(...sources) | Pair emissions by index |
Custom operators
import { createOperator, DONE, NEXT } from '@epikodelabs/streamix';
const onlyPrime = () =>
createOperator<number, number>('onlyPrime', source => ({
async next() {
while (true) {
const result = await source.next();
if (result.done) return DONE;
if (isPrime(result.value)) return NEXT(result.value);
}
}
}));Consume a single emission
for await (const first of pipe(interval(1000), take(1))) {
console.log(first);
}The pipeline completes after the first emitted value.
HTTP client
import { createHttpClient, readJson, useTimeout } from '@epikodelabs/streamix/networking';
const api = createHttpClient({
baseUrl: 'https://api.example.com',
middlewares: [useTimeout(5000)],
});
for await (const data of api.request('/items', readJson)) {
console.log(data);
}Why pull-based?
Most reactive libraries push values eagerly. streamix pulls: the consumer asks for the next value, and only then is it computed.
async function* primes() {
let n = 2;
while (true) {
if (isPrime(n)) yield n;
n++;
}
}
// Only 5 primes are ever computed
for await (const p of pipe(primes(), take(5))) {
console.log(p);
}This gives you on-demand computation, bounded memory, and consumer-driven backpressure without manual coordination.
streamix vs RxJS
| streamix | RxJS | |
|---|---|---|
| Execution model | Pull-based (lazy) | Push-based (eager) |
| Backpressure | Consumer-driven | Manual patterns required |
| Async/await | Native | Limited |
| Bundle size | Small | Larger |
| Reactive state | Atoms + derived | Manual stores |
Resources
License
MIT
API Reference
Check the detailed API Reference here.