Skip to main content

hermes_core/merge/
segment_manager.rs

1//! Segment manager — coordinates segment commit, background merging, and trained structures.
2//!
3//! Architecture:
4//! - **Single mutation queue**: All metadata mutations serialize through `tokio::sync::Mutex<ManagerState>`.
5//! - **Active-operation ownership**: Every segment that is being built, merged,
6//!   or reordered is registered before its first file is written and remains
7//!   registered until it is either published in metadata or abandoned.
8//! - **Concurrent merges**: Multiple non-overlapping merges can run in parallel.
9//!   New merges are rejected only if they share segments with an active operation.
10//! - **Auto-trigger**: Each completed merge re-evaluates the merge policy and spawns
11//!   new merges if eligible (cascading merges for higher tiers).
12//! - **ArcSwap for trained**: Lock-free reads of trained vector structures.
13//!
14//! # Segment lifecycle invariant
15//!
16//! Every on-disk `seg_*` ID must be protected by at least one of these owners:
17//!
18//! 1. `state.metadata` while the segment is live and searchable;
19//! 2. `active_operations` while an indexing/merge/reorder task owns it; or
20//! 3. `tracker` while a retired segment is still visible to a reader or its
21//!    filesystem deletion is scheduled.
22//!
23//! An ID with no owner is an orphan and may be swept. Transitions are ordered
24//! so the new owner is installed before the old owner is released.
25//!
26//! # Locking model (deadlock-free by construction)
27//!
28//! ```text
29//! Lock ordering (acquire in this order):
30//!   1. state               — tokio::sync::Mutex, held for mutations + disk I/O
31//!   2. active_operations   — parking_lot::Mutex (sync), sub-μs hold, RAII guard
32//!   3. tracker.inner       — parking_lot::Mutex (sync), sub-μs hold
33//!
34//! Independent bookkeeping (never held with `state`):
35//!   trained                — arc_swap::ArcSwapOption, lock-free
36//!   vector_artifact_update — atomic producer gate, held across manual training
37//!   merge_handles          — parking_lot::Mutex, synchronous short hold
38//!   lifecycle_handles      — parking_lot::Mutex, synchronous short hold
39//!   merge/reorder permits  — tokio semaphores shared by configuration
40//! ```
41//!
42//! **Rule:** Never hold a sync lock while `.await`-ing.
43
44use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwap;
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    PublishedIndexGeneration, SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot,
58    SegmentTracker, TrainedVectorStructures,
59};
60#[cfg(feature = "native")]
61use crate::segment::{SegmentMerger, SegmentReader};
62
63use super::{MergePolicy, SegmentInfo};
64
65const FORCE_MERGE_MAX_FAN_IN: usize = 64;
66
67#[derive(Debug)]
68struct ForceMergeGroup {
69    segments: Vec<(String, u32)>,
70    total_docs: u64,
71}
72
73/// Deterministically pack segments into near-minimal final outputs.
74///
75/// Best-fit decreasing avoids the pathological smallest-first behavior where,
76/// for example, `4 + 4` is merged before two `6 + 4` pairs under a 10-doc
77/// limit. That old order both left excess segments and made its intermediate
78/// outputs eligible for another full BP pass. Bin packing is NP-hard; BFD is a
79/// bounded, deterministic approximation that produces maximal groups (no two
80/// output groups can still fit together).
81fn plan_force_merge_groups(
82    mut segments: Vec<(String, u32)>,
83    max_docs: u64,
84) -> Vec<ForceMergeGroup> {
85    segments.sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
86        right_docs
87            .cmp(left_docs)
88            .then_with(|| left_id.cmp(right_id))
89    });
90
91    let mut groups: Vec<ForceMergeGroup> = Vec::new();
92    for segment in segments {
93        let docs = u64::from(segment.1);
94        let best_group = groups
95            .iter()
96            .enumerate()
97            .filter_map(|(index, group)| {
98                group
99                    .total_docs
100                    .checked_add(docs)
101                    .filter(|&total| total <= max_docs)
102                    .map(|_| (index, group.total_docs))
103            })
104            .max_by_key(|&(index, used)| (used, std::cmp::Reverse(index)))
105            .map(|(index, _)| index);
106
107        if let Some(index) = best_group {
108            groups[index].total_docs += docs;
109            groups[index].segments.push(segment);
110        } else {
111            groups.push(ForceMergeGroup {
112                segments: vec![segment],
113                total_docs: docs,
114            });
115        }
116    }
117
118    // Preserve the old small-to-large source order inside each output. Besides
119    // stable document ordering, BMP block-copy compression depends on source
120    // alignment at group boundaries.
121    for group in &mut groups {
122        group
123            .segments
124            .sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
125                left_docs
126                    .cmp(right_docs)
127                    .then_with(|| left_id.cmp(right_id))
128            });
129    }
130
131    // Smallest outputs first limit transient disk headroom. Tie-break on the
132    // first (deterministically ordered) segment ID for reproducible plans.
133    groups.sort_unstable_by(|left, right| {
134        left.total_docs
135            .cmp(&right.total_docs)
136            .then_with(|| left.segments[0].0.cmp(&right.segments[0].0))
137    });
138    groups
139}
140
141fn force_merge_output_count(source_count: usize) -> usize {
142    if source_count < 2 {
143        return 0;
144    }
145    (source_count - 1).div_ceil(FORCE_MERGE_MAX_FAN_IN - 1)
146}
147
148#[derive(Debug)]
149struct ForceMergeStep {
150    /// Source or earlier-step node IDs, in final document order.
151    inputs: Vec<usize>,
152}
153
154#[derive(Debug)]
155struct ForceMergeHierarchy {
156    steps: Vec<ForceMergeStep>,
157    root: usize,
158}
159
160/// Build a shallow ordered 64-ary reduction tree with the minimum possible
161/// number of merge outputs. All internal nodes have full fan-in except one
162/// deepest partial node; distributing internal nodes evenly keeps source
163/// rewrite depth balanced without changing concatenation order.
164fn plan_force_merge_hierarchy(source_count: usize) -> ForceMergeHierarchy {
165    debug_assert!(source_count >= 2);
166    let internal_count = force_merge_output_count(source_count);
167    let max_leaves = 1usize
168        .checked_add(internal_count.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
169        .expect("force-merge hierarchy size exceeds usize");
170    let deficit = max_leaves - source_count;
171    debug_assert!(deficit < FORCE_MERGE_MAX_FAN_IN - 1);
172
173    fn build(
174        leaf_start: usize,
175        leaf_count: usize,
176        internal_count: usize,
177        deficit: usize,
178        source_count: usize,
179        steps: &mut Vec<ForceMergeStep>,
180    ) -> usize {
181        debug_assert!(internal_count > 0);
182        if internal_count == 1 {
183            let arity = FORCE_MERGE_MAX_FAN_IN - deficit;
184            debug_assert_eq!(leaf_count, arity);
185            debug_assert!((2..=FORCE_MERGE_MAX_FAN_IN).contains(&arity));
186            let output = source_count + steps.len();
187            steps.push(ForceMergeStep {
188                inputs: (leaf_start..leaf_start + arity).collect(),
189            });
190            return output;
191        }
192
193        // Keep this node full and place the one partial arity, if any, in the
194        // deepest/largest child. This minimizes bytes rewritten versus making
195        // the root partial (65 inputs become 2 + 63 leaves, then one final
196        // merge, rather than rewriting a 64-input prefix).
197        let child_internal_total = internal_count - 1;
198        let base = child_internal_total / FORCE_MERGE_MAX_FAN_IN;
199        let extra = child_internal_total % FORCE_MERGE_MAX_FAN_IN;
200        let mut child_internal = vec![base; FORCE_MERGE_MAX_FAN_IN];
201        for count in &mut child_internal[..extra] {
202            *count += 1;
203        }
204        let partial_child = (deficit > 0).then(|| {
205            child_internal
206                .iter()
207                .position(|&count| count > 0)
208                .expect("a non-root partial node requires an internal child")
209        });
210
211        let mut cursor = leaf_start;
212        let mut inputs = Vec::with_capacity(FORCE_MERGE_MAX_FAN_IN);
213        for (child, &child_internals) in child_internal.iter().enumerate() {
214            if child_internals == 0 {
215                inputs.push(cursor);
216                cursor += 1;
217                continue;
218            }
219            let child_deficit = usize::from(partial_child == Some(child)) * deficit;
220            let child_leaves = 1usize
221                .checked_add(child_internals.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
222                .and_then(|maximum| maximum.checked_sub(child_deficit))
223                .expect("force-merge child size exceeds usize");
224            inputs.push(build(
225                cursor,
226                child_leaves,
227                child_internals,
228                child_deficit,
229                source_count,
230                steps,
231            ));
232            cursor += child_leaves;
233        }
234        debug_assert_eq!(cursor, leaf_start + leaf_count);
235        let output = source_count + steps.len();
236        steps.push(ForceMergeStep { inputs });
237        output
238    }
239
240    let mut steps = Vec::with_capacity(internal_count);
241    let root = build(
242        0,
243        source_count,
244        internal_count,
245        deficit,
246        source_count,
247        &mut steps,
248    );
249    debug_assert_eq!(steps.len(), internal_count);
250    ForceMergeHierarchy { steps, root }
251}
252
253// ============================================================================
254// RAII active-operation tracking
255// ============================================================================
256
257/// Tracks every segment ID owned by an in-flight lifecycle operation.
258///
259/// Merge/reorder guards include both sources and output, providing mutual
260/// exclusion as well as orphan-sweep protection. Indexing guards contain the
261/// new output only and live from before the first write through commit/abort.
262struct ActiveOperationState {
263    segment_ids: HashSet<String>,
264    operation_tokens: HashSet<u64>,
265    /// Subset of `operation_tokens` owned by indexing producers. Their guards
266    /// travel with built-but-uncommitted `PreparedSegment`s and are released
267    /// only by a later commit/abort, so drain barriers must not wait on them:
268    /// the commit that would release them can be blocked on the barrier's own
269    /// caller (writer write lock / `&mut self`).
270    indexing_tokens: HashSet<u64>,
271    next_operation_token: u64,
272    accepting: bool,
273    /// Retraining stages a complete replacement segment generation. Ordinary
274    /// merge/reorder work is paused so its source set cannot change midway;
275    /// indexing producers remain allowed and deliberately emit flat vectors.
276    non_indexing_paused: bool,
277}
278
279struct ActiveSegmentOperations {
280    inner: parking_lot::Mutex<ActiveOperationState>,
281    idle: Notify,
282    shutdown: Notify,
283    shutdown_requested: Arc<AtomicBool>,
284    /// Owning index, so lifecycle decisions are attributable when several
285    /// indexes register and defer operations concurrently.
286    index_label: Arc<str>,
287}
288
289impl ActiveSegmentOperations {
290    fn new(index_label: Arc<str>) -> Self {
291        Self {
292            inner: parking_lot::Mutex::new(ActiveOperationState {
293                segment_ids: HashSet::new(),
294                operation_tokens: HashSet::new(),
295                indexing_tokens: HashSet::new(),
296                next_operation_token: 0,
297                accepting: true,
298                non_indexing_paused: false,
299            }),
300            idle: Notify::new(),
301            shutdown: Notify::new(),
302            shutdown_requested: Arc::new(AtomicBool::new(false)),
303            index_label,
304        }
305    }
306
307    /// Try to claim IDs for a self-draining lifecycle operation (merge,
308    /// reorder, cleanup). Returns a guard on success, `None` if any requested
309    /// ID is already owned by another active operation.
310    fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
311        self.try_register_kind(segment_ids, false, false)
312    }
313
314    /// Try to claim IDs for an indexing producer whose guard is held until
315    /// metadata publication (commit) rather than task completion.
316    fn try_register_indexing(
317        self: &Arc<Self>,
318        segment_ids: Vec<String>,
319    ) -> Option<SegmentOperationGuard> {
320        self.try_register_kind(segment_ids, true, false)
321    }
322
323    /// Claim source/output IDs for the exclusive vector-generation updater,
324    /// which is the only non-indexing producer allowed through its own pause.
325    fn try_register_vector_update(
326        self: &Arc<Self>,
327        segment_ids: Vec<String>,
328    ) -> Option<SegmentOperationGuard> {
329        self.try_register_kind(segment_ids, false, true)
330    }
331
332    fn try_register_kind(
333        self: &Arc<Self>,
334        segment_ids: Vec<String>,
335        indexing: bool,
336        vector_update: bool,
337    ) -> Option<SegmentOperationGuard> {
338        let mut inner = self.inner.lock();
339        if !inner.accepting {
340            log::debug!(
341                "[segment_lifecycle] index={} rejected operation during shutdown",
342                self.index_label
343            );
344            return None;
345        }
346        if !indexing && !vector_update && inner.non_indexing_paused {
347            log::debug!(
348                "[segment_lifecycle] index={} deferred operation during dense vector retraining",
349                self.index_label
350            );
351            return None;
352        }
353        // Check for overlap with any active lifecycle operation.
354        for id in &segment_ids {
355            if inner.segment_ids.contains(id) {
356                log::debug!(
357                    "[segment_lifecycle] index={} rejected: {} overlaps with an active operation ({} active IDs)",
358                    self.index_label,
359                    id,
360                    inner.segment_ids.len()
361                );
362                return None;
363            }
364        }
365        log::debug!(
366            "[segment_lifecycle] index={} registered {} IDs (total active: {})",
367            self.index_label,
368            segment_ids.len(),
369            inner.segment_ids.len() + segment_ids.len()
370        );
371        let operation_token = inner.next_operation_token;
372        let next_operation_token = operation_token.checked_add(1)?;
373        for id in &segment_ids {
374            inner.segment_ids.insert(id.clone());
375        }
376        inner.next_operation_token = next_operation_token;
377        inner.operation_tokens.insert(operation_token);
378        if indexing {
379            inner.indexing_tokens.insert(operation_token);
380        }
381        Some(SegmentOperationGuard {
382            active_operations: Arc::clone(self),
383            segment_ids,
384            operation_token,
385        })
386    }
387
388    /// Snapshot of all IDs owned by active operations.
389    fn snapshot(&self) -> HashSet<String> {
390        self.inner.lock().segment_ids.clone()
391    }
392
393    /// Exact identities of self-draining operations (merge/reorder/cleanup)
394    /// active at one instant, plus the number of indexing tokens excluded.
395    /// Unlike segment IDs, tokens cannot be reused by a later retry, so an
396    /// artifact-update barrier can drain only pre-gate producers without being
397    /// starved by new flat producers.
398    ///
399    /// Indexing tokens are deliberately excluded: their guards are parked in
400    /// built-but-uncommitted `PreparedSegment`s and only a later commit — which
401    /// may be blocked on the barrier's caller — releases them, so waiting on
402    /// them deadlocks (see `begin_vector_artifact_update`).
403    fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
404        let inner = self.inner.lock();
405        let tokens = inner
406            .operation_tokens
407            .difference(&inner.indexing_tokens)
408            .copied()
409            .collect();
410        (tokens, inner.indexing_tokens.len())
411    }
412
413    /// Atomically prevent new lifecycle work from starting. Existing guards
414    /// remain valid and can be drained with [`Self::wait_until_idle`].
415    fn stop_accepting(&self) {
416        self.shutdown_requested.store(true, Ordering::Release);
417        let mut inner = self.inner.lock();
418        inner.accepting = false;
419        self.shutdown.notify_waiters();
420        if inner.segment_ids.is_empty() {
421            self.idle.notify_waiters();
422        }
423    }
424
425    fn pause_non_indexing(&self) {
426        self.inner.lock().non_indexing_paused = true;
427    }
428
429    fn resume_non_indexing(&self) {
430        self.inner.lock().non_indexing_paused = false;
431        self.idle.notify_waiters();
432    }
433
434    fn is_accepting(&self) -> bool {
435        self.inner.lock().accepting
436    }
437
438    fn cancellation_flag(&self) -> Arc<AtomicBool> {
439        Arc::clone(&self.shutdown_requested)
440    }
441
442    /// Wait until every operation that started before shutdown has released
443    /// its ownership. Register/check and notification are ordered to avoid a
444    /// missed wakeup between observing a non-empty set and awaiting.
445    async fn wait_until_idle(&self) {
446        loop {
447            let notified = self.idle.notified();
448            if self.inner.lock().segment_ids.is_empty() {
449                return;
450            }
451            notified.await;
452        }
453    }
454
455    async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
456        while !operations.is_empty() {
457            let notified = self.idle.notified();
458            if self.inner.lock().operation_tokens.is_disjoint(operations) {
459                return;
460            }
461            notified.await;
462        }
463    }
464
465    /// Resolve when shutdown starts, without missing a notification between
466    /// checking the state and registering the waiter.
467    async fn wait_for_shutdown(&self) {
468        loop {
469            let notified = self.shutdown.notified();
470            if !self.inner.lock().accepting {
471                return;
472            }
473            notified.await;
474        }
475    }
476}
477
478/// RAII ownership of segment IDs used by an active lifecycle operation.
479/// Dropping on success, error, cancellation, or panic makes abandoned outputs
480/// eligible for sweeping automatically.
481pub(crate) struct SegmentOperationGuard {
482    active_operations: Arc<ActiveSegmentOperations>,
483    segment_ids: Vec<String>,
484    operation_token: u64,
485}
486
487impl Drop for SegmentOperationGuard {
488    fn drop(&mut self) {
489        let mut inner = self.active_operations.inner.lock();
490        for id in &self.segment_ids {
491            inner.segment_ids.remove(id);
492        }
493        inner.operation_tokens.remove(&self.operation_token);
494        inner.indexing_tokens.remove(&self.operation_token);
495        // Token barriers need notification on every completion, not only the
496        // transition to complete global idleness.
497        self.active_operations.idle.notify_waiters();
498        if inner.segment_ids.is_empty() {
499            debug_assert!(inner.operation_tokens.is_empty());
500        }
501    }
502}
503
504/// Exclusive gate for an index-level trained-vector artifact update.
505///
506/// Segment producers consult this gate before capturing the current trained
507/// structures. Once the gate is raised, new producers deliberately emit flat
508/// vector data; waiting for already-active producers to drain then guarantees
509/// that the committed source set stays stable while replacements are staged.
510struct VectorArtifactUpdateLease {
511    updating: Arc<AtomicBool>,
512    active_operations: Arc<ActiveSegmentOperations>,
513}
514
515impl Drop for VectorArtifactUpdateLease {
516    fn drop(&mut self) {
517        self.updating.store(false, Ordering::Release);
518        self.active_operations.resume_non_indexing();
519    }
520}
521
522#[derive(Clone)]
523pub(crate) struct VectorArtifactUpdateGuard {
524    _lease: Arc<VectorArtifactUpdateLease>,
525}
526
527/// Merge-time/manual BP pools are shared by every index in this process.
528/// A pool per `SegmentManager` multiplied a 96-core host into two 48-thread
529/// merge pools plus the optimizer pool (200+ process threads in production).
530static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
531
532const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
533const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
534
535#[derive(Default)]
536struct MergeRetryState {
537    retry_after: Option<std::time::Instant>,
538    consecutive_failures: u32,
539}
540
541fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
542    let shift = consecutive_failures.saturating_sub(1).min(16);
543    MERGE_RETRY_BASE_DELAY
544        .checked_mul(1u32 << shift)
545        .unwrap_or(MERGE_RETRY_MAX_DELAY)
546        .min(MERGE_RETRY_MAX_DELAY)
547}
548
549/// Merge JoinHandles taken out of the shared list for draining.
550///
551/// Drain futures are awaited inline by RPC handlers (force_merge/reorder) and
552/// can be dropped at any await when a client disconnects. Handles are awaited
553/// through this guard and removed only after completion, so a cancelled drain
554/// returns every un-awaited (and possibly still-running) merge to the shared
555/// list instead of silently detaching it from shutdown, abort, and
556/// force-merge tracking.
557struct DrainedMergeHandles<'a> {
558    shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
559    drained: Vec<JoinHandle<()>>,
560}
561
562impl<'a> DrainedMergeHandles<'a> {
563    fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
564        let drained = std::mem::take(&mut *shared.lock());
565        Self { shared, drained }
566    }
567
568    fn is_empty(&self) -> bool {
569        self.drained.is_empty()
570    }
571
572    /// Await the next handle. It stays owned by this guard while being polled
573    /// and is discarded only once it has completed, so cancellation at the
574    /// await reinserts it via `Drop`.
575    async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
576        let handle = self.drained.last_mut()?;
577        let result = handle.await;
578        self.drained.pop();
579        Some(result)
580    }
581}
582
583impl Drop for DrainedMergeHandles<'_> {
584    fn drop(&mut self) {
585        if !self.drained.is_empty() {
586            self.shared.lock().append(&mut self.drained);
587        }
588    }
589}
590
591/// Spawn and register auxiliary lifecycle work as one synchronous operation.
592///
593/// Registering *after* `spawn` left a small deletion race: shutdown could
594/// observe an empty handle list while the newly spawned filesystem task was
595/// already running. Holding the handle-list mutex across `Handle::spawn`
596/// makes task creation visible to the drain before either side can proceed.
597fn try_spawn_lifecycle<F>(
598    handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
599    runtime: &tokio::runtime::Handle,
600    future: F,
601) -> bool
602where
603    F: std::future::Future<Output = ()> + Send + 'static,
604{
605    std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
606        let mut handles = handles.lock();
607        handles.retain(|handle| !handle.is_finished());
608        handles.push(runtime.spawn(future));
609    }))
610    .is_ok()
611}
612
613/// Deletes an uncommitted merge/reorder output if its task unwinds.
614///
615/// Normal `Result::Err` paths delete outputs synchronously so callers observe
616/// a clean directory before returning. This guard covers the path those
617/// branches cannot: a panic after output files have been created. The cleanup
618/// callback re-checks metadata before deleting, so a panic after a successful
619/// metadata commit cannot remove a live segment.
620struct OutputCleanupGuard {
621    segment_id: SegmentId,
622    cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
623}
624
625impl OutputCleanupGuard {
626    fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
627        Self {
628            segment_id,
629            cleanup: Some(cleanup),
630        }
631    }
632
633    fn disarm(&mut self) {
634        self.cleanup = None;
635    }
636}
637
638impl Drop for OutputCleanupGuard {
639    fn drop(&mut self) {
640        if let Some(cleanup) = self.cleanup.take() {
641            cleanup(self.segment_id);
642        }
643    }
644}
645
646/// All mutable state behind the single async Mutex.
647struct ManagerState {
648    metadata: IndexMetadata,
649    merge_policy: Box<dyn MergePolicy>,
650}
651
652type ReplacementRefresh = Arc<
653    dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
654        + Send
655        + Sync,
656>;
657
658/// Reconcile a long-lived consumer after a durable segment replacement.
659/// This runs inside the same lifecycle-owned task as metadata publication, so
660/// cancelling the caller cannot leave the replacement published but unannounced.
661async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>, index_label: &str) {
662    let Some(refresh) = refresh else {
663        return;
664    };
665    let mut last_error = None;
666    for attempt in 0..3 {
667        match refresh().await {
668            Ok(()) => return,
669            Err(error) => {
670                last_error = Some(error);
671                if attempt < 2 {
672                    tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
673                }
674            }
675        }
676    }
677    if let Some(error) = last_error {
678        log::warn!(
679            "[segment_lifecycle] index={index_label} replacement topology refresh failed after 3 attempts: {}",
680            error,
681        );
682    }
683}
684
685#[cfg(feature = "native")]
686struct MergeTaskError {
687    error: Error,
688    unavailable_segments: Vec<String>,
689}
690
691#[cfg(feature = "native")]
692impl MergeTaskError {
693    fn source(segment_id: String, error: Error) -> Self {
694        Self {
695            error,
696            unavailable_segments: vec![segment_id],
697        }
698    }
699
700    fn sources(segment_ids: Vec<String>, error: Error) -> Self {
701        Self {
702            error,
703            unavailable_segments: segment_ids,
704        }
705    }
706}
707
708#[cfg(feature = "native")]
709impl From<Error> for MergeTaskError {
710    fn from(error: Error) -> Self {
711        Self {
712            error,
713            unavailable_segments: Vec::new(),
714        }
715    }
716}
717
718#[cfg(feature = "native")]
719fn is_deterministic_source_error(error: &Error) -> bool {
720    matches!(error, Error::Corruption(_) | Error::Serialization(_))
721        || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
722}
723
724#[cfg(feature = "native")]
725fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
726    if is_deterministic_source_error(&error) {
727        MergeTaskError::source(segment_id, error)
728    } else {
729        // Timeouts, interrupted reads, permission changes, and other generic
730        // I/O failures may be transient. Back them off instead of quarantining
731        // a healthy metadata segment for the rest of the process lifetime.
732        MergeTaskError::from(error)
733    }
734}
735
736#[cfg(feature = "native")]
737type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
738
739#[derive(Clone, Copy)]
740enum ReplacementLayout {
741    /// No BP ran while combining these sources. The combined segment is not
742    /// globally reordered, but an interrupted BP lineage remains owed its
743    /// record-level deepening pass.
744    BlockCopy,
745    /// A BP pass produced the replacement layout.
746    BpReordered { converged: bool },
747    /// A vector-only rewrite leaves document order and sparse layout exactly
748    /// unchanged, so its persisted BP progress must not be reset or advanced.
749    PreserveSingleSource,
750}
751
752fn replacement_bp_state(
753    parent_has_debt: bool,
754    parent_unconverged_passes: u32,
755    layout: ReplacementLayout,
756) -> (bool, bool, u32) {
757    match layout {
758        ReplacementLayout::BlockCopy => (
759            false,
760            !parent_has_debt,
761            if parent_has_debt {
762                parent_unconverged_passes
763            } else {
764                0
765            },
766        ),
767        ReplacementLayout::BpReordered { converged } => (
768            true,
769            converged,
770            if converged {
771                0
772            } else {
773                parent_unconverged_passes.saturating_add(1)
774            },
775        ),
776        ReplacementLayout::PreserveSingleSource => {
777            unreachable!("preserved layouts retain the complete source metadata")
778        }
779    }
780}
781
782#[derive(Clone, Copy, Debug, Eq, PartialEq)]
783enum VectorSegmentRewriteOutcome {
784    Rewritten,
785    AlreadyCurrent,
786    SourceGone,
787    Conflict,
788    Deferred,
789}
790
791/// Complete but unpublished vector-only replacement. Its lifecycle claim and
792/// cleanup guard stay armed until the whole codebook generation commits.
793pub(crate) struct StagedVectorSegment {
794    source_id: String,
795    output_id: SegmentId,
796    doc_count: u32,
797    _operation: SegmentOperationGuard,
798    cleanup: OutputCleanupGuard,
799}
800
801/// Segment manager — coordinates segment commit, background merging, and trained structures.
802///
803/// SOLE owner of `metadata.json`. All metadata mutations go through `state` Mutex.
804pub struct SegmentManager<D: DirectoryWriter + 'static> {
805    /// Serializes ALL metadata mutations.
806    state: Arc<AsyncMutex<ManagerState>>,
807
808    /// RAII ownership for every in-flight segment lifecycle operation.
809    active_operations: Arc<ActiveSegmentOperations>,
810
811    /// Metadata-live segments involved in a deterministic source/corruption
812    /// failure. They stay searchable (and operator-visible) but are excluded
813    /// from merges for this process lifetime, preventing a bad candidate from
814    /// consuming full rewrite capacity on every retry.
815    quarantined_segments: parking_lot::Mutex<HashSet<String>>,
816
817    /// Generic merge failures pause scheduling briefly. Source-specific open
818    /// failures use `quarantined_segments` instead so healthy work can continue.
819    merge_retry: parking_lot::Mutex<MergeRetryState>,
820
821    /// Per-source backoff for non-deterministic standalone reorder failures.
822    /// Optimizer scans are periodic, but a pass can outlast the scan interval;
823    /// without completion-based backoff it would restart almost immediately.
824    reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
825
826    /// In-flight merge JoinHandles — supports multiple concurrent merges.
827    merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
828
829    /// At most one task per index waits for application-wide merge capacity.
830    /// Without this wakeup, an index denied by another index can remain idle
831    /// forever when no later commit happens to re-run merge policy evaluation.
832    global_merge_wakeup_pending: AtomicBool,
833
834    /// Non-zero while an explicit force merge is draining/running. Automatic
835    /// merges do not start during this window; otherwise a merge spawned
836    /// between the initial drain and foreground BP reservation could claim
837    /// source segments and then wait behind the foreground gate.
838    force_merge_active: AtomicUsize,
839
840    /// Times `force_merge` observed a conflicting active operation and retried.
841    /// Test-only observability for the conflict-retry backoff.
842    #[cfg(test)]
843    force_merge_conflict_retries: std::sync::atomic::AtomicU64,
844
845    /// Auxiliary lifecycle tasks: metadata transactions, deferred deletes,
846    /// and capacity wakeups. Handles registered here are drained before index
847    /// removal.
848    lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
849
850    /// Trained vector structures — lock-free reads via ArcSwap.
851    /// Wrapped in `Arc` so cancellation-safe metadata transactions can publish
852    /// the matching in-memory generation after their durable commit point.
853    published_generation: Arc<ArcSwap<PublishedIndexGeneration>>,
854
855    /// Raised while index-level trained artifacts and their metadata are being
856    /// replaced. Search readers keep using the last valid generation, while
857    /// segment producers fall back to flat output until publication completes.
858    vector_artifact_update: Arc<AtomicBool>,
859
860    /// Reference counting for safe segment deletion (sync Mutex for Drop).
861    tracker: Arc<SegmentTracker>,
862
863    /// Cached deletion callback for snapshots (avoids allocation per acquire_snapshot).
864    delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
865
866    /// Directory for segment I/O
867    directory: Arc<D>,
868    /// Schema for segment operations
869    schema: Arc<crate::dsl::Schema>,
870    /// Compression mode used for term dictionaries produced by merges.
871    optimization: crate::structures::IndexOptimization,
872    /// Posting codec used by new segments and merge re-encoding.
873    posting_codec: crate::structures::PostingCodec,
874    /// Term cache blocks for segment readers during merge
875    term_cache_blocks: usize,
876    /// Hard concurrency limit for background merges. A semaphore permit is
877    /// acquired before lifecycle ownership, closing the old handle-count race
878    /// where concurrent schedulers could exceed the configured maximum.
879    merge_permits: Arc<Semaphore>,
880    /// Application-wide merge limit shared across index managers.
881    global_merge_permits: Arc<Semaphore>,
882    /// Shared across every index opened from the same `IndexConfig`. This
883    /// bounds whole BP rewrites (optimizer + merge-time + manual) separately
884    /// from Rayon thread width, preventing N × memory-budget amplification.
885    reorder_permits: Arc<ReorderConcurrencyGate>,
886    /// Run BP reordering of `reorder`-attributed BMP fields inside merges.
887    /// Persisted index configuration (schema-level `reorder_on_merge: true`
888    /// in SDL); merged segments are marked `reordered` and skipped by the
889    /// standalone optimizer pass.
890    reorder_on_merge: bool,
891    /// Wall-clock budget for merge-time BP (from `IndexConfig`); truncated
892    /// passes mark the merged segment `bp_converged = false` so the
893    /// background optimizer deepens it later (warm-started).
894    merge_bp_time_budget: Option<std::time::Duration>,
895    /// Memory budget for the BP forward index (merge-time and background
896    /// reorder). Over-budget passes drop highest-df dims, logged loudly.
897    bp_memory_budget_bytes: usize,
898    /// Application-owned shared pool, when configured. This is the server
899    /// path and ensures optimizer and merge-time work use the same threads.
900    background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
901    /// Writer-owned topology reconciler. Every durable replacement invokes
902    /// it from tracked lifecycle work so background merges cannot leave
903    /// primary-key snapshots pinning retired sources indefinitely.
904    replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
905}
906
907struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
908
909impl Drop for ForceMergeActivityGuard<'_> {
910    fn drop(&mut self) {
911        self.0.fetch_sub(1, Ordering::AcqRel);
912    }
913}
914
915impl<D: DirectoryWriter + 'static> SegmentManager<D> {
916    /// Create a new segment manager with existing metadata
917    #[allow(clippy::too_many_arguments)]
918    pub fn new(
919        directory: Arc<D>,
920        schema: Arc<crate::dsl::Schema>,
921        metadata: IndexMetadata,
922        merge_policy: Box<dyn MergePolicy>,
923        term_cache_blocks: usize,
924        max_concurrent_merges: usize,
925        global_merge_permits: Arc<Semaphore>,
926        merge_bp_time_budget: Option<std::time::Duration>,
927        bp_memory_budget_bytes: usize,
928        reorder_permits: Arc<ReorderConcurrencyGate>,
929        background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
930    ) -> Self {
931        // Persisted index option: set via `reorder_on_merge: true` in the SDL
932        // at index creation. Absent = disabled (merges block-copy).
933        let reorder_on_merge = schema.reorder_on_merge();
934        if reorder_on_merge {
935            log::info!(
936                "[merge] index={} reorder-on-merge enabled by index schema",
937                schema.index_label()
938            );
939        }
940
941        let tracker = Arc::new(SegmentTracker::new());
942        for seg_id in metadata.segment_metas.keys() {
943            tracker.register(seg_id);
944        }
945
946        let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
947            Arc::new(parking_lot::Mutex::new(Vec::new()));
948        let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
949            let dir = Arc::clone(&directory);
950            let tracker = Arc::clone(&tracker);
951            let lifecycle_handles = Arc::clone(&lifecycle_handles);
952            let cleanup_index_label: Arc<str> = schema.index_label().into();
953            Arc::new(move |segment_ids| {
954                // Guard: if the tokio runtime is gone (program exit), skip async
955                // deletion. Segment files become orphans cleaned up on next startup.
956                let Ok(handle) = tokio::runtime::Handle::try_current() else {
957                    // Release in-process protection as well: if the process is
958                    // still alive, a later sweep must be able to retry.
959                    tracker.complete_deletion(&segment_ids);
960                    return;
961                };
962                let dir = Arc::clone(&dir);
963                let task_tracker = Arc::clone(&tracker);
964                let task_index_label = Arc::clone(&cleanup_index_label);
965                let cleanup_ids = segment_ids.clone();
966                let future = async move {
967                    for &segment_id in &segment_ids {
968                        log::info!(
969                            "[segment_cleanup] index={} deleting deferred segment {}",
970                            task_index_label,
971                            segment_id.to_hex()
972                        );
973                        if let Err(error) =
974                            crate::segment::delete_segment(dir.as_ref(), segment_id).await
975                        {
976                            log::warn!(
977                                "[segment_cleanup] index={} deferred delete failed for {}: {}",
978                                task_index_label,
979                                segment_id.to_hex(),
980                                error,
981                            );
982                        }
983                    }
984                    task_tracker.complete_deletion(&segment_ids);
985                };
986                if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
987                    // Spawning can fail only during runtime teardown. Release
988                    // the scheduled-deletion claim so an in-process sweep can
989                    // retry; crash recovery handles a process exit.
990                    tracker.complete_deletion(&cleanup_ids);
991                    log::warn!(
992                        "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
993                        cleanup_index_label
994                    );
995                }
996            })
997        };
998
999        let initial_generation = Arc::new(PublishedIndexGeneration {
1000            publication_id: metadata.publication_generation,
1001            schema: Arc::clone(&schema),
1002            trained_vectors: None,
1003        });
1004        Self {
1005            state: Arc::new(AsyncMutex::new(ManagerState {
1006                metadata,
1007                merge_policy,
1008            })),
1009            active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1010            quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1011            merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1012            reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1013            merge_handles: parking_lot::Mutex::new(Vec::new()),
1014            global_merge_wakeup_pending: AtomicBool::new(false),
1015            force_merge_active: AtomicUsize::new(0),
1016            #[cfg(test)]
1017            force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1018            lifecycle_handles,
1019            published_generation: Arc::new(ArcSwap::new(initial_generation)),
1020            vector_artifact_update: Arc::new(AtomicBool::new(false)),
1021            tracker,
1022            delete_fn,
1023            directory,
1024            schema,
1025            optimization: crate::structures::IndexOptimization::default(),
1026            posting_codec: crate::structures::PostingCodec::default(),
1027            term_cache_blocks,
1028            merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1029            global_merge_permits,
1030            reorder_permits,
1031            reorder_on_merge,
1032            merge_bp_time_budget,
1033            bp_memory_budget_bytes,
1034            background_reorder_pool,
1035            replacement_refresh: parking_lot::RwLock::new(None),
1036        }
1037    }
1038
1039    /// Set the persisted text-index layout selected by `IndexConfig` (or by
1040    /// an explicit `SegmentBuilderConfig` override).
1041    pub fn with_posting_config(
1042        mut self,
1043        optimization: crate::structures::IndexOptimization,
1044        posting_codec: crate::structures::PostingCodec,
1045    ) -> Self {
1046        self.optimization = optimization;
1047        self.posting_codec = posting_codec;
1048        self
1049    }
1050
1051    pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1052    where
1053        F: Fn() -> Fut + Send + Sync + 'static,
1054        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1055    {
1056        *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1057    }
1058
1059    /// Bounded rayon pool for background CPU (merge-time BP, manual reorder,
1060    /// and global vector-codebook training).
1061    /// Query scoring uses a dedicated search pool; keeping background work on
1062    /// this separate bounded pool prevents a merge or retrain from queueing
1063    /// every search behind its CPU passes.
1064    pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1065        if let Some(pool) = &self.background_reorder_pool {
1066            return Arc::clone(pool);
1067        }
1068        Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1069            let threads = (num_cpus::get() / 2).max(1);
1070            log::info!(
1071                "[merge] process-wide background CPU pool: {} thread(s)",
1072                threads
1073            );
1074            Arc::new(
1075                rayon::ThreadPoolBuilder::new()
1076                    .num_threads(threads)
1077                    .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1078                    .build()
1079                    .expect("failed to build background CPU pool"),
1080            )
1081        }))
1082    }
1083
1084    /// Stop new indexing/merge/reorder operations from claiming segment IDs.
1085    /// Used as the first half of index deletion; the writer then joins its
1086    /// workers before [`Self::wait_for_shutdown`] drains remaining ownership.
1087    pub fn begin_shutdown(&self) {
1088        self.active_operations.stop_accepting();
1089    }
1090
1091    /// Run a lifecycle mutation independently of its requesting future.
1092    ///
1093    /// Metadata writes contain an atomic rename. If an RPC is cancelled while
1094    /// awaiting that I/O, dropping the request must not abandon the matching
1095    /// in-memory/tracker transition. The spawned transaction is tracked for
1096    /// index shutdown; the oneshot only reports its result to a caller that is
1097    /// still interested.
1098    async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1099    where
1100        T: Send + 'static,
1101        F: std::future::Future<Output = Result<T>> + Send + 'static,
1102    {
1103        let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1104        let future = async move {
1105            let result = transaction.await;
1106            let _ = result_tx.send(result);
1107        };
1108        let runtime = tokio::runtime::Handle::current();
1109        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1110            return Err(Error::Internal(
1111                "runtime rejected lifecycle metadata transaction".into(),
1112            ));
1113        }
1114        result_rx.await.map_err(|_| {
1115            Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1116        })?
1117    }
1118
1119    /// Arm unwind cleanup for an output that is not visible in metadata yet.
1120    fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1121        let manager = Arc::clone(self);
1122        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1123            let Ok(handle) = tokio::runtime::Handle::try_current() else {
1124                log::warn!(
1125                    "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1126                    manager.schema.index_label(),
1127                    segment_id.to_hex(),
1128                );
1129                return;
1130            };
1131
1132            let cleanup_manager = Arc::clone(&manager);
1133            let future = async move {
1134                cleanup_manager
1135                    .delete_output_if_unregistered(segment_id, "task unwind")
1136                    .await;
1137            };
1138            if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1139                log::warn!(
1140                    "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1141                    manager.schema.index_label(),
1142                    segment_id.to_hex(),
1143                );
1144            }
1145        });
1146
1147        OutputCleanupGuard::new(output_id, cleanup)
1148    }
1149
1150    /// Delete an abandoned indexing output while retaining its lifecycle
1151    /// claim until the last file operation completes. The explicit runtime
1152    /// handle makes this safe from dedicated indexing OS threads, which are
1153    /// outside Tokio's entered context.
1154    pub(crate) fn schedule_unpublished_segment_cleanup(
1155        self: &Arc<Self>,
1156        output_id: SegmentId,
1157        operation: SegmentOperationGuard,
1158        runtime: tokio::runtime::Handle,
1159    ) {
1160        let manager = Arc::clone(self);
1161        let output_hex = output_id.to_hex();
1162        let future = async move {
1163            manager
1164                .delete_output_if_unregistered(output_id, "indexing abort or failure")
1165                .await;
1166            drop(operation);
1167        };
1168        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1169            // The dropped future releases operation ownership. Startup sweep
1170            // handles its output if the runtime is already tearing down.
1171            log::warn!(
1172                "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1173                self.schema.index_label(),
1174                output_hex,
1175            );
1176        }
1177    }
1178
1179    /// Claim a newly generated indexing segment before its first file write.
1180    ///
1181    /// The returned guard must travel with the built segment until metadata
1182    /// publication or abort. UUID collisions are treated as corruption rather
1183    /// than silently sharing lifecycle ownership.
1184    pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1185        match self
1186            .active_operations
1187            .try_register_indexing(vec![segment_id.clone()])
1188        {
1189            Some(operation) => Ok(operation),
1190            None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1191            None => Err(Error::Corruption(format!(
1192                "new segment ID {} is already owned by an active operation",
1193                segment_id
1194            ))),
1195        }
1196    }
1197
1198    /// Validate the small, mandatory core of a completed segment before it can
1199    /// become metadata-live. Optional vector/sparse/position/fast files are
1200    /// schema- and data-dependent and are validated by `SegmentReader` when used.
1201    async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1202        let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1203            Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1204        })?;
1205        let files = SegmentFiles::new(id.0);
1206
1207        for path in files.mandatory_paths() {
1208            if !self.directory.exists(path).await.map_err(Error::Io)? {
1209                return Err(Error::Corruption(format!(
1210                    "segment {} cannot be published: mandatory file {:?} is missing",
1211                    segment_id, path
1212                )));
1213            }
1214        }
1215
1216        let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1217            Error::Corruption(format!(
1218                "segment {} cannot be published: missing/unreadable {:?}: {}",
1219                segment_id, files.meta, e
1220            ))
1221        })?;
1222        let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1223            Error::Corruption(format!(
1224                "segment {} cannot be published: failed reading {:?}: {}",
1225                segment_id, files.meta, e
1226            ))
1227        })?;
1228        let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1229            Error::Corruption(format!(
1230                "segment {} cannot be published: invalid {:?}: {}",
1231                segment_id, files.meta, e
1232            ))
1233        })?;
1234
1235        if meta.id != id.0 || meta.num_docs != expected_docs {
1236            return Err(Error::Corruption(format!(
1237                "segment {} cannot be published: metadata identity/docs mismatch \
1238                 (id={:032x}, docs={}, expected_docs={})",
1239                segment_id, meta.id, meta.num_docs, expected_docs
1240            )));
1241        }
1242
1243        Ok(())
1244    }
1245
1246    fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1247        let inserted = self
1248            .quarantined_segments
1249            .lock()
1250            .insert(segment_id.to_string());
1251        if inserted {
1252            log::error!(
1253                "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1254                 It remains metadata-live for explicit repair but is excluded from merges until restart",
1255                self.schema.index_label(),
1256                segment_id,
1257                error,
1258            );
1259        }
1260    }
1261
1262    fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1263        let mut retry = self.merge_retry.lock();
1264        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1265        let delay = merge_retry_delay(retry.consecutive_failures);
1266        retry.retry_after = std::time::Instant::now().checked_add(delay);
1267        log::warn!(
1268            "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1269            self.schema.index_label(),
1270            delay.as_secs_f64(),
1271            retry.consecutive_failures,
1272            error,
1273        );
1274        delay
1275    }
1276
1277    fn clear_merge_retry_backoff(&self) {
1278        *self.merge_retry.lock() = MergeRetryState::default();
1279    }
1280
1281    fn merge_retry_is_paused(&self) -> bool {
1282        let mut retry = self.merge_retry.lock();
1283        match retry.retry_after {
1284            Some(deadline) if deadline > std::time::Instant::now() => true,
1285            Some(_) => {
1286                retry.retry_after = None;
1287                false
1288            }
1289            None => false,
1290        }
1291    }
1292
1293    fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1294        let mut retries = self.reorder_retries.lock();
1295        let retry = retries.entry(segment_id.to_string()).or_default();
1296        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1297        let delay = merge_retry_delay(retry.consecutive_failures);
1298        retry.retry_after = std::time::Instant::now().checked_add(delay);
1299        log::warn!(
1300            "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1301            self.schema.index_label(),
1302            segment_id,
1303            delay.as_secs_f64(),
1304            retry.consecutive_failures,
1305            error,
1306        );
1307    }
1308
1309    fn clear_reorder_retry(&self, segment_id: &str) {
1310        self.reorder_retries.lock().remove(segment_id);
1311    }
1312
1313    fn paused_reorder_segments(&self) -> HashSet<String> {
1314        let now = std::time::Instant::now();
1315        let mut retries = self.reorder_retries.lock();
1316        let mut paused = HashSet::new();
1317        for (segment_id, retry) in retries.iter_mut() {
1318            match retry.retry_after {
1319                Some(deadline) if deadline > now => {
1320                    paused.insert(segment_id.clone());
1321                }
1322                Some(_) => retry.retry_after = None,
1323                None => {}
1324            }
1325        }
1326        paused
1327    }
1328
1329    /// Re-evaluate this index when another index releases application-wide
1330    /// merge capacity. The atomic flag bounds this to one waiter per index and
1331    /// the tracked handle makes index shutdown drain it deterministically.
1332    fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1333        if self
1334            .global_merge_wakeup_pending
1335            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1336            .is_err()
1337        {
1338            return;
1339        }
1340
1341        let manager = Arc::clone(self);
1342        let future = async move {
1343            let capacity = tokio::select! {
1344                biased;
1345                () = manager.active_operations.wait_for_shutdown() => None,
1346                permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1347            };
1348
1349            manager
1350                .global_merge_wakeup_pending
1351                .store(false, Ordering::Release);
1352            if let Some(permit) = capacity {
1353                // This task is only a notification. The normal scheduler must
1354                // acquire both global and per-index permits atomically enough
1355                // for its own candidate selection.
1356                drop(permit);
1357                manager.maybe_merge().await;
1358            }
1359        };
1360        let runtime = tokio::runtime::Handle::current();
1361        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1362            self.global_merge_wakeup_pending
1363                .store(false, Ordering::Release);
1364            log::warn!(
1365                "[merge] index={} runtime rejected global-capacity wakeup task",
1366                self.schema.index_label()
1367            );
1368        }
1369    }
1370
1371    #[cfg(test)]
1372    pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1373        self.quarantined_segments.lock().contains(segment_id)
1374    }
1375
1376    /// Delete a failed output only if metadata did not make it live.
1377    ///
1378    /// Rechecking under `state` also makes unwind cleanup safe if it races
1379    /// successful publication of the same output.
1380    async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1381        let output_hex = output_id.to_hex();
1382        {
1383            let st = self.state.lock().await;
1384            if st.metadata.has_segment(&output_hex) {
1385                return;
1386            }
1387        }
1388
1389        // UUIDs are generated per producer and cannot be adopted by another
1390        // publisher after this check. Never hold the metadata mutex while a
1391        // multi-GB filesystem deletion runs.
1392        log::info!(
1393            "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1394            self.schema.index_label(),
1395            output_hex,
1396            reason,
1397        );
1398        if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1399        {
1400            log::warn!(
1401                "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1402                self.schema.index_label(),
1403                output_hex,
1404                error,
1405            );
1406        }
1407    }
1408
1409    // ========================================================================
1410    // Read path (brief lock or lock-free)
1411    // ========================================================================
1412
1413    /// Get the current segment IDs
1414    pub async fn get_segment_ids(&self) -> Vec<String> {
1415        self.state.lock().await.metadata.segment_ids()
1416    }
1417
1418    /// Get trained vector structures (lock-free via ArcSwap)
1419    pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1420        self.published_generation.load().trained_vectors.clone()
1421    }
1422
1423    /// Current schema/vector publication. Cloning the Arc gives an operation a
1424    /// stable configuration even if an ALTER publishes concurrently.
1425    pub(crate) fn published_generation(&self) -> Arc<PublishedIndexGeneration> {
1426        self.published_generation.load_full()
1427    }
1428
1429    pub(crate) fn publication_id(&self) -> u64 {
1430        self.published_generation.load().publication_id
1431    }
1432
1433    /// Capture trained structures for a segment producer.
1434    ///
1435    /// The second gate check closes the race where an update begins after the
1436    /// first check but before the ArcSwap load. Producers have lifecycle guards
1437    /// before calling this method, so an updater that raised the gate waits for
1438    /// any producer that successfully captured the previous generation.
1439    pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1440        if self.vector_artifact_update.load(Ordering::Acquire) {
1441            return None;
1442        }
1443        let trained = self.published_generation.load().trained_vectors.clone();
1444        if self.vector_artifact_update.load(Ordering::Acquire) {
1445            None
1446        } else {
1447            trained
1448        }
1449    }
1450
1451    /// Start an exclusive trained-artifact update and drain merge/reorder
1452    /// producers that may already hold the previous generation.
1453    ///
1454    /// New segment operations may continue while this waits, but they observe
1455    /// the gate through `trained_for_segment_build` and therefore emit flat
1456    /// vector data until the replacement generation is published. The guard
1457    /// is cancellation-safe: dropping the requesting future reopens ANN
1458    /// production without leaving the manager wedged.
1459    ///
1460    /// Indexing tokens cannot be waited on: their guards are parked inside
1461    /// built-but-uncommitted `PreparedSegment`s and are released only by a
1462    /// later commit. That commit typically needs the writer this update's
1463    /// caller already holds (server write lock / embedded `&mut self`), so
1464    /// waiting would permanently wedge vector-index finalization. They also
1465    /// cannot be ignored: such a segment may already contain ANN data bound to
1466    /// the previous artifact generation and could otherwise be committed after
1467    /// the new generation is published. Reject promptly and let the caller
1468    /// commit or abort the pending generation before retrying.
1469    pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1470        self.vector_artifact_update
1471            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1472            .map_err(|_| {
1473                Error::Internal("a trained-vector artifact update is already in progress".into())
1474            })?;
1475        self.active_operations.pause_non_indexing();
1476        let guard = VectorArtifactUpdateGuard {
1477            _lease: Arc::new(VectorArtifactUpdateLease {
1478                updating: Arc::clone(&self.vector_artifact_update),
1479                active_operations: Arc::clone(&self.active_operations),
1480            }),
1481        };
1482        let (preexisting, parked_indexing) =
1483            self.active_operations.draining_operation_tokens_snapshot();
1484        if parked_indexing > 0 {
1485            return Err(Error::Internal(format!(
1486                "cannot update trained-vector artifacts while {parked_indexing} indexing \
1487                 segment(s) are built but uncommitted; commit or abort the pending \
1488                 generation and retry"
1489            )));
1490        }
1491        self.active_operations
1492            .wait_until_operations_finish(&preexisting)
1493            .await;
1494        Ok(guard)
1495    }
1496
1497    /// Load trained structures from disk and publish to ArcSwap.
1498    /// Copies metadata under lock, releases lock, then does disk I/O.
1499    pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1500        // Copy vector_fields under lock (cheap clone of HashMap<u32, FieldMeta>)
1501        let (vector_fields, schema, publication_id) = {
1502            let st = self.state.lock().await;
1503            (
1504                st.metadata.vector_fields.clone(),
1505                Arc::new(st.metadata.schema.clone()),
1506                st.metadata.publication_generation,
1507            )
1508        };
1509        // Disk I/O happens WITHOUT holding the state lock
1510        let trained = IndexMetadata::try_load_trained_from_fields(
1511            &vector_fields,
1512            schema.as_ref(),
1513            self.directory.as_ref(),
1514        )
1515        .await?
1516        .map(Arc::new);
1517        // Publish exactly the validated snapshot, including None. Retaining a
1518        // previous map when metadata has no Built fields would let new segments
1519        // depend on artifacts no longer referenced durably.
1520        self.published_generation
1521            .store(Arc::new(PublishedIndexGeneration {
1522                publication_id,
1523                schema,
1524                trained_vectors: trained,
1525            }));
1526        Ok(())
1527    }
1528
1529    /// Atomically publish a fully staged vector generation.
1530    ///
1531    /// Every replacement segment and its global codebook is complete before
1532    /// this transaction starts. Metadata, tracker ownership, and the lock-free
1533    /// trained pointer advance under the same state lock; snapshots therefore
1534    /// observe either the complete old generation or the complete new one.
1535    pub(crate) async fn publish_vector_generation(
1536        self: &Arc<Self>,
1537        artifact_update: &VectorArtifactUpdateGuard,
1538        vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1539        next_trained: Arc<TrainedVectorStructures>,
1540        staged: Vec<StagedVectorSegment>,
1541    ) -> Result<()> {
1542        let schema = self.published_generation().schema.clone();
1543        self.publish_vector_generation_with_schema(
1544            artifact_update,
1545            schema,
1546            vector_fields,
1547            Some(next_trained),
1548            staged,
1549        )
1550        .await
1551    }
1552
1553    pub(crate) async fn publish_vector_generation_with_schema(
1554        self: &Arc<Self>,
1555        artifact_update: &VectorArtifactUpdateGuard,
1556        schema: Arc<crate::dsl::Schema>,
1557        vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1558        next_trained: Option<Arc<TrainedVectorStructures>>,
1559        mut staged: Vec<StagedVectorSegment>,
1560    ) -> Result<()> {
1561        if !self.vector_artifact_update.load(Ordering::Acquire) {
1562            return Err(Error::Internal(
1563                "vector generation publication lost its exclusive update lease".into(),
1564            ));
1565        }
1566
1567        for replacement in &staged {
1568            self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1569                .await?;
1570        }
1571
1572        let mut st = Arc::clone(&self.state).lock_owned().await;
1573        let mut next = st.metadata.clone();
1574        next.schema = (*schema).clone();
1575        next.vector_fields = vector_fields;
1576        next.refresh_total_vectors();
1577
1578        for replacement in &staged {
1579            let source_info = next
1580                .segment_metas
1581                .remove(&replacement.source_id)
1582                .ok_or_else(|| {
1583                    Error::Corruption(format!(
1584                        "vector generation source {} disappeared before publication",
1585                        replacement.source_id,
1586                    ))
1587                })?;
1588            let output_hex = replacement.output_id.to_hex();
1589            if next.segment_metas.contains_key(&output_hex) {
1590                return Err(Error::Corruption(format!(
1591                    "vector generation output {output_hex} is already metadata-live"
1592                )));
1593            }
1594            // A vector-only rewrite changes neither document order nor merge
1595            // lineage, so preserve the complete lifecycle record verbatim.
1596            next.add_segment_meta(output_hex, source_info);
1597        }
1598
1599        let directory = Arc::clone(&self.directory);
1600        let published_generation = Arc::clone(&self.published_generation);
1601        let tracker = Arc::clone(&self.tracker);
1602        let replacement_refresh = self.replacement_refresh.read().clone();
1603        // Keep the producer gate raised if the requesting future is cancelled
1604        // after the metadata transaction has been detached. The last guard
1605        // clone drops only after durable metadata and ArcSwap state agree.
1606        let artifact_update = artifact_update.clone();
1607        let index_label = self.schema.index_label().to_owned();
1608        next.publication_generation =
1609            next.publication_generation.checked_add(1).ok_or_else(|| {
1610                Error::Corruption("vector publication generation exhausted u64".into())
1611            })?;
1612        let next_schema = schema;
1613        let next_publication_id = next.publication_generation;
1614        self.run_lifecycle_transaction(async move {
1615            let _artifact_update = artifact_update;
1616            next.save(directory.as_ref()).await?;
1617
1618            for replacement in &staged {
1619                tracker.register(&replacement.output_id.to_hex());
1620            }
1621            st.metadata = next;
1622            published_generation.store(Arc::new(PublishedIndexGeneration {
1623                publication_id: next_publication_id,
1624                schema: next_schema,
1625                trained_vectors: next_trained,
1626            }));
1627
1628            // Outputs are now durably live. Disarm unwind cleanup before
1629            // retiring the old generation and releasing operation ownership.
1630            for replacement in &mut staged {
1631                replacement.cleanup.disarm();
1632            }
1633            let retired = staged
1634                .iter()
1635                .map(|replacement| replacement.source_id.clone())
1636                .collect::<Vec<_>>();
1637            let ready_to_delete = tracker.mark_for_deletion(&retired);
1638            drop(st);
1639            for &segment_id in &ready_to_delete {
1640                if let Err(error) =
1641                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
1642                {
1643                    log::warn!(
1644                        "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1645                        segment_id.to_hex(),
1646                        error,
1647                    );
1648                }
1649            }
1650            tracker.complete_deletion(&ready_to_delete);
1651            refresh_replacement_topology(replacement_refresh, &index_label).await;
1652            Ok(())
1653        })
1654        .await
1655    }
1656
1657    /// Publish query-only vector parameters without rebuilding segments or
1658    /// artifacts. The publication ID still advances so cached readers reload
1659    /// even though the segment ID set is unchanged.
1660    pub(crate) async fn publish_vector_schema_only(
1661        self: &Arc<Self>,
1662        artifact_update: &VectorArtifactUpdateGuard,
1663        schema: Arc<crate::dsl::Schema>,
1664    ) -> Result<()> {
1665        let vector_fields = self
1666            .read_metadata(|metadata| metadata.vector_fields.clone())
1667            .await;
1668        let trained = self.published_generation().trained_vectors.clone();
1669        self.publish_vector_generation_with_schema(
1670            artifact_update,
1671            schema,
1672            vector_fields,
1673            trained,
1674            Vec::new(),
1675        )
1676        .await
1677    }
1678
1679    /// Read metadata with a closure (no persist)
1680    pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1681    where
1682        F: FnOnce(&IndexMetadata) -> R,
1683    {
1684        let st = self.state.lock().await;
1685        f(&st.metadata)
1686    }
1687
1688    /// Update metadata with a closure and persist atomically
1689    pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1690    where
1691        F: FnOnce(&mut IndexMetadata),
1692    {
1693        let mut st = Arc::clone(&self.state).lock_owned().await;
1694        let mut next = st.metadata.clone();
1695        f(&mut next);
1696        let directory = Arc::clone(&self.directory);
1697        self.run_lifecycle_transaction(async move {
1698            next.save(directory.as_ref()).await?;
1699            st.metadata = next;
1700            Ok(())
1701        })
1702        .await
1703    }
1704
1705    /// Acquire a snapshot of current segments for reading.
1706    /// The snapshot holds references — segments won't be deleted while snapshot exists.
1707    pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1708        let (acquired, generation) = {
1709            let st = self.state.lock().await;
1710            let segment_ids = st.metadata.segment_ids();
1711            (
1712                self.tracker.acquire(&segment_ids),
1713                self.published_generation.load_full(),
1714            )
1715        };
1716
1717        SegmentSnapshot::with_generation(
1718            Arc::clone(&self.tracker),
1719            acquired,
1720            generation,
1721            Arc::clone(&self.delete_fn),
1722        )
1723    }
1724
1725    /// Get the segment tracker
1726    pub fn tracker(&self) -> Arc<SegmentTracker> {
1727        Arc::clone(&self.tracker)
1728    }
1729
1730    /// Get the directory
1731    pub fn directory(&self) -> Arc<D> {
1732        Arc::clone(&self.directory)
1733    }
1734}
1735
1736// ============================================================================
1737// Native-only: commit, merging, force_merge
1738// ============================================================================
1739
1740#[cfg(feature = "native")]
1741impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1742    /// Atomic commit: register new segments + persist metadata.
1743    pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1744        // Indexing guards still own these IDs here, so the orphan sweeper
1745        // cannot remove files between validation and metadata publication.
1746        for (segment_id, num_docs) in new_segments {
1747            self.validate_completed_segment(segment_id, *num_docs)
1748                .await?;
1749        }
1750
1751        let mut st = Arc::clone(&self.state).lock_owned().await;
1752        let mut next = st.metadata.clone();
1753        let mut added = Vec::new();
1754        for (segment_id, num_docs) in new_segments {
1755            if !next.has_segment(segment_id) {
1756                next.add_segment(segment_id.clone(), *num_docs);
1757                added.push(segment_id.clone());
1758            }
1759        }
1760
1761        // Durable-before-visible: a save failure leaves both in-memory metadata
1762        // and tracker unchanged, so callers can retry the prepared commit.
1763        // The tracked transaction continues if the requesting RPC is cancelled;
1764        // unpublished cleanup waits on this owned state guard before deciding
1765        // whether the files became metadata-live.
1766        let directory = Arc::clone(&self.directory);
1767        let tracker = Arc::clone(&self.tracker);
1768        self.run_lifecycle_transaction(async move {
1769            next.save(directory.as_ref()).await?;
1770            for segment_id in &added {
1771                tracker.register(segment_id);
1772            }
1773            st.metadata = next;
1774            Ok(())
1775        })
1776        .await
1777    }
1778
1779    /// Evaluate merge policy and spawn background merges for all eligible candidates.
1780    ///
1781    /// **Atomicity**: The entire filter → find_merges → spawn_merge sequence runs
1782    /// under the `state` lock to prevent a TOCTOU race where concurrent callers
1783    /// both see segments as eligible before either claims operation ownership.
1784    /// `spawn_merge` is non-blocking (just `try_register` + `tokio::spawn`), so
1785    /// holding the state lock through it is safe and sub-microsecond.
1786    ///
1787    /// The hard merge semaphore is acquired before lifecycle ownership, so
1788    /// concurrent triggers cannot exceed configured merge capacity.
1789    pub async fn maybe_merge(self: &Arc<Self>) {
1790        if !self.active_operations.is_accepting() {
1791            log::debug!(
1792                "[maybe_merge] index={} manager is shutting down, skipping",
1793                self.schema.index_label()
1794            );
1795            return;
1796        }
1797        if self.merge_retry_is_paused() {
1798            log::debug!(
1799                "[maybe_merge] index={} retry backoff active, skipping",
1800                self.schema.index_label()
1801            );
1802            return;
1803        }
1804
1805        // Finished handles no longer need to be retained. Concurrency itself
1806        // is enforced by `merge_permits`, not this bookkeeping vector.
1807        {
1808            let mut handles = self.merge_handles.lock();
1809            handles.retain(|h| !h.is_finished());
1810        }
1811        let local_slots = self.merge_permits.available_permits();
1812        let global_slots = self.global_merge_permits.available_permits();
1813        let slots_available = local_slots.min(global_slots);
1814
1815        // Hold state lock through spawn_merge to make filter + register atomic.
1816        // This closes the TOCTOU window where concurrent maybe_merge calls could
1817        // both see the same segments as eligible before either registers them.
1818        {
1819            let st = self.state.lock().await;
1820            let quarantined = self.quarantined_segments.lock().clone();
1821            let active_ids = self.active_operations.snapshot();
1822
1823            // Backlog pressure is based on every metadata-live segment, not
1824            // only currently available inputs. Otherwise a wave of in-flight
1825            // merges hides most of the topology and makes later candidates
1826            // look healthy while the original backlog still exists.
1827            let live_segments: Vec<SegmentInfo> = st
1828                .metadata
1829                .segment_metas
1830                .iter()
1831                .filter(|(id, _)| {
1832                    !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1833                })
1834                .map(|(id, info)| SegmentInfo {
1835                    id: id.clone(),
1836                    num_docs: info.num_docs,
1837                })
1838                .collect();
1839            let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1840
1841            // Exclude segments owned by another operation. Pending retirement
1842            // and quarantined segments were already removed above.
1843            let segments: Vec<SegmentInfo> = live_segments
1844                .iter()
1845                .filter(|segment| !active_ids.contains(&segment.id))
1846                .cloned()
1847                .collect();
1848
1849            log::debug!(
1850                "[maybe_merge] index={} {} eligible segments",
1851                self.schema.index_label(),
1852                segments.len()
1853            );
1854
1855            let candidates = st.merge_policy.find_merges(&segments);
1856
1857            if candidates.is_empty() {
1858                return;
1859            }
1860
1861            // Register a capacity waiter only for an index that actually has
1862            // eligible work. Scheduling one waiter for every idle index while
1863            // the process gate was full caused an avoidable wakeup stampede.
1864            if slots_available == 0 {
1865                if local_slots > 0 && global_slots == 0 {
1866                    self.schedule_global_merge_wakeup();
1867                }
1868                log::debug!(
1869                    "[maybe_merge] index={} at max concurrent merges, skipping",
1870                    self.schema.index_label()
1871                );
1872                return;
1873            }
1874
1875            log::debug!(
1876                "[maybe_merge] index={} {} merge candidates, {} slots available",
1877                self.schema.index_label(),
1878                candidates.len(),
1879                slots_available
1880            );
1881
1882            let mut handles = Vec::new();
1883            for c in candidates {
1884                if handles.len() >= slots_available {
1885                    break;
1886                }
1887                // Under severe topology pressure, retire segments first. BP
1888                // remains optional for correctness and the block-copy output
1889                // is explicitly marked unreordered, so the optimizer performs
1890                // one pass after the compaction wave instead of every merge
1891                // task serializing behind the whole-pass gate.
1892                let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1893                if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1894                    handles.push(h);
1895                }
1896            }
1897            if !handles.is_empty() {
1898                if severe_backlog && self.reorder_on_merge {
1899                    log::info!(
1900                        "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1901                         block-copy merge(s), deferring BP to the optimizer",
1902                        self.schema.index_label(),
1903                        live_segments.len(),
1904                        handles.len(),
1905                    );
1906                }
1907                // Publish handles before releasing `state`. A force merge
1908                // raises its admission barrier under the same lock, then
1909                // drains this list, so no spawned merge can fall into the
1910                // otherwise-deadlocking gap between spawn and bookkeeping.
1911                self.merge_handles.lock().extend(handles);
1912            }
1913        }
1914    }
1915
1916    /// Spawn a background merge task with RAII tracking.
1917    ///
1918    /// Pre-generates the output segment ID. The operation guard registers all segment IDs
1919    /// (old + output) in `active_operations`. When the task ends (success, failure, or
1920    /// panic), the guard drops and segments are automatically unregistered.
1921    ///
1922    /// On completion, the task auto-triggers `maybe_merge` to evaluate cascading merges.
1923    /// Returns the JoinHandle if the merge was spawned, None if it was skipped.
1924    fn spawn_merge(
1925        self: &Arc<Self>,
1926        segment_ids_to_merge: Vec<String>,
1927        reorder_bmp: bool,
1928    ) -> Option<JoinHandle<()>> {
1929        if self.force_merge_active.load(Ordering::Acquire) > 0 {
1930            log::debug!(
1931                "[spawn_merge] index={} skipped: explicit force merge has priority",
1932                self.schema.index_label()
1933            );
1934            return None;
1935        }
1936        let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1937            Ok(permit) => permit,
1938            Err(_) => {
1939                log::debug!(
1940                    "[spawn_merge] index={} skipped: global merge capacity is full",
1941                    self.schema.index_label()
1942                );
1943                self.schedule_global_merge_wakeup();
1944                return None;
1945            }
1946        };
1947        let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1948            Ok(permit) => permit,
1949            Err(_) => {
1950                log::debug!(
1951                    "[spawn_merge] index={} skipped: no merge permit available",
1952                    self.schema.index_label()
1953                );
1954                return None;
1955            }
1956        };
1957        let output_id = SegmentId::new();
1958        let output_hex = output_id.to_hex();
1959
1960        let mut all_ids = segment_ids_to_merge.clone();
1961        all_ids.push(output_hex);
1962
1963        let guard = match self.active_operations.try_register(all_ids) {
1964            Some(g) => g,
1965            None => {
1966                log::debug!(
1967                    "[spawn_merge] index={} skipped: segments overlap with an active operation",
1968                    self.schema.index_label()
1969                );
1970                return None;
1971            }
1972        };
1973
1974        let sm = Arc::clone(self);
1975        let ids = segment_ids_to_merge;
1976
1977        let index_label = self.schema.index_label().to_owned();
1978        Some(tokio::spawn(async move {
1979            let mut reevaluate = false;
1980            let mut retry_delay = None;
1981
1982            let result = sm
1983                .merge_and_replace_registered(
1984                    &ids,
1985                    output_id,
1986                    reorder_bmp,
1987                    ReorderPriority::AutomaticMerge,
1988                )
1989                .await;
1990
1991            match result {
1992                Ok(_) => {
1993                    sm.clear_merge_retry_backoff();
1994                    reevaluate = true;
1995                }
1996                Err(MergeTaskError {
1997                    error: Error::IndexClosed,
1998                    ..
1999                }) => {
2000                    log::debug!(
2001                        "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
2002                        ids,
2003                    );
2004                }
2005                Err(MergeTaskError {
2006                    error,
2007                    unavailable_segments,
2008                }) => {
2009                    log::error!(
2010                        "[merge] index={index_label} background merge failed for segments {:?}: {}",
2011                        ids,
2012                        error
2013                    );
2014                    if !unavailable_segments.is_empty() {
2015                        // Recompute without known-bad inputs. This is not a
2016                        // retry of the same candidate because policy filtering
2017                        // excludes every quarantined ID.
2018                        reevaluate = true;
2019                    } else {
2020                        retry_delay = Some(sm.pause_merge_retries(&error));
2021                    }
2022                }
2023            }
2024            // Release source/output ownership before re-evaluating policy, so
2025            // the completed operation cannot artificially hide candidates.
2026            drop(guard);
2027            // A failed merge must not reserve capacity during its retry delay.
2028            drop(merge_permit);
2029            drop(global_merge_permit);
2030
2031            if reevaluate {
2032                sm.maybe_merge().await;
2033            } else if let Some(retry_delay) = retry_delay {
2034                // A backoff without a wakeup can strand eligible segments
2035                // forever when no later commit happens. The sleep runs as a
2036                // tracked *lifecycle* task, not inside this merge JoinHandle:
2037                // wait_for_all_merges/force_merge/reorder drain merge handles,
2038                // and a pure backoff timer with no work in flight must not
2039                // stall them for up to MERGE_RETRY_MAX_DELAY.
2040                sm.schedule_merge_retry_wakeup(retry_delay);
2041            }
2042        }))
2043    }
2044
2045    /// Execute the common build → validate → durable replacement transaction
2046    /// for a batch whose source/output IDs are already lifecycle-owned.
2047    ///
2048    /// Scheduling, retry policy, and post-replacement snapshot refresh remain
2049    /// with the caller. Keeping the transaction here prevents background and
2050    /// force-merge paths from drifting on cleanup or BP metadata semantics.
2051    async fn merge_and_replace_registered(
2052        self: &Arc<Self>,
2053        ids: &[String],
2054        output_id: SegmentId,
2055        reorder_bmp: bool,
2056        priority: ReorderPriority,
2057    ) -> MergeTaskResult<(String, u32, bool)> {
2058        let mut output_cleanup = self.output_cleanup_guard(output_id);
2059        let generation = self.published_generation();
2060        let trained = self.trained_for_segment_build();
2061        let granularity = if reorder_bmp {
2062            self.merge_granularity(ids).await
2063        } else {
2064            crate::segment::reorder::BpGranularity::Auto
2065        };
2066        let result = Self::do_merge(
2067            self.directory.as_ref(),
2068            &generation.schema,
2069            ids,
2070            output_id,
2071            self.term_cache_blocks,
2072            self.optimization,
2073            self.posting_codec,
2074            trained.as_deref(),
2075            reorder_bmp,
2076            granularity,
2077            self.merge_bp_time_budget,
2078            self.bp_memory_budget_bytes,
2079            Arc::clone(&self.reorder_permits),
2080            priority,
2081            self.active_operations.cancellation_flag(),
2082            Some(self.background_cpu_pool()),
2083        )
2084        .await;
2085
2086        let (new_id, doc_count, bp_converged) = match result {
2087            Ok(value) => value,
2088            Err(error) => {
2089                for segment_id in &error.unavailable_segments {
2090                    self.quarantine_segment(segment_id, &error.error);
2091                }
2092                self.delete_output_if_unregistered(output_id, "merge failure")
2093                    .await;
2094                output_cleanup.disarm();
2095                return Err(error);
2096            }
2097        };
2098
2099        let layout = if reorder_bmp {
2100            ReplacementLayout::BpReordered {
2101                converged: bp_converged,
2102            }
2103        } else {
2104            ReplacementLayout::BlockCopy
2105        };
2106        if let Err(error) = self
2107            .replace_segments(ids, new_id.clone(), doc_count, layout)
2108            .await
2109        {
2110            self.delete_output_if_unregistered(output_id, "replacement failure")
2111                .await;
2112            output_cleanup.disarm();
2113            return Err(MergeTaskError::from(error));
2114        }
2115        output_cleanup.disarm();
2116        Ok((new_id, doc_count, bp_converged))
2117    }
2118
2119    /// Re-evaluate merge policy after a failure backoff, outside the tracked
2120    /// merge JoinHandles that merge waiters drain. Shutdown still drains this
2121    /// task deterministically (lifecycle handles) and interrupts its sleep.
2122    fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2123        let manager = Arc::clone(self);
2124        let future = async move {
2125            tokio::select! {
2126                () = tokio::time::sleep(retry_delay) => {
2127                    manager.maybe_merge().await;
2128                }
2129                () = manager.active_operations.wait_for_shutdown() => {}
2130            }
2131        };
2132        let runtime = tokio::runtime::Handle::current();
2133        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2134            log::warn!(
2135                "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2136                 unmerged until the next commit re-runs merge policy evaluation",
2137                self.schema.index_label()
2138            );
2139        }
2140    }
2141
2142    /// Atomically replace old segments with a new merged segment.
2143    /// Computes merge generation as max(parent gens) + 1 and records ancestors.
2144    /// `reordered` marks whether the new segment was BP-reordered.
2145    async fn replace_segments(
2146        self: &Arc<Self>,
2147        old_ids: &[String],
2148        new_id: String,
2149        doc_count: u32,
2150        layout: ReplacementLayout,
2151    ) -> Result<()> {
2152        // The operation guard owns the output during validation. Publication
2153        // below replaces that ownership with metadata + tracker atomically.
2154        self.validate_completed_segment(&new_id, doc_count).await?;
2155        let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2156            Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2157        })?;
2158        let output_reader = SegmentReader::open(
2159            self.directory.as_ref(),
2160            output_id,
2161            self.published_generation().schema.clone(),
2162            self.term_cache_blocks,
2163        )
2164        .await
2165        .map_err(|error| match error {
2166            // Preserve retryable storage failures as I/O. Structural failures
2167            // are deterministic for this completed output and get explicit
2168            // corruption context.
2169            Error::Io(_) | Error::IndexClosed => error,
2170            error => Error::Corruption(format!(
2171                "replacement segment {new_id} failed full reader validation: {error}"
2172            )),
2173        })?;
2174        if output_reader.num_docs() != doc_count {
2175            return Err(Error::Corruption(format!(
2176                "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2177                output_reader.num_docs(),
2178            )));
2179        }
2180        drop(output_reader);
2181
2182        let mut st = Arc::clone(&self.state).lock_owned().await;
2183        // Every source must still be live: callers hold operation ownership,
2184        // so a missing source means a stale merge/reorder whose input was
2185        // already replaced. Adding the output would duplicate its documents.
2186        let missing: Vec<&String> = old_ids
2187            .iter()
2188            .filter(|id| !st.metadata.has_segment(id))
2189            .collect();
2190        if !missing.is_empty() {
2191            return Err(Error::Corruption(format!(
2192                "replace_segments: source segment(s) {:?} not in metadata — \
2193                 refusing to add output {} (would duplicate documents)",
2194                missing, new_id
2195            )));
2196        }
2197
2198        let replacement_info = match layout {
2199            ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2200                let generation = old_ids
2201                    .iter()
2202                    .filter_map(|id| st.metadata.segment_metas.get(id))
2203                    .map(|info| info.generation)
2204                    .max()
2205                    .unwrap_or(0)
2206                    .checked_add(1)
2207                    .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2208                let parent_unconverged_passes = old_ids
2209                    .iter()
2210                    .filter_map(|id| st.metadata.segment_metas.get(id))
2211                    .map(|info| info.bp_unconverged_passes)
2212                    .max()
2213                    .unwrap_or(0);
2214                let parent_has_debt = old_ids
2215                    .iter()
2216                    .filter_map(|id| st.metadata.segment_metas.get(id))
2217                    .any(|info| !info.bp_converged);
2218                let (reordered, bp_converged, bp_unconverged_passes) =
2219                    replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2220                SegmentMetaInfo {
2221                    num_docs: doc_count,
2222                    ancestors: old_ids.to_vec(),
2223                    generation,
2224                    reordered,
2225                    bp_converged,
2226                    bp_unconverged_passes,
2227                }
2228            }
2229            ReplacementLayout::PreserveSingleSource => {
2230                let [source_id] = old_ids else {
2231                    return Err(Error::Internal(
2232                        "layout-preserving replacement requires exactly one source".into(),
2233                    ));
2234                };
2235                let mut source = st
2236                    .metadata
2237                    .segment_metas
2238                    .get(source_id)
2239                    .cloned()
2240                    .ok_or_else(|| {
2241                        Error::Corruption(format!(
2242                            "layout-preserving replacement source {source_id} disappeared"
2243                        ))
2244                    })?;
2245                source.num_docs = doc_count;
2246                source
2247            }
2248        };
2249        let retired_ids = old_ids.to_vec();
2250        let mut next = st.metadata.clone();
2251        for id in old_ids {
2252            next.remove_segment(id);
2253        }
2254        next.add_segment_meta(new_id.clone(), replacement_info);
2255
2256        let directory = Arc::clone(&self.directory);
2257        let tracker = Arc::clone(&self.tracker);
2258        let replacement_refresh = self.replacement_refresh.read().clone();
2259        let index_label = self.schema.index_label().to_owned();
2260        self.run_lifecycle_transaction(async move {
2261            // Durable-before-visible. If persistence fails, old metadata and
2262            // tracker ownership stay intact and source deletion is never armed.
2263            next.save(directory.as_ref()).await?;
2264            tracker.register(&new_id);
2265            st.metadata = next;
2266
2267            // Keep state locked until retired sources enter the tracker. The
2268            // transaction itself also performs deletion, so cancellation of
2269            // the requesting merge cannot strand pending-deletion ownership.
2270            let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2271            drop(st);
2272            for &segment_id in &ready_to_delete {
2273                if let Err(error) =
2274                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
2275                {
2276                    log::warn!(
2277                        "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2278                        segment_id.to_hex(),
2279                        error,
2280                    );
2281                }
2282            }
2283            tracker.complete_deletion(&ready_to_delete);
2284            refresh_replacement_topology(replacement_refresh, &index_label).await;
2285            Ok(())
2286        })
2287        .await
2288    }
2289
2290    /// Perform the actual merge operation (pure function — no shared state access).
2291    /// `output_segment_id` is pre-generated by the caller so active-operation ownership
2292    /// is installed before any output file is written.
2293    /// Returns (new_segment_id_hex, total_doc_count).
2294    #[allow(clippy::too_many_arguments)]
2295    async fn do_merge(
2296        directory: &D,
2297        schema: &Arc<crate::dsl::Schema>,
2298        segment_ids_to_merge: &[String],
2299        output_segment_id: SegmentId,
2300        term_cache_blocks: usize,
2301        optimization: crate::structures::IndexOptimization,
2302        posting_codec: crate::structures::PostingCodec,
2303        trained: Option<&TrainedVectorStructures>,
2304        reorder_bmp: bool,
2305        granularity: crate::segment::reorder::BpGranularity,
2306        merge_bp_time_budget: Option<std::time::Duration>,
2307        bp_memory_budget_bytes: usize,
2308        reorder_permits: Arc<ReorderConcurrencyGate>,
2309        reorder_priority: ReorderPriority,
2310        cancellation: Arc<AtomicBool>,
2311        bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2312    ) -> MergeTaskResult<(String, u32, bool)> {
2313        let output_hex = output_segment_id.to_hex();
2314        let load_start = std::time::Instant::now();
2315
2316        let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2317        for id_str in segment_ids_to_merge {
2318            let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2319                MergeTaskError::source(
2320                    id_str.clone(),
2321                    Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2322                )
2323            })?;
2324            segment_ids.push(id);
2325        }
2326
2327        // Cheap fail-fast before opening every reader. `join_all` otherwise
2328        // waits for all healthy multi-GB inputs to load even when one source's
2329        // `.meta` is already absent, turning a known-corrupt candidate into a
2330        // large CPU/IO spike before it can be quarantined.
2331        let mut unavailable_sources = Vec::new();
2332        let mut missing_files = Vec::new();
2333        for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2334            let files = SegmentFiles::new(id.0);
2335            let mut source_unavailable = false;
2336            for path in files.mandatory_paths() {
2337                let exists = directory
2338                    .exists(path)
2339                    .await
2340                    .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2341                if !exists {
2342                    source_unavailable = true;
2343                    missing_files.push(format!("{}:{:?}", id_str, path));
2344                }
2345            }
2346            if source_unavailable {
2347                unavailable_sources.push(id_str.clone());
2348            }
2349        }
2350        if !unavailable_sources.is_empty() {
2351            return Err(MergeTaskError::sources(
2352                unavailable_sources,
2353                Error::Corruption(format!(
2354                    "merge sources are missing mandatory files: {}",
2355                    missing_files.join(", ")
2356                )),
2357            ));
2358        }
2359
2360        let schema_arc = Arc::clone(schema);
2361        let futures: Vec<_> = segment_ids
2362            .iter()
2363            .map(|&sid| {
2364                let sch = Arc::clone(&schema_arc);
2365                async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2366            })
2367            .collect();
2368
2369        let results = futures::future::join_all(futures).await;
2370        let mut readers = Vec::with_capacity(results.len());
2371        let mut total_docs = 0u64;
2372        for (i, result) in results.into_iter().enumerate() {
2373            match result {
2374                Ok(r) => {
2375                    total_docs += r.meta().num_docs as u64;
2376                    readers.push(r);
2377                }
2378                Err(e) => {
2379                    log::error!(
2380                        "[merge] index={} Failed to open segment {}: {:?}",
2381                        schema.index_label(),
2382                        segment_ids_to_merge[i],
2383                        e
2384                    );
2385                    return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2386                }
2387            }
2388        }
2389        if total_docs > u32::MAX as u64 {
2390            return Err(Error::Internal(format!(
2391                "Merged segment doc count ({}) exceeds u32::MAX",
2392                total_docs
2393            ))
2394            .into());
2395        }
2396
2397        // Pre-merge validation: verify each source segment's store doc count
2398        // matches its metadata. Catching mismatches early avoids building a
2399        // corrupted merged segment and leaving orphan files on disk.
2400        for (i, reader) in readers.iter().enumerate() {
2401            let meta_docs = reader.meta().num_docs;
2402            let store_docs = reader.store().num_docs();
2403            if store_docs != meta_docs {
2404                return Err(MergeTaskError::source(
2405                    segment_ids_to_merge[i].clone(),
2406                    Error::Corruption(format!(
2407                        "pre-merge validation: segment {} store has {} docs but meta says {}",
2408                        segment_ids_to_merge[i], store_docs, meta_docs
2409                    )),
2410                ));
2411            }
2412        }
2413
2414        log::info!(
2415            "[merge] index={} loaded {} segment readers in {:.1}s",
2416            schema.index_label(),
2417            readers.len(),
2418            load_start.elapsed().as_secs_f64()
2419        );
2420
2421        let merger = SegmentMerger::new(Arc::clone(schema))
2422            .with_posting_config(optimization, posting_codec)
2423            .with_bmp_reorder(reorder_bmp)
2424            .with_granularity(granularity)
2425            .with_bp_budget(crate::segment::BpBudget {
2426                min_partition_docs: None,
2427                time_budget: merge_bp_time_budget,
2428            })
2429            .with_cancellation(cancellation)
2430            .with_bp_memory_budget(bp_memory_budget_bytes)
2431            .with_reorder_permits(reorder_permits)
2432            .with_reorder_priority(reorder_priority)
2433            .with_background_pool(bg_cpu_pool);
2434
2435        log::info!(
2436            "[merge] index={} {} segments -> {} (trained={})",
2437            schema.index_label(),
2438            segment_ids_to_merge.len(),
2439            output_hex,
2440            trained.map_or(0, |t| t.centroids.len()),
2441        );
2442
2443        let (_merged_meta, merge_stats) = merger
2444            .merge(directory, &readers, output_segment_id, trained)
2445            .await
2446            .map_err(|error| {
2447                if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2448                    // The merge has already opened every input successfully;
2449                    // a structural/serialization failure is deterministic for
2450                    // this candidate. Attribute all inputs rather than running
2451                    // the same multi-GB rewrite forever. This is deliberately
2452                    // not used for I/O errors, which may be transient/output-side.
2453                    MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2454                } else {
2455                    MergeTaskError::from(error)
2456                }
2457            })?;
2458        let bp_converged = merge_stats.bp_converged;
2459        if !bp_converged {
2460            log::info!(
2461                "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2462                 the background optimizer deepens it later",
2463                schema.index_label(),
2464            );
2465        }
2466
2467        log::info!(
2468            "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2469            schema.index_label(),
2470            load_start.elapsed().as_secs_f64(),
2471            readers.len(),
2472            total_docs,
2473        );
2474
2475        Ok((output_hex, total_docs as u32, bp_converged))
2476    }
2477
2478    /// Drain all in-flight merge tasks safely.
2479    ///
2480    /// Merge/reorder writers use synchronous block-in-place sections, so an
2481    /// abort request takes effect only after the owned writer reaches its next
2482    /// await. Draining the complete task remains necessary before index
2483    /// deletion or orphan cleanup can safely proceed.
2484    ///
2485    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2486    /// handles to `merge_handles` so later drains still see in-flight merges.
2487    pub async fn abort_merges(&self) {
2488        loop {
2489            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2490            if handles.is_empty() {
2491                return;
2492            }
2493            while let Some(result) = handles.join_next().await {
2494                if let Err(error) = result
2495                    && error.is_panic()
2496                {
2497                    log::error!(
2498                        "[merge] index={} background task panicked while draining: {}",
2499                        self.schema.index_label(),
2500                        error
2501                    );
2502                }
2503            }
2504        }
2505    }
2506
2507    /// Wait for all current in-flight merges to complete.
2508    ///
2509    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2510    /// handles to `merge_handles` so later drains still see in-flight merges.
2511    pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2512        let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2513        while handles.join_next().await.is_some() {}
2514    }
2515
2516    /// Wait for all eligible merges to complete, including cascading merges.
2517    ///
2518    /// Drains current handles, then loops. Each completed merge auto-triggers
2519    /// `maybe_merge` (which pushes new handles) before its JoinHandle resolves,
2520    /// so by the time `join_next` returns all cascading handles are registered.
2521    ///
2522    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2523    /// handles to `merge_handles` so later drains still see in-flight merges.
2524    pub async fn wait_for_all_merges(self: &Arc<Self>) {
2525        loop {
2526            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2527            if handles.is_empty() {
2528                break;
2529            }
2530            while handles.join_next().await.is_some() {}
2531        }
2532    }
2533
2534    /// Complete the second half of shutdown after the owning `IndexWriter`
2535    /// has been dropped. This drains tracked merges and then waits for every
2536    /// remaining guard, including optimizer reorders that are intentionally
2537    /// launched outside the writer lock.
2538    pub async fn wait_for_shutdown(self: &Arc<Self>) {
2539        self.wait_for_all_merges().await;
2540        self.active_operations.wait_until_idle().await;
2541        loop {
2542            let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2543            if handles.is_empty() {
2544                break;
2545            }
2546            for handle in handles {
2547                if let Err(error) = handle.await
2548                    && error.is_panic()
2549                {
2550                    log::error!(
2551                        "[segment_cleanup] index={} task panicked while draining: {}",
2552                        self.schema.index_label(),
2553                        error
2554                    );
2555                }
2556            }
2557        }
2558    }
2559
2560    /// Force merge all segments into one (packing up to the u32 document
2561    /// format limit per output).
2562    ///
2563    /// An explicit force merge is a full compaction: unlike background
2564    /// merges, it deliberately ignores the policy's `max_segment_docs` cap
2565    /// (which exists to bound background BP/merge cost, not to keep an index
2566    /// permanently split). Only the u32 doc-id format limit can force more
2567    /// than one output, and that outcome is logged loudly.
2568    ///
2569    /// Each batch is registered in `active_operations` via an RAII guard to prevent
2570    /// `maybe_merge` from spawning a conflicting background merge.
2571    pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2572        self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2573            .await
2574    }
2575
2576    /// Force merge while refreshing long-lived segment snapshots after the
2577    /// initial background-merge drain and every durable replacement.
2578    ///
2579    /// Long-lived consumers such as the primary-key index and cached
2580    /// `IndexReader` hold segment snapshots. Refreshing them only after the
2581    /// complete force merge retains every retired source file for the entire
2582    /// operation, which can temporarily double a large index on disk. The
2583    /// hook lets the owning `IndexWriter`/server advance those snapshots after
2584    /// each batch while the force merge remains otherwise memory bounded.
2585    pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2586        self: &Arc<Self>,
2587        mut refresh_snapshots: F,
2588    ) -> Result<()>
2589    where
2590        F: FnMut() -> Fut,
2591        Fut: std::future::Future<Output = Result<()>>,
2592    {
2593        // Conflicting owners that never appear in `merge_handles` (background
2594        // reorders, a concurrent force-merge) can hold a batch segment for
2595        // minutes to hours; retrying without parking would busy-spin a runtime
2596        // worker and hammer the state mutex for that whole window.
2597        const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2598            std::time::Duration::from_millis(100);
2599        // When every remaining mergeable segment is owned by another
2600        // operation (e.g. background BP reorders), there is nothing to do but
2601        // wait for a release; those passes run for minutes, so poll slowly.
2602        const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2603
2604        let (_force_merge_activity, policy_segment_docs) = {
2605            let st = self.state.lock().await;
2606            // `maybe_merge` registers and publishes every task handle under
2607            // this same state lock. Raising the barrier here therefore makes
2608            // the subsequent drain race-free.
2609            self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2610            (
2611                ForceMergeActivityGuard(&self.force_merge_active),
2612                st.merge_policy.max_segment_docs(),
2613            )
2614        };
2615
2616        // Wait for all in-flight background merges (including cascading)
2617        // before starting forced merges to avoid try_register conflicts.
2618        let background_merges = self
2619            .merge_handles
2620            .lock()
2621            .iter()
2622            .filter(|handle| !handle.is_finished())
2623            .count();
2624        if background_merges > 0 {
2625            log::info!(
2626                "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2627                self.schema.index_label(),
2628                background_merges,
2629            );
2630        }
2631        let drain_start = std::time::Instant::now();
2632        self.wait_for_all_merges().await;
2633        if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2634            log::info!(
2635                "[force_merge] index={} drained background merges in {:.1}s",
2636                self.schema.index_label(),
2637                drain_start.elapsed().as_secs_f64(),
2638            );
2639        }
2640
2641        // A background merge may have published after the writer's preceding
2642        // commit refresh but before the drain completed. Advance PK/search
2643        // snapshots now so its retired sources do not survive until the first
2644        // forced replacement (or forever when no batch is needed).
2645        let refresh_start = std::time::Instant::now();
2646        refresh_snapshots().await?;
2647        if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2648            log::info!(
2649                "[force_merge] index={} initial snapshot refresh took {:.1}s",
2650                self.schema.index_label(),
2651                refresh_start.elapsed().as_secs_f64(),
2652            );
2653        }
2654
2655        // Outputs produced by a completed bin are already maximal according
2656        // to the BFD plan. Excluding them from later replanning prevents an
2657        // expensive final BP output from being selected and reordered again.
2658        let mut completed_outputs = HashSet::new();
2659        // One INFO line per wait episode, not per 1s poll; DEBUG afterwards.
2660        let mut logged_held_wait = false;
2661
2662        loop {
2663            if !self.active_operations.is_accepting() {
2664                return Err(Error::IndexClosed);
2665            }
2666
2667            let segments: Vec<(String, u32)> = {
2668                let st = self.state.lock().await;
2669                st.metadata
2670                    .segment_metas
2671                    .iter()
2672                    .filter(|(id, _)| !completed_outputs.contains(*id))
2673                    .map(|(id, info)| (id.clone(), info.num_docs))
2674                    .collect()
2675            };
2676
2677            // Route around segments owned by active operations instead of
2678            // retrying the same conflicting batch. Group ownership is claimed
2679            // before foreground capacity, so an in-flight optimizer/retrain
2680            // can finish without a capacity/ownership lock inversion.
2681            let active_ids = self.active_operations.snapshot();
2682            let held = segments
2683                .iter()
2684                .filter(|(id, _)| active_ids.contains(id))
2685                .count();
2686            let free_segments: Vec<_> = segments
2687                .into_iter()
2688                .filter(|(id, _)| !active_ids.contains(id))
2689                .collect();
2690            // Every segment format and doc map uses u32 document IDs.
2691            // An explicit force merge packs up to that format limit: the
2692            // policy's `max_segment_docs` bounds *background* merge cost and
2693            // must not silently leave a forced compaction above one segment
2694            // (observed in prod: an 8.8M-doc index stuck at 2 segments under
2695            // a 5M-doc policy cap).
2696            let max_docs = u64::from(u32::MAX);
2697            let planned_groups = plan_force_merge_groups(free_segments, max_docs);
2698            if let Some(cap) = policy_segment_docs {
2699                for group in planned_groups
2700                    .iter()
2701                    .filter(|group| group.segments.len() >= 2)
2702                    .filter(|group| group.total_docs > u64::from(cap))
2703                {
2704                    log::warn!(
2705                        "[force_merge] index={} output of {} docs intentionally exceeds the \
2706                         background merge policy cap of {} docs (force merge compacts to the \
2707                         u32 format limit)",
2708                        self.schema.index_label(),
2709                        group.total_docs,
2710                        cap,
2711                    );
2712                }
2713            }
2714            let next_group = planned_groups
2715                .into_iter()
2716                .find(|group| group.segments.len() >= 2);
2717
2718            let Some(group) = next_group else {
2719                if held == 0 {
2720                    if !completed_outputs.is_empty() {
2721                        // A segment held while earlier groups ran may now fit
2722                        // with one of those outputs. Reconsider completed bins
2723                        // once before declaring convergence. In the normal
2724                        // no-conflict path BFD groups are pairwise maximal, so
2725                        // this extra planning round performs no merge.
2726                        completed_outputs.clear();
2727                        continue;
2728                    }
2729                    // Every remaining free segment is either already at the
2730                    // u32 format limit or cannot be paired without exceeding
2731                    // it. More than one leftover segment is an exceptional,
2732                    // loudly-reported outcome — never a silent policy effect.
2733                    let remaining = {
2734                        let st = self.state.lock().await;
2735                        st.metadata.segment_metas.len()
2736                    };
2737                    if remaining > 1 {
2738                        log::warn!(
2739                            "[force_merge] index={} finished with {} segments: combined \
2740                             document count exceeds the u32 segment format limit, so a \
2741                             single output is impossible",
2742                            self.schema.index_label(),
2743                            remaining,
2744                        );
2745                    }
2746                    // A reorder that was already running when foreground
2747                    // admission began may have published after the preceding
2748                    // callback. Reconcile external readers one final time so
2749                    // they cannot retain its retired source indefinitely.
2750                    refresh_snapshots().await?;
2751                    return Ok(());
2752                }
2753                if !logged_held_wait {
2754                    log::info!(
2755                        "[force_merge] index={} waiting: {} segment(s) held by active \
2756                         merge/reorder operations, no free group can merge",
2757                        self.schema.index_label(),
2758                        held
2759                    );
2760                    logged_held_wait = true;
2761                } else {
2762                    log::debug!(
2763                        "[force_merge] index={} still waiting on {} held segment(s)",
2764                        self.schema.index_label(),
2765                        held
2766                    );
2767                }
2768                #[cfg(test)]
2769                self.force_merge_conflict_retries
2770                    .fetch_add(1, Ordering::Relaxed);
2771                tokio::select! {
2772                    biased;
2773                    () = self.active_operations.wait_for_shutdown() => {
2774                        return Err(Error::IndexClosed);
2775                    }
2776                    () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2777                }
2778                continue;
2779            };
2780            logged_held_wait = false;
2781
2782            // Claim the complete final group and all of its future output IDs
2783            // up front. This freezes the hierarchy while allowing every
2784            // durable intermediate replacement to release its source files.
2785            let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2786            let output_ids: Vec<_> = (0..hierarchy.steps.len())
2787                .map(|_| SegmentId::new())
2788                .collect();
2789            let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2790            let mut all_ids = source_ids.clone();
2791            all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2792            let group_guard = {
2793                let st = self.state.lock().await;
2794                source_ids
2795                    .iter()
2796                    .all(|id| st.metadata.has_segment(id))
2797                    .then(|| self.active_operations.try_register(all_ids))
2798                    .flatten()
2799            };
2800            let _group_guard = match group_guard {
2801                Some(guard) => guard,
2802                None if !self.active_operations.is_accepting() => {
2803                    return Err(Error::IndexClosed);
2804                }
2805                None => {
2806                    #[cfg(test)]
2807                    self.force_merge_conflict_retries
2808                        .fetch_add(1, Ordering::Relaxed);
2809                    log::debug!(
2810                        "[force_merge] index={} group lost a registration race, replanning",
2811                        self.schema.index_label()
2812                    );
2813                    let had_tracked_merges = !self.merge_handles.lock().is_empty();
2814                    self.wait_for_merging_thread().await;
2815                    if !had_tracked_merges {
2816                        tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2817                    }
2818                    continue;
2819                }
2820            };
2821
2822            log::info!(
2823                "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
2824                self.schema.index_label(),
2825                group.segments.len(),
2826                group.total_docs,
2827                output_ids.len(),
2828            );
2829
2830            // Claiming the complete group before any capacity wait prevents a
2831            // vector-generation pause from waiting for a global slot retained
2832            // by force merge while force merge waits for that pause to end.
2833            //
2834            // Background merges acquire global merge capacity before the
2835            // shared BP gate. Preserve that order here: wait for one group
2836            // slot while background BP remains admitted, then pause new BP
2837            // passes. The retained slot guarantees this hierarchy can make
2838            // progress even if other indexes subsequently park at the gate.
2839            let group_global_merge_permit = if self.reorder_on_merge {
2840                let capacity_start = std::time::Instant::now();
2841                let permit = tokio::select! {
2842                    biased;
2843                    () = self.active_operations.wait_for_shutdown() => {
2844                        return Err(Error::IndexClosed);
2845                    }
2846                    permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2847                        permit.map_err(|_| {
2848                            Error::Internal(
2849                                "global background merge scheduler is closed".into(),
2850                            )
2851                        })?
2852                    }
2853                };
2854                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2855                    log::info!(
2856                        "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
2857                        self.schema.index_label(),
2858                        capacity_start.elapsed().as_secs_f64(),
2859                    );
2860                }
2861                Some(permit)
2862            } else {
2863                None
2864            };
2865            let _foreground_reorder = if self.reorder_on_merge {
2866                log::info!(
2867                    "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
2868                    self.schema.index_label(),
2869                    self.reorder_permits.limit(),
2870                );
2871                let admission_start = std::time::Instant::now();
2872                let guard = Arc::clone(&self.reorder_permits)
2873                    .begin_foreground()
2874                    .await
2875                    .map_err(|_| {
2876                        Error::Internal("background reorder scheduler is closed".into())
2877                    })?;
2878                if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2879                    log::info!(
2880                        "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
2881                        self.schema.index_label(),
2882                        admission_start.elapsed().as_secs_f64(),
2883                    );
2884                }
2885                Some(guard)
2886            } else {
2887                None
2888            };
2889
2890            let source_count = group.segments.len();
2891            let mut nodes: Vec<Option<(String, u32)>> =
2892                group.segments.into_iter().map(Some).collect();
2893            nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2894            for (step_index, step) in hierarchy.steps.iter().enumerate() {
2895                let final_pass = step_index + 1 == hierarchy.steps.len();
2896                let mut batch_entries = Vec::with_capacity(step.inputs.len());
2897                for &node in &step.inputs {
2898                    let entry = nodes
2899                        .get_mut(node)
2900                        .and_then(Option::take)
2901                        .expect("force-merge hierarchy must reference an available node");
2902                    batch_entries.push(entry);
2903                }
2904                let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2905                let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2906                let output_id = output_ids[step_index];
2907
2908                let capacity_start = std::time::Instant::now();
2909                let step_global_merge_permit = if group_global_merge_permit.is_none() {
2910                    Some(tokio::select! {
2911                        biased;
2912                        () = self.active_operations.wait_for_shutdown() => {
2913                            return Err(Error::IndexClosed);
2914                        }
2915                        permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2916                            permit.map_err(|_| {
2917                                Error::Internal(
2918                                    "global background merge scheduler is closed".into(),
2919                                )
2920                            })?
2921                        }
2922                    })
2923                } else {
2924                    None
2925                };
2926                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2927                    log::info!(
2928                        "[force_merge] index={} waited {:.1}s for global merge capacity",
2929                        self.schema.index_label(),
2930                        capacity_start.elapsed().as_secs_f64(),
2931                    );
2932                }
2933
2934                // Intermediate reductions are streaming block-copy merges.
2935                // Only the final pass pays BP, so a group with hundreds of
2936                // tiny sources never reorders the same documents repeatedly.
2937                let reorder_bmp = final_pass && self.reorder_on_merge;
2938                log::info!(
2939                    "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
2940                    self.schema.index_label(),
2941                    if final_pass {
2942                        "final"
2943                    } else {
2944                        "fan-in reduction"
2945                    },
2946                    batch.len(),
2947                    batch_docs,
2948                    reorder_bmp,
2949                );
2950                let (new_segment_id, total_docs, _) = self
2951                    .merge_and_replace_registered(
2952                        &batch,
2953                        output_id,
2954                        reorder_bmp,
2955                        ReorderPriority::Foreground,
2956                    )
2957                    .await
2958                    .map_err(|error| error.error)?;
2959                drop(step_global_merge_permit);
2960
2961                // Advance PK/read snapshots after every hierarchy level so
2962                // retired sources do not accumulate until the final output.
2963                let refresh_start = std::time::Instant::now();
2964                refresh_snapshots().await?;
2965                if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2966                    log::info!(
2967                        "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
2968                        self.schema.index_label(),
2969                        refresh_start.elapsed().as_secs_f64(),
2970                    );
2971                }
2972
2973                let output_node = source_count + step_index;
2974                debug_assert!(nodes[output_node].is_none());
2975                nodes[output_node] = Some((new_segment_id, total_docs));
2976            }
2977            let (root_id, _) = nodes[hierarchy.root]
2978                .take()
2979                .expect("force-merge hierarchy must produce its root");
2980            debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2981            completed_outputs.insert(root_id);
2982        }
2983    }
2984
2985    fn segment_needs_vector_rewrite(
2986        &self,
2987        schema: &crate::dsl::Schema,
2988        reader: &SegmentReader,
2989        field_ids: &[u32],
2990        trained: &TrainedVectorStructures,
2991        rewrite_existing: bool,
2992    ) -> Result<bool> {
2993        for &field_id in field_ids {
2994            let flat = reader.flat_vectors().get(&field_id);
2995            let ann = reader.vector_indexes().get(&field_id);
2996            if ann.is_some() && flat.is_none() {
2997                return Err(Error::Corruption(format!(
2998                    "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2999                    reader.meta().id,
3000                )));
3001            }
3002
3003            let Some(flat) = flat else {
3004                continue;
3005            };
3006            if flat.num_vectors == 0 {
3007                continue;
3008            }
3009            if rewrite_existing {
3010                return Ok(true);
3011            }
3012            let field = crate::dsl::Field(field_id);
3013            let entry = schema.get_field_entry(field).ok_or_else(|| {
3014                Error::Corruption(format!(
3015                    "segment {:032x} references unknown vector field {field_id}",
3016                    reader.meta().id,
3017                ))
3018            })?;
3019            let current = match entry.field_type {
3020                // TQ payloads carry no trained generation, so an existing one
3021                // is always current; a tq field must never be staged for a
3022                // vector-generation rewrite.
3023                crate::dsl::FieldType::DenseVector
3024                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3025                        config.index_type == crate::dsl::VectorIndexType::Tq
3026                    }) =>
3027                {
3028                    matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
3029                }
3030                crate::dsl::FieldType::DenseVector
3031                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3032                        config.index_type == crate::dsl::VectorIndexType::IvfTq
3033                    }) =>
3034                {
3035                    let config = entry
3036                        .dense_vector_config
3037                        .as_ref()
3038                        .expect("matched IVF-TQ configuration");
3039                    match (ann, trained.centroids.get(&field_id)) {
3040                        (
3041                            Some(crate::segment::VectorIndex::IvfTq { index, .. }),
3042                            Some(centroids),
3043                        ) => {
3044                            let header = index.get().header();
3045                            crate::structures::is_ivf_tq_cosine_generation(centroids.version)
3046                                && crate::structures::is_ivf_tq_cosine_generation(
3047                                    header.quantizer_version,
3048                                )
3049                                && header.dim == config.dim
3050                                && header.num_clusters == centroids.num_clusters
3051                                && header.quantizer_version == centroids.version
3052                                && header.codebook_version
3053                                    == crate::structures::vector::quantization::tq_expected_fingerprint(
3054                                        config.dim,
3055                                    )
3056                                && header.routing == config.ivf_routing
3057                        }
3058                        (None, None) => true,
3059                        _ => false,
3060                    }
3061                }
3062                crate::dsl::FieldType::DenseVector
3063                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3064                        config.index_type == crate::dsl::VectorIndexType::Scann
3065                    }) =>
3066                {
3067                    match (ann, trained.scann_artifacts.get(&field_id)) {
3068                        (Some(crate::segment::VectorIndex::ScannAh(index)), Some(artifact)) => {
3069                            index
3070                                .get()
3071                                .validate_scann_generation(
3072                                    artifact.config(),
3073                                    artifact.generation(),
3074                                    artifact.artifact_id(),
3075                                )
3076                                .is_ok()
3077                        }
3078                        (None, None) => true,
3079                        _ => false,
3080                    }
3081                }
3082                // Only ivf_tq dense fields are trainable; other dense
3083                // index types never reach a vector-generation rewrite.
3084                crate::dsl::FieldType::DenseVector => false,
3085                crate::dsl::FieldType::BinaryDenseVector
3086                    if entry
3087                        .binary_dense_vector_config
3088                        .as_ref()
3089                        .is_some_and(|config| {
3090                            config.index_type == crate::dsl::BinaryIndexType::Scann
3091                        }) =>
3092                {
3093                    match (ann, trained.scann_artifacts.get(&field_id)) {
3094                        (Some(crate::segment::VectorIndex::ScannBinary(index)), Some(artifact)) => {
3095                            index
3096                                .get()
3097                                .validate_scann_generation(
3098                                    artifact.config(),
3099                                    artifact.generation(),
3100                                    artifact.artifact_id(),
3101                                )
3102                                .is_ok()
3103                        }
3104                        (None, None) => true,
3105                        _ => false,
3106                    }
3107                }
3108                crate::dsl::FieldType::BinaryDenseVector => matches!(
3109                    (ann, trained.binary_quantizers.get(&field_id)),
3110                    (Some(crate::segment::VectorIndex::BinaryIvf(_)), Some(_)) | (None, None)
3111                ),
3112                _ => false,
3113            };
3114            if !current {
3115                return Ok(true);
3116            }
3117        }
3118        Ok(false)
3119    }
3120
3121    async fn acquire_vector_rewrite_capacity(
3122        &self,
3123    ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
3124        let global = tokio::select! {
3125            biased;
3126            () = self.active_operations.wait_for_shutdown() => {
3127                return Err(Error::IndexClosed);
3128            }
3129            permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3130                permit.map_err(|_| Error::Internal(
3131                    "global background merge scheduler is closed".into()
3132                ))?
3133            }
3134        };
3135        let local = tokio::select! {
3136            biased;
3137            () = self.active_operations.wait_for_shutdown() => {
3138                return Err(Error::IndexClosed);
3139            }
3140            permit = Arc::clone(&self.merge_permits).acquire_owned() => {
3141                permit.map_err(|_| Error::Internal(
3142                    "background merge scheduler is closed".into()
3143                ))?
3144            }
3145        };
3146        Ok((global, local))
3147    }
3148
3149    async fn build_vector_replacement(
3150        self: &Arc<Self>,
3151        schemas: (&Arc<crate::dsl::Schema>, &Arc<crate::dsl::Schema>),
3152        segment_id: &str,
3153        source_id: SegmentId,
3154        output_id: SegmentId,
3155        trained: &TrainedVectorStructures,
3156        failure_context: &'static str,
3157    ) -> Result<(String, u32, OutputCleanupGuard)> {
3158        let mut cleanup = self.output_cleanup_guard(output_id);
3159        match crate::segment::reorder::rewrite_vector_segment(
3160            self.directory.as_ref(),
3161            schemas,
3162            source_id,
3163            output_id,
3164            self.term_cache_blocks,
3165            trained,
3166            Some(self.background_cpu_pool()),
3167        )
3168        .await
3169        {
3170            Ok((new_id, doc_count)) => {
3171                self.validate_completed_segment(&new_id, doc_count).await?;
3172                Ok((new_id, doc_count, cleanup))
3173            }
3174            Err(error) => {
3175                self.delete_output_if_unregistered(output_id, failure_context)
3176                    .await;
3177                cleanup.disarm();
3178                if is_deterministic_source_error(&error) {
3179                    self.quarantine_segment(segment_id, &error);
3180                }
3181                Err(error)
3182            }
3183        }
3184    }
3185
3186    /// Build every required replacement segment without exposing any of them.
3187    /// The returned guards keep both source and output generations alive until
3188    /// [`Self::publish_vector_generation`] commits the complete set.
3189    pub(crate) async fn stage_vector_generation(
3190        self: &Arc<Self>,
3191        _artifact_update: &VectorArtifactUpdateGuard,
3192        segment_ids: &[String],
3193        field_ids: &[u32],
3194        trained: Arc<TrainedVectorStructures>,
3195        rewrite_existing: bool,
3196    ) -> Result<Vec<StagedVectorSegment>> {
3197        let schema = self.published_generation().schema.clone();
3198        self.stage_vector_generation_with_schema(
3199            _artifact_update,
3200            segment_ids,
3201            field_ids,
3202            trained,
3203            rewrite_existing,
3204            schema,
3205        )
3206        .await
3207    }
3208
3209    pub(crate) async fn stage_vector_generation_with_schema(
3210        self: &Arc<Self>,
3211        _artifact_update: &VectorArtifactUpdateGuard,
3212        segment_ids: &[String],
3213        field_ids: &[u32],
3214        trained: Arc<TrainedVectorStructures>,
3215        rewrite_existing: bool,
3216        schema: Arc<crate::dsl::Schema>,
3217    ) -> Result<Vec<StagedVectorSegment>> {
3218        if !self.vector_artifact_update.load(Ordering::Acquire) {
3219            return Err(Error::Internal(
3220                "cannot stage a vector generation without an exclusive update lease".into(),
3221            ));
3222        }
3223
3224        let source_schema = self.published_generation().schema.clone();
3225        let mut staged = Vec::new();
3226        for segment_id in segment_ids {
3227            if self.quarantined_segments.lock().contains(segment_id) {
3228                return Err(Error::Corruption(format!(
3229                    "segment {segment_id} is quarantined after a deterministic source failure"
3230                )));
3231            }
3232            let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3233                Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3234            })?;
3235
3236            // A single rewrite may hold several gigabytes while assigning all
3237            // vectors. Reuse the ordinary local and process-wide merge bounds.
3238            let _capacity = self.acquire_vector_rewrite_capacity().await?;
3239
3240            let output_id = SegmentId::new();
3241            let output_hex = output_id.to_hex();
3242            let operation = {
3243                let st = self.state.lock().await;
3244                if !st.metadata.has_segment(segment_id) {
3245                    return Err(Error::Corruption(format!(
3246                        "vector generation source {segment_id} disappeared while lifecycle work was paused"
3247                    )));
3248                }
3249                self.active_operations
3250                    .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3251            }
3252            .ok_or_else(|| {
3253                if self.active_operations.is_accepting() {
3254                    Error::Internal(format!(
3255                        "vector generation could not claim stable source {segment_id}"
3256                    ))
3257                } else {
3258                    Error::IndexClosed
3259                }
3260            })?;
3261
3262            let reader = SegmentReader::open(
3263                self.directory.as_ref(),
3264                source_id,
3265                Arc::clone(&source_schema),
3266                self.term_cache_blocks,
3267            )
3268            .await?;
3269            if !self.segment_needs_vector_rewrite(
3270                schema.as_ref(),
3271                &reader,
3272                field_ids,
3273                trained.as_ref(),
3274                rewrite_existing,
3275            )? {
3276                continue;
3277            }
3278            drop(reader);
3279
3280            let (new_id, doc_count, cleanup) = self
3281                .build_vector_replacement(
3282                    (&source_schema, &schema),
3283                    segment_id,
3284                    source_id,
3285                    output_id,
3286                    trained.as_ref(),
3287                    "vector generation staging failure",
3288                )
3289                .await?;
3290            debug_assert_eq!(new_id, output_hex);
3291            let output_reader = SegmentReader::open(
3292                self.directory.as_ref(),
3293                output_id,
3294                Arc::clone(&schema),
3295                self.term_cache_blocks,
3296            )
3297            .await?;
3298            if self.segment_needs_vector_rewrite(
3299                schema.as_ref(),
3300                &output_reader,
3301                field_ids,
3302                trained.as_ref(),
3303                false,
3304            )? {
3305                return Err(Error::Corruption(format!(
3306                    "staged vector segment {new_id} does not match its candidate codebook generation"
3307                )));
3308            }
3309
3310            staged.push(StagedVectorSegment {
3311                source_id: segment_id.clone(),
3312                output_id,
3313                doc_count,
3314                _operation: operation,
3315                cleanup,
3316            });
3317        }
3318        Ok(staged)
3319    }
3320
3321    async fn rewrite_vector_segment_once(
3322        self: &Arc<Self>,
3323        segment_id: &str,
3324        field_ids: &[u32],
3325    ) -> Result<VectorSegmentRewriteOutcome> {
3326        if self.quarantined_segments.lock().contains(segment_id) {
3327            return Err(Error::Corruption(format!(
3328                "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3329            )));
3330        }
3331        let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3332            Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3333        })?;
3334
3335        // Match ordinary merge lock ordering: capacity before lifecycle
3336        // ownership. A vector rewrite can hold several gigabytes while it
3337        // assigns vectors, so it participates in both local and process-wide
3338        // merge limits.
3339        let _capacity = self.acquire_vector_rewrite_capacity().await?;
3340
3341        let output_id = SegmentId::new();
3342        let output_hex = output_id.to_hex();
3343        let all_ids = vec![segment_id.to_owned(), output_hex];
3344        let operation = {
3345            let st = self.state.lock().await;
3346            if !st.metadata.has_segment(segment_id) {
3347                return Ok(VectorSegmentRewriteOutcome::SourceGone);
3348            }
3349            self.active_operations.try_register(all_ids)
3350        };
3351        let _operation = match operation {
3352            Some(operation) => operation,
3353            None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3354            None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3355        };
3356
3357        let Some(trained) = self.trained_for_segment_build() else {
3358            return Ok(VectorSegmentRewriteOutcome::Deferred);
3359        };
3360        let schema = self.published_generation().schema.clone();
3361
3362        let reader = SegmentReader::open(
3363            self.directory.as_ref(),
3364            source_id,
3365            Arc::clone(&schema),
3366            self.term_cache_blocks,
3367        )
3368        .await?;
3369        if !self.segment_needs_vector_rewrite(
3370            schema.as_ref(),
3371            &reader,
3372            field_ids,
3373            trained.as_ref(),
3374            false,
3375        )? {
3376            return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3377        }
3378        drop(reader);
3379
3380        let (new_id, doc_count, mut output_cleanup) = self
3381            .build_vector_replacement(
3382                (&schema, &schema),
3383                segment_id,
3384                source_id,
3385                output_id,
3386                trained.as_ref(),
3387                "vector rewrite failure",
3388            )
3389            .await?;
3390
3391        if let Err(error) = self
3392            .replace_segments(
3393                &[segment_id.to_owned()],
3394                new_id,
3395                doc_count,
3396                ReplacementLayout::PreserveSingleSource,
3397            )
3398            .await
3399        {
3400            self.delete_output_if_unregistered(output_id, "vector replacement failure")
3401                .await;
3402            output_cleanup.disarm();
3403            return Err(error);
3404        }
3405        output_cleanup.disarm();
3406        Ok(VectorSegmentRewriteOutcome::Rewritten)
3407    }
3408
3409    /// Finalize every committed flat vector segment against the published
3410    /// global ANN generation.
3411    /// Unlike force-merge this handles one segment and segments already at the
3412    /// merge policy's maximum size.
3413    pub(crate) async fn rewrite_vector_segments(
3414        self: &Arc<Self>,
3415        field_ids: &[u32],
3416    ) -> Result<usize> {
3417        if field_ids.is_empty() {
3418            return Ok(0);
3419        }
3420        let mut rewritten = 0usize;
3421        loop {
3422            let segment_ids = self.get_segment_ids().await;
3423            let mut conflicted = false;
3424            let mut changed = false;
3425            for segment_id in segment_ids {
3426                match self
3427                    .rewrite_vector_segment_once(&segment_id, field_ids)
3428                    .await?
3429                {
3430                    VectorSegmentRewriteOutcome::Rewritten => {
3431                        rewritten += 1;
3432                        changed = true;
3433                    }
3434                    VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3435                    VectorSegmentRewriteOutcome::Deferred => {
3436                        return Err(Error::Internal(
3437                            "ANN finalization lost the published trained generation".into(),
3438                        ));
3439                    }
3440                    VectorSegmentRewriteOutcome::AlreadyCurrent
3441                    | VectorSegmentRewriteOutcome::SourceGone => {}
3442                }
3443            }
3444            if !conflicted && !changed {
3445                log::info!(
3446                    "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3447                    self.schema.index_label(),
3448                    rewritten,
3449                );
3450                return Ok(rewritten);
3451            }
3452            tokio::select! {
3453                biased;
3454                () = self.active_operations.wait_for_shutdown() => {
3455                    return Err(Error::IndexClosed);
3456                }
3457                () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3458            }
3459        }
3460    }
3461
3462    /// A producer that started in the force-flat phase can commit after the
3463    /// main finalization snapshot. Upgrade exactly those new segments in a
3464    /// tracked background task; ordinary producers already using the current
3465    /// generation are detected and skipped without rewriting.
3466    pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3467        if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3468            return;
3469        }
3470        let manager = Arc::clone(self);
3471        let future = async move {
3472            let field_ids = manager
3473                .read_metadata(|metadata| {
3474                    metadata
3475                        .vector_fields
3476                        .keys()
3477                        .filter(|field_id| metadata.is_field_built(**field_id))
3478                        .copied()
3479                        .collect::<Vec<_>>()
3480                })
3481                .await;
3482            for segment_id in segment_ids {
3483                loop {
3484                    match manager
3485                        .rewrite_vector_segment_once(&segment_id, &field_ids)
3486                        .await
3487                    {
3488                        Ok(VectorSegmentRewriteOutcome::Conflict) => {
3489                            tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3490                        }
3491                        Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3492                        Ok(_) => break,
3493                        Err(error) => {
3494                            log::error!(
3495                                "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3496                                manager.schema.index_label(),
3497                                segment_id,
3498                                error,
3499                            );
3500                            break;
3501                        }
3502                    }
3503                }
3504            }
3505        };
3506        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3507            log::warn!(
3508                "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3509                self.schema.index_label()
3510            );
3511            return;
3512        };
3513        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3514            log::warn!(
3515                "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3516                self.schema.index_label()
3517            );
3518        }
3519    }
3520
3521    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
3522    ///
3523    /// Each segment is individually rebuilt with reordered BMP blocks.
3524    /// Non-BMP fields are copied unchanged via streaming file copy.
3525    ///
3526    /// Uses active-operation ownership to prevent concurrent work on the same segment.
3527    pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3528        self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3529            .await
3530    }
3531
3532    /// Reorder all segments while advancing long-lived snapshots after the
3533    /// background-merge drain and every durable replacement.
3534    pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3535        self: &Arc<Self>,
3536        mut refresh_snapshots: F,
3537    ) -> Result<()>
3538    where
3539        F: FnMut() -> Fut,
3540        Fut: std::future::Future<Output = Result<()>>,
3541    {
3542        self.wait_for_all_merges().await;
3543        refresh_snapshots().await?;
3544        let segment_ids = self.get_segment_ids().await;
3545
3546        if segment_ids.is_empty() {
3547            log::info!(
3548                "[reorder] index={} no segments to reorder",
3549                self.schema.index_label()
3550            );
3551            return Ok(());
3552        }
3553
3554        log::info!(
3555            "[reorder] index={} reordering {} segments",
3556            self.schema.index_label(),
3557            segment_ids.len()
3558        );
3559
3560        for seg_id in segment_ids {
3561            match self
3562                .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3563                .await
3564            {
3565                Ok(true) => refresh_snapshots().await?,
3566                Ok(false) => log::warn!(
3567                    "[reorder] index={} segment {} skipped (in merge)",
3568                    self.schema.index_label(),
3569                    seg_id
3570                ),
3571                Err(e) => return Err(e),
3572            }
3573        }
3574
3575        // A segment skipped because another lifecycle owner held it may have
3576        // been replaced after the preceding callback.
3577        refresh_snapshots().await?;
3578        log::info!(
3579            "[reorder] index={} all segments reordered",
3580            self.schema.index_label()
3581        );
3582        Ok(())
3583    }
3584
3585    /// Get segment IDs that have not been reordered yet.
3586    ///
3587    /// Excludes segments currently involved in a merge or reorder operation
3588    /// to avoid wasted work (the optimizer would skip them anyway).
3589    pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3590        self.unreordered_segments()
3591            .await
3592            .into_iter()
3593            .map(|(id, _)| id)
3594            .collect()
3595    }
3596
3597    /// Segments never reordered, with doc counts — for the optimizer to pick
3598    /// a size-appropriate BP budget.
3599    pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3600        let quarantined = self.quarantined_segments.lock().clone();
3601        let paused = self.paused_reorder_segments();
3602        let st = self.state.lock().await;
3603        let active_ids = self.active_operations.snapshot();
3604        st.metadata
3605            .segment_metas
3606            .iter()
3607            .filter(|(id, info)| {
3608                !info.reordered
3609                    && info.bp_converged
3610                    && !active_ids.contains(*id)
3611                    && !quarantined.contains(*id)
3612                    && !paused.contains(*id)
3613            })
3614            .map(|(id, info)| (id.clone(), info.num_docs))
3615            .collect()
3616    }
3617
3618    /// Segments whose last BP pass hit its wall-clock budget before finishing
3619    /// (`bp_converged == false`). A warm-started follow-up pass deepens the
3620    /// ordering; the optimizer revisits these at low priority.
3621    pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3622        self.unconverged_segments_below(u32::MAX)
3623            .await
3624            .into_iter()
3625            .map(|(id, docs, _)| (id, docs))
3626            .collect()
3627    }
3628
3629    /// Unconverged segments still below a hard replacement-lineage work
3630    /// bound. Includes the persisted attempt count for scheduler diagnostics.
3631    pub async fn unconverged_segments_below(
3632        &self,
3633        max_unconverged_passes: u32,
3634    ) -> Vec<(String, u32, u32)> {
3635        let quarantined = self.quarantined_segments.lock().clone();
3636        let paused = self.paused_reorder_segments();
3637        let st = self.state.lock().await;
3638        let active_ids = self.active_operations.snapshot();
3639        st.metadata
3640            .segment_metas
3641            .iter()
3642            .filter(|(id, info)| {
3643                !info.bp_converged
3644                    && info.bp_unconverged_passes < max_unconverged_passes
3645                    && !active_ids.contains(*id)
3646                    && !quarantined.contains(*id)
3647                    && !paused.contains(*id)
3648            })
3649            .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3650            .collect()
3651    }
3652
3653    /// Granularity for a BP pass whose sources are `ids`: `Records` when any
3654    /// source carries unconverged BP debt, `Auto` otherwise.
3655    ///
3656    /// Alignment with the depth budget (docs/block-level-reorder.md): an
3657    /// unconverged segment is owed a deepening pass, and the output of this
3658    /// pass will be marked `bp_converged`. `Auto` would measure the partial
3659    /// pass's residual coherence, potentially take the blockwise path — which
3660    /// cannot deepen record clustering — and end the cascade at partial
3661    /// quality. Only record-level BP discharges the debt.
3662    async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3663        let st = self.state.lock().await;
3664        let deepening = ids.iter().any(|id| {
3665            st.metadata
3666                .segment_metas
3667                .get(id)
3668                .is_some_and(|info| !info.bp_converged)
3669        });
3670        drop(st);
3671        if deepening {
3672            log::info!(
3673                "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3674                self.schema.index_label(),
3675            );
3676            crate::segment::reorder::BpGranularity::Records
3677        } else {
3678            crate::segment::reorder::BpGranularity::Auto
3679        }
3680    }
3681
3682    /// Reorder a single segment via BP. Returns Ok(true) if reordered, Ok(false) if skipped.
3683    ///
3684    /// Non-blocking: operation ownership prevents conflicts with background merges.
3685    /// Copies unchanged files and rebuilds only the sparse file with reordered BMP data.
3686    pub async fn reorder_single_segment(
3687        self: &Arc<Self>,
3688        seg_id: &str,
3689        rayon_pool: Option<Arc<rayon::ThreadPool>>,
3690        bp_budget: crate::segment::BpBudget,
3691    ) -> Result<bool> {
3692        let source_id = SegmentId::from_hex(seg_id)
3693            .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3694        if self.quarantined_segments.lock().contains(seg_id) {
3695            return Err(Error::Corruption(format!(
3696                "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3697                seg_id
3698            )));
3699        }
3700        if self.force_merge_active.load(Ordering::Acquire) > 0 {
3701            log::debug!(
3702                "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3703                self.schema.index_label(),
3704                seg_id,
3705            );
3706            return Ok(false);
3707        }
3708
3709        // Whole-pass concurrency is independent from Rayon width. One pass
3710        // can already use every configured BP worker; this permit bounds the
3711        // much larger forward-index and rewrite working set across indexes,
3712        // optimizer tasks, and merge-time BP.
3713        let reorder_gate = Arc::clone(&self.reorder_permits);
3714        let _reorder_permit = tokio::select! {
3715            biased;
3716            () = self.active_operations.wait_for_shutdown() => {
3717                return Err(Error::IndexClosed);
3718            }
3719            permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3720                permit.map_err(|_| {
3721                    Error::Internal("background reorder scheduler is closed".into())
3722                })?
3723            }
3724        };
3725
3726        let output_id = SegmentId::new();
3727        let output_hex = output_id.to_hex();
3728        let source_ids = [seg_id.to_string()];
3729        let granularity = self.merge_granularity(&source_ids).await;
3730
3731        // Register while holding `state`, matching orphan cleanup's deletion
3732        // barrier. Candidates are scanned ahead of time and can go stale: a
3733        // merge may have consumed this segment since. Its files may even still
3734        // be on disk (deferred deletion under a searcher snapshot) — reordering
3735        // them would re-insert a duplicate copy of docs the merge output holds.
3736        let all_ids = vec![seg_id.to_string(), output_hex];
3737        let (_guard, source_docs, schema) = {
3738            let st = self.state.lock().await;
3739            // Force merge raises this barrier under the same state lock, so
3740            // this second check closes the race with the cheap early check
3741            // above and prevents optimizer starvation between final groups.
3742            if self.force_merge_active.load(Ordering::Acquire) > 0 {
3743                log::debug!(
3744                    "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3745                    self.schema.index_label(),
3746                    seg_id,
3747                );
3748                return Ok(false);
3749            }
3750            let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3751                log::info!(
3752                    "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
3753                    self.schema.index_label(),
3754                    seg_id
3755                );
3756                self.clear_reorder_retry(seg_id);
3757                return Ok(false);
3758            };
3759
3760            let schema = self.published_generation().schema.clone();
3761            match self.active_operations.try_register(all_ids) {
3762                Some(guard) => (guard, source_meta.num_docs, schema),
3763                None if !self.active_operations.is_accepting() => {
3764                    return Err(Error::IndexClosed);
3765                }
3766                None => {
3767                    log::debug!(
3768                        "[optimizer] index={} segment {} in active merge, skipping",
3769                        self.schema.index_label(),
3770                        seg_id
3771                    );
3772                    return Ok(false);
3773                }
3774            }
3775        };
3776
3777        // Fail before allocating a forward index or creating output files.
3778        // Missing mandatory files are deterministic and should remove this
3779        // segment from future optimizer scans, not consume the same CPU every
3780        // interval. Other I/O failures remain retryable.
3781        if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3782            if is_deterministic_source_error(&error) {
3783                self.quarantine_segment(seg_id, &error);
3784            } else if !matches!(&error, Error::IndexClosed) {
3785                self.pause_reorder_retries(seg_id, &error);
3786            }
3787            return Err(error);
3788        }
3789
3790        let mut output_cleanup = self.output_cleanup_guard(output_id);
3791
3792        let reorder_result = crate::segment::reorder::reorder_segment(
3793            self.directory.as_ref(),
3794            &schema,
3795            source_id,
3796            output_id,
3797            self.term_cache_blocks,
3798            self.bp_memory_budget_bytes,
3799            bp_budget,
3800            granularity,
3801            self.optimization,
3802            self.posting_codec,
3803            rayon_pool,
3804            Some(self.active_operations.cancellation_flag()),
3805        )
3806        .await;
3807        let (new_id, total_docs, bp_converged) = match reorder_result {
3808            Ok(v) => v,
3809            Err(e) => {
3810                // A failed pass may have copied tens of GB before dying;
3811                // delete the uncommitted output before propagating.
3812                self.delete_output_if_unregistered(output_id, "reorder failure")
3813                    .await;
3814                output_cleanup.disarm();
3815                if is_deterministic_source_error(&e) {
3816                    self.quarantine_segment(seg_id, &e);
3817                } else if !matches!(&e, Error::IndexClosed) {
3818                    self.pause_reorder_retries(seg_id, &e);
3819                }
3820                return Err(e);
3821            }
3822        };
3823
3824        // A pass with a depth floor above block granularity has, by
3825        // definition, not converged to block-level order — record it as
3826        // unconverged so the optimizer's deepening ladder revisits it with a
3827        // full-depth (warm-started) pass. Depth caps are only used by the
3828        // optimizer's first pass on large segments.
3829        let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3830        if let Err(e) = self
3831            .replace_segments(
3832                &[seg_id.to_string()],
3833                new_id,
3834                total_docs,
3835                ReplacementLayout::BpReordered {
3836                    converged: ladder_converged,
3837                },
3838            )
3839            .await
3840        {
3841            self.delete_output_if_unregistered(output_id, "replacement failure")
3842                .await;
3843            output_cleanup.disarm();
3844            if !matches!(&e, Error::IndexClosed) {
3845                self.pause_reorder_retries(seg_id, &e);
3846            }
3847            return Err(e);
3848        }
3849        output_cleanup.disarm();
3850        self.clear_reorder_retry(seg_id);
3851
3852        Ok(true)
3853    }
3854
3855    /// Clean up orphan segment files not registered in metadata.
3856    ///
3857    /// Reads metadata, active-operation ownership, and snapshot-deferred
3858    /// deletions to determine which segments are legitimate. Filesystem
3859    /// deletion is asynchronous; in-flight outputs and retired sources still
3860    /// held by readers are both protected.
3861    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3862        let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3863
3864        if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3865            for entry in entries {
3866                let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3867                    continue;
3868                };
3869                let Some(rest) = filename.strip_prefix("seg_") else {
3870                    continue;
3871                };
3872                let Some(hex_id) = rest.get(..32) else {
3873                    continue;
3874                };
3875                if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3876                    continue;
3877                }
3878                orphan_files
3879                    .entry(hex_id.to_ascii_lowercase())
3880                    .or_default()
3881                    .push(entry);
3882            }
3883        }
3884
3885        let mut deleted = 0;
3886        for (hex_id, paths) in &orphan_files {
3887            // Revalidate and atomically claim deletion under the same
3888            // state -> active_operations -> tracker order used by publishers.
3889            // The claim lets us release `state` before filesystem I/O: deleting
3890            // a multi-GB orphan must not freeze commits and snapshot acquisition.
3891            let deletion_guard = {
3892                let st = self.state.lock().await;
3893                if st.metadata.has_segment(hex_id) {
3894                    continue;
3895                }
3896                let Some(guard) = self
3897                    .active_operations
3898                    .try_register(vec![hex_id.to_string()])
3899                else {
3900                    continue;
3901                };
3902                if self.tracker.is_deletion_protected(hex_id) {
3903                    drop(guard);
3904                    continue;
3905                }
3906                guard
3907            };
3908
3909            // Delete what was actually discovered, not only the currently
3910            // known SegmentFiles extensions. This also removes partial files
3911            // left by older formats instead of reporting the same orphan on
3912            // every startup forever.
3913            let results =
3914                futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3915                    .await;
3916            let removed = results.into_iter().all(|result| match result {
3917                Ok(()) => true,
3918                Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3919                Err(error) => {
3920                    log::warn!(
3921                        "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
3922                        self.schema.index_label(),
3923                        hex_id,
3924                        error,
3925                    );
3926                    false
3927                }
3928            });
3929            // Releasing this claim is the deletion barrier. No producer can
3930            // adopt the ID while its files are being removed.
3931            drop(deletion_guard);
3932            if removed {
3933                deleted += 1;
3934                log::info!(
3935                    "[segment_cleanup] index={} swept orphan segment {}",
3936                    self.schema.index_label(),
3937                    hex_id
3938                );
3939            }
3940        }
3941
3942        Ok(deleted)
3943    }
3944}
3945
3946#[cfg(test)]
3947mod tests {
3948    use super::*;
3949    use std::sync::atomic::{AtomicBool, Ordering};
3950
3951    fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3952        let schema = crate::dsl::SchemaBuilder::default().build();
3953        let metadata = IndexMetadata::new(schema.clone());
3954        Arc::new(SegmentManager::new(
3955            Arc::new(crate::directories::RamDirectory::new()),
3956            Arc::new(schema),
3957            metadata,
3958            Box::new(crate::merge::NoMergePolicy),
3959            0,
3960            1,
3961            Arc::new(Semaphore::new(1)),
3962            None,
3963            1024,
3964            Arc::new(ReorderConcurrencyGate::new(1)),
3965            None,
3966        ))
3967    }
3968
3969    #[test]
3970    fn force_merge_planner_pairs_large_and_small_segments() {
3971        let groups = plan_force_merge_groups(
3972            vec![
3973                ("a".into(), 6),
3974                ("b".into(), 6),
3975                ("c".into(), 4),
3976                ("d".into(), 4),
3977            ],
3978            10,
3979        );
3980
3981        assert_eq!(groups.len(), 2);
3982        assert!(groups.iter().all(|group| group.total_docs == 10));
3983        assert!(groups.iter().all(|group| group.segments.len() == 2));
3984    }
3985
3986    #[test]
3987    fn force_merge_planner_leaves_oversized_segments_alone() {
3988        let groups = plan_force_merge_groups(
3989            vec![
3990                ("oversized".into(), 11),
3991                ("small-a".into(), 5),
3992                ("small-b".into(), 5),
3993            ],
3994            10,
3995        );
3996
3997        assert_eq!(groups.len(), 2);
3998        assert_eq!(groups[0].total_docs, 10);
3999        assert_eq!(groups[0].segments.len(), 2);
4000        assert_eq!(groups[1].total_docs, 11);
4001        assert_eq!(groups[1].segments.len(), 1);
4002    }
4003
4004    #[test]
4005    fn force_merge_planner_never_exceeds_segment_format_limit() {
4006        let groups = plan_force_merge_groups(
4007            vec![
4008                ("large-a".into(), 3_000_000_000),
4009                ("large-b".into(), 2_000_000_000),
4010            ],
4011            u64::from(u32::MAX),
4012        );
4013        assert_eq!(groups.len(), 2);
4014        assert!(
4015            groups
4016                .iter()
4017                .all(|group| group.total_docs <= u64::from(u32::MAX))
4018        );
4019    }
4020
4021    #[test]
4022    fn force_merge_hierarchy_has_one_final_bp_pass() {
4023        assert_eq!(force_merge_output_count(1), 0);
4024        assert_eq!(force_merge_output_count(2), 1);
4025        assert_eq!(force_merge_output_count(64), 1);
4026        assert_eq!(force_merge_output_count(65), 2);
4027        assert_eq!(force_merge_output_count(127), 2);
4028        assert_eq!(force_merge_output_count(128), 3);
4029        assert_eq!(force_merge_output_count(1_000), 16);
4030    }
4031
4032    fn expand_force_merge_node(
4033        hierarchy: &ForceMergeHierarchy,
4034        source_count: usize,
4035        node: usize,
4036        sources: &mut Vec<usize>,
4037    ) {
4038        if node < source_count {
4039            sources.push(node);
4040            return;
4041        }
4042
4043        let step_index = node - source_count;
4044        let step = hierarchy
4045            .steps
4046            .get(step_index)
4047            .expect("merge input must refer to an existing source or output");
4048        for &input in &step.inputs {
4049            assert!(
4050                input < node,
4051                "merge step {step_index} refers to a future output node {input}"
4052            );
4053            expand_force_merge_node(hierarchy, source_count, input, sources);
4054        }
4055    }
4056
4057    #[test]
4058    fn force_merge_hierarchy_has_minimal_valid_arity() {
4059        let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
4060
4061        for source_count in source_counts {
4062            let hierarchy = plan_force_merge_hierarchy(source_count);
4063            let output_count = hierarchy.steps.len();
4064
4065            assert!(
4066                hierarchy
4067                    .steps
4068                    .iter()
4069                    .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
4070                "invalid merge arity for {source_count} sources"
4071            );
4072            assert!(
4073                source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
4074                "{output_count} outputs cannot reduce {source_count} sources"
4075            );
4076            assert!(
4077                output_count == 1
4078                    || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
4079                "{output_count} outputs are not minimal for {source_count} sources"
4080            );
4081            assert_eq!(output_count, force_merge_output_count(source_count));
4082        }
4083    }
4084
4085    #[test]
4086    fn force_merge_hierarchy_preserves_exact_source_order() {
4087        for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
4088            let hierarchy = plan_force_merge_hierarchy(source_count);
4089            let mut sources = Vec::with_capacity(source_count);
4090            expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
4091            assert_eq!(
4092                sources,
4093                (0..source_count).collect::<Vec<_>>(),
4094                "source order changed for {source_count} sources"
4095            );
4096        }
4097    }
4098
4099    fn force_merge_rewrite_cost(source_count: usize) -> usize {
4100        let hierarchy = plan_force_merge_hierarchy(source_count);
4101        let mut node_weights = vec![1usize; source_count];
4102        let mut rewrite_cost = 0usize;
4103
4104        for (step_index, step) in hierarchy.steps.iter().enumerate() {
4105            let output = source_count + step_index;
4106            let output_weight = step
4107                .inputs
4108                .iter()
4109                .map(|&input| {
4110                    assert!(
4111                        input < output,
4112                        "merge step {step_index} refers to future output {input}"
4113                    );
4114                    node_weights[input]
4115                })
4116                .sum::<usize>();
4117            rewrite_cost += output_weight;
4118            node_weights.push(output_weight);
4119        }
4120
4121        assert_eq!(node_weights[hierarchy.root], source_count);
4122        rewrite_cost
4123    }
4124
4125    #[test]
4126    fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
4127        assert_eq!(force_merge_rewrite_cost(65), 67);
4128        assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
4129    }
4130
4131    #[test]
4132    fn block_copy_carries_bp_debt_without_spending_an_attempt() {
4133        assert_eq!(
4134            replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
4135            (false, false, 3),
4136        );
4137        assert_eq!(
4138            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
4139            (true, false, 4),
4140        );
4141        assert_eq!(
4142            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
4143            (true, true, 0),
4144        );
4145    }
4146
4147    #[tokio::test]
4148    async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
4149        let mut schema_builder = crate::dsl::SchemaBuilder::default();
4150        let field = schema_builder.add_text_field("text", true, true);
4151        let schema = schema_builder.build();
4152        let directory = crate::directories::RamDirectory::new();
4153        let config = crate::index::IndexConfig {
4154            num_indexing_threads: 1,
4155            merge_policy: Box::new(crate::merge::NoMergePolicy),
4156            ..Default::default()
4157        };
4158        let mut writer = crate::index::IndexWriter::create(directory, schema, config)
4159            .await
4160            .unwrap();
4161        for value in ["one", "two", "three"] {
4162            let mut document = crate::dsl::Document::new();
4163            document.add_text(field, value);
4164            writer.add_document(document).unwrap();
4165            writer.commit().await.unwrap();
4166        }
4167
4168        let manager = Arc::clone(writer.segment_manager());
4169        let held_id = manager.get_segment_ids().await.pop().unwrap();
4170        let mut held = Some(
4171            manager
4172                .active_operations
4173                .try_register(vec![held_id])
4174                .unwrap(),
4175        );
4176        let batches = Arc::new(AtomicUsize::new(0));
4177        let batch_count = Arc::clone(&batches);
4178        writer
4179            .force_merge_with_snapshot_refresh(move || {
4180                let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
4181                // Refresh #1 follows the startup drain. Keep the source held
4182                // until refresh #2, after the two free sources were replaced.
4183                if refresh == 2 {
4184                    drop(held.take());
4185                }
4186                std::future::ready(Ok(()))
4187            })
4188            .await
4189            .unwrap();
4190
4191        assert_eq!(manager.get_segment_ids().await.len(), 1);
4192        assert_eq!(
4193            batches.load(Ordering::Relaxed),
4194            4,
4195            "initial/final refreshes plus two replacements are required after the held source releases"
4196        );
4197    }
4198
4199    #[tokio::test]
4200    async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
4201        let mut schema_builder = crate::dsl::SchemaBuilder::default();
4202        schema_builder.set_reorder_on_merge(true);
4203        let schema = schema_builder.build();
4204        let mut metadata = IndexMetadata::new(schema.clone());
4205        metadata.add_segment("00000000000000000000000000000001".into(), 1);
4206        metadata.add_segment("00000000000000000000000000000002".into(), 1);
4207
4208        let global_merge_permits = Arc::new(Semaphore::new(1));
4209        let manager = Arc::new(SegmentManager::new(
4210            Arc::new(crate::directories::RamDirectory::new()),
4211            Arc::new(schema),
4212            metadata,
4213            Box::new(crate::merge::NoMergePolicy),
4214            0,
4215            1,
4216            Arc::clone(&global_merge_permits),
4217            None,
4218            1024,
4219            Arc::new(ReorderConcurrencyGate::new(1)),
4220            None,
4221        ));
4222
4223        // Vector-generation staging rejects ordinary lifecycle claims while
4224        // acquiring the shared global merge permit separately for each source.
4225        // Force merge must wait before capacity admission; retaining the only
4226        // global slot here would deadlock both operations between sources.
4227        manager.active_operations.pause_non_indexing();
4228        let force_merge = {
4229            let manager = Arc::clone(&manager);
4230            tokio::spawn(async move { manager.force_merge().await })
4231        };
4232        tokio::time::timeout(std::time::Duration::from_secs(1), async {
4233            while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4234                tokio::task::yield_now().await;
4235            }
4236        })
4237        .await
4238        .expect("force merge never reached the paused group claim");
4239
4240        assert_eq!(
4241            global_merge_permits.available_permits(),
4242            1,
4243            "force merge retained global capacity while vector staging blocked group ownership"
4244        );
4245
4246        force_merge.abort();
4247        let _ = force_merge.await;
4248        manager.active_operations.resume_non_indexing();
4249    }
4250
4251    #[test]
4252    fn output_cleanup_guard_runs_during_panic_unwind() {
4253        let cleaned = Arc::new(AtomicBool::new(false));
4254        let cleaned_in_callback = Arc::clone(&cleaned);
4255        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4256            cleaned_in_callback.store(true, Ordering::SeqCst);
4257        });
4258
4259        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4260            let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4261            panic!("simulated reorder panic");
4262        }));
4263
4264        assert!(result.is_err());
4265        assert!(
4266            cleaned.load(Ordering::SeqCst),
4267            "partial output cleanup must run during unwind"
4268        );
4269    }
4270
4271    #[test]
4272    fn output_cleanup_guard_disarms_after_commit() {
4273        let cleaned = Arc::new(AtomicBool::new(false));
4274        let cleaned_in_callback = Arc::clone(&cleaned);
4275        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4276            cleaned_in_callback.store(true, Ordering::SeqCst);
4277        });
4278
4279        {
4280            let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4281            guard.disarm();
4282        }
4283
4284        assert!(!cleaned.load(Ordering::SeqCst));
4285    }
4286
4287    #[test]
4288    fn test_active_operation_guard_releases_ownership() {
4289        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4290        {
4291            let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4292            let snap = active.snapshot();
4293            assert!(snap.contains("a"));
4294            assert!(snap.contains("b"));
4295        }
4296        assert!(active.snapshot().is_empty());
4297    }
4298
4299    #[test]
4300    fn test_non_overlapping_operations_can_run_concurrently() {
4301        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4302        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4303        let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4304        let snap = active.snapshot();
4305        assert_eq!(snap.len(), 4);
4306
4307        drop(first);
4308        let snap = active.snapshot();
4309        assert_eq!(snap.len(), 2);
4310        assert!(snap.contains("c"));
4311        assert!(snap.contains("d"));
4312    }
4313
4314    #[test]
4315    fn test_overlapping_operation_is_rejected_until_release() {
4316        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4317        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4318        assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4319        drop(first);
4320        assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4321    }
4322
4323    #[test]
4324    fn test_active_operation_snapshot() {
4325        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4326        let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4327        let snap = active.snapshot();
4328        assert!(snap.contains("x"));
4329        assert!(snap.contains("y"));
4330        assert!(!snap.contains("z"));
4331    }
4332
4333    #[tokio::test]
4334    async fn operation_barrier_ignores_producers_started_after_snapshot() {
4335        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4336        let before_gate = active.try_register(vec!["old".into()]).unwrap();
4337        let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4338        assert_eq!(parked_indexing, 0);
4339        let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4340
4341        let waiter = {
4342            let active = Arc::clone(&active);
4343            tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4344        };
4345        tokio::task::yield_now().await;
4346        assert!(!waiter.is_finished());
4347
4348        drop(before_gate);
4349        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4350            .await
4351            .expect("pre-gate operation barrier was starved by a post-gate producer")
4352            .unwrap();
4353        assert!(active.snapshot().contains("new-flat"));
4354        drop(after_gate);
4355    }
4356
4357    #[tokio::test]
4358    async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4359        let manager = lifecycle_test_manager();
4360        let current = manager.published_generation();
4361        manager
4362            .published_generation
4363            .store(Arc::new(PublishedIndexGeneration {
4364                publication_id: current.publication_id,
4365                schema: current.schema.clone(),
4366                trained_vectors: Some(Arc::new(TrainedVectorStructures {
4367                    centroids: rustc_hash::FxHashMap::default(),
4368                    binary_quantizers: rustc_hash::FxHashMap::default(),
4369                    ..Default::default()
4370                })),
4371            }));
4372
4373        let guard = manager.begin_vector_artifact_update().await.unwrap();
4374        assert!(
4375            manager.trained().is_some(),
4376            "search readers keep the last fully validated generation"
4377        );
4378        assert!(
4379            manager.trained_for_segment_build().is_none(),
4380            "new segment producers must stay flat during an artifact update"
4381        );
4382
4383        let detached_transaction_guard = guard.clone();
4384        drop(guard);
4385        assert!(
4386            manager.trained_for_segment_build().is_none(),
4387            "a detached lifecycle transaction must retain the producer gate after request cancellation"
4388        );
4389        drop(detached_transaction_guard);
4390        assert!(manager.trained_for_segment_build().is_some());
4391    }
4392
4393    #[tokio::test]
4394    async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4395        let manager = lifecycle_test_manager();
4396        let guard = manager.begin_vector_artifact_update().await.unwrap();
4397        assert!(
4398            manager
4399                .active_operations
4400                .try_register(vec!["merge".into()])
4401                .is_none(),
4402            "ordinary merge/reorder work must not change staged sources"
4403        );
4404        let indexing = manager
4405            .active_operations
4406            .try_register_indexing(vec!["fresh".into()])
4407            .expect("indexing remains available in flat mode");
4408        drop(indexing);
4409
4410        drop(guard);
4411        assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4412        assert!(
4413            manager
4414                .active_operations
4415                .try_register(vec!["merge".into()])
4416                .is_some()
4417        );
4418    }
4419
4420    #[tokio::test]
4421    async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4422        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4423        let guard = active.try_register(vec!["live".into()]).unwrap();
4424        let cancellation = active.cancellation_flag();
4425        active.stop_accepting();
4426        assert!(cancellation.load(Ordering::Acquire));
4427        assert!(active.try_register(vec!["new".into()]).is_none());
4428
4429        let waiter = {
4430            let active = Arc::clone(&active);
4431            tokio::spawn(async move { active.wait_until_idle().await })
4432        };
4433        tokio::task::yield_now().await;
4434        assert!(!waiter.is_finished());
4435        drop(guard);
4436        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4437            .await
4438            .expect("shutdown waiter missed the final guard notification")
4439            .unwrap();
4440    }
4441
4442    #[tokio::test]
4443    async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4444        let manager = lifecycle_test_manager();
4445        let started = Arc::new(Semaphore::new(0));
4446        let release = Arc::new(Semaphore::new(0));
4447        let completed = Arc::new(AtomicBool::new(false));
4448
4449        let request = {
4450            let manager = Arc::clone(&manager);
4451            let started = Arc::clone(&started);
4452            let release = Arc::clone(&release);
4453            let completed = Arc::clone(&completed);
4454            tokio::spawn(async move {
4455                manager
4456                    .run_lifecycle_transaction(async move {
4457                        started.add_permits(1);
4458                        let _permit = release.acquire().await.unwrap();
4459                        completed.store(true, Ordering::Release);
4460                        Ok(())
4461                    })
4462                    .await
4463            })
4464        };
4465
4466        let _started = started.acquire().await.unwrap();
4467        request.abort();
4468        assert!(request.await.unwrap_err().is_cancelled());
4469        release.add_permits(1);
4470
4471        manager.begin_shutdown();
4472        tokio::time::timeout(
4473            std::time::Duration::from_secs(1),
4474            manager.wait_for_shutdown(),
4475        )
4476        .await
4477        .expect("shutdown did not drain detached lifecycle transaction");
4478        assert!(completed.load(Ordering::Acquire));
4479    }
4480
4481    #[tokio::test]
4482    async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4483        let manager = lifecycle_test_manager();
4484        {
4485            let mut state = manager.state.lock().await;
4486            state.metadata.add_segment_meta(
4487                "eligible".into(),
4488                SegmentMetaInfo {
4489                    num_docs: 10,
4490                    ancestors: Vec::new(),
4491                    generation: 1,
4492                    reordered: true,
4493                    bp_converged: false,
4494                    bp_unconverged_passes: 2,
4495                },
4496            );
4497            state.metadata.add_segment_meta(
4498                "at-limit".into(),
4499                SegmentMetaInfo {
4500                    num_docs: 20,
4501                    ancestors: Vec::new(),
4502                    generation: 1,
4503                    reordered: true,
4504                    bp_converged: false,
4505                    bp_unconverged_passes: 3,
4506                },
4507            );
4508            state.metadata.add_segment_meta(
4509                "carried-debt".into(),
4510                SegmentMetaInfo {
4511                    num_docs: 15,
4512                    ancestors: Vec::new(),
4513                    generation: 2,
4514                    reordered: false,
4515                    bp_converged: false,
4516                    bp_unconverged_passes: 2,
4517                },
4518            );
4519            state.metadata.add_segment_meta(
4520                "carried-debt-at-limit".into(),
4521                SegmentMetaInfo {
4522                    num_docs: 25,
4523                    ancestors: Vec::new(),
4524                    generation: 2,
4525                    reordered: false,
4526                    bp_converged: false,
4527                    bp_unconverged_passes: 3,
4528                },
4529            );
4530            state.metadata.add_segment_meta(
4531                "converged".into(),
4532                SegmentMetaInfo {
4533                    num_docs: 30,
4534                    ancestors: Vec::new(),
4535                    generation: 1,
4536                    reordered: true,
4537                    bp_converged: true,
4538                    bp_unconverged_passes: 0,
4539                },
4540            );
4541            state.metadata.add_segment("fresh".into(), 40);
4542        }
4543
4544        assert_eq!(
4545            manager.unreordered_segments().await,
4546            vec![("fresh".into(), 40)],
4547            "a block-copy output with BP debt is not a fresh first-pass candidate",
4548        );
4549        let mut eligible = manager.unconverged_segments_below(3).await;
4550        eligible.sort_unstable();
4551        assert_eq!(
4552            eligible,
4553            vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4554        );
4555        assert!(manager.unconverged_segments_below(0).await.is_empty());
4556    }
4557
4558    #[test]
4559    fn merge_retry_backoff_is_exponential_and_capped() {
4560        assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4561        assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4562        assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4563        assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4564    }
4565
4566    #[test]
4567    fn only_deterministic_source_errors_are_quarantined() {
4568        assert!(is_deterministic_source_error(&Error::Corruption(
4569            "bad footer".into()
4570        )));
4571        assert!(is_deterministic_source_error(&Error::Io(
4572            std::io::Error::from(std::io::ErrorKind::NotFound)
4573        )));
4574        assert!(!is_deterministic_source_error(&Error::Io(
4575            std::io::Error::from(std::io::ErrorKind::TimedOut)
4576        )));
4577        assert!(!is_deterministic_source_error(&Error::Io(
4578            std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4579        )));
4580    }
4581
4582    #[test]
4583    fn transient_reorder_failure_is_backed_off_until_cleared() {
4584        let manager = lifecycle_test_manager();
4585        manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4586        assert!(manager.paused_reorder_segments().contains("source"));
4587        manager.clear_reorder_retry("source");
4588        assert!(!manager.paused_reorder_segments().contains("source"));
4589    }
4590
4591    /// Fails `exists` with the transient I/O error class that sends a
4592    /// background merge into its generic retry backoff (not source quarantine).
4593    #[derive(Default)]
4594    struct FailingExistsDirectory(crate::directories::RamDirectory);
4595
4596    #[async_trait::async_trait]
4597    impl crate::directories::Directory for FailingExistsDirectory {
4598        async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4599            Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4600        }
4601
4602        async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4603            self.0.file_size(path).await
4604        }
4605
4606        async fn open_read(
4607            &self,
4608            path: &std::path::Path,
4609        ) -> std::io::Result<crate::directories::FileHandle> {
4610            self.0.open_read(path).await
4611        }
4612
4613        async fn read_range(
4614            &self,
4615            path: &std::path::Path,
4616            range: std::ops::Range<u64>,
4617        ) -> std::io::Result<crate::directories::OwnedBytes> {
4618            self.0.read_range(path, range).await
4619        }
4620
4621        async fn list_files(
4622            &self,
4623            prefix: &std::path::Path,
4624        ) -> std::io::Result<Vec<std::path::PathBuf>> {
4625            self.0.list_files(prefix).await
4626        }
4627
4628        async fn open_lazy(
4629            &self,
4630            path: &std::path::Path,
4631        ) -> std::io::Result<crate::directories::FileHandle> {
4632            self.0.open_lazy(path).await
4633        }
4634    }
4635
4636    #[async_trait::async_trait]
4637    impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4638        async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4639            self.0.write(path, data).await
4640        }
4641
4642        async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4643            self.0.delete(path).await
4644        }
4645
4646        async fn rename(
4647            &self,
4648            from: &std::path::Path,
4649            to: &std::path::Path,
4650        ) -> std::io::Result<()> {
4651            self.0.rename(from, to).await
4652        }
4653
4654        async fn sync(&self) -> std::io::Result<()> {
4655            self.0.sync().await
4656        }
4657
4658        async fn streaming_writer(
4659            &self,
4660            path: &std::path::Path,
4661        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4662            self.0.streaming_writer(path).await
4663        }
4664    }
4665
4666    #[derive(Debug, Clone)]
4667    struct MergeEverythingPolicy;
4668
4669    impl MergePolicy for MergeEverythingPolicy {
4670        fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4671            if segments.len() < 2 {
4672                return Vec::new();
4673            }
4674            vec![crate::merge::MergeCandidate {
4675                segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4676            }]
4677        }
4678
4679        fn clone_box(&self) -> Box<dyn MergePolicy> {
4680            Box::new(self.clone())
4681        }
4682    }
4683
4684    #[tokio::test]
4685    async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4686        let manager = lifecycle_test_manager();
4687        // Simulates a memory-budget mid-cycle segment build whose guard is
4688        // parked inside a PreparedSegment: only a later commit releases this
4689        // token, and that commit can be blocked on the very caller of the
4690        // artifact update (writer write lock / &mut self).
4691        let parked_indexing = manager
4692            .protect_new_segment("00000000000000000000000000000abc".into())
4693            .unwrap();
4694
4695        let error = tokio::time::timeout(
4696            std::time::Duration::from_secs(2),
4697            manager.begin_vector_artifact_update(),
4698        )
4699        .await
4700        .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4701        .err()
4702        .expect("an old-generation prepared segment must block artifact replacement")
4703        .to_string();
4704        assert!(error.contains("built but uncommitted"), "{error}");
4705        assert!(
4706            !manager.vector_artifact_update.load(Ordering::Acquire),
4707            "a rejected update must release the producer gate"
4708        );
4709
4710        drop(parked_indexing);
4711
4712        let guard = manager
4713            .begin_vector_artifact_update()
4714            .await
4715            .expect("artifact update should succeed after the pending generation is resolved");
4716        drop(guard);
4717    }
4718
4719    #[tokio::test]
4720    async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4721        let manager = lifecycle_test_manager();
4722        let merge_like = manager
4723            .active_operations
4724            .try_register(vec!["merge-source".into()])
4725            .unwrap();
4726
4727        let waiter = {
4728            let manager = Arc::clone(&manager);
4729            tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4730        };
4731        for _ in 0..8 {
4732            tokio::task::yield_now().await;
4733        }
4734        assert!(
4735            !waiter.is_finished(),
4736            "artifact update must drain merge/reorder producers that may hold the previous generation"
4737        );
4738
4739        drop(merge_like);
4740        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4741            .await
4742            .expect("artifact update missed the lifecycle guard release")
4743            .unwrap()
4744            .unwrap();
4745    }
4746
4747    #[tokio::test]
4748    async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4749        let manager = lifecycle_test_manager();
4750        let release = Arc::new(Semaphore::new(0));
4751        let merge_task = {
4752            let release = Arc::clone(&release);
4753            tokio::spawn(async move {
4754                let _permit = release.acquire().await.unwrap();
4755            })
4756        };
4757        manager.merge_handles.lock().push(merge_task);
4758
4759        let waiter = {
4760            let manager = Arc::clone(&manager);
4761            tokio::spawn(async move { manager.wait_for_all_merges().await })
4762        };
4763        for _ in 0..8 {
4764            tokio::task::yield_now().await;
4765        }
4766        assert!(!waiter.is_finished());
4767        // Simulates tonic dropping a force_merge/reorder RPC future at the
4768        // JoinHandle await when the client disconnects.
4769        waiter.abort();
4770        let join_error = waiter.await.unwrap_err();
4771        assert!(join_error.is_cancelled());
4772
4773        assert!(
4774            !manager.merge_handles.lock().is_empty(),
4775            "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4776        );
4777
4778        // A later drain must still see and await the real in-flight merge.
4779        release.add_permits(1);
4780        tokio::time::timeout(
4781            std::time::Duration::from_secs(1),
4782            manager.wait_for_all_merges(),
4783        )
4784        .await
4785        .expect("subsequent drain missed the reinserted merge handle");
4786        assert!(manager.merge_handles.lock().is_empty());
4787    }
4788
4789    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4790    async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4791        let manager = lifecycle_test_manager();
4792        {
4793            let mut state = manager.state.lock().await;
4794            state
4795                .metadata
4796                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4797            state
4798                .metadata
4799                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4800        }
4801        // A background reorder (or a concurrent force-merge) owns one segment
4802        // in the batch but never appears in merge_handles.
4803        let reorder_like = manager
4804            .active_operations
4805            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4806            .unwrap();
4807
4808        let force_merge = {
4809            let manager = Arc::clone(&manager);
4810            tokio::spawn(async move { manager.force_merge().await })
4811        };
4812
4813        tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4814        let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4815        assert!(
4816            retries >= 1,
4817            "force_merge never observed the conflicting owner (retries={retries})"
4818        );
4819        assert!(
4820            retries < 20,
4821            "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4822        );
4823
4824        drop(reorder_like);
4825        // With the conflict gone the loop proceeds; the batch then fails fast
4826        // in do_merge (the test IDs have no files), proving the loop exited.
4827        let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4828            .await
4829            .expect("force_merge kept spinning after the conflicting owner released")
4830            .unwrap();
4831        assert!(result.is_err());
4832    }
4833
4834    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4835    async fn force_merge_routes_around_segments_held_by_reorder() {
4836        let manager = lifecycle_test_manager();
4837        {
4838            let mut state = manager.state.lock().await;
4839            state
4840                .metadata
4841                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4842            state
4843                .metadata
4844                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4845            state
4846                .metadata
4847                .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4848        }
4849        // A background reorder owns one segment and holds it for the whole
4850        // test (in prod: a BP pass runs for minutes while force_merge spins).
4851        let _reorder_like = manager
4852            .active_operations
4853            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4854            .unwrap();
4855
4856        // Regression: force_merge used to rebuild the identical smallest-N
4857        // batch (including the held segment) every 100ms and retry-log
4858        // forever. It must instead skip the held segment and immediately
4859        // make progress on the two free ones — reaching do_merge (which
4860        // fails fast here: the test IDs have no files) proves the batch was
4861        // built without the held segment while the reorder is STILL active.
4862        let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4863            let manager = Arc::clone(&manager);
4864            async move { manager.force_merge().await }
4865        })
4866        .await
4867        .expect("force_merge livelocked on a segment held by an active reorder");
4868        assert!(result.is_err(), "fake segment files must fail the merge");
4869
4870        assert_eq!(
4871            manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4872            0,
4873            "batch built from the ownership snapshot must not collide with the held segment"
4874        );
4875    }
4876
4877    #[tokio::test]
4878    async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4879        let schema = crate::dsl::SchemaBuilder::default().build();
4880        let mut metadata = IndexMetadata::new(schema.clone());
4881        metadata.add_segment("00000000000000000000000000000001".into(), 10);
4882        metadata.add_segment("00000000000000000000000000000002".into(), 10);
4883        let manager = Arc::new(SegmentManager::new(
4884            Arc::new(FailingExistsDirectory::default()),
4885            Arc::new(schema),
4886            metadata,
4887            Box::new(MergeEverythingPolicy),
4888            0,
4889            1,
4890            Arc::new(Semaphore::new(1)),
4891            None,
4892            1024,
4893            Arc::new(ReorderConcurrencyGate::new(1)),
4894            None,
4895        ));
4896
4897        // Spawns a background merge that fails with a transient I/O error and
4898        // arms the 30s..30min retry backoff.
4899        manager.maybe_merge().await;
4900
4901        tokio::time::timeout(
4902            std::time::Duration::from_secs(5),
4903            manager.wait_for_all_merges(),
4904        )
4905        .await
4906        .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4907        assert!(
4908            manager.merge_retry_is_paused(),
4909            "the failed merge should have armed the retry backoff"
4910        );
4911
4912        // Shutdown still drains the pending backoff wakeup deterministically.
4913        manager.begin_shutdown();
4914        tokio::time::timeout(
4915            std::time::Duration::from_secs(5),
4916            manager.wait_for_shutdown(),
4917        )
4918        .await
4919        .expect("shutdown did not drain the merge retry wakeup task");
4920    }
4921}