1use std::sync::Arc;
33use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
34
35use rustc_hash::FxHashMap;
36
37use crate::directories::DirectoryWriter;
38use crate::dsl::{Document, Field, Schema};
39use crate::error::{Error, Result};
40use crate::segment::{SegmentBuilder, SegmentBuilderConfig, SegmentId};
41use crate::tokenizer::BoxedTokenizer;
42
43use super::IndexConfig;
44
45const PIPELINE_MAX_SIZE_IN_DOCS: usize = 10_000;
47
48pub struct IndexWriter<D: DirectoryWriter + 'static> {
60 pub(super) directory: Arc<D>,
61 pub(super) schema: Arc<Schema>,
62 pub(super) config: IndexConfig,
63 doc_sender: async_channel::Sender<Document>,
66 workers: Vec<std::thread::JoinHandle<()>>,
68 worker_state: Arc<WorkerState<D>>,
70 pub(super) segment_manager: Arc<crate::merge::SegmentManager<D>>,
72 flushed_segments: Vec<PreparedSegment>,
75 primary_key_index: Option<super::primary_key::PrimaryKeyIndex>,
77}
78
79struct WorkerState<D: DirectoryWriter + 'static> {
81 directory: Arc<D>,
82 schema: Arc<Schema>,
83 builder_config: SegmentBuilderConfig,
84 tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
85 memory_budget_per_worker: usize,
87 segment_manager: Arc<crate::merge::SegmentManager<D>>,
89 built_segments: parking_lot::Mutex<Vec<PreparedSegment>>,
92
93 flush_count: AtomicUsize,
100 flush_mutex: parking_lot::Mutex<()>,
102 flush_cvar: parking_lot::Condvar,
103 resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
105 resume_epoch: AtomicUsize,
108 resume_cvar: parking_lot::Condvar,
110 shutdown: AtomicBool,
112 num_workers: usize,
114}
115
116struct PreparedSegment {
122 id: String,
123 num_docs: u32,
124 _operation: crate::merge::SegmentOperationGuard,
125}
126
127impl PreparedSegment {
128 fn metadata_entry(&self) -> (String, u32) {
129 (self.id.clone(), self.num_docs)
130 }
131}
132
133impl<D: DirectoryWriter + 'static> IndexWriter<D> {
134 pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
136 Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
137 }
138
139 pub async fn create_with_config(
141 directory: D,
142 schema: Schema,
143 config: IndexConfig,
144 builder_config: SegmentBuilderConfig,
145 ) -> Result<Self> {
146 let directory = Arc::new(directory);
147 let schema = Arc::new(schema);
148 directory.set_index_label(schema.index_label());
150 let metadata = super::IndexMetadata::new((*schema).clone());
151
152 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
153 Arc::clone(&directory),
154 Arc::clone(&schema),
155 metadata,
156 config.merge_policy.clone_box(),
157 config.term_cache_blocks,
158 config.max_concurrent_merges,
159 config.merge_bp_time_budget,
160 config.bp_memory_budget_bytes,
161 ));
162 segment_manager.update_metadata(|_| {}).await?;
163
164 Ok(Self::new_with_parts(
165 directory,
166 schema,
167 config,
168 builder_config,
169 segment_manager,
170 ))
171 }
172
173 pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
175 Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
176 }
177
178 pub async fn open_with_config(
180 directory: D,
181 config: IndexConfig,
182 builder_config: SegmentBuilderConfig,
183 ) -> Result<Self> {
184 let directory = Arc::new(directory);
185 let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
186 let schema = Arc::new(metadata.schema.clone());
187 directory.set_index_label(schema.index_label());
189
190 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
191 Arc::clone(&directory),
192 Arc::clone(&schema),
193 metadata,
194 config.merge_policy.clone_box(),
195 config.term_cache_blocks,
196 config.max_concurrent_merges,
197 config.merge_bp_time_budget,
198 config.bp_memory_budget_bytes,
199 ));
200 segment_manager.load_and_publish_trained().await;
201
202 Ok(Self::new_with_parts(
203 directory,
204 schema,
205 config,
206 builder_config,
207 segment_manager,
208 ))
209 }
210
211 pub fn from_index(index: &super::Index<D>) -> Self {
214 Self::new_with_parts(
215 Arc::clone(&index.directory),
216 Arc::clone(&index.schema),
217 index.config.clone(),
218 SegmentBuilderConfig::default(),
219 Arc::clone(&index.segment_manager),
220 )
221 }
222
223 fn new_with_parts(
229 directory: Arc<D>,
230 schema: Arc<Schema>,
231 config: IndexConfig,
232 builder_config: SegmentBuilderConfig,
233 segment_manager: Arc<crate::merge::SegmentManager<D>>,
234 ) -> Self {
235 let registry = crate::tokenizer::TokenizerRegistry::new();
237 let mut tokenizers = FxHashMap::default();
238 for (field, entry) in schema.fields() {
239 if matches!(entry.field_type, crate::dsl::FieldType::Text)
240 && let Some(ref tok_name) = entry.tokenizer
241 && let Some(tok) = registry.get(tok_name)
242 {
243 tokenizers.insert(field, tok);
244 }
245 }
246
247 let num_workers = config.num_indexing_threads.max(1);
248 let worker_state = Arc::new(WorkerState {
249 directory: Arc::clone(&directory),
250 schema: Arc::clone(&schema),
251 builder_config,
252 tokenizers: parking_lot::RwLock::new(tokenizers),
253 memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
254 segment_manager: Arc::clone(&segment_manager),
255 built_segments: parking_lot::Mutex::new(Vec::new()),
256 flush_count: AtomicUsize::new(0),
257 flush_mutex: parking_lot::Mutex::new(()),
258 flush_cvar: parking_lot::Condvar::new(),
259 resume_receiver: parking_lot::Mutex::new(None),
260 resume_epoch: AtomicUsize::new(0),
261 resume_cvar: parking_lot::Condvar::new(),
262 shutdown: AtomicBool::new(false),
263 num_workers,
264 });
265 let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
266
267 Self {
268 directory,
269 schema,
270 config,
271 doc_sender,
272 workers,
273 worker_state,
274 segment_manager,
275 flushed_segments: Vec::new(),
276 primary_key_index: None,
277 }
278 }
279
280 fn spawn_workers(
281 worker_state: &Arc<WorkerState<D>>,
282 num_workers: usize,
283 ) -> (
284 async_channel::Sender<Document>,
285 Vec<std::thread::JoinHandle<()>>,
286 ) {
287 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
288 let handle = tokio::runtime::Handle::current();
289 let mut workers = Vec::with_capacity(num_workers);
290 for i in 0..num_workers {
291 let state = Arc::clone(worker_state);
292 let rx = receiver.clone();
293 let rt = handle.clone();
294 workers.push(
295 std::thread::Builder::new()
296 .name(format!("index-worker-{}", i))
297 .spawn(move || Self::worker_loop(state, rx, rt))
298 .expect("failed to spawn index worker thread"),
299 );
300 }
301 (sender, workers)
302 }
303
304 pub fn schema(&self) -> &Schema {
306 &self.schema
307 }
308
309 pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
312 self.worker_state
313 .tokenizers
314 .write()
315 .insert(field, Box::new(tokenizer));
316 }
317
318 pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
334 use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
335
336 let field = match self.schema.primary_field() {
337 Some(f) => f,
338 None => return Ok(()),
339 };
340
341 let snapshot = self.segment_manager.acquire_snapshot().await;
342 let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
343
344 let cached = match self
346 .directory
347 .open_read(std::path::Path::new(PK_BLOOM_FILE))
348 .await
349 {
350 Ok(handle) => {
351 let data = handle.read_bytes_range(0..handle.len()).await;
352 match data {
353 Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
354 Err(_) => None,
355 }
356 }
357 Err(_) => None,
358 };
359
360 let load_futures: Vec<_> = current_seg_ids
362 .iter()
363 .map(|seg_id_str| {
364 let seg_id_str = seg_id_str.clone();
365 let dir = self.directory.as_ref();
366 let schema = Arc::clone(&self.schema);
367 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
368 })
369 .collect();
370 let all_data = futures::future::try_join_all(load_futures).await?;
371
372 if let Some((persisted_seg_ids, bloom)) = cached {
373 let mut pk_data = Vec::with_capacity(all_data.len());
375 let mut new_data = Vec::new();
376 for d in all_data {
377 if persisted_seg_ids.contains(&d.segment_id) {
378 pk_data.push(d);
379 } else {
380 new_data.push(d);
381 }
382 }
383 let needs_persist = !new_data.is_empty();
384 let new_start = pk_data.len();
385 pk_data.extend(new_data);
386
387 let pk_index = if new_start == pk_data.len() {
388 super::primary_key::PrimaryKeyIndex::from_persisted(
390 field,
391 bloom,
392 pk_data,
393 &[],
394 snapshot,
395 )
396 } else {
397 tokio::task::spawn_blocking(move || {
399 let mut bloom = bloom;
402 let mut added = 0usize;
403 let num_new = pk_data.len() - new_start;
404 for data in &pk_data[new_start..] {
405 if let Some(ff) = data.fast_fields.get(&field.0)
406 && let Some(dict) = ff.text_dict()
407 {
408 for key in dict.iter() {
409 bloom.insert(key.as_bytes());
410 added += 1;
411 }
412 }
413 }
414 if added > 0 {
415 log::info!(
416 "[primary_key] bloom: added {} keys from {} new segment(s)",
417 added,
418 num_new,
419 );
420 }
421 super::primary_key::PrimaryKeyIndex::from_persisted(
422 field,
423 bloom,
424 pk_data,
425 &[],
426 snapshot,
427 )
428 })
429 .await
430 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
431 };
432
433 if needs_persist {
434 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
435 }
436
437 self.primary_key_index = Some(pk_index);
438 } else {
439 let pk_index = tokio::task::spawn_blocking(move || {
441 super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
442 })
443 .await
444 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
445
446 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
447 self.primary_key_index = Some(pk_index);
448 }
449
450 Ok(())
451 }
452
453 async fn persist_pk_bloom(
456 &self,
457 pk_index: &super::primary_key::PrimaryKeyIndex,
458 segment_ids: &[String],
459 ) {
460 use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
461
462 let bloom_bytes = pk_index.bloom_to_bytes();
463 let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
464 if let Err(e) = self
465 .directory
466 .write(std::path::Path::new(PK_BLOOM_FILE), &data)
467 .await
468 {
469 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
470 }
471 }
472
473 pub fn add_document(&self, doc: Document) -> Result<()> {
478 if let Some(ref pk_index) = self.primary_key_index {
479 pk_index.check_and_insert(&doc)?;
480 }
481 match self.doc_sender.try_send(doc) {
482 Ok(()) => Ok(()),
483 Err(async_channel::TrySendError::Full(doc)) => {
484 if let Some(ref pk_index) = self.primary_key_index {
486 pk_index.rollback_uncommitted_key(&doc);
487 }
488 Err(Error::QueueFull)
489 }
490 Err(async_channel::TrySendError::Closed(doc)) => {
491 if let Some(ref pk_index) = self.primary_key_index {
493 pk_index.rollback_uncommitted_key(&doc);
494 }
495 Err(Error::Internal("Document channel closed".into()))
496 }
497 }
498 }
499
500 pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
505 let total = documents.len();
506 for (i, doc) in documents.into_iter().enumerate() {
507 match self.add_document(doc) {
508 Ok(()) => {}
509 Err(Error::QueueFull) => return Ok(i),
510 Err(e) => return Err(e),
511 }
512 }
513 Ok(total)
514 }
515
516 fn worker_loop(
529 state: Arc<WorkerState<D>>,
530 initial_receiver: async_channel::Receiver<Document>,
531 handle: tokio::runtime::Handle,
532 ) {
533 let mut receiver = initial_receiver;
534 let mut my_epoch = 0usize;
535
536 loop {
537 let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
541 let mut builder: Option<SegmentBuilder> = None;
542
543 while let Ok(doc) = receiver.recv_blocking() {
544 if builder.is_none() {
546 match SegmentBuilder::new(
547 Arc::clone(&state.schema),
548 state.builder_config.clone(),
549 ) {
550 Ok(mut b) => {
551 for (field, tokenizer) in state.tokenizers.read().iter() {
552 b.set_tokenizer(*field, tokenizer.clone_box());
553 }
554 builder = Some(b);
555 }
556 Err(e) => {
557 log::error!("Failed to create segment builder: {:?}", e);
558 continue;
559 }
560 }
561 }
562
563 let b = builder.as_mut().unwrap();
564 if let Err(e) = b.add_document(doc) {
565 log::error!("Failed to index document: {:?}", e);
566 continue;
567 }
568
569 let builder_memory = b.estimated_memory_bytes();
570
571 if b.num_docs() & 0x3FFF == 0 {
572 log::debug!(
573 "[indexing] docs={}, memory={:.2} MB, budget={:.2} MB",
574 b.num_docs(),
575 builder_memory as f64 / (1024.0 * 1024.0),
576 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0)
577 );
578 }
579
580 const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
582
583 let effective_budget = state.memory_budget_per_worker * 4 / 5;
587
588 if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
589 log::info!(
590 "[indexing] memory budget reached, building segment: \
591 docs={}, memory={:.2} MB, budget={:.2} MB",
592 b.num_docs(),
593 builder_memory as f64 / (1024.0 * 1024.0),
594 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0),
595 );
596 let full_builder = builder.take().unwrap();
597 Self::build_segment_inline(&state, full_builder, &handle);
598 }
599 }
600
601 if let Some(b) = builder.take()
603 && b.num_docs() > 0
604 {
605 Self::build_segment_inline(&state, b, &handle);
606 }
607 }));
608
609 if build_result.is_err() {
610 log::error!(
611 "[worker] panic during indexing cycle — documents in this cycle may be lost"
612 );
613 }
614
615 let prev = state.flush_count.fetch_add(1, Ordering::Release);
618 if prev + 1 == state.num_workers {
619 let _lock = state.flush_mutex.lock();
621 state.flush_cvar.notify_one();
622 }
623
624 {
628 let mut lock = state.resume_receiver.lock();
629 loop {
630 if state.shutdown.load(Ordering::Acquire) {
631 return;
632 }
633 let current_epoch = state.resume_epoch.load(Ordering::Acquire);
634 if current_epoch > my_epoch
635 && let Some(rx) = lock.as_ref()
636 {
637 receiver = rx.clone();
638 my_epoch = current_epoch;
639 break;
640 }
641 state.resume_cvar.wait(&mut lock);
642 }
643 }
644 }
645 }
646
647 fn build_segment_inline(
651 state: &WorkerState<D>,
652 builder: SegmentBuilder,
653 handle: &tokio::runtime::Handle,
654 ) {
655 let segment_id = SegmentId::new();
656 let segment_hex = segment_id.to_hex();
657 let operation = match state
660 .segment_manager
661 .protect_new_segment(segment_hex.clone())
662 {
663 Ok(operation) => operation,
664 Err(e) => {
665 log::error!(
666 "[segment_build_failed] segment_id={} lifecycle_error={}",
667 segment_hex,
668 e,
669 );
670 return;
671 }
672 };
673 let trained = state.segment_manager.trained();
674 let doc_count = builder.num_docs();
675 let build_start = std::time::Instant::now();
676
677 log::info!(
678 "[segment_build] segment_id={} doc_count={} ann={}",
679 segment_hex,
680 doc_count,
681 trained.is_some()
682 );
683
684 match handle.block_on(builder.build(
685 state.directory.as_ref(),
686 segment_id,
687 trained.as_deref(),
688 )) {
689 Ok(meta) if meta.num_docs > 0 => {
690 let duration_ms = build_start.elapsed().as_millis() as u64;
691 log::info!(
692 "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
693 segment_hex,
694 meta.num_docs,
695 duration_ms,
696 );
697 state.built_segments.lock().push(PreparedSegment {
698 id: segment_hex,
699 num_docs: meta.num_docs,
700 _operation: operation,
701 });
702 }
703 Ok(_) => {}
704 Err(e) => {
705 log::error!(
706 "[segment_build_failed] segment_id={} error={:?}",
707 segment_hex,
708 e
709 );
710 if let Err(delete_error) = handle.block_on(crate::segment::delete_segment(
711 state.directory.as_ref(),
712 segment_id,
713 )) {
714 log::warn!(
715 "[segment_cleanup] failed deleting partial indexing segment {}: {}",
716 segment_hex,
717 delete_error,
718 );
719 }
720 }
721 }
722 }
723
724 pub async fn maybe_merge(&self) {
730 self.segment_manager.maybe_merge().await;
731 }
732
733 pub async fn abort_merges(&self) {
735 self.segment_manager.abort_merges().await;
736 }
737
738 pub async fn wait_for_merging_thread(&self) {
740 self.segment_manager.wait_for_merging_thread().await;
741 }
742
743 pub async fn wait_for_all_merges(&self) {
745 self.segment_manager.wait_for_all_merges().await;
746 }
747
748 pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
750 self.segment_manager.tracker()
751 }
752
753 pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
755 self.segment_manager.acquire_snapshot().await
756 }
757
758 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
760 self.segment_manager.cleanup_orphan_segments().await
761 }
762
763 pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
774 self.doc_sender.close();
776
777 self.worker_state.resume_cvar.notify_all();
781
782 let state = Arc::clone(&self.worker_state);
785 let all_flushed = tokio::task::spawn_blocking(move || {
786 let mut lock = state.flush_mutex.lock();
787 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
788 while state.flush_count.load(Ordering::Acquire) < state.num_workers {
789 let remaining = deadline.saturating_duration_since(std::time::Instant::now());
790 if remaining.is_zero() {
791 log::error!(
792 "[prepare_commit] timed out waiting for workers: {}/{} flushed",
793 state.flush_count.load(Ordering::Acquire),
794 state.num_workers
795 );
796 return false;
797 }
798 state.flush_cvar.wait_for(&mut lock, remaining);
799 }
800 true
801 })
802 .await
803 .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
804
805 if !all_flushed {
806 self.resume_workers();
808 return Err(Error::Internal(format!(
809 "prepare_commit timed out: {}/{} workers flushed",
810 self.worker_state.flush_count.load(Ordering::Acquire),
811 self.worker_state.num_workers
812 )));
813 }
814
815 let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
817 self.flushed_segments.extend(built);
818
819 Ok(PreparedCommit {
820 writer: self,
821 is_resolved: false,
822 is_published: false,
823 })
824 }
825
826 pub async fn commit(&mut self) -> Result<bool> {
831 self.prepare_commit().await?.commit().await
832 }
833
834 pub async fn force_merge(&mut self) -> Result<()> {
836 self.prepare_commit().await?.commit().await?;
837 self.segment_manager.force_merge().await
838 }
839
840 pub async fn reorder(&mut self) -> Result<()> {
845 self.prepare_commit().await?.commit().await?;
846 self.segment_manager.reorder_segments().await
847 }
848
849 pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
851 &self.segment_manager
852 }
853
854 fn resume_workers(&mut self) {
859 if tokio::runtime::Handle::try_current().is_err() {
860 self.worker_state.shutdown.store(true, Ordering::Release);
863 self.worker_state.resume_cvar.notify_all();
864 return;
865 }
866
867 self.worker_state.flush_count.store(0, Ordering::Release);
869
870 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
872 self.doc_sender = sender;
873
874 {
876 let mut lock = self.worker_state.resume_receiver.lock();
877 *lock = Some(receiver);
878 }
879 self.worker_state
880 .resume_epoch
881 .fetch_add(1, Ordering::Release);
882 self.worker_state.resume_cvar.notify_all();
883 }
884
885 }
887
888impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
889 fn drop(&mut self) {
890 self.worker_state.shutdown.store(true, Ordering::Release);
892 self.doc_sender.close();
894 self.worker_state.resume_cvar.notify_all();
896 for w in std::mem::take(&mut self.workers) {
898 let _ = w.join();
899 }
900 }
901}
902
903pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
910 writer: &'a mut IndexWriter<D>,
911 is_resolved: bool,
912 is_published: bool,
915}
916
917impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
918 pub async fn commit(mut self) -> Result<bool> {
922 let segments = std::mem::take(&mut self.writer.flushed_segments);
923
924 if segments.is_empty() {
926 log::debug!("[commit] no segments to commit, skipping");
927 self.is_resolved = true;
928 self.writer.resume_workers();
929 return Ok(false);
930 }
931
932 let metadata_entries: Vec<(String, u32)> = segments
933 .iter()
934 .map(PreparedSegment::metadata_entry)
935 .collect();
936 if let Err(error) = self.writer.segment_manager.commit(&metadata_entries).await {
937 self.writer.flushed_segments = segments;
941 self.is_resolved = true;
942 self.writer.resume_workers();
943 return Err(error);
944 }
945 self.is_published = true;
946
947 drop(segments);
950
951 let post_publish_result = async {
953 if let Some(ref mut pk_index) = self.writer.primary_key_index {
954 let snapshot = self.writer.segment_manager.acquire_snapshot().await;
955 let existing_ids: std::collections::HashSet<&str> =
956 pk_index.committed_segment_ids().collect();
957
958 let load_futures: Vec<_> = snapshot
960 .segment_ids()
961 .iter()
962 .filter(|id| !existing_ids.contains(id.as_str()))
963 .map(|seg_id_str| {
964 let seg_id_str = seg_id_str.clone();
965 let dir = self.writer.directory.as_ref();
966 let schema = Arc::clone(&self.writer.schema);
967 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
968 })
969 .collect();
970 let new_data = futures::future::try_join_all(load_futures).await?;
971
972 let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
973 pk_index.refresh_incremental(new_data, snapshot);
974
975 let bloom_bytes = pk_index.bloom_to_bytes();
977 let data = super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes);
978 if let Err(e) = self
979 .writer
980 .directory
981 .write(
982 std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
983 &data,
984 )
985 .await
986 {
987 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
988 }
989 }
990
991 self.writer.segment_manager.maybe_merge().await;
992 Ok(())
993 }
994 .await;
995
996 self.is_resolved = true;
999 self.writer.resume_workers();
1000 post_publish_result.map(|()| true)
1001 }
1002
1003 pub fn abort(mut self) {
1006 self.is_resolved = true;
1007 self.writer.flushed_segments.clear();
1008 if let Some(ref mut pk_index) = self.writer.primary_key_index {
1009 pk_index.clear_uncommitted();
1010 }
1011 self.writer.resume_workers();
1012 }
1013}
1014
1015impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
1016 fn drop(&mut self) {
1017 if !self.is_resolved {
1018 if self.is_published {
1019 log::warn!("PreparedCommit dropped after metadata publication — resuming workers");
1020 } else {
1021 log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
1022 self.writer.flushed_segments.clear();
1023 if let Some(ref mut pk_index) = self.writer.primary_key_index {
1024 pk_index.clear_uncommitted();
1025 }
1026 }
1027 self.writer.resume_workers();
1028 }
1029 }
1030}
1031
1032async fn load_pk_segment_data<D: crate::directories::Directory>(
1034 dir: &D,
1035 seg_id_str: &str,
1036 schema: &Arc<crate::dsl::Schema>,
1037) -> Result<super::primary_key::PkSegmentData> {
1038 let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
1039 .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
1040 let files = crate::segment::SegmentFiles::new(seg_id.0);
1041 let fast_fields =
1042 crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
1043 Ok(super::primary_key::PkSegmentData {
1044 segment_id: seg_id_str.to_string(),
1045 fast_fields,
1046 })
1047}