1use crate::error::FinError;
18use crate::types::NanoTimestamp;
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
22pub enum LatencyPhase {
23 SubmitToAck,
25 AckToFill,
27 FillToBookUpdate,
29}
30
31#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
33pub struct LatencySample {
34 pub phase: LatencyPhase,
36 pub latency_ns: i64,
38}
39
40#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
42pub struct PhaseStats {
43 pub count: usize,
45 pub min_ns: i64,
47 pub max_ns: i64,
49 pub mean_ns: f64,
51 pub p50_ns: i64,
53 pub p95_ns: i64,
55 pub p99_ns: i64,
57}
58
59#[derive(Debug, Default)]
82pub struct OrderLatencyTracker {
83 pending: std::collections::HashMap<String, OrderTimestamps>,
85 submit_to_ack: Vec<i64>,
87 ack_to_fill: Vec<i64>,
89 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 pub fn new() -> Self {
103 Self::default()
104 }
105
106 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 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 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 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 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 pub fn pending_count(&self) -> usize {
228 self.pending.len()
229 }
230
231 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
241fn 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 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 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}