Skip to content

streamix

Reactive flows built on async iterators.

License

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.

bash
npm install @epikodelabs/streamix

Core concepts ​

Atoms ​

Atoms are reactive values: readable, writable, composable with derived, and consumable as async iterables when you need pipelines.

ts
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.

ts
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 to Atom<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 ​

FactoryDescription
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
EMPTYCompletes 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 ​

ts
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 ​

ts
for await (const first of pipe(interval(1000), take(1))) {
  console.log(first);
}

The pipeline completes after the first emitted value.


HTTP client ​

ts
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.

ts
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 ​

streamixRxJS
Execution modelPull-based (lazy)Push-based (eager)
BackpressureConsumer-drivenManual patterns required
Async/awaitNativeLimited
Bundle sizeSmallLarger
Reactive stateAtoms + derivedManual stores

Resources ​


License ​

MIT

API Reference ​

Check the detailed API Reference here.

Released under the MIT License.