1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, AtomicUsize, 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, ReorderConcurrencyGate, ReorderPriority, 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 force_merge_active: AtomicUsize,
566
567 #[cfg(test)]
570 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
571
572 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
576
577 trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
581
582 vector_artifact_update: Arc<AtomicBool>,
586
587 tracker: Arc<SegmentTracker>,
589
590 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
592
593 directory: Arc<D>,
595 schema: Arc<crate::dsl::Schema>,
597 term_cache_blocks: usize,
599 merge_permits: Arc<Semaphore>,
603 global_merge_permits: Arc<Semaphore>,
605 reorder_permits: Arc<ReorderConcurrencyGate>,
609 reorder_on_merge: bool,
614 merge_bp_time_budget: Option<std::time::Duration>,
618 bp_memory_budget_bytes: usize,
621 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
624}
625
626struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
627
628impl Drop for ForceMergeActivityGuard<'_> {
629 fn drop(&mut self) {
630 self.0.fetch_sub(1, Ordering::AcqRel);
631 }
632}
633
634impl<D: DirectoryWriter + 'static> SegmentManager<D> {
635 #[allow(clippy::too_many_arguments)]
637 pub fn new(
638 directory: Arc<D>,
639 schema: Arc<crate::dsl::Schema>,
640 metadata: IndexMetadata,
641 merge_policy: Box<dyn MergePolicy>,
642 term_cache_blocks: usize,
643 max_concurrent_merges: usize,
644 global_merge_permits: Arc<Semaphore>,
645 merge_bp_time_budget: Option<std::time::Duration>,
646 bp_memory_budget_bytes: usize,
647 reorder_permits: Arc<ReorderConcurrencyGate>,
648 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
649 ) -> Self {
650 let reorder_on_merge = schema.reorder_on_merge();
653 if reorder_on_merge {
654 log::info!("[merge] reorder-on-merge enabled by index schema");
655 }
656
657 let tracker = Arc::new(SegmentTracker::new());
658 for seg_id in metadata.segment_metas.keys() {
659 tracker.register(seg_id);
660 }
661
662 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
663 Arc::new(parking_lot::Mutex::new(Vec::new()));
664 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
665 let dir = Arc::clone(&directory);
666 let tracker = Arc::clone(&tracker);
667 let lifecycle_handles = Arc::clone(&lifecycle_handles);
668 Arc::new(move |segment_ids| {
669 let Ok(handle) = tokio::runtime::Handle::try_current() else {
672 tracker.complete_deletion(&segment_ids);
675 return;
676 };
677 let dir = Arc::clone(&dir);
678 let task_tracker = Arc::clone(&tracker);
679 let cleanup_ids = segment_ids.clone();
680 let future = async move {
681 for &segment_id in &segment_ids {
682 log::info!(
683 "[segment_cleanup] deleting deferred segment {}",
684 segment_id.to_hex()
685 );
686 if let Err(error) =
687 crate::segment::delete_segment(dir.as_ref(), segment_id).await
688 {
689 log::warn!(
690 "[segment_cleanup] deferred delete failed for {}: {}",
691 segment_id.to_hex(),
692 error,
693 );
694 }
695 }
696 task_tracker.complete_deletion(&segment_ids);
697 };
698 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
699 tracker.complete_deletion(&cleanup_ids);
703 log::warn!(
704 "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
705 );
706 }
707 })
708 };
709
710 Self {
711 state: Arc::new(AsyncMutex::new(ManagerState {
712 metadata,
713 merge_policy,
714 })),
715 active_operations: Arc::new(ActiveSegmentOperations::new()),
716 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
717 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
718 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
719 merge_handles: parking_lot::Mutex::new(Vec::new()),
720 global_merge_wakeup_pending: AtomicBool::new(false),
721 force_merge_active: AtomicUsize::new(0),
722 #[cfg(test)]
723 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
724 lifecycle_handles,
725 trained: Arc::new(ArcSwapOption::new(None)),
726 vector_artifact_update: Arc::new(AtomicBool::new(false)),
727 tracker,
728 delete_fn,
729 directory,
730 schema,
731 term_cache_blocks,
732 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
733 global_merge_permits,
734 reorder_permits,
735 reorder_on_merge,
736 merge_bp_time_budget,
737 bp_memory_budget_bytes,
738 background_reorder_pool,
739 }
740 }
741
742 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
748 if let Some(pool) = &self.background_reorder_pool {
749 return Arc::clone(pool);
750 }
751 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
752 let threads = (num_cpus::get() / 2).max(1);
753 log::info!(
754 "[merge] process-wide background CPU pool: {} thread(s)",
755 threads
756 );
757 Arc::new(
758 rayon::ThreadPoolBuilder::new()
759 .num_threads(threads)
760 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
761 .build()
762 .expect("failed to build background CPU pool"),
763 )
764 }))
765 }
766
767 pub fn begin_shutdown(&self) {
771 self.active_operations.stop_accepting();
772 }
773
774 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
782 where
783 T: Send + 'static,
784 F: std::future::Future<Output = Result<T>> + Send + 'static,
785 {
786 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
787 let future = async move {
788 let result = transaction.await;
789 let _ = result_tx.send(result);
790 };
791 let runtime = tokio::runtime::Handle::current();
792 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
793 return Err(Error::Internal(
794 "runtime rejected lifecycle metadata transaction".into(),
795 ));
796 }
797 result_rx.await.map_err(|_| {
798 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
799 })?
800 }
801
802 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
804 let manager = Arc::clone(self);
805 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
806 let Ok(handle) = tokio::runtime::Handle::try_current() else {
807 log::warn!(
808 "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
809 segment_id.to_hex(),
810 );
811 return;
812 };
813
814 let cleanup_manager = Arc::clone(&manager);
815 let future = async move {
816 cleanup_manager
817 .delete_output_if_unregistered(segment_id, "task unwind")
818 .await;
819 };
820 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
821 log::warn!(
822 "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
823 segment_id.to_hex(),
824 );
825 }
826 });
827
828 OutputCleanupGuard::new(output_id, cleanup)
829 }
830
831 pub(crate) fn schedule_unpublished_segment_cleanup(
836 self: &Arc<Self>,
837 output_id: SegmentId,
838 operation: SegmentOperationGuard,
839 runtime: tokio::runtime::Handle,
840 ) {
841 let manager = Arc::clone(self);
842 let output_hex = output_id.to_hex();
843 let future = async move {
844 manager
845 .delete_output_if_unregistered(output_id, "indexing abort or failure")
846 .await;
847 drop(operation);
848 };
849 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
850 log::warn!(
853 "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
854 output_hex,
855 );
856 }
857 }
858
859 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
865 match self
866 .active_operations
867 .try_register_indexing(vec![segment_id.clone()])
868 {
869 Some(operation) => Ok(operation),
870 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
871 None => Err(Error::Corruption(format!(
872 "new segment ID {} is already owned by an active operation",
873 segment_id
874 ))),
875 }
876 }
877
878 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
882 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
883 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
884 })?;
885 let files = SegmentFiles::new(id.0);
886
887 for path in files.mandatory_paths() {
888 if !self.directory.exists(path).await.map_err(Error::Io)? {
889 return Err(Error::Corruption(format!(
890 "segment {} cannot be published: mandatory file {:?} is missing",
891 segment_id, path
892 )));
893 }
894 }
895
896 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
897 Error::Corruption(format!(
898 "segment {} cannot be published: missing/unreadable {:?}: {}",
899 segment_id, files.meta, e
900 ))
901 })?;
902 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
903 Error::Corruption(format!(
904 "segment {} cannot be published: failed reading {:?}: {}",
905 segment_id, files.meta, e
906 ))
907 })?;
908 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
909 Error::Corruption(format!(
910 "segment {} cannot be published: invalid {:?}: {}",
911 segment_id, files.meta, e
912 ))
913 })?;
914
915 if meta.id != id.0 || meta.num_docs != expected_docs {
916 return Err(Error::Corruption(format!(
917 "segment {} cannot be published: metadata identity/docs mismatch \
918 (id={:032x}, docs={}, expected_docs={})",
919 segment_id, meta.id, meta.num_docs, expected_docs
920 )));
921 }
922
923 Ok(())
924 }
925
926 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
927 let inserted = self
928 .quarantined_segments
929 .lock()
930 .insert(segment_id.to_string());
931 if inserted {
932 log::error!(
933 "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
934 It remains metadata-live for explicit repair but is excluded from merges until restart",
935 segment_id,
936 error,
937 );
938 }
939 }
940
941 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
942 let mut retry = self.merge_retry.lock();
943 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
944 let delay = merge_retry_delay(retry.consecutive_failures);
945 retry.retry_after = std::time::Instant::now().checked_add(delay);
946 log::warn!(
947 "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
948 delay.as_secs_f64(),
949 retry.consecutive_failures,
950 error,
951 );
952 delay
953 }
954
955 fn clear_merge_retry_backoff(&self) {
956 *self.merge_retry.lock() = MergeRetryState::default();
957 }
958
959 fn merge_retry_is_paused(&self) -> bool {
960 let mut retry = self.merge_retry.lock();
961 match retry.retry_after {
962 Some(deadline) if deadline > std::time::Instant::now() => true,
963 Some(_) => {
964 retry.retry_after = None;
965 false
966 }
967 None => false,
968 }
969 }
970
971 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
972 let mut retries = self.reorder_retries.lock();
973 let retry = retries.entry(segment_id.to_string()).or_default();
974 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
975 let delay = merge_retry_delay(retry.consecutive_failures);
976 retry.retry_after = std::time::Instant::now().checked_add(delay);
977 log::warn!(
978 "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
979 segment_id,
980 delay.as_secs_f64(),
981 retry.consecutive_failures,
982 error,
983 );
984 }
985
986 fn clear_reorder_retry(&self, segment_id: &str) {
987 self.reorder_retries.lock().remove(segment_id);
988 }
989
990 fn paused_reorder_segments(&self) -> HashSet<String> {
991 let now = std::time::Instant::now();
992 let mut retries = self.reorder_retries.lock();
993 let mut paused = HashSet::new();
994 for (segment_id, retry) in retries.iter_mut() {
995 match retry.retry_after {
996 Some(deadline) if deadline > now => {
997 paused.insert(segment_id.clone());
998 }
999 Some(_) => retry.retry_after = None,
1000 None => {}
1001 }
1002 }
1003 paused
1004 }
1005
1006 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1010 if self
1011 .global_merge_wakeup_pending
1012 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1013 .is_err()
1014 {
1015 return;
1016 }
1017
1018 let manager = Arc::clone(self);
1019 let future = async move {
1020 let capacity = tokio::select! {
1021 biased;
1022 () = manager.active_operations.wait_for_shutdown() => None,
1023 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1024 };
1025
1026 manager
1027 .global_merge_wakeup_pending
1028 .store(false, Ordering::Release);
1029 if let Some(permit) = capacity {
1030 drop(permit);
1034 manager.maybe_merge().await;
1035 }
1036 };
1037 let runtime = tokio::runtime::Handle::current();
1038 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1039 self.global_merge_wakeup_pending
1040 .store(false, Ordering::Release);
1041 log::warn!("[merge] runtime rejected global-capacity wakeup task");
1042 }
1043 }
1044
1045 #[cfg(test)]
1046 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1047 self.quarantined_segments.lock().contains(segment_id)
1048 }
1049
1050 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1055 let output_hex = output_id.to_hex();
1056 {
1057 let st = self.state.lock().await;
1058 if st.metadata.has_segment(&output_hex) {
1059 return;
1060 }
1061 }
1062
1063 log::info!(
1067 "[segment_cleanup] deleting uncommitted output {} after {}",
1068 output_hex,
1069 reason,
1070 );
1071 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1072 {
1073 log::warn!(
1074 "[segment_cleanup] failed deleting uncommitted output {}: {}",
1075 output_hex,
1076 error,
1077 );
1078 }
1079 }
1080
1081 pub async fn get_segment_ids(&self) -> Vec<String> {
1087 self.state.lock().await.metadata.segment_ids()
1088 }
1089
1090 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1092 self.trained.load_full()
1093 }
1094
1095 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1102 if self.vector_artifact_update.load(Ordering::Acquire) {
1103 return None;
1104 }
1105 let trained = self.trained.load_full();
1106 if self.vector_artifact_update.load(Ordering::Acquire) {
1107 None
1108 } else {
1109 trained
1110 }
1111 }
1112
1113 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1132 self.vector_artifact_update
1133 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1134 .map_err(|_| {
1135 Error::Internal("a trained-vector artifact update is already in progress".into())
1136 })?;
1137 self.active_operations.pause_non_indexing();
1138 let guard = VectorArtifactUpdateGuard {
1139 _lease: Arc::new(VectorArtifactUpdateLease {
1140 updating: Arc::clone(&self.vector_artifact_update),
1141 active_operations: Arc::clone(&self.active_operations),
1142 }),
1143 };
1144 let (preexisting, parked_indexing) =
1145 self.active_operations.draining_operation_tokens_snapshot();
1146 if parked_indexing > 0 {
1147 return Err(Error::Internal(format!(
1148 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1149 segment(s) are built but uncommitted; commit or abort the pending \
1150 generation and retry"
1151 )));
1152 }
1153 self.active_operations
1154 .wait_until_operations_finish(&preexisting)
1155 .await;
1156 Ok(guard)
1157 }
1158
1159 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1162 let vector_fields = {
1164 let st = self.state.lock().await;
1165 st.metadata.vector_fields.clone()
1166 };
1167 let trained = IndexMetadata::try_load_trained_from_fields(
1169 &vector_fields,
1170 self.schema.as_ref(),
1171 self.directory.as_ref(),
1172 )
1173 .await?
1174 .map(Arc::new);
1175 self.trained.store(trained);
1179 Ok(())
1180 }
1181
1182 pub(crate) async fn publish_vector_generation(
1189 self: &Arc<Self>,
1190 artifact_update: &VectorArtifactUpdateGuard,
1191 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1192 next_trained: Arc<TrainedVectorStructures>,
1193 mut staged: Vec<StagedVectorSegment>,
1194 ) -> Result<()> {
1195 if !self.vector_artifact_update.load(Ordering::Acquire) {
1196 return Err(Error::Internal(
1197 "vector generation publication lost its exclusive update lease".into(),
1198 ));
1199 }
1200
1201 for replacement in &staged {
1202 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1203 .await?;
1204 }
1205
1206 let mut st = Arc::clone(&self.state).lock_owned().await;
1207 let mut next = st.metadata.clone();
1208 next.vector_fields = vector_fields;
1209 next.refresh_total_vectors();
1210
1211 for replacement in &staged {
1212 let source_info = next
1213 .segment_metas
1214 .remove(&replacement.source_id)
1215 .ok_or_else(|| {
1216 Error::Corruption(format!(
1217 "vector generation source {} disappeared before publication",
1218 replacement.source_id,
1219 ))
1220 })?;
1221 let output_hex = replacement.output_id.to_hex();
1222 if next.segment_metas.contains_key(&output_hex) {
1223 return Err(Error::Corruption(format!(
1224 "vector generation output {output_hex} is already metadata-live"
1225 )));
1226 }
1227 next.add_segment_meta(output_hex, source_info);
1230 }
1231
1232 let directory = Arc::clone(&self.directory);
1233 let trained = Arc::clone(&self.trained);
1234 let tracker = Arc::clone(&self.tracker);
1235 let artifact_update = artifact_update.clone();
1239 self.run_lifecycle_transaction(async move {
1240 let _artifact_update = artifact_update;
1241 next.save(directory.as_ref()).await?;
1242
1243 for replacement in &staged {
1244 tracker.register(&replacement.output_id.to_hex());
1245 }
1246 st.metadata = next;
1247 trained.store(Some(next_trained));
1248
1249 for replacement in &mut staged {
1252 replacement.cleanup.disarm();
1253 }
1254 let retired = staged
1255 .iter()
1256 .map(|replacement| replacement.source_id.clone())
1257 .collect::<Vec<_>>();
1258 let ready_to_delete = tracker.mark_for_deletion(&retired);
1259 drop(st);
1260 for &segment_id in &ready_to_delete {
1261 if let Err(error) =
1262 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1263 {
1264 log::warn!(
1265 "[segment_cleanup] immediate dense-vector generation delete failed for {}: {}",
1266 segment_id.to_hex(),
1267 error,
1268 );
1269 }
1270 }
1271 tracker.complete_deletion(&ready_to_delete);
1272 Ok(())
1273 })
1274 .await
1275 }
1276
1277 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1279 where
1280 F: FnOnce(&IndexMetadata) -> R,
1281 {
1282 let st = self.state.lock().await;
1283 f(&st.metadata)
1284 }
1285
1286 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1288 where
1289 F: FnOnce(&mut IndexMetadata),
1290 {
1291 let mut st = Arc::clone(&self.state).lock_owned().await;
1292 let mut next = st.metadata.clone();
1293 f(&mut next);
1294 let directory = Arc::clone(&self.directory);
1295 self.run_lifecycle_transaction(async move {
1296 next.save(directory.as_ref()).await?;
1297 st.metadata = next;
1298 Ok(())
1299 })
1300 .await
1301 }
1302
1303 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1306 let (acquired, trained) = {
1307 let st = self.state.lock().await;
1308 let segment_ids = st.metadata.segment_ids();
1309 (self.tracker.acquire(&segment_ids), self.trained.load_full())
1310 };
1311
1312 SegmentSnapshot::with_generation(
1313 Arc::clone(&self.tracker),
1314 acquired,
1315 trained,
1316 Arc::clone(&self.delete_fn),
1317 )
1318 }
1319
1320 pub fn tracker(&self) -> Arc<SegmentTracker> {
1322 Arc::clone(&self.tracker)
1323 }
1324
1325 pub fn directory(&self) -> Arc<D> {
1327 Arc::clone(&self.directory)
1328 }
1329}
1330
1331#[cfg(feature = "native")]
1336impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1337 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1339 for (segment_id, num_docs) in new_segments {
1342 self.validate_completed_segment(segment_id, *num_docs)
1343 .await?;
1344 }
1345
1346 let mut st = Arc::clone(&self.state).lock_owned().await;
1347 let mut next = st.metadata.clone();
1348 let mut added = Vec::new();
1349 for (segment_id, num_docs) in new_segments {
1350 if !next.has_segment(segment_id) {
1351 next.add_segment(segment_id.clone(), *num_docs);
1352 added.push(segment_id.clone());
1353 }
1354 }
1355
1356 let directory = Arc::clone(&self.directory);
1362 let tracker = Arc::clone(&self.tracker);
1363 self.run_lifecycle_transaction(async move {
1364 next.save(directory.as_ref()).await?;
1365 for segment_id in &added {
1366 tracker.register(segment_id);
1367 }
1368 st.metadata = next;
1369 Ok(())
1370 })
1371 .await
1372 }
1373
1374 pub async fn maybe_merge(self: &Arc<Self>) {
1385 if !self.active_operations.is_accepting() {
1386 log::debug!("[maybe_merge] manager is shutting down, skipping");
1387 return;
1388 }
1389 if self.merge_retry_is_paused() {
1390 log::debug!("[maybe_merge] retry backoff active, skipping");
1391 return;
1392 }
1393
1394 {
1397 let mut handles = self.merge_handles.lock();
1398 handles.retain(|h| !h.is_finished());
1399 }
1400 let local_slots = self.merge_permits.available_permits();
1401 let global_slots = self.global_merge_permits.available_permits();
1402 let slots_available = local_slots.min(global_slots);
1403
1404 {
1408 let st = self.state.lock().await;
1409 let quarantined = self.quarantined_segments.lock().clone();
1410 let active_ids = self.active_operations.snapshot();
1411
1412 let segments: Vec<SegmentInfo> = st
1415 .metadata
1416 .segment_metas
1417 .iter()
1418 .filter(|(id, _)| {
1419 !self.tracker.is_pending_deletion(id)
1420 && !active_ids.contains(*id)
1421 && !quarantined.contains(*id)
1422 })
1423 .map(|(id, info)| SegmentInfo {
1424 id: id.clone(),
1425 num_docs: info.num_docs,
1426 })
1427 .collect();
1428
1429 log::debug!("[maybe_merge] {} eligible segments", segments.len());
1430
1431 let candidates = st.merge_policy.find_merges(&segments);
1432
1433 if candidates.is_empty() {
1434 return;
1435 }
1436
1437 if slots_available == 0 {
1441 if local_slots > 0 && global_slots == 0 {
1442 self.schedule_global_merge_wakeup();
1443 }
1444 log::debug!("[maybe_merge] at max concurrent merges, skipping");
1445 return;
1446 }
1447
1448 log::debug!(
1449 "[maybe_merge] {} merge candidates, {} slots available",
1450 candidates.len(),
1451 slots_available
1452 );
1453
1454 let mut handles = Vec::new();
1455 for c in candidates {
1456 if handles.len() >= slots_available {
1457 break;
1458 }
1459 if let Some(h) = self.spawn_merge(c.segment_ids) {
1460 handles.push(h);
1461 }
1462 }
1463 if !handles.is_empty() {
1464 self.merge_handles.lock().extend(handles);
1469 }
1470 }
1471 }
1472
1473 fn spawn_merge(self: &Arc<Self>, segment_ids_to_merge: Vec<String>) -> Option<JoinHandle<()>> {
1482 if self.force_merge_active.load(Ordering::Acquire) > 0 {
1483 log::debug!("[spawn_merge] skipped: explicit force merge has priority");
1484 return None;
1485 }
1486 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1487 Ok(permit) => permit,
1488 Err(_) => {
1489 log::debug!("[spawn_merge] skipped: global merge capacity is full");
1490 self.schedule_global_merge_wakeup();
1491 return None;
1492 }
1493 };
1494 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1495 Ok(permit) => permit,
1496 Err(_) => {
1497 log::debug!("[spawn_merge] skipped: no merge permit available");
1498 return None;
1499 }
1500 };
1501 let output_id = SegmentId::new();
1502 let output_hex = output_id.to_hex();
1503
1504 let mut all_ids = segment_ids_to_merge.clone();
1505 all_ids.push(output_hex);
1506
1507 let guard = match self.active_operations.try_register(all_ids) {
1508 Some(g) => g,
1509 None => {
1510 log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1511 return None;
1512 }
1513 };
1514
1515 let sm = Arc::clone(self);
1516 let ids = segment_ids_to_merge;
1517
1518 Some(tokio::spawn(async move {
1519 let mut output_cleanup = sm.output_cleanup_guard(output_id);
1520 let mut reevaluate = false;
1521 let mut retry_delay = None;
1522
1523 let trained_snap = sm.trained_for_segment_build();
1524 let granularity = sm.merge_granularity(&ids).await;
1525 let result = Self::do_merge(
1526 sm.directory.as_ref(),
1527 &sm.schema,
1528 &ids,
1529 output_id,
1530 sm.term_cache_blocks,
1531 trained_snap.as_deref(),
1532 sm.reorder_on_merge,
1533 granularity,
1534 sm.merge_bp_time_budget,
1535 sm.bp_memory_budget_bytes,
1536 Arc::clone(&sm.reorder_permits),
1537 ReorderPriority::Background,
1538 Some(sm.background_cpu_pool()),
1539 )
1540 .await;
1541
1542 match result {
1543 Ok((new_id, doc_count, bp_converged)) => {
1544 match sm
1545 .replace_segments(
1546 &ids,
1547 new_id,
1548 doc_count,
1549 ReplacementLayout::Recomputed {
1550 reordered: sm.reorder_on_merge,
1551 bp_converged,
1552 },
1553 )
1554 .await
1555 {
1556 Ok(()) => {
1557 output_cleanup.disarm();
1558 sm.clear_merge_retry_backoff();
1559 reevaluate = true;
1560 }
1561 Err(e) => {
1562 sm.delete_output_if_unregistered(output_id, "replacement failure")
1563 .await;
1564 output_cleanup.disarm();
1565 retry_delay = Some(sm.pause_merge_retries(&e));
1566 log::error!("[merge] failed to publish merged segment: {}", e);
1567 }
1568 }
1569 }
1570 Err(MergeTaskError {
1571 error,
1572 unavailable_segments,
1573 }) => {
1574 log::error!(
1575 "[merge] background merge failed for segments {:?}: {}",
1576 ids,
1577 error
1578 );
1579 if !unavailable_segments.is_empty() {
1580 for segment_id in &unavailable_segments {
1581 sm.quarantine_segment(segment_id, &error);
1582 }
1583 reevaluate = true;
1587 } else {
1588 retry_delay = Some(sm.pause_merge_retries(&error));
1589 }
1590 sm.delete_output_if_unregistered(output_id, "merge failure")
1591 .await;
1592 output_cleanup.disarm();
1593 }
1594 }
1595 drop(guard);
1598 drop(merge_permit);
1600 drop(global_merge_permit);
1601
1602 if reevaluate {
1603 sm.maybe_merge().await;
1604 } else if let Some(retry_delay) = retry_delay {
1605 sm.schedule_merge_retry_wakeup(retry_delay);
1612 }
1613 }))
1614 }
1615
1616 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
1620 let manager = Arc::clone(self);
1621 let future = async move {
1622 tokio::select! {
1623 () = tokio::time::sleep(retry_delay) => {
1624 manager.maybe_merge().await;
1625 }
1626 () = manager.active_operations.wait_for_shutdown() => {}
1627 }
1628 };
1629 let runtime = tokio::runtime::Handle::current();
1630 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1631 log::warn!(
1632 "[merge] runtime rejected merge-retry wakeup task; eligible segments may stay \
1633 unmerged until the next commit re-runs merge policy evaluation"
1634 );
1635 }
1636 }
1637
1638 async fn replace_segments(
1642 self: &Arc<Self>,
1643 old_ids: &[String],
1644 new_id: String,
1645 doc_count: u32,
1646 layout: ReplacementLayout,
1647 ) -> Result<()> {
1648 self.validate_completed_segment(&new_id, doc_count).await?;
1651 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1652 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1653 })?;
1654 let output_reader = SegmentReader::open(
1655 self.directory.as_ref(),
1656 output_id,
1657 Arc::clone(&self.schema),
1658 self.term_cache_blocks,
1659 )
1660 .await
1661 .map_err(|error| match error {
1662 Error::Io(_) | Error::IndexClosed => error,
1666 error => Error::Corruption(format!(
1667 "replacement segment {new_id} failed full reader validation: {error}"
1668 )),
1669 })?;
1670 if output_reader.num_docs() != doc_count {
1671 return Err(Error::Corruption(format!(
1672 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
1673 output_reader.num_docs(),
1674 )));
1675 }
1676 drop(output_reader);
1677
1678 let mut st = Arc::clone(&self.state).lock_owned().await;
1679 let missing: Vec<&String> = old_ids
1683 .iter()
1684 .filter(|id| !st.metadata.has_segment(id))
1685 .collect();
1686 if !missing.is_empty() {
1687 return Err(Error::Corruption(format!(
1688 "replace_segments: source segment(s) {:?} not in metadata — \
1689 refusing to add output {} (would duplicate documents)",
1690 missing, new_id
1691 )));
1692 }
1693
1694 let replacement_info = match layout {
1695 ReplacementLayout::Recomputed {
1696 reordered,
1697 bp_converged,
1698 } => {
1699 let generation = old_ids
1700 .iter()
1701 .filter_map(|id| st.metadata.segment_metas.get(id))
1702 .map(|info| info.generation)
1703 .max()
1704 .unwrap_or(0)
1705 .checked_add(1)
1706 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
1707 let parent_unconverged_passes = old_ids
1708 .iter()
1709 .filter_map(|id| st.metadata.segment_metas.get(id))
1710 .map(|info| info.bp_unconverged_passes)
1711 .max()
1712 .unwrap_or(0);
1713 let passes = if reordered && !bp_converged {
1714 parent_unconverged_passes.saturating_add(1)
1715 } else {
1716 0
1717 };
1718 SegmentMetaInfo {
1719 num_docs: doc_count,
1720 ancestors: old_ids.to_vec(),
1721 generation,
1722 reordered,
1723 bp_converged,
1724 bp_unconverged_passes: passes,
1725 }
1726 }
1727 ReplacementLayout::PreserveSingleSource => {
1728 let [source_id] = old_ids else {
1729 return Err(Error::Internal(
1730 "layout-preserving replacement requires exactly one source".into(),
1731 ));
1732 };
1733 let mut source = st
1734 .metadata
1735 .segment_metas
1736 .get(source_id)
1737 .cloned()
1738 .ok_or_else(|| {
1739 Error::Corruption(format!(
1740 "layout-preserving replacement source {source_id} disappeared"
1741 ))
1742 })?;
1743 source.num_docs = doc_count;
1744 source
1745 }
1746 };
1747 let retired_ids = old_ids.to_vec();
1748 let mut next = st.metadata.clone();
1749 for id in old_ids {
1750 next.remove_segment(id);
1751 }
1752 next.add_segment_meta(new_id.clone(), replacement_info);
1753
1754 let directory = Arc::clone(&self.directory);
1755 let tracker = Arc::clone(&self.tracker);
1756 self.run_lifecycle_transaction(async move {
1757 next.save(directory.as_ref()).await?;
1760 tracker.register(&new_id);
1761 st.metadata = next;
1762
1763 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
1767 drop(st);
1768 for &segment_id in &ready_to_delete {
1769 if let Err(error) =
1770 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1771 {
1772 log::warn!(
1773 "[segment_cleanup] immediate delete failed for {}: {}",
1774 segment_id.to_hex(),
1775 error,
1776 );
1777 }
1778 }
1779 tracker.complete_deletion(&ready_to_delete);
1780 Ok(())
1781 })
1782 .await
1783 }
1784
1785 #[allow(clippy::too_many_arguments)]
1790 async fn do_merge(
1791 directory: &D,
1792 schema: &Arc<crate::dsl::Schema>,
1793 segment_ids_to_merge: &[String],
1794 output_segment_id: SegmentId,
1795 term_cache_blocks: usize,
1796 trained: Option<&TrainedVectorStructures>,
1797 reorder_bmp: bool,
1798 granularity: crate::segment::reorder::BpGranularity,
1799 merge_bp_time_budget: Option<std::time::Duration>,
1800 bp_memory_budget_bytes: usize,
1801 reorder_permits: Arc<ReorderConcurrencyGate>,
1802 reorder_priority: ReorderPriority,
1803 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
1804 ) -> MergeTaskResult<(String, u32, bool)> {
1805 let output_hex = output_segment_id.to_hex();
1806 let load_start = std::time::Instant::now();
1807
1808 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
1809 for id_str in segment_ids_to_merge {
1810 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
1811 MergeTaskError::source(
1812 id_str.clone(),
1813 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
1814 )
1815 })?;
1816 segment_ids.push(id);
1817 }
1818
1819 let mut unavailable_sources = Vec::new();
1824 let mut missing_files = Vec::new();
1825 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
1826 let files = SegmentFiles::new(id.0);
1827 let mut source_unavailable = false;
1828 for path in files.mandatory_paths() {
1829 let exists = directory
1830 .exists(path)
1831 .await
1832 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
1833 if !exists {
1834 source_unavailable = true;
1835 missing_files.push(format!("{}:{:?}", id_str, path));
1836 }
1837 }
1838 if source_unavailable {
1839 unavailable_sources.push(id_str.clone());
1840 }
1841 }
1842 if !unavailable_sources.is_empty() {
1843 return Err(MergeTaskError::sources(
1844 unavailable_sources,
1845 Error::Corruption(format!(
1846 "merge sources are missing mandatory files: {}",
1847 missing_files.join(", ")
1848 )),
1849 ));
1850 }
1851
1852 let schema_arc = Arc::clone(schema);
1853 let futures: Vec<_> = segment_ids
1854 .iter()
1855 .map(|&sid| {
1856 let sch = Arc::clone(&schema_arc);
1857 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
1858 })
1859 .collect();
1860
1861 let results = futures::future::join_all(futures).await;
1862 let mut readers = Vec::with_capacity(results.len());
1863 let mut total_docs = 0u64;
1864 for (i, result) in results.into_iter().enumerate() {
1865 match result {
1866 Ok(r) => {
1867 total_docs += r.meta().num_docs as u64;
1868 readers.push(r);
1869 }
1870 Err(e) => {
1871 log::error!(
1872 "[merge] Failed to open segment {}: {:?}",
1873 segment_ids_to_merge[i],
1874 e
1875 );
1876 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
1877 }
1878 }
1879 }
1880 if total_docs > u32::MAX as u64 {
1881 return Err(Error::Internal(format!(
1882 "Merged segment doc count ({}) exceeds u32::MAX",
1883 total_docs
1884 ))
1885 .into());
1886 }
1887
1888 for (i, reader) in readers.iter().enumerate() {
1892 let meta_docs = reader.meta().num_docs;
1893 let store_docs = reader.store().num_docs();
1894 if store_docs != meta_docs {
1895 return Err(MergeTaskError::source(
1896 segment_ids_to_merge[i].clone(),
1897 Error::Corruption(format!(
1898 "pre-merge validation: segment {} store has {} docs but meta says {}",
1899 segment_ids_to_merge[i], store_docs, meta_docs
1900 )),
1901 ));
1902 }
1903 }
1904
1905 log::info!(
1906 "[merge] loaded {} segment readers in {:.1}s",
1907 readers.len(),
1908 load_start.elapsed().as_secs_f64()
1909 );
1910
1911 let merger = SegmentMerger::new(Arc::clone(schema))
1912 .with_bmp_reorder(reorder_bmp)
1913 .with_granularity(granularity)
1914 .with_bp_budget(crate::segment::BpBudget {
1915 min_partition_docs: None,
1916 time_budget: merge_bp_time_budget,
1917 })
1918 .with_bp_memory_budget(bp_memory_budget_bytes)
1919 .with_reorder_permits(reorder_permits)
1920 .with_reorder_priority(reorder_priority)
1921 .with_background_pool(bg_cpu_pool);
1922
1923 log::info!(
1924 "[merge] {} segments -> {} (trained={})",
1925 segment_ids_to_merge.len(),
1926 output_hex,
1927 trained.map_or(0, |t| t.centroids.len()),
1928 );
1929
1930 let (_merged_meta, merge_stats) = merger
1931 .merge(directory, &readers, output_segment_id, trained)
1932 .await
1933 .map_err(|error| {
1934 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
1935 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
1941 } else {
1942 MergeTaskError::from(error)
1943 }
1944 })?;
1945 let bp_converged = merge_stats.bp_converged;
1946 if !bp_converged {
1947 log::info!(
1948 "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
1949 the background optimizer deepens it later",
1950 );
1951 }
1952
1953 log::info!(
1954 "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
1955 load_start.elapsed().as_secs_f64(),
1956 readers.len(),
1957 total_docs,
1958 );
1959
1960 Ok((output_hex, total_docs as u32, bp_converged))
1961 }
1962
1963 pub async fn abort_merges(&self) {
1974 loop {
1975 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
1976 if handles.is_empty() {
1977 return;
1978 }
1979 while let Some(result) = handles.join_next().await {
1980 if let Err(error) = result
1981 && error.is_panic()
1982 {
1983 log::error!("[merge] background task panicked while draining: {}", error);
1984 }
1985 }
1986 }
1987 }
1988
1989 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
1994 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
1995 while handles.join_next().await.is_some() {}
1996 }
1997
1998 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2007 loop {
2008 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2009 if handles.is_empty() {
2010 break;
2011 }
2012 while handles.join_next().await.is_some() {}
2013 }
2014 }
2015
2016 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2021 self.wait_for_all_merges().await;
2022 self.active_operations.wait_until_idle().await;
2023 loop {
2024 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2025 if handles.is_empty() {
2026 break;
2027 }
2028 for handle in handles {
2029 if let Err(error) = handle.await
2030 && error.is_panic()
2031 {
2032 log::error!("[segment_cleanup] task panicked while draining: {}", error);
2033 }
2034 }
2035 }
2036 }
2037
2038 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2047 const FORCE_MERGE_BATCH: usize = 64;
2048 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2053 std::time::Duration::from_millis(100);
2054 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2058
2059 let (_force_merge_activity, max_segment_docs) = {
2060 let st = self.state.lock().await;
2061 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2065 (
2066 ForceMergeActivityGuard(&self.force_merge_active),
2067 st.merge_policy.max_segment_docs(),
2068 )
2069 };
2070
2071 self.wait_for_all_merges().await;
2074
2075 let _foreground_reorder = if self.reorder_on_merge {
2082 log::info!(
2083 "[force_merge] prioritizing BP capacity ({} total pass slot(s))",
2084 self.reorder_permits.limit(),
2085 );
2086 Some(
2087 Arc::clone(&self.reorder_permits)
2088 .begin_foreground()
2089 .await
2090 .map_err(|_| {
2091 Error::Internal("background reorder scheduler is closed".into())
2092 })?,
2093 )
2094 } else {
2095 None
2096 };
2097
2098 let mut logged_held_wait = false;
2100
2101 loop {
2102 if !self.active_operations.is_accepting() {
2103 return Err(Error::IndexClosed);
2104 }
2105 let mut segments: Vec<(String, u32)> = {
2107 let st = self.state.lock().await;
2108 st.metadata
2109 .segment_metas
2110 .iter()
2111 .map(|(id, info)| (id.clone(), info.num_docs))
2112 .collect()
2113 };
2114
2115 if segments.len() < 2 {
2116 return Ok(());
2117 }
2118
2119 segments.sort_by_key(|(_, docs)| *docs);
2120
2121 let active_ids = self.active_operations.snapshot();
2129 let held: usize = segments
2130 .iter()
2131 .filter(|(id, _)| active_ids.contains(id))
2132 .count();
2133
2134 let max_docs = max_segment_docs.map(|m| m as u64).unwrap_or(u64::MAX);
2136 let mut batch = Vec::new();
2137 let mut batch_docs = 0u64;
2138
2139 for (id, docs) in &segments {
2140 if active_ids.contains(id) {
2141 continue;
2142 }
2143 if batch.len() >= FORCE_MERGE_BATCH {
2144 break;
2145 }
2146 let next_total = batch_docs + *docs as u64;
2147 if next_total > max_docs && !batch.is_empty() {
2148 break;
2149 }
2150 batch.push(id.clone());
2151 batch_docs += *docs as u64;
2152 }
2153
2154 if batch.len() < 2 {
2155 if held == 0 {
2156 return Ok(());
2159 }
2160 if !logged_held_wait {
2164 log::info!(
2165 "[force_merge] waiting: {} segment(s) held by active \
2166 merge/reorder operations, none free to merge",
2167 held
2168 );
2169 logged_held_wait = true;
2170 } else {
2171 log::debug!("[force_merge] still waiting on {} held segment(s)", held);
2172 }
2173 #[cfg(test)]
2174 self.force_merge_conflict_retries
2175 .fetch_add(1, Ordering::Relaxed);
2176 tokio::select! {
2177 biased;
2178 () = self.active_operations.wait_for_shutdown() => {
2179 return Err(Error::IndexClosed);
2180 }
2181 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2182 }
2183 continue;
2184 }
2185 logged_held_wait = false;
2186
2187 let _global_merge_permit = tokio::select! {
2188 biased;
2189 () = self.active_operations.wait_for_shutdown() => {
2190 return Err(Error::IndexClosed);
2191 }
2192 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2193 permit.map_err(|_| {
2194 Error::Internal("global background merge scheduler is closed".into())
2195 })?
2196 }
2197 };
2198
2199 let output_id = SegmentId::new();
2200 let output_hex = output_id.to_hex();
2201
2202 let mut all_ids = batch.clone();
2205 all_ids.push(output_hex);
2206 let guard = {
2207 let st = self.state.lock().await;
2208 batch
2209 .iter()
2210 .all(|id| st.metadata.has_segment(id))
2211 .then(|| self.active_operations.try_register(all_ids))
2212 .flatten()
2213 };
2214 let _guard = match guard {
2215 Some(g) => g,
2216 None if !self.active_operations.is_accepting() => {
2217 return Err(Error::IndexClosed);
2218 }
2219 None => {
2220 #[cfg(test)]
2221 self.force_merge_conflict_retries
2222 .fetch_add(1, Ordering::Relaxed);
2223 drop(_global_merge_permit);
2226 log::debug!("[force_merge] batch lost a registration race, rebuilding");
2232 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2233 self.wait_for_merging_thread().await;
2234 if !had_tracked_merges {
2235 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2236 }
2237 continue;
2238 }
2239 };
2240 log::info!(
2244 "[force_merge] merging batch of {} segments ({} docs)",
2245 batch.len(),
2246 batch_docs
2247 );
2248 let mut output_cleanup = self.output_cleanup_guard(output_id);
2249
2250 let trained_snap = self.trained_for_segment_build();
2251 let granularity = self.merge_granularity(&batch).await;
2252 let merge_result = Self::do_merge(
2253 self.directory.as_ref(),
2254 &self.schema,
2255 &batch,
2256 output_id,
2257 self.term_cache_blocks,
2258 trained_snap.as_deref(),
2259 self.reorder_on_merge,
2260 granularity,
2261 self.merge_bp_time_budget,
2262 self.bp_memory_budget_bytes,
2263 Arc::clone(&self.reorder_permits),
2264 ReorderPriority::Foreground,
2265 Some(self.background_cpu_pool()),
2266 )
2267 .await;
2268 let (new_segment_id, total_docs, bp_converged) = match merge_result {
2269 Ok(v) => v,
2270 Err(MergeTaskError {
2271 error,
2272 unavailable_segments,
2273 }) => {
2274 for segment_id in &unavailable_segments {
2275 self.quarantine_segment(segment_id, &error);
2276 }
2277 self.delete_output_if_unregistered(output_id, "force-merge failure")
2278 .await;
2279 output_cleanup.disarm();
2280 return Err(error);
2281 }
2282 };
2283
2284 if let Err(e) = self
2285 .replace_segments(
2286 &batch,
2287 new_segment_id,
2288 total_docs,
2289 ReplacementLayout::Recomputed {
2290 reordered: self.reorder_on_merge,
2291 bp_converged,
2292 },
2293 )
2294 .await
2295 {
2296 self.delete_output_if_unregistered(output_id, "replacement failure")
2297 .await;
2298 output_cleanup.disarm();
2299 return Err(e);
2300 }
2301 output_cleanup.disarm();
2302
2303 }
2305 }
2306
2307 fn segment_needs_vector_rewrite(
2308 &self,
2309 reader: &SegmentReader,
2310 field_ids: &[u32],
2311 rewrite_existing: bool,
2312 ) -> Result<bool> {
2313 for &field_id in field_ids {
2314 let flat = reader.flat_vectors().get(&field_id);
2315 let ann = reader.vector_indexes().get(&field_id);
2316 if ann.is_some() && flat.is_none() {
2317 return Err(Error::Corruption(format!(
2318 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2319 reader.meta().id,
2320 )));
2321 }
2322
2323 let Some(flat) = flat else {
2324 continue;
2325 };
2326 if flat.num_vectors == 0 {
2327 continue;
2328 }
2329 if rewrite_existing {
2330 return Ok(true);
2331 }
2332 let field = crate::dsl::Field(field_id);
2333 let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2334 Error::Corruption(format!(
2335 "segment {:032x} references unknown vector field {field_id}",
2336 reader.meta().id,
2337 ))
2338 })?;
2339 let current = match entry.field_type {
2340 crate::dsl::FieldType::DenseVector
2344 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2345 config.index_type == crate::dsl::VectorIndexType::Tq
2346 }) =>
2347 {
2348 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2349 }
2350 crate::dsl::FieldType::DenseVector
2351 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2352 config.index_type == crate::dsl::VectorIndexType::IvfTq
2353 }) =>
2354 {
2355 matches!(ann, Some(crate::segment::VectorIndex::IvfTq { .. }))
2356 }
2357 crate::dsl::FieldType::DenseVector => false,
2360 crate::dsl::FieldType::BinaryDenseVector => {
2361 matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2362 }
2363 _ => false,
2364 };
2365 if !current {
2366 return Ok(true);
2367 }
2368 }
2369 Ok(false)
2370 }
2371
2372 async fn acquire_vector_rewrite_capacity(
2373 &self,
2374 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2375 let global = tokio::select! {
2376 biased;
2377 () = self.active_operations.wait_for_shutdown() => {
2378 return Err(Error::IndexClosed);
2379 }
2380 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2381 permit.map_err(|_| Error::Internal(
2382 "global background merge scheduler is closed".into()
2383 ))?
2384 }
2385 };
2386 let local = tokio::select! {
2387 biased;
2388 () = self.active_operations.wait_for_shutdown() => {
2389 return Err(Error::IndexClosed);
2390 }
2391 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2392 permit.map_err(|_| Error::Internal(
2393 "background merge scheduler is closed".into()
2394 ))?
2395 }
2396 };
2397 Ok((global, local))
2398 }
2399
2400 async fn build_vector_replacement(
2401 self: &Arc<Self>,
2402 segment_id: &str,
2403 source_id: SegmentId,
2404 output_id: SegmentId,
2405 trained: &TrainedVectorStructures,
2406 failure_context: &'static str,
2407 ) -> Result<(String, u32, OutputCleanupGuard)> {
2408 let mut cleanup = self.output_cleanup_guard(output_id);
2409 match crate::segment::reorder::rewrite_vector_segment(
2410 self.directory.as_ref(),
2411 &self.schema,
2412 source_id,
2413 output_id,
2414 self.term_cache_blocks,
2415 trained,
2416 Some(self.background_cpu_pool()),
2417 )
2418 .await
2419 {
2420 Ok((new_id, doc_count)) => {
2421 self.validate_completed_segment(&new_id, doc_count).await?;
2422 Ok((new_id, doc_count, cleanup))
2423 }
2424 Err(error) => {
2425 self.delete_output_if_unregistered(output_id, failure_context)
2426 .await;
2427 cleanup.disarm();
2428 if is_deterministic_source_error(&error) {
2429 self.quarantine_segment(segment_id, &error);
2430 }
2431 Err(error)
2432 }
2433 }
2434 }
2435
2436 pub(crate) async fn stage_vector_generation(
2440 self: &Arc<Self>,
2441 _artifact_update: &VectorArtifactUpdateGuard,
2442 segment_ids: &[String],
2443 field_ids: &[u32],
2444 trained: Arc<TrainedVectorStructures>,
2445 rewrite_existing: bool,
2446 ) -> Result<Vec<StagedVectorSegment>> {
2447 if !self.vector_artifact_update.load(Ordering::Acquire) {
2448 return Err(Error::Internal(
2449 "cannot stage a vector generation without an exclusive update lease".into(),
2450 ));
2451 }
2452
2453 let mut staged = Vec::new();
2454 for segment_id in segment_ids {
2455 if self.quarantined_segments.lock().contains(segment_id) {
2456 return Err(Error::Corruption(format!(
2457 "segment {segment_id} is quarantined after a deterministic source failure"
2458 )));
2459 }
2460 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2461 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2462 })?;
2463
2464 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2467
2468 let output_id = SegmentId::new();
2469 let output_hex = output_id.to_hex();
2470 let operation = {
2471 let st = self.state.lock().await;
2472 if !st.metadata.has_segment(segment_id) {
2473 return Err(Error::Corruption(format!(
2474 "vector generation source {segment_id} disappeared while lifecycle work was paused"
2475 )));
2476 }
2477 self.active_operations
2478 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2479 }
2480 .ok_or_else(|| {
2481 if self.active_operations.is_accepting() {
2482 Error::Internal(format!(
2483 "vector generation could not claim stable source {segment_id}"
2484 ))
2485 } else {
2486 Error::IndexClosed
2487 }
2488 })?;
2489
2490 let reader = SegmentReader::open(
2491 self.directory.as_ref(),
2492 source_id,
2493 Arc::clone(&self.schema),
2494 self.term_cache_blocks,
2495 )
2496 .await?;
2497 if !self.segment_needs_vector_rewrite(&reader, field_ids, rewrite_existing)? {
2498 continue;
2499 }
2500 drop(reader);
2501
2502 let (new_id, doc_count, cleanup) = self
2503 .build_vector_replacement(
2504 segment_id,
2505 source_id,
2506 output_id,
2507 trained.as_ref(),
2508 "vector generation staging failure",
2509 )
2510 .await?;
2511 debug_assert_eq!(new_id, output_hex);
2512 let output_reader = SegmentReader::open(
2513 self.directory.as_ref(),
2514 output_id,
2515 Arc::clone(&self.schema),
2516 self.term_cache_blocks,
2517 )
2518 .await?;
2519 if self.segment_needs_vector_rewrite(&output_reader, field_ids, false)? {
2520 return Err(Error::Corruption(format!(
2521 "staged vector segment {new_id} does not match its candidate codebook generation"
2522 )));
2523 }
2524
2525 staged.push(StagedVectorSegment {
2526 source_id: segment_id.clone(),
2527 output_id,
2528 doc_count,
2529 _operation: operation,
2530 cleanup,
2531 });
2532 }
2533 Ok(staged)
2534 }
2535
2536 async fn rewrite_vector_segment_once(
2537 self: &Arc<Self>,
2538 segment_id: &str,
2539 field_ids: &[u32],
2540 ) -> Result<VectorSegmentRewriteOutcome> {
2541 if self.quarantined_segments.lock().contains(segment_id) {
2542 return Err(Error::Corruption(format!(
2543 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
2544 )));
2545 }
2546 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2547 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2548 })?;
2549
2550 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2555
2556 let output_id = SegmentId::new();
2557 let output_hex = output_id.to_hex();
2558 let all_ids = vec![segment_id.to_owned(), output_hex];
2559 let operation = {
2560 let st = self.state.lock().await;
2561 if !st.metadata.has_segment(segment_id) {
2562 return Ok(VectorSegmentRewriteOutcome::SourceGone);
2563 }
2564 self.active_operations.try_register(all_ids)
2565 };
2566 let _operation = match operation {
2567 Some(operation) => operation,
2568 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
2569 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
2570 };
2571
2572 let Some(trained) = self.trained_for_segment_build() else {
2573 return Ok(VectorSegmentRewriteOutcome::Deferred);
2574 };
2575
2576 let reader = SegmentReader::open(
2577 self.directory.as_ref(),
2578 source_id,
2579 Arc::clone(&self.schema),
2580 self.term_cache_blocks,
2581 )
2582 .await?;
2583 if !self.segment_needs_vector_rewrite(&reader, field_ids, false)? {
2584 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
2585 }
2586 drop(reader);
2587
2588 let (new_id, doc_count, mut output_cleanup) = self
2589 .build_vector_replacement(
2590 segment_id,
2591 source_id,
2592 output_id,
2593 trained.as_ref(),
2594 "vector rewrite failure",
2595 )
2596 .await?;
2597
2598 if let Err(error) = self
2599 .replace_segments(
2600 &[segment_id.to_owned()],
2601 new_id,
2602 doc_count,
2603 ReplacementLayout::PreserveSingleSource,
2604 )
2605 .await
2606 {
2607 self.delete_output_if_unregistered(output_id, "vector replacement failure")
2608 .await;
2609 output_cleanup.disarm();
2610 return Err(error);
2611 }
2612 output_cleanup.disarm();
2613 Ok(VectorSegmentRewriteOutcome::Rewritten)
2614 }
2615
2616 pub(crate) async fn rewrite_vector_segments(
2621 self: &Arc<Self>,
2622 field_ids: &[u32],
2623 ) -> Result<usize> {
2624 if field_ids.is_empty() {
2625 return Ok(0);
2626 }
2627 let mut rewritten = 0usize;
2628 loop {
2629 let segment_ids = self.get_segment_ids().await;
2630 let mut conflicted = false;
2631 let mut changed = false;
2632 for segment_id in segment_ids {
2633 match self
2634 .rewrite_vector_segment_once(&segment_id, field_ids)
2635 .await?
2636 {
2637 VectorSegmentRewriteOutcome::Rewritten => {
2638 rewritten += 1;
2639 changed = true;
2640 }
2641 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
2642 VectorSegmentRewriteOutcome::Deferred => {
2643 return Err(Error::Internal(
2644 "ANN finalization lost the published trained generation".into(),
2645 ));
2646 }
2647 VectorSegmentRewriteOutcome::AlreadyCurrent
2648 | VectorSegmentRewriteOutcome::SourceGone => {}
2649 }
2650 }
2651 if !conflicted && !changed {
2652 log::info!(
2653 "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
2654 rewritten,
2655 );
2656 return Ok(rewritten);
2657 }
2658 tokio::select! {
2659 biased;
2660 () = self.active_operations.wait_for_shutdown() => {
2661 return Err(Error::IndexClosed);
2662 }
2663 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
2664 }
2665 }
2666 }
2667
2668 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
2673 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
2674 return;
2675 }
2676 let manager = Arc::clone(self);
2677 let future = async move {
2678 let field_ids = manager
2679 .read_metadata(|metadata| {
2680 metadata
2681 .vector_fields
2682 .keys()
2683 .filter(|field_id| metadata.is_field_built(**field_id))
2684 .copied()
2685 .collect::<Vec<_>>()
2686 })
2687 .await;
2688 for segment_id in segment_ids {
2689 loop {
2690 match manager
2691 .rewrite_vector_segment_once(&segment_id, &field_ids)
2692 .await
2693 {
2694 Ok(VectorSegmentRewriteOutcome::Conflict) => {
2695 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
2696 }
2697 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
2698 Ok(_) => break,
2699 Err(error) => {
2700 log::error!(
2701 "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
2702 segment_id,
2703 error,
2704 );
2705 break;
2706 }
2707 }
2708 }
2709 }
2710 };
2711 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
2712 log::warn!(
2713 "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
2714 );
2715 return;
2716 };
2717 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2718 log::warn!(
2719 "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
2720 );
2721 }
2722 }
2723
2724 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
2731 self.wait_for_all_merges().await;
2732 let segment_ids = self.get_segment_ids().await;
2733
2734 if segment_ids.is_empty() {
2735 log::info!("[reorder] no segments to reorder");
2736 return Ok(());
2737 }
2738
2739 log::info!("[reorder] reordering {} segments", segment_ids.len());
2740
2741 for seg_id in segment_ids {
2742 match self
2743 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
2744 .await
2745 {
2746 Ok(true) => {}
2747 Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
2748 Err(e) => return Err(e),
2749 }
2750 }
2751
2752 log::info!("[reorder] all segments reordered");
2753 Ok(())
2754 }
2755
2756 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
2761 self.unreordered_segments()
2762 .await
2763 .into_iter()
2764 .map(|(id, _)| id)
2765 .collect()
2766 }
2767
2768 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
2771 let quarantined = self.quarantined_segments.lock().clone();
2772 let paused = self.paused_reorder_segments();
2773 let st = self.state.lock().await;
2774 let active_ids = self.active_operations.snapshot();
2775 st.metadata
2776 .segment_metas
2777 .iter()
2778 .filter(|(id, info)| {
2779 !info.reordered
2780 && !active_ids.contains(*id)
2781 && !quarantined.contains(*id)
2782 && !paused.contains(*id)
2783 })
2784 .map(|(id, info)| (id.clone(), info.num_docs))
2785 .collect()
2786 }
2787
2788 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
2792 self.unconverged_segments_below(u32::MAX)
2793 .await
2794 .into_iter()
2795 .map(|(id, docs, _)| (id, docs))
2796 .collect()
2797 }
2798
2799 pub async fn unconverged_segments_below(
2802 &self,
2803 max_unconverged_passes: u32,
2804 ) -> Vec<(String, u32, u32)> {
2805 let quarantined = self.quarantined_segments.lock().clone();
2806 let paused = self.paused_reorder_segments();
2807 let st = self.state.lock().await;
2808 let active_ids = self.active_operations.snapshot();
2809 st.metadata
2810 .segment_metas
2811 .iter()
2812 .filter(|(id, info)| {
2813 info.reordered
2814 && !info.bp_converged
2815 && info.bp_unconverged_passes < max_unconverged_passes
2816 && !active_ids.contains(*id)
2817 && !quarantined.contains(*id)
2818 && !paused.contains(*id)
2819 })
2820 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
2821 .collect()
2822 }
2823
2824 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
2834 let st = self.state.lock().await;
2835 let deepening = ids.iter().any(|id| {
2836 st.metadata
2837 .segment_metas
2838 .get(id)
2839 .is_some_and(|info| info.reordered && !info.bp_converged)
2840 });
2841 drop(st);
2842 if deepening {
2843 log::info!(
2844 "[reorder] source segment(s) unconverged — forcing record-level BP (deepening pass)",
2845 );
2846 crate::segment::reorder::BpGranularity::Records
2847 } else {
2848 crate::segment::reorder::BpGranularity::Auto
2849 }
2850 }
2851
2852 pub async fn reorder_single_segment(
2857 self: &Arc<Self>,
2858 seg_id: &str,
2859 rayon_pool: Option<Arc<rayon::ThreadPool>>,
2860 bp_budget: crate::segment::BpBudget,
2861 ) -> Result<bool> {
2862 let source_id = SegmentId::from_hex(seg_id)
2863 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
2864 if self.quarantined_segments.lock().contains(seg_id) {
2865 return Err(Error::Corruption(format!(
2866 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
2867 seg_id
2868 )));
2869 }
2870
2871 let reorder_gate = Arc::clone(&self.reorder_permits);
2876 let _reorder_permit = tokio::select! {
2877 biased;
2878 () = self.active_operations.wait_for_shutdown() => {
2879 return Err(Error::IndexClosed);
2880 }
2881 permit = reorder_gate.acquire(ReorderPriority::Background) => {
2882 permit.map_err(|_| {
2883 Error::Internal("background reorder scheduler is closed".into())
2884 })?
2885 }
2886 };
2887
2888 let output_id = SegmentId::new();
2889 let output_hex = output_id.to_hex();
2890 let source_ids = [seg_id.to_string()];
2891 let granularity = self.merge_granularity(&source_ids).await;
2892
2893 let all_ids = vec![seg_id.to_string(), output_hex];
2899 let (_guard, source_docs) = {
2900 let st = self.state.lock().await;
2901 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
2902 log::info!(
2903 "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
2904 seg_id
2905 );
2906 self.clear_reorder_retry(seg_id);
2907 return Ok(false);
2908 };
2909
2910 match self.active_operations.try_register(all_ids) {
2911 Some(guard) => (guard, source_meta.num_docs),
2912 None if !self.active_operations.is_accepting() => {
2913 return Err(Error::IndexClosed);
2914 }
2915 None => {
2916 log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
2917 return Ok(false);
2918 }
2919 }
2920 };
2921
2922 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
2927 if is_deterministic_source_error(&error) {
2928 self.quarantine_segment(seg_id, &error);
2929 } else if !matches!(&error, Error::IndexClosed) {
2930 self.pause_reorder_retries(seg_id, &error);
2931 }
2932 return Err(error);
2933 }
2934
2935 let mut output_cleanup = self.output_cleanup_guard(output_id);
2936
2937 let reorder_result = crate::segment::reorder::reorder_segment(
2938 self.directory.as_ref(),
2939 &self.schema,
2940 source_id,
2941 output_id,
2942 self.term_cache_blocks,
2943 self.bp_memory_budget_bytes,
2944 bp_budget,
2945 granularity,
2946 rayon_pool,
2947 )
2948 .await;
2949 let (new_id, total_docs, bp_converged) = match reorder_result {
2950 Ok(v) => v,
2951 Err(e) => {
2952 self.delete_output_if_unregistered(output_id, "reorder failure")
2955 .await;
2956 output_cleanup.disarm();
2957 if is_deterministic_source_error(&e) {
2958 self.quarantine_segment(seg_id, &e);
2959 } else if !matches!(&e, Error::IndexClosed) {
2960 self.pause_reorder_retries(seg_id, &e);
2961 }
2962 return Err(e);
2963 }
2964 };
2965
2966 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
2972 if let Err(e) = self
2973 .replace_segments(
2974 &[seg_id.to_string()],
2975 new_id,
2976 total_docs,
2977 ReplacementLayout::Recomputed {
2978 reordered: true,
2979 bp_converged: ladder_converged,
2980 },
2981 )
2982 .await
2983 {
2984 self.delete_output_if_unregistered(output_id, "replacement failure")
2985 .await;
2986 output_cleanup.disarm();
2987 if !matches!(&e, Error::IndexClosed) {
2988 self.pause_reorder_retries(seg_id, &e);
2989 }
2990 return Err(e);
2991 }
2992 output_cleanup.disarm();
2993 self.clear_reorder_retry(seg_id);
2994
2995 Ok(true)
2996 }
2997
2998 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3005 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3006
3007 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3008 for entry in entries {
3009 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3010 continue;
3011 };
3012 let Some(rest) = filename.strip_prefix("seg_") else {
3013 continue;
3014 };
3015 let Some(hex_id) = rest.get(..32) else {
3016 continue;
3017 };
3018 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3019 continue;
3020 }
3021 orphan_files
3022 .entry(hex_id.to_ascii_lowercase())
3023 .or_default()
3024 .push(entry);
3025 }
3026 }
3027
3028 let mut deleted = 0;
3029 for (hex_id, paths) in &orphan_files {
3030 let deletion_guard = {
3035 let st = self.state.lock().await;
3036 if st.metadata.has_segment(hex_id) {
3037 continue;
3038 }
3039 let Some(guard) = self
3040 .active_operations
3041 .try_register(vec![hex_id.to_string()])
3042 else {
3043 continue;
3044 };
3045 if self.tracker.is_deletion_protected(hex_id) {
3046 drop(guard);
3047 continue;
3048 }
3049 guard
3050 };
3051
3052 let results =
3057 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3058 .await;
3059 let removed = results.into_iter().all(|result| match result {
3060 Ok(()) => true,
3061 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3062 Err(error) => {
3063 log::warn!(
3064 "[segment_cleanup] failed sweeping orphan segment {}: {}",
3065 hex_id,
3066 error,
3067 );
3068 false
3069 }
3070 });
3071 drop(deletion_guard);
3074 if removed {
3075 deleted += 1;
3076 log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3077 }
3078 }
3079
3080 Ok(deleted)
3081 }
3082}
3083
3084#[cfg(test)]
3085mod tests {
3086 use super::*;
3087 use std::sync::atomic::{AtomicBool, Ordering};
3088
3089 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3090 let schema = crate::dsl::SchemaBuilder::default().build();
3091 let metadata = IndexMetadata::new(schema.clone());
3092 Arc::new(SegmentManager::new(
3093 Arc::new(crate::directories::RamDirectory::new()),
3094 Arc::new(schema),
3095 metadata,
3096 Box::new(crate::merge::NoMergePolicy),
3097 0,
3098 1,
3099 Arc::new(Semaphore::new(1)),
3100 None,
3101 1024,
3102 Arc::new(ReorderConcurrencyGate::new(1)),
3103 None,
3104 ))
3105 }
3106
3107 #[test]
3108 fn output_cleanup_guard_runs_during_panic_unwind() {
3109 let cleaned = Arc::new(AtomicBool::new(false));
3110 let cleaned_in_callback = Arc::clone(&cleaned);
3111 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3112 cleaned_in_callback.store(true, Ordering::SeqCst);
3113 });
3114
3115 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3116 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3117 panic!("simulated reorder panic");
3118 }));
3119
3120 assert!(result.is_err());
3121 assert!(
3122 cleaned.load(Ordering::SeqCst),
3123 "partial output cleanup must run during unwind"
3124 );
3125 }
3126
3127 #[test]
3128 fn output_cleanup_guard_disarms_after_commit() {
3129 let cleaned = Arc::new(AtomicBool::new(false));
3130 let cleaned_in_callback = Arc::clone(&cleaned);
3131 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3132 cleaned_in_callback.store(true, Ordering::SeqCst);
3133 });
3134
3135 {
3136 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3137 guard.disarm();
3138 }
3139
3140 assert!(!cleaned.load(Ordering::SeqCst));
3141 }
3142
3143 #[test]
3144 fn test_active_operation_guard_releases_ownership() {
3145 let active = Arc::new(ActiveSegmentOperations::new());
3146 {
3147 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3148 let snap = active.snapshot();
3149 assert!(snap.contains("a"));
3150 assert!(snap.contains("b"));
3151 }
3152 assert!(active.snapshot().is_empty());
3153 }
3154
3155 #[test]
3156 fn test_non_overlapping_operations_can_run_concurrently() {
3157 let active = Arc::new(ActiveSegmentOperations::new());
3158 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3159 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3160 let snap = active.snapshot();
3161 assert_eq!(snap.len(), 4);
3162
3163 drop(first);
3164 let snap = active.snapshot();
3165 assert_eq!(snap.len(), 2);
3166 assert!(snap.contains("c"));
3167 assert!(snap.contains("d"));
3168 }
3169
3170 #[test]
3171 fn test_overlapping_operation_is_rejected_until_release() {
3172 let active = Arc::new(ActiveSegmentOperations::new());
3173 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3174 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3175 drop(first);
3176 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3177 }
3178
3179 #[test]
3180 fn test_active_operation_snapshot() {
3181 let active = Arc::new(ActiveSegmentOperations::new());
3182 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3183 let snap = active.snapshot();
3184 assert!(snap.contains("x"));
3185 assert!(snap.contains("y"));
3186 assert!(!snap.contains("z"));
3187 }
3188
3189 #[tokio::test]
3190 async fn operation_barrier_ignores_producers_started_after_snapshot() {
3191 let active = Arc::new(ActiveSegmentOperations::new());
3192 let before_gate = active.try_register(vec!["old".into()]).unwrap();
3193 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3194 assert_eq!(parked_indexing, 0);
3195 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3196
3197 let waiter = {
3198 let active = Arc::clone(&active);
3199 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3200 };
3201 tokio::task::yield_now().await;
3202 assert!(!waiter.is_finished());
3203
3204 drop(before_gate);
3205 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3206 .await
3207 .expect("pre-gate operation barrier was starved by a post-gate producer")
3208 .unwrap();
3209 assert!(active.snapshot().contains("new-flat"));
3210 drop(after_gate);
3211 }
3212
3213 #[tokio::test]
3214 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
3215 let manager = lifecycle_test_manager();
3216 manager
3217 .trained
3218 .store(Some(Arc::new(TrainedVectorStructures {
3219 centroids: rustc_hash::FxHashMap::default(),
3220 binary_quantizers: rustc_hash::FxHashMap::default(),
3221 ..Default::default()
3222 })));
3223
3224 let guard = manager.begin_vector_artifact_update().await.unwrap();
3225 assert!(
3226 manager.trained().is_some(),
3227 "search readers keep the last fully validated generation"
3228 );
3229 assert!(
3230 manager.trained_for_segment_build().is_none(),
3231 "new segment producers must stay flat during an artifact update"
3232 );
3233
3234 let detached_transaction_guard = guard.clone();
3235 drop(guard);
3236 assert!(
3237 manager.trained_for_segment_build().is_none(),
3238 "a detached lifecycle transaction must retain the producer gate after request cancellation"
3239 );
3240 drop(detached_transaction_guard);
3241 assert!(manager.trained_for_segment_build().is_some());
3242 }
3243
3244 #[tokio::test]
3245 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
3246 let manager = lifecycle_test_manager();
3247 let guard = manager.begin_vector_artifact_update().await.unwrap();
3248 assert!(
3249 manager
3250 .active_operations
3251 .try_register(vec!["merge".into()])
3252 .is_none(),
3253 "ordinary merge/reorder work must not change staged sources"
3254 );
3255 let indexing = manager
3256 .active_operations
3257 .try_register_indexing(vec!["fresh".into()])
3258 .expect("indexing remains available in flat mode");
3259 drop(indexing);
3260
3261 drop(guard);
3262 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
3263 assert!(
3264 manager
3265 .active_operations
3266 .try_register(vec!["merge".into()])
3267 .is_some()
3268 );
3269 }
3270
3271 #[tokio::test]
3272 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
3273 let active = Arc::new(ActiveSegmentOperations::new());
3274 let guard = active.try_register(vec!["live".into()]).unwrap();
3275 active.stop_accepting();
3276 assert!(active.try_register(vec!["new".into()]).is_none());
3277
3278 let waiter = {
3279 let active = Arc::clone(&active);
3280 tokio::spawn(async move { active.wait_until_idle().await })
3281 };
3282 tokio::task::yield_now().await;
3283 assert!(!waiter.is_finished());
3284 drop(guard);
3285 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3286 .await
3287 .expect("shutdown waiter missed the final guard notification")
3288 .unwrap();
3289 }
3290
3291 #[tokio::test]
3292 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
3293 let manager = lifecycle_test_manager();
3294 let started = Arc::new(Semaphore::new(0));
3295 let release = Arc::new(Semaphore::new(0));
3296 let completed = Arc::new(AtomicBool::new(false));
3297
3298 let request = {
3299 let manager = Arc::clone(&manager);
3300 let started = Arc::clone(&started);
3301 let release = Arc::clone(&release);
3302 let completed = Arc::clone(&completed);
3303 tokio::spawn(async move {
3304 manager
3305 .run_lifecycle_transaction(async move {
3306 started.add_permits(1);
3307 let _permit = release.acquire().await.unwrap();
3308 completed.store(true, Ordering::Release);
3309 Ok(())
3310 })
3311 .await
3312 })
3313 };
3314
3315 let _started = started.acquire().await.unwrap();
3316 request.abort();
3317 assert!(request.await.unwrap_err().is_cancelled());
3318 release.add_permits(1);
3319
3320 manager.begin_shutdown();
3321 tokio::time::timeout(
3322 std::time::Duration::from_secs(1),
3323 manager.wait_for_shutdown(),
3324 )
3325 .await
3326 .expect("shutdown did not drain detached lifecycle transaction");
3327 assert!(completed.load(Ordering::Acquire));
3328 }
3329
3330 #[tokio::test]
3331 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
3332 let manager = lifecycle_test_manager();
3333 {
3334 let mut state = manager.state.lock().await;
3335 state.metadata.add_segment_meta(
3336 "eligible".into(),
3337 SegmentMetaInfo {
3338 num_docs: 10,
3339 ancestors: Vec::new(),
3340 generation: 1,
3341 reordered: true,
3342 bp_converged: false,
3343 bp_unconverged_passes: 2,
3344 },
3345 );
3346 state.metadata.add_segment_meta(
3347 "at-limit".into(),
3348 SegmentMetaInfo {
3349 num_docs: 20,
3350 ancestors: Vec::new(),
3351 generation: 1,
3352 reordered: true,
3353 bp_converged: false,
3354 bp_unconverged_passes: 3,
3355 },
3356 );
3357 state.metadata.add_segment_meta(
3358 "converged".into(),
3359 SegmentMetaInfo {
3360 num_docs: 30,
3361 ancestors: Vec::new(),
3362 generation: 1,
3363 reordered: true,
3364 bp_converged: true,
3365 bp_unconverged_passes: 0,
3366 },
3367 );
3368 state.metadata.add_segment("fresh".into(), 40);
3369 }
3370
3371 assert_eq!(
3372 manager.unconverged_segments_below(3).await,
3373 vec![("eligible".into(), 10, 2)]
3374 );
3375 assert!(manager.unconverged_segments_below(0).await.is_empty());
3376 }
3377
3378 #[test]
3379 fn merge_retry_backoff_is_exponential_and_capped() {
3380 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
3381 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
3382 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
3383 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
3384 }
3385
3386 #[test]
3387 fn only_deterministic_source_errors_are_quarantined() {
3388 assert!(is_deterministic_source_error(&Error::Corruption(
3389 "bad footer".into()
3390 )));
3391 assert!(is_deterministic_source_error(&Error::Io(
3392 std::io::Error::from(std::io::ErrorKind::NotFound)
3393 )));
3394 assert!(!is_deterministic_source_error(&Error::Io(
3395 std::io::Error::from(std::io::ErrorKind::TimedOut)
3396 )));
3397 assert!(!is_deterministic_source_error(&Error::Io(
3398 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
3399 )));
3400 }
3401
3402 #[test]
3403 fn transient_reorder_failure_is_backed_off_until_cleared() {
3404 let manager = lifecycle_test_manager();
3405 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
3406 assert!(manager.paused_reorder_segments().contains("source"));
3407 manager.clear_reorder_retry("source");
3408 assert!(!manager.paused_reorder_segments().contains("source"));
3409 }
3410
3411 #[derive(Default)]
3414 struct FailingExistsDirectory(crate::directories::RamDirectory);
3415
3416 #[async_trait::async_trait]
3417 impl crate::directories::Directory for FailingExistsDirectory {
3418 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
3419 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
3420 }
3421
3422 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
3423 self.0.file_size(path).await
3424 }
3425
3426 async fn open_read(
3427 &self,
3428 path: &std::path::Path,
3429 ) -> std::io::Result<crate::directories::FileHandle> {
3430 self.0.open_read(path).await
3431 }
3432
3433 async fn read_range(
3434 &self,
3435 path: &std::path::Path,
3436 range: std::ops::Range<u64>,
3437 ) -> std::io::Result<crate::directories::OwnedBytes> {
3438 self.0.read_range(path, range).await
3439 }
3440
3441 async fn list_files(
3442 &self,
3443 prefix: &std::path::Path,
3444 ) -> std::io::Result<Vec<std::path::PathBuf>> {
3445 self.0.list_files(prefix).await
3446 }
3447
3448 async fn open_lazy(
3449 &self,
3450 path: &std::path::Path,
3451 ) -> std::io::Result<crate::directories::FileHandle> {
3452 self.0.open_lazy(path).await
3453 }
3454 }
3455
3456 #[async_trait::async_trait]
3457 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
3458 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
3459 self.0.write(path, data).await
3460 }
3461
3462 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
3463 self.0.delete(path).await
3464 }
3465
3466 async fn rename(
3467 &self,
3468 from: &std::path::Path,
3469 to: &std::path::Path,
3470 ) -> std::io::Result<()> {
3471 self.0.rename(from, to).await
3472 }
3473
3474 async fn sync(&self) -> std::io::Result<()> {
3475 self.0.sync().await
3476 }
3477
3478 async fn streaming_writer(
3479 &self,
3480 path: &std::path::Path,
3481 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
3482 self.0.streaming_writer(path).await
3483 }
3484 }
3485
3486 #[derive(Debug, Clone)]
3487 struct MergeEverythingPolicy;
3488
3489 impl MergePolicy for MergeEverythingPolicy {
3490 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
3491 if segments.len() < 2 {
3492 return Vec::new();
3493 }
3494 vec![crate::merge::MergeCandidate {
3495 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
3496 }]
3497 }
3498
3499 fn clone_box(&self) -> Box<dyn MergePolicy> {
3500 Box::new(self.clone())
3501 }
3502 }
3503
3504 #[tokio::test]
3505 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
3506 let manager = lifecycle_test_manager();
3507 let parked_indexing = manager
3512 .protect_new_segment("00000000000000000000000000000abc".into())
3513 .unwrap();
3514
3515 let error = tokio::time::timeout(
3516 std::time::Duration::from_secs(2),
3517 manager.begin_vector_artifact_update(),
3518 )
3519 .await
3520 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
3521 .err()
3522 .expect("an old-generation prepared segment must block artifact replacement")
3523 .to_string();
3524 assert!(error.contains("built but uncommitted"), "{error}");
3525 assert!(
3526 !manager.vector_artifact_update.load(Ordering::Acquire),
3527 "a rejected update must release the producer gate"
3528 );
3529
3530 drop(parked_indexing);
3531
3532 let guard = manager
3533 .begin_vector_artifact_update()
3534 .await
3535 .expect("artifact update should succeed after the pending generation is resolved");
3536 drop(guard);
3537 }
3538
3539 #[tokio::test]
3540 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
3541 let manager = lifecycle_test_manager();
3542 let merge_like = manager
3543 .active_operations
3544 .try_register(vec!["merge-source".into()])
3545 .unwrap();
3546
3547 let waiter = {
3548 let manager = Arc::clone(&manager);
3549 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
3550 };
3551 for _ in 0..8 {
3552 tokio::task::yield_now().await;
3553 }
3554 assert!(
3555 !waiter.is_finished(),
3556 "artifact update must drain merge/reorder producers that may hold the previous generation"
3557 );
3558
3559 drop(merge_like);
3560 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3561 .await
3562 .expect("artifact update missed the lifecycle guard release")
3563 .unwrap()
3564 .unwrap();
3565 }
3566
3567 #[tokio::test]
3568 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
3569 let manager = lifecycle_test_manager();
3570 let release = Arc::new(Semaphore::new(0));
3571 let merge_task = {
3572 let release = Arc::clone(&release);
3573 tokio::spawn(async move {
3574 let _permit = release.acquire().await.unwrap();
3575 })
3576 };
3577 manager.merge_handles.lock().push(merge_task);
3578
3579 let waiter = {
3580 let manager = Arc::clone(&manager);
3581 tokio::spawn(async move { manager.wait_for_all_merges().await })
3582 };
3583 for _ in 0..8 {
3584 tokio::task::yield_now().await;
3585 }
3586 assert!(!waiter.is_finished());
3587 waiter.abort();
3590 let join_error = waiter.await.unwrap_err();
3591 assert!(join_error.is_cancelled());
3592
3593 assert!(
3594 !manager.merge_handles.lock().is_empty(),
3595 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
3596 );
3597
3598 release.add_permits(1);
3600 tokio::time::timeout(
3601 std::time::Duration::from_secs(1),
3602 manager.wait_for_all_merges(),
3603 )
3604 .await
3605 .expect("subsequent drain missed the reinserted merge handle");
3606 assert!(manager.merge_handles.lock().is_empty());
3607 }
3608
3609 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3610 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
3611 let manager = lifecycle_test_manager();
3612 {
3613 let mut state = manager.state.lock().await;
3614 state
3615 .metadata
3616 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
3617 state
3618 .metadata
3619 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
3620 }
3621 let reorder_like = manager
3624 .active_operations
3625 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
3626 .unwrap();
3627
3628 let force_merge = {
3629 let manager = Arc::clone(&manager);
3630 tokio::spawn(async move { manager.force_merge().await })
3631 };
3632
3633 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
3634 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
3635 assert!(
3636 retries >= 1,
3637 "force_merge never observed the conflicting owner (retries={retries})"
3638 );
3639 assert!(
3640 retries < 20,
3641 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
3642 );
3643
3644 drop(reorder_like);
3645 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
3648 .await
3649 .expect("force_merge kept spinning after the conflicting owner released")
3650 .unwrap();
3651 assert!(result.is_err());
3652 }
3653
3654 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3655 async fn force_merge_routes_around_segments_held_by_reorder() {
3656 let manager = lifecycle_test_manager();
3657 {
3658 let mut state = manager.state.lock().await;
3659 state
3660 .metadata
3661 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
3662 state
3663 .metadata
3664 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
3665 state
3666 .metadata
3667 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
3668 }
3669 let _reorder_like = manager
3672 .active_operations
3673 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
3674 .unwrap();
3675
3676 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
3683 let manager = Arc::clone(&manager);
3684 async move { manager.force_merge().await }
3685 })
3686 .await
3687 .expect("force_merge livelocked on a segment held by an active reorder");
3688 assert!(result.is_err(), "fake segment files must fail the merge");
3689
3690 assert_eq!(
3691 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
3692 0,
3693 "batch built from the ownership snapshot must not collide with the held segment"
3694 );
3695 }
3696
3697 #[tokio::test]
3698 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
3699 let schema = crate::dsl::SchemaBuilder::default().build();
3700 let mut metadata = IndexMetadata::new(schema.clone());
3701 metadata.add_segment("00000000000000000000000000000001".into(), 10);
3702 metadata.add_segment("00000000000000000000000000000002".into(), 10);
3703 let manager = Arc::new(SegmentManager::new(
3704 Arc::new(FailingExistsDirectory::default()),
3705 Arc::new(schema),
3706 metadata,
3707 Box::new(MergeEverythingPolicy),
3708 0,
3709 1,
3710 Arc::new(Semaphore::new(1)),
3711 None,
3712 1024,
3713 Arc::new(ReorderConcurrencyGate::new(1)),
3714 None,
3715 ));
3716
3717 manager.maybe_merge().await;
3720
3721 tokio::time::timeout(
3722 std::time::Duration::from_secs(5),
3723 manager.wait_for_all_merges(),
3724 )
3725 .await
3726 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
3727 assert!(
3728 manager.merge_retry_is_paused(),
3729 "the failed merge should have armed the retry backoff"
3730 );
3731
3732 manager.begin_shutdown();
3734 tokio::time::timeout(
3735 std::time::Duration::from_secs(5),
3736 manager.wait_for_shutdown(),
3737 )
3738 .await
3739 .expect("shutdown did not drain the merge retry wakeup task");
3740 }
3741}