1use datum::StreamError;
2use thiserror::Error;
3
4pub type MqResult<T> = Result<T, MqError>;
6
7#[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}