use crate::cache::BoundedCache;
use crate::content::property::PropertyType;
use crate::error::{Error, Result};
use crate::hashing::{compare_utf16_strings, map_entry_hash};
use crate::segment::identifier::SegmentIdentifier;
use crate::segment::parsed_segment::MAXIMUM_SEGMENT_SIZE;
use crate::segment::record::{RecordIdentifier, RecordType};
use crate::writer::identifier_generator::{
new_bulk_segment_identifier, new_data_segment_identifier,
};
use crate::writer::segment_builder::{
BuiltSegment, GarbageCollectionGeneration, SegmentBufferBuilder,
};
mod collections;
mod nodes;
#[cfg(test)]
mod test_support;
mod values;
pub use nodes::*;
pub(crate) use values::*;
pub trait SegmentSink {
fn write_segment(&mut self, segment: BuiltSegment) -> Result<()>;
}
pub struct RecordWriter<Sink: SegmentSink> {
pub(crate) sink: Sink,
pub(crate) generation: GarbageCollectionGeneration,
pub(crate) writer_identifier: String,
pub(crate) segment_sequence: u32,
pub(crate) current: SegmentBufferBuilder,
pub(crate) value_cache: BoundedCache<Vec<u8>, RecordIdentifier>,
pub(crate) template_cache: BoundedCache<TemplateKey, RecordIdentifier>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum BulkBlockSharing {
WithinOneStore,
AcrossStores,
}
impl<Sink: SegmentSink> RecordWriter<Sink> {
#[must_use]
pub fn new(sink: Sink, generation: GarbageCollectionGeneration) -> Self {
Self::with_writer_identifier(sink, generation, "froe")
}
#[must_use]
pub fn with_writer_identifier(
sink: Sink,
generation: GarbageCollectionGeneration,
writer_identifier: &str,
) -> Self {
let mut writer = Self {
sink,
generation,
writer_identifier: writer_identifier.to_owned(),
segment_sequence: 0,
current: SegmentBufferBuilder::new(new_data_segment_identifier(), generation),
value_cache: BoundedCache::new(VALUE_DEDUP_BUDGET_BYTES),
template_cache: BoundedCache::new(TEMPLATE_DEDUP_BUDGET_BYTES),
};
writer.write_segment_info_record();
writer
}
#[must_use]
pub fn sink(&self) -> &Sink {
&self.sink
}
pub fn finish(mut self) -> Result<Sink> {
self.flush_current_segment()?;
Ok(self.sink)
}
pub fn flush_current_segment(&mut self) -> Result<()> {
if self.current.record_count() <= 1 {
return Ok(());
}
let mut fresh = SegmentBufferBuilder::new(new_data_segment_identifier(), self.generation);
std::mem::swap(&mut fresh, &mut self.current);
self.write_segment_info_record();
self.sink.write_segment(fresh.finish())
}
pub(crate) fn write_segment_info_record(&mut self) {
let milliseconds = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |duration| duration.as_millis());
let info = format!(
"{{\"wid\":\"{}\",\"sno\":{},\"t\":{milliseconds}}}",
self.writer_identifier, self.segment_sequence
);
self.segment_sequence = self.segment_sequence.wrapping_add(1);
let bytes = info.as_bytes();
if bytes.len() < SMALL_VALUE_LIMIT {
let record = self
.current
.allocate(RecordType::Value, 1 + bytes.len(), &[])
.expect("the segment-info record always fits an empty segment");
let target = self.current.record_bytes_mut(record);
target[0] = bytes.len() as u8;
target[1..=bytes.len()].copy_from_slice(bytes);
} else {
let truncated = &bytes[..bytes.len().min(MEDIUM_VALUE_LIMIT - 1)];
let record = self
.current
.allocate(RecordType::Value, 2 + truncated.len(), &[])
.expect("the segment-info record always fits an empty segment");
let stored = (truncated.len() - SMALL_VALUE_LIMIT) as u16 | 0x8000;
let target = self.current.record_bytes_mut(record);
target[0..2].copy_from_slice(&stored.to_be_bytes());
target[2..2 + truncated.len()].copy_from_slice(truncated);
}
}
pub(crate) fn allocate(
&mut self,
record_type: RecordType,
size: usize,
referenced: &[RecordIdentifier],
) -> Result<u32> {
let referenced_segments: Vec<SegmentIdentifier> = referenced
.iter()
.map(|identifier| identifier.segment)
.collect();
if let Ok(record_number) = self
.current
.allocate(record_type, size, &referenced_segments)
{
return Ok(record_number);
}
self.flush_current_segment()?;
self.current
.allocate(record_type, size, &referenced_segments)
.map_err(|_| Error::InvalidFormat {
details: format!("record of {size} bytes cannot fit even an empty segment"),
})
}
pub(crate) fn write_identifier_at(
&mut self,
record_number: u32,
offset: usize,
identifier: RecordIdentifier,
) {
let reference = self.current.reference_for(identifier.segment);
let bytes = self.current.record_bytes_mut(record_number);
let target: &mut [u8; 6] = (&mut bytes[offset..offset + 6])
.try_into()
.expect("the slice is exactly six bytes");
SegmentBufferBuilder::write_record_identifier_bytes(
reference,
identifier.record_number,
target,
);
}
pub(crate) fn identifier_of(&self, record_number: u32) -> RecordIdentifier {
RecordIdentifier::new(self.current.identifier(), record_number)
}
}
#[cfg(test)]
mod tests {
use super::RecordWriter;
use crate::content::value::read_string;
use crate::segment::record::RecordIdentifier;
use crate::writer::record_writer::nodes::ChildNodesToWrite;
use crate::writer::record_writer::test_support::{MemoryStore, new_writer};
#[test]
fn a_repeated_string_or_template_reuses_the_record_it_already_wrote() {
let mut writer = new_writer();
let first = writer.write_string("nt:unstructured").expect("first");
let second = writer.write_string("nt:unstructured").expect("second");
assert_eq!(first, second, "an identical value reuses its record");
let distinct = writer.write_string("cq:Page").expect("distinct");
assert_ne!(first, distinct, "a different value gets its own record");
let shape = |writer: &mut RecordWriter<MemoryStore>| {
writer
.write_template(
Some("nt:unstructured"),
&["mix:versionable".to_owned()],
&ChildNodesToWrite::Zero,
&[],
)
.expect("template")
};
let first_template = shape(&mut writer);
let second_template = shape(&mut writer);
assert_eq!(
first_template, second_template,
"an identical shape reuses its template record"
);
let other_template = writer
.write_template(Some("cq:Page"), &[], &ChildNodesToWrite::Zero, &[])
.expect("other template");
assert_ne!(
first_template, other_template,
"a different shape gets its own template record"
);
}
#[test]
fn writing_rolls_over_across_segments() {
let mut writer = new_writer();
let identifiers: Vec<(String, RecordIdentifier)> = (0..40)
.map(|index| {
let text = format!("{index:04}").repeat(4000);
let identifier = writer.write_string(&text).expect("write");
(text, identifier)
})
.collect();
let store = writer.finish().expect("finish");
assert!(
store.write_order.len() > 1,
"forty 16 KiB strings cannot fit one segment"
);
for (text, identifier) in identifiers {
assert_eq!(read_string(&store, identifier).expect("read"), text);
}
}
}