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<(String, u32)>,
74 primary_key_index: Option<super::primary_key::PrimaryKeyIndex>,
76}
77
78struct WorkerState<D: DirectoryWriter + 'static> {
80 directory: Arc<D>,
81 schema: Arc<Schema>,
82 builder_config: SegmentBuilderConfig,
83 tokenizers: parking_lot::RwLock<FxHashMap<Field, BoxedTokenizer>>,
84 memory_budget_per_worker: usize,
86 segment_manager: Arc<crate::merge::SegmentManager<D>>,
88 built_segments: parking_lot::Mutex<Vec<(String, u32)>>,
90
91 flush_count: AtomicUsize,
98 flush_mutex: parking_lot::Mutex<()>,
100 flush_cvar: parking_lot::Condvar,
101 resume_receiver: parking_lot::Mutex<Option<async_channel::Receiver<Document>>>,
103 resume_epoch: AtomicUsize,
106 resume_cvar: parking_lot::Condvar,
108 shutdown: AtomicBool,
110 num_workers: usize,
112}
113
114impl<D: DirectoryWriter + 'static> IndexWriter<D> {
115 pub async fn create(directory: D, schema: Schema, config: IndexConfig) -> Result<Self> {
117 Self::create_with_config(directory, schema, config, SegmentBuilderConfig::default()).await
118 }
119
120 pub async fn create_with_config(
122 directory: D,
123 schema: Schema,
124 config: IndexConfig,
125 builder_config: SegmentBuilderConfig,
126 ) -> Result<Self> {
127 let directory = Arc::new(directory);
128 let schema = Arc::new(schema);
129 directory.set_index_label(schema.index_label());
131 let metadata = super::IndexMetadata::new((*schema).clone());
132
133 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
134 Arc::clone(&directory),
135 Arc::clone(&schema),
136 metadata,
137 config.merge_policy.clone_box(),
138 config.term_cache_blocks,
139 config.max_concurrent_merges,
140 config.merge_bp_time_budget,
141 config.bp_memory_budget_bytes,
142 ));
143 segment_manager.update_metadata(|_| {}).await?;
144
145 Ok(Self::new_with_parts(
146 directory,
147 schema,
148 config,
149 builder_config,
150 segment_manager,
151 ))
152 }
153
154 pub async fn open(directory: D, config: IndexConfig) -> Result<Self> {
156 Self::open_with_config(directory, config, SegmentBuilderConfig::default()).await
157 }
158
159 pub async fn open_with_config(
161 directory: D,
162 config: IndexConfig,
163 builder_config: SegmentBuilderConfig,
164 ) -> Result<Self> {
165 let directory = Arc::new(directory);
166 let metadata = super::IndexMetadata::load(directory.as_ref()).await?;
167 let schema = Arc::new(metadata.schema.clone());
168 directory.set_index_label(schema.index_label());
170
171 let segment_manager = Arc::new(crate::merge::SegmentManager::new(
172 Arc::clone(&directory),
173 Arc::clone(&schema),
174 metadata,
175 config.merge_policy.clone_box(),
176 config.term_cache_blocks,
177 config.max_concurrent_merges,
178 config.merge_bp_time_budget,
179 config.bp_memory_budget_bytes,
180 ));
181 segment_manager.load_and_publish_trained().await;
182
183 Ok(Self::new_with_parts(
184 directory,
185 schema,
186 config,
187 builder_config,
188 segment_manager,
189 ))
190 }
191
192 pub fn from_index(index: &super::Index<D>) -> Self {
195 Self::new_with_parts(
196 Arc::clone(&index.directory),
197 Arc::clone(&index.schema),
198 index.config.clone(),
199 SegmentBuilderConfig::default(),
200 Arc::clone(&index.segment_manager),
201 )
202 }
203
204 fn new_with_parts(
210 directory: Arc<D>,
211 schema: Arc<Schema>,
212 config: IndexConfig,
213 builder_config: SegmentBuilderConfig,
214 segment_manager: Arc<crate::merge::SegmentManager<D>>,
215 ) -> Self {
216 let registry = crate::tokenizer::TokenizerRegistry::new();
218 let mut tokenizers = FxHashMap::default();
219 for (field, entry) in schema.fields() {
220 if matches!(entry.field_type, crate::dsl::FieldType::Text)
221 && let Some(ref tok_name) = entry.tokenizer
222 && let Some(tok) = registry.get(tok_name)
223 {
224 tokenizers.insert(field, tok);
225 }
226 }
227
228 let num_workers = config.num_indexing_threads.max(1);
229 let worker_state = Arc::new(WorkerState {
230 directory: Arc::clone(&directory),
231 schema: Arc::clone(&schema),
232 builder_config,
233 tokenizers: parking_lot::RwLock::new(tokenizers),
234 memory_budget_per_worker: config.max_indexing_memory_bytes / num_workers,
235 segment_manager: Arc::clone(&segment_manager),
236 built_segments: parking_lot::Mutex::new(Vec::new()),
237 flush_count: AtomicUsize::new(0),
238 flush_mutex: parking_lot::Mutex::new(()),
239 flush_cvar: parking_lot::Condvar::new(),
240 resume_receiver: parking_lot::Mutex::new(None),
241 resume_epoch: AtomicUsize::new(0),
242 resume_cvar: parking_lot::Condvar::new(),
243 shutdown: AtomicBool::new(false),
244 num_workers,
245 });
246 let (doc_sender, workers) = Self::spawn_workers(&worker_state, num_workers);
247
248 Self {
249 directory,
250 schema,
251 config,
252 doc_sender,
253 workers,
254 worker_state,
255 segment_manager,
256 flushed_segments: Vec::new(),
257 primary_key_index: None,
258 }
259 }
260
261 fn spawn_workers(
262 worker_state: &Arc<WorkerState<D>>,
263 num_workers: usize,
264 ) -> (
265 async_channel::Sender<Document>,
266 Vec<std::thread::JoinHandle<()>>,
267 ) {
268 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
269 let handle = tokio::runtime::Handle::current();
270 let mut workers = Vec::with_capacity(num_workers);
271 for i in 0..num_workers {
272 let state = Arc::clone(worker_state);
273 let rx = receiver.clone();
274 let rt = handle.clone();
275 workers.push(
276 std::thread::Builder::new()
277 .name(format!("index-worker-{}", i))
278 .spawn(move || Self::worker_loop(state, rx, rt))
279 .expect("failed to spawn index worker thread"),
280 );
281 }
282 (sender, workers)
283 }
284
285 pub fn schema(&self) -> &Schema {
287 &self.schema
288 }
289
290 pub fn set_tokenizer<T: crate::tokenizer::Tokenizer>(&mut self, field: Field, tokenizer: T) {
293 self.worker_state
294 .tokenizers
295 .write()
296 .insert(field, Box::new(tokenizer));
297 }
298
299 pub async fn init_primary_key_dedup(&mut self) -> Result<()> {
315 use super::primary_key::{PK_BLOOM_FILE, deserialize_pk_bloom};
316
317 let field = match self.schema.primary_field() {
318 Some(f) => f,
319 None => return Ok(()),
320 };
321
322 let snapshot = self.segment_manager.acquire_snapshot().await;
323 let current_seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
324
325 let cached = match self
327 .directory
328 .open_read(std::path::Path::new(PK_BLOOM_FILE))
329 .await
330 {
331 Ok(handle) => {
332 let data = handle.read_bytes_range(0..handle.len()).await;
333 match data {
334 Ok(bytes) => deserialize_pk_bloom(bytes.as_slice()),
335 Err(_) => None,
336 }
337 }
338 Err(_) => None,
339 };
340
341 let load_futures: Vec<_> = current_seg_ids
343 .iter()
344 .map(|seg_id_str| {
345 let seg_id_str = seg_id_str.clone();
346 let dir = self.directory.as_ref();
347 let schema = Arc::clone(&self.schema);
348 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
349 })
350 .collect();
351 let all_data = futures::future::try_join_all(load_futures).await?;
352
353 if let Some((persisted_seg_ids, bloom)) = cached {
354 let mut pk_data = Vec::with_capacity(all_data.len());
356 let mut new_data = Vec::new();
357 for d in all_data {
358 if persisted_seg_ids.contains(&d.segment_id) {
359 pk_data.push(d);
360 } else {
361 new_data.push(d);
362 }
363 }
364 let needs_persist = !new_data.is_empty();
365 let new_start = pk_data.len();
366 pk_data.extend(new_data);
367
368 let pk_index = if new_start == pk_data.len() {
369 super::primary_key::PrimaryKeyIndex::from_persisted(
371 field,
372 bloom,
373 pk_data,
374 &[],
375 snapshot,
376 )
377 } else {
378 tokio::task::spawn_blocking(move || {
380 let mut bloom = bloom;
383 let mut added = 0usize;
384 let num_new = pk_data.len() - new_start;
385 for data in &pk_data[new_start..] {
386 if let Some(ff) = data.fast_fields.get(&field.0)
387 && let Some(dict) = ff.text_dict()
388 {
389 for key in dict.iter() {
390 bloom.insert(key.as_bytes());
391 added += 1;
392 }
393 }
394 }
395 if added > 0 {
396 log::info!(
397 "[primary_key] bloom: added {} keys from {} new segment(s)",
398 added,
399 num_new,
400 );
401 }
402 super::primary_key::PrimaryKeyIndex::from_persisted(
403 field,
404 bloom,
405 pk_data,
406 &[],
407 snapshot,
408 )
409 })
410 .await
411 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?
412 };
413
414 if needs_persist {
415 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
416 }
417
418 self.primary_key_index = Some(pk_index);
419 } else {
420 let pk_index = tokio::task::spawn_blocking(move || {
422 super::primary_key::PrimaryKeyIndex::new(field, all_data, snapshot)
423 })
424 .await
425 .map_err(|e| Error::Internal(format!("spawn_blocking failed: {}", e)))?;
426
427 self.persist_pk_bloom(&pk_index, ¤t_seg_ids).await;
428 self.primary_key_index = Some(pk_index);
429 }
430
431 Ok(())
432 }
433
434 async fn persist_pk_bloom(
437 &self,
438 pk_index: &super::primary_key::PrimaryKeyIndex,
439 segment_ids: &[String],
440 ) {
441 use super::primary_key::{PK_BLOOM_FILE, serialize_pk_bloom};
442
443 let bloom_bytes = pk_index.bloom_to_bytes();
444 let data = serialize_pk_bloom(segment_ids, &bloom_bytes);
445 if let Err(e) = self
446 .directory
447 .write(std::path::Path::new(PK_BLOOM_FILE), &data)
448 .await
449 {
450 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
451 }
452 }
453
454 pub fn add_document(&self, doc: Document) -> Result<()> {
459 if let Some(ref pk_index) = self.primary_key_index {
460 pk_index.check_and_insert(&doc)?;
461 }
462 match self.doc_sender.try_send(doc) {
463 Ok(()) => Ok(()),
464 Err(async_channel::TrySendError::Full(doc)) => {
465 if let Some(ref pk_index) = self.primary_key_index {
467 pk_index.rollback_uncommitted_key(&doc);
468 }
469 Err(Error::QueueFull)
470 }
471 Err(async_channel::TrySendError::Closed(doc)) => {
472 if let Some(ref pk_index) = self.primary_key_index {
474 pk_index.rollback_uncommitted_key(&doc);
475 }
476 Err(Error::Internal("Document channel closed".into()))
477 }
478 }
479 }
480
481 pub fn add_documents(&self, documents: Vec<Document>) -> Result<usize> {
486 let total = documents.len();
487 for (i, doc) in documents.into_iter().enumerate() {
488 match self.add_document(doc) {
489 Ok(()) => {}
490 Err(Error::QueueFull) => return Ok(i),
491 Err(e) => return Err(e),
492 }
493 }
494 Ok(total)
495 }
496
497 fn worker_loop(
510 state: Arc<WorkerState<D>>,
511 initial_receiver: async_channel::Receiver<Document>,
512 handle: tokio::runtime::Handle,
513 ) {
514 let mut receiver = initial_receiver;
515 let mut my_epoch = 0usize;
516
517 loop {
518 let build_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
522 let mut builder: Option<SegmentBuilder> = None;
523
524 while let Ok(doc) = receiver.recv_blocking() {
525 if builder.is_none() {
527 match SegmentBuilder::new(
528 Arc::clone(&state.schema),
529 state.builder_config.clone(),
530 ) {
531 Ok(mut b) => {
532 for (field, tokenizer) in state.tokenizers.read().iter() {
533 b.set_tokenizer(*field, tokenizer.clone_box());
534 }
535 builder = Some(b);
536 }
537 Err(e) => {
538 log::error!("Failed to create segment builder: {:?}", e);
539 continue;
540 }
541 }
542 }
543
544 let b = builder.as_mut().unwrap();
545 if let Err(e) = b.add_document(doc) {
546 log::error!("Failed to index document: {:?}", e);
547 continue;
548 }
549
550 let builder_memory = b.estimated_memory_bytes();
551
552 if b.num_docs() & 0x3FFF == 0 {
553 log::debug!(
554 "[indexing] docs={}, memory={:.2} MB, budget={:.2} MB",
555 b.num_docs(),
556 builder_memory as f64 / (1024.0 * 1024.0),
557 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0)
558 );
559 }
560
561 const MIN_DOCS_BEFORE_FLUSH: u32 = 100;
563
564 let effective_budget = state.memory_budget_per_worker * 4 / 5;
568
569 if builder_memory >= effective_budget && b.num_docs() >= MIN_DOCS_BEFORE_FLUSH {
570 log::info!(
571 "[indexing] memory budget reached, building segment: \
572 docs={}, memory={:.2} MB, budget={:.2} MB",
573 b.num_docs(),
574 builder_memory as f64 / (1024.0 * 1024.0),
575 state.memory_budget_per_worker as f64 / (1024.0 * 1024.0),
576 );
577 let full_builder = builder.take().unwrap();
578 Self::build_segment_inline(&state, full_builder, &handle);
579 }
580 }
581
582 if let Some(b) = builder.take()
584 && b.num_docs() > 0
585 {
586 Self::build_segment_inline(&state, b, &handle);
587 }
588 }));
589
590 if build_result.is_err() {
591 log::error!(
592 "[worker] panic during indexing cycle — documents in this cycle may be lost"
593 );
594 }
595
596 let prev = state.flush_count.fetch_add(1, Ordering::Release);
599 if prev + 1 == state.num_workers {
600 let _lock = state.flush_mutex.lock();
602 state.flush_cvar.notify_one();
603 }
604
605 {
609 let mut lock = state.resume_receiver.lock();
610 loop {
611 if state.shutdown.load(Ordering::Acquire) {
612 return;
613 }
614 let current_epoch = state.resume_epoch.load(Ordering::Acquire);
615 if current_epoch > my_epoch
616 && let Some(rx) = lock.as_ref()
617 {
618 receiver = rx.clone();
619 my_epoch = current_epoch;
620 break;
621 }
622 state.resume_cvar.wait(&mut lock);
623 }
624 }
625 }
626 }
627
628 fn build_segment_inline(
632 state: &WorkerState<D>,
633 builder: SegmentBuilder,
634 handle: &tokio::runtime::Handle,
635 ) {
636 let segment_id = SegmentId::new();
637 let segment_hex = segment_id.to_hex();
638 let trained = state.segment_manager.trained();
639 let doc_count = builder.num_docs();
640 let build_start = std::time::Instant::now();
641
642 log::info!(
643 "[segment_build] segment_id={} doc_count={} ann={}",
644 segment_hex,
645 doc_count,
646 trained.is_some()
647 );
648
649 match handle.block_on(builder.build(
650 state.directory.as_ref(),
651 segment_id,
652 trained.as_deref(),
653 )) {
654 Ok(meta) if meta.num_docs > 0 => {
655 let duration_ms = build_start.elapsed().as_millis() as u64;
656 log::info!(
657 "[segment_build_done] segment_id={} doc_count={} duration_ms={}",
658 segment_hex,
659 meta.num_docs,
660 duration_ms,
661 );
662 state
663 .built_segments
664 .lock()
665 .push((segment_hex, meta.num_docs));
666 }
667 Ok(_) => {}
668 Err(e) => {
669 log::error!(
670 "[segment_build_failed] segment_id={} error={:?}",
671 segment_hex,
672 e
673 );
674 }
675 }
676 }
677
678 pub async fn maybe_merge(&self) {
684 self.segment_manager.maybe_merge().await;
685 }
686
687 pub async fn abort_merges(&self) {
689 self.segment_manager.abort_merges().await;
690 }
691
692 pub async fn wait_for_merging_thread(&self) {
694 self.segment_manager.wait_for_merging_thread().await;
695 }
696
697 pub async fn wait_for_all_merges(&self) {
699 self.segment_manager.wait_for_all_merges().await;
700 }
701
702 pub fn tracker(&self) -> std::sync::Arc<crate::segment::SegmentTracker> {
704 self.segment_manager.tracker()
705 }
706
707 pub async fn acquire_snapshot(&self) -> crate::segment::SegmentSnapshot {
709 self.segment_manager.acquire_snapshot().await
710 }
711
712 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
714 self.segment_manager.cleanup_orphan_segments().await
715 }
716
717 pub async fn prepare_commit(&mut self) -> Result<PreparedCommit<'_, D>> {
728 self.doc_sender.close();
730
731 self.worker_state.resume_cvar.notify_all();
735
736 let state = Arc::clone(&self.worker_state);
739 let all_flushed = tokio::task::spawn_blocking(move || {
740 let mut lock = state.flush_mutex.lock();
741 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(300);
742 while state.flush_count.load(Ordering::Acquire) < state.num_workers {
743 let remaining = deadline.saturating_duration_since(std::time::Instant::now());
744 if remaining.is_zero() {
745 log::error!(
746 "[prepare_commit] timed out waiting for workers: {}/{} flushed",
747 state.flush_count.load(Ordering::Acquire),
748 state.num_workers
749 );
750 return false;
751 }
752 state.flush_cvar.wait_for(&mut lock, remaining);
753 }
754 true
755 })
756 .await
757 .map_err(|e| Error::Internal(format!("Failed to wait for workers: {}", e)))?;
758
759 if !all_flushed {
760 self.resume_workers();
762 return Err(Error::Internal(format!(
763 "prepare_commit timed out: {}/{} workers flushed",
764 self.worker_state.flush_count.load(Ordering::Acquire),
765 self.worker_state.num_workers
766 )));
767 }
768
769 let built = std::mem::take(&mut *self.worker_state.built_segments.lock());
771 self.flushed_segments.extend(built);
772
773 Ok(PreparedCommit {
774 writer: self,
775 is_resolved: false,
776 })
777 }
778
779 pub async fn commit(&mut self) -> Result<bool> {
784 self.prepare_commit().await?.commit().await
785 }
786
787 pub async fn force_merge(&mut self) -> Result<()> {
789 self.prepare_commit().await?.commit().await?;
790 self.segment_manager.force_merge().await
791 }
792
793 pub async fn reorder(&mut self) -> Result<()> {
798 self.prepare_commit().await?.commit().await?;
799 self.segment_manager.reorder_segments().await
800 }
801
802 pub fn segment_manager(&self) -> &Arc<crate::merge::SegmentManager<D>> {
804 &self.segment_manager
805 }
806
807 fn resume_workers(&mut self) {
812 if tokio::runtime::Handle::try_current().is_err() {
813 self.worker_state.shutdown.store(true, Ordering::Release);
816 self.worker_state.resume_cvar.notify_all();
817 return;
818 }
819
820 self.worker_state.flush_count.store(0, Ordering::Release);
822
823 let (sender, receiver) = async_channel::bounded(PIPELINE_MAX_SIZE_IN_DOCS);
825 self.doc_sender = sender;
826
827 {
829 let mut lock = self.worker_state.resume_receiver.lock();
830 *lock = Some(receiver);
831 }
832 self.worker_state
833 .resume_epoch
834 .fetch_add(1, Ordering::Release);
835 self.worker_state.resume_cvar.notify_all();
836 }
837
838 }
840
841impl<D: DirectoryWriter + 'static> Drop for IndexWriter<D> {
842 fn drop(&mut self) {
843 self.worker_state.shutdown.store(true, Ordering::Release);
845 self.doc_sender.close();
847 self.worker_state.resume_cvar.notify_all();
849 for w in std::mem::take(&mut self.workers) {
851 let _ = w.join();
852 }
853 }
854}
855
856pub struct PreparedCommit<'a, D: DirectoryWriter + 'static> {
863 writer: &'a mut IndexWriter<D>,
864 is_resolved: bool,
865}
866
867impl<'a, D: DirectoryWriter + 'static> PreparedCommit<'a, D> {
868 pub async fn commit(mut self) -> Result<bool> {
872 self.is_resolved = true;
873 let segments = std::mem::take(&mut self.writer.flushed_segments);
874
875 if segments.is_empty() {
877 log::debug!("[commit] no segments to commit, skipping");
878 self.writer.resume_workers();
879 return Ok(false);
880 }
881
882 self.writer.segment_manager.commit(segments).await?;
883
884 if let Some(ref mut pk_index) = self.writer.primary_key_index {
886 let snapshot = self.writer.segment_manager.acquire_snapshot().await;
887 let existing_ids: std::collections::HashSet<&str> =
888 pk_index.committed_segment_ids().collect();
889
890 let load_futures: Vec<_> = snapshot
892 .segment_ids()
893 .iter()
894 .filter(|id| !existing_ids.contains(id.as_str()))
895 .map(|seg_id_str| {
896 let seg_id_str = seg_id_str.clone();
897 let dir = self.writer.directory.as_ref();
898 let schema = Arc::clone(&self.writer.schema);
899 async move { load_pk_segment_data(dir, &seg_id_str, &schema).await }
900 })
901 .collect();
902 let new_data = futures::future::try_join_all(load_futures).await?;
903
904 let seg_ids: Vec<String> = snapshot.segment_ids().to_vec();
905 pk_index.refresh_incremental(new_data, snapshot);
906
907 let bloom_bytes = pk_index.bloom_to_bytes();
909 let data = super::primary_key::serialize_pk_bloom(&seg_ids, &bloom_bytes);
910 if let Err(e) = self
911 .writer
912 .directory
913 .write(
914 std::path::Path::new(super::primary_key::PK_BLOOM_FILE),
915 &data,
916 )
917 .await
918 {
919 log::warn!("[primary_key] failed to persist bloom cache: {}", e);
920 }
921 }
922
923 self.writer.segment_manager.maybe_merge().await;
924 self.writer.resume_workers();
925 Ok(true)
926 }
927
928 pub fn abort(mut self) {
931 self.is_resolved = true;
932 self.writer.flushed_segments.clear();
933 if let Some(ref mut pk_index) = self.writer.primary_key_index {
934 pk_index.clear_uncommitted();
935 }
936 self.writer.resume_workers();
937 }
938}
939
940impl<D: DirectoryWriter + 'static> Drop for PreparedCommit<'_, D> {
941 fn drop(&mut self) {
942 if !self.is_resolved {
943 log::warn!("PreparedCommit dropped without commit/abort — auto-aborting");
944 self.writer.flushed_segments.clear();
945 if let Some(ref mut pk_index) = self.writer.primary_key_index {
946 pk_index.clear_uncommitted();
947 }
948 self.writer.resume_workers();
949 }
950 }
951}
952
953async fn load_pk_segment_data<D: crate::directories::Directory>(
955 dir: &D,
956 seg_id_str: &str,
957 schema: &Arc<crate::dsl::Schema>,
958) -> Result<super::primary_key::PkSegmentData> {
959 let seg_id = crate::segment::SegmentId::from_hex(seg_id_str)
960 .ok_or_else(|| Error::Internal(format!("Invalid segment id: {}", seg_id_str)))?;
961 let files = crate::segment::SegmentFiles::new(seg_id.0);
962 let fast_fields =
963 crate::segment::reader::loader::load_fast_fields_file(dir, &files, schema).await?;
964 Ok(super::primary_key::PkSegmentData {
965 segment_id: seg_id_str.to_string(),
966 fast_fields,
967 })
968}