1use std::collections::BTreeMap;
28
29use crate::datastore::Datastore;
30use crate::push_interp;
31use crate::QlogPointRtt;
32use crate::QlogPointu64;
33use netlog::h2::H2_DEFAULT_WINDOW_SIZE;
34
35use qlog::events::quic::QuicFrame;
36
37#[derive(Default)]
38pub struct SeriesStore {
39 pub local_cwnd: Vec<QlogPointu64>,
40 pub local_bytes_in_flight: Vec<QlogPointu64>,
41 pub local_ssthresh: Vec<QlogPointu64>,
42 pub local_pacing_rate: Vec<QlogPointu64>,
43 pub local_delivery_rate: Vec<QlogPointu64>,
44 pub local_send_rate: Vec<QlogPointu64>,
45 pub local_ack_rate: Vec<QlogPointu64>,
46
47 pub local_min_rtt: Vec<QlogPointRtt>,
48 pub local_latest_rtt: Vec<QlogPointRtt>,
49 pub local_smoothed_rtt: Vec<QlogPointRtt>,
50
51 pub onertt_packet_created: Vec<QlogPointu64>,
52 pub onertt_packet_sent: Vec<QlogPointu64>,
53 pub onertt_packet_sent_aggregate_count: Vec<QlogPointu64>,
54 pub onertt_packet_lost_hacky: Vec<QlogPointu64>,
55 pub onertt_packet_lost_aggregate_count: Vec<QlogPointu64>,
56 pub onertt_packet_delivered_aggregate_count: Vec<QlogPointu64>,
57
58 pub onertt_packet_received: Vec<QlogPointu64>,
59
60 pub netlog_missing_packets: Vec<f64>,
61
62 pub onertt_packet_created_sent_delta: Vec<(u64, f64)>,
64
65 pub sent_max_data: Vec<QlogPointu64>,
66 pub sent_stream_max_data: BTreeMap<u64, Vec<QlogPointu64>>,
67
68 pub received_max_data: Vec<QlogPointu64>,
69 pub received_stream_max_data: BTreeMap<u64, Vec<QlogPointu64>>,
70
71 pub stream_buffer_reads: BTreeMap<u64, Vec<QlogPointu64>>,
72 pub sum_stream_buffer_reads: Vec<QlogPointu64>,
73
74 pub stream_buffer_writes: BTreeMap<u64, Vec<QlogPointu64>>,
75 pub sum_stream_buffer_writes: Vec<QlogPointu64>,
76
77 pub stream_buffer_dropped: BTreeMap<u64, Vec<QlogPointu64>>,
78 pub sum_stream_buffer_dropped: Vec<QlogPointu64>,
79
80 pub sent_stream_frames_series: BTreeMap<u64, Vec<QlogPointu64>>,
81
82 pub received_stream_frames_series: BTreeMap<u64, Vec<QlogPointu64>>,
83
84 pub received_data_frames_series: BTreeMap<u64, Vec<QlogPointu64>>,
85 pub received_data_max: BTreeMap<u64, u64>,
86 pub sent_data_frames_series: BTreeMap<u64, Vec<QlogPointu64>>,
87 pub sent_data_max: BTreeMap<u64, u64>,
88
89 pub h2_send_window_series_balanced: BTreeMap<u32, Vec<(f64, i32)>>,
90 pub h2_send_window_balanced_max: BTreeMap<u32, i32>,
91 pub h2_send_window_series_absolute: BTreeMap<u32, Vec<(f64, u64)>>,
92 pub h2_send_window_absolute_max: BTreeMap<u32, u64>,
93
94 pub netlog_h2_stream_received_connection_cumulative: Vec<QlogPointu64>,
95 pub netlog_quic_stream_received_connection_cumulative: Vec<QlogPointu64>,
96
97 pub netlog_quic_client_side_window_updates: BTreeMap<i64, Vec<(f64, u64)>>,
98
99 pub sum_received_stream_max_data: Vec<QlogPointu64>,
100 pub sum_sent_stream_max_data: Vec<QlogPointu64>,
101
102 pub sent_x_min: f64,
103 pub sent_x_max: f64,
104
105 pub received_x_min: f64,
106 pub received_x_max: f64,
107
108 pub y_max_stream_send_plot: u64,
109 pub y_max_stream_recv_plot: u64,
110 pub y_max_congestion_plot: u64,
111 pub y_max_rtt_plot: f32,
112
113 pub y_max_onertt_pkt_sent_plot: u64,
114 pub y_min_onertt_packet_created_sent_delta: f64,
115 pub y_max_onertt_packet_created_sent_delta: f64,
116
117 pub y_max_onertt_pkt_received_plot: u64,
118
119 pub max_pacing_rate: u64,
120 pub max_delivery_rate: u64,
121 pub max_send_rate: u64,
122 pub max_ack_rate: u64,
123}
124
125impl SeriesStore {
126 pub fn from_datastore(data_store: &Datastore) -> Self {
127 let mut series_store = SeriesStore::default();
128
129 series_store.populate_series_values(data_store);
130
131 series_store
132 }
133
134 fn update_sent_x_axis_max(&mut self, x: f64) {
135 self.sent_x_max = self.sent_x_max.max(x);
136 }
137
138 fn update_received_x_axis_max(&mut self, x: f64) {
139 self.received_x_max = self.sent_x_max.max(x);
140 }
141
142 fn update_congestion_y_axis_max(&mut self, y: u64) {
143 self.y_max_congestion_plot = self.y_max_congestion_plot.max(y);
144 }
145
146 fn update_stream_send_y_axis_max(&mut self, y: u64) {
147 self.y_max_stream_send_plot = self.y_max_stream_send_plot.max(y);
148 }
149
150 fn update_stream_recv_y_axis_max(&mut self, y: u64) {
151 self.y_max_stream_recv_plot = self.y_max_stream_recv_plot.max(y);
152 }
153
154 fn update_rtt_y_axis_max(&mut self, y: f32) {
155 self.y_max_rtt_plot = self.y_max_rtt_plot.max(y);
156 }
157
158 fn cwnd(&mut self, data_store: &Datastore) {
159 for point in &data_store.local_cwnd {
160 self.update_sent_x_axis_max(point.0);
161 self.update_congestion_y_axis_max(point.1);
162
163 push_interp(&mut self.local_cwnd, *point);
164 }
165 }
166
167 fn bif(&mut self, data_store: &Datastore) {
168 for point in &data_store.local_bytes_in_flight {
169 self.update_sent_x_axis_max(point.0);
170 self.update_congestion_y_axis_max(point.1);
171
172 push_interp(&mut self.local_bytes_in_flight, *point);
173 }
174 }
175
176 fn min_rtt(&mut self, data_store: &Datastore) {
177 for point in &data_store.local_min_rtt {
178 self.update_sent_x_axis_max(point.0);
179 self.update_rtt_y_axis_max(point.1);
180
181 push_interp(&mut self.local_min_rtt, *point);
182 }
183 }
184
185 fn latest_rtt(&mut self, data_store: &Datastore) {
186 for point in &data_store.local_latest_rtt {
187 self.update_sent_x_axis_max(point.0);
188 self.update_rtt_y_axis_max(point.1);
189
190 push_interp(&mut self.local_latest_rtt, *point);
191 }
192 }
193
194 fn pacing_rate(&mut self, data_store: &Datastore) {
195 for point in &data_store.local_pacing_rate {
196 self.update_sent_x_axis_max(point.0);
197 self.max_pacing_rate = self.max_pacing_rate.max(point.1);
198
199 push_interp(&mut self.local_pacing_rate, *point);
200 }
201 }
202
203 fn delivery_rate(&mut self, data_store: &Datastore) {
204 for point in &data_store.local_delivery_rate {
205 self.update_sent_x_axis_max(point.0);
206 self.max_delivery_rate = self.max_delivery_rate.max(point.1);
207
208 push_interp(&mut self.local_delivery_rate, *point);
209 }
210 }
211
212 fn send_rate(&mut self, data_store: &Datastore) {
213 for point in &data_store.local_send_rate {
214 self.update_sent_x_axis_max(point.0);
215 self.max_send_rate = self.max_send_rate.max(point.1);
216
217 push_interp(&mut self.local_send_rate, *point);
218 }
219 }
220
221 fn ack_rate(&mut self, data_store: &Datastore) {
222 for point in &data_store.local_ack_rate {
223 self.update_sent_x_axis_max(point.0);
224 self.max_ack_rate = self.max_ack_rate.max(point.1);
225
226 push_interp(&mut self.local_ack_rate, *point);
227 }
228 }
229
230 fn ssthresh(&mut self, data_store: &Datastore) {
231 for point in &data_store.local_ssthresh {
232 self.update_sent_x_axis_max(point.0);
233 push_interp(&mut self.local_ssthresh, *point);
236 }
237 }
238
239 fn smoothed_rtt(&mut self, data_store: &Datastore) {
240 for point in &data_store.local_smoothed_rtt {
241 self.update_sent_x_axis_max(point.0);
242 push_interp(&mut self.local_smoothed_rtt, *point);
245 }
246 }
247
248 fn sent_max_data(&mut self, data_store: &Datastore) {
249 for point in &data_store.sent_max_data {
250 self.update_sent_x_axis_max(point.0);
251 self.update_stream_recv_y_axis_max(point.1);
252
253 push_interp(&mut self.sent_max_data, *point);
254 }
255 }
256
257 fn sum_sent_stream_max_data(&mut self, data_store: &Datastore) {
258 for point in &data_store.sent_stream_max_data_tracker.sum_series {
259 self.update_sent_x_axis_max(point.0);
260 self.update_stream_recv_y_axis_max(point.1);
261
262 push_interp(&mut self.sum_sent_stream_max_data, *point);
263 }
264 }
265
266 fn packet_sent(&mut self, data_store: &Datastore) {
267 if let Some(onertt_pkts) =
268 &data_store.packet_sent.get(&crate::PacketType::OneRtt)
269 {
270 let mut sent_count = 0;
272 let mut delivered_count = 0;
273 let mut lost_count = 0;
274
275 for (pkt_num, pkt_info) in onertt_pkts.iter() {
276 sent_count += 1;
277 push_interp(
278 &mut self.onertt_packet_created,
279 (pkt_info.created_time, *pkt_num),
280 );
281 push_interp(
282 &mut self.onertt_packet_sent_aggregate_count,
283 (pkt_info.created_time, sent_count),
284 );
285
286 if pkt_info.acked.is_none() {
289 self.onertt_packet_lost_hacky
290 .push((pkt_info.created_time, *pkt_num));
291 lost_count += 1;
292 push_interp(
293 &mut self.onertt_packet_lost_aggregate_count,
294 (pkt_info.created_time, lost_count),
295 );
296 } else {
297 delivered_count += 1;
298 push_interp(
299 &mut self.onertt_packet_delivered_aggregate_count,
300 (pkt_info.created_time, delivered_count),
301 );
302 }
303
304 self.y_max_onertt_pkt_sent_plot =
305 std::cmp::max(self.y_max_onertt_pkt_sent_plot, *pkt_num);
306
307 if let Some(packet_sent_at) = pkt_info.send_at_time {
309 push_interp(
310 &mut self.onertt_packet_sent,
311 (packet_sent_at, *pkt_num),
312 );
313
314 let delta = packet_sent_at - pkt_info.created_time;
315 push_interp(
316 &mut self.onertt_packet_created_sent_delta,
317 (*pkt_num, delta),
318 );
319
320 if delta > self.y_max_onertt_packet_created_sent_delta {
322 self.y_max_onertt_packet_created_sent_delta = delta;
323 }
324
325 if delta < self.y_min_onertt_packet_created_sent_delta {
326 self.y_min_onertt_packet_created_sent_delta = delta;
327 }
328 }
329 }
330 }
331 }
332
333 fn packet_recv(&mut self, data_store: &Datastore) {
334 if let Some(onertt_pkts) =
335 &data_store.packet_received.get(&crate::PacketType::OneRtt)
336 {
337 for (pkt_num, pkt_info) in onertt_pkts.iter() {
338 self.update_received_x_axis_max(pkt_info.created_time);
339 self.y_max_onertt_pkt_received_plot =
340 self.y_max_onertt_pkt_received_plot.max(*pkt_num);
341
342 self.onertt_packet_received
343 .push((pkt_info.created_time, *pkt_num));
344 }
345 }
346 }
347
348 fn missing_packets(&mut self, data_store: &Datastore) {
349 let mut last: Vec<u64> = vec![];
355
356 for (event_time, missing_pkts) in
357 &data_store.netlog_ack_sent_missing_packets_raw
358 {
359 if &last != missing_pkts {
360 last = missing_pkts.clone();
361
362 self.netlog_missing_packets.push(*event_time);
363 }
364 }
365 }
366
367 fn sent_stream_max_data(&mut self, data_store: &Datastore) {
368 for (stream, points) in
369 &data_store.sent_stream_max_data_tracker.per_stream
370 {
371 let mut series_points = vec![];
372
373 for point in points {
374 self.update_sent_x_axis_max(point.0);
375 self.update_stream_recv_y_axis_max(point.1);
376
377 push_interp(&mut series_points, *point);
378 }
379
380 self.sent_stream_max_data.insert(*stream, series_points);
381 }
382 }
383
384 fn received_stream_max_data(&mut self, data_store: &Datastore) {
385 for (stream, points) in
386 &data_store.received_stream_max_data_tracker.per_stream
387 {
388 let mut series_points = vec![];
389
390 for point in points {
391 self.update_sent_x_axis_max(point.0);
392 self.update_stream_send_y_axis_max(point.1);
393
394 push_interp(&mut series_points, *point);
395 }
396
397 self.received_stream_max_data.insert(*stream, series_points);
398 }
399 }
400
401 fn stream_buffer_reads(&mut self, data_store: &Datastore) {
402 for (stream, points) in &data_store.stream_buffer_reads_tracker.per_stream
403 {
404 let mut series_points = vec![];
405
406 for point in points {
407 let y = point.1.offset + point.1.length;
408
409 self.update_sent_x_axis_max(point.0);
410 self.update_stream_recv_y_axis_max(y);
411
412 push_interp(&mut series_points, (point.0, y));
413 }
414
415 self.stream_buffer_reads.insert(*stream, series_points);
416 }
417 }
418
419 fn sum_stream_buffer_reads(&mut self, data_store: &Datastore) {
420 for point in &data_store.stream_buffer_reads_tracker.sum_series {
421 push_interp(&mut self.sum_stream_buffer_reads, *point);
422 }
423 }
424
425 fn stream_buffer_writes(&mut self, data_store: &Datastore) {
426 for (stream, points) in
427 &data_store.stream_buffer_writes_tracker.per_stream
428 {
429 let mut series_points = vec![];
430
431 for point in points {
432 let y = point.1.offset + point.1.length;
433
434 self.update_sent_x_axis_max(point.0);
435 self.update_stream_send_y_axis_max(y);
436
437 push_interp(&mut series_points, (point.0, y));
438 }
439
440 self.stream_buffer_writes.insert(*stream, series_points);
441 }
442 }
443
444 fn sum_stream_buffer_writes(&mut self, data_store: &Datastore) {
445 for point in &data_store.stream_buffer_writes_tracker.sum_series {
446 push_interp(&mut self.sum_stream_buffer_writes, *point);
447 }
448
449 if let Some((_, y)) =
450 data_store.stream_buffer_writes_tracker.sum_series.last()
451 {
452 self.update_stream_send_y_axis_max(*y);
453 }
454 }
455
456 fn stream_buffer_dropped(&mut self, data_store: &Datastore) {
457 for (stream, points) in
458 &data_store.stream_buffer_dropped_tracker.per_stream
459 {
460 let mut series_points = vec![];
461
462 for point in points {
463 let y = point.1.offset + point.1.length;
464
465 self.update_sent_x_axis_max(point.0);
466 self.update_stream_send_y_axis_max(y);
467
468 push_interp(&mut series_points, (point.0, y));
469 }
470
471 self.stream_buffer_dropped.insert(*stream, series_points);
472 }
473 }
474
475 fn sum_stream_buffer_dropped(&mut self, data_store: &Datastore) {
476 for point in &data_store.stream_buffer_dropped_tracker.sum_series {
477 push_interp(&mut self.sum_stream_buffer_dropped, *point);
478 }
479
480 if let Some((_, y)) =
481 data_store.stream_buffer_dropped_tracker.sum_series.last()
482 {
483 self.update_stream_send_y_axis_max(*y);
484 }
485 }
486
487 fn sent_stream_frames(&mut self, data_store: &Datastore) {
488 for (stream, points) in &data_store.sent_stream_frames {
489 let mut series_points = vec![];
490
491 for point in points {
492 if let (_, QuicFrame::Stream { offset, raw, .. }) = point {
493 let offset = offset.unwrap_or_default();
494 let length = raw
495 .clone()
496 .unwrap_or_default()
497 .payload_length
498 .unwrap_or_default();
499 let y = offset + length;
500
501 self.update_sent_x_axis_max(point.0);
502 self.update_stream_send_y_axis_max(y);
503
504 series_points.push((point.0, y));
505 }
506 }
507
508 self.sent_stream_frames_series
509 .insert(*stream, series_points);
510 }
511 }
512
513 fn received_stream_frames(&mut self, data_store: &Datastore) {
514 for (stream, points) in &data_store.received_stream_frames {
515 let mut series_points = vec![];
516
517 for point in points {
518 let y = point.1.offset + point.1.length;
519
520 self.received_x_max = self.received_x_max.max(point.0);
521 self.update_stream_recv_y_axis_max(y);
522
523 series_points.push((point.0, y));
524 }
525
526 self.received_stream_frames_series
527 .insert(*stream, series_points);
528 }
529 }
530
531 fn received_data_frames(&mut self, data_store: &Datastore) {
532 for (stream, points) in &data_store.received_data_frames {
533 let mut series_points = vec![];
534
535 if let Some(first) = points.first() {
537 series_points.push((first.0, 0));
538 } else {
539 continue;
540 }
541
542 let mut last_y = series_points.first().unwrap().1;
543
544 for point in points {
545 let new_y = last_y + point.1;
546
547 self.received_x_max = self.received_x_max.max(point.0);
548 self.update_stream_recv_y_axis_max(new_y);
549
550 series_points.push((point.0, new_y));
551
552 last_y = new_y;
553 }
554
555 self.received_data_max
556 .insert(*stream, series_points.last().unwrap().1);
557
558 self.received_data_frames_series
559 .insert(*stream, series_points);
560 }
561 }
562
563 fn sent_data_frames(&mut self, data_store: &Datastore) {
564 for (stream, points) in &data_store.sent_data_frames {
565 let mut series_points = vec![];
566
567 if let Some(first) = points.first() {
569 series_points.push((first.0, 0));
570 } else {
571 continue;
572 }
573
574 let mut last_y = series_points.first().unwrap().1;
575
576 for point in points {
577 let new_y = last_y + point.1;
578
579 self.update_sent_x_axis_max(point.0);
580 series_points.push((point.0, new_y));
581
582 series_points.push((point.0, new_y));
583
584 last_y = new_y;
585 }
586
587 self.sent_data_max
588 .insert(*stream, series_points.last().unwrap().1);
589
590 self.sent_data_frames_series.insert(*stream, series_points);
591 }
592 }
593
594 fn initial_h2_fc_value(stream_id: u32, data_store: &Datastore) -> u32 {
599 if stream_id == 0 {
600 H2_DEFAULT_WINDOW_SIZE
603 } else {
604 data_store
605 .h2_server_settings
606 .initial_window_size
607 .unwrap_or(H2_DEFAULT_WINDOW_SIZE)
608 }
609 }
610
611 fn h2_fc_balanced(&mut self, data_store: &Datastore) {
613 for (stream, points) in &data_store.h2_send_window_updates_balanced {
614 let mut series_points = vec![];
615 let mut y_max = 0;
616
617 let initial_window =
618 Self::initial_h2_fc_value(*stream, data_store) as i32;
619
620 if let Some(first) = points.first() {
622 series_points.push((first.0, initial_window));
623 } else {
624 continue;
625 }
626
627 let mut last_y = series_points.first().unwrap().1;
628
629 for point in points {
630 let new_y = last_y + point.1;
631 self.update_sent_x_axis_max(point.0);
632 y_max = y_max.max(new_y);
633
634 push_interp(&mut series_points, (point.0, new_y));
635 last_y = new_y;
636 }
637
638 self.h2_send_window_series_balanced
639 .insert(*stream, series_points);
640 self.h2_send_window_balanced_max.insert(*stream, y_max);
641 }
642 }
643
644 fn h2_fc_absolute(&mut self, data_store: &Datastore) {
645 for (stream, points) in &data_store.h2_send_window_updates_absolute {
646 let mut series_points = vec![];
647 let mut y_max = 0;
648
649 let initial_window =
650 Self::initial_h2_fc_value(*stream, data_store) as u64;
651
652 if let Some(first) = points.first() {
654 series_points.push((first.0, initial_window));
655 } else {
656 continue;
657 }
658
659 let mut last_y = series_points.first().unwrap().1;
660
661 for point in points {
662 let new_y = last_y + point.1;
663 self.update_sent_x_axis_max(point.0);
664 y_max = y_max.max(new_y);
665
666 push_interp(&mut series_points, (point.0, new_y));
667
668 last_y = new_y;
669 }
670
671 self.h2_send_window_series_absolute
672 .insert(*stream, series_points);
673 self.h2_send_window_absolute_max.insert(*stream, y_max);
674 }
675 }
676
677 fn netlog_quic_client_side_window_updates(&mut self, data_store: &Datastore) {
678 for (stream, points) in &data_store.netlog_quic_client_side_window_updates
679 {
680 let s = self
681 .netlog_quic_client_side_window_updates
682 .entry(*stream)
683 .or_default();
684
685 for point in points {
686 push_interp(s, *point);
687 }
688 }
689 }
690
691 fn netlog_h2_stream_received_connection_cumulative(
692 &mut self, data_store: &Datastore,
693 ) {
694 for point in &data_store.netlog_h2_stream_received_connection_cumulative {
695 push_interp(
696 &mut self.netlog_h2_stream_received_connection_cumulative,
697 *point,
698 );
699 }
700 }
701
702 fn netlog_quic_stream_received_connection_cumulative(
703 &mut self, data_store: &Datastore,
704 ) {
705 for point in &data_store.netlog_quic_stream_received_connection_cumulative
706 {
707 push_interp(
708 &mut self.netlog_quic_stream_received_connection_cumulative,
709 *point,
710 );
711 }
712 }
713
714 fn received_max_data(&mut self, data_store: &Datastore) {
715 for point in &data_store.received_max_data {
716 push_interp(&mut self.received_max_data, *point);
717 }
718 }
719
720 fn sum_received_stream_max_data(&mut self, data_store: &Datastore) {
721 for point in &data_store.received_stream_max_data_tracker.sum_series {
722 push_interp(&mut self.sum_received_stream_max_data, *point);
723 }
724 }
725
726 fn populate_series_values(&mut self, data_store: &Datastore) {
727 self.cwnd(data_store);
728 self.bif(data_store);
729 self.min_rtt(data_store);
730 self.latest_rtt(data_store);
731 self.pacing_rate(data_store);
732 self.delivery_rate(data_store);
733 self.send_rate(data_store);
734 self.ack_rate(data_store);
735 self.ssthresh(data_store);
736 self.smoothed_rtt(data_store);
737
738 self.sent_max_data(data_store);
739 self.sum_sent_stream_max_data(data_store);
740
741 self.packet_sent(data_store);
742 self.packet_recv(data_store);
743 self.missing_packets(data_store);
744
745 self.sent_stream_max_data(data_store);
746 self.received_stream_max_data(data_store);
747
748 self.stream_buffer_reads(data_store);
749 self.sum_stream_buffer_reads(data_store);
750
751 self.stream_buffer_writes(data_store);
752 self.sum_stream_buffer_writes(data_store);
753
754 self.stream_buffer_dropped(data_store);
755 self.sum_stream_buffer_dropped(data_store);
756
757 self.sent_stream_frames(data_store);
758 self.received_stream_frames(data_store);
759
760 self.sent_data_frames(data_store);
761 self.received_data_frames(data_store);
762
763 self.h2_fc_balanced(data_store);
764 self.h2_fc_absolute(data_store);
765
766 self.netlog_h2_stream_received_connection_cumulative(data_store);
767 self.netlog_quic_stream_received_connection_cumulative(data_store);
768 self.netlog_quic_client_side_window_updates(data_store);
769
770 self.received_max_data(data_store);
771 self.sum_received_stream_max_data(data_store);
772 }
773}