use subms_merge_iterator::MergeIterator;
const VENUE_TAPES: [&[i64]; 3] = [
&[8_000, 9_100, 9_400, 9_800],
&[8_500, 9_300, 9_600],
&[9_050, 9_450, 9_900],
];
const BID_LADDERS: [&[i64]; 2] = [&[10_120, 10_105, 10_101, 10_095], &[10_118, 10_110, 10_099]];
const REFERENCE_LEVELS: [&[(&str, Option<&str>)]; 3] = [
&[
("AAPL", Some("listed")),
("ENRN", Some("listed")),
("MSFT", Some("listed")),
],
&[("AAPL", Some("listed-adr"))],
&[("ENRN", None)],
];
#[cfg(any(feature = "dedup", feature = "priority"))]
const PRICE_FLUSHED: &[(&str, i64)] = &[("AAPL", 150), ("MSFT", 300)];
#[cfg(any(feature = "dedup", feature = "priority"))]
const PRICE_MEMTABLE: &[(&str, i64)] = &[("AAPL", 152)];
fn main() {
println!(
"market-data store: {} venue tapes, {} bid ladders, {} reference levels",
VENUE_TAPES.len(),
BID_LADDERS.len(),
REFERENCE_LEVELS.len()
);
base_consolidated_tape();
#[cfg(feature = "seek-to")]
session_window_scan();
#[cfg(feature = "reverse")]
walk_bid_ladder_down();
#[cfg(feature = "tombstones")]
resolve_reference_rows();
#[cfg(feature = "dedup")]
compact_last_prices();
#[cfg(feature = "priority")]
memtable_wins_the_read();
}
fn tapes() -> Vec<std::iter::Copied<std::slice::Iter<'static, i64>>> {
VENUE_TAPES.iter().map(|t| t.iter().copied()).collect()
}
fn base_consolidated_tape() {
println!("\n== base: consolidated trade tape ==");
let merge = MergeIterator::new(tapes());
println!(" live venues: {}", merge.live_streams());
println!(" earliest trade: {:?}", merge.peek());
let tape: Vec<i64> = merge.collect();
println!(" {} trades in order: {tape:?}", tape.len());
assert_eq!(tape.len(), 10, "every trade appears once");
assert!(
tape.windows(2).all(|w| w[0] <= w[1]),
"the tape stays chronological"
);
}
#[cfg(feature = "seek-to")]
fn session_window_scan() {
use subms_merge_iterator::SeekableMergeIterator;
println!("\n== seek-to: one session window out of the tape ==");
let (open, close) = (9_300, 9_800);
let mut scan = SeekableMergeIterator::new(tapes());
scan.seek(&open);
scan.set_upper_bound(close);
let window: Vec<i64> = scan.collect();
println!(" window [{open}, {close}): {window:?}");
assert_eq!(
window,
vec![9_300, 9_400, 9_450, 9_600],
"half-open: the close tick is excluded"
);
}
#[cfg(feature = "reverse")]
fn walk_bid_ladder_down() {
use subms_merge_iterator::ReverseMergeIterator;
println!("\n== reverse: walk the consolidated bid ladder down ==");
let ladders: Vec<_> = BID_LADDERS.iter().map(|l| l.iter().copied()).collect();
let mut book = ReverseMergeIterator::new(ladders);
println!(" best bid across venues: {:?}", book.peek());
let limit = 10_100;
book.seek_for_prev(&10_110);
book.set_lower_bound(limit);
let fillable: Vec<i64> = book.collect();
println!(" levels from 10110 down to the {limit} limit: {fillable:?}");
assert_eq!(
fillable,
vec![10_110, 10_105, 10_101],
"descending, and the lower bound is inclusive"
);
}
#[cfg(feature = "tombstones")]
fn resolve_reference_rows() {
use subms_merge_iterator::{TombstoneEntry, TombstoneMergeIterator};
println!("\n== tombstones: resolve the reference rows ==");
let levels: Vec<std::vec::IntoIter<TombstoneEntry<&str, &str>>> = REFERENCE_LEVELS
.iter()
.map(|rows| {
rows.iter()
.map(|&(sym, status)| match status {
Some(s) => TombstoneEntry::live(sym, s),
None => TombstoneEntry::tombstone(sym),
})
.collect::<Vec<_>>()
.into_iter()
})
.collect();
let resolved: Vec<(&str, &str)> = TombstoneMergeIterator::new(levels)
.map(|e| (e.key, e.value.unwrap()))
.collect();
println!(" live instruments: {resolved:?}");
assert_eq!(
resolved,
vec![("AAPL", "listed-adr"), ("MSFT", "listed")],
"the delisted symbol is shadowed out, AAPL takes the newer row"
);
}
#[cfg(feature = "dedup")]
fn compact_last_prices() {
use subms_merge_iterator::{DedupEntry, DedupMergeIterator};
println!("\n== dedup: compact the last-price rows ==");
let flushed = PRICE_FLUSHED
.iter()
.map(|&(k, v)| DedupEntry::new(k, v))
.collect::<Vec<_>>();
let memtable = PRICE_MEMTABLE
.iter()
.map(|&(k, v)| DedupEntry::new(k, v))
.collect::<Vec<_>>();
let compacted: Vec<(&str, i64)> =
DedupMergeIterator::new([flushed.into_iter(), memtable.into_iter()])
.map(|e| (e.key, e.value))
.collect();
println!(" compacted last prices: {compacted:?}");
assert_eq!(compacted, vec![("AAPL", 152), ("MSFT", 300)]);
}
#[cfg(feature = "priority")]
fn memtable_wins_the_read() {
use subms_merge_iterator::{PriorityEntry, PriorityMergeIterator, PrioritySource};
println!("\n== priority: the memtable is authoritative ==");
let memtable = PrioritySource::new(
100,
PRICE_MEMTABLE
.iter()
.map(|&(k, v)| PriorityEntry::new(k, v))
.collect::<Vec<_>>()
.into_iter(),
);
let flushed = PrioritySource::new(
10,
PRICE_FLUSHED
.iter()
.map(|&(k, v)| PriorityEntry::new(k, v))
.collect::<Vec<_>>()
.into_iter(),
);
let view: Vec<(&str, i64)> = PriorityMergeIterator::new([memtable, flushed])
.map(|e| (e.key, e.value))
.collect();
println!(" resolved read view: {view:?}");
assert_eq!(
view,
vec![("AAPL", 152), ("MSFT", 300)],
"the memtable wins AAPL despite being registered first"
);
}