1use 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 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 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 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 pub fn is_retry_later(&self) -> bool {
1423 matches!(
1424 self,
1425 Error::RetryLater { .. } | Error::GetLatestCacheRetryExceeded { .. }
1426 )
1427 }
1428
1429 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 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}