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
101 changes: 101 additions & 0 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -568,6 +568,34 @@ pub async fn enqueue_task(
Ok((StatusCode::ACCEPTED, Json(task)))
}

/// Cursor for the read-only merge failure ledger.
#[derive(Debug, Deserialize)]
pub struct MergeFailureParams {
pub after: Option<String>,
}

/// Read the durable retry ledger without probing workers or opening WAL data.
pub async fn list_merge_failures(
State(state): State<Arc<MasterState>>,
Query(params): Query<MergeFailureParams>,
) -> Result<Json<serde_json::Value>, MasterError> {
let (failures, next) = state
.task_store
.merge_coordinator()
.failure_page(params.after.as_deref(), 256)
.await
.map_err(|error| {
if error == "invalid merge failure cursor" {
MasterError::InvalidRequest(error)
} else {
MasterError::Internal(error)
}
})?;
Ok(Json(
serde_json::json!({"failures": failures, "next": next}),
))
}

/// `GET /api/v1/scheduler/cooldowns` — targets the auto-sweeps are skipping
/// because they failed repeatedly, with the failure count, when the cooldown
/// lapses, and the last error. A store in this list is broken in a way that
Expand Down Expand Up @@ -652,6 +680,7 @@ pub fn api_router() -> Router<Arc<MasterState>> {
.route("/experiments/{name}/rescan", post(rescan_experiment))
.route("/tasks", post(enqueue_task).get(list_tasks))
.route("/scheduler/cooldowns", get(list_cooldowns))
.route("/scheduler/merge-failures", get(list_merge_failures))
.route("/scheduler/repairs", get(list_repairs))
.route("/tasks/{id}", get(get_task))
.route("/rescan", post(rescan))
Expand Down Expand Up @@ -805,6 +834,78 @@ mod tests {
}
}

#[tokio::test]
#[ignore = "requires ETCD_TEST_ENDPOINTS"]
async fn merge_failure_api_pages_metadata_without_resetting_retry_budgets() {
use lance_context_api::TaskKind;
let dir = TempDir::new().unwrap();
let state = MasterState::new(test_config(&dir)).await.unwrap();
state
.task_store
.enqueue(TaskKind::MergeWal, "hot", Vec::new())
.await
.unwrap();
let claim = state
.task_store
.claim_next_of_kinds(crate::task_store::TaskKinds::MERGE_WAL)
.await
.unwrap()
.unwrap();
let coordinator = state.task_store.merge_coordinator();
let proof = state.task_store.merge_claim(&claim);
for i in 0..257 {
coordinator
.record_failure(&proof, "hot", &format!("worker-{i:03}"), "storage timeout")
.await
.unwrap();
}
let before = coordinator
.failure("hot", "worker-000")
.await
.unwrap()
.unwrap();
let Json(first) = list_merge_failures(
State(state.clone()),
Query(MergeFailureParams { after: None }),
)
.await
.unwrap();
assert_eq!(first["failures"].as_array().unwrap().len(), 256);
let cursor = first["next"].as_str().unwrap().to_string();
let Json(second) = list_merge_failures(
State(state.clone()),
Query(MergeFailureParams {
after: Some(cursor),
}),
)
.await
.unwrap();
assert_eq!(second["failures"].as_array().unwrap().len(), 1);
assert!(second["next"].is_null());
let after = coordinator
.failure("hot", "worker-000")
.await
.unwrap()
.unwrap();
assert_eq!(after.consecutive_attempts, before.consecutive_attempts);
assert_eq!(after.next_retry_ms, before.next_retry_ms);
assert!(matches!(
list_merge_failures(
State(state.clone()),
Query(MergeFailureParams {
after: Some("/unrelated-prefix/key".into()),
})
)
.await,
Err(MasterError::InvalidRequest(_))
));
state
.task_store
.finish(claim, Err("test complete".into()))
.await
.unwrap();
}

/// Create N experiments via core, register them, scan, and assert the stats
/// table + list endpoint reflect them.
#[tokio::test]
Expand Down
17 changes: 17 additions & 0 deletions docs/merge-failures.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# Inspecting merge failures

`GET /api/v1/scheduler/merge-failures?after=<cursor>` reads the durable retry
ledger. It returns `failures` (at most 256 records) and `next` (an opaque cursor,
or null at the end). An invalid cursor returns HTTP 400. The endpoint only reads
etcd metadata; it never opens WAL payloads, probes workers, retries a merge, or
resets a failure budget. A ledger changing during pagination is not a snapshot.

Each record includes `target`, `endpoint`, `class`, `consecutive_attempts`,
`last_error`, `last_failure_ms`, `next_retry_ms`, and `needs_attention`.
Timestamps are Unix milliseconds. Errors are capped at 4096 characters in the
runtime ledger. `needs_attention` means automatic probes are insufficient to
resolve the incident; it does not authorize data deletion or bypass ownership.

Retry, storage-fencing, and rollout behavior are defined in
[merge-recovery.md](merge-recovery.md). This API adds observation only and does
not change those policies or provide an external notification integration.
Loading