diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index c9d7cd8..15ce912 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -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, +} + +/// Read the durable retry ledger without probing workers or opening WAL data. +pub async fn list_merge_failures( + State(state): State>, + Query(params): Query, +) -> Result, 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 @@ -652,6 +680,7 @@ pub fn api_router() -> Router> { .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)) @@ -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] diff --git a/docs/merge-failures.md b/docs/merge-failures.md new file mode 100644 index 0000000..6b867b7 --- /dev/null +++ b/docs/merge-failures.md @@ -0,0 +1,17 @@ +# Inspecting merge failures + +`GET /api/v1/scheduler/merge-failures?after=` 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.