From ba6f01094550aaa96169bcdd29ac7d2e28182f57 Mon Sep 17 00:00:00 2001 From: Beinan Date: Tue, 29 Sep 2026 21:54:21 +0000 Subject: [PATCH] feat(master): repair base tables whose manifest names missing files, automatically Six production stores have manifests that reference data files storage no longer has. Every scan, merge and compaction of them fails with "Not found: .../data/.lance"; no retry changes that, and the cooldown only spaced the failures out. The rows in those fragments are already gone. Meanwhile the stores' WAL generations pile up (one is at 7,985) because nothing can merge into the base table. StorageBase::repair_missing_fragments HEADs every data and deletion file the current manifest names and commits an Operation::Delete of the fragments with missing files, the same transaction a row delete that empties a fragment commits. Everything else, including every WAL generation, is untouched. The master enqueues a Repair task the first time any task fails with that error signature (no counting: the condition is deterministic) and re-enqueues the failed task behind it. Each repair is recorded in etcd with the dropped fragment ids, their row counts and the missing paths, exposed at GET /api/v1/scheduler/repairs and logged at WARN, so "why does this store have fewer rows" is answerable later. Co-Authored-By: Claude Fable 5 --- crates/lance-context-api/src/lib.rs | 30 +++ crates/lance-context-core/src/lib.rs | 2 +- .../lance-context-core/src/rollout_store.rs | 97 ++++++++ crates/lance-context-core/src/store.rs | 21 ++ crates/lance-context-core/src/store_base.rs | 104 ++++++++- crates/lance-context-master/src/routes.rs | 19 +- crates/lance-context-master/src/scheduler.rs | 213 ++++++++++++++++++ crates/lance-context-master/src/task_store.rs | 70 +++++- 8 files changed, 548 insertions(+), 8 deletions(-) diff --git a/crates/lance-context-api/src/lib.rs b/crates/lance-context-api/src/lib.rs index bfb3684..7530e11 100644 --- a/crates/lance-context-api/src/lib.rs +++ b/crates/lance-context-api/src/lib.rs @@ -1486,6 +1486,11 @@ pub enum TaskKind { /// Runs on the master; serialized per experiment against `Compact` because /// both mutate the shared base table. IndexId, + /// Drop base-table fragments whose data files are missing from storage, + /// so a manifest that names files that no longer exist stops failing + /// every merge and compaction. Enqueued automatically when a task fails + /// with that error; the rows in those fragments are already gone. + Repair, } /// Lifecycle state of a scheduled task, generalized from [`CompactJobStatus`] @@ -1504,6 +1509,31 @@ pub enum TaskState { } /// One unit of scheduled work plus its lifecycle, as surfaced to the queue UI. +/// One fragment a repair dropped from a base table. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct RepairedFragment { + pub id: u64, + pub physical_rows: Option, + pub missing_files: Vec, +} + +/// A base-table repair the master performed. Reported by +/// `GET /api/v1/scheduler/repairs`, most recent first, so anyone asking why a +/// store has fewer rows than expected can see when and what was dropped. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct RepairRecord { + pub target: String, + /// Unix ms when the repair committed. + pub repaired_at_ms: i64, + /// Manifest version the repair read and the one it committed. + pub read_version: u64, + pub committed_version: u64, + /// The task whose failure triggered the repair; `Repair` when it was + /// enqueued by hand. + pub triggered_by: TaskKind, + pub dropped: Vec, +} + /// A task target the master's sweeps are skipping because it has failed /// repeatedly. Reported by `GET /api/v1/scheduler/cooldowns`. #[derive(Debug, Clone, Serialize, Deserialize)] diff --git a/crates/lance-context-core/src/lib.rs b/crates/lance-context-core/src/lib.rs index 279b969..a0aec7d 100644 --- a/crates/lance-context-core/src/lib.rs +++ b/crates/lance-context-core/src/lib.rs @@ -74,7 +74,7 @@ pub use storage::{ }; pub use store::{ CompactionConfig, CompactionStats, ContextStore, ContextStoreOptions, DistanceMetric, - IdIndexType, ReadProjection, + DroppedFragment, IdIndexType, ReadProjection, RepairReport, }; // Re-export CompactionMetrics from lance for Python bindings diff --git a/crates/lance-context-core/src/rollout_store.rs b/crates/lance-context-core/src/rollout_store.rs index 4f650c9..a2fd2fc 100644 --- a/crates/lance-context-core/src/rollout_store.rs +++ b/crates/lance-context-core/src/rollout_store.rs @@ -691,6 +691,16 @@ impl RolloutStore { self.base.compact(options).await } + /// See `StorageBase::base_data_files`. + pub fn base_data_files(&self) -> Vec { + self.base.base_data_files() + } + + /// See `StorageBase::repair_missing_fragments`. + pub async fn repair_missing_fragments(&mut self) -> LanceResult { + self.base.repair_missing_fragments().await + } + /// Build a ZoneMap scalar index on the base table's `id` column. Idempotent. /// See `StorageBase::create_key_zonemap_index`. pub async fn create_id_zonemap_index(&mut self) -> LanceResult<()> { @@ -3587,6 +3597,93 @@ mod tests { }); } + /// A manifest that names a data file storage no longer has makes every + /// scan, merge and compaction fail with `Not found`, and no retry fixes + /// it. `repair_missing_fragments` drops exactly those fragments and + /// commits, after which the table works again minus the rows that were + /// already gone. + #[test] + fn repair_drops_fragments_whose_files_are_missing() { + let dir = TempDir::new().unwrap(); + let uri = dir.path().to_string_lossy().to_string(); + let runtime = tokio::runtime::Runtime::new().unwrap(); + runtime.block_on(async { + let mut store = RolloutStore::open_with_options( + &uri, + RolloutStoreOptions { + storage_options: None, + session: None, + shard_id: Some("rollout-0".to_string()), + merge_after_generations: Some(1), + merge_max_generations: None, + merge_max_bytes: None, + pending_generations_warn: None, + merge_budget: None, + }, + ) + .await + .unwrap(); + // Three fragments of one row each. + for i in 0..3 { + store + .add(&[assistant_record(&format!("a-{i}"))]) + .await + .unwrap(); + store.flush().await.unwrap(); + store.maybe_merge_own_shard().await.unwrap(); + } + assert_eq!(store.base.dataset.count_fragments(), 3); + assert_eq!(store.list(None, None).await.unwrap().len(), 3); + + // Nothing missing: no commit. + let v = store.version(); + let report = store.repair_missing_fragments().await.unwrap(); + assert!(report.dropped.is_empty()); + assert_eq!(report.committed_version, None); + assert_eq!(store.version(), v); + + // Lose the second fragment's data file behind Lance's back. + let victim = store.base.dataset.fragments()[1].clone(); + let victim_id = victim.id; + let victim_path = victim.files[0].path.clone(); + std::fs::remove_file(dir.path().join("data").join(&victim_path)).unwrap(); + + // The table is now broken the way production stores were. + let mut broken = + RolloutStore::open_existing_with_options(&uri, RolloutStoreOptions::default()) + .await + .unwrap(); + let err = broken + .compact(Some(CompactionConfig { + min_fragments: 1, + ..Default::default() + })) + .await + .unwrap_err() + .to_string(); + assert!(err.contains("Not found"), "{err}"); + + let report = store.repair_missing_fragments().await.unwrap(); + assert_eq!(report.dropped.len(), 1); + assert_eq!(report.dropped[0].id, victim_id); + assert_eq!(report.dropped[0].physical_rows, Some(1)); + assert_eq!(report.dropped[0].missing_files, vec![victim_path]); + assert!(report.committed_version.unwrap() > report.read_version); + + // Two rows remain and the table compacts again. + let rows = store.list(None, None).await.unwrap(); + assert_eq!(rows.len(), 2); + store + .compact(Some(CompactionConfig { + min_fragments: 1, + ..Default::default() + })) + .await + .unwrap(); + assert_eq!(store.list(None, None).await.unwrap().len(), 2); + }); + } + #[test] fn compact_reduces_fragments_and_preserves_reads() { // Each WAL merge appends a fragment to the base table, so several merges diff --git a/crates/lance-context-core/src/store.rs b/crates/lance-context-core/src/store.rs index 5a0acb9..137f2a9 100644 --- a/crates/lance-context-core/src/store.rs +++ b/crates/lance-context-core/src/store.rs @@ -175,6 +175,27 @@ impl DistanceMetric { } } +/// One fragment dropped by `repair_missing_fragments`. +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct DroppedFragment { + pub id: u64, + /// Rows the fragment held, if the manifest recorded it. These rows are + /// gone; the repair only makes the manifest say so. + pub physical_rows: Option, + /// The files the manifest named that storage does not have. + pub missing_files: Vec, +} + +/// What a `repair_missing_fragments` pass did. +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct RepairReport { + /// Manifest version the repair inspected. + pub read_version: u64, + /// Version the repair committed; `None` when nothing was missing. + pub committed_version: Option, + pub dropped: Vec, +} + /// Statistics about compaction status and history. #[derive(Debug, Clone)] pub struct CompactionStats { diff --git a/crates/lance-context-core/src/store_base.rs b/crates/lance-context-core/src/store_base.rs index 5bdbb30..148108a 100644 --- a/crates/lance-context-core/src/store_base.rs +++ b/crates/lance-context-core/src/store_base.rs @@ -71,7 +71,7 @@ use crate::merge_budget::{MergeMemoryBudget, MergeReservation}; use crate::metrics::{ count, observe_duration, observe_phase, observe_value, timer_elapsed, timer_start, }; -use crate::store::{CompactionConfig, CompactionStats}; +use crate::store::{CompactionConfig, CompactionStats, DroppedFragment, RepairReport}; /// Number of shard manifest files to scan per batch when discovering the latest /// shard state. @@ -1332,6 +1332,108 @@ impl StorageBase { } } + /// Relative paths (under `data/`) of every data file the base table's + /// current manifest references, in fragment order. For tests and + /// diagnostics. + pub fn base_data_files(&self) -> Vec { + self.dataset + .fragments() + .iter() + .flat_map(|f| f.files.iter().map(|d| d.path.clone())) + .collect() + } + + /// Drop every base-table fragment whose data or deletion file is missing + /// from storage, so the manifest stops naming files that do not exist. + /// + /// A manifest can outlive its files (a cleanup that ran against a stale + /// listing, an object-store soft-delete that expired, an interrupted + /// rewrite). From then on every scan, merge and compaction of the table + /// fails with `Not found: .../data/.lance` and nothing retries it + /// into existence: the rows in those fragments are already gone. This + /// commits a `Delete` of exactly those fragment ids, which is the same + /// transaction a row delete that empties a fragment would commit, and + /// leaves every other fragment and every WAL generation untouched. + /// + /// Returns what was dropped so the caller can record it. A dataset with + /// nothing missing commits nothing and returns an empty report. + pub async fn repair_missing_fragments(&mut self) -> LanceResult { + use lance::dataset::transaction::{Operation, Transaction}; + use lance::dataset::write::CommitBuilder; + + self.ensure_writable()?; + self.reload().await?; + let object_store = self.dataset.object_store(None).await?; + let data_dir = self.dataset.data_dir(); + let deletions_dir = self.dataset.deletions_dir(); + let read_version = self.dataset.version().version; + + let mut dropped = Vec::new(); + for fragment in self.dataset.fragments().iter() { + let mut missing = Vec::new(); + for file in &fragment.files { + let path = data_dir.clone().join(file.path.as_str()); + if !object_store.exists(&path).await? { + missing.push(file.path.clone()); + } + } + if let Some(deletion) = &fragment.deletion_file { + // Mirrors lance_table::io::deletion::deletion_file_path. + let name = format!( + "{}-{}-{}.{}", + fragment.id, + deletion.read_version, + deletion.id, + deletion.file_type.suffix() + ); + let path = deletions_dir.clone().join(name.as_str()); + if !object_store.exists(&path).await? { + missing.push(format!("_deletions/{name}")); + } + } + if !missing.is_empty() { + dropped.push(DroppedFragment { + id: fragment.id, + physical_rows: fragment.physical_rows, + missing_files: missing, + }); + } + } + if dropped.is_empty() { + return Ok(RepairReport { + read_version, + committed_version: None, + dropped, + }); + } + + let operation = Operation::Delete { + updated_fragments: Vec::new(), + deleted_fragment_ids: dropped.iter().map(|d| d.id).collect(), + predicate: "repair: fragment data file missing from storage".to_string(), + }; + let committed = CommitBuilder::new(Arc::new(self.dataset.clone())) + .execute(Transaction::new(read_version, operation, None)) + .await?; + let committed_version = committed.version().version; + self.reload().await?; + warn!( + read_version, + committed_version, + fragments = dropped.len(), + rows = dropped + .iter() + .map(|d| d.physical_rows.unwrap_or(0)) + .sum::(), + "repaired base table: dropped fragments whose files are missing from storage" + ); + Ok(RepairReport { + read_version, + committed_version: Some(committed_version), + dropped, + }) + } + /// Build a ZoneMap scalar index on the base table's key column. /// /// The key column is the table's (unenforced) primary key, so a lightweight diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index cf648ab..d3a3ebf 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -13,8 +13,8 @@ use serde::Deserialize; use lance_context_api::{ CompactJobStatus, EnqueueTaskRequest, ExperimentDetail, ExperimentListResponse, - ExperimentRecordsResponse, ExperimentSummary, SqlQueryRequest, SqlQueryResponse, TaskCooldown, - TaskKind, TaskListResponse, TaskRecord, TaskState, + ExperimentRecordsResponse, ExperimentSummary, RepairRecord, SqlQueryRequest, SqlQueryResponse, + TaskCooldown, TaskKind, TaskListResponse, TaskRecord, TaskState, }; use lance_context_core::{rollout_record_to_dto, ListSource, RolloutFilters, RolloutStore}; use tokio::sync::RwLock; @@ -548,6 +548,20 @@ pub async fn list_cooldowns( .map_err(MasterError::from_lance) } +/// `GET /api/v1/scheduler/repairs` — base-table repairs the master performed, +/// most recent first: which fragments were dropped from which store, when, +/// and how many rows they held. +pub async fn list_repairs( + State(state): State>, +) -> Result>, MasterError> { + state + .task_store + .list_repairs() + .await + .map(Json) + .map_err(MasterError::from_lance) +} + /// `GET /api/v1/tasks` — paginated tasks (queue + recent history), newest first. pub async fn list_tasks( State(state): State>, @@ -601,6 +615,7 @@ pub fn api_router() -> Router> { .route("/experiments/{name}/compact/status", get(compact_status)) .route("/tasks", post(enqueue_task).get(list_tasks)) .route("/scheduler/cooldowns", get(list_cooldowns)) + .route("/scheduler/repairs", get(list_repairs)) .route("/tasks/{id}", get(get_task)) .route("/rescan", post(rescan)) } diff --git a/crates/lance-context-master/src/scheduler.rs b/crates/lance-context-master/src/scheduler.rs index 0ae6cb0..2ffe22e 100644 --- a/crates/lance-context-master/src/scheduler.rs +++ b/crates/lance-context-master/src/scheduler.rs @@ -82,6 +82,7 @@ fn kind_label(kind: TaskKind) -> &'static str { TaskKind::Compact => "compact", TaskKind::MergeWal => "merge_wal", TaskKind::IndexId => "index_id", + TaskKind::Repair => "repair", } } @@ -146,6 +147,7 @@ async fn run_task(state: &Arc, claim: TaskClaim, timing: TaskClaimT TaskKind::Compact => run_compaction(state, &task).await, TaskKind::MergeWal => run_merge_wal(state, &task.target).await, TaskKind::IndexId => run_index_id(state, &task.target).await, + TaskKind::Repair => run_repair(state, &task).await, }; let work_elapsed = started.elapsed(); let result = if outcome.is_ok() { "success" } else { "failed" }; @@ -162,6 +164,12 @@ async fn run_task(state: &Arc, claim: TaskClaim, timing: TaskClaimT if let Err(error) = &outcome { tracing::warn!(task = %task.id, target = %task.target, error, "task failed"); + if task.kind != TaskKind::Repair && is_missing_fragment_error(error) { + // The manifest names a file storage does not have. No retry and + // no cooldown changes that; a repair does, so enqueue one now + // and re-run this task behind it. + schedule_repair(state, &task).await; + } } let commit_start = std::time::Instant::now(); let finished = state.task_store.finish(claim, outcome).await; @@ -172,6 +180,107 @@ async fn run_task(state: &Arc, claim: TaskClaim, timing: TaskClaimT } } +/// Whether a task error is the deterministic "manifest names a data file +/// that is not there" failure. Every scan of the base table then fails the +/// same way, from the master's compaction and from every worker's merge +/// alike, until the fragment is dropped. +fn is_missing_fragment_error(error: &str) -> bool { + error.contains("Not found") + && (error.contains("/data/") || error.contains("/_deletions/")) + && (error.contains(".lance") || error.contains(".arrow") || error.contains(".bin")) +} + +/// Enqueue a `Repair` for the failed task's target, then the original task +/// again depending on it. Both are best-effort: the failure is already +/// recorded, and the next sweep enqueues the original kind anyway. +async fn schedule_repair(state: &Arc, failed: &TaskRecord) { + let repair = match enqueue(state, TaskKind::Repair, &failed.target).await { + Ok(repair) => repair, + Err(error) => { + tracing::warn!(target = %failed.target, %error, "failed to enqueue repair"); + return; + } + }; + tracing::warn!( + target = %failed.target, + after = ?failed.kind, + repair = %repair.id, + "base table names missing files; repair enqueued" + ); + if let Err(error) = enqueue_with_deps(state, failed.kind, &failed.target, vec![repair.id]).await + { + tracing::warn!(target = %failed.target, %error, "failed to re-enqueue after repair"); + } +} + +/// Drop base-table fragments whose files are missing (see +/// `StorageBase::repair_missing_fragments`) and record what was dropped. +async fn run_repair(state: &Arc, task: &TaskRecord) -> Result { + let (kind, name) = parse_target(&task.target); + if kind != StoreKind::Rollout { + return Err(format!("repair is not implemented for {kind:?} stores")); + } + let uri = state.rollout_uri(name); + let mut store = RolloutStore::open_existing_with_options(&uri, state.rollout_store_options()) + .await + .map_err(|e| e.to_string())?; + let report = store + .repair_missing_fragments() + .await + .map_err(|e| e.to_string())?; + let Some(committed_version) = report.committed_version else { + return Ok("nothing missing; no repair needed".to_string()); + }; + let rows: usize = report + .dropped + .iter() + .map(|d| d.physical_rows.unwrap_or(0)) + .sum(); + // The task whose failure scheduled this repair is the one re-enqueued + // behind it (see `schedule_repair`); a manual repair has none. + let triggered_by = state + .task_store + .list() + .await + .ok() + .and_then(|tasks| { + tasks + .into_iter() + .find(|t| t.depends_on.contains(&task.id)) + .map(|t| t.kind) + }) + .unwrap_or(TaskKind::Repair); + let record = lance_context_api::RepairRecord { + target: task.target.clone(), + repaired_at_ms: chrono::Utc::now().timestamp_millis(), + read_version: report.read_version, + committed_version, + triggered_by, + dropped: report + .dropped + .iter() + .map(|d| lance_context_api::RepairedFragment { + id: d.id, + physical_rows: d.physical_rows, + missing_files: d.missing_files.clone(), + }) + .collect(), + }; + if let Err(error) = state.task_store.record_repair(&record).await { + tracing::warn!(target = %task.target, %error, "failed to record repair"); + } + metrics::counter!("master_repairs_total").increment(1); + metrics::counter!("master_repair_fragments_dropped_total") + .increment(report.dropped.len() as u64); + Ok(format!( + "dropped {} fragments ({} rows) whose files were missing; version {} -> {}", + report.dropped.len(), + rows, + report.read_version, + committed_version + )) +} + /// Compact one experiment. The task-store claim owns the per-experiment write /// lock for the full execution. /// @@ -921,6 +1030,110 @@ mod tests { worker.abort(); } + #[test] + fn missing_fragment_error_is_recognised() { + assert!(is_missing_fragment_error( + "Wrapped error: Not found: rocketkeep/x.rollout.lance/data/0101abcd.lance, /rustc/..." + )); + assert!(is_missing_fragment_error( + "LanceError(IO): Not found: rocketkeep/x.rollout.lance/_deletions/3-12-7.arrow" + )); + // A missing manifest or WAL generation is a different failure. + assert!(!is_missing_fragment_error( + "Not found: rocketkeep/x.rollout.lance/_versions/12.manifest" + )); + assert!(!is_missing_fragment_error( + "Not found: rocketkeep/x.rollout.lance/_mem_wal/shard/gen_5/_versions" + )); + assert!(!is_missing_fragment_error("HTTP 500 Internal Server Error")); + } + + /// A base table whose manifest names a missing data file fails every + /// compaction with `Not found`. That failure enqueues a `Repair` and the + /// compaction again behind it; the repair drops the dead fragment and is + /// recorded, and the re-run compaction succeeds. + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn missing_fragment_failure_triggers_repair_and_rerun() { + let dir = TempDir::new().unwrap(); + let state = MasterState::new(config(&dir)).await.unwrap(); + let worker = spawn_scheduler(&state); + + let name = "exp"; + let uri = state.rollout_uri(name); + let victim_path = { + let mut store = RolloutStore::open(&uri).await.unwrap(); + for i in 0..4 { + store + .add(&[rollout_record(&format!("r{i}"))]) + .await + .unwrap(); + store.cleanup_own_shard().await.unwrap(); + } + store.base_data_files()[1].clone() + }; + std::fs::remove_file( + dir.path() + .join(format!("{name}.rollout.lance/data/{victim_path}")), + ) + .unwrap(); + state + .registry + .write() + .await + .upsert(name, &uri) + .await + .unwrap(); + crate::scanner::scan_once(&state).await.unwrap(); + + let compact = enqueue(&state, TaskKind::Compact, name).await.unwrap(); + let failed = await_terminal(&state, &compact.id).await; + assert_eq!(failed.state, TaskState::Failed, "got {failed:?}"); + assert!( + failed.error.as_deref().unwrap_or("").contains("Not found"), + "{failed:?}" + ); + + // The failure enqueued a repair and a dependent compaction. + let tasks = state.task_store.list().await.unwrap(); + let repair = tasks + .iter() + .find(|t| t.kind == TaskKind::Repair && t.target == name) + .expect("repair enqueued") + .clone(); + let rerun = tasks + .iter() + .find(|t| t.kind == TaskKind::Compact && t.target == name && t.id != compact.id) + .expect("compaction re-enqueued") + .clone(); + assert_eq!(rerun.depends_on, vec![repair.id.clone()]); + + let repair = await_terminal(&state, &repair.id).await; + assert_eq!(repair.state, TaskState::Done, "got {repair:?}"); + assert!(repair + .detail + .as_deref() + .unwrap() + .starts_with("dropped 1 fragments (1 rows)")); + let rerun = await_terminal(&state, &rerun.id).await; + assert_eq!(rerun.state, TaskState::Done, "got {rerun:?}"); + + let repairs = state.task_store.list_repairs().await.unwrap(); + assert_eq!(repairs.len(), 1); + assert_eq!(repairs[0].target, name); + assert_eq!(repairs[0].triggered_by, TaskKind::Compact); + assert_eq!(repairs[0].dropped.len(), 1); + assert_eq!(repairs[0].dropped[0].missing_files, vec![victim_path]); + + // Three rows remain readable. + let store = RolloutStore::open_existing_with_options(&uri, Default::default()) + .await + .unwrap(); + assert_eq!(store.list(None, None).await.unwrap().len(), 3); + + worker.abort(); + } + /// Manual enqueue of an `IndexId` task -> dispatcher builds the ZoneMap /// index -> task reaches Done with the expected detail summary. #[tokio::test] diff --git a/crates/lance-context-master/src/task_store.rs b/crates/lance-context-master/src/task_store.rs index d3e5b43..d8296e3 100644 --- a/crates/lance-context-master/src/task_store.rs +++ b/crates/lance-context-master/src/task_store.rs @@ -16,7 +16,7 @@ use etcd_client::{ Certificate, Client, Compare, CompareOp, ConnectOptions, GetOptions, Identity, PutOptions, TlsOptions, Txn, TxnOp, }; -use lance_context_api::{TaskCooldown, TaskKind, TaskRecord, TaskState}; +use lance_context_api::{RepairRecord, TaskCooldown, TaskKind, TaskRecord, TaskState}; use lance_context_core::generate_id; use tokio::sync::oneshot; use tokio::task::JoinHandle; @@ -62,7 +62,8 @@ impl TaskKinds { match kind { TaskKind::Compact => self.compact, TaskKind::MergeWal => self.merge_wal, - TaskKind::IndexId => self.index_id, + // Repair is a base-table write like IndexId and shares its pool. + TaskKind::IndexId | TaskKind::Repair => self.index_id, } } } @@ -271,6 +272,19 @@ impl TaskStore { self.inner.list_cooldowns().await } + /// Persist what a repair dropped, keyed by target and time, so it can be + /// answered for later. Kept for `history_ttl_secs` like task history. + pub async fn record_repair(&self, record: &RepairRecord) -> lance::Result<()> { + self.inner + .record_repair(record, self.history_ttl_secs) + .await + } + + /// Every recorded repair, most recent first. + pub async fn list_repairs(&self) -> lance::Result> { + self.inner.list_repairs().await + } + /// Try to acquire a named coordination lock without waiting. This is used /// around the shared Lance stats table, whose mutations must have one writer /// across all master replicas. @@ -1102,6 +1116,50 @@ impl EtcdTaskStore { .collect() } + fn repairs_prefix(&self) -> String { + format!("{}/repairs/", self.prefix) + } + + async fn record_repair(&self, record: &RepairRecord, ttl_secs: u64) -> lance::Result<()> { + let mut client = self.client.clone(); + let lease = client + .lease_grant(ttl_secs.max(60) as i64, None) + .await + .map_err(etcd_error("grant repair lease"))? + .id(); + let key = format!( + "{}{}/{}", + self.repairs_prefix(), + encode_segment(&record.target), + record.repaired_at_ms + ); + let value = serde_json::to_vec(record) + .map_err(|e| lance::Error::io(format!("encode repair: {e}")))?; + client + .put(key, value, Some(PutOptions::new().with_lease(lease))) + .await + .map_err(etcd_error("put repair"))?; + Ok(()) + } + + async fn list_repairs(&self) -> lance::Result> { + let mut client = self.client.clone(); + let response = client + .get(self.repairs_prefix(), Some(GetOptions::new().with_prefix())) + .await + .map_err(etcd_error("list repairs"))?; + let mut out = response + .kvs() + .iter() + .map(|kv| { + serde_json::from_slice::(kv.value()) + .map_err(|e| lance::Error::io(format!("decode repair: {e}"))) + }) + .collect::>>()?; + out.sort_by_key(|r| std::cmp::Reverse(r.repaired_at_ms)); + Ok(out) + } + fn dedupe_key(&self, kind: TaskKind, target: &str, depends_on: &[String]) -> Option { should_dedupe(kind, depends_on).then(|| { format!( @@ -1174,12 +1232,15 @@ fn should_dedupe(kind: TaskKind, depends_on: &[String]) -> bool { depends_on.is_empty() && matches!( kind, - TaskKind::Compact | TaskKind::IndexId | TaskKind::MergeWal + TaskKind::Compact | TaskKind::IndexId | TaskKind::MergeWal | TaskKind::Repair ) } fn requires_target_lock(kind: TaskKind) -> bool { - matches!(kind, TaskKind::Compact | TaskKind::IndexId) + matches!( + kind, + TaskKind::Compact | TaskKind::IndexId | TaskKind::Repair + ) } fn kind_label(kind: TaskKind) -> &'static str { @@ -1187,6 +1248,7 @@ fn kind_label(kind: TaskKind) -> &'static str { TaskKind::Compact => "compact", TaskKind::MergeWal => "merge-wal", TaskKind::IndexId => "index-id", + TaskKind::Repair => "repair", } }