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