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