Skip to main content

fin_primitives/latency/
mod.rs

1//! Order latency tracking: measures submit→ack, ack→fill, fill→book-update phases.
2//!
3//! ## Responsibility
4//! Tracks order latency across the three phases of an order lifecycle:
5//! submit → acknowledge → fill → book-update.
6//! Provides percentile statistics (p50 / p95 / p99) for each phase.
7//!
8//! ## Guarantees
9//! - All latencies are stored as `i64` nanoseconds; no floating-point drift in storage
10//! - Percentiles are computed from sorted snapshots; no panics on empty sets
11//! - Phases are tracked independently; missing phases do not corrupt other phases
12//!
13//! ## NOT Responsible For
14//! - Clock synchronization
15//! - Order routing
16
17use crate::error::FinError;
18use crate::types::NanoTimestamp;
19
20/// Phase of an order's lifecycle.
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
23pub enum LatencyPhase {
24    /// Time from order submission to exchange acknowledgement (ns).
25    SubmitToAck,
26    /// Time from acknowledgement to first fill (ns).
27    AckToFill,
28    /// Time from fill to order book update reflecting the fill (ns).
29    FillToBookUpdate,
30}
31
32/// A latency measurement for one phase of a single order.
33#[derive(Debug, Clone)]
34#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
35pub struct LatencySample {
36    /// Which phase this measurement covers.
37    pub phase: LatencyPhase,
38    /// Latency in nanoseconds.
39    pub latency_ns: i64,
40}
41
42/// Per-phase statistics snapshot.
43#[derive(Debug, Clone)]
44#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
45pub struct PhaseStats {
46    /// Number of samples recorded.
47    pub count: usize,
48    /// Minimum observed latency (ns).
49    pub min_ns: i64,
50    /// Maximum observed latency (ns).
51    pub max_ns: i64,
52    /// Approximate mean latency (ns).
53    pub mean_ns: f64,
54    /// 50th percentile (median) latency (ns).
55    pub p50_ns: i64,
56    /// 95th percentile latency (ns).
57    pub p95_ns: i64,
58    /// 99th percentile latency (ns).
59    pub p99_ns: i64,
60}
61
62/// Records open orders and accumulates latency samples for each lifecycle phase.
63///
64/// # Example
65/// ```rust
66/// use fin_primitives::latency::{OrderLatencyTracker, LatencyPhase};
67/// use fin_primitives::types::NanoTimestamp;
68///
69/// let mut tracker = OrderLatencyTracker::new();
70/// let t0 = NanoTimestamp::new(1_000_000_000);
71/// let t1 = NanoTimestamp::new(1_000_001_000);
72/// let t2 = NanoTimestamp::new(1_000_002_500);
73/// let t3 = NanoTimestamp::new(1_000_003_000);
74///
75/// tracker.record_submit("ord1", t0);
76/// tracker.record_ack("ord1", t1).unwrap();
77/// tracker.record_fill("ord1", t2).unwrap();
78/// tracker.record_book_update("ord1", t3).unwrap();
79///
80/// let stats = tracker.stats(LatencyPhase::SubmitToAck).unwrap();
81/// assert_eq!(stats.count, 1);
82/// assert_eq!(stats.p50_ns, 1000);
83/// ```
84#[derive(Debug, Default)]
85pub struct OrderLatencyTracker {
86    /// Pending orders: order_id → lifecycle timestamps (submit, ack, fill, book_update).
87    pending: std::collections::HashMap<String, OrderTimestamps>,
88    /// Accumulated submit-to-ack latencies (ns).
89    submit_to_ack: Vec<i64>,
90    /// Accumulated ack-to-fill latencies (ns).
91    ack_to_fill: Vec<i64>,
92    /// Accumulated fill-to-book-update latencies (ns).
93    fill_to_book: Vec<i64>,
94}
95
96#[derive(Debug, Clone, Default)]
97struct OrderTimestamps {
98    submit: Option<i64>,
99    ack: Option<i64>,
100    fill: Option<i64>,
101}
102
103impl OrderLatencyTracker {
104    /// Creates a new empty `OrderLatencyTracker`.
105    pub fn new() -> Self {
106        Self::default()
107    }
108
109    /// Records the submission timestamp of `order_id`.
110    ///
111    /// If a record for `order_id` already exists it is reset.
112    pub fn record_submit(&mut self, order_id: impl Into<String>, ts: NanoTimestamp) {
113        let ots = OrderTimestamps { submit: Some(ts.nanos()), ..OrderTimestamps::default() };
114        self.pending.insert(order_id.into(), ots);
115    }
116
117    /// Records the acknowledgement timestamp of `order_id`.
118    ///
119    /// Stores the submit→ack latency sample.
120    ///
121    /// # Errors
122    /// - [`FinError::InvalidInput`] if `order_id` is unknown or submit was not recorded.
123    /// - [`FinError::InvalidInput`] if ack timestamp is before submit timestamp.
124    pub fn record_ack(&mut self, order_id: &str, ts: NanoTimestamp) -> Result<(), FinError> {
125        let rec = self
126            .pending
127            .get_mut(order_id)
128            .ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
129        let submit = rec.submit.ok_or_else(|| {
130            FinError::InvalidInput(format!("submit not recorded for '{order_id}'"))
131        })?;
132        let ack_ns = ts.nanos();
133        if ack_ns < submit {
134            return Err(FinError::InvalidInput(format!(
135                "ack timestamp before submit for '{order_id}'"
136            )));
137        }
138        rec.ack = Some(ack_ns);
139        self.submit_to_ack.push(ack_ns - submit);
140        Ok(())
141    }
142
143    /// Records the fill timestamp of `order_id`.
144    ///
145    /// Stores the ack→fill latency sample.
146    ///
147    /// # Errors
148    /// - [`FinError::InvalidInput`] if `order_id` is unknown or ack was not recorded.
149    /// - [`FinError::InvalidInput`] if fill timestamp is before ack timestamp.
150    pub fn record_fill(&mut self, order_id: &str, ts: NanoTimestamp) -> Result<(), FinError> {
151        let rec = self
152            .pending
153            .get_mut(order_id)
154            .ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
155        let ack = rec.ack.ok_or_else(|| {
156            FinError::InvalidInput(format!("ack not recorded for '{order_id}'"))
157        })?;
158        let fill_ns = ts.nanos();
159        if fill_ns < ack {
160            return Err(FinError::InvalidInput(format!(
161                "fill timestamp before ack for '{order_id}'"
162            )));
163        }
164        rec.fill = Some(fill_ns);
165        self.ack_to_fill.push(fill_ns - ack);
166        Ok(())
167    }
168
169    /// Records the book-update timestamp of `order_id` and removes it from pending.
170    ///
171    /// Stores the fill→book-update latency sample.
172    ///
173    /// # Errors
174    /// - [`FinError::InvalidInput`] if `order_id` is unknown or fill was not recorded.
175    /// - [`FinError::InvalidInput`] if book_update timestamp is before fill timestamp.
176    pub fn record_book_update(
177        &mut self,
178        order_id: &str,
179        ts: NanoTimestamp,
180    ) -> Result<(), FinError> {
181        let rec = self
182            .pending
183            .remove(order_id)
184            .ok_or_else(|| FinError::InvalidInput(format!("unknown order '{order_id}'")))?;
185        let fill = rec.fill.ok_or_else(|| {
186            FinError::InvalidInput(format!("fill not recorded for '{order_id}'"))
187        })?;
188        let book_ns = ts.nanos();
189        if book_ns < fill {
190            return Err(FinError::InvalidInput(format!(
191                "book_update timestamp before fill for '{order_id}'"
192            )));
193        }
194        self.fill_to_book.push(book_ns - fill);
195        Ok(())
196    }
197
198    /// Returns percentile statistics for `phase`, or `None` if no samples exist.
199    pub fn stats(&self, phase: LatencyPhase) -> Option<PhaseStats> {
200        let samples = match phase {
201            LatencyPhase::SubmitToAck => &self.submit_to_ack,
202            LatencyPhase::AckToFill => &self.ack_to_fill,
203            LatencyPhase::FillToBookUpdate => &self.fill_to_book,
204        };
205        if samples.is_empty() {
206            return None;
207        }
208        let mut sorted = samples.clone();
209        sorted.sort_unstable();
210        let n = sorted.len();
211        let min_ns = *sorted.first().unwrap_or(&0);
212        let max_ns = *sorted.last().unwrap_or(&0);
213        let mean_ns = sorted.iter().map(|&v| v as f64).sum::<f64>() / n as f64;
214        let p50_ns = percentile_ns(&sorted, 50);
215        let p95_ns = percentile_ns(&sorted, 95);
216        let p99_ns = percentile_ns(&sorted, 99);
217        Some(PhaseStats {
218            count: n,
219            min_ns,
220            max_ns,
221            mean_ns,
222            p50_ns,
223            p95_ns,
224            p99_ns,
225        })
226    }
227
228    /// Returns the number of orders still awaiting lifecycle completion.
229    pub fn pending_count(&self) -> usize {
230        self.pending.len()
231    }
232
233    /// Returns all accumulated samples for `phase`.
234    pub fn samples(&self, phase: LatencyPhase) -> &[i64] {
235        match phase {
236            LatencyPhase::SubmitToAck => &self.submit_to_ack,
237            LatencyPhase::AckToFill => &self.ack_to_fill,
238            LatencyPhase::FillToBookUpdate => &self.fill_to_book,
239        }
240    }
241}
242
243/// Returns the `p`th percentile (nearest-rank method) of a **sorted** slice.
244fn percentile_ns(sorted: &[i64], p: usize) -> i64 {
245    if sorted.is_empty() {
246        return 0;
247    }
248    let idx = ((p * sorted.len()) / 100).min(sorted.len() - 1);
249    sorted[idx]
250}
251
252#[cfg(test)]
253mod tests {
254    use super::*;
255
256    fn ts(n: i64) -> NanoTimestamp {
257        NanoTimestamp::new(n)
258    }
259
260    fn full_lifecycle(tracker: &mut OrderLatencyTracker, id: &str, t0: i64, t1: i64, t2: i64, t3: i64) {
261        tracker.record_submit(id, ts(t0));
262        tracker.record_ack(id, ts(t1)).unwrap();
263        tracker.record_fill(id, ts(t2)).unwrap();
264        tracker.record_book_update(id, ts(t3)).unwrap();
265    }
266
267    #[test]
268    fn test_single_order_lifecycle() {
269        let mut tracker = OrderLatencyTracker::new();
270        full_lifecycle(&mut tracker, "o1", 1000, 2000, 4000, 5000);
271
272        let s2a = tracker.stats(LatencyPhase::SubmitToAck).unwrap();
273        assert_eq!(s2a.count, 1);
274        assert_eq!(s2a.p50_ns, 1000);
275
276        let a2f = tracker.stats(LatencyPhase::AckToFill).unwrap();
277        assert_eq!(a2f.p50_ns, 2000);
278
279        let f2b = tracker.stats(LatencyPhase::FillToBookUpdate).unwrap();
280        assert_eq!(f2b.p50_ns, 1000);
281
282        assert_eq!(tracker.pending_count(), 0);
283    }
284
285    #[test]
286    fn test_percentiles_multiple_orders() {
287        let mut tracker = OrderLatencyTracker::new();
288        // 10 orders with submit→ack latencies 100ns..1000ns (step 100)
289        for i in 1..=10_i64 {
290            let id = format!("o{i}");
291            let base = i * 10_000;
292            tracker.record_submit(&id, ts(base));
293            tracker.record_ack(&id, ts(base + i * 100)).unwrap();
294            tracker.record_fill(&id, ts(base + i * 100 + 500)).unwrap();
295            tracker.record_book_update(&id, ts(base + i * 100 + 500 + 200)).unwrap();
296        }
297        let stats = tracker.stats(LatencyPhase::SubmitToAck).unwrap();
298        assert_eq!(stats.count, 10);
299        assert_eq!(stats.min_ns, 100);
300        assert_eq!(stats.max_ns, 1000);
301        // nearest-rank: idx = (50 * 10) / 100 = 5 → value at index 5 = 600
302        assert_eq!(stats.p50_ns, 600);
303        assert_eq!(stats.p99_ns, 1000);
304    }
305
306    #[test]
307    fn test_no_stats_before_samples() {
308        let tracker = OrderLatencyTracker::new();
309        assert!(tracker.stats(LatencyPhase::SubmitToAck).is_none());
310    }
311
312    #[test]
313    fn test_unknown_order_errors() {
314        let mut tracker = OrderLatencyTracker::new();
315        assert!(matches!(
316            tracker.record_ack("ghost", ts(1000)).unwrap_err(),
317            FinError::InvalidInput(_)
318        ));
319    }
320
321    #[test]
322    fn test_ack_before_submit_errors() {
323        let mut tracker = OrderLatencyTracker::new();
324        tracker.record_submit("o1", ts(5000));
325        assert!(matches!(
326            tracker.record_ack("o1", ts(4000)).unwrap_err(),
327            FinError::InvalidInput(_)
328        ));
329    }
330
331    #[test]
332    fn test_fill_before_ack_errors() {
333        let mut tracker = OrderLatencyTracker::new();
334        tracker.record_submit("o1", ts(1000));
335        tracker.record_ack("o1", ts(2000)).unwrap();
336        assert!(matches!(
337            tracker.record_fill("o1", ts(1500)).unwrap_err(),
338            FinError::InvalidInput(_)
339        ));
340    }
341}