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", } }