#![expect(dead_code)]
use std::collections::VecDeque;
use bytes::Bytes;
use crate::h3::qpack::error::QpackError;
use crate::h3::qpack::static_table;
use crate::h3::qpack::table::DynamicTable;
use crate::hpack::{huffman, integer, HpackError};
const SECTION_ACK: u8 = 0x80;
const STREAM_CANCELLATION: u8 = 0x40;
const INSERT_COUNT_INCREMENT: u8 = 0x00;
const SET_CAPACITY: u8 = 0b0010_0000;
const INSERT_WITH_NAME_REF: u8 = 0b1000_0000;
const INSERT_WITH_LITERAL_NAME: u8 = 0b0100_0000;
const DUPLICATE: u8 = 0b0000_0000;
const INDEXED: u8 = 0b1000_0000;
const INDEXED_POST_BASE: u8 = 0b0001_0000;
const LITERAL_NAME_REF: u8 = 0b0100_0000;
const LITERAL_POST_BASE_NAME_REF: u8 = 0b0000_0000;
const LITERAL_LITERAL_NAME: u8 = 0b0010_0000;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnblockedSection {
pub stream_id: u64,
pub headers: Vec<(Bytes, Bytes)>,
}
#[derive(Debug)]
struct BlockedSection {
stream_id: u64,
buf: Bytes,
ric: u64,
since: u64,
}
#[derive(Debug)]
pub struct Decoder {
dynamic: DynamicTable,
max_capacity: u64,
known_received: u64,
blocked: VecDeque<BlockedSection>,
blocked_by_stream: Vec<(u64, usize)>,
section_size_by_stream: Vec<(u64, usize)>,
max_blocked_streams: usize,
max_field_section_size: usize,
decoder_stream: Vec<u8>,
encoder_stream_pending: Vec<u8>,
decoder_stream_pending: Vec<u8>,
}
const MAX_DECODER_STREAM_PENDING: usize = 64;
fn string_len(buf: &[u8], prefix_bits: u8) -> Option<usize> {
let int_prefix = prefix_bits - 1;
let int_len = integer::encoded_len(buf, int_prefix)?;
let mut off = 1;
let len = integer::decode(buf, &mut off, int_prefix, buf[0]).ok()?;
let content = usize::try_from(len).ok()?;
let total = int_len.checked_add(content)?;
if buf.len() < total {
return None;
}
Some(total)
}
fn encoder_instruction_len(buf: &[u8]) -> Option<usize> {
let header = *buf.first()?;
match header & 0xC0 {
0x80 | 0xC0 => {
let index_len = integer::encoded_len(buf, 6)?;
let value_len = string_len(buf.get(index_len..)?, 8)?;
Some(index_len + value_len)
}
0x40 => {
let name_len = string_len(buf, 6)?;
let value_len = string_len(buf.get(name_len..)?, 8)?;
Some(name_len + value_len)
}
_ => integer::encoded_len(buf, 5),
}
}
fn decoder_instruction_len(buf: &[u8]) -> Option<usize> {
let header = *buf.first()?;
let prefix_bits: u8 = if header & 0x80 != 0 { 7 } else { 6 };
integer::encoded_len(buf, prefix_bits)
}
impl Decoder {
#[inline]
pub fn new(max_capacity: u64, max_blocked_streams: usize) -> Self {
Self {
dynamic: DynamicTable::new(0),
max_capacity,
known_received: 0,
blocked: VecDeque::new(),
blocked_by_stream: Vec::new(),
section_size_by_stream: Vec::new(),
max_blocked_streams,
max_field_section_size: usize::MAX,
decoder_stream: Vec::new(),
encoder_stream_pending: Vec::new(),
decoder_stream_pending: Vec::new(),
}
}
#[inline]
pub fn set_max_field_section_size(&mut self, size: usize) {
self.max_field_section_size = size;
}
#[inline]
pub fn inserted(&self) -> u64 {
self.dynamic.inserted()
}
#[inline]
pub fn known_received(&self) -> u64 {
self.known_received
}
#[inline]
pub fn pending_blocked(&self) -> usize {
self.blocked.len()
}
#[inline]
pub fn take_decoder_stream(&mut self) -> Bytes {
Bytes::from(std::mem::take(&mut self.decoder_stream))
}
#[inline]
pub fn feed_decoder_stream(&mut self, buf: &[u8]) -> Result<(), QpackError> {
self.decoder_stream_pending.extend_from_slice(buf);
if self.decoder_stream_pending.len() > MAX_DECODER_STREAM_PENDING {
return Err(QpackError::DecoderStream);
}
let mut consumed = 0;
while consumed < self.decoder_stream_pending.len() {
let Some(len) = decoder_instruction_len(&self.decoder_stream_pending[consumed..])
else {
break;
};
if consumed + len > self.decoder_stream_pending.len() {
break;
}
let instr = &self.decoder_stream_pending[consumed..consumed + len];
let mut off = 0;
let header = instr[0];
if header & 0x80 != 0 {
integer::decode(instr, &mut off, 7, header).map_err(dec_stream_err)?;
} else if header & 0x40 != 0 {
integer::decode(instr, &mut off, 6, header).map_err(dec_stream_err)?;
} else {
let increment =
integer::decode(instr, &mut off, 6, header).map_err(dec_stream_err)?;
if increment == 0 {
return Err(QpackError::DecoderStream);
}
}
consumed += len;
}
self.decoder_stream_pending.drain(..consumed);
Ok(())
}
#[inline]
pub fn feed_encoder_stream(&mut self, buf: &[u8]) -> Result<Vec<UnblockedSection>, QpackError> {
self.encoder_stream_pending.extend_from_slice(buf);
let cap = (self.max_capacity as usize).saturating_add(1024);
if self.encoder_stream_pending.len() > cap {
return Err(QpackError::EncoderStream);
}
let mut consumed = 0;
while consumed < self.encoder_stream_pending.len() {
let Some(len) = encoder_instruction_len(&self.encoder_stream_pending[consumed..])
else {
break;
};
if consumed + len > self.encoder_stream_pending.len() {
break;
}
let instr = self.encoder_stream_pending[consumed..consumed + len].to_vec();
self.parse_encoder_instruction(&instr)?;
consumed += len;
}
self.encoder_stream_pending.drain(..consumed);
let mut sections = Vec::new();
while let Some(front) = self.blocked.front() {
if front.ric > self.dynamic.inserted() {
break;
}
let front = self.blocked.pop_front().expect("front just inspected");
self.remove_blocked_section(front.stream_id);
let headers = self.decode_ready(&front.buf)?;
let size: usize = headers.iter().map(|(n, v)| n.len() + v.len()).sum();
if self.account_section(front.stream_id, size) > self.max_field_section_size {
return Err(QpackError::DecompressionFailed);
}
if front.ric > 0 {
self.acknowledge(front.ric);
self.emit_section_ack(front.stream_id);
}
sections.push(UnblockedSection {
stream_id: front.stream_id,
headers,
});
}
if self.dynamic.inserted() > self.known_received {
integer::encode(
&mut self.decoder_stream,
self.dynamic.inserted() - self.known_received,
6,
INSERT_COUNT_INCREMENT,
);
self.known_received = self.dynamic.inserted();
}
Ok(sections)
}
#[inline]
fn parse_encoder_instruction(&mut self, instr: &[u8]) -> Result<(), QpackError> {
let mut off = 0;
let header = instr[0];
off += 1;
match header & 0xC0 {
0x80 | 0xC0 => {
let index = integer::decode(instr, &mut off, 6, header).map_err(enc_stream_err)?;
let value = self
.read_value_string(instr, &mut off)
.map_err(enc_stream_err)?;
let name = if header & 0x40 != 0 {
let idx = usize::try_from(index).map_err(|_| QpackError::EncoderStream)?;
let (name, _) = static_table::get(idx).ok_or(QpackError::EncoderStream)?;
Bytes::from_static(name)
} else {
let (name, _) = self
.dynamic
.get_relative_bytes(index)
.ok_or(QpackError::EncoderStream)?;
name
};
self.insert_entry(name, value)?;
}
0x40 => {
let name = self
.read_string(instr, &mut off, 6, header)
.map_err(enc_stream_err)?;
let value = self
.read_value_string(instr, &mut off)
.map_err(enc_stream_err)?;
self.insert_entry(name, value)?;
}
_ => {
if header & 0x20 != 0 {
let capacity =
integer::decode(instr, &mut off, 5, header).map_err(enc_stream_err)?;
if capacity > self.max_capacity {
return Err(QpackError::EncoderStream);
}
let evicted = self.dynamic.evict_for_capacity(capacity);
if evicted > self.known_received {
return Err(QpackError::EncoderStream);
}
self.dynamic.set_capacity(capacity);
} else {
let index =
integer::decode(instr, &mut off, 5, header).map_err(enc_stream_err)?;
let (name, value) = self
.dynamic
.get_relative_bytes(index)
.ok_or(QpackError::EncoderStream)?;
self.insert_entry(name, value)?;
}
}
}
Ok(())
}
#[inline]
pub fn decode_block(
&mut self,
buf: &[u8],
stream_id: u64,
now: u64,
) -> Result<Option<Vec<(Bytes, Bytes)>>, QpackError> {
let (ric, _, _) = self.read_prefix(buf)?;
let stream_blocked = self.stream_has_blocked(stream_id);
if ric > self.dynamic.inserted() || stream_blocked {
if !stream_blocked && self.blocked_by_stream.len() >= self.max_blocked_streams {
return Err(QpackError::DecompressionFailed);
}
self.blocked.push_back(BlockedSection {
stream_id,
buf: Bytes::copy_from_slice(buf),
ric,
since: now,
});
self.add_blocked_section(stream_id);
return Ok(None);
}
let headers = self.decode_ready(buf)?;
let size: usize = headers.iter().map(|(n, v)| n.len() + v.len()).sum();
if self.account_section(stream_id, size) > self.max_field_section_size {
return Err(QpackError::DecompressionFailed);
}
if ric > 0 {
self.acknowledge(ric);
self.emit_section_ack(stream_id);
}
Ok(Some(headers))
}
#[inline]
pub fn stream_finished(&mut self, stream_id: u64) {
self.section_size_by_stream
.retain(|(id, _)| *id != stream_id);
}
#[inline]
pub fn stream_cancelled(&mut self, stream_id: u64) -> Bytes {
self.blocked.retain(|b| b.stream_id != stream_id);
self.blocked_by_stream.retain(|(id, _)| *id != stream_id);
self.section_size_by_stream
.retain(|(id, _)| *id != stream_id);
let mut out = Vec::new();
integer::encode(&mut out, stream_id, 6, STREAM_CANCELLATION);
self.decoder_stream.extend_from_slice(&out);
Bytes::from(out)
}
#[inline]
pub fn expire_blocked(&mut self, now: u64, max_age: u64) -> Bytes {
let mut out = Vec::new();
let mut cancelled = Vec::new();
for blocked in &self.blocked {
if now.saturating_sub(blocked.since) > max_age
&& !cancelled.contains(&blocked.stream_id)
{
cancelled.push(blocked.stream_id);
integer::encode(&mut out, blocked.stream_id, 6, STREAM_CANCELLATION);
}
}
if !cancelled.is_empty() {
self.blocked
.retain(|blocked| !cancelled.contains(&blocked.stream_id));
self.blocked_by_stream
.retain(|(stream_id, _)| !cancelled.contains(stream_id));
self.section_size_by_stream
.retain(|(stream_id, _)| !cancelled.contains(stream_id));
}
self.decoder_stream.extend_from_slice(&out);
Bytes::from(out)
}
#[inline]
fn decode_ready(&self, buf: &[u8]) -> Result<Vec<(Bytes, Bytes)>, QpackError> {
let (ric, base, mut off) = self.read_prefix(buf)?;
let mut headers = Vec::new();
let mut needed = 0u64;
while off < buf.len() {
let header = buf[off];
off += 1;
if header & 0x80 != 0 {
let index = integer::decode(buf, &mut off, 6, header).map_err(dec_failed)?;
if header & 0x40 != 0 {
let idx =
usize::try_from(index).map_err(|_| QpackError::DecompressionFailed)?;
let (name, value) =
static_table::get(idx).ok_or(QpackError::DecompressionFailed)?;
headers.push((Bytes::from_static(name), Bytes::from_static(value)));
} else {
let (name, value) = self
.dynamic
.get_base_relative_bytes(base, index)
.ok_or(QpackError::DecompressionFailed)?;
needed = needed.max(base - index);
headers.push((name, value));
}
} else if header & 0x40 != 0 {
let index = integer::decode(buf, &mut off, 4, header).map_err(dec_failed)?;
let value = self.read_value_string(buf, &mut off).map_err(dec_failed)?;
if header & 0x10 != 0 {
let idx =
usize::try_from(index).map_err(|_| QpackError::DecompressionFailed)?;
let (name, _) =
static_table::get(idx).ok_or(QpackError::DecompressionFailed)?;
headers.push((Bytes::from_static(name), value));
} else {
let (name, _) = self
.dynamic
.get_base_relative_bytes(base, index)
.ok_or(QpackError::DecompressionFailed)?;
needed = needed.max(base - index);
headers.push((name, value));
}
} else if header & 0x20 != 0 {
let name = self
.read_string(buf, &mut off, 4, header)
.map_err(dec_failed)?;
let value = self.read_value_string(buf, &mut off).map_err(dec_failed)?;
headers.push((name, value));
} else if header & 0x10 != 0 {
let index = integer::decode(buf, &mut off, 4, header).map_err(dec_failed)?;
let (name, value) = self
.dynamic
.get_post_base_bytes(base, index)
.ok_or(QpackError::DecompressionFailed)?;
needed = needed.max(base + index + 1);
headers.push((name, value));
} else {
let index = integer::decode(buf, &mut off, 3, header).map_err(dec_failed)?;
let (name, _) = self
.dynamic
.get_post_base_bytes(base, index)
.ok_or(QpackError::DecompressionFailed)?;
needed = needed.max(base + index + 1);
let value = self.read_value_string(buf, &mut off).map_err(dec_failed)?;
headers.push((name, value));
}
}
if ric != needed {
return Err(QpackError::DecompressionFailed);
}
Ok(headers)
}
#[inline]
fn read_prefix(&self, buf: &[u8]) -> Result<(u64, u64, usize), QpackError> {
let header = *buf.first().ok_or(QpackError::DecompressionFailed)?;
let mut off = 1;
let enc_ric = integer::decode(buf, &mut off, 8, header).map_err(dec_failed)?;
let max_entries = self.max_capacity / 32;
let ric = if enc_ric == 0 || max_entries == 0 {
0
} else {
let full_range = 2 * max_entries;
if enc_ric > full_range {
return Err(QpackError::DecompressionFailed);
}
let max_value = self.dynamic.inserted() + max_entries;
let max_wrapped = (max_value / full_range) * full_range;
let mut ric = max_wrapped + enc_ric - 1;
if ric > max_value {
if ric <= full_range {
return Err(QpackError::DecompressionFailed);
}
ric -= full_range;
}
if ric == 0 {
return Err(QpackError::DecompressionFailed);
}
ric
};
let header = *buf.get(off).ok_or(QpackError::DecompressionFailed)?;
off += 1;
let delta = integer::decode(buf, &mut off, 7, header).map_err(dec_failed)?;
let base = if header & 0x80 != 0 {
ric.checked_sub(delta + 1)
.ok_or(QpackError::DecompressionFailed)?
} else {
ric.checked_add(delta)
.ok_or(QpackError::DecompressionFailed)?
};
Ok((ric, base, off))
}
#[inline]
fn read_string(
&self,
buf: &[u8],
off: &mut usize,
prefix_bits: u8,
header: u8,
) -> Result<Bytes, HpackError> {
let huffman = header & (1 << (prefix_bits - 1)) != 0;
let len = integer::decode(buf, off, prefix_bits - 1, header)?;
let len = usize::try_from(len).map_err(|_| HpackError::InvalidString)?;
let end = (*off).checked_add(len).ok_or(HpackError::InvalidString)?;
let src = buf.get(*off..end).ok_or(HpackError::InvalidString)?;
*off = end;
if huffman {
let mut dst = Vec::with_capacity(len);
huffman::decode(src, &mut dst)?;
Ok(Bytes::from(dst))
} else {
Ok(Bytes::copy_from_slice(src))
}
}
#[inline]
fn read_value_string(&self, buf: &[u8], off: &mut usize) -> Result<Bytes, HpackError> {
let header = *buf.get(*off).ok_or(HpackError::InvalidString)?;
*off += 1;
self.read_string(buf, off, 8, header)
}
#[inline]
fn insert_entry(&mut self, name: Bytes, value: Bytes) -> Result<(), QpackError> {
let size = DynamicTable::entry_size(&name, &value);
let evicted = self.dynamic.would_evict(size);
if evicted > self.known_received {
return Err(QpackError::EncoderStream);
}
self.dynamic
.insert(name, value)
.map_err(|_| QpackError::EncoderStream)
}
#[inline]
fn acknowledge(&mut self, ric: u64) {
self.known_received = self.known_received.max(ric);
}
#[inline]
fn emit_section_ack(&mut self, stream_id: u64) {
integer::encode(&mut self.decoder_stream, stream_id, 7, SECTION_ACK);
}
#[inline]
fn stream_has_blocked(&self, stream_id: u64) -> bool {
self.blocked_by_stream
.iter()
.any(|(id, _)| *id == stream_id)
}
#[inline]
fn add_blocked_section(&mut self, stream_id: u64) {
if let Some((_, count)) = self
.blocked_by_stream
.iter_mut()
.find(|(id, _)| *id == stream_id)
{
*count += 1;
} else {
self.blocked_by_stream.push((stream_id, 1));
}
}
#[inline]
fn account_section(&mut self, stream_id: u64, size: usize) -> usize {
if let Some((_, total)) = self
.section_size_by_stream
.iter_mut()
.find(|(id, _)| *id == stream_id)
{
*total = total.saturating_add(size);
*total
} else {
self.section_size_by_stream.push((stream_id, size));
size
}
}
#[inline]
fn remove_blocked_section(&mut self, stream_id: u64) {
let Some(index) = self
.blocked_by_stream
.iter()
.position(|(id, _)| *id == stream_id)
else {
debug_assert!(false, "blocked section missing stream accounting");
return;
};
if self.blocked_by_stream[index].1 == 1 {
self.blocked_by_stream.swap_remove(index);
} else {
self.blocked_by_stream[index].1 -= 1;
}
}
}
#[inline]
fn enc_stream_err(_: HpackError) -> QpackError {
QpackError::EncoderStream
}
#[inline]
fn dec_failed(_: HpackError) -> QpackError {
QpackError::DecompressionFailed
}
#[inline]
fn dec_stream_err(_: HpackError) -> QpackError {
QpackError::DecoderStream
}
#[cfg(test)]
mod tests;