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: WriterLock,
166}
167
168#[derive(Default)]
169struct CommitFinalizationState {
170 in_progress: AtomicBool,
171 idle: tokio::sync::Notify,
172}
173
174impl CommitFinalizationState {
175 fn begin(&self) -> bool {
176 self.in_progress
177 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
178 .is_ok()
179 }
180
181 fn finish(&self) {
182 self.in_progress.store(false, Ordering::Release);
183 self.idle.notify_waiters();
184 }
185
186 async fn wait_until_idle(&self) {
187 while self.in_progress.load(Ordering::Acquire) {
188 let notified = self.idle.notified();
189 if !self.in_progress.load(Ordering::Acquire) {
190 break;
191 }
192 notified.await;
193 }
194 }
195}
196
197struct WorkerState<D: DirectoryWriter + 'static> {
199 directory: Arc<D>,
200 schema: Arc<Schema>,
201 builder_config: SegmentBuilderConfig,
202 tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
203 memory_budget_per_worker: usize,
205 segment_manager: Arc<crate::merge::SegmentManager<D>>,
207 built_segments: parking_lot::Mutex<Vec<PreparedSegment<D>>>,
210 cycle_error: parking_lot::Mutex<Option<String>>,
215 cycle_failed: AtomicBool,
216
217 flush_count: AtomicUsize,
224 flush_mutex: parking_lot::Mutex<()>,
226 flush_cvar: parking_lot::Condvar,
227 resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
229 resume_epoch: AtomicUsize,
232 resume_cvar: parking_lot::Condvar,
234 shutdown: AtomicBool,
236 num_workers: usize,
238}
239
240struct PreparedSegment<D: DirectoryWriter + 'static> {
246 id: String,
247 segment_id: SegmentId,
248 num_docs: u32,
249 segment_manager: Arc<crate::merge::SegmentManager<D>>,
250 operation: Option<crate::merge::SegmentOperationGuard>,
251 runtime: tokio::runtime::Handle,
252 published: bool,
253}
254
255impl<D: DirectoryWriter + 'static> PreparedSegment<D> {
256 fn metadata_entry(&self) -> (String, u32) {
257 (self.id.clone(), self.num_docs)
258 }
259
260 fn mark_published(&mut self) {
261 self.published = true;
262 drop(self.operation.take());
264 }
265}
266
267impl<D: DirectoryWriter + 'static> WorkerState<D> {
268 fn record_cycle_error(&self, error: impl Into<String>) {
269 let mut first_error = self.cycle_error.lock();
270 if first_error.is_none() {
271 *first_error = Some(error.into());
272 }
273 drop(first_error);
274 self.cycle_failed.store(true, Ordering::Release);
275 }
276}
277
278impl<D: DirectoryWriter + 'static> Drop for PreparedSegment<D> {
279 fn drop(&mut self) {
280 if self.published {
281 return;
282 }
283 let Some(operation) = self.operation.take() else {
284 return;
285 };
286 self.segment_manager.schedule_unpublished_segment_cleanup(
287 self.segment_id,
288 operation,
289 self.runtime.clone(),
290 );
291 }
292}
293
294impl<D: DirectoryWriter + 'static> IndexWriter<D> {
295 pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
297 Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
298 }
299
300 pub async fn create_with_config(
302 directory: D,
303 schema: Schema,
304 config: IndexConfig,
305 builder_config: SegmentBuilderConfig,
306 ) -> Result<Self> {
307 let directory = Arc::new(directory);
308 let schema = Arc::new(schema);
309 directory.set_index_label(schema.index_label());
311
312 let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
314 if let WriterLock::Unavailable { reason } = &writer_lock {
315 return Err(Error::Internal(reason.clone()));
316 }
317 if directory
321 .exists(std::path::Path::new(super::INDEX_META_FILENAME))
322 .await?
323 {
324 return Err(Error::Internal(format!(
325 "refusing to create index: {} already exists in this directory; \
326 use IndexWriter::open to open the existing index, or delete the \
327 directory first if you really want to start over",
328 super::INDEX_META_FILENAME
329 )));
330 }
331
332 let metadata = super::IndexMetadata::new((*schema).clone());
333
334 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
335 Arc::clone(&directory),
336 Arc::clone(&schema),
337 metadata,
338 config.merge_policy.clone_box(),
339 config.term_cache_blocks,
340 config.max_concurrent_merges,
341 Arc::clone(&config.background_merge_permits),
342 config.merge_bp_time_budget,
343 config.bp_memory_budget_bytes,
344 Arc::clone(&config.background_reorder_permits),
345 config.background_reorder_pool.clone(),
346 ));
347 segment_manager.update_metadata(|_| {}).await?;
348
349 Ok(Self::new_with_parts(
350 directory,
351 schema,
352 config,
353 builder_config,
354 segment_manager,
355 writer_lock,
356 ))
357 }
358
359 pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
368 Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
369 }
370
371 pub async fn open_with_config(
373 directory: D,
374 config: IndexConfig,
375 builder_config: SegmentBuilderConfig,
376 ) -> Result<Self> {
377 let directory = Arc::new(directory);
378
379 let writer_lock = try_acquire_writer_lock(directory.as_ref())?;
382 if let WriterLock::Unavailable { reason } = &writer_lock {
383 return Err(Error::Internal(reason.clone()));
384 }
385
386 let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
387 let schema = Arc::new(metadata.schema.clone());
388 directory.set_index_label(schema.index_label());
390
391 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
392 Arc::clone(&directory),
393 Arc::clone(&schema),
394 metadata,
395 config.merge_policy.clone_box(),
396 config.term_cache_blocks,
397 config.max_concurrent_merges,
398 Arc::clone(&config.background_merge_permits),
399 config.merge_bp_time_budget,
400 config.bp_memory_budget_bytes,
401 Arc::clone(&config.background_reorder_permits),
402 config.background_reorder_pool.clone(),
403 ));
404 let swept = segment_manager.cleanup_orphan_segments().await?;
405 if swept > 0 {
406 log::warn!(
407 "[segment_cleanup] swept {} orphan segment(s) while opening writer",
408 swept
409 );
410 }
411 segment_manager.try_load_and_publish_trained().await?;
412
413 Ok(Self::new_with_parts(
414 directory,
415 schema,
416 config,
417 builder_config,
418 segment_manager,
419 writer_lock,
420 ))
421 }
422
423 pub fn from_index(index: &super::Index<D>) -> Self {
430 let writer_lock = match try_acquire_writer_lock(index.directory.as_ref()) {
431 Ok(lock) => lock,
432 Err(error) => WriterLock::Unavailable {
433 reason: format!("failed to acquire the single-writer lock: {error}"),
434 },
435 };
436 if let WriterLock::Unavailable { reason } = &writer_lock {
437 log::error!("[writer_lock] {reason}");
438 }
439 Self::new_with_parts(
440 Arc::clone(&index.directory),
441 Arc::clone(&index.schema),
442 index.config.clone(),
443 SegmentBuilderConfig::default(),
444 Arc::clone(&index.segment_manager),
445 writer_lock,
446 )
447 }
448
449 fn new_with_parts(
455 directory: Arc<D>,
456 schema: Arc<Schema>,
457 config: IndexConfig,
458 builder_config: SegmentBuilderConfig,
459 segment_manager: Arc<crate::merge::SegmentManager<D>>,
460 writer_lock: WriterLock,
461 ) -> Self {
462 let registry = crate::tokenizer::TokenizerRegistry::new();
464 let mut tokenizers = FxHashMap::default();
465 for (field, entry) in schema.fields() {
466 if matches!(entry.field_type, crate::dsl::FieldType::Text)
467 && let Some(ref tok_name) = entry.tokenizer
468 && let Some(tok) = registry.get(tok_name)
469 {
470 tokenizers.insert(field, tok);
471 }
472 }
473
474 let num_workers = config.num_indexing_threads.max(1);
475 let worker_state = Arc::new(WorkerState {
476 directory: Arc::clone(&directory),
477 schema: Arc::clone(&schema),
478 builder_config,
479 tokenizers: parking_lot::RwLock::new(tokenizers),
480 memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
481 segment_manager: Arc::clone(&segment_manager),
482 built_segments: parking_lot::Mutex::new(Vec::new()),
483 cycle_error: parking_lot::Mutex::new(None),
484 cycle_failed: AtomicBool::new(false),
485 flush_count: AtomicUsize::new(0),
486 flush_mutex: parking_lot::Mutex::new(()),
487 flush_cvar: parking_lot::Condvar::new(),
488 resume_receiver: parking_lot::Mutex::new(None),
489 resume_epoch: AtomicUsize::new(0),
490 resume_cvar: parking_lot::Condvar::new(),
491 shutdown: AtomicBool::new(false),
492 num_workers,
493 });
494 let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
495
496 Self {
497 directory,
498 schema,
499 config,
500 doc_sender: Arc::new(parking_lot::RwLock::new(doc_sender)),
501 workers,
502 worker_state,
503 segment_manager,
504 flushed_segments: Arc::new(parking_lot::Mutex::new(Vec::new())),
505 primary_key_index: Arc::new(parking_lot::RwLock::new(None)),
506 commit_finalization: Arc::new(CommitFinalizationState::default()),
507 pk_reservations_retained: Arc::new(AtomicBool::new(false)),
508 writer_lock,
509 }
510 }
511
512 fn ensure_writer_lock(&self) -> Result<()> {
514 if let WriterLock::Unavailable { reason } = &self.writer_lock {
515 return Err(Error::Internal(reason.clone()));
516 }
517 Ok(())
518 }
519
520 fn clear_uncommitted_pk_reservations(&self) {
529 if self.pk_reservations_retained.load(Ordering::Acquire) {
530 log::warn!(
531 "[primary_key] keeping uncommitted reservations through abort: a \
532 failed post-commit refresh left them as the only record of \
533 committed keys; they are cleared by the next successful commit"
534 );
535 return;
536 }
537 if let Some(pk_index) = self.primary_key_index.write().as_mut() {
538 pk_index.clear_uncommitted();
539 }
540 }
541
542 fn spawn_workers(
543 worker_state: &Arc<WorkerState<D>>,
544 num_workers: usize,
545 ) -> (
546 async_channel::Sender<Document>,
547 Vec<std::thread::JoinHandle<()>>,
548 ) {
549 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
550 let handle = tokio::runtime::Handle::current();
551 let mut workers = Vec::with_capacity(num_workers);
552 for i in 0..num_workers {
553 let state = Arc::clone(worker_state);
554 let rx = receiver.clone();
555 let rt = handle.clone();
556 workers.push(
557 std::thread::Builder::new()
558 .name(format!("index-worker-{}", i))
559 .spawn(move || Self::worker_loop(state, rx, rt))
560 .expect("failed to spawn index worker thread"),
561 );
562 }
563 (sender, workers)
564 }
565
566 pub fn schema(&self) -> &Schema {
568 &self.schema
569 }
570
571 pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
574 self.worker_state
575 .tokenizers
576 .write()
577 .insert(field, Box::new(tokenizer));
578 }
579
580 pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
596 use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
597
598 self.commit_finalization.wait_until_idle().await;
599
600 let field = match self.schema.primary_field() {
601 Some(f) => f,
602 None => return Ok(()),
603 };
604
605 let snapshot = self.segment_manager.acquire_snapshot().await;
606 let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
607
608 let cached = match self
610 .directory
611 .open_read(std::path::Path::new(PK_BLOOM_FILE))
612 .await
613 {
614 Ok(handle) => {
615 let data = handle.read_bytes_range(0..handle.len()).await;
616 match data {
617 Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
618 Err(_) => None,
619 }
620 }
621 Err(_) => None,
622 };
623
624 let load_futures: Vec<_> = current_seg_ids
626 .iter()
627 .map(|seg_id_str| {
628 let seg_id_str = seg_id_str.clone();
629 let dir = self.directory.as_ref();
630 let schema = Arc::clone(&self.schema);
631 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
632 })
633 .collect();
634 let all_data = futures::future::try_join_all(load_futures).await?;
635
636 if let Some((persisted_seg_ids, bloom)) = cached {
637 let mut pk_data = Vec::with_capacity(all_data.len());
639 let mut new_data = Vec::new();
640 for d in all_data {
641 if persisted_seg_ids.contains(&d.segment_id) {
642 pk_data.push(d);
643 } else {
644 new_data.push(d);
645 }
646 }
647 let needs_persist = !new_data.is_empty();
648 let new_start = pk_data.len();
649 pk_data.extend(new_data);
650
651 let pk_index = if new_start == pk_data.len() {
652 super::primary_key::PrimaryKeyIndex::from_persisted(
654 field,
655 bloom,
656 pk_data,
657 &[],
658 snapshot,
659 )
660 } else {
661 tokio::task::spawn_blocking(move || {
663 let mut bloom = bloom;
666 let mut added = 0usize;
667 let num_new = pk_data.len() - new_start;
668 for data in &pk_data[new_start..] {
669 if let Some(ff) = data.fast_fields.get(&field.0)
670 && let Some(dict) = ff.text_dict()
671 {
672 for key in dict.iter() {
673 bloom.insert(key.as_bytes());
674 added += 1;
675 }
676 }
677 }
678 if added > 0 {
679 log::info!(
680 "[primary_key] bloom: added {} keys from {} new segment(s)",
681 added,
682 num_new,
683 );
684 }
685 super::primary_key::PrimaryKeyIndex::from_persisted(
686 field,
687 bloom,
688 pk_data,
689 &[],
690 snapshot,
691 )
692 })
693 .await
694 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
695 };
696
697 if needs_persist {
698 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
699 }
700
701 *self.primary_key_index.write() = Some(pk_index);
702 } else {
703 let pk_index = tokio::task::spawn_blocking(move || {
705 super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
706 })
707 .await
708 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
709
710 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
711 *self.primary_key_index.write() = Some(pk_index);
712 }
713
714 self.pk_reservations_retained
718 .store(false, Ordering::Release);
719
720 Ok(())
721 }
722
723 async fn persist_pk_bloom(
726 &self,
727 pk_index: &super::primary_key::PrimaryKeyIndex,
728 segment_ids: &[String],
729 ) {
730 use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
731
732 let bloom_bytes = pk_index.bloom_to_bytes();
733 let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
734 if let Err(e) = self
735 .directory
736 .write(std::path::Path::new(PK_BLOOM_FILE), &data)
737 .await
738 {
739 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
740 }
741 }
742
743 pub fn add_document(&self, doc: Document) -> Result<()> {
749 self.ensure_writer_lock()?;
750 if self.worker_state.shutdown.load(Ordering::Acquire) {
751 return Err(Error::IndexClosed);
752 }
753 if self.commit_finalization.in_progress.load(Ordering::Acquire) {
754 return Err(Error::CommitInProgress);
755 }
756 let sender = self.doc_sender.read().clone();
757 if sender.is_closed() {
761 return Err(Error::CommitInProgress);
762 }
763 let primary_key_index = self.primary_key_index.read();
764 if let Some(ref pk_index) = *primary_key_index {
765 pk_index.check_and_insert(&doc)?;
766 }
767 match sender.try_send(doc) {
768 Ok(()) => Ok(()),
769 Err(async_channel::TrySendError::Full(doc)) => {
770 if let Some(ref pk_index) = *primary_key_index {
772 pk_index.rollback_uncommitted_key(&doc);
773 }
774 Err(Error::QueueFull)
775 }
776 Err(async_channel::TrySendError::Closed(doc)) => {
777 if let Some(ref pk_index) = *primary_key_index {
779 pk_index.rollback_uncommitted_key(&doc);
780 }
781 Err(Error::CommitInProgress)
782 }
783 }
784 }
785
786 pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
791 let total = documents.len();
792 for (i, doc) in documents.into_iter().enumerate() {
793 match self.add_document(doc) {
794 Ok(()) => {}
795 Err(Error::QueueFull | Error::CommitInProgress) => return Ok(i),
796 Err(e) => return Err(e),
797 }
798 }
799 Ok(total)
800 }
801
802 fn worker_loop(
815 state: Arc<WorkerState<D>>,
816 initial_receiver: async_channel::Receiver<Document>,
817 handle: tokio::runtime::Handle,
818 ) {
819 let mut receiver = initial_receiver;
820 let mut my_epoch = 0usize;
821
822 loop {
823 let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
827 let mut builder: Option<SegmentBuilder> = None;
828
829 while let Ok(doc) = receiver.recv_blocking() {
830 if state.shutdown.load(Ordering::Acquire) {
831 break;
832 }
833 if state.cycle_failed.load(Ordering::Acquire) {
838 continue;
839 }
840 if builder.is_none() {
842 match SegmentBuilder::new(
843 Arc::clone(&state.schema),
844 state.builder_config.clone(),
845 ) {
846 Ok(mut b) => {
847 for (field, tokenizer) in state.tokenizers.read().iter() {
848 b.set_tokenizer(*field, tokenizer.clone_box());
849 }
850 builder = Some(b);
851 }
852 Err(e) => {
853 log::error!("Failed to create segment builder: {:?}", e);
854 state.record_cycle_error(format!(
855 "failed to create segment builder: {e}"
856 ));
857 continue;
858 }
859 }
860 }
861
862 let b = builder.as_mut().unwrap();
863 if let Err(e) = b.add_document(doc) {
864 log::error!("Failed to index document: {:?}", e);
865 state.record_cycle_error(format!("failed to index document: {e}"));
866 continue;
867 }
868
869 let builder_memory = b.estimated_memory_bytes();
870
871 if b.num_docs() & 0x3FFF == 0 {
872 log::debug!(
873 "[indexing] docs={}, memory={:.2} MB, budget={:.2} MB",
874 b.num_docs(),
875 builder_memory as f64 / (1024.0 * 1024.0),
876 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0)
877 );
878 }
879
880 const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
882
883 let effective_budget = state.memory_budget_per_worker * 4 / 5;
887
888 if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
889 log::info!(
890 "[indexing] memory budget reached, building segment: \
891 docs={}, memory={:.2} MB, budget={:.2} MB",
892 b.num_docs(),
893 builder_memory as f64 / (1024.0 * 1024.0),
894 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0),
895 );
896 let full_builder = builder.take().unwrap();
897 Self::build_segment_inline(&state, full_builder, &handle);
898 }
899 }
900
901 if !state.cycle_failed.load(Ordering::Acquire)
903 && let Some(b) = builder.take()
904 && b.num_docs() > 0
905 {
906 Self::build_segment_inline(&state, b, &handle);
907 }
908 }));
909
910 if build_result.is_err() {
911 log::error!(
912 "[worker] panic during indexing cycle — documents in this cycle may be lost"
913 );
914 state.record_cycle_error("indexing worker panicked while building the batch");
915 }
916
917 let prev = state.flush_count.fetch_add(1, Ordering::Release);
920 if prev + 1 == state.num_workers {
921 let _lock = state.flush_mutex.lock();
928 state.flush_cvar.notify_all();
929 }
930
931 {
935 let mut lock = state.resume_receiver.lock();
936 loop {
937 if state.shutdown.load(Ordering::Acquire) {
938 return;
939 }
940 let current_epoch = state.resume_epoch.load(Ordering::Acquire);
941 if current_epoch > my_epoch
942 && let Some(rx) = lock.as_ref()
943 {
944 receiver = rx.clone();
945 my_epoch = current_epoch;
946 break;
947 }
948 state.resume_cvar.wait(&mut lock);
949 }
950 }
951 }
952 }
953
954 fn build_segment_inline(
958 state: &WorkerState<D>,
959 builder: SegmentBuilder,
960 handle: &tokio::runtime::Handle,
961 ) {
962 let segment_id = SegmentId::new();
963 let segment_hex = segment_id.to_hex();
964 let operation = match state
967 .segment_manager
968 .protect_new_segment(segment_hex.clone())
969 {
970 Ok(operation) => operation,
971 Err(e) => {
972 log::error!(
973 "[segment_build_failed] segment_id={} lifecycle_error={}",
974 segment_hex,
975 e,
976 );
977 state.record_cycle_error(format!(
978 "failed to claim segment {segment_hex} for building: {e}"
979 ));
980 return;
981 }
982 };
983 let trained = state.segment_manager.trained_for_segment_build();
984 let doc_count = builder.num_docs();
985 let build_start = std::time::Instant::now();
986
987 log::info!(
988 "[segment_build] segment_id={} doc_count={} ann={}",
989 segment_hex,
990 doc_count,
991 trained.is_some()
992 );
993
994 let mut prepared = PreparedSegment {
998 id: segment_hex.clone(),
999 segment_id,
1000 num_docs: doc_count,
1001 segment_manager: Arc::clone(&state.segment_manager),
1002 operation: Some(operation),
1003 runtime: handle.clone(),
1004 published: false,
1005 };
1006
1007 match handle.block_on(builder.build(
1008 state.directory.as_ref(),
1009 segment_id,
1010 trained.as_deref(),
1011 )) {
1012 Ok(meta) if meta.num_docs == doc_count && meta.num_docs > 0 => {
1013 let duration_ms = build_start.elapsed().as_millis() as u64;
1014 log::info!(
1015 "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
1016 segment_hex,
1017 meta.num_docs,
1018 duration_ms,
1019 );
1020 prepared.num_docs = meta.num_docs;
1021 state.built_segments.lock().push(prepared);
1022 }
1023 Ok(meta) => {
1024 let error = format!(
1025 "segment {segment_hex} built {} docs from a {doc_count}-document builder",
1026 meta.num_docs
1027 );
1028 log::error!("[segment_build_failed] {error}");
1029 state.record_cycle_error(error);
1030 }
1031 Err(e) => {
1032 log::error!(
1033 "[segment_build_failed] segment_id={} error={:?}",
1034 segment_hex,
1035 e
1036 );
1037 state.record_cycle_error(format!("failed to build segment {segment_hex}: {e}"));
1040 }
1041 }
1042 }
1043
1044 pub async fn maybe_merge(&self) {
1050 self.segment_manager.maybe_merge().await;
1051 }
1052
1053 pub async fn abort_merges(&self) {
1056 self.segment_manager.abort_merges().await;
1057 }
1058
1059 pub async fn shutdown(&mut self) -> Result<()> {
1064 self.segment_manager.begin_shutdown();
1065 self.signal_worker_shutdown();
1066
1067 self.commit_finalization.wait_until_idle().await;
1072
1073 let workers = std::mem::take(&mut self.workers);
1074 let panicked = tokio::task::spawn_blocking(move || {
1075 workers
1076 .into_iter()
1077 .map(|worker| worker.join().is_err())
1078 .filter(|panicked| *panicked)
1079 .count()
1080 })
1081 .await
1082 .map_err(|error| Error::Internal(format!("failed to join index workers: {}", error)))?;
1083 if panicked > 0 {
1084 log::error!("[index_shutdown] {} indexing worker(s) panicked", panicked);
1085 }
1086
1087 self.flushed_segments.lock().clear();
1090 self.worker_state.built_segments.lock().clear();
1091 if let Some(pk_index) = self.primary_key_index.write().as_mut() {
1092 pk_index.clear_uncommitted();
1093 }
1094 Ok(())
1095 }
1096
1097 pub async fn wait_for_merging_thread(&self) {
1099 self.segment_manager.wait_for_merging_thread().await;
1100 }
1101
1102 pub async fn wait_for_all_merges(&self) {
1104 self.segment_manager.wait_for_all_merges().await;
1105 }
1106
1107 pub async fn wait_for_commit_finalization(&self) {
1112 self.commit_finalization.wait_until_idle().await;
1113 }
1114
1115 pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
1117 self.segment_manager.tracker()
1118 }
1119
1120 pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
1122 self.segment_manager.acquire_snapshot().await
1123 }
1124
1125 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
1130 self.ensure_writer_lock()?;
1131 self.segment_manager.cleanup_orphan_segments().await
1132 }
1133
1134 pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
1145 self.ensure_writer_lock()?;
1146 if self.worker_state.shutdown.load(Ordering::Acquire) {
1147 return Err(Error::IndexClosed);
1148 }
1149 if self.commit_finalization.in_progress.load(Ordering::Acquire) {
1150 return Err(Error::CommitInProgress);
1151 }
1152 self.doc_sender.read().close();
1154
1155 self.worker_state.resume_cvar.notify_all();
1159
1160 let state = Arc::clone(&self.worker_state);
1163 let all_flushed = tokio::task::spawn_blocking(move || {
1164 let mut lock = state.flush_mutex.lock();
1165 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
1166 while state.flush_count.load(Ordering::Acquire) < state.num_workers {
1167 let remaining = deadline.saturating_duration_since(std::time::Instant::now());
1168 if remaining.is_zero() {
1169 log::error!(
1170 "[prepare_commit] timed out waiting for workers: {}/{} flushed",
1171 state.flush_count.load(Ordering::Acquire),
1172 state.num_workers
1173 );
1174 return false;
1175 }
1176 state.flush_cvar.wait_for(&mut lock, remaining);
1177 }
1178 true
1179 })
1180 .await
1181 .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
1182
1183 if !all_flushed {
1184 return Err(Error::Internal(format!(
1192 "prepare_commit timed out: {}/{} workers flushed; writer remains paused, retry commit",
1193 self.worker_state.flush_count.load(Ordering::Acquire),
1194 self.worker_state.num_workers
1195 )));
1196 }
1197
1198 let cycle_error = { self.worker_state.cycle_error.lock().take() };
1199 if let Some(error) = cycle_error {
1200 self.flushed_segments.lock().clear();
1205 self.worker_state.built_segments.lock().clear();
1206 self.clear_uncommitted_pk_reservations();
1207 self.resume_workers();
1208 return Err(Error::Internal(format!(
1209 "indexing generation failed; no documents from this batch were committed: {error}"
1210 )));
1211 }
1212
1213 let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
1215 self.flushed_segments.lock().extend(built);
1216
1217 Ok(PreparedCommit {
1218 writer: self,
1219 is_resolved: false,
1220 })
1221 }
1222
1223 pub async fn commit(&mut self) -> Result<bool> {
1228 self.prepare_commit().await?.commit().await
1229 }
1230
1231 pub async fn force_merge(&mut self) -> Result<()> {
1233 self.prepare_commit().await?.commit().await?;
1234 self.segment_manager.force_merge().await
1235 }
1236
1237 pub async fn reorder(&mut self) -> Result<()> {
1242 self.prepare_commit().await?.commit().await?;
1243 self.segment_manager.reorder_segments().await
1244 }
1245
1246 pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
1248 &self.segment_manager
1249 }
1250
1251 fn resume_workers(&mut self) {
1256 Self::resume_workers_shared(&self.worker_state, &self.doc_sender);
1257 }
1258
1259 fn resume_workers_shared(
1260 worker_state: &Arc<WorkerState<D>>,
1261 doc_sender: &Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1262 ) {
1263 if worker_state.shutdown.load(Ordering::Acquire) {
1264 return;
1265 }
1266 if tokio::runtime::Handle::try_current().is_err() {
1267 worker_state.shutdown.store(true, Ordering::Release);
1270 worker_state.resume_cvar.notify_all();
1271 return;
1272 }
1273
1274 worker_state.flush_count.store(0, Ordering::Release);
1276 *worker_state.cycle_error.lock() = None;
1277 worker_state.cycle_failed.store(false, Ordering::Release);
1278
1279 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
1281 *doc_sender.write() = sender;
1282
1283 {
1285 let mut lock = worker_state.resume_receiver.lock();
1286 *lock = Some(receiver);
1287 }
1288 worker_state.resume_epoch.fetch_add(1, Ordering::Release);
1289 worker_state.resume_cvar.notify_all();
1290 }
1291
1292 fn signal_worker_shutdown(&self) {
1293 self.worker_state.shutdown.store(true, Ordering::Release);
1294 self.doc_sender.read().close();
1295 self.worker_state.resume_cvar.notify_all();
1296 }
1297
1298 }
1300
1301impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
1302 fn drop(&mut self) {
1303 self.signal_worker_shutdown();
1304 for w in std::mem::take(&mut self.workers) {
1305 let _ = w.join();
1306 }
1307 }
1308}
1309
1310pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
1317 writer: &'a mut IndexWriter<D>,
1318 is_resolved: bool,
1319}
1320
1321struct PreparedSegmentsGuard<D: DirectoryWriter + 'static> {
1326 segments: Option<Vec<PreparedSegment<D>>>,
1327 retry_slot: Arc<parking_lot::Mutex<Vec<PreparedSegment<D>>>>,
1328}
1329
1330impl<D: DirectoryWriter + 'static> PreparedSegmentsGuard<D> {
1331 fn metadata_entries(&self) -> Vec<(String, u32)> {
1332 self.segments
1333 .as_deref()
1334 .unwrap_or_default()
1335 .iter()
1336 .map(PreparedSegment::metadata_entry)
1337 .collect()
1338 }
1339
1340 fn take_published(&mut self) -> Vec<PreparedSegment<D>> {
1341 self.segments.take().unwrap_or_default()
1342 }
1343}
1344
1345impl<D: DirectoryWriter + 'static> Drop for PreparedSegmentsGuard<D> {
1346 fn drop(&mut self) {
1347 if let Some(segments) = self.segments.take() {
1348 self.retry_slot.lock().extend(segments);
1349 }
1350 }
1351}
1352
1353struct CommitFinalizationGuard<D: DirectoryWriter + 'static> {
1358 state: Arc<CommitFinalizationState>,
1359 worker_state: Arc<WorkerState<D>>,
1360 doc_sender: Arc<parking_lot::RwLock<async_channel::Sender<Document>>>,
1361 resume_workers: bool,
1362}
1363
1364impl<D: DirectoryWriter + 'static> CommitFinalizationGuard<D> {
1365 fn resume_on_drop(&mut self) {
1366 self.resume_workers = true;
1367 }
1368}
1369
1370impl<D: DirectoryWriter + 'static> Drop for CommitFinalizationGuard<D> {
1371 fn drop(&mut self) {
1372 if self.resume_workers {
1373 IndexWriter::<D>::resume_workers_shared(&self.worker_state, &self.doc_sender);
1374 }
1375 self.state.finish();
1376 }
1377}
1378
1379struct OwnedCommitFinalization<D: DirectoryWriter + 'static> {
1384 directory: Arc<D>,
1385 schema: Arc<Schema>,
1386 segment_manager: Arc<crate::merge::SegmentManager<D>>,
1387 primary_key_index: Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1388 prepared: PreparedSegmentsGuard<D>,
1389 finalization: Option<CommitFinalizationGuard<D>>,
1390 publication_observed: Arc<AtomicBool>,
1391 pk_reservations_retained: Arc<AtomicBool>,
1392}
1393
1394async fn refresh_primary_key_after_commit<D: DirectoryWriter + 'static>(
1395 directory: &Arc<D>,
1396 schema: &Arc<Schema>,
1397 segment_manager: &Arc<crate::merge::SegmentManager<D>>,
1398 primary_key_index: &Arc<parking_lot::RwLock<Option<super::primary_key::PrimaryKeyIndex>>>,
1399) -> Result<()> {
1400 let existing_ids: std::collections::HashSet<String> = {
1401 let guard = primary_key_index.read();
1402 let Some(pk_index) = guard.as_ref() else {
1403 return Ok(());
1404 };
1405 pk_index
1406 .committed_segment_ids()
1407 .map(ToOwned::to_owned)
1408 .collect()
1409 };
1410
1411 let snapshot = segment_manager.acquire_snapshot().await;
1412 let load_futures: Vec<_> = snapshot
1413 .segment_ids()
1414 .iter()
1415 .filter(|id| !existing_ids.contains(id.as_str()))
1416 .map(|seg_id_str| {
1417 let seg_id_str = seg_id_str.clone();
1418 let dir = directory.as_ref();
1419 let schema = Arc::clone(schema);
1420 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
1421 })
1422 .collect();
1423 let new_data = futures::future::try_join_all(load_futures).await?;
1424 let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
1425
1426 let bloom_file = {
1427 let mut guard = primary_key_index.write();
1428 let Some(pk_index) = guard.as_mut() else {
1429 return Ok(());
1430 };
1431 pk_index.refresh_incremental(new_data, snapshot);
1432 let bloom_bytes = pk_index.bloom_to_bytes();
1433 super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes)
1434 };
1435
1436 if let Err(error) = directory
1437 .write(
1438 std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
1439 &bloom_file,
1440 )
1441 .await
1442 {
1443 log::warn!("[primary_key] failed to persist bloom cache: {}", error);
1444 }
1445 Ok(())
1446}
1447
1448async fn finalize_prepared_commit<D: DirectoryWriter + 'static>(
1449 mut commit: OwnedCommitFinalization<D>,
1450) -> Result<bool> {
1451 let metadata_entries = commit.prepared.metadata_entries();
1452
1453 commit.segment_manager.commit(&metadata_entries).await?;
1457 commit.publication_observed.store(true, Ordering::Release);
1458
1459 let mut published = commit.prepared.take_published();
1460 for segment in &mut published {
1461 segment.mark_published();
1462 }
1463 drop(published);
1464 if let Some(finalization) = commit.finalization.as_mut() {
1468 finalization.resume_on_drop();
1469 } else {
1470 log::error!("owned commit finalization guard was already released after publication");
1471 }
1472
1473 match refresh_primary_key_after_commit(
1478 &commit.directory,
1479 &commit.schema,
1480 &commit.segment_manager,
1481 &commit.primary_key_index,
1482 )
1483 .await
1484 {
1485 Ok(()) => commit
1488 .pk_reservations_retained
1489 .store(false, Ordering::Release),
1490 Err(error) => {
1491 commit
1496 .pk_reservations_retained
1497 .store(true, Ordering::Release);
1498 log::error!(
1499 "[primary_key] committed metadata but failed to refresh dedup state; \
1500 retaining reservations until a later successful commit: {}",
1501 error,
1502 );
1503 }
1504 }
1505
1506 drop(commit.finalization.take());
1510 commit.segment_manager.maybe_merge().await;
1511 Ok(true)
1512}
1513
1514impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
1515 pub async fn commit(mut self) -> Result<bool> {
1519 let segments = std::mem::take(&mut *self.writer.flushed_segments.lock());
1520
1521 if segments.is_empty() {
1523 log::debug!("[commit] no segments to commit, skipping");
1524 self.is_resolved = true;
1525 self.writer.resume_workers();
1526 return Ok(false);
1527 }
1528
1529 if !self.writer.commit_finalization.begin() {
1530 self.writer.flushed_segments.lock().extend(segments);
1531 self.is_resolved = true;
1535 return Err(Error::CommitInProgress);
1536 }
1537
1538 let publication_observed = Arc::new(AtomicBool::new(false));
1539 let owned = OwnedCommitFinalization {
1540 directory: Arc::clone(&self.writer.directory),
1541 schema: Arc::clone(&self.writer.schema),
1542 segment_manager: Arc::clone(&self.writer.segment_manager),
1543 primary_key_index: Arc::clone(&self.writer.primary_key_index),
1544 prepared: PreparedSegmentsGuard {
1545 segments: Some(segments),
1546 retry_slot: Arc::clone(&self.writer.flushed_segments),
1547 },
1548 finalization: Some(CommitFinalizationGuard {
1549 state: Arc::clone(&self.writer.commit_finalization),
1550 worker_state: Arc::clone(&self.writer.worker_state),
1551 doc_sender: Arc::clone(&self.writer.doc_sender),
1552 resume_workers: false,
1553 }),
1554 publication_observed: Arc::clone(&publication_observed),
1555 pk_reservations_retained: Arc::clone(&self.writer.pk_reservations_retained),
1556 };
1557
1558 self.is_resolved = true;
1563 let task_publication = Arc::clone(&publication_observed);
1564 let task = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
1565 tokio::spawn(async move {
1566 match std::panic::AssertUnwindSafe(finalize_prepared_commit(owned))
1567 .catch_unwind()
1568 .await
1569 {
1570 Ok(result) => result,
1571 Err(_) if task_publication.load(Ordering::Acquire) => {
1572 log::error!(
1573 "owned commit finalizer panicked after metadata publication; \
1574 treating the durable generation as committed"
1575 );
1576 Ok(true)
1577 }
1578 Err(_) => Err(Error::Internal(
1579 "owned commit finalizer panicked before metadata publication".into(),
1580 )),
1581 }
1582 })
1583 }))
1584 .map_err(|_| Error::Internal("runtime rejected owned commit finalizer".into()))?;
1585
1586 match task.await {
1587 Ok(result) => result,
1588 Err(error) if publication_observed.load(Ordering::Acquire) => {
1589 log::error!(
1590 "owned commit finalizer terminated after metadata publication: {}; \
1591 treating the durable generation as committed",
1592 error,
1593 );
1594 Ok(true)
1595 }
1596 Err(error) => Err(Error::Internal(format!(
1597 "owned commit finalizer terminated unexpectedly: {error}"
1598 ))),
1599 }
1600 }
1601
1602 pub fn abort(mut self) {
1605 self.is_resolved = true;
1606 self.writer.flushed_segments.lock().clear();
1607 self.writer.clear_uncommitted_pk_reservations();
1608 self.writer.resume_workers();
1609 }
1610}
1611
1612impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
1613 fn drop(&mut self) {
1614 if !self.is_resolved {
1615 log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
1616 self.writer.flushed_segments.lock().clear();
1617 self.writer.clear_uncommitted_pk_reservations();
1618 self.writer.resume_workers();
1619 }
1620 }
1621}
1622
1623async fn load_pk_segment_data<D: crate::directories::Directory>(
1625 dir: &D,
1626 seg_id_str: &str,
1627 schema: &Arc<crate::dsl::Schema>,
1628) -> Result<super::primary_key::PkSegmentData> {
1629 let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
1630 .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
1631 let files = crate::segment::SegmentFiles::new(seg_id.0);
1632 let fast_fields =
1633 crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
1634 Ok(super::primary_key::PkSegmentData {
1635 segment_id: seg_id_str.to_string(),
1636 fast_fields,
1637 })
1638}