use std::io::Write;
use crate::error::{Error, Result};
use crate::external_sort::{SortedPasses, SpillRecord};
use crate::index::lucene::codec::data_output::CodecOutput;
use crate::index::lucene::codec::packed::{
PACKED_VERSION_CURRENT, bits_required, write_block_packed, write_monotonic_block_packed,
write_packed,
};
const DATA_CODEC: &str = "Lucene45DocValuesData";
const META_CODEC: &str = "Lucene45ValuesMetadata";
const VERSION_CURRENT: i32 = 1;
const NUMERIC: u8 = 0;
const BINARY: u8 = 1;
const SORTED: u8 = 2;
const SORTED_SET: u8 = 3;
pub const BLOCK_SIZE: usize = 16384;
pub const ADDRESS_INTERVAL: usize = 16;
const DELTA_COMPRESSED: i32 = 0;
const GCD_COMPRESSED: i32 = 1;
const TABLE_COMPRESSED: i32 = 2;
const BINARY_FIXED_UNCOMPRESSED: i32 = 0;
const BINARY_PREFIX_COMPRESSED: i32 = 2;
const SORTED_SET_WITH_ADDRESSES: i32 = 0;
const SORTED_SET_SINGLE_VALUED_SORTED: i32 = 1;
pub const MISSING_ORD: i64 = -1;
const TABLE_LIMIT: usize = 256;
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct NumericRecord {
pub document: i32,
pub value: Option<i64>,
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct DictionaryRecord {
pub ordinal: i64,
pub value: Vec<u8>,
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct OrdinalRecord {
pub document: i32,
pub ordinal: i64,
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub struct SetOrdinalRecord {
pub document: i32,
pub ordinal: i64,
}
macro_rules! numeric_spill_record {
($record:ident, $first:ident, $second:ident) => {
impl SpillRecord for $record {
fn encode(&self, buffer: &mut Vec<u8>) {
buffer.extend_from_slice(&self.$first.to_be_bytes());
buffer.extend_from_slice(&self.$second.to_be_bytes());
}
fn decode(bytes: &[u8]) -> Result<Self> {
if bytes.len() != 12 {
return Err(Error::InvalidFormat {
details: format!("a doc-value record is twelve bytes, not {}", bytes.len()),
});
}
Ok(Self {
$first: i32::from_be_bytes(bytes[..4].try_into().expect("four bytes")),
$second: i64::from_be_bytes(bytes[4..].try_into().expect("eight bytes")),
})
}
fn resident_size(&self) -> usize {
12
}
}
};
}
numeric_spill_record!(OrdinalRecord, document, ordinal);
numeric_spill_record!(SetOrdinalRecord, document, ordinal);
impl SpillRecord for NumericRecord {
fn encode(&self, buffer: &mut Vec<u8>) {
buffer.extend_from_slice(&self.document.to_be_bytes());
buffer.push(u8::from(self.value.is_some()));
buffer.extend_from_slice(&self.value.unwrap_or(0).to_be_bytes());
}
fn decode(bytes: &[u8]) -> Result<Self> {
if bytes.len() != 13 {
return Err(Error::InvalidFormat {
details: format!("a numeric record is thirteen bytes, not {}", bytes.len()),
});
}
let value = i64::from_be_bytes(bytes[5..].try_into().expect("eight bytes"));
Ok(Self {
document: i32::from_be_bytes(bytes[..4].try_into().expect("four bytes")),
value: (bytes[4] == 1).then_some(value),
})
}
fn resident_size(&self) -> usize {
13
}
}
impl SpillRecord for DictionaryRecord {
fn encode(&self, buffer: &mut Vec<u8>) {
buffer.extend_from_slice(&self.ordinal.to_be_bytes());
buffer.extend_from_slice(&self.value);
}
fn decode(bytes: &[u8]) -> Result<Self> {
if bytes.len() < 8 {
return Err(Error::InvalidFormat {
details: format!(
"a dictionary record is at least eight bytes, not {}",
bytes.len()
),
});
}
Ok(Self {
ordinal: i64::from_be_bytes(bytes[..8].try_into().expect("eight bytes")),
value: bytes[8..].to_vec(),
})
}
fn resident_size(&self) -> usize {
8 + self.value.len()
}
}
#[derive(Debug)]
pub struct DocValuesFiles<Sink> {
pub data: Sink,
pub metadata: Sink,
}
pub trait DictionaryStream {
fn walk(&mut self, visit: &mut dyn FnMut(&[u8]) -> Result<()>) -> Result<()>;
}
impl DictionaryStream for SortedPasses<DictionaryRecord> {
fn walk(&mut self, visit: &mut dyn FnMut(&[u8]) -> Result<()>) -> Result<()> {
for record in self.pass()? {
visit(&record?.value)?;
}
Ok(())
}
}
trait NumericStream {
fn walk(&mut self, visit: &mut dyn FnMut(Option<i64>) -> Result<()>) -> Result<()>;
}
struct MappedStream<'source, Record: SpillRecord, Map> {
source: &'source mut SortedPasses<Record>,
map: Map,
}
impl<Record: SpillRecord, Map: FnMut(&Record) -> Option<i64>> NumericStream
for MappedStream<'_, Record, Map>
{
fn walk(&mut self, visit: &mut dyn FnMut(Option<i64>) -> Result<()>) -> Result<()> {
for record in self.source.pass()? {
let record = record?;
visit((self.map)(&record))?;
}
Ok(())
}
}
struct SingleOrdinalStream<'source> {
source: &'source mut SortedPasses<SetOrdinalRecord>,
document_count: i64,
}
impl NumericStream for SingleOrdinalStream<'_> {
fn walk(&mut self, visit: &mut dyn FnMut(Option<i64>) -> Result<()>) -> Result<()> {
let mut document = 0i64;
for record in self.source.pass()? {
let record = record?;
let at = i64::from(record.document);
while document < at {
visit(Some(MISSING_ORD))?;
document += 1;
}
visit(Some(record.ordinal))?;
document += 1;
}
while document < self.document_count {
visit(Some(MISSING_ORD))?;
document += 1;
}
Ok(())
}
}
fn greatest_common_divisor(first: i64, second: i64) -> i64 {
let mut left = first.unsigned_abs();
let mut right = second.unsigned_abs();
if left == 0 {
return right as i64;
}
if right == 0 {
return left as i64;
}
let common = (left | right).trailing_zeros();
left >>= left.trailing_zeros();
loop {
right >>= right.trailing_zeros();
if left == right {
break;
}
if left > right {
std::mem::swap(&mut left, &mut right);
}
right -= left;
}
(left << common) as i64
}
fn common_prefix_length(left: &[u8], right: &[u8]) -> usize {
left.iter()
.zip(right.iter())
.take_while(|(first, second)| first == second)
.count()
}
fn write_block_packed_stream<Sink: Write>(
data: &mut CodecOutput<Sink>,
stream: &mut dyn NumericStream,
transform: &dyn Fn(i64) -> i64,
) -> Result<()> {
let mut block: Vec<i64> = Vec::with_capacity(BLOCK_SIZE);
let mut failure = None;
stream.walk(&mut |value| {
block.push(transform(value.unwrap_or(0)));
if block.len() == BLOCK_SIZE {
if let Err(error) = write_block_packed(data, &block) {
failure = Some(error);
}
block.clear();
}
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
if !block.is_empty() {
write_block_packed(data, &block)?;
}
Ok(())
}
fn write_monotonic_blocks<Sink: Write>(data: &mut CodecOutput<Sink>, values: &[i64]) -> Result<()> {
for block in values.chunks(BLOCK_SIZE) {
write_monotonic_block_packed(data, block)?;
}
Ok(())
}
fn write_packed_stream<Sink: Write>(
data: &mut CodecOutput<Sink>,
stream: &mut dyn NumericStream,
encode: &dyn Fn(i64) -> Result<u64>,
bits: u32,
) -> Result<()> {
let mut group: Vec<u64> = Vec::with_capacity(8);
let mut failure = None;
stream.walk(&mut |value| {
group.push(encode(value.unwrap_or(0))?);
if group.len() == 8 {
if let Err(error) = write_packed(data, &group, bits) {
failure = Some(error);
}
group.clear();
}
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
if !group.is_empty() {
write_packed(data, &group, bits)?;
}
Ok(())
}
pub struct DocValuesConsumer<Sink: Write> {
data: CodecOutput<Sink>,
metadata: CodecOutput<Sink>,
document_count: i64,
}
impl<Sink: Write> DocValuesConsumer<Sink> {
pub fn new(data: Sink, metadata: Sink, document_count: i64) -> Result<Self> {
let mut data = CodecOutput::new(data);
let mut metadata = CodecOutput::new(metadata);
data.write_header(DATA_CODEC, VERSION_CURRENT)?;
metadata.write_header(META_CODEC, VERSION_CURRENT)?;
Ok(Self {
data,
metadata,
document_count,
})
}
pub fn add_numeric(
&mut self,
field_number: i32,
values: &mut SortedPasses<NumericRecord>,
) -> Result<()> {
let mut expected = 0i32;
for record in values.pass()? {
let record = record?;
if record.document != expected {
return Err(Error::InvalidFormat {
details: format!(
"a numeric doc-value stream is dense and in document order: document \
{expected} was expected and {} came",
record.document
),
});
}
expected += 1;
}
if i64::from(expected) != self.document_count {
return Err(Error::InvalidFormat {
details: format!(
"a numeric doc-value stream holds one value per document: {expected} for a \
segment of {} documents",
self.document_count
),
});
}
let mut stream = MappedStream {
source: values,
map: |record: &NumericRecord| record.value,
};
self.write_numeric_entry(field_number, true, &mut stream)
}
pub fn add_sorted(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
ordinals: &mut SortedPasses<OrdinalRecord>,
) -> Result<()> {
let mut stream = MappedStream {
source: ordinals,
map: |record: &OrdinalRecord| Some(record.ordinal),
};
self.write_sorted(field_number, values, &mut stream)
}
pub fn add_sorted_set(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
document_ordinals: &mut SortedPasses<SetOrdinalRecord>,
) -> Result<()> {
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(SORTED_SET)?;
let mut single_valued = true;
Self::walk_counts(document_ordinals, self.document_count, &mut |count| {
if count > 1 {
single_valued = false;
}
Ok(())
})?;
if single_valued {
self.metadata.write_vint(SORTED_SET_SINGLE_VALUED_SORTED)?;
let mut stream = SingleOrdinalStream {
source: document_ordinals,
document_count: self.document_count,
};
return self.write_sorted(field_number, values, &mut stream);
}
self.metadata.write_vint(SORTED_SET_WITH_ADDRESSES)?;
self.add_terms_dictionary(field_number, values)?;
let mut stream = MappedStream {
source: document_ordinals,
map: |record: &SetOrdinalRecord| Some(record.ordinal),
};
self.write_numeric_entry(field_number, false, &mut stream)?;
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(NUMERIC)?;
self.metadata.write_vint(DELTA_COMPRESSED)?;
self.metadata.write_long(-1)?;
self.metadata.write_vint(PACKED_VERSION_CURRENT)?;
self.metadata.write_long(self.data.position() as i64)?;
self.metadata.write_vlong(self.document_count)?;
self.metadata.write_vint(BLOCK_SIZE as i32)?;
let mut addresses: Vec<i64> = Vec::new();
let mut running = 0i64;
let mut failure = None;
Self::walk_counts(document_ordinals, self.document_count, &mut |count| {
running += count;
addresses.push(running);
if addresses.len() == BLOCK_SIZE {
if let Err(error) = write_monotonic_block_packed(&mut self.data, &addresses) {
failure = Some(error);
}
addresses.clear();
}
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
if !addresses.is_empty() {
write_monotonic_block_packed(&mut self.data, &addresses)?;
}
Ok(())
}
pub fn finish(mut self) -> Result<DocValuesFiles<Sink>> {
self.metadata.write_vint(-1)?;
Ok(DocValuesFiles {
data: self.data.into_inner(),
metadata: self.metadata.into_inner(),
})
}
fn walk_counts(
source: &mut SortedPasses<SetOrdinalRecord>,
document_count: i64,
visit: &mut dyn FnMut(i64) -> Result<()>,
) -> Result<()> {
let mut document = 0i64;
let mut count = 0i64;
for record in source.pass()? {
let record = record?;
let at = i64::from(record.document);
if at < document || at >= document_count {
return Err(Error::InvalidFormat {
details: format!(
"a sorted-set stream is in document order and inside the segment: \
document {at} follows {document} of {document_count}"
),
});
}
while document < at {
visit(count)?;
count = 0;
document += 1;
}
count += 1;
}
while document < document_count {
visit(count)?;
count = 0;
document += 1;
}
Ok(())
}
fn write_sorted(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
ordinals: &mut dyn NumericStream,
) -> Result<()> {
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(SORTED)?;
self.add_terms_dictionary(field_number, values)?;
self.write_numeric_entry(field_number, false, ordinals)
}
}
impl<Sink: Write> DocValuesConsumer<Sink> {
fn write_numeric_entry(
&mut self,
field_number: i32,
optimize_storage: bool,
stream: &mut dyn NumericStream,
) -> Result<()> {
let statistics = Self::measure(optimize_storage, stream)?;
let format = statistics.format();
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(NUMERIC)?;
self.metadata.write_vint(format)?;
if statistics.missing {
self.metadata.write_long(self.data.position() as i64)?;
Self::write_missing_bitset(&mut self.data, stream)?;
} else {
self.metadata.write_long(-1)?;
}
self.metadata.write_vint(PACKED_VERSION_CURRENT)?;
self.metadata.write_long(self.data.position() as i64)?;
self.metadata.write_vlong(statistics.count as i64)?;
self.metadata.write_vint(BLOCK_SIZE as i32)?;
match format {
GCD_COMPRESSED => {
self.metadata.write_long(statistics.minimum)?;
self.metadata.write_long(statistics.divisor)?;
let minimum = statistics.minimum;
let divisor = statistics.divisor;
write_block_packed_stream(&mut self.data, stream, &move |value| {
(value - minimum) / divisor
})
}
TABLE_COMPRESSED => {
let table: Vec<i64> = statistics.unique.clone().unwrap_or_default();
self.metadata.write_vint(table.len() as i32)?;
for value in &table {
self.metadata.write_long(*value)?;
}
let bits = bits_required((table.len() as u64).wrapping_sub(1));
let lookup = table;
write_packed_stream(
&mut self.data,
stream,
&move |value| {
lookup
.binary_search(&value)
.map(|at| at as u64)
.map_err(|_| Error::InvalidFormat {
details: format!(
"{value} is not in the table the statistics pass built; \
the two passes saw different values"
),
})
},
bits,
)
}
_ => write_block_packed_stream(&mut self.data, stream, &|value| value),
}
}
fn measure(optimize_storage: bool, stream: &mut dyn NumericStream) -> Result<Statistics> {
let mut statistics = Statistics::default();
if !optimize_storage {
stream.walk(&mut |_| {
statistics.count += 1;
Ok(())
})?;
return Ok(statistics);
}
let mut unique = Some(Vec::<i64>::new());
let mut disable_table = false;
stream.walk(&mut |value| {
let value = value.unwrap_or_else(|| {
statistics.missing = true;
0
});
if statistics.divisor != 1 {
if !(i64::MIN / 2..=i64::MAX / 2).contains(&value) {
statistics.divisor = 1;
} else if statistics.count != 0 {
statistics.divisor =
greatest_common_divisor(statistics.divisor, value - statistics.minimum);
}
}
statistics.minimum = statistics.minimum.min(value);
statistics.maximum = statistics.maximum.max(value);
if let Some(values) = unique.as_mut()
&& let Err(at) = values.binary_search(&value)
{
values.insert(at, value);
disable_table = values.len() > TABLE_LIMIT;
}
if disable_table {
unique = None;
}
statistics.count += 1;
Ok(())
})?;
statistics.unique = unique;
Ok(statistics)
}
fn write_missing_bitset(
data: &mut CodecOutput<Sink>,
stream: &mut dyn NumericStream,
) -> Result<()> {
let mut bits = 0u8;
let mut count = 0usize;
let mut failure = None;
stream.walk(&mut |value| {
if count == 8 {
if let Err(error) = data.write_byte(bits) {
failure = Some(error);
}
count = 0;
bits = 0;
}
if value.is_some() {
bits |= 1 << (count & 7);
}
count += 1;
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
if count > 0 {
data.write_byte(bits)?;
}
Ok(())
}
fn add_terms_dictionary(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
) -> Result<()> {
let mut minimum = i32::MAX;
let mut maximum = i32::MIN;
values.walk(&mut |value| {
let length = value.len() as i32;
minimum = minimum.min(length);
maximum = maximum.max(length);
Ok(())
})?;
if minimum == maximum {
return self.add_fixed_length_dictionary(field_number, values, minimum);
}
self.add_prefix_compressed_dictionary(field_number, values, minimum, maximum)
}
fn add_fixed_length_dictionary(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
length: i32,
) -> Result<()> {
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(BINARY)?;
let start = self.data.position();
let mut count = 0u64;
let mut failure = None;
let data = &mut self.data;
values.walk(&mut |value| {
if let Err(error) = data.write_bytes(value) {
failure = Some(error);
}
count += 1;
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
self.metadata.write_vint(BINARY_FIXED_UNCOMPRESSED)?;
self.metadata.write_long(-1)?;
self.metadata.write_vint(length)?;
self.metadata.write_vint(length)?;
self.metadata.write_vlong(count as i64)?;
self.metadata.write_long(start as i64)?;
Ok(())
}
fn add_prefix_compressed_dictionary(
&mut self,
field_number: i32,
values: &mut impl DictionaryStream,
minimum: i32,
maximum: i32,
) -> Result<()> {
self.metadata.write_vint(field_number)?;
self.metadata.write_byte(BINARY)?;
self.metadata.write_vint(BINARY_PREFIX_COMPRESSED)?;
self.metadata.write_long(-1)?;
let start = self.data.position();
let mut addresses: Vec<i64> = Vec::new();
let mut last: Vec<u8> = Vec::new();
let mut count = 0u64;
let mut failure = None;
let data = &mut self.data;
values.walk(&mut |value| {
let mut write = || -> Result<()> {
if count.is_multiple_of(ADDRESS_INTERVAL as u64) {
addresses.push((data.position() - start) as i64);
last.clear();
}
let shared = common_prefix_length(&last, value);
data.write_vint(shared as i32)?;
data.write_vint((value.len() - shared) as i32)?;
data.write_bytes(&value[shared..])?;
last.clear();
last.extend_from_slice(value);
count += 1;
Ok(())
};
if let Err(error) = write() {
failure = Some(error);
}
Ok(())
})?;
if let Some(error) = failure {
return Err(error);
}
let index_start = self.data.position();
write_monotonic_blocks(&mut self.data, &addresses)?;
self.metadata.write_vint(minimum)?;
self.metadata.write_vint(maximum)?;
self.metadata.write_vlong(count as i64)?;
self.metadata.write_long(start as i64)?;
self.metadata.write_vint(ADDRESS_INTERVAL as i32)?;
self.metadata.write_long(index_start as i64)?;
self.metadata.write_vint(PACKED_VERSION_CURRENT)?;
self.metadata.write_vint(BLOCK_SIZE as i32)?;
Ok(())
}
}
struct Statistics {
count: u64,
minimum: i64,
maximum: i64,
divisor: i64,
missing: bool,
unique: Option<Vec<i64>>,
}
impl Default for Statistics {
fn default() -> Self {
Self {
count: 0,
minimum: i64::MAX,
maximum: i64::MIN,
divisor: 0,
missing: false,
unique: None,
}
}
}
impl Statistics {
fn format(&self) -> i32 {
let delta = self.maximum.wrapping_sub(self.minimum);
if let Some(values) = self.unique.as_ref() {
let narrower =
bits_required((values.len() as u64).wrapping_sub(1)) < bits_required(delta as u64);
if (delta < 0 || narrower) && i32::try_from(self.count).is_ok() {
return TABLE_COMPRESSED;
}
}
if self.divisor != 0 && self.divisor != 1 {
GCD_COMPRESSED
} else {
DELTA_COMPRESSED
}
}
}