fin-stream
Real-time market data plumbing for Rust. Four exchange formats in, one exact-decimal tick out, through a lock-free ring into bars, books and health checks.
OrderBook·HealthMonitor·TickReplayer
$ cargo add fin-streamcrates.io has 2.4.2; main (2.11.0) has the examples and newer modules. Clone it next to fin-primitives to build.
cargo run --example tape at its recorded pace.Four examples, no keys, no network
Each example is deterministic: seeded generators or a recorded file. In a terminal they stream in real time; piped, they print at once. What you see below is their actual output, captured from the main branch.
The tape
A feed thread builds each trade in its venue's wire format, normalizes it with TickNormalizer and pushes it into SpscRing<NormalizedTick, 64>. The main thread pops, prints the tape with the price colored by tick direction, and rolls the same ticks into 2-second bars with OhlcvAggregator. Alpaca and Polygon do not report the aggressor, so their side is n/a, not a guess.
$ cargo run --example tapeBTC-USD time and sales, 4 venues, 2s bars venue JSON -> TickNormalizer -> SpscRing<NormalizedTick, 64> -> tape + OhlcvAggregator time (exch) venue side price size latency 14:30:00.005 coinbase buy · 64,250.19 0.1112 █ 12ms 14:30:00.242 coinbase sell ▲ 64,256.87 0.2008 █▊ 17ms 14:30:00.285 binance sell ▼ 64,255.51 0.1375 █▎ 45ms 14:30:00.398 binance sell ▼ 64,255.22 0.7251 ██████▌ 39ms 14:30:00.414 binance buy ▲ 64,259.81 0.3319 ███ 43ms 14:30:00.722 alpaca n/a ▲ 64,267.52 0.0583 ▌ 27ms 14:30:01.051 binance buy ▲ 64,270.47 0.0148 ▏ 45ms 14:30:01.171 binance buy ▼ 64,264.65 0.0458 ▍ 43ms 14:30:01.347 alpaca n/a ▲ 64,268.54 0.0491 ▌ 26ms 14:30:01.470 coinbase buy ▲ 64,270.66 0.1478 █▍ 17ms 14:30:01.585 polygon n/a ▲ 64,272.65 0.0088 ▏ 25ms 14:30:01.596 polygon n/a ▼ 64,261.83 0.0025 ▏ 30ms 14:30:01.730 binance buy ▲ 64,267.25 0.1924 █▊ 44ms 14:30:01.751 binance buy ▼ 64,262.73 0.1248 █▏ 43ms 14:30:01.805 polygon n/a ▲ 64,267.78 0.4513 ████ 30ms 14:30:01.985 coinbase sell ▼ 64,254.78 0.1446 █▎ 15ms 14:30:00 bar o 64,250.19 h 64,272.65 l 64,250.19 c 64,254.78 2.75 BTC +4.59 14:30:02.076 coinbase sell ▼ 64,254.51 0.1404 █▎ 15ms 14:30:02.272 binance buy ▲ 64,256.12 0.3661 ███▎ 39ms 14:30:02.299 binance buy ▲ 64,260.83 0.9429 ████████▌ 44ms 14:30:02.335 polygon n/a ▲ 64,264.61 0.2904 ██▋ 25ms 14:30:02.340 binance buy ▼ 64,263.10 0.2889 ██▋ 46ms 14:30:02.665 binance sell ▼ 64,253.12 0.2262 ██ 42ms 14:30:02.896 alpaca n/a ▲ 64,257.10 0.1227 █▏ 26ms 14:30:03.057 binance buy ▲ 64,258.46 0.0495 ▌ 46ms 14:30:03.082 polygon n/a ▲ 64,259.47 0.7773 ███████ 30ms 14:30:03.287 alpaca n/a ▲ 64,259.67 0.0884 ▊ 19ms 14:30:03.488 binance buy ▲ 64,266.70 0.0470 ▍ 38ms 14:30:03.497 alpaca n/a ▼ 64,263.99 0.2778 ██▌ 22ms 14:30:03.607 polygon n/a ▼ 64,260.43 0.0992 ▉ 25ms 14:30:03.665 binance buy ▲ 64,261.87 0.4262 ███▉ 44ms 14:30:03.863 binance sell ▲ 64,263.44 0.1553 █▍ 45ms 14:30:03.932 binance buy ▲ 64,266.13 0.2023 █▉ 38ms 14:30:03.989 polygon n/a ▲ 64,272.26 0.0251 ▎ 30ms 14:30:03.998 polygon n/a ▼ 64,267.94 0.1328 █▎ 26ms 14:30:02 bar o 64,254.51 h 64,272.26 l 64,253.12 c 64,267.94 4.66 BTC +13.43 14:30:04.006 binance buy ▲ 64,269.08 0.6293 █████▋ 40ms 14:30:04.248 coinbase sell ▼ 64,264.41 0.0345 ▎ 18ms 14:30:04.364 binance buy ▼ 64,260.37 0.1389 █▎ 38ms 14:30:04.383 coinbase buy ▲ 64,266.46 0.1183 █▏ 19ms 14:30:04.762 coinbase buy ▲ 64,269.99 0.1807 █▋ 18ms 14:30:04.772 binance sell ▼ 64,269.81 0.1393 █▎ 40ms 14:30:04 bar o 64,269.08 h 64,269.99 l 64,260.37 c 64,269.81 1.24 BTC +0.73 40 ticks alpaca 5 binance 19 coinbase 8 polygon 8 3 bars range 64,250.19 .. 64,272.65 vwap 64,261.76 volume 8.6460 BTC ring peak 40/64 slots in use, 0 full-ring retries, 1 ms wall clock (instant mode)
One trade, four wire formats
Binance sends strings and a maker flag, Coinbase spells out the side with an ISO-8601 time, Alpaca sends JSON numbers, Polygon a nanosecond epoch. All four become the same NormalizedTick; malformed payloads come back as a typed StreamError.
$ cargo run --example normalizeTickNormalizer the same BTC-USD trade as four venues send it binance strings for numbers; m = buyer is maker {"T":1790260200112,"e":"trade","m":false,"p":"64250.10","q":"0.01200","t":4120001} -> price 64250.10 qty 0.01200 side Buy exch_ts 14:30:00.112 coinbase side spelled out; ISO-8601 time {"price":"64250.10","side":"buy","size":"0.01200000","time":"2026-09-24T14:30:00.112Z","trade_id":"91200337"} -> price 64250.10 qty 0.01200000 side Buy exch_ts 14:30:00.112 alpaca JSON numbers; no side {"T":"t","i":88410023,"p":64250.1,"s":0.012,"t":"2026-09-24T14:30:00.112Z"} -> price 64250.1 qty 0.012 side None exch_ts 14:30:00.112 polygon nanosecond epoch; no side {"ev":"XT","i":"55c1e07a","p":64250.1,"s":0.012,"t":1790260200112000000} -> price 64250.1 qty 0.012 side None exch_ts 14:30:00.112 rejected typed errors, never a panic binance {"T":1790260200112,"q":"0.5"} -> Tick parse error from Binance: missing field 'p' coinbase {"price":"-3.10","side":"buy","size":"1"} -> Invalid tick: price must be positive, got -3.10 alpaca {"p":"sixty-four thousand","s":1} -> Tick parse error from Alpaca: field 'p' parse error: Invalid decimal: unknown character
A feed goes quiet
HealthMonitor gets a heartbeat per tick and a sweep every 500 ms of simulated time. Coinbase stalls, turns stale after 2 s, opens its circuit after three stale checks, and closes it on the next heartbeat. Polygon beats slowly but has its own 5 s threshold, so it never alarms.
$ cargo run --example feed_healthHealthMonitor stale after 2 s quiet (polygon: 5 s), circuit opens after 3 stale checks binance ││││││││││││││││││││││││││││││││││││││││││││││││ healthy 96 beats, last 0.0s ago coinbase │││││││││····▒▒█████████││││││││││││││││││││││││ healthy 33 beats, last 0.0s ago alpaca ·││·││·││·││·││·││·││·││·││·││·│···▒▒███████████ open 21 beats, last 8.2s ago polygon ·····│·····│·····│·····│·····│·····│·····│·····│ healthy 8 beats, last 0.0s ago 0s 5s 10s 15s 20s │ heartbeat · quiet ▒ stale █ circuit open t+ 7.0s coinbase Feed 'coinbase' is stale: last tick was 2500ms ago (threshold: 2000ms) t+ 8.0s coinbase circuit OPEN after 3 stale checks t+12.5s coinbase heartbeat received, circuit closed t+18.0s alpaca Feed 'alpaca' is stale: last tick was 2250ms ago (threshold: 2000ms) t+19.0s alpaca circuit OPEN after 3 stale checks end state 3 healthy, 1 stale, 0 unknown; ratio_healthy 0.75
Recorded data, live code path
TickReplayer streams a 600-tick NDJSON file through the same TickSource trait a live feed implements, honoring the recorded gaps at a speed multiplier. Each 30-second bar is drawn as it closes, on a fixed price axis.
$ cargo run --example replayTickReplayer examples/data/btc-usd-10m.ndjson at full speed, 30s bars bar 63,945 64,202 64,459 close vs open volume ┬────────────────────────────────────────────┬ 14:30:00 ──────██── 64,193.68 -0.01% 7.45 14:30:30 ████████████──── 64,306.78 +0.16% 8.72 14:31:00 ────██████████████ 64,148.94 -0.08% 8.61 14:31:30 ────███████████████─── 64,321.89 +0.19% 7.22 14:32:00 ────████─────── 64,353.62 +0.24% 6.63 14:32:30 ────────██── 64,371.11 +0.26% 6.75 14:33:00 ─█████──── 64,413.24 +0.33% 7.95 14:33:30 ────█████ 64,337.80 +0.21% 6.49 14:34:00 ███████─── 64,276.51 +0.12% 6.93 14:34:30 ────████████── 64,357.60 +0.24% 6.81 14:35:00 ██████████─ 64,430.49 +0.36% 7.86 14:35:30 ─████████████████── 64,268.38 +0.10% 8.37 14:36:00 ─████████████████ 64,407.62 +0.32% 8.54 14:36:30 █████████─── 64,326.93 +0.19% 8.36 14:37:00 ─────███── 64,378.12 +0.27% 8.18 14:37:30 ───█████████── 64,273.49 +0.11% 6.86 14:38:00 ───────────█████ 64,186.34 -0.02% 7.74 14:38:30 ─────███████─── 64,266.58 +0.10% 7.00 14:39:00 ─████──── 64,325.51 +0.19% 7.58 14:39:30 ████████ 64,222.57 +0.03% 6.68 600 ticks 20 bars, 0 parse errors, 9 ms wall clock for 600 s of market time
How the pieces fit
Every stage takes and returns plain values, so they compose in whatever order your system needs. TickReplayer implements the same TickSource trait a live feed would, so strategy code runs unchanged on recorded data.
Exact prices, typed errors
Prices and sizes are rust_decimal::Decimal. Every fallible call, constructors included, returns Result<_, StreamError>; the crate denies unwrap, expect and panic in Clippy.
A ring that does not allocate
SpscRing<T, N> allocates its slots once. Push and pop are atomic loads and stores with no lock; a full ring returns Err instead of blocking. Unsafe code lives only here, behind a safe API.
Measured, not promised
Criterion medians from benches/tick_hot_path.rs on an Intel Core i7-13700KF (Windows 11, rustc 1.91), single-threaded, 2026-09-25. Reproduce with cargo bench --bench tick_hot_path.
| Benchmark | Median | One iteration |
|---|---|---|
ring_push_pop_u64 | 2.3 ns | push a u64 into the ring and pop it back |
order_book_best_levels | 11.7 ns | best bid and ask from a 10-level book |
ring_push_pop_normalized_tick | 52.8 ns | build a NormalizedTick, push, pop |
order_book_apply_delta | 64.8 ns | apply one bid-level update |
ohlcv_feed_same_window | 101 ns | build a tick and feed an open 1-minute bar |
ohlcv_feed_bar_completion | 200 ns | build a tick that closes a 1-second bar |
tick_normalize_binance | 714 ns | clone a Binance JSON payload and normalize it |
tick_normalize_coinbase | 831 ns | build a Coinbase JSON payload and normalize it |
Normalization dominates, mostly JSON handling. At about 0.7 to 0.8 µs per tick that is over a million normalized ticks per second on one core. These are microbenchmarks on one thread, not an end-to-end cross-thread throughput test.
Running in two minutes
main depends on fin-primitives through a relative path, so clone both side by side.
Clone
git clone https://github.com/Mattbusel/fin-primitives git clone https://github.com/Mattbusel/fin-stream cd fin-stream
Or cargo add fin-stream for the 2.4.2 release on crates.io.
Watch the tape
cargo run --example tape
Then normalize, feed_health and replay. NO_COLOR=1 turns color off.
Check it
cargo test --doc cargo test --test '*'
The doctests compile and run every Rust block in the README as well as the crate docs.