diff --git a/crates/lance-context-master/src/routes.rs b/crates/lance-context-master/src/routes.rs index e0e8eb7..859eff1 100644 --- a/crates/lance-context-master/src/routes.rs +++ b/crates/lance-context-master/src/routes.rs @@ -411,6 +411,42 @@ pub async fn rescan( Ok(Json(serde_json::json!({ "scanned": n }))) } +/// `POST /api/v1/experiments/{name}/rescan` +/// +/// Re-observe one store now and persist its stats row, so the next merge +/// sweep sees its real pending-generation count instead of waiting up to a +/// scan interval. Same work as `GET /experiments/{name}?fresh=true`, exposed +/// as a POST because operators reach for it to *change* scheduler input. +pub async fn rescan_experiment( + State(state): State>, + Path(name): Path, +) -> Result, MasterError> { + let entry = state + .registry + .write() + .await + .get(&name) + .await + .map_err(MasterError::from_lance)? + .ok_or_else(|| MasterError::NotFound(format!("experiment '{}' does not exist", name)))?; + scanner::refresh_one(&state, &entry.name, &entry.uri) + .await + .map_err(MasterError::from_lance)?; + let row = state + .stats + .lock() + .await + .get(&name) + .await + .map_err(MasterError::from_lance)? + .ok_or_else(|| { + MasterError::Internal(format!("stats row for '{}' missing after refresh", name)) + })?; + Ok(Json(ExperimentDetail { + summary: row.into_summary(), + })) +} + /// Map a task's terminal `detail` string back into the fragment counts the /// legacy `CompactJobStatus::Done` shape carries. Best-effort: unparsable /// details yield zeros (the UI only needs the state to stop polling). @@ -613,6 +649,7 @@ pub fn api_router() -> Router> { ) .route("/experiments/{name}/compact", post(compact_experiment)) .route("/experiments/{name}/compact/status", get(compact_status)) + .route("/experiments/{name}/rescan", post(rescan_experiment)) .route("/tasks", post(enqueue_task).get(list_tasks)) .route("/scheduler/cooldowns", get(list_cooldowns)) .route("/scheduler/repairs", get(list_repairs)) @@ -967,6 +1004,43 @@ mod tests { assert!(matches!(res, Err(MasterError::NotFound(_)))); } + /// The merge sweep reads the stats row; a rescan must make a store's + /// real pending count visible without waiting for the next scan round. + #[tokio::test] + #[ignore = "requires ETCD_TEST_ENDPOINTS"] + async fn rescan_experiment_refreshes_pending_generations() { + let dir = TempDir::new().unwrap(); + let state = MasterState::new(test_config(&dir)).await.unwrap(); + let name = "behind"; + let uri = state.rollout_uri(name); + let store = RolloutStore::open(&uri).await.unwrap(); + state + .registry + .write() + .await + .upsert(name, &uri) + .await + .unwrap(); + scanner::scan_once(&state).await.unwrap(); + let before = state.stats.lock().await.get(name).await.unwrap().unwrap(); + assert_eq!(before.pending_wal_generations, 0); + + store.add(&[test_record("a", false)]).await.unwrap(); + store.flush().await.unwrap(); + store.add(&[test_record("b", false)]).await.unwrap(); + store.flush().await.unwrap(); + + let Json(detail) = rescan_experiment(State(state.clone()), Path(name.to_string())) + .await + .unwrap(); + assert_eq!(detail.summary.pending_wal_generations, 2); + let after = state.stats.lock().await.get(name).await.unwrap().unwrap(); + assert_eq!(after.pending_wal_generations, 2); + + let res = rescan_experiment(State(state), Path("missing".to_string())).await; + assert!(matches!(res, Err(MasterError::NotFound(_)))); + } + #[tokio::test] #[ignore = "requires ETCD_TEST_ENDPOINTS"] async fn fresh_detail_refreshes_only_requested_experiment() {