From 53c44dc81b1e664b7ff45c754f6a4e6b1d330b5b Mon Sep 17 00:00:00 2001 From: Milind Srivastava Date: Mon, 21 Sep 2026 00:30:35 -0400 Subject: [PATCH 1/4] test: add query forwarding disable mode --- .../src/drivers/query/adapters/config.rs | 25 ++++ data_plane/src/drivers/query/servers/http.rs | 18 +++ .../src/drivers/query/servers/metrics.rs | 14 ++ data_plane/src/main.rs | 120 +++++++++++++++++- .../query_engines/asap_query_engine/engine.rs | 37 +++++- .../asap_query_engine/exact_subqueries.rs | 18 +++ data_plane/src/query_engines/mod.rs | 2 + .../src/query_engines/query_forwarding.rs | 25 ++++ .../src/tests/prometheus_forwarding_tests.rs | 85 ++++++++++++- .../disable-query-forwarding-decisions.md | 25 ++++ 10 files changed, 363 insertions(+), 6 deletions(-) create mode 100644 data_plane/src/query_engines/query_forwarding.rs create mode 100644 docs/design_docs/disable-query-forwarding-decisions.md diff --git a/data_plane/src/drivers/query/adapters/config.rs b/data_plane/src/drivers/query/adapters/config.rs index 4e2976e12..044b3ef72 100644 --- a/data_plane/src/drivers/query/adapters/config.rs +++ b/data_plane/src/drivers/query/adapters/config.rs @@ -1,4 +1,5 @@ use crate::drivers::query::fallback::FallbackClient; +use crate::query_engines::QueryForwardingPolicy; use crate::storage_engines::types::enums::{QueryLanguage, QueryProtocol}; use std::sync::Arc; @@ -13,6 +14,9 @@ pub struct AdapterConfig { /// Optional fallback client for unsupported queries pub fallback: Option>, + + /// Whether this adapter may issue query requests to its fallback backend. + pub query_forwarding_policy: QueryForwardingPolicy, } impl std::fmt::Debug for AdapterConfig { @@ -24,6 +28,7 @@ impl std::fmt::Debug for AdapterConfig { "fallback", &self.fallback.as_ref().map(|_| "Some(FallbackClient)"), ) + .field("query_forwarding_policy", &self.query_forwarding_policy) .finish() } } @@ -39,9 +44,18 @@ impl AdapterConfig { protocol, language, fallback, + query_forwarding_policy: QueryForwardingPolicy::Enabled, } } + pub fn with_query_forwarding_policy(mut self, policy: QueryForwardingPolicy) -> Self { + self.query_forwarding_policy = policy; + if !policy.allows_external_queries() { + self.fallback = None; + } + self + } + /// Create a configuration for Prometheus HTTP with PromQL /// Convenience constructor for backward compatibility pub fn prometheus_promql(fallback_url: String, forward_unsupported: bool) -> Self { @@ -88,4 +102,15 @@ mod tests { QueryLanguage::MetricsQl ); } + + #[test] + fn disabled_query_forwarding_removes_the_fallback_client() { + let config = AdapterConfig::prometheus_promql("http://prom:9090".into(), true) + .with_query_forwarding_policy(QueryForwardingPolicy::Disabled); + assert!(config.fallback.is_none()); + assert_eq!( + config.query_forwarding_policy, + QueryForwardingPolicy::Disabled + ); + } } diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index 55f365cf9..fd3cb978c 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -1236,6 +1236,7 @@ async fn process_via_simple_engine( Err(status) => status.into_response(), } } else { + record_disabled_forwarding(state, "prometheus", "instant"); debug!("Query not supported and forwarding disabled, returning error"); // Adapter formats the unsupported query error for its protocol. // We still annotate `data_source: asap_query` so callers @@ -2034,6 +2035,21 @@ fn query_status_label(response: &Response) -> &'static str { } } +fn record_disabled_forwarding(state: &AppState, backend: &str, path: &str) { + if !state + .config + .adapter_config + .query_forwarding_policy + .allows_external_queries() + { + debug!( + backend, + path, "query forwarding disabled; external query blocked" + ); + srv_metrics::record_query_forwarding_blocked(backend, path); + } +} + // ============================================================ // Metrics Handler // ============================================================ @@ -2248,6 +2264,7 @@ async fn process_range_query_request( )})), ) .into_response() + } Err(EngineRouterError::AllFailed { last }) => { use crate::query_engines::EngineError; @@ -2273,6 +2290,7 @@ async fn process_range_query_request( Err(status) => status.into_response(), } } else { + record_disabled_forwarding(state, "prometheus", "range"); match state.adapter.format_unsupported_query_response().await { Ok(json) => json.into_response(), Err(status) => status.into_response(), diff --git a/data_plane/src/drivers/query/servers/metrics.rs b/data_plane/src/drivers/query/servers/metrics.rs index 46cbeeecc..9af5d7ffe 100644 --- a/data_plane/src/drivers/query/servers/metrics.rs +++ b/data_plane/src/drivers/query/servers/metrics.rs @@ -34,6 +34,13 @@ lazy_static! { ) .unwrap(); + pub static ref QUERY_FORWARDING_BLOCKED_TOTAL: CounterVec = register_counter_vec!( + "asap_query_forwarding_blocked_total", + "External query attempts blocked by the query-forwarding policy", + &["backend", "path"] + ) + .unwrap(); + pub static ref INGEST_SAMPLES_TOTAL: CounterVec = register_counter_vec!( "asap_ingest_samples_total", "Raw samples accepted by the ingest server, labelled by wire protocol", @@ -63,6 +70,7 @@ lazy_static! { pub fn register_all() { lazy_static::initialize(&QUERY_REQUESTS_TOTAL); lazy_static::initialize(&QUERY_DURATION_SECONDS); + lazy_static::initialize(&QUERY_FORWARDING_BLOCKED_TOTAL); lazy_static::initialize(&INGEST_SAMPLES_TOTAL); lazy_static::initialize(&INGEST_BATCH_DURATION_SECONDS); lazy_static::initialize(&INGEST_DECODE_ERRORS_TOTAL); @@ -79,3 +87,9 @@ pub fn record_query_outcome(query_type: &str, status: &str) { .with_label_values(&[query_type, status]) .inc(); } + +pub fn record_query_forwarding_blocked(backend: &str, path: &str) { + QUERY_FORWARDING_BLOCKED_TOTAL + .with_label_values(&[backend, path]) + .inc(); +} diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 08e7bdc62..28110fa6a 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -1,6 +1,7 @@ use clap::{Parser, ValueEnum}; use std::fs; use std::sync::Arc; +use thiserror::Error; use tokio::signal; use tracing::{error, info, warn}; @@ -155,6 +156,10 @@ struct Args { #[arg(long)] forward_unsupported_queries: bool, + /// Disable all external query forwarding for isolated tests. + #[arg(long)] + disable_query_forwarding: bool, + /// Database path (currently unused, kept for compatibility) #[arg(long, default_value = "sketchdb.db")] db_path: String, @@ -382,7 +387,40 @@ struct Args { backend_storage_routing: Option, } +#[derive(Debug, Error)] +enum QueryForwardingConfigError { + #[error("--disable-query-forwarding conflicts with --forward-unsupported-queries")] + ConflictingFlags, + #[error("--profile asapquery requires query forwarding and cannot be combined with --disable-query-forwarding")] + AsapqueryRequiresForwarding, + #[error("--disable-query-forwarding cannot be combined with --victoriametrics-http-port")] + VictoriaMetricsListenerConfigured, + #[error("--disable-query-forwarding cannot be combined with --clickhouse-http-port")] + ClickHouseListenerConfigured, +} + +fn validate_query_forwarding_configuration( + args: &Args, +) -> std::result::Result<(), QueryForwardingConfigError> { + if args.disable_query_forwarding { + if args.forward_unsupported_queries { + return Err(QueryForwardingConfigError::ConflictingFlags); + } + if args.profile == RuntimeProfile::Asapquery { + return Err(QueryForwardingConfigError::AsapqueryRequiresForwarding); + } + if args.victoriametrics_http_port.is_some() { + return Err(QueryForwardingConfigError::VictoriaMetricsListenerConfigured); + } + if args.clickhouse_http_port.is_some() { + return Err(QueryForwardingConfigError::ClickHouseListenerConfigured); + } + } + Ok(()) +} + fn validate_profile(args: &Args) -> Result<()> { + validate_query_forwarding_configuration(args)?; if args.profile != RuntimeProfile::Asapquery { if args.streaming_config.is_none() { return Err("the distributed profile requires --streaming-config".into()); @@ -762,9 +800,18 @@ async fn main() -> Result<()> { // query engine so SeriesLookup classification drives the Phase 6 archive // failover via EngineError::CapabilityMiss when the ASAP tier is empty / // ghost / unknown. + let query_forwarding_policy = if args.disable_query_forwarding { + data_plane::query_engines::QueryForwardingPolicy::Disabled + } else { + data_plane::query_engines::QueryForwardingPolicy::Enabled + }; + if !query_forwarding_policy.allows_external_queries() { + info!("query forwarding disabled for this process"); + } let engine = ASAPQueryEngine::new(args.prometheus_scrape_interval) .with_sketch_index(summary_store.clone()) .with_active_physical_plan(active_physical_plan.clone()) + .with_query_forwarding_policy(query_forwarding_policy) .with_exact_subquery_endpoint(args.prometheus_server.clone()) .with_metricsql_exact_subquery_endpoint(args.victoriametrics_url.clone()); @@ -1018,7 +1065,8 @@ async fn main() -> Result<()> { let adapter_config = AdapterConfig::prometheus_promql( args.prometheus_server.clone(), args.forward_unsupported_queries, - ); + ) + .with_query_forwarding_policy(query_forwarding_policy); let http_config = HttpServerConfig { port: args.http_port, @@ -1436,7 +1484,7 @@ fn setup_logging( #[cfg(test)] mod tests { - use super::{validate_profile, Args}; + use super::{validate_profile, validate_query_forwarding_configuration, Args}; use clap::Parser; use data_plane::drivers::AdapterConfig; @@ -1520,4 +1568,72 @@ mod tests { "forward_unsupported=true must install the Prom fallback", ); } + + #[test] + fn disable_query_forwarding_rejects_conflicting_flags() { + let args = Args::try_parse_from([ + "data_plane", + "--streaming-config", + "streaming.yaml", + "--disable-query-forwarding", + "--forward-unsupported-queries", + ]) + .unwrap(); + assert!(validate_profile(&args) + .unwrap_err() + .to_string() + .contains("conflicts")); + } + + #[test] + fn disable_query_forwarding_rejects_forwarding_listeners() { + let args = Args::try_parse_from([ + "data_plane", + "--streaming-config", + "streaming.yaml", + "--disable-query-forwarding", + "--victoriametrics-http-port", + "8429", + ]) + .unwrap(); + assert!(validate_profile(&args) + .unwrap_err() + .to_string() + .contains("victoriametrics-http-port")); + } + + #[test] + fn disable_query_forwarding_rejects_clickhouse_listener() { + let clickhouse = Args::try_parse_from([ + "data_plane", + "--streaming-config", + "streaming.yaml", + "--disable-query-forwarding", + "--clickhouse-http-port", + "8124", + ]) + .unwrap(); + assert!(validate_profile(&clickhouse) + .unwrap_err() + .to_string() + .contains("clickhouse-http-port")); + + } + + #[test] + fn disable_query_forwarding_rejects_asapquery_profile() { + let args = Args::try_parse_from([ + "data_plane", + "--profile", + "asapquery", + "--physical-plan", + "plan.json", + "--disable-query-forwarding", + ]) + .unwrap(); + assert!(validate_profile(&args) + .unwrap_err() + .to_string() + .contains("requires query forwarding")); + } } diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index c7fe6c8be..37f0fde60 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -1,4 +1,5 @@ use std::sync::Arc; +use tracing::debug; use asap_types::query_requirements::QueryRequirements; use asap_types::KeyByLabelNames; @@ -89,6 +90,7 @@ pub struct ASAPQueryEngine { active_physical_plan: Option, exact_subquery_endpoint: Option, metricsql_exact_subquery_endpoint: Option, + query_forwarding_policy: crate::query_engines::QueryForwardingPolicy, exact_subquery_client: reqwest::Client, } @@ -155,6 +157,7 @@ impl ASAPQueryEngine { active_physical_plan: None, exact_subquery_endpoint: None, metricsql_exact_subquery_endpoint: None, + query_forwarding_policy: crate::query_engines::QueryForwardingPolicy::Enabled, exact_subquery_client: reqwest::Client::builder() .timeout(std::time::Duration::from_secs(60)) .build() @@ -171,6 +174,14 @@ impl ASAPQueryEngine { self.metricsql_exact_subquery_endpoint = Some(endpoint); self } + + pub fn with_query_forwarding_policy( + mut self, + policy: crate::query_engines::QueryForwardingPolicy, + ) -> Self { + self.query_forwarding_policy = policy; + self + } async fn prepare_query_inputs( &self, physical: &crate::storage_engines::types::RuntimePhysicalPlan, @@ -225,11 +236,33 @@ impl ASAPQueryEngine { }, ); } + if !self.query_forwarding_policy.allows_external_queries() + && entry.nodes.values().any(|node| { + matches!( + node, + asap_types::query_plan::QueryPlanNode::ExternalExact { .. } + ) + }) + { + debug!( + language = ?entry.language, + query_id = %entry.query_id, + "query forwarding disabled; external exact subquery blocked" + ); + } super::exact_subqueries::prepare_external( entry, times, - self.exact_subquery_endpoint.as_deref(), - self.metricsql_exact_subquery_endpoint.as_deref(), + if self.query_forwarding_policy.allows_external_queries() { + self.exact_subquery_endpoint.as_deref() + } else { + None + }, + if self.query_forwarding_policy.allows_external_queries() { + self.metricsql_exact_subquery_endpoint.as_deref() + } else { + None + }, &self.exact_subquery_client, prepared, ) diff --git a/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs b/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs index 1630051e2..2a451d6a2 100644 --- a/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs +++ b/data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs @@ -631,6 +631,24 @@ mod tests { assert_eq!((exact.remote_evaluations, exact.remote_rpcs), (0, 0)); } + #[tokio::test] + async fn missing_external_endpoint_fails_closed_before_any_rpc() { + let result = prepare_external( + &candidate_entry("sum by (job) (rate(m[5m]))"), + &[1_000], + None, + None, + &reqwest::Client::new(), + candidate_rows(["api".into()]), + ) + .await; + assert!(matches!( + result, + Err(EngineError::CapabilityMiss { ref detail, .. }) + if detail.contains("Prometheus exact endpoint unavailable") + )); + } + #[tokio::test] async fn candidate_exact_is_discovered_and_prepared_behind_candidate_topk_root() { use asap_types::query_plan::{logical::Grouping, CandidateCompleteness}; diff --git a/data_plane/src/query_engines/mod.rs b/data_plane/src/query_engines/mod.rs index 7e6bf6fcc..ece33e1d8 100644 --- a/data_plane/src/query_engines/mod.rs +++ b/data_plane/src/query_engines/mod.rs @@ -17,6 +17,7 @@ pub mod asap_clickhouse_query_engine; pub mod asap_query_engine; pub mod canonical; +pub mod query_forwarding; pub mod query_result; pub mod routing; @@ -30,6 +31,7 @@ pub use crate::storage_engines::sketch_db::query::timeline_dispatch; pub use crate::storage_engines::sketch_db::query::window_merger; pub use asap_query_engine::ASAPQueryEngine; +pub use query_forwarding::{QueryForwardingError, QueryForwardingPolicy}; pub use query_result::{InstantVector, QueryResult, RangeVector, RangeVectorElement, Sample}; pub use timeline_dispatch::{combine_statistic, CombinedResult}; pub use window_merger::{create_window_merger, NaiveMerger, WindowMerger}; diff --git a/data_plane/src/query_engines/query_forwarding.rs b/data_plane/src/query_engines/query_forwarding.rs new file mode 100644 index 000000000..25512b748 --- /dev/null +++ b/data_plane/src/query_engines/query_forwarding.rs @@ -0,0 +1,25 @@ +use thiserror::Error; + +/// Controls whether query-serving code may contact an external backend. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum QueryForwardingPolicy { + #[default] + Enabled, + Disabled, +} + +impl QueryForwardingPolicy { + pub const fn allows_external_queries(self) -> bool { + matches!(self, Self::Enabled) + } +} + +/// Typed error for an external query blocked by the test policy. +#[derive(Debug, Error, Clone, PartialEq, Eq)] +pub enum QueryForwardingError { + #[error("external query forwarding is disabled for backend {backend} ({path})")] + Disabled { + backend: &'static str, + path: &'static str, + }, +} diff --git a/data_plane/src/tests/prometheus_forwarding_tests.rs b/data_plane/src/tests/prometheus_forwarding_tests.rs index e35c227fc..ff205f4e0 100644 --- a/data_plane/src/tests/prometheus_forwarding_tests.rs +++ b/data_plane/src/tests/prometheus_forwarding_tests.rs @@ -1,11 +1,14 @@ use crate::drivers::query::adapters::AdapterConfig; use crate::drivers::query::servers::http::{HttpServer, HttpServerConfig}; -use crate::query_engines::ASAPQueryEngine; +use crate::query_engines::{ASAPQueryEngine, QueryForwardingPolicy}; #[cfg(test)] use crate::storage_engines::types::{QueryLanguage, StreamingConfig}; use reqwest::Client; use serde_json::Value; -use std::sync::Arc; +use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; use tokio::net::TcpListener; use tokio::time::{sleep, Duration}; @@ -63,6 +66,40 @@ async fn start_mock_prometheus_server() -> Result, +) -> Result> { + use axum::{routing::get, Router}; + let app = Router::new() + .route( + "/api/v1/query", + get({ + let calls = calls.clone(); + move || async move { + calls.fetch_add(1, Ordering::SeqCst); + axum::Json(serde_json::json!({"status":"success","data":{"resultType":"vector","result":[]}})) + } + }), + ) + .route( + "/api/v1/query_range", + get({ + let calls = calls.clone(); + move || async move { + calls.fetch_add(1, Ordering::SeqCst); + axum::Json(serde_json::json!({"status":"success","data":{"resultType":"matrix","result":[]}})) + } + }), + ); + let listener = TcpListener::bind("127.0.0.1:0").await?; + let port = listener.local_addr()?.port(); + tokio::spawn(async move { + axum::serve(listener, app).await.unwrap(); + }); + sleep(Duration::from_millis(100)).await; + Ok(port) +} + async fn setup_test_server(prometheus_port: u16) -> (HttpServer, u16) { let config = HttpServerConfig { port: 0, // Use random port @@ -185,6 +222,50 @@ async fn test_forwarding_disabled() { assert_eq!(response_json["status"], "error"); } +#[tokio::test] +async fn disable_query_forwarding_makes_zero_instant_or_range_requests() { + let calls = Arc::new(AtomicUsize::new(0)); + let prometheus_port = start_counting_prometheus_server(calls.clone()) + .await + .unwrap(); + let config = HttpServerConfig { + port: 0, + handle_http_requests: true, + adapter_config: AdapterConfig::prometheus_promql( + format!("http://127.0.0.1:{prometheus_port}"), + true, + ) + .with_query_forwarding_policy(QueryForwardingPolicy::Disabled), + }; + let query_engine = Arc::new(ASAPQueryEngine::new(15000)); + let idx = Arc::new(crate::storage_engines::sketch_db::index::SketchStore::new()); + let server = HttpServer::new(config, query_engine, idx); + let server_port = server.start_test_server().await.unwrap(); + let client = Client::new(); + + let instant = client + .get(format!("http://127.0.0.1:{server_port}/api/v1/query")) + .query(&[("query", "rate(unsupported_metric[5m])")]) + .send() + .await + .unwrap(); + assert!(instant.status().is_success()); + + let range = client + .get(format!("http://127.0.0.1:{server_port}/api/v1/query_range")) + .query(&[ + ("query", "rate(unsupported_metric[5m])"), + ("start", "100"), + ("end", "200"), + ("step", "15"), + ]) + .send() + .await + .unwrap(); + assert!(range.status().is_success() || range.status().is_client_error()); + assert_eq!(calls.load(Ordering::SeqCst), 0); +} + #[tokio::test] async fn test_prometheus_server_unreachable() { // Use an unreachable port for Prometheus diff --git a/docs/design_docs/disable-query-forwarding-decisions.md b/docs/design_docs/disable-query-forwarding-decisions.md new file mode 100644 index 000000000..65c83c6c8 --- /dev/null +++ b/docs/design_docs/disable-query-forwarding-decisions.md @@ -0,0 +1,25 @@ +# Disable query forwarding: decisions + +Q: What does the mode prohibit? +A: No Prometheus, VictoriaMetrics, ClickHouse, or exact-subquery query requests. Health and runtime metadata remain allowed. + +Q: How is it enabled? +A: CLI-only `--disable-query-forwarding`. + +Q: What conflicts are errors? +A: Combining it with `--forward-unsupported-queries`, `--profile asapquery`, `--victoriametrics-http-port`, or `--clickhouse-http-port` is a startup error. + +Q: What happens when a plan needs an external exact subquery? +A: Fail closed with the normal protocol-specific unsupported/capability-miss response; never return an empty or partial result. + +Q: How is the behavior represented in code? +A: A centralized forwarding policy with typed errors/enums, not scattered booleans or ad-hoc error strings. + +Q: What is the default? +A: No flag means existing forwarding behavior is unchanged. + +Q: How is it tested? +A: Request-capturing mocks cover instant and range paths, proving zero external query requests. + +Q: What observability is required? +A: Startup INFO announces the mode; blocked attempts emit DEBUG logs and increment a backend/path-labeled counter. Existing protocol response shapes remain unchanged. From 460629d4403d07cf6ee2232b249a80cae22f5213 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 21 Sep 2026 16:02:00 +0000 Subject: [PATCH 2/4] test(query): verify no-forwarding mode in production process --- data_plane/src/drivers/query/servers/http.rs | 1 - data_plane/src/main.rs | 3 +- .../query_engines/asap_query_engine/engine.rs | 7 + data_plane/src/query_engines/mod.rs | 2 +- .../src/query_engines/query_forwarding.rs | 12 -- .../disable_query_forwarding_process_e2e.rs | 132 ++++++++++++++++++ .../disable-query-forwarding-decisions.md | 6 +- 7 files changed, 144 insertions(+), 19 deletions(-) create mode 100644 data_plane/tests/disable_query_forwarding_process_e2e.rs diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index fd3cb978c..f62a9f092 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2264,7 +2264,6 @@ async fn process_range_query_request( )})), ) .into_response() - } Err(EngineRouterError::AllFailed { last }) => { use crate::query_engines::EngineError; diff --git a/data_plane/src/main.rs b/data_plane/src/main.rs index 28110fa6a..91c109e81 100644 --- a/data_plane/src/main.rs +++ b/data_plane/src/main.rs @@ -1484,7 +1484,7 @@ fn setup_logging( #[cfg(test)] mod tests { - use super::{validate_profile, validate_query_forwarding_configuration, Args}; + use super::{validate_profile, Args}; use clap::Parser; use data_plane::drivers::AdapterConfig; @@ -1617,7 +1617,6 @@ mod tests { .unwrap_err() .to_string() .contains("clickhouse-http-port")); - } #[test] diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index 37f0fde60..3de7bceb5 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -249,6 +249,13 @@ impl ASAPQueryEngine { query_id = %entry.query_id, "query forwarding disabled; external exact subquery blocked" ); + crate::drivers::query::servers::metrics::record_query_forwarding_blocked( + match entry.language { + asap_types::QueryLanguage::MetricsQl => "victoriametrics", + _ => "prometheus", + }, + "exact_subquery", + ); } super::exact_subqueries::prepare_external( entry, diff --git a/data_plane/src/query_engines/mod.rs b/data_plane/src/query_engines/mod.rs index ece33e1d8..0470cb7db 100644 --- a/data_plane/src/query_engines/mod.rs +++ b/data_plane/src/query_engines/mod.rs @@ -31,7 +31,7 @@ pub use crate::storage_engines::sketch_db::query::timeline_dispatch; pub use crate::storage_engines::sketch_db::query::window_merger; pub use asap_query_engine::ASAPQueryEngine; -pub use query_forwarding::{QueryForwardingError, QueryForwardingPolicy}; +pub use query_forwarding::QueryForwardingPolicy; pub use query_result::{InstantVector, QueryResult, RangeVector, RangeVectorElement, Sample}; pub use timeline_dispatch::{combine_statistic, CombinedResult}; pub use window_merger::{create_window_merger, NaiveMerger, WindowMerger}; diff --git a/data_plane/src/query_engines/query_forwarding.rs b/data_plane/src/query_engines/query_forwarding.rs index 25512b748..a1eb2d24c 100644 --- a/data_plane/src/query_engines/query_forwarding.rs +++ b/data_plane/src/query_engines/query_forwarding.rs @@ -1,5 +1,3 @@ -use thiserror::Error; - /// Controls whether query-serving code may contact an external backend. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub enum QueryForwardingPolicy { @@ -13,13 +11,3 @@ impl QueryForwardingPolicy { matches!(self, Self::Enabled) } } - -/// Typed error for an external query blocked by the test policy. -#[derive(Debug, Error, Clone, PartialEq, Eq)] -pub enum QueryForwardingError { - #[error("external query forwarding is disabled for backend {backend} ({path})")] - Disabled { - backend: &'static str, - path: &'static str, - }, -} diff --git a/data_plane/tests/disable_query_forwarding_process_e2e.rs b/data_plane/tests/disable_query_forwarding_process_e2e.rs new file mode 100644 index 000000000..ca05cd93a --- /dev/null +++ b/data_plane/tests/disable_query_forwarding_process_e2e.rs @@ -0,0 +1,132 @@ +//! The production CLI's no-forwarding mode keeps query traffic inside the backend. + +use std::io::Write; +use std::net::TcpListener; +use std::process::{Child, Command, Stdio}; +use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, +}; +use std::time::Duration; + +use axum::{routing::get, Router}; + +struct ChildGuard(Child); + +impl Drop for ChildGuard { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +fn unused_port() -> u16 { + TcpListener::bind("127.0.0.1:0") + .unwrap() + .local_addr() + .unwrap() + .port() +} + +#[tokio::test] +async fn cli_mode_blocks_instant_and_range_forwarding() { + let calls = Arc::new(AtomicUsize::new(0)); + let count = calls.clone(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let upstream = format!("http://{}", listener.local_addr().unwrap()); + let capture = tokio::spawn(async move { + axum::serve( + listener, + Router::new() + .route( + "/api/v1/query", + get({ + let count = count.clone(); + move || async move { + count.fetch_add(1, Ordering::SeqCst); + "unexpected instant query" + } + }), + ) + .route( + "/api/v1/query_range", + get(move || async move { + count.fetch_add(1, Ordering::SeqCst); + "unexpected range query" + }), + ), + ) + .await + .unwrap(); + }); + + let mut config = tempfile::NamedTempFile::new().unwrap(); + write!(config, "aggregations: []\n").unwrap(); + let output = tempfile::tempdir().unwrap(); + let port = unused_port(); + let mut child = ChildGuard( + Command::new(env!("CARGO_BIN_EXE_data_plane")) + .arg("--streaming-config") + .arg(config.path()) + .args([ + "--disable-query-forwarding", + "--prometheus-server", + &upstream, + ]) + .args(["--http-port", &port.to_string()]) + .arg("--output-dir") + .arg(output.path()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .unwrap(), + ); + let client = reqwest::Client::new(); + let base = format!("http://127.0.0.1:{port}"); + let mut ready = false; + for _ in 0..100 { + if let Ok(response) = client.get(format!("{base}/api/v1/health")).send().await { + if response.status().is_success() { + ready = true; + break; + } + } + assert!( + child.0.try_wait().unwrap().is_none(), + "backend exited before ready" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } + assert!(ready, "backend did not become ready"); + + let instant = client + .get(format!("{base}/api/v1/query")) + .query(&[("query", "rate(unplanned_metric[5m])")]) + .send() + .await + .unwrap(); + assert_eq!(instant.status(), reqwest::StatusCode::OK); + assert_eq!( + instant.json::().await.unwrap()["status"], + "error" + ); + + let range = client + .get(format!("{base}/api/v1/query_range")) + .query(&[ + ("query", "rate(unplanned_metric[5m])"), + ("start", "100"), + ("end", "200"), + ("step", "15"), + ]) + .send() + .await + .unwrap(); + assert_eq!(range.status(), reqwest::StatusCode::OK); + assert_eq!( + range.json::().await.unwrap()["status"], + "error" + ); + assert_eq!(calls.load(Ordering::SeqCst), 0); + capture.abort(); +} diff --git a/docs/design_docs/disable-query-forwarding-decisions.md b/docs/design_docs/disable-query-forwarding-decisions.md index 65c83c6c8..5faf2493a 100644 --- a/docs/design_docs/disable-query-forwarding-decisions.md +++ b/docs/design_docs/disable-query-forwarding-decisions.md @@ -13,13 +13,13 @@ Q: What happens when a plan needs an external exact subquery? A: Fail closed with the normal protocol-specific unsupported/capability-miss response; never return an empty or partial result. Q: How is the behavior represented in code? -A: A centralized forwarding policy with typed errors/enums, not scattered booleans or ad-hoc error strings. +A: A shared forwarding-policy enum gates adapter fallback and planned exact subqueries. Q: What is the default? A: No flag means existing forwarding behavior is unchanged. Q: How is it tested? -A: Request-capturing mocks cover instant and range paths, proving zero external query requests. +A: A production-process request-capture test covers instant and range paths; CLI validation tests cover conflicting settings. Q: What observability is required? -A: Startup INFO announces the mode; blocked attempts emit DEBUG logs and increment a backend/path-labeled counter. Existing protocol response shapes remain unchanged. +A: Startup INFO announces the mode; blocked HTTP fallbacks and planned exact subqueries emit DEBUG logs and increment a backend/path-labeled counter. Existing protocol response shapes remain unchanged. From 0c137a250e5321b302f95b20c1dbcb982d9abc0d Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 21 Sep 2026 16:56:27 +0000 Subject: [PATCH 3/4] test(query): satisfy Clippy in forwarding process test --- data_plane/tests/disable_query_forwarding_process_e2e.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/data_plane/tests/disable_query_forwarding_process_e2e.rs b/data_plane/tests/disable_query_forwarding_process_e2e.rs index ca05cd93a..7c8c65855 100644 --- a/data_plane/tests/disable_query_forwarding_process_e2e.rs +++ b/data_plane/tests/disable_query_forwarding_process_e2e.rs @@ -61,7 +61,7 @@ async fn cli_mode_blocks_instant_and_range_forwarding() { }); let mut config = tempfile::NamedTempFile::new().unwrap(); - write!(config, "aggregations: []\n").unwrap(); + writeln!(config, "aggregations: []").unwrap(); let output = tempfile::tempdir().unwrap(); let port = unused_port(); let mut child = ChildGuard( From 97d872afe163d6b5bb4462e2393ab70b21753a94 Mon Sep 17 00:00:00 2001 From: zz_y Date: Mon, 21 Sep 2026 17:17:04 +0000 Subject: [PATCH 4/4] fix(query): count only configured blocked fallbacks --- .../src/drivers/query/adapters/config.rs | 19 ++++ data_plane/src/drivers/query/servers/http.rs | 11 +-- .../query_engines/asap_query_engine/engine.rs | 89 +++++++++++++++++++ 3 files changed, 114 insertions(+), 5 deletions(-) diff --git a/data_plane/src/drivers/query/adapters/config.rs b/data_plane/src/drivers/query/adapters/config.rs index 044b3ef72..36cff467c 100644 --- a/data_plane/src/drivers/query/adapters/config.rs +++ b/data_plane/src/drivers/query/adapters/config.rs @@ -17,6 +17,9 @@ pub struct AdapterConfig { /// Whether this adapter may issue query requests to its fallback backend. pub query_forwarding_policy: QueryForwardingPolicy, + + /// A fallback was configured but removed by the forwarding policy. + pub fallback_blocked_by_policy: bool, } impl std::fmt::Debug for AdapterConfig { @@ -29,6 +32,10 @@ impl std::fmt::Debug for AdapterConfig { &self.fallback.as_ref().map(|_| "Some(FallbackClient)"), ) .field("query_forwarding_policy", &self.query_forwarding_policy) + .field( + "fallback_blocked_by_policy", + &self.fallback_blocked_by_policy, + ) .finish() } } @@ -45,13 +52,17 @@ impl AdapterConfig { language, fallback, query_forwarding_policy: QueryForwardingPolicy::Enabled, + fallback_blocked_by_policy: false, } } pub fn with_query_forwarding_policy(mut self, policy: QueryForwardingPolicy) -> Self { self.query_forwarding_policy = policy; if !policy.allows_external_queries() { + self.fallback_blocked_by_policy |= self.fallback.is_some(); self.fallback = None; + } else { + self.fallback_blocked_by_policy = false; } self } @@ -108,9 +119,17 @@ mod tests { let config = AdapterConfig::prometheus_promql("http://prom:9090".into(), true) .with_query_forwarding_policy(QueryForwardingPolicy::Disabled); assert!(config.fallback.is_none()); + assert!(config.fallback_blocked_by_policy); assert_eq!( config.query_forwarding_policy, QueryForwardingPolicy::Disabled ); } + + #[test] + fn disabling_without_a_fallback_does_not_claim_a_blocked_request() { + let config = AdapterConfig::prometheus_promql("http://prom:9090".into(), false) + .with_query_forwarding_policy(QueryForwardingPolicy::Disabled); + assert!(!config.fallback_blocked_by_policy); + } } diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index f62a9f092..382ec0b43 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -2036,11 +2036,12 @@ fn query_status_label(response: &Response) -> &'static str { } fn record_disabled_forwarding(state: &AppState, backend: &str, path: &str) { - if !state - .config - .adapter_config - .query_forwarding_policy - .allows_external_queries() + if state.config.adapter_config.fallback_blocked_by_policy + && !state + .config + .adapter_config + .query_forwarding_policy + .allows_external_queries() { debug!( backend, diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index 3de7bceb5..cae8ff229 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -64,6 +64,95 @@ mod readiness_coverage_tests { } } +#[cfg(test)] +mod forwarding_policy_tests { + use super::ASAPQueryEngine; + use crate::query_engines::asap_query_engine::test_plan; + use crate::query_engines::routing::query_engine_routing::QueryEngine; + use crate::query_engines::{EngineError, QueryForwardingPolicy}; + use crate::storage_engines::sketch_db::index::SketchStore; + use asap_types::query_plan::{ + ExternalExactOutput, ExternalExactRequest, FallbackPolicy, InstantExecution, QueryLanguage, + QueryNodeId, QueryPlanEntry, QueryPlanNode, + }; + use std::collections::BTreeMap; + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }; + + #[tokio::test] + async fn disabled_policy_blocks_an_installed_exact_plan_before_http() { + let calls = Arc::new(AtomicUsize::new(0)); + let captured = calls.clone(); + let app = axum::Router::new().route( + "/api/v1/query", + axum::routing::get(move || { + let captured = captured.clone(); + async move { + captured.fetch_add(1, Ordering::SeqCst); + axum::Json(serde_json::json!({ + "status":"success", + "data":{"resultType":"vector","result":[]} + })) + } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let endpoint = format!("http://{}", listener.local_addr().unwrap()); + let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let query = "sum(rate(m[5m]))"; + let canonical = asap_types::query_plan::canonical_promql(query).unwrap(); + let entry = QueryPlanEntry { + language: QueryLanguage::PromQl, + query_id: canonical.clone(), + canonical_query: canonical, + fixed_evaluation: None, + root: QueryNodeId(0), + nodes: BTreeMap::from([( + QueryNodeId(0), + QueryPlanNode::ExternalExact { + request: ExternalExactRequest { + language: QueryLanguage::PromQl, + expression: query.into(), + output: ExternalExactOutput::InstantVector, + parameters: BTreeMap::new(), + start_parameter: None, + end_parameter: None, + input_contracts: vec![], + }, + inputs: vec![], + }, + )]), + instant: InstantExecution { + lookback_ms: 300_000, + full_history: false, + cumulative_readout: true, + }, + fallback: FallbackPolicy::ExactBackend, + }; + let index = Arc::new(SketchStore::new()); + let active = test_plan::install(&index, &[], vec![entry]); + let enabled = ASAPQueryEngine::new(15_000) + .with_active_physical_plan(active.clone()) + .with_exact_subquery_endpoint(endpoint.clone()); + enabled.execute_at(query, 1_000).await.unwrap(); + assert_eq!(calls.load(Ordering::SeqCst), 1); + let engine = ASAPQueryEngine::new(15_000) + .with_active_physical_plan(active) + .with_exact_subquery_endpoint(endpoint) + .with_query_forwarding_policy(QueryForwardingPolicy::Disabled); + let result = engine.execute_at(query, 1_000).await; + assert!( + matches!(result, Err(EngineError::CapabilityMiss { .. })), + "{result:?}" + ); + assert_eq!(calls.load(Ordering::SeqCst), 1); + server.abort(); + } +} + #[cfg(test)] use crate::storage_engines::types::KeyByLabelValues; #[cfg(test)]