Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions crates/lance-context-api/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`]
Expand All @@ -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<usize>,
pub missing_files: Vec<String>,
}

/// 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<RepairedFragment>,
}

/// 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)]
Expand Down
2 changes: 1 addition & 1 deletion crates/lance-context-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
97 changes: 97 additions & 0 deletions crates/lance-context-core/src/rollout_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -691,6 +691,16 @@ impl RolloutStore {
self.base.compact(options).await
}

/// See `StorageBase::base_data_files`.
pub fn base_data_files(&self) -> Vec<String> {
self.base.base_data_files()
}

/// See `StorageBase::repair_missing_fragments`.
pub async fn repair_missing_fragments(&mut self) -> LanceResult<crate::RepairReport> {
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<()> {
Expand Down Expand Up @@ -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
Expand Down
21 changes: 21 additions & 0 deletions crates/lance-context-core/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<usize>,
/// The files the manifest named that storage does not have.
pub missing_files: Vec<String>,
}

/// 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<u64>,
pub dropped: Vec<DroppedFragment>,
}

/// Statistics about compaction status and history.
#[derive(Debug, Clone)]
pub struct CompactionStats {
Expand Down
104 changes: 103 additions & 1 deletion crates/lance-context-core/src/store_base.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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<String> {
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/<file>.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<RepairReport> {
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::<usize>(),
"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
Expand Down
19 changes: 17 additions & 2 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Arc<MasterState>>,
) -> Result<Json<Vec<RepairRecord>>, 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<Arc<MasterState>>,
Expand Down Expand Up @@ -601,6 +615,7 @@ pub fn api_router() -> Router<Arc<MasterState>> {
.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))
}
Expand Down
Loading
Loading