1use 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 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
128pub 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
145pub 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 let mut poll = mio::Poll::new().unwrap();
163 let mut events = mio::Events::with_capacity(1024);
164
165 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 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 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 let (write, send_info) = conn.send(&mut out).expect("initial send failed");
221
222 let mut client = SyncClient::new(close_trigger_frames);
223 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 let new = wait - timeout;
272 wait_duration = Some(new);
273 Some(timeout)
274 } else if wait < timeout {
275 Some(wait)
276 } else {
277 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 events.is_empty() {
292 log::debug!("timed out");
293
294 conn.on_timeout();
295 }
296
297 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 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 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 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 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 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 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 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
549pub 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 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}