Skip to main content

h3i/client/
sync_client.rs

1// Copyright (C) 2024, 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
27//! Responsible for creating a [quiche::Connection] and managing I/O.
28
29use std::slice::Iter;
30use std::time::Duration;
31use std::time::Instant;
32
33use ring::rand::*;
34
35use crate::client::QUIC_VERSION;
36use crate::frame::H3iFrame;
37use crate::quiche;
38
39use crate::actions::h3::Action;
40use crate::actions::h3::StreamEventType;
41use crate::actions::h3::WaitType;
42use crate::actions::h3::WaitingFor;
43use crate::client::execute_action;
44use crate::client::parse_streams;
45use crate::client::ClientError;
46use crate::client::ConnectionCloseDetails;
47use crate::client::MAX_DATAGRAM_SIZE;
48use crate::config::Config;
49
50use super::parse_args;
51use super::Client;
52use super::CloseTriggerFrames;
53use super::ConnectionSummary;
54use super::ParsedArgs;
55use super::StreamMap;
56use super::StreamParserMap;
57
58#[derive(Default)]
59struct SyncClient {
60    streams: StreamMap,
61    stream_parsers: StreamParserMap,
62}
63
64impl SyncClient {
65    fn new(close_trigger_frames: Option<CloseTriggerFrames>) -> Self {
66        Self {
67            streams: StreamMap::new(close_trigger_frames),
68            ..Default::default()
69        }
70    }
71}
72
73impl Client for SyncClient {
74    fn stream_parsers_mut(&mut self) -> &mut StreamParserMap {
75        &mut self.stream_parsers
76    }
77
78    fn handle_response_frame(&mut self, stream_id: u64, frame: H3iFrame) {
79        self.streams.insert(stream_id, frame);
80    }
81}
82
83fn create_config(args: &Config, should_log_keys: bool) -> quiche::Config {
84    // Create the configuration for the QUIC connection.
85    let mut config = quiche::Config::new(QUIC_VERSION).unwrap();
86
87    config.verify_peer(args.verify_peer);
88    config.set_application_protos(&[b"h3"]).unwrap();
89    config.set_max_idle_timeout(args.idle_timeout);
90    config.set_send_capacity_factor(args.send_capacity_factor);
91    config.set_max_recv_udp_payload_size(MAX_DATAGRAM_SIZE);
92    config.set_max_send_udp_payload_size(MAX_DATAGRAM_SIZE);
93    config.set_initial_max_data(10_000_000);
94    config
95        .set_initial_max_stream_data_bidi_local(args.max_stream_data_bidi_local);
96    config.set_initial_max_stream_data_bidi_remote(
97        args.max_stream_data_bidi_remote,
98    );
99    config.set_initial_max_stream_data_uni(args.max_stream_data_uni);
100    config.set_initial_max_streams_bidi(args.max_streams_bidi);
101    config.set_initial_max_streams_uni(args.max_streams_uni);
102    config.set_disable_active_migration(true);
103    config.set_active_connection_id_limit(0);
104
105    config.set_max_connection_window(args.max_window);
106    config.set_max_stream_window(args.max_stream_window);
107    config.set_enable_send_streams_blocked(true);
108    config.grease(false);
109
110    if args.enable_early_data {
111        config.enable_early_data();
112    }
113
114    if args.enable_dgram {
115        config.enable_dgram(
116            true,
117            args.dgram_recv_queue_len,
118            args.dgram_send_queue_len,
119        );
120    }
121    if should_log_keys {
122        config.log_keys()
123    }
124
125    config
126}
127
128/// Connect to a server and execute provided actions.
129///
130/// Constructs a socket and [quiche::Connection] based on the provided `args`,
131/// then iterates over `actions`.
132///
133/// If `close_trigger_frames` is specified, h3i will close the connection
134/// immediately upon receiving all of the supplied frames rather than waiting
135/// for the idle timeout. See [`CloseTriggerFrames`] for details.
136///
137/// Returns a [ConnectionSummary] on success, [ClientError] on failure.
138pub fn connect(
139    args: Config, actions: Vec<Action>,
140    close_trigger_frames: Option<CloseTriggerFrames>,
141) -> std::result::Result<ConnectionSummary, ClientError> {
142    connect_with_early_data(args, None, actions, close_trigger_frames)
143}
144
145/// Connect to a server and execute provided early_action and actions.
146///
147/// See `connect` for additional documentation.
148pub fn connect_with_early_data(
149    args: Config, early_actions: Option<Vec<Action>>, actions: Vec<Action>,
150    close_trigger_frames: Option<CloseTriggerFrames>,
151) -> std::result::Result<ConnectionSummary, ClientError> {
152    let mut buf = [0; 65535];
153    let mut out = [0; MAX_DATAGRAM_SIZE];
154
155    let ParsedArgs {
156        connect_url,
157        bind_addr,
158        peer_addr,
159    } = parse_args(&args);
160
161    // Setup the event loop.
162    let mut poll = mio::Poll::new().unwrap();
163    let mut events = mio::Events::with_capacity(1024);
164
165    // Create the UDP socket backing the QUIC connection, and register it with
166    // the event loop.
167    let mut socket = mio::net::UdpSocket::bind(bind_addr).unwrap();
168    poll.registry()
169        .register(&mut socket, mio::Token(0), mio::Interest::READABLE)
170        .unwrap();
171
172    let mut keylog = None;
173    if let Some(keylog_path) = std::env::var_os("SSLKEYLOGFILE") {
174        let file = std::fs::OpenOptions::new()
175            .create(true)
176            .append(true)
177            .open(keylog_path)
178            .unwrap();
179
180        keylog = Some(file);
181    }
182
183    let mut config = create_config(&args, keylog.is_some());
184
185    // Generate a random source connection ID for the connection.
186    let mut scid = [0; quiche::MAX_CONN_ID_LEN];
187
188    let rng = SystemRandom::new();
189    rng.fill(&mut scid[..]).unwrap();
190
191    let scid = quiche::ConnectionId::from_ref(&scid);
192
193    let Ok(local_addr) = socket.local_addr() else {
194        return Err(ClientError::Other("invalid socket".to_string()));
195    };
196
197    // Create a new client-side QUIC connection.
198    let mut conn =
199        quiche::connect(connect_url, &scid, local_addr, peer_addr, &mut config)
200            .map_err(|e| ClientError::Other(e.to_string()))?;
201
202    if let Some(session) = &args.session {
203        conn.set_session(session)
204            .map_err(|error| ClientError::Other(error.to_string()))?;
205    }
206
207    if let Some(keylog) = &mut keylog {
208        if let Ok(keylog) = keylog.try_clone() {
209            conn.set_keylog(Box::new(keylog));
210        }
211    }
212
213    log::info!(
214        "connecting to {peer_addr:} from {local_addr:} with scid {scid:?}",
215    );
216
217    let mut app_proto_selected = false;
218
219    // Send ClientHello and initiate the handshake.
220    let (write, send_info) = conn.send(&mut out).expect("initial send failed");
221
222    let mut client = SyncClient::new(close_trigger_frames);
223    // Send early data if connection is_in_early_data (resumption with 0-RTT was
224    // successful) and if we have early_actions.
225    if conn.is_in_early_data() {
226        if let Some(early_actions) = early_actions {
227            let mut early_action_iter = early_actions.iter();
228            let mut wait_duration = None;
229            let mut wait_instant = None;
230            let mut waiting_for = WaitingFor::default();
231
232            check_duration_and_do_actions(
233                &mut wait_duration,
234                &mut wait_instant,
235                &mut early_action_iter,
236                &mut conn,
237                &mut waiting_for,
238                client.stream_parsers_mut(),
239            );
240        }
241    }
242
243    while let Err(e) = socket.send_to(&out[..write], send_info.to) {
244        if e.kind() == std::io::ErrorKind::WouldBlock {
245            log::debug!(
246                "{} -> {}: send() would block",
247                socket.local_addr().unwrap(),
248                send_info.to
249            );
250            continue;
251        }
252
253        return Err(ClientError::Other(format!("send() failed: {e:?}")));
254    }
255
256    let app_data_start = std::time::Instant::now();
257
258    let mut action_iter = actions.iter();
259    let mut wait_duration = None;
260    let mut wait_instant = None;
261
262    let mut waiting_for = WaitingFor::default();
263
264    loop {
265        let actual_sleep = match (wait_duration, conn.timeout()) {
266            (Some(wait), Some(timeout)) => {
267                #[allow(clippy::comparison_chain)]
268                if timeout < wait {
269                    // shave some off the wait time so it doesn't go longer
270                    // than user really wanted.
271                    let new = wait - timeout;
272                    wait_duration = Some(new);
273                    Some(timeout)
274                } else if wait < timeout {
275                    Some(wait)
276                } else {
277                    // same, so picking either doesn't matter
278                    Some(timeout)
279                }
280            },
281            (None, Some(timeout)) => Some(timeout),
282            (Some(wait), None) => Some(wait),
283            _ => None,
284        };
285
286        log::debug!("actual sleep is {actual_sleep:?}");
287        poll.poll(&mut events, actual_sleep).unwrap();
288
289        // If the event loop reported no events, run a belt and braces check on
290        // the quiche connection's timeouts.
291        if events.is_empty() {
292            log::debug!("timed out");
293
294            conn.on_timeout();
295        }
296
297        // Read incoming UDP packets from the socket and feed them to quiche,
298        // until there are no more packets to read.
299        for event in &events {
300            let socket = match event.token() {
301                mio::Token(0) => &socket,
302
303                _ => unreachable!(),
304            };
305
306            let local_addr = socket.local_addr().unwrap();
307            'read: loop {
308                let (len, from) = match socket.recv_from(&mut buf) {
309                    Ok(v) => v,
310
311                    Err(e) => {
312                        // There are no more UDP packets to read on this socket.
313                        // Process subsequent events.
314                        if e.kind() == std::io::ErrorKind::WouldBlock {
315                            break 'read;
316                        }
317
318                        return Err(ClientError::Other(format!(
319                            "{local_addr}: recv() failed: {e:?}"
320                        )));
321                    },
322                };
323
324                let recv_info = quiche::RecvInfo {
325                    to: local_addr,
326                    from,
327                };
328
329                // Process potentially coalesced packets.
330                let _read = match conn.recv(&mut buf[..len], recv_info) {
331                    Ok(v) => v,
332
333                    Err(e) => {
334                        log::debug!("{local_addr}: recv failed: {e:?}");
335                        continue 'read;
336                    },
337                };
338            }
339        }
340
341        log::debug!("done reading");
342
343        if conn.is_closed() {
344            log::info!(
345                "connection closed with error={:?} did_idle_timeout={}, stats={:?} path_stats={:?}",
346                conn.peer_error(),
347                conn.is_timed_out(),
348                conn.stats(),
349                conn.path_stats().collect::<Vec<quiche::PathStats>>(),
350            );
351
352            if !conn.is_established() {
353                log::info!(
354                    "connection timed out after {:?}",
355                    app_data_start.elapsed(),
356                );
357
358                return Err(ClientError::HandshakeFail);
359            }
360
361            break;
362        }
363
364        // Create a new application protocol session once the QUIC connection is
365        // established.
366        if (conn.is_established() || conn.is_in_early_data()) &&
367            !app_proto_selected
368        {
369            app_proto_selected = true;
370        }
371
372        if app_proto_selected {
373            check_duration_and_do_actions(
374                &mut wait_duration,
375                &mut wait_instant,
376                &mut action_iter,
377                &mut conn,
378                &mut waiting_for,
379                client.stream_parsers_mut(),
380            );
381
382            let mut wait_cleared = false;
383            for response in parse_streams(&mut conn, &mut client) {
384                let stream_id = response.stream_id;
385
386                if let StreamEventType::Finished = response.event_type {
387                    waiting_for.clear_waits_on_stream(stream_id);
388                } else {
389                    waiting_for.remove_wait(response);
390                }
391
392                wait_cleared = true;
393            }
394
395            // Check if a CanOpenNumStreams wait is satisfied.
396            let before = waiting_for.is_empty();
397            waiting_for.check_can_open_num_streams(&conn);
398            if !before && waiting_for.is_empty() {
399                wait_cleared = true;
400            }
401
402            if client.streams.all_close_trigger_frames_seen() {
403                client.streams.close_due_to_trigger_frames(&mut conn);
404            }
405
406            if wait_cleared {
407                check_duration_and_do_actions(
408                    &mut wait_duration,
409                    &mut wait_instant,
410                    &mut action_iter,
411                    &mut conn,
412                    &mut waiting_for,
413                    client.stream_parsers_mut(),
414                );
415            }
416        }
417
418        // Provides as many CIDs as possible.
419        while conn.scids_left() > 0 {
420            let (scid, reset_token) = generate_cid_and_reset_token();
421
422            if conn.new_scid(&scid, reset_token, false).is_err() {
423                break;
424            }
425        }
426
427        // Generate outgoing QUIC packets and send them on the UDP socket, until
428        // quiche reports that there are no more packets to be sent.
429        let sockets = vec![&socket];
430
431        for socket in sockets {
432            let local_addr = socket.local_addr().unwrap();
433
434            for peer_addr in conn.paths_iter(local_addr) {
435                loop {
436                    let (write, send_info) = match conn.send_on_path(
437                        &mut out,
438                        Some(local_addr),
439                        Some(peer_addr),
440                    ) {
441                        Ok(v) => v,
442
443                        Err(quiche::Error::Done) => {
444                            break;
445                        },
446
447                        Err(e) => {
448                            log::error!(
449                                "{local_addr} -> {peer_addr}: send failed: {e:?}"
450                            );
451
452                            conn.close(false, 0x1, b"fail").ok();
453                            break;
454                        },
455                    };
456
457                    if let Err(e) = socket.send_to(&out[..write], send_info.to) {
458                        if e.kind() == std::io::ErrorKind::WouldBlock {
459                            log::debug!(
460                                "{} -> {}: send() would block",
461                                local_addr,
462                                send_info.to
463                            );
464                            break;
465                        }
466
467                        return Err(ClientError::Other(format!(
468                            "{} -> {}: send() failed: {:?}",
469                            local_addr, send_info.to, e
470                        )));
471                    }
472                }
473            }
474        }
475
476        if conn.is_closed() {
477            log::info!(
478                "connection closed, {:?} {:?}",
479                conn.stats(),
480                conn.path_stats().collect::<Vec<quiche::PathStats>>()
481            );
482
483            if !conn.is_established() {
484                log::info!(
485                    "connection timed out after {:?}",
486                    app_data_start.elapsed(),
487                );
488
489                return Err(ClientError::HandshakeFail);
490            }
491
492            break;
493        }
494    }
495
496    Ok(ConnectionSummary {
497        stream_map: client.streams,
498        stats: Some(conn.stats()),
499        path_stats: conn.path_stats().collect(),
500        conn_close_details: ConnectionCloseDetails::new(&conn),
501    })
502}
503
504fn check_duration_and_do_actions(
505    wait_duration: &mut Option<Duration>, wait_instant: &mut Option<Instant>,
506    action_iter: &mut Iter<Action>, conn: &mut quiche::Connection,
507    waiting_for: &mut WaitingFor, stream_parsers: &mut StreamParserMap,
508) {
509    match wait_duration.as_ref() {
510        None => {
511            if let Some(idle_wait) =
512                handle_actions(action_iter, conn, waiting_for, stream_parsers)
513            {
514                *wait_duration = Some(idle_wait);
515                *wait_instant = Some(Instant::now());
516
517                // TODO: the wait period could still be larger than the
518                // negotiated idle timeout.
519                // We could in theory check quiche's idle_timeout value if
520                // it was public.
521                log::info!(
522                    "waiting for {idle_wait:?} before executing more actions"
523                );
524            }
525        },
526
527        Some(period) => {
528            let now = Instant::now();
529            let then = wait_instant.unwrap();
530            log::debug!(
531                "checking if actions wait period elapsed {:?} > {:?}",
532                now.duration_since(then),
533                wait_duration
534            );
535            if now.duration_since(then) >= *period {
536                log::debug!("yup!");
537                *wait_duration = None;
538
539                if let Some(idle_wait) =
540                    handle_actions(action_iter, conn, waiting_for, stream_parsers)
541                {
542                    *wait_duration = Some(idle_wait);
543                }
544            }
545        },
546    }
547}
548
549/// Generate a new pair of Source Connection ID and reset token.
550pub fn generate_cid_and_reset_token() -> (quiche::ConnectionId<'static>, u128) {
551    let rng = SystemRandom::new();
552
553    let mut scid = [0; quiche::MAX_CONN_ID_LEN];
554    rng.fill(&mut scid[..]).unwrap();
555    let scid = scid.to_vec().into();
556
557    let mut reset_token = [0; 16];
558    rng.fill(&mut reset_token[..]).unwrap();
559
560    let reset_token = u128::from_be_bytes(reset_token);
561    (scid, reset_token)
562}
563
564fn handle_actions<'a, I>(
565    iter: &mut I, conn: &mut quiche::Connection, waiting_for: &mut WaitingFor,
566    stream_parsers: &mut StreamParserMap,
567) -> Option<Duration>
568where
569    I: Iterator<Item = &'a Action>,
570{
571    if !waiting_for.is_empty() {
572        log::debug!(
573            "won't fire an action due to waiting for responses: {waiting_for:?}"
574        );
575        return None;
576    }
577
578    // Send actions
579    for action in iter {
580        match action {
581            Action::FlushPackets => return None,
582            Action::Wait { wait_type } => match wait_type {
583                WaitType::WaitDuration(period) => return Some(*period),
584                WaitType::StreamEvent(response) => {
585                    log::info!(
586                        "waiting for {response:?} before executing more actions"
587                    );
588                    waiting_for.add_wait(response);
589                    return None;
590                },
591                WaitType::CanOpenNumStreams(required_streams) => {
592                    log::info!(
593                        "h3i: waiting for peer_streams_left_bidi >= {required_streams:?}"
594                    );
595                    waiting_for.set_required_stream_quota(*required_streams);
596                    return None;
597                },
598            },
599            action => execute_action(action, conn, stream_parsers),
600        }
601    }
602
603    None
604}