Skip to main content

datum_mq/
error.rs

1use datum::StreamError;
2use thiserror::Error;
3
4/// Result type used by `datum-mq`.
5pub type MqResult<T> = Result<T, MqError>;
6
7/// Kafka connector errors.
8#[derive(Debug, Error)]
9pub enum MqError {
10    #[error("invalid Kafka configuration: {0}")]
11    InvalidConfig(String),
12
13    #[error("Kafka control handle is not available")]
14    ControlUnavailable,
15
16    #[error("Kafka assignment was lost for {topic}:{partition}; refusing to commit offset")]
17    AssignmentLost { topic: String, partition: i32 },
18
19    #[error("Kafka drain timed out")]
20    DrainTimeout,
21
22    #[error("Kafka producer delivery failed: {0}")]
23    Delivery(String),
24
25    #[error(
26        "native Kafka producer delivery failed for {topic}:{partition}: {message} (broker code {code})"
27    )]
28    NativeDelivery {
29        topic: String,
30        partition: i32,
31        code: i16,
32        message: String,
33    },
34
35    #[error("native Kafka TLS handshake with broker {broker} failed: {message}")]
36    NativeTls { broker: String, message: String },
37
38    #[error("native Kafka SASL {mechanism} authentication with broker {broker} failed: {message}")]
39    NativeSasl {
40        broker: String,
41        mechanism: String,
42        message: String,
43    },
44
45    #[error("Kafka operation failed: {0}")]
46    Failed(String),
47}
48
49impl From<MqError> for StreamError {
50    fn from(error: MqError) -> Self {
51        StreamError::Failed(error.to_string())
52    }
53}