diff --git a/README.md b/README.md index 7c4403c4..1fadc3a6 100644 --- a/README.md +++ b/README.md @@ -137,16 +137,15 @@ To build and run just this backend, see **Building from source** below. ## Building from source This repo path-deps `asap-precompute-rs` and `asap-gorilla` from -[ASAPCollector](https://github.com/ProjectASAP/ASAPCollector), and -`asap_sketchlib` from -[asap_sketchlib](https://github.com/ProjectASAP/asap_sketchlib). -Clone all three repos as siblings under `~/repos/`: +[ASAPCollector](https://github.com/ProjectASAP/ASAPCollector). +`asap_sketchlib` is resolved by Cargo from its pinned Git dependency, +so it does not need to be cloned as a local build prerequisite. Clone +the backend and collector as siblings under `~/repos/`: ``` ~/repos/ ├── ASAPCollector/ # path-dep'd by this repo -├── ASAPQuery-backend/ # this repo -└── asap_sketchlib/ # path-dep'd by both +└── ASAPQuery-backend/ # this repo ``` Then: diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 744eb1c7..ab18e4d7 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -1704,14 +1704,52 @@ async fn process_range_query_request( } } Err(_) => { - debug!( - "Modern range-query path returned CapabilityMiss for query='{}', \ - falling through to unsupported", - parsed_request.query - ); - match state.adapter.format_unsupported_query_response().await { - Ok(json) => json.into_response(), - Err(status) => status.into_response()} + // ASAP sketch tier returned CapabilityMiss — try ThanosQueryEngine if registered. + let thanos_engine = state + .query_router + .engine_by_id(crate::query_engines::thanos_query_engine::DATA_SOURCE_THANOS_QUERY_ID) + .cloned(); + if let Some(engine) = thanos_engine { + debug!( + "Range-query ASAP miss — forwarding to ThanosQueryEngine: query='{}'", + parsed_request.query + ); + match engine + .execute_range(&parsed_request.query, start_ms, end_ms, step_ms) + .await + { + Ok(query_result) => match state + .adapter + .format_range_success_response( + &query_result, + &promql_utilities::data_model::KeyByLabelNames::default(), + ) + .await + { + Ok(response) => response.into_response(), + Err(status) => status.into_response(), + }, + Err(e) => { + warn!( + "ThanosQueryEngine range query failed: {e}" + ); + match state.adapter.format_unsupported_query_response().await { + Ok(json) => json.into_response(), + Err(status) => status.into_response(), + } + } + } + } else { + debug!( + "Modern range-query path returned CapabilityMiss for query='{}', \ + no ThanosQueryEngine registered — falling through to unsupported", + parsed_request.query + ); + match state.adapter.format_unsupported_query_response().await { + Ok(json) => json.into_response(), + Err(status) => status.into_response(), + } + } } } } @@ -4255,6 +4293,64 @@ aggregations: ); } + // ── Path A2 range-query e2e test ──────────────────────────────────── + + /// GET /api/v1/query_range returns success when ASAP sketch misses + /// and ThanosQueryEngine is registered. Covers the full HTTP path: + /// browser → ASAP backend → ThanosQueryEngine → mock thanos → matrix response + #[tokio::test] + async fn http_query_range_forwards_to_thanos_when_asap_misses() { + use crate::query_engines::thanos_query_engine::forward::test_support::CANNED_MATRIX_BODY; + use crate::query_engines::thanos_query_engine::forward::test_support::spawn_mock_thanos_capture_range; + use crate::query_engines::thanos_query_engine::{ThanosQueryConfig, ThanosQueryEngine}; + + let (mock_url, _captured, _mock_handle) = + spawn_mock_thanos_capture_range(CANNED_MATRIX_BODY).await; + let cfg = ThanosQueryConfig { + base_url: mock_url, + request_timeout: std::time::Duration::from_secs(5), + }; + let engine = ThanosQueryEngine::new(cfg).expect("engine"); + let arc_engine: Arc = Arc::new(engine); + + // Use GorillaObjectStore so ASAP-tier misses and router falls through to Thanos. + let server_port = setup_test_server_with_named_router( + StorageBackend::GorillaObjectStore, + vec![arc_engine], + ) + .await; + + let client = Client::new(); + let resp = client + .get(format!("http://127.0.0.1:{server_port}/api/v1/query_range")) + .query(&[ + ("query", "rate(http_requests_total[1m])"), + ("start", "1700000000"), + ("end", "1700003600"), + ("step", "15"), + ]) + .send() + .await + .expect("request"); + + assert!( + resp.status().is_success(), + "expected 2xx from range query forwarded to Thanos; got {}", + resp.status() + ); + let body: serde_json::Value = resp.json().await.unwrap(); + assert_eq!( + body["status"].as_str().unwrap_or(""), + "success", + "wire response must carry status=success; got {body}" + ); + assert_eq!( + body["data"]["resultType"].as_str().unwrap_or(""), + "matrix", + "wire response must carry resultType=matrix; got {body}" + ); + } + #[tokio::test] async fn http_engine_override_can_target_thanos_query_id() { // X-ASAP-Engine: thanos_query must reach the forwarder diff --git a/data_plane/src/query_engines/routing/query_engine_routing.rs b/data_plane/src/query_engines/routing/query_engine_routing.rs index d0dfda8b..0e212539 100644 --- a/data_plane/src/query_engines/routing/query_engine_routing.rs +++ b/data_plane/src/query_engines/routing/query_engine_routing.rs @@ -69,6 +69,25 @@ pub trait QueryEngine: Send + Sync { /// Answer `query` against this engine's storage tier. async fn execute(&self, query: &str) -> Result; + /// Forward a range query to this engine's storage tier. + /// + /// Params are milliseconds since epoch. The default returns `CapabilityMiss` + /// so existing impls need not change; override when the engine has native + /// range-query support (e.g. Thanos). + async fn execute_range( + &self, + query: &str, + start_ms: u64, + end_ms: u64, + step_ms: u64, + ) -> Result { + let _ = (query, start_ms, end_ms, step_ms); + Err(EngineError::capability_miss( + self.capabilities().data_source_id, + "execute_range not implemented for this engine", + )) + } + /// What this engine can serve. Cheap; the router calls it on every /// `register` and may re-call to refresh cost estimates. fn capabilities(&self) -> EngineCapabilities; diff --git a/data_plane/src/query_engines/thanos_query_engine/forward.rs b/data_plane/src/query_engines/thanos_query_engine/forward.rs index 3abc4770..8016d479 100644 --- a/data_plane/src/query_engines/thanos_query_engine/forward.rs +++ b/data_plane/src/query_engines/thanos_query_engine/forward.rs @@ -129,6 +129,10 @@ impl ThanosQueryConfig { fn instant_endpoint(&self) -> String { format!("{}/api/v1/query", self.base_url) } + + fn range_endpoint(&self) -> String { + format!("{}/api/v1/query_range", self.base_url) + } } /// Forwards PromQL queries to an upstream `thanos-query` sidecar @@ -254,6 +258,67 @@ impl ThanosQueryEngine { .map_err(ThanosQueryError::ParseError)?; Ok(result) } + + /// Forward `query` to `${base_url}/api/v1/query_range` with + /// the given time bounds and return the matrix as a [`QueryResult`]. + pub async fn query_range( + &self, + query: &str, + start_ms: u64, + end_ms: u64, + step_ms: u64, + ) -> Result { + let started = Instant::now(); + let url = self.config.range_endpoint(); + let start_s = format!("{:.3}", start_ms as f64 / 1000.0); + let end_s = format!("{:.3}", end_ms as f64 / 1000.0); + let step_s = format!("{:.3}", step_ms as f64 / 1000.0); + debug!( + url = %url, + query = query, + start = %start_s, + end = %end_s, + step = %step_s, + "thanos-forward: issuing range query", + ); + + let resp = self + .client + .post(&url) + .form(&[ + ("query", query), + ("start", start_s.as_str()), + ("end", end_s.as_str()), + ("step", step_s.as_str()), + ]) + .send() + .await + .map_err(|e| ThanosQueryError::Unreachable(e.to_string()))?; + + let status = resp.status(); + if status.is_server_error() { + return Err(ThanosQueryError::Unreachable(format!( + "upstream returned {status}", + ))); + } + if !status.is_success() { + let body = resp.text().await.unwrap_or_default(); + return Err(ThanosQueryError::BadQuery { + status: status.as_u16(), + body, + }); + } + + let payload: ThanosResponse = resp + .json() + .await + .map_err(|e| ThanosQueryError::ParseError(e.to_string()))?; + + let elapsed_ms = started.elapsed().as_millis(); + let result = build_result_from_thanos_payload(payload, elapsed_ms) + .map_err(ThanosQueryError::ParseError)?; + Ok(result) + } } #[async_trait] @@ -295,6 +360,43 @@ impl QueryEngine for ThanosQueryEngine { } } + async fn execute_range( + &self, + query: &str, + start_ms: u64, + end_ms: u64, + step_ms: u64, + ) -> Result { + match self.query_range(query, start_ms, end_ms, step_ms).await { + Ok(result) => Ok(result), + Err(ThanosQueryError::Unreachable(reason)) => { + warn!( + engine = self.data_source_id, + error = %reason, + "thanos-forward: upstream unreachable (range query)", + ); + Err(crate::query_engines::EngineError::backend( + self.data_source_id, + format!("thanos_unreachable: {reason}"), + )) + } + Err(ThanosQueryError::BadQuery { status, body }) => { + Err(crate::query_engines::EngineError::capability_miss( + self.data_source_id, + format!("thanos rejected range query (status {status}): {body}"), + )) + } + Err(ThanosQueryError::ParseError(msg)) => Err(crate::query_engines::EngineError::backend( + self.data_source_id, + format!("thanos range response parse error: {msg}"), + )), + Err(ThanosQueryError::ConfigInvalid(msg)) => Err(crate::query_engines::EngineError::backend( + self.data_source_id, + format!("thanos client misconfigured: {msg}"), + )), + } + } + fn capabilities(&self) -> EngineCapabilities { EngineCapabilities { data_source_id: self.data_source_id, @@ -662,6 +764,97 @@ pub mod test_support { ] } }"#; + + /// A minimal canned matrix response — one series with two samples. + pub const CANNED_MATRIX_BODY: &str = r#"{ + "status": "success", + "data": { + "resultType": "matrix", + "result": [ + { + "metric": {"__name__": "http_requests_total", "job": "api"}, + "values": [[1700000000.0, "42"], [1700000015.0, "43"]] + } + ] + } + }"#; + + use std::collections::HashMap; + use std::sync::{Arc, Mutex}; + + #[derive(Debug, Clone, Default)] + pub struct CapturedParams(pub HashMap); + + /// Spawn a mock thanos that captures form params of the first POST + /// to `/api/v1/query_range`, then returns `canned_body`. + pub async fn spawn_mock_thanos_capture_range( + canned_body: &'static str, + ) -> (String, Arc>>, tokio::task::JoinHandle<()>) { + use axum::extract::Form; + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let local: SocketAddr = listener.local_addr().expect("local_addr"); + let base_url = format!("http://{local}"); + let captured: Arc>> = Arc::new(Mutex::new(None)); + let cap2 = captured.clone(); + let app = axum::Router::new() + .route( + "/api/v1/query_range", + axum::routing::post(move |Form(params): Form>| { + let captured = cap2.clone(); + async move { + *captured.lock().unwrap() = Some(CapturedParams(params)); + canned_body + } + }), + ) + .route("/api/v1/query", axum::routing::post(move || async move { canned_body })); + let handle = tokio::spawn(async move { + axum::serve(listener, app).await.expect("capture serve"); + }); + tokio::task::yield_now().await; + (base_url, captured, handle) + } + + /// Spawn a mock that sleeps 10 s before answering `/api/v1/query_range`. + pub async fn spawn_mock_thanos_slow_range() -> (String, JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let local: SocketAddr = listener.local_addr().expect("local_addr"); + let base_url = format!("http://{local}"); + let app: Router = Router::new() + .route( + "/api/v1/query_range", + axum::routing::post(|| async { + tokio::time::sleep(std::time::Duration::from_secs(10)).await; + "{}" + }), + ) + .route("/api/v1/query", axum::routing::post(|| async { "{}" })); + let handle = tokio::spawn(async move { + axum::serve(listener, app).await.expect("slow serve"); + }); + tokio::task::yield_now().await; + (base_url, handle) + } + + /// Spawn a mock that returns `status_code` on `/api/v1/query_range`. + pub async fn spawn_mock_thanos_status_range(status_code: u16) -> (String, JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let local: SocketAddr = listener.local_addr().expect("local_addr"); + let base_url = format!("http://{local}"); + let sc = axum::http::StatusCode::from_u16(status_code) + .unwrap_or(axum::http::StatusCode::INTERNAL_SERVER_ERROR); + let app: Router = Router::new() + .route( + "/api/v1/query_range", + axum::routing::post(move || async move { sc }), + ) + .route("/api/v1/query", axum::routing::post(|| async { axum::http::StatusCode::OK })); + let handle = tokio::spawn(async move { + axum::serve(listener, app).await.expect("status serve"); + }); + tokio::task::yield_now().await; + (base_url, handle) + } } // --------------------------------------------------------------------------- @@ -853,4 +1046,122 @@ mod tests { other => panic!("expected Vector, got {other:?}"), } } + + // ── query_range unit tests ────────────────────────────────────────── + + #[tokio::test] + async fn query_range_calls_range_endpoint_not_instant() { + use test_support::{spawn_mock_thanos_capture_range, CANNED_MATRIX_BODY}; + let (url, _captured, _handle) = + spawn_mock_thanos_capture_range(CANNED_MATRIX_BODY).await; + let engine = ThanosQueryEngine::new(config_for(&url)).expect("engine"); + let result = engine + .query_range("rate(http_requests_total[1m])", 1_700_000_000_000, 1_700_003_600_000, 15_000) + .await; + assert!(result.is_ok(), "expected Ok; got {result:?}"); + match result.unwrap() { + QueryResult::Matrix(_) => {} + other => panic!("expected Matrix result, got {other:?}"), + } + } + + #[tokio::test] + async fn query_range_passes_params_correctly() { + use test_support::{spawn_mock_thanos_capture_range, CANNED_MATRIX_BODY}; + let (url, captured, _handle) = + spawn_mock_thanos_capture_range(CANNED_MATRIX_BODY).await; + let engine = ThanosQueryEngine::new(config_for(&url)).expect("engine"); + engine + .query_range("up", 1_700_000_000_000, 1_700_003_600_000, 15_000) + .await + .expect("ok"); + let params = captured.lock().unwrap().clone().expect("params captured"); + assert_eq!(params.0.get("query").map(|s| s.as_str()), Some("up"), + "query param must be forwarded verbatim"); + let start: f64 = params.0["start"].parse().unwrap(); + assert!((start - 1_700_000_000.0).abs() < 0.1, "start must be unix seconds: {start}"); + let end: f64 = params.0["end"].parse().unwrap(); + assert!((end - 1_700_003_600.0).abs() < 0.1, "end must be unix seconds: {end}"); + let step: f64 = params.0["step"].parse().unwrap(); + assert!((step - 15.0).abs() < 0.1, "step must be seconds: {step}"); + } + + #[tokio::test] + async fn query_range_parses_matrix_response() { + use test_support::{spawn_mock_thanos_capture_range, CANNED_MATRIX_BODY}; + let (url, _captured, _handle) = + spawn_mock_thanos_capture_range(CANNED_MATRIX_BODY).await; + let engine = ThanosQueryEngine::new(config_for(&url)).expect("engine"); + let result = engine + .query_range("up", 1_700_000_000_000, 1_700_003_600_000, 15_000) + .await + .expect("ok"); + match result { + QueryResult::Matrix(mv) => { + assert_eq!(mv.values.len(), 1, "one series in canned body"); + assert_eq!(mv.values[0].samples.len(), 2, "two samples in canned body"); + assert!((mv.values[0].samples[0].value - 42.0).abs() < 1e-9); + assert!((mv.values[0].samples[1].value - 43.0).abs() < 1e-9); + } + other => panic!("expected Matrix result, got {other:?}"), + } + } + + #[tokio::test] + async fn query_range_4xx_returns_capability_miss() { + use test_support::spawn_mock_thanos_status_range; + let (url, _handle) = spawn_mock_thanos_status_range(400).await; + let engine = ThanosQueryEngine::new(config_for(&url)).expect("engine"); + let raw = engine.query_range("bad[", 0, 1_000, 1_000).await; + assert!( + matches!(raw, Err(ThanosQueryError::BadQuery { .. })), + "4xx must surface as BadQuery; got {raw:?}", + ); + let trait_err = QueryEngine::execute_range(&engine, "bad[", 0, 1_000, 1_000) + .await.unwrap_err(); + assert!( + matches!(trait_err, crate::query_engines::EngineError::CapabilityMiss { .. }), + "4xx must fold into CapabilityMiss via trait; got {trait_err:?}", + ); + } + + #[tokio::test] + async fn query_range_5xx_returns_backend_error() { + use test_support::spawn_mock_thanos_status_range; + let (url, _handle) = spawn_mock_thanos_status_range(503).await; + let engine = ThanosQueryEngine::new(config_for(&url)).expect("engine"); + let raw = engine.query_range("up", 0, 1_000, 1_000).await; + assert!( + matches!(raw, Err(ThanosQueryError::Unreachable(_))), + "5xx must surface as Unreachable; got {raw:?}", + ); + let trait_err = QueryEngine::execute_range(&engine, "up", 0, 1_000, 1_000) + .await.unwrap_err(); + assert!( + matches!(trait_err, crate::query_engines::EngineError::Backend { .. }), + "5xx must fold into Backend via trait; got {trait_err:?}", + ); + } + + #[tokio::test] + async fn query_range_timeout_returns_backend_error() { + use test_support::spawn_mock_thanos_slow_range; + let (url, _handle) = spawn_mock_thanos_slow_range().await; + let cfg = ThanosQueryConfig { + base_url: url.trim_end_matches('/').to_string(), + request_timeout: Duration::from_millis(100), + }; + let engine = ThanosQueryEngine::new(cfg).expect("engine"); + let raw = engine.query_range("up", 0, 1_000_000_000, 1_000).await; + assert!( + matches!(raw, Err(ThanosQueryError::Unreachable(_))), + "timeout must surface as Unreachable; got {raw:?}", + ); + let trait_err = QueryEngine::execute_range(&engine, "up", 0, 1_000_000_000, 1_000) + .await.unwrap_err(); + assert!( + matches!(trait_err, crate::query_engines::EngineError::Backend { .. }), + "timeout must fold into Backend via trait; got {trait_err:?}", + ); + } }