From 8af3a04ed7ec225e9b3fa41bda8fd6459f77179f Mon Sep 17 00:00:00 2001 From: discord9 Date: Wed, 16 Sep 2026 06:28:15 +0000 Subject: [PATCH] fix: preserve structured query errors through distributed execution (#9161) * fix: preserve structured query errors through distributed execution Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> * test: update SQL expectations for preserved query error codes Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --------- Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- src/common/query/src/error.rs | 148 +++++++- src/common/recordbatch/src/error.rs | 124 ++++++- src/flow/src/batching_mode/task/ckpt.rs | 45 --- src/flow/src/batching_mode/task/test.rs | 214 ++--------- src/query/src/datafusion.rs | 347 +++++++++++++++++- src/query/src/dist_plan/merge_scan.rs | 117 ++++-- src/query/src/error.rs | 15 + src/servers/src/error.rs | 36 +- tests-integration/src/grpc/flight.rs | 8 +- .../optimizer/filter_push_down.result | 2 +- .../standalone/common/order/limit.result | 18 +- .../optimizer/filter_push_down.result | 2 +- 12 files changed, 791 insertions(+), 285 deletions(-) diff --git a/src/common/query/src/error.rs b/src/common/query/src/error.rs index dd2b29adf3a..4c53be86b6e 100644 --- a/src/common/query/src/error.rs +++ b/src/common/query/src/error.rs @@ -270,18 +270,158 @@ pub fn datafusion_status_code( e: &DataFusionError, default_status: Option, ) -> StatusCode { - match e { + let mut error = e; + loop { + error = match error { + DataFusionError::Shared(inner) => inner, + DataFusionError::Context(_, inner) | DataFusionError::Diagnostic(_, inner) => inner, + _ => break, + }; + } + + match error { DataFusionError::Internal(_) => StatusCode::Internal, DataFusionError::NotImplemented(_) => StatusCode::Unsupported, DataFusionError::Plan(_) => StatusCode::PlanQuery, - DataFusionError::External(e) => { - if let Some(ext) = (*e).downcast_ref::() { + DataFusionError::External(error) => { + if let Some(ext) = (*error).downcast_ref::() { + ext.status_code() + } else if let Some(ext) = (*error).downcast_ref::() { ext.status_code() } else { default_status.unwrap_or(StatusCode::EngineExecuteQuery) } } - DataFusionError::Diagnostic(_, e) => datafusion_status_code::(e, default_status), _ => default_status.unwrap_or(StatusCode::EngineExecuteQuery), } } + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use common_error::ext::PlainError; + + use super::*; + + #[test] + fn test_datafusion_status_code_preserves_external_errors_through_wrappers() { + let boxed_error = |status| { + DataFusionError::External(Box::new(BoxedError::new(PlainError::new( + "neutral error".to_string(), + status, + )))) + }; + + for status in [ + StatusCode::RequestOutdated, + StatusCode::Unknown, + StatusCode::Unsupported, + ] { + let errors = [ + boxed_error(status), + DataFusionError::Shared(Arc::new(boxed_error(status))), + DataFusionError::Context("context".to_string(), Box::new(boxed_error(status))), + DataFusionError::Diagnostic( + Box::new(datafusion_common::Diagnostic::new_error("diagnostic", None)), + Box::new(boxed_error(status)), + ), + DataFusionError::Shared(Arc::new(DataFusionError::Context( + "context".to_string(), + Box::new(DataFusionError::Diagnostic( + Box::new(datafusion_common::Diagnostic::new_error("diagnostic", None)), + Box::new(boxed_error(status)), + )), + ))), + ]; + + for error in errors { + assert_eq!(datafusion_status_code::(&error, None), status); + assert_eq!( + datafusion_status_code::(&error, Some(StatusCode::PlanQuery)), + status + ); + } + } + + let direct_error = DataFusionError::External(Box::new(Error::DynFilterPayloadTooLarge { + payload_size_bytes: 2, + max_payload_bytes: 1, + location: Location::default(), + })); + assert_eq!( + datafusion_status_code::(&direct_error, Some(StatusCode::Internal)), + StatusCode::PlanQuery + ); + + let boundary_error = + DataFusionError::Shared(Arc::new(DataFusionError::External(Box::new( + BoxedError::new(common_recordbatch::error::Error::PhysicalExpr { + error: DataFusionError::NotImplemented("inner error".to_string()), + location: Location::default(), + }), + )))); + assert_eq!( + datafusion_status_code::(&boundary_error, Some(StatusCode::PlanQuery)), + StatusCode::Internal + ); + } + + #[test] + fn test_datafusion_status_code_uses_default_for_untyped_errors() { + let wrap = |error| { + DataFusionError::Shared(Arc::new(DataFusionError::Context( + "context".to_string(), + Box::new(DataFusionError::Diagnostic( + Box::new(datafusion_common::Diagnostic::new_error("diagnostic", None)), + Box::new(error), + )), + ))) + }; + let errors = || { + [ + ( + DataFusionError::External(Box::new(std::io::Error::other("neutral io error"))), + StatusCode::EngineExecuteQuery, + StatusCode::Unknown, + ), + ( + DataFusionError::Internal("neutral internal error".to_string()), + StatusCode::Internal, + StatusCode::Internal, + ), + ( + DataFusionError::NotImplemented("neutral not implemented error".to_string()), + StatusCode::Unsupported, + StatusCode::Unsupported, + ), + ( + DataFusionError::Plan("neutral plan error".to_string()), + StatusCode::PlanQuery, + StatusCode::PlanQuery, + ), + ( + DataFusionError::External(Box::new(DataFusionError::Internal( + "inner error".to_string(), + ))), + StatusCode::EngineExecuteQuery, + StatusCode::Unknown, + ), + ] + }; + + for (error, none_expected, default_expected) in + errors() + .into_iter() + .chain(errors().map(|(error, none_expected, default_expected)| { + (wrap(error), none_expected, default_expected) + })) + { + assert_eq!(datafusion_status_code::(&error, None), none_expected); + assert_eq!( + datafusion_status_code::(&error, Some(StatusCode::Unknown)), + default_expected + ); + } + } +} diff --git a/src/common/recordbatch/src/error.rs b/src/common/recordbatch/src/error.rs index 469320c19cf..bddb8ae8032 100644 --- a/src/common/recordbatch/src/error.rs +++ b/src/common/recordbatch/src/error.rs @@ -205,7 +205,26 @@ impl ErrorExt for Error { | Error::PhysicalExpr { .. } | Error::RecordBatchSliceIndexOverflow { .. } => StatusCode::Internal, - Error::PollStream { .. } => StatusCode::EngineExecuteQuery, + Error::PollStream { error, .. } => { + let mut error = error; + loop { + error = match error { + datafusion::error::DataFusionError::Shared(inner) => inner, + datafusion::error::DataFusionError::Context(_, inner) + | datafusion::error::DataFusionError::Diagnostic(_, inner) => inner, + _ => break, + }; + } + + match error { + datafusion::error::DataFusionError::External(source) => source + .downcast_ref::() + .map_or(StatusCode::EngineExecuteQuery, |source| { + source.status_code() + }), + _ => StatusCode::EngineExecuteQuery, + } + } Error::ArrowCompute { .. } => StatusCode::IllegalState, @@ -242,3 +261,106 @@ impl ErrorExt for Error { } } } + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use common_error::ext::PlainError; + use datafusion::error::DataFusionError; + + use super::*; + + #[test] + fn poll_stream_status_code_preserves_boxed_error_through_wrappers() { + let boxed_error = |status| { + DataFusionError::External(Box::new(BoxedError::new(PlainError::new( + "neutral error".to_string(), + status, + )))) + }; + + for status in [ + StatusCode::RequestOutdated, + StatusCode::Unknown, + StatusCode::Unsupported, + ] { + let errors = [ + boxed_error(status), + DataFusionError::Shared(Arc::new(boxed_error(status))), + DataFusionError::Context("context".to_string(), Box::new(boxed_error(status))), + DataFusionError::Diagnostic( + Box::new(datafusion::common::Diagnostic::new_error( + "diagnostic", + None, + )), + Box::new(boxed_error(status)), + ), + DataFusionError::Shared(Arc::new(DataFusionError::Context( + "context".to_string(), + Box::new(DataFusionError::Diagnostic( + Box::new(datafusion::common::Diagnostic::new_error( + "diagnostic", + None, + )), + Box::new(boxed_error(status)), + )), + ))), + ]; + + for error in errors { + let error = Error::PollStream { + error, + location: Location::default(), + }; + assert_eq!(error.status_code(), status); + } + } + + let error = Error::PollStream { + error: DataFusionError::Shared(Arc::new(DataFusionError::External(Box::new( + BoxedError::new(Error::PhysicalExpr { + error: DataFusionError::NotImplemented("inner error".to_string()), + location: Location::default(), + }), + )))), + location: Location::default(), + }; + assert_eq!(error.status_code(), StatusCode::Internal); + } + + #[test] + fn poll_stream_status_code_defaults_for_other_datafusion_errors() { + let wrap = |error| { + DataFusionError::Shared(Arc::new(DataFusionError::Context( + "context".to_string(), + Box::new(DataFusionError::Diagnostic( + Box::new(datafusion::common::Diagnostic::new_error( + "diagnostic", + None, + )), + Box::new(error), + )), + ))) + }; + let errors = || { + [ + DataFusionError::External(Box::new(std::io::Error::other("neutral io error"))), + DataFusionError::Internal("neutral internal error".to_string()), + DataFusionError::NotImplemented("neutral not implemented error".to_string()), + DataFusionError::Plan("neutral plan error".to_string()), + DataFusionError::External(Box::new(DataFusionError::Internal( + "inner error".to_string(), + ))), + ] + }; + + for error in errors().into_iter().chain(errors().map(wrap)) { + let error = Error::PollStream { + error, + location: Location::default(), + }; + assert_eq!(error.status_code(), StatusCode::EngineExecuteQuery); + } + } +} diff --git a/src/flow/src/batching_mode/task/ckpt.rs b/src/flow/src/batching_mode/task/ckpt.rs index 960e541c73c..c736a0fdbdf 100644 --- a/src/flow/src/batching_mode/task/ckpt.rs +++ b/src/flow/src/batching_mode/task/ckpt.rs @@ -12,7 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::error::Error as StdError; use std::time::Duration; use client::OutputWithMetrics; @@ -31,42 +30,6 @@ use crate::metrics::{ }; use crate::{Error, FlowId}; -/// Liveness guard: when a fenced repair query fails with a wrapped error whose -/// text indicates a stale snapshot fence (even when `StatusCode::RequestOutdated` -/// was lost through client layers), classify it as `SnapshotFenceExpired` to -/// break the retry loop and force a rebind of the fence high `H`. -/// -/// Long-term the structured `StatusCode` / retry hint path should be preserved -/// end-to-end; this text fallback is a narrow safety measure. -fn matches_stale_snapshot_fence_text(err: &Error) -> bool { - let markers = [ - "STALE_SNAPSHOT_FENCE", - "REBIND_SNAPSHOT_FENCE", - "snapshot upper bound stale", - ]; - // Check the top-level error Display and Debug. - let debug_str = format!("{:?}", err); - let display_str = err.to_string(); - for marker in &markers { - if debug_str.contains(marker) || display_str.contains(marker) { - return true; - } - } - // Walk the error source chain. - let mut source = err.source(); - while let Some(s) = source { - let debug_str = format!("{:?}", s); - let display_str = s.to_string(); - for marker in &markers { - if debug_str.contains(marker) || display_str.contains(marker) { - return true; - } - } - source = s.source(); - } - false -} - impl BatchingTask { /// Classify execution errors into checkpoint fallback reasons. A stale /// snapshot fence is special only for fenced repair chunks. @@ -80,14 +43,6 @@ impl BatchingTask { } else { FlowQueryFallbackReason::StaleCursor } - } else if matches!(coverage, QueryCoverage::FencedRepairChunk { .. }) - && matches_stale_snapshot_fence_text(err) - { - // Narrow text-based fallback for wrapped errors where the - // structured StatusCode::RequestOutdated was lost through - // frontend/client layers. Without this fenced repair will - // retry the same stale `given_seq` every refresh tick. - FlowQueryFallbackReason::SnapshotFenceExpired } else if matches!(coverage, QueryCoverage::IncrementalDelta) { FlowQueryFallbackReason::IncrementalQueryFailure } else { diff --git a/src/flow/src/batching_mode/task/test.rs b/src/flow/src/batching_mode/task/test.rs index 24cdbf883e0..acb872d8fea 100644 --- a/src/flow/src/batching_mode/task/test.rs +++ b/src/flow/src/batching_mode/task/test.rs @@ -13,6 +13,7 @@ // limitations under the License. use std::collections::{BTreeMap, BTreeSet, HashMap}; +use std::sync::Arc; use catalog::RegisterTableRequest; use catalog::memory::MemoryCatalogManager; @@ -391,50 +392,6 @@ fn flow_error_with_status(status_code: StatusCode) -> Error { .unwrap_err() } -/// Test-only error that carries a non-RequestOutdated status code but -/// displays a stale-snapshot-fence marker string, simulating the real-world -/// scenario where the structured status code is lost through frontend/client -/// wrapping layers. -#[derive(Debug)] -struct StaleFenceTextError { - code: StatusCode, - message: String, -} - -impl std::fmt::Display for StaleFenceTextError { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - write!(f, "{}", self.message) - } -} - -impl std::error::Error for StaleFenceTextError {} - -impl common_error::ext::ErrorExt for StaleFenceTextError { - fn status_code(&self) -> StatusCode { - self.code - } - fn as_any(&self) -> &dyn std::any::Any { - self - } -} - -impl common_error::ext::StackError for StaleFenceTextError { - fn debug_fmt(&self, _: usize, _: &mut Vec) {} - fn next(&self) -> Option<&dyn common_error::ext::StackError> { - None - } -} - -fn flow_error_with_code_and_text(code: StatusCode, text: &str) -> Error { - let inner = StaleFenceTextError { - code, - message: text.to_string(), - }; - Err::<(), _>(BoxedError::new(inner)) - .context(crate::error::ExternalSnafu) - .unwrap_err() -} - fn dirty_range(start: i64, end: i64) -> DirtyTimeWindows { let mut dirty = DirtyTimeWindows::default(); dirty.add_window( @@ -1447,78 +1404,6 @@ fn test_query_failure_reason_distinguishes_fenced_repair_stale_fence() { ); } -/// Wrapped errors carrying stale snapshot fence marker text in their -/// Display/Debug chain should be classified as `SnapshotFenceExpired` on -/// fenced repair coverage, even when the structured `StatusCode::RequestOutdated` -/// was lost through client layering. This prevents an infinite retry loop -/// where the fenced chunk re-sends the same stale `given_seq` every tick. -#[test] -fn test_query_failure_reason_text_fallback_stale_snapshot_fence() { - let high = BTreeMap::new(); - let fenced = QueryCoverage::FencedRepairChunk { high: high.clone() }; - - // STALE_SNAPSHOT_FENCE marker with a non-RequestOutdated status code - let err = flow_error_with_code_and_text( - StatusCode::Internal, - "gRPC error: STALE_SNAPSHOT_FENCE: snapshot upper bound stale, region: 1024/0", - ); - assert_eq!( - BatchingTask::query_failure_reason(&err, &fenced), - FlowQueryFallbackReason::SnapshotFenceExpired - ); - - // REBIND_SNAPSHOT_FENCE marker - let err = flow_error_with_code_and_text( - StatusCode::Internal, - "STALE_SNAPSHOT_FENCE ... retry_hint: REBIND_SNAPSHOT_FENCE", - ); - assert_eq!( - BatchingTask::query_failure_reason(&err, &fenced), - FlowQueryFallbackReason::SnapshotFenceExpired - ); - - // snapshot upper bound stale marker (the natural-language fragment) - let err = flow_error_with_code_and_text( - StatusCode::Internal, - "query failed: snapshot upper bound stale, consider rebinding", - ); - assert_eq!( - BatchingTask::query_failure_reason(&err, &fenced), - FlowQueryFallbackReason::SnapshotFenceExpired - ); - - // Fenced coverage with a generic wrapped error (no stale-fence marker) → - // still QueryFailure - let generic_err = - flow_error_with_code_and_text(StatusCode::Internal, "some transient network error"); - assert_eq!( - BatchingTask::query_failure_reason(&generic_err, &fenced), - FlowQueryFallbackReason::QueryFailure - ); - - // Non-fenced incremental coverage with stale-fence marker text must NOT - // classify as SnapshotFenceExpired; it should remain IncrementalQueryFailure. - let err = flow_error_with_code_and_text( - StatusCode::Internal, - "STALE_SNAPSHOT_FENCE blob in unexpected context", - ); - assert_eq!( - BatchingTask::query_failure_reason(&err, &QueryCoverage::IncrementalDelta), - FlowQueryFallbackReason::IncrementalQueryFailure - ); - - // Existing RequestOutdated behavior is unchanged. - let outdated_err = flow_error_with_status(StatusCode::RequestOutdated); - assert_eq!( - BatchingTask::query_failure_reason(&outdated_err, &fenced), - FlowQueryFallbackReason::SnapshotFenceExpired - ); - assert_eq!( - BatchingTask::query_failure_reason(&outdated_err, &QueryCoverage::IncrementalDelta), - FlowQueryFallbackReason::StaleCursor - ); -} - #[test] fn test_fenced_repair_coverage_produces_snapshot_seq_map_for_distributed_metadata_path() { // Covers the metadata boundary between QueryCoverage and the @@ -1563,11 +1448,28 @@ async fn test_fenced_repair_stale_fence_next_plan_is_scoped_base_repair() { { let mut state = task.state.write().unwrap(); + let error = Err::<(), _>(BoxedError::new( + common_recordbatch::error::Error::PollStream { + error: datafusion::error::DataFusionError::Shared(Arc::new( + datafusion::error::DataFusionError::External(Box::new(BoxedError::new( + MockError::new(StatusCode::RequestOutdated), + ))), + )), + location: snafu::Location::default(), + }, + )) + .context(crate::error::ExternalSnafu) + .unwrap_err(); + let reason = BatchingTask::query_failure_reason( + &error, + &QueryCoverage::FencedRepairChunk { high: high.clone() }, + ); + assert_eq!(reason, FlowQueryFallbackReason::SnapshotFenceExpired); let decision = BatchingTask::apply_query_failure_to_state( &mut state, std::time::Duration::from_millis(1), &QueryCoverage::FencedRepairChunk { high }, - FlowQueryFallbackReason::SnapshotFenceExpired, + reason, ); assert_eq!( decision, @@ -1577,9 +1479,11 @@ async fn test_fenced_repair_stale_fence_next_plan_is_scoped_base_repair() { }) ); assert!(state.pending_fenced_repair().is_none()); + assert_eq!(state.dirty_time_windows.len(), 1); // Simulate the outer execution failure restore for the in-flight chunk. state.restore_scoped_windows(&filter); + assert_eq!(state.dirty_time_windows.len(), 2); } let plan = task @@ -1666,82 +1570,6 @@ fn test_fenced_repair_transient_non_stale_failure_retries_same_high() { ); } -/// When `query_failure_reason` classifies a wrapped error as -/// `SnapshotFenceExpired` via the text-marker fallback (not via -/// `StatusCode::RequestOutdated`), the state machine must still -/// abandon the fenced repair and produce a `ScopedBaseRepair` plan -/// next, exactly like the structured-code path. -#[tokio::test] -async fn test_text_fallback_stale_fence_produces_scoped_base_repair() { - let TestTaskParts { - task, - query_engine, - .. - } = new_time_window_test_task_with_query( - "SELECT number, date_bin(INTERVAL '5 second', ts) AS time_window FROM numbers_with_ts GROUP BY time_window, number", - ) - .await; - let high = BTreeMap::from([(1_u64, 10_u64), (2_u64, 20_u64)]); - let filter = { - let mut state = task.state.write().unwrap(); - state - .dirty_time_windows - .add_window(Timestamp::new_second(10), Some(Timestamp::new_second(15))); - state - .dirty_time_windows - .add_window(Timestamp::new_second(100), Some(Timestamp::new_second(105))); - state.start_fenced_repair(high.clone()).unwrap(); - next_fenced_repair_filter(&mut state, 1) - }; - - // Construct a wrapped error that hits the text fallback (non-RequestOutdated - // status code with STALE_SNAPSHOT_FENCE marker text). - let err = flow_error_with_code_and_text( - StatusCode::Internal, - "STALE_SNAPSHOT_FENCE: snapshot upper bound stale, retry_hint: REBIND_SNAPSHOT_FENCE", - ); - let coverage = QueryCoverage::FencedRepairChunk { high }; - let reason = BatchingTask::query_failure_reason(&err, &coverage); - assert_eq!(reason, FlowQueryFallbackReason::SnapshotFenceExpired); - - { - let mut state = task.state.write().unwrap(); - let decision = BatchingTask::apply_query_failure_to_state( - &mut state, - std::time::Duration::from_millis(1), - &coverage, - reason, - ); - assert_eq!( - decision, - Some(FlowCheckpointDecision::FallbackToFullSnapshot { - previous_mode: CheckpointMode::FullSnapshot, - reason: FlowQueryFallbackReason::SnapshotFenceExpired, - }) - ); - assert!(state.pending_fenced_repair().is_none()); - - // Simulate the outer execution failure restore for the in-flight chunk. - state.restore_scoped_windows(&filter); - } - - let plan = task - .gen_query_with_time_window( - query_engine, - &aggregate_time_window_sink_schema(), - &[], - false, - Some(1), - ) - .await - .unwrap() - .expect("text-fallback stale fence should restore dirty windows for a fresh scoped repair"); - assert!( - matches!(plan.coverage, QueryCoverage::ScopedBaseRepair), - "next plan after text-fallback stale fence should be ScopedBaseRepair" - ); -} - #[test] fn test_checkpoint_decision_labels_are_stable() { let advance = FlowCheckpointDecision::AdvancedIncremental { diff --git a/src/query/src/datafusion.rs b/src/query/src/datafusion.rs index 63b834bed55..7df99f059ca 100644 --- a/src/query/src/datafusion.rs +++ b/src/query/src/datafusion.rs @@ -814,17 +814,20 @@ mod tests { use api::v1::SemanticType; use arrow::array::{ArrayRef, UInt64Array}; use arrow_schema::SortOptions; + use async_trait::async_trait; use catalog::RegisterTableRequest; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, NUMBERS_TABLE_ID}; use common_error::ext::BoxedError; - use common_recordbatch::{EmptyRecordBatchStream, SendableRecordBatchStream, util}; + use common_recordbatch::{ + EmptyRecordBatchStream, RecordBatch, SendableRecordBatchStream, util, + }; use datafusion::physical_plan::display::{DisplayAs, DisplayFormatType}; use datafusion::physical_plan::expressions::PhysicalSortExpr; use datafusion::physical_plan::joins::{HashJoinExec, JoinOn, PartitionMode}; use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet; use datafusion::physical_plan::{ExecutionPlan, PhysicalExpr}; use datafusion::prelude::{col, lit}; - use datafusion_common::{JoinType, NullEquality}; + use datafusion_common::{JoinType, NullEquality, ScalarValue}; use datafusion_physical_expr::expressions::Column; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, SchemaRef}; @@ -835,12 +838,13 @@ mod tests { PartitionRange, PrepareRequest, QueryScanContext, RegionScanner, ScannerProperties, }; use store_api::storage::{RegionId, ScanRequest}; + use table::metadata::{TableInfoBuilder, TableMetaBuilder}; use table::table::numbers::{NUMBERS_TABLE_NAME, NumbersTable}; use table::table::scan::RegionScanExec; use super::*; use crate::options::QueryOptions; - use crate::parser::QueryLanguageParser; + use crate::parser::{QueryLanguageParser, QueryStatement}; use crate::part_sort::PartSortExec; use crate::query_engine::{QueryEngineFactory, QueryEngineRef}; @@ -1422,4 +1426,341 @@ mod tests { assert!(right_update_calls.load(Ordering::Relaxed) > 0); assert!(right_last_filter_len.load(Ordering::Relaxed) > 0); } + #[derive(Default)] + struct RecordingMutationHandler { + inserts: std::sync::Mutex>, + } + + #[async_trait] + impl common_function::handlers::TableMutationHandler for RecordingMutationHandler { + async fn insert( + &self, + request: table::requests::InsertRequest, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + self.inserts.lock().unwrap().push(request); + Ok(common_query::Output::new_with_affected_rows(1)) + } + + async fn delete( + &self, + _request: table::requests::DeleteRequest, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected delete") + } + + async fn flush( + &self, + _request: table::requests::FlushTableRequest, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected flush") + } + + async fn compact( + &self, + _request: table::requests::CompactTableRequest, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected compact") + } + + async fn build_index( + &self, + _request: table::requests::BuildIndexTableRequest, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected build_index") + } + + async fn flush_region( + &self, + _region_id: store_api::storage::RegionId, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected flush_region") + } + + async fn compact_region( + &self, + _region_id: store_api::storage::RegionId, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected compact_region") + } + + async fn discard_unflushed_data( + &self, + _region_id: store_api::storage::RegionId, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected discard_unflushed_data") + } + + async fn discard_unflushed_data_by_table( + &self, + _table_name: table::table_name::TableName, + _ctx: session::context::QueryContextRef, + ) -> common_query::error::Result { + unimplemented!("unexpected discard_unflushed_data_by_table") + } + } + + fn native_schema() -> Arc { + Arc::new(Schema::new(vec![ + ColumnSchema::new("dim", ConcreteDataType::date_datatype(), true), + ColumnSchema::new("amount", ConcreteDataType::decimal128_datatype(30, 2), true), + ColumnSchema::new( + "elapsed", + ConcreteDataType::duration_millisecond_datatype(), + true, + ), + ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ColumnSchema::new( + "updated_at", + ConcreteDataType::timestamp_millisecond_datatype(), + true, + ), + ColumnSchema::new("marker", ConcreteDataType::uint8_datatype(), true), + ColumnSchema::new("payload", ConcreteDataType::binary_datatype(), true), + ColumnSchema::new("epoch", ConcreteDataType::uint64_datatype(), true), + ])) + } + + fn register_native_tables(catalog: &catalog::memory::MemoryCatalogManager) { + let schema = native_schema(); + let meta = TableMetaBuilder::empty() + .schema(schema.clone()) + .primary_key_indices(vec![]) + .value_indices((0..schema.num_columns()).collect()) + .next_column_id(8) + .build() + .unwrap(); + let info = TableInfoBuilder::default() + .name("native_regression") + .table_id(9001) + .table_version(0) + .meta(meta) + .build() + .unwrap(); + catalog + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "native_regression".to_string(), + table_id: 9001, + table: table::test_util::EmptyTable::from_table_info(&info), + }) + .unwrap(); + + let schema = Arc::new(Schema::new(vec![ + ColumnSchema::new("dim", ConcreteDataType::date_datatype(), true), + ColumnSchema::new("amount", ConcreteDataType::decimal128_datatype(20, 2), true), + ColumnSchema::new( + "elapsed", + ConcreteDataType::duration_millisecond_datatype(), + true, + ), + ColumnSchema::new( + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + false, + ) + .with_time_index(true), + ])); + let rows = RecordBatch::new( + schema, + vec![ + Arc::new(datatypes::vectors::DateVector::from_slice([0, 2])) as VectorRef, + Arc::new( + datatypes::vectors::Decimal128Vector::from_slice([10000, 20000]) + .with_precision_and_scale(20, 2) + .unwrap(), + ) as VectorRef, + Arc::new(datatypes::vectors::DurationMillisecondVector::from_values( + [10, 20], + )) as VectorRef, + Arc::new(datatypes::vectors::TimestampMillisecondVector::from_slice( + [1, 2], + )) as VectorRef, + ], + ) + .unwrap(); + catalog + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: "native_aggregate".to_string(), + table_id: 9002, + table: table::test_util::MemTable::table("native_aggregate", rows), + }) + .unwrap(); + } + + async fn run_sql(engine: &QueryEngineRef, sql: &str) -> Vec { + let stmt = QueryLanguageParser::parse_sql(sql, &QueryContext::arc()).unwrap(); + let plan = engine + .planner() + .plan(&stmt, QueryContext::arc()) + .await + .unwrap(); + match engine + .execute(plan, QueryContext::arc()) + .await + .unwrap() + .data + { + OutputData::Stream(stream) => util::collect(stream).await.unwrap(), + _ => unreachable!(), + } + } + + #[tokio::test] + async fn test_native_executor_null_insert_and_aggregates() { + let catalog = catalog::memory::new_memory_catalog_manager().unwrap(); + register_native_tables(&catalog); + let handler = Arc::new(RecordingMutationHandler::default()); + let engine = QueryEngineFactory::new( + catalog, + None, + Some(handler.clone()), + None, + None, + false, + QueryOptions::default(), + ) + .query_engine(); + + // The CASTs force Insert::can_extract_values false. This records the + // QueryEngine DML path selected by operator/src/statement/dml.rs. + let insert_sql = "INSERT INTO native_regression (dim, amount, elapsed, ts, updated_at, marker, payload, epoch) VALUES (NULL, NULL, NULL, CAST(-62135596799999 AS TIMESTAMP(3)), CAST(NULL AS TIMESTAMP(3)), CAST(1 AS UInt8), X'0102', CAST(1 AS UInt64))"; + let stmt = QueryLanguageParser::parse_sql(insert_sql, &QueryContext::arc()).unwrap(); + let QueryStatement::Sql(sql::statements::statement::Statement::Insert(insert)) = &stmt + else { + unreachable!() + }; + assert!(!insert.can_extract_values()); + let plan = engine + .planner() + .plan(&stmt, QueryContext::arc()) + .await + .unwrap(); + assert!(matches!( + engine + .execute(plan, QueryContext::arc()) + .await + .unwrap() + .data, + OutputData::AffectedRows(1) + )); + + let request = handler.inserts.lock().unwrap().pop().unwrap(); + for (name, data_type) in [ + ("dim", ConcreteDataType::date_datatype()), + ("amount", ConcreteDataType::decimal128_datatype(30, 2)), + ("elapsed", ConcreteDataType::duration_millisecond_datatype()), + ( + "updated_at", + ConcreteDataType::timestamp_millisecond_datatype(), + ), + ] { + let vector = request.columns_values.get(name).unwrap(); + assert_eq!(vector.len(), 1, "{name}"); + assert_eq!(vector.data_type(), data_type, "{name}"); + assert!(vector.is_null(0), "{name}"); + } + assert!(!request.columns_values["ts"].is_null(0)); + for (name, data_type) in [ + ("marker", ConcreteDataType::uint8_datatype()), + ("payload", ConcreteDataType::binary_datatype()), + ("epoch", ConcreteDataType::uint64_datatype()), + ] { + assert_eq!( + request.columns_values[name].data_type(), + data_type, + "{name}" + ); + } + + let batches = run_sql( + &engine, + "SELECT MIN(dim), MAX(dim), SUM(amount), SUM(elapsed) FROM native_aggregate", + ) + .await; + let batch = &batches[0]; + assert_eq!(batch.num_rows(), 1); + for (index, (data_type, value)) in [ + ( + ConcreteDataType::date_datatype(), + ScalarValue::Date32(Some(0)), + ), + ( + ConcreteDataType::date_datatype(), + ScalarValue::Date32(Some(2)), + ), + ( + ConcreteDataType::decimal128_datatype(30, 2), + ScalarValue::Decimal128(Some(30000), 30, 2), + ), + ( + ConcreteDataType::duration_millisecond_datatype(), + ScalarValue::DurationMillisecond(Some(30)), + ), + ] + .into_iter() + .enumerate() + { + assert_eq!(batch.schema.column_schemas()[index].data_type, data_type); + assert_eq!( + ScalarValue::try_from_array(batch.column(index).as_ref(), 0).unwrap(), + value + ); + } + + let batches = run_sql( + &engine, + "SELECT aggregate.amount_sum + aggregate.amount_sum AS amount_add, \ + aggregate.elapsed_sum + aggregate.elapsed_sum AS elapsed_add, \ + CASE WHEN aggregate.min_dim < CAST('1970-01-02' AS DATE) THEN true ELSE false END AS min_before, \ + CASE WHEN aggregate.max_dim > CAST('1970-01-01' AS DATE) THEN true ELSE false END AS max_after \ + FROM (SELECT MIN(dim) AS min_dim, MAX(dim) AS max_dim, SUM(amount) AS amount_sum, SUM(elapsed) AS elapsed_sum \ + FROM native_aggregate) AS aggregate", + ) + .await; + let batch = &batches[0]; + assert_eq!(batch.num_rows(), 1); + for (index, (data_type, value)) in [ + ( + ConcreteDataType::decimal128_datatype(31, 2), + ScalarValue::Decimal128(Some(60000), 31, 2), + ), + ( + ConcreteDataType::duration_millisecond_datatype(), + ScalarValue::DurationMillisecond(Some(60)), + ), + ( + ConcreteDataType::boolean_datatype(), + ScalarValue::Boolean(Some(true)), + ), + ( + ConcreteDataType::boolean_datatype(), + ScalarValue::Boolean(Some(true)), + ), + ] + .into_iter() + .enumerate() + { + assert_eq!(batch.schema.column_schemas()[index].data_type, data_type); + assert_eq!( + ScalarValue::try_from_array(batch.column(index).as_ref(), 0).unwrap(), + value + ); + } + } } diff --git a/src/query/src/dist_plan/merge_scan.rs b/src/query/src/dist_plan/merge_scan.rs index e7e78bfec61..2efdd7b1fa2 100644 --- a/src/query/src/dist_plan/merge_scan.rs +++ b/src/query/src/dist_plan/merge_scan.rs @@ -24,6 +24,7 @@ use arrow_schema::{ }; use async_stream::stream; use common_catalog::parse_catalog_and_schema_from_db_string; +use common_error::ext::BoxedError; use common_plugins::GREPTIME_EXEC_READ_COST; use common_query::request::QueryRequest; use common_recordbatch::adapter::{RecordBatchMetrics, region_scan_output_bytes}; @@ -746,7 +747,7 @@ impl MergeScanExec { } let mut stream = do_get_result.map_err(|e| { MERGE_SCAN_ERRORS_TOTAL.inc(); - DataFusionError::External(Box::new(e)) + DataFusionError::External(Box::new(BoxedError::new(e))) })?; if let Some(subscriber_rollback) = subscriber_rollback.as_mut() { @@ -815,7 +816,8 @@ impl MergeScanExec { let poll_elapsed = poll_timer.elapsed(); poll_duration += poll_elapsed; - let batch = batch.map_err(|e| DataFusionError::External(Box::new(e)))?; + let batch = batch + .map_err(|e| DataFusionError::External(Box::new(BoxedError::new(e))))?; let df_batch = batch.into_df_record_batch(); if !Arc::ptr_eq(&advertised_schema, df_batch.schema_ref()) { validate_remote_schema( @@ -1423,6 +1425,8 @@ mod tests { use arrow_schema::{DataType as TestArrowDataType, Field, TimeUnit}; use async_trait::async_trait; use common_base::Plugins; + use common_error::ext::{ErrorExt, PlainError}; + use common_error::status_code::StatusCode; use common_meta::peer::Peer; use common_query::request::{ INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY, InitialDynFilterRegs, @@ -1434,6 +1438,7 @@ mod tests { use datafusion::config::ConfigOptions; use datafusion::execution::SessionStateBuilder; use datafusion::physical_plan::filter_pushdown::ChildFilterPushdownResult; + use datafusion::physical_plan::repartition::RepartitionExec; use datafusion::physical_plan::{StatisticsArgs, StatisticsContext}; use datafusion_common::TableReference; use datafusion_expr::{LogicalPlanBuilder, col, lit}; @@ -1448,6 +1453,7 @@ mod tests { use session::ReadPreference; use session::context::QueryContext; use session::query_id::QueryId; + use snafu::IntoError; use table::table::scan::REGION_SCAN_EXEC_NAME; use table::table_name::TableName; use tokio::sync::{Notify, oneshot}; @@ -1830,7 +1836,7 @@ mod tests { } #[tokio::test] - async fn failed_do_get_rolls_back_new_subscriber_without_starting_fanout() { + async fn failed_do_get_preserves_status_code_and_rolls_back_subscriber() { let handler = Arc::new(FailingRegionQueryHandler::default()); let query_ctx = QueryContext::arc(); let state = Arc::new(QueryEngineState::new( @@ -1886,8 +1892,12 @@ mod tests { ) .unwrap(); - let mut stream = exec.to_stream(task_ctx, 0).unwrap(); - assert!(stream.next().await.unwrap().is_err()); + let mut stream = common_recordbatch::adapter::RecordBatchStreamAdapter::try_new( + exec.to_stream(task_ctx, 0).unwrap(), + ) + .unwrap(); + let error = stream.next().await.unwrap().unwrap_err(); + assert_eq!(error.status_code(), StatusCode::RequestOutdated); assert_eq!(handler.do_get_calls.load(Ordering::SeqCst), 1); assert!(handler.saw_subscriber.load(Ordering::SeqCst)); @@ -1897,6 +1907,46 @@ mod tests { assert!(!entries[0].fanout_started_for_test()); } + #[tokio::test] + async fn repartitioned_merge_scan_later_stream_error_preserves_status_code() { + let region_id = RegionId::new(1024, 1); + let handler = Arc::new(TestRegionQueryHandler::with_responses(vec![( + region_id, + int64_schema(&["a", "b"]), + vec![Err(common_recordbatch::error::ExternalSnafu.into_error( + BoxedError::new(PlainError::new( + "neutral stream error".to_string(), + StatusCode::RequestOutdated, + )), + ))], + )])); + let merge_scan = Arc::new(merge_scan_exec_with_handler( + vec![region_id], + expected_int64_schema(), + handler, + 1, + )); + let repartition = + RepartitionExec::try_new(merge_scan, Partitioning::RoundRobinBatch(2)).unwrap(); + assert_eq!( + repartition + .properties() + .output_partitioning() + .partition_count(), + 2 + ); + + let mut stream = common_recordbatch::adapter::RecordBatchStreamAdapter::try_new( + repartition + .execute(0, Arc::new(TaskContext::default())) + .unwrap(), + ) + .unwrap(); + + let error = stream.next().await.unwrap().unwrap_err(); + assert_eq!(error.status_code(), StatusCode::RequestOutdated); + } + #[tokio::test] async fn aborting_pending_do_get_poll_rolls_back_subscriber_without_starting_fanout() { let handler = Arc::new(PendingDoGetHandler::default()); @@ -2097,10 +2147,9 @@ mod tests { assert_eq!(registry_manager.registry_count(), 0); } - #[derive(Clone)] struct TestRegionResponse { advertised_schema: Arc, - batches: Vec, + batches: Vec>, } #[derive(Default)] @@ -2117,7 +2166,7 @@ mod tests { region_id, TestRegionResponse { advertised_schema: batch.schema.clone(), - batches: vec![batch], + batches: vec![Ok(batch)], }, ) }) @@ -2126,7 +2175,13 @@ mod tests { } fn with_responses( - responses: impl IntoIterator, Vec)>, + responses: impl IntoIterator< + Item = ( + RegionId, + Arc, + Vec>, + ), + >, ) -> Self { let responses = responses .into_iter() @@ -2146,19 +2201,17 @@ mod tests { struct TestRecordBatchStream { schema: Arc, - batches: Vec, - index: usize, + batches: Vec>, } impl Stream for TestRecordBatchStream { type Item = common_recordbatch::error::Result; fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { - if let Some(batch) = self.batches.get(self.index).cloned() { - self.index += 1; - Poll::Ready(Some(Ok(batch))) - } else { + if self.batches.is_empty() { Poll::Ready(None) + } else { + Poll::Ready(Some(self.batches.remove(0))) } } } @@ -2219,10 +2272,13 @@ mod tests { }), Ordering::SeqCst, ); - crate::error::UnimplementedSnafu { - operation: "test do_get failure", - } - .fail() + Err(crate::error::Error::QueryExecution { + source: BoxedError::new(PlainError::new( + "neutral do_get error".to_string(), + StatusCode::RequestOutdated, + )), + location: snafu::Location::default(), + }) } async fn handle_remote_dyn_filter_update( @@ -2540,8 +2596,19 @@ mod tests { .expect("test handler needs a response for every requested region"); Ok(Box::pin(TestRecordBatchStream { schema: response.advertised_schema.clone(), - batches: response.batches.clone(), - index: 0, + batches: response + .batches + .iter() + .map(|batch| match batch { + Ok(batch) => Ok(batch.clone()), + Err(error) => Err(common_recordbatch::error::ExternalSnafu.into_error( + BoxedError::new(PlainError::new( + error.to_string(), + error.status_code(), + )), + )), + }) + .collect(), })) } @@ -3102,7 +3169,7 @@ mod tests { Arc::new(TestRegionQueryHandler::with_responses(vec![( RegionId::new(1024, 1), remote_schema, - vec![batch], + vec![Ok(batch)], )])), 1, )) @@ -3139,7 +3206,7 @@ mod tests { Arc::new(TestRegionQueryHandler::with_responses(vec![( region_id, advertised_schema, - vec![batch], + vec![Ok(batch)], )])), 1, ); @@ -3172,7 +3239,7 @@ mod tests { Arc::new(TestRegionQueryHandler::with_responses(vec![( region_id, advertised_schema, - vec![batch], + vec![Ok(batch)], )])), 1, ); @@ -3207,7 +3274,7 @@ mod tests { Arc::new(TestRegionQueryHandler::with_responses(vec![( region_id, advertised_schema, - vec![batch], + vec![Ok(batch)], )])), 1, )) diff --git a/src/query/src/error.rs b/src/query/src/error.rs index fdfeaafde97..a2946415570 100644 --- a/src/query/src/error.rs +++ b/src/query/src/error.rs @@ -498,8 +498,23 @@ impl From for DataFusionError { #[cfg(test)] mod tests { + use common_error::ext::PlainError; + use super::*; + #[test] + fn test_datafusion_external_boxed_error_status_code() { + let error = Error::DataFusion { + error: DataFusionError::External(Box::new(BoxedError::new(PlainError::new( + "neutral error".to_string(), + StatusCode::RequestOutdated, + )))), + location: Location::default(), + }; + + assert_eq!(error.status_code(), StatusCode::RequestOutdated); + } + #[test] fn test_build_backend_delegates_error_metadata() { let source = common_datasource::error::LocalFileAccessDisabledSnafu { diff --git a/src/servers/src/error.rs b/src/servers/src/error.rs index fa5957759c7..fe1d226a10c 100644 --- a/src/servers/src/error.rs +++ b/src/servers/src/error.rs @@ -745,7 +745,7 @@ impl ErrorExt for Error { #[cfg(not(windows))] UpdateJemallocMetrics { .. } => StatusCode::Internal, - CollectRecordbatch { .. } => StatusCode::EngineExecuteQuery, + CollectRecordbatch { source, .. } => source.status_code(), ExecuteQuery { source, .. } | ExecutePlan { source, .. } @@ -986,3 +986,37 @@ pub fn status_code_to_http_status(status_code: &StatusCode) -> HttpStatusCode { | StatusCode::EngineExecuteQuery => HttpStatusCode::INTERNAL_SERVER_ERROR, } } + +#[cfg(test)] +mod tests { + use common_error::GREPTIME_DB_HEADER_ERROR_CODE; + use common_error::ext::PlainError; + + use super::*; + + #[test] + fn collect_recordbatch_preserves_poll_stream_status_in_tonic_status() { + let error = Error::CollectRecordbatch { + source: common_recordbatch::error::Error::PollStream { + error: DataFusionError::External(Box::new(BoxedError::new(PlainError::new( + "neutral error".to_string(), + StatusCode::RequestOutdated, + )))), + location: Location::default(), + }, + location: Location::default(), + }; + + let status: tonic::Status = error.into(); + assert_eq!(status.code(), tonic::Code::InvalidArgument); + assert_eq!( + status + .metadata() + .get(GREPTIME_DB_HEADER_ERROR_CODE) + .unwrap() + .to_str() + .unwrap(), + (StatusCode::RequestOutdated as u32).to_string() + ); + } +} diff --git a/tests-integration/src/grpc/flight.rs b/tests-integration/src/grpc/flight.rs index 2486dddee05..fb27a273516 100644 --- a/tests-integration/src/grpc/flight.rs +++ b/tests-integration/src/grpc/flight.rs @@ -33,6 +33,8 @@ mod test { use client::region::RegionRequester; use client::{Client, Database}; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; + use common_error::ext::ErrorExt; + use common_error::status_code::StatusCode; use common_grpc::channel_manager::{ChannelConfig, ChannelManager}; use common_grpc::flight::do_put::{DoPutMetadata, DoPutResponse}; use common_grpc::flight::{FlightDecoder, FlightEncoder, FlightMessage}; @@ -784,12 +786,14 @@ mod test { let mut stale_fence_error = None; while let Some(batch) = stream.next().await { if let Err(err) = batch { - stale_fence_error = Some(format!("{err:?}")); + stale_fence_error = Some(err); break; } } - let err_msg = stale_fence_error.expect("expected stale snapshot fence rejection"); + let stale_fence_error = stale_fence_error.expect("expected stale snapshot fence rejection"); + assert_eq!(stale_fence_error.status_code(), StatusCode::RequestOutdated); + let err_msg = format!("{stale_fence_error:?}"); assert!( err_msg.contains("STALE_SNAPSHOT_FENCE") || err_msg.contains("RequestOutdated") diff --git a/tests/cases/distributed/optimizer/filter_push_down.result b/tests/cases/distributed/optimizer/filter_push_down.result index 83c438de4e9..c4eab6b1e09 100644 --- a/tests/cases/distributed/optimizer/filter_push_down.result +++ b/tests/cases/distributed/optimizer/filter_push_down.result @@ -157,7 +157,7 @@ SELECT * FROM integers i1 WHERE NOT EXISTS(SELECT i FROM integers WHERE i=i1.i) SELECT i1.i,i2.i FROM integers i1, integers i2 WHERE i1.i=(SELECT i FROM integers WHERE i1.i=i) AND i1.i=i2.i ORDER BY i1.i; -Error: 3001(EngineExecuteQuery), Error during planning: Correlated scalar subquery must be aggregated to return at most one row +Error: 3000(PlanQuery), Error during planning: Correlated scalar subquery must be aggregated to return at most one row SELECT * FROM (SELECT i1.i AS a, i2.i AS b FROM integers i1, integers i2) a1 WHERE a=b ORDER BY 1; diff --git a/tests/cases/standalone/common/order/limit.result b/tests/cases/standalone/common/order/limit.result index 558f8f6a0fa..b69f44e0465 100644 --- a/tests/cases/standalone/common/order/limit.result +++ b/tests/cases/standalone/common/order/limit.result @@ -25,7 +25,7 @@ SELECT b FROM test ORDER BY b LIMIT 2 OFFSET 0; SELECT a FROM test LIMIT 1.25; -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Float64 +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Float64 SELECT a FROM test LIMIT 2-1; @@ -49,7 +49,7 @@ Error: 1001(Unsupported), This feature is not implemented: Unsupported LIMIT exp SELECT a FROM test LIMIT row_number() OVER (); -Error: 3001(EngineExecuteQuery), This feature is not implemented: Unsupported LIMIT expression: Some(Cast(Cast { expr: WindowFunction(WindowFunction { fun: WindowUDF(WindowUDF { inner: RowNumber { signature: Signature { type_signature: Nullary, volatility: Immutable, parameter_names: None } } }), params: WindowFunctionParams { args: [], partition_by: [], order_by: [], window_frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }, filter: None, null_treatment: None, distinct: false } }), field: Field { name: "", data_type: Int64, nullable: true } })) +Error: 1001(Unsupported), This feature is not implemented: Unsupported LIMIT expression: Some(Cast(Cast { expr: WindowFunction(WindowFunction { fun: WindowUDF(WindowUDF { inner: RowNumber { signature: Signature { type_signature: Nullary, volatility: Immutable, parameter_names: None } } }), params: WindowFunctionParams { args: [], partition_by: [], order_by: [], window_frame: WindowFrame { units: Rows, start_bound: Preceding(UInt64(NULL)), end_bound: Following(UInt64(NULL)), is_causal: false }, filter: None, null_treatment: None, distinct: false } }), field: Field { name: "", data_type: Int64, nullable: true } })) CREATE TABLE test2 (a STRING, ts TIMESTAMP TIME INDEX); @@ -69,7 +69,7 @@ SELECT * FROM test2 LIMIT 3; select 1 limit date '1992-01-01'; -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Date32 +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Date32 CREATE TABLE integers(i TIMESTAMP TIME INDEX); @@ -102,23 +102,23 @@ SELECT * FROM integers LIMIT 4; SELECT * FROM integers as int LIMIT (SELECT MIN(integers.i) FROM integers); -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) SELECT * FROM integers as int OFFSET (SELECT MIN(integers.i) FROM integers); -Error: 3001(EngineExecuteQuery), Error during planning: Expected OFFSET to be an integer or null, but got Timestamp(ms) +Error: 3000(PlanQuery), Error during planning: Expected OFFSET to be an integer or null, but got Timestamp(ms) SELECT * FROM integers as int LIMIT (SELECT MAX(integers.i) FROM integers) OFFSET (SELECT MIN(integers.i) FROM integers); -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) SELECT * FROM integers as int LIMIT (SELECT max(integers.i) FROM integers where i > 5); -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) SELECT * FROM integers as int LIMIT (SELECT max(integers.i) FROM integers where i > 5); -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Timestamp(ms) SELECT * FROM integers as int LIMIT (SELECT NULL); @@ -130,7 +130,7 @@ Error: 1001(Unsupported), This feature is not implemented: Unsupported LIMIT exp SELECT * FROM integers as int LIMIT (SELECT 'ab'); -Error: 3001(EngineExecuteQuery), Error during planning: Expected LIMIT to be an integer or null, but got Utf8 +Error: 3000(PlanQuery), Error during planning: Expected LIMIT to be an integer or null, but got Utf8 DROP TABLE integers; diff --git a/tests/cases/standalone/optimizer/filter_push_down.result b/tests/cases/standalone/optimizer/filter_push_down.result index 6705bb3c83f..242cebd1167 100644 --- a/tests/cases/standalone/optimizer/filter_push_down.result +++ b/tests/cases/standalone/optimizer/filter_push_down.result @@ -139,7 +139,7 @@ SELECT * FROM integers i1 WHERE NOT EXISTS(SELECT i FROM integers WHERE i=i1.i) SELECT i1.i,i2.i FROM integers i1, integers i2 WHERE i1.i=(SELECT i FROM integers WHERE i1.i=i) AND i1.i=i2.i ORDER BY i1.i; -Error: 3001(EngineExecuteQuery), Error during planning: Correlated scalar subquery must be aggregated to return at most one row +Error: 3000(PlanQuery), Error during planning: Correlated scalar subquery must be aggregated to return at most one row SELECT * FROM (SELECT i1.i AS a, i2.i AS b FROM integers i1, integers i2) a1 WHERE a=b ORDER BY 1;