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 rustc_hash::FxHashMap;
36
37use crate::directories::DirectoryWriter;
38use crate::dsl::{Document, Field, Schema};
39use crate::error::{Error, Result};
40use crate::segment::{SegmentBuilder, SegmentBuilderConfig, SegmentId};
41use crate::tokenizer::BoxedTokenizer;
42
43use super::IndexConfig;
44
45/// Total pipeline capacity (in documents).
46const PIPELINE_MAX_SIZE_IN_DOCS: usize = 10_000;
47
48/// Async IndexWriter for adding documents and committing segments.
49///
50/// **Backpressure:** `add_document()` is sync, O(1). Returns `Error::QueueFull`
51/// when the shared queue is at capacity — caller must back off.
52///
53/// **Two-phase commit:**
54/// - `prepare_commit()` → `PreparedCommit::commit()` or `PreparedCommit::abort()`
55/// - `commit()` is a convenience that does both phases.
56/// - Between prepare and commit, the caller can do external work (WAL, sync, etc.)
57///   knowing that abort is possible if something fails.
58/// - Dropping `PreparedCommit` without calling commit/abort auto-aborts.
59pub struct IndexWriter<D: DirectoryWriter + 'static> {
60    pub(super) directory: Arc<D>,
61    pub(super) schema: Arc<Schema>,
62    pub(super) config: IndexConfig,
63    /// MPMC sender — `try_send(&self)` is thread-safe, no lock needed.
64    /// Replaced on each commit cycle (workers get new receiver via resume).
65    doc_sender: async_channel::Sender<Document>,
66    /// Worker OS thread handles — long-lived, survive across commits.
67    workers: Vec<std::thread::JoinHandle<()>>,
68    /// Shared worker state (immutable config + mutable segment output + sync)
69    worker_state: Arc<WorkerState<D>>,
70    /// Segment manager — owns metadata.json, handles segments and background merging
71    pub(super) segment_manager: Arc<crate::merge::SegmentManager<D>>,
72    /// Segments flushed to disk but not yet registered in metadata. Each item
73    /// owns an active-operation guard, so orphan sweeping cannot delete it.
74    flushed_segments: Vec<PreparedSegment>,
75    /// Primary key dedup index (None if schema has no primary field)
76    primary_key_index: Option<super::primary_key::PrimaryKeyIndex>,
77}
78
79/// Shared state for worker threads.
80struct WorkerState<D: DirectoryWriter + 'static> {
81    directory: Arc<D>,
82    schema: Arc<Schema>,
83    builder_config: SegmentBuilderConfig,
84    tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
85    /// Fixed per-worker memory budget (bytes). When a builder exceeds this, segment is built.
86    memory_budget_per_worker: usize,
87    /// Segment manager — workers read trained structures from its ArcSwap (lock-free).
88    segment_manager: Arc<crate::merge::SegmentManager<D>>,
89    /// Segments built by workers, collected by `prepare_commit()`. Their RAII
90    /// guards protect both in-progress and completed-uncommitted files.
91    built_segments: parking_lot::Mutex<Vec<PreparedSegment>>,
92
93    // === Worker lifecycle synchronization ===
94    // Workers survive across commits. On prepare_commit the channel is closed;
95    // workers flush their builders, increment flush_count, then wait on
96    // resume_cvar for a new receiver. commit/abort creates a fresh channel
97    // and wakes them.
98    /// Number of workers that have completed their flush.
99    flush_count: AtomicUsize,
100    /// Mutex + condvar for prepare_commit to wait on all workers flushed.
101    flush_mutex: parking_lot::Mutex<()>,
102    flush_cvar: parking_lot::Condvar,
103    /// Holds the new channel receiver after commit/abort. Workers clone from this.
104    resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
105    /// Monotonically increasing epoch, bumped by each resume_workers call.
106    /// Workers compare against their local epoch to avoid re-cloning a stale receiver.
107    resume_epoch: AtomicUsize,
108    /// Condvar for workers to wait for resume (new channel) or shutdown.
109    resume_cvar: parking_lot::Condvar,
110    /// When true, workers should exit permanently (IndexWriter dropped).
111    shutdown: AtomicBool,
112    /// Total number of worker threads.
113    num_workers: usize,
114}
115
116/// A completed indexing segment that has not been published in metadata yet.
117///
118/// `_operation` is intentionally data, not a side-channel set update: moving
119/// this value through worker → prepared commit → commit/abort moves lifecycle
120/// ownership with it, and every unwind/drop path releases ownership safely.
121struct PreparedSegment {
122    id: String,
123    num_docs: u32,
124    _operation: crate::merge::SegmentOperationGuard,
125}
126
127impl PreparedSegment {
128    fn metadata_entry(&self) -> (String, u32) {
129        (self.id.clone(), self.num_docs)
130    }
131}
132
133impl<D: DirectoryWriter + 'static> IndexWriter<D> {
134    /// Create a new index in the directory
135    pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
136        Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
137    }
138
139    /// Create a new index with custom builder config
140    pub async fn create_with_config(
141        directory: D,
142        schema: Schema,
143        config: IndexConfig,
144        builder_config: SegmentBuilderConfig,
145    ) -> Result<Self> {
146        let directory = Arc::new(directory);
147        let schema = Arc::new(schema);
148        // Directory-layer metrics (cold writes, lazy reads) carry the index label
149        directory.set_index_label(schema.index_label());
150        let metadata = super::IndexMetadata::new((*schema).clone());
151
152        let segment_manager = Arc::new(crate::merge::SegmentManager::new(
153            Arc::clone(&directory),
154            Arc::clone(&schema),
155            metadata,
156            config.merge_policy.clone_box(),
157            config.term_cache_blocks,
158            config.max_concurrent_merges,
159            config.merge_bp_time_budget,
160            config.bp_memory_budget_bytes,
161        ));
162        segment_manager.update_metadata(|_| {}).await?;
163
164        Ok(Self::new_with_parts(
165            directory,
166            schema,
167            config,
168            builder_config,
169            segment_manager,
170        ))
171    }
172
173    /// Open an existing index for writing
174    pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
175        Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
176    }
177
178    /// Open an existing index with custom builder config
179    pub async fn open_with_config(
180        directory: D,
181        config: IndexConfig,
182        builder_config: SegmentBuilderConfig,
183    ) -> Result<Self> {
184        let directory = Arc::new(directory);
185        let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
186        let schema = Arc::new(metadata.schema.clone());
187        // Directory-layer metrics (cold writes, lazy reads) carry the index label
188        directory.set_index_label(schema.index_label());
189
190        let segment_manager = Arc::new(crate::merge::SegmentManager::new(
191            Arc::clone(&directory),
192            Arc::clone(&schema),
193            metadata,
194            config.merge_policy.clone_box(),
195            config.term_cache_blocks,
196            config.max_concurrent_merges,
197            config.merge_bp_time_budget,
198            config.bp_memory_budget_bytes,
199        ));
200        segment_manager.load_and_publish_trained().await;
201
202        Ok(Self::new_with_parts(
203            directory,
204            schema,
205            config,
206            builder_config,
207            segment_manager,
208        ))
209    }
210
211    /// Create an IndexWriter from an existing Index.
212    /// Shares the SegmentManager for consistent segment lifecycle management.
213    pub fn from_index(index: &super::Index<D>) -> Self {
214        Self::new_with_parts(
215            Arc::clone(&index.directory),
216            Arc::clone(&index.schema),
217            index.config.clone(),
218            SegmentBuilderConfig::default(),
219            Arc::clone(&index.segment_manager),
220        )
221    }
222
223    // ========================================================================
224    // Construction + pipeline management
225    // ========================================================================
226
227    /// Common construction: creates worker state, spawns workers, assembles `Self`.
228    fn new_with_parts(
229        directory: Arc<D>,
230        schema: Arc<Schema>,
231        config: IndexConfig,
232        builder_config: SegmentBuilderConfig,
233        segment_manager: Arc<crate::merge::SegmentManager<D>>,
234    ) -> Self {
235        // Auto-configure tokenizers from schema for all text fields
236        let registry = crate::tokenizer::TokenizerRegistry::new();
237        let mut tokenizers = FxHashMap::default();
238        for (field, entry) in schema.fields() {
239            if matches!(entry.field_type, crate::dsl::FieldType::Text)
240                && let Some(ref tok_name) = entry.tokenizer
241                && let Some(tok) = registry.get(tok_name)
242            {
243                tokenizers.insert(field, tok);
244            }
245        }
246
247        let num_workers = config.num_indexing_threads.max(1);
248        let worker_state = Arc::new(WorkerState {
249            directory: Arc::clone(&directory),
250            schema: Arc::clone(&schema),
251            builder_config,
252            tokenizers: parking_lot::RwLock::new(tokenizers),
253            memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
254            segment_manager: Arc::clone(&segment_manager),
255            built_segments: parking_lot::Mutex::new(Vec::new()),
256            flush_count: AtomicUsize::new(0),
257            flush_mutex: parking_lot::Mutex::new(()),
258            flush_cvar: parking_lot::Condvar::new(),
259            resume_receiver: parking_lot::Mutex::new(None),
260            resume_epoch: AtomicUsize::new(0),
261            resume_cvar: parking_lot::Condvar::new(),
262            shutdown: AtomicBool::new(false),
263            num_workers,
264        });
265        let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
266
267        Self {
268            directory,
269            schema,
270            config,
271            doc_sender,
272            workers,
273            worker_state,
274            segment_manager,
275            flushed_segments: Vec::new(),
276            primary_key_index: None,
277        }
278    }
279
280    fn spawn_workers(
281        worker_state: &Arc<WorkerState<D>>,
282        num_workers: usize,
283    ) -> (
284        async_channel::Sender<Document>,
285        Vec<std::thread::JoinHandle<()>>,
286    ) {
287        let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
288        let handle = tokio::runtime::Handle::current();
289        let mut workers = Vec::with_capacity(num_workers);
290        for i in 0..num_workers {
291            let state = Arc::clone(worker_state);
292            let rx = receiver.clone();
293            let rt = handle.clone();
294            workers.push(
295                std::thread::Builder::new()
296                    .name(format!("index-worker-{}", i))
297                    .spawn(move || Self::worker_loop(state, rx, rt))
298                    .expect("failed to spawn index worker thread"),
299            );
300        }
301        (sender, workers)
302    }
303
304    /// Get the schema
305    pub fn schema(&self) -> &Schema {
306        &self.schema
307    }
308
309    /// Set tokenizer for a field.
310    /// Propagated to worker threads — takes effect for the next SegmentBuilder they create.
311    pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
312        self.worker_state
313            .tokenizers
314            .write()
315            .insert(field, Box::new(tokenizer));
316    }
317
318    /// Initialize primary key deduplication from committed segments.
319    ///
320    /// Tries to load a cached bloom filter from `pk_bloom.bin` first. If the
321    /// cache covers all current segments, the bloom is reused directly (fast
322    /// path). If new segments appeared since the cache was written, only their
323    /// keys are iterated (incremental). Falls back to a full rebuild when no
324    /// cache exists.
325    ///
326    /// Only loads fast-field data (text dictionaries) per segment — NOT full
327    /// `SegmentReader`s — to avoid duplicating dense/sparse index memory.
328    ///
329    /// The CPU-intensive bloom build is offloaded via `spawn_blocking` so it
330    /// does not block the tokio runtime.
331    ///
332    /// No-op if schema has no primary field.
333    pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
334        use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
335
336        let field = match self.schema.primary_field() {
337            Some(f) => f,
338            None => return Ok(()),
339        };
340
341        let snapshot = self.segment_manager.acquire_snapshot().await;
342        let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
343
344        // Try to load persisted bloom filter.
345        let cached = match self
346            .directory
347            .open_read(std::path::Path::new(PK_BLOOM_FILE))
348            .await
349        {
350            Ok(handle) => {
351                let data = handle.read_bytes_range(0..handle.len()).await;
352                match data {
353                    Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
354                    Err(_) => None,
355                }
356            }
357            Err(_) => None,
358        };
359
360        // Load lightweight fast-field data for all segments concurrently.
361        let load_futures: Vec<_> = current_seg_ids
362            .iter()
363            .map(|seg_id_str| {
364                let seg_id_str = seg_id_str.clone();
365                let dir = self.directory.as_ref();
366                let schema = Arc::clone(&self.schema);
367                async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
368            })
369            .collect();
370        let all_data = futures::future::try_join_all(load_futures).await?;
371
372        if let Some((persisted_seg_ids, bloom)) = cached {
373            // Partition: old segments (covered by bloom) first, new segments at end.
374            let mut pk_data = Vec::with_capacity(all_data.len());
375            let mut new_data = Vec::new();
376            for d in all_data {
377                if persisted_seg_ids.contains(&d.segment_id) {
378                    pk_data.push(d);
379                } else {
380                    new_data.push(d);
381                }
382            }
383            let needs_persist = !new_data.is_empty();
384            let new_start = pk_data.len();
385            pk_data.extend(new_data);
386
387            let pk_index = if new_start == pk_data.len() {
388                // Fast path: all segments covered by cache.
389                super::primary_key::PrimaryKeyIndex::from_persisted(
390                    field,
391                    bloom,
392                    pk_data,
393                    &[],
394                    snapshot,
395                )
396            } else {
397                // Incremental: only iterate new segments' keys.
398                tokio::task::spawn_blocking(move || {
399                    // Insert new segments' keys into the bloom, then construct
400                    // PrimaryKeyIndex with the pre-populated bloom.
401                    let mut bloom = bloom;
402                    let mut added = 0usize;
403                    let num_new = pk_data.len() - new_start;
404                    for data in &pk_data[new_start..] {
405                        if let Some(ff) = data.fast_fields.get(&field.0)
406                            && let Some(dict) = ff.text_dict()
407                        {
408                            for key in dict.iter() {
409                                bloom.insert(key.as_bytes());
410                                added += 1;
411                            }
412                        }
413                    }
414                    if added > 0 {
415                        log::info!(
416                            "[primary_key] bloom: added {} keys from {} new segment(s)",
417                            added,
418                            num_new,
419                        );
420                    }
421                    super::primary_key::PrimaryKeyIndex::from_persisted(
422                        field,
423                        bloom,
424                        pk_data,
425                        &[],
426                        snapshot,
427                    )
428                })
429                .await
430                .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
431            };
432
433            if needs_persist {
434                self.persist_pk_bloom(&pk_index, &current_seg_ids).await;
435            }
436
437            self.primary_key_index = Some(pk_index);
438        } else {
439            // No cache — full rebuild, offloaded to blocking thread.
440            let pk_index = tokio::task::spawn_blocking(move || {
441                super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
442            })
443            .await
444            .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
445
446            self.persist_pk_bloom(&pk_index, &current_seg_ids).await;
447            self.primary_key_index = Some(pk_index);
448        }
449
450        Ok(())
451    }
452
453    /// Persist the primary-key bloom filter to `pk_bloom.bin`.
454    /// Best-effort: errors are logged but not propagated.
455    async fn persist_pk_bloom(
456        &self,
457        pk_index: &super::primary_key::PrimaryKeyIndex,
458        segment_ids: &[String],
459    ) {
460        use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
461
462        let bloom_bytes = pk_index.bloom_to_bytes();
463        let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
464        if let Err(e) = self
465            .directory
466            .write(std::path::Path::new(PK_BLOOM_FILE), &data)
467            .await
468        {
469            log::warn!("[primary_key] failed to persist bloom cache: {}", e);
470        }
471    }
472
473    /// Add a document to the indexing queue (sync, O(1), lock-free).
474    ///
475    /// `Document` is moved into the channel (zero-copy). Workers compete to pull it.
476    /// Returns `Error::QueueFull` when the queue is at capacity — caller must back off.
477    pub fn add_document(&self, doc: Document) -> Result<()> {
478        if let Some(ref pk_index) = self.primary_key_index {
479            pk_index.check_and_insert(&doc)?;
480        }
481        match self.doc_sender.try_send(doc) {
482            Ok(()) => Ok(()),
483            Err(async_channel::TrySendError::Full(doc)) => {
484                // Roll back PK registration so the caller can retry later
485                if let Some(ref pk_index) = self.primary_key_index {
486                    pk_index.rollback_uncommitted_key(&doc);
487                }
488                Err(Error::QueueFull)
489            }
490            Err(async_channel::TrySendError::Closed(doc)) => {
491                // Roll back PK registration for defense-in-depth
492                if let Some(ref pk_index) = self.primary_key_index {
493                    pk_index.rollback_uncommitted_key(&doc);
494                }
495                Err(Error::Internal("Document channel closed".into()))
496            }
497        }
498    }
499
500    /// Add multiple documents to the indexing queue.
501    ///
502    /// Returns the number of documents successfully queued. Stops at the first
503    /// `QueueFull` and returns the count queued so far.
504    pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
505        let total = documents.len();
506        for (i, doc) in documents.into_iter().enumerate() {
507            match self.add_document(doc) {
508                Ok(()) => {}
509                Err(Error::QueueFull) => return Ok(i),
510                Err(e) => return Err(e),
511            }
512        }
513        Ok(total)
514    }
515
516    // ========================================================================
517    // Worker loop
518    // ========================================================================
519
520    /// Worker loop — runs on a dedicated OS thread, survives across commits.
521    ///
522    /// Outer loop: each iteration processes one commit cycle.
523    ///   Inner loop: pull documents from MPMC queue, index them, build segments
524    ///   when memory budget is exceeded.
525    ///   On channel close (prepare_commit): flush current builder, signal
526    ///   flush_count, wait for resume with new receiver.
527    ///   On shutdown (Drop): exit permanently.
528    fn worker_loop(
529        state: Arc<WorkerState<D>>,
530        initial_receiver: async_channel::Receiver<Document>,
531        handle: tokio::runtime::Handle,
532    ) {
533        let mut receiver = initial_receiver;
534        let mut my_epoch = 0usize;
535
536        loop {
537            // Wrap the recv+build phase in catch_unwind so a panic doesn't
538            // prevent flush_count from being signaled (which would hang
539            // prepare_commit forever).
540            let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
541                let mut builder: Option<SegmentBuilder> = None;
542
543                while let Ok(doc) = receiver.recv_blocking() {
544                    // Initialize builder if needed
545                    if builder.is_none() {
546                        match SegmentBuilder::new(
547                            Arc::clone(&state.schema),
548                            state.builder_config.clone(),
549                        ) {
550                            Ok(mut b) => {
551                                for (field, tokenizer) in state.tokenizers.read().iter() {
552                                    b.set_tokenizer(*field, tokenizer.clone_box());
553                                }
554                                builder = Some(b);
555                            }
556                            Err(e) => {
557                                log::error!("Failed to create segment builder: {:?}", e);
558                                continue;
559                            }
560                        }
561                    }
562
563                    let b = builder.as_mut().unwrap();
564                    if let Err(e) = b.add_document(doc) {
565                        log::error!("Failed to index document: {:?}", e);
566                        continue;
567                    }
568
569                    let builder_memory = b.estimated_memory_bytes();
570
571                    if b.num_docs() & 0x3FFF == 0 {
572                        log::debug!(
573                            "[indexing] docs={}, memory={:.2} MB, budget={:.2} MB",
574                            b.num_docs(),
575                            builder_memory as f64 / (1024.0 * 1024.0),
576                            state.memory_budget_per_worker as f64 / (1024.0 * 1024.0)
577                        );
578                    }
579
580                    // Require minimum 100 docs before flushing to avoid tiny segments
581                    const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
582
583                    // Reserve 20% headroom for segment build overhead (vid_set,
584                    // VidLookup, postings_flat, grid_entries). These temporary
585                    // allocations exist alongside the builder's data during build.
586                    let effective_budget = state.memory_budget_per_worker * 4 / 5;
587
588                    if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
589                        log::info!(
590                            "[indexing] memory budget reached, building segment: \
591                             docs={}, memory={:.2} MB, budget={:.2} MB",
592                            b.num_docs(),
593                            builder_memory as f64 / (1024.0 * 1024.0),
594                            state.memory_budget_per_worker as f64 / (1024.0 * 1024.0),
595                        );
596                        let full_builder = builder.take().unwrap();
597                        Self::build_segment_inline(&state, full_builder, &handle);
598                    }
599                }
600
601                // Channel closed — flush current builder
602                if let Some(b) = builder.take()
603                    && b.num_docs() > 0
604                {
605                    Self::build_segment_inline(&state, b, &handle);
606                }
607            }));
608
609            if build_result.is_err() {
610                log::error!(
611                    "[worker] panic during indexing cycle — documents in this cycle may be lost"
612                );
613            }
614
615            // Signal flush completion (always, even after panic — prevents
616            // prepare_commit from hanging)
617            let prev = state.flush_count.fetch_add(1, Ordering::Release);
618            if prev + 1 == state.num_workers {
619                // Last worker — wake prepare_commit
620                let _lock = state.flush_mutex.lock();
621                state.flush_cvar.notify_one();
622            }
623
624            // Wait for resume (new channel) or shutdown.
625            // Check resume_epoch to avoid re-cloning a stale receiver from
626            // a previous cycle.
627            {
628                let mut lock = state.resume_receiver.lock();
629                loop {
630                    if state.shutdown.load(Ordering::Acquire) {
631                        return;
632                    }
633                    let current_epoch = state.resume_epoch.load(Ordering::Acquire);
634                    if current_epoch > my_epoch
635                        && let Some(rx) = lock.as_ref()
636                    {
637                        receiver = rx.clone();
638                        my_epoch = current_epoch;
639                        break;
640                    }
641                    state.resume_cvar.wait(&mut lock);
642                }
643            }
644        }
645    }
646
647    /// Build a segment on the worker thread. Uses `Handle::block_on()` to bridge
648    /// into async context for I/O (streaming writers). CPU work (rayon) stays on
649    /// the worker thread / rayon pool.
650    fn build_segment_inline(
651        state: &WorkerState<D>,
652        builder: SegmentBuilder,
653        handle: &tokio::runtime::Handle,
654    ) {
655        let segment_id = SegmentId::new();
656        let segment_hex = segment_id.to_hex();
657        // Claim the ID before the first file write. The guard is moved into
658        // `PreparedSegment` on success and otherwise releases automatically.
659        let operation = match state
660            .segment_manager
661            .protect_new_segment(segment_hex.clone())
662        {
663            Ok(operation) => operation,
664            Err(e) => {
665                log::error!(
666                    "[segment_build_failed] segment_id={} lifecycle_error={}",
667                    segment_hex,
668                    e,
669                );
670                return;
671            }
672        };
673        let trained = state.segment_manager.trained();
674        let doc_count = builder.num_docs();
675        let build_start = std::time::Instant::now();
676
677        log::info!(
678            "[segment_build] segment_id={} doc_count={} ann={}",
679            segment_hex,
680            doc_count,
681            trained.is_some()
682        );
683
684        match handle.block_on(builder.build(
685            state.directory.as_ref(),
686            segment_id,
687            trained.as_deref(),
688        )) {
689            Ok(meta) if meta.num_docs > 0 => {
690                let duration_ms = build_start.elapsed().as_millis() as u64;
691                log::info!(
692                    "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
693                    segment_hex,
694                    meta.num_docs,
695                    duration_ms,
696                );
697                state.built_segments.lock().push(PreparedSegment {
698                    id: segment_hex,
699                    num_docs: meta.num_docs,
700                    _operation: operation,
701                });
702            }
703            Ok(_) => {}
704            Err(e) => {
705                log::error!(
706                    "[segment_build_failed] segment_id={} error={:?}",
707                    segment_hex,
708                    e
709                );
710                if let Err(delete_error) = handle.block_on(crate::segment::delete_segment(
711                    state.directory.as_ref(),
712                    segment_id,
713                )) {
714                    log::warn!(
715                        "[segment_cleanup] failed deleting partial indexing segment {}: {}",
716                        segment_hex,
717                        delete_error,
718                    );
719                }
720            }
721        }
722    }
723
724    // ========================================================================
725    // Public API — commit, merge, etc.
726    // ========================================================================
727
728    /// Check merge policy and spawn a background merge if needed.
729    pub async fn maybe_merge(&self) {
730        self.segment_manager.maybe_merge().await;
731    }
732
733    /// Abort all in-flight merge tasks without waiting for completion.
734    pub async fn abort_merges(&self) {
735        self.segment_manager.abort_merges().await;
736    }
737
738    /// Wait for the in-flight background merge to complete (if any).
739    pub async fn wait_for_merging_thread(&self) {
740        self.segment_manager.wait_for_merging_thread().await;
741    }
742
743    /// Wait for all eligible merges to complete, including cascading merges.
744    pub async fn wait_for_all_merges(&self) {
745        self.segment_manager.wait_for_all_merges().await;
746    }
747
748    /// Get the segment tracker for sharing with readers.
749    pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
750        self.segment_manager.tracker()
751    }
752
753    /// Acquire a snapshot of current segments for reading.
754    pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
755        self.segment_manager.acquire_snapshot().await
756    }
757
758    /// Clean up orphan segment files not registered in metadata.
759    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
760        self.segment_manager.cleanup_orphan_segments().await
761    }
762
763    /// Prepare commit — signal workers to flush, wait for completion, collect segments.
764    ///
765    /// All documents sent via `add_document` before this call are guaranteed
766    /// to be written to segment files on disk. Segments are NOT yet registered
767    /// in metadata — call `PreparedCommit::commit()` for that.
768    ///
769    /// Workers are NOT destroyed — they flush their builders and wait for
770    /// `resume_workers()` to give them a new channel.
771    ///
772    /// `add_document` will return `Closed` error until commit/abort resumes workers.
773    pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
774        // 1. Close channel → workers drain remaining docs and flush builders
775        self.doc_sender.close();
776
777        // Wake any workers still waiting on resume_cvar from previous cycle.
778        // They'll clone the stale receiver, enter recv_blocking, get Err
779        // immediately (sender already closed), flush, and signal completion.
780        self.worker_state.resume_cvar.notify_all();
781
782        // 2. Wait for all workers to complete their flush (via spawn_blocking
783        //    to avoid blocking the tokio runtime)
784        let state = Arc::clone(&self.worker_state);
785        let all_flushed = tokio::task::spawn_blocking(move || {
786            let mut lock = state.flush_mutex.lock();
787            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
788            while state.flush_count.load(Ordering::Acquire) < state.num_workers {
789                let remaining = deadline.saturating_duration_since(std::time::Instant::now());
790                if remaining.is_zero() {
791                    log::error!(
792                        "[prepare_commit] timed out waiting for workers: {}/{} flushed",
793                        state.flush_count.load(Ordering::Acquire),
794                        state.num_workers
795                    );
796                    return false;
797                }
798                state.flush_cvar.wait_for(&mut lock, remaining);
799            }
800            true
801        })
802        .await
803        .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
804
805        if !all_flushed {
806            // Resume workers so the system isn't stuck, then return error
807            self.resume_workers();
808            return Err(Error::Internal(format!(
809                "prepare_commit timed out: {}/{} workers flushed",
810                self.worker_state.flush_count.load(Ordering::Acquire),
811                self.worker_state.num_workers
812            )));
813        }
814
815        // 3. Collect built segments
816        let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
817        self.flushed_segments.extend(built);
818
819        Ok(PreparedCommit {
820            writer: self,
821            is_resolved: false,
822            is_published: false,
823        })
824    }
825
826    /// Commit (convenience): prepare_commit + commit in one call.
827    ///
828    /// Guarantees all prior `add_document` calls are committed.
829    /// Vector training is decoupled — call `build_vector_index()` manually.
830    pub async fn commit(&mut self) -> Result<bool> {
831        self.prepare_commit().await?.commit().await
832    }
833
834    /// Force merge all segments into one.
835    pub async fn force_merge(&mut self) -> Result<()> {
836        self.prepare_commit().await?.commit().await?;
837        self.segment_manager.force_merge().await
838    }
839
840    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
841    ///
842    /// Each segment is individually rebuilt with record-level BP reordering:
843    /// ordinals are shuffled across blocks so that similar content clusters tightly.
844    pub async fn reorder(&mut self) -> Result<()> {
845        self.prepare_commit().await?.commit().await?;
846        self.segment_manager.reorder_segments().await
847    }
848
849    /// Get the segment manager (for background optimizer access).
850    pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
851        &self.segment_manager
852    }
853
854    /// Resume workers with a fresh channel. Called after commit or abort.
855    ///
856    /// Workers are already alive — just give them a new channel and wake them.
857    /// If the tokio runtime has shut down (e.g., program exit), this is a no-op.
858    fn resume_workers(&mut self) {
859        if tokio::runtime::Handle::try_current().is_err() {
860            // Runtime is gone — signal permanent shutdown so workers don't
861            // hang forever on resume_cvar.
862            self.worker_state.shutdown.store(true, Ordering::Release);
863            self.worker_state.resume_cvar.notify_all();
864            return;
865        }
866
867        // Reset flush count for next cycle
868        self.worker_state.flush_count.store(0, Ordering::Release);
869
870        // Create new channel
871        let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
872        self.doc_sender = sender;
873
874        // Set new receiver, bump epoch, and wake all workers
875        {
876            let mut lock = self.worker_state.resume_receiver.lock();
877            *lock = Some(receiver);
878        }
879        self.worker_state
880            .resume_epoch
881            .fetch_add(1, Ordering::Release);
882        self.worker_state.resume_cvar.notify_all();
883    }
884
885    // Vector index methods (build_vector_index, etc.) are in vector_builder.rs
886}
887
888impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
889    fn drop(&mut self) {
890        // 1. Signal permanent shutdown
891        self.worker_state.shutdown.store(true, Ordering::Release);
892        // 2. Close channel to wake workers blocked on recv_blocking
893        self.doc_sender.close();
894        // 3. Wake workers that might be waiting on resume_cvar
895        self.worker_state.resume_cvar.notify_all();
896        // 4. Join worker threads
897        for w in std::mem::take(&mut self.workers) {
898            let _ = w.join();
899        }
900    }
901}
902
903/// A prepared commit that can be finalized or aborted.
904///
905/// Two-phase commit guard. Between `prepare_commit()` and
906/// `commit()`/`abort()`, segments are on disk but NOT in metadata.
907/// Dropping without calling either will auto-abort (discard segments,
908/// respawn workers).
909pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
910    writer: &'a mut IndexWriter<D>,
911    is_resolved: bool,
912    /// Metadata publication is the commit point. Post-publication cache work
913    /// may still fail or be cancelled, but must never be described as an abort.
914    is_published: bool,
915}
916
917impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
918    /// Finalize: register segments in metadata, evaluate merge policy, resume workers.
919    ///
920    /// Returns `true` if new segments were committed, `false` if nothing changed.
921    pub async fn commit(mut self) -> Result<bool> {
922        let segments = std::mem::take(&mut self.writer.flushed_segments);
923
924        // Fast path: nothing to commit
925        if segments.is_empty() {
926            log::debug!("[commit] no segments to commit, skipping");
927            self.is_resolved = true;
928            self.writer.resume_workers();
929            return Ok(false);
930        }
931
932        let metadata_entries: Vec<(String, u32)> = segments
933            .iter()
934            .map(PreparedSegment::metadata_entry)
935            .collect();
936        if let Err(error) = self.writer.segment_manager.commit(&metadata_entries).await {
937            // Publication is transactional, so these guarded files remain a
938            // valid prepared commit. Preserve them for the caller's next
939            // commit instead of leaking ownership or wedging the workers.
940            self.writer.flushed_segments = segments;
941            self.is_resolved = true;
942            self.writer.resume_workers();
943            return Err(error);
944        }
945        self.is_published = true;
946
947        // Metadata + tracker now own the IDs; releasing indexing ownership is
948        // the final transition from prepared to live.
949        drop(segments);
950
951        // Refresh primary key index: only load fast fields for NEW segments.
952        let post_publish_result = async {
953            if let Some(ref mut pk_index) = self.writer.primary_key_index {
954                let snapshot = self.writer.segment_manager.acquire_snapshot().await;
955                let existing_ids: std::collections::HashSet<&str> =
956                    pk_index.committed_segment_ids().collect();
957
958                // Only load fast fields for segments not already held.
959                let load_futures: Vec<_> = snapshot
960                    .segment_ids()
961                    .iter()
962                    .filter(|id| !existing_ids.contains(id.as_str()))
963                    .map(|seg_id_str| {
964                        let seg_id_str = seg_id_str.clone();
965                        let dir = self.writer.directory.as_ref();
966                        let schema = Arc::clone(&self.writer.schema);
967                        async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
968                    })
969                    .collect();
970                let new_data = futures::future::try_join_all(load_futures).await?;
971
972                let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
973                pk_index.refresh_incremental(new_data, snapshot);
974
975                // Persist bloom cache (extract bytes to avoid borrow conflict).
976                let bloom_bytes = pk_index.bloom_to_bytes();
977                let data = super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes);
978                if let Err(e) = self
979                    .writer
980                    .directory
981                    .write(
982                        std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
983                        &data,
984                    )
985                    .await
986                {
987                    log::warn!("[primary_key] failed to persist bloom cache: {}", e);
988                }
989            }
990
991            self.writer.segment_manager.maybe_merge().await;
992            Ok(())
993        }
994        .await;
995
996        // Worker availability must not depend on optional post-publication
997        // cache refresh or merge scheduling.
998        self.is_resolved = true;
999        self.writer.resume_workers();
1000        post_publish_result.map(|()| true)
1001    }
1002
1003    /// Abort: discard prepared segments, resume workers.
1004    /// Segment files become orphans (cleaned up by `cleanup_orphan_segments`).
1005    pub fn abort(mut self) {
1006        self.is_resolved = true;
1007        self.writer.flushed_segments.clear();
1008        if let Some(ref mut pk_index) = self.writer.primary_key_index {
1009            pk_index.clear_uncommitted();
1010        }
1011        self.writer.resume_workers();
1012    }
1013}
1014
1015impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
1016    fn drop(&mut self) {
1017        if !self.is_resolved {
1018            if self.is_published {
1019                log::warn!("PreparedCommit dropped after metadata publication — resuming workers");
1020            } else {
1021                log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
1022                self.writer.flushed_segments.clear();
1023                if let Some(ref mut pk_index) = self.writer.primary_key_index {
1024                    pk_index.clear_uncommitted();
1025                }
1026            }
1027            self.writer.resume_workers();
1028        }
1029    }
1030}
1031
1032/// Load only fast-field data for a segment (lightweight alternative to full SegmentReader).
1033async fn load_pk_segment_data<D: crate::directories::Directory>(
1034    dir: &D,
1035    seg_id_str: &str,
1036    schema: &Arc<crate::dsl::Schema>,
1037) -> Result<super::primary_key::PkSegmentData> {
1038    let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
1039        .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
1040    let files = crate::segment::SegmentFiles::new(seg_id.0);
1041    let fast_fields =
1042        crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
1043    Ok(super::primary_key::PkSegmentData {
1044        segment_id: seg_id_str.to_string(),
1045        fast_fields,
1046    })
1047}