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