Skip to main content

hermes_core/merge/
segment_manager.rs

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