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