Streaming Live Data
Websocket subscriptions, provisional bars, reconnection and backfilling the gap.
Streaming turns a one-shot run into a live evaluation. Call stream on a runtime instance to get an emitter, attach handlers, and connect. The stream emits open, bar, tick, error and close. Bar events carry a meta object with a closed flag; ticks carry the raw trade update when the provider supplies one.
Always warm the buffer before connecting. An RSI needs at least its length in bars before it means anything, and an EMA needs several multiples of its length to converge. Calling load first fills the history buffer so the first live value is correct rather than a number that drifts into correctness over the next hour.
Handle provisional bars deliberately. The runtime will evaluate on every update if you let it, which is right for a live readout and wrong for a decision. Push with the provisional flag to update a display value, and act only when the closed flag is true. Any trading logic that fires on a provisional bar will eventually fire on a wick that unwinds before the close.
Reconnection is built in but bounded. The reconnect option takes a retry count, a base backoff and a jitter flag, and exponential backoff with jitter is the correct default when a venue drops every client at once. After a successful reconnect the runtime backfills the bars it missed and replays them through the indicator buffer in order, so a thirty second outage does not leave a hole in the series.
Shut down cleanly. Call disconnect on SIGINT and SIGTERM so the socket closes and any provider-side subscription is released. A process that exits without disconnecting will often be rate-limited on its next connection attempt, which is a confusing failure to debug the first time it happens.
Example
Stream live bars over a websocket
import { AlgoBeamTS, Provider, indicators, type Bar } from 'algobeam-ts'
const algobeam = new AlgoBeamTS(Provider.Binance, 'SOLUSDT', '1m', 500)
await algobeam.load() // warm the buffer so indicators are not cold-started
const rsi = algobeam.series(indicators.rsi({ length: 14 }))
const stream = algobeam.stream({
transport: 'websocket',
reconnect: { retries: 8, backoffMs: 750, jitter: true },
})
stream.on('open', () => console.log('[algobeam] socket open'))
stream.on('bar', (bar: Bar, meta: { closed: boolean }) => {
const value = rsi.push(bar, { provisional: !meta.closed })
if (!meta.closed) return
const stamp = new Date(bar.time * 1000).toISOString()
console.log(`${bar.symbol} ${stamp} close=${bar.close} rsi=${value.toFixed(1)}`)
if (value > 70) console.warn(`[algobeam] ${bar.symbol} stretched at ${value.toFixed(1)}`)
})
stream.on('error', (error: unknown) => console.error('[algobeam] stream error', error))
stream.on('close', ({ code }: { code: number }) => console.warn('[algobeam] closed', code))
await stream.connect()
process.on('SIGINT', () => void stream.disconnect())Copy it, change one input, and run it again — the numbers are deterministic, so a difference in the output is always a difference you made.