///|
/// 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
}