Skip to main content

log_store/
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
15use 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            // Object store related errors
524            CreateWriter { .. }
525            | WriteIndex { .. }
526            | ReadIndex { .. }
527            | WalObjectStore { .. }
528            | StaleWalObject { .. }
529            | UnconfirmedWalEpochStart { .. }
530            | Io { .. } => StatusCode::StorageUnavailable,
531            // Raft engine
532            FetchEntry { .. } | RaftEngine { .. } | AddEntryLogBatch { .. } => {
533                StatusCode::StorageUnavailable
534            }
535            // Kafka producer related errors
536            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}