mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-08 13:02:42 +00:00
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>
This commit is contained in:
@@ -270,18 +270,158 @@ pub fn datafusion_status_code<T: ErrorExt + 'static>(
|
||||
e: &DataFusionError,
|
||||
default_status: Option<StatusCode>,
|
||||
) -> 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::<T>() {
|
||||
DataFusionError::External(error) => {
|
||||
if let Some(ext) = (*error).downcast_ref::<T>() {
|
||||
ext.status_code()
|
||||
} else if let Some(ext) = (*error).downcast_ref::<BoxedError>() {
|
||||
ext.status_code()
|
||||
} else {
|
||||
default_status.unwrap_or(StatusCode::EngineExecuteQuery)
|
||||
}
|
||||
}
|
||||
DataFusionError::Diagnostic(_, e) => datafusion_status_code::<T>(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>(&error, None), status);
|
||||
assert_eq!(
|
||||
datafusion_status_code::<Error>(&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::<Error>(&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::<Error>(&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>(&error, None), none_expected);
|
||||
assert_eq!(
|
||||
datafusion_status_code::<Error>(&error, Some(StatusCode::Unknown)),
|
||||
default_expected
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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::<BoxedError>()
|
||||
.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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<String>) {}
|
||||
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 {
|
||||
|
||||
+344
-3
@@ -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<Vec<table::requests::InsertRequest>>,
|
||||
}
|
||||
|
||||
#[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<common_query::Output> {
|
||||
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<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected delete")
|
||||
}
|
||||
|
||||
async fn flush(
|
||||
&self,
|
||||
_request: table::requests::FlushTableRequest,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected flush")
|
||||
}
|
||||
|
||||
async fn compact(
|
||||
&self,
|
||||
_request: table::requests::CompactTableRequest,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected compact")
|
||||
}
|
||||
|
||||
async fn build_index(
|
||||
&self,
|
||||
_request: table::requests::BuildIndexTableRequest,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected build_index")
|
||||
}
|
||||
|
||||
async fn flush_region(
|
||||
&self,
|
||||
_region_id: store_api::storage::RegionId,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected flush_region")
|
||||
}
|
||||
|
||||
async fn compact_region(
|
||||
&self,
|
||||
_region_id: store_api::storage::RegionId,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected compact_region")
|
||||
}
|
||||
|
||||
async fn discard_unflushed_data(
|
||||
&self,
|
||||
_region_id: store_api::storage::RegionId,
|
||||
_ctx: session::context::QueryContextRef,
|
||||
) -> common_query::error::Result<common_base::AffectedRows> {
|
||||
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<common_base::AffectedRows> {
|
||||
unimplemented!("unexpected discard_unflushed_data_by_table")
|
||||
}
|
||||
}
|
||||
|
||||
fn native_schema() -> Arc<Schema> {
|
||||
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<RecordBatch> {
|
||||
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
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<Schema>,
|
||||
batches: Vec<RecordBatch>,
|
||||
batches: Vec<common_recordbatch::error::Result<RecordBatch>>,
|
||||
}
|
||||
|
||||
#[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<Item = (RegionId, Arc<Schema>, Vec<RecordBatch>)>,
|
||||
responses: impl IntoIterator<
|
||||
Item = (
|
||||
RegionId,
|
||||
Arc<Schema>,
|
||||
Vec<common_recordbatch::error::Result<RecordBatch>>,
|
||||
),
|
||||
>,
|
||||
) -> Self {
|
||||
let responses = responses
|
||||
.into_iter()
|
||||
@@ -2146,19 +2201,17 @@ mod tests {
|
||||
|
||||
struct TestRecordBatchStream {
|
||||
schema: Arc<Schema>,
|
||||
batches: Vec<RecordBatch>,
|
||||
index: usize,
|
||||
batches: Vec<common_recordbatch::error::Result<RecordBatch>>,
|
||||
}
|
||||
|
||||
impl Stream for TestRecordBatchStream {
|
||||
type Item = common_recordbatch::error::Result<RecordBatch>;
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
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,
|
||||
))
|
||||
|
||||
@@ -498,8 +498,23 @@ impl From<Error> 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 {
|
||||
|
||||
@@ -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()
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user