1use crate::error::FinError;
18use crate::types::NanoTimestamp;
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
23pub enum LatencyPhase {
24 SubmitToAck,
26 AckToFill,
28 FillToBookUpdate,
30}
31
32#[derive(Debug, Clone)]
34#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
35pub struct LatencySample {
36 pub phase: LatencyPhase,
38 pub latency_ns: i64,
40}
41
42#[derive(Debug, Clone)]
44#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
45pub struct PhaseStats {
46 pub count: usize,
48 pub min_ns: i64,
50 pub max_ns: i64,
52 pub mean_ns: f64,
54 pub p50_ns: i64,
56 pub p95_ns: i64,
58 pub p99_ns: i64,
60}
61
62#[derive(Debug, Default)]
85pub struct OrderLatencyTracker {
86 pending: std::collections::HashMap<String, OrderTimestamps>,
88 submit_to_ack: Vec<i64>,
90 ack_to_fill: Vec<i64>,
92 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 pub fn new() -> Self {
106 Self::default()
107 }
108
109 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 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 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 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 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 pub fn pending_count(&self) -> usize {
230 self.pending.len()
231 }
232
233 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
243fn 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 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 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}