laminar_connectors/schema/
error.rs1use thiserror::Error;
7
8use crate::error::ConnectorError;
9
10pub type SchemaResult<T> = Result<T, SchemaError>;
12
13#[derive(Debug, Error)]
15pub enum SchemaError {
16 #[error("inference failed: {0}")]
18 InferenceFailed(String),
19
20 #[error("incompatible schemas: {0}")]
22 Incompatible(String),
23
24 #[error("registry error: {0}")]
26 RegistryError(String),
27
28 #[error("decode error: {0}")]
30 DecodeError(String),
31
32 #[error("evolution rejected: {0}")]
34 EvolutionRejected(String),
35
36 #[error("missing config: {0}")]
38 MissingConfig(String),
39
40 #[error("invalid config key '{key}': {message}")]
42 InvalidConfig {
43 key: String,
45 message: String,
47 },
48
49 #[error("duplicate wildcard: only one `*` is allowed in the column list")]
51 DuplicateWildcard,
52
53 #[error(
55 "wildcard without resolution: `*` requires a connector with a schema provider or registry"
56 )]
57 WildcardWithoutResolution,
58
59 #[error("wildcard prefix collision: prefixed column '{0}' collides with a declared column")]
61 WildcardPrefixCollision(String),
62
63 #[error("wildcard expanded to zero new columns: all source columns are already declared")]
66 WildcardNoNewFields,
67
68 #[error("arrow error: {0}")]
70 Arrow(#[from] arrow_schema::ArrowError),
71
72 #[error(transparent)]
74 Other(Box<dyn std::error::Error + Send + Sync>),
75}
76
77impl From<ConnectorError> for SchemaError {
78 fn from(err: ConnectorError) -> Self {
79 match err {
80 ConnectorError::ConfigurationError(msg) => SchemaError::InvalidConfig {
83 key: String::new(),
84 message: msg,
85 },
86 ConnectorError::SchemaMismatch(msg) => SchemaError::Incompatible(msg),
87 other => SchemaError::Other(Box::new(other)),
88 }
89 }
90}
91
92impl From<SchemaError> for ConnectorError {
93 fn from(err: SchemaError) -> Self {
94 match err {
95 SchemaError::MissingConfig(msg) => ConnectorError::missing_config(msg),
96 SchemaError::InvalidConfig { key, message } => {
97 ConnectorError::ConfigurationError(format!("invalid config key '{key}': {message}"))
98 }
99 SchemaError::Incompatible(msg) => ConnectorError::SchemaMismatch(msg),
100 SchemaError::DecodeError(msg) => ConnectorError::ReadError(msg),
101 other => ConnectorError::Internal(other.to_string()),
102 }
103 }
104}
105
106#[cfg(test)]
107mod tests {
108 use super::*;
109
110 #[test]
111 fn test_schema_error_display() {
112 let err = SchemaError::InferenceFailed("too few samples".into());
113 assert_eq!(err.to_string(), "inference failed: too few samples");
114 }
115
116 #[test]
117 fn test_schema_error_invalid_config() {
118 let err = SchemaError::InvalidConfig {
119 key: "format".into(),
120 message: "unknown format 'xml'".into(),
121 };
122 assert!(err.to_string().contains("format"));
123 assert!(err.to_string().contains("unknown format"));
124 }
125
126 #[test]
127 fn test_connector_to_schema_error() {
128 let ce = ConnectorError::missing_config("topic");
129 let se: SchemaError = ce.into();
130 assert!(
133 matches!(&se, SchemaError::InvalidConfig { message, .. } if message.contains("topic"))
134 );
135 }
136
137 #[test]
138 fn test_schema_to_connector_error() {
139 let se = SchemaError::Incompatible("field type mismatch".into());
140 let ce: ConnectorError = se.into();
141 assert!(matches!(ce, ConnectorError::SchemaMismatch(_)));
142 }
143
144 #[test]
145 fn invalid_schema_config_remains_a_connector_configuration_error() {
146 let error: ConnectorError = SchemaError::InvalidConfig {
147 key: "json.column.ts.epoch_unit".into(),
148 message: "invalid epoch unit".into(),
149 }
150 .into();
151 assert!(matches!(error, ConnectorError::ConfigurationError(_)));
152 }
153
154 #[test]
155 fn test_schema_error_from_arrow() {
156 let arrow_err = arrow_schema::ArrowError::SchemaError("bad schema".into());
157 let se: SchemaError = arrow_err.into();
158 assert!(matches!(se, SchemaError::Arrow(_)));
159 assert!(se.to_string().contains("bad schema"));
160 }
161
162 #[test]
163 fn test_other_connector_error_wraps() {
164 let ce = ConnectorError::ConnectionFailed("host down".into());
165 let se: SchemaError = ce.into();
166 assert!(matches!(se, SchemaError::Other(_)));
167 assert!(se.to_string().contains("host down"));
168 }
169
170 #[test]
171 fn test_wildcard_errors_display() {
172 let e1 = SchemaError::DuplicateWildcard;
173 assert!(e1.to_string().contains("duplicate wildcard"));
174
175 let e2 = SchemaError::WildcardWithoutResolution;
176 assert!(e2.to_string().contains("schema provider"));
177
178 let e3 = SchemaError::WildcardPrefixCollision("src_id".into());
179 assert!(e3.to_string().contains("src_id"));
180
181 let e4 = SchemaError::WildcardNoNewFields;
182 assert!(e4.to_string().contains("zero new columns"));
183 }
184}