1use std::any::Any;
16use std::sync::Arc;
17
18use common_error::ext::{ErrorExt, RetryHint, retry_hint_from_io_error};
19use common_error::status_code::StatusCode;
20use common_macro::stack_trace_debug;
21use common_runtime::error::Error as RuntimeError;
22use common_wal::kafka::rskafka_client_error_to_retry_hint;
23use object_store::error::retry_hint_from_opendal_error;
24use serde_json::error::Error as JsonError;
25use snafu::{Location, Snafu};
26use store_api::storage::RegionId;
27
28#[derive(Snafu)]
29#[snafu(visibility(pub))]
30#[stack_trace_debug]
31pub enum Error {
32 #[snafu(display("Failed to create TLS Config"))]
33 TlsConfig {
34 #[snafu(implicit)]
35 location: Location,
36 source: common_wal::error::Error,
37 },
38
39 #[snafu(display("Invalid provider type, expected: {}, actual: {}", expected, actual))]
40 InvalidProvider {
41 #[snafu(implicit)]
42 location: Location,
43 expected: String,
44 actual: String,
45 },
46
47 #[snafu(display("Failed to start log store task: {}", name))]
48 StartWalTask {
49 name: String,
50 #[snafu(implicit)]
51 location: Location,
52 source: RuntimeError,
53 },
54
55 #[snafu(display("Failed to stop log store task: {}", name))]
56 StopWalTask {
57 name: String,
58 #[snafu(implicit)]
59 location: Location,
60 source: RuntimeError,
61 },
62
63 #[snafu(display("Failed to add entry to LogBatch"))]
64 AddEntryLogBatch {
65 #[snafu(source)]
66 error: raft_engine::Error,
67 #[snafu(implicit)]
68 location: Location,
69 },
70
71 #[snafu(display("Failed to perform raft-engine operation"))]
72 RaftEngine {
73 #[snafu(source)]
74 error: raft_engine::Error,
75 #[snafu(implicit)]
76 location: Location,
77 },
78
79 #[snafu(display("Failed to perform IO on path: {}", path))]
80 Io {
81 path: String,
82 #[snafu(source)]
83 error: std::io::Error,
84 #[snafu(implicit)]
85 location: Location,
86 },
87
88 #[snafu(display("Log store not started yet"))]
89 IllegalState {
90 #[snafu(implicit)]
91 location: Location,
92 },
93
94 #[snafu(display("Namespace is illegal: {}", ns))]
95 IllegalNamespace {
96 ns: u64,
97 #[snafu(implicit)]
98 location: Location,
99 },
100
101 #[snafu(display(
102 "Failed to fetch entries from namespace: {}, start: {}, end: {}, max size: {}",
103 ns,
104 start,
105 end,
106 max_size,
107 ))]
108 FetchEntry {
109 ns: u64,
110 start: u64,
111 end: u64,
112 max_size: usize,
113 #[snafu(source)]
114 error: raft_engine::Error,
115 #[snafu(implicit)]
116 location: Location,
117 },
118
119 #[snafu(display(
120 "Cannot override compacted entry, namespace: {}, first index: {}, attempt index: {}",
121 namespace,
122 first_index,
123 attempt_index
124 ))]
125 OverrideCompactedEntry {
126 namespace: u64,
127 first_index: u64,
128 attempt_index: u64,
129 #[snafu(implicit)]
130 location: Location,
131 },
132
133 #[snafu(display(
134 "Failed to build a Kafka client, broker endpoints: {:?}",
135 broker_endpoints
136 ))]
137 BuildClient {
138 broker_endpoints: Vec<String>,
139 #[snafu(implicit)]
140 location: Location,
141 #[snafu(source)]
142 error: rskafka::client::error::Error,
143 },
144
145 #[snafu(display(
146 "Failed to build a Kafka partition client, topic: {}, partition: {}",
147 topic,
148 partition
149 ))]
150 BuildPartitionClient {
151 topic: String,
152 partition: i32,
153 #[snafu(implicit)]
154 location: Location,
155 #[snafu(source)]
156 error: rskafka::client::error::Error,
157 },
158
159 #[snafu(display("Missing required key in a record"))]
160 MissingKey {
161 #[snafu(implicit)]
162 location: Location,
163 },
164
165 #[snafu(display("Missing required value in a record"))]
166 MissingValue {
167 #[snafu(implicit)]
168 location: Location,
169 },
170
171 #[snafu(display("Failed to produce records to Kafka, topic: {}, size: {}", topic, size))]
172 ProduceRecord {
173 topic: String,
174 size: usize,
175 #[snafu(implicit)]
176 location: Location,
177 #[snafu(source)]
178 error: rskafka::client::producer::Error,
179 },
180
181 #[snafu(display("Failed to produce batch records to Kafka"))]
182 BatchProduce {
183 #[snafu(implicit)]
184 location: Location,
185 #[snafu(source)]
186 error: rskafka::client::error::Error,
187 },
188
189 #[snafu(display("Failed to read a record from Kafka, topic: {}", topic))]
190 ConsumeRecord {
191 topic: String,
192 #[snafu(implicit)]
193 location: Location,
194 #[snafu(source)]
195 error: rskafka::client::error::Error,
196 },
197
198 #[snafu(display("Failed to get the latest offset, topic: {}", topic))]
199 GetOffset {
200 topic: String,
201 #[snafu(implicit)]
202 location: Location,
203 #[snafu(source)]
204 error: rskafka::client::error::Error,
205 },
206
207 #[snafu(display("Failed to do a cast"))]
208 Cast {
209 #[snafu(implicit)]
210 location: Location,
211 },
212
213 #[snafu(display("Failed to encode object into json"))]
214 EncodeJson {
215 #[snafu(implicit)]
216 location: Location,
217 #[snafu(source)]
218 error: JsonError,
219 },
220
221 #[snafu(display("Failed to decode object from json"))]
222 DecodeJson {
223 #[snafu(implicit)]
224 location: Location,
225 #[snafu(source)]
226 error: JsonError,
227 },
228
229 #[snafu(display("The record sequence is not legal, error: {}", error))]
230 IllegalSequence {
231 #[snafu(implicit)]
232 location: Location,
233 error: String,
234 },
235
236 #[snafu(display(
237 "Attempt to append discontinuous log entry, region: {}, last index: {}, attempt index: {}",
238 region_id,
239 last_index,
240 attempt_index
241 ))]
242 DiscontinuousLogIndex {
243 region_id: RegionId,
244 last_index: u64,
245 attempt_index: u64,
246 },
247
248 #[snafu(display("OrderedBatchProducer is stopped",))]
249 OrderedBatchProducerStopped {
250 #[snafu(implicit)]
251 location: Location,
252 },
253
254 #[snafu(display("Failed to wait for ProduceResultReceiver"))]
255 WaitProduceResultReceiver {
256 #[snafu(implicit)]
257 location: Location,
258 #[snafu(source)]
259 error: tokio::sync::oneshot::error::RecvError,
260 },
261
262 #[snafu(display("Failed to wait for result of DumpIndex"))]
263 WaitDumpIndex {
264 #[snafu(implicit)]
265 location: Location,
266 #[snafu(source)]
267 error: tokio::sync::oneshot::error::RecvError,
268 },
269
270 #[snafu(display("Failed to create writer"))]
271 CreateWriter {
272 #[snafu(implicit)]
273 location: Location,
274 #[snafu(source)]
275 error: object_store::Error,
276 },
277
278 #[snafu(display("Failed to write index"))]
279 WriteIndex {
280 #[snafu(implicit)]
281 location: Location,
282 #[snafu(source)]
283 error: object_store::Error,
284 },
285
286 #[snafu(display("Failed to read index, path: {path}"))]
287 ReadIndex {
288 #[snafu(implicit)]
289 location: Location,
290 #[snafu(source)]
291 error: object_store::Error,
292 path: String,
293 },
294
295 #[snafu(display(
296 "The length of meta if exceeded the limit: {}, actual: {}",
297 limit,
298 actual
299 ))]
300 MetaLengthExceededLimit {
301 #[snafu(implicit)]
302 location: Location,
303 limit: usize,
304 actual: usize,
305 },
306
307 #[snafu(display("No max value"))]
308 NoMaxValue {
309 #[snafu(implicit)]
310 location: Location,
311 },
312
313 #[snafu(display("Corrupted WAL object, {}", reason))]
314 CorruptedWalObject {
315 reason: String,
316 #[snafu(implicit)]
317 location: Location,
318 },
319
320 #[snafu(display("Invalid WAL object store, {}", reason))]
321 InvalidWalObjectStore {
322 reason: String,
323 #[snafu(implicit)]
324 location: Location,
325 },
326
327 #[snafu(display(
328 "Object store WAL region mismatch, supplied region: {}, {}",
329 region_id,
330 reason
331 ))]
332 MismatchedWalRegion {
333 region_id: RegionId,
334 reason: String,
335 #[snafu(implicit)]
336 location: Location,
337 },
338
339 #[snafu(display(
340 "Invalid WAL entry range, region: {}, start: {}, end: {}",
341 region_id,
342 start_entry_id,
343 end_entry_id
344 ))]
345 InvalidWalEntryRange {
346 region_id: RegionId,
347 start_entry_id: u64,
348 end_entry_id: u64,
349 #[snafu(implicit)]
350 location: Location,
351 },
352
353 #[snafu(display("Incomplete multipart WAL entry of region {}", region_id))]
354 IncompleteWalEntry {
355 region_id: RegionId,
356 #[snafu(implicit)]
357 location: Location,
358 },
359
360 #[snafu(display("WAL object sequence is exhausted, last sequence: {}", last_object_seq))]
361 WalObjectSequenceExhausted {
362 last_object_seq: u64,
363 #[snafu(implicit)]
364 location: Location,
365 },
366
367 #[snafu(display(
368 "WAL object sequence {} is not settled: the open batch has assigned entry ids under it",
369 object_seq
370 ))]
371 WalObjectSequenceUnsettled {
372 object_seq: u64,
373 #[snafu(implicit)]
374 location: Location,
375 },
376
377 #[snafu(display(
378 "WAL entry positions of region {} in one object are exhausted",
379 region_id
380 ))]
381 WalEntryPositionExhausted {
382 region_id: RegionId,
383 #[snafu(implicit)]
384 location: Location,
385 },
386
387 #[snafu(display("WAL object already exists with different content, path: {}", path))]
388 WalObjectConflict {
389 path: String,
390 #[snafu(implicit)]
391 location: Location,
392 },
393
394 #[snafu(display(
395 "WAL object {} was written by the earlier epoch {}, this store writes epoch {}",
396 path,
397 existing_epoch,
398 epoch
399 ))]
400 StaleWalObject {
401 path: String,
402 existing_epoch: u64,
403 epoch: u64,
404 #[snafu(implicit)]
405 location: Location,
406 },
407
408 #[snafu(display(
409 "WAL object {} that starts epoch {} was already present, so the open cannot tell whether it wrote it",
410 path,
411 epoch
412 ))]
413 UnconfirmedWalEpochStart {
414 path: String,
415 epoch: u64,
416 #[snafu(implicit)]
417 location: Location,
418 },
419
420 #[snafu(display("Failed to {} WAL object, path: {}", operation, path))]
421 WalObjectStore {
422 operation: &'static str,
423 path: String,
424 #[snafu(source)]
425 error: object_store::Error,
426 #[snafu(implicit)]
427 location: Location,
428 },
429
430 #[snafu(display("Invalid WAL object, path: {}", path))]
431 InvalidWalObject {
432 path: String,
433 #[snafu(source(from(Error, Box::new)))]
434 source: Box<Error>,
435 #[snafu(implicit)]
436 location: Location,
437 },
438
439 #[snafu(display(
440 "Object store WAL prefix mismatch, expected: {}, actual: {}",
441 expected,
442 actual
443 ))]
444 MismatchedWalPrefix {
445 expected: String,
446 actual: String,
447 #[snafu(implicit)]
448 location: Location,
449 },
450
451 #[snafu(display("Object store WAL log store is stopped"))]
452 ObjectStoreWalStopped {
453 #[snafu(implicit)]
454 location: Location,
455 },
456
457 #[snafu(display("Object store WAL operation failed"))]
458 ObjectStoreWal {
459 source: Arc<Error>,
460 #[snafu(implicit)]
461 location: Location,
462 },
463}
464
465pub type Result<T> = std::result::Result<T, Error>;
466
467fn rskafka_client_error_to_status_code(error: &rskafka::client::error::Error) -> StatusCode {
468 match error {
469 rskafka::client::error::Error::Connection(_)
470 | rskafka::client::error::Error::Request(_)
471 | rskafka::client::error::Error::InvalidResponse(_)
472 | rskafka::client::error::Error::ServerError { .. }
473 | rskafka::client::error::Error::RetryFailed(_) => StatusCode::Internal,
474 rskafka::client::error::Error::Timeout => StatusCode::StorageUnavailable,
475 _ => StatusCode::Internal,
476 }
477}
478
479impl ErrorExt for Error {
480 fn as_any(&self) -> &dyn Any {
481 self
482 }
483
484 fn status_code(&self) -> StatusCode {
485 use Error::*;
486
487 match self {
488 TlsConfig { .. }
489 | InvalidProvider { .. }
490 | IllegalNamespace { .. }
491 | MissingKey { .. }
492 | MissingValue { .. }
493 | OverrideCompactedEntry { .. }
494 | InvalidWalObjectStore { .. }
495 | MismatchedWalPrefix { .. }
496 | MismatchedWalRegion { .. }
497 | IncompleteWalEntry { .. }
498 | InvalidWalEntryRange { .. } => StatusCode::InvalidArguments,
499 StartWalTask { .. }
500 | StopWalTask { .. }
501 | IllegalState { .. }
502 | NoMaxValue { .. }
503 | Cast { .. }
504 | EncodeJson { .. }
505 | DecodeJson { .. }
506 | IllegalSequence { .. }
507 | DiscontinuousLogIndex { .. }
508 | OrderedBatchProducerStopped { .. }
509 | WaitProduceResultReceiver { .. }
510 | WaitDumpIndex { .. }
511 | MetaLengthExceededLimit { .. }
512 | ObjectStoreWalStopped { .. } => StatusCode::Internal,
513
514 CorruptedWalObject { .. }
515 | WalObjectConflict { .. }
516 | WalObjectSequenceExhausted { .. }
517 | WalEntryPositionExhausted { .. } => StatusCode::Unexpected,
518 WalObjectSequenceUnsettled { .. } => StatusCode::IllegalState,
519
520 InvalidWalObject { source, .. } => source.status_code(),
521 ObjectStoreWal { source, .. } => source.status_code(),
522
523 CreateWriter { .. }
525 | WriteIndex { .. }
526 | ReadIndex { .. }
527 | WalObjectStore { .. }
528 | StaleWalObject { .. }
529 | UnconfirmedWalEpochStart { .. }
530 | Io { .. } => StatusCode::StorageUnavailable,
531 FetchEntry { .. } | RaftEngine { .. } | AddEntryLogBatch { .. } => {
533 StatusCode::StorageUnavailable
534 }
535 ProduceRecord { error, .. } => match error {
537 rskafka::client::producer::Error::Client(error) => {
538 rskafka_client_error_to_status_code(error)
539 }
540 rskafka::client::producer::Error::Aggregator(_)
541 | rskafka::client::producer::Error::FlushError(_)
542 | rskafka::client::producer::Error::TooLarge => StatusCode::Internal,
543 },
544 BuildClient { error, .. }
545 | BuildPartitionClient { error, .. }
546 | BatchProduce { error, .. }
547 | GetOffset { error, .. }
548 | ConsumeRecord { error, .. } => rskafka_client_error_to_status_code(error),
549 }
550 }
551
552 fn retry_hint(&self) -> RetryHint {
553 use Error::*;
554
555 match self {
556 CreateWriter { error, .. }
557 | WriteIndex { error, .. }
558 | ReadIndex { error, .. }
559 | WalObjectStore { error, .. } => retry_hint_from_opendal_error(error),
560 ObjectStoreWal { source, .. } => source.retry_hint(),
561 Io { error, .. } => retry_hint_from_io_error(error),
562 FetchEntry { .. }
563 | RaftEngine { .. }
564 | AddEntryLogBatch { .. }
565 | WalObjectSequenceUnsettled { .. }
566 | StaleWalObject { .. }
567 | UnconfirmedWalEpochStart { .. } => RetryHint::Retryable,
568 ProduceRecord { error, .. } => match error {
569 rskafka::client::producer::Error::Client(error) => {
570 rskafka_client_error_to_retry_hint(error)
571 }
572 rskafka::client::producer::Error::Aggregator(_)
573 | rskafka::client::producer::Error::FlushError(_)
574 | rskafka::client::producer::Error::TooLarge => RetryHint::NonRetryable,
575 },
576 BuildClient { error, .. }
577 | BuildPartitionClient { error, .. }
578 | BatchProduce { error, .. }
579 | GetOffset { error, .. }
580 | ConsumeRecord { error, .. } => rskafka_client_error_to_retry_hint(error),
581 _ => RetryHint::NonRetryable,
582 }
583 }
584}