use crate::{
h3::{H3Error, H3ErrorCode},
headers::{
entry_name::EntryName,
qpack::{
ConnectionAccumulator, FieldLineValue,
instruction::encoder::{
encode_duplicate, encode_insert_with_literal_name, encode_insert_with_name_ref,
encode_set_capacity,
},
static_table::first_match,
},
recent_pairs::RecentPairs,
},
};
use hashbrown::HashMap;
use std::{
borrow::Cow,
collections::VecDeque,
fmt::{self, Debug},
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[cfg_attr(not(test), allow(dead_code))]
pub(in crate::headers) enum InsertPolicy {
Never,
SeenK(u8),
}
impl Default for InsertPolicy {
fn default() -> Self {
Self::SeenK(2)
}
}
pub(super) struct TableState {
pub(super) entries: VecDeque<Entry>,
pub(super) max_capacity: usize,
pub(super) capacity: usize,
pub(super) current_size: usize,
pub(super) insert_count: u64,
pub(super) known_received_count: u64,
pub(super) pending_ops: VecDeque<Vec<u8>>,
pub(super) outstanding_sections: HashMap<u64, VecDeque<SectionRefs>>,
pub(super) failed: Option<H3ErrorCode>,
pub(super) max_blocked_streams: usize,
pub(super) by_name: HashMap<EntryName<'static>, NameIndex>,
pub(super) recent_pairs: RecentPairs,
pub(super) recent_names: RecentPairs,
pub(super) name_warm: bool,
pub(super) name_k: u8,
pub(super) insert_policy: InsertPolicy,
pub(super) primed_bytes: usize,
pub(super) accum: ConnectionAccumulator,
}
impl Debug for TableState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TableState")
.field(
"entries",
&fmt::from_fn(|f| {
let mut f = f.debug_map();
for Entry { name, value, .. } in &self.entries {
f.entry(name, &format_args!("{}", String::from_utf8_lossy(value)));
}
f.finish()
}),
)
.field("max_capacity", &self.max_capacity)
.field("capacity", &self.capacity)
.field("current_size", &self.current_size)
.field("insert_count", &self.insert_count)
.field("known_received_count", &self.known_received_count)
.field("pending_ops", &self.pending_ops)
.field("outstanding_sections", &self.outstanding_sections)
.field("failed", &self.failed)
.field("max_blocked_streams", &self.max_blocked_streams)
.field("by_name", &self.by_name)
.field("recent_pairs", &self.recent_pairs)
.field("recent_names", &self.recent_names)
.field("name_warm", &self.name_warm)
.field("name_k", &self.name_k)
.field("insert_policy", &self.insert_policy)
.field("primed_bytes", &self.primed_bytes)
.field("accum", &self.accum)
.finish()
}
}
#[derive(Default)]
pub(super) struct NameIndex {
pub(super) by_value: HashMap<Cow<'static, [u8]>, u64>,
pub(super) latest_any: u64,
}
impl Debug for NameIndex {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("NameIndex")
.field(
"by_value",
&fmt::from_fn(|f| {
let mut map = f.debug_map();
for (k, v) in &self.by_value {
map.entry(&format_args!("{}", String::from_utf8_lossy(k)), v);
}
map.finish()
}),
)
.field("latest_any", &self.latest_any)
.finish()
}
}
#[derive(Clone)]
pub(super) struct Entry {
pub(super) name: EntryName<'static>,
pub(super) value: Cow<'static, [u8]>,
pub(super) size: usize,
}
impl Debug for Entry {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Entry")
.field("name", &self.name)
.field(
"value",
&format_args!("{}", String::from_utf8_lossy(&self.value)),
)
.field("size", &self.size)
.finish()
}
}
#[derive(Debug, Clone, Copy)]
pub(in crate::headers) struct SectionRefs {
pub(in crate::headers) required_insert_count: u64,
pub(in crate::headers) min_ref_abs_idx: Option<u64>,
}
impl TableState {
pub(super) fn new(recent_pairs: RecentPairs) -> Self {
let recent_names = RecentPairs::with_size(recent_pairs.size());
Self {
entries: VecDeque::new(),
max_capacity: 0,
capacity: 0,
current_size: 0,
insert_count: 0,
known_received_count: 0,
pending_ops: VecDeque::new(),
outstanding_sections: HashMap::new(),
failed: None,
max_blocked_streams: 0,
by_name: HashMap::new(),
recent_pairs,
recent_names,
name_warm: false,
name_k: 3,
insert_policy: InsertPolicy::default(),
primed_bytes: 0,
accum: ConnectionAccumulator::default(),
}
}
pub(super) fn insert(
&mut self,
name: EntryName<'_>,
value: FieldLineValue<'_>,
extra_floor: Option<u64>,
) -> Result<u64, H3Error> {
if let Some(abs_idx) = self
.by_name
.get(&name)
.and_then(|idx| idx.by_value.get(value.as_bytes()).copied())
{
return self.duplicate(abs_idx, extra_floor);
}
let entry_size = name.len() + value.len() + 32;
let (wire, variant_floor) = if let Some(static_idx) = first_match(&name) {
(
encode_insert_with_name_ref(usize::from(static_idx), true, &value),
None,
)
} else if let Some(name_abs_idx) = self.by_name.get(&name).map(|idx| idx.latest_any) {
let relative_index = self.insert_count - 1 - name_abs_idx;
let wire = encode_insert_with_name_ref(
usize::try_from(relative_index).unwrap_or(usize::MAX),
false,
&value,
);
(wire, Some(name_abs_idx))
} else {
(
encode_insert_with_literal_name(name.as_bytes(), &value),
None,
)
};
self.make_room_for(entry_size, combine_floor(variant_floor, extra_floor))?;
let value = value.into_static();
Ok(self.insert_entry(name, value, entry_size, wire))
}
pub(in crate::headers::qpack::encoder_dynamic_table) fn duplicate(
&mut self,
abs_idx: u64,
extra_floor: Option<u64>,
) -> Result<u64, H3Error> {
let entry_size = self
.entry_at_abs(abs_idx)
.expect("insert's by_value lookup guarantees abs_idx is live")
.size;
let relative_index = self.insert_count - 1 - abs_idx;
let wire = encode_duplicate(usize::try_from(relative_index).unwrap_or(usize::MAX));
self.make_room_for(entry_size, combine_floor(Some(abs_idx), extra_floor))?;
let entry = self
.entry_at_abs(abs_idx)
.expect("preserved by make_room_for floor");
let name = entry.name.clone();
let value = entry.value.clone();
Ok(self.insert_entry(name, value, entry_size, wire))
}
pub(super) fn set_capacity(&mut self, new_capacity: usize) -> Result<(), H3Error> {
if new_capacity > self.max_capacity {
log::error!(
"qpack encoder: set_capacity {} exceeds max_capacity {}",
new_capacity,
self.max_capacity,
);
return Err(H3ErrorCode::QpackEncoderStreamError.into());
}
self.evict_down_to(new_capacity)?;
self.capacity = new_capacity;
self.pending_ops
.push_back(encode_set_capacity(new_capacity));
Ok(())
}
fn make_room_for(
&mut self,
entry_size: usize,
extra_floor: Option<u64>,
) -> Result<(), H3Error> {
if entry_size > self.capacity {
return Err(H3ErrorCode::QpackEncoderStreamError.into());
}
let target = self.capacity - entry_size;
let combined_floor = combine_floor(self.eviction_floor(), extra_floor);
self.evict_down_to_with_floor(target, combined_floor)
}
fn insert_entry(
&mut self,
name: EntryName<'_>,
value: Cow<'static, [u8]>,
entry_size: usize,
wire: Vec<u8>,
) -> u64 {
let name = name.into_owned();
let abs_idx = self.insert_count;
let name_index = self.by_name.entry(name.clone()).or_default();
name_index.by_value.insert(value.clone(), abs_idx);
name_index.latest_any = abs_idx;
self.entries.push_front(Entry {
name,
value,
size: entry_size,
});
self.current_size += entry_size;
self.insert_count += 1;
self.pending_ops.push_back(wire);
log::trace!(
"qpack encoder: inserted entry abs_idx={abs_idx} size={entry_size} current_size={} \
insert_count={}",
self.current_size,
self.insert_count,
);
abs_idx
}
pub(super) fn entry_at_abs(&self, abs_idx: u64) -> Option<&Entry> {
let oldest_abs = self.insert_count.checked_sub(self.entries.len() as u64)?;
if abs_idx < oldest_abs || abs_idx >= self.insert_count {
return None;
}
let pos = usize::try_from(self.insert_count - 1 - abs_idx).ok()?;
self.entries.get(pos)
}
pub(super) fn is_stream_blocking(&self, stream_id: u64) -> bool {
self.outstanding_sections
.get(&stream_id)
.is_some_and(|sections| {
sections
.iter()
.any(|s| s.required_insert_count > self.known_received_count)
})
}
pub(super) fn currently_blocked_streams(&self) -> usize {
let krc = self.known_received_count;
self.outstanding_sections
.iter()
.filter(|(_, sections)| sections.iter().any(|s| s.required_insert_count > krc))
.count()
}
fn eviction_floor(&self) -> Option<u64> {
self.outstanding_sections
.values()
.flat_map(|sections| sections.iter())
.filter_map(|s| s.min_ref_abs_idx)
.min()
}
fn evict_down_to(&mut self, target_size: usize) -> Result<(), H3Error> {
let floor = self.eviction_floor();
self.evict_down_to_with_floor(target_size, floor)
}
fn evict_down_to_with_floor(
&mut self,
target_size: usize,
floor: Option<u64>,
) -> Result<(), H3Error> {
let mut result = Ok(());
while self.current_size > target_size {
let evicted_abs = self.insert_count - self.entries.len() as u64;
if let Some(pin) = floor
&& evicted_abs >= pin
{
if evicted_abs > pin {
log::error!(
"qpack encoder: eviction proceeded past pinned entry (current_size={}, \
target_size={target_size}, evicted_abs={evicted_abs}, pin={pin})",
self.current_size,
);
} else {
log::trace!(
"qpack encoder: eviction blocked by pinned section (current_size={}, \
target_size={target_size}, evicted_abs={evicted_abs}, pin={pin})",
self.current_size,
);
}
result = Err(H3ErrorCode::QpackEncoderStreamError.into());
break;
}
let Entry { name, value, size } = self.entries.pop_back().expect("current_size > 0");
self.current_size -= size;
self.remove_from_reverse_index(&name, value.as_ref(), evicted_abs);
log::trace!("qpack encoder: evicted entry abs_idx={evicted_abs} size={size}");
}
result
}
fn remove_from_reverse_index(
&mut self,
name: &EntryName<'static>,
value: &[u8],
evicted_abs: u64,
) {
let Some(name_index) = self.by_name.get_mut(name) else {
return;
};
if name_index.by_value.get(value) == Some(&evicted_abs) {
name_index.by_value.remove(value);
}
let drop_name_entry = if name_index.latest_any == evicted_abs {
match name_index.by_value.values().copied().max() {
Some(newest) => {
name_index.latest_any = newest;
false
}
None => true,
}
} else {
false
};
if drop_name_entry {
self.by_name.remove(name);
}
}
}
fn combine_floor(a: Option<u64>, b: Option<u64>) -> Option<u64> {
match (a, b) {
(Some(a), Some(b)) => Some(a.min(b)),
(a, b) => a.or(b),
}
}