use crate::observability::read_metrics;
use crate::observability::read_phase::ReadPhase;
use crate::storage::cache::{ChunkKey, DecompressedChunkCache};
use crate::storage::sstable::compression::{Compression, CompressionAlgorithm};
use crate::storage::sstable::compression_info::CompressionInfo;
use crate::storage::sstable::reader::block_io::read_compressed_chunk_at;
use crate::storage::sstable::reader::data_access::DECOMPRESS_CALLS;
use crate::storage::sstable::reader::read_at::ReadAt;
use crate::{Error, Result};
use bytes::Bytes;
use std::sync::atomic::Ordering;
pub(crate) struct ChunkSource<'a> {
source: &'a dyn ReadAt,
comp_info: &'a CompressionInfo,
compression: Option<&'a Compression>,
cache: &'a DecompressedChunkCache,
file_size: u64,
header_offset: u64,
namespace: u64,
cache_id: u64,
}
pub(crate) fn compression_label_of(algorithm: Option<&CompressionAlgorithm>) -> &'static str {
match algorithm {
Some(a) => read_metrics::compression_attr(a),
None => read_metrics::COMPRESSION_NONE,
}
}
pub(crate) fn count_raw_chunk(bytes: &[u8], algorithm: Option<&CompressionAlgorithm>) {
read_metrics::record_decompressed_bytes(bytes.len(), compression_label_of(algorithm));
}
pub(crate) fn counted_raw_chunk(
bytes: Vec<u8>,
algorithm: Option<&CompressionAlgorithm>,
) -> Vec<u8> {
count_raw_chunk(&bytes, algorithm);
bytes
}
pub(crate) fn count_uncompressed_block(
compression_info: &Option<std::sync::Arc<CompressionInfo>>,
read: Result<Option<Vec<u8>>>,
) -> Result<Option<Vec<u8>>> {
if compression_info.is_none() {
if let Ok(Some(block)) = &read {
count_raw_chunk(block, None);
}
}
read
}
impl<'a> ChunkSource<'a> {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
source: &'a dyn ReadAt,
comp_info: &'a CompressionInfo,
compression: Option<&'a Compression>,
cache: &'a DecompressedChunkCache,
file_size: u64,
header_offset: u64,
namespace: u64,
cache_id: u64,
) -> Self {
Self {
source,
comp_info,
compression,
cache,
file_size,
header_offset,
namespace,
cache_id,
}
}
fn compression_label(&self) -> &'static str {
compression_label_of(self.compression.map(|c| c.algorithm()))
}
pub(crate) fn chunk(&self, index: usize) -> Result<Option<Bytes>> {
let key = ChunkKey::new(self.cache_id ^ self.namespace, index as u64);
if let Some(hit) = self.cache.get(&key) {
return Ok(Some(hit));
}
let compressed = match read_compressed_chunk_at(
self.source,
self.comp_info,
index,
self.file_size,
self.header_offset,
)? {
Some(c) => c,
None => return Ok(None), };
let incompressible = compressed.len() >= self.comp_info.max_compressed_length as usize;
Ok(Some(self.decode_and_cache(
key,
compressed,
incompressible,
)?))
}
pub(crate) fn decode_and_cache(
&self,
key: ChunkKey,
compressed: Vec<u8>,
incompressible: bool,
) -> Result<Bytes> {
if incompressible {
count_raw_chunk(&compressed, self.compression.map(|c| c.algorithm()));
return Ok(self.cache.insert(key, compressed));
}
let decompressed = if let Some(compression) = self.compression {
let d = crate::observability::read_phase::timed(ReadPhase::Decompress, || {
compression.decompress(&compressed)
})
.map_err(|e| {
Error::corruption(format!(
"ChunkSource: failed to decompress chunk (key={:?}): {}",
key, e
))
})?;
DECOMPRESS_CALLS.fetch_add(1, Ordering::Relaxed);
crate::storage::sstable::read_work_counters::record_chunk_path_alloc();
d
} else {
compressed
};
read_metrics::record_decompressed_bytes(decompressed.len(), self.compression_label());
Ok(self.cache.insert(key, decompressed))
}
pub(crate) fn decode_borrowed(&self, key: ChunkKey, compressed: &[u8]) -> Result<Bytes> {
let decompressed = if let Some(compression) = self.compression {
let d = crate::observability::read_phase::timed(ReadPhase::Decompress, || {
compression.decompress(compressed)
})
.map_err(|e| {
Error::corruption(format!(
"ChunkSource: failed to decompress chunk (key={:?}): {}",
key, e
))
})?;
DECOMPRESS_CALLS.fetch_add(1, Ordering::Relaxed);
crate::storage::sstable::read_work_counters::record_chunk_path_alloc();
d
} else {
compressed.to_vec()
};
read_metrics::record_decompressed_bytes(decompressed.len(), self.compression_label());
Ok(self.cache.insert(key, decompressed))
}
pub(crate) fn decompress_only(
compression: Option<&Compression>,
compressed: Vec<u8>,
) -> Result<Vec<u8>> {
if let Some(c) = compression {
let compressed_len = compressed.len();
let decompressed =
crate::observability::read_phase::timed(ReadPhase::Decompress, || {
c.decompress(&compressed)
})
.map_err(|e| {
Error::corruption(format!(
"ChunkSource: decompress failed ({} compressed bytes): {}",
compressed_len, e
))
})?;
read_metrics::record_decompressed_bytes(
decompressed.len(),
read_metrics::compression_attr(c.algorithm()),
);
Ok(decompressed)
} else {
read_metrics::record_decompressed_bytes(
compressed.len(),
read_metrics::COMPRESSION_NONE,
);
Ok(compressed)
}
}
}