1use std::collections::BTreeMap;
28use std::collections::BTreeSet;
29use std::collections::HashMap;
30use std::collections::HashSet;
31use std::convert::TryFrom;
32use std::fmt::Debug;
33use std::fmt::Display;
34
35use log::error;
36use log::trace;
37use qlog::events::http3::Http3Frame;
38use qlog::events::quic::QuicFrame;
39use qlog::events::quic::TransportInitiator;
40use qlog::events::EventData;
41use qlog::events::RawInfo;
42
43use regex::Regex;
44use tabled::Tabled;
45
46use crate::request_stub::find_header_value;
47use crate::request_stub::HttpRequestStub;
48use crate::request_stub::NaOption;
49use crate::trackers::StreamBufferTracker;
50use crate::trackers::StreamMaxTracker;
51use crate::LogFileData;
52use crate::PacketType;
53use crate::QlogPointRtt;
54use crate::QlogPointu64;
55use crate::RawLogEvents::Netlog;
56use netlog;
57use netlog::h2;
58use netlog::h2::Event::*;
59use netlog::h2::*;
60use netlog::h3;
61use netlog::h3::Event::*;
62use netlog::http;
63use netlog::quic;
64use netlog::quic::Event::*;
65use netlog::quic::*;
66use netlog::read_netlog_record;
67
68pub type ParseResult<T> = Result<T, serde_json::Error>;
69#[derive(Debug, Clone)]
70pub struct PacketInfoStub {
71 pub acked: Option<bool>,
72 pub raw: Option<RawInfo>,
73 pub created_time: f64,
74 pub send_at_time: Option<f64>,
75 pub ty: PacketType,
76 pub number: u64,
77}
78
79#[derive(Clone, Debug)]
80pub struct StreamAccess {
81 pub offset: u64,
82 pub length: u64,
83}
84
85#[derive(Default, Debug)]
86pub struct PrintStatsConfig {
87 pub rx_flow_control: bool,
88 pub tx_flow_control: bool,
89 pub reset_streams: bool,
90 pub stream_buffering: bool,
91 pub tx_stream_frames: bool,
92 pub packet_stats: bool,
93}
94
95#[derive(Clone, Copy, Debug, Default)]
96pub struct RequestAtServerDeltas {
97 pub rx_hdr_tx_hdr: NaOption<f64>,
98 pub rx_hdr_tx_first_data: NaOption<f64>,
99 pub rx_hdr_tx_last_data: NaOption<f64>,
100 pub tx_first_data_tx_last_data: NaOption<f64>,
101}
102
103#[derive(Clone, Copy, Debug, Default)]
104pub struct RequestAtClientDeltas {
105 pub discover_tx_hdr: NaOption<f64>,
106 pub tx_hdr_rx_hdr: NaOption<f64>,
107 pub tx_hdr_rx_first_data: NaOption<f64>,
108 pub tx_hdr_rx_last_data: NaOption<f64>,
109 pub tx_first_data_tx_last_data: NaOption<f64>,
110 pub rx_first_data_rx_last_data: NaOption<f64>,
111 pub rx_hdr_rx_last_data: NaOption<f64>,
112}
113
114#[derive(Debug, Default)]
115pub enum RequestActor {
116 #[default]
117 Client,
118 Server,
119}
120
121impl From<VantagePoint> for RequestActor {
122 fn from(value: VantagePoint) -> Self {
123 match value {
124 VantagePoint::Client => RequestActor::Client,
125 VantagePoint::Server => RequestActor::Server,
126 }
127 }
128}
129
130#[derive(Debug, Default, Clone, Copy)]
131pub enum VantagePoint {
132 #[default]
133 Client,
134 Server,
135}
136
137#[derive(Debug, Default)]
138pub struct StreamDatapoint {
139 pub offset: u64,
140 pub length: u64,
141}
142
143#[derive(Default, Clone, Copy, PartialEq)]
144pub enum ApplicationProto {
145 Http2,
146 #[default]
147 Http3,
148}
149
150impl Debug for ApplicationProto {
151 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
152 Display::fmt(self, f)
153 }
154}
155
156impl Display for ApplicationProto {
157 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
158 let v = match self {
159 ApplicationProto::Http2 => "HTTP/2",
160 ApplicationProto::Http3 => "HTTP/3",
161 };
162
163 write!(f, "{}", v)
164 }
165}
166
167#[derive(Debug, Default, Tabled)]
168pub struct QuicSessionClose {
169 #[tabled(rename = "ID")]
170 pub session_id: i64,
171 #[tabled(rename = "SNI")]
172 pub sni: String,
173 #[tabled(rename = "Error")]
174 pub quic_error: i64,
175 #[tabled(rename = "Description")]
176 pub quic_error_pretty: NaOption<String>,
177 #[tabled(rename = "From peer")]
178 pub from_peer: bool,
179 #[tabled(rename = "Additional Details")]
180 pub details: String,
181}
182
183#[derive(Debug, Default, Tabled)]
184pub struct H2SessionClose {
185 #[tabled(rename = "ID")]
186 pub session_id: i64,
187 #[tabled(rename = "SNI")]
188 pub sni: String,
189 #[tabled(rename = "Error")]
190 pub net_err: i64,
191 #[tabled(rename = "Description")]
192 pub net_err_pretty: NaOption<String>,
193 #[tabled(rename = "Additional Details")]
194 pub details: String,
195}
196
197#[derive(Default, Debug)]
198pub struct QuicStreamStopSending {
199 pub quic_rst_stream_error: u64,
200 pub quic_rst_stream_error_friendly: Option<String>,
201}
202
203impl Display for QuicStreamStopSending {
204 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
205 write!(
206 f,
207 "code={}, {}",
208 self.quic_rst_stream_error,
209 self.quic_rst_stream_error_friendly.as_deref().unwrap_or("")
210 )
211 }
212}
213
214#[derive(Default, Debug)]
215pub struct QuicStreamReset {
216 pub offset: u64,
217 pub quic_rst_stream_error: u64,
218 pub quic_rst_stream_error_friendly: Option<String>,
219}
220
221impl Display for QuicStreamReset {
222 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
223 write!(
224 f,
225 "offset={}, code={}, {}",
226 self.offset,
227 self.quic_rst_stream_error,
228 self.quic_rst_stream_error_friendly.as_deref().unwrap_or("")
229 )
230 }
231}
232
233#[derive(Default, Debug)]
234pub struct H2StreamReset {
235 pub error: String,
236 pub description: String,
237}
238
239impl Display for H2StreamReset {
240 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
241 write!(f, "code={}, {}", self.error, self.description,)
242 }
243}
244
245#[derive(Default, Debug)]
246pub struct Datastore {
247 pub vantage_point: VantagePoint,
248 pub application_proto: ApplicationProto,
249 pub session_id: Option<i64>,
250 pub host: Option<String>,
251
252 pub h2_client_settings: Http2Settings,
253 pub h2_server_settings: Http2Settings,
254
255 pub client_quic_tps: TransportParameters,
256
257 pub last_event_time: f64,
258
259 pub packet_sent: HashMap<PacketType, BTreeMap<u64, PacketInfoStub>>,
263 pub packet_received: HashMap<PacketType, BTreeMap<u64, PacketInfoStub>>,
264
265 pub netlog_ack_sent_missing_packet: BTreeMap<PacketType, BTreeSet<u64>>,
269 pub netlog_ack_received_missing_packet: BTreeMap<PacketType, BTreeSet<u64>>,
270
271 pub packet_acked: Vec<QlogPointu64>,
272
273 pub local_cwnd: Vec<QlogPointu64>,
274 pub local_bytes_in_flight: Vec<QlogPointu64>,
275 pub local_ssthresh: Vec<QlogPointu64>,
276 pub local_pacing_rate: Vec<QlogPointu64>,
277 pub local_delivery_rate: Vec<QlogPointu64>,
278 pub local_send_rate: Vec<QlogPointu64>,
279 pub local_ack_rate: Vec<QlogPointu64>,
280
281 pub local_min_rtt: Vec<QlogPointRtt>,
282 pub local_latest_rtt: Vec<QlogPointRtt>,
283 pub local_smoothed_rtt: Vec<QlogPointRtt>,
284
285 pub congestion_state_updates: Vec<(f64, u64, String)>,
286
287 pub received_max_data: Vec<QlogPointu64>,
288
289 pub received_stream_max_data_tracker: StreamMaxTracker,
292
293 pub sent_max_data: Vec<QlogPointu64>,
294
295 pub sent_stream_max_data_tracker: StreamMaxTracker,
298
299 pub stream_buffer_reads_tracker: StreamBufferTracker,
302
303 pub stream_buffer_writes_tracker: StreamBufferTracker,
306
307 pub stream_buffer_dropped_tracker: StreamBufferTracker,
310
311 pub received_reset_stream: BTreeMap<u64, Vec<QuicFrame>>,
312 pub sent_reset_stream: BTreeMap<u64, Vec<QuicFrame>>,
313
314 pub received_stream_frames: BTreeMap<u64, Vec<(f64, StreamDatapoint)>>,
315 pub received_stream_frames_count_based:
316 BTreeMap<u64, Vec<(u64, StreamDatapoint)>>,
317 pub total_received_stream_frame_count: u64,
318
319 pub sent_stream_frames: BTreeMap<u64, Vec<(f64, QuicFrame)>>,
320 pub sent_stream_frames_count_based: BTreeMap<u64, Vec<(u64, QuicFrame)>>,
321 pub total_sent_stream_frame_count: u64,
322
323 pub received_data_frames: BTreeMap<u64, Vec<(f64, u64)>>,
324 pub received_data_frames_count_based: BTreeMap<u64, Vec<(u64, u64)>>,
325 pub total_received_data_frame_count: u64,
326 pub received_data_cumulative: BTreeMap<u64, Vec<(f64, u64)>>,
327 pub received_data_cumulative_max: BTreeMap<u64, u64>,
328
329 pub sent_data_frames: BTreeMap<u64, Vec<(f64, u64)>>,
330 pub sent_data_frames_count_based: BTreeMap<u64, Vec<(u64, u64)>>,
331 pub total_sent_data_frame_count: u64,
332 pub sent_data_cumulative: BTreeMap<u64, Vec<(f64, u64)>>,
333 pub sent_data_cumulative_max: BTreeMap<u64, u64>,
334
335 pub http_requests: BTreeMap<u64, HttpRequestStub>,
336 pub largest_data_frame_rx_length_global: u64,
337 pub largest_data_frame_tx_length_global: u64,
338
339 pub local_init_max_stream_data_bidi_local: u64,
340 pub local_init_max_stream_data_uni: u64,
341 pub peer_init_max_stream_data_bidi_local: u64,
342 pub peer_init_max_stream_data_bidi_remote: u64,
343 pub peer_init_max_stream_data_uni: u64,
344
345 pub h2_recv_window_updates: BTreeMap<u32, Vec<(f64, i32)>>,
346
347 pub h2_send_window_updates_balanced: BTreeMap<u32, Vec<(f64, i32)>>,
349
350 pub h2_send_window_updates_absolute: BTreeMap<u32, Vec<(f64, u64)>>,
352
353 pub netlog_quic_server_window_blocked: BTreeMap<i64, Vec<f64>>,
354 pub netlog_quic_client_side_window_updates: BTreeMap<i64, Vec<(f64, u64)>>,
355
356 pub netlog_h2_stream_received_connection_cumulative: Vec<QlogPointu64>,
357 pub netlog_quic_stream_received_connection_cumulative: Vec<QlogPointu64>,
358
359 pub received_packets_netlog: Vec<(f64, PacketInfoStub)>,
360 pub discontinuous_packet_number_count: u64,
361
362 pub netlog_ack_sent_missing_packets_raw: Vec<(f64, Vec<u64>)>,
363
364 pub total_tx_ack: usize,
365 pub max_ack_sent_missing_packets_size: usize,
366
367 pub total_rx_ack: usize,
368 pub max_ack_received_missing_packets_size: usize,
369
370 pub quic_session_close: Option<QuicSessionClose>,
371 pub h2_session_close: Option<H2SessionClose>,
372
373 pub h2_concurrent_requests: u64,
374}
375
376fn is_bidi(stream_id: u64) -> bool {
377 (stream_id & 0x2) == 0
378}
379
380impl Datastore {
381 pub fn consume_netlog_event(
382 &mut self, session_start_time: u64, ev_hdr: &netlog::EventHeader,
383 event: &netlog::Event, constants: &netlog::constants::Constants,
384 stream_bind: &StreamBindingMap,
385 h3_session_requests: Option<&Vec<ReqOverH3>>,
386 ) {
387 match event {
388 netlog::Event::Http(_e) => (),
390 netlog::Event::H2(e) => self.consume_netlog_h2(
391 session_start_time,
392 ev_hdr,
393 e,
394 constants,
395 stream_bind,
396 ),
397
398 netlog::Event::H3(e) => self.consume_netlog_h3(
399 session_start_time,
400 ev_hdr,
401 e,
402 h3_session_requests,
403 ),
404 netlog::Event::Quic(e) =>
405 self.consume_netlog_quic(session_start_time, ev_hdr, e, constants),
406 }
407 }
408
409 fn consume_netlog_quic(
410 &mut self, session_start_time: u64, ev_hdr: &netlog::EventHeader,
411 ev: &quic::Event, constants: &netlog::constants::Constants,
412 ) {
413 let rel_event_time = (ev_hdr.time_num - session_start_time) as f64;
414
415 match ev {
416 QuicSessionStreamFrameReceived(e) => {
417 let s = self
418 .received_stream_frames
419 .entry(e.params.stream_id)
420 .or_default();
421 s.push((rel_event_time, StreamDatapoint {
422 offset: e.params.offset,
423 length: e.params.length,
424 }));
425
426 let s = self
427 .received_stream_frames_count_based
428 .entry(e.params.stream_id)
429 .or_default();
430 s.push((
431 self.total_received_stream_frame_count,
432 StreamDatapoint {
433 offset: e.params.offset,
434 length: e.params.length,
435 },
436 ));
437
438 self.total_received_stream_frame_count += 1;
439
440 if e.params.fin {
441 if let Some(req) =
442 self.http_requests.get_mut(&e.params.stream_id)
443 {
444 req.time_fin_rx = Some(rel_event_time);
445 }
446 }
447
448 let cumulative = if let Some(last) = self
449 .netlog_quic_stream_received_connection_cumulative
450 .last()
451 {
452 last.1 + e.params.length
453 } else {
454 e.params.length
456 };
457
458 self.netlog_quic_stream_received_connection_cumulative
459 .push((rel_event_time, cumulative));
460 },
461
462 QuicSessionUnauthenticatedPacketHeaderReceived(e) => {
463 if let Some((_, last)) = self.received_packets_netlog.last() {
464 let gap = if e.params.packet_number > last.number {
465 e.params.packet_number - last.number
466 } else {
467 last.number.saturating_sub(e.params.packet_number)
468 };
469
470 if gap > 1 {
471 self.discontinuous_packet_number_count += 1;
472 }
473 }
474
475 let packet_type = PacketType::from_netlog_packet_header(
476 &e.params.header_format,
477 &e.params.long_header_type,
478 );
479
480 let packet_info = PacketInfoStub {
481 acked: None,
482 raw: None,
483 created_time: rel_event_time,
484 send_at_time: None,
485 ty: packet_type,
486 number: e.params.packet_number,
487 };
488
489 self.received_packets_netlog
490 .push((rel_event_time, packet_info.clone()));
491
492 let s = self.packet_received.entry(packet_type).or_default();
493
494 s.insert(e.params.packet_number, packet_info);
495 },
496
497 QuicSessionPacketSent(e) => {
498 let packet_type = PacketType::from_netlog_encryption_level(
499 &e.params.encryption_level,
500 );
501
502 let packet_info = PacketInfoStub {
503 acked: None,
504 raw: None,
505 created_time: rel_event_time,
506 send_at_time: None,
507 ty: packet_type,
508 number: e.params.packet_number,
509 };
510
511 let s = self.packet_sent.entry(packet_type).or_default();
512
513 s.insert(e.params.packet_number, packet_info);
514
515 if let Some(pkts) = self
517 .netlog_ack_sent_missing_packet
518 .get_mut(&PacketType::Unknown)
519 {
520 if !pkts.is_empty() {
521 let old = std::mem::take(pkts);
522
523 let s = self
524 .netlog_ack_sent_missing_packet
525 .entry(packet_type)
526 .or_default();
527
528 for num in old {
529 s.insert(num);
530 }
531 }
532 }
533 },
534
535 QuicSessionAckFrameSent(e) => {
536 self.total_tx_ack += 1;
537 self.max_ack_sent_missing_packets_size = std::cmp::max(
538 e.params.missing_packets.len(),
539 self.max_ack_sent_missing_packets_size,
540 );
541
542 if !e.params.missing_packets.is_empty() {
543 self.netlog_ack_sent_missing_packets_raw
544 .push((rel_event_time, e.params.missing_packets.clone()));
545 }
546
547 let s = self
548 .netlog_ack_sent_missing_packet
549 .entry(PacketType::Unknown)
550 .or_default();
551
552 for pn in &e.params.missing_packets {
553 s.insert(*pn);
558 }
559 },
560
561 QuicSessionAckFrameReceived(e) => {
562 self.total_rx_ack += 1;
563 self.max_ack_received_missing_packets_size = std::cmp::max(
564 e.params.missing_packets.len(),
565 self.max_ack_sent_missing_packets_size,
566 );
567
568 if !e.params.missing_packets.is_empty() {
569 let parent_packet = self.received_packets_netlog.last();
572 if let Some((_, pkt_info)) = parent_packet {
573 let s = self
574 .netlog_ack_received_missing_packet
575 .entry(pkt_info.ty)
576 .or_default();
577
578 for missing in &e.params.missing_packets {
579 s.insert(*missing);
580 }
581 }
582 }
583 },
584
585 QuicSessionClosed(e) => {
586 self.quic_session_close = Some(QuicSessionClose {
587 session_id: self.session_id.unwrap_or(-1),
588 sni: self.host.clone().unwrap_or("ERROR UNKNOWN".to_string()),
589 details: e.params.details.clone(),
590 from_peer: e.params.from_peer,
591 quic_error: e.params.quic_error,
592 quic_error_pretty: NaOption::new(
593 constants
594 .quic_error_id_keyed
595 .get(&e.params.quic_error)
596 .cloned(),
597 ),
598 });
599 },
600
601 QuicSessionRstStreamFrameReceived(e) => {
602 if let Some(req) = self.http_requests.get_mut(&e.params.stream_id)
605 {
606 req.quic_stream_reset_received = Some(QuicStreamReset {
607 offset: e.params.offset,
608 quic_rst_stream_error: e.params.quic_rst_stream_error,
609 quic_rst_stream_error_friendly: constants
610 .quic_rst_stream_error_id_keyed
611 .get(&(e.params.quic_rst_stream_error as i64))
612 .cloned(),
613 });
614 }
615 },
616
617 QuicSessionRstStreamFrameSent(e) => {
618 if let Some(req) = self.http_requests.get_mut(&e.params.stream_id)
621 {
622 req.quic_stream_reset_sent = Some(QuicStreamReset {
623 offset: e.params.offset,
624 quic_rst_stream_error: e.params.quic_rst_stream_error,
625 quic_rst_stream_error_friendly: constants
626 .quic_rst_stream_error_id_keyed
627 .get(&(e.params.quic_rst_stream_error as i64))
628 .cloned(),
629 });
630 }
631 },
632
633 QuicSessionStopSendingFrameSent(e) => {
634 if let Some(req) = self.http_requests.get_mut(&e.params.stream_id)
637 {
638 req.quic_stream_stop_sending_sent =
639 Some(QuicStreamStopSending {
640 quic_rst_stream_error: e.params.quic_rst_stream_error,
641 quic_rst_stream_error_friendly: constants
642 .quic_rst_stream_error_id_keyed
643 .get(&(e.params.quic_rst_stream_error as i64))
644 .cloned(),
645 });
646 }
647 },
648
649 QuicSessionBlockedFrameReceived(e) => {
650 let s = self
651 .netlog_quic_server_window_blocked
652 .entry(e.params.stream_id)
653 .or_default();
654 s.push(rel_event_time);
655 },
656
657 QuicSessionWindowUpdateFrameSent(e) => {
658 let s = self
659 .netlog_quic_client_side_window_updates
660 .entry(e.params.stream_id)
661 .or_default();
662
663 s.push((rel_event_time, e.params.byte_offset));
664 },
665
666 QuicSessionTransportParametersSent(e) => {
667 self.client_quic_tps =
668 e.params.quic_transport_parameters.clone().into();
669
670 let s = self
671 .netlog_quic_client_side_window_updates
672 .entry(-1)
673 .or_default();
674
675 s.push((
676 rel_event_time,
677 self.client_quic_tps.initial_max_data.unwrap_or_default(),
678 ));
679 },
680
681 QuicSession(_) | QuicSessionTransportParametersReceived(_) => (),
683 }
684 }
685
686 fn consume_netlog_h3(
687 &mut self, session_start_time: u64, ev_hdr: &netlog::EventHeader,
688 ev: &h3::Event, h3_session_requests: Option<&Vec<ReqOverH3>>,
689 ) {
690 let rel_event_time = (ev_hdr.time_num - session_start_time) as f64;
691
692 match ev {
693 Http3PriorityUpdateSent(e) => {
694 let req =
695 self.get_or_insert_http_req(e.params.prioritized_element_id);
696 req.priority_updates
697 .push(e.params.priority_field_value.clone());
698 },
699
700 Http3HeadersSent(e) => {
701 let req = self.get_or_insert_http_req(e.params.stream_id);
702 req.time_first_headers_tx.get_or_insert(rel_event_time);
703 req.set_request_info_from_netlog(&e.params.headers);
704
705 if let Some(reqs) = h3_session_requests {
706 for r in reqs {
707 if let Some(stream_id) = r.quic_stream_id {
708 if stream_id == e.params.stream_id {
709 req.time_discovery = Some(
713 (r.discover_time
714 .wrapping_sub(session_start_time)
715 as i64)
716 as f64,
717 );
718 break;
719 }
720 }
721 }
722 }
723 },
724
725 Http3HeadersReceived(_) => (),
728
729 Http3HeadersDecoded(e) => {
730 let req = self.get_or_insert_http_req(e.params.stream_id);
731 req.time_first_headers_rx.get_or_insert(rel_event_time);
732 req.set_response_info_from_netlog(&e.params.headers);
733 },
734
735 Http3DataFrameReceived(e) => {
736 let req = self.get_or_insert_http_req(e.params.stream_id);
737
738 req.time_first_data_rx.get_or_insert(rel_event_time);
739
740 let _ = req.time_last_data_rx.insert(rel_event_time);
741
742 let length = e.params.payload_length;
743 req.time_data_rx_set.push((rel_event_time, length));
744 self.largest_data_frame_rx_length_global = std::cmp::max(
745 self.largest_data_frame_rx_length_global,
746 length,
747 );
748
749 let s = self
750 .received_data_frames
751 .entry(e.params.stream_id)
752 .or_default();
753
754 s.push((rel_event_time, e.params.payload_length));
755
756 let s = self
757 .received_data_frames_count_based
758 .entry(e.params.stream_id)
759 .or_default();
760 s.push((
761 self.total_received_data_frame_count,
762 e.params.payload_length,
763 ));
764
765 self.total_received_data_frame_count += 1;
766
767 let s = self
768 .received_data_cumulative
769 .entry(e.params.stream_id)
770 .or_default();
771
772 let received_data_cumulative = if let Some(last) = s.last() {
773 last.1 + e.params.payload_length
774 } else {
775 e.params.payload_length
777 };
778
779 s.push((rel_event_time, received_data_cumulative));
780
781 let s = self
782 .received_data_cumulative_max
783 .entry(e.params.stream_id)
784 .or_default();
785
786 *s = std::cmp::max(*s, received_data_cumulative);
787
788 let s = self.http_requests.entry(e.params.stream_id).or_default();
789
790 s.server_transferred_bytes = std::cmp::max(
791 s.server_transferred_bytes,
792 NaOption::new(Some(received_data_cumulative)),
793 )
794 },
795
796 Http3DataSent(_) => (),
798 }
799 }
800
801 fn consume_netlog_h2(
802 &mut self, session_start_time: u64, ev_hdr: &netlog::EventHeader,
803 ev: &h2::Event, constants: &netlog::constants::Constants,
804 stream_bind: &StreamBindingMap,
805 ) {
806 let rel_event_time = (ev_hdr.time_num - session_start_time) as f64;
807
808 match ev {
809 Http2SessionSendSettings(e) => {
810 match Http2Settings::try_from(e.params.settings.as_slice()) {
811 Ok(v) => self.h2_client_settings = v,
812
813 Err(e) => error!("{}", e),
814 }
815 },
816
817 Http2SessionRecvSetting(e) => {
818 let re = Regex::new(H2_RECV_SETTING_PATTERN).unwrap();
819
820 if let Some(m) =
821 re.captures(&e.params.id).and_then(|caps| caps.get(1))
822 {
823 if let Ok(id) = m.as_str().parse::<u16>() {
824 self.h2_server_settings.set_from_wire(id, e.params.value);
825 } else {
826 error!("parsing H2 setting {:?}", e.params);
827 }
828 } else {
829 error!("parsing H2 setting {:?}", e.params);
830 }
831 },
832 Http2SessionSendHeaders(e) => {
833 let req = self.get_or_insert_http_req(e.params.stream_id as u64);
834 req.time_first_headers_tx.get_or_insert(rel_event_time);
835 req.set_request_info_from_netlog(&e.params.headers);
836
837 if let Some(sb) = stream_bind.get(&e.params.source_dependency.id)
838 {
839 req.time_discovery = Some(
842 (sb.request_discovery_time
843 .wrapping_sub(session_start_time)
844 as i64) as f64,
845 );
846 }
847
848 self.h2_concurrent_requests += 1;
849 },
850
851 Http2SessionRecvHeaders(e) => {
852 let req = self.get_or_insert_http_req(e.params.stream_id as u64);
853 req.time_first_headers_rx.get_or_insert(rel_event_time);
854 req.set_response_info_from_netlog(&e.params.headers);
855
856 if e.params.fin {
857 self.h2_concurrent_requests -= 1;
858 }
859 },
860
861 Http2SessionRecvData(e) => {
862 let req = self.get_or_insert_http_req(e.params.stream_id as u64);
863
864 req.time_first_data_rx.get_or_insert(rel_event_time);
865
866 let _ = req.time_last_data_rx.insert(rel_event_time);
867
868 let length = e.params.size as u64;
869 req.time_data_rx_set.push((rel_event_time, length));
870 self.largest_data_frame_rx_length_global = std::cmp::max(
871 self.largest_data_frame_rx_length_global,
872 length,
873 );
874
875 let s = self
876 .received_data_frames
877 .entry(e.params.stream_id as u64)
878 .or_default();
879
880 s.push((rel_event_time, e.params.size as u64));
881
882 let s = self
883 .received_data_frames_count_based
884 .entry(e.params.stream_id as u64)
885 .or_default();
886 s.push((
887 self.total_received_data_frame_count,
888 e.params.size as u64,
889 ));
890
891 self.total_received_data_frame_count += 1;
892
893 let s = self
894 .received_data_cumulative
895 .entry(e.params.stream_id as u64)
896 .or_default();
897
898 let received_data_cumulative = if let Some(last) = s.last() {
899 last.1 + e.params.size as u64
900 } else {
901 e.params.size as u64
903 };
904
905 s.push((rel_event_time, received_data_cumulative));
906
907 let s = self
908 .received_data_cumulative_max
909 .entry(e.params.stream_id as u64)
910 .or_default();
911
912 *s = std::cmp::max(*s, received_data_cumulative);
913
914 let s = self
915 .http_requests
916 .entry(e.params.stream_id as u64)
917 .or_default();
918
919 s.server_transferred_bytes = std::cmp::max(
920 s.server_transferred_bytes,
921 NaOption::new(Some(received_data_cumulative)),
922 );
923
924 if e.params.fin {
925 self.h2_concurrent_requests -= 1;
926 }
927
928 let cumulative = if let Some(last) =
929 self.netlog_h2_stream_received_connection_cumulative.last()
930 {
931 last.1 + e.params.size as u64
932 } else {
933 e.params.size as u64
935 };
936
937 self.netlog_h2_stream_received_connection_cumulative
938 .push((rel_event_time, cumulative));
939 },
940
941 Http2SessionSendData(e) => {
942 let req = self.get_or_insert_http_req(e.params.stream_id as u64);
943
944 req.time_first_data_tx.get_or_insert(rel_event_time);
945
946 let _ = req.time_last_data_tx.insert(rel_event_time);
947
948 let length = e.params.size as u64;
949 req.time_data_tx_set.push((rel_event_time, length));
950 self.largest_data_frame_tx_length_global = std::cmp::max(
951 self.largest_data_frame_tx_length_global,
952 length,
953 );
954
955 let s = self
956 .sent_data_frames
957 .entry(e.params.stream_id as u64)
958 .or_default();
959
960 s.push((rel_event_time, e.params.size as u64));
961
962 let s = self
963 .sent_data_frames_count_based
964 .entry(e.params.stream_id as u64)
965 .or_default();
966 s.push((self.total_sent_data_frame_count, e.params.size as u64));
967
968 self.total_sent_data_frame_count += 1;
969
970 let s = self
972 .h2_send_window_updates_balanced
973 .entry(e.params.stream_id)
974 .or_default();
975 s.push((rel_event_time, -(e.params.size as i32)));
976
977 let s = self
978 .sent_data_cumulative
979 .entry(e.params.stream_id as u64)
980 .or_default();
981
982 let sent_data_cumulative = if let Some(last) = s.last() {
983 last.1 + e.params.size as u64
984 } else {
985 e.params.size as u64
987 };
988
989 s.push((rel_event_time, sent_data_cumulative));
990
991 let s = self
992 .sent_data_cumulative_max
993 .entry(e.params.stream_id as u64)
994 .or_default();
995
996 *s = std::cmp::max(*s, sent_data_cumulative);
997
998 let s = self
999 .http_requests
1000 .entry(e.params.stream_id as u64)
1001 .or_default();
1002
1003 s.client_transferred_bytes = std::cmp::max(
1004 s.client_transferred_bytes,
1005 NaOption::new(Some(sent_data_cumulative)),
1006 );
1007 },
1008
1009 Http2SessionRecvWindowUpdate(e) => {
1011 let s = self
1012 .h2_send_window_updates_balanced
1013 .entry(e.params.stream_id)
1014 .or_default();
1015 s.push((rel_event_time, e.params.delta));
1016 },
1017
1018 Http2SessionSendWindowUpdate(e) => {
1020 let s = self
1021 .h2_send_window_updates_balanced
1022 .entry(e.params.stream_id)
1023 .or_default();
1024 s.push((rel_event_time, e.params.delta));
1025
1026 let s = self
1027 .h2_send_window_updates_absolute
1028 .entry(e.params.stream_id)
1029 .or_default();
1030 s.push((rel_event_time, e.params.delta as u64));
1032 },
1033
1034 Http2SessionClose(e) => {
1035 self.h2_session_close = Some(H2SessionClose {
1036 session_id: self.session_id.unwrap_or(-1),
1037 sni: self.host.clone().unwrap_or("ERROR UNKNOWN".to_string()),
1038 details: e.params.description.clone(),
1039 net_err: e.params.net_error,
1040 net_err_pretty: NaOption::new(
1041 constants
1042 .net_error_id_keyed
1043 .get(&e.params.net_error)
1044 .cloned(),
1045 ),
1046 });
1047 },
1048
1049 Http2SessionSendRstStream(e) => {
1050 if let Some(req) =
1051 self.http_requests.get_mut(&(e.params.stream_id as u64))
1052 {
1053 req.h2_stream_reset_sent = Some(H2StreamReset {
1054 error: e.params.error_code.clone(),
1055 description: e.params.description.clone(),
1056 });
1057 }
1058 },
1059
1060 Http2SessionRecvRstStream(e) => {
1061 if let Some(req) =
1062 self.http_requests.get_mut(&(e.params.stream_id as u64))
1063 {
1064 req.h2_stream_reset_receive = Some(H2StreamReset {
1065 error: e.params.error_code.clone(),
1066 description: "".to_string(),
1067 });
1068 }
1069 self.h2_concurrent_requests -= 1;
1070 },
1071
1072 Http2Session(_) => (),
1074 Http2SessionInitialized(_) => (),
1075 Http2SessionUpdateRecvWindow(_) => (),
1076 Http2SessionUpdateSendWindow(_) => (),
1077 Http2SessionUpdateStreamsSendWindowSize(_) => (),
1078 Http2SessionStalledMaxStreams(_) => (),
1079
1080 Http2StreamUpdateSendWindow(_) => (),
1081 Http2StreamUpdateRecvWindow(_) => (),
1082 Http2StreamStalledByStreamSendWindow(_) => (),
1083 Http2SessionPing(_) => (),
1084
1085 Http2SessionRecvGoaway(_) => (),
1086 }
1087 }
1088
1089 pub fn consume_qlog_event(
1090 &mut self, event: &qlog::events::Event, process_acks: bool,
1091 ) {
1092 let ev_time = event.time;
1093
1094 if ev_time > self.last_event_time {
1095 self.last_event_time = ev_time;
1096 }
1097
1098 match &event.data {
1099 EventData::QuicParametersSet(v) =>
1100 self.consume_qlog_transport_parameters_set(v),
1101
1102 EventData::QuicPacketReceived(v) =>
1103 self.consume_qlog_packet_received(v, ev_time, process_acks),
1104
1105 EventData::QuicPacketSent(v) =>
1106 self.consume_qlog_packet_sent(v, ev_time),
1107
1108 EventData::QuicStreamDataMoved(v) =>
1109 self.consume_qlog_stream_data_moved(v, ev_time),
1110
1111 EventData::QuicMetricsUpdated(v) =>
1112 self.consume_qlog_metrics_updated(v, ev_time),
1113
1114 EventData::QuicCongestionStateUpdated(v) =>
1115 self.consume_qlog_congestion_state_updated(v, ev_time),
1116
1117 EventData::Http3FrameCreated(v) => match self.vantage_point {
1118 VantagePoint::Client =>
1119 self.consume_qlog_h3_frame_created_client(v, ev_time),
1120 VantagePoint::Server =>
1121 self.consume_qlog_h3_frame_created_server(v, ev_time),
1122 },
1123
1124 EventData::Http3FrameParsed(v) => match self.vantage_point {
1125 VantagePoint::Client =>
1126 self.consume_qlog_h3_frame_parsed_client(v, ev_time),
1127 VantagePoint::Server =>
1128 self.consume_qlog_h3_frame_parsed_server(v, ev_time),
1129 },
1130
1131 _ => (), }
1133 }
1134
1135 pub fn with_qlog_events(
1136 events: &[qlog::events::Event], vantage_point: &qlog::VantagePointType,
1137 process_acks: bool,
1138 ) -> Self {
1139 let vp = match vantage_point {
1140 qlog::VantagePointType::Client => VantagePoint::Client,
1141 qlog::VantagePointType::Server => VantagePoint::Server,
1142 _ => panic!("unknown vantage point type"),
1143 };
1144
1145 let mut ds = Datastore {
1146 vantage_point: vp,
1147 ..Default::default()
1148 };
1149
1150 for event in events {
1151 ds.consume_qlog_event(event, process_acks);
1152 }
1153
1154 ds.hydrate_http_requests();
1155 ds.finalize();
1156
1157 ds
1158 }
1159
1160 pub fn finalize(&mut self) {
1161 if let Some(last) = self.local_cwnd.last().cloned() {
1162 self.local_cwnd.push((self.last_event_time, last.1));
1163 }
1164
1165 if let Some(last) = self.local_pacing_rate.last().cloned() {
1166 trace!("pushing last {:?}", last);
1167 self.local_pacing_rate.push((self.last_event_time, last.1));
1168 }
1169 }
1170
1171 pub fn hydrate_http_requests(&mut self) {
1172 for req in self.http_requests.values_mut() {
1173 req.calculate_deltas();
1174 req.calculate_upload_download_rate();
1175 }
1176 }
1177
1178 fn consume_qlog_transport_parameters_set(
1179 &mut self, tp: &qlog::events::quic::ParametersSet,
1180 ) {
1181 match tp.initiator {
1182 Some(TransportInitiator::Local) => {
1183 if let Some(max_data) = tp.initial_max_data {
1184 self.sent_max_data.push((0.0, max_data));
1185 }
1186
1187 if let Some(max_stream_data) =
1188 tp.initial_max_stream_data_bidi_local
1189 {
1190 self.local_init_max_stream_data_bidi_local = max_stream_data;
1191 }
1192
1193 if let Some(max_stream_data) = tp.initial_max_stream_data_uni {
1194 self.local_init_max_stream_data_uni = max_stream_data;
1195 }
1196 },
1197
1198 Some(TransportInitiator::Remote) => {
1199 if let Some(max_data) = tp.initial_max_data {
1200 self.received_max_data.push((0.0, max_data));
1201 }
1202
1203 if let Some(max_stream_data) =
1204 tp.initial_max_stream_data_bidi_local
1205 {
1206 self.peer_init_max_stream_data_bidi_local = max_stream_data;
1207 }
1208
1209 if let Some(max_stream_data) =
1210 tp.initial_max_stream_data_bidi_remote
1211 {
1212 self.peer_init_max_stream_data_bidi_remote = max_stream_data;
1213 }
1214
1215 if let Some(max_stream_data) = tp.initial_max_stream_data_uni {
1216 self.peer_init_max_stream_data_uni = max_stream_data;
1217 }
1218 },
1219
1220 _ => unimplemented!(),
1221 }
1222 }
1223
1224 fn consume_qlog_packet_received(
1225 &mut self, pr: &qlog::events::quic::PacketReceived, ev_time: f64,
1226 process_acks: bool,
1227 ) {
1228 if let Some(frames) = &pr.frames {
1229 for frame in frames {
1230 match frame {
1231 QuicFrame::Ack { acked_ranges, .. } => {
1232 if process_acks {
1233 if let Some(ack_ranges) = acked_ranges {
1234 let ty = PacketType::from_qlog_packet_type(
1235 &pr.header.packet_type,
1236 );
1237 if let Some(pkt_space) =
1238 self.packet_sent.get_mut(&ty)
1239 {
1240 for range in ack_ranges {
1241 pkt_space
1244 .range_mut(range.as_range_inclusive())
1245 .for_each(|e| e.1.acked = Some(true));
1246 }
1247 }
1248 }
1249 }
1250 },
1251
1252 QuicFrame::MaxData { maximum, .. } => {
1253 self.received_max_data.push((ev_time, *maximum));
1254 },
1255
1256 QuicFrame::MaxStreamData {
1257 stream_id, maximum, ..
1258 } => {
1259 let init_val = if is_bidi(*stream_id) {
1260 self.peer_init_max_stream_data_bidi_remote
1261 } else {
1262 self.peer_init_max_stream_data_uni
1263 };
1264 self.received_stream_max_data_tracker
1265 .update(*stream_id, *maximum, ev_time, init_val);
1266 },
1267
1268 QuicFrame::ResetStream { stream_id, .. } => {
1269 let s = self
1270 .received_reset_stream
1271 .entry(*stream_id)
1272 .or_default();
1273 s.push(frame.clone());
1274 },
1275
1276 QuicFrame::Stream {
1277 stream_id,
1278 offset,
1279 raw,
1280 ..
1281 } => {
1282 let length = raw
1283 .clone()
1284 .unwrap_or_default()
1285 .payload_length
1286 .unwrap_or_default();
1287 let s = self
1288 .received_stream_frames
1289 .entry(*stream_id)
1290 .or_default();
1291 s.push((ev_time, StreamDatapoint {
1292 length,
1293 offset: offset.unwrap_or_default(),
1294 }));
1295
1296 let s = self
1297 .received_stream_frames_count_based
1298 .entry(*stream_id)
1299 .or_default();
1300 s.push((
1301 self.total_received_stream_frame_count,
1302 StreamDatapoint {
1303 length,
1304 offset: offset.unwrap_or_default(),
1305 },
1306 ));
1307
1308 self.total_received_stream_frame_count += 1;
1309 },
1310
1311 _ => (),
1312 }
1313 }
1314 }
1315 }
1316
1317 fn consume_qlog_packet_sent(
1318 &mut self, ps: &qlog::events::quic::PacketSent, ev_time: f64,
1319 ) {
1320 if ps.header.packet_number.is_none() {
1322 return;
1323 }
1324
1325 let packet_type =
1326 PacketType::from_qlog_packet_type(&ps.header.packet_type);
1327 let packet_info = PacketInfoStub {
1328 acked: None,
1329 raw: ps.raw.clone(),
1330 created_time: ev_time,
1331 send_at_time: ps.send_at_time,
1332 ty: packet_type,
1333 number: ps.header.packet_number.unwrap(),
1334 };
1335
1336 let s = self.packet_sent.entry(packet_type).or_default();
1337
1338 s.insert(ps.header.packet_number.unwrap(), packet_info);
1339
1340 let event_time = ps.send_at_time.unwrap_or(ev_time);
1343
1344 if let Some(frames) = &ps.frames {
1345 for frame in frames {
1346 match frame {
1347 QuicFrame::Ack { .. } => {
1348 },
1350
1351 QuicFrame::MaxData { maximum, .. } => {
1352 self.sent_max_data.push((event_time, *maximum));
1353 },
1354
1355 QuicFrame::MaxStreamData {
1356 stream_id, maximum, ..
1357 } => {
1358 let init_val = if is_bidi(*stream_id) {
1359 self.local_init_max_stream_data_bidi_local
1360 } else {
1361 self.local_init_max_stream_data_uni
1362 };
1363 self.sent_stream_max_data_tracker
1364 .update(*stream_id, *maximum, ev_time, init_val);
1365 },
1366
1367 QuicFrame::ResetStream { stream_id, .. } => {
1368 let s =
1369 self.sent_reset_stream.entry(*stream_id).or_default();
1370 s.push(frame.clone());
1371 },
1372
1373 QuicFrame::Stream { stream_id, .. } => {
1374 let s = self
1375 .sent_stream_frames
1376 .entry(*stream_id)
1377 .or_default();
1378 s.push((event_time, frame.clone()));
1379
1380 let s = self
1381 .sent_stream_frames_count_based
1382 .entry(*stream_id)
1383 .or_default();
1384 s.push((
1385 self.total_sent_stream_frame_count,
1386 frame.clone(),
1387 ));
1388
1389 self.total_sent_stream_frame_count += 1;
1390 },
1391
1392 QuicFrame::DataBlocked { limit, .. } => {
1393 trace!(
1394 "todo DATA_BLOCKED t={} limit={}",
1395 event_time,
1396 limit
1397 );
1398 },
1399
1400 QuicFrame::StreamDataBlocked {
1401 stream_id, limit, ..
1402 } => {
1403 trace!(
1404 "todo STREAM_DATA_BLOCKED t={} stream={} limit={}",
1405 event_time,
1406 stream_id,
1407 limit
1408 );
1409 },
1410
1411 _ => (),
1412 }
1413 }
1414 }
1415 }
1416
1417 fn consume_qlog_stream_data_moved(
1418 &mut self, dm: &qlog::events::quic::StreamDataMoved, ev_time: f64,
1419 ) {
1420 if let Some(recipient) = &dm.to {
1421 let tracker = match recipient {
1422 qlog::events::DataRecipient::Application =>
1423 &mut self.stream_buffer_reads_tracker,
1424 qlog::events::DataRecipient::Transport =>
1425 &mut self.stream_buffer_writes_tracker,
1426 qlog::events::DataRecipient::Dropped =>
1427 &mut self.stream_buffer_dropped_tracker,
1428 _ => todo!(),
1429 };
1430
1431 if let Some(stream_id) = dm.stream_id {
1432 if let Some(raw) = &dm.raw {
1433 if let (Some(offset), Some(length)) = (dm.offset, raw.length)
1434 {
1435 tracker.update(
1436 stream_id,
1437 StreamAccess { offset, length },
1438 ev_time,
1439 );
1440 }
1441 }
1442 }
1443 }
1444 }
1445
1446 fn consume_qlog_metrics_updated(
1447 &mut self, mu: &qlog::events::quic::RecoveryMetricsUpdated, ev_time: f64,
1448 ) {
1449 if let Some(cwnd) = mu.congestion_window {
1450 self.local_cwnd.push((ev_time, cwnd));
1451 }
1452
1453 if let Some(bif) = mu.bytes_in_flight {
1454 self.local_bytes_in_flight.push((ev_time, bif));
1455 }
1456
1457 if let Some(rtt) = mu.min_rtt {
1458 self.local_min_rtt.push((ev_time, rtt));
1459 }
1460
1461 if let Some(rtt) = mu.latest_rtt {
1462 self.local_latest_rtt.push((ev_time, rtt));
1463 }
1464
1465 if let Some(rtt) = mu.smoothed_rtt {
1466 self.local_smoothed_rtt.push((ev_time, rtt));
1467 }
1468
1469 if let Some(thresh) = mu.ssthresh {
1470 self.local_ssthresh.push((ev_time, thresh));
1471 }
1472
1473 if let Some(pacing_rate) = mu.pacing_rate {
1474 self.local_pacing_rate.push((ev_time, pacing_rate));
1475 }
1476
1477 if let Some(rate) =
1479 mu.ex_data.get("cf_delivery_rate").and_then(|v| v.as_u64())
1480 {
1481 self.local_delivery_rate.push((ev_time, rate));
1482 }
1483
1484 if let Some(rate) =
1485 mu.ex_data.get("cf_send_rate").and_then(|v| v.as_u64())
1486 {
1487 self.local_send_rate.push((ev_time, rate));
1488 }
1489
1490 if let Some(rate) = mu.ex_data.get("cf_ack_rate").and_then(|v| v.as_u64())
1491 {
1492 self.local_ack_rate.push((ev_time, rate));
1493 }
1494 }
1495
1496 fn consume_qlog_congestion_state_updated(
1497 &mut self, csu: &qlog::events::quic::CongestionStateUpdated, ev_time: f64,
1498 ) {
1499 if let Some(point) = self.local_cwnd.last() {
1500 self.congestion_state_updates.push((
1503 ev_time,
1504 point.1,
1505 csu.new.clone(),
1506 ));
1507 }
1508 }
1509
1510 fn get_or_insert_http_req(&mut self, stream_id: u64) -> &mut HttpRequestStub {
1511 self.http_requests
1512 .entry(stream_id)
1513 .or_insert(HttpRequestStub {
1514 stream_id,
1515 request_actor: self.vantage_point.into(),
1516 ..Default::default()
1517 })
1518 }
1519
1520 fn consume_qlog_h3_frame_created_client(
1521 &mut self, fc: &qlog::events::http3::FrameCreated, ev_time: f64,
1522 ) {
1523 match &fc.frame {
1524 Http3Frame::Headers { headers, .. } => {
1525 let req = self.get_or_insert_http_req(fc.stream_id);
1526 req.time_first_headers_tx.get_or_insert(ev_time);
1527 req.set_request_info_from_qlog(headers);
1528 },
1529
1530 Http3Frame::Data { .. } => {
1531 let req = self.get_or_insert_http_req(fc.stream_id);
1532
1533 req.time_first_data_tx.get_or_insert(ev_time);
1534
1535 let _ = req.time_last_data_tx.insert(ev_time);
1536
1537 let length = fc.length.unwrap_or_default();
1538 req.time_data_tx_set.push((ev_time, length));
1539 self.largest_data_frame_tx_length_global = std::cmp::max(
1540 self.largest_data_frame_tx_length_global,
1541 length,
1542 );
1543 },
1544
1545 Http3Frame::PriorityUpdate {
1546 stream_id: Some(stream_id),
1547 priority_field_value,
1548 ..
1549 } => {
1550 let req: &mut HttpRequestStub =
1551 self.get_or_insert_http_req(*stream_id);
1552 req.priority_updates.push(priority_field_value.clone());
1553 },
1554
1555 _ => (),
1557 }
1558 }
1559
1560 fn consume_qlog_h3_frame_created_server(
1561 &mut self, fc: &qlog::events::http3::FrameCreated, ev_time: f64,
1562 ) {
1563 match &fc.frame {
1564 Http3Frame::Headers { headers, .. } => {
1565 let req = self.get_or_insert_http_req(fc.stream_id);
1566 req.time_first_headers_tx.get_or_insert(ev_time);
1567 req.set_response_info_from_qlog(headers);
1568 },
1569
1570 Http3Frame::Data { .. } => {
1571 let req = self.get_or_insert_http_req(fc.stream_id);
1572
1573 req.time_first_data_tx.get_or_insert(ev_time);
1574
1575 let _ = req.time_last_data_tx.insert(ev_time);
1576
1577 let length = fc.length.unwrap_or_default();
1578 req.time_data_tx_set.push((ev_time, length));
1579 self.largest_data_frame_tx_length_global = std::cmp::max(
1580 self.largest_data_frame_tx_length_global,
1581 length,
1582 );
1583 },
1584
1585 _ => (),
1587 }
1588 }
1589
1590 fn consume_qlog_h3_frame_parsed_client(
1591 &mut self, fp: &qlog::events::http3::FrameParsed, ev_time: f64,
1592 ) {
1593 match &fp.frame {
1594 Http3Frame::Headers { headers, .. } => {
1595 let req = self.get_or_insert_http_req(fp.stream_id);
1596 req.time_first_headers_rx.get_or_insert(ev_time);
1597
1598 req.set_response_info_from_qlog(headers);
1599 },
1600
1601 Http3Frame::Data { .. } => {
1602 let req = self.get_or_insert_http_req(fp.stream_id);
1603
1604 req.time_first_data_rx.get_or_insert(ev_time);
1605
1606 let _ = req.time_last_data_rx.insert(ev_time);
1607
1608 let length = fp.length.unwrap_or_default();
1610 req.time_data_rx_set.push((ev_time, length));
1611 self.largest_data_frame_rx_length_global = std::cmp::max(
1612 self.largest_data_frame_rx_length_global,
1613 length,
1614 );
1615
1616 let s =
1617 self.received_data_frames.entry(fp.stream_id).or_default();
1618
1619 s.push((ev_time, length));
1620 },
1621
1622 _ => (),
1624 }
1625 }
1626
1627 fn consume_qlog_h3_frame_parsed_server(
1628 &mut self, fp: &qlog::events::http3::FrameParsed, ev_time: f64,
1629 ) {
1630 match &fp.frame {
1631 Http3Frame::Headers { headers, .. } => {
1632 let req = self.get_or_insert_http_req(fp.stream_id);
1633 req.time_first_headers_rx.get_or_insert(ev_time);
1634 req.path = NaOption::new(find_header_value(headers, ":path"));
1635 req.client_pri_hdr =
1636 NaOption::new(find_header_value(headers, "priority"));
1637 },
1638
1639 Http3Frame::Data { .. } => {
1640 let req = self.get_or_insert_http_req(fp.stream_id);
1641
1642 req.time_first_data_rx.get_or_insert(ev_time);
1643
1644 let _ = req.time_last_data_rx.insert(ev_time);
1645
1646 let length = fp.length.unwrap_or_default();
1647 req.time_data_rx_set.push((ev_time, length));
1648 self.largest_data_frame_rx_length_global = std::cmp::max(
1649 self.largest_data_frame_rx_length_global,
1650 length,
1651 );
1652 },
1653
1654 Http3Frame::PriorityUpdate {
1655 stream_id: Some(stream_id),
1656 priority_field_value,
1657 ..
1658 } => {
1659 let req = self.get_or_insert_http_req(*stream_id);
1660 req.priority_updates.push(priority_field_value.clone());
1661 },
1662
1663 _ => (),
1665 }
1666 }
1667
1668 pub fn with_sqlog_reader_events(
1669 events: &[qlog::reader::Event], vantage_point: &qlog::VantagePointType,
1670 process_acks: bool,
1671 ) -> Self {
1672 let vp = match vantage_point {
1673 qlog::VantagePointType::Client => VantagePoint::Client,
1674 qlog::VantagePointType::Server => VantagePoint::Server,
1675 _ => panic!("unknown vantage point type"),
1676 };
1677
1678 let mut ds = Datastore {
1679 total_sent_stream_frame_count: 0,
1680 vantage_point: vp,
1681 ..Default::default()
1682 };
1683
1684 for event in events {
1685 match event {
1686 qlog::reader::Event::Qlog(ev) => {
1687 ds.consume_qlog_event(ev, process_acks);
1688 },
1689
1690 qlog::reader::Event::Json(ev) => {
1691 error!("unhandled Json event {:?}", ev);
1693 },
1694 }
1695 }
1696
1697 ds.hydrate_http_requests();
1698 ds.finalize();
1699
1700 ds
1701 }
1702}
1703
1704#[derive(Tabled)]
1705pub struct NetlogSession {
1706 #[tabled(rename = "ID")]
1707 session_id: i64,
1708 #[tabled(rename = "Protocol")]
1709 application_proto: ApplicationProto,
1710 #[tabled(rename = "SNI")]
1711 host: String,
1712 start_time: u64,
1713}
1714
1715impl Debug for NetlogSession {
1716 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1717 write!(
1718 f,
1719 "app proto={:?}, host={}",
1720 self.application_proto, self.host
1721 )
1722 }
1723}
1724
1725pub struct RequestDiscovery {
1726 pub time: u64,
1727 pub stream_job_id: Option<i64>,
1728}
1729
1730pub struct StreamBind {
1731 pub time: u64,
1732 pub request_discovery_id: i64,
1733 pub request_discovery_time: u64,
1734}
1735
1736pub type RequestDiscoveryMap = BTreeMap<i64, RequestDiscovery>;
1737pub type StreamBindingMap = BTreeMap<i64, StreamBind>;
1738
1739#[derive(Debug)]
1740pub struct ReqOverH3 {
1741 pub id: i64,
1742 pub discover_time: u64,
1743 pub session_id: Option<i64>,
1744 pub quic_stream_id: Option<u64>,
1745}
1746
1747pub type RequestOverH3Map = BTreeMap<i64, Vec<ReqOverH3>>;
1748
1749pub fn with_netlog_reader<R: std::io::BufRead>(
1750 reader: &mut R, hostname_filter: HashSet<String>,
1751 constants: &netlog::constants::Constants,
1752) -> (Vec<LogFileData>, BTreeMap<i64, NetlogSession>) {
1753 read_netlog_record(reader);
1755
1756 let mut sessions: BTreeMap<i64, NetlogSession> = BTreeMap::new();
1757 let mut session_events: BTreeMap<
1758 i64,
1759 Vec<(netlog::EventHeader, netlog::Event)>,
1760 > = BTreeMap::new();
1761
1762 let mut h3_session_requests: RequestOverH3Map = BTreeMap::new();
1763
1764 let mut req_id_to_session_id: BTreeMap<i64, i64> = BTreeMap::new();
1765
1766 let mut request_discovery: RequestDiscoveryMap = BTreeMap::new();
1767 let mut stream_bind: StreamBindingMap = BTreeMap::new();
1768
1769 while let Some(event) = read_netlog_record(reader) {
1770 let res: Result<netlog::EventHeader, serde_json::Error> =
1771 serde_json::from_slice(&event);
1772
1773 match res {
1774 Ok(mut event_hdr) => {
1775 event_hdr.populate_strings(constants);
1776 event_hdr.time_num = event_hdr.time.parse::<u64>().unwrap();
1777
1778 if event_hdr.phase_string == "PHASE_BEGIN" {
1784 match event_hdr.ty_string.as_str() {
1785 "QUIC_SESSION" => {
1786 let ev: QuicSessionEvent =
1787 serde_json::from_slice(&event).unwrap();
1788
1789 let host =
1792 format!("{}:{}", ev.params.host, ev.params.port,);
1793
1794 sessions.insert(event_hdr.source.id, NetlogSession {
1795 session_id: event_hdr.source.id,
1796 application_proto: ApplicationProto::Http3,
1797 host: host.clone(),
1798 start_time: event_hdr
1799 .time
1800 .parse::<u64>()
1801 .unwrap(),
1802 });
1803
1804 let do_insert = hostname_filter.is_empty() ||
1805 hostname_filter.contains(&host);
1806
1807 if do_insert {
1808 session_events
1809 .insert(event_hdr.source.id, Vec::new());
1810
1811 h3_session_requests
1812 .insert(event_hdr.source.id, Vec::new());
1813 }
1814 },
1815
1816 "HTTP2_SESSION" => {
1817 let ev: Http2SessionEvent =
1818 serde_json::from_slice(&event).unwrap();
1819
1820 sessions.insert(event_hdr.source.id, NetlogSession {
1821 session_id: event_hdr.source.id,
1822 application_proto: ApplicationProto::Http2,
1823 host: ev.params.host.clone(),
1824 start_time: event_hdr
1825 .time
1826 .parse::<u64>()
1827 .unwrap(),
1828 });
1829
1830 let do_insert = hostname_filter.is_empty() ||
1831 hostname_filter.contains(&ev.params.host);
1832
1833 if do_insert {
1834 session_events
1835 .insert(event_hdr.source.id, Vec::new());
1836 }
1837 },
1838
1839 "CORS_REQUEST" => {
1840 request_discovery.insert(
1843 event_hdr.source.id,
1844 RequestDiscovery {
1845 time: event_hdr.time_num,
1846 stream_job_id: None,
1847 },
1848 );
1849 },
1850
1851 _ => (),
1852 }
1853 }
1854
1855 if event_hdr.ty_string.starts_with("HTTP_") {
1856 let event = netlog::http::parse_event(&event_hdr, &event);
1857
1858 #[allow(clippy::single_match)]
1861 match event {
1862 Some(netlog::Event::Http(e)) => {
1863 match e {
1864 http::Event::HttpStreamJobBoundToRequest(v) => {
1865 let request_discovery_id = v.params.source_dependency.id;
1866 if let Some(rd) = request_discovery.get_mut(&request_discovery_id) {
1867 stream_bind.insert(event_hdr.source.id, StreamBind{time: event_hdr.time_num, request_discovery_id, request_discovery_time: rd.time});
1868 rd.stream_job_id = Some(event_hdr.source.id);
1869 }
1870 },
1871
1872 http::Event::HttpStreamRequestBoundToQuicSession(v) => {
1873 let request_discovery_id = event_hdr.source.id;
1874 if let Some(rd) = request_discovery.get(&request_discovery_id) {
1875
1876 if let Some(session_requests) = h3_session_requests.get_mut(&v.params.source_dependency.id) {
1877 let req = ReqOverH3{id: event_hdr.source.id, discover_time: rd.time, session_id: Some(v.params.source_dependency.id), quic_stream_id: None };
1878 session_requests.push(req);
1879
1880 req_id_to_session_id.insert(event_hdr.source.id, v.params.source_dependency.id);
1882 }
1883 }
1884 }
1885
1886 http::Event::HttpTransactionQuicSendRequestHeaders(v) => {
1887 let req_id = event_hdr.source.id;
1888 if let Some(session_id) = req_id_to_session_id.get(&req_id) {
1889 if let Some(reqs) = h3_session_requests.get_mut(session_id) {
1890 for req in reqs {
1892 if req.id == req_id {
1893 req.quic_stream_id = Some(v.params.quic_stream_id);
1894 break;
1895 }
1896 }
1897 }
1898 }
1899 }
1900
1901 _ => (),
1902 }
1903 },
1904
1905 _ => (),
1907 }
1908 }
1909
1910 if let Some(session) =
1911 session_events.get_mut(&event_hdr.source.id)
1912 {
1913 if let Some(ev) = netlog::parse_event(&event_hdr, &event) {
1914 session.push((event_hdr, ev));
1915 }
1916 }
1917 },
1918
1919 Err(e) => {
1920 error!("Error deserializing: {}", e);
1921 error!("input value {}", String::from_utf8_lossy(&event));
1922
1923 },
1925 }
1926 }
1927
1928 println!("All sessions in this netlog = {:#?}", sessions);
1929
1930 let mut log_file_data = Vec::new();
1931
1932 for (session_id, details) in &sessions {
1933 if let Some(events) = session_events.get(session_id) {
1934 let mut ds = Datastore {
1935 session_id: Some(*session_id),
1936 application_proto: details.application_proto,
1937 host: Some(details.host.clone()),
1938 total_sent_stream_frame_count: 0,
1939 ..Default::default()
1940 };
1941
1942 for (ev_hdr, event) in events {
1943 ds.consume_netlog_event(
1944 details.start_time,
1945 ev_hdr,
1946 event,
1947 constants,
1948 &stream_bind,
1949 h3_session_requests.get(session_id),
1950 );
1951 }
1952
1953 ds.hydrate_http_requests();
1954 ds.finalize();
1955 log_file_data.push(LogFileData {
1956 datastore: ds,
1957 raw: Netlog,
1958 });
1959 }
1960 }
1961
1962 (log_file_data, sessions)
1963}