fin-stream
Rust · Tokio · MIT

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.

TickNormalizer→SpscRing→OhlcvAggregator
OrderBook·HealthMonitor·TickReplayer
$ cargo add fin-stream

crates.io has 2.4.2; main (2.11.0) has the examples and newer modules. Clone it next to fin-primitives to build.

BTC-USD time and sales · 4 venues · 2s bars14:30:00.000
Replaying the real output of 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 tape
examples/tape.rs
$ cargo run --example tape
  BTC-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 normalize
examples/normalize.rs
$ cargo run --example normalize
  TickNormalizer  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_health
examples/feed_health.rs
$ cargo run --example feed_health
  HealthMonitor  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 replay
examples/replay.rs
$ cargo run --example replay
  TickReplayer  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.

Architecture: WsManager, TickReplayer and SyntheticMarketGenerator feed TickNormalizer and FeedAggregator, an SpscRing hands ticks to OhlcvAggregator, OrderBook, the normalizers and analytics, and HealthMonitor watches every feed.

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.

BenchmarkMedianOne iteration
ring_push_pop_u642.3 nspush a u64 into the ring and pop it back
order_book_best_levels11.7 nsbest bid and ask from a 10-level book
ring_push_pop_normalized_tick52.8 nsbuild a NormalizedTick, push, pop
order_book_apply_delta64.8 nsapply one bid-level update
ohlcv_feed_same_window101 nsbuild a tick and feed an open 1-minute bar
ohlcv_feed_bar_completion200 nsbuild a tick that closes a 1-second bar
tick_normalize_binance714 nsclone a Binance JSON payload and normalize it
tick_normalize_coinbase831 nsbuild 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.