Skip to main content

hermes_core/merge/
segment_manager.rs

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