diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index 14cb482..7bd8253 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -28,11 +28,10 @@ use lance_context_api::ExperimentSummary; /// scan round. const OBSERVE_TIMEOUT: Duration = Duration::from_secs(30); -/// Bound on one steady-state `_stats` maintenance pass so a slow object store -/// cannot wedge the scanner loop. +/// Bound on the compaction half of one `_stats` maintenance pass, which runs +/// under the store lock, so a slow object store cannot wedge the scanner loop. /// -/// Deliberately not applied to the first pass after startup: see -/// [`maintain_stats`]. +/// The cleanup half is not bounded: see [`maintain_stats`]. const MAINTENANCE_TIMEOUT: Duration = Duration::from_secs(300); /// Bound on the merge and compaction steps of retiring one experiment. Larger @@ -47,50 +46,61 @@ const RETIRE_TIMEOUT: Duration = Duration::from_secs(600); /// one per round and this pass is cheap. It remains necessary to reclaim those /// versions, and to recover deployments that ran the old per-row path. /// -/// # The first pass runs without a timeout +/// # The store lock is held only for the compaction /// -/// A deployment upgraded from the per-row path can arrive with a chain -/// hundreds of thousands of versions long (246k+ observed). Compacting and -/// pruning that cannot finish inside `MAINTENANCE_TIMEOUT`, so every pass timed -/// out, rolled back, and left the table exactly as bloated as before — the -/// bound guaranteed the table could never recover. The first pass after startup -/// therefore runs unbounded, and subsequent passes take the bound: by then the -/// backlog is gone and any pass exceeding it is a genuine fault. +/// Compaction rewrites the table and must exclude other writers, so it runs +/// under `state.stats`. The cleanup that follows deletes objects no live +/// version references; it needs no exclusive access, only time -- a backlog of +/// tens of thousands of versions is tens of thousands of deletes against a +/// remote object store. It used to run under the same lock, and every sweep +/// and every post-compaction stats refresh on every master waited behind it +/// (36 minutes observed after a routine restart). It now runs on a detached +/// handle with the lock released; the store is re-opened afterwards so its +/// in-memory handle never points at a pruned manifest. +/// +/// The cleanup is deliberately unbounded. Lance deletes as it goes, so a +/// timeout would only abandon the pass mid-way, not undo it, and with the lock +/// released nothing waits on it any more. Only the compaction is bounded by +/// [`MAINTENANCE_TIMEOUT`]. /// /// Callers must hold the `stats-writer` coordination lock so only one replica /// ever rewrites the dataset. pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { let ttl = Duration::from_secs(state.config.stats_history_ttl_secs); let start = std::time::Instant::now(); - let mut stats = state.stats.lock().await; - // `swap` so exactly one pass per process is unbounded, even if several - // scanner ticks race here. - let first_pass = state.stats_maintenance_done.swap(true, Ordering::SeqCst); - let outcome = if first_pass { - tokio::time::timeout(MAINTENANCE_TIMEOUT, stats.maintain(ttl)) + let compacted = { + let mut stats = state.stats.lock().await; + tokio::time::timeout(MAINTENANCE_TIMEOUT, stats.compact()) .await - .unwrap_or_else(|_| Err(lance::Error::io("stats maintenance timed out"))) - } else { - tracing::info!( - version = stats.version(), - "running first stats maintenance pass without a timeout; \ - a table carried over from the per-row write path can take a while to reclaim" - ); - stats.maintain(ttl).await + .unwrap_or_else(|_| Err(lance::Error::io("stats compaction timed out"))) }; + let outcome = match compacted { + Ok((compaction, cleaner)) => { + let removal = cleaner.cleanup(ttl).await; + // Re-open regardless of the cleanup's outcome: a failed cleanup + // may still have deleted manifests the current handle references. + let reload = state.stats.lock().await.reload().await; + match (removal, reload) { + (Ok(removal), Ok(())) => Ok((compaction, removal)), + (Err(e), _) | (Ok(_), Err(e)) => Err(e), + } + } + Err(e) => Err(e), + }; + let version = state.stats.lock().await.version(); match outcome { Ok((compaction, removal)) => { state.stats_maintenance_failures.store(0, Ordering::Relaxed); state .stats_last_reclaimed_version - .store(stats.version(), Ordering::Relaxed); + .store(version, Ordering::Relaxed); metrics::histogram!("master_stats_maintenance_duration_seconds") .record(start.elapsed().as_secs_f64()); metrics::counter!("master_stats_versions_removed_total") .increment(removal.old_versions); - metrics::gauge!("master_stats_version").set(stats.version() as f64); + metrics::gauge!("master_stats_version").set(version as f64); metrics::gauge!("master_stats_maintenance_consecutive_failures").set(0.0); metrics::gauge!("master_stats_unreclaimed_versions").set(0.0); tracing::info!( @@ -98,7 +108,8 @@ pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { fragments_added = compaction.fragments_added, old_versions_removed = removal.old_versions, bytes_removed = removal.bytes_removed, - version = stats.version(), + version, + elapsed_secs = start.elapsed().as_secs(), "stats maintenance complete" ); Ok(()) @@ -121,7 +132,7 @@ pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { .fetch_add(1, Ordering::Relaxed) + 1; let last_reclaimed = state.stats_last_reclaimed_version.load(Ordering::Relaxed); - let unreclaimed = stats.version().saturating_sub(last_reclaimed); + let unreclaimed = version.saturating_sub(last_reclaimed); metrics::counter!("master_stats_maintenance_failures_total").increment(1); metrics::gauge!("master_stats_maintenance_consecutive_failures").set(failures as f64); @@ -131,7 +142,7 @@ pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { error = %e, consecutive_failures = failures, unreclaimed_versions = unreclaimed, - version = stats.version(), + version, "stats maintenance failed; old manifests are not being reclaimed" ); Err(e) diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 19904b5..1d5b197 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -44,6 +44,10 @@ use crate::task_store::{TaskClaim, TaskKinds}; /// anything still over the threshold is picked up by the next tick. const MAX_SWEEP_ENQUEUE: usize = 256; +/// How long a finished compaction waits for the `stats-writer` lock to refresh +/// its stats row before giving up and leaving it to the next scan round. +const STATS_REFRESH_LOCK_WAIT: Duration = Duration::from_secs(10); + /// Task-target prefix marking a generic store. Store names match /// `[A-Za-z0-9_][A-Za-z0-9._-]*`, so a `:` can never appear in a bare name /// and the prefix is unambiguous. Carrying the kind in the target string -- @@ -383,13 +387,30 @@ async fn merge_wal_one(http: &reqwest::Client, url: &str) -> Result, name: &str, store: &RolloutStore) { - let guard = match state.task_store.coordination_lock("stats-writer").await { - Ok(guard) => guard, - Err(e) => { + let guard = match tokio::time::timeout( + STATS_REFRESH_LOCK_WAIT, + state.task_store.coordination_lock("stats-writer"), + ) + .await + { + Ok(Ok(guard)) => guard, + Ok(Err(e)) => { tracing::warn!(store = %name, error = %e, "stats writer lock failed"); return; } + Err(_) => { + tracing::info!( + store = %name, + "stats writer busy; leaving the post-compaction refresh to the next scan" + ); + return; + } }; let obs = match store.observe().await { Ok(obs) => obs, diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index 6fe54d5..dba767d 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -104,13 +104,6 @@ pub struct MasterState { pub http: reqwest::Client, /// Process-wide compaction permits shared by scheduler and retirement work. pub(crate) compaction_permits: Arc, - /// Whether this process has already run one `_stats` maintenance pass. - /// - /// The first pass runs without a timeout so a deployment carrying a version - /// chain from the old per-row write path can actually reclaim it; a bounded - /// first pass times out forever and never recovers. See - /// [`crate::scanner::maintain_stats`]. - pub stats_maintenance_done: std::sync::atomic::AtomicBool, /// Consecutive `_stats` maintenance failures, for alerting. /// /// Failure was previously silent: the success counter simply stopped @@ -167,7 +160,6 @@ impl MasterState { task_store, http: reqwest::Client::new(), compaction_permits: Arc::new(Semaphore::new(compaction_concurrency)), - stats_maintenance_done: std::sync::atomic::AtomicBool::new(false), stats_maintenance_failures: std::sync::atomic::AtomicU64::new(0), stats_last_reclaimed_version: std::sync::atomic::AtomicU64::new(0), }); diff --git a/crates/lance-context-master/src/stats_store.rs b/crates/lance-context-master/src/stats_store.rs index a493408..4bd7ffb 100644 --- a/crates/lance-context-master/src/stats_store.rs +++ b/crates/lance-context-master/src/stats_store.rs @@ -544,10 +544,23 @@ impl StatsStore { /// /// Callers must serialize this with other mutations (it takes `&mut self`) /// and, across master replicas, hold the `stats-writer` coordination lock. + /// Callers that share this store behind a lock should prefer + /// [`Self::compact`] followed by [`Self::cleanup`] on the returned handle, + /// so the lock is not held for the cleanup's many object-store deletes. pub async fn maintain( &mut self, older_than: Duration, ) -> LanceResult<(CompactionMetrics, RemovalStats)> { + let (compaction, cleaner) = self.compact().await?; + let removal = cleaner.cleanup(older_than).await?; + self.reload().await?; + Ok((compaction, removal)) + } + + /// The compaction half of [`Self::maintain`]: rewrite the table's fragments + /// into one. Returns a [`StatsCleaner`] that prunes old versions *without* + /// borrowing the store, so a caller can release its lock first. + pub async fn compact(&mut self) -> LanceResult<(CompactionMetrics, StatsCleaner)> { self.dataset.checkout_latest().await?; let options = CompactionOptions { @@ -560,16 +573,45 @@ impl StatsStore { }; let compaction = compact_files(&mut self.dataset, options, None).await?; - // Re-open so the handle (and the cleanup below) sees the rewritten - // version rather than the pre-compaction manifest. + // Re-open so the handle (and the cleanup) sees the rewritten version + // rather than the pre-compaction manifest. + self.reload().await?; + Ok(( + compaction, + StatsCleaner { + dataset: self.dataset.clone(), + }, + )) + } + + /// Re-open the dataset at its latest version. Call after a + /// [`StatsCleaner::cleanup`] so the in-memory handle does not reference + /// manifests the cleanup removed. + pub async fn reload(&mut self) -> LanceResult<()> { self.dataset = Self::load(&self.uri, self.storage_options.clone()).await?; + Ok(()) + } +} +/// A detached handle for the old-version cleanup half of maintenance. +/// +/// Cleanup only reads manifests and deletes objects no live version references, +/// so it needs no exclusive access to the store: readers on the current +/// version, and even writers committing new versions, are unaffected. What it +/// does need is time -- tens of thousands of deletes against a remote object +/// store -- and holding the store's lock for that stalled every sweep and +/// every post-compaction stats refresh on every master for the duration. +pub struct StatsCleaner { + dataset: Dataset, +} + +impl StatsCleaner { + /// Drop manifest versions older than `older_than` and the files only they + /// referenced. + pub async fn cleanup(self, older_than: Duration) -> LanceResult { let grace = chrono::TimeDelta::from_std(older_than) .map_err(|e| LanceError::io(format!("invalid stats history TTL: {e}")))?; - let removal = self.dataset.cleanup_old_versions(grace, None, None).await?; - self.dataset = Self::load(&self.uri, self.storage_options.clone()).await?; - - Ok((compaction, removal)) + self.dataset.cleanup_old_versions(grace, None, None).await } } @@ -695,6 +737,37 @@ mod tests { assert_eq!(s.list(None, 10, 0).await.unwrap().len(), 1); } + /// The cleanup half of maintenance runs on a detached handle so the master + /// can release its store lock first. The store must stay fully usable -- + /// reads and writes -- while that cleanup runs, and be intact after it. + #[tokio::test] + async fn cleanup_runs_detached_from_the_store() { + let dir = TempDir::new().unwrap(); + let mut s = new_store(&dir).await; + for i in 0..20 { + s.upsert(&sample("exp-a", i)).await.unwrap(); + s.upsert(&sample("exp-b", i)).await.unwrap(); + } + let (compaction, cleaner) = s.compact().await.unwrap(); + assert!(compaction.fragments_removed > 0, "nothing compacted"); + + // The store is free while the cleaner exists: write through it. + s.upsert(&sample("exp-c", 1)).await.unwrap(); + assert_eq!(s.count(None).await.unwrap(), 3); + + let removal = cleaner.cleanup(Duration::from_secs(0)).await.unwrap(); + assert!(removal.old_versions > 0, "no versions reclaimed"); + + // The write that landed during cleanup is the live version; reloading + // onto it keeps every row. + s.reload().await.unwrap(); + let rows = s.list(None, 10, 0).await.unwrap(); + assert_eq!(rows.len(), 3); + assert_eq!(rows[0].row_count, 19); + s.upsert(&sample("exp-d", 1)).await.unwrap(); + assert_eq!(s.count(None).await.unwrap(), 4); + } + #[tokio::test] async fn no_compaction_sentinel_maps_to_none() { let row = sample("x", 1);