use super::{
BuiltSegment, BulkBlockSharing, Error, GarbageCollectionGeneration, MAXIMUM_SEGMENT_SIZE,
RecordIdentifier, RecordType, RecordWriter, Result, SegmentSink, new_bulk_segment_identifier,
};
pub(crate) const BLOCK_SIZE: usize = 4096;
pub(crate) const VALUE_DEDUP_BUDGET_BYTES: usize = 32 * 1024 * 1024;
pub(crate) const MAXIMUM_DEDUPLICATED_VALUE_BYTES: usize = 1024;
pub(crate) const SMALL_VALUE_LIMIT: usize = 128;
pub(crate) const MEDIUM_VALUE_LIMIT: usize = (1 << 14) + 128;
pub(crate) const BLOB_IDENTIFIER_SMALL_LIMIT: usize = 4096;
impl<Sink: SegmentSink> RecordWriter<Sink> {
pub fn write_string(&mut self, text: &str) -> Result<RecordIdentifier> {
self.write_deduplicated_value(text.as_bytes())
}
pub(super) fn write_deduplicated_value(&mut self, bytes: &[u8]) -> Result<RecordIdentifier> {
if bytes.len() > MAXIMUM_DEDUPLICATED_VALUE_BYTES {
return self.write_value_bytes(bytes);
}
if let Some(existing) = self.value_cache.get(&bytes.to_vec()) {
return Ok(existing);
}
let written = self.write_value_bytes(bytes)?;
self.value_cache.insert(bytes.to_vec(), written);
Ok(written)
}
pub fn write_binary_content(&mut self, content: &[u8]) -> Result<RecordIdentifier> {
self.write_value_bytes(content)
}
pub fn copy_binary_value(
&mut self,
source: &dyn crate::content::provider::SegmentProvider,
source_value: RecordIdentifier,
sharing: BulkBlockSharing,
) -> Result<RecordIdentifier> {
let length = crate::content::value::read_value_length(source, source_value)?;
if length < MEDIUM_VALUE_LIMIT as u64 {
let content = crate::content::value::read_binary_content(source, source_value)?;
return self.write_binary_content(&content);
}
let block_count = length.div_ceil(BLOCK_SIZE as u64);
let source_view = source.segment(source_value.segment)?;
let list_identifier =
source_view.read_record_identifier(source_value.record_number, 8, 0)?;
let source_blocks =
crate::content::list::uncounted_list_entries(source, list_identifier, block_count)?;
let mut block_identifiers = Vec::with_capacity(source_blocks.len());
let mut remaining = length;
for source_block in source_blocks {
let block_length = remaining.min(BLOCK_SIZE as u64) as usize;
remaining -= block_length as u64;
if source_block.segment.is_bulk_segment() && sharing == BulkBlockSharing::WithinOneStore
{
block_identifiers.push(source_block);
continue;
}
let block_view = source.segment(source_block.segment)?;
let block_bytes = block_view.read_bytes(source_block.record_number, 0, block_length)?;
let record = self.allocate(RecordType::Block, block_length, &[])?;
self.current.record_bytes_mut(record)[..block_length].copy_from_slice(block_bytes);
block_identifiers.push(self.identifier_of(record));
}
let body =
self.write_list_body(&block_identifiers)?
.ok_or_else(|| Error::InvalidFormat {
details: "a long binary always has at least one block".to_owned(),
})?;
let record = self.allocate(RecordType::Value, 8 + 6, &[body])?;
let stored = (length - MEDIUM_VALUE_LIMIT as u64) | (0b11 << 62);
self.current.record_bytes_mut(record)[0..8].copy_from_slice(&stored.to_be_bytes());
self.write_identifier_at(record, 8, body);
Ok(self.identifier_of(record))
}
pub(super) fn write_value_bytes(&mut self, content: &[u8]) -> Result<RecordIdentifier> {
if content.len() < SMALL_VALUE_LIMIT {
let record = self.allocate(RecordType::Value, 1 + content.len(), &[])?;
let bytes = self.current.record_bytes_mut(record);
bytes[0] = content.len() as u8;
bytes[1..=content.len()].copy_from_slice(content);
return Ok(self.identifier_of(record));
}
if content.len() < MEDIUM_VALUE_LIMIT {
let record = self.allocate(RecordType::Value, 2 + content.len(), &[])?;
let stored = (content.len() - SMALL_VALUE_LIMIT) as u16 | 0x8000;
let bytes = self.current.record_bytes_mut(record);
bytes[0..2].copy_from_slice(&stored.to_be_bytes());
bytes[2..2 + content.len()].copy_from_slice(content);
return Ok(self.identifier_of(record));
}
let block_identifiers = self.write_blocks(content)?;
let list_body =
self.write_list_body(&block_identifiers)?
.ok_or_else(|| Error::InvalidFormat {
details: "a long value always has at least one block".to_owned(),
})?;
let record = self.allocate(RecordType::Value, 8 + 6, &[list_body])?;
let stored = (content.len() - MEDIUM_VALUE_LIMIT) as u64 | (0b11 << 62);
let header = stored.to_be_bytes();
self.current.record_bytes_mut(record)[0..8].copy_from_slice(&header);
self.write_identifier_at(record, 8, list_body);
Ok(self.identifier_of(record))
}
pub(super) fn write_blocks(&mut self, content: &[u8]) -> Result<Vec<RecordIdentifier>> {
let mut block_identifiers = Vec::with_capacity(content.len().div_ceil(BLOCK_SIZE));
let mut remaining = content;
while remaining.len() >= MAXIMUM_SEGMENT_SIZE {
let bulk_identifier = new_bulk_segment_identifier();
let (bulk_content, rest) = remaining.split_at(MAXIMUM_SEGMENT_SIZE);
self.sink.write_segment(BuiltSegment {
identifier: bulk_identifier,
bytes: bulk_content.to_vec(),
generation: GarbageCollectionGeneration {
generation: 0,
full_generation: 0,
is_compacted: false,
},
referenced_segments: Vec::new(),
binary_reference_identifiers: Vec::new(),
})?;
for block_offset in (0..MAXIMUM_SEGMENT_SIZE).step_by(BLOCK_SIZE) {
block_identifiers.push(RecordIdentifier::new(bulk_identifier, block_offset as u32));
}
remaining = rest;
}
for chunk in remaining.chunks(BLOCK_SIZE) {
let record = self.allocate(RecordType::Block, chunk.len(), &[])?;
self.current.record_bytes_mut(record)[..chunk.len()].copy_from_slice(chunk);
block_identifiers.push(self.identifier_of(record));
}
Ok(block_identifiers)
}
pub fn write_external_binary_identifier(
&mut self,
blob_identifier: &str,
) -> Result<RecordIdentifier> {
let encoded = blob_identifier.as_bytes();
if encoded.len() < BLOB_IDENTIFIER_SMALL_LIMIT {
let record =
self.allocate(RecordType::ExternalBlobIdentifier, 2 + encoded.len(), &[])?;
self.current
.register_binary_reference(blob_identifier.to_owned());
let stored = 0xE000u16 | encoded.len() as u16;
let bytes = self.current.record_bytes_mut(record);
bytes[0..2].copy_from_slice(&stored.to_be_bytes());
bytes[2..2 + encoded.len()].copy_from_slice(encoded);
return Ok(self.identifier_of(record));
}
let string_identifier = self.write_string(blob_identifier)?;
let record = self.allocate(
RecordType::ExternalBlobIdentifier,
1 + 6,
&[string_identifier],
)?;
self.current
.register_binary_reference(blob_identifier.to_owned());
self.current.record_bytes_mut(record)[0] = 0xF0;
self.write_identifier_at(record, 1, string_identifier);
Ok(self.identifier_of(record))
}
}
#[cfg(test)]
mod tests {
use crate::content::value::{BinaryValue, read_binary_value, read_string};
use crate::segment::record::RecordIdentifier;
use crate::writer::record_writer::test_support::new_writer;
#[test]
fn strings_round_trip_at_every_size_class() {
let mut writer = new_writer();
let small = "hello".to_owned();
let boundary_small = "x".repeat(127);
let medium = "y".repeat(4000);
let boundary_medium = "m".repeat(16511);
let long = "z".repeat(20_000);
let identifiers: Vec<(String, RecordIdentifier)> =
[small, boundary_small, medium, boundary_medium, long]
.into_iter()
.map(|text| {
let identifier = writer.write_string(&text).expect("write");
(text, identifier)
})
.collect();
let store = writer.finish().expect("finish");
for (text, identifier) in identifiers {
assert_eq!(read_string(&store, identifier).expect("read"), text);
}
}
#[test]
fn huge_values_use_bulk_segments() {
let mut writer = new_writer();
let content: Vec<u8> = (0..600 * 1024).map(|index| (index % 251) as u8).collect();
let identifier = writer.write_binary_content(&content).expect("write");
let store = writer.finish().expect("finish");
let bulk_count = store
.write_order
.iter()
.filter(|segment| segment.is_bulk_segment())
.count();
assert_eq!(bulk_count, 2, "two full 256 KiB runs become bulk segments");
match read_binary_value(&store, identifier).expect("classify") {
BinaryValue::Inline { length, .. } => assert_eq!(length, content.len() as u64),
BinaryValue::External { .. } => panic!("inline binary expected"),
}
assert_eq!(
crate::content::value::read_binary_content(&store, identifier).expect("content"),
content
);
}
#[test]
fn external_binary_identifiers_round_trip() {
let mut writer = new_writer();
let short = writer
.write_external_binary_identifier("datastore-0001")
.expect("short");
let long_identifier_text = "reference-".repeat(500);
let long = writer
.write_external_binary_identifier(&long_identifier_text)
.expect("long");
let store = writer.finish().expect("finish");
assert_eq!(
read_binary_value(&store, short).expect("read"),
BinaryValue::External {
blob_identifier: "datastore-0001".to_owned()
}
);
assert_eq!(
read_binary_value(&store, long).expect("read"),
BinaryValue::External {
blob_identifier: long_identifier_text
}
);
}
}