1use 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
46const PIPELINE_MAX_SIZE_IN_DOCS: usize = 10_000;
48
49pub const WRITER_LOCK_FILENAME: &str = ".hermes_writer.lock";
51
52enum WriterLock {
61 Held { _file: std::fs::File },
63 NotApplicable,
66 Unavailable { reason: String },
69}
70
71fn 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 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
94fn 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
124pub struct IndexWriter<D: DirectoryWriter + 'static> {
138 pub(super) directory: Arc<D>,
139 pub(super) schema: Arc<Schema>,
140 pub(super) config: IndexConfig,
141 doc_sender: Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
144 workers: Vec<std::thread::JoinHandle<()>>,
146 worker_state: Arc<WorkerState<D>>,
148 pub(super) segment_manager: Arc<crate::merge::SegmentManager<D>>,
150 flushed_segments: Arc<parking_lot::Mutex<Vec<PreparedSegment<D>>>>,
153 primary_key_index: Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
155 commit_finalization: Arc<CommitFinalizationState>,
159 pk_reservations_retained: Arc<AtomicBool>,
164 writer_lock: parking_lot::RwLock<WriterLock>,
169}
170
171#[derive(Default)]
172struct CommitFinalizationState {
173 in_progress: AtomicBool,
174 idle: tokio::sync::Notify,
175}
176
177impl CommitFinalizationState {
178 fn begin(&self) -> bool {
179 self.in_progress
180 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
181 .is_ok()
182 }
183
184 fn finish(&self) {
185 self.in_progress.store(false, Ordering::Release);
186 self.idle.notify_waiters();
187 }
188
189 async fn wait_until_idle(&self) {
190 while self.in_progress.load(Ordering::Acquire) {
191 let notified = self.idle.notified();
192 if !self.in_progress.load(Ordering::Acquire) {
193 break;
194 }
195 notified.await;
196 }
197 }
198}
199
200struct WorkerState<D: DirectoryWriter + 'static> {
202 directory: Arc<D>,
203 schema: Arc<Schema>,
204 builder_config: SegmentBuilderConfig,
205 tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
206 memory_budget_per_worker: usize,
208 segment_manager: Arc<crate::merge::SegmentManager<D>>,
210 built_segments: parking_lot::Mutex<Vec<PreparedSegment<D>>>,
213 cycle_error: parking_lot::Mutex<Option<String>>,
218 cycle_failed: AtomicBool,
219
220 flush_count: AtomicUsize,
227 flush_mutex: parking_lot::Mutex<()>,
229 flush_cvar: parking_lot::Condvar,
230 resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
232 resume_epoch: AtomicUsize,
235 resume_cvar: parking_lot::Condvar,
237 shutdown: AtomicBool,
239 num_workers: usize,
241}
242
243struct PreparedSegment<D: DirectoryWriter + 'static> {
249 id: String,
250 segment_id: SegmentId,
251 num_docs: u32,
252 segment_manager: Arc<crate::merge::SegmentManager<D>>,
253 operation: Option<crate::merge::SegmentOperationGuard>,
254 runtime: tokio::runtime::Handle,
255 needs_vector_upgrade: bool,
256 published: bool,
257}
258
259impl<D: DirectoryWriter + 'static> PreparedSegment<D> {
260 fn metadata_entry(&self) -> (String, u32) {
261 (self.id.clone(), self.num_docs)
262 }
263
264 fn mark_published(&mut self) {
265 self.published = true;
266 drop(self.operation.take());
268 }
269}
270
271impl<D: DirectoryWriter + 'static> WorkerState<D> {
272 fn record_cycle_error(&self, error: impl Into<String>) {
273 let mut first_error = self.cycle_error.lock();
274 if first_error.is_none() {
275 *first_error = Some(error.into());
276 }
277 drop(first_error);
278 self.cycle_failed.store(true, Ordering::Release);
279 }
280}
281
282impl<D: DirectoryWriter + 'static> Drop for PreparedSegment<D> {
283 fn drop(&mut self) {
284 if self.published {
285 return;
286 }
287 let Some(operation) = self.operation.take() else {
288 return;
289 };
290 self.segment_manager.schedule_unpublished_segment_cleanup(
291 self.segment_id,
292 operation,
293 self.runtime.clone(),
294 );
295 }
296}
297
298impl<D: DirectoryWriter + 'static> IndexWriter<D> {
299 pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
301 Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
302 }
303
304 pub async fn create_with_config(
306 directory: D,
307 schema: Schema,
308 config: IndexConfig,
309 builder_config: SegmentBuilderConfig,
310 ) -> Result<Self> {
311 crate::dsl::reject_removed_vector_index_types(&schema).map_err(Error::Schema)?;
312 let directory = Arc::new(directory);
313 let schema = Arc::new(schema);
314 directory.set_index_label(schema.index_label());
316
317 let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
319 if let WriterLock::Unavailable { reason } = &writer_lock {
320 return Err(Error::Internal(reason.clone()));
321 }
322 if directory
326 .exists(std::path::Path::new(super::INDEX_META_FILENAME))
327 .await?
328 {
329 return Err(Error::Internal(format!(
330 "refusing to create index: {} already exists in this directory; \
331 use IndexWriter::open to open the existing index, or delete the \
332 directory first if you really want to start over",
333 super::INDEX_META_FILENAME
334 )));
335 }
336
337 let metadata = super::IndexMetadata::new((*schema).clone());
338
339 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
340 Arc::clone(&directory),
341 Arc::clone(&schema),
342 metadata,
343 config.merge_policy.clone_box(),
344 config.term_cache_blocks,
345 config.max_concurrent_merges,
346 Arc::clone(&config.background_merge_permits),
347 config.merge_bp_time_budget,
348 config.bp_memory_budget_bytes,
349 Arc::clone(&config.background_reorder_permits),
350 config.background_reorder_pool.clone(),
351 ));
352 segment_manager.update_metadata(|_| {}).await?;
353
354 Ok(Self::new_with_parts(
355 directory,
356 schema,
357 config,
358 builder_config,
359 segment_manager,
360 writer_lock,
361 ))
362 }
363
364 pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
373 Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
374 }
375
376 pub async fn open_with_config(
378 directory: D,
379 config: IndexConfig,
380 builder_config: SegmentBuilderConfig,
381 ) -> Result<Self> {
382 let directory = Arc::new(directory);
383
384 let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
387 if let WriterLock::Unavailable { reason } = &writer_lock {
388 return Err(Error::Internal(reason.clone()));
389 }
390
391 let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
392 let schema = Arc::new(metadata.schema.clone());
393 directory.set_index_label(schema.index_label());
395
396 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
397 Arc::clone(&directory),
398 Arc::clone(&schema),
399 metadata,
400 config.merge_policy.clone_box(),
401 config.term_cache_blocks,
402 config.max_concurrent_merges,
403 Arc::clone(&config.background_merge_permits),
404 config.merge_bp_time_budget,
405 config.bp_memory_budget_bytes,
406 Arc::clone(&config.background_reorder_permits),
407 config.background_reorder_pool.clone(),
408 ));
409 let swept = segment_manager.cleanup_orphan_segments().await?;
410 if swept > 0 {
411 log::warn!(
412 "[segment_cleanup] swept {} orphan segment(s) while opening writer",
413 swept
414 );
415 }
416 segment_manager.try_load_and_publish_trained().await?;
417
418 Ok(Self::new_with_parts(
419 directory,
420 schema,
421 config,
422 builder_config,
423 segment_manager,
424 writer_lock,
425 ))
426 }
427
428 pub fn from_index(index: &super::Index<D>) -> Self {
435 let writer_lock = match try_acquire_writer_lock(index.directory.as_ref()) {
436 Ok(lock) => lock,
437 Err(error) => WriterLock::Unavailable {
438 reason: format!("failed to acquire the single-writer lock: {error}"),
439 },
440 };
441 if let WriterLock::Unavailable { reason } = &writer_lock {
442 log::error!("[writer_lock] {reason}");
443 }
444 Self::new_with_parts(
445 Arc::clone(&index.directory),
446 Arc::clone(&index.schema),
447 index.config.clone(),
448 SegmentBuilderConfig::default(),
449 Arc::clone(&index.segment_manager),
450 writer_lock,
451 )
452 }
453
454 fn new_with_parts(
460 directory: Arc<D>,
461 schema: Arc<Schema>,
462 config: IndexConfig,
463 builder_config: SegmentBuilderConfig,
464 segment_manager: Arc<crate::merge::SegmentManager<D>>,
465 writer_lock: WriterLock,
466 ) -> Self {
467 let registry = crate::tokenizer::TokenizerRegistry::new();
469 let mut tokenizers = FxHashMap::default();
470 for (field, entry) in schema.fields() {
471 if matches!(entry.field_type, crate::dsl::FieldType::Text)
472 && let Some(ref tok_name) = entry.tokenizer
473 && let Some(tok) = registry.get(tok_name)
474 {
475 tokenizers.insert(field, tok);
476 }
477 }
478
479 let num_workers = config.num_indexing_threads.max(1);
480 let worker_state = Arc::new(WorkerState {
481 directory: Arc::clone(&directory),
482 schema: Arc::clone(&schema),
483 builder_config,
484 tokenizers: parking_lot::RwLock::new(tokenizers),
485 memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
486 segment_manager: Arc::clone(&segment_manager),
487 built_segments: parking_lot::Mutex::new(Vec::new()),
488 cycle_error: parking_lot::Mutex::new(None),
489 cycle_failed: AtomicBool::new(false),
490 flush_count: AtomicUsize::new(0),
491 flush_mutex: parking_lot::Mutex::new(()),
492 flush_cvar: parking_lot::Condvar::new(),
493 resume_receiver: parking_lot::Mutex::new(None),
494 resume_epoch: AtomicUsize::new(0),
495 resume_cvar: parking_lot::Condvar::new(),
496 shutdown: AtomicBool::new(false),
497 num_workers,
498 });
499 let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
500
501 Self {
502 directory,
503 schema,
504 config,
505 doc_sender: Arc::new(parking_lot::RwLock::new(doc_sender)),
506 workers,
507 worker_state,
508 segment_manager,
509 flushed_segments: Arc::new(parking_lot::Mutex::new(Vec::new())),
510 primary_key_index: Arc::new(parking_lot::RwLock::new(None)),
511 commit_finalization: Arc::new(CommitFinalizationState::default()),
512 pk_reservations_retained: Arc::new(AtomicBool::new(false)),
513 writer_lock: parking_lot::RwLock::new(writer_lock),
514 }
515 }
516
517 fn ensure_writer_lock(&self) -> Result<()> {
525 if !matches!(&*self.writer_lock.read(), WriterLock::Unavailable { .. }) {
527 return Ok(());
528 }
529
530 let mut lock = self.writer_lock.write();
531 if !matches!(&*lock, WriterLock::Unavailable { .. }) {
533 return Ok(());
534 }
535 match try_acquire_writer_lock(self.directory.as_ref())? {
536 acquired @ (WriterLock::Held { .. } | WriterLock::NotApplicable) => {
537 log::info!(
538 "[writer_lock] single-writer lock acquired after retry; \
539 the previous holder has released it — resuming writes"
540 );
541 *lock = acquired;
542 Ok(())
543 }
544 WriterLock::Unavailable { reason } => {
545 let err = Error::Internal(reason.clone());
546 *lock = WriterLock::Unavailable { reason };
547 Err(err)
548 }
549 }
550 }
551
552 fn clear_uncommitted_pk_reservations(&self) {
561 if self.pk_reservations_retained.load(Ordering::Acquire) {
562 log::warn!(
563 "[primary_key] keeping uncommitted reservations through abort: a \
564 failed post-commit refresh left them as the only record of \
565 committed keys; they are cleared by the next successful commit"
566 );
567 return;
568 }
569 if let Some(pk_index) = self.primary_key_index.write().as_mut() {
570 pk_index.clear_uncommitted();
571 }
572 }
573
574 fn spawn_workers(
575 worker_state: &Arc<WorkerState<D>>,
576 num_workers: usize,
577 ) -> (
578 async_channel::Sender<Document>,
579 Vec<std::thread::JoinHandle<()>>,
580 ) {
581 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
582 let handle = tokio::runtime::Handle::current();
583 let mut workers = Vec::with_capacity(num_workers);
584 for i in 0..num_workers {
585 let state = Arc::clone(worker_state);
586 let rx = receiver.clone();
587 let rt = handle.clone();
588 workers.push(
589 std::thread::Builder::new()
590 .name(format!("index-worker-{}", i))
591 .spawn(move || Self::worker_loop(state, rx, rt))
592 .expect("failed to spawn index worker thread"),
593 );
594 }
595 (sender, workers)
596 }
597
598 pub fn schema(&self) -> &Schema {
600 &self.schema
601 }
602
603 pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
606 self.worker_state
607 .tokenizers
608 .write()
609 .insert(field, Box::new(tokenizer));
610 }
611
612 pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
628 use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
629
630 self.commit_finalization.wait_until_idle().await;
631
632 let field = match self.schema.primary_field() {
633 Some(f) => f,
634 None => return Ok(()),
635 };
636
637 let snapshot = self.segment_manager.acquire_snapshot().await;
638 let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
639
640 let cached = match self
642 .directory
643 .open_read(std::path::Path::new(PK_BLOOM_FILE))
644 .await
645 {
646 Ok(handle) => {
647 let data = handle.read_bytes_range(0..handle.len()).await;
648 match data {
649 Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
650 Err(_) => None,
651 }
652 }
653 Err(_) => None,
654 };
655
656 let load_futures: Vec<_> = current_seg_ids
658 .iter()
659 .map(|seg_id_str| {
660 let seg_id_str = seg_id_str.clone();
661 let dir = self.directory.as_ref();
662 let schema = Arc::clone(&self.schema);
663 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
664 })
665 .collect();
666 let all_data = futures::future::try_join_all(load_futures).await?;
667
668 if let Some((persisted_seg_ids, bloom)) = cached {
669 let mut pk_data = Vec::with_capacity(all_data.len());
671 let mut new_data = Vec::new();
672 for d in all_data {
673 if persisted_seg_ids.contains(&d.segment_id) {
674 pk_data.push(d);
675 } else {
676 new_data.push(d);
677 }
678 }
679 let needs_persist = !new_data.is_empty();
680 let new_start = pk_data.len();
681 pk_data.extend(new_data);
682
683 let pk_index = if new_start == pk_data.len() {
684 super::primary_key::PrimaryKeyIndex::from_persisted(
686 field,
687 bloom,
688 pk_data,
689 &[],
690 snapshot,
691 )
692 } else {
693 tokio::task::spawn_blocking(move || {
695 let mut bloom = bloom;
698 let mut added = 0usize;
699 let num_new = pk_data.len() - new_start;
700 for data in &pk_data[new_start..] {
701 if let Some(ff) = data.fast_fields.get(&field.0)
702 && let Some(dict) = ff.text_dict()
703 {
704 for key in dict.iter() {
705 bloom.insert(key.as_bytes());
706 added += 1;
707 }
708 }
709 }
710 if added > 0 {
711 log::info!(
712 "[primary_key] bloom: added {} keys from {} new segment(s)",
713 added,
714 num_new,
715 );
716 }
717 super::primary_key::PrimaryKeyIndex::from_persisted(
718 field,
719 bloom,
720 pk_data,
721 &[],
722 snapshot,
723 )
724 })
725 .await
726 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
727 };
728
729 if needs_persist {
730 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
731 }
732
733 *self.primary_key_index.write() = Some(pk_index);
734 } else {
735 let pk_index = tokio::task::spawn_blocking(move || {
737 super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
738 })
739 .await
740 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
741
742 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
743 *self.primary_key_index.write() = Some(pk_index);
744 }
745
746 self.pk_reservations_retained
750 .store(false, Ordering::Release);
751
752 Ok(())
753 }
754
755 async fn persist_pk_bloom(
758 &self,
759 pk_index: &super::primary_key::PrimaryKeyIndex,
760 segment_ids: &[String],
761 ) {
762 use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
763
764 let bloom_bytes = pk_index.bloom_to_bytes();
765 let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
766 if let Err(e) = self
767 .directory
768 .write(std::path::Path::new(PK_BLOOM_FILE), &data)
769 .await
770 {
771 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
772 }
773 }
774
775 pub fn add_document(&self, doc: Document) -> Result<()> {
781 self.ensure_writer_lock()?;
782 if self.worker_state.shutdown.load(Ordering::Acquire) {
783 return Err(Error::IndexClosed);
784 }
785 if self.commit_finalization.in_progress.load(Ordering::Acquire) {
786 return Err(Error::CommitInProgress);
787 }
788 let sender = self.doc_sender.read().clone();
789 if sender.is_closed() {
793 return Err(Error::CommitInProgress);
794 }
795 let primary_key_index = self.primary_key_index.read();
796 if let Some(ref pk_index) = *primary_key_index {
797 pk_index.check_and_insert(&doc)?;
798 }
799 match sender.try_send(doc) {
800 Ok(()) => Ok(()),
801 Err(async_channel::TrySendError::Full(doc)) => {
802 if let Some(ref pk_index) = *primary_key_index {
804 pk_index.rollback_uncommitted_key(&doc);
805 }
806 Err(Error::QueueFull)
807 }
808 Err(async_channel::TrySendError::Closed(doc)) => {
809 if let Some(ref pk_index) = *primary_key_index {
811 pk_index.rollback_uncommitted_key(&doc);
812 }
813 Err(Error::CommitInProgress)
814 }
815 }
816 }
817
818 pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
823 let total = documents.len();
824 for (i, doc) in documents.into_iter().enumerate() {
825 match self.add_document(doc) {
826 Ok(()) => {}
827 Err(Error::QueueFull | Error::CommitInProgress) => return Ok(i),
828 Err(e) => return Err(e),
829 }
830 }
831 Ok(total)
832 }
833
834 fn worker_loop(
847 state: Arc<WorkerState<D>>,
848 initial_receiver: async_channel::Receiver<Document>,
849 handle: tokio::runtime::Handle,
850 ) {
851 let mut receiver = initial_receiver;
852 let mut my_epoch = 0usize;
853
854 loop {
855 let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
859 let mut builder: Option<SegmentBuilder> = None;
860
861 while let Ok(doc) = receiver.recv_blocking() {
862 if state.shutdown.load(Ordering::Acquire) {
863 break;
864 }
865 if state.cycle_failed.load(Ordering::Acquire) {
870 continue;
871 }
872 if builder.is_none() {
874 match SegmentBuilder::new(
875 Arc::clone(&state.schema),
876 state.builder_config.clone(),
877 ) {
878 Ok(mut b) => {
879 for (field, tokenizer) in state.tokenizers.read().iter() {
880 b.set_tokenizer(*field, tokenizer.clone_box());
881 }
882 builder = Some(b);
883 }
884 Err(e) => {
885 log::error!("Failed to create segment builder: {:?}", e);
886 state.record_cycle_error(format!(
887 "failed to create segment builder: {e}"
888 ));
889 continue;
890 }
891 }
892 }
893
894 let b = builder.as_mut().unwrap();
895 if let Err(e) = b.add_document(doc) {
896 log::error!("Failed to index document: {:?}", e);
897 state.record_cycle_error(format!("failed to index document: {e}"));
898 continue;
899 }
900
901 let builder_memory = b.estimated_memory_bytes();
902
903 if b.num_docs() & 0x3FFF == 0 {
904 log::debug!(
905 "[indexing] docs={}, memory={}, budget={}",
906 b.num_docs(),
907 crate::format_bytes(builder_memory as u64),
908 crate::format_bytes(state.memory_budget_per_worker as u64)
909 );
910 }
911
912 const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
914
915 let effective_budget = state.memory_budget_per_worker * 4 / 5;
919
920 if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
921 log::info!(
922 "[indexing] memory budget reached, building segment: \
923 docs={}, memory={}, budget={}",
924 b.num_docs(),
925 crate::format_bytes(builder_memory as u64),
926 crate::format_bytes(state.memory_budget_per_worker as u64),
927 );
928 let full_builder = builder.take().unwrap();
929 Self::build_segment_inline(&state, full_builder, &handle);
930 }
931 }
932
933 if !state.cycle_failed.load(Ordering::Acquire)
935 && let Some(b) = builder.take()
936 && b.num_docs() > 0
937 {
938 Self::build_segment_inline(&state, b, &handle);
939 }
940 }));
941
942 if build_result.is_err() {
943 log::error!(
944 "[worker] panic during indexing cycle — documents in this cycle may be lost"
945 );
946 state.record_cycle_error("indexing worker panicked while building the batch");
947 }
948
949 let prev = state.flush_count.fetch_add(1, Ordering::Release);
952 if prev + 1 == state.num_workers {
953 let _lock = state.flush_mutex.lock();
960 state.flush_cvar.notify_all();
961 }
962
963 {
967 let mut lock = state.resume_receiver.lock();
968 loop {
969 if state.shutdown.load(Ordering::Acquire) {
970 return;
971 }
972 let current_epoch = state.resume_epoch.load(Ordering::Acquire);
973 if current_epoch > my_epoch
974 && let Some(rx) = lock.as_ref()
975 {
976 receiver = rx.clone();
977 my_epoch = current_epoch;
978 break;
979 }
980 state.resume_cvar.wait(&mut lock);
981 }
982 }
983 }
984 }
985
986 fn build_segment_inline(
990 state: &WorkerState<D>,
991 builder: SegmentBuilder,
992 handle: &tokio::runtime::Handle,
993 ) {
994 let segment_id = SegmentId::new();
995 let segment_hex = segment_id.to_hex();
996 let operation = match state
999 .segment_manager
1000 .protect_new_segment(segment_hex.clone())
1001 {
1002 Ok(operation) => operation,
1003 Err(e) => {
1004 log::error!(
1005 "[segment_build_failed] segment_id={} lifecycle_error={}",
1006 segment_hex,
1007 e,
1008 );
1009 state.record_cycle_error(format!(
1010 "failed to claim segment {segment_hex} for building: {e}"
1011 ));
1012 return;
1013 }
1014 };
1015 let trained = state.segment_manager.trained_for_segment_build();
1016 let doc_count = builder.num_docs();
1017 let build_start = std::time::Instant::now();
1018
1019 log::info!(
1020 "[segment_build] segment_id={} doc_count={} ann={}",
1021 segment_hex,
1022 doc_count,
1023 trained.is_some()
1024 );
1025
1026 let mut prepared = PreparedSegment {
1030 id: segment_hex.clone(),
1031 segment_id,
1032 num_docs: doc_count,
1033 segment_manager: Arc::clone(&state.segment_manager),
1034 operation: Some(operation),
1035 runtime: handle.clone(),
1036 needs_vector_upgrade: trained.is_none(),
1037 published: false,
1038 };
1039
1040 match handle.block_on(builder.build(
1041 state.directory.as_ref(),
1042 segment_id,
1043 trained.as_deref(),
1044 )) {
1045 Ok(meta) if meta.num_docs == doc_count && meta.num_docs > 0 => {
1046 let duration_ms = build_start.elapsed().as_millis() as u64;
1047 log::info!(
1048 "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
1049 segment_hex,
1050 meta.num_docs,
1051 duration_ms,
1052 );
1053 prepared.num_docs = meta.num_docs;
1054 state.built_segments.lock().push(prepared);
1055 }
1056 Ok(meta) => {
1057 let error = format!(
1058 "segment {segment_hex} built {} docs from a {doc_count}-document builder",
1059 meta.num_docs
1060 );
1061 log::error!("[segment_build_failed] {error}");
1062 state.record_cycle_error(error);
1063 }
1064 Err(e) => {
1065 log::error!(
1066 "[segment_build_failed] segment_id={} error={:?}",
1067 segment_hex,
1068 e
1069 );
1070 state.record_cycle_error(format!("failed to build segment {segment_hex}: {e}"));
1073 }
1074 }
1075 }
1076
1077 pub async fn maybe_merge(&self) {
1083 self.segment_manager.maybe_merge().await;
1084 }
1085
1086 pub async fn abort_merges(&self) {
1089 self.segment_manager.abort_merges().await;
1090 }
1091
1092 pub async fn shutdown(&mut self) -> Result<()> {
1097 self.segment_manager.begin_shutdown();
1098 self.signal_worker_shutdown();
1099
1100 self.commit_finalization.wait_until_idle().await;
1105
1106 let workers = std::mem::take(&mut self.workers);
1107 let panicked = tokio::task::spawn_blocking(move || {
1108 workers
1109 .into_iter()
1110 .map(|worker| worker.join().is_err())
1111 .filter(|panicked| *panicked)
1112 .count()
1113 })
1114 .await
1115 .map_err(|error| Error::Internal(format!("failed to join index workers: {}", error)))?;
1116 if panicked > 0 {
1117 log::error!("[index_shutdown] {} indexing worker(s) panicked", panicked);
1118 }
1119
1120 self.flushed_segments.lock().clear();
1123 self.worker_state.built_segments.lock().clear();
1124 if let Some(pk_index) = self.primary_key_index.write().as_mut() {
1125 pk_index.clear_uncommitted();
1126 }
1127 Ok(())
1128 }
1129
1130 pub async fn wait_for_merging_thread(&self) {
1132 self.segment_manager.wait_for_merging_thread().await;
1133 }
1134
1135 pub async fn wait_for_all_merges(&self) {
1137 self.segment_manager.wait_for_all_merges().await;
1138 }
1139
1140 pub async fn wait_for_commit_finalization(&self) {
1145 self.commit_finalization.wait_until_idle().await;
1146 }
1147
1148 pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
1150 self.segment_manager.tracker()
1151 }
1152
1153 pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
1155 self.segment_manager.acquire_snapshot().await
1156 }
1157
1158 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
1163 self.ensure_writer_lock()?;
1164 self.segment_manager.cleanup_orphan_segments().await
1165 }
1166
1167 pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
1178 self.ensure_writer_lock()?;
1179 if self.worker_state.shutdown.load(Ordering::Acquire) {
1180 return Err(Error::IndexClosed);
1181 }
1182 if self.commit_finalization.in_progress.load(Ordering::Acquire) {
1183 return Err(Error::CommitInProgress);
1184 }
1185 self.doc_sender.read().close();
1187
1188 self.worker_state.resume_cvar.notify_all();
1192
1193 let state = Arc::clone(&self.worker_state);
1196 let all_flushed = tokio::task::spawn_blocking(move || {
1197 let mut lock = state.flush_mutex.lock();
1198 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
1199 while state.flush_count.load(Ordering::Acquire) < state.num_workers {
1200 let remaining = deadline.saturating_duration_since(std::time::Instant::now());
1201 if remaining.is_zero() {
1202 log::error!(
1203 "[prepare_commit] timed out waiting for workers: {}/{} flushed",
1204 state.flush_count.load(Ordering::Acquire),
1205 state.num_workers
1206 );
1207 return false;
1208 }
1209 state.flush_cvar.wait_for(&mut lock, remaining);
1210 }
1211 true
1212 })
1213 .await
1214 .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
1215
1216 if !all_flushed {
1217 return Err(Error::Internal(format!(
1225 "prepare_commit timed out: {}/{} workers flushed; writer remains paused, retry commit",
1226 self.worker_state.flush_count.load(Ordering::Acquire),
1227 self.worker_state.num_workers
1228 )));
1229 }
1230
1231 let cycle_error = { self.worker_state.cycle_error.lock().take() };
1232 if let Some(error) = cycle_error {
1233 self.flushed_segments.lock().clear();
1238 self.worker_state.built_segments.lock().clear();
1239 self.clear_uncommitted_pk_reservations();
1240 self.resume_workers();
1241 return Err(Error::Internal(format!(
1242 "indexing generation failed; no documents from this batch were committed: {error}"
1243 )));
1244 }
1245
1246 let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
1248 self.flushed_segments.lock().extend(built);
1249
1250 Ok(PreparedCommit {
1251 writer: self,
1252 is_resolved: false,
1253 })
1254 }
1255
1256 pub async fn commit(&mut self) -> Result<bool> {
1261 self.prepare_commit().await?.commit().await
1262 }
1263
1264 pub async fn force_merge(&mut self) -> Result<()> {
1266 self.prepare_commit().await?.commit().await?;
1267 self.segment_manager.force_merge().await
1268 }
1269
1270 pub async fn reorder(&mut self) -> Result<()> {
1275 self.prepare_commit().await?.commit().await?;
1276 self.segment_manager.reorder_segments().await
1277 }
1278
1279 pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
1281 &self.segment_manager
1282 }
1283
1284 fn resume_workers(&mut self) {
1289 Self::resume_workers_shared(&self.worker_state, &self.doc_sender);
1290 }
1291
1292 fn resume_workers_shared(
1293 worker_state: &Arc<WorkerState<D>>,
1294 doc_sender: &Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1295 ) {
1296 if worker_state.shutdown.load(Ordering::Acquire) {
1297 return;
1298 }
1299 if tokio::runtime::Handle::try_current().is_err() {
1300 worker_state.shutdown.store(true, Ordering::Release);
1303 worker_state.resume_cvar.notify_all();
1304 return;
1305 }
1306
1307 worker_state.flush_count.store(0, Ordering::Release);
1309 *worker_state.cycle_error.lock() = None;
1310 worker_state.cycle_failed.store(false, Ordering::Release);
1311
1312 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
1314 *doc_sender.write() = sender;
1315
1316 {
1318 let mut lock = worker_state.resume_receiver.lock();
1319 *lock = Some(receiver);
1320 }
1321 worker_state.resume_epoch.fetch_add(1, Ordering::Release);
1322 worker_state.resume_cvar.notify_all();
1323 }
1324
1325 fn signal_worker_shutdown(&self) {
1326 self.worker_state.shutdown.store(true, Ordering::Release);
1327 self.doc_sender.read().close();
1328 self.worker_state.resume_cvar.notify_all();
1329 }
1330
1331 }
1333
1334impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
1335 fn drop(&mut self) {
1336 self.signal_worker_shutdown();
1337 for w in std::mem::take(&mut self.workers) {
1338 let _ = w.join();
1339 }
1340 }
1341}
1342
1343pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
1350 writer: &'a mut IndexWriter<D>,
1351 is_resolved: bool,
1352}
1353
1354struct PreparedSegmentsGuard<D: DirectoryWriter + 'static> {
1359 segments: Option<Vec<PreparedSegment<D>>>,
1360 retry_slot: Arc<parking_lot::Mutex<Vec<PreparedSegment<D>>>>,
1361}
1362
1363impl<D: DirectoryWriter + 'static> PreparedSegmentsGuard<D> {
1364 fn metadata_entries(&self) -> Vec<(String, u32)> {
1365 self.segments
1366 .as_deref()
1367 .unwrap_or_default()
1368 .iter()
1369 .map(PreparedSegment::metadata_entry)
1370 .collect()
1371 }
1372
1373 fn take_published(&mut self) -> Vec<PreparedSegment<D>> {
1374 self.segments.take().unwrap_or_default()
1375 }
1376
1377 fn vector_upgrade_segment_ids(&self) -> Vec<String> {
1378 self.segments
1379 .as_deref()
1380 .unwrap_or_default()
1381 .iter()
1382 .filter(|segment| segment.needs_vector_upgrade)
1383 .map(|segment| segment.id.clone())
1384 .collect()
1385 }
1386}
1387
1388impl<D: DirectoryWriter + 'static> Drop for PreparedSegmentsGuard<D> {
1389 fn drop(&mut self) {
1390 if let Some(segments) = self.segments.take() {
1391 self.retry_slot.lock().extend(segments);
1392 }
1393 }
1394}
1395
1396struct CommitFinalizationGuard<D: DirectoryWriter + 'static> {
1401 state: Arc<CommitFinalizationState>,
1402 worker_state: Arc<WorkerState<D>>,
1403 doc_sender: Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1404 resume_workers: bool,
1405}
1406
1407impl<D: DirectoryWriter + 'static> CommitFinalizationGuard<D> {
1408 fn resume_on_drop(&mut self) {
1409 self.resume_workers = true;
1410 }
1411}
1412
1413impl<D: DirectoryWriter + 'static> Drop for CommitFinalizationGuard<D> {
1414 fn drop(&mut self) {
1415 if self.resume_workers {
1416 IndexWriter::<D>::resume_workers_shared(&self.worker_state, &self.doc_sender);
1417 }
1418 self.state.finish();
1419 }
1420}
1421
1422struct OwnedCommitFinalization<D: DirectoryWriter + 'static> {
1427 directory: Arc<D>,
1428 schema: Arc<Schema>,
1429 segment_manager: Arc<crate::merge::SegmentManager<D>>,
1430 primary_key_index: Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1431 prepared: PreparedSegmentsGuard<D>,
1432 finalization: Option<CommitFinalizationGuard<D>>,
1433 publication_observed: Arc<AtomicBool>,
1434 pk_reservations_retained: Arc<AtomicBool>,
1435}
1436
1437async fn refresh_primary_key_after_commit<D: DirectoryWriter + 'static>(
1438 directory: &Arc<D>,
1439 schema: &Arc<Schema>,
1440 segment_manager: &Arc<crate::merge::SegmentManager<D>>,
1441 primary_key_index: &Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1442) -> Result<()> {
1443 let existing_ids: std::collections::HashSet<String> = {
1444 let guard = primary_key_index.read();
1445 let Some(pk_index) = guard.as_ref() else {
1446 return Ok(());
1447 };
1448 pk_index
1449 .committed_segment_ids()
1450 .map(ToOwned::to_owned)
1451 .collect()
1452 };
1453
1454 let snapshot = segment_manager.acquire_snapshot().await;
1455 let load_futures: Vec<_> = snapshot
1456 .segment_ids()
1457 .iter()
1458 .filter(|id| !existing_ids.contains(id.as_str()))
1459 .map(|seg_id_str| {
1460 let seg_id_str = seg_id_str.clone();
1461 let dir = directory.as_ref();
1462 let schema = Arc::clone(schema);
1463 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
1464 })
1465 .collect();
1466 let new_data = futures::future::try_join_all(load_futures).await?;
1467 let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
1468
1469 let bloom_file = {
1470 let mut guard = primary_key_index.write();
1471 let Some(pk_index) = guard.as_mut() else {
1472 return Ok(());
1473 };
1474 pk_index.refresh_incremental(new_data, snapshot);
1475 let bloom_bytes = pk_index.bloom_to_bytes();
1476 super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes)
1477 };
1478
1479 if let Err(error) = directory
1480 .write(
1481 std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
1482 &bloom_file,
1483 )
1484 .await
1485 {
1486 log::warn!("[primary_key] failed to persist bloom cache: {}", error);
1487 }
1488 Ok(())
1489}
1490
1491async fn finalize_prepared_commit<D: DirectoryWriter + 'static>(
1492 mut commit: OwnedCommitFinalization<D>,
1493) -> Result<bool> {
1494 let metadata_entries = commit.prepared.metadata_entries();
1495 let published_segment_ids = commit.prepared.vector_upgrade_segment_ids();
1496
1497 commit.segment_manager.commit(&metadata_entries).await?;
1501 commit.publication_observed.store(true, Ordering::Release);
1502
1503 let mut published = commit.prepared.take_published();
1504 for segment in &mut published {
1505 segment.mark_published();
1506 }
1507 drop(published);
1508 commit
1509 .segment_manager
1510 .schedule_vector_segment_upgrades(published_segment_ids);
1511 if let Some(finalization) = commit.finalization.as_mut() {
1515 finalization.resume_on_drop();
1516 } else {
1517 log::error!("owned commit finalization guard was already released after publication");
1518 }
1519
1520 match refresh_primary_key_after_commit(
1525 &commit.directory,
1526 &commit.schema,
1527 &commit.segment_manager,
1528 &commit.primary_key_index,
1529 )
1530 .await
1531 {
1532 Ok(()) => commit
1535 .pk_reservations_retained
1536 .store(false, Ordering::Release),
1537 Err(error) => {
1538 commit
1543 .pk_reservations_retained
1544 .store(true, Ordering::Release);
1545 log::error!(
1546 "[primary_key] committed metadata but failed to refresh dedup state; \
1547 retaining reservations until a later successful commit: {}",
1548 error,
1549 );
1550 }
1551 }
1552
1553 drop(commit.finalization.take());
1557 commit.segment_manager.maybe_merge().await;
1558 Ok(true)
1559}
1560
1561impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
1562 pub async fn commit(mut self) -> Result<bool> {
1566 let segments = std::mem::take(&mut *self.writer.flushed_segments.lock());
1567
1568 if segments.is_empty() {
1570 log::debug!("[commit] no segments to commit, skipping");
1571 self.is_resolved = true;
1572 self.writer.resume_workers();
1573 return Ok(false);
1574 }
1575
1576 if !self.writer.commit_finalization.begin() {
1577 self.writer.flushed_segments.lock().extend(segments);
1578 self.is_resolved = true;
1582 return Err(Error::CommitInProgress);
1583 }
1584
1585 let publication_observed = Arc::new(AtomicBool::new(false));
1586 let owned = OwnedCommitFinalization {
1587 directory: Arc::clone(&self.writer.directory),
1588 schema: Arc::clone(&self.writer.schema),
1589 segment_manager: Arc::clone(&self.writer.segment_manager),
1590 primary_key_index: Arc::clone(&self.writer.primary_key_index),
1591 prepared: PreparedSegmentsGuard {
1592 segments: Some(segments),
1593 retry_slot: Arc::clone(&self.writer.flushed_segments),
1594 },
1595 finalization: Some(CommitFinalizationGuard {
1596 state: Arc::clone(&self.writer.commit_finalization),
1597 worker_state: Arc::clone(&self.writer.worker_state),
1598 doc_sender: Arc::clone(&self.writer.doc_sender),
1599 resume_workers: false,
1600 }),
1601 publication_observed: Arc::clone(&publication_observed),
1602 pk_reservations_retained: Arc::clone(&self.writer.pk_reservations_retained),
1603 };
1604
1605 self.is_resolved = true;
1610 let task_publication = Arc::clone(&publication_observed);
1611 let task = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1612 tokio::spawn(async move {
1613 match std::panic::AssertUnwindSafe(finalize_prepared_commit(owned))
1614 .catch_unwind()
1615 .await
1616 {
1617 Ok(result) => result,
1618 Err(_) if task_publication.load(Ordering::Acquire) => {
1619 log::error!(
1620 "owned commit finalizer panicked after metadata publication; \
1621 treating the durable generation as committed"
1622 );
1623 Ok(true)
1624 }
1625 Err(_) => Err(Error::Internal(
1626 "owned commit finalizer panicked before metadata publication".into(),
1627 )),
1628 }
1629 })
1630 }))
1631 .map_err(|_| Error::Internal("runtime rejected owned commit finalizer".into()))?;
1632
1633 match task.await {
1634 Ok(result) => result,
1635 Err(error) if publication_observed.load(Ordering::Acquire) => {
1636 log::error!(
1637 "owned commit finalizer terminated after metadata publication: {}; \
1638 treating the durable generation as committed",
1639 error,
1640 );
1641 Ok(true)
1642 }
1643 Err(error) => Err(Error::Internal(format!(
1644 "owned commit finalizer terminated unexpectedly: {error}"
1645 ))),
1646 }
1647 }
1648
1649 pub fn abort(mut self) {
1652 self.is_resolved = true;
1653 self.writer.flushed_segments.lock().clear();
1654 self.writer.clear_uncommitted_pk_reservations();
1655 self.writer.resume_workers();
1656 }
1657}
1658
1659impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
1660 fn drop(&mut self) {
1661 if !self.is_resolved {
1662 log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
1663 self.writer.flushed_segments.lock().clear();
1664 self.writer.clear_uncommitted_pk_reservations();
1665 self.writer.resume_workers();
1666 }
1667 }
1668}
1669
1670async fn load_pk_segment_data<D: crate::directories::Directory>(
1672 dir: &D,
1673 seg_id_str: &str,
1674 schema: &Arc<crate::dsl::Schema>,
1675) -> Result<super::primary_key::PkSegmentData> {
1676 let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
1677 .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
1678 let files = crate::segment::SegmentFiles::new(seg_id.0);
1679 let fast_fields =
1680 crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
1681 Ok(super::primary_key::PkSegmentData {
1682 segment_id: seg_id_str.to_string(),
1683 fast_fields,
1684 })
1685}