diff --git a/crates/lance-context-master/src/scanner.rs b/crates/lance-context-master/src/scanner.rs index e6ed92c..806d020 100644 --- a/crates/lance-context-master/src/scanner.rs +++ b/crates/lance-context-master/src/scanner.rs @@ -507,6 +507,20 @@ struct ScanEntry { kind: ScanKind, } +/// Count pending generations under the observe timeout; shard manifests are +/// one small object read per shard, but object storage still has to answer. +async fn pending_or_timeout( + name: &str, + pending: impl std::future::Future>, +) -> lance::Result { + match tokio::time::timeout(OBSERVE_TIMEOUT, pending).await { + Ok(result) => result, + Err(_) => Err(lance::Error::io(format!( + "pending-generation count timed out for store '{name}'" + ))), + } +} + /// Observe one generic store. Only the fields the WAL-merge sweep and the UI /// need: version (for the skip-if-unchanged shortcut), base row count, /// fragment count and pending MemWAL generations. Compaction counters stay at @@ -529,11 +543,19 @@ async fn observe_generic( } }; + // WAL flushes never bump the base-table version, so the unchanged- + // version shortcut must still recount pending generations: a store that + // no worker or master has merged keeps its old version while its shards + // accumulate thousands of generations, and the merge sweep reads this + // row to decide whether to merge it at all. Only the row count (a + // base-table scan) is carried over. let current_version = store.version() as i64; if let Some(prev) = previous { if prev.version != StatRow::UNKNOWN_VERSION && prev.version == current_version { + let pending = pending_or_timeout(name, store.pending_wal_generations()).await?; let mut row = prev.clone(); row.uri = uri.to_string(); + row.pending_wal_generations = pending as i64; row.scanned_at = Utc::now().timestamp_millis(); return Ok((row, true)); } @@ -593,12 +615,16 @@ async fn observe_one( } }; - // Free: the manifest is already in hand from the open above. + // Free: the manifest is already in hand from the open above. Pending + // generations are not: see `observe_generic` for why they are recounted + // even when the base-table version is unchanged. let current_version = store.version() as i64; if let Some(prev) = previous { if prev.version != StatRow::UNKNOWN_VERSION && prev.version == current_version { + let pending = pending_or_timeout(name, store.pending_wal_generations()).await?; let mut row = prev.clone(); row.uri = uri.to_string(); + row.pending_wal_generations = pending as i64; row.scanned_at = Utc::now().timestamp_millis(); return Ok((row, true)); } @@ -1053,6 +1079,72 @@ mod incremental_scan_tests { .unwrap(); assert!(skipped); assert_eq!(again.pending_wal_generations, 3); + + // More flushes do not move the base version, but the merge sweep + // reads this count: the skipped row must still carry the new one. + { + let store = GenericStore::open( + &uri, + SchemaSpec::new(vec![( + ID_COLUMN.to_string(), + ColumnSpec::required(ColumnType::String { large: false }), + )]), + GenericStoreOptions { + seal_on_add: true, + ..Default::default() + }, + ) + .await + .unwrap(); + for i in 3..5 { + store + .add(&[serde_json::json!({"id": format!("r{i}")}) + .as_object() + .unwrap() + .clone()]) + .await + .unwrap(); + } + } + let (recounted, skipped) = + observe_generic(&name, &uri, Some(&again), GenericStoreOptions::default()) + .await + .unwrap(); + assert!(skipped, "base version unchanged"); + assert_eq!(recounted.version, again.version); + assert_eq!(recounted.pending_wal_generations, 5); + } + + /// WAL flushes never bump the base-table version. A store that no one + /// has merged keeps its version forever while its shards fill up; the + /// unchanged-version shortcut must still recount pending generations or + /// the merge sweep never learns about it (17k generations sat behind a + /// stale "3" in production). + #[tokio::test] + async fn skipped_round_recounts_pending_generations() { + 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 (prev, _) = observe_one("e", &uri, None, RolloutStoreOptions::default()) + .await + .unwrap(); + assert_eq!(prev.pending_wal_generations, 1); + + for id in ["b", "c", "d"] { + store.add(&[rec(id)]).await.unwrap(); + store.flush().await.unwrap(); + } + let (next, skipped) = observe_one("e", &uri, Some(&prev), RolloutStoreOptions::default()) + .await + .unwrap(); + assert!(skipped, "base version unchanged"); + assert_eq!(next.version, prev.version); + assert_eq!(next.pending_wal_generations, 4); } /// A write moves the base version, so the next scan must observe fully and