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