diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index d7d64538b..d66126061 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -673,6 +673,7 @@ pub fn router(state: AppState) -> Router { "/xrpc/sh.tangled.actor.searchActorsTypeahead", get(search_actors_typeahead), ) + .route("/_ready", get(get_ready)) .route("/xrpc/sh.tangled.bobbin.getCoverage", get(get_coverage)) .route("/xrpc/sh.tangled.bobbin.awaitRecord", get(await_record)) .route( @@ -3820,6 +3821,19 @@ async fn resolve_mini_doc( Ok(Json(doc).into_response()) } +async fn get_ready(State(state): State) -> (StatusCode, Json) { + let mut envelope = CoverageEnvelope::from(state.coverage.snapshot()); + envelope.stream = state + .stream_health + .as_deref() + .map(|health| health.snapshot().into()); + let status = envelope + .ready + .then_some(StatusCode::OK) + .unwrap_or(StatusCode::SERVICE_UNAVAILABLE); + (status, Json(envelope)) +} + async fn get_coverage(State(state): State) -> Json { let mut envelope = CoverageEnvelope::from(state.coverage.snapshot()); envelope.stream = state diff --git a/bobbin/crates/xrpc/tests/coverage.rs b/bobbin/crates/xrpc/tests/coverage.rs index 00e2981f8..e860d5d17 100644 --- a/bobbin/crates/xrpc/tests/coverage.rs +++ b/bobbin/crates/xrpc/tests/coverage.rs @@ -331,3 +331,39 @@ async fn any_origin_allows_preflight_and_get() { "*" ); } + +fn ready_request(path: &str) -> Request { + Request::builder() + .uri(path) + .body(Body::empty()) + .unwrap() +} + +#[tokio::test] +async fn ready_returns_503_when_warming() { + let h = Harness::new().await; + for path in ["/_ready", "/xrpc/_ready"] { + let app = router(h.state.clone()); + let (status, body) = json_response(app.oneshot(ready_request(path)).await.unwrap()).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + assert_eq!(body["ready"], json!(false)); + assert_eq!(body["eventsProcessed"], json!(0)); + } +} + +#[tokio::test] +async fn ready_returns_200_when_promoted() { + let h = Harness::new().await; + h.coverage.update(|_| Coverage::Ready { + events_processed: 42, + last_cursor: HydrantCursor::new(100), + }); + for path in ["/_ready", "/xrpc/_ready"] { + let app = router(h.state.clone()); + let (status, body) = json_response(app.oneshot(ready_request(path)).await.unwrap()).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["ready"], json!(true)); + assert_eq!(body["eventsProcessed"], json!(42)); + assert_eq!(body["lastCursor"], json!(100)); + } +}