Skip to main content

hermes_core/index/
writer.rs

1//! IndexWriter — async document indexing with parallel segment building.
2//!
3//! This module is only compiled with the "native" feature.
4//!
5//! # Architecture
6//!
7//! ```text
8//! add_document() ──try_send──► [shared bounded MPMC] ◄──recv── worker 0
9//!                                                     ◄──recv── worker 1
10//!                                                     ◄──recv── worker N
11//! ```
12//!
13//! - **Shared MPMC queue** (`async_channel`): all workers compete for documents.
14//!   Busy workers (building segments) naturally stop pulling; free workers pick up slack.
15//! - **Zero-copy pipeline**: `Document` is moved (never cloned) through every stage:
16//!   `add_document()` → channel → `recv_blocking()` → `SegmentBuilder::add_document()`.
17//! - `add_document` returns `QueueFull` when the queue is at capacity.
18//! - **Workers are OS threads**: CPU-intensive work (tokenization, posting list building)
19//!   runs on dedicated threads, never blocking the tokio async runtime.
20//!   Async I/O (segment file writes) is bridged via `Handle::block_on()`.
21//! - **Fixed per-worker memory budget**: `max_indexing_memory_bytes / num_workers`.
22//! - **Two-phase commit**:
23//!   1. `prepare_commit()` — closes queue, workers flush builders to disk.
24//!      Returns a `PreparedCommit` guard. No new documents accepted until resolved.
25//!   2. `PreparedCommit::commit()` — registers segments in metadata, resumes workers.
26//!   3. `PreparedCommit::abort()` — discards prepared segments, resumes workers.
27//!   4. `commit()` — convenience: `prepare_commit().await?.commit().await`.
28//!
29//! Since `prepare_commit`/`commit` take `&mut self`, Rust’s borrow checker
30//! guarantees no concurrent `add_document` calls during the commit window.
31
32use std::sync::Arc;
33use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
34
35use futures::FutureExt;
36use rustc_hash::FxHashMap;
37
38use crate::directories::DirectoryWriter;
39use crate::dsl::{Document, Field, Schema};
40use crate::error::{Error, Result};
41use crate::segment::{SegmentBuilder, SegmentBuilderConfig, SegmentId};
42use crate::tokenizer::BoxedTokenizer;
43
44use super::IndexConfig;
45
46/// Total pipeline capacity (in documents).
47const PIPELINE_MAX_SIZE_IN_DOCS: usize = 10_000;
48
49/// File name of the advisory single-writer lock inside the index directory.
50pub const WRITER_LOCK_FILENAME: &str = ".hermes_writer.lock";
51
52/// Advisory single-writer lock state.
53///
54/// Two independent writers on one index directory silently destroy each
55/// other's data: the orphan sweep at writer open deletes the other process's
56/// unpublished segment files, and metadata saves are last-writer-wins. For
57/// directories rooted on a local filesystem the writer therefore holds an OS
58/// advisory lock for its whole lifetime; the kernel releases it automatically
59/// when the process dies.
60enum WriterLock {
61    /// Lock acquired. Closing the file (writer drop) releases it.
62    Held { _file: std::fs::File },
63    /// The directory has no lockable local filesystem root (e.g. RAM or
64    /// remote directories) — cross-process locking is not applicable.
65    NotApplicable,
66    /// Another writer holds the lock. Every mutating operation fails loudly
67    /// with this message instead of silently double-writing.
68    Unavailable { reason: String },
69}
70
71/// Local filesystem root of the index directory, when the directory type
72/// exposes one.
73fn writer_lock_root<D: DirectoryWriter + 'static>(directory: &D) -> Option<std::path::PathBuf> {
74    let any: &dyn std::any::Any = directory;
75    if let Some(mmap) = any.downcast_ref::<crate::directories::MmapDirectory>() {
76        return Some(mmap.root().to_path_buf());
77    }
78    // FsDirectory does not expose its root path, so the single-writer lock
79    // cannot be enforced for it yet. Say so loudly instead of silently
80    // skipping protection for a filesystem-backed writer.
81    if any
82        .downcast_ref::<crate::directories::FsDirectory>()
83        .is_some()
84    {
85        log::warn!(
86            "[writer_lock] FsDirectory exposes no root path; single-writer locking \
87             is not enforced for this writer — do not open a second writer for the \
88             same index directory"
89        );
90    }
91    None
92}
93
94/// Try to take the exclusive single-writer lock for `directory`.
95///
96/// Returns `WriterLock::Unavailable` (not `Err`) on conflict so infallible
97/// constructors can defer the failure to their first mutating operation.
98fn try_acquire_writer_lock<D: DirectoryWriter + 'static>(directory: &D) -> Result<WriterLock> {
99    let Some(root) = writer_lock_root(directory) else {
100        return Ok(WriterLock::NotApplicable);
101    };
102    std::fs::create_dir_all(&root)?;
103    let lock_path = root.join(WRITER_LOCK_FILENAME);
104    let file = std::fs::OpenOptions::new()
105        .create(true)
106        .truncate(false)
107        .write(true)
108        .open(&lock_path)?;
109    match file.try_lock() {
110        Ok(()) => Ok(WriterLock::Held { _file: file }),
111        Err(std::fs::TryLockError::WouldBlock) => Ok(WriterLock::Unavailable {
112            reason: format!(
113                "another IndexWriter already holds the single-writer lock for this \
114                 index ({}); Hermes supports one writer per index directory — stop \
115                 the other writer (e.g. a running hermes-server or hermes-tool) \
116                 before opening this one",
117                lock_path.display()
118            ),
119        }),
120        Err(std::fs::TryLockError::Error(error)) => Err(Error::Io(error)),
121    }
122}
123
124/// Async IndexWriter for adding documents and committing segments.
125///
126/// **Backpressure:** `add_document()` is sync and O(1). It returns
127/// `Error::QueueFull` when the shared queue is full and
128/// `Error::CommitInProgress` while a generation is publishing or awaiting
129/// retry; callers must back off.
130///
131/// **Two-phase commit:**
132/// - `prepare_commit()` → `PreparedCommit::commit()` or `PreparedCommit::abort()`
133/// - `commit()` is a convenience that does both phases.
134/// - Between prepare and commit, the caller can do external work (WAL, sync, etc.)
135///   knowing that abort is possible if something fails.
136/// - Dropping `PreparedCommit` without calling commit/abort auto-aborts.
137pub struct IndexWriter<D: DirectoryWriter + 'static> {
138    pub(super) directory: Arc<D>,
139    pub(super) schema: Arc<Schema>,
140    pub(super) config: IndexConfig,
141    /// MPMC sender, replaced under a brief lock on each commit cycle (workers
142    /// get the corresponding new receiver via resume).
143    doc_sender: Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
144    /// Worker OS thread handles — long-lived, survive across commits.
145    workers: Vec<std::thread::JoinHandle<()>>,
146    /// Shared worker state (immutable config + mutable segment output + sync)
147    worker_state: Arc<WorkerState<D>>,
148    /// Segment manager — owns metadata.json, handles segments and background merging
149    pub(super) segment_manager: Arc<crate::merge::SegmentManager<D>>,
150    /// Segments flushed to disk but not yet registered in metadata. Each item
151    /// owns an active-operation guard, so orphan sweeping cannot delete it.
152    flushed_segments: Arc<parking_lot::Mutex<Vec<PreparedSegment<D>>>>,
153    /// Primary key dedup index (None if schema has no primary field)
154    primary_key_index: Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
155    /// Tracks the owned finalizer spawned by `PreparedCommit::commit`. The
156    /// requesting future may disappear, but a second commit generation must
157    /// not start until this one has made publication and worker state agree.
158    commit_finalization: Arc<CommitFinalizationState>,
159    /// True while a failed post-commit PK refresh has left the uncommitted
160    /// reservations as the ONLY record of already-committed keys (fail-closed,
161    /// see `finalize_prepared_commit`). While set, abort paths must NOT clear
162    /// the reservations or duplicate primary keys could be admitted.
163    pk_reservations_retained: Arc<AtomicBool>,
164    /// Advisory single-writer lock, held for the writer's lifetime.
165    /// `Unavailable` is retryable: the conflicting holder may exit at any
166    /// time (the kernel then releases its lock), so `ensure_writer_lock`
167    /// re-attempts acquisition instead of caching the conflict forever.
168    writer_lock: parking_lot::RwLock<WriterLock>,
169}
170
171#[derive(Default)]
172struct CommitFinalizationState {
173    in_progress: AtomicBool,
174    idle: tokio::sync::Notify,
175}
176
177impl CommitFinalizationState {
178    fn begin(&self) -> bool {
179        self.in_progress
180            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
181            .is_ok()
182    }
183
184    fn finish(&self) {
185        self.in_progress.store(false, Ordering::Release);
186        self.idle.notify_waiters();
187    }
188
189    async fn wait_until_idle(&self) {
190        while self.in_progress.load(Ordering::Acquire) {
191            let notified = self.idle.notified();
192            if !self.in_progress.load(Ordering::Acquire) {
193                break;
194            }
195            notified.await;
196        }
197    }
198}
199
200/// Shared state for worker threads.
201struct WorkerState<D: DirectoryWriter + 'static> {
202    directory: Arc<D>,
203    schema: Arc<Schema>,
204    builder_config: SegmentBuilderConfig,
205    tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
206    /// Fixed per-worker memory budget (bytes). When a builder exceeds this, segment is built.
207    memory_budget_per_worker: usize,
208    /// Segment manager — workers read trained structures from its ArcSwap (lock-free).
209    segment_manager: Arc<crate::merge::SegmentManager<D>>,
210    /// Segments built by workers, collected by `prepare_commit()`. Their RAII
211    /// guards protect both in-progress and completed-uncommitted files.
212    built_segments: parking_lot::Mutex<Vec<PreparedSegment<D>>>,
213    /// First failure in the current flush generation. Worker-side indexing is
214    /// asynchronous, so `prepare_commit` is the only sound place to surface
215    /// it to the caller. A failed generation is aborted as a unit; publishing
216    /// only its successful segments would silently lose documents.
217    cycle_error: parking_lot::Mutex<Option<String>>,
218    cycle_failed: AtomicBool,
219
220    // === Worker lifecycle synchronization ===
221    // Workers survive across commits. On prepare_commit the channel is closed;
222    // workers flush their builders, increment flush_count, then wait on
223    // resume_cvar for a new receiver. commit/abort creates a fresh channel
224    // and wakes them.
225    /// Number of workers that have completed their flush.
226    flush_count: AtomicUsize,
227    /// Mutex + condvar for prepare_commit to wait on all workers flushed.
228    flush_mutex: parking_lot::Mutex<()>,
229    flush_cvar: parking_lot::Condvar,
230    /// Holds the new channel receiver after commit/abort. Workers clone from this.
231    resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
232    /// Monotonically increasing epoch, bumped by each resume_workers call.
233    /// Workers compare against their local epoch to avoid re-cloning a stale receiver.
234    resume_epoch: AtomicUsize,
235    /// Condvar for workers to wait for resume (new channel) or shutdown.
236    resume_cvar: parking_lot::Condvar,
237    /// When true, workers should exit permanently (IndexWriter dropped).
238    shutdown: AtomicBool,
239    /// Total number of worker threads.
240    num_workers: usize,
241}
242
243/// A completed indexing segment that has not been published in metadata yet.
244///
245/// `operation` is intentionally data, not a side-channel set update: moving
246/// this value through worker → prepared commit → commit/abort moves lifecycle
247/// ownership with it, and every unwind/drop path releases ownership safely.
248struct PreparedSegment<D: DirectoryWriter + 'static> {
249    id: String,
250    segment_id: SegmentId,
251    num_docs: u32,
252    segment_manager: Arc<crate::merge::SegmentManager<D>>,
253    operation: Option<crate::merge::SegmentOperationGuard>,
254    runtime: tokio::runtime::Handle,
255    needs_vector_upgrade: bool,
256    published: bool,
257}
258
259impl<D: DirectoryWriter + 'static> PreparedSegment<D> {
260    fn metadata_entry(&self) -> (String, u32) {
261        (self.id.clone(), self.num_docs)
262    }
263
264    fn mark_published(&mut self) {
265        self.published = true;
266        // Metadata + SegmentTracker are now the durable lifecycle owners.
267        drop(self.operation.take());
268    }
269}
270
271impl<D: DirectoryWriter + 'static> WorkerState<D> {
272    fn record_cycle_error(&self, error: impl Into<String>) {
273        let mut first_error = self.cycle_error.lock();
274        if first_error.is_none() {
275            *first_error = Some(error.into());
276        }
277        drop(first_error);
278        self.cycle_failed.store(true, Ordering::Release);
279    }
280}
281
282impl<D: DirectoryWriter + 'static> Drop for PreparedSegment<D> {
283    fn drop(&mut self) {
284        if self.published {
285            return;
286        }
287        let Some(operation) = self.operation.take() else {
288            return;
289        };
290        self.segment_manager.schedule_unpublished_segment_cleanup(
291            self.segment_id,
292            operation,
293            self.runtime.clone(),
294        );
295    }
296}
297
298impl<D: DirectoryWriter + 'static> IndexWriter<D> {
299    /// Create a new index in the directory
300    pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
301        Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
302    }
303
304    /// Create a new index with custom builder config
305    pub async fn create_with_config(
306        directory: D,
307        schema: Schema,
308        config: IndexConfig,
309        builder_config: SegmentBuilderConfig,
310    ) -> Result<Self> {
311        let directory = Arc::new(directory);
312        let schema = Arc::new(schema);
313        // Directory-layer metrics (cold writes, lazy reads) carry the index label
314        directory.set_index_label(schema.index_label());
315
316        // Refuse a second writer before touching any index state.
317        let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
318        if let WriterLock::Unavailable { reason } = &writer_lock {
319            return Err(Error::Internal(reason.clone()));
320        }
321        // Refuse to clobber an existing index: persisting a fresh empty
322        // metadata.json would orphan every committed segment, and the next
323        // writer open's orphan sweep would permanently delete them.
324        if directory
325            .exists(std::path::Path::new(super::INDEX_META_FILENAME))
326            .await?
327        {
328            return Err(Error::Internal(format!(
329                "refusing to create index: {} already exists in this directory; \
330                 use IndexWriter::open to open the existing index, or delete the \
331                 directory first if you really want to start over",
332                super::INDEX_META_FILENAME
333            )));
334        }
335
336        let metadata = super::IndexMetadata::new((*schema).clone());
337
338        let segment_manager = Arc::new(crate::merge::SegmentManager::new(
339            Arc::clone(&directory),
340            Arc::clone(&schema),
341            metadata,
342            config.merge_policy.clone_box(),
343            config.term_cache_blocks,
344            config.max_concurrent_merges,
345            Arc::clone(&config.background_merge_permits),
346            config.merge_bp_time_budget,
347            config.bp_memory_budget_bytes,
348            Arc::clone(&config.background_reorder_permits),
349            config.background_reorder_pool.clone(),
350        ));
351        segment_manager.update_metadata(|_| {}).await?;
352
353        Ok(Self::new_with_parts(
354            directory,
355            schema,
356            config,
357            builder_config,
358            segment_manager,
359            writer_lock,
360        ))
361    }
362
363    /// Open an existing index for exclusive writing.
364    ///
365    /// Multiple independent writers for the same directory are unsupported;
366    /// for filesystem-rooted directories this is enforced with an advisory
367    /// single-writer lock ([`WRITER_LOCK_FILENAME`]) held for the writer's
368    /// lifetime. This path removes crash-leftover outputs before starting its
369    /// workers. Use [`Index::writer`](super::Index::writer) to share lifecycle
370    /// state with an already-open search index.
371    pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
372        Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
373    }
374
375    /// Open an existing index with custom builder config
376    pub async fn open_with_config(
377        directory: D,
378        config: IndexConfig,
379        builder_config: SegmentBuilderConfig,
380    ) -> Result<Self> {
381        let directory = Arc::new(directory);
382
383        // The lock must be held before the orphan sweep below: sweeping while
384        // another process's writer is live deletes its in-flight outputs.
385        let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
386        if let WriterLock::Unavailable { reason } = &writer_lock {
387            return Err(Error::Internal(reason.clone()));
388        }
389
390        let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
391        let schema = Arc::new(metadata.schema.clone());
392        // Directory-layer metrics (cold writes, lazy reads) carry the index label
393        directory.set_index_label(schema.index_label());
394
395        let segment_manager = Arc::new(crate::merge::SegmentManager::new(
396            Arc::clone(&directory),
397            Arc::clone(&schema),
398            metadata,
399            config.merge_policy.clone_box(),
400            config.term_cache_blocks,
401            config.max_concurrent_merges,
402            Arc::clone(&config.background_merge_permits),
403            config.merge_bp_time_budget,
404            config.bp_memory_budget_bytes,
405            Arc::clone(&config.background_reorder_permits),
406            config.background_reorder_pool.clone(),
407        ));
408        let swept = segment_manager.cleanup_orphan_segments().await?;
409        if swept > 0 {
410            log::warn!(
411                "[segment_cleanup] swept {} orphan segment(s) while opening writer",
412                swept
413            );
414        }
415        segment_manager.try_load_and_publish_trained().await?;
416
417        Ok(Self::new_with_parts(
418            directory,
419            schema,
420            config,
421            builder_config,
422            segment_manager,
423            writer_lock,
424        ))
425    }
426
427    /// Create an IndexWriter from an existing Index.
428    /// Shares the SegmentManager for consistent segment lifecycle management.
429    ///
430    /// This constructor is infallible, so a single-writer lock conflict is
431    /// deferred: the returned writer fails loudly on its first mutating
432    /// operation instead of silently double-writing next to another writer.
433    pub fn from_index(index: &super::Index<D>) -> Self {
434        let writer_lock = match try_acquire_writer_lock(index.directory.as_ref()) {
435            Ok(lock) => lock,
436            Err(error) => WriterLock::Unavailable {
437                reason: format!("failed to acquire the single-writer lock: {error}"),
438            },
439        };
440        if let WriterLock::Unavailable { reason } = &writer_lock {
441            log::error!("[writer_lock] {reason}");
442        }
443        Self::new_with_parts(
444            Arc::clone(&index.directory),
445            Arc::clone(&index.schema),
446            index.config.clone(),
447            SegmentBuilderConfig::default(),
448            Arc::clone(&index.segment_manager),
449            writer_lock,
450        )
451    }
452
453    // ========================================================================
454    // Construction + pipeline management
455    // ========================================================================
456
457    /// Common construction: creates worker state, spawns workers, assembles `Self`.
458    fn new_with_parts(
459        directory: Arc<D>,
460        schema: Arc<Schema>,
461        config: IndexConfig,
462        builder_config: SegmentBuilderConfig,
463        segment_manager: Arc<crate::merge::SegmentManager<D>>,
464        writer_lock: WriterLock,
465    ) -> Self {
466        // Auto-configure tokenizers from schema for all text fields
467        let registry = crate::tokenizer::TokenizerRegistry::new();
468        let mut tokenizers = FxHashMap::default();
469        for (field, entry) in schema.fields() {
470            if matches!(entry.field_type, crate::dsl::FieldType::Text)
471                && let Some(ref tok_name) = entry.tokenizer
472                && let Some(tok) = registry.get(tok_name)
473            {
474                tokenizers.insert(field, tok);
475            }
476        }
477
478        let num_workers = config.num_indexing_threads.max(1);
479        let worker_state = Arc::new(WorkerState {
480            directory: Arc::clone(&directory),
481            schema: Arc::clone(&schema),
482            builder_config,
483            tokenizers: parking_lot::RwLock::new(tokenizers),
484            memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
485            segment_manager: Arc::clone(&segment_manager),
486            built_segments: parking_lot::Mutex::new(Vec::new()),
487            cycle_error: parking_lot::Mutex::new(None),
488            cycle_failed: AtomicBool::new(false),
489            flush_count: AtomicUsize::new(0),
490            flush_mutex: parking_lot::Mutex::new(()),
491            flush_cvar: parking_lot::Condvar::new(),
492            resume_receiver: parking_lot::Mutex::new(None),
493            resume_epoch: AtomicUsize::new(0),
494            resume_cvar: parking_lot::Condvar::new(),
495            shutdown: AtomicBool::new(false),
496            num_workers,
497        });
498        let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
499
500        Self {
501            directory,
502            schema,
503            config,
504            doc_sender: Arc::new(parking_lot::RwLock::new(doc_sender)),
505            workers,
506            worker_state,
507            segment_manager,
508            flushed_segments: Arc::new(parking_lot::Mutex::new(Vec::new())),
509            primary_key_index: Arc::new(parking_lot::RwLock::new(None)),
510            commit_finalization: Arc::new(CommitFinalizationState::default()),
511            pk_reservations_retained: Arc::new(AtomicBool::new(false)),
512            writer_lock: parking_lot::RwLock::new(writer_lock),
513        }
514    }
515
516    /// Fail loudly when another writer owns the single-writer lock.
517    ///
518    /// A deferred conflict (`from_index` during a writer handover, e.g. a
519    /// rolling pod restart) is not permanent: the holder exits and the kernel
520    /// releases its advisory lock. Re-attempt acquisition on every call in
521    /// the `Unavailable` state so the writer recovers as soon as the lock
522    /// frees, instead of rejecting all writes for its lifetime.
523    fn ensure_writer_lock(&self) -> Result<()> {
524        // Fast path: uncontended read on the healthy states.
525        if !matches!(&*self.writer_lock.read(), WriterLock::Unavailable { .. }) {
526            return Ok(());
527        }
528
529        let mut lock = self.writer_lock.write();
530        // Another thread may have recovered while we waited for the write lock.
531        if !matches!(&*lock, WriterLock::Unavailable { .. }) {
532            return Ok(());
533        }
534        match try_acquire_writer_lock(self.directory.as_ref())? {
535            acquired @ (WriterLock::Held { .. } | WriterLock::NotApplicable) => {
536                log::info!(
537                    "[writer_lock] single-writer lock acquired after retry; \
538                     the previous holder has released it — resuming writes"
539                );
540                *lock = acquired;
541                Ok(())
542            }
543            WriterLock::Unavailable { reason } => {
544                let err = Error::Internal(reason.clone());
545                *lock = WriterLock::Unavailable { reason };
546                Err(err)
547            }
548        }
549    }
550
551    /// Clear primary-key reservations after an aborted or failed generation.
552    ///
553    /// Skipped while a failed post-commit PK refresh has left the uncommitted
554    /// reservations as the ONLY record of already-committed keys (fail-closed,
555    /// see `finalize_prepared_commit`): wiping them would admit duplicate
556    /// primary keys. Retaining the aborted generation's keys as well is
557    /// deliberately conservative — they clear on the next successful commit's
558    /// refresh.
559    fn clear_uncommitted_pk_reservations(&self) {
560        if self.pk_reservations_retained.load(Ordering::Acquire) {
561            log::warn!(
562                "[primary_key] keeping uncommitted reservations through abort: a \
563                 failed post-commit refresh left them as the only record of \
564                 committed keys; they are cleared by the next successful commit"
565            );
566            return;
567        }
568        if let Some(pk_index) = self.primary_key_index.write().as_mut() {
569            pk_index.clear_uncommitted();
570        }
571    }
572
573    fn spawn_workers(
574        worker_state: &Arc<WorkerState<D>>,
575        num_workers: usize,
576    ) -> (
577        async_channel::Sender<Document>,
578        Vec<std::thread::JoinHandle<()>>,
579    ) {
580        let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
581        let handle = tokio::runtime::Handle::current();
582        let mut workers = Vec::with_capacity(num_workers);
583        for i in 0..num_workers {
584            let state = Arc::clone(worker_state);
585            let rx = receiver.clone();
586            let rt = handle.clone();
587            workers.push(
588                std::thread::Builder::new()
589                    .name(format!("index-worker-{}", i))
590                    .spawn(move || Self::worker_loop(state, rx, rt))
591                    .expect("failed to spawn index worker thread"),
592            );
593        }
594        (sender, workers)
595    }
596
597    /// Get the schema
598    pub fn schema(&self) -> &Schema {
599        &self.schema
600    }
601
602    /// Set tokenizer for a field.
603    /// Propagated to worker threads — takes effect for the next SegmentBuilder they create.
604    pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
605        self.worker_state
606            .tokenizers
607            .write()
608            .insert(field, Box::new(tokenizer));
609    }
610
611    /// Initialize primary key deduplication from committed segments.
612    ///
613    /// Tries to load a cached bloom filter from `pk_bloom.bin` first. If the
614    /// cache covers all current segments, the bloom is reused directly (fast
615    /// path). If new segments appeared since the cache was written, only their
616    /// keys are iterated (incremental). Falls back to a full rebuild when no
617    /// cache exists.
618    ///
619    /// Only loads fast-field data (text dictionaries) per segment — NOT full
620    /// `SegmentReader`s — to avoid duplicating dense/sparse index memory.
621    ///
622    /// The CPU-intensive bloom build is offloaded via `spawn_blocking` so it
623    /// does not block the tokio runtime.
624    ///
625    /// No-op if schema has no primary field.
626    pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
627        use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
628
629        self.commit_finalization.wait_until_idle().await;
630
631        let field = match self.schema.primary_field() {
632            Some(f) => f,
633            None => return Ok(()),
634        };
635
636        let snapshot = self.segment_manager.acquire_snapshot().await;
637        let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
638
639        // Try to load persisted bloom filter.
640        let cached = match self
641            .directory
642            .open_read(std::path::Path::new(PK_BLOOM_FILE))
643            .await
644        {
645            Ok(handle) => {
646                let data = handle.read_bytes_range(0..handle.len()).await;
647                match data {
648                    Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
649                    Err(_) => None,
650                }
651            }
652            Err(_) => None,
653        };
654
655        // Load lightweight fast-field data for all segments concurrently.
656        let load_futures: Vec<_> = current_seg_ids
657            .iter()
658            .map(|seg_id_str| {
659                let seg_id_str = seg_id_str.clone();
660                let dir = self.directory.as_ref();
661                let schema = Arc::clone(&self.schema);
662                async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
663            })
664            .collect();
665        let all_data = futures::future::try_join_all(load_futures).await?;
666
667        if let Some((persisted_seg_ids, bloom)) = cached {
668            // Partition: old segments (covered by bloom) first, new segments at end.
669            let mut pk_data = Vec::with_capacity(all_data.len());
670            let mut new_data = Vec::new();
671            for d in all_data {
672                if persisted_seg_ids.contains(&d.segment_id) {
673                    pk_data.push(d);
674                } else {
675                    new_data.push(d);
676                }
677            }
678            let needs_persist = !new_data.is_empty();
679            let new_start = pk_data.len();
680            pk_data.extend(new_data);
681
682            let pk_index = if new_start == pk_data.len() {
683                // Fast path: all segments covered by cache.
684                super::primary_key::PrimaryKeyIndex::from_persisted(
685                    field,
686                    bloom,
687                    pk_data,
688                    &[],
689                    snapshot,
690                )
691            } else {
692                // Incremental: only iterate new segments' keys.
693                tokio::task::spawn_blocking(move || {
694                    // Insert new segments' keys into the bloom, then construct
695                    // PrimaryKeyIndex with the pre-populated bloom.
696                    let mut bloom = bloom;
697                    let mut added = 0usize;
698                    let num_new = pk_data.len() - new_start;
699                    for data in &pk_data[new_start..] {
700                        if let Some(ff) = data.fast_fields.get(&field.0)
701                            && let Some(dict) = ff.text_dict()
702                        {
703                            for key in dict.iter() {
704                                bloom.insert(key.as_bytes());
705                                added += 1;
706                            }
707                        }
708                    }
709                    if added > 0 {
710                        log::info!(
711                            "[primary_key] bloom: added {} keys from {} new segment(s)",
712                            added,
713                            num_new,
714                        );
715                    }
716                    super::primary_key::PrimaryKeyIndex::from_persisted(
717                        field,
718                        bloom,
719                        pk_data,
720                        &[],
721                        snapshot,
722                    )
723                })
724                .await
725                .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
726            };
727
728            if needs_persist {
729                self.persist_pk_bloom(&pk_index, &current_seg_ids).await;
730            }
731
732            *self.primary_key_index.write() = Some(pk_index);
733        } else {
734            // No cache — full rebuild, offloaded to blocking thread.
735            let pk_index = tokio::task::spawn_blocking(move || {
736                super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
737            })
738            .await
739            .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
740
741            self.persist_pk_bloom(&pk_index, &current_seg_ids).await;
742            *self.primary_key_index.write() = Some(pk_index);
743        }
744
745        // The freshly built index covers every committed segment, so any
746        // reservations retained after a failed post-commit refresh are
747        // superseded by committed_data.
748        self.pk_reservations_retained
749            .store(false, Ordering::Release);
750
751        Ok(())
752    }
753
754    /// Persist the primary-key bloom filter to `pk_bloom.bin`.
755    /// Best-effort: errors are logged but not propagated.
756    async fn persist_pk_bloom(
757        &self,
758        pk_index: &super::primary_key::PrimaryKeyIndex,
759        segment_ids: &[String],
760    ) {
761        use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
762
763        let bloom_bytes = pk_index.bloom_to_bytes();
764        let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
765        if let Err(e) = self
766            .directory
767            .write(std::path::Path::new(PK_BLOOM_FILE), &data)
768            .await
769        {
770            log::warn!("[primary_key] failed to persist bloom cache: {}", e);
771        }
772    }
773
774    /// Add a document to the indexing queue (sync, O(1)).
775    ///
776    /// `Document` is moved into the channel (zero-copy). Workers compete to pull it.
777    /// Returns an explicit backpressure error when the queue is at capacity or
778    /// a prepared commit generation is not yet resolved.
779    pub fn add_document(&self, doc: Document) -> Result<()> {
780        self.ensure_writer_lock()?;
781        if self.worker_state.shutdown.load(Ordering::Acquire) {
782            return Err(Error::IndexClosed);
783        }
784        if self.commit_finalization.in_progress.load(Ordering::Acquire) {
785            return Err(Error::CommitInProgress);
786        }
787        let sender = self.doc_sender.read().clone();
788        // A publication error deliberately leaves the prepared generation and
789        // its workers paused for a lossless retry. Report this as backpressure
790        // instead of inserting/rolling back a PK key against a closed channel.
791        if sender.is_closed() {
792            return Err(Error::CommitInProgress);
793        }
794        let primary_key_index = self.primary_key_index.read();
795        if let Some(ref pk_index) = *primary_key_index {
796            pk_index.check_and_insert(&doc)?;
797        }
798        match sender.try_send(doc) {
799            Ok(()) => Ok(()),
800            Err(async_channel::TrySendError::Full(doc)) => {
801                // Roll back PK registration so the caller can retry later
802                if let Some(ref pk_index) = *primary_key_index {
803                    pk_index.rollback_uncommitted_key(&doc);
804                }
805                Err(Error::QueueFull)
806            }
807            Err(async_channel::TrySendError::Closed(doc)) => {
808                // Roll back PK registration for defense-in-depth
809                if let Some(ref pk_index) = *primary_key_index {
810                    pk_index.rollback_uncommitted_key(&doc);
811                }
812                Err(Error::CommitInProgress)
813            }
814        }
815    }
816
817    /// Add multiple documents to the indexing queue.
818    ///
819    /// Returns the number of documents successfully queued. Stops at the first
820    /// backpressure error and returns the count queued so far.
821    pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
822        let total = documents.len();
823        for (i, doc) in documents.into_iter().enumerate() {
824            match self.add_document(doc) {
825                Ok(()) => {}
826                Err(Error::QueueFull | Error::CommitInProgress) => return Ok(i),
827                Err(e) => return Err(e),
828            }
829        }
830        Ok(total)
831    }
832
833    // ========================================================================
834    // Worker loop
835    // ========================================================================
836
837    /// Worker loop — runs on a dedicated OS thread, survives across commits.
838    ///
839    /// Outer loop: each iteration processes one commit cycle.
840    ///   Inner loop: pull documents from MPMC queue, index them, build segments
841    ///   when memory budget is exceeded.
842    ///   On channel close (prepare_commit): flush current builder, signal
843    ///   flush_count, wait for resume with new receiver.
844    ///   On shutdown (Drop): exit permanently.
845    fn worker_loop(
846        state: Arc<WorkerState<D>>,
847        initial_receiver: async_channel::Receiver<Document>,
848        handle: tokio::runtime::Handle,
849    ) {
850        let mut receiver = initial_receiver;
851        let mut my_epoch = 0usize;
852
853        loop {
854            // Wrap the recv+build phase in catch_unwind so a panic doesn't
855            // prevent flush_count from being signaled (which would hang
856            // prepare_commit forever).
857            let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
858                let mut builder: Option<SegmentBuilder> = None;
859
860                while let Ok(doc) = receiver.recv_blocking() {
861                    if state.shutdown.load(Ordering::Acquire) {
862                        break;
863                    }
864                    // Another worker already invalidated this generation.
865                    // Drain the shared queue so prepare_commit can complete,
866                    // but do not spend CPU/RAM building outputs that must be
867                    // discarded transactionally.
868                    if state.cycle_failed.load(Ordering::Acquire) {
869                        continue;
870                    }
871                    // Initialize builder if needed
872                    if builder.is_none() {
873                        match SegmentBuilder::new(
874                            Arc::clone(&state.schema),
875                            state.builder_config.clone(),
876                        ) {
877                            Ok(mut b) => {
878                                for (field, tokenizer) in state.tokenizers.read().iter() {
879                                    b.set_tokenizer(*field, tokenizer.clone_box());
880                                }
881                                builder = Some(b);
882                            }
883                            Err(e) => {
884                                log::error!("Failed to create segment builder: {:?}", e);
885                                state.record_cycle_error(format!(
886                                    "failed to create segment builder: {e}"
887                                ));
888                                continue;
889                            }
890                        }
891                    }
892
893                    let b = builder.as_mut().unwrap();
894                    if let Err(e) = b.add_document(doc) {
895                        log::error!("Failed to index document: {:?}", e);
896                        state.record_cycle_error(format!("failed to index document: {e}"));
897                        continue;
898                    }
899
900                    let builder_memory = b.estimated_memory_bytes();
901
902                    if b.num_docs() & 0x3FFF == 0 {
903                        log::debug!(
904                            "[indexing] docs={}, memory={:.2} MB, budget={:.2} MB",
905                            b.num_docs(),
906                            builder_memory as f64 / (1024.0 * 1024.0),
907                            state.memory_budget_per_worker as f64 / (1024.0 * 1024.0)
908                        );
909                    }
910
911                    // Require minimum 100 docs before flushing to avoid tiny segments
912                    const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
913
914                    // Reserve 20% headroom for segment build overhead (vid_set,
915                    // VidLookup, postings_flat, grid_entries). These temporary
916                    // allocations exist alongside the builder's data during build.
917                    let effective_budget = state.memory_budget_per_worker * 4 / 5;
918
919                    if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
920                        log::info!(
921                            "[indexing] memory budget reached, building segment: \
922                             docs={}, memory={:.2} MB, budget={:.2} MB",
923                            b.num_docs(),
924                            builder_memory as f64 / (1024.0 * 1024.0),
925                            state.memory_budget_per_worker as f64 / (1024.0 * 1024.0),
926                        );
927                        let full_builder = builder.take().unwrap();
928                        Self::build_segment_inline(&state, full_builder, &handle);
929                    }
930                }
931
932                // Channel closed — flush current builder
933                if !state.cycle_failed.load(Ordering::Acquire)
934                    && let Some(b) = builder.take()
935                    && b.num_docs() > 0
936                {
937                    Self::build_segment_inline(&state, b, &handle);
938                }
939            }));
940
941            if build_result.is_err() {
942                log::error!(
943                    "[worker] panic during indexing cycle — documents in this cycle may be lost"
944                );
945                state.record_cycle_error("indexing worker panicked while building the batch");
946            }
947
948            // Signal flush completion (always, even after panic — prevents
949            // prepare_commit from hanging)
950            let prev = state.flush_count.fetch_add(1, Ordering::Release);
951            if prev + 1 == state.num_workers {
952                // Last worker — wake prepare_commit. notify_all, not
953                // notify_one: a cancelled commit leaves its detached
954                // spawn_blocking waiter parked on this condvar, and with a
955                // single notification that dead waiter would consume the
956                // only wakeup, stalling a retried prepare_commit for its
957                // full deadline.
958                let _lock = state.flush_mutex.lock();
959                state.flush_cvar.notify_all();
960            }
961
962            // Wait for resume (new channel) or shutdown.
963            // Check resume_epoch to avoid re-cloning a stale receiver from
964            // a previous cycle.
965            {
966                let mut lock = state.resume_receiver.lock();
967                loop {
968                    if state.shutdown.load(Ordering::Acquire) {
969                        return;
970                    }
971                    let current_epoch = state.resume_epoch.load(Ordering::Acquire);
972                    if current_epoch > my_epoch
973                        && let Some(rx) = lock.as_ref()
974                    {
975                        receiver = rx.clone();
976                        my_epoch = current_epoch;
977                        break;
978                    }
979                    state.resume_cvar.wait(&mut lock);
980                }
981            }
982        }
983    }
984
985    /// Build a segment on the worker thread. Uses `Handle::block_on()` to bridge
986    /// into async context for I/O (streaming writers). CPU work (rayon) stays on
987    /// the worker thread / rayon pool.
988    fn build_segment_inline(
989        state: &WorkerState<D>,
990        builder: SegmentBuilder,
991        handle: &tokio::runtime::Handle,
992    ) {
993        let segment_id = SegmentId::new();
994        let segment_hex = segment_id.to_hex();
995        // Claim the ID before the first file write. The guard is moved into
996        // `PreparedSegment` on success and otherwise releases automatically.
997        let operation = match state
998            .segment_manager
999            .protect_new_segment(segment_hex.clone())
1000        {
1001            Ok(operation) => operation,
1002            Err(e) => {
1003                log::error!(
1004                    "[segment_build_failed] segment_id={} lifecycle_error={}",
1005                    segment_hex,
1006                    e,
1007                );
1008                state.record_cycle_error(format!(
1009                    "failed to claim segment {segment_hex} for building: {e}"
1010                ));
1011                return;
1012            }
1013        };
1014        let trained = state.segment_manager.trained_for_segment_build();
1015        let doc_count = builder.num_docs();
1016        let build_start = std::time::Instant::now();
1017
1018        log::info!(
1019            "[segment_build] segment_id={} doc_count={} ann={}",
1020            segment_hex,
1021            doc_count,
1022            trained.is_some()
1023        );
1024
1025        // Construct the cleanup owner before building. It keeps lifecycle
1026        // ownership through async deletion on ordinary error, abort, and
1027        // panic unwind; crash recovery is the only path left to the sweeper.
1028        let mut prepared = PreparedSegment {
1029            id: segment_hex.clone(),
1030            segment_id,
1031            num_docs: doc_count,
1032            segment_manager: Arc::clone(&state.segment_manager),
1033            operation: Some(operation),
1034            runtime: handle.clone(),
1035            needs_vector_upgrade: trained.is_none(),
1036            published: false,
1037        };
1038
1039        match handle.block_on(builder.build(
1040            state.directory.as_ref(),
1041            segment_id,
1042            trained.as_deref(),
1043        )) {
1044            Ok(meta) if meta.num_docs == doc_count && meta.num_docs > 0 => {
1045                let duration_ms = build_start.elapsed().as_millis() as u64;
1046                log::info!(
1047                    "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
1048                    segment_hex,
1049                    meta.num_docs,
1050                    duration_ms,
1051                );
1052                prepared.num_docs = meta.num_docs;
1053                state.built_segments.lock().push(prepared);
1054            }
1055            Ok(meta) => {
1056                let error = format!(
1057                    "segment {segment_hex} built {} docs from a {doc_count}-document builder",
1058                    meta.num_docs
1059                );
1060                log::error!("[segment_build_failed] {error}");
1061                state.record_cycle_error(error);
1062            }
1063            Err(e) => {
1064                log::error!(
1065                    "[segment_build_failed] segment_id={} error={:?}",
1066                    segment_hex,
1067                    e
1068                );
1069                // `prepared` owns the lifecycle claim and schedules one
1070                // tracked, idempotent cleanup pass when this scope ends.
1071                state.record_cycle_error(format!("failed to build segment {segment_hex}: {e}"));
1072            }
1073        }
1074    }
1075
1076    // ========================================================================
1077    // Public API — commit, merge, etc.
1078    // ========================================================================
1079
1080    /// Check merge policy and spawn a background merge if needed.
1081    pub async fn maybe_merge(&self) {
1082        self.segment_manager.maybe_merge().await;
1083    }
1084
1085    /// Drain all in-flight merge tasks.
1086    /// Blocking merge phases cannot be cancelled safely once started.
1087    pub async fn abort_merges(&self) {
1088        self.segment_manager.abort_merges().await;
1089    }
1090
1091    /// Stop accepting lifecycle work, stop and join indexing workers, and
1092    /// discard unpublished segments. Index deletion calls this while holding
1093    /// the registry writer lock so in-flight requests finish first and stale
1094    /// writer Arcs cannot restart work afterward.
1095    pub async fn shutdown(&mut self) -> Result<()> {
1096        self.segment_manager.begin_shutdown();
1097        self.signal_worker_shutdown();
1098
1099        // A cancelled commit request leaves its owned finalizer running. Do not
1100        // clear shared PK/prepared state while that task may still publish or
1101        // refresh it. Worker shutdown is signalled first, so a successful
1102        // finalizer cannot restart ingestion while deletion is waiting.
1103        self.commit_finalization.wait_until_idle().await;
1104
1105        let workers = std::mem::take(&mut self.workers);
1106        let panicked = tokio::task::spawn_blocking(move || {
1107            workers
1108                .into_iter()
1109                .map(|worker| worker.join().is_err())
1110                .filter(|panicked| *panicked)
1111                .count()
1112        })
1113        .await
1114        .map_err(|error| Error::Internal(format!("failed to join index workers: {}", error)))?;
1115        if panicked > 0 {
1116            log::error!("[index_shutdown] {} indexing worker(s) panicked", panicked);
1117        }
1118
1119        // No commit is possible after shutdown. Dropping these RAII values
1120        // releases their lifecycle ownership before directory deletion.
1121        self.flushed_segments.lock().clear();
1122        self.worker_state.built_segments.lock().clear();
1123        if let Some(pk_index) = self.primary_key_index.write().as_mut() {
1124            pk_index.clear_uncommitted();
1125        }
1126        Ok(())
1127    }
1128
1129    /// Wait for the in-flight background merge to complete (if any).
1130    pub async fn wait_for_merging_thread(&self) {
1131        self.segment_manager.wait_for_merging_thread().await;
1132    }
1133
1134    /// Wait for all eligible merges to complete, including cascading merges.
1135    pub async fn wait_for_all_merges(&self) {
1136        self.segment_manager.wait_for_all_merges().await;
1137    }
1138
1139    /// Wait until an owned commit finalizer has reconciled durable metadata,
1140    /// primary-key state, and worker availability. Normally callers need not
1141    /// use this: it exists for orderly shutdown and request supervisors that
1142    /// want to observe completion after cancelling their original waiter.
1143    pub async fn wait_for_commit_finalization(&self) {
1144        self.commit_finalization.wait_until_idle().await;
1145    }
1146
1147    /// Get the segment tracker for sharing with readers.
1148    pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
1149        self.segment_manager.tracker()
1150    }
1151
1152    /// Acquire a snapshot of current segments for reading.
1153    pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
1154        self.segment_manager.acquire_snapshot().await
1155    }
1156
1157    /// Clean up orphan segment files not registered in metadata.
1158    ///
1159    /// Requires the single-writer lock: sweeping while another process's
1160    /// writer is live would delete its in-flight segment outputs.
1161    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
1162        self.ensure_writer_lock()?;
1163        self.segment_manager.cleanup_orphan_segments().await
1164    }
1165
1166    /// Prepare commit — signal workers to flush, wait for completion, collect segments.
1167    ///
1168    /// All documents sent via `add_document` before this call are guaranteed
1169    /// to be written to segment files on disk. Segments are NOT yet registered
1170    /// in metadata — call `PreparedCommit::commit()` for that.
1171    ///
1172    /// Workers are NOT destroyed — they flush their builders and wait for
1173    /// `resume_workers()` to give them a new channel.
1174    ///
1175    /// `add_document` returns `CommitInProgress` until commit/abort resumes workers.
1176    pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
1177        self.ensure_writer_lock()?;
1178        if self.worker_state.shutdown.load(Ordering::Acquire) {
1179            return Err(Error::IndexClosed);
1180        }
1181        if self.commit_finalization.in_progress.load(Ordering::Acquire) {
1182            return Err(Error::CommitInProgress);
1183        }
1184        // 1. Close channel → workers drain remaining docs and flush builders
1185        self.doc_sender.read().close();
1186
1187        // Wake any workers still waiting on resume_cvar from previous cycle.
1188        // They'll clone the stale receiver, enter recv_blocking, get Err
1189        // immediately (sender already closed), flush, and signal completion.
1190        self.worker_state.resume_cvar.notify_all();
1191
1192        // 2. Wait for all workers to complete their flush (via spawn_blocking
1193        //    to avoid blocking the tokio runtime)
1194        let state = Arc::clone(&self.worker_state);
1195        let all_flushed = tokio::task::spawn_blocking(move || {
1196            let mut lock = state.flush_mutex.lock();
1197            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
1198            while state.flush_count.load(Ordering::Acquire) < state.num_workers {
1199                let remaining = deadline.saturating_duration_since(std::time::Instant::now());
1200                if remaining.is_zero() {
1201                    log::error!(
1202                        "[prepare_commit] timed out waiting for workers: {}/{} flushed",
1203                        state.flush_count.load(Ordering::Acquire),
1204                        state.num_workers
1205                    );
1206                    return false;
1207                }
1208                state.flush_cvar.wait_for(&mut lock, remaining);
1209            }
1210            true
1211        })
1212        .await
1213        .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
1214
1215        if !all_flushed {
1216            // Keep this commit cycle paused. Resetting flush_count and handing
1217            // out a new receiver while an old worker is still building lets
1218            // that late worker increment the *next* cycle's counter. A later
1219            // prepare can then return before all of its workers flushed and
1220            // publish an incomplete set of segments. The caller may retry
1221            // prepare_commit; it will observe the same generation and collect
1222            // every completed output once the lagging worker finishes.
1223            return Err(Error::Internal(format!(
1224                "prepare_commit timed out: {}/{} workers flushed; writer remains paused, retry commit",
1225                self.worker_state.flush_count.load(Ordering::Acquire),
1226                self.worker_state.num_workers
1227            )));
1228        }
1229
1230        let cycle_error = { self.worker_state.cycle_error.lock().take() };
1231        if let Some(error) = cycle_error {
1232            // No partial publication: some documents in this generation no
1233            // longer exist in a worker builder, so successful sibling outputs
1234            // cannot be committed without violating commit's all-prior-docs
1235            // guarantee. Their RAII drops retain ownership through deletion.
1236            self.flushed_segments.lock().clear();
1237            self.worker_state.built_segments.lock().clear();
1238            self.clear_uncommitted_pk_reservations();
1239            self.resume_workers();
1240            return Err(Error::Internal(format!(
1241                "indexing generation failed; no documents from this batch were committed: {error}"
1242            )));
1243        }
1244
1245        // 3. Collect built segments
1246        let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
1247        self.flushed_segments.lock().extend(built);
1248
1249        Ok(PreparedCommit {
1250            writer: self,
1251            is_resolved: false,
1252        })
1253    }
1254
1255    /// Commit (convenience): prepare_commit + commit in one call.
1256    ///
1257    /// Guarantees all prior `add_document` calls are committed.
1258    /// Vector training is decoupled — call `build_vector_index()` manually.
1259    pub async fn commit(&mut self) -> Result<bool> {
1260        self.prepare_commit().await?.commit().await
1261    }
1262
1263    /// Force merge all segments into one.
1264    pub async fn force_merge(&mut self) -> Result<()> {
1265        self.prepare_commit().await?.commit().await?;
1266        self.segment_manager.force_merge().await
1267    }
1268
1269    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
1270    ///
1271    /// Each segment is individually rebuilt with record-level BP reordering:
1272    /// ordinals are shuffled across blocks so that similar content clusters tightly.
1273    pub async fn reorder(&mut self) -> Result<()> {
1274        self.prepare_commit().await?.commit().await?;
1275        self.segment_manager.reorder_segments().await
1276    }
1277
1278    /// Get the segment manager (for background optimizer access).
1279    pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
1280        &self.segment_manager
1281    }
1282
1283    /// Resume workers with a fresh channel. Called after commit or abort.
1284    ///
1285    /// Workers are already alive — just give them a new channel and wake them.
1286    /// If the tokio runtime has shut down (e.g., program exit), this is a no-op.
1287    fn resume_workers(&mut self) {
1288        Self::resume_workers_shared(&self.worker_state, &self.doc_sender);
1289    }
1290
1291    fn resume_workers_shared(
1292        worker_state: &Arc<WorkerState<D>>,
1293        doc_sender: &Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1294    ) {
1295        if worker_state.shutdown.load(Ordering::Acquire) {
1296            return;
1297        }
1298        if tokio::runtime::Handle::try_current().is_err() {
1299            // Runtime is gone — signal permanent shutdown so workers don't
1300            // hang forever on resume_cvar.
1301            worker_state.shutdown.store(true, Ordering::Release);
1302            worker_state.resume_cvar.notify_all();
1303            return;
1304        }
1305
1306        // Reset flush count for next cycle
1307        worker_state.flush_count.store(0, Ordering::Release);
1308        *worker_state.cycle_error.lock() = None;
1309        worker_state.cycle_failed.store(false, Ordering::Release);
1310
1311        // Create new channel
1312        let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
1313        *doc_sender.write() = sender;
1314
1315        // Set new receiver, bump epoch, and wake all workers
1316        {
1317            let mut lock = worker_state.resume_receiver.lock();
1318            *lock = Some(receiver);
1319        }
1320        worker_state.resume_epoch.fetch_add(1, Ordering::Release);
1321        worker_state.resume_cvar.notify_all();
1322    }
1323
1324    fn signal_worker_shutdown(&self) {
1325        self.worker_state.shutdown.store(true, Ordering::Release);
1326        self.doc_sender.read().close();
1327        self.worker_state.resume_cvar.notify_all();
1328    }
1329
1330    // Vector index methods (build_vector_index, etc.) are in vector_builder.rs
1331}
1332
1333impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
1334    fn drop(&mut self) {
1335        self.signal_worker_shutdown();
1336        for w in std::mem::take(&mut self.workers) {
1337            let _ = w.join();
1338        }
1339    }
1340}
1341
1342/// A prepared commit that can be finalized or aborted.
1343///
1344/// Two-phase commit guard. Between `prepare_commit()` and
1345/// `commit()`/`abort()`, segments are on disk but NOT in metadata.
1346/// Dropping without calling either will auto-abort (discard segments,
1347/// respawn workers).
1348pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
1349    writer: &'a mut IndexWriter<D>,
1350    is_resolved: bool,
1351}
1352
1353/// Returns prepared segments to the writer if an owned commit finalizer fails
1354/// or unwinds before it can establish that metadata owns them. Retrying commit
1355/// is safe even when publication actually won the race: `SegmentManager::commit`
1356/// is idempotent and the operation guards keep the files protected meanwhile.
1357struct PreparedSegmentsGuard<D: DirectoryWriter + 'static> {
1358    segments: Option<Vec<PreparedSegment<D>>>,
1359    retry_slot: Arc<parking_lot::Mutex<Vec<PreparedSegment<D>>>>,
1360}
1361
1362impl<D: DirectoryWriter + 'static> PreparedSegmentsGuard<D> {
1363    fn metadata_entries(&self) -> Vec<(String, u32)> {
1364        self.segments
1365            .as_deref()
1366            .unwrap_or_default()
1367            .iter()
1368            .map(PreparedSegment::metadata_entry)
1369            .collect()
1370    }
1371
1372    fn take_published(&mut self) -> Vec<PreparedSegment<D>> {
1373        self.segments.take().unwrap_or_default()
1374    }
1375
1376    fn vector_upgrade_segment_ids(&self) -> Vec<String> {
1377        self.segments
1378            .as_deref()
1379            .unwrap_or_default()
1380            .iter()
1381            .filter(|segment| segment.needs_vector_upgrade)
1382            .map(|segment| segment.id.clone())
1383            .collect()
1384    }
1385}
1386
1387impl<D: DirectoryWriter + 'static> Drop for PreparedSegmentsGuard<D> {
1388    fn drop(&mut self) {
1389        if let Some(segments) = self.segments.take() {
1390            self.retry_slot.lock().extend(segments);
1391        }
1392    }
1393}
1394
1395/// Couples completion of the owned commit task to writer availability. The
1396/// default is deliberately fail-closed: a pre-publication error or panic keeps
1397/// workers paused so the retained prepared generation can be retried. Only the
1398/// normal published path arms resumption.
1399struct CommitFinalizationGuard<D: DirectoryWriter + 'static> {
1400    state: Arc<CommitFinalizationState>,
1401    worker_state: Arc<WorkerState<D>>,
1402    doc_sender: Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1403    resume_workers: bool,
1404}
1405
1406impl<D: DirectoryWriter + 'static> CommitFinalizationGuard<D> {
1407    fn resume_on_drop(&mut self) {
1408        self.resume_workers = true;
1409    }
1410}
1411
1412impl<D: DirectoryWriter + 'static> Drop for CommitFinalizationGuard<D> {
1413    fn drop(&mut self) {
1414        if self.resume_workers {
1415            IndexWriter::<D>::resume_workers_shared(&self.worker_state, &self.doc_sender);
1416        }
1417        self.state.finish();
1418    }
1419}
1420
1421/// Everything needed to finish one prepared generation is moved into this
1422/// value before spawning. Its two guards therefore reconcile segment
1423/// ownership and writer availability even if Tokio drops the task before its
1424/// first poll.
1425struct OwnedCommitFinalization<D: DirectoryWriter + 'static> {
1426    directory: Arc<D>,
1427    schema: Arc<Schema>,
1428    segment_manager: Arc<crate::merge::SegmentManager<D>>,
1429    primary_key_index: Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1430    prepared: PreparedSegmentsGuard<D>,
1431    finalization: Option<CommitFinalizationGuard<D>>,
1432    publication_observed: Arc<AtomicBool>,
1433    pk_reservations_retained: Arc<AtomicBool>,
1434}
1435
1436async fn refresh_primary_key_after_commit<D: DirectoryWriter + 'static>(
1437    directory: &Arc<D>,
1438    schema: &Arc<Schema>,
1439    segment_manager: &Arc<crate::merge::SegmentManager<D>>,
1440    primary_key_index: &Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1441) -> Result<()> {
1442    let existing_ids: std::collections::HashSet<String> = {
1443        let guard = primary_key_index.read();
1444        let Some(pk_index) = guard.as_ref() else {
1445            return Ok(());
1446        };
1447        pk_index
1448            .committed_segment_ids()
1449            .map(ToOwned::to_owned)
1450            .collect()
1451    };
1452
1453    let snapshot = segment_manager.acquire_snapshot().await;
1454    let load_futures: Vec<_> = snapshot
1455        .segment_ids()
1456        .iter()
1457        .filter(|id| !existing_ids.contains(id.as_str()))
1458        .map(|seg_id_str| {
1459            let seg_id_str = seg_id_str.clone();
1460            let dir = directory.as_ref();
1461            let schema = Arc::clone(schema);
1462            async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
1463        })
1464        .collect();
1465    let new_data = futures::future::try_join_all(load_futures).await?;
1466    let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
1467
1468    let bloom_file = {
1469        let mut guard = primary_key_index.write();
1470        let Some(pk_index) = guard.as_mut() else {
1471            return Ok(());
1472        };
1473        pk_index.refresh_incremental(new_data, snapshot);
1474        let bloom_bytes = pk_index.bloom_to_bytes();
1475        super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes)
1476    };
1477
1478    if let Err(error) = directory
1479        .write(
1480            std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
1481            &bloom_file,
1482        )
1483        .await
1484    {
1485        log::warn!("[primary_key] failed to persist bloom cache: {}", error);
1486    }
1487    Ok(())
1488}
1489
1490async fn finalize_prepared_commit<D: DirectoryWriter + 'static>(
1491    mut commit: OwnedCommitFinalization<D>,
1492) -> Result<bool> {
1493    let metadata_entries = commit.prepared.metadata_entries();
1494    let published_segment_ids = commit.prepared.vector_upgrade_segment_ids();
1495
1496    // This entire future is owned by a Tokio task. Cancelling the RPC only
1497    // drops its JoinHandle; it cannot split durable metadata publication from
1498    // PK reservations or worker resumption.
1499    commit.segment_manager.commit(&metadata_entries).await?;
1500    commit.publication_observed.store(true, Ordering::Release);
1501
1502    let mut published = commit.prepared.take_published();
1503    for segment in &mut published {
1504        segment.mark_published();
1505    }
1506    drop(published);
1507    commit
1508        .segment_manager
1509        .schedule_vector_segment_upgrades(published_segment_ids);
1510    // Publication is irreversible. From here onward every exit path, including
1511    // panic unwind, must make the writer available again while PK reservations
1512    // remain fail-closed until refresh succeeds.
1513    if let Some(finalization) = commit.finalization.as_mut() {
1514        finalization.resume_on_drop();
1515    } else {
1516        log::error!("owned commit finalization guard was already released after publication");
1517    }
1518
1519    // Metadata publication is the commit point. Cache refresh is fail-closed:
1520    // retaining the generation's uncommitted keys may cause conservative
1521    // duplicate rejections, but can never admit a duplicate or turn a durable
1522    // commit into an API error.
1523    match refresh_primary_key_after_commit(
1524        &commit.directory,
1525        &commit.schema,
1526        &commit.segment_manager,
1527        &commit.primary_key_index,
1528    )
1529    .await
1530    {
1531        // A successful refresh folded every committed key into committed_data
1532        // and cleared the reservations — nothing retained anymore.
1533        Ok(()) => commit
1534            .pk_reservations_retained
1535            .store(false, Ordering::Release),
1536        Err(error) => {
1537            // The retained reservations are now the ONLY record of the
1538            // published segments' keys. Abort paths must not clear them
1539            // (see clear_uncommitted_pk_reservations) or duplicates would
1540            // be admitted.
1541            commit
1542                .pk_reservations_retained
1543                .store(true, Ordering::Release);
1544            log::error!(
1545                "[primary_key] committed metadata but failed to refresh dedup state; \
1546                 retaining reservations until a later successful commit: {}",
1547                error,
1548            );
1549        }
1550    }
1551
1552    // Merge scheduling is optional post-commit work and may briefly wait on
1553    // manager state. Reconcile worker availability first so it cannot extend
1554    // ingestion backpressure after metadata and PK state already agree.
1555    drop(commit.finalization.take());
1556    commit.segment_manager.maybe_merge().await;
1557    Ok(true)
1558}
1559
1560impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
1561    /// Finalize: register segments in metadata, evaluate merge policy, resume workers.
1562    ///
1563    /// Returns `true` if new segments were committed, `false` if nothing changed.
1564    pub async fn commit(mut self) -> Result<bool> {
1565        let segments = std::mem::take(&mut *self.writer.flushed_segments.lock());
1566
1567        // Fast path: nothing to commit
1568        if segments.is_empty() {
1569            log::debug!("[commit] no segments to commit, skipping");
1570            self.is_resolved = true;
1571            self.writer.resume_workers();
1572            return Ok(false);
1573        }
1574
1575        if !self.writer.commit_finalization.begin() {
1576            self.writer.flushed_segments.lock().extend(segments);
1577            // Keep the prepared generation paused. Letting `Drop` auto-abort
1578            // here would delete the retryable segments owned by another
1579            // finalization state transition.
1580            self.is_resolved = true;
1581            return Err(Error::CommitInProgress);
1582        }
1583
1584        let publication_observed = Arc::new(AtomicBool::new(false));
1585        let owned = OwnedCommitFinalization {
1586            directory: Arc::clone(&self.writer.directory),
1587            schema: Arc::clone(&self.writer.schema),
1588            segment_manager: Arc::clone(&self.writer.segment_manager),
1589            primary_key_index: Arc::clone(&self.writer.primary_key_index),
1590            prepared: PreparedSegmentsGuard {
1591                segments: Some(segments),
1592                retry_slot: Arc::clone(&self.writer.flushed_segments),
1593            },
1594            finalization: Some(CommitFinalizationGuard {
1595                state: Arc::clone(&self.writer.commit_finalization),
1596                worker_state: Arc::clone(&self.writer.worker_state),
1597                doc_sender: Arc::clone(&self.writer.doc_sender),
1598                resume_workers: false,
1599            }),
1600            publication_observed: Arc::clone(&publication_observed),
1601            pk_reservations_retained: Arc::clone(&self.writer.pk_reservations_retained),
1602        };
1603
1604        // From this point the owned value, not this cancel-sensitive guard,
1605        // controls every segment and the paused worker generation. Resolve the
1606        // local guard before spawning so even a runtime-spawn panic cannot
1607        // auto-abort the retryable generation during unwind.
1608        self.is_resolved = true;
1609        let task_publication = Arc::clone(&publication_observed);
1610        let task = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1611            tokio::spawn(async move {
1612                match std::panic::AssertUnwindSafe(finalize_prepared_commit(owned))
1613                    .catch_unwind()
1614                    .await
1615                {
1616                    Ok(result) => result,
1617                    Err(_) if task_publication.load(Ordering::Acquire) => {
1618                        log::error!(
1619                            "owned commit finalizer panicked after metadata publication; \
1620                             treating the durable generation as committed"
1621                        );
1622                        Ok(true)
1623                    }
1624                    Err(_) => Err(Error::Internal(
1625                        "owned commit finalizer panicked before metadata publication".into(),
1626                    )),
1627                }
1628            })
1629        }))
1630        .map_err(|_| Error::Internal("runtime rejected owned commit finalizer".into()))?;
1631
1632        match task.await {
1633            Ok(result) => result,
1634            Err(error) if publication_observed.load(Ordering::Acquire) => {
1635                log::error!(
1636                    "owned commit finalizer terminated after metadata publication: {}; \
1637                     treating the durable generation as committed",
1638                    error,
1639                );
1640                Ok(true)
1641            }
1642            Err(error) => Err(Error::Internal(format!(
1643                "owned commit finalizer terminated unexpectedly: {error}"
1644            ))),
1645        }
1646    }
1647
1648    /// Abort: discard prepared segments, delete their files asynchronously,
1649    /// and resume workers. Lifecycle ownership is held until deletion ends.
1650    pub fn abort(mut self) {
1651        self.is_resolved = true;
1652        self.writer.flushed_segments.lock().clear();
1653        self.writer.clear_uncommitted_pk_reservations();
1654        self.writer.resume_workers();
1655    }
1656}
1657
1658impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
1659    fn drop(&mut self) {
1660        if !self.is_resolved {
1661            log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
1662            self.writer.flushed_segments.lock().clear();
1663            self.writer.clear_uncommitted_pk_reservations();
1664            self.writer.resume_workers();
1665        }
1666    }
1667}
1668
1669/// Load only fast-field data for a segment (lightweight alternative to full SegmentReader).
1670async fn load_pk_segment_data<D: crate::directories::Directory>(
1671    dir: &D,
1672    seg_id_str: &str,
1673    schema: &Arc<crate::dsl::Schema>,
1674) -> Result<super::primary_key::PkSegmentData> {
1675    let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
1676        .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
1677    let files = crate::segment::SegmentFiles::new(seg_id.0);
1678    let fast_fields =
1679        crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
1680    Ok(super::primary_key::PkSegmentData {
1681        segment_id: seg_id_str.to_string(),
1682        fast_fields,
1683    })
1684}