Skip to content
Merged
Show file tree
Hide file tree
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
44 changes: 44 additions & 0 deletions data_plane/src/drivers/query/adapters/config.rs
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -13,6 +14,12 @@ pub struct AdapterConfig {

/// Optional fallback client for unsupported queries
pub fallback: Option<Arc<dyn FallbackClient>>,

/// 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 {
Expand All @@ -24,6 +31,11 @@ impl std::fmt::Debug for AdapterConfig {
"fallback",
&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()
}
}
Expand All @@ -39,7 +51,20 @@ impl AdapterConfig {
protocol,
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
}

/// Create a configuration for Prometheus HTTP with PromQL
Expand Down Expand Up @@ -88,4 +113,23 @@ 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!(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);
}
}
18 changes: 18 additions & 0 deletions data_plane/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -2034,6 +2035,22 @@ fn query_status_label(response: &Response) -> &'static str {
}
}

fn record_disabled_forwarding(state: &AppState, backend: &str, path: &str) {
if state.config.adapter_config.fallback_blocked_by_policy
&& !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
// ============================================================
Expand Down Expand Up @@ -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(),
Expand Down
14 changes: 14 additions & 0 deletions data_plane/src/drivers/query/servers/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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);
Expand All @@ -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();
}
117 changes: 116 additions & 1 deletion data_plane/src/main.rs
Original file line number Diff line number Diff line change
@@ -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};

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -382,7 +387,40 @@ struct Args {
backend_storage_routing: Option<std::path::PathBuf>,
}

#[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());
Expand Down Expand Up @@ -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());

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -1520,4 +1568,71 @@ 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"));
}
}
Loading
Loading