Skip to main content

qlog_dancer/
datastore.rs

1// Copyright (C) 2025, Cloudflare, Inc.
2// All rights reserved.
3//
4// Redistribution and use in source and binary forms, with or without
5// modification, are permitted provided that the following conditions are
6// met:
7//
8//     * Redistributions of source code must retain the above copyright notice,
9//       this list of conditions and the following disclaimer.
10//
11//     * Redistributions in binary form must reproduce the above copyright
12//       notice, this list of conditions and the following disclaimer in the
13//       documentation and/or other materials provided with the distribution.
14//
15// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS
16// IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO,
17// THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR
18// PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR
19// CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL,
20// EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO,
21// PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR
22// PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF
23// LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING
24// NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS
25// SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
26
27use 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    // There are several packet spaces, so store a map of all packets sent
260    // according to packet space. Each space then contains a map of packet
261    // header info keyed off the packet number.
262    pub packet_sent: HashMap<PacketType, BTreeMap<u64, PacketInfoStub>>,
263    pub packet_received: HashMap<PacketType, BTreeMap<u64, PacketInfoStub>>,
264
265    // TODO: netlog packet sent happens after frame, so we can't detect the
266    // packet type properly. Stick all in one bucket for now and accept we'll
267    // alias packet numbers that overlap between spaced
268    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    /// Tracks per-stream max data: full history, current max, and cumulative
290    /// sum.
291    pub received_stream_max_data_tracker: StreamMaxTracker,
292
293    pub sent_max_data: Vec<QlogPointu64>,
294
295    /// Tracks per-stream max data: full history, current max, and cumulative
296    /// sum.
297    pub sent_stream_max_data_tracker: StreamMaxTracker,
298
299    /// Tracks stream buffer reads: per-stream history, current max, and running
300    /// sum.
301    pub stream_buffer_reads_tracker: StreamBufferTracker,
302
303    /// Tracks stream buffer writes: per-stream history, current max, and
304    /// running sum.
305    pub stream_buffer_writes_tracker: StreamBufferTracker,
306
307    /// Tracks stream buffer dropped: per-stream history, current max, and
308    /// running sum.
309    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    // Balance against incoming data to make it easier to plot in some cases
348    pub h2_send_window_updates_balanced: BTreeMap<u32, Vec<(f64, i32)>>,
349
350    // Just store raw updates for clear absolute values
351    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            // nothing to do for this type just now
389            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                    // insert a 0'th point
455                    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                // Go back and update the Ack frame type if there was one
516                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                    // At this stage, we don't know the packet type we sent the
554                    // ACK in, because it comes later in the netlog. Insert with
555                    // a placeholder now, and we'll update later in
556                    // QuicSessionPacketSent handler. Ugly but functional.
557                    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                    // For netlogs, it is assumed that the last packet received
570                    // relates to this event.
571                    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                // Non-request streams can be reset, we don't care about them
603                // right now
604                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                // Non-request streams can be reset, we don't care about them
619                // right now
620                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                // Non-request streams can be stopped, we don't care about them
635                // right now
636                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            // ignore the other events for now
682            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                                // Hat-tip Olivia Trewin: this one cool trick
710                                // allows u64's to be substracted into a correct
711                                // i64.
712                                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            // TODO: this is reception of headers frame, before the field
726            // section is decoded. Ignore for now and just use the decoded event.
727            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                    // insert a 0'th point
776                    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            // TODO: add support for logging HTTP/2 sending
797            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                    // Wrapping subtraction before casting preserves the signed
840                    // difference between the two `u64` timestamps.
841                    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                    // insert a 0'th point
902                    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                    // insert a 0'th point
934                    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                // counterintuitively, reduces our local send window
971                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                    // insert a 0'th point
986                    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            // counterintuitively, updates our local send window
1010            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            // counterintuitively, updates our local receive window
1019            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                // Window updates always have a positive delta, so this is fine.
1031                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            // ignore the other events for now
1073            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            _ => (), // trace!("skipping {:?}", event.data),
1132        }
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                                        // TODO: check ack ranges and rust
1242                                        // Range mapping is correct
1243                                        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 there's no packet number we'll have to skip processing.
1321        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        // Prefer to use the packet_sent send_at_time if it exists. Otherwise
1341        // fallback to the event time.
1342        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                        // TODO
1349                    },
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        // Extract rate metrics from ex_data
1478        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            // give this a virtual y-value of the last cwnd value recorded, we
1501            // can choose to use it or not later.
1502            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            // ignore other frames
1556            _ => (),
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            // ignore other frames
1586            _ => (),
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                // TODO: is default length sensible here?
1609                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            // ignore other frames
1623            _ => (),
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            // ignore other frames
1664            _ => (),
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                    // Just swallow the failure and move on
1692                    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    // second line in a netlog is always `"events": [` so skip it
1754    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 this is a session creation, store the session, so we can
1779                // link events with it. The source ID of these events is the
1780                // unique value that we will use to link things together.
1781                // This assumes events belonging to a session do not occur
1782                // before the session is created.
1783                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                            // QUIC sessions split host and port, which
1790                            // interferes with filter expression, so merge them
1791                            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                            // Seems to be the earliest netlog event related to
1841                            // any request.
1842                            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                    // This will eventually deal with other events, and having
1859                    // to refactor back and forth is a waste.
1860                    #[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                                            // populate reverse mapping, each unique request ID has a session ID
1881                                            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                                            // todo replace vec with map?
1891                                            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                        // ignore other events
1906                        _ => (),
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                // Just swallow the failure and move on
1924            },
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}