Skip to main content

qlog_dancer/
wirefilter.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 qlog::events::quic::QuicFrame;
28use qlog::events::EventData;
29use qlog::reader::Event;
30use std::iter::FromIterator;
31use std::vec;
32use wirefilter::ExecutionContext;
33use wirefilter::Scheme;
34use wirefilter::TypedArray;
35
36use crate::category_and_type_from_event;
37
38fn stream_ids(event: &Event) -> TypedArray<'_, i64> {
39    let mut ids: TypedArray<i64> = TypedArray::new();
40
41    match event {
42        Event::Qlog(event) => match &event.data {
43            EventData::QuicStreamDataMoved(v) =>
44                if let Some(id) = v.stream_id {
45                    ids.push(id as i64);
46                },
47            EventData::QuicPacketSent(v) => {
48                if let Some(frames) = &v.frames {
49                    for frame in frames {
50                        match frame {
51                            QuicFrame::ResetStream { stream_id, .. } =>
52                                ids.push(*stream_id as i64),
53                            QuicFrame::StopSending { stream_id, .. } =>
54                                ids.push(*stream_id as i64),
55                            QuicFrame::Stream { stream_id, .. } =>
56                                ids.push(*stream_id as i64),
57                            QuicFrame::MaxStreamData { stream_id, .. } =>
58                                ids.push(*stream_id as i64),
59                            QuicFrame::StreamDataBlocked {
60                                stream_id, ..
61                            } => ids.push(*stream_id as i64),
62
63                            // other frames are not related to streams
64                            _ => (),
65                        }
66                    }
67                }
68            },
69            EventData::Http3StreamTypeSet(v) => {
70                ids.push(v.stream_id as i64);
71            },
72            EventData::Http3FrameCreated(v) => {
73                ids.push(v.stream_id as i64);
74            },
75            EventData::Http3FrameParsed(v) => {
76                ids.push(v.stream_id as i64);
77            },
78
79            // other events are not related to streams
80            _ => (),
81        },
82
83        // TODO: try and fuzzy extract stream id
84        Event::Json(_event) => {},
85    }
86
87    ids
88}
89
90pub fn filter_sqlog_events(mut events: Vec<Event>, filter: &str) -> Vec<Event> {
91    let mut ret = vec![];
92
93    let mut builder = Scheme! {
94        category: Bytes,
95        name: Bytes,
96        stream_id: Array(Int),
97    };
98
99    builder
100        .add_function("any", wirefilter::AnyFunction {})
101        .unwrap();
102
103    let scheme = builder.build();
104    let ast = scheme.parse(filter).unwrap();
105    let filter = ast.compile();
106
107    // TODO: smarter filtering rather then drain / recreate
108    for event in events.drain(..) {
109        // Recreate context each time to appease borrow checker
110        let mut ctx = ExecutionContext::new(&scheme);
111
112        let filter_match = match &event {
113            Event::Qlog(ev) => {
114                let (cat, ty) = category_and_type_from_event(&ev);
115
116                ctx.set_field_value(
117                    scheme.get_field("category").unwrap(),
118                    cat.clone(),
119                )
120                .unwrap();
121                ctx.set_field_value(
122                    scheme.get_field("name").unwrap(),
123                    ty.clone(),
124                )
125                .unwrap();
126
127                ctx.set_field_value(
128                    scheme.get_field("stream_id").unwrap(),
129                    stream_ids(&event),
130                )
131                .unwrap();
132
133                filter.execute(&ctx).unwrap()
134            },
135            Event::Json(ev) => {
136                let (cat, ty) = category_and_type_from_event(&ev);
137                ctx.set_field_value(
138                    scheme.get_field("category").unwrap(),
139                    cat.clone(),
140                )
141                .unwrap();
142                ctx.set_field_value(
143                    scheme.get_field("name").unwrap(),
144                    ty.clone(),
145                )
146                .unwrap();
147
148                ctx.set_field_value(
149                    scheme.get_field("stream_id").unwrap(),
150                    stream_ids(&event),
151                )
152                .unwrap();
153                filter.execute(&ctx).unwrap()
154            },
155        };
156
157        if filter_match {
158            ret.push(event);
159        }
160    }
161
162    ret
163}
164
165#[cfg(test)]
166mod tests {
167    use crate::wirefilter::filter_sqlog_events;
168    use qlog::events::quic::PacketHeader;
169    use qlog::events::quic::PacketSent;
170    use qlog::events::quic::PacketType::Initial;
171    use qlog::events::quic::QuicFrame;
172    use qlog::events::EventData::QuicPacketSent;
173    use qlog::events::RawInfo;
174    use qlog::reader::Event;
175
176    fn stream_frame(stream_id: u64) -> QuicFrame {
177        QuicFrame::Stream {
178            stream_id,
179            offset: Some(0),
180            fin: Some(true),
181            raw: Some(Box::new(RawInfo {
182                length: None,
183                payload_length: Some(10),
184                data: None,
185            })),
186        }
187    }
188
189    // Events are not cloneable in this qlog version, so use a helper.
190    fn events() -> Vec<Event> {
191        let mut events = vec![];
192        let scid = [0x7e, 0x37, 0xe4, 0xdc, 0xc6, 0x68, 0x2d, 0xa8];
193        let dcid = [0x36, 0xce, 0x10, 0x4e, 0xee, 0x50, 0x10, 0x1c];
194        let pkt_hdr = PacketHeader::new(
195            Initial,
196            Some(0),
197            None,
198            None,
199            Some(1),
200            Some(&scid),
201            Some(&dcid),
202        );
203        let raw = RawInfo {
204            length: None,
205            payload_length: Some(0),
206            data: None,
207        };
208
209        let frames = vec![
210            QuicFrame::Crypto {
211                offset: 0,
212                raw: Some(Box::new(raw)),
213            },
214            stream_frame(1),
215            stream_frame(2),
216            stream_frame(3),
217            stream_frame(4),
218            stream_frame(5),
219        ];
220
221        let raw = RawInfo {
222            length: Some(1251),
223            payload_length: Some(1224),
224            data: None,
225        };
226
227        let event_data = QuicPacketSent(PacketSent {
228            header: pkt_hdr.clone(),
229            frames: Some(frames),
230            stateless_reset_token: None,
231            supported_versions: None,
232            raw: Some(raw.clone()),
233            datagram_id: None,
234            is_mtu_probe_packet: None,
235            send_at_time: None,
236            trigger: None,
237        });
238
239        events.push(Event::Qlog(qlog::events::Event::with_time(0.0, event_data)));
240
241        let frames = vec![
242            stream_frame(0),
243            stream_frame(100),
244            stream_frame(200),
245            stream_frame(300),
246            stream_frame(400),
247        ];
248
249        let event_data = QuicPacketSent(PacketSent {
250            header: pkt_hdr.clone(),
251            frames: Some(frames),
252            stateless_reset_token: None,
253            supported_versions: None,
254            raw: Some(raw.clone()),
255            datagram_id: None,
256            is_mtu_probe_packet: None,
257            send_at_time: None,
258            trigger: None,
259        });
260
261        events.push(Event::Qlog(qlog::events::Event::with_time(0.0, event_data)));
262
263        let frames = vec![
264            stream_frame(1),
265            stream_frame(100),
266            stream_frame(2),
267            stream_frame(200),
268        ];
269
270        let event_data = QuicPacketSent(PacketSent {
271            header: pkt_hdr,
272            frames: Some(frames),
273            stateless_reset_token: None,
274            supported_versions: None,
275            raw: Some(raw),
276            datagram_id: None,
277            is_mtu_probe_packet: None,
278            send_at_time: None,
279            trigger: None,
280        });
281
282        events.push(Event::Qlog(qlog::events::Event::with_time(0.0, event_data)));
283
284        events
285    }
286
287    #[test]
288    fn test_stream_id_filter_no_match() {
289        let events = events();
290        assert_eq!(events.len(), 3);
291
292        let filter = "any(stream_id[*]==13)";
293        let filtered_events = filter_sqlog_events(events, filter);
294        assert!(filtered_events.is_empty());
295    }
296
297    #[test]
298    fn test_stream_id_filter_stream0() {
299        let events = events();
300        assert_eq!(events.len(), 3);
301
302        let filter = "any(stream_id[*]==0)";
303        let filtered_events = filter_sqlog_events(events, filter);
304        assert_eq!(filtered_events.len(), 1);
305
306        let ev = &filtered_events[0];
307        match ev {
308            Event::Qlog(event) => {
309                // assert_eq!
310                match &event.data {
311                    QuicPacketSent(packet_sent) => {
312                        assert_eq!(
313                            packet_sent.frames,
314                            Some(vec![
315                                stream_frame(0),
316                                stream_frame(100),
317                                stream_frame(200),
318                                stream_frame(300),
319                                stream_frame(400),
320                            ])
321                        );
322                    },
323                    _ => panic!("unexpected event data"),
324                }
325            },
326            Event::Json(_json_event) => panic!("unexpected type"),
327        }
328    }
329
330    #[test]
331    fn test_stream_id_filter_stream1() {
332        let events = events();
333        assert_eq!(events.len(), 3);
334
335        let filter = "any(stream_id[*]==1)";
336        let filtered_events = filter_sqlog_events(events, filter);
337        assert_eq!(filtered_events.len(), 2);
338
339        let ev = &filtered_events[0];
340
341        let raw = RawInfo {
342            length: None,
343            payload_length: Some(0),
344            data: None,
345        };
346
347        match ev {
348            Event::Qlog(event) => match &event.data {
349                QuicPacketSent(packet_sent) => {
350                    assert_eq!(
351                        packet_sent.frames,
352                        Some(vec![
353                            QuicFrame::Crypto {
354                                offset: 0,
355                                raw: Some(Box::new(raw)),
356                            },
357                            stream_frame(1),
358                            stream_frame(2),
359                            stream_frame(3),
360                            stream_frame(4),
361                            stream_frame(5),
362                        ])
363                    );
364                },
365                _ => panic!("unexpected event data"),
366            },
367            Event::Json(_json_event) => panic!("unexpected type"),
368        }
369
370        let ev = &filtered_events[1];
371        match ev {
372            Event::Qlog(event) => {
373                // assert_eq!
374                match &event.data {
375                    QuicPacketSent(packet_sent) => {
376                        assert_eq!(
377                            packet_sent.frames,
378                            Some(vec![
379                                stream_frame(1),
380                                stream_frame(100),
381                                stream_frame(2),
382                                stream_frame(200),
383                            ])
384                        );
385                    },
386                    _ => panic!("unexpected event data"),
387                }
388            },
389            Event::Json(_json_event) => panic!("unexpected type"),
390        }
391    }
392
393    #[test]
394    fn test_stream_id_filter_stream0_and_stream3() {
395        let events = events();
396        assert_eq!(events.len(), 3);
397
398        let filter = "any(stream_id[*]==0) || any(stream_id[*]==3)";
399        let filtered_events = filter_sqlog_events(events, filter);
400        assert_eq!(filtered_events.len(), 2);
401
402        let raw = RawInfo {
403            length: None,
404            payload_length: Some(0),
405            data: None,
406        };
407
408        let ev = &filtered_events[0];
409        match ev {
410            Event::Qlog(event) => {
411                // assert_eq!
412                match &event.data {
413                    QuicPacketSent(packet_sent) => {
414                        assert_eq!(
415                            packet_sent.frames,
416                            Some(vec![
417                                QuicFrame::Crypto {
418                                    offset: 0,
419                                    raw: Some(Box::new(raw)),
420                                },
421                                stream_frame(1),
422                                stream_frame(2),
423                                stream_frame(3),
424                                stream_frame(4),
425                                stream_frame(5),
426                            ])
427                        );
428                    },
429                    _ => panic!("unexpected event data"),
430                }
431            },
432            Event::Json(_json_event) => panic!("unexpected type"),
433        }
434
435        let ev = &filtered_events[1];
436        match ev {
437            Event::Qlog(event) => {
438                // assert_eq!
439                match &event.data {
440                    QuicPacketSent(packet_sent) => {
441                        assert_eq!(
442                            packet_sent.frames,
443                            Some(vec![
444                                stream_frame(0),
445                                stream_frame(100),
446                                stream_frame(200),
447                                stream_frame(300),
448                                stream_frame(400),
449                            ])
450                        );
451                    },
452                    _ => panic!("unexpected event data"),
453                }
454            },
455            Event::Json(_json_event) => panic!("unexpected type"),
456        }
457    }
458}