Skip to main content

h2ts_client/
frames.rs

1//! HTTP/2 frame model + codec (RFC 7540 §6) — the Rust port of `frames/types.ts`
2//! and `frames/codec.ts`. A [`Frame`] is a tagged enum; [`serialize_frame`] writes
3//! the 9-byte header + payload, and [`FrameDecoder`] streams complete frames out
4//! of arbitrary byte chunks, skipping unknown frame types (RFC 7540 §4.1).
5
6use crate::bytes::{ByteReader, ByteWriter};
7use crate::errors::{ErrorCode, H2Error};
8
9pub const FRAME_HEADER_SIZE: usize = 9;
10pub const DEFAULT_MAX_FRAME_SIZE: usize = 16384;
11
12/// Frame type identifiers (RFC 7540 §6).
13pub mod frame_type {
14    pub const DATA: u8 = 0x0;
15    pub const HEADERS: u8 = 0x1;
16    pub const PRIORITY: u8 = 0x2;
17    pub const RST_STREAM: u8 = 0x3;
18    pub const SETTINGS: u8 = 0x4;
19    pub const PUSH_PROMISE: u8 = 0x5;
20    pub const PING: u8 = 0x6;
21    pub const GOAWAY: u8 = 0x7;
22    pub const WINDOW_UPDATE: u8 = 0x8;
23    pub const CONTINUATION: u8 = 0x9;
24}
25
26/// Flag bits per frame type.
27mod flags {
28    pub const DATA_END_STREAM: u8 = 0x1;
29    pub const DATA_PADDED: u8 = 0x8;
30    pub const HEADERS_END_STREAM: u8 = 0x1;
31    pub const HEADERS_END_HEADERS: u8 = 0x4;
32    pub const HEADERS_PADDED: u8 = 0x8;
33    pub const HEADERS_PRIORITY: u8 = 0x20;
34    pub const SETTINGS_ACK: u8 = 0x1;
35    pub const PING_ACK: u8 = 0x1;
36    pub const PUSH_PROMISE_END_HEADERS: u8 = 0x4;
37    pub const PUSH_PROMISE_PADDED: u8 = 0x8;
38    pub const CONTINUATION_END_HEADERS: u8 = 0x4;
39}
40
41/// SETTINGS parameters (RFC 7540 §6.5.2). Absent fields are `None`.
42#[derive(Debug, Clone, Default, PartialEq, Eq)]
43pub struct Settings {
44    pub header_table_size: Option<u32>,      // 0x1
45    pub enable_push: Option<bool>,           // 0x2
46    pub max_concurrent_streams: Option<u32>, // 0x3
47    pub initial_window_size: Option<u32>,    // 0x4
48    pub max_frame_size: Option<u32>,         // 0x5
49    pub max_header_list_size: Option<u32>,   // 0x6
50}
51
52/// Stream priority (RFC 7540 §6.3). `weight` is the human value `1..=256`.
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub struct Priority {
55    pub stream_dependency: u32,
56    pub weight: u16,
57    pub exclusive: bool,
58}
59
60/// An HTTP/2 frame. SETTINGS/PING/GOAWAY always carry stream id 0, so it is implicit.
61#[derive(Debug, Clone, PartialEq, Eq)]
62pub enum Frame {
63    Data {
64        stream_id: u32,
65        data: Vec<u8>,
66        end_stream: bool,
67    },
68    Headers {
69        stream_id: u32,
70        header_block_fragment: Vec<u8>,
71        end_stream: bool,
72        end_headers: bool,
73        priority: Option<Priority>,
74    },
75    Priority {
76        stream_id: u32,
77        priority: Priority,
78    },
79    RstStream {
80        stream_id: u32,
81        error_code: u32,
82    },
83    Settings {
84        ack: bool,
85        settings: Settings,
86    },
87    PushPromise {
88        stream_id: u32,
89        promised_stream_id: u32,
90        header_block_fragment: Vec<u8>,
91        end_headers: bool,
92    },
93    Ping {
94        ack: bool,
95        opaque_data: [u8; 8],
96    },
97    Goaway {
98        last_stream_id: u32,
99        error_code: u32,
100        debug_data: Vec<u8>,
101    },
102    WindowUpdate {
103        stream_id: u32,
104        window_size_increment: u32,
105    },
106    Continuation {
107        stream_id: u32,
108        header_block_fragment: Vec<u8>,
109        end_headers: bool,
110    },
111}
112
113struct Encoded {
114    type_id: u8,
115    flags: u8,
116    stream_id: u32,
117    payload: Vec<u8>,
118}
119
120fn encode_body(frame: &Frame) -> Encoded {
121    match frame {
122        Frame::Data {
123            stream_id,
124            data,
125            end_stream,
126        } => Encoded {
127            type_id: frame_type::DATA,
128            flags: if *end_stream {
129                flags::DATA_END_STREAM
130            } else {
131                0
132            },
133            stream_id: *stream_id,
134            payload: data.clone(),
135        },
136        Frame::Headers {
137            stream_id,
138            header_block_fragment,
139            end_stream,
140            end_headers,
141            priority,
142        } => {
143            let mut f = if *end_headers {
144                flags::HEADERS_END_HEADERS
145            } else {
146                0
147            };
148            if *end_stream {
149                f |= flags::HEADERS_END_STREAM;
150            }
151            let payload = if let Some(p) = priority {
152                f |= flags::HEADERS_PRIORITY;
153                let mut w = ByteWriter::with_capacity(5 + header_block_fragment.len());
154                let dep = p.stream_dependency & 0x7fff_ffff;
155                w.u32(if p.exclusive { dep | 0x8000_0000 } else { dep });
156                w.u8((p.weight - 1) as u8);
157                w.bytes(header_block_fragment);
158                w.into_vec()
159            } else {
160                header_block_fragment.clone()
161            };
162            Encoded {
163                type_id: frame_type::HEADERS,
164                flags: f,
165                stream_id: *stream_id,
166                payload,
167            }
168        }
169        Frame::Priority {
170            stream_id,
171            priority,
172        } => {
173            let mut w = ByteWriter::with_capacity(5);
174            let dep = priority.stream_dependency & 0x7fff_ffff;
175            w.u32(if priority.exclusive {
176                dep | 0x8000_0000
177            } else {
178                dep
179            });
180            w.u8((priority.weight - 1) as u8);
181            Encoded {
182                type_id: frame_type::PRIORITY,
183                flags: 0,
184                stream_id: *stream_id,
185                payload: w.into_vec(),
186            }
187        }
188        Frame::RstStream {
189            stream_id,
190            error_code,
191        } => {
192            let mut w = ByteWriter::with_capacity(4);
193            w.u32(*error_code);
194            Encoded {
195                type_id: frame_type::RST_STREAM,
196                flags: 0,
197                stream_id: *stream_id,
198                payload: w.into_vec(),
199            }
200        }
201        Frame::Settings { ack, settings } => {
202            let mut w = ByteWriter::with_capacity(48);
203            if let Some(v) = settings.header_table_size {
204                w.u16(0x1);
205                w.u32(v);
206            }
207            if let Some(v) = settings.enable_push {
208                w.u16(0x2);
209                w.u32(u32::from(v));
210            }
211            if let Some(v) = settings.max_concurrent_streams {
212                w.u16(0x3);
213                w.u32(v);
214            }
215            if let Some(v) = settings.initial_window_size {
216                w.u16(0x4);
217                w.u32(v);
218            }
219            if let Some(v) = settings.max_frame_size {
220                w.u16(0x5);
221                w.u32(v);
222            }
223            if let Some(v) = settings.max_header_list_size {
224                w.u16(0x6);
225                w.u32(v);
226            }
227            Encoded {
228                type_id: frame_type::SETTINGS,
229                flags: if *ack { flags::SETTINGS_ACK } else { 0 },
230                stream_id: 0,
231                payload: w.into_vec(),
232            }
233        }
234        Frame::PushPromise {
235            stream_id,
236            promised_stream_id,
237            header_block_fragment,
238            end_headers,
239        } => {
240            let mut w = ByteWriter::with_capacity(4 + header_block_fragment.len());
241            w.u32(promised_stream_id & 0x7fff_ffff);
242            w.bytes(header_block_fragment);
243            Encoded {
244                type_id: frame_type::PUSH_PROMISE,
245                flags: if *end_headers {
246                    flags::PUSH_PROMISE_END_HEADERS
247                } else {
248                    0
249                },
250                stream_id: *stream_id,
251                payload: w.into_vec(),
252            }
253        }
254        Frame::Ping { ack, opaque_data } => Encoded {
255            type_id: frame_type::PING,
256            flags: if *ack { flags::PING_ACK } else { 0 },
257            stream_id: 0,
258            payload: opaque_data.to_vec(),
259        },
260        Frame::Goaway {
261            last_stream_id,
262            error_code,
263            debug_data,
264        } => {
265            let mut w = ByteWriter::with_capacity(8 + debug_data.len());
266            w.u32(last_stream_id & 0x7fff_ffff);
267            w.u32(*error_code);
268            w.bytes(debug_data);
269            Encoded {
270                type_id: frame_type::GOAWAY,
271                flags: 0,
272                stream_id: 0,
273                payload: w.into_vec(),
274            }
275        }
276        Frame::WindowUpdate {
277            stream_id,
278            window_size_increment,
279        } => {
280            let mut w = ByteWriter::with_capacity(4);
281            w.u32(window_size_increment & 0x7fff_ffff);
282            Encoded {
283                type_id: frame_type::WINDOW_UPDATE,
284                flags: 0,
285                stream_id: *stream_id,
286                payload: w.into_vec(),
287            }
288        }
289        Frame::Continuation {
290            stream_id,
291            header_block_fragment,
292            end_headers,
293        } => Encoded {
294            type_id: frame_type::CONTINUATION,
295            flags: if *end_headers {
296                flags::CONTINUATION_END_HEADERS
297            } else {
298                0
299            },
300            stream_id: *stream_id,
301            payload: header_block_fragment.clone(),
302        },
303    }
304}
305
306/// Serialize a frame to bytes (9-byte header + payload).
307pub fn serialize_frame(frame: &Frame) -> Vec<u8> {
308    let e = encode_body(frame);
309    let mut w = ByteWriter::with_capacity(FRAME_HEADER_SIZE + e.payload.len());
310    w.u24(e.payload.len() as u32);
311    w.u8(e.type_id);
312    w.u8(e.flags);
313    w.u32(e.stream_id & 0x7fff_ffff);
314    w.bytes(&e.payload);
315    w.into_vec()
316}
317
318fn truncated() -> H2Error {
319    H2Error::new(ErrorCode::ProtocolError, "truncated frame payload")
320}
321
322fn read_padded<'a>(r: &mut ByteReader<'a>, len: usize, padded: bool) -> Result<&'a [u8], H2Error> {
323    if !padded {
324        return r.bytes(len).ok_or_else(truncated);
325    }
326    if len < 1 {
327        return Err(H2Error::new(
328            ErrorCode::ProtocolError,
329            "padded frame missing pad length",
330        ));
331    }
332    let pad_length = r.u8() as usize;
333    let data_length = (len - 1).checked_sub(pad_length).ok_or_else(|| {
334        H2Error::new(ErrorCode::ProtocolError, "pad length exceeds frame payload")
335    })?;
336    let data = r.bytes(data_length).ok_or_else(truncated)?;
337    r.bytes(pad_length); // discard padding
338    Ok(data)
339}
340
341/// Parse one frame payload. `Ok(None)` means an unknown frame type to skip.
342fn parse_payload(
343    type_id: u8,
344    flag_bits: u8,
345    stream_id: u32,
346    payload: &[u8],
347) -> Result<Option<Frame>, H2Error> {
348    let mut r = ByteReader::new(payload);
349    let len = payload.len();
350
351    let frame = match type_id {
352        frame_type::DATA => {
353            let data = read_padded(&mut r, len, flag_bits & flags::DATA_PADDED != 0)?;
354            Frame::Data {
355                stream_id,
356                data: data.to_vec(),
357                end_stream: flag_bits & flags::DATA_END_STREAM != 0,
358            }
359        }
360        frame_type::HEADERS => {
361            let padded = flag_bits & flags::HEADERS_PADDED != 0;
362            let has_priority = flag_bits & flags::HEADERS_PRIORITY != 0;
363            let pad_length = if padded {
364                if len < 1 {
365                    return Err(H2Error::new(
366                        ErrorCode::ProtocolError,
367                        "invalid HEADERS padding",
368                    ));
369                }
370                r.u8() as usize
371            } else {
372                0
373            };
374            let priority = if has_priority {
375                if r.remaining() < 5 {
376                    return Err(H2Error::new(
377                        ErrorCode::ProtocolError,
378                        "HEADERS priority truncated",
379                    ));
380                }
381                let dep = r.u32();
382                let weight = r.u8() as u16 + 1;
383                Some(Priority {
384                    stream_dependency: dep & 0x7fff_ffff,
385                    exclusive: dep & 0x8000_0000 != 0,
386                    weight,
387                })
388            } else {
389                None
390            };
391            let overhead = usize::from(padded) + if has_priority { 5 } else { 0 } + pad_length;
392            let frag_length = len
393                .checked_sub(overhead)
394                .ok_or_else(|| H2Error::new(ErrorCode::ProtocolError, "invalid HEADERS padding"))?;
395            let fragment = r.bytes(frag_length).ok_or_else(truncated)?;
396            Frame::Headers {
397                stream_id,
398                header_block_fragment: fragment.to_vec(),
399                end_stream: flag_bits & flags::HEADERS_END_STREAM != 0,
400                end_headers: flag_bits & flags::HEADERS_END_HEADERS != 0,
401                priority,
402            }
403        }
404        frame_type::PRIORITY => {
405            if len != 5 {
406                return Err(H2Error::stream(
407                    ErrorCode::FrameSizeError,
408                    "PRIORITY must be 5 bytes",
409                    stream_id,
410                ));
411            }
412            let dep = r.u32();
413            let weight = r.u8() as u16 + 1;
414            Frame::Priority {
415                stream_id,
416                priority: Priority {
417                    stream_dependency: dep & 0x7fff_ffff,
418                    exclusive: dep & 0x8000_0000 != 0,
419                    weight,
420                },
421            }
422        }
423        frame_type::RST_STREAM => {
424            if len != 4 {
425                return Err(H2Error::new(
426                    ErrorCode::FrameSizeError,
427                    "RST_STREAM must be 4 bytes",
428                ));
429            }
430            Frame::RstStream {
431                stream_id,
432                error_code: r.u32(),
433            }
434        }
435        frame_type::SETTINGS => {
436            let ack = flag_bits & flags::SETTINGS_ACK != 0;
437            if ack && len != 0 {
438                return Err(H2Error::new(
439                    ErrorCode::FrameSizeError,
440                    "SETTINGS ACK must be empty",
441                ));
442            }
443            if !len.is_multiple_of(6) {
444                return Err(H2Error::new(
445                    ErrorCode::FrameSizeError,
446                    "SETTINGS length not multiple of 6",
447                ));
448            }
449            let mut settings = Settings::default();
450            for _ in 0..len / 6 {
451                let id = r.u16();
452                let value = r.u32();
453                match id {
454                    0x1 => settings.header_table_size = Some(value),
455                    0x2 => settings.enable_push = Some(value != 0),
456                    0x3 => settings.max_concurrent_streams = Some(value),
457                    0x4 => settings.initial_window_size = Some(value),
458                    0x5 => settings.max_frame_size = Some(value),
459                    0x6 => settings.max_header_list_size = Some(value),
460                    _ => {} // unknown settings ignored (RFC 7540 §6.5.2)
461                }
462            }
463            Frame::Settings { ack, settings }
464        }
465        frame_type::PUSH_PROMISE => {
466            let padded = flag_bits & flags::PUSH_PROMISE_PADDED != 0;
467            let pad_length = if padded {
468                if len < 1 {
469                    return Err(H2Error::new(
470                        ErrorCode::ProtocolError,
471                        "invalid PUSH_PROMISE padding",
472                    ));
473                }
474                r.u8() as usize
475            } else {
476                0
477            };
478            if r.remaining() < 4 {
479                return Err(H2Error::new(
480                    ErrorCode::ProtocolError,
481                    "PUSH_PROMISE truncated",
482                ));
483            }
484            let promised = r.u32() & 0x7fff_ffff;
485            let overhead = usize::from(padded) + 4 + pad_length;
486            let frag_length = len.checked_sub(overhead).ok_or_else(|| {
487                H2Error::new(ErrorCode::ProtocolError, "invalid PUSH_PROMISE padding")
488            })?;
489            let fragment = r.bytes(frag_length).ok_or_else(truncated)?;
490            Frame::PushPromise {
491                stream_id,
492                promised_stream_id: promised,
493                header_block_fragment: fragment.to_vec(),
494                end_headers: flag_bits & flags::PUSH_PROMISE_END_HEADERS != 0,
495            }
496        }
497        frame_type::PING => {
498            if len != 8 {
499                return Err(H2Error::new(
500                    ErrorCode::FrameSizeError,
501                    "PING must be 8 bytes",
502                ));
503            }
504            let mut opaque = [0u8; 8];
505            opaque.copy_from_slice(r.bytes(8).ok_or_else(truncated)?);
506            Frame::Ping {
507                ack: flag_bits & flags::PING_ACK != 0,
508                opaque_data: opaque,
509            }
510        }
511        frame_type::GOAWAY => {
512            if len < 8 {
513                return Err(H2Error::new(ErrorCode::FrameSizeError, "GOAWAY too short"));
514            }
515            let last_stream_id = r.u32() & 0x7fff_ffff;
516            let error_code = r.u32();
517            let debug_data = r.bytes(len - 8).ok_or_else(truncated)?.to_vec();
518            Frame::Goaway {
519                last_stream_id,
520                error_code,
521                debug_data,
522            }
523        }
524        frame_type::WINDOW_UPDATE => {
525            if len != 4 {
526                return Err(H2Error::new(
527                    ErrorCode::FrameSizeError,
528                    "WINDOW_UPDATE must be 4 bytes",
529                ));
530            }
531            Frame::WindowUpdate {
532                stream_id,
533                window_size_increment: r.u32() & 0x7fff_ffff,
534            }
535        }
536        frame_type::CONTINUATION => Frame::Continuation {
537            stream_id,
538            header_block_fragment: r.bytes(len).ok_or_else(truncated)?.to_vec(),
539            end_headers: flag_bits & flags::CONTINUATION_END_HEADERS != 0,
540        },
541        _ => return Ok(None), // unknown frame type: skip (RFC 7540 §4.1)
542    };
543    Ok(Some(frame))
544}
545
546/// Streaming frame decoder. Feed it arbitrary byte chunks; it returns whatever
547/// complete frames are now available, buffering any partial frame internally.
548pub struct FrameDecoder {
549    pending: Vec<u8>,
550    max_frame_size: usize,
551}
552
553impl Default for FrameDecoder {
554    fn default() -> Self {
555        Self::new(DEFAULT_MAX_FRAME_SIZE)
556    }
557}
558
559impl FrameDecoder {
560    pub fn new(max_frame_size: usize) -> Self {
561        Self {
562            pending: Vec::new(),
563            max_frame_size,
564        }
565    }
566
567    pub fn push(&mut self, chunk: &[u8]) -> Result<Vec<Frame>, H2Error> {
568        let buf: Vec<u8> = if self.pending.is_empty() {
569            chunk.to_vec()
570        } else {
571            let mut v = core::mem::take(&mut self.pending);
572            v.extend_from_slice(chunk);
573            v
574        };
575
576        let mut frames = Vec::new();
577        let mut offset = 0;
578
579        while buf.len() - offset >= FRAME_HEADER_SIZE {
580            let length = ((buf[offset] as usize) << 16)
581                | ((buf[offset + 1] as usize) << 8)
582                | (buf[offset + 2] as usize);
583            if length > self.max_frame_size {
584                return Err(H2Error::new(
585                    ErrorCode::FrameSizeError,
586                    format!("frame length {length} exceeds max {}", self.max_frame_size),
587                ));
588            }
589            let total = FRAME_HEADER_SIZE + length;
590            if buf.len() - offset < total {
591                break; // wait for more bytes
592            }
593
594            let type_id = buf[offset + 3];
595            let flag_bits = buf[offset + 4];
596            let stream_id = u32::from_be_bytes([
597                buf[offset + 5],
598                buf[offset + 6],
599                buf[offset + 7],
600                buf[offset + 8],
601            ]) & 0x7fff_ffff;
602            let payload = &buf[offset + FRAME_HEADER_SIZE..offset + total];
603
604            if let Some(frame) = parse_payload(type_id, flag_bits, stream_id, payload)? {
605                frames.push(frame);
606            }
607            offset += total;
608        }
609
610        self.pending = if offset == buf.len() {
611            Vec::new()
612        } else {
613            buf[offset..].to_vec()
614        };
615        Ok(frames)
616    }
617}