1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwapOption;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, SegmentMetaInfo};
56use crate::segment::{
57 SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot, SegmentTracker, TrainedVectorStructures,
58};
59#[cfg(feature = "native")]
60use crate::segment::{SegmentMerger, SegmentReader};
61
62use super::{MergePolicy, SegmentInfo};
63
64struct ActiveOperationState {
74 segment_ids: HashSet<String>,
75 operation_tokens: HashSet<u64>,
76 indexing_tokens: HashSet<u64>,
82 next_operation_token: u64,
83 accepting: bool,
84 non_indexing_paused: bool,
88}
89
90struct ActiveSegmentOperations {
91 inner: parking_lot::Mutex<ActiveOperationState>,
92 idle: Notify,
93 shutdown: Notify,
94}
95
96impl ActiveSegmentOperations {
97 fn new() -> Self {
98 Self {
99 inner: parking_lot::Mutex::new(ActiveOperationState {
100 segment_ids: HashSet::new(),
101 operation_tokens: HashSet::new(),
102 indexing_tokens: HashSet::new(),
103 next_operation_token: 0,
104 accepting: true,
105 non_indexing_paused: false,
106 }),
107 idle: Notify::new(),
108 shutdown: Notify::new(),
109 }
110 }
111
112 fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
116 self.try_register_kind(segment_ids, false, false)
117 }
118
119 fn try_register_indexing(
122 self: &Arc<Self>,
123 segment_ids: Vec<String>,
124 ) -> Option<SegmentOperationGuard> {
125 self.try_register_kind(segment_ids, true, false)
126 }
127
128 fn try_register_vector_update(
131 self: &Arc<Self>,
132 segment_ids: Vec<String>,
133 ) -> Option<SegmentOperationGuard> {
134 self.try_register_kind(segment_ids, false, true)
135 }
136
137 fn try_register_kind(
138 self: &Arc<Self>,
139 segment_ids: Vec<String>,
140 indexing: bool,
141 vector_update: bool,
142 ) -> Option<SegmentOperationGuard> {
143 let mut inner = self.inner.lock();
144 if !inner.accepting {
145 log::debug!("[segment_lifecycle] rejected operation during shutdown");
146 return None;
147 }
148 if !indexing && !vector_update && inner.non_indexing_paused {
149 log::debug!("[segment_lifecycle] deferred operation during dense vector retraining");
150 return None;
151 }
152 for id in &segment_ids {
154 if inner.segment_ids.contains(id) {
155 log::debug!(
156 "[segment_lifecycle] rejected: {} overlaps with an active operation ({} active IDs)",
157 id,
158 inner.segment_ids.len()
159 );
160 return None;
161 }
162 }
163 log::debug!(
164 "[segment_lifecycle] registered {} IDs (total active: {})",
165 segment_ids.len(),
166 inner.segment_ids.len() + segment_ids.len()
167 );
168 let operation_token = inner.next_operation_token;
169 let next_operation_token = operation_token.checked_add(1)?;
170 for id in &segment_ids {
171 inner.segment_ids.insert(id.clone());
172 }
173 inner.next_operation_token = next_operation_token;
174 inner.operation_tokens.insert(operation_token);
175 if indexing {
176 inner.indexing_tokens.insert(operation_token);
177 }
178 Some(SegmentOperationGuard {
179 active_operations: Arc::clone(self),
180 segment_ids,
181 operation_token,
182 })
183 }
184
185 fn snapshot(&self) -> HashSet<String> {
187 self.inner.lock().segment_ids.clone()
188 }
189
190 fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
201 let inner = self.inner.lock();
202 let tokens = inner
203 .operation_tokens
204 .difference(&inner.indexing_tokens)
205 .copied()
206 .collect();
207 (tokens, inner.indexing_tokens.len())
208 }
209
210 fn stop_accepting(&self) {
213 let mut inner = self.inner.lock();
214 inner.accepting = false;
215 self.shutdown.notify_waiters();
216 if inner.segment_ids.is_empty() {
217 self.idle.notify_waiters();
218 }
219 }
220
221 fn pause_non_indexing(&self) {
222 self.inner.lock().non_indexing_paused = true;
223 }
224
225 fn resume_non_indexing(&self) {
226 self.inner.lock().non_indexing_paused = false;
227 self.idle.notify_waiters();
228 }
229
230 fn is_accepting(&self) -> bool {
231 self.inner.lock().accepting
232 }
233
234 async fn wait_until_idle(&self) {
238 loop {
239 let notified = self.idle.notified();
240 if self.inner.lock().segment_ids.is_empty() {
241 return;
242 }
243 notified.await;
244 }
245 }
246
247 async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
248 while !operations.is_empty() {
249 let notified = self.idle.notified();
250 if self.inner.lock().operation_tokens.is_disjoint(operations) {
251 return;
252 }
253 notified.await;
254 }
255 }
256
257 async fn wait_for_shutdown(&self) {
260 loop {
261 let notified = self.shutdown.notified();
262 if !self.inner.lock().accepting {
263 return;
264 }
265 notified.await;
266 }
267 }
268}
269
270pub(crate) struct SegmentOperationGuard {
274 active_operations: Arc<ActiveSegmentOperations>,
275 segment_ids: Vec<String>,
276 operation_token: u64,
277}
278
279impl Drop for SegmentOperationGuard {
280 fn drop(&mut self) {
281 let mut inner = self.active_operations.inner.lock();
282 for id in &self.segment_ids {
283 inner.segment_ids.remove(id);
284 }
285 inner.operation_tokens.remove(&self.operation_token);
286 inner.indexing_tokens.remove(&self.operation_token);
287 self.active_operations.idle.notify_waiters();
290 if inner.segment_ids.is_empty() {
291 debug_assert!(inner.operation_tokens.is_empty());
292 }
293 }
294}
295
296struct VectorArtifactUpdateLease {
303 updating: Arc<AtomicBool>,
304 active_operations: Arc<ActiveSegmentOperations>,
305}
306
307impl Drop for VectorArtifactUpdateLease {
308 fn drop(&mut self) {
309 self.updating.store(false, Ordering::Release);
310 self.active_operations.resume_non_indexing();
311 }
312}
313
314#[derive(Clone)]
315pub(crate) struct VectorArtifactUpdateGuard {
316 _lease: Arc<VectorArtifactUpdateLease>,
317}
318
319static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
323
324const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
325const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
326
327#[derive(Default)]
328struct MergeRetryState {
329 retry_after: Option<std::time::Instant>,
330 consecutive_failures: u32,
331}
332
333fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
334 let shift = consecutive_failures.saturating_sub(1).min(16);
335 MERGE_RETRY_BASE_DELAY
336 .checked_mul(1u32 << shift)
337 .unwrap_or(MERGE_RETRY_MAX_DELAY)
338 .min(MERGE_RETRY_MAX_DELAY)
339}
340
341struct DrainedMergeHandles<'a> {
350 shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
351 drained: Vec<JoinHandle<()>>,
352}
353
354impl<'a> DrainedMergeHandles<'a> {
355 fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
356 let drained = std::mem::take(&mut *shared.lock());
357 Self { shared, drained }
358 }
359
360 fn is_empty(&self) -> bool {
361 self.drained.is_empty()
362 }
363
364 async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
368 let handle = self.drained.last_mut()?;
369 let result = handle.await;
370 self.drained.pop();
371 Some(result)
372 }
373}
374
375impl Drop for DrainedMergeHandles<'_> {
376 fn drop(&mut self) {
377 if !self.drained.is_empty() {
378 self.shared.lock().append(&mut self.drained);
379 }
380 }
381}
382
383fn try_spawn_lifecycle<F>(
390 handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
391 runtime: &tokio::runtime::Handle,
392 future: F,
393) -> bool
394where
395 F: std::future::Future<Output = ()> + Send + 'static,
396{
397 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
398 let mut handles = handles.lock();
399 handles.retain(|handle| !handle.is_finished());
400 handles.push(runtime.spawn(future));
401 }))
402 .is_ok()
403}
404
405struct OutputCleanupGuard {
413 segment_id: SegmentId,
414 cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
415}
416
417impl OutputCleanupGuard {
418 fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
419 Self {
420 segment_id,
421 cleanup: Some(cleanup),
422 }
423 }
424
425 fn disarm(&mut self) {
426 self.cleanup = None;
427 }
428}
429
430impl Drop for OutputCleanupGuard {
431 fn drop(&mut self) {
432 if let Some(cleanup) = self.cleanup.take() {
433 cleanup(self.segment_id);
434 }
435 }
436}
437
438struct ManagerState {
440 metadata: IndexMetadata,
441 merge_policy: Box<dyn MergePolicy>,
442}
443
444#[cfg(feature = "native")]
445struct MergeTaskError {
446 error: Error,
447 unavailable_segments: Vec<String>,
448}
449
450#[cfg(feature = "native")]
451impl MergeTaskError {
452 fn source(segment_id: String, error: Error) -> Self {
453 Self {
454 error,
455 unavailable_segments: vec![segment_id],
456 }
457 }
458
459 fn sources(segment_ids: Vec<String>, error: Error) -> Self {
460 Self {
461 error,
462 unavailable_segments: segment_ids,
463 }
464 }
465}
466
467#[cfg(feature = "native")]
468impl From<Error> for MergeTaskError {
469 fn from(error: Error) -> Self {
470 Self {
471 error,
472 unavailable_segments: Vec::new(),
473 }
474 }
475}
476
477#[cfg(feature = "native")]
478fn is_deterministic_source_error(error: &Error) -> bool {
479 matches!(error, Error::Corruption(_) | Error::Serialization(_))
480 || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
481}
482
483#[cfg(feature = "native")]
484fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
485 if is_deterministic_source_error(&error) {
486 MergeTaskError::source(segment_id, error)
487 } else {
488 MergeTaskError::from(error)
492 }
493}
494
495#[cfg(feature = "native")]
496type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
497
498#[derive(Clone, Copy)]
499enum ReplacementLayout {
500 Recomputed {
501 reordered: bool,
502 bp_converged: bool,
503 },
504 PreserveSingleSource,
507}
508
509#[derive(Clone, Copy, Debug, Eq, PartialEq)]
510enum VectorSegmentRewriteOutcome {
511 Rewritten,
512 AlreadyCurrent,
513 SourceGone,
514 Conflict,
515 Deferred,
516}
517
518pub(crate) struct StagedVectorSegment {
521 source_id: String,
522 output_id: SegmentId,
523 doc_count: u32,
524 _operation: SegmentOperationGuard,
525 cleanup: OutputCleanupGuard,
526}
527
528pub struct SegmentManager<D: DirectoryWriter + 'static> {
532 state: Arc<AsyncMutex<ManagerState>>,
534
535 active_operations: Arc<ActiveSegmentOperations>,
537
538 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
543
544 merge_retry: parking_lot::Mutex<MergeRetryState>,
547
548 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
552
553 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
555
556 global_merge_wakeup_pending: AtomicBool,
560
561 #[cfg(test)]
564 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
565
566 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
570
571 trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
575
576 vector_artifact_update: Arc<AtomicBool>,
580
581 tracker: Arc<SegmentTracker>,
583
584 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
586
587 directory: Arc<D>,
589 schema: Arc<crate::dsl::Schema>,
591 term_cache_blocks: usize,
593 merge_permits: Arc<Semaphore>,
597 global_merge_permits: Arc<Semaphore>,
599 reorder_permits: Arc<Semaphore>,
603 reorder_on_merge: bool,
608 merge_bp_time_budget: Option<std::time::Duration>,
612 bp_memory_budget_bytes: usize,
615 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
618}
619
620impl<D: DirectoryWriter + 'static> SegmentManager<D> {
621 #[allow(clippy::too_many_arguments)]
623 pub fn new(
624 directory: Arc<D>,
625 schema: Arc<crate::dsl::Schema>,
626 metadata: IndexMetadata,
627 merge_policy: Box<dyn MergePolicy>,
628 term_cache_blocks: usize,
629 max_concurrent_merges: usize,
630 global_merge_permits: Arc<Semaphore>,
631 merge_bp_time_budget: Option<std::time::Duration>,
632 bp_memory_budget_bytes: usize,
633 reorder_permits: Arc<Semaphore>,
634 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
635 ) -> Self {
636 let reorder_on_merge = schema.reorder_on_merge();
639 if reorder_on_merge {
640 log::info!("[merge] reorder-on-merge enabled by index schema");
641 }
642
643 let tracker = Arc::new(SegmentTracker::new());
644 for seg_id in metadata.segment_metas.keys() {
645 tracker.register(seg_id);
646 }
647
648 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
649 Arc::new(parking_lot::Mutex::new(Vec::new()));
650 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
651 let dir = Arc::clone(&directory);
652 let tracker = Arc::clone(&tracker);
653 let lifecycle_handles = Arc::clone(&lifecycle_handles);
654 Arc::new(move |segment_ids| {
655 let Ok(handle) = tokio::runtime::Handle::try_current() else {
658 tracker.complete_deletion(&segment_ids);
661 return;
662 };
663 let dir = Arc::clone(&dir);
664 let task_tracker = Arc::clone(&tracker);
665 let cleanup_ids = segment_ids.clone();
666 let future = async move {
667 for &segment_id in &segment_ids {
668 log::info!(
669 "[segment_cleanup] deleting deferred segment {}",
670 segment_id.to_hex()
671 );
672 if let Err(error) =
673 crate::segment::delete_segment(dir.as_ref(), segment_id).await
674 {
675 log::warn!(
676 "[segment_cleanup] deferred delete failed for {}: {}",
677 segment_id.to_hex(),
678 error,
679 );
680 }
681 }
682 task_tracker.complete_deletion(&segment_ids);
683 };
684 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
685 tracker.complete_deletion(&cleanup_ids);
689 log::warn!(
690 "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
691 );
692 }
693 })
694 };
695
696 Self {
697 state: Arc::new(AsyncMutex::new(ManagerState {
698 metadata,
699 merge_policy,
700 })),
701 active_operations: Arc::new(ActiveSegmentOperations::new()),
702 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
703 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
704 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
705 merge_handles: parking_lot::Mutex::new(Vec::new()),
706 global_merge_wakeup_pending: AtomicBool::new(false),
707 #[cfg(test)]
708 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
709 lifecycle_handles,
710 trained: Arc::new(ArcSwapOption::new(None)),
711 vector_artifact_update: Arc::new(AtomicBool::new(false)),
712 tracker,
713 delete_fn,
714 directory,
715 schema,
716 term_cache_blocks,
717 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
718 global_merge_permits,
719 reorder_permits,
720 reorder_on_merge,
721 merge_bp_time_budget,
722 bp_memory_budget_bytes,
723 background_reorder_pool,
724 }
725 }
726
727 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
733 if let Some(pool) = &self.background_reorder_pool {
734 return Arc::clone(pool);
735 }
736 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
737 let threads = (num_cpus::get() / 2).max(1);
738 log::info!(
739 "[merge] process-wide background CPU pool: {} thread(s)",
740 threads
741 );
742 Arc::new(
743 rayon::ThreadPoolBuilder::new()
744 .num_threads(threads)
745 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
746 .build()
747 .expect("failed to build background CPU pool"),
748 )
749 }))
750 }
751
752 pub fn begin_shutdown(&self) {
756 self.active_operations.stop_accepting();
757 }
758
759 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
767 where
768 T: Send + 'static,
769 F: std::future::Future<Output = Result<T>> + Send + 'static,
770 {
771 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
772 let future = async move {
773 let result = transaction.await;
774 let _ = result_tx.send(result);
775 };
776 let runtime = tokio::runtime::Handle::current();
777 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
778 return Err(Error::Internal(
779 "runtime rejected lifecycle metadata transaction".into(),
780 ));
781 }
782 result_rx.await.map_err(|_| {
783 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
784 })?
785 }
786
787 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
789 let manager = Arc::clone(self);
790 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
791 let Ok(handle) = tokio::runtime::Handle::try_current() else {
792 log::warn!(
793 "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
794 segment_id.to_hex(),
795 );
796 return;
797 };
798
799 let cleanup_manager = Arc::clone(&manager);
800 let future = async move {
801 cleanup_manager
802 .delete_output_if_unregistered(segment_id, "task unwind")
803 .await;
804 };
805 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
806 log::warn!(
807 "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
808 segment_id.to_hex(),
809 );
810 }
811 });
812
813 OutputCleanupGuard::new(output_id, cleanup)
814 }
815
816 pub(crate) fn schedule_unpublished_segment_cleanup(
821 self: &Arc<Self>,
822 output_id: SegmentId,
823 operation: SegmentOperationGuard,
824 runtime: tokio::runtime::Handle,
825 ) {
826 let manager = Arc::clone(self);
827 let output_hex = output_id.to_hex();
828 let future = async move {
829 manager
830 .delete_output_if_unregistered(output_id, "indexing abort or failure")
831 .await;
832 drop(operation);
833 };
834 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
835 log::warn!(
838 "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
839 output_hex,
840 );
841 }
842 }
843
844 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
850 match self
851 .active_operations
852 .try_register_indexing(vec![segment_id.clone()])
853 {
854 Some(operation) => Ok(operation),
855 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
856 None => Err(Error::Corruption(format!(
857 "new segment ID {} is already owned by an active operation",
858 segment_id
859 ))),
860 }
861 }
862
863 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
867 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
868 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
869 })?;
870 let files = SegmentFiles::new(id.0);
871
872 for path in files.mandatory_paths() {
873 if !self.directory.exists(path).await.map_err(Error::Io)? {
874 return Err(Error::Corruption(format!(
875 "segment {} cannot be published: mandatory file {:?} is missing",
876 segment_id, path
877 )));
878 }
879 }
880
881 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
882 Error::Corruption(format!(
883 "segment {} cannot be published: missing/unreadable {:?}: {}",
884 segment_id, files.meta, e
885 ))
886 })?;
887 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
888 Error::Corruption(format!(
889 "segment {} cannot be published: failed reading {:?}: {}",
890 segment_id, files.meta, e
891 ))
892 })?;
893 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
894 Error::Corruption(format!(
895 "segment {} cannot be published: invalid {:?}: {}",
896 segment_id, files.meta, e
897 ))
898 })?;
899
900 if meta.id != id.0 || meta.num_docs != expected_docs {
901 return Err(Error::Corruption(format!(
902 "segment {} cannot be published: metadata identity/docs mismatch \
903 (id={:032x}, docs={}, expected_docs={})",
904 segment_id, meta.id, meta.num_docs, expected_docs
905 )));
906 }
907
908 Ok(())
909 }
910
911 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
912 let inserted = self
913 .quarantined_segments
914 .lock()
915 .insert(segment_id.to_string());
916 if inserted {
917 log::error!(
918 "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
919 It remains metadata-live for explicit repair but is excluded from merges until restart",
920 segment_id,
921 error,
922 );
923 }
924 }
925
926 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
927 let mut retry = self.merge_retry.lock();
928 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
929 let delay = merge_retry_delay(retry.consecutive_failures);
930 retry.retry_after = std::time::Instant::now().checked_add(delay);
931 log::warn!(
932 "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
933 delay.as_secs_f64(),
934 retry.consecutive_failures,
935 error,
936 );
937 delay
938 }
939
940 fn clear_merge_retry_backoff(&self) {
941 *self.merge_retry.lock() = MergeRetryState::default();
942 }
943
944 fn merge_retry_is_paused(&self) -> bool {
945 let mut retry = self.merge_retry.lock();
946 match retry.retry_after {
947 Some(deadline) if deadline > std::time::Instant::now() => true,
948 Some(_) => {
949 retry.retry_after = None;
950 false
951 }
952 None => false,
953 }
954 }
955
956 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
957 let mut retries = self.reorder_retries.lock();
958 let retry = retries.entry(segment_id.to_string()).or_default();
959 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
960 let delay = merge_retry_delay(retry.consecutive_failures);
961 retry.retry_after = std::time::Instant::now().checked_add(delay);
962 log::warn!(
963 "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
964 segment_id,
965 delay.as_secs_f64(),
966 retry.consecutive_failures,
967 error,
968 );
969 }
970
971 fn clear_reorder_retry(&self, segment_id: &str) {
972 self.reorder_retries.lock().remove(segment_id);
973 }
974
975 fn paused_reorder_segments(&self) -> HashSet<String> {
976 let now = std::time::Instant::now();
977 let mut retries = self.reorder_retries.lock();
978 let mut paused = HashSet::new();
979 for (segment_id, retry) in retries.iter_mut() {
980 match retry.retry_after {
981 Some(deadline) if deadline > now => {
982 paused.insert(segment_id.clone());
983 }
984 Some(_) => retry.retry_after = None,
985 None => {}
986 }
987 }
988 paused
989 }
990
991 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
995 if self
996 .global_merge_wakeup_pending
997 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
998 .is_err()
999 {
1000 return;
1001 }
1002
1003 let manager = Arc::clone(self);
1004 let future = async move {
1005 let capacity = tokio::select! {
1006 biased;
1007 () = manager.active_operations.wait_for_shutdown() => None,
1008 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1009 };
1010
1011 manager
1012 .global_merge_wakeup_pending
1013 .store(false, Ordering::Release);
1014 if let Some(permit) = capacity {
1015 drop(permit);
1019 manager.maybe_merge().await;
1020 }
1021 };
1022 let runtime = tokio::runtime::Handle::current();
1023 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1024 self.global_merge_wakeup_pending
1025 .store(false, Ordering::Release);
1026 log::warn!("[merge] runtime rejected global-capacity wakeup task");
1027 }
1028 }
1029
1030 #[cfg(test)]
1031 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1032 self.quarantined_segments.lock().contains(segment_id)
1033 }
1034
1035 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1040 let output_hex = output_id.to_hex();
1041 {
1042 let st = self.state.lock().await;
1043 if st.metadata.has_segment(&output_hex) {
1044 return;
1045 }
1046 }
1047
1048 log::info!(
1052 "[segment_cleanup] deleting uncommitted output {} after {}",
1053 output_hex,
1054 reason,
1055 );
1056 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1057 {
1058 log::warn!(
1059 "[segment_cleanup] failed deleting uncommitted output {}: {}",
1060 output_hex,
1061 error,
1062 );
1063 }
1064 }
1065
1066 pub async fn get_segment_ids(&self) -> Vec<String> {
1072 self.state.lock().await.metadata.segment_ids()
1073 }
1074
1075 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1077 self.trained.load_full()
1078 }
1079
1080 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1087 if self.vector_artifact_update.load(Ordering::Acquire) {
1088 return None;
1089 }
1090 let trained = self.trained.load_full();
1091 if self.vector_artifact_update.load(Ordering::Acquire) {
1092 None
1093 } else {
1094 trained
1095 }
1096 }
1097
1098 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1117 self.vector_artifact_update
1118 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1119 .map_err(|_| {
1120 Error::Internal("a trained-vector artifact update is already in progress".into())
1121 })?;
1122 self.active_operations.pause_non_indexing();
1123 let guard = VectorArtifactUpdateGuard {
1124 _lease: Arc::new(VectorArtifactUpdateLease {
1125 updating: Arc::clone(&self.vector_artifact_update),
1126 active_operations: Arc::clone(&self.active_operations),
1127 }),
1128 };
1129 let (preexisting, parked_indexing) =
1130 self.active_operations.draining_operation_tokens_snapshot();
1131 if parked_indexing > 0 {
1132 return Err(Error::Internal(format!(
1133 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1134 segment(s) are built but uncommitted; commit or abort the pending \
1135 generation and retry"
1136 )));
1137 }
1138 self.active_operations
1139 .wait_until_operations_finish(&preexisting)
1140 .await;
1141 Ok(guard)
1142 }
1143
1144 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1147 let vector_fields = {
1149 let st = self.state.lock().await;
1150 st.metadata.vector_fields.clone()
1151 };
1152 let trained = IndexMetadata::try_load_trained_from_fields(
1154 &vector_fields,
1155 self.schema.as_ref(),
1156 self.directory.as_ref(),
1157 )
1158 .await?
1159 .map(Arc::new);
1160 self.trained.store(trained);
1164 Ok(())
1165 }
1166
1167 pub(crate) async fn publish_vector_generation(
1174 self: &Arc<Self>,
1175 artifact_update: &VectorArtifactUpdateGuard,
1176 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1177 next_trained: Arc<TrainedVectorStructures>,
1178 mut staged: Vec<StagedVectorSegment>,
1179 ) -> Result<()> {
1180 if !self.vector_artifact_update.load(Ordering::Acquire) {
1181 return Err(Error::Internal(
1182 "vector generation publication lost its exclusive update lease".into(),
1183 ));
1184 }
1185
1186 for replacement in &staged {
1187 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1188 .await?;
1189 }
1190
1191 let mut st = Arc::clone(&self.state).lock_owned().await;
1192 let mut next = st.metadata.clone();
1193 next.vector_fields = vector_fields;
1194 next.refresh_total_vectors();
1195
1196 for replacement in &staged {
1197 let source_info = next
1198 .segment_metas
1199 .remove(&replacement.source_id)
1200 .ok_or_else(|| {
1201 Error::Corruption(format!(
1202 "vector generation source {} disappeared before publication",
1203 replacement.source_id,
1204 ))
1205 })?;
1206 let output_hex = replacement.output_id.to_hex();
1207 if next.segment_metas.contains_key(&output_hex) {
1208 return Err(Error::Corruption(format!(
1209 "vector generation output {output_hex} is already metadata-live"
1210 )));
1211 }
1212 next.add_segment_meta(output_hex, source_info);
1215 }
1216
1217 let directory = Arc::clone(&self.directory);
1218 let trained = Arc::clone(&self.trained);
1219 let tracker = Arc::clone(&self.tracker);
1220 let artifact_update = artifact_update.clone();
1224 self.run_lifecycle_transaction(async move {
1225 let _artifact_update = artifact_update;
1226 next.save(directory.as_ref()).await?;
1227
1228 for replacement in &staged {
1229 tracker.register(&replacement.output_id.to_hex());
1230 }
1231 st.metadata = next;
1232 trained.store(Some(next_trained));
1233
1234 for replacement in &mut staged {
1237 replacement.cleanup.disarm();
1238 }
1239 let retired = staged
1240 .iter()
1241 .map(|replacement| replacement.source_id.clone())
1242 .collect::<Vec<_>>();
1243 let ready_to_delete = tracker.mark_for_deletion(&retired);
1244 drop(st);
1245 for &segment_id in &ready_to_delete {
1246 if let Err(error) =
1247 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1248 {
1249 log::warn!(
1250 "[segment_cleanup] immediate dense-vector generation delete failed for {}: {}",
1251 segment_id.to_hex(),
1252 error,
1253 );
1254 }
1255 }
1256 tracker.complete_deletion(&ready_to_delete);
1257 Ok(())
1258 })
1259 .await
1260 }
1261
1262 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1264 where
1265 F: FnOnce(&IndexMetadata) -> R,
1266 {
1267 let st = self.state.lock().await;
1268 f(&st.metadata)
1269 }
1270
1271 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1273 where
1274 F: FnOnce(&mut IndexMetadata),
1275 {
1276 let mut st = Arc::clone(&self.state).lock_owned().await;
1277 let mut next = st.metadata.clone();
1278 f(&mut next);
1279 let directory = Arc::clone(&self.directory);
1280 self.run_lifecycle_transaction(async move {
1281 next.save(directory.as_ref()).await?;
1282 st.metadata = next;
1283 Ok(())
1284 })
1285 .await
1286 }
1287
1288 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1291 let (acquired, trained) = {
1292 let st = self.state.lock().await;
1293 let segment_ids = st.metadata.segment_ids();
1294 (self.tracker.acquire(&segment_ids), self.trained.load_full())
1295 };
1296
1297 SegmentSnapshot::with_generation(
1298 Arc::clone(&self.tracker),
1299 acquired,
1300 trained,
1301 Arc::clone(&self.delete_fn),
1302 )
1303 }
1304
1305 pub fn tracker(&self) -> Arc<SegmentTracker> {
1307 Arc::clone(&self.tracker)
1308 }
1309
1310 pub fn directory(&self) -> Arc<D> {
1312 Arc::clone(&self.directory)
1313 }
1314}
1315
1316#[cfg(feature = "native")]
1321impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1322 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1324 for (segment_id, num_docs) in new_segments {
1327 self.validate_completed_segment(segment_id, *num_docs)
1328 .await?;
1329 }
1330
1331 let mut st = Arc::clone(&self.state).lock_owned().await;
1332 let mut next = st.metadata.clone();
1333 let mut added = Vec::new();
1334 for (segment_id, num_docs) in new_segments {
1335 if !next.has_segment(segment_id) {
1336 next.add_segment(segment_id.clone(), *num_docs);
1337 added.push(segment_id.clone());
1338 }
1339 }
1340
1341 let directory = Arc::clone(&self.directory);
1347 let tracker = Arc::clone(&self.tracker);
1348 self.run_lifecycle_transaction(async move {
1349 next.save(directory.as_ref()).await?;
1350 for segment_id in &added {
1351 tracker.register(segment_id);
1352 }
1353 st.metadata = next;
1354 Ok(())
1355 })
1356 .await
1357 }
1358
1359 pub async fn maybe_merge(self: &Arc<Self>) {
1370 if !self.active_operations.is_accepting() {
1371 log::debug!("[maybe_merge] manager is shutting down, skipping");
1372 return;
1373 }
1374 if self.merge_retry_is_paused() {
1375 log::debug!("[maybe_merge] retry backoff active, skipping");
1376 return;
1377 }
1378
1379 {
1382 let mut handles = self.merge_handles.lock();
1383 handles.retain(|h| !h.is_finished());
1384 }
1385 let local_slots = self.merge_permits.available_permits();
1386 let global_slots = self.global_merge_permits.available_permits();
1387 let slots_available = local_slots.min(global_slots);
1388
1389 let new_handles = {
1393 let st = self.state.lock().await;
1394 let quarantined = self.quarantined_segments.lock().clone();
1395 let active_ids = self.active_operations.snapshot();
1396
1397 let segments: Vec<SegmentInfo> = st
1400 .metadata
1401 .segment_metas
1402 .iter()
1403 .filter(|(id, _)| {
1404 !self.tracker.is_pending_deletion(id)
1405 && !active_ids.contains(*id)
1406 && !quarantined.contains(*id)
1407 })
1408 .map(|(id, info)| SegmentInfo {
1409 id: id.clone(),
1410 num_docs: info.num_docs,
1411 })
1412 .collect();
1413
1414 log::debug!("[maybe_merge] {} eligible segments", segments.len());
1415
1416 let candidates = st.merge_policy.find_merges(&segments);
1417
1418 if candidates.is_empty() {
1419 return;
1420 }
1421
1422 if slots_available == 0 {
1426 if local_slots > 0 && global_slots == 0 {
1427 self.schedule_global_merge_wakeup();
1428 }
1429 log::debug!("[maybe_merge] at max concurrent merges, skipping");
1430 return;
1431 }
1432
1433 log::debug!(
1434 "[maybe_merge] {} merge candidates, {} slots available",
1435 candidates.len(),
1436 slots_available
1437 );
1438
1439 let mut handles = Vec::new();
1440 for c in candidates {
1441 if handles.len() >= slots_available {
1442 break;
1443 }
1444 if let Some(h) = self.spawn_merge(c.segment_ids) {
1445 handles.push(h);
1446 }
1447 }
1448 handles
1449 };
1451
1452 if !new_handles.is_empty() {
1453 self.merge_handles.lock().extend(new_handles);
1457 }
1458 }
1459
1460 fn spawn_merge(self: &Arc<Self>, segment_ids_to_merge: Vec<String>) -> Option<JoinHandle<()>> {
1469 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1470 Ok(permit) => permit,
1471 Err(_) => {
1472 log::debug!("[spawn_merge] skipped: global merge capacity is full");
1473 self.schedule_global_merge_wakeup();
1474 return None;
1475 }
1476 };
1477 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1478 Ok(permit) => permit,
1479 Err(_) => {
1480 log::debug!("[spawn_merge] skipped: no merge permit available");
1481 return None;
1482 }
1483 };
1484 let output_id = SegmentId::new();
1485 let output_hex = output_id.to_hex();
1486
1487 let mut all_ids = segment_ids_to_merge.clone();
1488 all_ids.push(output_hex);
1489
1490 let guard = match self.active_operations.try_register(all_ids) {
1491 Some(g) => g,
1492 None => {
1493 log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1494 return None;
1495 }
1496 };
1497
1498 let sm = Arc::clone(self);
1499 let ids = segment_ids_to_merge;
1500
1501 Some(tokio::spawn(async move {
1502 let mut output_cleanup = sm.output_cleanup_guard(output_id);
1503 let mut reevaluate = false;
1504 let mut retry_delay = None;
1505
1506 let trained_snap = sm.trained_for_segment_build();
1507 let granularity = sm.merge_granularity(&ids).await;
1508 let result = Self::do_merge(
1509 sm.directory.as_ref(),
1510 &sm.schema,
1511 &ids,
1512 output_id,
1513 sm.term_cache_blocks,
1514 trained_snap.as_deref(),
1515 sm.reorder_on_merge,
1516 granularity,
1517 sm.merge_bp_time_budget,
1518 sm.bp_memory_budget_bytes,
1519 Arc::clone(&sm.reorder_permits),
1520 Some(sm.background_cpu_pool()),
1521 )
1522 .await;
1523
1524 match result {
1525 Ok((new_id, doc_count, bp_converged)) => {
1526 match sm
1527 .replace_segments(
1528 &ids,
1529 new_id,
1530 doc_count,
1531 ReplacementLayout::Recomputed {
1532 reordered: sm.reorder_on_merge,
1533 bp_converged,
1534 },
1535 )
1536 .await
1537 {
1538 Ok(()) => {
1539 output_cleanup.disarm();
1540 sm.clear_merge_retry_backoff();
1541 reevaluate = true;
1542 }
1543 Err(e) => {
1544 sm.delete_output_if_unregistered(output_id, "replacement failure")
1545 .await;
1546 output_cleanup.disarm();
1547 retry_delay = Some(sm.pause_merge_retries(&e));
1548 log::error!("[merge] failed to publish merged segment: {}", e);
1549 }
1550 }
1551 }
1552 Err(MergeTaskError {
1553 error,
1554 unavailable_segments,
1555 }) => {
1556 log::error!(
1557 "[merge] background merge failed for segments {:?}: {}",
1558 ids,
1559 error
1560 );
1561 if !unavailable_segments.is_empty() {
1562 for segment_id in &unavailable_segments {
1563 sm.quarantine_segment(segment_id, &error);
1564 }
1565 reevaluate = true;
1569 } else {
1570 retry_delay = Some(sm.pause_merge_retries(&error));
1571 }
1572 sm.delete_output_if_unregistered(output_id, "merge failure")
1573 .await;
1574 output_cleanup.disarm();
1575 }
1576 }
1577 drop(guard);
1580 drop(merge_permit);
1582 drop(global_merge_permit);
1583
1584 if reevaluate {
1585 sm.maybe_merge().await;
1586 } else if let Some(retry_delay) = retry_delay {
1587 sm.schedule_merge_retry_wakeup(retry_delay);
1594 }
1595 }))
1596 }
1597
1598 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
1602 let manager = Arc::clone(self);
1603 let future = async move {
1604 tokio::select! {
1605 () = tokio::time::sleep(retry_delay) => {
1606 manager.maybe_merge().await;
1607 }
1608 () = manager.active_operations.wait_for_shutdown() => {}
1609 }
1610 };
1611 let runtime = tokio::runtime::Handle::current();
1612 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1613 log::warn!(
1614 "[merge] runtime rejected merge-retry wakeup task; eligible segments may stay \
1615 unmerged until the next commit re-runs merge policy evaluation"
1616 );
1617 }
1618 }
1619
1620 async fn replace_segments(
1624 self: &Arc<Self>,
1625 old_ids: &[String],
1626 new_id: String,
1627 doc_count: u32,
1628 layout: ReplacementLayout,
1629 ) -> Result<()> {
1630 self.validate_completed_segment(&new_id, doc_count).await?;
1633 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1634 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1635 })?;
1636 let output_reader = SegmentReader::open(
1637 self.directory.as_ref(),
1638 output_id,
1639 Arc::clone(&self.schema),
1640 self.term_cache_blocks,
1641 )
1642 .await
1643 .map_err(|error| match error {
1644 Error::Io(_) | Error::IndexClosed => error,
1648 error => Error::Corruption(format!(
1649 "replacement segment {new_id} failed full reader validation: {error}"
1650 )),
1651 })?;
1652 if output_reader.num_docs() != doc_count {
1653 return Err(Error::Corruption(format!(
1654 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
1655 output_reader.num_docs(),
1656 )));
1657 }
1658 drop(output_reader);
1659
1660 let mut st = Arc::clone(&self.state).lock_owned().await;
1661 let missing: Vec<&String> = old_ids
1665 .iter()
1666 .filter(|id| !st.metadata.has_segment(id))
1667 .collect();
1668 if !missing.is_empty() {
1669 return Err(Error::Corruption(format!(
1670 "replace_segments: source segment(s) {:?} not in metadata — \
1671 refusing to add output {} (would duplicate documents)",
1672 missing, new_id
1673 )));
1674 }
1675
1676 let replacement_info = match layout {
1677 ReplacementLayout::Recomputed {
1678 reordered,
1679 bp_converged,
1680 } => {
1681 let generation = old_ids
1682 .iter()
1683 .filter_map(|id| st.metadata.segment_metas.get(id))
1684 .map(|info| info.generation)
1685 .max()
1686 .unwrap_or(0)
1687 .checked_add(1)
1688 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
1689 let parent_unconverged_passes = old_ids
1690 .iter()
1691 .filter_map(|id| st.metadata.segment_metas.get(id))
1692 .map(|info| info.bp_unconverged_passes)
1693 .max()
1694 .unwrap_or(0);
1695 let passes = if reordered && !bp_converged {
1696 parent_unconverged_passes.saturating_add(1)
1697 } else {
1698 0
1699 };
1700 SegmentMetaInfo {
1701 num_docs: doc_count,
1702 ancestors: old_ids.to_vec(),
1703 generation,
1704 reordered,
1705 bp_converged,
1706 bp_unconverged_passes: passes,
1707 }
1708 }
1709 ReplacementLayout::PreserveSingleSource => {
1710 let [source_id] = old_ids else {
1711 return Err(Error::Internal(
1712 "layout-preserving replacement requires exactly one source".into(),
1713 ));
1714 };
1715 let mut source = st
1716 .metadata
1717 .segment_metas
1718 .get(source_id)
1719 .cloned()
1720 .ok_or_else(|| {
1721 Error::Corruption(format!(
1722 "layout-preserving replacement source {source_id} disappeared"
1723 ))
1724 })?;
1725 source.num_docs = doc_count;
1726 source
1727 }
1728 };
1729 let retired_ids = old_ids.to_vec();
1730 let mut next = st.metadata.clone();
1731 for id in old_ids {
1732 next.remove_segment(id);
1733 }
1734 next.add_segment_meta(new_id.clone(), replacement_info);
1735
1736 let directory = Arc::clone(&self.directory);
1737 let tracker = Arc::clone(&self.tracker);
1738 self.run_lifecycle_transaction(async move {
1739 next.save(directory.as_ref()).await?;
1742 tracker.register(&new_id);
1743 st.metadata = next;
1744
1745 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
1749 drop(st);
1750 for &segment_id in &ready_to_delete {
1751 if let Err(error) =
1752 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1753 {
1754 log::warn!(
1755 "[segment_cleanup] immediate delete failed for {}: {}",
1756 segment_id.to_hex(),
1757 error,
1758 );
1759 }
1760 }
1761 tracker.complete_deletion(&ready_to_delete);
1762 Ok(())
1763 })
1764 .await
1765 }
1766
1767 #[allow(clippy::too_many_arguments)]
1772 async fn do_merge(
1773 directory: &D,
1774 schema: &Arc<crate::dsl::Schema>,
1775 segment_ids_to_merge: &[String],
1776 output_segment_id: SegmentId,
1777 term_cache_blocks: usize,
1778 trained: Option<&TrainedVectorStructures>,
1779 reorder_bmp: bool,
1780 granularity: crate::segment::reorder::BpGranularity,
1781 merge_bp_time_budget: Option<std::time::Duration>,
1782 bp_memory_budget_bytes: usize,
1783 reorder_permits: Arc<Semaphore>,
1784 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
1785 ) -> MergeTaskResult<(String, u32, bool)> {
1786 let output_hex = output_segment_id.to_hex();
1787 let load_start = std::time::Instant::now();
1788
1789 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
1790 for id_str in segment_ids_to_merge {
1791 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
1792 MergeTaskError::source(
1793 id_str.clone(),
1794 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
1795 )
1796 })?;
1797 segment_ids.push(id);
1798 }
1799
1800 let mut unavailable_sources = Vec::new();
1805 let mut missing_files = Vec::new();
1806 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
1807 let files = SegmentFiles::new(id.0);
1808 let mut source_unavailable = false;
1809 for path in files.mandatory_paths() {
1810 let exists = directory
1811 .exists(path)
1812 .await
1813 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
1814 if !exists {
1815 source_unavailable = true;
1816 missing_files.push(format!("{}:{:?}", id_str, path));
1817 }
1818 }
1819 if source_unavailable {
1820 unavailable_sources.push(id_str.clone());
1821 }
1822 }
1823 if !unavailable_sources.is_empty() {
1824 return Err(MergeTaskError::sources(
1825 unavailable_sources,
1826 Error::Corruption(format!(
1827 "merge sources are missing mandatory files: {}",
1828 missing_files.join(", ")
1829 )),
1830 ));
1831 }
1832
1833 let schema_arc = Arc::clone(schema);
1834 let futures: Vec<_> = segment_ids
1835 .iter()
1836 .map(|&sid| {
1837 let sch = Arc::clone(&schema_arc);
1838 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
1839 })
1840 .collect();
1841
1842 let results = futures::future::join_all(futures).await;
1843 let mut readers = Vec::with_capacity(results.len());
1844 let mut total_docs = 0u64;
1845 for (i, result) in results.into_iter().enumerate() {
1846 match result {
1847 Ok(r) => {
1848 total_docs += r.meta().num_docs as u64;
1849 readers.push(r);
1850 }
1851 Err(e) => {
1852 log::error!(
1853 "[merge] Failed to open segment {}: {:?}",
1854 segment_ids_to_merge[i],
1855 e
1856 );
1857 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
1858 }
1859 }
1860 }
1861 if total_docs > u32::MAX as u64 {
1862 return Err(Error::Internal(format!(
1863 "Merged segment doc count ({}) exceeds u32::MAX",
1864 total_docs
1865 ))
1866 .into());
1867 }
1868
1869 for (i, reader) in readers.iter().enumerate() {
1873 let meta_docs = reader.meta().num_docs;
1874 let store_docs = reader.store().num_docs();
1875 if store_docs != meta_docs {
1876 return Err(MergeTaskError::source(
1877 segment_ids_to_merge[i].clone(),
1878 Error::Corruption(format!(
1879 "pre-merge validation: segment {} store has {} docs but meta says {}",
1880 segment_ids_to_merge[i], store_docs, meta_docs
1881 )),
1882 ));
1883 }
1884 }
1885
1886 log::info!(
1887 "[merge] loaded {} segment readers in {:.1}s",
1888 readers.len(),
1889 load_start.elapsed().as_secs_f64()
1890 );
1891
1892 let merger = SegmentMerger::new(Arc::clone(schema))
1893 .with_bmp_reorder(reorder_bmp)
1894 .with_granularity(granularity)
1895 .with_bp_budget(crate::segment::BpBudget {
1896 min_partition_docs: None,
1897 time_budget: merge_bp_time_budget,
1898 })
1899 .with_bp_memory_budget(bp_memory_budget_bytes)
1900 .with_reorder_permits(reorder_permits)
1901 .with_background_pool(bg_cpu_pool);
1902
1903 log::info!(
1904 "[merge] {} segments -> {} (trained={})",
1905 segment_ids_to_merge.len(),
1906 output_hex,
1907 trained.map_or(0, |t| t.centroids.len()),
1908 );
1909
1910 let (_merged_meta, merge_stats) = merger
1911 .merge(directory, &readers, output_segment_id, trained)
1912 .await
1913 .map_err(|error| {
1914 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
1915 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
1921 } else {
1922 MergeTaskError::from(error)
1923 }
1924 })?;
1925 let bp_converged = merge_stats.bp_converged;
1926 if !bp_converged {
1927 log::info!(
1928 "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
1929 the background optimizer deepens it later",
1930 );
1931 }
1932
1933 log::info!(
1934 "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
1935 load_start.elapsed().as_secs_f64(),
1936 readers.len(),
1937 total_docs,
1938 );
1939
1940 Ok((output_hex, total_docs as u32, bp_converged))
1941 }
1942
1943 pub async fn abort_merges(&self) {
1954 loop {
1955 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
1956 if handles.is_empty() {
1957 return;
1958 }
1959 while let Some(result) = handles.join_next().await {
1960 if let Err(error) = result
1961 && error.is_panic()
1962 {
1963 log::error!("[merge] background task panicked while draining: {}", error);
1964 }
1965 }
1966 }
1967 }
1968
1969 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
1974 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
1975 while handles.join_next().await.is_some() {}
1976 }
1977
1978 pub async fn wait_for_all_merges(self: &Arc<Self>) {
1987 loop {
1988 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
1989 if handles.is_empty() {
1990 break;
1991 }
1992 while handles.join_next().await.is_some() {}
1993 }
1994 }
1995
1996 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2001 self.wait_for_all_merges().await;
2002 self.active_operations.wait_until_idle().await;
2003 loop {
2004 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2005 if handles.is_empty() {
2006 break;
2007 }
2008 for handle in handles {
2009 if let Err(error) = handle.await
2010 && error.is_panic()
2011 {
2012 log::error!("[segment_cleanup] task panicked while draining: {}", error);
2013 }
2014 }
2015 }
2016 }
2017
2018 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2027 const FORCE_MERGE_BATCH: usize = 64;
2028 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2033 std::time::Duration::from_millis(100);
2034 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2038
2039 let max_segment_docs = {
2040 let st = self.state.lock().await;
2041 st.merge_policy.max_segment_docs()
2042 };
2043
2044 self.wait_for_all_merges().await;
2047
2048 let mut logged_held_wait = false;
2050
2051 loop {
2052 if !self.active_operations.is_accepting() {
2053 return Err(Error::IndexClosed);
2054 }
2055 let mut segments: Vec<(String, u32)> = {
2057 let st = self.state.lock().await;
2058 st.metadata
2059 .segment_metas
2060 .iter()
2061 .map(|(id, info)| (id.clone(), info.num_docs))
2062 .collect()
2063 };
2064
2065 if segments.len() < 2 {
2066 return Ok(());
2067 }
2068
2069 segments.sort_by_key(|(_, docs)| *docs);
2070
2071 let active_ids = self.active_operations.snapshot();
2079 let held: usize = segments
2080 .iter()
2081 .filter(|(id, _)| active_ids.contains(id))
2082 .count();
2083
2084 let max_docs = max_segment_docs.map(|m| m as u64).unwrap_or(u64::MAX);
2086 let mut batch = Vec::new();
2087 let mut batch_docs = 0u64;
2088
2089 for (id, docs) in &segments {
2090 if active_ids.contains(id) {
2091 continue;
2092 }
2093 if batch.len() >= FORCE_MERGE_BATCH {
2094 break;
2095 }
2096 let next_total = batch_docs + *docs as u64;
2097 if next_total > max_docs && !batch.is_empty() {
2098 break;
2099 }
2100 batch.push(id.clone());
2101 batch_docs += *docs as u64;
2102 }
2103
2104 if batch.len() < 2 {
2105 if held == 0 {
2106 return Ok(());
2109 }
2110 if !logged_held_wait {
2114 log::info!(
2115 "[force_merge] waiting: {} segment(s) held by active \
2116 merge/reorder operations, none free to merge",
2117 held
2118 );
2119 logged_held_wait = true;
2120 } else {
2121 log::debug!("[force_merge] still waiting on {} held segment(s)", held);
2122 }
2123 #[cfg(test)]
2124 self.force_merge_conflict_retries
2125 .fetch_add(1, Ordering::Relaxed);
2126 tokio::select! {
2127 biased;
2128 () = self.active_operations.wait_for_shutdown() => {
2129 return Err(Error::IndexClosed);
2130 }
2131 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2132 }
2133 continue;
2134 }
2135 logged_held_wait = false;
2136
2137 let _global_merge_permit = tokio::select! {
2138 biased;
2139 () = self.active_operations.wait_for_shutdown() => {
2140 return Err(Error::IndexClosed);
2141 }
2142 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2143 permit.map_err(|_| {
2144 Error::Internal("global background merge scheduler is closed".into())
2145 })?
2146 }
2147 };
2148
2149 let output_id = SegmentId::new();
2150 let output_hex = output_id.to_hex();
2151
2152 let mut all_ids = batch.clone();
2155 all_ids.push(output_hex);
2156 let guard = {
2157 let st = self.state.lock().await;
2158 batch
2159 .iter()
2160 .all(|id| st.metadata.has_segment(id))
2161 .then(|| self.active_operations.try_register(all_ids))
2162 .flatten()
2163 };
2164 let _guard = match guard {
2165 Some(g) => g,
2166 None if !self.active_operations.is_accepting() => {
2167 return Err(Error::IndexClosed);
2168 }
2169 None => {
2170 #[cfg(test)]
2171 self.force_merge_conflict_retries
2172 .fetch_add(1, Ordering::Relaxed);
2173 drop(_global_merge_permit);
2176 log::debug!("[force_merge] batch lost a registration race, rebuilding");
2182 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2183 self.wait_for_merging_thread().await;
2184 if !had_tracked_merges {
2185 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2186 }
2187 continue;
2188 }
2189 };
2190 log::info!(
2194 "[force_merge] merging batch of {} segments ({} docs)",
2195 batch.len(),
2196 batch_docs
2197 );
2198 let mut output_cleanup = self.output_cleanup_guard(output_id);
2199
2200 let trained_snap = self.trained_for_segment_build();
2201 let granularity = self.merge_granularity(&batch).await;
2202 let merge_result = Self::do_merge(
2203 self.directory.as_ref(),
2204 &self.schema,
2205 &batch,
2206 output_id,
2207 self.term_cache_blocks,
2208 trained_snap.as_deref(),
2209 self.reorder_on_merge,
2210 granularity,
2211 self.merge_bp_time_budget,
2212 self.bp_memory_budget_bytes,
2213 Arc::clone(&self.reorder_permits),
2214 Some(self.background_cpu_pool()),
2215 )
2216 .await;
2217 let (new_segment_id, total_docs, bp_converged) = match merge_result {
2218 Ok(v) => v,
2219 Err(MergeTaskError {
2220 error,
2221 unavailable_segments,
2222 }) => {
2223 for segment_id in &unavailable_segments {
2224 self.quarantine_segment(segment_id, &error);
2225 }
2226 self.delete_output_if_unregistered(output_id, "force-merge failure")
2227 .await;
2228 output_cleanup.disarm();
2229 return Err(error);
2230 }
2231 };
2232
2233 if let Err(e) = self
2234 .replace_segments(
2235 &batch,
2236 new_segment_id,
2237 total_docs,
2238 ReplacementLayout::Recomputed {
2239 reordered: self.reorder_on_merge,
2240 bp_converged,
2241 },
2242 )
2243 .await
2244 {
2245 self.delete_output_if_unregistered(output_id, "replacement failure")
2246 .await;
2247 output_cleanup.disarm();
2248 return Err(e);
2249 }
2250 output_cleanup.disarm();
2251
2252 }
2254 }
2255
2256 fn segment_needs_vector_rewrite(
2257 &self,
2258 reader: &SegmentReader,
2259 field_ids: &[u32],
2260 rewrite_existing: bool,
2261 ) -> Result<bool> {
2262 for &field_id in field_ids {
2263 let flat = reader.flat_vectors().get(&field_id);
2264 let ann = reader.vector_indexes().get(&field_id);
2265 if ann.is_some() && flat.is_none() {
2266 return Err(Error::Corruption(format!(
2267 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2268 reader.meta().id,
2269 )));
2270 }
2271
2272 let Some(flat) = flat else {
2273 continue;
2274 };
2275 if flat.num_vectors == 0 {
2276 continue;
2277 }
2278 if rewrite_existing {
2279 return Ok(true);
2280 }
2281 let field = crate::dsl::Field(field_id);
2282 let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2283 Error::Corruption(format!(
2284 "segment {:032x} references unknown vector field {field_id}",
2285 reader.meta().id,
2286 ))
2287 })?;
2288 let current = match entry.field_type {
2289 crate::dsl::FieldType::DenseVector => {
2290 matches!(ann, Some(crate::segment::VectorIndex::IvfPq(_)))
2291 }
2292 crate::dsl::FieldType::BinaryDenseVector => {
2293 matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2294 }
2295 _ => false,
2296 };
2297 if !current {
2298 return Ok(true);
2299 }
2300 }
2301 Ok(false)
2302 }
2303
2304 async fn acquire_vector_rewrite_capacity(
2305 &self,
2306 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2307 let global = tokio::select! {
2308 biased;
2309 () = self.active_operations.wait_for_shutdown() => {
2310 return Err(Error::IndexClosed);
2311 }
2312 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2313 permit.map_err(|_| Error::Internal(
2314 "global background merge scheduler is closed".into()
2315 ))?
2316 }
2317 };
2318 let local = tokio::select! {
2319 biased;
2320 () = self.active_operations.wait_for_shutdown() => {
2321 return Err(Error::IndexClosed);
2322 }
2323 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2324 permit.map_err(|_| Error::Internal(
2325 "background merge scheduler is closed".into()
2326 ))?
2327 }
2328 };
2329 Ok((global, local))
2330 }
2331
2332 async fn build_vector_replacement(
2333 self: &Arc<Self>,
2334 segment_id: &str,
2335 source_id: SegmentId,
2336 output_id: SegmentId,
2337 trained: &TrainedVectorStructures,
2338 failure_context: &'static str,
2339 ) -> Result<(String, u32, OutputCleanupGuard)> {
2340 let mut cleanup = self.output_cleanup_guard(output_id);
2341 match crate::segment::reorder::rewrite_vector_segment(
2342 self.directory.as_ref(),
2343 &self.schema,
2344 source_id,
2345 output_id,
2346 self.term_cache_blocks,
2347 trained,
2348 Some(self.background_cpu_pool()),
2349 )
2350 .await
2351 {
2352 Ok((new_id, doc_count)) => {
2353 self.validate_completed_segment(&new_id, doc_count).await?;
2354 Ok((new_id, doc_count, cleanup))
2355 }
2356 Err(error) => {
2357 self.delete_output_if_unregistered(output_id, failure_context)
2358 .await;
2359 cleanup.disarm();
2360 if is_deterministic_source_error(&error) {
2361 self.quarantine_segment(segment_id, &error);
2362 }
2363 Err(error)
2364 }
2365 }
2366 }
2367
2368 pub(crate) async fn stage_vector_generation(
2372 self: &Arc<Self>,
2373 _artifact_update: &VectorArtifactUpdateGuard,
2374 segment_ids: &[String],
2375 field_ids: &[u32],
2376 trained: Arc<TrainedVectorStructures>,
2377 rewrite_existing: bool,
2378 ) -> Result<Vec<StagedVectorSegment>> {
2379 if !self.vector_artifact_update.load(Ordering::Acquire) {
2380 return Err(Error::Internal(
2381 "cannot stage a vector generation without an exclusive update lease".into(),
2382 ));
2383 }
2384
2385 let mut staged = Vec::new();
2386 for segment_id in segment_ids {
2387 if self.quarantined_segments.lock().contains(segment_id) {
2388 return Err(Error::Corruption(format!(
2389 "segment {segment_id} is quarantined after a deterministic source failure"
2390 )));
2391 }
2392 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2393 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2394 })?;
2395
2396 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2399
2400 let output_id = SegmentId::new();
2401 let output_hex = output_id.to_hex();
2402 let operation = {
2403 let st = self.state.lock().await;
2404 if !st.metadata.has_segment(segment_id) {
2405 return Err(Error::Corruption(format!(
2406 "vector generation source {segment_id} disappeared while lifecycle work was paused"
2407 )));
2408 }
2409 self.active_operations
2410 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2411 }
2412 .ok_or_else(|| {
2413 if self.active_operations.is_accepting() {
2414 Error::Internal(format!(
2415 "vector generation could not claim stable source {segment_id}"
2416 ))
2417 } else {
2418 Error::IndexClosed
2419 }
2420 })?;
2421
2422 let reader = SegmentReader::open(
2423 self.directory.as_ref(),
2424 source_id,
2425 Arc::clone(&self.schema),
2426 self.term_cache_blocks,
2427 )
2428 .await?;
2429 if !self.segment_needs_vector_rewrite(&reader, field_ids, rewrite_existing)? {
2430 continue;
2431 }
2432 drop(reader);
2433
2434 let (new_id, doc_count, cleanup) = self
2435 .build_vector_replacement(
2436 segment_id,
2437 source_id,
2438 output_id,
2439 trained.as_ref(),
2440 "vector generation staging failure",
2441 )
2442 .await?;
2443 debug_assert_eq!(new_id, output_hex);
2444 let output_reader = SegmentReader::open(
2445 self.directory.as_ref(),
2446 output_id,
2447 Arc::clone(&self.schema),
2448 self.term_cache_blocks,
2449 )
2450 .await?;
2451 if self.segment_needs_vector_rewrite(&output_reader, field_ids, false)? {
2452 return Err(Error::Corruption(format!(
2453 "staged vector segment {new_id} does not match its candidate codebook generation"
2454 )));
2455 }
2456
2457 staged.push(StagedVectorSegment {
2458 source_id: segment_id.clone(),
2459 output_id,
2460 doc_count,
2461 _operation: operation,
2462 cleanup,
2463 });
2464 }
2465 Ok(staged)
2466 }
2467
2468 async fn rewrite_vector_segment_once(
2469 self: &Arc<Self>,
2470 segment_id: &str,
2471 field_ids: &[u32],
2472 ) -> Result<VectorSegmentRewriteOutcome> {
2473 if self.quarantined_segments.lock().contains(segment_id) {
2474 return Err(Error::Corruption(format!(
2475 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
2476 )));
2477 }
2478 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2479 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2480 })?;
2481
2482 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2487
2488 let output_id = SegmentId::new();
2489 let output_hex = output_id.to_hex();
2490 let all_ids = vec![segment_id.to_owned(), output_hex];
2491 let operation = {
2492 let st = self.state.lock().await;
2493 if !st.metadata.has_segment(segment_id) {
2494 return Ok(VectorSegmentRewriteOutcome::SourceGone);
2495 }
2496 self.active_operations.try_register(all_ids)
2497 };
2498 let _operation = match operation {
2499 Some(operation) => operation,
2500 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
2501 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
2502 };
2503
2504 let Some(trained) = self.trained_for_segment_build() else {
2505 return Ok(VectorSegmentRewriteOutcome::Deferred);
2506 };
2507
2508 let reader = SegmentReader::open(
2509 self.directory.as_ref(),
2510 source_id,
2511 Arc::clone(&self.schema),
2512 self.term_cache_blocks,
2513 )
2514 .await?;
2515 if !self.segment_needs_vector_rewrite(&reader, field_ids, false)? {
2516 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
2517 }
2518 drop(reader);
2519
2520 let (new_id, doc_count, mut output_cleanup) = self
2521 .build_vector_replacement(
2522 segment_id,
2523 source_id,
2524 output_id,
2525 trained.as_ref(),
2526 "vector rewrite failure",
2527 )
2528 .await?;
2529
2530 if let Err(error) = self
2531 .replace_segments(
2532 &[segment_id.to_owned()],
2533 new_id,
2534 doc_count,
2535 ReplacementLayout::PreserveSingleSource,
2536 )
2537 .await
2538 {
2539 self.delete_output_if_unregistered(output_id, "vector replacement failure")
2540 .await;
2541 output_cleanup.disarm();
2542 return Err(error);
2543 }
2544 output_cleanup.disarm();
2545 Ok(VectorSegmentRewriteOutcome::Rewritten)
2546 }
2547
2548 pub(crate) async fn rewrite_vector_segments(
2553 self: &Arc<Self>,
2554 field_ids: &[u32],
2555 ) -> Result<usize> {
2556 if field_ids.is_empty() {
2557 return Ok(0);
2558 }
2559 let mut rewritten = 0usize;
2560 loop {
2561 let segment_ids = self.get_segment_ids().await;
2562 let mut conflicted = false;
2563 let mut changed = false;
2564 for segment_id in segment_ids {
2565 match self
2566 .rewrite_vector_segment_once(&segment_id, field_ids)
2567 .await?
2568 {
2569 VectorSegmentRewriteOutcome::Rewritten => {
2570 rewritten += 1;
2571 changed = true;
2572 }
2573 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
2574 VectorSegmentRewriteOutcome::Deferred => {
2575 return Err(Error::Internal(
2576 "ANN finalization lost the published trained generation".into(),
2577 ));
2578 }
2579 VectorSegmentRewriteOutcome::AlreadyCurrent
2580 | VectorSegmentRewriteOutcome::SourceGone => {}
2581 }
2582 }
2583 if !conflicted && !changed {
2584 log::info!(
2585 "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
2586 rewritten,
2587 );
2588 return Ok(rewritten);
2589 }
2590 tokio::select! {
2591 biased;
2592 () = self.active_operations.wait_for_shutdown() => {
2593 return Err(Error::IndexClosed);
2594 }
2595 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
2596 }
2597 }
2598 }
2599
2600 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
2605 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
2606 return;
2607 }
2608 let manager = Arc::clone(self);
2609 let future = async move {
2610 let field_ids = manager
2611 .read_metadata(|metadata| {
2612 metadata
2613 .vector_fields
2614 .keys()
2615 .filter(|field_id| metadata.is_field_built(**field_id))
2616 .copied()
2617 .collect::<Vec<_>>()
2618 })
2619 .await;
2620 for segment_id in segment_ids {
2621 loop {
2622 match manager
2623 .rewrite_vector_segment_once(&segment_id, &field_ids)
2624 .await
2625 {
2626 Ok(VectorSegmentRewriteOutcome::Conflict) => {
2627 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
2628 }
2629 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
2630 Ok(_) => break,
2631 Err(error) => {
2632 log::error!(
2633 "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
2634 segment_id,
2635 error,
2636 );
2637 break;
2638 }
2639 }
2640 }
2641 }
2642 };
2643 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
2644 log::warn!(
2645 "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
2646 );
2647 return;
2648 };
2649 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2650 log::warn!(
2651 "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
2652 );
2653 }
2654 }
2655
2656 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
2663 self.wait_for_all_merges().await;
2664 let segment_ids = self.get_segment_ids().await;
2665
2666 if segment_ids.is_empty() {
2667 log::info!("[reorder] no segments to reorder");
2668 return Ok(());
2669 }
2670
2671 log::info!("[reorder] reordering {} segments", segment_ids.len());
2672
2673 for seg_id in segment_ids {
2674 match self
2675 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
2676 .await
2677 {
2678 Ok(true) => {}
2679 Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
2680 Err(e) => return Err(e),
2681 }
2682 }
2683
2684 log::info!("[reorder] all segments reordered");
2685 Ok(())
2686 }
2687
2688 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
2693 self.unreordered_segments()
2694 .await
2695 .into_iter()
2696 .map(|(id, _)| id)
2697 .collect()
2698 }
2699
2700 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
2703 let quarantined = self.quarantined_segments.lock().clone();
2704 let paused = self.paused_reorder_segments();
2705 let st = self.state.lock().await;
2706 let active_ids = self.active_operations.snapshot();
2707 st.metadata
2708 .segment_metas
2709 .iter()
2710 .filter(|(id, info)| {
2711 !info.reordered
2712 && !active_ids.contains(*id)
2713 && !quarantined.contains(*id)
2714 && !paused.contains(*id)
2715 })
2716 .map(|(id, info)| (id.clone(), info.num_docs))
2717 .collect()
2718 }
2719
2720 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
2724 self.unconverged_segments_below(u32::MAX)
2725 .await
2726 .into_iter()
2727 .map(|(id, docs, _)| (id, docs))
2728 .collect()
2729 }
2730
2731 pub async fn unconverged_segments_below(
2734 &self,
2735 max_unconverged_passes: u32,
2736 ) -> Vec<(String, u32, u32)> {
2737 let quarantined = self.quarantined_segments.lock().clone();
2738 let paused = self.paused_reorder_segments();
2739 let st = self.state.lock().await;
2740 let active_ids = self.active_operations.snapshot();
2741 st.metadata
2742 .segment_metas
2743 .iter()
2744 .filter(|(id, info)| {
2745 info.reordered
2746 && !info.bp_converged
2747 && info.bp_unconverged_passes < max_unconverged_passes
2748 && !active_ids.contains(*id)
2749 && !quarantined.contains(*id)
2750 && !paused.contains(*id)
2751 })
2752 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
2753 .collect()
2754 }
2755
2756 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
2766 let st = self.state.lock().await;
2767 let deepening = ids.iter().any(|id| {
2768 st.metadata
2769 .segment_metas
2770 .get(id)
2771 .is_some_and(|info| info.reordered && !info.bp_converged)
2772 });
2773 drop(st);
2774 if deepening {
2775 log::info!(
2776 "[reorder] source segment(s) unconverged — forcing record-level BP (deepening pass)",
2777 );
2778 crate::segment::reorder::BpGranularity::Records
2779 } else {
2780 crate::segment::reorder::BpGranularity::Auto
2781 }
2782 }
2783
2784 pub async fn reorder_single_segment(
2789 self: &Arc<Self>,
2790 seg_id: &str,
2791 rayon_pool: Option<Arc<rayon::ThreadPool>>,
2792 bp_budget: crate::segment::BpBudget,
2793 ) -> Result<bool> {
2794 let source_id = SegmentId::from_hex(seg_id)
2795 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
2796 if self.quarantined_segments.lock().contains(seg_id) {
2797 return Err(Error::Corruption(format!(
2798 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
2799 seg_id
2800 )));
2801 }
2802
2803 let _reorder_permit = tokio::select! {
2808 biased;
2809 () = self.active_operations.wait_for_shutdown() => {
2810 return Err(Error::IndexClosed);
2811 }
2812 permit = Arc::clone(&self.reorder_permits).acquire_owned() => {
2813 permit.map_err(|_| {
2814 Error::Internal("background reorder scheduler is closed".into())
2815 })?
2816 }
2817 };
2818
2819 let output_id = SegmentId::new();
2820 let output_hex = output_id.to_hex();
2821 let source_ids = [seg_id.to_string()];
2822 let granularity = self.merge_granularity(&source_ids).await;
2823
2824 let all_ids = vec![seg_id.to_string(), output_hex];
2830 let (_guard, source_docs) = {
2831 let st = self.state.lock().await;
2832 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
2833 log::info!(
2834 "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
2835 seg_id
2836 );
2837 self.clear_reorder_retry(seg_id);
2838 return Ok(false);
2839 };
2840
2841 match self.active_operations.try_register(all_ids) {
2842 Some(guard) => (guard, source_meta.num_docs),
2843 None if !self.active_operations.is_accepting() => {
2844 return Err(Error::IndexClosed);
2845 }
2846 None => {
2847 log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
2848 return Ok(false);
2849 }
2850 }
2851 };
2852
2853 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
2858 if is_deterministic_source_error(&error) {
2859 self.quarantine_segment(seg_id, &error);
2860 } else if !matches!(&error, Error::IndexClosed) {
2861 self.pause_reorder_retries(seg_id, &error);
2862 }
2863 return Err(error);
2864 }
2865
2866 let mut output_cleanup = self.output_cleanup_guard(output_id);
2867
2868 let reorder_result = crate::segment::reorder::reorder_segment(
2869 self.directory.as_ref(),
2870 &self.schema,
2871 source_id,
2872 output_id,
2873 self.term_cache_blocks,
2874 self.bp_memory_budget_bytes,
2875 bp_budget,
2876 granularity,
2877 rayon_pool,
2878 )
2879 .await;
2880 let (new_id, total_docs, bp_converged) = match reorder_result {
2881 Ok(v) => v,
2882 Err(e) => {
2883 self.delete_output_if_unregistered(output_id, "reorder failure")
2886 .await;
2887 output_cleanup.disarm();
2888 if is_deterministic_source_error(&e) {
2889 self.quarantine_segment(seg_id, &e);
2890 } else if !matches!(&e, Error::IndexClosed) {
2891 self.pause_reorder_retries(seg_id, &e);
2892 }
2893 return Err(e);
2894 }
2895 };
2896
2897 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
2903 if let Err(e) = self
2904 .replace_segments(
2905 &[seg_id.to_string()],
2906 new_id,
2907 total_docs,
2908 ReplacementLayout::Recomputed {
2909 reordered: true,
2910 bp_converged: ladder_converged,
2911 },
2912 )
2913 .await
2914 {
2915 self.delete_output_if_unregistered(output_id, "replacement failure")
2916 .await;
2917 output_cleanup.disarm();
2918 if !matches!(&e, Error::IndexClosed) {
2919 self.pause_reorder_retries(seg_id, &e);
2920 }
2921 return Err(e);
2922 }
2923 output_cleanup.disarm();
2924 self.clear_reorder_retry(seg_id);
2925
2926 Ok(true)
2927 }
2928
2929 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
2936 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
2937
2938 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
2939 for entry in entries {
2940 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
2941 continue;
2942 };
2943 let Some(rest) = filename.strip_prefix("seg_") else {
2944 continue;
2945 };
2946 let Some(hex_id) = rest.get(..32) else {
2947 continue;
2948 };
2949 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
2950 continue;
2951 }
2952 orphan_files
2953 .entry(hex_id.to_ascii_lowercase())
2954 .or_default()
2955 .push(entry);
2956 }
2957 }
2958
2959 let mut deleted = 0;
2960 for (hex_id, paths) in &orphan_files {
2961 let deletion_guard = {
2966 let st = self.state.lock().await;
2967 if st.metadata.has_segment(hex_id) {
2968 continue;
2969 }
2970 let Some(guard) = self
2971 .active_operations
2972 .try_register(vec![hex_id.to_string()])
2973 else {
2974 continue;
2975 };
2976 if self.tracker.is_deletion_protected(hex_id) {
2977 drop(guard);
2978 continue;
2979 }
2980 guard
2981 };
2982
2983 let results =
2988 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
2989 .await;
2990 let removed = results.into_iter().all(|result| match result {
2991 Ok(()) => true,
2992 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
2993 Err(error) => {
2994 log::warn!(
2995 "[segment_cleanup] failed sweeping orphan segment {}: {}",
2996 hex_id,
2997 error,
2998 );
2999 false
3000 }
3001 });
3002 drop(deletion_guard);
3005 if removed {
3006 deleted += 1;
3007 log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3008 }
3009 }
3010
3011 Ok(deleted)
3012 }
3013}
3014
3015#[cfg(test)]
3016mod tests {
3017 use super::*;
3018 use std::sync::atomic::{AtomicBool, Ordering};
3019
3020 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3021 let schema = crate::dsl::SchemaBuilder::default().build();
3022 let metadata = IndexMetadata::new(schema.clone());
3023 Arc::new(SegmentManager::new(
3024 Arc::new(crate::directories::RamDirectory::new()),
3025 Arc::new(schema),
3026 metadata,
3027 Box::new(crate::merge::NoMergePolicy),
3028 0,
3029 1,
3030 Arc::new(Semaphore::new(1)),
3031 None,
3032 1024,
3033 Arc::new(Semaphore::new(1)),
3034 None,
3035 ))
3036 }
3037
3038 #[test]
3039 fn output_cleanup_guard_runs_during_panic_unwind() {
3040 let cleaned = Arc::new(AtomicBool::new(false));
3041 let cleaned_in_callback = Arc::clone(&cleaned);
3042 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3043 cleaned_in_callback.store(true, Ordering::SeqCst);
3044 });
3045
3046 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3047 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3048 panic!("simulated reorder panic");
3049 }));
3050
3051 assert!(result.is_err());
3052 assert!(
3053 cleaned.load(Ordering::SeqCst),
3054 "partial output cleanup must run during unwind"
3055 );
3056 }
3057
3058 #[test]
3059 fn output_cleanup_guard_disarms_after_commit() {
3060 let cleaned = Arc::new(AtomicBool::new(false));
3061 let cleaned_in_callback = Arc::clone(&cleaned);
3062 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3063 cleaned_in_callback.store(true, Ordering::SeqCst);
3064 });
3065
3066 {
3067 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3068 guard.disarm();
3069 }
3070
3071 assert!(!cleaned.load(Ordering::SeqCst));
3072 }
3073
3074 #[test]
3075 fn test_active_operation_guard_releases_ownership() {
3076 let active = Arc::new(ActiveSegmentOperations::new());
3077 {
3078 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3079 let snap = active.snapshot();
3080 assert!(snap.contains("a"));
3081 assert!(snap.contains("b"));
3082 }
3083 assert!(active.snapshot().is_empty());
3084 }
3085
3086 #[test]
3087 fn test_non_overlapping_operations_can_run_concurrently() {
3088 let active = Arc::new(ActiveSegmentOperations::new());
3089 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3090 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3091 let snap = active.snapshot();
3092 assert_eq!(snap.len(), 4);
3093
3094 drop(first);
3095 let snap = active.snapshot();
3096 assert_eq!(snap.len(), 2);
3097 assert!(snap.contains("c"));
3098 assert!(snap.contains("d"));
3099 }
3100
3101 #[test]
3102 fn test_overlapping_operation_is_rejected_until_release() {
3103 let active = Arc::new(ActiveSegmentOperations::new());
3104 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3105 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3106 drop(first);
3107 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3108 }
3109
3110 #[test]
3111 fn test_active_operation_snapshot() {
3112 let active = Arc::new(ActiveSegmentOperations::new());
3113 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3114 let snap = active.snapshot();
3115 assert!(snap.contains("x"));
3116 assert!(snap.contains("y"));
3117 assert!(!snap.contains("z"));
3118 }
3119
3120 #[tokio::test]
3121 async fn operation_barrier_ignores_producers_started_after_snapshot() {
3122 let active = Arc::new(ActiveSegmentOperations::new());
3123 let before_gate = active.try_register(vec!["old".into()]).unwrap();
3124 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3125 assert_eq!(parked_indexing, 0);
3126 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3127
3128 let waiter = {
3129 let active = Arc::clone(&active);
3130 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3131 };
3132 tokio::task::yield_now().await;
3133 assert!(!waiter.is_finished());
3134
3135 drop(before_gate);
3136 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3137 .await
3138 .expect("pre-gate operation barrier was starved by a post-gate producer")
3139 .unwrap();
3140 assert!(active.snapshot().contains("new-flat"));
3141 drop(after_gate);
3142 }
3143
3144 #[tokio::test]
3145 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
3146 let manager = lifecycle_test_manager();
3147 manager
3148 .trained
3149 .store(Some(Arc::new(TrainedVectorStructures {
3150 centroids: rustc_hash::FxHashMap::default(),
3151 binary_quantizers: rustc_hash::FxHashMap::default(),
3152 codebooks: rustc_hash::FxHashMap::default(),
3153 ..Default::default()
3154 })));
3155
3156 let guard = manager.begin_vector_artifact_update().await.unwrap();
3157 assert!(
3158 manager.trained().is_some(),
3159 "search readers keep the last fully validated generation"
3160 );
3161 assert!(
3162 manager.trained_for_segment_build().is_none(),
3163 "new segment producers must stay flat during an artifact update"
3164 );
3165
3166 let detached_transaction_guard = guard.clone();
3167 drop(guard);
3168 assert!(
3169 manager.trained_for_segment_build().is_none(),
3170 "a detached lifecycle transaction must retain the producer gate after request cancellation"
3171 );
3172 drop(detached_transaction_guard);
3173 assert!(manager.trained_for_segment_build().is_some());
3174 }
3175
3176 #[tokio::test]
3177 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
3178 let manager = lifecycle_test_manager();
3179 let guard = manager.begin_vector_artifact_update().await.unwrap();
3180 assert!(
3181 manager
3182 .active_operations
3183 .try_register(vec!["merge".into()])
3184 .is_none(),
3185 "ordinary merge/reorder work must not change staged sources"
3186 );
3187 let indexing = manager
3188 .active_operations
3189 .try_register_indexing(vec!["fresh".into()])
3190 .expect("indexing remains available in flat mode");
3191 drop(indexing);
3192
3193 drop(guard);
3194 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
3195 assert!(
3196 manager
3197 .active_operations
3198 .try_register(vec!["merge".into()])
3199 .is_some()
3200 );
3201 }
3202
3203 #[tokio::test]
3204 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
3205 let active = Arc::new(ActiveSegmentOperations::new());
3206 let guard = active.try_register(vec!["live".into()]).unwrap();
3207 active.stop_accepting();
3208 assert!(active.try_register(vec!["new".into()]).is_none());
3209
3210 let waiter = {
3211 let active = Arc::clone(&active);
3212 tokio::spawn(async move { active.wait_until_idle().await })
3213 };
3214 tokio::task::yield_now().await;
3215 assert!(!waiter.is_finished());
3216 drop(guard);
3217 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3218 .await
3219 .expect("shutdown waiter missed the final guard notification")
3220 .unwrap();
3221 }
3222
3223 #[tokio::test]
3224 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
3225 let manager = lifecycle_test_manager();
3226 let started = Arc::new(Semaphore::new(0));
3227 let release = Arc::new(Semaphore::new(0));
3228 let completed = Arc::new(AtomicBool::new(false));
3229
3230 let request = {
3231 let manager = Arc::clone(&manager);
3232 let started = Arc::clone(&started);
3233 let release = Arc::clone(&release);
3234 let completed = Arc::clone(&completed);
3235 tokio::spawn(async move {
3236 manager
3237 .run_lifecycle_transaction(async move {
3238 started.add_permits(1);
3239 let _permit = release.acquire().await.unwrap();
3240 completed.store(true, Ordering::Release);
3241 Ok(())
3242 })
3243 .await
3244 })
3245 };
3246
3247 let _started = started.acquire().await.unwrap();
3248 request.abort();
3249 assert!(request.await.unwrap_err().is_cancelled());
3250 release.add_permits(1);
3251
3252 manager.begin_shutdown();
3253 tokio::time::timeout(
3254 std::time::Duration::from_secs(1),
3255 manager.wait_for_shutdown(),
3256 )
3257 .await
3258 .expect("shutdown did not drain detached lifecycle transaction");
3259 assert!(completed.load(Ordering::Acquire));
3260 }
3261
3262 #[tokio::test]
3263 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
3264 let manager = lifecycle_test_manager();
3265 {
3266 let mut state = manager.state.lock().await;
3267 state.metadata.add_segment_meta(
3268 "eligible".into(),
3269 SegmentMetaInfo {
3270 num_docs: 10,
3271 ancestors: Vec::new(),
3272 generation: 1,
3273 reordered: true,
3274 bp_converged: false,
3275 bp_unconverged_passes: 2,
3276 },
3277 );
3278 state.metadata.add_segment_meta(
3279 "at-limit".into(),
3280 SegmentMetaInfo {
3281 num_docs: 20,
3282 ancestors: Vec::new(),
3283 generation: 1,
3284 reordered: true,
3285 bp_converged: false,
3286 bp_unconverged_passes: 3,
3287 },
3288 );
3289 state.metadata.add_segment_meta(
3290 "converged".into(),
3291 SegmentMetaInfo {
3292 num_docs: 30,
3293 ancestors: Vec::new(),
3294 generation: 1,
3295 reordered: true,
3296 bp_converged: true,
3297 bp_unconverged_passes: 0,
3298 },
3299 );
3300 state.metadata.add_segment("fresh".into(), 40);
3301 }
3302
3303 assert_eq!(
3304 manager.unconverged_segments_below(3).await,
3305 vec![("eligible".into(), 10, 2)]
3306 );
3307 assert!(manager.unconverged_segments_below(0).await.is_empty());
3308 }
3309
3310 #[test]
3311 fn merge_retry_backoff_is_exponential_and_capped() {
3312 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
3313 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
3314 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
3315 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
3316 }
3317
3318 #[test]
3319 fn only_deterministic_source_errors_are_quarantined() {
3320 assert!(is_deterministic_source_error(&Error::Corruption(
3321 "bad footer".into()
3322 )));
3323 assert!(is_deterministic_source_error(&Error::Io(
3324 std::io::Error::from(std::io::ErrorKind::NotFound)
3325 )));
3326 assert!(!is_deterministic_source_error(&Error::Io(
3327 std::io::Error::from(std::io::ErrorKind::TimedOut)
3328 )));
3329 assert!(!is_deterministic_source_error(&Error::Io(
3330 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
3331 )));
3332 }
3333
3334 #[test]
3335 fn transient_reorder_failure_is_backed_off_until_cleared() {
3336 let manager = lifecycle_test_manager();
3337 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
3338 assert!(manager.paused_reorder_segments().contains("source"));
3339 manager.clear_reorder_retry("source");
3340 assert!(!manager.paused_reorder_segments().contains("source"));
3341 }
3342
3343 #[derive(Default)]
3346 struct FailingExistsDirectory(crate::directories::RamDirectory);
3347
3348 #[async_trait::async_trait]
3349 impl crate::directories::Directory for FailingExistsDirectory {
3350 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
3351 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
3352 }
3353
3354 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
3355 self.0.file_size(path).await
3356 }
3357
3358 async fn open_read(
3359 &self,
3360 path: &std::path::Path,
3361 ) -> std::io::Result<crate::directories::FileHandle> {
3362 self.0.open_read(path).await
3363 }
3364
3365 async fn read_range(
3366 &self,
3367 path: &std::path::Path,
3368 range: std::ops::Range<u64>,
3369 ) -> std::io::Result<crate::directories::OwnedBytes> {
3370 self.0.read_range(path, range).await
3371 }
3372
3373 async fn list_files(
3374 &self,
3375 prefix: &std::path::Path,
3376 ) -> std::io::Result<Vec<std::path::PathBuf>> {
3377 self.0.list_files(prefix).await
3378 }
3379
3380 async fn open_lazy(
3381 &self,
3382 path: &std::path::Path,
3383 ) -> std::io::Result<crate::directories::FileHandle> {
3384 self.0.open_lazy(path).await
3385 }
3386 }
3387
3388 #[async_trait::async_trait]
3389 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
3390 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
3391 self.0.write(path, data).await
3392 }
3393
3394 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
3395 self.0.delete(path).await
3396 }
3397
3398 async fn rename(
3399 &self,
3400 from: &std::path::Path,
3401 to: &std::path::Path,
3402 ) -> std::io::Result<()> {
3403 self.0.rename(from, to).await
3404 }
3405
3406 async fn sync(&self) -> std::io::Result<()> {
3407 self.0.sync().await
3408 }
3409
3410 async fn streaming_writer(
3411 &self,
3412 path: &std::path::Path,
3413 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
3414 self.0.streaming_writer(path).await
3415 }
3416 }
3417
3418 #[derive(Debug, Clone)]
3419 struct MergeEverythingPolicy;
3420
3421 impl MergePolicy for MergeEverythingPolicy {
3422 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
3423 if segments.len() < 2 {
3424 return Vec::new();
3425 }
3426 vec![crate::merge::MergeCandidate {
3427 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
3428 }]
3429 }
3430
3431 fn clone_box(&self) -> Box<dyn MergePolicy> {
3432 Box::new(self.clone())
3433 }
3434 }
3435
3436 #[tokio::test]
3437 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
3438 let manager = lifecycle_test_manager();
3439 let parked_indexing = manager
3444 .protect_new_segment("00000000000000000000000000000abc".into())
3445 .unwrap();
3446
3447 let error = tokio::time::timeout(
3448 std::time::Duration::from_secs(2),
3449 manager.begin_vector_artifact_update(),
3450 )
3451 .await
3452 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
3453 .err()
3454 .expect("an old-generation prepared segment must block artifact replacement")
3455 .to_string();
3456 assert!(error.contains("built but uncommitted"), "{error}");
3457 assert!(
3458 !manager.vector_artifact_update.load(Ordering::Acquire),
3459 "a rejected update must release the producer gate"
3460 );
3461
3462 drop(parked_indexing);
3463
3464 let guard = manager
3465 .begin_vector_artifact_update()
3466 .await
3467 .expect("artifact update should succeed after the pending generation is resolved");
3468 drop(guard);
3469 }
3470
3471 #[tokio::test]
3472 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
3473 let manager = lifecycle_test_manager();
3474 let merge_like = manager
3475 .active_operations
3476 .try_register(vec!["merge-source".into()])
3477 .unwrap();
3478
3479 let waiter = {
3480 let manager = Arc::clone(&manager);
3481 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
3482 };
3483 for _ in 0..8 {
3484 tokio::task::yield_now().await;
3485 }
3486 assert!(
3487 !waiter.is_finished(),
3488 "artifact update must drain merge/reorder producers that may hold the previous generation"
3489 );
3490
3491 drop(merge_like);
3492 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3493 .await
3494 .expect("artifact update missed the lifecycle guard release")
3495 .unwrap()
3496 .unwrap();
3497 }
3498
3499 #[tokio::test]
3500 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
3501 let manager = lifecycle_test_manager();
3502 let release = Arc::new(Semaphore::new(0));
3503 let merge_task = {
3504 let release = Arc::clone(&release);
3505 tokio::spawn(async move {
3506 let _permit = release.acquire().await.unwrap();
3507 })
3508 };
3509 manager.merge_handles.lock().push(merge_task);
3510
3511 let waiter = {
3512 let manager = Arc::clone(&manager);
3513 tokio::spawn(async move { manager.wait_for_all_merges().await })
3514 };
3515 for _ in 0..8 {
3516 tokio::task::yield_now().await;
3517 }
3518 assert!(!waiter.is_finished());
3519 waiter.abort();
3522 let join_error = waiter.await.unwrap_err();
3523 assert!(join_error.is_cancelled());
3524
3525 assert!(
3526 !manager.merge_handles.lock().is_empty(),
3527 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
3528 );
3529
3530 release.add_permits(1);
3532 tokio::time::timeout(
3533 std::time::Duration::from_secs(1),
3534 manager.wait_for_all_merges(),
3535 )
3536 .await
3537 .expect("subsequent drain missed the reinserted merge handle");
3538 assert!(manager.merge_handles.lock().is_empty());
3539 }
3540
3541 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3542 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
3543 let manager = lifecycle_test_manager();
3544 {
3545 let mut state = manager.state.lock().await;
3546 state
3547 .metadata
3548 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
3549 state
3550 .metadata
3551 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
3552 }
3553 let reorder_like = manager
3556 .active_operations
3557 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
3558 .unwrap();
3559
3560 let force_merge = {
3561 let manager = Arc::clone(&manager);
3562 tokio::spawn(async move { manager.force_merge().await })
3563 };
3564
3565 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
3566 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
3567 assert!(
3568 retries >= 1,
3569 "force_merge never observed the conflicting owner (retries={retries})"
3570 );
3571 assert!(
3572 retries < 20,
3573 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
3574 );
3575
3576 drop(reorder_like);
3577 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
3580 .await
3581 .expect("force_merge kept spinning after the conflicting owner released")
3582 .unwrap();
3583 assert!(result.is_err());
3584 }
3585
3586 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3587 async fn force_merge_routes_around_segments_held_by_reorder() {
3588 let manager = lifecycle_test_manager();
3589 {
3590 let mut state = manager.state.lock().await;
3591 state
3592 .metadata
3593 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
3594 state
3595 .metadata
3596 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
3597 state
3598 .metadata
3599 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
3600 }
3601 let _reorder_like = manager
3604 .active_operations
3605 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
3606 .unwrap();
3607
3608 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
3615 let manager = Arc::clone(&manager);
3616 async move { manager.force_merge().await }
3617 })
3618 .await
3619 .expect("force_merge livelocked on a segment held by an active reorder");
3620 assert!(result.is_err(), "fake segment files must fail the merge");
3621
3622 assert_eq!(
3623 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
3624 0,
3625 "batch built from the ownership snapshot must not collide with the held segment"
3626 );
3627 }
3628
3629 #[tokio::test]
3630 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
3631 let schema = crate::dsl::SchemaBuilder::default().build();
3632 let mut metadata = IndexMetadata::new(schema.clone());
3633 metadata.add_segment("00000000000000000000000000000001".into(), 10);
3634 metadata.add_segment("00000000000000000000000000000002".into(), 10);
3635 let manager = Arc::new(SegmentManager::new(
3636 Arc::new(FailingExistsDirectory::default()),
3637 Arc::new(schema),
3638 metadata,
3639 Box::new(MergeEverythingPolicy),
3640 0,
3641 1,
3642 Arc::new(Semaphore::new(1)),
3643 None,
3644 1024,
3645 Arc::new(Semaphore::new(1)),
3646 None,
3647 ));
3648
3649 manager.maybe_merge().await;
3652
3653 tokio::time::timeout(
3654 std::time::Duration::from_secs(5),
3655 manager.wait_for_all_merges(),
3656 )
3657 .await
3658 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
3659 assert!(
3660 manager.merge_retry_is_paused(),
3661 "the failed merge should have armed the retry backoff"
3662 );
3663
3664 manager.begin_shutdown();
3666 tokio::time::timeout(
3667 std::time::Duration::from_secs(5),
3668 manager.wait_for_shutdown(),
3669 )
3670 .await
3671 .expect("shutdown did not drain the merge retry wakeup task");
3672 }
3673}