use std::io::Write;
use crate::error::{Error, Result};
use crate::index::lucene::codec::data_output::CodecOutput;
use crate::index::lucene::codec::packed::{
FORMAT_PACKED, PACKED_VERSION_CURRENT, bits_required, write_packed,
};
pub const BLOCK_SIZE: usize = 128;
pub const VERSION_CURRENT: i32 = 1;
pub const TERMS_CODEC: &str = "Lucene41PostingsWriterTerms";
const DOC_CODEC: &str = "Lucene41PostingsWriterDoc";
const POS_CODEC: &str = "Lucene41PostingsWriterPos";
const PAY_CODEC: &str = "Lucene41PostingsWriterPay";
const ALL_VALUES_EQUAL: u8 = 0;
const MAX_SKIP_LEVELS: usize = 10;
const SKIP_MULTIPLIER: usize = 8;
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord)]
pub enum IndexOptions {
Documents,
DocumentsAndFrequencies,
DocumentsAndFrequenciesAndPositions,
DocumentsAndFrequenciesAndPositionsAndOffsets,
}
impl IndexOptions {
#[must_use]
pub fn has_frequencies(self) -> bool {
self >= Self::DocumentsAndFrequencies
}
#[must_use]
pub fn has_positions(self) -> bool {
self >= Self::DocumentsAndFrequenciesAndPositions
}
#[must_use]
pub fn has_offsets(self) -> bool {
self == Self::DocumentsAndFrequenciesAndPositionsAndOffsets
}
}
#[derive(Clone, Copy, Debug)]
pub struct SegmentShape {
pub document_count: i32,
pub any_field_has_positions: bool,
pub any_field_has_offsets: bool,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct LongCount(u8);
impl LongCount {
#[must_use]
pub const fn value(self) -> usize {
self.0 as usize
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct TermMetadata {
pub document_start: u64,
pub position_start: u64,
pub payload_start: u64,
pub singleton_document: Option<i32>,
pub last_position_block_offset: Option<u64>,
pub skip_offset: Option<u64>,
pub document_frequency: i32,
pub total_term_frequency: i64,
}
#[derive(Debug)]
pub struct PostingsFiles<Sink> {
pub document: Sink,
pub position: Option<Sink>,
pub payload: Option<Sink>,
}
pub fn write_terms_header<Sink: Write>(output: &mut CodecOutput<Sink>) -> Result<()> {
output.write_header(TERMS_CODEC, VERSION_CURRENT)?;
output.write_vint(BLOCK_SIZE as i32)
}
fn write_format_table<Sink: Write>(output: &mut CodecOutput<Sink>) -> Result<()> {
output.write_vint(PACKED_VERSION_CURRENT)?;
for bits_per_value in 1..=32i32 {
output.write_vint((FORMAT_PACKED << 5) | (bits_per_value - 1))?;
}
Ok(())
}
fn write_block<Sink: Write>(output: &mut CodecOutput<Sink>, values: &[i32]) -> Result<()> {
debug_assert_eq!(values.len(), BLOCK_SIZE);
let first = values[0];
if values.iter().all(|value| *value == first) {
output.write_byte(ALL_VALUES_EQUAL)?;
return output.write_vint(first);
}
let mut combined = 0u64;
let mut packed = Vec::with_capacity(BLOCK_SIZE);
for value in values {
let value = u64::from(*value as u32);
combined |= value;
packed.push(value);
}
let bits_per_value = bits_required(combined);
output.write_byte(bits_per_value as u8)?;
write_packed(output, &packed, bits_per_value)
}
#[derive(Clone, Copy)]
struct SkipPoint {
document: i32,
document_pointer: u64,
position_pointer: u64,
payload_pointer: u64,
position_buffer_upto: usize,
}
struct SkipWriter {
level_count: usize,
has_positions: bool,
has_offsets: bool,
buffers: Vec<Vec<u8>>,
last_document: Vec<i32>,
last_document_pointer: Vec<u64>,
last_position_pointer: Vec<u64>,
last_payload_pointer: Vec<u64>,
}
impl SkipWriter {
fn new(document_count: i32) -> Self {
let level_count = Self::level_count(document_count);
Self {
level_count,
has_positions: false,
has_offsets: false,
buffers: vec![Vec::new(); level_count],
last_document: vec![0; level_count],
last_document_pointer: vec![0; level_count],
last_position_pointer: vec![0; level_count],
last_payload_pointer: vec![0; level_count],
}
}
fn level_count(document_count: i32) -> usize {
if document_count <= BLOCK_SIZE as i32 {
return 1;
}
let mut remaining = document_count as usize / BLOCK_SIZE;
let mut levels = 1;
while remaining >= SKIP_MULTIPLIER {
remaining /= SKIP_MULTIPLIER;
levels += 1;
}
levels.min(MAX_SKIP_LEVELS)
}
fn set_field(&mut self, has_positions: bool, has_offsets: bool) {
self.has_positions = has_positions;
self.has_offsets = has_offsets;
}
fn reset(&mut self, document_pointer: u64, position_pointer: u64, payload_pointer: u64) {
for buffer in &mut self.buffers {
buffer.clear();
}
self.last_document.fill(0);
self.last_document_pointer.fill(document_pointer);
if self.has_positions {
self.last_position_pointer.fill(position_pointer);
if self.has_offsets {
self.last_payload_pointer.fill(payload_pointer);
}
}
}
fn buffer_skip(&mut self, document_count: i32, point: &SkipPoint) -> Result<()> {
debug_assert_eq!(
document_count as usize % BLOCK_SIZE,
0,
"a skip point falls on a block boundary"
);
let mut levels = 1usize;
let mut remaining = document_count as usize / BLOCK_SIZE;
while remaining.is_multiple_of(SKIP_MULTIPLIER) && levels < self.level_count {
levels += 1;
remaining /= SKIP_MULTIPLIER;
}
let mut child_pointer = 0u64;
for level in 0..levels {
self.write_skip_data(level, point)?;
let new_child_pointer = self.buffers[level].len() as u64;
if level != 0 {
let mut output = CodecOutput::new(&mut self.buffers[level]);
output.write_vlong(child_pointer as i64)?;
}
child_pointer = new_child_pointer;
}
Ok(())
}
fn write_skip_data(&mut self, level: usize, point: &SkipPoint) -> Result<()> {
let document_delta = point.document - self.last_document[level];
self.last_document[level] = point.document;
let document_pointer_delta = point.document_pointer - self.last_document_pointer[level];
self.last_document_pointer[level] = point.document_pointer;
let position_pointer_delta = point.position_pointer - self.last_position_pointer[level];
let payload_pointer_delta = point.payload_pointer - self.last_payload_pointer[level];
if self.has_positions {
self.last_position_pointer[level] = point.position_pointer;
if self.has_offsets {
self.last_payload_pointer[level] = point.payload_pointer;
}
}
let has_positions = self.has_positions;
let has_offsets = self.has_offsets;
let position_buffer_upto = point.position_buffer_upto as i32;
let mut output = CodecOutput::new(&mut self.buffers[level]);
output.write_vint(document_delta)?;
output.write_vint(document_pointer_delta as i32)?;
if has_positions {
output.write_vint(position_pointer_delta as i32)?;
output.write_vint(position_buffer_upto)?;
if has_offsets {
output.write_vint(payload_pointer_delta as i32)?;
}
}
Ok(())
}
fn write_skip<Sink: Write>(&self, output: &mut CodecOutput<Sink>) -> Result<u64> {
let pointer = output.position();
for level in (1..self.level_count).rev() {
let length = self.buffers[level].len();
if length > 0 {
output.write_vlong(length as i64)?;
output.write_bytes(&self.buffers[level])?;
}
}
output.write_bytes(&self.buffers[0])?;
Ok(pointer)
}
}
pub struct PostingsWriter<Sink: Write> {
document_output: CodecOutput<Sink>,
position_output: Option<CodecOutput<Sink>>,
payload_output: Option<CodecOutput<Sink>>,
options: IndexOptions,
document_deltas: Vec<i32>,
frequencies: Vec<i32>,
document_buffer_upto: usize,
position_deltas: Vec<i32>,
offset_start_deltas: Vec<i32>,
offset_lengths: Vec<i32>,
position_buffer_upto: usize,
last_document: i32,
last_position: i32,
last_start_offset: i32,
document_count: i32,
frequency_total: i64,
position_count: i64,
last_block_document: i32,
last_block_position_pointer: u64,
last_block_payload_pointer: u64,
last_block_position_buffer_upto: usize,
document_start: u64,
position_start: u64,
payload_start: u64,
skip: SkipWriter,
}
impl<Sink: Write> PostingsWriter<Sink> {
pub fn new(
shape: SegmentShape,
document: Sink,
position: Option<Sink>,
payload: Option<Sink>,
) -> Result<Self> {
if shape.any_field_has_positions != position.is_some() {
return Err(Error::InvalidFormat {
details: format!(
"a segment with positions = {} needs a .pos sink and only then; \
one was {}given",
shape.any_field_has_positions,
if position.is_some() { "" } else { "not " }
),
});
}
if shape.any_field_has_offsets != payload.is_some() {
return Err(Error::InvalidFormat {
details: format!(
"a segment with offsets = {} needs a .pay sink and only then; \
one was {}given",
shape.any_field_has_offsets,
if payload.is_some() { "" } else { "not " }
),
});
}
if shape.any_field_has_offsets && !shape.any_field_has_positions {
return Err(Error::InvalidFormat {
details: "offsets without positions is not an index option Lucene has; \
.pay is created inside the branch that creates .pos"
.to_owned(),
});
}
let mut document_output = CodecOutput::new(document);
document_output.write_header(DOC_CODEC, VERSION_CURRENT)?;
write_format_table(&mut document_output)?;
let position_output = match position {
Some(sink) => {
let mut output = CodecOutput::new(sink);
output.write_header(POS_CODEC, VERSION_CURRENT)?;
Some(output)
}
None => None,
};
let payload_output = match payload {
Some(sink) => {
let mut output = CodecOutput::new(sink);
output.write_header(PAY_CODEC, VERSION_CURRENT)?;
Some(output)
}
None => None,
};
Ok(Self {
document_output,
position_output,
payload_output,
options: IndexOptions::Documents,
document_deltas: vec![0; BLOCK_SIZE],
frequencies: vec![0; BLOCK_SIZE],
document_buffer_upto: 0,
position_deltas: vec![0; BLOCK_SIZE],
offset_start_deltas: vec![0; BLOCK_SIZE],
offset_lengths: vec![0; BLOCK_SIZE],
position_buffer_upto: 0,
last_document: 0,
last_position: 0,
last_start_offset: 0,
document_count: 0,
frequency_total: 0,
position_count: 0,
last_block_document: -1,
last_block_position_pointer: 0,
last_block_payload_pointer: 0,
last_block_position_buffer_upto: 0,
document_start: 0,
position_start: 0,
payload_start: 0,
skip: SkipWriter::new(shape.document_count),
})
}
pub fn set_field(&mut self, options: IndexOptions) -> Result<LongCount> {
if options.has_positions() && self.position_output.is_none() {
return Err(Error::InvalidFormat {
details: "a field with positions in a segment whose field infos said it had \
none: there is no .pos to write into"
.to_owned(),
});
}
if options.has_offsets() && self.payload_output.is_none() {
return Err(Error::InvalidFormat {
details: "a field with offsets in a segment whose field infos said it had \
none: there is no .pay to write into"
.to_owned(),
});
}
self.options = options;
self.skip
.set_field(options.has_positions(), options.has_offsets());
Ok(LongCount(if options.has_positions() {
if options.has_offsets() { 3 } else { 2 }
} else {
1
}))
}
pub fn start_term(&mut self) {
self.document_start = self.document_output.position();
if let Some(output) = self.position_output.as_ref() {
self.position_start = output.position();
}
if let Some(output) = self.payload_output.as_ref() {
self.payload_start = output.position();
}
self.last_document = 0;
self.last_block_document = -1;
self.skip
.reset(self.document_start, self.position_start, self.payload_start);
}
pub fn start_document(&mut self, document: i32, frequency: i32) -> Result<()> {
if self.last_block_document != -1 && self.document_buffer_upto == 0 {
let point = SkipPoint {
document: self.last_block_document,
document_pointer: self.document_output.position(),
position_pointer: self.last_block_position_pointer,
payload_pointer: self.last_block_payload_pointer,
position_buffer_upto: self.last_block_position_buffer_upto,
};
self.skip.buffer_skip(self.document_count, &point)?;
}
let delta = document - self.last_document;
if document < 0 || (self.document_count > 0 && delta <= 0) {
return Err(Error::InvalidFormat {
details: format!(
"postings are written in increasing document order and {document} does \
not follow {}",
self.last_document
),
});
}
if self.options.has_frequencies() && frequency < 1 {
return Err(Error::InvalidFormat {
details: format!(
"a document in the postings of a term occurs at least once, so a \
frequency of {frequency} cannot be encoded"
),
});
}
self.document_deltas[self.document_buffer_upto] = delta;
if self.options.has_frequencies() {
self.frequencies[self.document_buffer_upto] = frequency;
self.frequency_total += i64::from(frequency);
}
self.document_buffer_upto += 1;
self.document_count += 1;
if self.document_buffer_upto == BLOCK_SIZE {
write_block(&mut self.document_output, &self.document_deltas)?;
if self.options.has_frequencies() {
write_block(&mut self.document_output, &self.frequencies)?;
}
}
self.last_document = document;
self.last_position = 0;
self.last_start_offset = 0;
Ok(())
}
pub fn add_position(
&mut self,
position: i32,
start_offset: i32,
end_offset: i32,
) -> Result<()> {
if !self.options.has_positions() {
return Err(Error::InvalidFormat {
details: "a position in a field whose index options have none".to_owned(),
});
}
let position_delta = position - self.last_position;
if position_delta < 0 {
return Err(Error::InvalidFormat {
details: format!(
"positions within a document increase, and {position} does not follow {}",
self.last_position
),
});
}
self.position_deltas[self.position_buffer_upto] = position_delta;
if self.options.has_offsets() {
if start_offset < self.last_start_offset || end_offset < start_offset {
return Err(Error::InvalidFormat {
details: format!(
"offsets increase and enclose their term: {start_offset}..{end_offset} \
does not follow a start of {}",
self.last_start_offset
),
});
}
self.offset_start_deltas[self.position_buffer_upto] =
start_offset - self.last_start_offset;
self.offset_lengths[self.position_buffer_upto] = end_offset - start_offset;
self.last_start_offset = start_offset;
}
self.position_buffer_upto += 1;
self.position_count += 1;
self.last_position = position;
if self.position_buffer_upto == BLOCK_SIZE {
let position_output =
self.position_output
.as_mut()
.ok_or_else(|| Error::InvalidFormat {
details: "a field with positions but no .pos sink".to_owned(),
})?;
write_block(position_output, &self.position_deltas)?;
if self.options.has_offsets() {
let payload_output =
self.payload_output
.as_mut()
.ok_or_else(|| Error::InvalidFormat {
details: "a field with offsets but no .pay sink".to_owned(),
})?;
write_block(payload_output, &self.offset_start_deltas)?;
write_block(payload_output, &self.offset_lengths)?;
}
self.position_buffer_upto = 0;
}
Ok(())
}
pub fn finish_document(&mut self) {
if self.document_buffer_upto == BLOCK_SIZE {
self.last_block_document = self.last_document;
if let Some(output) = self.position_output.as_ref() {
if let Some(payload) = self.payload_output.as_ref() {
self.last_block_payload_pointer = payload.position();
}
self.last_block_position_pointer = output.position();
self.last_block_position_buffer_upto = self.position_buffer_upto;
}
self.document_buffer_upto = 0;
}
}
pub fn finish_term(&mut self) -> Result<TermMetadata> {
if self.document_count == 0 {
return Err(Error::InvalidFormat {
details: "a term with no documents has no postings to finish".to_owned(),
});
}
let singleton_document = if self.document_count == 1 {
Some(self.document_deltas[0])
} else {
self.write_document_tail()?;
None
};
let total_term_frequency = if self.options.has_frequencies() {
self.frequency_total
} else {
-1
};
let last_position_block_offset = self.finish_positions(total_term_frequency)?;
let skip_offset = if self.document_count > BLOCK_SIZE as i32 {
Some(self.skip.write_skip(&mut self.document_output)? - self.document_start)
} else {
None
};
let metadata = TermMetadata {
document_start: self.document_start,
position_start: self.position_start,
payload_start: self.payload_start,
singleton_document,
last_position_block_offset,
skip_offset,
document_frequency: self.document_count,
total_term_frequency,
};
self.document_buffer_upto = 0;
self.position_buffer_upto = 0;
self.last_document = 0;
self.document_count = 0;
self.frequency_total = 0;
self.position_count = 0;
Ok(metadata)
}
fn write_document_tail(&mut self) -> Result<()> {
for index in 0..self.document_buffer_upto {
let delta = self.document_deltas[index];
if !self.options.has_frequencies() {
self.document_output.write_vint(delta)?;
} else if self.frequencies[index] == 1 {
self.document_output.write_vint((delta << 1) | 1)?;
} else {
self.document_output.write_vint(delta << 1)?;
self.document_output.write_vint(self.frequencies[index])?;
}
}
Ok(())
}
fn finish_positions(&mut self, total_term_frequency: i64) -> Result<Option<u64>> {
if !self.options.has_positions() {
return Ok(None);
}
if self.position_count != total_term_frequency {
return Err(Error::InvalidFormat {
details: format!(
"a field with positions records one position per occurrence, so \
{} positions cannot belong to a term of total frequency \
{total_term_frequency}",
self.position_count
),
});
}
let has_offsets = self.options.has_offsets();
let position_start = self.position_start;
let buffer_upto = self.position_buffer_upto;
let position_deltas = &self.position_deltas;
let offset_start_deltas = &self.offset_start_deltas;
let offset_lengths = &self.offset_lengths;
let output = self
.position_output
.as_mut()
.ok_or_else(|| Error::InvalidFormat {
details: "a field with positions but no .pos sink".to_owned(),
})?;
let last_position_block_offset = if total_term_frequency > BLOCK_SIZE as i64 {
Some(output.position() - position_start)
} else {
None
};
let mut last_offset_length = -1i32;
for index in 0..buffer_upto {
output.write_vint(position_deltas[index])?;
if has_offsets {
let start_delta = offset_start_deltas[index];
let length = offset_lengths[index];
if length == last_offset_length {
output.write_vint(start_delta << 1)?;
} else {
output.write_vint((start_delta << 1) | 1)?;
output.write_vint(length)?;
last_offset_length = length;
}
}
}
Ok(last_position_block_offset)
}
pub fn finish(self) -> PostingsFiles<Sink> {
PostingsFiles {
document: self.document_output.into_inner(),
position: self.position_output.map(CodecOutput::into_inner),
payload: self.payload_output.map(CodecOutput::into_inner),
}
}
}