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
48 changes: 42 additions & 6 deletions asap-query-engine/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,10 @@ use std::collections::HashMap;
use std::sync::Arc;
use std::time::Instant;
use tokio::net::TcpListener;
use tracing::{debug, info};
use tracing::{debug, info, warn};

use crate::drivers::query::adapters::{create_http_adapter, AdapterConfig, HttpProtocolAdapter};
use crate::engines::SimpleEngine;
use crate::engines::{QueryExecutionError, SimpleEngine};
use crate::query_tracker::QueryTracker;
use crate::stores::Store;

Expand Down Expand Up @@ -44,6 +44,22 @@ struct AppState {
fallback: Option<Arc<dyn crate::drivers::query::fallback::FallbackClient>>,
}

async fn format_native_execution_error(state: &AppState, error: QueryExecutionError) -> Response {
match state
.adapter
.format_error_response(
&crate::drivers::query::adapters::AdapterError::ProtocolError(error.to_string()),
)
.await
{
Ok(mut response) => {
*response.status_mut() = StatusCode::INTERNAL_SERVER_ERROR;
response
}
Err(status) => status.into_response(),
}
}

impl HttpServer {
pub fn new(
config: HttpServerConfig,
Expand Down Expand Up @@ -193,7 +209,7 @@ async fn process_query_request(
.query_engine
.handle_query(parsed_request.query.clone(), parsed_request.time)
{
Some((query_output_labels, query_result)) => {
Ok(Some((query_output_labels, query_result))) => {
let query_duration = query_start_time.elapsed();
debug!("=== QUERY ENGINE SUCCESS ===");
debug!(
Expand Down Expand Up @@ -233,7 +249,7 @@ async fn process_query_request(
Err(status) => status.into_response(),
}
}
None => {
Ok(None) => {
let total_duration = start_time.elapsed();
debug!("=== QUERY ENGINE RETURNED NONE ===");
debug!(
Expand Down Expand Up @@ -271,6 +287,16 @@ async fn process_query_request(
}
}
}
Err(error) => {
let total_duration = start_time.elapsed();
warn!(query = %parsed_request.query, error = %error, "Native query execution failed");
info!(
"query='{}' destination=none_native_error total_latency_ms={:.2}",
parsed_request.query,
total_duration.as_secs_f64() * 1000.0
);
format_native_execution_error(state, error).await
}
}
}

Expand Down Expand Up @@ -524,7 +550,7 @@ async fn process_range_query_request(
parsed_request.end,
parsed_request.step,
) {
Some((query_output_labels, query_result)) => {
Ok(Some((query_output_labels, query_result))) => {
let query_duration = query_start_time.elapsed();
debug!(
"Range query execution took: {:.2}ms",
Expand All @@ -547,7 +573,7 @@ async fn process_range_query_request(
Err(status) => status.into_response(),
}
}
None => {
Ok(None) => {
let total_duration = start_time.elapsed();
debug!("Range query returned None - query not supported");

Expand Down Expand Up @@ -578,6 +604,16 @@ async fn process_range_query_request(
}
}
}
Err(error) => {
let total_duration = start_time.elapsed();
warn!(query = %parsed_request.query, error = %error, "Native range query execution failed");
info!(
"query='{}' destination=none_native_error total_latency_ms={:.2}",
parsed_request.query,
total_duration.as_secs_f64() * 1000.0
);
format_native_execution_error(state, error).await
}
}
}

Expand Down
2 changes: 1 addition & 1 deletion asap-query-engine/src/engines/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,5 +6,5 @@ pub(crate) mod sliding_window_composition;
pub mod window_merger;

pub use query_result::{InstantVector, QueryResult, RangeVector, RangeVectorElement, Sample};
pub use simple_engine::SimpleEngine;
pub use simple_engine::{QueryExecutionError, SimpleEngine};
pub use window_merger::{create_window_merger, NaiveMerger, WindowMerger};
Loading
Loading