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