Skip to main content

flow/
error.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Error definition for flow module
16
17use std::any::Any;
18
19use api::v1::CreateTableExpr;
20use arrow_schema::ArrowError;
21use common_error::ext::{BoxedError, RetryHint};
22use common_error::{
23    GREPTIME_DB_HEADER_ERROR_RETRY_HINT, define_into_tonic_status, from_err_code_msg_to_header,
24};
25use common_macro::stack_trace_debug;
26use common_telemetry::common_error::ext::ErrorExt;
27use common_telemetry::common_error::status_code::StatusCode;
28use snafu::{Location, Snafu};
29use tonic::codegen::http::HeaderValue;
30use tonic::metadata::MetadataMap;
31
32use crate::FlowId;
33
34/// This error is used to represent all possible errors that can occur in the flow module.
35#[derive(Snafu)]
36#[snafu(visibility(pub))]
37#[stack_trace_debug]
38pub enum Error {
39    #[snafu(display(
40        "Failed to insert into flow: region_id={}, flow_ids={:?}",
41        region_id,
42        flow_ids
43    ))]
44    InsertIntoFlow {
45        region_id: u64,
46        flow_ids: Vec<u64>,
47        source: BoxedError,
48        #[snafu(implicit)]
49        location: Location,
50    },
51
52    #[snafu(display("Flow engine is still recovering"))]
53    FlowNotRecovered {
54        #[snafu(implicit)]
55        location: Location,
56    },
57
58    #[snafu(display("Error encountered while creating flow: {sql}"))]
59    CreateFlow {
60        sql: String,
61        source: BoxedError,
62        #[snafu(implicit)]
63        location: Location,
64    },
65
66    #[snafu(display("Error encountered while creating sink table for flow: {create:?}"))]
67    CreateSinkTable {
68        create: CreateTableExpr,
69        source: BoxedError,
70        #[snafu(implicit)]
71        location: Location,
72    },
73
74    #[snafu(display("Time error"))]
75    Time {
76        source: common_time::error::Error,
77        #[snafu(implicit)]
78        location: Location,
79    },
80
81    #[snafu(display("No available frontend found after timeout: {timeout:?}, context: {context}"))]
82    NoAvailableFrontend {
83        timeout: std::time::Duration,
84        context: String,
85        #[snafu(implicit)]
86        location: Location,
87    },
88
89    #[snafu(display("External error"))]
90    External {
91        source: BoxedError,
92        #[snafu(implicit)]
93        location: Location,
94    },
95
96    #[snafu(display("Internal error"))]
97    Internal {
98        reason: String,
99        #[snafu(implicit)]
100        location: Location,
101    },
102
103    #[snafu(display("Table not found: {name}"))]
104    TableNotFound {
105        name: String,
106        #[snafu(implicit)]
107        location: Location,
108    },
109
110    #[snafu(display("Table not found: {msg}, meta error: {source}"))]
111    TableNotFoundMeta {
112        source: common_meta::error::Error,
113        msg: String,
114        #[snafu(implicit)]
115        location: Location,
116    },
117
118    #[snafu(display("Flow not found, id={id}"))]
119    FlowNotFound {
120        id: FlowId,
121        #[snafu(implicit)]
122        location: Location,
123    },
124
125    #[snafu(display("Failed to list flows in flownode={id:?}"))]
126    ListFlows {
127        id: Option<common_meta::FlownodeId>,
128        source: common_meta::error::Error,
129        #[snafu(implicit)]
130        location: Location,
131    },
132
133    #[snafu(display("Flow already exist, id={id}"))]
134    FlowAlreadyExist {
135        id: FlowId,
136        #[snafu(implicit)]
137        location: Location,
138    },
139
140    #[snafu(display("Failed to join task"))]
141    JoinTask {
142        #[snafu(source)]
143        error: tokio::task::JoinError,
144        #[snafu(implicit)]
145        location: Location,
146    },
147
148    #[snafu(display("Invalid query: {reason}"))]
149    InvalidQuery {
150        reason: String,
151        #[snafu(implicit)]
152        location: Location,
153    },
154
155    #[snafu(display("Flow plan error: {reason}"))]
156    Plan {
157        reason: String,
158        #[snafu(implicit)]
159        location: Location,
160    },
161
162    #[snafu(display("Unsupported: {reason}"))]
163    Unsupported {
164        reason: String,
165        #[snafu(implicit)]
166        location: Location,
167    },
168
169    #[snafu(display("Datatypes error: {source} with extra message: {extra}"))]
170    Datatypes {
171        source: datatypes::Error,
172        extra: String,
173        #[snafu(implicit)]
174        location: Location,
175    },
176
177    #[snafu(display("Arrow error: {raw:?} in context: {context}"))]
178    Arrow {
179        #[snafu(source)]
180        raw: ArrowError,
181        context: String,
182        #[snafu(implicit)]
183        location: Location,
184    },
185
186    #[snafu(display("Datafusion error: {raw:?} in context: {context}"))]
187    Datafusion {
188        #[snafu(source)]
189        raw: datafusion_common::DataFusionError,
190        context: String,
191        #[snafu(implicit)]
192        location: Location,
193    },
194
195    #[snafu(display("Unexpected: {reason}"))]
196    Unexpected {
197        reason: String,
198        #[snafu(implicit)]
199        location: Location,
200    },
201
202    #[snafu(display("Illegal check task state: {reason}"))]
203    IllegalCheckTaskState {
204        reason: String,
205        #[snafu(implicit)]
206        location: Location,
207    },
208
209    #[snafu(display(
210        "Failed to sync with check task for flow {} with allow_drop={}",
211        flow_id,
212        allow_drop
213    ))]
214    SyncCheckTask {
215        flow_id: FlowId,
216        allow_drop: bool,
217        #[snafu(implicit)]
218        location: Location,
219    },
220
221    #[snafu(display("Failed to start server"))]
222    StartServer {
223        #[snafu(implicit)]
224        location: Location,
225        source: servers::error::Error,
226    },
227
228    #[snafu(display("Failed to shutdown server"))]
229    ShutdownServer {
230        #[snafu(implicit)]
231        location: Location,
232        source: servers::error::Error,
233    },
234
235    #[snafu(display("Failed to initialize meta client"))]
236    MetaClientInit {
237        #[snafu(implicit)]
238        location: Location,
239        source: meta_client::error::Error,
240    },
241
242    #[snafu(display("Failed to parse address {}", addr))]
243    ParseAddr {
244        addr: String,
245        #[snafu(source)]
246        error: std::net::AddrParseError,
247    },
248
249    #[snafu(display("Failed to get cache from cache registry: {}", name))]
250    CacheRequired {
251        #[snafu(implicit)]
252        location: Location,
253        name: String,
254    },
255
256    #[snafu(display("Invalid request: {context}"))]
257    InvalidRequest {
258        context: String,
259        source: client::Error,
260        #[snafu(implicit)]
261        location: Location,
262    },
263
264    #[snafu(display("Failed to encode logical plan in substrait"))]
265    SubstraitEncodeLogicalPlan {
266        #[snafu(implicit)]
267        location: Location,
268        source: substrait::error::Error,
269    },
270
271    #[snafu(display("Failed to convert column schema to proto column def"))]
272    ConvertColumnSchema {
273        #[snafu(implicit)]
274        location: Location,
275        source: operator::error::Error,
276    },
277
278    #[snafu(display("Failed to create channel manager for gRPC client"))]
279    InvalidClientConfig {
280        #[snafu(implicit)]
281        location: Location,
282        source: common_grpc::error::Error,
283    },
284}
285
286/// the outer message is the full error stack, and inner message in header is the last error message that can be show directly to user
287pub fn to_status_with_last_err(err: impl ErrorExt) -> tonic::Status {
288    let msg = err.to_string();
289    let last_err_msg = common_error::ext::StackError::last(&err).to_string();
290    let code = err.status_code() as u32;
291    let mut header = from_err_code_msg_to_header(code, &last_err_msg);
292    header.insert(
293        GREPTIME_DB_HEADER_ERROR_RETRY_HINT,
294        HeaderValue::from_static(err.retry_hint().as_str()),
295    );
296
297    tonic::Status::with_metadata(
298        tonic::Code::InvalidArgument,
299        msg,
300        MetadataMap::from_headers(header),
301    )
302}
303
304/// Result type for flow module
305pub type Result<T> = std::result::Result<T, Error>;
306
307impl ErrorExt for Error {
308    fn status_code(&self) -> StatusCode {
309        match self {
310            Self::JoinTask { .. }
311            | Self::Datafusion { .. }
312            | Self::InsertIntoFlow { .. }
313            | Self::NoAvailableFrontend { .. }
314            | Self::FlowNotRecovered { .. } => StatusCode::Internal,
315            Self::FlowAlreadyExist { .. } => StatusCode::TableAlreadyExists,
316            Self::TableNotFound { .. }
317            | Self::TableNotFoundMeta { .. }
318            | Self::ListFlows { .. } => StatusCode::TableNotFound,
319            Self::FlowNotFound { .. } => StatusCode::FlowNotFound,
320            Self::Plan { .. } | Self::Datatypes { .. } => StatusCode::PlanQuery,
321            Self::CreateFlow { .. }
322            | Self::CreateSinkTable { .. }
323            | Self::Arrow { .. }
324            | Self::Time { .. } => StatusCode::EngineExecuteQuery,
325            Self::Unexpected { .. }
326            | Self::SyncCheckTask { .. }
327            | Self::IllegalCheckTaskState { .. } => StatusCode::Unexpected,
328            Self::Unsupported { .. } => StatusCode::Unsupported,
329            Self::External { source, .. } => source.status_code(),
330            Self::Internal { .. } | Self::CacheRequired { .. } => StatusCode::Internal,
331            Self::StartServer { source, .. } | Self::ShutdownServer { source, .. } => {
332                source.status_code()
333            }
334            Self::MetaClientInit { source, .. } => source.status_code(),
335
336            Self::InvalidQuery { .. }
337            | Self::InvalidRequest { .. }
338            | Self::ParseAddr { .. }
339            | Self::InvalidClientConfig { .. } => StatusCode::InvalidArguments,
340
341            Error::SubstraitEncodeLogicalPlan { source, .. } => source.status_code(),
342
343            Error::ConvertColumnSchema { source, .. } => source.status_code(),
344        }
345    }
346
347    fn as_any(&self) -> &dyn Any {
348        self
349    }
350
351    fn retry_hint(&self) -> RetryHint {
352        match self {
353            Self::FlowNotRecovered { .. } | Self::NoAvailableFrontend { .. } => {
354                RetryHint::Retryable
355            }
356
357            Self::InsertIntoFlow { source, .. }
358            | Self::CreateFlow { source, .. }
359            | Self::CreateSinkTable { source, .. }
360            | Self::External { source, .. } => source.retry_hint(),
361
362            Self::Time { source, .. } => source.retry_hint(),
363            Self::TableNotFoundMeta { source, .. } | Self::ListFlows { source, .. } => {
364                source.retry_hint()
365            }
366            Self::Datatypes { source, .. } => source.retry_hint(),
367            Self::StartServer { source, .. } | Self::ShutdownServer { source, .. } => {
368                source.retry_hint()
369            }
370            Self::MetaClientInit { source, .. } => source.retry_hint(),
371            Self::InvalidRequest { source, .. } => source.retry_hint(),
372            Self::SubstraitEncodeLogicalPlan { source, .. } => source.retry_hint(),
373            Self::ConvertColumnSchema { source, .. } => source.retry_hint(),
374            Self::InvalidClientConfig { source, .. } => source.retry_hint(),
375
376            _ => RetryHint::NonRetryable,
377        }
378    }
379}
380
381define_into_tonic_status!(Error);