1use 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#[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
286pub 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
304pub 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);