1use super::error::{H2Error, Reason};
20use bytes::{Bytes, BytesMut};
21
22pub const FRAME_HEADER_LEN: usize = 9;
24pub const DEFAULT_MAX_FRAME_SIZE: usize = 16_384;
26pub const MAX_FRAME_SIZE_LIMIT: usize = 16_777_215;
28pub const DEFAULT_INITIAL_WINDOW_SIZE: u32 = 1_048_576;
30pub const MAX_WINDOW_SIZE: u32 = 2_147_483_647;
32pub const CLIENT_PREFACE: &[u8] = b"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n";
34
35const DATA_TYPE: u8 = 0x00;
36const HEADERS_TYPE: u8 = 0x01;
37const PRIORITY_TYPE: u8 = 0x02;
38const RST_STREAM_TYPE: u8 = 0x03;
39const SETTINGS_TYPE: u8 = 0x04;
40const PUSH_PROMISE_TYPE: u8 = 0x05;
41const PING_TYPE: u8 = 0x06;
42const GOAWAY_TYPE: u8 = 0x07;
43const WINDOW_UPDATE_TYPE: u8 = 0x08;
44const CONTINUATION_TYPE: u8 = 0x09;
45
46const FLAG_END_STREAM: u8 = 0x01;
48const FLAG_ACK: u8 = 0x01;
49const FLAG_END_HEADERS: u8 = 0x04;
50const FLAG_PADDED: u8 = 0x08;
51const FLAG_PRIORITY: u8 = 0x20;
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
55pub struct Priority {
56 pub exclusive: bool,
57 pub dependency: u32,
58 pub weight: u8,
59}
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub struct Setting {
64 pub id: u16,
65 pub value: u32,
66}
67
68#[derive(Debug, Clone, PartialEq, Eq)]
74pub enum Frame {
75 Data {
76 stream_id: u32,
77 end_stream: bool,
78 data: Bytes,
79 },
80 Headers {
81 stream_id: u32,
82 end_stream: bool,
83 end_headers: bool,
84 priority: Option<Priority>,
85 block: Bytes,
86 },
87 Priority {
88 stream_id: u32,
89 priority: Priority,
90 },
91 Reset {
92 stream_id: u32,
93 error_code: u32,
94 },
95 Settings {
96 ack: bool,
97 settings: Vec<Setting>,
98 },
99 PushPromise {
100 stream_id: u32,
101 end_headers: bool,
102 promised_stream_id: u32,
103 block: Bytes,
104 },
105 Ping {
106 ack: bool,
107 payload: [u8; 8],
108 },
109 GoAway {
110 last_stream_id: u32,
111 error_code: u32,
112 debug: Bytes,
113 },
114 WindowUpdate {
115 stream_id: u32,
116 increment: u32,
117 },
118 Continuation {
119 stream_id: u32,
120 end_headers: bool,
121 block: Bytes,
122 },
123 Unknown {
126 typ: u8,
127 flags: u8,
128 stream_id: u32,
129 payload: Bytes,
130 },
131}
132
133#[derive(Debug)]
135pub struct FrameDecoder {
136 buf: BytesMut,
137 max_frame_size: usize,
138 block_stream: Option<u32>,
142}
143
144impl FrameDecoder {
145 #[inline]
148 pub fn new(max_frame_size: usize) -> FrameDecoder {
149 FrameDecoder {
150 buf: BytesMut::new(),
151 max_frame_size,
152 block_stream: None,
153 }
154 }
155
156 #[inline]
158 pub fn extend(&mut self, bytes: &[u8]) {
159 self.buf.extend_from_slice(bytes);
160 }
161
162 #[inline]
164 pub fn max_frame_size(&self) -> usize {
165 self.max_frame_size
166 }
167
168 #[inline]
170 pub fn set_max_frame_size(&mut self, max_frame_size: usize) {
171 self.max_frame_size = max_frame_size;
172 }
173
174 #[inline]
176 pub fn block_stream(&self) -> Option<u32> {
177 self.block_stream
178 }
179
180 #[inline]
183 pub fn next_frame(&mut self) -> Result<Option<Frame>, H2Error> {
184 if self.buf.len() < FRAME_HEADER_LEN {
185 return Ok(None);
186 }
187
188 let payload_len = u32::from_be_bytes([0, self.buf[0], self.buf[1], self.buf[2]]) as usize;
191 let typ = self.buf[3];
192 let flags = self.buf[4];
193 let raw_stream_id =
194 u32::from_be_bytes([self.buf[5], self.buf[6], self.buf[7], self.buf[8]]);
195 let stream_id = raw_stream_id & 0x7fff_ffff;
199
200 if payload_len > self.max_frame_size {
201 return Err(H2Error::frame_size(
202 "frame payload exceeds SETTINGS_MAX_FRAME_SIZE",
203 ));
204 }
205
206 let total = FRAME_HEADER_LEN + payload_len;
207 if self.buf.len() < total {
208 return Ok(None);
209 }
210
211 let _header = self.buf.split_to(FRAME_HEADER_LEN);
212 let payload = self.buf.split_to(payload_len).freeze();
213 let frame = parse_frame(typ, flags, stream_id, payload, self)?;
214 Ok(Some(frame))
215 }
216}
217
218#[inline]
220fn parse_frame(
221 typ: u8,
222 flags: u8,
223 stream_id: u32,
224 payload: Bytes,
225 decoder: &mut FrameDecoder,
226) -> Result<Frame, H2Error> {
227 let mut body = &payload[..];
228
229 if let Some(open_stream) = decoder.block_stream {
231 if typ != CONTINUATION_TYPE {
232 return Err(H2Error::protocol(
233 "frame received while field block was open",
234 ));
235 }
236 if stream_id != open_stream {
237 return Err(H2Error::protocol(
238 "CONTINUATION frame on a different stream than the open field block",
239 ));
240 }
241 } else if typ == CONTINUATION_TYPE {
242 return Err(H2Error::protocol(
243 "CONTINUATION frame without a preceding HEADERS or PUSH_PROMISE",
244 ));
245 }
246
247 let frame = match typ {
248 DATA_TYPE => {
249 require_stream(stream_id)?;
250 let end_stream = flags & FLAG_END_STREAM != 0;
251 let pad_len = take_padding(&mut body, flags & FLAG_PADDED != 0)?;
252 Frame::Data {
253 stream_id,
254 end_stream,
255 data: payload_slice(&payload, body, pad_len),
256 }
257 }
258 HEADERS_TYPE => {
259 require_stream(stream_id)?;
260 let end_stream = flags & FLAG_END_STREAM != 0;
261 let end_headers = flags & FLAG_END_HEADERS != 0;
262 let pad_len = take_padding(&mut body, flags & FLAG_PADDED != 0)?;
263 let priority = if flags & FLAG_PRIORITY != 0 {
264 Some(read_priority(&mut body, stream_id)?)
265 } else {
266 None
267 };
268 if !end_headers {
269 decoder.block_stream = Some(stream_id);
270 }
271 Frame::Headers {
272 stream_id,
273 end_stream,
274 end_headers,
275 priority,
276 block: payload_slice(&payload, body, pad_len),
277 }
278 }
279 PRIORITY_TYPE => {
280 require_stream(stream_id)?;
281 if body.len() != 5 {
282 return Err(H2Error::frame_size(
283 "PRIORITY frame payload must be exactly 5 octets",
284 ));
285 }
286 let priority = read_priority(&mut body, stream_id)?;
287 Frame::Priority {
288 stream_id,
289 priority,
290 }
291 }
292 RST_STREAM_TYPE => {
293 require_stream(stream_id)?;
294 if body.len() != 4 {
295 return Err(H2Error::frame_size(
296 "RST_STREAM frame payload must be exactly 4 octets",
297 ));
298 }
299 Frame::Reset {
300 stream_id,
301 error_code: read_u32(body),
302 }
303 }
304 SETTINGS_TYPE => {
305 if stream_id != 0 {
306 return Err(H2Error::protocol(
307 "SETTINGS frame must use stream identifier 0",
308 ));
309 }
310 let ack = flags & FLAG_ACK != 0;
311 if ack && !body.is_empty() {
312 return Err(H2Error::frame_size(
313 "SETTINGS frame with ACK flag must have an empty payload",
314 ));
315 }
316 if !body.len().is_multiple_of(6) {
317 return Err(H2Error::frame_size(
318 "SETTINGS frame payload must be a multiple of 6 octets",
319 ));
320 }
321 let mut settings = Vec::with_capacity(body.len() / 6);
322 while !body.is_empty() {
323 let id = u16::from_be_bytes([body[0], body[1]]);
324 let value = u32::from_be_bytes([body[2], body[3], body[4], body[5]]);
325 validate_setting(id, value)?;
326 settings.push(Setting { id, value });
327 body = &body[6..];
328 }
329 if !ack {
333 for setting in &settings {
334 if setting.id == 0x05 {
335 decoder.set_max_frame_size(setting.value as usize);
336 }
337 }
338 }
339 Frame::Settings { ack, settings }
340 }
341 PUSH_PROMISE_TYPE => {
342 require_stream(stream_id)?;
343 let end_headers = flags & FLAG_END_HEADERS != 0;
344 let pad_len = take_padding(&mut body, flags & FLAG_PADDED != 0)?;
345 if body.len() < 4 {
346 return Err(H2Error::frame_size(
347 "PUSH_PROMISE frame payload must be at least 4 octets",
348 ));
349 }
350 let promised_stream_id = read_u32(&body[..4]) & 0x7fff_ffff;
351 if promised_stream_id == 0 {
352 return Err(H2Error::protocol(
353 "PUSH_PROMISE promised stream identifier is 0",
354 ));
355 }
356 if !end_headers {
357 decoder.block_stream = Some(stream_id);
358 }
359 body = &body[4..];
360 Frame::PushPromise {
361 stream_id,
362 end_headers,
363 promised_stream_id,
364 block: payload_slice(&payload, body, pad_len),
365 }
366 }
367 PING_TYPE => {
368 if stream_id != 0 {
369 return Err(H2Error::protocol("PING frame must use stream identifier 0"));
370 }
371 if body.len() != 8 {
372 return Err(H2Error::frame_size(
373 "PING frame payload must be exactly 8 octets",
374 ));
375 }
376 Frame::Ping {
377 ack: flags & FLAG_ACK != 0,
378 payload: body[..8]
379 .try_into()
380 .expect("PING payload is exactly 8 octets (validated above)"),
381 }
382 }
383 GOAWAY_TYPE => {
384 if stream_id != 0 {
385 return Err(H2Error::protocol(
386 "GOAWAY frame must use stream identifier 0",
387 ));
388 }
389 if body.len() < 8 {
390 return Err(H2Error::frame_size(
391 "GOAWAY frame payload must be at least 8 octets",
392 ));
393 }
394 let last_stream_id = read_u32(&body[..4]) & 0x7fff_ffff;
395 let error_code = read_u32(&body[4..8]);
396 Frame::GoAway {
397 last_stream_id,
398 error_code,
399 debug: payload.slice(8..),
400 }
401 }
402 WINDOW_UPDATE_TYPE => {
403 if body.len() != 4 {
404 return Err(H2Error::frame_size(
405 "WINDOW_UPDATE frame payload must be exactly 4 octets",
406 ));
407 }
408 let increment = read_u32(body) & 0x7fff_ffff;
409 if increment == 0 {
410 return Err(H2Error::protocol("WINDOW_UPDATE frame with zero increment"));
411 }
412 Frame::WindowUpdate {
413 stream_id,
414 increment,
415 }
416 }
417 CONTINUATION_TYPE => {
418 let end_headers = flags & FLAG_END_HEADERS != 0;
419 if end_headers {
420 decoder.block_stream = None;
421 }
422 Frame::Continuation {
423 stream_id,
424 end_headers,
425 block: payload,
426 }
427 }
428 _ => Frame::Unknown {
429 typ,
430 flags,
431 stream_id,
432 payload,
433 },
434 };
435
436 Ok(frame)
437}
438
439#[inline]
443fn payload_slice(payload: &Bytes, body: &[u8], pad_len: usize) -> Bytes {
444 let start = payload.len() - pad_len - body.len();
445 payload.slice(start..start + body.len())
446}
447
448#[inline]
454fn take_padding(body: &mut &[u8], padded: bool) -> Result<usize, H2Error> {
455 if !padded {
456 return Ok(0);
457 }
458 let pad_len = *body
459 .first()
460 .ok_or_else(|| H2Error::frame_size("PADDED frame payload too short for pad length"))?
461 as usize;
462 if pad_len >= body.len() {
463 return Err(H2Error::protocol(
464 "padding length is the length of the frame payload or greater",
465 ));
466 }
467 *body = &body[1..body.len() - pad_len];
468 Ok(pad_len)
469}
470
471#[inline]
474fn read_priority(body: &mut &[u8], stream_id: u32) -> Result<Priority, H2Error> {
475 if body.len() < 5 {
476 return Err(H2Error::frame_size(
477 "frame payload too short for priority fields",
478 ));
479 }
480 let exclusive = body[0] & 0x80 != 0;
481 let dependency = u32::from_be_bytes([body[0], body[1], body[2], body[3]]) & 0x7fff_ffff;
482 if dependency == stream_id {
483 return Err(H2Error::protocol(
484 "stream priority depends on its own stream identifier",
485 ));
486 }
487 let weight = body[4];
488 *body = &body[5..];
489 Ok(Priority {
490 exclusive,
491 dependency,
492 weight,
493 })
494}
495
496#[inline]
497fn require_stream(stream_id: u32) -> Result<(), H2Error> {
498 if stream_id == 0 {
499 Err(H2Error::new(
500 Reason::ProtocolError,
501 "frame must use a non-zero stream identifier",
502 ))
503 } else {
504 Ok(())
505 }
506}
507
508#[inline]
512fn validate_setting(id: u16, value: u32) -> Result<(), H2Error> {
513 match id {
514 0x02 => {
515 if value > 1 {
517 return Err(H2Error::protocol("SETTINGS_ENABLE_PUSH must be 0 or 1"));
518 }
519 }
520 0x04 => {
521 if value > MAX_WINDOW_SIZE {
523 return Err(H2Error::new(
524 Reason::FlowControlError,
525 "SETTINGS_INITIAL_WINDOW_SIZE exceeds 2^31-1",
526 ));
527 }
528 }
529 0x05 if !(DEFAULT_MAX_FRAME_SIZE..=MAX_FRAME_SIZE_LIMIT).contains(&(value as usize)) => {
530 return Err(H2Error::protocol(
532 "SETTINGS_MAX_FRAME_SIZE outside 16384..16777215",
533 ));
534 }
535 _ => {}
536 }
537 Ok(())
538}
539
540#[inline]
541fn read_u32(bytes: &[u8]) -> u32 {
542 u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]])
543}
544
545#[derive(Debug, Clone, Copy, Default)]
549pub struct FrameWriter {
550 pub max_frame_size: usize,
553}
554
555impl FrameWriter {
556 pub const fn new(max_frame_size: usize) -> FrameWriter {
557 FrameWriter { max_frame_size }
558 }
559
560 #[inline]
561 fn header(out: &mut Vec<u8>, payload_len: usize, typ: u8, flags: u8, stream_id: u32) {
562 out.push(((payload_len >> 16) & 0xff) as u8);
563 out.push(((payload_len >> 8) & 0xff) as u8);
564 out.push((payload_len & 0xff) as u8);
565 out.push(typ);
566 out.push(flags);
567 out.push(((stream_id >> 24) & 0x7f) as u8);
568 out.push(((stream_id >> 16) & 0xff) as u8);
569 out.push(((stream_id >> 8) & 0xff) as u8);
570 out.push((stream_id & 0xff) as u8);
571 }
572
573 #[inline]
574 pub fn write_data(&self, out: &mut Vec<u8>, stream_id: u32, end_stream: bool, data: &[u8]) {
575 let flags = if end_stream { FLAG_END_STREAM } else { 0 };
576 FrameWriter::header(out, data.len(), DATA_TYPE, flags, stream_id);
577 out.extend_from_slice(data);
578 }
579
580 #[inline]
582 pub fn write_headers(
583 &self,
584 out: &mut Vec<u8>,
585 stream_id: u32,
586 end_stream: bool,
587 end_headers: bool,
588 priority: Option<Priority>,
589 block: &[u8],
590 ) {
591 let mut flags = 0;
592 if end_stream {
593 flags |= FLAG_END_STREAM;
594 }
595 if end_headers {
596 flags |= FLAG_END_HEADERS;
597 }
598 let extra = if priority.is_some() { 5 } else { 0 };
599 if extra != 0 {
600 flags |= FLAG_PRIORITY;
601 }
602 FrameWriter::header(out, block.len() + extra, HEADERS_TYPE, flags, stream_id);
603 if let Some(p) = priority {
604 let dep = p.dependency & 0x7fff_ffff | if p.exclusive { 0x8000_0000 } else { 0 };
605 out.extend_from_slice(&dep.to_be_bytes());
606 out.push(p.weight);
607 }
608 out.extend_from_slice(block);
609 }
610
611 #[inline]
615 pub fn write_field_block(
616 &self,
617 out: &mut Vec<u8>,
618 stream_id: u32,
619 end_stream: bool,
620 block: &[u8],
621 ) {
622 let capacity = self.max_frame_size;
623 if block.len() <= capacity {
624 self.write_headers(out, stream_id, end_stream, true, None, block);
625 return;
626 }
627 let first = &block[..capacity];
628 self.write_headers(out, stream_id, end_stream, false, None, first);
629 let mut rest = &block[capacity..];
630 while rest.len() > capacity {
631 self.write_continuation(out, stream_id, false, &rest[..capacity]);
632 rest = &rest[capacity..];
633 }
634 self.write_continuation(out, stream_id, true, rest);
635 }
636
637 #[inline]
638 pub fn write_priority(&self, out: &mut Vec<u8>, stream_id: u32, priority: Priority) {
639 let dep =
640 priority.dependency & 0x7fff_ffff | if priority.exclusive { 0x8000_0000 } else { 0 };
641 FrameWriter::header(out, 5, PRIORITY_TYPE, 0, stream_id);
642 out.extend_from_slice(&dep.to_be_bytes());
643 out.push(priority.weight);
644 }
645
646 #[inline]
647 pub fn write_reset(&self, out: &mut Vec<u8>, stream_id: u32, error_code: u32) {
648 FrameWriter::header(out, 4, RST_STREAM_TYPE, 0, stream_id);
649 out.extend_from_slice(&error_code.to_be_bytes());
650 }
651
652 #[inline]
653 pub fn write_settings(&self, out: &mut Vec<u8>, settings: &[Setting]) {
654 FrameWriter::header(out, settings.len() * 6, SETTINGS_TYPE, 0, 0);
655 for setting in settings {
656 out.extend_from_slice(&setting.id.to_be_bytes());
657 out.extend_from_slice(&setting.value.to_be_bytes());
658 }
659 }
660
661 #[inline]
662 pub fn write_settings_ack(&self, out: &mut Vec<u8>) {
663 FrameWriter::header(out, 0, SETTINGS_TYPE, FLAG_ACK, 0);
664 }
665
666 #[inline]
667 pub fn write_push_promise(
668 &self,
669 out: &mut Vec<u8>,
670 stream_id: u32,
671 promised_stream_id: u32,
672 block: &[u8],
673 ) {
674 FrameWriter::header(
675 out,
676 4 + block.len(),
677 PUSH_PROMISE_TYPE,
678 FLAG_END_HEADERS,
679 stream_id,
680 );
681 out.extend_from_slice(&(promised_stream_id & 0x7fff_ffff).to_be_bytes());
682 out.extend_from_slice(block);
683 }
684
685 #[inline]
686 pub fn write_ping(&self, out: &mut Vec<u8>, payload: &[u8; 8]) {
687 FrameWriter::header(out, 8, PING_TYPE, 0, 0);
688 out.extend_from_slice(payload);
689 }
690
691 #[inline]
692 pub fn write_ping_ack(&self, out: &mut Vec<u8>, payload: &[u8; 8]) {
693 FrameWriter::header(out, 8, PING_TYPE, FLAG_ACK, 0);
694 out.extend_from_slice(payload);
695 }
696
697 #[inline]
698 pub fn write_goaway(
699 &self,
700 out: &mut Vec<u8>,
701 last_stream_id: u32,
702 error_code: u32,
703 debug: &[u8],
704 ) {
705 FrameWriter::header(out, 8 + debug.len(), GOAWAY_TYPE, 0, 0);
706 out.extend_from_slice(&(last_stream_id & 0x7fff_ffff).to_be_bytes());
707 out.extend_from_slice(&error_code.to_be_bytes());
708 out.extend_from_slice(debug);
709 }
710
711 #[inline]
712 pub fn write_window_update(&self, out: &mut Vec<u8>, stream_id: u32, increment: u32) {
713 FrameWriter::header(out, 4, WINDOW_UPDATE_TYPE, 0, stream_id);
714 out.extend_from_slice(&(increment & 0x7fff_ffff).to_be_bytes());
715 }
716
717 #[inline]
718 pub fn write_continuation(
719 &self,
720 out: &mut Vec<u8>,
721 stream_id: u32,
722 end_headers: bool,
723 block: &[u8],
724 ) {
725 let flags = if end_headers { FLAG_END_HEADERS } else { 0 };
726 FrameWriter::header(out, block.len(), CONTINUATION_TYPE, flags, stream_id);
727 out.extend_from_slice(block);
728 }
729}
730
731#[cfg(test)]
732mod tests;