diff --git a/crates/lance-context-core/src/registry.rs b/crates/lance-context-core/src/registry.rs index d6e9bb3..87c787f 100644 --- a/crates/lance-context-core/src/registry.rs +++ b/crates/lance-context-core/src/registry.rs @@ -19,11 +19,14 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; +use std::time::Duration; use arrow_array::{Int64Array, RecordBatch, RecordBatchIterator, StringArray}; use arrow_schema::{ArrowError, DataType, Field, Schema}; use chrono::Utc; use futures::TryStreamExt; +use lance::dataset::cleanup::RemovalStats; +use lance::dataset::optimize::{compact_files, CompactionMetrics, CompactionOptions}; use lance::dataset::{builder::DatasetBuilder, Dataset, WriteMode, WriteParams}; use lance::io::{ObjectStoreParams, StorageOptionsAccessor}; use lance::{Error as LanceError, Result as LanceResult}; @@ -347,6 +350,73 @@ impl RolloutRegistry { pub fn uri(&self) -> &str { &self.uri } + + /// Current dataset version (manifest chain head). + pub fn version(&self) -> u64 { + self.dataset.version().version + } + + /// Fold the one-row append fragments produced by [`Self::upsert`] into a + /// few, then drop manifest versions older than `older_than`. + /// + /// Every `create` is a delete plus an append, so the registry gains two + /// versions and one fragment per store. Lance keeps every manifest until + /// cleaned, and every manifest lists every fragment, so an unmaintained + /// registry grows quadratically: 34k versions of ~840 KB each (29 GB of + /// manifests) for a 10k-row, three-column table was observed in + /// production, and each `contains`/`get` re-reads the head manifest. + /// + /// Callers that share this registry behind a lock should prefer + /// [`Self::compact`] followed by [`RegistryCleaner::cleanup`] so the lock + /// is not held across the cleanup's 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`]. Returns a + /// [`RegistryCleaner`] that prunes old versions without borrowing the + /// registry. + pub async fn compact(&mut self) -> LanceResult<(CompactionMetrics, RegistryCleaner)> { + self.reload().await?; + let options = CompactionOptions { + target_rows_per_fragment: 1_048_576, + materialize_deletions: true, + materialize_deletions_threshold: 0.0, + ..Default::default() + }; + let compaction = compact_files(&mut self.dataset, options, None).await?; + self.reload().await?; + Ok(( + compaction, + RegistryCleaner { + dataset: self.dataset.clone(), + }, + )) + } +} + +/// A detached handle for the old-version cleanup half of registry +/// maintenance. See `StatsCleaner` in the master for why this is split off: +/// cleanup deletes objects no live version references, so it needs no +/// exclusive access, only time. +pub struct RegistryCleaner { + dataset: Dataset, +} + +impl RegistryCleaner { + /// 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 registry history TTL: {e}")))?; + self.dataset.cleanup_old_versions(grace, None, None).await + } } #[cfg(test)] @@ -361,6 +431,90 @@ mod tests { .unwrap() } + /// Maintenance folds the per-upsert fragments and prunes old manifests + /// while keeping every live row; the registry keeps working afterwards. + #[tokio::test] + async fn maintain_bounds_versions_and_preserves_rows() { + let dir = TempDir::new().unwrap(); + let mut r = new_registry(&dir).await; + for i in 0..20 { + r.upsert(&format!("exp-{i}"), &format!("/data/exp-{i}.lance")) + .await + .unwrap(); + } + let before = r.version(); + assert!( + before >= 40, + "expected two versions per upsert, got {before}" + ); + + let (compaction, removal) = r.maintain(Duration::from_secs(0)).await.unwrap(); + assert!(compaction.fragments_removed > 0, "nothing compacted"); + assert!(removal.old_versions > 0, "no versions reclaimed"); + + let mut names: Vec = r + .list() + .await + .unwrap() + .into_iter() + .map(|e| e.name) + .collect(); + names.sort(); + assert_eq!(names.len(), 20); + assert_eq!(names[0], "exp-0"); + assert!(r.contains("exp-19").await.unwrap()); + r.upsert("exp-new", "/data/exp-new.lance").await.unwrap(); + assert!(r.contains("exp-new").await.unwrap()); + } + + /// A second handle on the same URI (a worker) keeps creating stores while + /// another (the master) maintains. Nothing the worker wrote is lost and + /// both handles see every row afterwards. + #[tokio::test] + async fn maintenance_does_not_lose_concurrent_worker_upserts() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().join("_registry.lance"); + let uri = uri.to_str().unwrap(); + let mut master = RolloutRegistry::open_or_create(uri, None).await.unwrap(); + let mut worker = RolloutRegistry::open_or_create(uri, None).await.unwrap(); + for i in 0..20 { + worker.upsert(&format!("old-{i}"), "/x").await.unwrap(); + } + + let (_, cleaner) = master.compact().await.unwrap(); + // Worker writes land between the compaction and the cleanup ... + for i in 0..5 { + worker.upsert(&format!("mid-{i}"), "/x").await.unwrap(); + } + cleaner.cleanup(Duration::from_secs(0)).await.unwrap(); + master.reload().await.unwrap(); + // ... and after it. + worker.upsert("late", "/x").await.unwrap(); + + assert_eq!(master.list().await.unwrap().len(), 26); + assert_eq!(worker.list().await.unwrap().len(), 26); + assert!(master.contains("mid-3").await.unwrap()); + assert!(worker.contains("late").await.unwrap()); + } + + /// The cleanup half runs on a detached handle: the registry stays usable + /// (reads and writes) while the cleaner is outstanding. + #[tokio::test] + async fn cleanup_runs_detached_from_the_registry() { + let dir = TempDir::new().unwrap(); + let mut r = new_registry(&dir).await; + for i in 0..20 { + r.upsert(&format!("exp-{i}"), "/x").await.unwrap(); + } + let (_, cleaner) = r.compact().await.unwrap(); + r.upsert("during", "/x").await.unwrap(); + assert!(r.contains("during").await.unwrap()); + let removal = cleaner.cleanup(Duration::from_secs(0)).await.unwrap(); + assert!(removal.old_versions > 0); + r.reload().await.unwrap(); + assert_eq!(r.list().await.unwrap().len(), 21); + } + #[tokio::test] async fn create_or_load_recovers_when_another_caller_wins() { let dir = TempDir::new().unwrap(); diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index 7bd8253..e6ed92c 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -14,7 +14,8 @@ use std::time::Duration; use chrono::Utc; use futures::stream::{self, StreamExt}; use lance_context_core::{ - CompactionConfig, GenericStore, GenericStoreOptions, RolloutStore, RolloutStoreOptions, + CompactionConfig, GenericStore, GenericStoreOptions, RolloutRegistry, RolloutStore, + RolloutStoreOptions, }; use tokio::sync::Semaphore; use tokio::task::JoinHandle; @@ -150,6 +151,55 @@ pub async fn maintain_stats(state: &Arc) -> lance::Result<()> { } } +/// Compact a store registry and prune its old manifest versions. +/// +/// Workers write the registry on every store create/delete as a delete plus an +/// append, so it gains two versions and one fragment per store and, without +/// this, every manifest lists every fragment: 34k versions of ~840 KB each +/// were observed for a 10k-row table, and every `create`/`get` re-read one. +/// +/// Same shape as [`maintain_stats`]: the compaction (a `Rewrite` that must not +/// race another master's) runs under `state.stats`-style exclusion via the +/// registry lock, bounded by [`MAINTENANCE_TIMEOUT`]; the cleanup runs on a +/// detached handle with the lock released. Workers keep appending in between: +/// Lance retries their commits past the rewrite, and the grace window keeps +/// any version a worker may still be reading. +async fn maintain_registry( + state: &Arc, + label: &'static str, + registry: &tokio::sync::RwLock, +) -> lance::Result<()> { + let ttl = Duration::from_secs(state.config.stats_history_ttl_secs); + let start = std::time::Instant::now(); + let (compaction, cleaner) = { + let mut registry = registry.write().await; + tokio::time::timeout(MAINTENANCE_TIMEOUT, registry.compact()) + .await + .unwrap_or_else(|_| Err(lance::Error::io("registry compaction timed out")))? + }; + let removal = cleaner.cleanup(ttl).await; + let reload = registry.write().await.reload().await; + let removal = match (removal, reload) { + (Ok(removal), Ok(())) => removal, + (Err(e), _) | (Ok(_), Err(e)) => return Err(e), + }; + let version = registry.read().await.version(); + metrics::counter!("master_registry_versions_removed_total", "registry" => label) + .increment(removal.old_versions); + metrics::gauge!("master_registry_version", "registry" => label).set(version as f64); + tracing::info!( + registry = label, + fragments_removed = compaction.fragments_removed, + fragments_added = compaction.fragments_added, + old_versions_removed = removal.old_versions, + bytes_removed = removal.bytes_removed, + version, + elapsed_secs = start.elapsed().as_secs(), + "registry maintenance complete" + ); + Ok(()) +} + /// Run a single scan pass: refresh every experiment's stats row and drop rows /// for experiments no longer in the registry. Returns the number of /// experiments successfully observed. @@ -179,6 +229,14 @@ async fn try_scan_once(state: &Arc, maintain: bool) -> lance::Resul if let Err(e) = maintain_stats(state).await { tracing::warn!(error = %e, "stats maintenance failed"); } + for (label, registry) in [ + ("rollout", &state.registry), + ("generic", &state.generic_registry), + ] { + if let Err(e) = maintain_registry(state, label, registry).await { + tracing::warn!(registry = label, error = %e, "registry maintenance failed"); + } + } } let release = state.task_store.release_coordination_lock(guard).await; match (result, release) {