transferum
v1.8.0
Published
A language for describing interactions between components — type-safe primitives that compose into interaction graphs with compile-time-checked contracts
Maintainers
Readme
Transferum

A language for describing interactions between components.
Type-safe primitives — transfers (nodes), bridges (edges), builders (chain constructors), and operators (data transformers) — compose into interaction graphs with compile-time-checked contracts.
Table of Contents
- Why Transferum?
- Domain-Specific Applications
- Comparison with Alternatives
- Installation & Import
- Core Concepts
- Transfers
- Transfer Comparison Table
- PushChannelTransfer
- DelayedPushChannelTransfer
- DebounceTransfer
- ThrottleTransfer
- PushStoredChannelTransfer
- BufferTransfer
- ManualBufferTransfer
- ManualFlowTransfer
- GateTransfer
- MergeTransfer
- SplitTransfer
- PollingSourceTransfer
- PollingProxyTransfer
- PollingFlowTransfer
- IdlePollingTransfer
- ChannelTransfer
- StoredChannelTransfer
- SinkTransfer
- WriteTransfer
- ReadTransfer
- ConvertTransfer
- ConditionTransfer
- DisplaceTransfer
- UniversalCompositeTransfer
- Async Transfer Comparison Table
- Async Transfers
- Backpressure
- Operators
- Storages
- Tickers
- Helpers
- Bridges
- Pipeline Builders
- Factories
- Utilities
- Guards
- Types
- Configurations
- Running Tests
- License
Why Transferum?
Problems It Solves
As data flows grow in complexity, connecting heterogeneous sources becomes a major bottleneck. Mixing push-based streams, pull-based APIs, polling loops, and async operations demands endless bespoke glue code, and this integration cost compounds with every new stage.
Transferum addresses the structural issues that emerge at scale:
- No unified contract between stages — connecting a subscribable source to a pushable sink, a pullable buffer to a polling proxy, or a sync stream to an async consumer shouldn't require hand-written adapters every time.
- Type erosion across boundaries — as data passes through transformation, filtering, and async stages, input/output types are lost. Mismatches surface at runtime, not at compile time.
- Sync/async impedance mismatch — mixing synchronous reactive streams with async operations typically forces a full migration to one model or manual bridging at every boundary.
- Flow control as an afterthought — gating, routing, rate limiting, and fallback behavior are scattered across conditional flags and ad-hoc timers rather than composed from reusable primitives.
- Fragmented lifecycle management — subscriptions, timers, and intermediate nodes are cleaned up inconsistently. A uniform
destroy()contract makes resource ownership explicit and leaks detectable.
Transferum provides composable, type-safe building blocks with a uniform capability system. You declare what each stage does (push, pull, transform, filter, poll) — the library handles how data moves between stages, including sync/async bridging, flow control, and resource cleanup.
Not a stream library — a data transfer graph library. Unlike RxJS, where everything is an
Observableand composition happens inside a single object, Transferum treats transfers as nodes and bridges as edges in a data transfer graph.linkTransfersconnects nodes by inspecting their capability contracts — not their class names. Builders assemble subgraphs fluently. Operators are local graph transformations. This makes Transferum architecturally closer to dataflow systems and component graphs than to classical reactive stream libraries.
Transfer ──── Bridge ──── Transfer
\ │
\ │
─────── Bridge ──── TransferThis is a graph.
Key Benefits
| Benefit | How |
|-------------------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| Type-safe pipelines | Each transfer and operator carries its input/output types. Builders enforce type compatibility at compile time — a mismatch is a compile error, not a runtime crash. |
| Uniform capability model | Every transfer declares its capabilities via flags (isPushable, isSubscribable, isGate, …). linkTransfers automatically selects the correct wiring strategy — no manual glue code. Flags also define the transfer's TypeScript interface at compile time — methods like push() or subscribe() are type-guaranteed, not runtime-guessed (see Capability Flags System). |
| Sync + async in one system | Sync and async transfers coexist. linkTransfers prefers sync when possible and falls back to async strategies when needed. No separate "async world." |
| Composable architecture | Transfers link into chains, bridges connect chains with gate control, builders assemble chains fluently, operators transform data — all orthogonal and reusable. |
| Explicit lifecycle | Every resource has an explicit cleanup method — destroy() for transfers, bridges, and subscription managers; stop() for tickers; unsubscribe() for subscriptions. Builders track owned resources and clean them up in one call. No leaked timers or subscriptions. |
| Reactive by default, pull when needed | Most transfers are subscribable (push-based reactivity). Polling transfers add pull-based data acquisition on the same foundation. Use the right model per stage without switching libraries. |
| Local, fail-safe error handling | Errors are local to each transfer — one stage's failure doesn't kill the pipeline. With onError — suppressed, stream continues. Without — visible (exception/rejection), and polling stops (no zombie tickers). Per-stage granularity (onAcceptError/onEmitError, onDestroyError). Typed ErrorHandler<TSource> passes the transfer instance. No silent swallowing. |
| Undefined suppression | undefined never propagates through the chain of transfers — it means "no data", not "empty value." Use null as an explicit empty marker when needed. This eliminates an entire class of null-check bugs in downstream consumers. |
| Built-in backpressure | Four async transfers (AsyncSinkTransfer, AsyncWriteTransfer, AsyncConvertTransfer, AsyncConditionTransfer) support optional maxConcurrency, bufferSize, and onBufferOverflow — limiting parallel async operations, queuing excess data, and handling overflow gracefully. Defaults are backward-compatible (unlimited). See Backpressure. |
| Ordered async execution | AsyncSinkTransfer and AsyncWriteTransfer support optional ordered: true — callback/write invocations are executed sequentially in data-arrival order, regardless of their async duration. AsyncConvertTransfer and AsyncConditionTransfer automatically enforce ordered emission when maxConcurrency > 1 via an internal Sequence Guard (no config needed). See Backpressure. |
| No god-objects, no utility sprawl | BaseTransfer stays minimal — only capability declarations, no logic. Each transfer class models exactly one behavioral concept (buffering, gating, merging, polling...). No central object knows about every other component. Complexity is concentrated in the type layer; runtime code stays compact and readable. |
| Pluggable linking | LinkStrategyInterface lets you override how transfers are wired together — implement link() for custom logic (logging, validation, serialization, inter-process bridging) and inject it into CompositeTransferBuilder.start(transfer, { linkStrategy }). The default DefaultLinkStrategy delegates to linkTransfers() — zero-config for existing code, plug-and-play for custom needs. See Linking. |
Use Cases
Real-time UI updates from API polling
import {
CompositeTransferBuilder, createAsyncPollingSourceTransfer, createConvertTransfer,
createMapOperator, createPushStoredChannelTransfer,
} from 'transferum';
// Poll an API every 5 seconds, transform the response, update subscribers
const polling = createAsyncPollingSourceTransfer<ServerState>({
fetcher: async (): Promise<ServerState> => await fetchApi('/api/state'),
interval: 5000,
activated: true,
});
const pipeline = CompositeTransferBuilder
.start(polling)
.to(createConvertTransfer<ServerState, ViewModel>({
operator: createMapOperator((state) => toViewModel(state)),
}))
.finish(createPushStoredChannelTransfer<ViewModel>());
pipeline.subscribe((vm) => renderUI(vm));Debounced user input with async validation
import {
CompositeTransferBuilder, createPushStoredChannelTransfer, createDebounceTransfer, createAsyncConditionTransfer,
createAsyncConvertTransfer, createAsyncSinkTransfer, createAsyncMapOperator,
} from 'transferum';
// Debounce input → validate → transform → send to async sink
const input = createPushStoredChannelTransfer<string>();
const pipeline = CompositeTransferBuilder
.start(input)
.to(createDebounceTransfer<string>({ delay: 300 }))
.to(createAsyncConditionTransfer<string>({ shouldAccept: async (s) => s.length > 0 }))
.to(createAsyncConvertTransfer<string, ValidationResult>({
operator: createAsyncMapOperator(async (s) => await validate(s)),
}), { onLinkError: (e) => console.error(e) })
.finish(createAsyncSinkTransfer<ValidationResult>({
callback: async (result) => await saveResult(result),
}), { onLinkError: (e) => console.error(e) });
input.push('[email protected]'); // debounced → validated → savedMerging multiple data sources into a single view
import { createAsyncPollingSourceTransfer, createPushStoredChannelTransfer, createMergeTransfer } from 'transferum';
// Merge multiple sensor streams into one (all sources must have the same type)
const tempSensor = createAsyncPollingSourceTransfer<SensorData>({
fetcher: () => Promise.resolve({ sensor: 'temperature', value: 25 }),
interval: 1000,
activated: true,
});
const humiditySensor = createAsyncPollingSourceTransfer<SensorData>({
fetcher: () => Promise.resolve({ sensor: 'humidity', value: 55 }),
interval: 1000,
activated: true,
});
const merge = createMergeTransfer<SensorData>({ sources: [tempSensor, humiditySensor] });
merge.subscribe((data) => updateDashboard(data)); // receives data from both sensorsConditional routing with bridges
import { createBridgeSelector, createPassBridge } from 'transferum';
// Route data to different processing pipelines based on a selector
const fastBridge = createPassBridge({ source, target: fastPipeline, activated: false });
const slowBridge = createPassBridge({ source, target: slowPipeline, activated: false });
const router = createBridgeSelector({
bridges: { fast: fastBridge, slow: slowBridge },
initialKey: 'fast',
activated: true,
owned: false,
});
// Switch route at runtime
router.select('slow');Idle fallback polling
import { createAsyncIdlePollingTransfer } from 'transferum';
// When sensor push API stops interacting, fall back to polling for fresh data
const channel = createAsyncIdlePollingTransfer<SensorState>({
fetcher: async () => await fetchSensorState(),
timeout: 10000, // 10 s of inactivity → start polling
interval: 2000, // poll every 2 s
activated: true,
onError: (e, currentChannel) => {
console.warn('Cannot get sensor state');
showAlert(currentChannel);
},
});
channel.subscribe((item) => appendToChart(item));
// Sensor pushes items (e.g. by websocket) → real-time. User goes idle → automatic polling kicks in.
channel.push({ temperature: 25 });Async data pipeline with storage
import {
CompositeTransferBuilder, createPushStoredChannelTransfer, createAsyncConvertTransfer, createAsyncMapOperator,
} from 'transferum';
// Push data → async transform → notify subscribers
const source = createPushStoredChannelTransfer<RawData>();
const pipeline = CompositeTransferBuilder
.start(source)
.to(createAsyncConvertTransfer<RawData, ProcessedData>({
operator: createAsyncMapOperator(async (raw) => await process(raw)),
}), { onLinkError: (e) => console.error(e) })
.finish(createPushStoredChannelTransfer<ProcessedData>(), {
onLinkError: (e) => console.error(e),
});
pipeline.subscribe((data) => console.log('processed data', data));
source.push(rawData); // → async transformation → processed dataBroadcast to multiple consumers
import { createPushStoredChannelTransfer, createSplitTransfer, createSinkTransfer, createWriteTransfer, linkTransfers } from 'transferum';
// One source → multiple independent consumers
const source = createPushStoredChannelTransfer<Telemetry>();
const split = createSplitTransfer<Telemetry>({
targets: [
createSinkTransfer({ callback: (t) => logTelemetry(t) }),
createSinkTransfer({ callback: (t) => updateChart(t) }),
createWriteTransfer({ flow: telemetryStorage }),
],
});
linkTransfers(source, split);
source.push(telemetry); // → logged, charted, and stored simultaneouslyDomain-Specific Applications
Transferum is designed for building complex, predictable data processing systems across various domains:
Game Development
Input Processing Pipeline
Collect events from keyboard → filter (e.g., DebounceTransfer to prevent spam) → transform into game commands → route to appropriate systems.
import {
createDebounceTransfer, createConvertTransfer, createMapOperator,
createBridgeSelector, createPassBridge,
} from 'transferum';
// Raw input source (for keystrokes)
// Delays handling by 16ms (~1 frame at 60 FPS) to debounce rapid inputs
const rawInputSource = createDebounceTransfer<KeyboardEvent>({ delay: 16 });
// Converter: transforms raw KeyboardEvent into a simple string action
const inputConverter = createConvertTransfer<KeyboardEvent, string>({
operator: createMapOperator((event) => event.code), // Extracts the key code
});
// Target subsystems (consumers that process the final commands)
const carSystem = { push: (cmd: string) => console.log(`🚗 Car executing: ${cmd}`) };
const planeSystem = { push: (cmd: string) => console.log(`✈️ Plane executing: ${cmd}`) };
const menuSystem = { push: (cmd: string) => console.log(`📋 Menu processing: ${cmd}`) };
// Router: manages which input context/subsystem is currently active
const gameplayRouter = createBridgeSelector({
bridges: {
driving: createPassBridge({ source: inputConverter, target: carSystem }),
flying: createPassBridge({ source: inputConverter, target: planeSystem }),
ui: createPassBridge({ source: inputConverter, target: menuSystem }),
},
initialKey: 'driving', // The player starts inside a car by default
activated: true,
owned: true,
});
// Pipe data from the raw source to the converter (proper subscription without loops)
rawInputSource.subscribe((event) => inputConverter.push(event));
// Gameplay Simulation:
async function runSimulation() {
// 1. Player presses 'KeyW' while driving the car
rawInputSource.push(new KeyboardEvent('keydown', { code: 'KeyW' }));
// Expected Output after debounce: 🚗 Car executing: KeyW
// Wait for debounce timer (16ms) to fire and deliver data to carSystem
await sleep(20);
// 2. Player boards a plane — switch the control context
gameplayRouter.select('flying');
// 3. Player presses the exact same 'KeyW' key
rawInputSource.push(new KeyboardEvent('keydown', { code: 'KeyW' }));
// Expected Output after debounce: ✈️ Plane executing: KeyW
// Wait for the second debounce timer to fire
await sleep(20);
}
runSimulation();Particle & Sound Effects
Use BridgeMultiSelector to activate multiple effects simultaneously on events (explosion, hit).
import { createBridgeMultiSelector, createPassBridge } from 'transferum';
const effects = createBridgeMultiSelector({
bridges: {
explosion: createPassBridge({ source: trigger, target: particleSystem, activated: false }),
sound: createPassBridge({ source: trigger, target: audioSystem, activated: false }),
shake: createPassBridge({ source: trigger, target: cameraShake, activated: false }),
},
initialKeys: [],
activated: true,
owned: true,
});
// On explosion event
effects.check('explosion');
effects.check('sound');
effects.check('shake');IoT & Automation
Sensor Data Aggregation
Read data from multiple sensors (temperature, humidity, motion) via AsyncPollingSourceTransfer → filter (ConditionTransfer) → aggregate → send to cloud or local storage.
import {
CompositeTransferBuilder, createAsyncPollingSourceTransfer, createMergeTransfer,
createConditionTransfer, createAsyncWriteTransfer,
} from 'transferum';
const sensor1 = createAsyncPollingSourceTransfer<SensorData>({
fetcher: () => Promise.resolve({ temperature: 25, humidity: 50 }),
interval: 50,
activated: true,
});
const sensor2 = createAsyncPollingSourceTransfer<SensorData>({
fetcher: () => Promise.resolve({ temperature: 26, humidity: 55 }),
interval: 50,
activated: true,
});
const aggregator = createMergeTransfer<SensorData>({
sources: [sensor1, sensor2],
});
const pipeline = CompositeTransferBuilder
.start(aggregator)
.to(createConditionTransfer<SensorData>({ shouldAccept: (d) => d.temperature > 0 && d.humidity >= 0 }))
.finish(createAsyncWriteTransfer<SensorData>({ flow: cloudStorage }));Device Control
Process commands from users or external systems → route to specific actuators (BridgeSelector) → receive feedback.
import { createBridgeSelector, createPassBridge } from 'transferum';
const commandRouter = createBridgeSelector({
bridges: {
light: createPassBridge({ source: commandChannel, target: lightController, activated: false }),
thermostat: createPassBridge({ source: commandChannel, target: thermostatController, activated: false }),
lock: createPassBridge({ source: commandChannel, target: lockController, activated: false }),
},
initialKey: 'light',
activated: true,
owned: true,
});
commandRouter.select('thermostat'); // switch to thermostat controlMonitoring & Alerts
Poll device temperature via AsyncPollingSourceTransfer → filter by threshold (ConditionTransfer) → throttle alerts (ThrottleTransfer) → transform into alert (ConvertTransfer) → send notification.
import {
CompositeTransferBuilder, createConditionTransfer, createThrottleTransfer,
createConvertTransfer, createMapOperator, createAsyncPollingSourceTransfer,
} from 'transferum';
const TEMPERATURE_THRESHOLD = 95;
const tempMonitor = createAsyncPollingSourceTransfer<number>({
fetcher: async () => await readTemperature(),
interval: 1000,
activated: true,
});
const alertPipeline = CompositeTransferBuilder
.start(tempMonitor)
.to(createConditionTransfer<number>({ shouldAccept: (temp) => temp > TEMPERATURE_THRESHOLD }))
.to(createThrottleTransfer<number>({ interval: 5000 }))
.finish(createConvertTransfer<number, Alert>({ operator: createMapOperator((temp): Alert => ({ type: 'HIGH_TEMP', value: temp })) }));
alertPipeline.subscribe((alert) => sendNotification(alert));UI/UX Applications
Reactive Forms
Process user input in form fields → DebounceTransfer for autosave or live search → validate (ConditionTransfer) → async transform (MapOperator) → store results.
import {
CompositeTransferBuilder, createDebounceTransfer, createAsyncConditionTransfer,
createAsyncConvertTransfer, createPushStoredChannelTransfer, createAsyncMapOperator,
} from 'transferum';
const searchInput = createDebounceTransfer<string>({ delay: 300 });
const pipeline = CompositeTransferBuilder
.start(searchInput)
.to(createAsyncConditionTransfer<string>({ shouldAccept: async (s) => s.length >= 3 }))
.to(createAsyncConvertTransfer<string, SearchResult[]>({
operator: createAsyncMapOperator(async (query) => await searchAPI(query)),
}), { onLinkError: (e) => console.error(e) })
.finish(createPushStoredChannelTransfer<SearchResult[]>(), {
onLinkError: (e) => console.error(e),
});
pipeline.subscribe((results) => renderSuggestions(results));
searchInput.push('user query');Monitoring & Logging Systems
Metrics Collection
Collect metrics from various sources (server logs, client events) → filter → transform → send to multiple monitoring systems (BridgeMultiSelector for Prometheus, ELK, Sentry simultaneously).
import { createPushStoredChannelTransfer, createBridgeMultiSelector, createPassBridge } from 'transferum';
const metricsChannel = createPushStoredChannelTransfer<Metric>();
const destinations = createBridgeMultiSelector({
bridges: {
prometheus: createPassBridge({ source: metricsChannel, target: prometheusWriter, activated: false }),
elk: createPassBridge({ source: metricsChannel, target: elkWriter, activated: false }),
sentry: createPassBridge({ source: metricsChannel, target: sentryWriter, activated: false }),
},
initialKeys: ['prometheus', 'elk'],
activated: true,
owned: true,
});
metricsChannel.push({ name: 'request_latency', value: 150 });Financial Applications
Stock Market Data Processing
Receive streaming quotes → calculate indicators (MapOperator) → filter by conditions (ConditionTransfer) → transform into trading signals (ConvertTransfer) → execute trades.
import {
CompositeTransferBuilder, createPushChannelTransfer, createPushStoredChannelTransfer, createConvertTransfer,
createConditionTransfer, createAsyncSinkTransfer, createMapOperator,
} from 'transferum';
const quoteStream = createPushChannelTransfer<Quote[]>();
const thresholdChannel = createPushStoredChannelTransfer<number>({ initialValue: 100 });
const indicatorPipeline = CompositeTransferBuilder
.start(quoteStream)
.to(createConvertTransfer<Quote[], TechnicalIndicator>({
operator: createMapOperator((quotes): TechnicalIndicator => ({
value: quotes.reduce((sum, q) => sum + q.price, 0),
symbols: quotes.map((q) => q.symbol),
threshold: thresholdChannel.pull() ?? 0,
})),
}))
.to(createConditionTransfer<TechnicalIndicator>({ shouldAccept: (ind) => ind.value > ind.threshold }))
.to(createConvertTransfer<TechnicalIndicator, TradingSignal>({
operator: createMapOperator((ind) => ({
action: 'BUY',
symbols: ind.symbols,
targetPrice: ind.value,
})),
}))
.finish(createAsyncSinkTransfer<TradingSignal>({
callback: async (signal) => await executeTrade(signal),
}), { onLinkError: (e) => console.error(e) });
quoteStream.push([{ symbol: 'AAPL', price: 150, timestamp: Date.now() }]);
thresholdChannel.push(200);Portfolio Management
Use BridgeSelector to switch between different strategies or data sources.
import { createBridgeSelector, createPassBridge } from 'transferum';
const strategyRouter = createBridgeSelector({
bridges: {
conservative: createPassBridge({ source: marketData, target: conservativeStrategy, activated: false }),
aggressive: createPassBridge({ source: marketData, target: aggressiveStrategy, activated: false }),
balanced: createPassBridge({ source: marketData, target: balancedStrategy, activated: true }),
},
initialKey: 'balanced',
activated: true,
owned: true,
});
// Switch strategy based on market conditions
if (marketVolatility > HIGH_THRESHOLD) {
strategyRouter.select('conservative');
}Comparison with Alternatives
Transferum exists in a rich ecosystem of reactive and stream-processing libraries. This section compares it with popular alternatives to help you make an informed choice.
Transferum vs RxJS
RxJS is the most widely adopted reactive programming library for JavaScript/TypeScript.
| Aspect | Transferum | RxJS |
|----------------------------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|-----------------------------------------------------------------------------------------|
| Bundle size | ~15 KB minified | ~35 KB minified (full), <5 KB (selective imports) |
| Dependencies | Zero | Zero (v7+) |
| Learning curve | Moderate — explicit primitives | Steep — 100+ operators, complex concepts |
| Type inference | Strong — tuple-based pipeline types | Strong — but complex generic chains |
| Sync/Async unify | Built-in — linkTransfers handles both | Manual — from(), toPromise(), firstValueFrom() |
| Pull-based | Native — Pullable, PollingProxy | Limited — mostly push-based |
| Gate/Flow control | Built-in — GateTransfer, BridgeSelector, DisplaceTransfer | Manual — takeUntil(), switchMap(), subjects |
| Resource cleanup | Explicit — destroy() on every transfer | Subscription-based — subscription.unsubscribe() |
| Error handling | Local, non-fatal — onError per transfer, fail-safe polling, typed source | Stream-level — catchError(), retry(), errors terminate stream |
| Undefined handling | Suppressed — undefined never propagates | Propagated — undefined is a valid value |
| Operators | ~10 pure operators (stateless transforms only) | 100+ operators (creation, transformation, filtering, combination, utility) |
| Flow control / Rate limiting | Built-in transfers — DebounceTransfer, ThrottleTransfer, BufferTransfer, GateTransfer (with explicit lifecycle) | Built-in operators — throttle(), buffer(), sample() (state lives in subscription) |
| Operator-equivalent coverage | Many RxJS operators are transfers: debounceTime→DebounceTransfer, filter→ConditionTransfer, merge→MergeTransfer, share→SplitTransfer, takeUntil→GateTransfer, delay→DelayedPushChannelTransfer, switchMap→DisplaceTransfer | All flow control is operator-based — no separate node lifecycle |
| Testing | Simple — fake timers, direct method calls | Complex — TestScheduler, marble diagrams |
| Community | Small — single maintainer | Large — Google, widespread adoption |
Key differences:
- Architecture: RxJS uses
Observable+Operator+Subscriptionmodel. Transferum usesTransfer+Bridge+Builderwith capability flags. - Models, not abstractions: RxJS builds around a single abstraction (
Observable) that serves as source, connection, and handler simultaneously. Transferum builds around distinct models — transfers are nodes, bridges are edges, operators are transforms. Each has a clear role; none tries to be all of them. - Composability: RxJS operators are functions that transform observables. Transferum transfers are objects that can be linked via
linkTransfers()automatically. - Error handling: RxJS errors propagate through the stream and terminate it unless caught with
catchError()/retry(). Transferum errors are local to each transfer — withonError, suppressed and the stream continues; without, the error is visible (exception/rejection) and polling stops (fail-safe). One stage's failure doesn't kill the pipeline. See Error Handling. - Scheduling: RxJS has
Schedulerabstraction (async, asyncSchedule, animationFrame). Transferum hasTicker(RAFTicker, IntervalTicker) for polling.
Code comparison — Conditional routing with runtime switching:
// RxJS
import { Subject, filter, mergeMap, from, catchError, EMPTY, Subscription } from 'rxjs';
const transports = {
sentry: {
level: 'ERROR',
send: (l) => sentryAPI.send(l),
},
elk: {
level: 'WARN',
send: (l) => elkAPI.send(l),
},
prometheus: {
level: 'INFO',
send: (l) => prometheusAPI.send(l),
},
};
// Subscription infrastructure
const source$ = new Subject<LogEntry>();
const subscriptions = new Map<string, Subscription>();
function activateRoute(key: string) {
if (subscriptions.has(key)) return;
const config = transports[key];
if (!config) return; // Guard against unknown keys
const sub = source$.pipe(
// Filter logs by level from config
filter((l) => l.level === config.level),
// Wrap async send in Observable, swallow errors to keep source$ alive
mergeMap((l) =>
from(config.send(l)).pipe(
catchError((e) => {
console.error(`Error in ${key}:`, e);
return EMPTY;
})
)
)
).subscribe();
subscriptions.set(key, sub);
}
function deactivateRoute(key: string) {
subscriptions.get(key)?.unsubscribe();
subscriptions.delete(key);
}
// Initial routes
activateRoute('sentry');
activateRoute('elk');
activateRoute('prometheus');
// Runtime: leave only Sentry and Prometheus
deactivateRoute('elk');
// Activate / deactivate single route
activateRoute('elk');
deactivateRoute('elk');// Transferum
import {
createBridgeMultiSelector, createPassBridge, createPushStoredChannelTransfer,
createConditionTransfer, createAsyncSinkTransfer, createSplitTransfer, linkTransfers,
} from 'transferum';
const source = createPushStoredChannelTransfer<LogEntry>();
const sentryFilter = createConditionTransfer<LogEntry>({ shouldAccept: (l) => l.level === 'ERROR' });
const elkFilter = createConditionTransfer<LogEntry>({ shouldAccept: (l) => l.level === 'WARN' });
const prometheusFilter = createConditionTransfer<LogEntry>({ shouldAccept: (l) => l.level === 'INFO' });
linkTransfers(source, createSplitTransfer<LogEntry>({ targets: [sentryFilter, elkFilter, prometheusFilter] }));
const routes = createBridgeMultiSelector({
bridges: {
sentry: createPassBridge({
source: sentryFilter,
target: createAsyncSinkTransfer<LogEntry>({ callback: async (l) => sentryAPI.send(l), onError: (e) => console.error(e) }),
activated: false,
}),
elk: createPassBridge({
source: elkFilter,
target: createAsyncSinkTransfer<LogEntry>({ callback: async (l) => elkAPI.send(l), onError: (e) => console.error(e) }),
activated: false,
}),
prometheus: createPassBridge({
source: prometheusFilter,
target: createAsyncSinkTransfer<LogEntry>({ callback: async (l) => prometheusAPI.send(l), onError: (e) => console.error(e) }),
activated: false,
}),
},
initialKeys: ['sentry', 'elk', 'prometheus'],
activated: true,
owned: true,
});
// Runtime: leave only Sentry and Prometheus
routes.select(['sentry', 'prometheus']);
// Activate / deactivate single route
routes.check('elk');
routes.uncheck('elk');RxJS handles routing via manual subscription management — a Map to track active subscriptions and explicit
activateRoute/deactivateRoute functions. Transferum's BridgeMultiSelector is a first-class routing object:
declarative bridge list, built-in subscription management (owned: true), and select() / check() / uncheck()
to switch routes at runtime.
Code comparison — Debounced search with error handling and empty-result suppression:
The API call can fail — errors must be logged without killing the stream. Empty result sets must not reach the renderer.
// RxJS
import { fromEvent, from, EMPTY } from 'rxjs';
import { debounceTime, switchMap, filter, map, catchError } from 'rxjs/operators';
fromEvent(searchInput, 'input')
.pipe(
debounceTime(300),
map((e: Event) => (e.target as HTMLInputElement).value),
filter(query => query.length >= 3),
switchMap(query =>
from(searchAPI(query)).pipe(
catchError(e => {
console.error(e);
return EMPTY;
})
)
),
filter(results => results.length > 0)
)
.subscribe(results => render(results));// Transferum
import {
CompositeTransferBuilder, createPushChannelTransfer, createDebounceTransfer, createConditionTransfer,
createDisplaceTransfer, createAsyncConvertTransfer, createAsyncMapOperator, createSinkTransfer,
} from 'transferum';
const input = createPushChannelTransfer<string>();
const pipeline = CompositeTransferBuilder
.start(input)
.to(createDebounceTransfer<string>({ delay: 300 }))
.to(createConditionTransfer<string>({ shouldAccept: q => q.length >= 3 }))
.to(createDisplaceTransfer<string, SearchResult[]>({
factory: () => createAsyncConvertTransfer<string, SearchResult[]>({
operator: createAsyncMapOperator(async (query) => await searchAPI(query)),
onError: (e) => console.error(e),
}),
}))
.to(createConditionTransfer<SearchResult[]>({ shouldAccept: results => results.length > 0 }))
.finish(createSinkTransfer<SearchResult[]>({
callback: results => render(results),
}), { owned: true });
input.push(query); // manually push, or integrate with DOM eventTransferum vs Most.js
Most.js is a lightweight, high-performance FRP library.
| Aspect | Transferum | Most.js |
|-------------------|---------------------------------------------------|--------------------------|
| Bundle size | ~15 KB | ~7 KB |
| Status | Active (2026) | Maintenance mode (2020+) |
| Async support | Built-in async transfers | Native async event loop |
| Pull-based | Yes — Pullable, PollingProxy | No — push-only |
| TypeScript | First-class, strict types | Community typings |
| Operators | ~10 (pure transforms; flow control via transfers) | ~40 |
Most.js excels in raw performance for push-based streams but lacks Transferum's pull-based primitives and unified sync/async model.
Transferum vs Bacon.js / Kefir
Bacon.js and Kefir are Functional Reactive Programming (FRP) libraries with Property (stateful) and EventStream (stateless) abstractions.
| Aspect | Transferum | Bacon.js / Kefir |
|--------------------|------------------------------------------------------|--------------------------------------|
| State model | Explicit — ProxyReference<T> per transfer | Implicit — Property holds state |
| Stream types | Capability flags (isSubscribable, isPullable) | Two types: EventStream, Property |
| Error handling | Local, non-fatal — onError per transfer, fail-safe | Error events terminate stream |
| Async | First-class async transfers | Via fromPromise() |
| Bundle size | ~15 KB | ~12 KB (Bacon), ~8 KB (Kefir) |
| Status | Active | Bacon: maintenance, Kefir: archived |
Transferum's capability flags provide more granularity than the two-type model, allowing fine-grained control over data flow mechanics. Unlike Bacon.js/Kefir's runtime-only EventStream / Property distinction, Transferum's flags are compile-time type literals — TypeScript knows which methods each transfer exposes without runtime checks.
Transferum vs Node Streams / WHATWG Streams
Node Streams and WHATWG Streams (the web standard) both address data piping through typed stream objects.
| Aspect | Transferum | Node Streams / WHATWG Streams |
|--------------------|-------------------------------------------------------------------------|--------------------------------------------------------------------------------------------------------------------------|
| Type system | Capabilities are composable — mix push + pull + gate + poll in one type | 4 fixed roles (Readable / Writable / Duplex / Transform) — no composition of capabilities within a single stream |
| Model | Capabilities, not stream types | Fixed stream types — a new role means a new class |
| Push + Pull | Both first-class — isPushable / isPullable as flags | Readable = pull-oriented, Writable = push-oriented |
| Gate / Routing | Built-in — GateTransfer, BridgeSelector | Manual — pipe chains, no routing primitive |
| Error handling | Local, non-fatal — per-transfer onError | Stream-level — errors propagate and can destroy stream |
| Lifecycle | Explicit destroy() on every transfer | destroy() / cancel() — inconsistent across impls |
| TypeScript | First-class — computed interfaces from flags | Community typings (Node), limited generics (WHATWG) |
Node Streams and WHATWG Streams solve piping well, but they introduce a new stream type for each role (Readable, Writable, Duplex, Transform). Transferum introduces capabilities instead — a transfer declares what it can do via flags, and the type system computes the correct interface. Adding a new capability doesn't require a new stream class; it requires a new flag. This scales combinatorially better.
Quick Comparison Table
| Feature | Transferum | RxJS | Most.js | Bacon.js | Kefir | Node Streams | WHATWG Streams |
|----------------------------|-------------------------------|--------------------------------------|-------------------|-------------------|-------------------|-----------------|-----------------|
| Bundle size (minified) | ~15 KB | ~35 KB | ~7 KB | ~12 KB | ~8 KB | Built-in | Built-in |
| Dependencies | 0 | 0 | 0 | 0 | 0 | 0 | 0 |
| Pull-based | ✓ | ✗ | ✗ | ✗ | ✗ | ✓ (Readable) | ✓ (Readable) |
| Push-based | ✓ | ✓ | ✓ | ✓ | ✓ | ✓ (Writable) | ✓ (Writable) |
| Sync/Async unify | ✓ | Partial | Partial | Partial | Partial | Partial | Partial |
| Built-in polling | ✓ (Ticker) | ✗ (manual) | ✗ | ✗ | ✗ | ✗ | ✗ |
| Gate/Flow control | ✓ (GateTransfer, Bridge) | Manual | Manual | Manual | Manual | Manual | Manual |
| Error handling | Local, non-fatal, fail-safe | Stream-level (catchError, retry) | Stream terminates | Stream terminates | Stream terminates | Stream-level | Stream-level |
| Undefined suppression | ✓ | ✗ | ✗ | ✗ | ✗ | ✗ | ✗ |
| Operators count | ~10 (pure transforms) | 100+ | ~40 | ~60 | ~50 | ~15 (Transform) | ~10 (Transform) |
| Flow-control as nodes | ✓ (transfers with lifecycle) | ✗ (operators only) | ✗ | ✗ | ✗ | ✗ | ✗ |
| TypeScript support | Excellent | Excellent | Good | Fair | Fair | Fair | Fair |
| Community size | Small | Very large | Medium | Small | Small | Large | Medium |
| Maintenance status | Active | Active | Maintenance | Maintenance | Archived | Active | Active |
When to Choose Transferum
Transferum is ideal for:
- TypeScript-first projects — Strict type inference. Capability flags are compile-time literals (
isPushable: true, notboolean), so TypeScript knows which methods each transfer exposes —push(),subscribe(),pull()are type-guaranteed without casts or runtime checks. - Mixed sync/async pipelines — Unified model without manual conversion.
- Pull-based data acquisition — Polling APIs, sensors, or storage with
PollingProxy. - Explicit flow control — Gates, bridges, and selectors for runtime routing.
- Resource-conscious environments — Smaller bundle than RxJS, zero dependencies.
- Resilient error handling — Local, non-fatal error handling: one stage's failure doesn't kill the pipeline, broken pollers stop (fail-safe), no silent swallowing. Per-stage granularity with typed
ErrorHandler<TSource>. - Undefined-free data flow —
undefinednever propagates through the chain — it means "no data," not "empty value," eliminating null-check defects in consumers. - Game development / IoT — Frame-aligned tickers, idle polling, sensor aggregation.
- Infrastructure libraries — Where components interact in multiple ways (push, pull, polling, trigger, read, write) and the interaction contract matters more than the data format.
- Integration layers — Connecting heterogeneous components or services where a typed, capability-based contract is more valuable than a uniform stream interface.
- Complex pipelines with heterogeneous nodes — Where different stages have fundamentally different capabilities (some push-only, some pull-only, some gated, some polling) and these differences should be statically checked, not discovered at runtime.
- Embedded DSLs — Where the library serves as a language for describing an interaction graph with compile-time safety.
Example fit: Real-time dashboard with API polling, debounced user input, conditional routing to multiple visualizations, and unified sync/async data flows.
When to Consider Alternatives
Consider RxJS if:
- You need operators not covered by Transferum's transfers: stream combinators (
combineLatest,zip,withLatestFrom), windowing (bufferCount,windowTime), or retry policies (retryWhen). Many common operators (debounceTime,throttleTime,filter,merge,delay,takeUntil,switchMap) have transfer equivalents — see the comparison table above. - Your team already has RxJS expertise.
- You need advanced scheduling (test scheduler, virtual time).
- You need complex stream combination operators (e.g., combining 5+ distinct data sources with intricate join logic).
- Community support and long-term stability are top priorities.
Consider Most.js if:
- Raw performance is the top priority.
- You need a small, push-only stream library.
- You don't need pull-based or polling primitives.
Consider Bacon.js (or maintaining Kefir) only if:
- You are working on a legacy codebase already locked into these specific FRP abstractions.
Consider Node Streams / WHATWG Streams if:
- You need the web standard or Node.js native streaming protocol for interop.
- Your pipeline is purely I/O-oriented (file, network, compression) and doesn't need push+pull+gate+poll in one system.
- You don't need compile-time capability guarantees.
Installation & Import
npm i transferumAll components are exported from a single entry point:
import {
// Transfers
PushChannelTransfer,
DelayedPushChannelTransfer,
DebounceTransfer,
ThrottleTransfer,
PushStoredChannelTransfer,
BufferTransfer,
ManualBufferTransfer,
ManualFlowTransfer,
GateTransfer,
MergeTransfer,
SplitTransfer,
PollingSourceTransfer,
PollingProxyTransfer,
PollingFlowTransfer,
IdlePollingTransfer,
ChannelTransfer,
StoredChannelTransfer,
SinkTransfer,
WriteTransfer,
ReadTransfer,
ConvertTransfer,
ConditionTransfer,
DisplaceTransfer,
UniversalCompositeTransfer,
// Async transfers
AsyncSinkTransfer,
AsyncWriteTransfer,
AsyncReadTransfer,
AsyncConvertTransfer,
AsyncConditionTransfer,
AsyncPollingSourceTransfer,
AsyncPollingProxyTransfer,
AsyncPollingFlowTransfer,
AsyncIdlePollingTransfer,
AsyncStoredChannelTransfer,
// Operators
TransparentOperator,
MapOperator,
FilterOperator,
ReducerOperator,
GuardOperator,
PipelineOperator,
// Async operators
AsyncMapOperator,
AsyncGuardOperator,
AsyncPipelineOperator,
// Storages
LatestStorage,
QueueStorage,
StackStorage,
// Tickers
RAFTicker,
IntervalTicker,
TickerInterface,
// Helpers
Subscriber,
SubscriptionManager,
StateSubscriptionManager,
ProxyReference,
DisposableSubscriberAdapter,
// Bridges
PassBridge,
TransformBridge,
TransferBridge,
BridgeAggregator,
BridgeSelector,
BridgeMultiSelector,
AsyncTransformBridge,
// Linking
DefaultLinkStrategy,
LinkStrategyInterface,
// Builders
InputPipelineBuilder,
OutputPipelineBuilder,
DuplexPipelineBuilder,
OperatorPipelineBuilder,
// Async builders
AsyncInputPipelineBuilder,
AsyncOutputPipelineBuilder,
AsyncDuplexPipelineBuilder,
AsyncOperatorPipelineBuilder,
// Factories
createPushChannelTransfer,
createDelayedPushChannelTransfer,
createDebounceTransfer,
createThrottleTransfer,
createPushStoredChannelTransfer,
createDefaultLinkStrategy,
// ... all create* functions (including createAsync*)
// Utilities
linkTransfers,
handleError,
// Guards
isPushable,
isPullable,
isSubscribable,
isPollingProxy,
isTriggerable,
isGate,
isAsyncPushable,
isAsyncPullable,
isAsyncPollingProxy,
isAsyncTriggerable,
} from 'transferum';Core Concepts
Architectural Invariants
Transferum is built on a single idea:
Behavior can be described as a composition of independent capabilities that simultaneously determine the type, the implementation, and the rules of interaction.
Everything else — transfers, bridges, builders, operators — follows from this principle. They are its consequences, not the idea itself. The invariants below are the shape these consequences take in code.
One idea, one direction. A single concept — capability — flows through every layer of the library:
Capability flags ↓ Type system (computed interfaces) ↓ Transfer (declares capabilities, implements behavior) ↓ Bridge (inspects capabilities, wires transfers) ↓ Builder (assembles transfers, enforces type compatibility)Operators are not a separate layer — they are stateless transforms used inside transfers. Adding a new capability propagates automatically through all layers. Adding a new transfer class requires only declaring its flags — the type system,
linkTransfers, and builders adapt without changes.
The invariants:
| Invariant | What it means | Where enforced |
|-------------------------------------------------|-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------------------------------------------------------------------------------------------|
| Transfers don't know their neighbors | A transfer defines its own behavior (push, pull, subscribe, gate…) but never references, imports, or checks the class of another transfer. It doesn't know what's upstream or downstream — it only fulfills its own contract. | BaseTransfer and all descendants — no cross-references between transfer classes |
| Bridges don't know concrete implementations | A bridge inspects capability flags, never class names. There is no instanceof chain, no RTTI, no class-name switching. Any transfer with the right capabilities is bridgeable — including ones that don't exist yet. | linkTransfers, LinkStrategyInterface, PassBridge, BridgeSelector, all bridge classes |
| Capabilities define the contract | Flags are the single source of truth. They determine the TypeScript interface (compile-time), the available methods (runtime), the linking strategy (linkTransfers), and the builder type compatibility. One piece of metadata, many readers. | interfaces.ts, types.ts, linkTransfers, LinkStrategyInterface, builders |
| BaseTransfer is minimal | The base class contains only capability declarations — no state, no logic. State lives in BaseStateTransfer<T> (one level down). This prevents the god-object pattern where a base accumulates knowledge of all descendants. | BaseTransfer class hierarchy |
| One class, one behavior | Each transfer models exactly one behavioral concept. There are no combinatorial mega-classes (BufferedStoredGateTransfer). Complex behavior emerges from composition, not from inheritance depth. | All transfer classes |
| Operators are stateless transforms | Operators transform data (apply(input) → output) and hold no pipeline state. They live inside transfers (ConvertTransfer, TransformBridge), not as standalone pipeline nodes. The transfer owns the behavioral contract; the operator owns the data transformation. | OperatorInterface, AsyncOperatorInterface |
| destroy() is universal | Every resource has an explicit cleanup path — destroy() for transfers, bridges, and subscription managers; stop() for tickers; unsubscribe() for subscriptions. Builders track owned resources and clean them up in one destroy() call. No resource lacks a cleanup path. | DisposableInterface, all transfers, all bridges, builders; TickerInterface.stop() |
| undefined never propagates | undefined means "no data," not "empty value." It is suppressed at SubscriptionManager.sendState() — subscribers are never notified with undefined. Use null for explicit empty markers. | SubscriptionManager, all transfers |
Why these matter. A new transfer, bridge, or operator that respects these invariants integrates without touching existing code — the architecture stays coherent as it grows.
Capability Flags System
Each transfer implements CommunicationContractInterface — a set of boolean flags. Each flag is both a runtime value and a compile-time guarantee: when a flag is true, the corresponding method is part of the transfer's TypeScript type — TypeScript knows it exists without any casts or runtime checks.
| Flag | Methods | Description |
|-------------------|----------------------------------------------------------------------------------|--------------------------------------------------------------|
| isInput | — | Can accept data from outside (acts as an input) |
| isOutput | — | Can yield data (acts as an output) |
| isDuplex | — | Both input and output simultaneously (isInput && isOutput) |
| isPushable | push(data) | Data can be pushed into the transfer |
| isPullable | pull() | Data can be read from the transfer |
| isSubscribable | subscribe(handler) | The transfer can be subscribed to |
| isTriggerable | trigger() | Has a manual emission trigger |
| isGate | activate() / deactivate() / toggle() / active / onStateChange(handler) | Flow control (on/off) + state subscription |
| isPollingSource | — | Has internal source polling |
| isPollingProxy | setFetcher() / clearFetcher() | Polls the previous node in the chain |
Asynchronous flags:
| Flag
