1use std::collections::hash_map;
286use std::collections::HashSet;
287use std::collections::VecDeque;
288
289#[cfg(feature = "sfv")]
290use std::convert::TryFrom;
291use std::fmt;
292use std::fmt::Write;
293
294#[cfg(feature = "qlog")]
295use qlog::events::http3::FrameCreated;
296#[cfg(feature = "qlog")]
297use qlog::events::http3::FrameParsed;
298#[cfg(feature = "qlog")]
299use qlog::events::http3::Http3EventType;
300#[cfg(feature = "qlog")]
301use qlog::events::http3::Http3Frame;
302#[cfg(feature = "qlog")]
303use qlog::events::http3::Initiator;
304#[cfg(feature = "qlog")]
305use qlog::events::http3::StreamType;
306#[cfg(feature = "qlog")]
307use qlog::events::http3::StreamTypeSet;
308#[cfg(feature = "qlog")]
309use qlog::events::EventData;
310#[cfg(feature = "qlog")]
311use qlog::events::EventImportance;
312#[cfg(feature = "qlog")]
313use qlog::events::EventType;
314
315use crate::buffers::BufFactory;
316use crate::BufSplit;
317
318pub const APPLICATION_PROTOCOL: &[&[u8]] = &[b"h3"];
326
327const PRIORITY_URGENCY_OFFSET: u8 = 124;
329
330const PRIORITY_URGENCY_LOWER_BOUND: u8 = 0;
334const PRIORITY_URGENCY_UPPER_BOUND: u8 = 7;
335const PRIORITY_URGENCY_DEFAULT: u8 = 3;
336const PRIORITY_INCREMENTAL_DEFAULT: bool = false;
337
338pub const PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT: u64 = 256;
343
344pub const SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT: u64 = 32_768;
348
349#[cfg(feature = "qlog")]
350const QLOG_FRAME_CREATED: EventType =
351 EventType::Http3EventType(Http3EventType::FrameCreated);
352#[cfg(feature = "qlog")]
353const QLOG_FRAME_PARSED: EventType =
354 EventType::Http3EventType(Http3EventType::FrameParsed);
355#[cfg(feature = "qlog")]
356const QLOG_STREAM_TYPE_SET: EventType =
357 EventType::Http3EventType(Http3EventType::StreamTypeSet);
358
359pub type Result<T> = std::result::Result<T, Error>;
366
367#[derive(Clone, Copy, Debug, PartialEq, Eq)]
369pub enum Error {
370 Done,
372
373 BufferTooShort,
375
376 InternalError,
378
379 ExcessiveLoad,
382
383 IdError,
387
388 StreamCreationError,
391
392 ClosedCriticalStream,
394
395 MissingSettings,
397
398 FrameUnexpected,
400
401 FrameError,
403
404 QpackDecompressionFailed,
406
407 TransportError(crate::Error),
409
410 StreamBlocked,
413
414 SettingsError,
416
417 RequestRejected,
419
420 RequestCancelled,
422
423 RequestIncomplete,
426
427 MessageError,
429
430 ConnectError,
433
434 VersionFallback,
437}
438
439#[derive(Copy, Clone, Debug, Eq, PartialEq)]
443pub enum WireErrorCode {
444 NoError = 0x100,
447 GeneralProtocolError = 0x101,
451 InternalError = 0x102,
453 StreamCreationError = 0x103,
456 ClosedCriticalStream = 0x104,
458 FrameUnexpected = 0x105,
461 FrameError = 0x106,
464 ExcessiveLoad = 0x107,
467 IdError = 0x108,
470 SettingsError = 0x109,
472 MissingSettings = 0x10a,
474 RequestRejected = 0x10b,
477 RequestCancelled = 0x10c,
479 RequestIncomplete = 0x10d,
482 MessageError = 0x10e,
484 ConnectError = 0x10f,
487 VersionFallback = 0x110,
490}
491
492impl Error {
493 fn to_wire(self) -> u64 {
494 match self {
495 Error::Done => WireErrorCode::NoError as u64,
496 Error::InternalError => WireErrorCode::InternalError as u64,
497 Error::StreamCreationError =>
498 WireErrorCode::StreamCreationError as u64,
499 Error::ClosedCriticalStream =>
500 WireErrorCode::ClosedCriticalStream as u64,
501 Error::FrameUnexpected => WireErrorCode::FrameUnexpected as u64,
502 Error::FrameError => WireErrorCode::FrameError as u64,
503 Error::ExcessiveLoad => WireErrorCode::ExcessiveLoad as u64,
504 Error::IdError => WireErrorCode::IdError as u64,
505 Error::MissingSettings => WireErrorCode::MissingSettings as u64,
506 Error::QpackDecompressionFailed => 0x200,
507 Error::BufferTooShort => 0x999,
508 Error::TransportError { .. } | Error::StreamBlocked => 0xFF,
509 Error::SettingsError => WireErrorCode::SettingsError as u64,
510 Error::RequestRejected => WireErrorCode::RequestRejected as u64,
511 Error::RequestCancelled => WireErrorCode::RequestCancelled as u64,
512 Error::RequestIncomplete => WireErrorCode::RequestIncomplete as u64,
513 Error::MessageError => WireErrorCode::MessageError as u64,
514 Error::ConnectError => WireErrorCode::ConnectError as u64,
515 Error::VersionFallback => WireErrorCode::VersionFallback as u64,
516 }
517 }
518
519 #[cfg(feature = "ffi")]
520 fn to_c(self) -> libc::ssize_t {
521 match self {
522 Error::Done => -1,
523 Error::BufferTooShort => -2,
524 Error::InternalError => -3,
525 Error::ExcessiveLoad => -4,
526 Error::IdError => -5,
527 Error::StreamCreationError => -6,
528 Error::ClosedCriticalStream => -7,
529 Error::MissingSettings => -8,
530 Error::FrameUnexpected => -9,
531 Error::FrameError => -10,
532 Error::QpackDecompressionFailed => -11,
533 Error::StreamBlocked => -13,
535 Error::SettingsError => -14,
536 Error::RequestRejected => -15,
537 Error::RequestCancelled => -16,
538 Error::RequestIncomplete => -17,
539 Error::MessageError => -18,
540 Error::ConnectError => -19,
541 Error::VersionFallback => -20,
542
543 Error::TransportError(quic_error) => quic_error.to_c() - 1000,
544 }
545 }
546}
547
548impl fmt::Display for Error {
549 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
550 write!(f, "{self:?}")
551 }
552}
553
554impl std::error::Error for Error {
555 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
556 None
557 }
558}
559
560impl From<super::Error> for Error {
561 fn from(err: super::Error) -> Self {
562 match err {
563 super::Error::Done => Error::Done,
564
565 _ => Error::TransportError(err),
566 }
567 }
568}
569
570impl From<octets::BufferTooShortError> for Error {
571 fn from(_err: octets::BufferTooShortError) -> Self {
572 Error::BufferTooShort
573 }
574}
575
576pub struct Config {
578 max_field_section_size: Option<u64>,
579 qpack_max_table_capacity: Option<u64>,
580 qpack_blocked_streams: Option<u64>,
581 connect_protocol_enabled: Option<u64>,
582 additional_settings: Option<Vec<(u64, u64)>>,
585
586 max_priority_update_size: u64,
587}
588
589impl Config {
590 pub const fn new() -> Result<Config> {
592 Ok(Config {
593 max_field_section_size: Some(SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT),
594 qpack_max_table_capacity: None,
595 qpack_blocked_streams: None,
596 connect_protocol_enabled: None,
597 additional_settings: None,
598 max_priority_update_size:
599 PRIORITY_UPDATE_FRAME_PAYLOAD_MAX_SIZE_DEFAULT,
600 })
601 }
602
603 pub fn set_max_field_section_size(&mut self, v: u64) {
620 self.max_field_section_size = Some(v);
621 }
622
623 pub fn set_qpack_max_table_capacity(&mut self, v: u64) {
627 self.qpack_max_table_capacity = Some(v);
628 }
629
630 pub fn set_qpack_blocked_streams(&mut self, v: u64) {
634 self.qpack_blocked_streams = Some(v);
635 }
636
637 pub fn enable_extended_connect(&mut self, enabled: bool) {
641 if enabled {
642 self.connect_protocol_enabled = Some(1);
643 } else {
644 self.connect_protocol_enabled = None;
645 }
646 }
647
648 pub fn set_additional_settings(
668 &mut self, additional_settings: Vec<(u64, u64)>,
669 ) -> Result<()> {
670 let explicit_quiche_settings = HashSet::from([
671 frame::SETTINGS_QPACK_MAX_TABLE_CAPACITY,
672 frame::SETTINGS_MAX_FIELD_SECTION_SIZE,
673 frame::SETTINGS_QPACK_BLOCKED_STREAMS,
674 frame::SETTINGS_ENABLE_CONNECT_PROTOCOL,
675 frame::SETTINGS_H3_DATAGRAM,
676 frame::SETTINGS_H3_DATAGRAM_00,
677 ]);
678
679 let dedup_settings: HashSet<u64> =
680 additional_settings.iter().map(|(key, _)| *key).collect();
681
682 if dedup_settings.len() != additional_settings.len() ||
683 !explicit_quiche_settings.is_disjoint(&dedup_settings)
684 {
685 return Err(Error::SettingsError);
686 }
687 self.additional_settings = Some(additional_settings);
688 Ok(())
689 }
690
691 pub fn set_max_priority_update_size(&mut self, v: u64) {
703 self.max_priority_update_size = v;
704 }
705}
706
707pub trait NameValue {
709 fn name(&self) -> &[u8];
711
712 fn value(&self) -> &[u8];
714}
715
716impl<N, V> NameValue for (N, V)
717where
718 N: AsRef<[u8]>,
719 V: AsRef<[u8]>,
720{
721 fn name(&self) -> &[u8] {
722 self.0.as_ref()
723 }
724
725 fn value(&self) -> &[u8] {
726 self.1.as_ref()
727 }
728}
729
730#[derive(Clone, PartialEq, Eq)]
732pub struct Header(Vec<u8>, Vec<u8>);
733
734fn try_print_as_readable(hdr: &[u8], f: &mut fmt::Formatter) -> fmt::Result {
735 match std::str::from_utf8(hdr) {
736 Ok(s) => f.write_str(&s.escape_default().to_string()),
737 Err(_) => write!(f, "{hdr:?}"),
738 }
739}
740
741impl fmt::Debug for Header {
742 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
743 f.write_char('"')?;
744 try_print_as_readable(&self.0, f)?;
745 f.write_str(": ")?;
746 try_print_as_readable(&self.1, f)?;
747 f.write_char('"')
748 }
749}
750
751impl Header {
752 pub fn new(name: &[u8], value: &[u8]) -> Self {
756 Self(name.to_vec(), value.to_vec())
757 }
758}
759
760impl NameValue for Header {
761 fn name(&self) -> &[u8] {
762 &self.0
763 }
764
765 fn value(&self) -> &[u8] {
766 &self.1
767 }
768}
769
770#[derive(Clone, Debug, PartialEq, Eq)]
772pub struct HeaderRef<'a>(&'a [u8], &'a [u8]);
773
774impl<'a> HeaderRef<'a> {
775 pub const fn new(name: &'a [u8], value: &'a [u8]) -> Self {
777 Self(name, value)
778 }
779}
780
781impl NameValue for HeaderRef<'_> {
782 fn name(&self) -> &[u8] {
783 self.0
784 }
785
786 fn value(&self) -> &[u8] {
787 self.1
788 }
789}
790
791#[derive(Clone, Debug, PartialEq, Eq)]
793pub enum Event {
794 Headers {
796 list: Vec<Header>,
799
800 more_frames: bool,
802 },
803
804 Data,
816
817 Finished,
819
820 Reset(u64),
824
825 PriorityUpdate,
839
840 GoAway,
842}
843
844#[derive(Clone, Copy, Debug, PartialEq, Eq)]
851#[repr(C)]
852pub struct Priority {
853 urgency: u8,
854 incremental: bool,
855}
856
857impl Default for Priority {
858 fn default() -> Self {
859 Priority {
860 urgency: PRIORITY_URGENCY_DEFAULT,
861 incremental: PRIORITY_INCREMENTAL_DEFAULT,
862 }
863 }
864}
865
866impl Priority {
867 pub const fn new(urgency: u8, incremental: bool) -> Self {
869 Priority {
870 urgency,
871 incremental,
872 }
873 }
874}
875
876#[cfg(feature = "sfv")]
877#[cfg_attr(docsrs, doc(cfg(feature = "sfv")))]
878impl TryFrom<&[u8]> for Priority {
879 type Error = Error;
880
881 fn try_from(value: &[u8]) -> std::result::Result<Self, Self::Error> {
897 let dict = match sfv::Parser::parse_dictionary(value) {
898 Ok(v) => v,
899
900 Err(_) => return Err(Error::Done),
901 };
902
903 let urgency = match dict.get("u") {
904 Some(sfv::ListEntry::Item(item)) => match item.bare_item.as_int() {
910 Some(v) => {
911 if !(PRIORITY_URGENCY_LOWER_BOUND as i64..=
912 PRIORITY_URGENCY_UPPER_BOUND as i64)
913 .contains(&v)
914 {
915 PRIORITY_URGENCY_UPPER_BOUND
916 } else {
917 v as u8
918 }
919 },
920
921 None => return Err(Error::Done),
922 },
923
924 Some(sfv::ListEntry::InnerList(_)) => return Err(Error::Done),
925
926 None => PRIORITY_URGENCY_DEFAULT,
928 };
929
930 let incremental = match dict.get("i") {
931 Some(sfv::ListEntry::Item(item)) =>
932 item.bare_item.as_bool().ok_or(Error::Done)?,
933
934 _ => false,
936 };
937
938 Ok(Priority::new(urgency, incremental))
939 }
940}
941
942struct ConnectionSettings {
943 pub max_field_section_size: Option<u64>,
944 pub qpack_max_table_capacity: Option<u64>,
945 pub qpack_blocked_streams: Option<u64>,
946 pub connect_protocol_enabled: Option<u64>,
947 pub h3_datagram: Option<u64>,
948 pub additional_settings: Option<Vec<(u64, u64)>>,
949 pub raw: Option<Vec<(u64, u64)>>,
950}
951
952#[derive(Default)]
953struct QpackStreams {
954 pub encoder_stream_id: Option<u64>,
955 pub encoder_stream_bytes: u64,
956 pub decoder_stream_id: Option<u64>,
957 pub decoder_stream_bytes: u64,
958}
959
960#[derive(Clone, Default)]
966#[non_exhaustive]
967pub struct Stats {
968 pub qpack_encoder_stream_recv_bytes: u64,
970 pub qpack_decoder_stream_recv_bytes: u64,
972}
973
974fn close_conn_critical_stream<F: BufFactory>(
975 conn: &mut super::Connection<F>,
976) -> Result<()> {
977 conn.close(
978 true,
979 Error::ClosedCriticalStream.to_wire(),
980 b"Critical stream closed.",
981 )?;
982
983 Err(Error::ClosedCriticalStream)
984}
985
986fn close_conn_if_critical_stream_finished<F: BufFactory>(
987 conn: &mut super::Connection<F>, stream_id: u64,
988) -> Result<()> {
989 if conn.stream_finished(stream_id) {
990 close_conn_critical_stream(conn)?;
991 }
992
993 Ok(())
994}
995
996pub struct Connection {
998 is_server: bool,
999
1000 next_request_stream_id: u64,
1001 next_uni_stream_id: u64,
1002
1003 streams: crate::stream::StreamIdHashMap<stream::Stream>,
1004
1005 local_settings: ConnectionSettings,
1006 peer_settings: ConnectionSettings,
1007
1008 control_stream_id: Option<u64>,
1009 peer_control_stream_id: Option<u64>,
1010
1011 qpack_encoder: qpack::Encoder,
1012 qpack_decoder: qpack::Decoder,
1013
1014 local_qpack_streams: QpackStreams,
1015 peer_qpack_streams: QpackStreams,
1016
1017 max_push_id: u64,
1018
1019 finished_streams: VecDeque<u64>,
1024
1025 frames_greased: bool,
1026
1027 local_goaway_id: Option<u64>,
1028 peer_goaway_id: Option<u64>,
1029
1030 max_priority_update_size: u64,
1031}
1032
1033impl Connection {
1034 fn new(
1035 config: &Config, is_server: bool, enable_dgram: bool,
1036 ) -> Result<Connection> {
1037 let initial_uni_stream_id = if is_server { 0x3 } else { 0x2 };
1038 let h3_datagram = if enable_dgram { Some(1) } else { None };
1039
1040 Ok(Connection {
1041 is_server,
1042
1043 next_request_stream_id: 0,
1044
1045 next_uni_stream_id: initial_uni_stream_id,
1046
1047 streams: Default::default(),
1048
1049 local_settings: ConnectionSettings {
1050 max_field_section_size: config.max_field_section_size,
1051 qpack_max_table_capacity: config.qpack_max_table_capacity,
1052 qpack_blocked_streams: config.qpack_blocked_streams,
1053 connect_protocol_enabled: config.connect_protocol_enabled,
1054 h3_datagram,
1055 additional_settings: config.additional_settings.clone(),
1056 raw: Default::default(),
1057 },
1058
1059 peer_settings: ConnectionSettings {
1060 max_field_section_size: None,
1061 qpack_max_table_capacity: None,
1062 qpack_blocked_streams: None,
1063 h3_datagram: None,
1064 connect_protocol_enabled: None,
1065 additional_settings: Default::default(),
1066 raw: Default::default(),
1067 },
1068
1069 control_stream_id: None,
1070 peer_control_stream_id: None,
1071
1072 qpack_encoder: qpack::Encoder::new(),
1073 qpack_decoder: qpack::Decoder::new(),
1074
1075 local_qpack_streams: Default::default(),
1076 peer_qpack_streams: Default::default(),
1077
1078 max_push_id: 0,
1079
1080 finished_streams: VecDeque::new(),
1081
1082 frames_greased: false,
1083
1084 local_goaway_id: None,
1085 peer_goaway_id: None,
1086
1087 max_priority_update_size: config.max_priority_update_size,
1088 })
1089 }
1090
1091 pub fn with_transport<F: BufFactory>(
1108 conn: &mut super::Connection<F>, config: &Config,
1109 ) -> Result<Connection> {
1110 let is_client = !conn.is_server;
1111 if is_client && !(conn.is_established() || conn.is_in_early_data()) {
1112 trace!("{} QUIC connection must be established or in early data before creating an HTTP/3 connection", conn.trace_id());
1113 return Err(Error::InternalError);
1114 }
1115
1116 let mut http3_conn =
1117 Connection::new(config, conn.is_server, conn.dgram_enabled())?;
1118
1119 match http3_conn.send_settings(conn) {
1120 Ok(_) => (),
1121
1122 Err(e) => {
1123 conn.close(true, e.to_wire(), b"Error opening control stream")?;
1124 return Err(e);
1125 },
1126 };
1127
1128 http3_conn.open_qpack_encoder_stream(conn).ok();
1131 http3_conn.open_qpack_decoder_stream(conn).ok();
1132
1133 if conn.grease {
1134 http3_conn.open_grease_stream(conn).ok();
1137 }
1138
1139 Ok(http3_conn)
1140 }
1141
1142 pub fn send_request<T: NameValue, F: BufFactory>(
1161 &mut self, conn: &mut super::Connection<F>, headers: &[T], fin: bool,
1162 ) -> Result<u64> {
1163 if self.peer_goaway_id.is_some() {
1166 return Err(Error::FrameUnexpected);
1167 }
1168
1169 let stream_id = self.next_request_stream_id;
1170
1171 self.streams.insert(
1172 stream_id,
1173 <stream::Stream>::new(
1174 stream_id,
1175 true,
1176 self.local_settings
1177 .max_field_section_size
1178 .unwrap_or(SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT),
1179 self.max_priority_update_size,
1180 ),
1181 );
1182
1183 if let Err(e) = conn.stream_send(stream_id, b"", false) {
1188 self.streams.remove(&stream_id);
1189
1190 if e == super::Error::Done {
1191 return Err(Error::StreamBlocked);
1192 }
1193
1194 return Err(e.into());
1195 };
1196
1197 if let Err(e) = self.send_headers(conn, stream_id, headers, fin) {
1198 if e == Error::StreamBlocked {
1204 self.streams.remove(&stream_id);
1205 }
1206
1207 return Err(e);
1208 }
1209
1210 self.next_request_stream_id = self
1213 .next_request_stream_id
1214 .checked_add(4)
1215 .ok_or(Error::IdError)?;
1216
1217 Ok(stream_id)
1218 }
1219
1220 pub fn send_response<T: NameValue, F: BufFactory>(
1258 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1259 headers: &[T], fin: bool,
1260 ) -> Result<()> {
1261 let priority = Default::default();
1262
1263 self.send_response_with_priority(
1264 conn, stream_id, headers, &priority, fin,
1265 )?;
1266
1267 Ok(())
1268 }
1269
1270 pub fn send_response_with_priority<T: NameValue, F: BufFactory>(
1314 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1315 headers: &[T], priority: &Priority, fin: bool,
1316 ) -> Result<()> {
1317 match self.streams.get(&stream_id) {
1318 Some(s) => {
1319 if s.local_initialized() {
1321 return Err(Error::FrameUnexpected);
1322 }
1323
1324 s
1325 },
1326
1327 None => return Err(Error::FrameUnexpected),
1328 };
1329
1330 self.send_headers(conn, stream_id, headers, fin)?;
1331
1332 let urgency = priority
1334 .urgency
1335 .clamp(PRIORITY_URGENCY_LOWER_BOUND, PRIORITY_URGENCY_UPPER_BOUND) +
1336 PRIORITY_URGENCY_OFFSET;
1337
1338 conn.stream_priority(stream_id, urgency, priority.incremental)?;
1339
1340 Ok(())
1341 }
1342
1343 pub fn send_additional_headers<T: NameValue, F: BufFactory>(
1370 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1371 headers: &[T], is_trailer_section: bool, fin: bool,
1372 ) -> Result<()> {
1373 if !self.is_server && !is_trailer_section {
1375 return Err(Error::FrameUnexpected);
1376 }
1377
1378 match self.streams.get(&stream_id) {
1379 Some(s) => {
1380 if !s.local_initialized() {
1382 return Err(Error::FrameUnexpected);
1383 }
1384
1385 if s.trailers_sent() {
1387 return Err(Error::FrameUnexpected);
1388 }
1389
1390 s
1391 },
1392
1393 None => return Err(Error::FrameUnexpected),
1394 };
1395
1396 self.send_headers(conn, stream_id, headers, fin)?;
1397
1398 if is_trailer_section {
1399 if let Some(s) = self.streams.get_mut(&stream_id) {
1402 s.mark_trailers_sent();
1403 }
1404 }
1405
1406 Ok(())
1407 }
1408
1409 pub fn send_additional_headers_with_priority<T: NameValue, F: BufFactory>(
1441 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1442 headers: &[T], priority: &Priority, is_trailer_section: bool, fin: bool,
1443 ) -> Result<()> {
1444 self.send_additional_headers(
1445 conn,
1446 stream_id,
1447 headers,
1448 is_trailer_section,
1449 fin,
1450 )?;
1451
1452 let urgency = priority
1454 .urgency
1455 .clamp(PRIORITY_URGENCY_LOWER_BOUND, PRIORITY_URGENCY_UPPER_BOUND) +
1456 PRIORITY_URGENCY_OFFSET;
1457
1458 conn.stream_priority(stream_id, urgency, priority.incremental)?;
1459
1460 Ok(())
1461 }
1462
1463 fn encode_header_block<T: NameValue>(
1464 &mut self, headers: &[T],
1465 ) -> Result<Vec<u8>> {
1466 let headers_len = headers
1467 .iter()
1468 .fold(0, |acc, h| acc + h.value().len() + h.name().len() + 32);
1469
1470 let mut header_block = vec![0; headers_len];
1471 let len = self
1472 .qpack_encoder
1473 .encode(headers, &mut header_block)
1474 .map_err(|_| Error::InternalError)?;
1475
1476 header_block.truncate(len);
1477
1478 Ok(header_block)
1479 }
1480
1481 fn send_headers<T: NameValue, F: BufFactory>(
1482 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1483 headers: &[T], fin: bool,
1484 ) -> Result<()> {
1485 let mut d = [42; 10];
1486 let mut b = octets::OctetsMut::with_slice(&mut d);
1487
1488 if !self.frames_greased && conn.grease {
1489 self.send_grease_frames(conn, stream_id)?;
1490 self.frames_greased = true;
1491 }
1492
1493 let header_block = self.encode_header_block(headers)?;
1494
1495 let overhead = octets::varint_len(frame::HEADERS_FRAME_TYPE_ID) +
1496 octets::varint_len(header_block.len() as u64);
1497
1498 match conn.stream_writable(stream_id, overhead + header_block.len()) {
1501 Ok(true) => (),
1502
1503 Ok(false) => return Err(Error::StreamBlocked),
1504
1505 Err(e) => {
1506 if conn.stream_finished(stream_id) {
1507 self.streams.remove(&stream_id);
1508 }
1509
1510 return Err(e.into());
1511 },
1512 };
1513
1514 b.put_varint(frame::HEADERS_FRAME_TYPE_ID)?;
1515 b.put_varint(header_block.len() as u64)?;
1516 let off = b.off();
1517 conn.stream_send(stream_id, &d[..off], false)?;
1518
1519 conn.stream_send(stream_id, &header_block, fin)?;
1521
1522 trace!(
1523 "{} tx frm HEADERS stream={} len={} fin={}",
1524 conn.trace_id(),
1525 stream_id,
1526 header_block.len(),
1527 fin
1528 );
1529
1530 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
1531 let qlog_headers = headers
1532 .iter()
1533 .map(|h| qlog::events::http3::HttpHeader {
1534 name: Some(String::from_utf8_lossy(h.name()).into_owned()),
1535 name_bytes: None,
1536 value: Some(String::from_utf8_lossy(h.value()).into_owned()),
1537 value_bytes: None,
1538 })
1539 .collect();
1540
1541 let frame = Http3Frame::Headers {
1542 headers: qlog_headers,
1543 raw: None,
1544 };
1545 let ev_data = EventData::Http3FrameCreated(FrameCreated {
1546 stream_id,
1547 length: Some(header_block.len() as u64),
1548 frame,
1549 ..Default::default()
1550 });
1551
1552 q.add_event_data_now(ev_data).ok();
1553 });
1554
1555 if fin {
1556 self.finish_local_stream(conn, stream_id, true);
1557 } else if let Some(s) = self.streams.get_mut(&stream_id) {
1558 s.initialize_local();
1559 }
1560
1561 Ok(())
1562 }
1563
1564 pub fn send_body<F: BufFactory>(
1579 &mut self, conn: &mut super::Connection<F>, stream_id: u64, body: &[u8],
1580 fin: bool,
1581 ) -> Result<usize> {
1582 self.do_send_body(
1583 conn,
1584 stream_id,
1585 body,
1586 fin,
1587 |conn: &mut super::Connection<F>,
1588 header: &[u8],
1589 stream_id: u64,
1590 body: &[u8],
1591 body_len: usize,
1592 fin: bool| {
1593 conn.stream_send(stream_id, header, false)?;
1594 Ok(conn
1595 .stream_send(stream_id, &body[..body_len], fin)
1596 .map(|v| (v, v))?)
1597 },
1598 )
1599 }
1600
1601 pub fn send_body_zc<F>(
1620 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1621 body: &mut F::Buf, fin: bool,
1622 ) -> Result<usize>
1623 where
1624 F: BufFactory,
1625 F::Buf: BufSplit,
1626 {
1627 self.do_send_body(
1628 conn,
1629 stream_id,
1630 body,
1631 fin,
1632 |conn: &mut super::Connection<F>,
1633 header: &[u8],
1634 stream_id: u64,
1635 body: &mut F::Buf,
1636 mut body_len: usize,
1637 fin: bool| {
1638 let with_prefix = body.try_add_prefix(header);
1639 if !with_prefix {
1640 conn.stream_send(stream_id, header, false)?;
1641 } else {
1642 body_len += header.len();
1643 }
1644
1645 let remainder = body.split_at(body_len);
1646 debug_assert_eq!(body.as_ref().len(), body_len);
1649
1650 let (mut n, rem) =
1651 conn.stream_send_zc(stream_id, body.clone(), fin)?;
1652 if rem.as_ref().is_some_and(|v| !v.as_ref().is_empty()) {
1653 debug_assert!(false);
1656 return Err(Error::InternalError);
1657 }
1658
1659 if with_prefix {
1660 n -= header.len();
1661 }
1662
1663 if !remainder.as_ref().is_empty() {
1664 let _ = std::mem::replace(body, remainder);
1665 }
1666
1667 Ok((n, n))
1668 },
1669 )
1670 }
1671
1672 fn do_send_body<F, B, R, SND>(
1673 &mut self, conn: &mut super::Connection<F>, stream_id: u64, body: B,
1674 fin: bool, write_fn: SND,
1675 ) -> Result<R>
1676 where
1677 F: BufFactory,
1678 B: AsRef<[u8]>,
1679 SND: FnOnce(
1680 &mut super::Connection<F>,
1681 &[u8],
1682 u64,
1683 B,
1684 usize,
1685 bool,
1686 ) -> Result<(usize, R)>,
1687 {
1688 let mut d = [42; 10];
1689 let mut b = octets::OctetsMut::with_slice(&mut d);
1690
1691 let len = body.as_ref().len();
1692
1693 if !stream_id.is_multiple_of(4) {
1695 return Err(Error::FrameUnexpected);
1696 }
1697
1698 match self.streams.get_mut(&stream_id) {
1699 Some(s) => {
1700 if !s.local_initialized() {
1701 return Err(Error::FrameUnexpected);
1702 }
1703
1704 if s.trailers_sent() {
1705 return Err(Error::FrameUnexpected);
1706 }
1707 },
1708
1709 None => {
1710 return Err(Error::FrameUnexpected);
1711 },
1712 };
1713
1714 if len == 0 && !fin {
1716 return Err(Error::Done);
1717 }
1718
1719 let overhead = octets::varint_len(frame::DATA_FRAME_TYPE_ID) +
1720 octets::varint_len(len as u64);
1721
1722 let stream_cap = match conn.stream_capacity(stream_id) {
1723 Ok(v) => v,
1724
1725 Err(e) => {
1726 if conn.stream_finished(stream_id) {
1727 self.streams.remove(&stream_id);
1728 }
1729
1730 return Err(e.into());
1731 },
1732 };
1733
1734 if stream_cap < overhead {
1736 let _ = conn.stream_writable(stream_id, overhead + 1);
1737 return Err(Error::Done);
1738 }
1739
1740 let body_len = std::cmp::min(len, stream_cap - overhead);
1742
1743 let fin = if body_len != len { false } else { fin };
1746
1747 if body_len == 0 && !fin {
1749 let _ = conn.stream_writable(stream_id, overhead + 1);
1750 return Err(Error::Done);
1751 }
1752
1753 b.put_varint(frame::DATA_FRAME_TYPE_ID)?;
1754 b.put_varint(body_len as u64)?;
1755 let off = b.off();
1756
1757 let (written, ret) =
1760 write_fn(conn, &d[..off], stream_id, body, body_len, fin)?;
1761 if written != body_len {
1762 debug_assert!(false);
1765 return Err(Error::InternalError);
1766 }
1767
1768 trace!(
1769 "{} tx frm DATA stream={} len={} fin={}",
1770 conn.trace_id(),
1771 stream_id,
1772 written,
1773 fin
1774 );
1775
1776 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
1777 let frame = Http3Frame::Data { raw: None };
1778 let ev_data = EventData::Http3FrameCreated(FrameCreated {
1779 stream_id,
1780 length: Some(written as u64),
1781 frame,
1782 ..Default::default()
1783 });
1784
1785 q.add_event_data_now(ev_data).ok();
1786 });
1787
1788 if written < len {
1789 let _ = conn.stream_writable(stream_id, overhead + 1);
1796 }
1797
1798 if fin && written == len {
1799 self.finish_local_stream(conn, stream_id, false);
1800 }
1801
1802 Ok(ret)
1803 }
1804
1805 pub fn dgram_enabled_by_peer<F: BufFactory>(
1813 &self, conn: &super::Connection<F>,
1814 ) -> bool {
1815 self.peer_settings.h3_datagram == Some(1) &&
1816 conn.dgram_max_writable_len().is_some()
1817 }
1818
1819 pub fn extended_connect_enabled_by_peer(&self) -> bool {
1827 self.peer_settings.connect_protocol_enabled == Some(1)
1828 }
1829
1830 pub fn recv_body<F: BufFactory>(
1842 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1843 out: &mut [u8],
1844 ) -> Result<usize> {
1845 self.recv_body_buf(conn, stream_id, out)
1846 }
1847
1848 pub fn recv_body_buf<F: BufFactory, OUT: bytes::BufMut>(
1885 &mut self, conn: &mut super::Connection<F>, stream_id: u64, mut out: OUT,
1886 ) -> Result<usize> {
1887 let mut total = 0;
1888
1889 while out.has_remaining_mut() {
1895 let stream = self.streams.get_mut(&stream_id).ok_or(Error::Done)?;
1896
1897 if stream.state() != stream::State::Data {
1898 break;
1899 }
1900
1901 let (read, fin) = match stream.try_consume_data(conn, &mut out) {
1902 Ok(v) => v,
1903
1904 Err(Error::Done) => break,
1905
1906 Err(e) => return Err(e),
1907 };
1908
1909 total += read;
1910
1911 if read == 0 || fin {
1913 break;
1914 }
1915
1916 match self.process_readable_stream(conn, stream_id, false) {
1921 Ok(_) => unreachable!(),
1922
1923 Err(Error::Done) => (),
1924
1925 Err(e) => return Err(e),
1926 };
1927
1928 if conn.stream_finished(stream_id) {
1929 break;
1930 }
1931 }
1932
1933 if conn.stream_finished(stream_id) {
1936 self.process_finished_stream(stream_id);
1937 }
1938
1939 if total == 0 {
1940 return Err(Error::Done);
1941 }
1942
1943 Ok(total)
1944 }
1945
1946 pub fn send_priority_update_for_request<F: BufFactory>(
1961 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
1962 priority: &Priority,
1963 ) -> Result<()> {
1964 let mut d = [42; 20];
1965 let mut b = octets::OctetsMut::with_slice(&mut d);
1966
1967 if self.is_server {
1969 return Err(Error::FrameUnexpected);
1970 }
1971
1972 if !stream_id.is_multiple_of(4) {
1973 return Err(Error::FrameUnexpected);
1974 }
1975
1976 let control_stream_id =
1977 self.control_stream_id.ok_or(Error::FrameUnexpected)?;
1978
1979 let urgency = priority
1980 .urgency
1981 .clamp(PRIORITY_URGENCY_LOWER_BOUND, PRIORITY_URGENCY_UPPER_BOUND);
1982
1983 let mut field_value = format!("u={urgency}");
1984
1985 if priority.incremental {
1986 field_value.push_str(",i");
1987 }
1988
1989 let priority_field_value = field_value.as_bytes();
1990 let frame_payload_len =
1991 octets::varint_len(stream_id) + priority_field_value.len();
1992
1993 let overhead =
1994 octets::varint_len(frame::PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID) +
1995 octets::varint_len(stream_id) +
1996 octets::varint_len(frame_payload_len as u64);
1997
1998 match conn.stream_writable(
2000 control_stream_id,
2001 overhead + priority_field_value.len(),
2002 ) {
2003 Ok(true) => (),
2004
2005 Ok(false) => return Err(Error::StreamBlocked),
2006
2007 Err(e) => {
2008 return Err(e.into());
2009 },
2010 }
2011
2012 b.put_varint(frame::PRIORITY_UPDATE_FRAME_REQUEST_TYPE_ID)?;
2013 b.put_varint(frame_payload_len as u64)?;
2014 b.put_varint(stream_id)?;
2015 let off = b.off();
2016 conn.stream_send(control_stream_id, &d[..off], false)?;
2017
2018 conn.stream_send(control_stream_id, priority_field_value, false)?;
2020
2021 trace!(
2022 "{} tx frm PRIORITY_UPDATE request_stream={} priority_field_value={}",
2023 conn.trace_id(),
2024 stream_id,
2025 field_value,
2026 );
2027
2028 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
2029 let frame = Http3Frame::PriorityUpdate {
2030 stream_id: Some(stream_id),
2031 push_id: None,
2032 priority_field_value: field_value.clone(),
2033 raw: None,
2034 };
2035
2036 let ev_data = EventData::Http3FrameCreated(FrameCreated {
2037 stream_id,
2038 length: Some(priority_field_value.len() as u64),
2039 frame,
2040 ..Default::default()
2041 });
2042
2043 q.add_event_data_now(ev_data).ok();
2044 });
2045
2046 Ok(())
2047 }
2048
2049 pub fn take_last_priority_update(
2064 &mut self, prioritized_element_id: u64,
2065 ) -> Result<Vec<u8>> {
2066 if let Some(stream) = self.streams.get_mut(&prioritized_element_id) {
2067 return stream.take_last_priority_update().ok_or(Error::Done);
2068 }
2069
2070 Err(Error::Done)
2071 }
2072
2073 pub fn poll<F: BufFactory>(
2111 &mut self, conn: &mut super::Connection<F>,
2112 ) -> Result<(u64, Event)> {
2113 if conn.local_error.is_some() {
2117 return Err(Error::Done);
2118 }
2119
2120 if let Some(stream_id) = self.peer_control_stream_id {
2122 match self.process_control_stream(conn, stream_id) {
2123 Ok(ev) => return Ok(ev),
2124
2125 Err(Error::Done) => (),
2126
2127 Err(e) => return Err(e),
2128 };
2129 }
2130
2131 if let Some(stream_id) = self.peer_qpack_streams.encoder_stream_id {
2132 match self.process_control_stream(conn, stream_id) {
2133 Ok(ev) => return Ok(ev),
2134
2135 Err(Error::Done) => (),
2136
2137 Err(e) => return Err(e),
2138 };
2139 }
2140
2141 if let Some(stream_id) = self.peer_qpack_streams.decoder_stream_id {
2142 match self.process_control_stream(conn, stream_id) {
2143 Ok(ev) => return Ok(ev),
2144
2145 Err(Error::Done) => (),
2146
2147 Err(e) => return Err(e),
2148 };
2149 }
2150
2151 if let Some(ev) = self.pop_finished_stream(conn) {
2153 return Ok(ev);
2154 }
2155
2156 for s in conn.readable() {
2158 trace!("{} stream id {} is readable", conn.trace_id(), s);
2159
2160 let ev = match self.process_readable_stream(conn, s, true) {
2161 Ok(v) => Some(v),
2162
2163 Err(Error::Done) => None,
2164
2165 Err(Error::TransportError(crate::Error::StreamReset(e))) => {
2168 self.remove_local_finished_stream(s);
2169
2170 return Ok((s, Event::Reset(e)));
2171 },
2172
2173 Err(e) => return Err(e),
2174 };
2175
2176 if conn.stream_finished(s) {
2177 self.process_finished_stream(s);
2178 }
2179
2180 if let Some(ev) = ev {
2182 return Ok(ev);
2183 }
2184 }
2185
2186 if let Some(ev) = self.pop_finished_stream(conn) {
2190 return Ok(ev);
2191 }
2192
2193 Err(Error::Done)
2194 }
2195
2196 pub fn send_goaway<F: BufFactory>(
2208 &mut self, conn: &mut super::Connection<F>, id: u64,
2209 ) -> Result<()> {
2210 let mut id = id;
2211
2212 if !self.is_server {
2216 id = 0;
2217 }
2218
2219 if self.is_server && !id.is_multiple_of(4) {
2220 return Err(Error::IdError);
2221 }
2222
2223 if let Some(sent_id) = self.local_goaway_id {
2224 if id > sent_id {
2225 return Err(Error::IdError);
2226 }
2227 }
2228
2229 if let Some(stream_id) = self.control_stream_id {
2230 let mut d = [42; 10];
2231 let mut b = octets::OctetsMut::with_slice(&mut d);
2232
2233 let frame = frame::Frame::GoAway { id };
2234
2235 let wire_len = frame.to_bytes(&mut b)?;
2236 let stream_cap = conn.stream_capacity(stream_id)?;
2237
2238 if stream_cap < wire_len {
2239 return Err(Error::StreamBlocked);
2240 }
2241
2242 trace!("{} tx frm {:?}", conn.trace_id(), frame);
2243
2244 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
2245 let ev_data = EventData::Http3FrameCreated(FrameCreated {
2246 stream_id,
2247 length: Some(octets::varint_len(id) as u64),
2248 frame: frame.to_qlog(),
2249 ..Default::default()
2250 });
2251
2252 q.add_event_data_now(ev_data).ok();
2253 });
2254
2255 let off = b.off();
2256 conn.stream_send(stream_id, &d[..off], false)?;
2257
2258 self.local_goaway_id = Some(id);
2259 }
2260
2261 Ok(())
2262 }
2263
2264 pub fn peer_settings_raw(&self) -> Option<&[(u64, u64)]> {
2268 self.peer_settings.raw.as_deref()
2269 }
2270
2271 fn open_uni_stream<F: BufFactory>(
2272 &mut self, conn: &mut super::Connection<F>, ty: u64,
2273 ) -> Result<u64> {
2274 let stream_id = self.next_uni_stream_id;
2275
2276 let mut d = [0; 8];
2277 let mut b = octets::OctetsMut::with_slice(&mut d);
2278
2279 match ty {
2280 stream::HTTP3_CONTROL_STREAM_TYPE_ID |
2282 stream::QPACK_ENCODER_STREAM_TYPE_ID |
2283 stream::QPACK_DECODER_STREAM_TYPE_ID => {
2284 conn.stream_priority(stream_id, 0, false)?;
2285 },
2286
2287 stream::HTTP3_PUSH_STREAM_TYPE_ID => (),
2289
2290 _ => {
2292 conn.stream_priority(stream_id, 255, false)?;
2293 },
2294 }
2295
2296 conn.stream_send(stream_id, b.put_varint(ty)?, false)?;
2297
2298 self.next_uni_stream_id = self
2301 .next_uni_stream_id
2302 .checked_add(4)
2303 .ok_or(Error::IdError)?;
2304
2305 Ok(stream_id)
2306 }
2307
2308 fn open_qpack_encoder_stream<F: BufFactory>(
2309 &mut self, conn: &mut super::Connection<F>,
2310 ) -> Result<()> {
2311 let stream_id =
2312 self.open_uni_stream(conn, stream::QPACK_ENCODER_STREAM_TYPE_ID)?;
2313
2314 self.local_qpack_streams.encoder_stream_id = Some(stream_id);
2315
2316 qlog_with_type!(QLOG_STREAM_TYPE_SET, conn.qlog, q, {
2317 let ev_data = EventData::Http3StreamTypeSet(StreamTypeSet {
2318 stream_id,
2319 initiator: Some(Initiator::Local),
2320 stream_type: StreamType::QpackEncode,
2321 ..Default::default()
2322 });
2323
2324 q.add_event_data_now(ev_data).ok();
2325 });
2326
2327 Ok(())
2328 }
2329
2330 fn open_qpack_decoder_stream<F: BufFactory>(
2331 &mut self, conn: &mut super::Connection<F>,
2332 ) -> Result<()> {
2333 let stream_id =
2334 self.open_uni_stream(conn, stream::QPACK_DECODER_STREAM_TYPE_ID)?;
2335
2336 self.local_qpack_streams.decoder_stream_id = Some(stream_id);
2337
2338 qlog_with_type!(QLOG_STREAM_TYPE_SET, conn.qlog, q, {
2339 let ev_data = EventData::Http3StreamTypeSet(StreamTypeSet {
2340 stream_id,
2341 initiator: Some(Initiator::Local),
2342 stream_type: StreamType::QpackDecode,
2343 ..Default::default()
2344 });
2345
2346 q.add_event_data_now(ev_data).ok();
2347 });
2348
2349 Ok(())
2350 }
2351
2352 fn send_grease_frames<F: BufFactory>(
2354 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
2355 ) -> Result<()> {
2356 let mut d = [0; 8];
2357
2358 let stream_cap = match conn.stream_capacity(stream_id) {
2359 Ok(v) => v,
2360
2361 Err(e) => {
2362 if conn.stream_finished(stream_id) {
2363 self.streams.remove(&stream_id);
2364 }
2365
2366 return Err(e.into());
2367 },
2368 };
2369
2370 let grease_frame1 = grease_value();
2371 let grease_frame2 = grease_value();
2372 let grease_payload = b"GREASE is the word";
2373
2374 let overhead = octets::varint_len(grease_frame1) + 1 + octets::varint_len(grease_frame2) + 1 + grease_payload.len(); if stream_cap < overhead {
2383 return Ok(());
2384 }
2385
2386 let mut b = octets::OctetsMut::with_slice(&mut d);
2388 conn.stream_send(stream_id, b.put_varint(grease_frame1)?, false)?;
2389
2390 let mut b = octets::OctetsMut::with_slice(&mut d);
2391 conn.stream_send(stream_id, b.put_varint(0)?, false)?;
2392
2393 trace!(
2394 "{} tx frm GREASE stream={} len=0",
2395 conn.trace_id(),
2396 stream_id
2397 );
2398
2399 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
2400 let frame = Http3Frame::Reserved {
2401 frame_type_bytes: grease_frame1,
2402 raw: None,
2403 };
2404 let ev_data = EventData::Http3FrameCreated(FrameCreated {
2405 stream_id,
2406 length: Some(0),
2407 frame,
2408 ..Default::default()
2409 });
2410
2411 q.add_event_data_now(ev_data).ok();
2412 });
2413
2414 let mut b = octets::OctetsMut::with_slice(&mut d);
2416 conn.stream_send(stream_id, b.put_varint(grease_frame2)?, false)?;
2417
2418 let mut b = octets::OctetsMut::with_slice(&mut d);
2419 conn.stream_send(stream_id, b.put_varint(18)?, false)?;
2420
2421 conn.stream_send(stream_id, grease_payload, false)?;
2422
2423 trace!(
2424 "{} tx frm GREASE stream={} len={}",
2425 conn.trace_id(),
2426 stream_id,
2427 grease_payload.len()
2428 );
2429
2430 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
2431 let frame = Http3Frame::Reserved {
2432 frame_type_bytes: grease_frame2,
2433 raw: None,
2434 };
2435 let ev_data = EventData::Http3FrameCreated(FrameCreated {
2436 stream_id,
2437 length: Some(grease_payload.len() as u64),
2438 frame,
2439 ..Default::default()
2440 });
2441
2442 q.add_event_data_now(ev_data).ok();
2443 });
2444
2445 Ok(())
2446 }
2447
2448 fn open_grease_stream<F: BufFactory>(
2451 &mut self, conn: &mut super::Connection<F>,
2452 ) -> Result<()> {
2453 let ty = grease_value();
2454 match self.open_uni_stream(conn, ty) {
2455 Ok(stream_id) => {
2456 conn.stream_send(stream_id, b"GREASE is the word", true)?;
2457
2458 trace!("{} open GREASE stream {}", conn.trace_id(), stream_id);
2459
2460 qlog_with_type!(QLOG_STREAM_TYPE_SET, conn.qlog, q, {
2461 let ev_data = EventData::Http3StreamTypeSet(StreamTypeSet {
2462 stream_id,
2463 initiator: Some(Initiator::Local),
2464 stream_type: StreamType::Unknown,
2465 stream_type_bytes: Some(ty),
2466 ..Default::default()
2467 });
2468
2469 q.add_event_data_now(ev_data).ok();
2470 });
2471 },
2472
2473 Err(Error::IdError) => {
2474 trace!("{} GREASE stream blocked", conn.trace_id(),);
2475
2476 return Ok(());
2477 },
2478
2479 Err(e) => return Err(e),
2480 };
2481
2482 Ok(())
2483 }
2484
2485 fn send_settings<F: BufFactory>(
2487 &mut self, conn: &mut super::Connection<F>,
2488 ) -> Result<()> {
2489 let stream_id = match self
2490 .open_uni_stream(conn, stream::HTTP3_CONTROL_STREAM_TYPE_ID)
2491 {
2492 Ok(v) => v,
2493
2494 Err(e) => {
2495 trace!("{} Control stream blocked", conn.trace_id(),);
2496
2497 if e == Error::Done {
2498 return Err(Error::InternalError);
2499 }
2500
2501 return Err(e);
2502 },
2503 };
2504
2505 self.control_stream_id = Some(stream_id);
2506
2507 qlog_with_type!(QLOG_STREAM_TYPE_SET, conn.qlog, q, {
2508 let ev_data = EventData::Http3StreamTypeSet(StreamTypeSet {
2509 stream_id,
2510 initiator: Some(Initiator::Local),
2511 stream_type: StreamType::Control,
2512 ..Default::default()
2513 });
2514
2515 q.add_event_data_now(ev_data).ok();
2516 });
2517
2518 let grease = if conn.grease {
2519 Some((grease_value(), grease_value()))
2520 } else {
2521 None
2522 };
2523
2524 let frame = frame::Frame::Settings {
2525 max_field_section_size: self.local_settings.max_field_section_size,
2526 qpack_max_table_capacity: self
2527 .local_settings
2528 .qpack_max_table_capacity,
2529 qpack_blocked_streams: self.local_settings.qpack_blocked_streams,
2530 connect_protocol_enabled: self
2531 .local_settings
2532 .connect_protocol_enabled,
2533 h3_datagram: self.local_settings.h3_datagram,
2534 grease,
2535 additional_settings: self.local_settings.additional_settings.clone(),
2536 raw: Default::default(),
2537 };
2538
2539 let mut d = [42; 128];
2540 let mut b = octets::OctetsMut::with_slice(&mut d);
2541
2542 frame.to_bytes(&mut b)?;
2543
2544 let off = b.off();
2545
2546 if let Some(id) = self.control_stream_id {
2547 conn.stream_send(id, &d[..off], false)?;
2548
2549 trace!(
2550 "{} tx frm SETTINGS stream={} len={}",
2551 conn.trace_id(),
2552 id,
2553 off
2554 );
2555
2556 qlog_with_type!(QLOG_FRAME_CREATED, conn.qlog, q, {
2557 let frame = frame.to_qlog();
2558 let ev_data = EventData::Http3FrameCreated(FrameCreated {
2559 stream_id: id,
2560 length: Some(off as u64),
2561 frame,
2562 ..Default::default()
2563 });
2564
2565 q.add_event_data_now(ev_data).ok();
2566 });
2567 }
2568
2569 Ok(())
2570 }
2571
2572 fn process_control_stream<F: BufFactory>(
2573 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
2574 ) -> Result<(u64, Event)> {
2575 close_conn_if_critical_stream_finished(conn, stream_id)?;
2576
2577 if !conn.stream_readable(stream_id) {
2578 return Err(Error::Done);
2579 }
2580
2581 match self.process_readable_stream(conn, stream_id, true) {
2582 Ok(ev) => return Ok(ev),
2583
2584 Err(Error::Done) => (),
2585
2586 Err(e) => return Err(e),
2587 };
2588
2589 close_conn_if_critical_stream_finished(conn, stream_id)?;
2590
2591 Err(Error::Done)
2592 }
2593
2594 fn process_readable_stream<F: BufFactory>(
2595 &mut self, conn: &mut super::Connection<F>, stream_id: u64, polling: bool,
2596 ) -> Result<(u64, Event)> {
2597 self.streams.entry(stream_id).or_insert_with(|| {
2598 <stream::Stream>::new(
2599 stream_id,
2600 false,
2601 self.local_settings
2602 .max_field_section_size
2603 .unwrap_or(SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT),
2604 self.max_priority_update_size,
2605 )
2606 });
2607
2608 while let Some(stream) = self.streams.get_mut(&stream_id) {
2613 match stream.state() {
2614 stream::State::StreamType => {
2615 stream.try_fill_buffer(conn)?;
2616
2617 let varint = match stream.try_consume_varint() {
2618 Ok(v) => v,
2619
2620 Err(_) => continue,
2621 };
2622
2623 let ty = stream::Type::deserialize(varint)?;
2624
2625 if let Err(e) = stream.set_ty(ty) {
2626 conn.close(true, e.to_wire(), b"")?;
2627 return Err(e);
2628 }
2629
2630 qlog_with_type!(QLOG_STREAM_TYPE_SET, conn.qlog, q, {
2631 let ty_val = if matches!(ty, stream::Type::Unknown) {
2632 Some(varint)
2633 } else {
2634 None
2635 };
2636
2637 let ev_data =
2638 EventData::Http3StreamTypeSet(StreamTypeSet {
2639 stream_id,
2640 initiator: Some(Initiator::Remote),
2641 stream_type: ty.to_qlog(),
2642 stream_type_bytes: ty_val,
2643 ..Default::default()
2644 });
2645
2646 q.add_event_data_now(ev_data).ok();
2647 });
2648
2649 match &ty {
2650 stream::Type::Control => {
2651 if self.peer_control_stream_id.is_some() {
2653 conn.close(
2654 true,
2655 Error::StreamCreationError.to_wire(),
2656 b"Received multiple control streams",
2657 )?;
2658
2659 return Err(Error::StreamCreationError);
2660 }
2661
2662 trace!(
2663 "{} open peer's control stream {}",
2664 conn.trace_id(),
2665 stream_id
2666 );
2667
2668 close_conn_if_critical_stream_finished(
2669 conn, stream_id,
2670 )?;
2671
2672 self.peer_control_stream_id = Some(stream_id);
2673 },
2674
2675 stream::Type::Push => {
2676 conn.close(
2680 true,
2681 Error::StreamCreationError.to_wire(),
2682 b"Received push stream.",
2683 )?;
2684
2685 return Err(Error::StreamCreationError);
2686 },
2687
2688 stream::Type::QpackEncoder => {
2689 if self.peer_qpack_streams.encoder_stream_id.is_some()
2691 {
2692 conn.close(
2693 true,
2694 Error::StreamCreationError.to_wire(),
2695 b"Received multiple QPACK encoder streams",
2696 )?;
2697
2698 return Err(Error::StreamCreationError);
2699 }
2700
2701 close_conn_if_critical_stream_finished(
2702 conn, stream_id,
2703 )?;
2704
2705 self.peer_qpack_streams.encoder_stream_id =
2706 Some(stream_id);
2707 },
2708
2709 stream::Type::QpackDecoder => {
2710 if self.peer_qpack_streams.decoder_stream_id.is_some()
2712 {
2713 conn.close(
2714 true,
2715 Error::StreamCreationError.to_wire(),
2716 b"Received multiple QPACK decoder streams",
2717 )?;
2718
2719 return Err(Error::StreamCreationError);
2720 }
2721
2722 close_conn_if_critical_stream_finished(
2723 conn, stream_id,
2724 )?;
2725
2726 self.peer_qpack_streams.decoder_stream_id =
2727 Some(stream_id);
2728 },
2729
2730 stream::Type::Unknown => {
2731 },
2734
2735 stream::Type::Request => unreachable!(),
2736 }
2737 },
2738
2739 stream::State::PushId => {
2740 stream.try_fill_buffer(conn)?;
2741
2742 let varint = match stream.try_consume_varint() {
2743 Ok(v) => v,
2744
2745 Err(_) => continue,
2746 };
2747
2748 if let Err(e) = stream.set_push_id(varint) {
2749 conn.close(true, e.to_wire(), b"")?;
2750 return Err(e);
2751 }
2752 },
2753
2754 stream::State::FrameType => {
2755 stream.try_fill_buffer(conn)?;
2756
2757 let varint = match stream.try_consume_varint() {
2758 Ok(v) => v,
2759
2760 Err(_) => continue,
2761 };
2762
2763 match stream.set_frame_type(varint) {
2764 Err(Error::FrameUnexpected) => {
2765 let msg = format!("Unexpected frame type {varint}");
2766
2767 conn.close(
2768 true,
2769 Error::FrameUnexpected.to_wire(),
2770 msg.as_bytes(),
2771 )?;
2772
2773 return Err(Error::FrameUnexpected);
2774 },
2775
2776 Err(e) => {
2777 conn.close(
2778 true,
2779 e.to_wire(),
2780 b"Error handling frame.",
2781 )?;
2782
2783 return Err(e);
2784 },
2785
2786 _ => (),
2787 }
2788 },
2789
2790 stream::State::FramePayloadLen => {
2791 stream.try_fill_buffer(conn)?;
2792
2793 let payload_len = match stream.try_consume_varint() {
2794 Ok(v) => v,
2795
2796 Err(_) => continue,
2797 };
2798
2799 if Some(frame::DATA_FRAME_TYPE_ID) == stream.frame_type() {
2802 trace!(
2803 "{} rx frm DATA stream={} wire_payload_len={}",
2804 conn.trace_id(),
2805 stream_id,
2806 payload_len
2807 );
2808
2809 qlog_with_type!(QLOG_FRAME_PARSED, conn.qlog, q, {
2810 let frame = Http3Frame::Data { raw: None };
2811
2812 let ev_data =
2813 EventData::Http3FrameParsed(FrameParsed {
2814 stream_id,
2815 length: Some(payload_len),
2816 frame,
2817 ..Default::default()
2818 });
2819
2820 q.add_event_data_now(ev_data).ok();
2821 });
2822 }
2823
2824 let res = stream.set_frame_payload_len(payload_len);
2825
2826 if let Err(e) = res {
2827 conn.close(true, e.to_wire(), b"")?;
2828 return Err(e);
2829 }
2830 },
2831
2832 stream::State::FramePayload => {
2833 if !polling {
2835 break;
2836 }
2837
2838 stream.try_fill_buffer(conn)?;
2839
2840 let (frame, payload_len) = match stream.try_consume_frame() {
2841 Ok(frame) => frame,
2842
2843 Err(Error::Done) => return Err(Error::Done),
2844
2845 Err(e) => {
2846 conn.close(
2847 true,
2848 e.to_wire(),
2849 b"Error handling frame.",
2850 )?;
2851
2852 return Err(e);
2853 },
2854 };
2855
2856 match self.process_frame(conn, stream_id, frame, payload_len)
2857 {
2858 Ok(ev) => return Ok(ev),
2859
2860 Err(Error::Done) => {
2861 if conn.stream_finished(stream_id) {
2864 break;
2865 }
2866 },
2867
2868 Err(e) => return Err(e),
2869 };
2870 },
2871
2872 stream::State::Data => {
2873 if !polling {
2875 break;
2876 }
2877
2878 if !stream.try_trigger_data_event() {
2879 break;
2880 }
2881
2882 return Ok((stream_id, Event::Data));
2883 },
2884
2885 stream::State::QpackInstruction => {
2886 let mut d = [0; 4096];
2887
2888 loop {
2890 let (recv, fin) = conn.stream_recv(stream_id, &mut d)?;
2891
2892 match stream.ty() {
2893 Some(stream::Type::QpackEncoder) =>
2894 self.peer_qpack_streams.encoder_stream_bytes +=
2895 recv as u64,
2896 Some(stream::Type::QpackDecoder) =>
2897 self.peer_qpack_streams.decoder_stream_bytes +=
2898 recv as u64,
2899 _ => unreachable!(),
2900 };
2901
2902 if fin {
2903 close_conn_critical_stream(conn)?;
2904 }
2905 }
2906 },
2907
2908 stream::State::SkipFramePayload => {
2909 stream.try_skip_frame(conn)?;
2910
2911 if conn.stream_finished(stream_id) {
2914 break;
2915 }
2916 },
2917
2918 stream::State::Drain => {
2919 conn.stream_shutdown(
2921 stream_id,
2922 crate::Shutdown::Read,
2923 0x100,
2924 )?;
2925
2926 break;
2927 },
2928
2929 stream::State::Finished => break,
2930 }
2931 }
2932
2933 Err(Error::Done)
2934 }
2935
2936 fn process_finished_stream(&mut self, stream_id: u64) {
2937 let stream = match self.streams.get_mut(&stream_id) {
2938 Some(v) => v,
2939
2940 None => return,
2941 };
2942
2943 if stream.state() == stream::State::Finished {
2944 return;
2945 }
2946
2947 match stream.ty() {
2948 Some(stream::Type::Request) | Some(stream::Type::Push) => {
2949 stream.finished();
2950
2951 self.finished_streams.push_back(stream_id);
2952 },
2953 Some(stream::Type::Unknown) | None => {
2954 self.streams.remove(&stream_id);
2955 },
2956 Some(stream::Type::Control) |
2959 Some(stream::Type::QpackEncoder) |
2960 Some(stream::Type::QpackDecoder) => (),
2961 };
2962 }
2963
2964 fn finish_local_stream<F: BufFactory>(
2965 &mut self, conn: &super::Connection<F>, stream_id: u64,
2966 initialize_local: bool,
2967 ) {
2968 let hash_map::Entry::Occupied(mut stream) = self.streams.entry(stream_id)
2969 else {
2970 return;
2971 };
2972
2973 {
2974 let stream = stream.get_mut();
2975
2976 if initialize_local {
2977 stream.initialize_local();
2978 }
2979
2980 stream.finish_local();
2981 }
2982
2983 if conn.stream_finished(stream_id) {
2984 stream.remove();
2985 }
2986 }
2987
2988 fn remove_local_finished_stream(&mut self, stream_id: u64) {
2989 if let hash_map::Entry::Occupied(stream) = self.streams.entry(stream_id) {
2990 if stream.get().local_finished() {
2991 stream.remove();
2992 }
2993 }
2994 }
2995
2996 fn pop_finished_stream<F: BufFactory>(
2997 &mut self, conn: &mut super::Connection<F>,
2998 ) -> Option<(u64, Event)> {
2999 let finished = self.finished_streams.pop_front()?;
3000
3001 self.remove_local_finished_stream(finished);
3002
3003 if conn.stream_readable(finished) {
3004 if let Err(crate::Error::StreamReset(e)) =
3007 conn.stream_recv(finished, &mut [])
3008 {
3009 return Some((finished, Event::Reset(e)));
3010 }
3011 }
3012
3013 Some((finished, Event::Finished))
3014 }
3015
3016 fn process_frame<F: BufFactory>(
3017 &mut self, conn: &mut super::Connection<F>, stream_id: u64,
3018 frame: frame::Frame, payload_len: u64,
3019 ) -> Result<(u64, Event)> {
3020 trace!(
3021 "{} rx frm {:?} stream={} payload_len={}",
3022 conn.trace_id(),
3023 frame,
3024 stream_id,
3025 payload_len
3026 );
3027
3028 qlog_with_type!(QLOG_FRAME_PARSED, conn.qlog, q, {
3029 if !matches!(frame, frame::Frame::Headers { .. }) {
3031 let frame = frame.to_qlog();
3032 let ev_data = EventData::Http3FrameParsed(FrameParsed {
3033 stream_id,
3034 length: Some(payload_len),
3035 frame,
3036 ..Default::default()
3037 });
3038
3039 q.add_event_data_now(ev_data).ok();
3040 }
3041 });
3042
3043 match frame {
3044 frame::Frame::Settings {
3045 max_field_section_size,
3046 qpack_max_table_capacity,
3047 qpack_blocked_streams,
3048 connect_protocol_enabled,
3049 h3_datagram,
3050 additional_settings,
3051 raw,
3052 ..
3053 } => {
3054 self.peer_settings = ConnectionSettings {
3055 max_field_section_size,
3056 qpack_max_table_capacity,
3057 qpack_blocked_streams,
3058 connect_protocol_enabled,
3059 h3_datagram,
3060 additional_settings,
3061 raw,
3062 };
3063
3064 if let Some(1) = h3_datagram {
3065 if conn.dgram_max_writable_len().is_none() {
3067 conn.close(
3068 true,
3069 Error::SettingsError.to_wire(),
3070 b"H3_DATAGRAM sent with value 1 but max_datagram_frame_size TP not set.",
3071 )?;
3072
3073 return Err(Error::SettingsError);
3074 }
3075 }
3076 },
3077
3078 frame::Frame::Headers { header_block } => {
3079 if let Some(s) = self.streams.get_mut(&stream_id) {
3081 if self.is_server && s.headers_received_count() == 2 {
3082 conn.close(
3083 true,
3084 Error::FrameUnexpected.to_wire(),
3085 b"Too many HEADERS frames",
3086 )?;
3087 return Err(Error::FrameUnexpected);
3088 }
3089
3090 s.increment_headers_received();
3091 }
3092
3093 let max_size = self
3096 .local_settings
3097 .max_field_section_size
3098 .unwrap_or(u64::MAX);
3099
3100 let headers = match self
3101 .qpack_decoder
3102 .decode(&header_block[..], max_size)
3103 {
3104 Ok(v) => v,
3105
3106 Err(e) => {
3107 let e = match e {
3108 qpack::Error::HeaderListTooLarge =>
3109 Error::ExcessiveLoad,
3110
3111 _ => Error::QpackDecompressionFailed,
3112 };
3113
3114 conn.close(true, e.to_wire(), b"Error parsing headers.")?;
3115
3116 return Err(e);
3117 },
3118 };
3119
3120 qlog_with_type!(QLOG_FRAME_PARSED, conn.qlog, q, {
3121 let qlog_headers = headers
3122 .iter()
3123 .map(|h| qlog::events::http3::HttpHeader {
3124 name: Some(
3125 String::from_utf8_lossy(h.name()).into_owned(),
3126 ),
3127 name_bytes: None,
3128 value: Some(
3129 String::from_utf8_lossy(h.value()).into_owned(),
3130 ),
3131 value_bytes: None,
3132 })
3133 .collect();
3134
3135 let frame = Http3Frame::Headers {
3136 headers: qlog_headers,
3137 raw: None,
3138 };
3139
3140 let ev_data = EventData::Http3FrameParsed(FrameParsed {
3141 stream_id,
3142 length: Some(payload_len),
3143 frame,
3144 ..Default::default()
3145 });
3146
3147 q.add_event_data_now(ev_data).ok();
3148 });
3149
3150 let more_frames = !conn.stream_finished(stream_id);
3151
3152 return Ok((stream_id, Event::Headers {
3153 list: headers,
3154 more_frames,
3155 }));
3156 },
3157
3158 frame::Frame::Data { .. } => {
3159 },
3161
3162 frame::Frame::GoAway { id } => {
3163 if !self.is_server && id % 4 != 0 {
3164 conn.close(
3165 true,
3166 Error::FrameUnexpected.to_wire(),
3167 b"GOAWAY received with ID of non-request stream",
3168 )?;
3169
3170 return Err(Error::IdError);
3171 }
3172
3173 if let Some(received_id) = self.peer_goaway_id {
3174 if id > received_id {
3175 conn.close(
3176 true,
3177 Error::IdError.to_wire(),
3178 b"GOAWAY received with ID larger than previously received",
3179 )?;
3180
3181 return Err(Error::IdError);
3182 }
3183 }
3184
3185 self.peer_goaway_id = Some(id);
3186
3187 return Ok((id, Event::GoAway));
3188 },
3189
3190 frame::Frame::MaxPushId { push_id } => {
3191 if !self.is_server {
3192 conn.close(
3193 true,
3194 Error::FrameUnexpected.to_wire(),
3195 b"MAX_PUSH_ID received by client",
3196 )?;
3197
3198 return Err(Error::FrameUnexpected);
3199 }
3200
3201 if push_id < self.max_push_id {
3202 conn.close(
3203 true,
3204 Error::IdError.to_wire(),
3205 b"MAX_PUSH_ID reduced limit",
3206 )?;
3207
3208 return Err(Error::IdError);
3209 }
3210
3211 self.max_push_id = push_id;
3212 },
3213
3214 frame::Frame::PushPromise { .. } => {
3215 if self.is_server {
3216 conn.close(
3217 true,
3218 Error::FrameUnexpected.to_wire(),
3219 b"PUSH_PROMISE received by server",
3220 )?;
3221
3222 return Err(Error::FrameUnexpected);
3223 }
3224
3225 if !stream_id.is_multiple_of(4) {
3226 conn.close(
3227 true,
3228 Error::FrameUnexpected.to_wire(),
3229 b"PUSH_PROMISE received on non-request stream",
3230 )?;
3231
3232 return Err(Error::FrameUnexpected);
3233 }
3234
3235 },
3237
3238 frame::Frame::CancelPush { .. } => {
3239 },
3241
3242 frame::Frame::PriorityUpdateRequest {
3243 prioritized_element_id,
3244 priority_field_value,
3245 } => {
3246 if !self.is_server {
3247 conn.close(
3248 true,
3249 Error::FrameUnexpected.to_wire(),
3250 b"PRIORITY_UPDATE received by client",
3251 )?;
3252
3253 return Err(Error::FrameUnexpected);
3254 }
3255
3256 if prioritized_element_id % 4 != 0 {
3257 conn.close(
3258 true,
3259 Error::FrameUnexpected.to_wire(),
3260 b"PRIORITY_UPDATE for request stream type with wrong ID",
3261 )?;
3262
3263 return Err(Error::FrameUnexpected);
3264 }
3265
3266 if prioritized_element_id > conn.streams.max_streams_bidi() * 4 {
3267 conn.close(
3268 true,
3269 Error::IdError.to_wire(),
3270 b"PRIORITY_UPDATE for request stream beyond max streams limit",
3271 )?;
3272
3273 return Err(Error::IdError);
3274 }
3275
3276 if conn.stream_closed(prioritized_element_id) {
3281 return Err(Error::Done);
3282 }
3283
3284 let stream = self
3286 .streams
3287 .entry(prioritized_element_id)
3288 .or_insert_with(|| {
3289 <stream::Stream>::new(
3290 prioritized_element_id,
3291 false,
3292 self.local_settings.max_field_section_size.unwrap_or(
3293 SETTINGS_MAX_FIELD_SECTION_SIZE_DEFAULT,
3294 ),
3295 self.max_priority_update_size,
3296 )
3297 });
3298
3299 let had_priority_update = stream.has_last_priority_update();
3300 stream.set_last_priority_update(Some(priority_field_value));
3301
3302 if !had_priority_update {
3305 return Ok((prioritized_element_id, Event::PriorityUpdate));
3306 } else {
3307 return Err(Error::Done);
3308 }
3309 },
3310
3311 frame::Frame::PriorityUpdatePush {
3312 prioritized_element_id,
3313 ..
3314 } => {
3315 if !self.is_server {
3316 conn.close(
3317 true,
3318 Error::FrameUnexpected.to_wire(),
3319 b"PRIORITY_UPDATE received by client",
3320 )?;
3321
3322 return Err(Error::FrameUnexpected);
3323 }
3324
3325 if prioritized_element_id % 3 != 0 {
3326 conn.close(
3327 true,
3328 Error::FrameUnexpected.to_wire(),
3329 b"PRIORITY_UPDATE for push stream type with wrong ID",
3330 )?;
3331
3332 return Err(Error::FrameUnexpected);
3333 }
3334
3335 },
3337
3338 frame::Frame::Unknown { .. } => (),
3339 }
3340
3341 Err(Error::Done)
3342 }
3343
3344 #[inline]
3346 pub fn stats(&self) -> Stats {
3347 Stats {
3348 qpack_encoder_stream_recv_bytes: self
3349 .peer_qpack_streams
3350 .encoder_stream_bytes,
3351 qpack_decoder_stream_recv_bytes: self
3352 .peer_qpack_streams
3353 .decoder_stream_bytes,
3354 }
3355 }
3356}
3357
3358pub fn grease_value() -> u64 {
3360 let n = super::rand::rand_u64_uniform(148_764_065_110_560_899);
3361 31 * n + 33
3362}
3363
3364#[doc(hidden)]
3365#[cfg(any(test, feature = "internal"))]
3366pub mod testing {
3367 use super::*;
3368
3369 use crate::test_utils;
3370 use crate::DefaultBufFactory;
3371
3372 pub struct Session<F = DefaultBufFactory>
3387 where
3388 F: BufFactory,
3389 {
3390 pub pipe: test_utils::Pipe<F>,
3391 pub client: Connection,
3392 pub server: Connection,
3393 }
3394
3395 impl Session {
3396 pub fn new() -> Result<Session> {
3397 Session::<DefaultBufFactory>::new_with_buf()
3398 }
3399
3400 pub fn with_configs(
3401 config: &mut crate::Config, h3_config: &Config,
3402 ) -> Result<Session> {
3403 Session::<DefaultBufFactory>::with_configs_and_buf(config, h3_config)
3404 }
3405
3406 pub fn default_configs() -> Result<(crate::Config, Config)> {
3407 fn path_relative_to_manifest_dir(path: &str) -> String {
3408 std::fs::canonicalize(
3409 std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join(path),
3410 )
3411 .unwrap()
3412 .to_string_lossy()
3413 .into_owned()
3414 }
3415
3416 let mut config = crate::Config::new(crate::PROTOCOL_VERSION)?;
3417 config.load_cert_chain_from_pem_file(
3418 &path_relative_to_manifest_dir("examples/cert.crt"),
3419 )?;
3420 config.load_priv_key_from_pem_file(
3421 &path_relative_to_manifest_dir("examples/cert.key"),
3422 )?;
3423 config.set_application_protos(&[b"h3"])?;
3424 config.set_initial_max_data(1500);
3425 config.set_initial_max_stream_data_bidi_local(150);
3426 config.set_initial_max_stream_data_bidi_remote(150);
3427 config.set_initial_max_stream_data_uni(150);
3428 config.set_initial_max_streams_bidi(5);
3429 config.set_initial_max_streams_uni(5);
3430 config.verify_peer(false);
3431 config.enable_dgram(true, 3, 3);
3432 config.set_ack_delay_exponent(8);
3433
3434 let h3_config = Config::new()?;
3435 Ok((config, h3_config))
3436 }
3437 }
3438
3439 impl<F: BufFactory> Session<F> {
3440 pub fn new_with_buf() -> Result<Session<F>> {
3441 let (mut config, h3_config) = Session::default_configs()?;
3442 Session::with_configs_and_buf(&mut config, &h3_config)
3443 }
3444
3445 pub fn with_configs_and_buf(
3446 config: &mut crate::Config, h3_config: &Config,
3447 ) -> Result<Session<F>> {
3448 let pipe = test_utils::Pipe::with_config_and_buf(config)?;
3449 let client_dgram = pipe.client.dgram_enabled();
3450 let server_dgram = pipe.server.dgram_enabled();
3451 Ok(Session {
3452 pipe,
3453 client: Connection::new(h3_config, false, client_dgram)?,
3454 server: Connection::new(h3_config, true, server_dgram)?,
3455 })
3456 }
3457
3458 pub fn handshake(&mut self) -> Result<()> {
3460 self.pipe.handshake()?;
3461
3462 self.client.send_settings(&mut self.pipe.client)?;
3464 self.pipe.advance().ok();
3465
3466 self.client
3467 .open_qpack_encoder_stream(&mut self.pipe.client)?;
3468 self.pipe.advance().ok();
3469
3470 self.client
3471 .open_qpack_decoder_stream(&mut self.pipe.client)?;
3472 self.pipe.advance().ok();
3473
3474 if self.pipe.client.grease {
3475 self.client.open_grease_stream(&mut self.pipe.client)?;
3476 }
3477
3478 self.pipe.advance().ok();
3479
3480 self.server.send_settings(&mut self.pipe.server)?;
3482 self.pipe.advance().ok();
3483
3484 self.server
3485 .open_qpack_encoder_stream(&mut self.pipe.server)?;
3486 self.pipe.advance().ok();
3487
3488 self.server
3489 .open_qpack_decoder_stream(&mut self.pipe.server)?;
3490 self.pipe.advance().ok();
3491
3492 if self.pipe.server.grease {
3493 self.server.open_grease_stream(&mut self.pipe.server)?;
3494 }
3495
3496 self.advance().ok();
3497
3498 while self.client.poll(&mut self.pipe.client).is_ok() {
3499 }
3501
3502 while self.server.poll(&mut self.pipe.server).is_ok() {
3503 }
3505
3506 Ok(())
3507 }
3508
3509 pub fn advance(&mut self) -> crate::Result<()> {
3511 self.pipe.advance()
3512 }
3513
3514 pub fn poll_client(&mut self) -> Result<(u64, Event)> {
3516 self.client.poll(&mut self.pipe.client)
3517 }
3518
3519 pub fn poll_server(&mut self) -> Result<(u64, Event)> {
3521 self.server.poll(&mut self.pipe.server)
3522 }
3523
3524 pub fn send_request(&mut self, fin: bool) -> Result<(u64, Vec<Header>)> {
3528 let req = vec![
3529 Header::new(b":method", b"GET"),
3530 Header::new(b":scheme", b"https"),
3531 Header::new(b":authority", b"quic.tech"),
3532 Header::new(b":path", b"/test"),
3533 Header::new(b"user-agent", b"quiche-test"),
3534 ];
3535
3536 let stream =
3537 self.client.send_request(&mut self.pipe.client, &req, fin)?;
3538
3539 self.advance().ok();
3540
3541 Ok((stream, req))
3542 }
3543
3544 pub fn send_response(
3548 &mut self, stream: u64, fin: bool,
3549 ) -> Result<Vec<Header>> {
3550 let resp = vec![
3551 Header::new(b":status", b"200"),
3552 Header::new(b"server", b"quiche-test"),
3553 ];
3554
3555 self.server.send_response(
3556 &mut self.pipe.server,
3557 stream,
3558 &resp,
3559 fin,
3560 )?;
3561
3562 self.advance().ok();
3563
3564 Ok(resp)
3565 }
3566
3567 pub fn send_body_client(
3571 &mut self, stream: u64, fin: bool,
3572 ) -> Result<Vec<u8>> {
3573 let bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3574
3575 self.client
3576 .send_body(&mut self.pipe.client, stream, &bytes, fin)?;
3577
3578 self.advance().ok();
3579
3580 Ok(bytes)
3581 }
3582
3583 pub fn recv_body_client(
3587 &mut self, stream: u64, buf: &mut [u8],
3588 ) -> Result<usize> {
3589 self.client.recv_body(&mut self.pipe.client, stream, buf)
3590 }
3591
3592 pub fn recv_body_buf_client<B: bytes::BufMut>(
3596 &mut self, stream: u64, buf: B,
3597 ) -> Result<usize> {
3598 self.client
3599 .recv_body_buf(&mut self.pipe.client, stream, buf)
3600 }
3601
3602 pub fn send_body_server(
3606 &mut self, stream: u64, fin: bool,
3607 ) -> Result<Vec<u8>> {
3608 let bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3609
3610 self.server
3611 .send_body(&mut self.pipe.server, stream, &bytes, fin)?;
3612
3613 self.advance().ok();
3614
3615 Ok(bytes)
3616 }
3617
3618 pub fn recv_body_server(
3622 &mut self, stream: u64, buf: &mut [u8],
3623 ) -> Result<usize> {
3624 self.server.recv_body(&mut self.pipe.server, stream, buf)
3625 }
3626
3627 pub fn recv_body_buf_server<B: bytes::BufMut>(
3631 &mut self, stream: u64, buf: B,
3632 ) -> Result<usize> {
3633 self.server
3634 .recv_body_buf(&mut self.pipe.server, stream, buf)
3635 }
3636
3637 pub fn send_frame_client(
3639 &mut self, frame: frame::Frame, stream_id: u64, fin: bool,
3640 ) -> Result<()> {
3641 let mut d = [42; 65535];
3642
3643 let mut b = octets::OctetsMut::with_slice(&mut d);
3644
3645 frame.to_bytes(&mut b)?;
3646
3647 let off = b.off();
3648 self.pipe.client.stream_send(stream_id, &d[..off], fin)?;
3649
3650 self.advance().ok();
3651
3652 Ok(())
3653 }
3654
3655 pub fn send_dgram_client(&mut self, flow_id: u64) -> Result<Vec<u8>> {
3659 let bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3660 let len = octets::varint_len(flow_id) + bytes.len();
3661 let mut d = vec![0; len];
3662 let mut b = octets::OctetsMut::with_slice(&mut d);
3663
3664 b.put_varint(flow_id)?;
3665 b.put_bytes(&bytes)?;
3666
3667 self.pipe.client.dgram_send(&d)?;
3668
3669 self.advance().ok();
3670
3671 Ok(bytes)
3672 }
3673
3674 pub fn recv_dgram_client(
3679 &mut self, buf: &mut [u8],
3680 ) -> Result<(usize, u64, usize)> {
3681 let len = self.pipe.client.dgram_recv(buf)?;
3682 let mut b = octets::Octets::with_slice(buf);
3683 let flow_id = b.get_varint()?;
3684
3685 Ok((len, flow_id, b.off()))
3686 }
3687
3688 pub fn send_dgram_server(&mut self, flow_id: u64) -> Result<Vec<u8>> {
3692 let bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
3693 let len = octets::varint_len(flow_id) + bytes.len();
3694 let mut d = vec![0; len];
3695 let mut b = octets::OctetsMut::with_slice(&mut d);
3696
3697 b.put_varint(flow_id)?;
3698 b.put_bytes(&bytes)?;
3699
3700 self.pipe.server.dgram_send(&d)?;
3701
3702 self.advance().ok();
3703
3704 Ok(bytes)
3705 }
3706
3707 pub fn recv_dgram_server(
3712 &mut self, buf: &mut [u8],
3713 ) -> Result<(usize, u64, usize)> {
3714 let len = self.pipe.server.dgram_recv(buf)?;
3715 let mut b = octets::Octets::with_slice(buf);
3716 let flow_id = b.get_varint()?;
3717
3718 Ok((len, flow_id, b.off()))
3719 }
3720
3721 pub fn send_frame_server(
3723 &mut self, frame: frame::Frame, stream_id: u64, fin: bool,
3724 ) -> Result<()> {
3725 let mut d = [42; 65535];
3726
3727 let mut b = octets::OctetsMut::with_slice(&mut d);
3728
3729 frame.to_bytes(&mut b)?;
3730
3731 let off = b.off();
3732 self.pipe.server.stream_send(stream_id, &d[..off], fin)?;
3733
3734 self.advance().ok();
3735
3736 Ok(())
3737 }
3738
3739 pub fn send_arbitrary_stream_data_client(
3741 &mut self, data: &[u8], stream_id: u64, fin: bool,
3742 ) -> Result<()> {
3743 self.pipe.client.stream_send(stream_id, data, fin)?;
3744
3745 self.advance().ok();
3746
3747 Ok(())
3748 }
3749
3750 pub fn send_arbitrary_stream_data_server(
3752 &mut self, data: &[u8], stream_id: u64, fin: bool,
3753 ) -> Result<()> {
3754 self.pipe.server.stream_send(stream_id, data, fin)?;
3755
3756 self.advance().ok();
3757
3758 Ok(())
3759 }
3760 }
3761}
3762
3763#[cfg(test)]
3764mod tests {
3765 use bytes::BufMut as _;
3766
3767 use super::*;
3768
3769 use super::testing::*;
3770
3771 #[test]
3772 fn grease_value_in_varint_limit() {
3774 assert!(grease_value() < 2u64.pow(62) - 1);
3775 }
3776
3777 #[test]
3778 fn h3_handshake_0rtt() {
3779 let mut buf = [0; 65535];
3780
3781 let mut config =
3784 crate::test_utils::config_no_pq(crate::PROTOCOL_VERSION).unwrap();
3785 config
3786 .load_cert_chain_from_pem_file("examples/cert.crt")
3787 .unwrap();
3788 config
3789 .load_priv_key_from_pem_file("examples/cert.key")
3790 .unwrap();
3791 config
3792 .set_application_protos(&[b"proto1", b"proto2"])
3793 .unwrap();
3794 config.set_initial_max_data(30);
3795 config.set_initial_max_stream_data_bidi_local(15);
3796 config.set_initial_max_stream_data_bidi_remote(15);
3797 config.set_initial_max_stream_data_uni(15);
3798 config.set_initial_max_streams_bidi(3);
3799 config.set_initial_max_streams_uni(3);
3800 config.enable_early_data();
3801 config.verify_peer(false);
3802
3803 let h3_config = Config::new().unwrap();
3804
3805 let mut pipe = crate::test_utils::Pipe::with_config(&mut config).unwrap();
3807 assert_eq!(pipe.handshake(), Ok(()));
3808
3809 let session = pipe.client.session().unwrap();
3811
3812 let mut pipe = crate::test_utils::Pipe::with_config(&mut config).unwrap();
3814 assert_eq!(pipe.client.set_session(session), Ok(()));
3815
3816 assert!(matches!(
3819 Connection::with_transport(&mut pipe.client, &h3_config),
3820 Err(Error::InternalError)
3821 ));
3822
3823 let (len, _) = pipe.client.send(&mut buf).unwrap();
3825
3826 assert!(Connection::with_transport(&mut pipe.client, &h3_config).is_ok());
3828 assert_eq!(pipe.server_recv(&mut buf[..len]), Ok(len));
3829
3830 let pkt_type = crate::packet::Type::ZeroRTT;
3832
3833 let frames = [crate::frame::Frame::Stream {
3834 stream_id: 6,
3835 data: <crate::range_buf::RangeBuf>::from(b"aaaaa", 0, true),
3836 }];
3837
3838 assert_eq!(
3839 pipe.send_pkt_to_server(pkt_type, &frames, &mut buf),
3840 Ok(1200)
3841 );
3842
3843 assert_eq!(pipe.server.undecryptable_pkts.len(), 0);
3844
3845 let mut r = pipe.server.readable();
3847 assert_eq!(r.next(), Some(6));
3848 assert_eq!(r.next(), None);
3849
3850 let mut b = [0; 15];
3851 assert_eq!(pipe.server.stream_recv(6, &mut b), Ok((5, true)));
3852 assert_eq!(&b[..5], b"aaaaa");
3853 }
3854
3855 #[test]
3856 fn request_no_body_response_no_body() {
3858 let mut s = Session::new().unwrap();
3859 s.handshake().unwrap();
3860
3861 let (stream, req) = s.send_request(true).unwrap();
3862
3863 assert_eq!(stream, 0);
3864
3865 let ev_headers = Event::Headers {
3866 list: req,
3867 more_frames: false,
3868 };
3869
3870 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
3871 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
3872
3873 let resp = s.send_response(stream, true).unwrap();
3874
3875 let ev_headers = Event::Headers {
3876 list: resp,
3877 more_frames: false,
3878 };
3879
3880 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
3881 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
3882 assert_eq!(s.poll_client(), Err(Error::Done));
3883 }
3884
3885 #[test]
3886 fn request_no_body_response_one_chunk() {
3888 let mut s = Session::new().unwrap();
3889 s.handshake().unwrap();
3890
3891 let (stream, req) = s.send_request(true).unwrap();
3892 assert_eq!(stream, 0);
3893
3894 let ev_headers = Event::Headers {
3895 list: req,
3896 more_frames: false,
3897 };
3898
3899 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
3900
3901 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
3902
3903 let resp = s.send_response(stream, false).unwrap();
3904
3905 let body = s.send_body_server(stream, true).unwrap();
3906
3907 let mut recv_buf = vec![0; body.len()];
3908
3909 let ev_headers = Event::Headers {
3910 list: resp,
3911 more_frames: true,
3912 };
3913
3914 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
3915
3916 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
3917 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
3918
3919 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
3920 assert_eq!(s.poll_client(), Err(Error::Done));
3921 }
3922
3923 #[test]
3924 fn request_no_body_response_many_chunks() {
3926 let mut s = Session::new().unwrap();
3927 s.handshake().unwrap();
3928
3929 let (stream, req) = s.send_request(true).unwrap();
3930
3931 let ev_headers = Event::Headers {
3932 list: req,
3933 more_frames: false,
3934 };
3935
3936 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
3937 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
3938
3939 let total_data_frames = 4;
3940
3941 let resp = s.send_response(stream, false).unwrap();
3942
3943 for _ in 0..total_data_frames - 1 {
3944 s.send_body_server(stream, false).unwrap();
3945 }
3946
3947 let body = s.send_body_server(stream, true).unwrap();
3948
3949 let mut recv_buf = vec![0; body.len()];
3950
3951 let ev_headers = Event::Headers {
3952 list: resp,
3953 more_frames: true,
3954 };
3955
3956 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
3957 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
3958 assert_eq!(s.poll_client(), Err(Error::Done));
3959
3960 for _ in 0..total_data_frames {
3961 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
3962 }
3963
3964 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
3965 assert_eq!(s.poll_client(), Err(Error::Done));
3966 }
3967
3968 #[test]
3969 fn request_no_body_response_many_chunks_with_buf() {
3971 let (mut config, h3_config) = Session::default_configs().unwrap();
3972 config.set_initial_congestion_window_packets(100);
3974 config.set_initial_max_data(200_000);
3975 config.set_initial_max_stream_data_bidi_local(200_000);
3976 config.set_initial_max_stream_data_bidi_remote(200_000);
3977 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
3978 s.handshake().unwrap();
3979
3980 let (stream, req) = s.send_request(true).unwrap();
3981
3982 let ev_headers = Event::Headers {
3983 list: req,
3984 more_frames: false,
3985 };
3986
3987 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
3988 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
3989
3990 let total_data_frames = 4;
3991
3992 let data = vec![0xab_u8; 16 * 1024];
3994
3995 let resp = s.send_response(stream, false).unwrap();
3996
3997 for _ in 0..total_data_frames - 1 {
3998 assert_eq!(
3999 s.server.send_body(&mut s.pipe.server, stream, &data, false),
4000 Ok(data.len())
4001 );
4002 s.advance().ok();
4003 }
4004
4005 s.server
4006 .send_body(&mut s.pipe.server, stream, &data, true)
4007 .unwrap();
4008 s.advance().ok();
4009
4010 let ev_headers = Event::Headers {
4011 list: resp,
4012 more_frames: true,
4013 };
4014
4015 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4016 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
4017 assert_eq!(s.poll_client(), Err(Error::Done));
4018
4019 let how_much_to_read_per_call = data.len() * 2 / 3;
4022 let mut remaining_to_read = total_data_frames * data.len();
4023 let mut recv_buf = Vec::new().limit(how_much_to_read_per_call);
4024 assert_eq!(
4025 s.recv_body_buf_client(stream, &mut recv_buf),
4026 Ok(how_much_to_read_per_call)
4027 );
4028 remaining_to_read -= how_much_to_read_per_call;
4029 assert_eq!(recv_buf.get_ref().len(), how_much_to_read_per_call);
4030
4031 while remaining_to_read > 0 {
4032 recv_buf.set_limit(data.len());
4034 let expected = std::cmp::min(data.len(), remaining_to_read);
4037 assert_eq!(
4038 s.recv_body_buf_client(stream, &mut recv_buf),
4039 Ok(expected)
4040 );
4041 remaining_to_read -= expected;
4042 }
4043 assert_eq!(recv_buf.get_ref().len(), total_data_frames * data.len());
4045
4046 assert_eq!(
4048 s.recv_body_buf_client(stream, &mut recv_buf),
4049 Err(Error::Done)
4050 );
4051
4052 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4053 assert_eq!(s.poll_client(), Err(Error::Done));
4054 }
4055
4056 #[test]
4057 fn request_one_chunk_response_no_body() {
4059 let mut s = Session::new().unwrap();
4060 s.handshake().unwrap();
4061
4062 let (stream, req) = s.send_request(false).unwrap();
4063
4064 let body = s.send_body_client(stream, true).unwrap();
4065
4066 let mut recv_buf = vec![0; body.len()];
4067
4068 let ev_headers = Event::Headers {
4069 list: req,
4070 more_frames: true,
4071 };
4072
4073 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4074
4075 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
4076 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
4077
4078 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4079
4080 let resp = s.send_response(stream, true).unwrap();
4081
4082 let ev_headers = Event::Headers {
4083 list: resp,
4084 more_frames: false,
4085 };
4086
4087 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4088 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4089 }
4090
4091 #[test]
4092 fn request_many_chunks_response_no_body() {
4094 let mut s = Session::new().unwrap();
4095 s.handshake().unwrap();
4096
4097 let (stream, req) = s.send_request(false).unwrap();
4098
4099 let total_data_frames = 4;
4100
4101 for _ in 0..total_data_frames - 1 {
4102 s.send_body_client(stream, false).unwrap();
4103 }
4104
4105 let body = s.send_body_client(stream, true).unwrap();
4106
4107 let mut recv_buf = vec![0; body.len()];
4108
4109 let ev_headers = Event::Headers {
4110 list: req,
4111 more_frames: true,
4112 };
4113
4114 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4115 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
4116 assert_eq!(s.poll_server(), Err(Error::Done));
4117
4118 for _ in 0..total_data_frames {
4119 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
4120 }
4121
4122 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4123
4124 let resp = s.send_response(stream, true).unwrap();
4125
4126 let ev_headers = Event::Headers {
4127 list: resp,
4128 more_frames: false,
4129 };
4130
4131 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4132 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4133 }
4134
4135 #[test]
4136 fn many_requests_many_chunks_response_one_chunk() {
4139 let mut s = Session::new().unwrap();
4140 s.handshake().unwrap();
4141
4142 let mut reqs = Vec::new();
4143
4144 let (stream1, req1) = s.send_request(false).unwrap();
4145 assert_eq!(stream1, 0);
4146 reqs.push(req1);
4147
4148 let (stream2, req2) = s.send_request(false).unwrap();
4149 assert_eq!(stream2, 4);
4150 reqs.push(req2);
4151
4152 let (stream3, req3) = s.send_request(false).unwrap();
4153 assert_eq!(stream3, 8);
4154 reqs.push(req3);
4155
4156 let body = s.send_body_client(stream1, false).unwrap();
4157 s.send_body_client(stream2, false).unwrap();
4158 s.send_body_client(stream3, false).unwrap();
4159
4160 let mut recv_buf = vec![0; body.len()];
4161
4162 s.send_body_client(stream3, true).unwrap();
4165 s.send_body_client(stream2, true).unwrap();
4166 s.send_body_client(stream1, true).unwrap();
4167
4168 let (_, ev) = s.poll_server().unwrap();
4169 let ev_headers = Event::Headers {
4170 list: reqs[0].clone(),
4171 more_frames: true,
4172 };
4173 assert_eq!(ev, ev_headers);
4174
4175 let (_, ev) = s.poll_server().unwrap();
4176 let ev_headers = Event::Headers {
4177 list: reqs[1].clone(),
4178 more_frames: true,
4179 };
4180 assert_eq!(ev, ev_headers);
4181
4182 let (_, ev) = s.poll_server().unwrap();
4183 let ev_headers = Event::Headers {
4184 list: reqs[2].clone(),
4185 more_frames: true,
4186 };
4187 assert_eq!(ev, ev_headers);
4188
4189 assert_eq!(s.poll_server(), Ok((0, Event::Data)));
4190 assert_eq!(s.recv_body_server(0, &mut recv_buf), Ok(body.len()));
4191 assert_eq!(s.poll_client(), Err(Error::Done));
4192 assert_eq!(s.recv_body_server(0, &mut recv_buf), Ok(body.len()));
4193 assert_eq!(s.poll_server(), Ok((0, Event::Finished)));
4194
4195 assert_eq!(s.poll_server(), Ok((4, Event::Data)));
4196 assert_eq!(s.recv_body_server(4, &mut recv_buf), Ok(body.len()));
4197 assert_eq!(s.poll_client(), Err(Error::Done));
4198 assert_eq!(s.recv_body_server(4, &mut recv_buf), Ok(body.len()));
4199 assert_eq!(s.poll_server(), Ok((4, Event::Finished)));
4200
4201 assert_eq!(s.poll_server(), Ok((8, Event::Data)));
4202 assert_eq!(s.recv_body_server(8, &mut recv_buf), Ok(body.len()));
4203 assert_eq!(s.poll_client(), Err(Error::Done));
4204 assert_eq!(s.recv_body_server(8, &mut recv_buf), Ok(body.len()));
4205 assert_eq!(s.poll_server(), Ok((8, Event::Finished)));
4206
4207 assert_eq!(s.poll_server(), Err(Error::Done));
4208
4209 let mut resps = Vec::new();
4210
4211 let resp1 = s.send_response(stream1, true).unwrap();
4212 resps.push(resp1);
4213
4214 let resp2 = s.send_response(stream2, true).unwrap();
4215 resps.push(resp2);
4216
4217 let resp3 = s.send_response(stream3, true).unwrap();
4218 resps.push(resp3);
4219
4220 for _ in 0..resps.len() {
4221 let (stream, ev) = s.poll_client().unwrap();
4222 let ev_headers = Event::Headers {
4223 list: resps[(stream / 4) as usize].clone(),
4224 more_frames: false,
4225 };
4226 assert_eq!(ev, ev_headers);
4227 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4228 }
4229
4230 assert_eq!(s.poll_client(), Err(Error::Done));
4231 }
4232
4233 #[test]
4234 fn request_no_body_response_one_chunk_empty_fin() {
4237 let mut s = Session::new().unwrap();
4238 s.handshake().unwrap();
4239
4240 let (stream, req) = s.send_request(true).unwrap();
4241
4242 let ev_headers = Event::Headers {
4243 list: req,
4244 more_frames: false,
4245 };
4246
4247 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4248 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4249
4250 let resp = s.send_response(stream, false).unwrap();
4251
4252 let body = s.send_body_server(stream, false).unwrap();
4253
4254 let mut recv_buf = vec![0; body.len()];
4255
4256 let ev_headers = Event::Headers {
4257 list: resp,
4258 more_frames: true,
4259 };
4260
4261 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4262
4263 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
4264 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
4265
4266 assert_eq!(s.pipe.server.stream_send(stream, &[], true), Ok(0));
4267 s.advance().ok();
4268
4269 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4270 assert_eq!(s.poll_client(), Err(Error::Done));
4271 }
4272
4273 #[test]
4274 fn request_no_body_response_no_body_with_grease() {
4277 let mut s = Session::new().unwrap();
4278 s.handshake().unwrap();
4279
4280 let (stream, req) = s.send_request(true).unwrap();
4281
4282 assert_eq!(stream, 0);
4283
4284 let ev_headers = Event::Headers {
4285 list: req,
4286 more_frames: false,
4287 };
4288
4289 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4290 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4291
4292 let resp = s.send_response(stream, false).unwrap();
4293
4294 let ev_headers = Event::Headers {
4295 list: resp,
4296 more_frames: true,
4297 };
4298
4299 let mut d = [42; 10];
4301 let mut b = octets::OctetsMut::with_slice(&mut d);
4302
4303 let frame_type = b.put_varint(148_764_065_110_560_899).unwrap();
4304 s.pipe.server.stream_send(0, frame_type, false).unwrap();
4305
4306 let frame_len = b.put_varint(10).unwrap();
4307 s.pipe.server.stream_send(0, frame_len, false).unwrap();
4308
4309 s.pipe.server.stream_send(0, &d, true).unwrap();
4310
4311 s.advance().ok();
4312
4313 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4314 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4315 assert_eq!(s.poll_client(), Err(Error::Done));
4316 }
4317
4318 #[test]
4319 fn body_response_before_headers() {
4321 let mut s = Session::new().unwrap();
4322 s.handshake().unwrap();
4323
4324 let (stream, req) = s.send_request(true).unwrap();
4325 assert_eq!(stream, 0);
4326
4327 let ev_headers = Event::Headers {
4328 list: req,
4329 more_frames: false,
4330 };
4331
4332 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4333
4334 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4335
4336 assert_eq!(
4337 s.send_body_server(stream, true),
4338 Err(Error::FrameUnexpected)
4339 );
4340
4341 assert_eq!(s.poll_client(), Err(Error::Done));
4342 }
4343
4344 #[test]
4345 fn send_body_invalid_client_stream() {
4348 let mut s = Session::new().unwrap();
4349 s.handshake().unwrap();
4350
4351 assert_eq!(s.send_body_client(0, true), Err(Error::FrameUnexpected));
4352
4353 assert_eq!(
4354 s.send_body_client(s.client.control_stream_id.unwrap(), true),
4355 Err(Error::FrameUnexpected)
4356 );
4357
4358 assert_eq!(
4359 s.send_body_client(
4360 s.client.local_qpack_streams.encoder_stream_id.unwrap(),
4361 true
4362 ),
4363 Err(Error::FrameUnexpected)
4364 );
4365
4366 assert_eq!(
4367 s.send_body_client(
4368 s.client.local_qpack_streams.decoder_stream_id.unwrap(),
4369 true
4370 ),
4371 Err(Error::FrameUnexpected)
4372 );
4373
4374 assert_eq!(
4375 s.send_body_client(s.client.peer_control_stream_id.unwrap(), true),
4376 Err(Error::FrameUnexpected)
4377 );
4378
4379 assert_eq!(
4380 s.send_body_client(
4381 s.client.peer_qpack_streams.encoder_stream_id.unwrap(),
4382 true
4383 ),
4384 Err(Error::FrameUnexpected)
4385 );
4386
4387 assert_eq!(
4388 s.send_body_client(
4389 s.client.peer_qpack_streams.decoder_stream_id.unwrap(),
4390 true
4391 ),
4392 Err(Error::FrameUnexpected)
4393 );
4394 }
4395
4396 #[test]
4397 fn send_body_invalid_server_stream() {
4400 let mut s = Session::new().unwrap();
4401 s.handshake().unwrap();
4402
4403 assert_eq!(s.send_body_server(0, true), Err(Error::FrameUnexpected));
4404
4405 assert_eq!(
4406 s.send_body_server(s.server.control_stream_id.unwrap(), true),
4407 Err(Error::FrameUnexpected)
4408 );
4409
4410 assert_eq!(
4411 s.send_body_server(
4412 s.server.local_qpack_streams.encoder_stream_id.unwrap(),
4413 true
4414 ),
4415 Err(Error::FrameUnexpected)
4416 );
4417
4418 assert_eq!(
4419 s.send_body_server(
4420 s.server.local_qpack_streams.decoder_stream_id.unwrap(),
4421 true
4422 ),
4423 Err(Error::FrameUnexpected)
4424 );
4425
4426 assert_eq!(
4427 s.send_body_server(s.server.peer_control_stream_id.unwrap(), true),
4428 Err(Error::FrameUnexpected)
4429 );
4430
4431 assert_eq!(
4432 s.send_body_server(
4433 s.server.peer_qpack_streams.encoder_stream_id.unwrap(),
4434 true
4435 ),
4436 Err(Error::FrameUnexpected)
4437 );
4438
4439 assert_eq!(
4440 s.send_body_server(
4441 s.server.peer_qpack_streams.decoder_stream_id.unwrap(),
4442 true
4443 ),
4444 Err(Error::FrameUnexpected)
4445 );
4446 }
4447
4448 #[test]
4449 fn trailers() {
4451 let mut s = Session::new().unwrap();
4452 s.handshake().unwrap();
4453
4454 let (stream, req) = s.send_request(false).unwrap();
4455
4456 let body = s.send_body_client(stream, false).unwrap();
4457
4458 let mut recv_buf = vec![0; body.len()];
4459
4460 let req_trailers = vec![Header::new(b"foo", b"bar")];
4461
4462 s.client
4463 .send_additional_headers(
4464 &mut s.pipe.client,
4465 stream,
4466 &req_trailers,
4467 true,
4468 true,
4469 )
4470 .unwrap();
4471
4472 s.advance().ok();
4473
4474 let ev_headers = Event::Headers {
4475 list: req,
4476 more_frames: true,
4477 };
4478
4479 let ev_trailers = Event::Headers {
4480 list: req_trailers,
4481 more_frames: false,
4482 };
4483
4484 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4485
4486 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
4487 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
4488
4489 assert_eq!(s.poll_server(), Ok((stream, ev_trailers)));
4490 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4491 }
4492
4493 #[test]
4494 fn informational_response() {
4496 let mut s = Session::new().unwrap();
4497 s.handshake().unwrap();
4498
4499 let (stream, req) = s.send_request(true).unwrap();
4500
4501 assert_eq!(stream, 0);
4502
4503 let ev_headers = Event::Headers {
4504 list: req,
4505 more_frames: false,
4506 };
4507
4508 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4509 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4510
4511 let info_resp = vec![
4512 Header::new(b":status", b"103"),
4513 Header::new(b"link", b"<https://example.com>; rel=\"preconnect\""),
4514 ];
4515
4516 let resp = vec![
4517 Header::new(b":status", b"200"),
4518 Header::new(b"server", b"quiche-test"),
4519 ];
4520
4521 s.server
4522 .send_response(&mut s.pipe.server, stream, &info_resp, false)
4523 .unwrap();
4524
4525 s.server
4526 .send_additional_headers(
4527 &mut s.pipe.server,
4528 stream,
4529 &resp,
4530 false,
4531 true,
4532 )
4533 .unwrap();
4534
4535 s.advance().ok();
4536
4537 let ev_info_headers = Event::Headers {
4538 list: info_resp,
4539 more_frames: true,
4540 };
4541
4542 let ev_headers = Event::Headers {
4543 list: resp,
4544 more_frames: false,
4545 };
4546
4547 assert_eq!(s.poll_client(), Ok((stream, ev_info_headers)));
4548 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
4549 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
4550 assert_eq!(s.poll_client(), Err(Error::Done));
4551 }
4552
4553 #[test]
4554 fn no_multiple_response() {
4557 let mut s = Session::new().unwrap();
4558 s.handshake().unwrap();
4559
4560 let (stream, req) = s.send_request(true).unwrap();
4561
4562 assert_eq!(stream, 0);
4563
4564 let ev_headers = Event::Headers {
4565 list: req,
4566 more_frames: false,
4567 };
4568
4569 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4570 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4571
4572 let info_resp = vec![
4573 Header::new(b":status", b"103"),
4574 Header::new(b"link", b"<https://example.com>; rel=\"preconnect\""),
4575 ];
4576
4577 let resp = vec![
4578 Header::new(b":status", b"200"),
4579 Header::new(b"server", b"quiche-test"),
4580 ];
4581
4582 s.server
4583 .send_response(&mut s.pipe.server, stream, &info_resp, false)
4584 .unwrap();
4585
4586 assert_eq!(
4587 Err(Error::FrameUnexpected),
4588 s.server
4589 .send_response(&mut s.pipe.server, stream, &resp, true)
4590 );
4591
4592 s.advance().ok();
4593
4594 let ev_info_headers = Event::Headers {
4595 list: info_resp,
4596 more_frames: true,
4597 };
4598
4599 assert_eq!(s.poll_client(), Ok((stream, ev_info_headers)));
4600 assert_eq!(s.poll_client(), Err(Error::Done));
4601 }
4602
4603 #[test]
4604 fn no_send_additional_before_initial_response() {
4606 let mut s = Session::new().unwrap();
4607 s.handshake().unwrap();
4608
4609 let (stream, req) = s.send_request(true).unwrap();
4610
4611 assert_eq!(stream, 0);
4612
4613 let ev_headers = Event::Headers {
4614 list: req,
4615 more_frames: false,
4616 };
4617
4618 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4619 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
4620
4621 let info_resp = vec![
4622 Header::new(b":status", b"103"),
4623 Header::new(b"link", b"<https://example.com>; rel=\"preconnect\""),
4624 ];
4625
4626 assert_eq!(
4627 Err(Error::FrameUnexpected),
4628 s.server.send_additional_headers(
4629 &mut s.pipe.server,
4630 stream,
4631 &info_resp,
4632 false,
4633 false
4634 )
4635 );
4636
4637 s.advance().ok();
4638
4639 assert_eq!(s.poll_client(), Err(Error::Done));
4640 }
4641
4642 #[test]
4643 fn additional_headers_before_data_client() {
4645 let mut s = Session::new().unwrap();
4646 s.handshake().unwrap();
4647
4648 let (stream, req) = s.send_request(false).unwrap();
4649
4650 let req_trailer = vec![Header::new(b"goodbye", b"world")];
4651
4652 assert_eq!(
4653 s.client.send_additional_headers(
4654 &mut s.pipe.client,
4655 stream,
4656 &req_trailer,
4657 true,
4658 false
4659 ),
4660 Ok(())
4661 );
4662
4663 s.advance().ok();
4664
4665 let ev_initial_headers = Event::Headers {
4666 list: req,
4667 more_frames: true,
4668 };
4669
4670 let ev_trailing_headers = Event::Headers {
4671 list: req_trailer,
4672 more_frames: true,
4673 };
4674
4675 assert_eq!(s.poll_server(), Ok((stream, ev_initial_headers)));
4676 assert_eq!(s.poll_server(), Ok((stream, ev_trailing_headers)));
4677 assert_eq!(s.poll_server(), Err(Error::Done));
4678 }
4679
4680 #[test]
4681 fn data_after_trailers_client() {
4683 let mut s = Session::new().unwrap();
4684 s.handshake().unwrap();
4685
4686 let (stream, req) = s.send_request(false).unwrap();
4687
4688 let body = s.send_body_client(stream, false).unwrap();
4689
4690 let mut recv_buf = vec![0; body.len()];
4691
4692 let req_trailers = vec![Header::new(b"foo", b"bar")];
4693
4694 s.client
4695 .send_additional_headers(
4696 &mut s.pipe.client,
4697 stream,
4698 &req_trailers,
4699 true,
4700 false,
4701 )
4702 .unwrap();
4703
4704 s.advance().ok();
4705
4706 s.send_frame_client(
4707 frame::Frame::Data {
4708 payload: vec![1, 2, 3, 4],
4709 },
4710 stream,
4711 true,
4712 )
4713 .unwrap();
4714
4715 let ev_headers = Event::Headers {
4716 list: req,
4717 more_frames: true,
4718 };
4719
4720 let ev_trailers = Event::Headers {
4721 list: req_trailers,
4722 more_frames: true,
4723 };
4724
4725 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4726 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
4727 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
4728 assert_eq!(s.poll_server(), Ok((stream, ev_trailers)));
4729 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
4730 }
4731
4732 #[test]
4733 fn max_push_id_from_client_good() {
4735 let mut s = Session::new().unwrap();
4736 s.handshake().unwrap();
4737
4738 s.send_frame_client(
4739 frame::Frame::MaxPushId { push_id: 1 },
4740 s.client.control_stream_id.unwrap(),
4741 false,
4742 )
4743 .unwrap();
4744
4745 assert_eq!(s.poll_server(), Err(Error::Done));
4746 }
4747
4748 #[test]
4749 fn max_push_id_from_client_bad_stream() {
4751 let mut s = Session::new().unwrap();
4752 s.handshake().unwrap();
4753
4754 let (stream, req) = s.send_request(false).unwrap();
4755
4756 s.send_frame_client(
4757 frame::Frame::MaxPushId { push_id: 2 },
4758 stream,
4759 false,
4760 )
4761 .unwrap();
4762
4763 let ev_headers = Event::Headers {
4764 list: req,
4765 more_frames: true,
4766 };
4767
4768 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4769 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
4770 }
4771
4772 #[test]
4773 fn max_push_id_from_client_limit_reduction() {
4776 let mut s = Session::new().unwrap();
4777 s.handshake().unwrap();
4778
4779 s.send_frame_client(
4780 frame::Frame::MaxPushId { push_id: 2 },
4781 s.client.control_stream_id.unwrap(),
4782 false,
4783 )
4784 .unwrap();
4785
4786 s.send_frame_client(
4787 frame::Frame::MaxPushId { push_id: 1 },
4788 s.client.control_stream_id.unwrap(),
4789 false,
4790 )
4791 .unwrap();
4792
4793 assert_eq!(s.poll_server(), Err(Error::IdError));
4794 }
4795
4796 #[test]
4797 fn max_push_id_from_server() {
4799 let mut s = Session::new().unwrap();
4800 s.handshake().unwrap();
4801
4802 s.send_frame_server(
4803 frame::Frame::MaxPushId { push_id: 1 },
4804 s.server.control_stream_id.unwrap(),
4805 false,
4806 )
4807 .unwrap();
4808
4809 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
4810 }
4811
4812 #[test]
4813 fn push_promise_from_client() {
4815 let mut s = Session::new().unwrap();
4816 s.handshake().unwrap();
4817
4818 let (stream, req) = s.send_request(false).unwrap();
4819
4820 let header_block = s.client.encode_header_block(&req).unwrap();
4821
4822 s.send_frame_client(
4823 frame::Frame::PushPromise {
4824 push_id: 1,
4825 header_block,
4826 },
4827 stream,
4828 false,
4829 )
4830 .unwrap();
4831
4832 let ev_headers = Event::Headers {
4833 list: req,
4834 more_frames: true,
4835 };
4836
4837 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4838 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
4839 }
4840
4841 #[test]
4842 fn push_stream_from_client() {
4844 let mut s = Session::new().unwrap();
4845 s.handshake().unwrap();
4846
4847 s.client
4848 .open_uni_stream(
4849 &mut s.pipe.client,
4850 stream::HTTP3_PUSH_STREAM_TYPE_ID,
4851 )
4852 .unwrap();
4853
4854 s.advance().ok();
4855
4856 assert_eq!(s.poll_server(), Err(Error::StreamCreationError));
4857 }
4858
4859 #[test]
4860 fn push_stream_from_server() {
4863 let mut s = Session::new().unwrap();
4864 s.handshake().unwrap();
4865
4866 s.server
4867 .open_uni_stream(
4868 &mut s.pipe.server,
4869 stream::HTTP3_PUSH_STREAM_TYPE_ID,
4870 )
4871 .unwrap();
4872
4873 s.advance().ok();
4874
4875 assert_eq!(s.poll_client(), Err(Error::StreamCreationError));
4876 }
4877
4878 #[test]
4879 fn cancel_push_from_client() {
4881 let mut s = Session::new().unwrap();
4882 s.handshake().unwrap();
4883
4884 s.send_frame_client(
4885 frame::Frame::CancelPush { push_id: 1 },
4886 s.client.control_stream_id.unwrap(),
4887 false,
4888 )
4889 .unwrap();
4890
4891 assert_eq!(s.poll_server(), Err(Error::Done));
4892 }
4893
4894 #[test]
4895 fn cancel_push_from_client_bad_stream() {
4897 let mut s = Session::new().unwrap();
4898 s.handshake().unwrap();
4899
4900 let (stream, req) = s.send_request(false).unwrap();
4901
4902 s.send_frame_client(
4903 frame::Frame::CancelPush { push_id: 2 },
4904 stream,
4905 false,
4906 )
4907 .unwrap();
4908
4909 let ev_headers = Event::Headers {
4910 list: req,
4911 more_frames: true,
4912 };
4913
4914 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
4915 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
4916 }
4917
4918 #[test]
4919 fn cancel_push_from_server() {
4921 let mut s = Session::new().unwrap();
4922 s.handshake().unwrap();
4923
4924 s.send_frame_server(
4925 frame::Frame::CancelPush { push_id: 1 },
4926 s.server.control_stream_id.unwrap(),
4927 false,
4928 )
4929 .unwrap();
4930
4931 assert_eq!(s.poll_client(), Err(Error::Done));
4932 }
4933
4934 #[test]
4935 fn goaway_from_client_good() {
4937 let mut s = Session::new().unwrap();
4938 s.handshake().unwrap();
4939
4940 s.client.send_goaway(&mut s.pipe.client, 100).unwrap();
4941
4942 s.advance().ok();
4943
4944 assert_eq!(s.poll_server(), Ok((0, Event::GoAway)));
4946 }
4947
4948 #[test]
4949 fn goaway_from_server_good() {
4951 let mut s = Session::new().unwrap();
4952 s.handshake().unwrap();
4953
4954 s.server.send_goaway(&mut s.pipe.server, 4000).unwrap();
4955
4956 s.advance().ok();
4957
4958 assert_eq!(s.poll_client(), Ok((4000, Event::GoAway)));
4959 }
4960
4961 #[test]
4962 fn client_request_after_goaway() {
4964 let mut s = Session::new().unwrap();
4965 s.handshake().unwrap();
4966
4967 s.server.send_goaway(&mut s.pipe.server, 4000).unwrap();
4968
4969 s.advance().ok();
4970
4971 assert_eq!(s.poll_client(), Ok((4000, Event::GoAway)));
4972
4973 assert_eq!(s.send_request(true), Err(Error::FrameUnexpected));
4974 }
4975
4976 #[test]
4977 fn goaway_from_server_invalid_id() {
4979 let mut s = Session::new().unwrap();
4980 s.handshake().unwrap();
4981
4982 s.send_frame_server(
4983 frame::Frame::GoAway { id: 1 },
4984 s.server.control_stream_id.unwrap(),
4985 false,
4986 )
4987 .unwrap();
4988
4989 assert_eq!(s.poll_client(), Err(Error::IdError));
4990 }
4991
4992 #[test]
4993 fn goaway_from_server_increase_id() {
4996 let mut s = Session::new().unwrap();
4997 s.handshake().unwrap();
4998
4999 s.send_frame_server(
5000 frame::Frame::GoAway { id: 0 },
5001 s.server.control_stream_id.unwrap(),
5002 false,
5003 )
5004 .unwrap();
5005
5006 s.send_frame_server(
5007 frame::Frame::GoAway { id: 4 },
5008 s.server.control_stream_id.unwrap(),
5009 false,
5010 )
5011 .unwrap();
5012
5013 assert_eq!(s.poll_client(), Ok((0, Event::GoAway)));
5014
5015 assert_eq!(s.poll_client(), Err(Error::IdError));
5016 }
5017
5018 #[test]
5019 #[cfg(feature = "sfv")]
5020 fn parse_priority_field_value() {
5021 assert_eq!(
5023 Ok(Priority::new(0, false)),
5024 Priority::try_from(b"u=0".as_slice())
5025 );
5026 assert_eq!(
5027 Ok(Priority::new(3, false)),
5028 Priority::try_from(b"u=3".as_slice())
5029 );
5030 assert_eq!(
5031 Ok(Priority::new(7, false)),
5032 Priority::try_from(b"u=7".as_slice())
5033 );
5034
5035 assert_eq!(
5036 Ok(Priority::new(0, true)),
5037 Priority::try_from(b"u=0, i".as_slice())
5038 );
5039 assert_eq!(
5040 Ok(Priority::new(3, true)),
5041 Priority::try_from(b"u=3, i".as_slice())
5042 );
5043 assert_eq!(
5044 Ok(Priority::new(7, true)),
5045 Priority::try_from(b"u=7, i".as_slice())
5046 );
5047
5048 assert_eq!(
5049 Ok(Priority::new(0, true)),
5050 Priority::try_from(b"u=0, i=?1".as_slice())
5051 );
5052 assert_eq!(
5053 Ok(Priority::new(3, true)),
5054 Priority::try_from(b"u=3, i=?1".as_slice())
5055 );
5056 assert_eq!(
5057 Ok(Priority::new(7, true)),
5058 Priority::try_from(b"u=7, i=?1".as_slice())
5059 );
5060
5061 assert_eq!(
5062 Ok(Priority::new(3, false)),
5063 Priority::try_from(b"".as_slice())
5064 );
5065
5066 assert_eq!(
5067 Ok(Priority::new(0, true)),
5068 Priority::try_from(b"u=0;foo, i;bar".as_slice())
5069 );
5070 assert_eq!(
5071 Ok(Priority::new(3, true)),
5072 Priority::try_from(b"u=3;hello, i;world".as_slice())
5073 );
5074 assert_eq!(
5075 Ok(Priority::new(7, true)),
5076 Priority::try_from(b"u=7;croeso, i;gymru".as_slice())
5077 );
5078
5079 assert_eq!(
5080 Ok(Priority::new(0, true)),
5081 Priority::try_from(b"u=0, i, spinaltap=11".as_slice())
5082 );
5083
5084 assert_eq!(Err(Error::Done), Priority::try_from(b"0".as_slice()));
5086 assert_eq!(
5087 Ok(Priority::new(7, false)),
5088 Priority::try_from(b"u=-1".as_slice())
5089 );
5090 assert_eq!(Err(Error::Done), Priority::try_from(b"u=0.2".as_slice()));
5091 assert_eq!(
5092 Ok(Priority::new(7, false)),
5093 Priority::try_from(b"u=100".as_slice())
5094 );
5095 assert_eq!(
5096 Err(Error::Done),
5097 Priority::try_from(b"u=3, i=true".as_slice())
5098 );
5099
5100 assert_eq!(Err(Error::Done), Priority::try_from(b"u=7, ".as_slice()));
5102 }
5103
5104 #[test]
5105 fn priority_update_request() {
5107 let mut s = Session::new().unwrap();
5108 s.handshake().unwrap();
5109
5110 s.client
5111 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5112 urgency: 3,
5113 incremental: false,
5114 })
5115 .unwrap();
5116 s.advance().ok();
5117
5118 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5119 assert_eq!(s.poll_server(), Err(Error::Done));
5120 }
5121
5122 #[test]
5123 fn priority_update_request_max_size_limit_default() {
5126 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
5127 config
5128 .load_cert_chain_from_pem_file("examples/cert.crt")
5129 .unwrap();
5130 config
5131 .load_priv_key_from_pem_file("examples/cert.key")
5132 .unwrap();
5133 config.set_application_protos(&[b"h3"]).unwrap();
5134 config.set_initial_max_data(1500);
5135 config.set_initial_max_stream_data_bidi_local(1500);
5136 config.set_initial_max_stream_data_bidi_remote(1500);
5137 config.set_initial_max_stream_data_uni(1500);
5138 config.set_initial_max_streams_bidi(5);
5139 config.set_initial_max_streams_uni(5);
5140 config.verify_peer(false);
5141
5142 let h3_config = Config::new().unwrap();
5143
5144 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
5145
5146 s.handshake().unwrap();
5147
5148 let mut d = vec![42; 600];
5149 let mut b = octets::OctetsMut::with_slice(&mut d);
5150
5151 let pu = frame::Frame::PriorityUpdateRequest {
5152 prioritized_element_id: 0,
5153 priority_field_value: vec![0; 512],
5154 };
5155
5156 pu.to_bytes(&mut b).unwrap();
5157
5158 s.pipe.client.stream_send(2, &d, true).unwrap();
5159
5160 s.advance().ok();
5161
5162 assert_eq!(s.poll_server(), Err(Error::ExcessiveLoad));
5163
5164 assert_eq!(
5165 s.pipe.server.local_error.as_ref().unwrap().error_code,
5166 Error::to_wire(Error::ExcessiveLoad)
5167 );
5168 }
5169
5170 #[test]
5171 fn priority_update_single_stream_rearm() {
5173 let mut s = Session::new().unwrap();
5174 s.handshake().unwrap();
5175
5176 s.client
5177 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5178 urgency: 3,
5179 incremental: false,
5180 })
5181 .unwrap();
5182 s.advance().ok();
5183
5184 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5185 assert_eq!(s.poll_server(), Err(Error::Done));
5186
5187 s.client
5188 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5189 urgency: 5,
5190 incremental: false,
5191 })
5192 .unwrap();
5193 s.advance().ok();
5194
5195 assert_eq!(s.poll_server(), Err(Error::Done));
5196
5197 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=5".to_vec()));
5200 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5201
5202 s.client
5203 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5204 urgency: 7,
5205 incremental: false,
5206 })
5207 .unwrap();
5208 s.advance().ok();
5209
5210 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5211 assert_eq!(s.poll_server(), Err(Error::Done));
5212
5213 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=7".to_vec()));
5214 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5215 }
5216
5217 #[test]
5218 fn priority_update_request_multiple_stream_arm_multiple_flights() {
5221 let mut s = Session::new().unwrap();
5222 s.handshake().unwrap();
5223
5224 s.client
5225 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5226 urgency: 3,
5227 incremental: false,
5228 })
5229 .unwrap();
5230 s.advance().ok();
5231
5232 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5233 assert_eq!(s.poll_server(), Err(Error::Done));
5234
5235 s.client
5236 .send_priority_update_for_request(&mut s.pipe.client, 4, &Priority {
5237 urgency: 1,
5238 incremental: false,
5239 })
5240 .unwrap();
5241 s.advance().ok();
5242
5243 assert_eq!(s.poll_server(), Ok((4, Event::PriorityUpdate)));
5244 assert_eq!(s.poll_server(), Err(Error::Done));
5245
5246 s.client
5247 .send_priority_update_for_request(&mut s.pipe.client, 8, &Priority {
5248 urgency: 2,
5249 incremental: false,
5250 })
5251 .unwrap();
5252 s.advance().ok();
5253
5254 assert_eq!(s.poll_server(), Ok((8, Event::PriorityUpdate)));
5255 assert_eq!(s.poll_server(), Err(Error::Done));
5256
5257 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=3".to_vec()));
5258 assert_eq!(s.server.take_last_priority_update(4), Ok(b"u=1".to_vec()));
5259 assert_eq!(s.server.take_last_priority_update(8), Ok(b"u=2".to_vec()));
5260 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5261 }
5262
5263 #[test]
5264 fn priority_update_request_multiple_stream_arm_single_flight() {
5267 let mut s = Session::new().unwrap();
5268 s.handshake().unwrap();
5269
5270 let mut d = [42; 65535];
5271
5272 let mut b = octets::OctetsMut::with_slice(&mut d);
5273
5274 let p1 = frame::Frame::PriorityUpdateRequest {
5275 prioritized_element_id: 0,
5276 priority_field_value: b"u=3".to_vec(),
5277 };
5278
5279 let p2 = frame::Frame::PriorityUpdateRequest {
5280 prioritized_element_id: 4,
5281 priority_field_value: b"u=3".to_vec(),
5282 };
5283
5284 let p3 = frame::Frame::PriorityUpdateRequest {
5285 prioritized_element_id: 8,
5286 priority_field_value: b"u=3".to_vec(),
5287 };
5288
5289 p1.to_bytes(&mut b).unwrap();
5290 p2.to_bytes(&mut b).unwrap();
5291 p3.to_bytes(&mut b).unwrap();
5292
5293 let off = b.off();
5294 s.pipe
5295 .client
5296 .stream_send(s.client.control_stream_id.unwrap(), &d[..off], false)
5297 .unwrap();
5298
5299 s.advance().ok();
5300
5301 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5302 assert_eq!(s.poll_server(), Ok((4, Event::PriorityUpdate)));
5303 assert_eq!(s.poll_server(), Ok((8, Event::PriorityUpdate)));
5304 assert_eq!(s.poll_server(), Err(Error::Done));
5305
5306 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=3".to_vec()));
5307 assert_eq!(s.server.take_last_priority_update(4), Ok(b"u=3".to_vec()));
5308 assert_eq!(s.server.take_last_priority_update(8), Ok(b"u=3".to_vec()));
5309
5310 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5311 }
5312
5313 #[test]
5314 fn priority_update_request_collected_completed() {
5317 let mut s = Session::new().unwrap();
5318 s.handshake().unwrap();
5319
5320 s.client
5321 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5322 urgency: 3,
5323 incremental: false,
5324 })
5325 .unwrap();
5326 s.advance().ok();
5327
5328 let (stream, req) = s.send_request(true).unwrap();
5329 let ev_headers = Event::Headers {
5330 list: req,
5331 more_frames: false,
5332 };
5333
5334 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5336 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
5337 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
5338 assert_eq!(s.poll_server(), Err(Error::Done));
5339
5340 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=3".to_vec()));
5341 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5342
5343 let resp = s.send_response(stream, true).unwrap();
5344
5345 let ev_headers = Event::Headers {
5346 list: resp,
5347 more_frames: false,
5348 };
5349
5350 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
5351 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
5352 assert_eq!(s.poll_client(), Err(Error::Done));
5353
5354 s.client
5356 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5357 urgency: 3,
5358 incremental: false,
5359 })
5360 .unwrap();
5361 s.advance().ok();
5362
5363 assert_eq!(s.poll_server(), Err(Error::Done));
5365 }
5366
5367 #[test]
5368 fn priority_update_request_after_h3_collection() {
5371 let mut s = Session::new().unwrap();
5372 s.handshake().unwrap();
5373
5374 let init_streams_server = s.server.streams.len();
5375
5376 let (stream, req) = s.send_request(true).unwrap();
5377 let ev_headers = Event::Headers {
5378 list: req,
5379 more_frames: false,
5380 };
5381
5382 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
5383 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
5384 assert_eq!(s.poll_server(), Err(Error::Done));
5385
5386 let resp = vec![
5387 Header::new(b":status", b"200"),
5388 Header::new(b"server", b"quiche-test"),
5389 ];
5390
5391 s.server
5392 .send_response(&mut s.pipe.server, stream, &resp, true)
5393 .unwrap();
5394
5395 assert_eq!(s.server.streams.len(), init_streams_server);
5399 assert!(s.pipe.server.stream_finished(stream));
5400 assert!(s.pipe.server.stream_closed(stream));
5401
5402 let stream_state = s.pipe.server.streams.get(stream).unwrap();
5403 assert!(stream_state.recv.is_fin());
5404 assert!(stream_state.send.is_fin());
5405 assert!(!s.pipe.server.streams.is_collected(stream));
5406
5407 s.client
5408 .send_priority_update_for_request(
5409 &mut s.pipe.client,
5410 stream,
5411 &Priority {
5412 urgency: 3,
5413 incremental: false,
5414 },
5415 )
5416 .unwrap();
5417
5418 let flight = crate::test_utils::emit_flight(&mut s.pipe.client).unwrap();
5419 crate::test_utils::process_flight(&mut s.pipe.server, flight).unwrap();
5420
5421 assert_eq!(s.poll_server(), Err(Error::Done));
5422 assert_eq!(s.server.streams.len(), init_streams_server);
5423 }
5424
5425 #[test]
5426 fn priority_update_request_collected_stopped() {
5429 let mut s = Session::new().unwrap();
5430 s.handshake().unwrap();
5431
5432 s.client
5433 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5434 urgency: 3,
5435 incremental: false,
5436 })
5437 .unwrap();
5438 s.advance().ok();
5439
5440 let (stream, req) = s.send_request(false).unwrap();
5441 let ev_headers = Event::Headers {
5442 list: req,
5443 more_frames: true,
5444 };
5445
5446 assert_eq!(s.poll_server(), Ok((0, Event::PriorityUpdate)));
5448 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
5449 assert_eq!(s.poll_server(), Err(Error::Done));
5450
5451 assert_eq!(s.server.take_last_priority_update(0), Ok(b"u=3".to_vec()));
5452 assert_eq!(s.server.take_last_priority_update(0), Err(Error::Done));
5453
5454 s.pipe
5455 .client
5456 .stream_shutdown(stream, crate::Shutdown::Write, 0x100)
5457 .unwrap();
5458 s.pipe
5459 .client
5460 .stream_shutdown(stream, crate::Shutdown::Read, 0x100)
5461 .unwrap();
5462
5463 s.advance().ok();
5464
5465 assert_eq!(s.poll_server(), Ok((0, Event::Reset(0x100))));
5466 assert_eq!(s.poll_server(), Err(Error::Done));
5467
5468 s.client
5470 .send_priority_update_for_request(&mut s.pipe.client, 0, &Priority {
5471 urgency: 3,
5472 incremental: false,
5473 })
5474 .unwrap();
5475 s.advance().ok();
5476
5477 assert_eq!(s.poll_server(), Err(Error::Done));
5479
5480 assert!(!s.pipe.server.streams.is_collected(0));
5482 assert_eq!(
5483 s.pipe.server.stream_capacity(0),
5484 Err(crate::Error::StreamStopped(0x100))
5485 );
5486 assert!(s.pipe.server.streams.is_collected(0));
5487 assert!(s.pipe.client.streams.is_collected(0));
5488 }
5489
5490 #[test]
5491 fn priority_update_push() {
5493 let mut s = Session::new().unwrap();
5494 s.handshake().unwrap();
5495
5496 s.send_frame_client(
5497 frame::Frame::PriorityUpdatePush {
5498 prioritized_element_id: 3,
5499 priority_field_value: b"u=3".to_vec(),
5500 },
5501 s.client.control_stream_id.unwrap(),
5502 false,
5503 )
5504 .unwrap();
5505
5506 assert_eq!(s.poll_server(), Err(Error::Done));
5507 }
5508
5509 #[test]
5510 fn priority_update_request_bad_stream() {
5513 let mut s = Session::new().unwrap();
5514 s.handshake().unwrap();
5515
5516 s.send_frame_client(
5517 frame::Frame::PriorityUpdateRequest {
5518 prioritized_element_id: 5,
5519 priority_field_value: b"u=3".to_vec(),
5520 },
5521 s.client.control_stream_id.unwrap(),
5522 false,
5523 )
5524 .unwrap();
5525
5526 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
5527 }
5528
5529 #[test]
5530 fn priority_update_push_bad_stream() {
5533 let mut s = Session::new().unwrap();
5534 s.handshake().unwrap();
5535
5536 s.send_frame_client(
5537 frame::Frame::PriorityUpdatePush {
5538 prioritized_element_id: 5,
5539 priority_field_value: b"u=3".to_vec(),
5540 },
5541 s.client.control_stream_id.unwrap(),
5542 false,
5543 )
5544 .unwrap();
5545
5546 assert_eq!(s.poll_server(), Err(Error::FrameUnexpected));
5547 }
5548
5549 #[test]
5550 fn priority_update_request_from_server() {
5552 let mut s = Session::new().unwrap();
5553 s.handshake().unwrap();
5554
5555 s.send_frame_server(
5556 frame::Frame::PriorityUpdateRequest {
5557 prioritized_element_id: 0,
5558 priority_field_value: b"u=3".to_vec(),
5559 },
5560 s.server.control_stream_id.unwrap(),
5561 false,
5562 )
5563 .unwrap();
5564
5565 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
5566 }
5567
5568 #[test]
5569 fn priority_update_push_from_server() {
5571 let mut s = Session::new().unwrap();
5572 s.handshake().unwrap();
5573
5574 s.send_frame_server(
5575 frame::Frame::PriorityUpdatePush {
5576 prioritized_element_id: 0,
5577 priority_field_value: b"u=3".to_vec(),
5578 },
5579 s.server.control_stream_id.unwrap(),
5580 false,
5581 )
5582 .unwrap();
5583
5584 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
5585 }
5586
5587 #[test]
5588 fn uni_stream_local_counting() {
5590 let config = Config::new().unwrap();
5591
5592 let h3_cln = Connection::new(&config, false, false).unwrap();
5593 assert_eq!(h3_cln.next_uni_stream_id, 2);
5594
5595 let h3_srv = Connection::new(&config, true, false).unwrap();
5596 assert_eq!(h3_srv.next_uni_stream_id, 3);
5597 }
5598
5599 #[test]
5600 fn open_multiple_control_streams() {
5602 let mut s = Session::new().unwrap();
5603 s.handshake().unwrap();
5604
5605 let stream_id = s.client.next_uni_stream_id;
5606
5607 let mut d = [42; 8];
5608 let mut b = octets::OctetsMut::with_slice(&mut d);
5609
5610 s.pipe
5611 .client
5612 .stream_send(
5613 stream_id,
5614 b.put_varint(stream::HTTP3_CONTROL_STREAM_TYPE_ID).unwrap(),
5615 false,
5616 )
5617 .unwrap();
5618
5619 s.advance().ok();
5620
5621 assert_eq!(s.poll_server(), Err(Error::StreamCreationError));
5622 }
5623
5624 #[test]
5625 fn close_control_stream_after_type() {
5627 let mut s = Session::new().unwrap();
5628 s.handshake().unwrap();
5629
5630 s.pipe
5631 .client
5632 .stream_send(s.client.control_stream_id.unwrap(), &[], true)
5633 .unwrap();
5634
5635 s.advance().ok();
5636
5637 assert_eq!(
5638 Err(Error::ClosedCriticalStream),
5639 s.server.poll(&mut s.pipe.server)
5640 );
5641 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5642 }
5643
5644 #[test]
5645 fn close_control_stream_after_frame() {
5648 let mut s = Session::new().unwrap();
5649 s.handshake().unwrap();
5650
5651 s.send_frame_client(
5652 frame::Frame::MaxPushId { push_id: 1 },
5653 s.client.control_stream_id.unwrap(),
5654 true,
5655 )
5656 .unwrap();
5657
5658 assert_eq!(
5659 Err(Error::ClosedCriticalStream),
5660 s.server.poll(&mut s.pipe.server)
5661 );
5662 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5663 }
5664
5665 #[test]
5666 fn reset_control_stream_after_type() {
5668 let mut s = Session::new().unwrap();
5669 s.handshake().unwrap();
5670
5671 s.pipe
5672 .client
5673 .stream_shutdown(
5674 s.client.control_stream_id.unwrap(),
5675 crate::Shutdown::Write,
5676 0,
5677 )
5678 .unwrap();
5679
5680 s.advance().ok();
5681
5682 assert_eq!(
5683 Err(Error::ClosedCriticalStream),
5684 s.server.poll(&mut s.pipe.server)
5685 );
5686 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5687 }
5688
5689 #[test]
5690 fn reset_control_stream_after_frame() {
5693 let mut s = Session::new().unwrap();
5694 s.handshake().unwrap();
5695
5696 s.send_frame_client(
5697 frame::Frame::MaxPushId { push_id: 1 },
5698 s.client.control_stream_id.unwrap(),
5699 false,
5700 )
5701 .unwrap();
5702
5703 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5704
5705 s.pipe
5706 .client
5707 .stream_shutdown(
5708 s.client.control_stream_id.unwrap(),
5709 crate::Shutdown::Write,
5710 0,
5711 )
5712 .unwrap();
5713
5714 s.advance().ok();
5715
5716 assert_eq!(
5717 Err(Error::ClosedCriticalStream),
5718 s.server.poll(&mut s.pipe.server)
5719 );
5720 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5721 }
5722
5723 #[test]
5724 fn close_qpack_stream_after_type() {
5726 let mut s = Session::new().unwrap();
5727 s.handshake().unwrap();
5728
5729 s.pipe
5730 .client
5731 .stream_send(
5732 s.client.local_qpack_streams.encoder_stream_id.unwrap(),
5733 &[],
5734 true,
5735 )
5736 .unwrap();
5737
5738 s.advance().ok();
5739
5740 assert_eq!(
5741 Err(Error::ClosedCriticalStream),
5742 s.server.poll(&mut s.pipe.server)
5743 );
5744 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5745 }
5746
5747 #[test]
5748 fn close_qpack_stream_after_data() {
5750 let mut s = Session::new().unwrap();
5751 s.handshake().unwrap();
5752
5753 let stream_id = s.client.local_qpack_streams.encoder_stream_id.unwrap();
5754 let d = [0; 1];
5755
5756 s.pipe.client.stream_send(stream_id, &d, false).unwrap();
5757 s.pipe.client.stream_send(stream_id, &d, true).unwrap();
5758
5759 s.advance().ok();
5760
5761 assert_eq!(
5762 Err(Error::ClosedCriticalStream),
5763 s.server.poll(&mut s.pipe.server)
5764 );
5765 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5766 }
5767
5768 #[test]
5769 fn reset_qpack_stream_after_type() {
5771 let mut s = Session::new().unwrap();
5772 s.handshake().unwrap();
5773
5774 s.pipe
5775 .client
5776 .stream_shutdown(
5777 s.client.local_qpack_streams.encoder_stream_id.unwrap(),
5778 crate::Shutdown::Write,
5779 0,
5780 )
5781 .unwrap();
5782
5783 s.advance().ok();
5784
5785 assert_eq!(
5786 Err(Error::ClosedCriticalStream),
5787 s.server.poll(&mut s.pipe.server)
5788 );
5789 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5790 }
5791
5792 #[test]
5793 fn reset_qpack_stream_after_data() {
5795 let mut s = Session::new().unwrap();
5796 s.handshake().unwrap();
5797
5798 let stream_id = s.client.local_qpack_streams.encoder_stream_id.unwrap();
5799 let d = [0; 1];
5800
5801 s.pipe.client.stream_send(stream_id, &d, false).unwrap();
5802 s.pipe.client.stream_send(stream_id, &d, false).unwrap();
5803
5804 s.advance().ok();
5805
5806 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5807
5808 s.pipe
5809 .client
5810 .stream_shutdown(stream_id, crate::Shutdown::Write, 0)
5811 .unwrap();
5812
5813 s.advance().ok();
5814
5815 assert_eq!(
5816 Err(Error::ClosedCriticalStream),
5817 s.server.poll(&mut s.pipe.server)
5818 );
5819 assert_eq!(Err(Error::Done), s.server.poll(&mut s.pipe.server));
5820 }
5821
5822 #[test]
5823 fn qpack_data() {
5825 let mut s = Session::new().unwrap();
5828 s.handshake().unwrap();
5829
5830 let e_stream_id = s.client.local_qpack_streams.encoder_stream_id.unwrap();
5831 let d_stream_id = s.client.local_qpack_streams.decoder_stream_id.unwrap();
5832 let d = [0; 20];
5833
5834 s.pipe.client.stream_send(e_stream_id, &d, false).unwrap();
5835 s.advance().ok();
5836
5837 s.pipe.client.stream_send(d_stream_id, &d, false).unwrap();
5838 s.advance().ok();
5839
5840 match s.server.poll(&mut s.pipe.server) {
5841 Ok(_) => panic!(),
5842
5843 Err(Error::Done) => {
5844 assert_eq!(s.server.peer_qpack_streams.encoder_stream_bytes, 20);
5845 assert_eq!(s.server.peer_qpack_streams.decoder_stream_bytes, 20);
5846 },
5847
5848 Err(_) => {
5849 panic!();
5850 },
5851 }
5852
5853 let stats = s.server.stats();
5854 assert_eq!(stats.qpack_encoder_stream_recv_bytes, 20);
5855 assert_eq!(stats.qpack_decoder_stream_recv_bytes, 20);
5856 }
5857
5858 #[test]
5859 fn max_state_buf_size() {
5861 let mut s = Session::new().unwrap();
5862 s.handshake().unwrap();
5863
5864 let req = vec![
5865 Header::new(b":method", b"GET"),
5866 Header::new(b":scheme", b"https"),
5867 Header::new(b":authority", b"quic.tech"),
5868 Header::new(b":path", b"/test"),
5869 Header::new(b"user-agent", b"quiche-test"),
5870 ];
5871
5872 assert_eq!(
5873 s.client.send_request(&mut s.pipe.client, &req, false),
5874 Ok(0)
5875 );
5876
5877 s.advance().ok();
5878
5879 let ev_headers = Event::Headers {
5880 list: req,
5881 more_frames: true,
5882 };
5883
5884 assert_eq!(s.server.poll(&mut s.pipe.server), Ok((0, ev_headers)));
5885
5886 let mut d = [42; 128];
5888 let mut b = octets::OctetsMut::with_slice(&mut d);
5889
5890 let frame_type = b.put_varint(frame::DATA_FRAME_TYPE_ID).unwrap();
5891 s.pipe.client.stream_send(0, frame_type, false).unwrap();
5892
5893 let frame_len = b.put_varint(1 << 24).unwrap();
5894 s.pipe.client.stream_send(0, frame_len, false).unwrap();
5895
5896 s.pipe.client.stream_send(0, &d, false).unwrap();
5897
5898 s.advance().ok();
5899
5900 assert_eq!(s.server.poll(&mut s.pipe.server), Ok((0, Event::Data)));
5901
5902 let mut s = Session::new().unwrap();
5904 s.handshake().unwrap();
5905
5906 let mut d = [42; 128];
5907 let mut b = octets::OctetsMut::with_slice(&mut d);
5908
5909 let frame_type = b.put_varint(148_764_065_110_560_899).unwrap();
5910 s.pipe.client.stream_send(0, frame_type, false).unwrap();
5911
5912 let frame_len = b.put_varint(1 << 24).unwrap();
5913 s.pipe.client.stream_send(0, frame_len, false).unwrap();
5914
5915 s.pipe.client.stream_send(0, &d, false).unwrap();
5916
5917 s.advance().ok();
5918
5919 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::ExcessiveLoad));
5920 }
5921
5922 #[test]
5923 fn stream_backpressure() {
5926 let bytes = vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
5927
5928 let mut s = Session::new().unwrap();
5929 s.handshake().unwrap();
5930
5931 let (stream, req) = s.send_request(false).unwrap();
5932
5933 let total_data_frames = 6;
5934
5935 for _ in 0..total_data_frames {
5936 assert_eq!(
5937 s.client
5938 .send_body(&mut s.pipe.client, stream, &bytes, false),
5939 Ok(bytes.len())
5940 );
5941
5942 s.advance().ok();
5943 }
5944
5945 assert_eq!(
5946 s.client.send_body(&mut s.pipe.client, stream, &bytes, true),
5947 Ok(bytes.len() - 2)
5948 );
5949
5950 s.advance().ok();
5951
5952 let mut recv_buf = vec![0; bytes.len()];
5953
5954 let ev_headers = Event::Headers {
5955 list: req,
5956 more_frames: true,
5957 };
5958
5959 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
5960 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
5961 assert_eq!(s.poll_server(), Err(Error::Done));
5962
5963 for _ in 0..total_data_frames {
5964 assert_eq!(
5965 s.recv_body_server(stream, &mut recv_buf),
5966 Ok(bytes.len())
5967 );
5968 }
5969
5970 assert_eq!(
5971 s.recv_body_server(stream, &mut recv_buf),
5972 Ok(bytes.len() - 2)
5973 );
5974
5975 assert_eq!(s.poll_server(), Err(Error::Done));
5978
5979 assert_eq!(s.pipe.server.data_blocked_sent_count, 0);
5980 assert_eq!(s.pipe.server.stream_data_blocked_sent_count, 0);
5981 assert_eq!(s.pipe.server.data_blocked_recv_count, 0);
5982 assert_eq!(s.pipe.server.stream_data_blocked_recv_count, 1);
5983
5984 assert_eq!(s.pipe.client.data_blocked_sent_count, 0);
5985 assert_eq!(s.pipe.client.stream_data_blocked_sent_count, 1);
5986 assert_eq!(s.pipe.client.data_blocked_recv_count, 0);
5987 assert_eq!(s.pipe.client.stream_data_blocked_recv_count, 0);
5988 }
5989
5990 #[test]
5991 fn request_max_header_size_limit_accepts_large_headers() {
5994 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
5995 config
5996 .load_cert_chain_from_pem_file("examples/cert.crt")
5997 .unwrap();
5998 config
5999 .load_priv_key_from_pem_file("examples/cert.key")
6000 .unwrap();
6001 config.set_application_protos(&[b"h3"]).unwrap();
6002 config.set_initial_max_data(150000000);
6003 config.set_initial_max_stream_data_bidi_local(150000000);
6004 config.set_initial_max_stream_data_bidi_remote(150000000);
6005 config.set_initial_max_stream_data_uni(150000000);
6006 config.set_initial_max_streams_bidi(5);
6007 config.set_initial_max_streams_uni(5);
6008 config.verify_peer(false);
6009 config.set_initial_congestion_window_packets(100);
6010
6011 let mut h3_config = Config::new().unwrap();
6012 h3_config.set_max_field_section_size(256 * 1024);
6013
6014 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6015
6016 s.handshake().unwrap();
6017
6018 let mut req = vec![
6019 Header::new(b":method", b"GET"),
6020 Header::new(b":scheme", b"https"),
6021 Header::new(b":authority", b"quic.tech"),
6022 Header::new(b":path", b"/test"),
6023 ];
6024
6025 for _ in 1..5000 {
6026 req.push(Header::new(b"aaaaaaaaaa", b"aaaaaaaaa"));
6027 }
6028
6029 let ev_headers = Event::Headers {
6030 list: req.clone(),
6031 more_frames: false,
6032 };
6033
6034 let stream = s
6035 .client
6036 .send_request(&mut s.pipe.client, &req, true)
6037 .unwrap();
6038
6039 s.advance().ok();
6040
6041 assert_eq!(stream, 0);
6042
6043 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
6044 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
6045 assert_eq!(s.poll_server(), Err(Error::Done));
6046 }
6047
6048 #[test]
6049 fn request_max_header_size_limit_decoded_field_section() {
6051 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6052 config
6053 .load_cert_chain_from_pem_file("examples/cert.crt")
6054 .unwrap();
6055 config
6056 .load_priv_key_from_pem_file("examples/cert.key")
6057 .unwrap();
6058 config.set_application_protos(&[b"h3"]).unwrap();
6059 config.set_initial_max_data(1500);
6060 config.set_initial_max_stream_data_bidi_local(150);
6061 config.set_initial_max_stream_data_bidi_remote(150);
6062 config.set_initial_max_stream_data_uni(150);
6063 config.set_initial_max_streams_bidi(5);
6064 config.set_initial_max_streams_uni(5);
6065 config.verify_peer(false);
6066
6067 let mut h3_config = Config::new().unwrap();
6068 h3_config.set_max_field_section_size(65);
6069
6070 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6071
6072 s.handshake().unwrap();
6073
6074 let req = vec![
6075 Header::new(b":method", b"GET"),
6076 Header::new(b":scheme", b"https"),
6077 Header::new(b":authority", b"quic.tech"),
6078 Header::new(b":path", b"/test"),
6079 Header::new(b"aaaaaaa", b"aaaaaaaa"),
6080 ];
6081
6082 let stream = s
6083 .client
6084 .send_request(&mut s.pipe.client, &req, true)
6085 .unwrap();
6086
6087 s.advance().ok();
6088
6089 assert_eq!(stream, 0);
6090
6091 assert_eq!(s.poll_server(), Err(Error::ExcessiveLoad));
6092
6093 assert_eq!(
6094 s.pipe.server.local_error.as_ref().unwrap().error_code,
6095 Error::to_wire(Error::ExcessiveLoad)
6096 );
6097 }
6098
6099 #[test]
6100 fn request_max_header_size_limit_default_abort_before_decode() {
6103 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6104 config
6105 .load_cert_chain_from_pem_file("examples/cert.crt")
6106 .unwrap();
6107 config
6108 .load_priv_key_from_pem_file("examples/cert.key")
6109 .unwrap();
6110 config.set_application_protos(&[b"h3"]).unwrap();
6111 config.set_initial_max_data(150000);
6112 config.set_initial_max_stream_data_bidi_local(150000);
6113 config.set_initial_max_stream_data_bidi_remote(150000);
6114 config.set_initial_max_stream_data_uni(150000);
6115 config.set_initial_max_streams_bidi(5);
6116 config.set_initial_max_streams_uni(5);
6117 config.verify_peer(false);
6118
6119 let h3_config = Config::new().unwrap();
6120
6121 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6122
6123 s.handshake().unwrap();
6124
6125 let mut d = vec![42; 200000];
6126 let mut b = octets::OctetsMut::with_slice(&mut d);
6127
6128 let hdrs = frame::Frame::Headers {
6129 header_block: vec![0; 65536],
6130 };
6131
6132 hdrs.to_bytes(&mut b).unwrap();
6133
6134 s.pipe.client.stream_send(0, &d, true).unwrap();
6135
6136 s.advance().ok();
6137
6138 assert_eq!(s.poll_server(), Err(Error::ExcessiveLoad));
6139
6140 assert_eq!(
6141 s.pipe.server.local_error.as_ref().unwrap().error_code,
6142 Error::to_wire(Error::ExcessiveLoad)
6143 );
6144 }
6145
6146 #[test]
6147 fn transport_error() {
6149 let mut s = Session::new().unwrap();
6150 s.handshake().unwrap();
6151
6152 let req = vec![
6153 Header::new(b":method", b"GET"),
6154 Header::new(b":scheme", b"https"),
6155 Header::new(b":authority", b"quic.tech"),
6156 Header::new(b":path", b"/test"),
6157 Header::new(b"user-agent", b"quiche-test"),
6158 ];
6159
6160 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(0));
6165 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(4));
6166 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(8));
6167 assert_eq!(
6168 s.client.send_request(&mut s.pipe.client, &req, true),
6169 Ok(12)
6170 );
6171 assert_eq!(
6172 s.client.send_request(&mut s.pipe.client, &req, true),
6173 Ok(16)
6174 );
6175
6176 assert_eq!(
6177 s.client.send_request(&mut s.pipe.client, &req, true),
6178 Err(Error::TransportError(crate::Error::StreamLimit))
6179 );
6180 }
6181
6182 #[test]
6183 fn data_before_headers() {
6185 let mut s = Session::new().unwrap();
6186 s.handshake().unwrap();
6187
6188 let mut d = [42; 128];
6189 let mut b = octets::OctetsMut::with_slice(&mut d);
6190
6191 let frame_type = b.put_varint(frame::DATA_FRAME_TYPE_ID).unwrap();
6192 s.pipe.client.stream_send(0, frame_type, false).unwrap();
6193
6194 let frame_len = b.put_varint(5).unwrap();
6195 s.pipe.client.stream_send(0, frame_len, false).unwrap();
6196
6197 s.pipe.client.stream_send(0, b"hello", false).unwrap();
6198
6199 s.advance().ok();
6200
6201 assert_eq!(
6202 s.server.poll(&mut s.pipe.server),
6203 Err(Error::FrameUnexpected)
6204 );
6205 }
6206
6207 #[test]
6208 fn poll_after_error() {
6210 let mut s = Session::new().unwrap();
6211 s.handshake().unwrap();
6212
6213 let mut d = [42; 128];
6214 let mut b = octets::OctetsMut::with_slice(&mut d);
6215
6216 let frame_type = b.put_varint(148_764_065_110_560_899).unwrap();
6217 s.pipe.client.stream_send(0, frame_type, false).unwrap();
6218
6219 let frame_len = b.put_varint(1 << 24).unwrap();
6220 s.pipe.client.stream_send(0, frame_len, false).unwrap();
6221
6222 s.pipe.client.stream_send(0, &d, false).unwrap();
6223
6224 s.advance().ok();
6225
6226 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::ExcessiveLoad));
6227
6228 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::Done));
6230 }
6231
6232 #[test]
6233 fn headers_blocked() {
6235 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6236 config
6237 .load_cert_chain_from_pem_file("examples/cert.crt")
6238 .unwrap();
6239 config
6240 .load_priv_key_from_pem_file("examples/cert.key")
6241 .unwrap();
6242 config.set_application_protos(&[b"h3"]).unwrap();
6243 config.set_initial_max_data(75);
6244 config.set_initial_max_stream_data_bidi_local(150);
6245 config.set_initial_max_stream_data_bidi_remote(150);
6246 config.set_initial_max_stream_data_uni(150);
6247 config.set_initial_max_streams_bidi(100);
6248 config.set_initial_max_streams_uni(5);
6249 config.verify_peer(false);
6250
6251 let h3_config = Config::new().unwrap();
6252
6253 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6254
6255 s.handshake().unwrap();
6256
6257 let req = vec![
6258 Header::new(b":method", b"GET"),
6259 Header::new(b":scheme", b"https"),
6260 Header::new(b":authority", b"quic.tech"),
6261 Header::new(b":path", b"/test"),
6262 ];
6263
6264 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(0));
6265
6266 assert_eq!(
6267 s.client.send_request(&mut s.pipe.client, &req, true),
6268 Err(Error::StreamBlocked)
6269 );
6270
6271 assert_eq!(s.pipe.client.stream_writable_next(), Some(2));
6273 assert_eq!(s.pipe.client.stream_writable_next(), Some(6));
6274 assert_eq!(s.pipe.client.stream_writable_next(), Some(10));
6275 assert_eq!(s.pipe.client.stream_writable_next(), None);
6276
6277 s.advance().ok();
6278
6279 assert_eq!(s.pipe.client.stream_writable_next(), Some(4));
6282 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(4));
6283
6284 assert_eq!(s.pipe.server.data_blocked_sent_count, 0);
6285 assert_eq!(s.pipe.server.stream_data_blocked_sent_count, 0);
6286 assert_eq!(s.pipe.server.data_blocked_recv_count, 1);
6287 assert_eq!(s.pipe.server.stream_data_blocked_recv_count, 0);
6288
6289 assert_eq!(s.pipe.client.data_blocked_sent_count, 1);
6290 assert_eq!(s.pipe.client.stream_data_blocked_sent_count, 0);
6291 assert_eq!(s.pipe.client.data_blocked_recv_count, 0);
6292 assert_eq!(s.pipe.client.stream_data_blocked_recv_count, 0);
6293 }
6294
6295 #[test]
6296 fn headers_blocked_on_conn() {
6298 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6299 config
6300 .load_cert_chain_from_pem_file("examples/cert.crt")
6301 .unwrap();
6302 config
6303 .load_priv_key_from_pem_file("examples/cert.key")
6304 .unwrap();
6305 config.set_application_protos(&[b"h3"]).unwrap();
6306 config.set_initial_max_data(75);
6307 config.set_initial_max_stream_data_bidi_local(150);
6308 config.set_initial_max_stream_data_bidi_remote(150);
6309 config.set_initial_max_stream_data_uni(150);
6310 config.set_initial_max_streams_bidi(100);
6311 config.set_initial_max_streams_uni(5);
6312 config.verify_peer(false);
6313
6314 let h3_config = Config::new().unwrap();
6315
6316 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6317
6318 s.handshake().unwrap();
6319
6320 let d = [42; 28];
6324 assert_eq!(s.pipe.client.stream_send(2, &d, false), Ok(23));
6325
6326 let req = vec![
6327 Header::new(b":method", b"GET"),
6328 Header::new(b":scheme", b"https"),
6329 Header::new(b":authority", b"quic.tech"),
6330 Header::new(b":path", b"/test"),
6331 ];
6332
6333 assert_eq!(
6336 s.client.send_request(&mut s.pipe.client, &req, true),
6337 Err(Error::StreamBlocked)
6338 );
6339 assert_eq!(s.pipe.client.stream_writable_next(), None);
6340
6341 s.advance().ok();
6344 assert_eq!(s.poll_server(), Err(Error::Done));
6345 s.advance().ok();
6346
6347 assert_eq!(s.pipe.client.stream_writable_next(), Some(2));
6349 assert_eq!(s.pipe.client.stream_writable_next(), Some(6));
6350 assert_eq!(s.client.send_request(&mut s.pipe.client, &req, true), Ok(0));
6351
6352 assert_eq!(s.pipe.server.data_blocked_sent_count, 0);
6353 assert_eq!(s.pipe.server.stream_data_blocked_sent_count, 0);
6354 assert_eq!(s.pipe.server.data_blocked_recv_count, 1);
6355 assert_eq!(s.pipe.server.stream_data_blocked_recv_count, 0);
6356
6357 assert_eq!(s.pipe.client.data_blocked_sent_count, 1);
6358 assert_eq!(s.pipe.client.stream_data_blocked_sent_count, 0);
6359 assert_eq!(s.pipe.client.data_blocked_recv_count, 0);
6360 assert_eq!(s.pipe.client.stream_data_blocked_recv_count, 0);
6361 }
6362
6363 #[test]
6364 fn headers_blocked_by_max_data_success_on_retry() {
6369 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6370 config
6371 .load_cert_chain_from_pem_file("examples/cert.crt")
6372 .unwrap();
6373 config
6374 .load_priv_key_from_pem_file("examples/cert.key")
6375 .unwrap();
6376 config.set_application_protos(&[b"h3"]).unwrap();
6377 config.set_initial_max_data(70);
6378 config.set_initial_max_stream_data_bidi_local(150);
6379 config.set_initial_max_stream_data_bidi_remote(150);
6380 config.set_initial_max_stream_data_uni(150);
6381 config.set_initial_max_streams_bidi(100);
6382 config.set_initial_max_streams_uni(5);
6383 config.verify_peer(false);
6384
6385 let h3_config = Config::new().unwrap();
6386
6387 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6388
6389 s.handshake().unwrap();
6390
6391 let req = vec![
6392 Header::new(b":method", b"GET"),
6393 Header::new(b":scheme", b"https"),
6394 Header::new(b":authority", b"quic.tech"),
6395 Header::new(b":path", b"/test/with/long/url"),
6396 ];
6397
6398 assert_eq!(
6403 s.client.send_request(&mut s.pipe.client, &req, true),
6404 Err(Error::StreamBlocked)
6405 );
6406
6407 assert!(!s.client.streams.contains_key(&0));
6410
6411 s.advance().ok();
6414 assert_eq!(s.poll_server(), Err(Error::Done));
6415 s.advance().ok();
6416
6417 let stream_id = s.client.send_request(&mut s.pipe.client, &req, true);
6420 assert_eq!(stream_id, Ok(0));
6421 assert!(s.client.streams.contains_key(&0));
6422 assert!(!s.client.streams.contains_key(&4));
6423
6424 let stream_id2 = s.client.send_request(&mut s.pipe.client, &req, true);
6426 assert_eq!(stream_id2, Ok(4));
6427 assert!(s.client.streams.contains_key(&0));
6428 assert!(s.client.streams.contains_key(&4));
6429
6430 s.advance().ok();
6431 }
6432
6433 #[test]
6434 fn send_body_truncation_stream_blocked() {
6437 use crate::test_utils::decode_pkt;
6438
6439 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6440 config
6441 .load_cert_chain_from_pem_file("examples/cert.crt")
6442 .unwrap();
6443 config
6444 .load_priv_key_from_pem_file("examples/cert.key")
6445 .unwrap();
6446 config.set_application_protos(&[b"h3"]).unwrap();
6447 config.set_initial_max_data(10000);
6449 config.set_initial_max_stream_data_bidi_local(80);
6450 config.set_initial_max_stream_data_bidi_remote(80);
6451 config.set_initial_max_stream_data_uni(150);
6452 config.set_initial_max_streams_bidi(100);
6453 config.set_initial_max_streams_uni(5);
6454 config.verify_peer(false);
6455
6456 let h3_config = Config::new().unwrap();
6457
6458 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6459
6460 s.handshake().unwrap();
6461
6462 let (stream, req) = s.send_request(true).unwrap();
6463
6464 let ev_headers = Event::Headers {
6465 list: req,
6466 more_frames: false,
6467 };
6468
6469 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
6470 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
6471
6472 let _ = s.send_response(stream, false).unwrap();
6473
6474 assert_eq!(s.pipe.server.streams.blocked().len(), 0);
6475
6476 let d = [42; 500];
6478 let mut off = 0;
6479
6480 let sent = s
6481 .server
6482 .send_body(&mut s.pipe.server, stream, &d, true)
6483 .unwrap();
6484 assert_eq!(sent, 25);
6485 off += sent;
6486
6487 assert_eq!(s.pipe.server.streams.blocked().len(), 1);
6489 assert_eq!(
6490 s.server
6491 .send_body(&mut s.pipe.server, stream, &d[off..], true),
6492 Err(Error::Done)
6493 );
6494 assert_eq!(s.pipe.server.streams.blocked().len(), 1);
6495
6496 let mut buf = [0; 65535];
6498 let (len, _) = s.pipe.server.send(&mut buf).unwrap();
6499
6500 let frames = decode_pkt(&mut s.pipe.client, &mut buf[..len]).unwrap();
6501
6502 let mut iter = frames.iter();
6503
6504 assert_eq!(
6505 iter.next(),
6506 Some(&crate::frame::Frame::StreamDataBlocked {
6507 stream_id: 0,
6508 limit: 80,
6509 })
6510 );
6511
6512 assert_eq!(s.pipe.server.streams.blocked().len(), 0);
6515
6516 assert_eq!(
6522 s.server
6523 .send_body(&mut s.pipe.server, stream, &d[off..], true),
6524 Err(Error::Done)
6525 );
6526 assert_eq!(s.pipe.server.streams.blocked().len(), 0);
6527 assert_eq!(s.pipe.server.send(&mut buf), Err(crate::Error::Done));
6528
6529 let frames = [crate::frame::Frame::MaxStreamData {
6531 stream_id: 0,
6532 max: 100,
6533 }];
6534
6535 let pkt_type = crate::packet::Type::Short;
6536 assert_eq!(
6537 s.pipe.send_pkt_to_server(pkt_type, &frames, &mut buf),
6538 Ok(39),
6539 );
6540
6541 let sent = s
6542 .server
6543 .send_body(&mut s.pipe.server, stream, &d[off..], true)
6544 .unwrap();
6545 assert_eq!(sent, 18);
6546
6547 assert_eq!(s.pipe.server.streams.blocked().len(), 1);
6549 assert_eq!(
6550 s.server
6551 .send_body(&mut s.pipe.server, stream, &d[off..], true),
6552 Err(Error::Done)
6553 );
6554 assert_eq!(s.pipe.server.streams.blocked().len(), 1);
6555
6556 let (len, _) = s.pipe.server.send(&mut buf).unwrap();
6557
6558 let frames = decode_pkt(&mut s.pipe.client, &mut buf[..len]).unwrap();
6559
6560 let mut iter = frames.iter();
6561
6562 assert_eq!(
6563 iter.next(),
6564 Some(&crate::frame::Frame::StreamDataBlocked {
6565 stream_id: 0,
6566 limit: 100,
6567 })
6568 );
6569 }
6570
6571 #[test]
6572 fn send_body_stream_blocked_by_small_cwnd() {
6574 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6575 config
6576 .load_cert_chain_from_pem_file("examples/cert.crt")
6577 .unwrap();
6578 config
6579 .load_priv_key_from_pem_file("examples/cert.key")
6580 .unwrap();
6581 config.set_application_protos(&[b"h3"]).unwrap();
6582 config.set_initial_max_data(100000);
6584 config.set_initial_max_stream_data_bidi_local(100000);
6585 config.set_initial_max_stream_data_bidi_remote(50000);
6586 config.set_initial_max_stream_data_uni(150);
6587 config.set_initial_max_streams_bidi(100);
6588 config.set_initial_max_streams_uni(5);
6589 config.verify_peer(false);
6590
6591 let h3_config = Config::new().unwrap();
6592
6593 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6594
6595 s.handshake().unwrap();
6596
6597 let (stream, req) = s.send_request(true).unwrap();
6598
6599 let ev_headers = Event::Headers {
6600 list: req,
6601 more_frames: false,
6602 };
6603
6604 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
6605 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
6606
6607 let _ = s.send_response(stream, false).unwrap();
6608
6609 assert_eq!(s.pipe.server.stream_writable_next(), Some(3));
6611 assert_eq!(s.pipe.server.stream_writable_next(), Some(7));
6612 assert_eq!(s.pipe.server.stream_writable_next(), Some(11));
6613 assert_eq!(s.pipe.server.stream_writable_next(), Some(stream));
6614 assert_eq!(s.pipe.server.stream_writable_next(), None);
6615
6616 let send_buf = [42; 80000];
6618
6619 let sent = s
6620 .server
6621 .send_body(&mut s.pipe.server, stream, &send_buf, true)
6622 .unwrap();
6623
6624 assert_eq!(sent, 11995);
6626
6627 s.advance().ok();
6628
6629 let mut recv_buf = [42; 80000];
6631 assert!(s.poll_client().is_ok());
6632 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
6633 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(11995));
6634
6635 s.advance().ok();
6636
6637 assert!(s.pipe.server.tx_cap < send_buf.len() - sent);
6639
6640 assert_eq!(s.pipe.server.stream_writable_next(), Some(0));
6642 }
6643
6644 #[test]
6645 fn send_body_stream_blocked_zero_length() {
6647 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6648 config
6649 .load_cert_chain_from_pem_file("examples/cert.crt")
6650 .unwrap();
6651 config
6652 .load_priv_key_from_pem_file("examples/cert.key")
6653 .unwrap();
6654 config.set_application_protos(&[b"h3"]).unwrap();
6655 config.set_initial_max_data(100000);
6657 config.set_initial_max_stream_data_bidi_local(100000);
6658 config.set_initial_max_stream_data_bidi_remote(50000);
6659 config.set_initial_max_stream_data_uni(150);
6660 config.set_initial_max_streams_bidi(100);
6661 config.set_initial_max_streams_uni(5);
6662 config.verify_peer(false);
6663
6664 let h3_config = Config::new().unwrap();
6665
6666 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6667
6668 s.handshake().unwrap();
6669
6670 let (stream, req) = s.send_request(true).unwrap();
6671
6672 let ev_headers = Event::Headers {
6673 list: req,
6674 more_frames: false,
6675 };
6676
6677 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
6678 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
6679
6680 let _ = s.send_response(stream, false).unwrap();
6681
6682 assert_eq!(s.pipe.server.stream_writable_next(), Some(3));
6684 assert_eq!(s.pipe.server.stream_writable_next(), Some(7));
6685 assert_eq!(s.pipe.server.stream_writable_next(), Some(11));
6686 assert_eq!(s.pipe.server.stream_writable_next(), Some(stream));
6687 assert_eq!(s.pipe.server.stream_writable_next(), None);
6688
6689 let send_buf = [42; 11994];
6692
6693 let sent = s
6694 .server
6695 .send_body(&mut s.pipe.server, stream, &send_buf, false)
6696 .unwrap();
6697
6698 assert_eq!(sent, 11994);
6699
6700 assert_eq!(s.pipe.server.stream_capacity(stream).unwrap(), 3);
6703 assert_eq!(
6704 s.server
6705 .send_body(&mut s.pipe.server, stream, &send_buf, false),
6706 Err(Error::Done)
6707 );
6708
6709 s.advance().ok();
6710
6711 let mut recv_buf = [42; 80000];
6713 assert!(s.poll_client().is_ok());
6714 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
6715 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(11994));
6716
6717 s.advance().ok();
6718
6719 assert_eq!(s.pipe.server.stream_writable_next(), Some(0));
6721 }
6722
6723 #[test]
6724 fn zero_length_data() {
6726 let mut s = Session::new().unwrap();
6727 s.handshake().unwrap();
6728
6729 let (stream, req) = s.send_request(false).unwrap();
6730
6731 assert_eq!(
6732 s.client.send_body(&mut s.pipe.client, 0, b"", false),
6733 Err(Error::Done)
6734 );
6735 assert_eq!(s.client.send_body(&mut s.pipe.client, 0, b"", true), Ok(0));
6736
6737 s.advance().ok();
6738
6739 let mut recv_buf = vec![0; 100];
6740
6741 let ev_headers = Event::Headers {
6742 list: req,
6743 more_frames: true,
6744 };
6745
6746 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
6747
6748 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
6749 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Err(Error::Done));
6750
6751 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
6752 assert_eq!(s.poll_server(), Err(Error::Done));
6753
6754 let resp = s.send_response(stream, false).unwrap();
6755
6756 assert_eq!(
6757 s.server.send_body(&mut s.pipe.server, 0, b"", false),
6758 Err(Error::Done)
6759 );
6760 assert_eq!(s.server.send_body(&mut s.pipe.server, 0, b"", true), Ok(0));
6761
6762 s.advance().ok();
6763
6764 let ev_headers = Event::Headers {
6765 list: resp,
6766 more_frames: true,
6767 };
6768
6769 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
6770
6771 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
6772 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Err(Error::Done));
6773
6774 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
6775 assert_eq!(s.poll_client(), Err(Error::Done));
6776 }
6777
6778 #[test]
6779 fn zero_length_data_blocked() {
6781 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6782 config
6783 .load_cert_chain_from_pem_file("examples/cert.crt")
6784 .unwrap();
6785 config
6786 .load_priv_key_from_pem_file("examples/cert.key")
6787 .unwrap();
6788 config.set_application_protos(&[b"h3"]).unwrap();
6789 config.set_initial_max_data(74);
6790 config.set_initial_max_stream_data_bidi_local(150);
6791 config.set_initial_max_stream_data_bidi_remote(150);
6792 config.set_initial_max_stream_data_uni(150);
6793 config.set_initial_max_streams_bidi(100);
6794 config.set_initial_max_streams_uni(5);
6795 config.verify_peer(false);
6796
6797 let h3_config = Config::new().unwrap();
6798
6799 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6800
6801 s.handshake().unwrap();
6802
6803 let req = vec![
6804 Header::new(b":method", b"GET"),
6805 Header::new(b":scheme", b"https"),
6806 Header::new(b":authority", b"quic.tech"),
6807 Header::new(b":path", b"/test"),
6808 ];
6809
6810 assert_eq!(
6811 s.client.send_request(&mut s.pipe.client, &req, false),
6812 Ok(0)
6813 );
6814
6815 assert_eq!(
6816 s.client.send_body(&mut s.pipe.client, 0, b"", true),
6817 Err(Error::Done)
6818 );
6819
6820 assert_eq!(s.pipe.client.stream_writable_next(), Some(2));
6822 assert_eq!(s.pipe.client.stream_writable_next(), Some(6));
6823 assert_eq!(s.pipe.client.stream_writable_next(), Some(10));
6824 assert_eq!(s.pipe.client.stream_writable_next(), None);
6825
6826 s.advance().ok();
6827
6828 assert_eq!(s.pipe.client.stream_writable_next(), Some(0));
6830 assert_eq!(s.client.send_body(&mut s.pipe.client, 0, b"", true), Ok(0));
6831 }
6832
6833 #[test]
6834 fn empty_settings() {
6836 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6837 config
6838 .load_cert_chain_from_pem_file("examples/cert.crt")
6839 .unwrap();
6840 config
6841 .load_priv_key_from_pem_file("examples/cert.key")
6842 .unwrap();
6843 config.set_application_protos(&[b"h3"]).unwrap();
6844 config.set_initial_max_data(1500);
6845 config.set_initial_max_stream_data_bidi_local(150);
6846 config.set_initial_max_stream_data_bidi_remote(150);
6847 config.set_initial_max_stream_data_uni(150);
6848 config.set_initial_max_streams_bidi(5);
6849 config.set_initial_max_streams_uni(5);
6850 config.verify_peer(false);
6851 config.set_ack_delay_exponent(8);
6852 config.grease(false);
6853
6854 let h3_config = Config::new().unwrap();
6855 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6856
6857 s.handshake().unwrap();
6858
6859 assert!(s.client.peer_settings_raw().is_some());
6860 assert!(s.server.peer_settings_raw().is_some());
6861 }
6862
6863 #[test]
6864 fn dgram_setting() {
6866 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6867 config
6868 .load_cert_chain_from_pem_file("examples/cert.crt")
6869 .unwrap();
6870 config
6871 .load_priv_key_from_pem_file("examples/cert.key")
6872 .unwrap();
6873 config.set_application_protos(&[b"h3"]).unwrap();
6874 config.set_initial_max_data(70);
6875 config.set_initial_max_stream_data_bidi_local(150);
6876 config.set_initial_max_stream_data_bidi_remote(150);
6877 config.set_initial_max_stream_data_uni(150);
6878 config.set_initial_max_streams_bidi(100);
6879 config.set_initial_max_streams_uni(5);
6880 config.enable_dgram(true, 1000, 1000);
6881 config.verify_peer(false);
6882
6883 let h3_config = Config::new().unwrap();
6884
6885 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6886 assert_eq!(s.pipe.handshake(), Ok(()));
6887
6888 s.client.send_settings(&mut s.pipe.client).unwrap();
6889 assert_eq!(s.pipe.advance(), Ok(()));
6890
6891 assert!(!s.server.dgram_enabled_by_peer(&s.pipe.server));
6894
6895 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::Done));
6897 assert!(s.server.dgram_enabled_by_peer(&s.pipe.server));
6898
6899 s.server.send_settings(&mut s.pipe.server).unwrap();
6901 assert_eq!(s.pipe.advance(), Ok(()));
6902 assert!(!s.client.dgram_enabled_by_peer(&s.pipe.client));
6903 assert_eq!(s.client.poll(&mut s.pipe.client), Err(Error::Done));
6904 assert!(s.client.dgram_enabled_by_peer(&s.pipe.client));
6905 }
6906
6907 #[test]
6908 fn dgram_setting_no_tp() {
6911 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6912 config
6913 .load_cert_chain_from_pem_file("examples/cert.crt")
6914 .unwrap();
6915 config
6916 .load_priv_key_from_pem_file("examples/cert.key")
6917 .unwrap();
6918 config.set_application_protos(&[b"h3"]).unwrap();
6919 config.set_initial_max_data(70);
6920 config.set_initial_max_stream_data_bidi_local(150);
6921 config.set_initial_max_stream_data_bidi_remote(150);
6922 config.set_initial_max_stream_data_uni(150);
6923 config.set_initial_max_streams_bidi(100);
6924 config.set_initial_max_streams_uni(5);
6925 config.verify_peer(false);
6926
6927 let h3_config = Config::new().unwrap();
6928
6929 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6930 assert_eq!(s.pipe.handshake(), Ok(()));
6931
6932 s.client.control_stream_id = Some(
6933 s.client
6934 .open_uni_stream(
6935 &mut s.pipe.client,
6936 stream::HTTP3_CONTROL_STREAM_TYPE_ID,
6937 )
6938 .unwrap(),
6939 );
6940
6941 let settings = frame::Frame::Settings {
6942 max_field_section_size: None,
6943 qpack_max_table_capacity: None,
6944 qpack_blocked_streams: None,
6945 connect_protocol_enabled: None,
6946 h3_datagram: Some(1),
6947 grease: None,
6948 additional_settings: Default::default(),
6949 raw: Default::default(),
6950 };
6951
6952 s.send_frame_client(settings, s.client.control_stream_id.unwrap(), false)
6953 .unwrap();
6954
6955 assert_eq!(s.pipe.advance(), Ok(()));
6956
6957 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::SettingsError));
6958 }
6959
6960 #[test]
6961 fn settings_h2_prohibited() {
6963 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
6964 config
6965 .load_cert_chain_from_pem_file("examples/cert.crt")
6966 .unwrap();
6967 config
6968 .load_priv_key_from_pem_file("examples/cert.key")
6969 .unwrap();
6970 config.set_application_protos(&[b"h3"]).unwrap();
6971 config.set_initial_max_data(70);
6972 config.set_initial_max_stream_data_bidi_local(150);
6973 config.set_initial_max_stream_data_bidi_remote(150);
6974 config.set_initial_max_stream_data_uni(150);
6975 config.set_initial_max_streams_bidi(100);
6976 config.set_initial_max_streams_uni(5);
6977 config.verify_peer(false);
6978
6979 let h3_config = Config::new().unwrap();
6980
6981 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
6982 assert_eq!(s.pipe.handshake(), Ok(()));
6983
6984 s.client.control_stream_id = Some(
6985 s.client
6986 .open_uni_stream(
6987 &mut s.pipe.client,
6988 stream::HTTP3_CONTROL_STREAM_TYPE_ID,
6989 )
6990 .unwrap(),
6991 );
6992
6993 s.server.control_stream_id = Some(
6994 s.server
6995 .open_uni_stream(
6996 &mut s.pipe.server,
6997 stream::HTTP3_CONTROL_STREAM_TYPE_ID,
6998 )
6999 .unwrap(),
7000 );
7001
7002 let frame_payload_len = 2u64;
7003 let settings = [
7004 frame::SETTINGS_FRAME_TYPE_ID as u8,
7005 frame_payload_len as u8,
7006 0x2, 1,
7008 ];
7009
7010 s.send_arbitrary_stream_data_client(
7011 &settings,
7012 s.client.control_stream_id.unwrap(),
7013 false,
7014 )
7015 .unwrap();
7016
7017 s.send_arbitrary_stream_data_server(
7018 &settings,
7019 s.server.control_stream_id.unwrap(),
7020 false,
7021 )
7022 .unwrap();
7023
7024 assert_eq!(s.pipe.advance(), Ok(()));
7025
7026 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::SettingsError));
7027
7028 assert_eq!(s.client.poll(&mut s.pipe.client), Err(Error::SettingsError));
7029 }
7030
7031 #[test]
7032 fn set_prohibited_additional_settings() {
7034 let mut h3_config = Config::new().unwrap();
7035 assert_eq!(
7036 h3_config.set_additional_settings(vec![(
7037 frame::SETTINGS_QPACK_MAX_TABLE_CAPACITY,
7038 43
7039 )]),
7040 Err(Error::SettingsError)
7041 );
7042 assert_eq!(
7043 h3_config.set_additional_settings(vec![(
7044 frame::SETTINGS_MAX_FIELD_SECTION_SIZE,
7045 43
7046 )]),
7047 Err(Error::SettingsError)
7048 );
7049 assert_eq!(
7050 h3_config.set_additional_settings(vec![(
7051 frame::SETTINGS_QPACK_BLOCKED_STREAMS,
7052 43
7053 )]),
7054 Err(Error::SettingsError)
7055 );
7056 assert_eq!(
7057 h3_config.set_additional_settings(vec![(
7058 frame::SETTINGS_ENABLE_CONNECT_PROTOCOL,
7059 43
7060 )]),
7061 Err(Error::SettingsError)
7062 );
7063 assert_eq!(
7064 h3_config
7065 .set_additional_settings(vec![(frame::SETTINGS_H3_DATAGRAM, 43)]),
7066 Err(Error::SettingsError)
7067 );
7068 }
7069
7070 #[test]
7071 fn settings_on_request_stream_client() {
7074 let mut s = Session::new().unwrap();
7075 s.handshake().unwrap();
7076
7077 let (stream, _req) = s.send_request(true).unwrap();
7078
7079 let settings = frame::Frame::Settings {
7080 max_field_section_size: None,
7081 qpack_max_table_capacity: None,
7082 qpack_blocked_streams: None,
7083 connect_protocol_enabled: None,
7084 h3_datagram: None,
7085 grease: None,
7086 additional_settings: Default::default(),
7087 raw: Default::default(),
7088 };
7089
7090 s.send_frame_server(settings, stream, false).unwrap();
7091
7092 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
7094 assert_eq!(
7095 s.pipe.client.local_error(),
7096 Some(&crate::ConnectionError {
7097 is_app: true,
7098 error_code: WireErrorCode::FrameUnexpected as u64,
7099 reason: format!(
7100 "Unexpected frame type {}",
7101 frame::SETTINGS_FRAME_TYPE_ID
7102 )
7103 .into_bytes(),
7104 })
7105 );
7106 }
7107
7108 #[test]
7109 fn cancel_push_on_request_stream_client() {
7112 let mut s = Session::new().unwrap();
7113 s.handshake().unwrap();
7114
7115 let (stream, _req) = s.send_request(true).unwrap();
7116 let cancel_push = frame::Frame::CancelPush { push_id: 0 };
7117 s.send_frame_server(cancel_push, stream, false).unwrap();
7118
7119 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
7121 assert_eq!(
7122 s.pipe.client.local_error(),
7123 Some(&crate::ConnectionError {
7124 is_app: true,
7125 error_code: WireErrorCode::FrameUnexpected as u64,
7126 reason: format!(
7127 "Unexpected frame type {}",
7128 frame::CANCEL_PUSH_FRAME_TYPE_ID
7129 )
7130 .into_bytes(),
7131 })
7132 );
7133 }
7134
7135 #[test]
7136 fn goaway_on_request_stream_client() {
7139 let mut s = Session::new().unwrap();
7140 s.handshake().unwrap();
7141
7142 let (stream, _req) = s.send_request(true).unwrap();
7143 let goaway = frame::Frame::GoAway { id: 0 };
7144
7145 s.send_frame_server(goaway, stream, false).unwrap();
7146
7147 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
7149 assert_eq!(
7150 s.pipe.client.local_error(),
7151 Some(&crate::ConnectionError {
7152 is_app: true,
7153 error_code: WireErrorCode::FrameUnexpected as u64,
7154 reason: format!(
7155 "Unexpected frame type {}",
7156 frame::GOAWAY_FRAME_TYPE_ID
7157 )
7158 .into_bytes(),
7159 })
7160 );
7161 }
7162
7163 #[test]
7164 fn max_push_id_on_request_stream_client() {
7167 let mut s = Session::new().unwrap();
7168 s.handshake().unwrap();
7169
7170 let (stream, _req) = s.send_request(true).unwrap();
7171 let max_push_id = frame::Frame::MaxPushId { push_id: 0 };
7172
7173 s.send_frame_server(max_push_id, stream, false).unwrap();
7174
7175 assert_eq!(s.poll_client(), Err(Error::FrameUnexpected));
7177 assert_eq!(
7178 s.pipe.client.local_error(),
7179 Some(&crate::ConnectionError {
7180 is_app: true,
7181 error_code: WireErrorCode::FrameUnexpected as u64,
7182 reason: format!(
7183 "Unexpected frame type {}",
7184 frame::MAX_PUSH_FRAME_TYPE_ID
7185 )
7186 .into_bytes(),
7187 })
7188 );
7189 }
7190
7191 #[test]
7192 fn set_additional_settings() {
7194 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
7195 config
7196 .load_cert_chain_from_pem_file("examples/cert.crt")
7197 .unwrap();
7198 config
7199 .load_priv_key_from_pem_file("examples/cert.key")
7200 .unwrap();
7201 config.set_application_protos(&[b"h3"]).unwrap();
7202 config.set_initial_max_data(70);
7203 config.set_initial_max_stream_data_bidi_local(150);
7204 config.set_initial_max_stream_data_bidi_remote(150);
7205 config.set_initial_max_stream_data_uni(150);
7206 config.set_initial_max_streams_bidi(100);
7207 config.set_initial_max_streams_uni(5);
7208 config.verify_peer(false);
7209 config.grease(false);
7210
7211 let mut h3_config = Config::new().unwrap();
7212 h3_config
7213 .set_additional_settings(vec![(42, 43), (44, 45)])
7214 .unwrap();
7215
7216 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7217 assert_eq!(s.pipe.handshake(), Ok(()));
7218
7219 assert_eq!(s.pipe.advance(), Ok(()));
7220
7221 s.client.send_settings(&mut s.pipe.client).unwrap();
7222 assert_eq!(s.pipe.advance(), Ok(()));
7223 assert_eq!(s.server.poll(&mut s.pipe.server), Err(Error::Done));
7224
7225 s.server.send_settings(&mut s.pipe.server).unwrap();
7226 assert_eq!(s.pipe.advance(), Ok(()));
7227 assert_eq!(s.client.poll(&mut s.pipe.client), Err(Error::Done));
7228
7229 assert_eq!(
7230 s.server.peer_settings_raw(),
7231 Some(&[(6, 32_768), (42, 43), (44, 45)][..])
7232 );
7233 assert_eq!(
7234 s.client.peer_settings_raw(),
7235 Some(&[(6, 32_768), (42, 43), (44, 45)][..])
7236 );
7237 }
7238
7239 #[test]
7240 fn single_dgram() {
7242 let mut buf = [0; 65535];
7243 let mut s = Session::new().unwrap();
7244 s.handshake().unwrap();
7245
7246 let result = (11, 0, 1);
7248
7249 s.send_dgram_client(0).unwrap();
7250
7251 assert_eq!(s.poll_server(), Err(Error::Done));
7252 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7253
7254 s.send_dgram_server(0).unwrap();
7255 assert_eq!(s.poll_client(), Err(Error::Done));
7256 assert_eq!(s.recv_dgram_client(&mut buf), Ok(result));
7257 }
7258
7259 #[test]
7260 fn multiple_dgram() {
7262 let mut buf = [0; 65535];
7263 let mut s = Session::new().unwrap();
7264 s.handshake().unwrap();
7265
7266 let result = (11, 0, 1);
7268
7269 s.send_dgram_client(0).unwrap();
7270 s.send_dgram_client(0).unwrap();
7271 s.send_dgram_client(0).unwrap();
7272
7273 assert_eq!(s.poll_server(), Err(Error::Done));
7274 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7275 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7276 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7277 assert_eq!(s.recv_dgram_server(&mut buf), Err(Error::Done));
7278
7279 s.send_dgram_server(0).unwrap();
7280 s.send_dgram_server(0).unwrap();
7281 s.send_dgram_server(0).unwrap();
7282
7283 assert_eq!(s.poll_client(), Err(Error::Done));
7284 assert_eq!(s.recv_dgram_client(&mut buf), Ok(result));
7285 assert_eq!(s.recv_dgram_client(&mut buf), Ok(result));
7286 assert_eq!(s.recv_dgram_client(&mut buf), Ok(result));
7287 assert_eq!(s.recv_dgram_client(&mut buf), Err(Error::Done));
7288 }
7289
7290 #[test]
7291 fn multiple_dgram_overflow() {
7293 let mut buf = [0; 65535];
7294 let mut s = Session::new().unwrap();
7295 s.handshake().unwrap();
7296
7297 let result = (11, 0, 1);
7299
7300 s.send_dgram_client(0).unwrap();
7302 s.send_dgram_client(0).unwrap();
7303 s.send_dgram_client(0).unwrap();
7304 s.send_dgram_client(0).unwrap();
7305 s.send_dgram_client(0).unwrap();
7306
7307 assert_eq!(s.poll_server(), Err(Error::Done));
7309 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7310 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7311 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7312 assert_eq!(s.recv_dgram_server(&mut buf), Err(Error::Done));
7313 }
7314
7315 #[test]
7316 fn poll_datagram_cycling_no_read() {
7318 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
7319 config
7320 .load_cert_chain_from_pem_file("examples/cert.crt")
7321 .unwrap();
7322 config
7323 .load_priv_key_from_pem_file("examples/cert.key")
7324 .unwrap();
7325 config.set_application_protos(&[b"h3"]).unwrap();
7326 config.set_initial_max_data(1500);
7327 config.set_initial_max_stream_data_bidi_local(150);
7328 config.set_initial_max_stream_data_bidi_remote(150);
7329 config.set_initial_max_stream_data_uni(150);
7330 config.set_initial_max_streams_bidi(100);
7331 config.set_initial_max_streams_uni(5);
7332 config.verify_peer(false);
7333 config.enable_dgram(true, 100, 100);
7334
7335 let h3_config = Config::new().unwrap();
7336 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7337 s.handshake().unwrap();
7338
7339 let (stream, req) = s.send_request(false).unwrap();
7341
7342 s.send_body_client(stream, true).unwrap();
7343
7344 let ev_headers = Event::Headers {
7345 list: req,
7346 more_frames: true,
7347 };
7348
7349 s.send_dgram_client(0).unwrap();
7350
7351 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
7352 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
7353
7354 assert_eq!(s.poll_server(), Err(Error::Done));
7355 }
7356
7357 #[test]
7358 fn poll_datagram_single_read() {
7360 let mut buf = [0; 65535];
7361
7362 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
7363 config
7364 .load_cert_chain_from_pem_file("examples/cert.crt")
7365 .unwrap();
7366 config
7367 .load_priv_key_from_pem_file("examples/cert.key")
7368 .unwrap();
7369 config.set_application_protos(&[b"h3"]).unwrap();
7370 config.set_initial_max_data(1500);
7371 config.set_initial_max_stream_data_bidi_local(150);
7372 config.set_initial_max_stream_data_bidi_remote(150);
7373 config.set_initial_max_stream_data_uni(150);
7374 config.set_initial_max_streams_bidi(100);
7375 config.set_initial_max_streams_uni(5);
7376 config.verify_peer(false);
7377 config.enable_dgram(true, 100, 100);
7378
7379 let h3_config = Config::new().unwrap();
7380 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7381 s.handshake().unwrap();
7382
7383 let result = (11, 0, 1);
7385
7386 let (stream, req) = s.send_request(false).unwrap();
7388
7389 let body = s.send_body_client(stream, true).unwrap();
7390
7391 let mut recv_buf = vec![0; body.len()];
7392
7393 let ev_headers = Event::Headers {
7394 list: req,
7395 more_frames: true,
7396 };
7397
7398 s.send_dgram_client(0).unwrap();
7399
7400 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
7401 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
7402
7403 assert_eq!(s.poll_server(), Err(Error::Done));
7404
7405 assert_eq!(s.recv_dgram_server(&mut buf), Ok(result));
7406
7407 assert_eq!(s.poll_server(), Err(Error::Done));
7408
7409 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
7410 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
7411 assert_eq!(s.poll_server(), Err(Error::Done));
7412
7413 let resp = s.send_response(stream, false).unwrap();
7415
7416 let body = s.send_body_server(stream, true).unwrap();
7417
7418 let mut recv_buf = vec![0; body.len()];
7419
7420 let ev_headers = Event::Headers {
7421 list: resp,
7422 more_frames: true,
7423 };
7424
7425 s.send_dgram_server(0).unwrap();
7426
7427 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
7428 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
7429
7430 assert_eq!(s.poll_client(), Err(Error::Done));
7431
7432 assert_eq!(s.recv_dgram_client(&mut buf), Ok(result));
7433
7434 assert_eq!(s.poll_client(), Err(Error::Done));
7435
7436 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
7437
7438 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
7439 assert_eq!(s.poll_client(), Err(Error::Done));
7440 }
7441
7442 #[test]
7443 fn poll_datagram_multi_read() {
7445 let mut buf = [0; 65535];
7446
7447 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
7448 config
7449 .load_cert_chain_from_pem_file("examples/cert.crt")
7450 .unwrap();
7451 config
7452 .load_priv_key_from_pem_file("examples/cert.key")
7453 .unwrap();
7454 config.set_application_protos(&[b"h3"]).unwrap();
7455 config.set_initial_max_data(1500);
7456 config.set_initial_max_stream_data_bidi_local(150);
7457 config.set_initial_max_stream_data_bidi_remote(150);
7458 config.set_initial_max_stream_data_uni(150);
7459 config.set_initial_max_streams_bidi(100);
7460 config.set_initial_max_streams_uni(5);
7461 config.verify_peer(false);
7462 config.enable_dgram(true, 100, 100);
7463
7464 let h3_config = Config::new().unwrap();
7465 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7466 s.handshake().unwrap();
7467
7468 let flow_0_result = (11, 0, 1);
7470 let flow_2_result = (11, 2, 1);
7471
7472 let (stream, req) = s.send_request(false).unwrap();
7474
7475 let body = s.send_body_client(stream, true).unwrap();
7476
7477 let mut recv_buf = vec![0; body.len()];
7478
7479 let ev_headers = Event::Headers {
7480 list: req,
7481 more_frames: true,
7482 };
7483
7484 s.send_dgram_client(0).unwrap();
7485 s.send_dgram_client(0).unwrap();
7486 s.send_dgram_client(0).unwrap();
7487 s.send_dgram_client(0).unwrap();
7488 s.send_dgram_client(0).unwrap();
7489 s.send_dgram_client(2).unwrap();
7490 s.send_dgram_client(2).unwrap();
7491 s.send_dgram_client(2).unwrap();
7492 s.send_dgram_client(2).unwrap();
7493 s.send_dgram_client(2).unwrap();
7494
7495 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
7496 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
7497
7498 assert_eq!(s.poll_server(), Err(Error::Done));
7499
7500 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7502 assert_eq!(s.poll_server(), Err(Error::Done));
7503 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7504 assert_eq!(s.poll_server(), Err(Error::Done));
7505 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7506 assert_eq!(s.poll_server(), Err(Error::Done));
7507
7508 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
7509 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
7510
7511 assert_eq!(s.poll_server(), Err(Error::Done));
7512
7513 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7515 assert_eq!(s.poll_server(), Err(Error::Done));
7516 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7517 assert_eq!(s.poll_server(), Err(Error::Done));
7518 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7519 assert_eq!(s.poll_server(), Err(Error::Done));
7520 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7521 assert_eq!(s.poll_server(), Err(Error::Done));
7522 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7523 assert_eq!(s.poll_server(), Err(Error::Done));
7524 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7525 assert_eq!(s.poll_server(), Err(Error::Done));
7526 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7527 assert_eq!(s.poll_server(), Err(Error::Done));
7528
7529 let resp = s.send_response(stream, false).unwrap();
7531
7532 let body = s.send_body_server(stream, true).unwrap();
7533
7534 let mut recv_buf = vec![0; body.len()];
7535
7536 let ev_headers = Event::Headers {
7537 list: resp,
7538 more_frames: true,
7539 };
7540
7541 s.send_dgram_server(0).unwrap();
7542 s.send_dgram_server(0).unwrap();
7543 s.send_dgram_server(0).unwrap();
7544 s.send_dgram_server(0).unwrap();
7545 s.send_dgram_server(0).unwrap();
7546 s.send_dgram_server(2).unwrap();
7547 s.send_dgram_server(2).unwrap();
7548 s.send_dgram_server(2).unwrap();
7549 s.send_dgram_server(2).unwrap();
7550 s.send_dgram_server(2).unwrap();
7551
7552 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
7553 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
7554
7555 assert_eq!(s.poll_client(), Err(Error::Done));
7556
7557 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_0_result));
7559 assert_eq!(s.poll_client(), Err(Error::Done));
7560 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_0_result));
7561 assert_eq!(s.poll_client(), Err(Error::Done));
7562 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_0_result));
7563 assert_eq!(s.poll_client(), Err(Error::Done));
7564
7565 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
7566 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
7567
7568 assert_eq!(s.poll_client(), Err(Error::Done));
7569
7570 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_0_result));
7572 assert_eq!(s.poll_client(), Err(Error::Done));
7573 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_0_result));
7574 assert_eq!(s.poll_client(), Err(Error::Done));
7575 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_2_result));
7576 assert_eq!(s.poll_client(), Err(Error::Done));
7577 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_2_result));
7578 assert_eq!(s.poll_client(), Err(Error::Done));
7579 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_2_result));
7580 assert_eq!(s.poll_client(), Err(Error::Done));
7581 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_2_result));
7582 assert_eq!(s.poll_client(), Err(Error::Done));
7583 assert_eq!(s.recv_dgram_client(&mut buf), Ok(flow_2_result));
7584 assert_eq!(s.poll_client(), Err(Error::Done));
7585 }
7586
7587 #[test]
7588 fn finished_is_for_requests() {
7591 let mut s = Session::new().unwrap();
7592 s.handshake().unwrap();
7593
7594 assert_eq!(s.poll_client(), Err(Error::Done));
7595 assert_eq!(s.poll_server(), Err(Error::Done));
7596
7597 assert_eq!(s.client.open_grease_stream(&mut s.pipe.client), Ok(()));
7598 assert_eq!(s.pipe.advance(), Ok(()));
7599
7600 assert_eq!(s.poll_client(), Err(Error::Done));
7601 assert_eq!(s.poll_server(), Err(Error::Done));
7602 }
7603
7604 #[test]
7605 fn unknown_uni_stream_leaks_past_max_streams_uni() {
7606 let (mut config, h3_config) = Session::default_configs().unwrap();
7607 config.set_initial_max_data(100_000);
7608 config.set_initial_max_stream_data_uni(100_000);
7609 config.set_initial_max_streams_uni(5);
7610 config.grease(false);
7611
7612 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7613 s.handshake().unwrap();
7614
7615 let baseline = s.server.streams.len();
7616 let mut leaked_streams = Vec::new();
7617 let mut id = 14;
7618
7619 let n = 500;
7620 for i in 0..n {
7621 s.pipe
7622 .client
7623 .stream_send(id, &[0x21], true)
7624 .unwrap_or_else(|e| {
7625 panic!(
7626 "stream credit was not re-issued after {} streams: {:?}",
7627 i, e
7628 )
7629 });
7630
7631 s.pipe.advance().unwrap();
7632 assert_eq!(s.poll_server(), Err(Error::Done));
7633 assert!(s.pipe.server.streams.is_collected(id));
7634 s.pipe.advance().unwrap();
7635
7636 leaked_streams.push(id);
7637 id += 4;
7638 }
7639
7640 for id in leaked_streams {
7641 assert!(s.pipe.server.streams.is_collected(id));
7642 }
7643
7644 assert_eq!(
7645 s.server.streams.len(),
7646 baseline,
7647 "expected {} allocated streams in stream map after {} streams with unknown type",
7648 baseline,
7649 n,
7650 );
7651 }
7652
7653 #[test]
7654 fn empty_uni_stream_leaks_past_max_streams_uni() {
7655 let (mut config, h3_config) = Session::default_configs().unwrap();
7656 config.set_initial_max_data(100_000);
7657 config.set_initial_max_stream_data_uni(100_000);
7658 config.set_initial_max_streams_uni(5);
7659 config.grease(false);
7660
7661 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7662 s.handshake().unwrap();
7663
7664 let baseline = s.server.streams.len();
7665 let mut leaked_streams = Vec::new();
7666 let mut id = 14;
7667
7668 let n = 500;
7669 for i in 0..n {
7670 s.pipe
7671 .client
7672 .stream_send(id, &[], true)
7673 .unwrap_or_else(|e| {
7674 panic!(
7675 "stream credit was not re-issued after {} streams: {:?}",
7676 i, e
7677 )
7678 });
7679
7680 s.pipe.advance().unwrap();
7681 assert_eq!(s.poll_server(), Err(Error::Done));
7682 assert!(s.pipe.server.streams.is_collected(id));
7683 s.pipe.advance().unwrap();
7684
7685 leaked_streams.push(id);
7686 id += 4;
7687 }
7688
7689 for id in leaked_streams {
7690 assert!(s.pipe.server.streams.is_collected(id));
7691 }
7692
7693 assert_eq!(
7694 s.server.streams.len(),
7695 baseline,
7696 "expected {} allocated streams in stream map after {} streams without stream type",
7697 baseline,
7698 n,
7699 );
7700 }
7701
7702 #[test]
7703 fn finished_once() {
7705 let mut s = Session::new().unwrap();
7706 s.handshake().unwrap();
7707
7708 let (stream, req) = s.send_request(false).unwrap();
7709 let body = s.send_body_client(stream, true).unwrap();
7710
7711 let mut recv_buf = vec![0; body.len()];
7712
7713 let ev_headers = Event::Headers {
7714 list: req,
7715 more_frames: true,
7716 };
7717
7718 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
7719 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
7720
7721 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
7722 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
7723
7724 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Err(Error::Done));
7725 assert_eq!(s.poll_server(), Err(Error::Done));
7726 }
7727
7728 #[test]
7729 fn data_event_rearm() {
7731 let bytes = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
7732
7733 let mut s = Session::new().unwrap();
7734 s.handshake().unwrap();
7735
7736 let (r1_id, r1_hdrs) = s.send_request(false).unwrap();
7737
7738 let mut recv_buf = vec![0; bytes.len()];
7739
7740 let r1_ev_headers = Event::Headers {
7741 list: r1_hdrs,
7742 more_frames: true,
7743 };
7744
7745 {
7748 let mut d = [42; 10];
7749 let mut b = octets::OctetsMut::with_slice(&mut d);
7750
7751 b.put_varint(frame::DATA_FRAME_TYPE_ID).unwrap();
7752 b.put_varint(bytes.len() as u64).unwrap();
7753 let off = b.off();
7754 s.pipe.client.stream_send(r1_id, &d[..off], false).unwrap();
7755
7756 assert_eq!(
7757 s.pipe.client.stream_send(r1_id, &bytes[..5], false),
7758 Ok(5)
7759 );
7760
7761 s.advance().ok();
7762 }
7763
7764 assert_eq!(s.poll_server(), Ok((r1_id, r1_ev_headers)));
7765 assert_eq!(s.poll_server(), Ok((r1_id, Event::Data)));
7766 assert_eq!(s.poll_server(), Err(Error::Done));
7767
7768 assert_eq!(s.recv_body_server(r1_id, &mut recv_buf), Ok(5));
7770
7771 assert_eq!(s.pipe.client.stream_send(r1_id, &bytes[5..], false), Ok(5));
7773 s.advance().ok();
7774
7775 assert_eq!(s.poll_server(), Ok((r1_id, Event::Data)));
7776 assert_eq!(s.poll_server(), Err(Error::Done));
7777
7778 assert_eq!(s.recv_body_server(r1_id, &mut recv_buf), Ok(5));
7780 assert_eq!(s.poll_server(), Err(Error::Done));
7781
7782 let r1_body = s.send_body_client(r1_id, false).unwrap();
7784
7785 assert_eq!(s.poll_server(), Ok((r1_id, Event::Data)));
7786 assert_eq!(s.poll_server(), Err(Error::Done));
7787
7788 assert_eq!(s.recv_body_server(r1_id, &mut recv_buf), Ok(r1_body.len()));
7789
7790 let (r2_id, r2_hdrs) = s.send_request(false).unwrap();
7792 let r2_ev_headers = Event::Headers {
7793 list: r2_hdrs,
7794 more_frames: true,
7795 };
7796 let r2_body = s.send_body_client(r2_id, false).unwrap();
7797
7798 s.advance().ok();
7799
7800 assert_eq!(s.poll_server(), Ok((r2_id, r2_ev_headers)));
7801 assert_eq!(s.poll_server(), Ok((r2_id, Event::Data)));
7802 assert_eq!(s.recv_body_server(r2_id, &mut recv_buf), Ok(r2_body.len()));
7803 assert_eq!(s.poll_server(), Err(Error::Done));
7804
7805 let r1_body = s.send_body_client(r1_id, false).unwrap();
7807
7808 let trailers = vec![Header::new(b"hello", b"world")];
7809
7810 s.client
7811 .send_headers(&mut s.pipe.client, r1_id, &trailers, true)
7812 .unwrap();
7813
7814 let r1_ev_trailers = Event::Headers {
7815 list: trailers.clone(),
7816 more_frames: false,
7817 };
7818
7819 s.advance().ok();
7820
7821 assert_eq!(s.poll_server(), Ok((r1_id, Event::Data)));
7822 assert_eq!(s.recv_body_server(r1_id, &mut recv_buf), Ok(r1_body.len()));
7823
7824 assert_eq!(s.poll_server(), Ok((r1_id, r1_ev_trailers)));
7825 assert_eq!(s.poll_server(), Ok((r1_id, Event::Finished)));
7826 assert_eq!(s.poll_server(), Err(Error::Done));
7827
7828 let r2_body = s.send_body_client(r2_id, false).unwrap();
7830
7831 s.client
7832 .send_headers(&mut s.pipe.client, r2_id, &trailers, false)
7833 .unwrap();
7834
7835 let r2_ev_trailers = Event::Headers {
7836 list: trailers,
7837 more_frames: true,
7838 };
7839
7840 s.advance().ok();
7841
7842 assert_eq!(s.poll_server(), Ok((r2_id, Event::Data)));
7843 assert_eq!(s.recv_body_server(r2_id, &mut recv_buf), Ok(r2_body.len()));
7844 assert_eq!(s.poll_server(), Ok((r2_id, r2_ev_trailers)));
7845 assert_eq!(s.poll_server(), Err(Error::Done));
7846
7847 let (r3_id, r3_hdrs) = s.send_request(false).unwrap();
7848
7849 let r3_ev_headers = Event::Headers {
7850 list: r3_hdrs,
7851 more_frames: true,
7852 };
7853
7854 {
7856 let mut d = [42; 10];
7857 let mut b = octets::OctetsMut::with_slice(&mut d);
7858
7859 b.put_varint(frame::DATA_FRAME_TYPE_ID).unwrap();
7860 b.put_varint(bytes.len() as u64).unwrap();
7861 let off = b.off();
7862 s.pipe.client.stream_send(r3_id, &d[..off], false).unwrap();
7863
7864 s.advance().ok();
7865 }
7866
7867 assert_eq!(s.poll_server(), Ok((r3_id, r3_ev_headers)));
7868 assert_eq!(s.poll_server(), Ok((r3_id, Event::Data)));
7869 assert_eq!(s.poll_server(), Err(Error::Done));
7870
7871 assert_eq!(s.recv_body_server(r3_id, &mut recv_buf), Err(Error::Done));
7872
7873 assert_eq!(s.pipe.client.stream_send(r3_id, &bytes[..5], false), Ok(5));
7874
7875 s.advance().ok();
7876
7877 assert_eq!(s.poll_server(), Ok((r3_id, Event::Data)));
7878 assert_eq!(s.poll_server(), Err(Error::Done));
7879
7880 assert_eq!(s.recv_body_server(r3_id, &mut recv_buf), Ok(5));
7881
7882 assert_eq!(s.pipe.client.stream_send(r3_id, &bytes[5..], false), Ok(5));
7883 s.advance().ok();
7884
7885 assert_eq!(s.poll_server(), Ok((r3_id, Event::Data)));
7886 assert_eq!(s.poll_server(), Err(Error::Done));
7887
7888 assert_eq!(s.recv_body_server(r3_id, &mut recv_buf), Ok(5));
7889
7890 let body = s.send_body_client(r3_id, false).unwrap();
7892 s.send_body_client(r3_id, false).unwrap();
7893 s.send_body_client(r3_id, false).unwrap();
7894
7895 assert_eq!(s.poll_server(), Ok((r3_id, Event::Data)));
7896 assert_eq!(s.poll_server(), Err(Error::Done));
7897
7898 {
7899 let mut d = [42; 10];
7900 let mut b = octets::OctetsMut::with_slice(&mut d);
7901
7902 b.put_varint(frame::DATA_FRAME_TYPE_ID).unwrap();
7903 b.put_varint(0).unwrap();
7904 let off = b.off();
7905 s.pipe.client.stream_send(r3_id, &d[..off], true).unwrap();
7906
7907 s.advance().ok();
7908 }
7909
7910 let mut recv_buf = vec![0; bytes.len() * 3];
7911
7912 assert_eq!(s.recv_body_server(r3_id, &mut recv_buf), Ok(body.len() * 3));
7913 }
7914
7915 #[test]
7916 fn dgram_event_rearm() {
7918 let mut buf = [0; 65535];
7919
7920 let mut config = crate::Config::new(crate::PROTOCOL_VERSION).unwrap();
7921 config
7922 .load_cert_chain_from_pem_file("examples/cert.crt")
7923 .unwrap();
7924 config
7925 .load_priv_key_from_pem_file("examples/cert.key")
7926 .unwrap();
7927 config.set_application_protos(&[b"h3"]).unwrap();
7928 config.set_initial_max_data(1500);
7929 config.set_initial_max_stream_data_bidi_local(150);
7930 config.set_initial_max_stream_data_bidi_remote(150);
7931 config.set_initial_max_stream_data_uni(150);
7932 config.set_initial_max_streams_bidi(100);
7933 config.set_initial_max_streams_uni(5);
7934 config.verify_peer(false);
7935 config.enable_dgram(true, 100, 100);
7936
7937 let h3_config = Config::new().unwrap();
7938 let mut s = Session::with_configs(&mut config, &h3_config).unwrap();
7939 s.handshake().unwrap();
7940
7941 let flow_0_result = (11, 0, 1);
7943 let flow_2_result = (11, 2, 1);
7944
7945 let (stream, req) = s.send_request(false).unwrap();
7947
7948 let body = s.send_body_client(stream, true).unwrap();
7949
7950 let mut recv_buf = vec![0; body.len()];
7951
7952 let ev_headers = Event::Headers {
7953 list: req,
7954 more_frames: true,
7955 };
7956
7957 s.send_dgram_client(0).unwrap();
7958 s.send_dgram_client(0).unwrap();
7959 s.send_dgram_client(2).unwrap();
7960 s.send_dgram_client(2).unwrap();
7961
7962 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
7963 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
7964
7965 assert_eq!(s.poll_server(), Err(Error::Done));
7966 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7967
7968 assert_eq!(s.poll_server(), Err(Error::Done));
7969 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7970
7971 assert_eq!(s.poll_server(), Err(Error::Done));
7972 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7973
7974 assert_eq!(s.poll_server(), Err(Error::Done));
7975 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7976
7977 assert_eq!(s.poll_server(), Err(Error::Done));
7978
7979 s.send_dgram_client(0).unwrap();
7980 s.send_dgram_client(2).unwrap();
7981
7982 assert_eq!(s.poll_server(), Err(Error::Done));
7983
7984 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_0_result));
7985 assert_eq!(s.poll_server(), Err(Error::Done));
7986
7987 assert_eq!(s.recv_dgram_server(&mut buf), Ok(flow_2_result));
7988 assert_eq!(s.poll_server(), Err(Error::Done));
7989
7990 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
7991 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
7992
7993 assert_eq!(s.pipe.client.dgram_sent_count, 6);
7995 assert_eq!(s.pipe.client.dgram_recv_count, 0);
7996 assert_eq!(s.pipe.server.dgram_sent_count, 0);
7997 assert_eq!(s.pipe.server.dgram_recv_count, 6);
7998
7999 let server_path = s.pipe.server.paths.get_active().expect("no active");
8000 let client_path = s.pipe.client.paths.get_active().expect("no active");
8001 assert_eq!(client_path.dgram_sent_count, 6);
8002 assert_eq!(client_path.dgram_recv_count, 0);
8003 assert_eq!(server_path.dgram_sent_count, 0);
8004 assert_eq!(server_path.dgram_recv_count, 6);
8005 }
8006
8007 #[test]
8008 fn reset_stream() {
8009 let mut buf = [0; 65535];
8010
8011 let mut s = Session::new().unwrap();
8012 s.handshake().unwrap();
8013
8014 let (stream, req) = s.send_request(false).unwrap();
8016
8017 let ev_headers = Event::Headers {
8018 list: req,
8019 more_frames: true,
8020 };
8021
8022 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8024 assert_eq!(s.poll_server(), Err(Error::Done));
8025
8026 let resp = s.send_response(stream, true).unwrap();
8027
8028 let ev_headers = Event::Headers {
8029 list: resp,
8030 more_frames: false,
8031 };
8032
8033 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8034 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
8035 assert_eq!(s.poll_client(), Err(Error::Done));
8036
8037 let frames = [crate::frame::Frame::ResetStream {
8039 stream_id: stream,
8040 error_code: 42,
8041 final_size: 68,
8042 }];
8043
8044 let pkt_type = crate::packet::Type::Short;
8045 assert_eq!(
8046 s.pipe.send_pkt_to_server(pkt_type, &frames, &mut buf),
8047 Ok(39)
8048 );
8049
8050 assert_eq!(s.poll_server(), Ok((stream, Event::Reset(42))));
8052 assert_eq!(s.poll_server(), Err(Error::Done));
8053
8054 assert_eq!(
8056 s.pipe.send_pkt_to_server(pkt_type, &frames, &mut buf),
8057 Ok(39)
8058 );
8059
8060 assert_eq!(s.poll_server(), Err(Error::Done));
8061 }
8062
8063 #[test]
8066 fn client_shutdown_write_server_fin() {
8067 let mut buf = [0; 65535];
8068 let mut s = Session::new().unwrap();
8069 s.handshake().unwrap();
8070
8071 let (stream, req) = s.send_request(false).unwrap();
8073
8074 let ev_headers = Event::Headers {
8075 list: req,
8076 more_frames: true,
8077 };
8078
8079 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8081 assert_eq!(s.poll_server(), Err(Error::Done));
8082
8083 let resp = s.send_response(stream, true).unwrap();
8084
8085 let ev_headers = Event::Headers {
8086 list: resp,
8087 more_frames: false,
8088 };
8089
8090 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8091 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
8092 assert_eq!(s.poll_client(), Err(Error::Done));
8093
8094 assert_eq!(
8096 s.pipe
8097 .client
8098 .stream_shutdown(stream, crate::Shutdown::Write, 42),
8099 Ok(())
8100 );
8101 assert_eq!(s.advance(), Ok(()));
8102
8103 assert_eq!(s.poll_server(), Ok((stream, Event::Reset(42))));
8105 assert_eq!(s.poll_server(), Err(Error::Done));
8106
8107 assert!(s.pipe.server.streams.is_collected(stream));
8109 assert!(s.pipe.client.streams.is_collected(stream));
8110
8111 let (stream, req) = s.send_request(false).unwrap();
8114
8115 let ev_headers = Event::Headers {
8116 list: req,
8117 more_frames: true,
8118 };
8119
8120 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8122 assert_eq!(s.poll_server(), Err(Error::Done));
8123
8124 let resp = s.send_response(stream, false).unwrap();
8126
8127 let ev_headers = Event::Headers {
8128 list: resp,
8129 more_frames: true,
8130 };
8131
8132 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8133 assert_eq!(s.poll_client(), Err(Error::Done));
8134
8135 assert_eq!(
8137 s.pipe
8138 .client
8139 .stream_shutdown(stream, crate::Shutdown::Write, 42),
8140 Ok(())
8141 );
8142 assert_eq!(s.advance(), Ok(()));
8143
8144 assert_eq!(s.poll_server(), Ok((stream, Event::Reset(42))));
8146 assert_eq!(s.poll_server(), Err(Error::Done));
8147
8148 s.send_body_server(stream, true).unwrap();
8150
8151 assert!(s.pipe.server.streams.is_collected(stream));
8153 assert!(!s.pipe.client.streams.is_collected(stream));
8156 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
8157 s.recv_body_client(stream, &mut buf).unwrap();
8158 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
8159 assert_eq!(s.poll_client(), Err(Error::Done));
8160 assert!(s.pipe.client.streams.is_collected(stream));
8161 }
8162
8163 #[test]
8164 fn client_shutdown_read() {
8165 let mut buf = [0; 65535];
8166 let mut s = Session::new().unwrap();
8167 s.handshake().unwrap();
8168
8169 let (stream, req) = s.send_request(false).unwrap();
8171
8172 let ev_headers = Event::Headers {
8173 list: req,
8174 more_frames: true,
8175 };
8176
8177 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8179 assert_eq!(s.poll_server(), Err(Error::Done));
8180
8181 let resp = s.send_response(stream, false).unwrap();
8182
8183 let ev_headers = Event::Headers {
8184 list: resp,
8185 more_frames: true,
8186 };
8187
8188 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8189 assert_eq!(s.poll_client(), Err(Error::Done));
8190 assert_eq!(
8192 s.pipe
8193 .client
8194 .stream_shutdown(stream, crate::Shutdown::Read, 42),
8195 Ok(())
8196 );
8197 assert_eq!(s.advance(), Ok(()));
8198
8199 assert_eq!(s.poll_server(), Err(Error::Done));
8201 let writables: Vec<u64> = s.pipe.server.writable().collect();
8202 assert!(writables.contains(&stream));
8203 assert_eq!(
8204 s.send_body_server(stream, false),
8205 Err(Error::TransportError(crate::Error::StreamStopped(42)))
8206 );
8207
8208 assert_eq!(
8210 s.client.send_body(&mut s.pipe.client, stream, &[], true),
8211 Ok(0)
8212 );
8213 assert_eq!(s.advance(), Ok(()));
8214 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
8217 assert_eq!(s.recv_body_server(stream, &mut buf), Err(Error::Done));
8218 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
8219 assert_eq!(s.poll_server(), Err(Error::Done));
8220
8221 assert!(s.pipe.client.streams.is_collected(stream));
8224 assert!(s.pipe.server.streams.is_collected(stream));
8225 }
8226
8227 #[test]
8228 fn reset_finished_at_server() {
8229 let mut s = Session::new().unwrap();
8230 s.handshake().unwrap();
8231
8232 let (stream, _req) = s.send_request(false).unwrap();
8234
8235 assert_eq!(
8237 s.pipe.client.stream_shutdown(0, crate::Shutdown::Write, 0),
8238 Ok(())
8239 );
8240
8241 assert_eq!(s.pipe.advance(), Ok(()));
8242
8243 assert_eq!(s.poll_server(), Ok((stream, Event::Reset(0))));
8245 assert_eq!(s.poll_server(), Err(Error::Done));
8246
8247 let (stream, req) = s.send_request(true).unwrap();
8249
8250 assert_eq!(
8252 s.pipe.client.stream_shutdown(4, crate::Shutdown::Write, 0),
8253 Ok(())
8254 );
8255
8256 let ev_headers = Event::Headers {
8257 list: req,
8258 more_frames: false,
8259 };
8260
8261 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8263 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
8264 assert_eq!(s.poll_server(), Err(Error::Done));
8265 }
8266
8267 #[test]
8268 fn reset_finished_at_server_with_data_pending() {
8269 let mut s = Session::new().unwrap();
8270 s.handshake().unwrap();
8271
8272 let (stream, req) = s.send_request(false).unwrap();
8274
8275 assert!(s.send_body_client(stream, false).is_ok());
8276
8277 assert_eq!(s.pipe.advance(), Ok(()));
8278
8279 let ev_headers = Event::Headers {
8280 list: req,
8281 more_frames: true,
8282 };
8283
8284 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8286 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
8287
8288 assert_eq!(
8290 s.pipe
8291 .client
8292 .stream_shutdown(stream, crate::Shutdown::Write, 0),
8293 Ok(())
8294 );
8295
8296 assert_eq!(s.pipe.advance(), Ok(()));
8297
8298 assert_eq!(s.poll_server(), Ok((stream, Event::Reset(0))));
8302 assert_eq!(s.poll_server(), Err(Error::Done));
8303 assert_eq!(s.pipe.server.readable().len(), 0);
8304 }
8305
8306 #[test]
8307 fn reset_finished_at_server_with_data_pending_2() {
8308 let mut s = Session::new().unwrap();
8309 s.handshake().unwrap();
8310
8311 let (stream, req) = s.send_request(false).unwrap();
8313
8314 assert!(s.send_body_client(stream, false).is_ok());
8315
8316 assert_eq!(s.pipe.advance(), Ok(()));
8317
8318 let ev_headers = Event::Headers {
8319 list: req,
8320 more_frames: true,
8321 };
8322
8323 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8325 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
8326
8327 assert_eq!(
8329 s.pipe
8330 .client
8331 .stream_shutdown(stream, crate::Shutdown::Write, 0),
8332 Ok(())
8333 );
8334
8335 assert_eq!(s.pipe.advance(), Ok(()));
8336
8337 assert_eq!(
8340 s.recv_body_server(stream, &mut [0; 100]),
8341 Err(Error::TransportError(crate::Error::StreamReset(0)))
8342 );
8343
8344 assert_eq!(s.poll_server(), Err(Error::Done));
8346 assert_eq!(s.pipe.server.readable().len(), 0);
8347 }
8348
8349 #[test]
8350 fn reset_finished_at_client() {
8351 let mut buf = [0; 65535];
8352 let mut s = Session::new().unwrap();
8353 s.handshake().unwrap();
8354
8355 let (stream, req) = s.send_request(false).unwrap();
8357
8358 let ev_headers = Event::Headers {
8359 list: req,
8360 more_frames: true,
8361 };
8362
8363 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8365 assert_eq!(s.poll_server(), Err(Error::Done));
8366
8367 s.send_response(stream, false).unwrap();
8369
8370 assert_eq!(s.pipe.advance(), Ok(()));
8371
8372 assert_eq!(
8374 s.pipe
8375 .server
8376 .stream_shutdown(stream, crate::Shutdown::Write, 0),
8377 Ok(())
8378 );
8379
8380 assert_eq!(s.pipe.advance(), Ok(()));
8381
8382 assert_eq!(s.poll_client(), Ok((stream, Event::Reset(0))));
8384 assert_eq!(s.poll_server(), Err(Error::Done));
8385
8386 let (stream, req) = s.send_request(true).unwrap();
8388
8389 let ev_headers = Event::Headers {
8390 list: req,
8391 more_frames: false,
8392 };
8393
8394 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8396 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
8397 assert_eq!(s.poll_server(), Err(Error::Done));
8398
8399 let resp = s.send_response(stream, true).unwrap();
8401
8402 assert_eq!(s.pipe.advance(), Ok(()));
8403
8404 let frames = [crate::frame::Frame::ResetStream {
8406 stream_id: stream,
8407 error_code: 42,
8408 final_size: 68,
8409 }];
8410
8411 let pkt_type = crate::packet::Type::Short;
8412 assert_eq!(
8413 s.pipe.send_pkt_to_server(pkt_type, &frames, &mut buf),
8414 Ok(39)
8415 );
8416
8417 assert_eq!(s.pipe.advance(), Ok(()));
8418
8419 let ev_headers = Event::Headers {
8420 list: resp,
8421 more_frames: false,
8422 };
8423
8424 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8426 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
8427 assert_eq!(s.poll_client(), Err(Error::Done));
8428 }
8429
8430 #[test]
8431 fn collect_completed_streams() {
8432 let mut s = Session::new().unwrap();
8433 s.handshake().unwrap();
8434
8435 let init_streams_client = s.client.streams.len();
8436 let init_streams_server = s.server.streams.len();
8437
8438 let (stream, req) = s.send_request(false).unwrap();
8440
8441 let ev_headers = Event::Headers {
8442 list: req,
8443 more_frames: true,
8444 };
8445
8446 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8448 assert_eq!(s.poll_server(), Err(Error::Done));
8449
8450 assert_eq!(s.client.streams.len(), init_streams_client + 1);
8451 assert_eq!(s.server.streams.len(), init_streams_server + 1);
8452
8453 let body = s.send_body_client(stream, true).unwrap();
8455
8456 let mut recv_buf = vec![0; body.len()];
8457
8458 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
8459 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
8460
8461 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
8462
8463 assert_eq!(s.client.streams.len(), init_streams_client + 1);
8464 assert_eq!(s.server.streams.len(), init_streams_server + 1);
8465
8466 let resp_headers = s.send_response(stream, false).unwrap();
8468 s.send_body_server(stream, true).unwrap();
8469
8470 let ev_headers = Event::Headers {
8471 list: resp_headers,
8472 more_frames: true,
8473 };
8474
8475 assert_eq!(s.poll_client(), Ok((stream, ev_headers)));
8476 assert_eq!(s.poll_client(), Ok((stream, Event::Data)));
8477 assert_eq!(s.recv_body_client(stream, &mut recv_buf), Ok(body.len()));
8478
8479 assert_eq!(s.client.streams.len(), init_streams_client + 1);
8481 assert_eq!(s.server.streams.len(), init_streams_server);
8482
8483 assert_eq!(s.poll_client(), Ok((stream, Event::Finished)));
8485 assert_eq!(s.poll_client(), Err(Error::Done));
8486
8487 assert_eq!(s.client.streams.len(), init_streams_client);
8488 }
8489
8490 #[test]
8491 fn collect_reset_streams() {
8492 let mut s = Session::new().unwrap();
8493 s.handshake().unwrap();
8494
8495 let init_streams_client = s.client.streams.len();
8496 let init_streams_server = s.server.streams.len();
8497
8498 let (stream, req) = s.send_request(false).unwrap();
8500
8501 let ev_headers = Event::Headers {
8502 list: req,
8503 more_frames: true,
8504 };
8505
8506 assert_eq!(s.poll_server(), Ok((stream, ev_headers)));
8508 assert_eq!(s.poll_server(), Err(Error::Done));
8509
8510 assert_eq!(s.client.streams.len(), init_streams_client + 1);
8511 assert_eq!(s.server.streams.len(), init_streams_server + 1);
8512
8513 let body = s.send_body_client(stream, true).unwrap();
8515
8516 let mut recv_buf = vec![0; body.len()];
8517
8518 assert_eq!(s.poll_server(), Ok((stream, Event::Data)));
8519 assert_eq!(s.recv_body_server(stream, &mut recv_buf), Ok(body.len()));
8520
8521 assert_eq!(s.poll_server(), Ok((stream, Event::Finished)));
8522
8523 assert_eq!(s.client.streams.len(), init_streams_client + 1);
8524 assert_eq!(s.server.streams.len(), init_streams_server + 1);
8525
8526 s.send_response(stream, false).unwrap();
8528 s.pipe
8529 .server
8530 .stream_shutdown(stream, crate::Shutdown::Write, 0)
8531 .unwrap();
8532
8533 s.advance().ok();
8534
8535 let _ = s.send_body_server(stream, true);
8544
8545 assert_eq!(s.poll_server(), Err(Error::Done));
8546
8547 assert_eq!(s.poll_client(), Ok((stream, Event::Reset(0))));
8548 assert_eq!(s.poll_client(), Err(Error::Done));
8549
8550 assert_eq!(s.client.streams.len(), init_streams_client);
8552 assert_eq!(s.server.streams.len(), init_streams_server);
8553 }
8554}
8555
8556#[cfg(feature = "ffi")]
8557mod ffi;
8558#[cfg(feature = "internal")]
8559#[doc(hidden)]
8560pub mod frame;
8561#[cfg(not(feature = "internal"))]
8562mod frame;
8563#[doc(hidden)]
8564pub mod qpack;
8565mod stream;