///|
/// Incremental writer that publishes one new immutable Segment per commit.
pub struct IndexWriter[D] {
directory : D
mut pending_delete_terms : Array[Term]
mut closed : Bool
}
///|
pub fn[D] IndexWriter::new(directory : D) -> IndexWriter[D] {
{ directory, pending_delete_terms: [], closed: false }
}
///|
fn[D : Directory] load_manifest_or_empty(
directory : D,
) -> IndexManifest raise PersistenceError {
if directory.exists(manifest_file_name) {
decode_manifest(directory.read(manifest_file_name))
} else {
empty_manifest()
}
}
///|
fn[D : Directory] load_segment(
directory : D,
meta : SegmentMeta,
) -> Segment raise PersistenceError {
decode_segment(directory.read(meta.file_name))
}
///|
fn[D : Directory] load_snapshot_segment(
directory : D,
meta : SegmentMeta,
) -> SnapshotSegment raise PersistenceError {
let segment = load_segment(directory, meta)
for doc_id in meta.deleted_docs {
guard doc_id.value < segment.doc_count() else {
raise PersistenceError::InvalidFormat(
"manifest tombstone exceeds segment document count",
)
}
}
SnapshotSegment::new(segment, meta.deleted_docs)
}
///|
fn[D : Directory] apply_delete_terms(
directory : D,
manifest : IndexManifest,
terms : Array[Term],
) -> ReadOnlyArray[SegmentMeta] raise PersistenceError {
let metas : Array[SegmentMeta] = []
for meta in manifest.segments {
let segment = load_segment(directory, meta)
let deleted = Array::make(segment.doc_count(), false)
for doc_id in meta.deleted_docs {
guard doc_id.value < segment.doc_count() else {
raise PersistenceError::InvalidFormat(
"manifest tombstone exceeds segment document count",
)
}
deleted[doc_id.value] = true
}
for term in terms {
for posting in segment.postings_for(term) {
deleted[posting.doc_id.value] = true
}
}
let deleted_docs : Array[DocId] = []
for doc_index in 0.. IndexManifest {
{
generation: manifest.generation + 1,
next_segment_id: manifest.next_segment_id,
segments,
}
}
///|
fn[D : Directory] validate_incremental_schema(
directory : D,
manifest : IndexManifest,
segment : Segment,
) -> Unit raise PersistenceError {
if manifest.segments.length() == 0 {
return
}
let existing = load_segment(directory, manifest.segments[0])
guard schemas_equal(existing.schema(), segment.schema()) else {
raise PersistenceError::InvalidFormat(
"incremental segment schema does not match the index",
)
}
}
///|
/// Writes the new Segment first and publishes it by replacing the manifest
/// last. Readers that already opened a manifest keep their previous snapshot.
pub fn[D : Directory] IndexWriter::commit(
self : IndexWriter[D],
segment : Segment,
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
let current = load_manifest_or_empty(self.directory)
validate_incremental_schema(self.directory, current, segment)
let current_with_deletes : IndexManifest = if self.pending_delete_terms.length() ==
0 {
current
} else {
{
generation: current.generation,
next_segment_id: current.next_segment_id,
segments: apply_delete_terms(
self.directory,
current,
self.pending_delete_terms,
),
}
}
let next = append_manifest_segment(current_with_deletes)
let new_meta = next.segments[next.segments.length() - 1]
self.directory.write(new_meta.file_name, encode_segment(segment))
self.directory.write(manifest_file_name, encode_manifest(next))
self.pending_delete_terms = []
}
///|
/// Queues an exact field-qualified Term deletion for the next writer commit.
pub fn[D] IndexWriter::delete_term(
self : IndexWriter[D],
term : Term,
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
if self.pending_delete_terms.search_by(queued => queued == term) is None {
self.pending_delete_terms.push(term)
}
}
///|
/// Applies queued Term deletions to currently published Segments and publishes
/// a new manifest generation. With no queued terms this is a no-op.
pub fn[D : Directory] IndexWriter::commit_deletes(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
if self.pending_delete_terms.length() == 0 {
return
}
guard self.directory.exists(manifest_file_name) else {
raise PersistenceError::NotFound(manifest_file_name)
}
let current = decode_manifest(self.directory.read(manifest_file_name))
let next = manifest_with_segments(
current,
apply_delete_terms(self.directory, current, self.pending_delete_terms),
)
self.directory.write(manifest_file_name, encode_manifest(next))
self.pending_delete_terms = []
}
///|
/// Merges every live document into one new immutable Segment and publishes it
/// as a new generation. Existing files remain available to older Readers.
pub fn[D : Directory] IndexWriter::merge(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
guard self.directory.exists(manifest_file_name) else {
raise PersistenceError::NotFound(manifest_file_name)
}
let current = decode_manifest(self.directory.read(manifest_file_name))
let metas = if self.pending_delete_terms.length() == 0 {
current.segments
} else {
apply_delete_terms(self.directory, current, self.pending_delete_terms)
}
let snapshots : Array[SnapshotSegment] = []
let mut schema : Schema? = None
for index in 0.. Unit {
self.closed = true
}
///|
/// Immutable snapshot of the Segment set named by one manifest generation.
pub struct IndexReader[D] {
directory : D
mut segments : ReadOnlyArray[SnapshotSegment]
schema_snapshot : Schema?
generation_snapshot : Int
mut closed : Bool
}
///|
pub fn[D : Directory] IndexReader::open(
directory : D,
) -> IndexReader[D] raise PersistenceError {
guard directory.exists(manifest_file_name) else {
raise PersistenceError::NotFound(manifest_file_name)
}
let manifest = decode_manifest(directory.read(manifest_file_name))
let segments : Array[SnapshotSegment] = []
let mut schema : Schema? = None
for index in 0.. Searcher raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
Searcher::from_segments(self.segments)
}
///|
pub fn[D] IndexReader::schema(
self : IndexReader[D],
) -> Schema? raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
self.schema_snapshot
}
///|
pub fn[D] IndexReader::segment_count(
self : IndexReader[D],
) -> Int raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
self.segments.length()
}
///|
pub fn[D] IndexReader::generation(
self : IndexReader[D],
) -> Int raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
self.generation_snapshot
}
///|
pub fn[D] IndexReader::close(self : IndexReader[D]) -> Unit {
self.segments = []
self.closed = true
}