Skip to main content

common_meta/
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::str::Utf8Error;
16use std::sync::Arc;
17
18use common_error::ext::{BoxedError, ErrorExt, RetryHint};
19use common_error::status_code::StatusCode;
20use common_macro::stack_trace_debug;
21use common_procedure::ProcedureId;
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;
27use table::metadata::TableId;
28
29use crate::DatanodeId;
30use crate::peer::Peer;
31
32mod retry_hint;
33
34pub use retry_hint::retry_hint_from_etcd_error;
35#[cfg(feature = "mysql_kvbackend")]
36pub use retry_hint::retry_hint_from_sqlx_error;
37#[cfg(feature = "pg_kvbackend")]
38pub use retry_hint::{retry_hint_from_postgres_error, retry_hint_from_postgres_pool_error};
39
40#[derive(Snafu)]
41#[snafu(visibility(pub))]
42#[stack_trace_debug]
43pub enum Error {
44    #[snafu(display("Empty key is not allowed"))]
45    EmptyKey {
46        #[snafu(implicit)]
47        location: Location,
48    },
49
50    #[snafu(display(
51        "Another procedure is operating the region: {} on peer: {}",
52        region_id,
53        peer_id
54    ))]
55    RegionOperatingRace {
56        #[snafu(implicit)]
57        location: Location,
58        peer_id: DatanodeId,
59        region_id: RegionId,
60    },
61
62    #[snafu(display("Failed to connect to Etcd"))]
63    ConnectEtcd {
64        #[snafu(source)]
65        error: etcd_client::Error,
66        #[snafu(implicit)]
67        location: Location,
68    },
69
70    #[snafu(display("Failed to execute via Etcd"))]
71    EtcdFailed {
72        #[snafu(source)]
73        error: etcd_client::Error,
74        #[snafu(implicit)]
75        location: Location,
76    },
77
78    #[snafu(display("Failed to execute {} txn operations via Etcd", max_operations))]
79    EtcdTxnFailed {
80        max_operations: usize,
81        #[snafu(source)]
82        error: etcd_client::Error,
83        #[snafu(implicit)]
84        location: Location,
85    },
86
87    #[snafu(display("Failed to get sequence: {}", err_msg))]
88    NextSequence {
89        err_msg: String,
90        #[snafu(implicit)]
91        location: Location,
92    },
93
94    #[snafu(display("Unexpected sequence value: {}", err_msg))]
95    UnexpectedSequenceValue {
96        err_msg: String,
97        #[snafu(implicit)]
98        location: Location,
99    },
100
101    #[snafu(display("Table info not found: {}", table))]
102    TableInfoNotFound {
103        table: String,
104        #[snafu(implicit)]
105        location: Location,
106    },
107
108    #[snafu(display("Failed to register procedure loader, type name: {}", type_name))]
109    RegisterProcedureLoader {
110        type_name: String,
111        #[snafu(implicit)]
112        location: Location,
113        source: common_procedure::error::Error,
114    },
115
116    #[snafu(display("Failed to register repartition procedure loader"))]
117    RegisterRepartitionProcedureLoader {
118        #[snafu(implicit)]
119        location: Location,
120        source: BoxedError,
121    },
122
123    #[snafu(display("Failed to create repartition procedure"))]
124    CreateRepartitionProcedure {
125        source: BoxedError,
126        #[snafu(implicit)]
127        location: Location,
128    },
129
130    #[snafu(display("Failed to persist the repartition GC requirement"))]
131    PersistRepartitionGcRequirement {
132        source: BoxedError,
133        #[snafu(implicit)]
134        location: Location,
135    },
136
137    #[snafu(display("Failed to submit procedure"))]
138    SubmitProcedure {
139        #[snafu(implicit)]
140        location: Location,
141        source: common_procedure::Error,
142    },
143
144    #[snafu(display("Failed to query procedure"))]
145    QueryProcedure {
146        #[snafu(implicit)]
147        location: Location,
148        source: common_procedure::Error,
149    },
150
151    #[snafu(display("Procedure not found: {pid}"))]
152    ProcedureNotFound {
153        #[snafu(implicit)]
154        location: Location,
155        pid: String,
156    },
157
158    #[snafu(display("Failed to parse procedure id: {key}"))]
159    ParseProcedureId {
160        #[snafu(implicit)]
161        location: Location,
162        key: String,
163        #[snafu(source)]
164        error: common_procedure::ParseIdError,
165    },
166
167    #[snafu(display("Unsupported operation {}", operation))]
168    Unsupported {
169        operation: String,
170        #[snafu(implicit)]
171        location: Location,
172    },
173
174    #[snafu(display("Trying to write to a read-only kv backend: {}", name))]
175    ReadOnlyKvBackend {
176        name: String,
177        #[snafu(implicit)]
178        location: Location,
179    },
180
181    #[snafu(display("Failed to get procedure state receiver, procedure id: {procedure_id}"))]
182    ProcedureStateReceiver {
183        procedure_id: ProcedureId,
184        #[snafu(implicit)]
185        location: Location,
186        source: common_procedure::Error,
187    },
188
189    #[snafu(display("Procedure state receiver not found: {procedure_id}"))]
190    ProcedureStateReceiverNotFound {
191        procedure_id: ProcedureId,
192        #[snafu(implicit)]
193        location: Location,
194    },
195
196    #[snafu(display("Failed to wait procedure done"))]
197    WaitProcedure {
198        #[snafu(implicit)]
199        location: Location,
200        source: common_procedure::Error,
201    },
202
203    #[snafu(display("Failed to start procedure manager"))]
204    StartProcedureManager {
205        #[snafu(implicit)]
206        location: Location,
207        source: common_procedure::Error,
208    },
209
210    #[snafu(display("Failed to stop procedure manager"))]
211    StopProcedureManager {
212        #[snafu(implicit)]
213        location: Location,
214        source: common_procedure::Error,
215    },
216
217    #[snafu(display(
218        "Failed to get procedure output, procedure id: {procedure_id}, error: {err_msg}"
219    ))]
220    ProcedureOutput {
221        procedure_id: String,
222        err_msg: String,
223        #[snafu(implicit)]
224        location: Location,
225    },
226
227    #[snafu(display("Primary key '{key}' not found when creating region request"))]
228    PrimaryKeyNotFound {
229        key: String,
230        #[snafu(implicit)]
231        location: Location,
232    },
233
234    #[snafu(display("Failed to build table meta for table: {}", table_name))]
235    BuildTableMeta {
236        table_name: String,
237        #[snafu(source)]
238        error: table::metadata::TableMetaBuilderError,
239        #[snafu(implicit)]
240        location: Location,
241    },
242
243    #[snafu(display("Table occurs error"))]
244    Table {
245        #[snafu(implicit)]
246        location: Location,
247        source: table::error::Error,
248    },
249
250    #[snafu(display("Failed to find table route for table id {}", table_id))]
251    TableRouteNotFound {
252        table_id: TableId,
253        #[snafu(implicit)]
254        location: Location,
255    },
256
257    #[snafu(display("Failed to find table repartition metadata for table id {}", table_id))]
258    TableRepartNotFound {
259        table_id: TableId,
260        #[snafu(implicit)]
261        location: Location,
262    },
263
264    #[snafu(display("Failed to decode protobuf"))]
265    DecodeProto {
266        #[snafu(implicit)]
267        location: Location,
268        #[snafu(source)]
269        error: prost::DecodeError,
270    },
271
272    #[snafu(display("Failed to decode packed file references"))]
273    DecodePackedFileRefs {
274        #[snafu(implicit)]
275        location: Location,
276        #[snafu(source)]
277        error: base64::DecodeError,
278    },
279
280    #[snafu(display("Invalid packed file references framing"))]
281    InvalidPackedFileRefs {
282        #[snafu(implicit)]
283        location: Location,
284    },
285
286    #[snafu(display("Failed to encode object into json"))]
287    EncodeJson {
288        #[snafu(implicit)]
289        location: Location,
290        #[snafu(source)]
291        error: JsonError,
292    },
293
294    #[snafu(display("Failed to decode object from json"))]
295    DecodeJson {
296        #[snafu(implicit)]
297        location: Location,
298        #[snafu(source)]
299        error: JsonError,
300    },
301
302    #[snafu(display("Failed to serialize to json: {}", input))]
303    SerializeToJson {
304        input: String,
305        #[snafu(source)]
306        error: serde_json::error::Error,
307        #[snafu(implicit)]
308        location: Location,
309    },
310
311    #[snafu(display("Failed to deserialize from json: {}", input))]
312    DeserializeFromJson {
313        input: String,
314        #[snafu(source)]
315        error: serde_json::error::Error,
316        #[snafu(implicit)]
317        location: Location,
318    },
319
320    #[snafu(display("Payload not exist"))]
321    PayloadNotExist {
322        #[snafu(implicit)]
323        location: Location,
324    },
325
326    #[snafu(display("Failed to serde json"))]
327    SerdeJson {
328        #[snafu(source)]
329        error: serde_json::error::Error,
330        #[snafu(implicit)]
331        location: Location,
332    },
333
334    #[snafu(display("Failed to parse value {} into key {}", value, key))]
335    ParseOption {
336        key: String,
337        value: String,
338        #[snafu(implicit)]
339        location: Location,
340    },
341
342    #[snafu(display(
343        "Conflicting schema options: {}={} and {}={}",
344        first_key,
345        first_value,
346        second_key,
347        second_value
348    ))]
349    ConflictingSchemaOptions {
350        first_key: String,
351        first_value: String,
352        second_key: String,
353        second_value: String,
354        #[snafu(implicit)]
355        location: Location,
356    },
357
358    #[snafu(display("Illegal state from server, code: {}, error: {}", code, err_msg))]
359    IllegalServerState {
360        code: i32,
361        err_msg: String,
362        #[snafu(implicit)]
363        location: Location,
364    },
365
366    #[snafu(display("Failed to convert alter table request"))]
367    ConvertAlterTableRequest {
368        source: common_grpc_expr::error::Error,
369        #[snafu(implicit)]
370        location: Location,
371    },
372
373    #[snafu(display("Invalid protobuf message: {err_msg}"))]
374    InvalidProtoMsg {
375        err_msg: String,
376        #[snafu(implicit)]
377        location: Location,
378    },
379
380    #[snafu(display("Unexpected: {err_msg}"))]
381    Unexpected {
382        err_msg: String,
383        #[snafu(implicit)]
384        location: Location,
385    },
386
387    #[snafu(display("Metasrv election has no leader at this moment"))]
388    ElectionNoLeader {
389        #[snafu(implicit)]
390        location: Location,
391    },
392
393    #[snafu(display("Metasrv election leader lease expired"))]
394    ElectionLeaderLeaseExpired {
395        #[snafu(implicit)]
396        location: Location,
397    },
398
399    #[snafu(display("Metasrv election leader lease changed during election"))]
400    ElectionLeaderLeaseChanged {
401        #[snafu(implicit)]
402        location: Location,
403    },
404
405    #[snafu(display("Table already exists, table: {}", table_name))]
406    TableAlreadyExists {
407        table_name: String,
408        #[snafu(implicit)]
409        location: Location,
410    },
411
412    #[snafu(display(
413        "Cannot drop table '{}': an older tombstone already uses the same full name",
414        table_name
415    ))]
416    /// Raised when a live table is recreated with a name still reserved by an older tombstone.
417    TableNameTombstoneConflict {
418        table_name: String,
419        existing_table_id: TableId,
420        dropping_table_id: TableId,
421        #[snafu(implicit)]
422        location: Location,
423    },
424
425    #[snafu(display("View already exists, view: {}", view_name))]
426    ViewAlreadyExists {
427        view_name: String,
428        #[snafu(implicit)]
429        location: Location,
430    },
431
432    #[snafu(display("Flow already exists: {}", flow_name))]
433    FlowAlreadyExists {
434        flow_name: String,
435        #[snafu(implicit)]
436        location: Location,
437    },
438
439    #[snafu(display("Schema already exists, catalog:{}, schema: {}", catalog, schema))]
440    SchemaAlreadyExists {
441        catalog: String,
442        schema: String,
443        #[snafu(implicit)]
444        location: Location,
445    },
446
447    #[snafu(display("Failed to convert raw key to str"))]
448    ConvertRawKey {
449        #[snafu(implicit)]
450        location: Location,
451        #[snafu(source)]
452        error: Utf8Error,
453    },
454
455    #[snafu(display("Table not found: '{}'", table_name))]
456    TableNotFound {
457        table_name: String,
458        #[snafu(implicit)]
459        location: Location,
460    },
461
462    #[snafu(display("Region not found: {}", region_id))]
463    RegionNotFound {
464        region_id: RegionId,
465        #[snafu(implicit)]
466        location: Location,
467    },
468
469    #[snafu(display("View not found: '{}'", view_name))]
470    ViewNotFound {
471        view_name: String,
472        #[snafu(implicit)]
473        location: Location,
474    },
475
476    #[snafu(display("Flow not found: '{}'", flow_name))]
477    FlowNotFound {
478        flow_name: String,
479        #[snafu(implicit)]
480        location: Location,
481    },
482
483    #[snafu(display("Flow route not found: '{}'", flow_name))]
484    FlowRouteNotFound {
485        flow_name: String,
486        #[snafu(implicit)]
487        location: Location,
488    },
489
490    #[snafu(display("Schema nod found, schema: {}", table_schema))]
491    SchemaNotFound {
492        table_schema: String,
493        #[snafu(implicit)]
494        location: Location,
495    },
496
497    #[snafu(display("Catalog not found, catalog: {}", catalog))]
498    CatalogNotFound {
499        catalog: String,
500        #[snafu(implicit)]
501        location: Location,
502    },
503
504    #[snafu(display("Invalid metadata, err: {}", err_msg))]
505    InvalidMetadata {
506        err_msg: String,
507        #[snafu(implicit)]
508        location: Location,
509    },
510
511    #[snafu(display("Invalid view info, err: {}", err_msg))]
512    InvalidViewInfo {
513        err_msg: String,
514        #[snafu(implicit)]
515        location: Location,
516    },
517
518    #[snafu(display("Invalid flow request body: {:?}", body))]
519    InvalidFlowRequestBody {
520        body: Box<Option<api::v1::flow::flow_request::Body>>,
521        #[snafu(implicit)]
522        location: Location,
523    },
524
525    #[snafu(display("Failed to get kv cache, err: {}", err_msg))]
526    GetKvCache { err_msg: String },
527
528    #[snafu(display("Get null from cache, key: {}", key))]
529    CacheNotGet {
530        key: String,
531        #[snafu(implicit)]
532        location: Location,
533    },
534
535    #[snafu(display("Etcd txn error: {err_msg}"))]
536    EtcdTxnOpResponse {
537        err_msg: String,
538        #[snafu(implicit)]
539        location: Location,
540    },
541
542    #[snafu(display("External error"))]
543    External {
544        #[snafu(implicit)]
545        location: Location,
546        source: BoxedError,
547    },
548
549    #[snafu(display("The response exceeded size limit"))]
550    ResponseExceededSizeLimit {
551        #[snafu(implicit)]
552        location: Location,
553        source: BoxedError,
554    },
555
556    #[snafu(display("Invalid heartbeat response"))]
557    InvalidHeartbeatResponse {
558        #[snafu(implicit)]
559        location: Location,
560    },
561
562    #[snafu(display("Failed to operate on datanode: {}", peer))]
563    OperateDatanode {
564        #[snafu(implicit)]
565        location: Location,
566        peer: Peer,
567        source: BoxedError,
568    },
569
570    #[snafu(display("Retry later"))]
571    RetryLater {
572        source: BoxedError,
573        clean_poisons: bool,
574    },
575
576    #[snafu(display("Abort procedure"))]
577    AbortProcedure {
578        #[snafu(implicit)]
579        location: Location,
580        source: BoxedError,
581        clean_poisons: bool,
582    },
583
584    #[snafu(display("Failed to serialize WAL options for region: {region_id}"))]
585    SerializeWalOptions {
586        region_id: RegionId,
587        #[snafu(source)]
588        error: serde_json::Error,
589        #[snafu(implicit)]
590        location: Location,
591    },
592
593    #[snafu(display("Invalid number of topics {}", num_topics))]
594    InvalidNumTopics {
595        num_topics: usize,
596        #[snafu(implicit)]
597        location: Location,
598    },
599
600    #[snafu(display(
601        "Failed to build a Kafka client, broker endpoints: {:?}",
602        broker_endpoints
603    ))]
604    BuildKafkaClient {
605        broker_endpoints: Vec<String>,
606        #[snafu(implicit)]
607        location: Location,
608        #[snafu(source)]
609        error: rskafka::client::error::Error,
610    },
611
612    #[snafu(display("Failed to create TLS Config"))]
613    TlsConfig {
614        #[snafu(implicit)]
615        location: Location,
616        source: common_wal::error::Error,
617    },
618
619    #[snafu(display("Failed to build a Kafka controller client"))]
620    BuildKafkaCtrlClient {
621        #[snafu(implicit)]
622        location: Location,
623        #[snafu(source)]
624        error: rskafka::client::error::Error,
625    },
626
627    #[snafu(display(
628        "Failed to get a Kafka partition client, topic: {}, partition: {}",
629        topic,
630        partition
631    ))]
632    KafkaPartitionClient {
633        topic: String,
634        partition: i32,
635        #[snafu(implicit)]
636        location: Location,
637        #[snafu(source)]
638        error: rskafka::client::error::Error,
639    },
640
641    #[snafu(display(
642        "Failed to get offset from Kafka, topic: {}, partition: {}",
643        topic,
644        partition
645    ))]
646    KafkaGetOffset {
647        topic: String,
648        partition: i32,
649        #[snafu(implicit)]
650        location: Location,
651        #[snafu(source)]
652        error: rskafka::client::error::Error,
653    },
654
655    #[snafu(display("Failed to produce records to Kafka, topic: {}", topic))]
656    ProduceRecord {
657        topic: String,
658        #[snafu(implicit)]
659        location: Location,
660        #[snafu(source)]
661        error: rskafka::client::error::Error,
662    },
663
664    #[snafu(display("Failed to create a Kafka wal topic"))]
665    CreateKafkaWalTopic {
666        #[snafu(implicit)]
667        location: Location,
668        #[snafu(source)]
669        error: rskafka::client::error::Error,
670    },
671
672    #[snafu(display("The topic pool is empty"))]
673    EmptyTopicPool {
674        #[snafu(implicit)]
675        location: Location,
676    },
677
678    #[snafu(display("Unexpected table route type: {}", err_msg))]
679    UnexpectedLogicalRouteTable {
680        #[snafu(implicit)]
681        location: Location,
682        err_msg: String,
683    },
684
685    #[snafu(display("The tasks of {} cannot be empty", name))]
686    EmptyDdlTasks {
687        name: String,
688        #[snafu(implicit)]
689        location: Location,
690    },
691
692    #[snafu(display("Metadata corruption: {}", err_msg))]
693    MetadataCorruption {
694        err_msg: String,
695        #[snafu(implicit)]
696        location: Location,
697    },
698
699    #[snafu(display("Alter logical tables invalid arguments: {}", err_msg))]
700    AlterLogicalTablesInvalidArguments {
701        err_msg: String,
702        #[snafu(implicit)]
703        location: Location,
704    },
705
706    #[snafu(display("Create logical tables invalid arguments: {}", err_msg))]
707    CreateLogicalTablesInvalidArguments {
708        err_msg: String,
709        #[snafu(implicit)]
710        location: Location,
711    },
712
713    #[snafu(display("Invalid node info key: {}", key))]
714    InvalidNodeInfoKey {
715        key: String,
716        #[snafu(implicit)]
717        location: Location,
718    },
719
720    #[snafu(display("Invalid node stat key: {}", key))]
721    InvalidStatKey {
722        key: String,
723        #[snafu(implicit)]
724        location: Location,
725    },
726
727    #[snafu(display("Failed to parse number: {}", err_msg))]
728    ParseNum {
729        err_msg: String,
730        #[snafu(source)]
731        error: std::num::ParseIntError,
732        #[snafu(implicit)]
733        location: Location,
734    },
735
736    #[snafu(display("Invalid role: {}", role))]
737    InvalidRole {
738        role: i32,
739        #[snafu(implicit)]
740        location: Location,
741    },
742
743    #[snafu(display("Invalid set database option, key: {}, value: {}", key, value))]
744    InvalidSetDatabaseOption {
745        key: String,
746        value: String,
747        #[snafu(implicit)]
748        location: Location,
749    },
750
751    #[snafu(display("Invalid unset database option, key: {}", key))]
752    InvalidUnsetDatabaseOption {
753        key: String,
754        #[snafu(implicit)]
755        location: Location,
756    },
757
758    #[snafu(display("Invalid prefix: {}, key: {}", prefix, key))]
759    MismatchPrefix {
760        prefix: String,
761        key: String,
762        #[snafu(implicit)]
763        location: Location,
764    },
765
766    #[snafu(display("Failed to move values: {err_msg}"))]
767    MoveValues {
768        err_msg: String,
769        #[snafu(implicit)]
770        location: Location,
771    },
772
773    #[snafu(display("Failed to restore tombstone, target key already exists: {key}"))]
774    TombstoneTargetAlreadyExists {
775        key: String,
776        #[snafu(implicit)]
777        location: Location,
778    },
779
780    #[snafu(display("Failed to parse {} from utf8", name))]
781    FromUtf8 {
782        name: String,
783        #[snafu(source)]
784        error: std::string::FromUtf8Error,
785        #[snafu(implicit)]
786        location: Location,
787    },
788
789    #[snafu(display("Value not exists"))]
790    ValueNotExist {
791        #[snafu(implicit)]
792        location: Location,
793    },
794
795    #[snafu(display("Failed to get cache"))]
796    GetCache { source: Arc<Error> },
797
798    #[snafu(display(
799        "Failed to get latest cache value after {} attempts due to concurrent invalidation",
800        attempts
801    ))]
802    GetLatestCacheRetryExceeded {
803        attempts: usize,
804        #[snafu(implicit)]
805        location: Location,
806    },
807
808    #[cfg(feature = "pg_kvbackend")]
809    #[snafu(display("Failed to execute via Postgres, sql: {}", sql))]
810    PostgresExecution {
811        sql: String,
812        #[snafu(source)]
813        error: tokio_postgres::Error,
814        #[snafu(implicit)]
815        location: Location,
816    },
817
818    #[cfg(feature = "pg_kvbackend")]
819    #[snafu(display("Failed to create connection pool for Postgres"))]
820    CreatePostgresPool {
821        #[snafu(source)]
822        error: deadpool_postgres::CreatePoolError,
823        #[snafu(implicit)]
824        location: Location,
825    },
826
827    #[cfg(feature = "pg_kvbackend")]
828    #[snafu(display("Failed to get Postgres connection from pool: {}", reason))]
829    GetPostgresConnection {
830        reason: String,
831        #[snafu(implicit)]
832        location: Location,
833    },
834
835    #[cfg(feature = "pg_kvbackend")]
836    #[snafu(display("Failed to get Postgres client"))]
837    GetPostgresClient {
838        #[snafu(source)]
839        error: deadpool::managed::PoolError<tokio_postgres::Error>,
840        #[snafu(implicit)]
841        location: Location,
842    },
843
844    #[cfg(feature = "pg_kvbackend")]
845    #[snafu(display("Failed to {} Postgres transaction", operation))]
846    PostgresTransaction {
847        #[snafu(source)]
848        error: tokio_postgres::Error,
849        #[snafu(implicit)]
850        location: Location,
851        operation: String,
852    },
853
854    #[cfg(feature = "pg_kvbackend")]
855    #[snafu(display("Failed to setup PostgreSQL TLS configuration: {}", reason))]
856    PostgresTlsConfig {
857        reason: String,
858        #[snafu(implicit)]
859        location: Location,
860    },
861
862    #[snafu(display("Failed to load TLS certificate from path: {}", path))]
863    LoadTlsCertificate {
864        path: String,
865        #[snafu(source)]
866        error: std::io::Error,
867        #[snafu(implicit)]
868        location: Location,
869    },
870
871    #[cfg(feature = "pg_kvbackend")]
872    #[snafu(display("Invalid TLS configuration: {}", reason))]
873    InvalidTlsConfig {
874        reason: String,
875        #[snafu(implicit)]
876        location: Location,
877    },
878
879    #[cfg(feature = "mysql_kvbackend")]
880    #[snafu(display("Failed to execute via MySql, sql: {}", sql))]
881    MySqlExecution {
882        sql: String,
883        #[snafu(source)]
884        error: sqlx::Error,
885        #[snafu(implicit)]
886        location: Location,
887    },
888
889    #[cfg(feature = "mysql_kvbackend")]
890    #[snafu(display("Failed to create connection pool for MySql"))]
891    CreateMySqlPool {
892        #[snafu(source)]
893        error: sqlx::Error,
894        #[snafu(implicit)]
895        location: Location,
896    },
897
898    #[cfg(feature = "mysql_kvbackend")]
899    #[snafu(display("Failed to decode sql value"))]
900    DecodeSqlValue {
901        #[snafu(source)]
902        error: sqlx::error::Error,
903        #[snafu(implicit)]
904        location: Location,
905    },
906
907    #[cfg(feature = "mysql_kvbackend")]
908    #[snafu(display("Failed to acquire mysql client from pool"))]
909    AcquireMySqlClient {
910        #[snafu(source)]
911        error: sqlx::Error,
912        #[snafu(implicit)]
913        location: Location,
914    },
915
916    #[cfg(feature = "mysql_kvbackend")]
917    #[snafu(display("Failed to {} MySql transaction", operation))]
918    MySqlTransaction {
919        #[snafu(source)]
920        error: sqlx::Error,
921        #[snafu(implicit)]
922        location: Location,
923        operation: String,
924    },
925
926    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
927    #[snafu(display("Rds transaction retry failed"))]
928    RdsTransactionRetryFailed {
929        #[snafu(implicit)]
930        location: Location,
931    },
932
933    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
934    #[snafu(display("Sql execution timeout, sql: {}, duration: {:?}", sql, duration))]
935    SqlExecutionTimeout {
936        sql: String,
937        duration: std::time::Duration,
938        #[snafu(implicit)]
939        location: Location,
940    },
941
942    #[snafu(display(
943        "Datanode table info not found, table id: {}, datanode id: {}",
944        table_id,
945        datanode_id
946    ))]
947    DatanodeTableInfoNotFound {
948        datanode_id: DatanodeId,
949        table_id: TableId,
950        #[snafu(implicit)]
951        location: Location,
952    },
953
954    #[snafu(display("Invalid topic name prefix: {}", prefix))]
955    InvalidTopicNamePrefix {
956        prefix: String,
957        #[snafu(implicit)]
958        location: Location,
959    },
960
961    #[snafu(display("No leader found for table_id: {}", table_id))]
962    NoLeader {
963        table_id: TableId,
964        #[snafu(implicit)]
965        location: Location,
966    },
967
968    #[snafu(display(
969        "Procedure poison key already exists with a different value, key: {}, value: {}",
970        key,
971        value
972    ))]
973    ProcedurePoisonConflict {
974        key: String,
975        value: String,
976        #[snafu(implicit)]
977        location: Location,
978    },
979
980    #[snafu(display("Failed to put poison, table metadata may be corrupted"))]
981    PutPoison {
982        #[snafu(implicit)]
983        location: Location,
984        #[snafu(source)]
985        source: common_procedure::error::Error,
986    },
987
988    #[snafu(display("Invalid file path: {}", file_path))]
989    InvalidFilePath {
990        #[snafu(implicit)]
991        location: Location,
992        file_path: String,
993    },
994
995    #[snafu(display("Failed to serialize flexbuffers"))]
996    SerializeFlexbuffers {
997        #[snafu(implicit)]
998        location: Location,
999        #[snafu(source)]
1000        error: flexbuffers::SerializationError,
1001    },
1002
1003    #[snafu(display("Failed to deserialize flexbuffers"))]
1004    DeserializeFlexbuffers {
1005        #[snafu(implicit)]
1006        location: Location,
1007        #[snafu(source)]
1008        error: flexbuffers::DeserializationError,
1009    },
1010
1011    #[snafu(display("Failed to read flexbuffers"))]
1012    ReadFlexbuffers {
1013        #[snafu(implicit)]
1014        location: Location,
1015        #[snafu(source)]
1016        error: flexbuffers::ReaderError,
1017    },
1018
1019    #[snafu(display("Invalid file name: {}", reason))]
1020    InvalidFileName {
1021        #[snafu(implicit)]
1022        location: Location,
1023        reason: String,
1024    },
1025
1026    #[snafu(display("Invalid file extension: {}", reason))]
1027    InvalidFileExtension {
1028        #[snafu(implicit)]
1029        location: Location,
1030        reason: String,
1031    },
1032
1033    #[snafu(display("Failed to write object, file path: {}", file_path))]
1034    WriteObject {
1035        #[snafu(implicit)]
1036        location: Location,
1037        file_path: String,
1038        #[snafu(source)]
1039        error: object_store::Error,
1040    },
1041
1042    #[snafu(display("Failed to read object, file path: {}", file_path))]
1043    ReadObject {
1044        #[snafu(implicit)]
1045        location: Location,
1046        file_path: String,
1047        #[snafu(source)]
1048        error: object_store::Error,
1049    },
1050
1051    #[snafu(display("Missing column ids"))]
1052    MissingColumnIds {
1053        #[snafu(implicit)]
1054        location: Location,
1055    },
1056
1057    #[snafu(display(
1058        "Missing column in column metadata: {}, table: {}, table_id: {}",
1059        column_name,
1060        table_name,
1061        table_id,
1062    ))]
1063    MissingColumnInColumnMetadata {
1064        column_name: String,
1065        #[snafu(implicit)]
1066        location: Location,
1067        table_name: String,
1068        table_id: TableId,
1069    },
1070
1071    #[snafu(display(
1072        "Mismatch column id: column_name: {}, column_id: {}, table: {}, table_id: {}",
1073        column_name,
1074        column_id,
1075        table_name,
1076        table_id,
1077    ))]
1078    MismatchColumnId {
1079        column_name: String,
1080        column_id: u32,
1081        #[snafu(implicit)]
1082        location: Location,
1083        table_name: String,
1084        table_id: TableId,
1085    },
1086
1087    #[snafu(display("Failed to convert column def, column: {}", column))]
1088    ConvertColumnDef {
1089        column: String,
1090        #[snafu(implicit)]
1091        location: Location,
1092        source: api::error::Error,
1093    },
1094
1095    #[snafu(display("Failed to convert time ranges"))]
1096    ConvertTimeRanges {
1097        #[snafu(implicit)]
1098        location: Location,
1099        source: api::error::Error,
1100    },
1101
1102    #[snafu(display(
1103        "Column metadata inconsistencies found in table: {}, table_id: {}",
1104        table_name,
1105        table_id
1106    ))]
1107    ColumnMetadataConflicts {
1108        table_name: String,
1109        table_id: TableId,
1110    },
1111
1112    #[snafu(display(
1113        "Column not found in column metadata, column_name: {}, column_id: {}",
1114        column_name,
1115        column_id
1116    ))]
1117    ColumnNotFound { column_name: String, column_id: u32 },
1118
1119    #[snafu(display(
1120        "Column id mismatch, column_name: {}, expected column_id: {}, actual column_id: {}",
1121        column_name,
1122        expected_column_id,
1123        actual_column_id
1124    ))]
1125    ColumnIdMismatch {
1126        column_name: String,
1127        expected_column_id: u32,
1128        actual_column_id: u32,
1129    },
1130
1131    #[snafu(display(
1132        "Timestamp column mismatch, expected column_name: {}, expected column_id: {}, actual column_name: {}, actual column_id: {}",
1133        expected_column_name,
1134        expected_column_id,
1135        actual_column_name,
1136        actual_column_id,
1137    ))]
1138    TimestampMismatch {
1139        expected_column_name: String,
1140        expected_column_id: u32,
1141        actual_column_name: String,
1142        actual_column_id: u32,
1143    },
1144
1145    #[cfg(feature = "enterprise")]
1146    #[snafu(display("Too large duration"))]
1147    TooLargeDuration {
1148        #[snafu(source)]
1149        error: prost_types::DurationError,
1150        #[snafu(implicit)]
1151        location: Location,
1152    },
1153
1154    #[cfg(feature = "enterprise")]
1155    #[snafu(display("Negative duration"))]
1156    NegativeDuration {
1157        #[snafu(source)]
1158        error: prost_types::DurationError,
1159        #[snafu(implicit)]
1160        location: Location,
1161    },
1162
1163    #[cfg(feature = "enterprise")]
1164    #[snafu(display("Missing interval field"))]
1165    MissingInterval {
1166        #[snafu(implicit)]
1167        location: Location,
1168    },
1169}
1170
1171pub type Result<T> = std::result::Result<T, Error>;
1172
1173impl ErrorExt for Error {
1174    fn status_code(&self) -> StatusCode {
1175        use Error::*;
1176        match self {
1177            IllegalServerState { .. }
1178            | EtcdTxnOpResponse { .. }
1179            | EtcdFailed { .. }
1180            | EtcdTxnFailed { .. }
1181            | ConnectEtcd { .. }
1182            | MoveValues { .. }
1183            | TombstoneTargetAlreadyExists { .. }
1184            | GetCache { .. }
1185            | GetLatestCacheRetryExceeded { .. }
1186            | SerializeToJson { .. }
1187            | DeserializeFromJson { .. }
1188            | ElectionNoLeader { .. }
1189            | ElectionLeaderLeaseExpired { .. }
1190            | ElectionLeaderLeaseChanged { .. } => StatusCode::Internal,
1191
1192            NoLeader { .. } => StatusCode::TableUnavailable,
1193            ValueNotExist { .. }
1194            | ProcedurePoisonConflict { .. }
1195            | ProcedureStateReceiverNotFound { .. }
1196            | MissingColumnIds { .. }
1197            | MissingColumnInColumnMetadata { .. }
1198            | MismatchColumnId { .. }
1199            | ColumnMetadataConflicts { .. }
1200            | ColumnNotFound { .. }
1201            | ColumnIdMismatch { .. }
1202            | TimestampMismatch { .. } => StatusCode::Unexpected,
1203
1204            Unsupported { .. } | ReadOnlyKvBackend { .. } => StatusCode::Unsupported,
1205            WriteObject { .. } | ReadObject { .. } => StatusCode::StorageUnavailable,
1206
1207            SerdeJson { .. }
1208            | ParseOption { .. }
1209            | InvalidProtoMsg { .. }
1210            | InvalidMetadata { .. }
1211            | Unexpected { .. }
1212            | TableInfoNotFound { .. }
1213            | NextSequence { .. }
1214            | UnexpectedSequenceValue { .. }
1215            | InvalidHeartbeatResponse { .. }
1216            | EncodeJson { .. }
1217            | DecodeJson { .. }
1218            | PayloadNotExist { .. }
1219            | ConvertRawKey { .. }
1220            | DecodeProto { .. }
1221            | DecodePackedFileRefs { .. }
1222            | InvalidPackedFileRefs { .. }
1223            | BuildTableMeta { .. }
1224            | TableRouteNotFound { .. }
1225            | TableRepartNotFound { .. }
1226            | RegionOperatingRace { .. }
1227            | SerializeWalOptions { .. }
1228            | BuildKafkaClient { .. }
1229            | BuildKafkaCtrlClient { .. }
1230            | KafkaPartitionClient { .. }
1231            | ProduceRecord { .. }
1232            | CreateKafkaWalTopic { .. }
1233            | EmptyTopicPool { .. }
1234            | UnexpectedLogicalRouteTable { .. }
1235            | ProcedureOutput { .. }
1236            | FromUtf8 { .. }
1237            | MetadataCorruption { .. }
1238            | KafkaGetOffset { .. }
1239            | ReadFlexbuffers { .. }
1240            | SerializeFlexbuffers { .. }
1241            | DeserializeFlexbuffers { .. }
1242            | ConvertTimeRanges { .. } => StatusCode::Unexpected,
1243
1244            GetKvCache { .. } | CacheNotGet { .. } => StatusCode::Internal,
1245
1246            SchemaAlreadyExists { .. } => StatusCode::DatabaseAlreadyExists,
1247
1248            ProcedureNotFound { .. }
1249            | InvalidViewInfo { .. }
1250            | PrimaryKeyNotFound { .. }
1251            | EmptyKey { .. }
1252            | AlterLogicalTablesInvalidArguments { .. }
1253            | CreateLogicalTablesInvalidArguments { .. }
1254            | MismatchPrefix { .. }
1255            | TlsConfig { .. }
1256            | InvalidSetDatabaseOption { .. }
1257            | InvalidUnsetDatabaseOption { .. }
1258            | InvalidTopicNamePrefix { .. }
1259            | InvalidFileExtension { .. }
1260            | InvalidFileName { .. }
1261            | InvalidFlowRequestBody { .. }
1262            | InvalidFilePath { .. }
1263            | ConflictingSchemaOptions { .. } => StatusCode::InvalidArguments,
1264
1265            #[cfg(feature = "enterprise")]
1266            MissingInterval { .. } | NegativeDuration { .. } | TooLargeDuration { .. } => {
1267                StatusCode::InvalidArguments
1268            }
1269
1270            FlowNotFound { .. } => StatusCode::FlowNotFound,
1271            FlowRouteNotFound { .. } => StatusCode::Unexpected,
1272            FlowAlreadyExists { .. } => StatusCode::FlowAlreadyExists,
1273
1274            ViewNotFound { .. } | TableNotFound { .. } | RegionNotFound { .. } => {
1275                StatusCode::TableNotFound
1276            }
1277            ViewAlreadyExists { .. }
1278            | TableAlreadyExists { .. }
1279            | TableNameTombstoneConflict { .. } => StatusCode::TableAlreadyExists,
1280
1281            SubmitProcedure { source, .. }
1282            | QueryProcedure { source, .. }
1283            | WaitProcedure { source, .. }
1284            | StartProcedureManager { source, .. }
1285            | StopProcedureManager { source, .. } => source.status_code(),
1286            RegisterProcedureLoader { source, .. } => source.status_code(),
1287            External { source, .. } => source.status_code(),
1288            ResponseExceededSizeLimit { source, .. } => source.status_code(),
1289            OperateDatanode { source, .. } => source.status_code(),
1290            Table { source, .. } => source.status_code(),
1291            RetryLater { source, .. } => source.status_code(),
1292            AbortProcedure { source, .. } => source.status_code(),
1293            ConvertAlterTableRequest { source, .. } => source.status_code(),
1294            PutPoison { source, .. } => source.status_code(),
1295            ConvertColumnDef { source, .. } => source.status_code(),
1296            ProcedureStateReceiver { source, .. } => source.status_code(),
1297            RegisterRepartitionProcedureLoader { source, .. } => source.status_code(),
1298            CreateRepartitionProcedure { source, .. } => source.status_code(),
1299            PersistRepartitionGcRequirement { source, .. } => source.status_code(),
1300
1301            ParseProcedureId { .. }
1302            | InvalidNumTopics { .. }
1303            | SchemaNotFound { .. }
1304            | CatalogNotFound { .. }
1305            | InvalidNodeInfoKey { .. }
1306            | InvalidStatKey { .. }
1307            | ParseNum { .. }
1308            | InvalidRole { .. }
1309            | EmptyDdlTasks { .. } => StatusCode::InvalidArguments,
1310
1311            LoadTlsCertificate { .. } => StatusCode::Internal,
1312
1313            #[cfg(feature = "pg_kvbackend")]
1314            PostgresExecution { .. }
1315            | CreatePostgresPool { .. }
1316            | GetPostgresConnection { .. }
1317            | GetPostgresClient { .. }
1318            | PostgresTransaction { .. }
1319            | PostgresTlsConfig { .. }
1320            | InvalidTlsConfig { .. } => StatusCode::Internal,
1321            #[cfg(feature = "mysql_kvbackend")]
1322            MySqlExecution { .. }
1323            | CreateMySqlPool { .. }
1324            | DecodeSqlValue { .. }
1325            | AcquireMySqlClient { .. }
1326            | MySqlTransaction { .. } => StatusCode::Internal,
1327            #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
1328            RdsTransactionRetryFailed { .. } | SqlExecutionTimeout { .. } => StatusCode::Internal,
1329            DatanodeTableInfoNotFound { .. } => StatusCode::Internal,
1330        }
1331    }
1332
1333    fn as_any(&self) -> &dyn std::any::Any {
1334        self
1335    }
1336
1337    fn retry_hint(&self) -> RetryHint {
1338        use Error::*;
1339
1340        match self {
1341            RetryLater { .. }
1342            | GetLatestCacheRetryExceeded { .. }
1343            | NoLeader { .. }
1344            | ElectionNoLeader { .. }
1345            | ElectionLeaderLeaseExpired { .. }
1346            | ElectionLeaderLeaseChanged { .. } => RetryHint::Retryable,
1347            ConnectEtcd { error, .. } | EtcdFailed { error, .. } | EtcdTxnFailed { error, .. } => {
1348                retry_hint_from_etcd_error(error)
1349            }
1350            WriteObject { error, .. } | ReadObject { error, .. } => {
1351                retry_hint_from_opendal_error(error)
1352            }
1353            BuildKafkaClient { error, .. }
1354            | BuildKafkaCtrlClient { error, .. }
1355            | KafkaPartitionClient { error, .. }
1356            | KafkaGetOffset { error, .. }
1357            | ProduceRecord { error, .. }
1358            | CreateKafkaWalTopic { error, .. } => rskafka_client_error_to_retry_hint(error),
1359            SubmitProcedure { source, .. }
1360            | QueryProcedure { source, .. }
1361            | WaitProcedure { source, .. }
1362            | StartProcedureManager { source, .. }
1363            | StopProcedureManager { source, .. }
1364            | RegisterProcedureLoader { source, .. }
1365            | PutPoison { source, .. }
1366            | ProcedureStateReceiver { source, .. } => source.retry_hint(),
1367            External { source, .. }
1368            | ResponseExceededSizeLimit { source, .. }
1369            | OperateDatanode { source, .. }
1370            | AbortProcedure { source, .. }
1371            | RegisterRepartitionProcedureLoader { source, .. }
1372            | CreateRepartitionProcedure { source, .. }
1373            | PersistRepartitionGcRequirement { source, .. } => source.retry_hint(),
1374            Table { source, .. } => source.retry_hint(),
1375            ConvertAlterTableRequest { source, .. } => source.retry_hint(),
1376            ConvertColumnDef { source, .. } => source.retry_hint(),
1377            GetCache { source, .. } => source.retry_hint(),
1378            #[cfg(feature = "pg_kvbackend")]
1379            PostgresExecution { error, .. } | PostgresTransaction { error, .. } => {
1380                retry_hint_from_postgres_error(error)
1381            }
1382            #[cfg(feature = "pg_kvbackend")]
1383            GetPostgresClient { error, .. } => retry_hint_from_postgres_pool_error(error),
1384            #[cfg(feature = "mysql_kvbackend")]
1385            MySqlExecution { error, .. }
1386            | CreateMySqlPool { error, .. }
1387            | AcquireMySqlClient { error, .. }
1388            | MySqlTransaction { error, .. } => retry_hint_from_sqlx_error(error),
1389            #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
1390            RdsTransactionRetryFailed { .. } | SqlExecutionTimeout { .. } => RetryHint::Retryable,
1391            _ => RetryHint::NonRetryable,
1392        }
1393    }
1394}
1395
1396impl Error {
1397    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
1398    /// Check if the error is a serialization error.
1399    pub fn is_serialization_error(&self) -> bool {
1400        match self {
1401            #[cfg(feature = "pg_kvbackend")]
1402            Error::PostgresExecution { error, .. } | Error::PostgresTransaction { error, .. } => {
1403                retry_hint::is_postgres_serialization_error(error)
1404            }
1405            #[cfg(feature = "mysql_kvbackend")]
1406            Error::MySqlExecution { error, .. } | Error::MySqlTransaction { error, .. } => {
1407                retry_hint::is_mysql_serialization_error(error)
1408            }
1409            _ => false,
1410        }
1411    }
1412
1413    /// Creates a new [Error::RetryLater] error from source `err`.
1414    pub fn retry_later<E: ErrorExt + Send + Sync + 'static>(err: E) -> Error {
1415        Error::RetryLater {
1416            source: BoxedError::new(err),
1417            clean_poisons: false,
1418        }
1419    }
1420
1421    /// Determine whether it is a retry later type through [StatusCode]
1422    pub fn is_retry_later(&self) -> bool {
1423        matches!(
1424            self,
1425            Error::RetryLater { .. } | Error::GetLatestCacheRetryExceeded { .. }
1426        )
1427    }
1428
1429    /// Determine whether it needs to clean poisons.
1430    pub fn need_clean_poisons(&self) -> bool {
1431        matches!(
1432            self,
1433            Error::AbortProcedure { clean_poisons, .. } if *clean_poisons
1434        ) || matches!(
1435            self,
1436            Error::RetryLater { clean_poisons, .. } if *clean_poisons
1437        )
1438    }
1439
1440    /// Returns true if the response exceeds the size limit.
1441    pub fn is_exceeded_size_limit(&self) -> bool {
1442        match self {
1443            Error::EtcdFailed {
1444                error: etcd_client::Error::GRpcStatus(status),
1445                ..
1446            } => status.code() == tonic::Code::OutOfRange,
1447            Error::ResponseExceededSizeLimit { .. } => true,
1448            _ => false,
1449        }
1450    }
1451}
1452
1453#[cfg(test)]
1454mod retry_hint_tests {
1455    use std::sync::Arc;
1456
1457    use common_error::mock::MockError;
1458
1459    use super::*;
1460
1461    #[test]
1462    fn test_retry_later_hint_is_retryable() {
1463        let err = Error::retry_later(MockError::new(StatusCode::Internal));
1464
1465        assert_eq!(err.retry_hint(), RetryHint::Retryable);
1466    }
1467
1468    #[test]
1469    fn test_latest_cache_retry_exceeded_hint_is_retryable() {
1470        let err = GetLatestCacheRetryExceededSnafu { attempts: 3_usize }.build();
1471
1472        assert_eq!(err.retry_hint(), RetryHint::Retryable);
1473    }
1474
1475    #[test]
1476    fn test_get_cache_forwards_retry_hint() {
1477        let source = Arc::new(Error::retry_later(MockError::new(StatusCode::Internal)));
1478        let err = Error::GetCache { source };
1479
1480        assert_eq!(err.retry_hint(), RetryHint::Retryable);
1481    }
1482
1483    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
1484    #[test]
1485    fn test_sql_execution_timeout_hint_is_retryable() {
1486        let err = SqlExecutionTimeoutSnafu {
1487            sql: "SELECT 1".to_string(),
1488            duration: std::time::Duration::from_secs(1),
1489        }
1490        .build();
1491
1492        assert_eq!(err.retry_hint(), RetryHint::Retryable);
1493    }
1494
1495    #[cfg(any(feature = "pg_kvbackend", feature = "mysql_kvbackend"))]
1496    #[test]
1497    fn test_rds_transaction_retry_failed_hint_is_retryable() {
1498        let err = RdsTransactionRetryFailedSnafu.build();
1499
1500        assert_eq!(err.retry_hint(), RetryHint::Retryable);
1501    }
1502
1503    #[test]
1504    fn test_default_hint_is_non_retryable() {
1505        let err = UnexpectedSnafu {
1506            err_msg: "mock error",
1507        }
1508        .build();
1509
1510        assert_eq!(err.retry_hint(), RetryHint::NonRetryable);
1511    }
1512}