use std::cell::Cell;
use std::collections::VecDeque;
use super::batch_pool::{is_tight, tls_pool};
use super::merge::MemBatch;
use gnitz_wire::RowSource;
use gnitz_wire::TypeCode;
use rustc_hash::FxHashMap;
type SpanKey = (usize, usize);
type SpanMap = FxHashMap<SpanKey, usize>;
#[inline]
pub(super) fn blob_span_key(content: &[u8]) -> SpanKey {
(content.as_ptr() as usize, content.len())
}
pub(crate) fn prorated_blob_cap(src_blob: usize, src_rows: usize, out_rows: usize) -> usize {
if src_blob == 0 || src_rows == 0 {
return 0;
}
let per_row = src_blob.div_ceil(src_rows) as u128;
(per_row * out_rows as u128).min(src_blob as u128) as usize
}
const RELOCATE_CELL_COST_BYTES: usize = 500;
pub(crate) fn should_relocate_blob(src_blob: usize, src_rows: usize, out_rows: usize) -> bool {
src_rows > 0 && src_blob > out_rows.saturating_mul(RELOCATE_CELL_COST_BYTES + src_blob / src_rows)
}
pub(crate) fn heap_is_wasteful(dead: usize, heap: usize) -> bool {
dead.saturating_mul(4) > heap
}
pub(super) fn carried_dead(
heap: usize,
dead: usize,
src_rows: usize,
kept_rows: usize,
excluded: impl FnOnce() -> usize,
) -> Option<usize> {
if kept_rows == 0 {
return None;
}
if heap == 0 {
return Some(0);
}
let left = src_rows - kept_rows;
if should_relocate_blob(heap, src_rows, kept_rows)
|| heap_is_wasteful(dead + prorated_blob_cap(heap, src_rows, left), heap)
{
return None;
}
let dead = if left == 0 { dead } else { dead + excluded() };
(!heap_is_wasteful(dead, heap)).then_some(dead)
}
#[inline]
pub(super) fn cell_long_bytes(cell: &[u8]) -> usize {
gnitz_wire::german_string_heap(cell, usize::MAX).map_or(0, |span| span.len())
}
pub(super) fn long_bytes_outside<S: RowSource>(src: &S, mask: u64, kept: &[(usize, usize)]) -> usize {
let bounds = kept.iter().copied().chain([(src.row_count(), src.row_count())]);
let gaps = bounds.scan(0, |next, (start, end)| {
let gap = *next..start;
*next = end;
Some(gap)
});
gaps.flatten().map(|row| row_long_bytes(src, mask, row)).sum()
}
#[inline]
pub(super) fn row_long_bytes<S: RowSource>(src: &S, mask: u64, row: usize) -> usize {
gnitz_wire::BitIter(mask)
.map(|pi| cell_long_bytes(src.get_col_ptr(row, pi, 16)))
.sum()
}
#[inline]
pub(super) fn rebase_string_cells(
cells: &mut [u8],
src_blob: &[u8],
dst_blob: &mut Vec<u8>,
heap_at: Option<usize>,
cache: Option<&mut BlobCache>,
) {
match heap_at {
Some(base) => gnitz_wire::shift_german_string_heaps(cells, base),
None => relocate_cells(cells.as_chunks_mut::<16>().0, src_blob, dst_blob, cache),
}
}
#[inline]
pub(super) fn rebase_string_cell(
cell: &mut [u8; 16],
src_blob: &[u8],
dst_blob: &mut Vec<u8>,
heap_at: Option<usize>,
cache: Option<&mut BlobCache>,
) {
match heap_at {
Some(base) => gnitz_wire::shift_german_string_heaps(cell, base),
None => relocate_german_string(cell, src_blob, dst_blob, cache),
}
}
#[inline(never)]
fn relocate_cells(cells: &mut [[u8; 16]], src_blob: &[u8], dst_blob: &mut Vec<u8>, cache: Option<&mut BlobCache>) {
match cache {
Some(cache) => {
for cell in cells {
relocate_german_string(cell, src_blob, dst_blob, Some(&mut *cache));
}
}
None => {
for cell in cells {
relocate_german_string(cell, src_blob, dst_blob, None);
}
}
}
}
#[inline]
pub(crate) fn relocate_german_string_vec(
src_cell: &[u8],
src_blob: &[u8],
dst_blob: &mut Vec<u8>,
cache: Option<&mut BlobCache>,
) -> [u8; 16] {
let mut cell: [u8; 16] = src_cell[..16]
.try_into()
.expect("relocate_german_string_vec: src must be a 16-byte German string cell");
relocate_german_string(&mut cell, src_blob, dst_blob, cache);
cell
}
#[inline]
fn relocate_german_string(cell: &mut [u8; 16], src_blob: &[u8], dst_blob: &mut Vec<u8>, cache: Option<&mut BlobCache>) {
*cell = gnitz_wire::relocate_german_string_with(cell, src_blob, |content| place_span(content, dst_blob, cache));
}
fn place_span(content: &[u8], dst_blob: &mut Vec<u8>, cache: Option<&mut BlobCache>) -> usize {
let at = dst_blob.len();
match cache {
Some(cache) => *cache.map().entry(blob_span_key(content)).or_insert_with(|| {
dst_blob.extend_from_slice(content);
at
}),
None => {
dst_blob.extend_from_slice(content);
at
}
}
}
const BLOB_CACHE_RESERVE_CAP: usize = 4096;
thread_local! {
static BLOB_CACHE_POOL: Cell<VecDeque<SpanMap>> = const { Cell::new(VecDeque::new()) };
}
pub(crate) struct BlobCache {
map: Option<SpanMap>,
cells: usize,
}
impl BlobCache {
pub(crate) fn new(cells: usize) -> Self {
BlobCache { map: None, cells }
}
#[inline]
pub(super) fn map(&mut self) -> &mut SpanMap {
let cells = self.cells;
self.map.get_or_insert_with(|| {
let reserve = cells.min(BLOB_CACHE_RESERVE_CAP);
let fits = |m: &SpanMap| m.capacity() >= reserve && is_tight(m.capacity(), cells.max(4));
let mut map = tls_pool::take(&BLOB_CACHE_POOL, fits).unwrap_or_default();
map.reserve(reserve);
map
})
}
}
impl Drop for BlobCache {
fn drop(&mut self) {
if let Some(mut map) = self.map.take() {
let bytes = map.capacity() * std::mem::size_of::<(SpanKey, usize)>();
map.clear();
tls_pool::recycle(&BLOB_CACHE_POOL, map, bytes);
}
}
}
pub(super) fn measure_dead_heap(mb: &MemBatch<'_>) -> usize {
walk_heap_spans(mb, |_, _| true).expect("an accepting walk")
}
pub(super) fn walk_heap_spans(mb: &MemBatch<'_>, mut accept: impl FnMut(&[u8; 16], TypeCode) -> bool) -> Option<usize> {
let schema = mb.schema;
let heap = mb.blob.len();
if !schema.has_german_string() {
return Some(heap);
}
let mut live = vec![0u64; heap.div_ceil(64)];
for (pi, col) in schema.payload_columns() {
if !col.type_code.is_german_string() {
continue;
}
for cell in mb.col_data(pi, 16).as_chunks::<16>().0 {
if !accept(cell, col.type_code) {
return None;
}
if let Some(span) = gnitz_wire::german_string_heap(cell, heap) {
mark_bits(&mut live, span);
}
}
}
let marked: usize = live.iter().map(|w| w.count_ones() as usize).sum();
Some(heap - marked)
}
fn mark_bits(bits: &mut [u64], span: std::ops::Range<usize>) {
let (mut at, end) = (span.start, span.end);
while at < end {
let (word, bit) = (at / 64, at % 64);
let n = (64 - bit).min(end - at);
bits[word] |= gnitz_wire::low_bits_mask(n) << bit;
at += n;
}
}
#[cfg(test)]
#[path = "tests/string_heap.rs"]
mod tests;