use crate::error::FinError;
use crate::types::NanoTimestamp;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum LatencyPhase {
SubmitToAck,
AckToFill,
FillToBookUpdate,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct LatencySample {
pub phase: LatencyPhase,
pub latency_ns: i64,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PhaseStats {
pub count: usize,
pub min_ns: i64,
pub max_ns: i64,
pub mean_ns: f64,
pub p50_ns: i64,
pub p95_ns: i64,
pub p99_ns: i64,
}
#[derive(Debug, Default)]
pub struct OrderLatencyTracker {
pending: std::collections::HashMap<String, OrderTimestamps>,
submit_to_ack: Vec<i64>,
ack_to_fill: Vec<i64>,
fill_to_book: Vec<i64>,
}
#[derive(Debug, Clone, Default)]
struct OrderTimestamps {
submit: Option<i64>,
ack: Option<i64>,
fill: Option<i64>,
}
impl OrderLatencyTracker {
pub fn new() -> Self {
Self::default()
}
pub fn record_submit(&mut self, order_id: impl Into<String>, ts: NanoTimestamp) {
let mut ots = OrderTimestamps::default();
ots.submit = Some(ts.nanos());
self.pending.insert(order_id.into(), ots);
}
pub fn record_ack(&mut self, order_id: &str, ts: NanoTimestamp) -> Result<(), FinError> {
let rec = self
.pending
.get_mut(order_id)
.ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
let submit = rec.submit.ok_or_else(|| {
FinError::InvalidInput(format!("submit not recorded for '{order_id}'"))
})?;
let ack_ns = ts.nanos();
if ack_ns < submit {
return Err(FinError::InvalidInput(format!(
"ack timestamp before submit for '{order_id}'"
)));
}
rec.ack = Some(ack_ns);
self.submit_to_ack.push(ack_ns - submit);
Ok(())
}
pub fn record_fill(&mut self, order_id: &str, ts: NanoTimestamp) -> Result<(), FinError> {
let rec = self
.pending
.get_mut(order_id)
.ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
let ack = rec.ack.ok_or_else(|| {
FinError::InvalidInput(format!("ack not recorded for '{order_id}'"))
})?;
let fill_ns = ts.nanos();
if fill_ns < ack {
return Err(FinError::InvalidInput(format!(
"fill timestamp before ack for '{order_id}'"
)));
}
rec.fill = Some(fill_ns);
self.ack_to_fill.push(fill_ns - ack);
Ok(())
}
pub fn record_book_update(
&mut self,
order_id: &str,
ts: NanoTimestamp,
) -> Result<(), FinError> {
let rec = self
.pending
.remove(order_id)
.ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
let fill = rec.fill.ok_or_else(|| {
FinError::InvalidInput(format!("fill not recorded for '{order_id}'"))
})?;
let book_ns = ts.nanos();
if book_ns < fill {
return Err(FinError::InvalidInput(format!(
"book_update timestamp before fill for '{order_id}'"
)));
}
self.fill_to_book.push(book_ns - fill);
Ok(())
}
pub fn stats(&self, phase: LatencyPhase) -> Option<PhaseStats> {
let samples = match phase {
LatencyPhase::SubmitToAck => &self.submit_to_ack,
LatencyPhase::AckToFill => &self.ack_to_fill,
LatencyPhase::FillToBookUpdate => &self.fill_to_book,
};
if samples.is_empty() {
return None;
}
let mut sorted = samples.clone();
sorted.sort_unstable();
let n = sorted.len();
let min_ns = *sorted.first().unwrap_or(&0);
let max_ns = *sorted.last().unwrap_or(&0);
let mean_ns = sorted.iter().map(|&v| v as f64).sum::<f64>() / n as f64;
let p50_ns = percentile_ns(&sorted, 50);
let p95_ns = percentile_ns(&sorted, 95);
let p99_ns = percentile_ns(&sorted, 99);
Some(PhaseStats {
count: n,
min_ns,
max_ns,
mean_ns,
p50_ns,
p95_ns,
p99_ns,
})
}
pub fn pending_count(&self) -> usize {
self.pending.len()
}
pub fn samples(&self, phase: LatencyPhase) -> &[i64] {
match phase {
LatencyPhase::SubmitToAck => &self.submit_to_ack,
LatencyPhase::AckToFill => &self.ack_to_fill,
LatencyPhase::FillToBookUpdate => &self.fill_to_book,
}
}
}
fn percentile_ns(sorted: &[i64], p: usize) -> i64 {
if sorted.is_empty() {
return 0;
}
let idx = ((p * sorted.len()) / 100).min(sorted.len() - 1);
sorted[idx]
}
#[cfg(test)]
mod tests {
use super::*;
fn ts(n: i64) -> NanoTimestamp {
NanoTimestamp::new(n)
}
fn full_lifecycle(tracker: &mut OrderLatencyTracker, id: &str, t0: i64, t1: i64, t2: i64, t3: i64) {
tracker.record_submit(id, ts(t0));
tracker.record_ack(id, ts(t1)).unwrap();
tracker.record_fill(id, ts(t2)).unwrap();
tracker.record_book_update(id, ts(t3)).unwrap();
}
#[test]
fn test_single_order_lifecycle() {
let mut tracker = OrderLatencyTracker::new();
full_lifecycle(&mut tracker, "o1", 1000, 2000, 4000, 5000);
let s2a = tracker.stats(LatencyPhase::SubmitToAck).unwrap();
assert_eq!(s2a.count, 1);
assert_eq!(s2a.p50_ns, 1000);
let a2f = tracker.stats(LatencyPhase::AckToFill).unwrap();
assert_eq!(a2f.p50_ns, 2000);
let f2b = tracker.stats(LatencyPhase::FillToBookUpdate).unwrap();
assert_eq!(f2b.p50_ns, 1000);
assert_eq!(tracker.pending_count(), 0);
}
#[test]
fn test_percentiles_multiple_orders() {
let mut tracker = OrderLatencyTracker::new();
for i in 1..=10_i64 {
let id = format!("o{i}");
let base = i * 10_000;
tracker.record_submit(&id, ts(base));
tracker.record_ack(&id, ts(base + i * 100)).unwrap();
tracker.record_fill(&id, ts(base + i * 100 + 500)).unwrap();
tracker.record_book_update(&id, ts(base + i * 100 + 500 + 200)).unwrap();
}
let stats = tracker.stats(LatencyPhase::SubmitToAck).unwrap();
assert_eq!(stats.count, 10);
assert_eq!(stats.min_ns, 100);
assert_eq!(stats.max_ns, 1000);
assert_eq!(stats.p50_ns, 600);
assert_eq!(stats.p99_ns, 1000);
}
#[test]
fn test_no_stats_before_samples() {
let tracker = OrderLatencyTracker::new();
assert!(tracker.stats(LatencyPhase::SubmitToAck).is_none());
}
#[test]
fn test_unknown_order_errors() {
let mut tracker = OrderLatencyTracker::new();
assert!(matches!(
tracker.record_ack("ghost", ts(1000)).unwrap_err(),
FinError::InvalidInput(_)
));
}
#[test]
fn test_ack_before_submit_errors() {
let mut tracker = OrderLatencyTracker::new();
tracker.record_submit("o1", ts(5000));
assert!(matches!(
tracker.record_ack("o1", ts(4000)).unwrap_err(),
FinError::InvalidInput(_)
));
}
#[test]
fn test_fill_before_ack_errors() {
let mut tracker = OrderLatencyTracker::new();
tracker.record_submit("o1", ts(1000));
tracker.record_ack("o1", ts(2000)).unwrap();
assert!(matches!(
tracker.record_fill("o1", ts(1500)).unwrap_err(),
FinError::InvalidInput(_)
));
}
}