From 1ad7c3f03d2aa412e24dd5c5bc2268c6b152119c Mon Sep 17 00:00:00 2001 From: Beinan Date: Mon, 27 Jul 2026 09:03:34 +0000 Subject: [PATCH] feat(master): retire experiments idle for a week from the stats table MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Stacked on #209. The stats table currently holds a row for every experiment that has ever existed, and each one costs an open on every scan round. At the scale this deployment is heading for — tens of thousands, the overwhelming majority written once and never again — that is the remaining reason a round cannot stay proportional to actual work. An experiment with no writes for `STATS_COLD_RETIRE_SECS` (default 7 days) is now dropped from the table. It stays in the registry and can still be observed on demand; it simply stops costing a row and an open every round. Retirement order is a correctness property, not an optimisation. Both auto-sweeps read the stats table and nothing else, so an experiment absent from it is invisible to them permanently. Dropping one that still has un-merged MemWAL generations would strand them: the rows stay readable because the read path unions the generations, but they are never folded into the base table, read amplification never recovers, the `_mem_wal/{shard}/` directories are never reclaimed, and no process is left that would notice. So retirement is: merge the WAL -> compact -> verify pending_wal_generations == 0 -> drop the row `compact_files` deliberately leaves MemWAL generations alone, so compaction on its own would not have been sufficient. Compaction is still done, because a retired table is one nothing will come back to tidy. Any step failing leaves the experiment in the table for the next round. Refusing to retire is always safe; retiring early is not. `STATS_COLD_RETIRE_SECS=0` disables retirement entirely. Every existing test config sets it to 0, so they exercise the unchanged path. Four tests, the important one asserting the shard is genuinely drained after retirement rather than that the function returned success. Verified it is not vacuous: removing the WAL-merge step makes it fail with exactly the stranding it exists to prevent. New: `master_stats_hot_experiments` (table size, now bounded rather than tracking the registry), `master_stats_experiments_retired_total`, `master_stats_retire_failures_total`. Co-Authored-By: Claude --- crates/lance-context-master/src/config.rs | 16 + crates/lance-context-master/src/routes.rs | 1 + crates/lance-context-master/src/scanner.rs | 288 ++++++++++++++++++ crates/lance-context-master/src/scheduler.rs | 1 + crates/lance-context-master/src/state.rs | 1 + crates/lance-context-master/src/task_store.rs | 1 + crates/lance-context-metrics/src/lib.rs | 23 ++ 7 files changed, 331 insertions(+) diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index e127068..50bf5bf 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -41,6 +41,22 @@ pub struct MasterConfig { #[arg(long, env = "STATS_HISTORY_TTL_SECS", default_value_t = 3_600)] pub stats_history_ttl_secs: u64, + /// Age, in seconds, after which an experiment with no writes is retired + /// from the stats table. `0` disables retirement (every known experiment + /// stays in the table forever). + /// + /// Retirement is what keeps the stats table proportional to *active* work + /// rather than to everything ever created. A retired experiment is still + /// listed in the registry and is observed on demand; it simply stops + /// costing a row and an open on every scan round. + /// + /// Before an experiment is retired its MemWAL generations are merged and + /// its fragments compacted, so it is left in a state that needs no further + /// maintenance — which is what makes it safe for the sweeps to stop looking + /// at it. See `scanner::retire_cold_experiments`. + #[arg(long, env = "STATS_COLD_RETIRE_SECS", default_value_t = 604_800)] + pub stats_cold_retire_secs: u64, + /// Interval, in seconds, between automatic compaction sweeps. `0` disables /// automatic compaction (manual triggers still work). #[arg(long, env = "COMPACTION_INTERVAL_SECS", default_value_t = 600)] diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index ad4a3b7..32b20db 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -540,6 +540,7 @@ mod tests { scan_concurrency: 4, stats_maintenance_every_n_scans: 0, stats_history_ttl_secs: 3_600, + stats_cold_retire_secs: 0, compaction_interval_secs: 0, min_fragments: 16, target_rows_per_fragment: 1_048_576, diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index 6bd0543..e466b0e 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -30,6 +30,11 @@ const OBSERVE_TIMEOUT: Duration = Duration::from_secs(30); /// [`maintain_stats`]. const MAINTENANCE_TIMEOUT: Duration = Duration::from_secs(300); +/// Bound on the merge and compaction steps of retiring one experiment. Larger +/// than `OBSERVE_TIMEOUT` because both genuinely rewrite data; a retirement that +/// exceeds it is abandoned and retried next round, which is harmless. +const RETIRE_TIMEOUT: Duration = Duration::from_secs(600); + /// Compact `_stats` and prune its old manifest versions. /// /// Scans now write the table as one `Overwrite` snapshot per round rather than @@ -241,6 +246,21 @@ async fn scan_once_inner(state: &Arc) -> lance::Result { } snapshot.sort_by(|a, b| a.name.cmp(&b.name)); + // Retire experiments with no writes for the configured window. Each is + // merged, compacted and verified quiescent first; only then is its row + // dropped, because a row absent from this table is invisible to both + // auto-sweeps forever. See `retire_cold_experiments`. + // + // Done before the snapshot is written so a retirement takes effect in the + // same commit rather than leaving a round where the row is stale. + let retire_after = Duration::from_secs(state.config.stats_cold_retire_secs); + let retired = + retire_cold_experiments(&snapshot, retire_after, Utc::now().timestamp_millis()).await; + if !retired.is_empty() { + snapshot.retain(|row| !retired.contains(&row.name)); + } + metrics::gauge!("master_stats_hot_experiments").set(snapshot.len() as f64); + let mut total_rows: i64 = 0; let mut total_fragments: i64 = 0; let mut live_count: usize = 0; @@ -367,6 +387,113 @@ async fn observe_one( )) } +/// Decide which of `rows` are cold enough to retire, and leave each one in a +/// state that needs no further maintenance before dropping it. +/// +/// Returns the names that were successfully retired. +/// +/// # Ordering is a correctness property, not an optimisation +/// +/// Both auto-sweeps read the stats table and nothing else, so an experiment +/// absent from it is invisible to them forever. Retiring one that still has +/// un-merged MemWAL generations would strand those generations permanently: +/// their rows stay readable (the read path unions them) but they are never +/// folded into the base table, read amplification never goes down, and the +/// `_mem_wal/{shard}/` directories are never reclaimed. That is a storage leak +/// with no process left to notice it. +/// +/// So retirement is: merge the WAL, compact, **verify the WAL actually drained**, +/// and only then drop the row. `compact_files` deliberately does not touch WAL +/// generations, so compaction alone would not have been enough. +/// +/// Any step failing leaves the experiment in the table to be retried next +/// round. Refusing to retire is always safe; retiring early is not. +async fn retire_cold_experiments( + rows: &[StatRow], + retire_after: Duration, + now_ms: i64, +) -> HashSet { + if retire_after.is_zero() { + return HashSet::new(); + } + let cutoff_ms = now_ms.saturating_sub(retire_after.as_millis() as i64); + + let mut retired = HashSet::new(); + for row in rows { + if row.last_updated > cutoff_ms { + continue; + } + match prepare_for_retirement(&row.name, &row.uri).await { + Ok(true) => { + retired.insert(row.name.clone()); + metrics::counter!("master_stats_experiments_retired_total").increment(1); + tracing::info!( + store = %row.name, + idle_ms = now_ms.saturating_sub(row.last_updated), + "retiring cold experiment from the stats table" + ); + } + Ok(false) => { + // Still had pending generations after the merge, so something + // else is writing or the merge did not fully drain. Keep it. + tracing::debug!( + store = %row.name, + "cold experiment not retired: WAL still pending after merge" + ); + } + Err(e) => { + metrics::counter!("master_stats_retire_failures_total").increment(1); + tracing::warn!( + store = %row.name, + error = %e, + "failed to prepare cold experiment for retirement; keeping it" + ); + } + } + } + retired +} + +/// Merge, compact, and verify one experiment is quiescent. +/// +/// `Ok(true)` means it is safe to drop from the stats table. +async fn prepare_for_retirement(name: &str, uri: &str) -> lance::Result { + let opts = RolloutStoreOptions::default(); + let mut store = match tokio::time::timeout( + OBSERVE_TIMEOUT, + RolloutStore::open_existing_with_options(uri, opts), + ) + .await + { + Ok(Ok(store)) => store, + Ok(Err(e)) => return Err(e), + Err(_) => { + return Err(lance::Error::io(format!( + "open timed out retiring store '{name}'" + ))) + } + }; + + // 1. Drain the WAL. Must precede compaction: `compact_files` rewrites base + // fragments and leaves MemWAL generations untouched. + tokio::time::timeout(RETIRE_TIMEOUT, store.cleanup_own_shard()) + .await + .map_err(|_| lance::Error::io(format!("WAL merge timed out retiring '{name}'")))??; + + // 2. Compact, so the retired table is not left as many small fragments that + // nothing will ever come back to tidy. + tokio::time::timeout(RETIRE_TIMEOUT, store.compact(None)) + .await + .map_err(|_| lance::Error::io(format!("compaction timed out retiring '{name}'")))??; + + // 3. Verify. Only a genuinely drained shard may leave the sweeps' view. + let obs = tokio::time::timeout(OBSERVE_TIMEOUT, store.observe()) + .await + .map_err(|_| lance::Error::io(format!("observe timed out retiring '{name}'")))??; + + Ok(obs.pending_wal_generations == 0) +} + /// Refresh one experiment immediately and persist its new stats row. /// /// Passes `None` as the previous row so this always does a full observation: @@ -644,3 +771,164 @@ mod incremental_scan_tests { assert_eq!(next.total_compactions, 7); } } + +#[cfg(test)] +mod retirement_tests { + use super::*; + use lance_context_core::{RolloutRecord, RolloutStore, RolloutStoreOptions, ROLE_ASSISTANT}; + use tempfile::TempDir; + + fn rec(id: &str) -> RolloutRecord { + RolloutRecord { + id: id.to_string(), + rollout_id: "r".to_string(), + problem_id: "p".to_string(), + dataset: None, + sequence_order: 0, + role: ROLE_ASSISTANT.to_string(), + created_at: Utc::now(), + content: Some("x".to_string()), + content_type: "text/plain".to_string(), + model_input_string: None, + model_output_string: None, + rationale: None, + problem_text: None, + user_metadata: None, + input_tokens: None, + output_tokens: None, + num_input_tokens: None, + num_output_tokens: None, + output_logprobs: None, + input_logprobs: None, + ref_logprobs: None, + loss_mask: None, + advantage: None, + reward: None, + raw_reward: None, + grader_id: None, + score: None, + include_in_training: None, + exclude_reason: None, + policy_version: None, + relationships: vec![], + binary_payload: None, + payload_size: None, + payload_checksum: None, + artifact_type: None, + metadata: None, + } + } + + fn row(name: &str, uri: &str, last_updated: i64) -> StatRow { + StatRow { + name: name.to_string(), + uri: uri.to_string(), + row_count: 1, + fragment_count: 1, + last_updated, + pending_wal_generations: 0, + last_compaction: StatRow::NO_COMPACTION, + total_compactions: 0, + scanned_at: last_updated, + version: 1, + } + } + + /// Retirement must drain the WAL before dropping the row. + /// + /// The sweeps read the stats table and nothing else, so a retired + /// experiment is invisible to them forever. Dropping one with pending + /// generations would strand them: never merged, read amplification never + /// recovers, `_mem_wal/` never reclaimed, and no process left to notice. + #[tokio::test] + async fn retirement_drains_the_wal_before_dropping_the_row() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().join("e.lance").to_string_lossy().to_string(); + { + let store = RolloutStore::open_with_options(&uri, RolloutStoreOptions::default()) + .await + .unwrap(); + for i in 0..3 { + store.add(&[rec(&format!("r{i}"))]).await.unwrap(); + store.flush().await.unwrap(); + } + // Pending generations exist at this point. + let obs = store.observe().await.unwrap(); + assert!( + obs.pending_wal_generations > 0, + "test setup must leave un-merged generations" + ); + } + + let now = Utc::now().timestamp_millis(); + let old = now - Duration::from_secs(30 * 86_400).as_millis() as i64; + let retired = + retire_cold_experiments(&[row("e", &uri, old)], Duration::from_secs(7 * 86_400), now) + .await; + + assert!(retired.contains("e"), "a cold experiment should retire"); + + // The decisive assertion: nothing was left behind for a sweep that will + // never look again. + let store = RolloutStore::open_existing_with_options(&uri, RolloutStoreOptions::default()) + .await + .unwrap(); + let obs = store.observe().await.unwrap(); + assert_eq!( + obs.pending_wal_generations, 0, + "retirement must leave zero pending generations, or they are stranded forever" + ); + assert_eq!(obs.row_count, 3, "no rows may be lost by retirement"); + } + + /// An experiment written recently must not be retired. + #[tokio::test] + async fn recently_written_experiment_is_not_retired() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().join("e.lance").to_string_lossy().to_string(); + { + let store = RolloutStore::open_with_options(&uri, RolloutStoreOptions::default()) + .await + .unwrap(); + store.add(&[rec("a")]).await.unwrap(); + store.flush().await.unwrap(); + } + + let now = Utc::now().timestamp_millis(); + let retired = + retire_cold_experiments(&[row("e", &uri, now)], Duration::from_secs(7 * 86_400), now) + .await; + assert!( + retired.is_empty(), + "a hot experiment must stay in the table" + ); + } + + /// Retirement disabled means nothing is ever retired, however old. + #[tokio::test] + async fn zero_window_disables_retirement() { + let now = Utc::now().timestamp_millis(); + let retired = + retire_cold_experiments(&[row("e", "/nonexistent", 0)], Duration::from_secs(0), now) + .await; + assert!(retired.is_empty()); + } + + /// An experiment that cannot be prepared stays in the table rather than + /// being dropped. Refusing to retire is always safe; retiring early is not. + #[tokio::test] + async fn unpreparable_experiment_is_kept() { + let now = Utc::now().timestamp_millis(); + let old = now - Duration::from_secs(30 * 86_400).as_millis() as i64; + let retired = retire_cold_experiments( + &[row("gone", "/no/such/dataset.lance", old)], + Duration::from_secs(7 * 86_400), + now, + ) + .await; + assert!( + retired.is_empty(), + "an experiment that failed to prepare must not be retired" + ); + } +} diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 19c7b0b..5bf2117 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -555,6 +555,7 @@ mod tests { scan_concurrency: 4, stats_maintenance_every_n_scans: 0, stats_history_ttl_secs: 3_600, + stats_cold_retire_secs: 0, compaction_interval_secs: 0, // Low threshold so a handful of appends crosses it. min_fragments: 2, diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index 726d68d..9d8c772 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -154,6 +154,7 @@ mod tests { scan_concurrency: 4, stats_maintenance_every_n_scans: 0, stats_history_ttl_secs: 3_600, + stats_cold_retire_secs: 0, compaction_interval_secs: 0, min_fragments: 16, target_rows_per_fragment: 1_048_576, diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index 1924397..5109ece 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -912,6 +912,7 @@ mod tests { scan_concurrency: 4, stats_maintenance_every_n_scans: 0, stats_history_ttl_secs: 3_600, + stats_cold_retire_secs: 0, compaction_interval_secs: 0, min_fragments: 16, target_rows_per_fragment: 1_048_576, diff --git a/crates/lance-context-metrics/src/lib.rs b/crates/lance-context-metrics/src/lib.rs index 65d4ad5..3c54d9c 100644 --- a/crates/lance-context-metrics/src/lib.rs +++ b/crates/lance-context-metrics/src/lib.rs @@ -250,6 +250,29 @@ fn describe_metrics() { "_stats versions created since the last successful maintenance pass. \ Sustained growth means old manifests are accumulating on storage." ); + describe_gauge!( + "master_stats_hot_experiments", + "Experiments currently held in the stats table. Bounded by retirement, \ + unlike the registry, which lists every experiment ever created." + ); + describe_gauge!( + "master_scan_experiments_skipped", + "Experiments a scan round skipped because their base version had not moved. \ + A ratio near 1 of total is the healthy state at scale." + ); + describe_gauge!( + "master_scan_experiments_observed", + "Experiments a scan round observed in full." + ); + describe_counter!( + "master_stats_experiments_retired_total", + "Cold experiments merged, compacted, verified quiescent, and dropped from \ + the stats table." + ); + describe_counter!( + "master_stats_retire_failures_total", + "Retirement attempts that failed; the experiment stays in the table and retries." + ); describe_gauge!( "master_rollout_rows", "Total rollout rows across all experiments as of the last scan."