From 7d337ca0aa7f0686b170dc854b23b2b815bfec18 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Tue, 6 Oct 2026 06:15:56 +0000 Subject: [PATCH 1/2] feat: schedule low-count WAL tails with a durable bounded cursor --- crates/lance-context-master/src/config.rs | 2 + crates/lance-context-master/src/lib.rs | 1 + crates/lance-context-master/src/routes.rs | 1 + crates/lance-context-master/src/scheduler.rs | 2 + crates/lance-context-master/src/state.rs | 1 + crates/lance-context-master/src/task_store.rs | 1 + crates/lance-context-master/src/wal_tail.rs | 357 ++++++++++++++++++ docs/wal-tail-sweep.md | 41 ++ 8 files changed, 406 insertions(+) create mode 100644 crates/lance-context-master/src/wal_tail.rs create mode 100644 docs/wal-tail-sweep.md diff --git a/crates/lance-context-master/src/config.rs b/crates/lance-context-master/src/config.rs index 84ccf2a2..fe639a54 100644 --- a/crates/lance-context-master/src/config.rs +++ b/crates/lance-context-master/src/config.rs @@ -14,6 +14,8 @@ pub struct MasterConfig { pub catchup: crate::catchup::CatchupConfig, #[command(flatten)] pub append: crate::rollout_append::AppendConfig, + #[command(flatten)] + pub wal_tail: crate::wal_tail::WalTailConfig, /// Data directory / object-store prefix shared with the data-plane server. #[arg(long, env = "DATA_DIR", default_value = "./data")] pub data_dir: String, diff --git a/crates/lance-context-master/src/lib.rs b/crates/lance-context-master/src/lib.rs index 38e9fede..6157d7a5 100644 --- a/crates/lance-context-master/src/lib.rs +++ b/crates/lance-context-master/src/lib.rs @@ -14,3 +14,4 @@ pub mod scheduler; pub mod state; pub mod stats_store; pub mod task_store; +pub mod wal_tail; diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index 58eb0c49..aa4ac1b7 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -803,6 +803,7 @@ mod tests { MasterConfig { append: Default::default(), catchup: Default::default(), + wal_tail: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), data_dir: dir.path().to_string_lossy().to_string(), diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index fa59a29b..7137f7c8 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -889,6 +889,7 @@ async fn enqueue_merge_request( /// not every scheduler task -- adequate for tests, which drop the whole /// `MasterState` immediately after. pub fn spawn_scheduler(state: &Arc) -> JoinHandle<()> { + crate::wal_tail::spawn(state); // Retry metadata independently of coarse stats sweeps. Every replica may // enqueue; the existing task dedupe/claim transaction elects one executor. let retry_state = state.clone(); @@ -1129,6 +1130,7 @@ mod tests { MasterConfig { append: Default::default(), catchup: Default::default(), + wal_tail: Default::default(), maintenance: Default::default(), merge_rollout: lance_context_merge::rollout::MergeRollout { owned_targets: ["exp", "generic:gs", "broken"] diff --git a/crates/lance-context-master/src/state.rs b/crates/lance-context-master/src/state.rs index 2e9a86d0..b13cee85 100644 --- a/crates/lance-context-master/src/state.rs +++ b/crates/lance-context-master/src/state.rs @@ -309,6 +309,7 @@ mod tests { MasterConfig { append: Default::default(), catchup: Default::default(), + wal_tail: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), data_dir: dir.path().to_string_lossy().to_string(), diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index ddeb320f..6ecffc18 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -2216,6 +2216,7 @@ mod tests { MasterConfig { append: Default::default(), catchup: Default::default(), + wal_tail: Default::default(), maintenance: Default::default(), merge_rollout: Default::default(), data_dir: dir.path().to_string_lossy().to_string(), diff --git a/crates/lance-context-master/src/wal_tail.rs b/crates/lance-context-master/src/wal_tail.rs new file mode 100644 index 00000000..c279ab24 --- /dev/null +++ b/crates/lance-context-master/src/wal_tail.rs @@ -0,0 +1,357 @@ +//! Bounded metadata-only service for WAL tails below the normal count trigger. +//! A service interval is not a claim about the creation age of a generation. +use crate::{config::MasterConfig, scheduler, state::MasterState, stats_store::StatRow}; +use etcd_client::{Compare, CompareOp, Txn, TxnOp}; +use lance_context_api::TaskKind; +use serde::{Deserialize, Serialize}; +use std::{sync::Arc, time::Duration}; + +#[derive(Clone, Debug, clap::Args)] +pub struct WalTailConfig { + /// Revisit low-count WAL tails after this interval; 0 disables this sweep. + /// Requires ordinary automatic WAL merge and worker endpoints as well. + #[arg(long, env = "MERGE_WAL_TAIL_INTERVAL_SECS", default_value_t = 0)] + pub wal_tail_interval_secs: u64, + /// At most this many candidates per fleet-wide 30-second batch (hard cap 64). + #[arg(long, env = "MERGE_WAL_TAIL_BATCH_SIZE", default_value_t = 16)] + pub wal_tail_batch_size: usize, + #[arg(long, env = "MERGE_WAL_TAIL_STATS_MAX_AGE_SECS", default_value_t = 900)] + pub wal_tail_stats_max_age_secs: u64, +} + +impl Default for WalTailConfig { + fn default() -> Self { + Self { + wal_tail_interval_secs: 0, + wal_tail_batch_size: 16, + wal_tail_stats_max_age_secs: 900, + } + } +} + +#[derive(Default, Serialize, Deserialize)] +struct Cursor { + after: Option, + next_batch_ms: i64, +} + +fn enabled(config: &MasterConfig) -> bool { + config.wal_tail.wal_tail_interval_secs > 0 + && config.merge_wal_interval_secs > 0 + && !config.worker_endpoints.is_empty() + && config.wal_tail.wal_tail_stats_max_age_secs > 0 +} + +fn candidates( + config: &MasterConfig, + snapshot: &[StatRow], + after: Option<&str>, + now: i64, +) -> Vec { + let age = config.wal_tail.wal_tail_stats_max_age_secs.min(86_400) as i64 * 1000; + let mut names: Vec<_> = snapshot + .iter() + .filter(|row| { + row.pending_wal_generations > 0 + && row.pending_wal_generations < config.merge_wal_min_generations + && row.scanned_at <= now + && now.saturating_sub(row.scanned_at) <= age + && !config.merge_rollout.draining(&row.name) + && after.is_none_or(|after| row.name.as_str() > after) + }) + .map(|row| row.name.clone()) + .collect(); + names.sort_unstable(); + names.dedup(); + names.truncate(config.wal_tail.wal_tail_batch_size.clamp(1, 64) + 1); + names +} + +pub(crate) fn spawn(state: &Arc) { + if !enabled(&state.config) { + return; + } + let weak = Arc::downgrade(state); + tokio::spawn(async move { + loop { + let Some(state) = weak.upgrade() else { + return; + }; + if let Err(error) = tick(&state, chrono::Utc::now().timestamp_millis()).await { + tracing::warn!(%error, "WAL tail sweep failed; ordinary scheduling continues"); + } + drop(state); + tokio::time::sleep(Duration::from_secs(30)).await; + } + }); +} + +async fn tick(state: &Arc, now: i64) -> Result { + if !enabled(&state.config) { + return Ok(0); + } + let Some(_operation) = state.admission.try_admit() else { + return Ok(0); + }; + let key = format!( + "{}/wal-tail-cursor", + state.config.etcd.etcd_prefix.trim_end_matches('/') + ); + let mut client = state.task_store.etcd_client().clone(); + let prior = client + .get(key.as_str(), None) + .await + .map_err(|e| e.to_string())?; + let old = prior.kvs().first(); + let cursor: Cursor = old + .map(|kv| serde_json::from_slice(kv.value())) + .transpose() + .map_err(|e| e.to_string())? + .unwrap_or_default(); + if cursor.next_batch_ms > now { + return Ok(0); + } + // Existing in-memory scalar stats only: no datasets, payloads, WAL listing + // or stats-writer lock. Followers without a fresh snapshot do no work. + let snapshot = state.stats_cache.read().await.clone(); + let mut names = candidates(&state.config, &snapshot, cursor.after.as_deref(), now); + if names.is_empty() && cursor.after.is_none() { + return Ok(0); + } + let limit = state.config.wal_tail.wal_tail_batch_size.clamp(1, 64); + let more = names.len() > limit; + names.truncate(limit); + let next = Cursor { + after: if more { names.last().cloned() } else { None }, + next_batch_ms: now.saturating_add(if more { + 30_000 + } else { + state + .config + .wal_tail + .wal_tail_interval_secs + .clamp(30, 86_400) as i64 + * 1000 + }), + }; + // Reserve the bounded page BEFORE any enqueues. A lost response or crash + // may defer this page until the next rotation; it never authorizes replay + // or a second replica's burst. The cursor survives master replacement. + let compare = old.map_or_else( + || Compare::version(key.as_str(), CompareOp::Equal, 0), + |kv| Compare::mod_revision(key.as_str(), CompareOp::Equal, kv.mod_revision()), + ); + if !client + .txn(Txn::new().when([compare]).and_then([TxnOp::put( + key.as_str(), + serde_json::to_vec(&next).map_err(|e| e.to_string())?, + None, + )])) + .await + .map_err(|e| e.to_string())? + .succeeded() + { + return Ok(0); + } + let mut queued = 0; + for target in names { + let result = enqueue_tail(state, &target).await; + match result { + Ok(true) => queued += 1, + Ok(false) => {} + Err(error) => { + tracing::warn!(%target, %error, "WAL tail enqueue unresolved; no immediate replay") + } + } + } + metrics::counter!("master_wal_tail_enqueue_requests_total").increment(queued as u64); + Ok(queued) +} + +async fn enqueue_tail(state: &Arc, target: &str) -> Result { + if state + .task_store + .is_cooling_down(TaskKind::MergeWal, target) + .await + .map_err(|e| e.to_string())? + || state + .task_store + .get_active_id(TaskKind::MergeWal, target) + .await + .map_err(|e| e.to_string())? + .is_some() + { + return Ok(false); + } + // Persistent publishers and catch-up Jobs already cover their own tails. + // This loop must never create a second publisher or change their records. + let owner = crate::catchup::store::active_key(&state.config.etcd.etcd_prefix, target); + if !state + .task_store + .etcd_client() + .clone() + .get(owner, None) + .await + .map_err(|e| e.to_string())? + .kvs() + .is_empty() + { + return Ok(false); + } + scheduler::enqueue(state, TaskKind::MergeWal, target) + .await + .map_err(|e| e.to_string())?; + Ok(true) +} + +#[cfg(test)] +mod tests { + use super::*; + use clap::Parser; + + fn config() -> MasterConfig { + let mut cfg = MasterConfig::parse_from(["master"]); + cfg.wal_tail.wal_tail_interval_secs = 3600; + cfg.wal_tail.wal_tail_batch_size = 2; + cfg.worker_endpoints = vec!["http://unused-test-worker".into()]; + cfg + } + + fn row(name: &str, pending: i64, at: i64) -> StatRow { + StatRow { + name: name.into(), + uri: format!("/unused/{name}"), + row_count: 0, + fragment_count: 0, + last_updated: at, + pending_wal_generations: pending, + last_compaction: -1, + total_compactions: 0, + scanned_at: at, + version: -1, + } + } + + #[test] + fn only_fresh_low_count_tails_enter_the_bounded_rotation() { + let mut cfg = config(); + cfg.merge_rollout.drain_targets = vec!["draining".into()]; + let now = 2_000_000; + let rows = vec![ + row("d", 4, now), + row("a", 1, now), + row("b", 2, now), + row("c", 3, now), + row("empty", 0, now), + row("hot", 100, now), + row("stale", 1, now - 901_000), + row("future", 1, now + 1), + row("unknown", -1, now), + row("draining", 1, now), + ]; + assert_eq!(candidates(&cfg, &rows, None, now), ["a", "b", "c"]); + assert_eq!(candidates(&cfg, &rows, Some("b"), now), ["c", "d"]); + assert!(enabled(&cfg)); + cfg.merge_wal_interval_secs = 0; + assert!( + !enabled(&cfg), + "maintenance-only masters must stay disabled" + ); + cfg.merge_wal_interval_secs = 600; + cfg.worker_endpoints.clear(); + assert!(!enabled(&cfg)); + } + + async fn fixture() -> (tempfile::TempDir, Arc) { + let dir = tempfile::tempdir().unwrap(); + let mut cfg = config(); + cfg.data_dir = dir.path().to_string_lossy().into(); + cfg.etcd.etcd_endpoints = std::env::var("ETCD_TEST_ENDPOINTS") + .expect("ETCD_TEST_ENDPOINTS required") + .split(',') + .map(str::to_string) + .collect(); + cfg.etcd.etcd_prefix = format!("/wal-tail-test/{}", lance_context_core::generate_id()); + (dir, MasterState::new(cfg).await.unwrap()) + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn replicas_share_a_page_and_restart_continues_the_cursor() { + let (_dir, state) = fixture().await; + let now = chrono::Utc::now().timestamp_millis(); + *state.stats_cache.write().await = + Arc::new(vec![row("c", 1, now), row("a", 1, now), row("b", 1, now)]); + let (one, two) = tokio::join!(tick(&state, now), tick(&state, now)); + assert_eq!(one.unwrap() + two.unwrap(), 2); + assert!(state + .task_store + .get_active_id(TaskKind::MergeWal, "c") + .await + .unwrap() + .is_none()); + let successor = MasterState::new(state.config.clone()).await.unwrap(); + *successor.stats_cache.write().await = state.stats_cache.read().await.clone(); + assert_eq!(tick(&successor, now + 1000).await.unwrap(), 0); + assert_eq!(tick(&successor, now + 30_000).await.unwrap(), 1); + assert!(successor + .task_store + .get_active_id(TaskKind::MergeWal, "c") + .await + .unwrap() + .is_some()); + assert_eq!(tick(&successor, now + 60_000).await.unwrap(), 0); + // The sweep never opens these intentionally nonexistent table paths. + assert!(!std::path::Path::new("/unused/a").exists()); + } + + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn external_owners_and_already_queued_work_are_not_replaced() { + let (_dir, state) = fixture().await; + let now = chrono::Utc::now().timestamp_millis(); + *state.stats_cache.write().await = + Arc::new(vec![row("a", 1, now), row("b", 1, now), row("c", 1, now)]); + let owner_key = crate::catchup::store::active_key(&state.config.etcd.etcd_prefix, "a"); + state + .task_store + .etcd_client() + .clone() + .put(owner_key.as_str(), "external-publisher", None) + .await + .unwrap(); + let queued = scheduler::enqueue(&state, TaskKind::MergeWal, "b") + .await + .unwrap(); + assert_eq!(tick(&state, now).await.unwrap(), 0); + assert!(state + .task_store + .get_active_id(TaskKind::MergeWal, "a") + .await + .unwrap() + .is_none()); + assert_eq!( + state + .task_store + .get_active_id(TaskKind::MergeWal, "b") + .await + .unwrap() + .as_deref(), + Some(queued.id.as_str()) + ); + assert_eq!( + state + .task_store + .etcd_client() + .clone() + .get(owner_key, None) + .await + .unwrap() + .kvs()[0] + .value(), + b"external-publisher" + ); + // Busy rows use a page position, not the whole fleet's remaining budget. + assert_eq!(tick(&state, now + 30_000).await.unwrap(), 1); + } +} diff --git a/docs/wal-tail-sweep.md b/docs/wal-tail-sweep.md new file mode 100644 index 00000000..6db7a096 --- /dev/null +++ b/docs/wal-tail-sweep.md @@ -0,0 +1,41 @@ +# Automatic service for low-count WAL tails + +A count-only merge trigger can leave one to seven sealed generations indefinitely +when ingestion stops below the default eight-generation threshold. The optional +master tail sweep periodically enqueues this work through the existing durable +scheduler. It does not use table names or table modification timestamps as a +proxy for the actual age of a WAL generation. + +Set `MERGE_WAL_TAIL_INTERVAL_SECS=3600` on ordinary merge-scheduling masters to +revisit low-count tails after each rotation. The default is zero (disabled). +`MERGE_WAL_INTERVAL_SECS` must also be nonzero and workers must be configured; +maintenance-only replicas remain disabled. This does not enable the Kubernetes +catch-up controller or change any table's merge protocol. + +The sweep uses the existing in-memory scalar stats snapshot, with +`MERGE_WAL_TAIL_STATS_MAX_AGE_SECS=900` by default. It never opens a table, scans +payloads or lists WAL objects. Missing, future or stale observations are excluded. +Only positive counts below `MERGE_WAL_MIN_GENERATIONS` enter this rotation; hot +work keeps the normal sweep. An ingestion/write/read hot path gains no new I/O. + +An etcd CAS reserves at most `MERGE_WAL_TAIL_BATCH_SIZE=16` candidates per +fleet-wide 30-second batch, hard capped at 64. The durable name cursor prevents +busy early targets from hiding later targets and survives master replacement. +The interval applies after the last page of a rotation; it is not an exact age +measurement or a guaranteed completion deadline for every table. Existing task +cooldowns, active tasks, draining targets and dedicated catch-up owners are +respected. No lock, owner, attempt record or failure ledger is deleted or reset. + +Page reservation precedes enqueue. A lost response or process crash may defer a +page to the next rotation, but never causes immediate replay of an uncertain +mutation or one replica's workload to be multiplied by every other replica. +An enqueue error is isolated to its target; other candidates continue. A new +rotation evaluates current fresh demand through normal scheduler deduplication. +`master_wal_tail_enqueue_requests_total` counts accepted enqueue requests, not +completed generations or necessarily newly created tasks during concurrent +normal scheduling. + +This sweep covers low-count ordinary-worker WAL tails. It deliberately leaves +externally registered persistent publisher lifecycle to its current owner; those +tables need the separate desired-coverage migration before external babysitting +can be retired. Bounded rolling metadata coverage remains required. From 37b51e0fcfcf4821e47ab6c23fcefbfb5be42364 Mon Sep 17 00:00:00 2001 From: Beinan Wang <> Date: Tue, 6 Oct 2026 06:20:42 +0000 Subject: [PATCH 2/2] fix: prevent empty followers from advancing WAL tail cursor --- crates/lance-context-master/src/wal_tail.rs | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/crates/lance-context-master/src/wal_tail.rs b/crates/lance-context-master/src/wal_tail.rs index c279ab24..9d60f018 100644 --- a/crates/lance-context-master/src/wal_tail.rs +++ b/crates/lance-context-master/src/wal_tail.rs @@ -114,6 +114,20 @@ async fn tick(state: &Arc, now: i64) -> Result { // Existing in-memory scalar stats only: no datasets, payloads, WAL listing // or stats-writer lock. Followers without a fresh snapshot do no work. let snapshot = state.stats_cache.read().await.clone(); + let max_age = state + .config + .wal_tail + .wal_tail_stats_max_age_secs + .min(86_400) as i64 + * 1000; + // An uninitialized/stale follower must not consume or reset the shared + // cursor before a replica with a usable snapshot can finish the rotation. + if !snapshot + .iter() + .any(|row| row.scanned_at <= now && now.saturating_sub(row.scanned_at) <= max_age) + { + return Ok(0); + } let mut names = candidates(&state.config, &snapshot, cursor.after.as_deref(), now); if names.is_empty() && cursor.after.is_none() { return Ok(0); @@ -291,6 +305,11 @@ mod tests { .unwrap() .is_none()); let successor = MasterState::new(state.config.clone()).await.unwrap(); + assert_eq!( + tick(&successor, now + 30_000).await.unwrap(), + 0, + "empty follower must not advance the shared cursor" + ); *successor.stats_cache.write().await = state.stats_cache.read().await.clone(); assert_eq!(tick(&successor, now + 1000).await.unwrap(), 0); assert_eq!(tick(&successor, now + 30_000).await.unwrap(), 1);