Skip to content
Merged
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
74 changes: 74 additions & 0 deletions crates/lance-context-master/src/routes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Arc<MasterState>>,
Path(name): Path<String>,
) -> Result<Json<ExperimentDetail>, 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).
Expand Down Expand Up @@ -613,6 +649,7 @@ pub fn api_router() -> Router<Arc<MasterState>> {
)
.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))
Expand Down Expand Up @@ -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() {
Expand Down
Loading