///|
/// Incremental writer that publishes one new immutable Segment per commit.
pub struct IndexWriter[D] {
directory : D
mut pending_delete_terms : Array[Term]
mut prepared_manifest_bytes : Bytes?
mut prepared_segment_name : String?
mut lock_token : Int64?
mut closed : Bool
}
///|
pub fn[D] IndexWriter::new(directory : D) -> IndexWriter[D] {
{
directory,
pending_delete_terms: [],
prepared_manifest_bytes: None,
prepared_segment_name: None,
lock_token: None,
closed: false,
}
}
///|
let writer_lock_file_name : String = "write.lock"
///|
let prepared_manifest_file_name : String = ".pending-meta.msi"
///|
fn manifest_references_file(manifest : IndexManifest, name : String) -> Bool {
manifest.segments.search_by(segment => segment.file_name == name) is Some(_)
}
///|
fn[D : DirectoryV2] IndexWriter::garbage_collect_unlocked(
self : IndexWriter[D],
) -> Int raise PersistenceError {
let manifest = load_manifest_or_empty(self.directory)
let mut removed = 0
for name in self.directory.list() {
let stale_temporary = name.has_suffix(".tmp") ||
name == prepared_manifest_file_name
let orphan_segment = name.has_prefix("segment-") &&
name.has_suffix(".msi") &&
!manifest_references_file(manifest, name)
if stale_temporary || orphan_segment {
self.directory.remove(name)
removed += 1
}
}
removed
}
///|
fn[D : DirectoryV2] IndexWriter::ensure_lock(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
if self.lock_token is Some(_) {
return
}
match self.directory.try_acquire_lock(writer_lock_file_name) {
Some(token) => self.lock_token = Some(token)
None => raise PersistenceError::LockBusy(writer_lock_file_name)
}
ignore(self.garbage_collect_unlocked())
}
///|
fn[D : DirectoryV2] IndexWriter::prepare_publication(
self : IndexWriter[D],
segment : Segment?,
segment_name : String?,
manifest : IndexManifest,
) -> Unit raise PersistenceError {
guard self.prepared_manifest_bytes is None else {
raise PersistenceError::AlreadyPrepared
}
match (segment, segment_name) {
(Some(segment), Some(name)) => {
self.directory.write_atomic(name, encode_segment(segment))
self.directory.sync(name)
self.prepared_segment_name = Some(name)
}
(None, None) => self.prepared_segment_name = None
_ => abort("prepared segment and name must be supplied together")
}
let manifest_bytes = encode_manifest(manifest)
self.directory.write_atomic(prepared_manifest_file_name, manifest_bytes)
self.directory.sync(prepared_manifest_file_name)
self.prepared_manifest_bytes = Some(manifest_bytes)
}
///|
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 and syncs immutable data plus a private prepared manifest without
/// changing the reader-visible manifest generation.
pub fn[D : DirectoryV2] IndexWriter::prepare_commit(
self : IndexWriter[D],
segment : Segment,
) -> Unit raise PersistenceError {
self.ensure_lock()
guard self.prepared_manifest_bytes is None else {
raise PersistenceError::AlreadyPrepared
}
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.prepare_publication(Some(segment), Some(new_meta.file_name), next)
}
///|
/// Atomically publishes the previously prepared manifest.
pub fn[D : DirectoryV2] IndexWriter::commit_prepared(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
self.ensure_lock()
let manifest_bytes = match self.prepared_manifest_bytes {
Some(bytes) => bytes
None => raise PersistenceError::NotPrepared
}
self.directory.write_atomic(manifest_file_name, manifest_bytes)
self.directory.sync(manifest_file_name)
self.directory.remove(prepared_manifest_file_name)
self.prepared_manifest_bytes = None
self.prepared_segment_name = None
self.pending_delete_terms = []
}
///|
/// Convenience one-shot commit built on prepare_commit + commit_prepared.
pub fn[D : DirectoryV2] IndexWriter::commit(
self : IndexWriter[D],
segment : Segment,
) -> Unit raise PersistenceError {
self.prepare_commit(segment)
self.commit_prepared()
}
///|
/// Discards a prepared publication and queued deletes. Already published
/// generations are never changed.
pub fn[D : DirectoryV2] IndexWriter::rollback(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
match self.prepared_segment_name {
Some(name) => {
let current = load_manifest_or_empty(self.directory)
if !manifest_references_file(current, name) {
self.directory.remove(name)
}
}
None => ()
}
self.directory.remove(prepared_manifest_file_name)
self.prepared_manifest_bytes = None
self.prepared_segment_name = None
self.pending_delete_terms = []
}
///|
pub fn[D] IndexWriter::is_prepared(self : IndexWriter[D]) -> Bool {
self.prepared_manifest_bytes is Some(_)
}
///|
/// Queues an exact field-qualified Term deletion for the next writer commit.
pub fn[D : DirectoryV2] IndexWriter::delete_term(
self : IndexWriter[D],
term : Term,
) -> Unit raise PersistenceError {
self.ensure_lock()
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 : DirectoryV2] IndexWriter::prepare_deletes(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
self.ensure_lock()
guard self.prepared_manifest_bytes is None else {
raise PersistenceError::AlreadyPrepared
}
if self.pending_delete_terms.length() == 0 {
return
}
guard Directory::exists(self.directory, manifest_file_name) else {
raise PersistenceError::NotFound(manifest_file_name)
}
let current = decode_manifest(
Directory::read(self.directory, manifest_file_name),
)
let next = manifest_with_segments(
current,
apply_delete_terms(self.directory, current, self.pending_delete_terms),
)
self.prepare_publication(None, None, next)
}
///|
pub fn[D : DirectoryV2] IndexWriter::commit_deletes(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
self.prepare_deletes()
if self.is_prepared() {
self.commit_prepared()
}
}
///|
/// 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 : DirectoryV2] IndexWriter::merge(
self : IndexWriter[D],
) -> Unit raise PersistenceError {
self.ensure_lock()
guard self.prepared_manifest_bytes is None else {
raise PersistenceError::AlreadyPrepared
}
guard Directory::exists(self.directory, manifest_file_name) else {
raise PersistenceError::NotFound(manifest_file_name)
}
let current = decode_manifest(
Directory::read(self.directory, 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 raise PersistenceError {
self.delete_term(term)
self.commit(replacement)
}
///|
/// Removes orphan Segments and stale temporary files not referenced by the
/// current manifest. Open readers are safe because they own decoded snapshots.
pub fn[D : DirectoryV2] IndexWriter::garbage_collect(
self : IndexWriter[D],
) -> Int raise PersistenceError {
self.ensure_lock()
guard self.prepared_manifest_bytes is None else {
raise PersistenceError::AlreadyPrepared
}
self.garbage_collect_unlocked()
}
///|
pub fn[D : DirectoryV2] IndexWriter::close(self : IndexWriter[D]) -> Unit {
if self.closed {
return
}
if self.prepared_manifest_bytes is Some(_) ||
self.prepared_segment_name is Some(_) {
self.rollback() catch {
_ => ()
}
}
match self.lock_token {
Some(token) => {
self.directory.release_lock(writer_lock_file_name, token) catch {
_ => ()
}
self.lock_token = None
}
None => ()
}
self.closed = true
}
///|
/// First deterministic merge policy: merge all segments once the configured
/// segment-count threshold is reached.
pub(all) struct MergePolicy {
max_segment_count : Int
} derive(Eq, @debug.Debug)
///|
pub fn MergePolicy::new(max_segment_count : Int) -> MergePolicy {
guard max_segment_count >= 2 else {
abort("merge threshold must be at least two segments")
}
{ max_segment_count, }
}
///|
pub fn MergePolicy::should_merge(
self : MergePolicy,
segment_count : Int,
) -> Bool {
segment_count >= self.max_segment_count
}
///|
/// Synchronous first scheduler with a hard per-merge input-byte budget. A
/// non-positive budget disables throttling.
pub(all) struct MergeScheduler {
policy : MergePolicy
max_input_bytes : Int
max_concurrent_merges : Int
mut in_flight_merges : Int
mut in_flight_bytes : Int
} derive(Eq, @debug.Debug)
///|
pub fn MergeScheduler::new(
policy : MergePolicy,
max_input_bytes : Int,
) -> MergeScheduler {
{
policy,
max_input_bytes,
max_concurrent_merges: 1,
in_flight_merges: 0,
in_flight_bytes: 0,
}
}
///|
/// Configures merge concurrency and aggregate in-flight bytes. A non-positive
/// byte budget disables the aggregate byte limit.
pub fn MergeScheduler::with_backpressure(
self : MergeScheduler,
max_concurrent_merges : Int,
max_in_flight_bytes : Int,
) -> MergeScheduler {
guard max_concurrent_merges > 0 else {
abort("maximum concurrent merges must be positive")
}
{
policy: self.policy,
max_input_bytes: if max_in_flight_bytes > 0 {
max_in_flight_bytes
} else {
self.max_input_bytes
},
max_concurrent_merges,
in_flight_merges: 0,
in_flight_bytes: 0,
}
}
///|
pub fn MergeScheduler::is_backpressured(
self : MergeScheduler,
next_input_bytes : Int,
) -> Bool {
self.in_flight_merges >= self.max_concurrent_merges ||
(
self.max_input_bytes > 0 &&
self.in_flight_bytes + next_input_bytes > self.max_input_bytes
)
}
///|
pub fn MergeScheduler::in_flight_merges(self : MergeScheduler) -> Int {
self.in_flight_merges
}
///|
pub fn MergeScheduler::in_flight_bytes(self : MergeScheduler) -> Int {
self.in_flight_bytes
}
///|
pub fn[D : DirectoryV2] MergeScheduler::run(
self : MergeScheduler,
writer : IndexWriter[D],
) -> Bool raise PersistenceError {
writer.ensure_lock()
let manifest = load_manifest_or_empty(writer.directory)
if !self.policy.should_merge(manifest.segments.length()) {
return false
}
let mut input_bytes = 0
for meta in manifest.segments {
input_bytes += Directory::read(writer.directory, meta.file_name).length()
}
if self.is_backpressured(input_bytes) {
return false
}
self.in_flight_merges += 1
self.in_flight_bytes += input_bytes
errdefer {
self.in_flight_merges -= 1
self.in_flight_bytes -= input_bytes
}
writer.merge()
self.in_flight_merges -= 1
self.in_flight_bytes -= input_bytes
true
}
///|
/// Immutable snapshot of the Segment set named by one manifest generation.
pub struct IndexReader[D] {
directory : D
mut segments : ReadOnlyArray[SnapshotSegment]
mut schema_snapshot : Schema?
mut generation_snapshot : Int
mut closed : Bool
}
///|
fn[D : Directory] load_reader_snapshot(
directory : D,
) -> (ReadOnlyArray[SnapshotSegment], Schema?, Int) 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.. IndexReader[D] raise PersistenceError {
let (segments, schema, generation) = load_reader_snapshot(directory)
{
directory,
segments,
schema_snapshot: schema,
generation_snapshot: generation,
closed: false,
}
}
///|
/// Opens a reader and immediately verifies analyzer resource fingerprints
/// persisted in its Schema. Use this entry point when fields pin dictionary or
/// synonym revisions.
pub fn[D : Directory] IndexReader::open_with_tokenizers(
directory : D,
tokenizers : TokenizerManager,
) -> IndexReader[D] raise {
let reader = IndexReader::open(directory)
match reader.schema_snapshot {
Some(schema) => tokenizers.validate_schema(schema)
None => ()
}
reader
}
///|
/// Reloads this reader if the current manifest generation changed. Searchers
/// created before reload remain immutable snapshots.
pub fn[D : Directory] IndexReader::reload(
self : IndexReader[D],
) -> Bool raise PersistenceError {
guard !self.closed else { raise PersistenceError::Closed }
let (segments, schema, generation) = load_reader_snapshot(self.directory)
if generation == self.generation_snapshot {
return false
}
self.segments = segments
self.schema_snapshot = schema
self.generation_snapshot = generation
true
}
///|
pub fn[D : Directory] IndexReader::open_if_changed(
self : IndexReader[D],
) -> Bool raise PersistenceError {
self.reload()
}
///|
pub fn[D] IndexReader::searcher(
self : IndexReader[D],
) -> 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
}