1use std::time::Duration;
4use thiserror::Error;
5
6#[derive(Debug, Error)]
14#[non_exhaustive]
15pub enum FaucetError {
16 #[error("HTTP error: {0}")]
17 Http(#[from] reqwest::Error),
18
19 #[error("HTTP {status} from {url}: {body}")]
25 HttpStatus {
26 status: u16,
27 url: String,
28 body: String,
29 },
30
31 #[error("JSON error: {0}")]
32 Json(#[from] serde_json::Error),
33
34 #[error("JSONPath error: {0}")]
35 JsonPath(String),
36
37 #[error("Auth error: {0}")]
38 Auth(String),
39
40 #[error("Rate limited: retry after {0:?}")]
44 RateLimited(Duration),
45
46 #[error("URL error: {0}")]
48 Url(String),
49
50 #[error("Transform error: {0}")]
52 Transform(String),
53
54 #[error("Config error: {0}")]
56 Config(String),
57
58 #[error("Source error: {0}")]
60 Source(String),
61
62 #[error("Sink error: {0}")]
64 Sink(String),
65
66 #[error("Quality check '{check}' failed: {message}")]
68 QualityFailure { check: String, message: String },
69
70 #[error("Schema drift on columns {columns:?}: {message}")]
73 SchemaDrift {
74 columns: Vec<String>,
75 message: String,
76 },
77
78 #[error("Policy `{rule}` violated on column `{column}`: {message}")]
82 PolicyViolation {
83 rule: String,
84 column: String,
85 message: String,
86 },
87
88 #[error("Budget `{budget}` exceeded: limit {limit}, actual {actual}")]
95 BudgetExceeded {
96 budget: String,
97 limit: u64,
98 actual: u64,
99 },
100
101 #[error("Profile drift on columns {columns:?}: {message}")]
105 ProfileDrift {
106 columns: Vec<String>,
107 message: String,
108 },
109
110 #[error("Contract v{version} violated: {message}")]
114 ContractViolation { version: String, message: String },
115
116 #[error("State error: {0}")]
119 State(String),
120
121 #[error(
126 "state '{key}' is incompatible: found {found}, expected {expected} — it was written by a \
127 newer faucet or by a different source; run the release that wrote it, or reset the row \
128 (`faucet state reset`) after confirming where it should resume"
129 )]
130 StateIncompatible {
131 key: String,
133 found: String,
135 expected: String,
137 },
138
139 #[error("Circuit open after {failures} consecutive failures; cooldown {cooldown:?}")]
143 CircuitOpen { failures: u32, cooldown: Duration },
144
145 #[error("Connector error: {0}")]
156 Custom(#[from] Box<dyn std::error::Error + Send + Sync>),
157}
158
159impl FaucetError {
160 pub fn sink_status(status: Option<u16>, url: impl Into<String>, message: String) -> Self {
166 match status {
167 Some(status) if status == 429 || status >= 500 => FaucetError::HttpStatus {
168 status,
169 url: url.into(),
170 body: message,
171 },
172 _ => FaucetError::Sink(message),
173 }
174 }
175
176 pub fn is_retriable(&self) -> bool {
187 match self {
188 FaucetError::Http(e) => {
190 if let Some(status) = e.status() {
192 status.is_server_error()
193 } else {
194 true
196 }
197 }
198 FaucetError::HttpStatus { status, .. } => *status >= 500 || *status == 429,
203 FaucetError::RateLimited(_) => true,
204 _ => false,
205 }
206 }
207}
208
209#[cfg(test)]
210mod tests {
211 use super::*;
212
213 #[test]
214 fn sink_status_types_only_retryable_statuses() {
215 for status in [429, 500, 503] {
216 let e = FaucetError::sink_status(Some(status), "s3://b/k", "m".into());
217 assert!(
218 matches!(e, FaucetError::HttpStatus { status: s, ref url, ref body }
219 if s == status && url == "s3://b/k" && body == "m")
220 );
221 assert!(e.is_retriable());
222 }
223 for status in [None, Some(403), Some(404)] {
224 let e = FaucetError::sink_status(status, "u", "m".into());
225 assert!(matches!(e, FaucetError::Sink(ref m) if m == "m"), "{e:?}");
226 }
227 }
228
229 #[test]
230 fn http_status_5xx_is_retriable() {
231 let err = FaucetError::HttpStatus {
232 status: 500,
233 url: "https://example.com".into(),
234 body: "Internal Server Error".into(),
235 };
236 assert!(err.is_retriable());
237
238 let err = FaucetError::HttpStatus {
239 status: 503,
240 url: "https://example.com".into(),
241 body: "".into(),
242 };
243 assert!(err.is_retriable());
244 }
245
246 #[test]
247 fn http_status_4xx_is_not_retriable() {
248 let err = FaucetError::HttpStatus {
249 status: 400,
250 url: "https://example.com".into(),
251 body: "Bad Request".into(),
252 };
253 assert!(!err.is_retriable());
254
255 let err = FaucetError::HttpStatus {
256 status: 404,
257 url: "https://example.com".into(),
258 body: "".into(),
259 };
260 assert!(!err.is_retriable());
261 }
262
263 #[test]
264 fn http_status_429_is_retriable() {
265 let err = FaucetError::HttpStatus {
268 status: 429,
269 url: "https://example.com".into(),
270 body: "Too Many Requests".into(),
271 };
272 assert!(err.is_retriable());
273 }
274
275 #[test]
276 fn rate_limited_is_retriable() {
277 let err = FaucetError::RateLimited(Duration::from_secs(30));
278 assert!(err.is_retriable());
279 }
280
281 #[test]
282 fn json_error_is_not_retriable() {
283 let serde_err = serde_json::from_str::<serde_json::Value>("not json").unwrap_err();
284 let err = FaucetError::Json(serde_err);
285 assert!(!err.is_retriable());
286 }
287
288 #[test]
289 fn jsonpath_error_is_not_retriable() {
290 let err = FaucetError::JsonPath("bad path".into());
291 assert!(!err.is_retriable());
292 }
293
294 #[test]
295 fn auth_error_is_not_retriable() {
296 let err = FaucetError::Auth("invalid token".into());
297 assert!(!err.is_retriable());
298 }
299
300 #[test]
301 fn url_error_is_not_retriable() {
302 let err = FaucetError::Url("bad url".into());
303 assert!(!err.is_retriable());
304 }
305
306 #[test]
307 fn transform_error_is_not_retriable() {
308 let err = FaucetError::Transform("bad regex".into());
309 assert!(!err.is_retriable());
310 }
311
312 #[test]
313 fn http_status_display_includes_url_and_body() {
314 let err = FaucetError::HttpStatus {
315 status: 422,
316 url: "https://api.example.com/test".into(),
317 body: "Unprocessable Entity".into(),
318 };
319 let msg = err.to_string();
320 assert!(msg.contains("422"));
321 assert!(msg.contains("https://api.example.com/test"));
322 assert!(msg.contains("Unprocessable Entity"));
323 }
324
325 #[test]
326 fn config_error_is_not_retriable() {
327 let err = FaucetError::Config("bad endpoint".into());
328 assert!(!err.is_retriable());
329 }
330
331 #[test]
332 fn config_error_display() {
333 let err = FaucetError::Config("missing descriptor".into());
334 assert_eq!(err.to_string(), "Config error: missing descriptor");
335 }
336
337 #[test]
338 fn source_error_is_not_retriable() {
339 let err = FaucetError::Source("query failed".into());
340 assert!(!err.is_retriable());
341 }
342
343 #[test]
344 fn source_error_display() {
345 let err = FaucetError::Source("connection refused".into());
346 assert_eq!(err.to_string(), "Source error: connection refused");
347 }
348
349 #[test]
350 fn custom_error_is_not_retriable() {
351 let err = FaucetError::Custom(Box::new(std::io::Error::other("custom failure")));
352 assert!(!err.is_retriable());
353 }
354
355 #[test]
356 fn custom_error_display() {
357 let err = FaucetError::Custom(Box::new(std::io::Error::other("custom failure")));
358 assert_eq!(err.to_string(), "Connector error: custom failure");
359 }
360
361 #[test]
362 fn custom_error_from_boxed() {
363 let io_err = std::io::Error::other("file missing");
364 let boxed: Box<dyn std::error::Error + Send + Sync> = Box::new(io_err);
365 let err: FaucetError = boxed.into();
366 assert!(matches!(err, FaucetError::Custom(_)));
367 }
368
369 #[test]
370 fn sink_error_is_not_retriable() {
371 let err = FaucetError::Sink("BigQuery insert failed".into());
372 assert!(!err.is_retriable());
373 }
374
375 #[test]
376 fn sink_error_display() {
377 let err = FaucetError::Sink("connection refused".into());
378 assert_eq!(err.to_string(), "Sink error: connection refused");
379 }
380
381 #[test]
382 fn quality_failure_is_not_retriable_and_displays() {
383 let err = FaucetError::QualityFailure {
384 check: "not_null".into(),
385 message: "field 'user_id' was null".into(),
386 };
387 assert!(!err.is_retriable());
388 let s = err.to_string();
389 assert!(s.contains("not_null"));
390 assert!(s.contains("user_id"));
391 }
392
393 #[test]
394 fn schema_drift_is_not_retriable_and_displays() {
395 let err = FaucetError::SchemaDrift {
396 columns: vec!["email".into(), "score".into()],
397 message: "2 new columns".into(),
398 };
399 assert!(!err.is_retriable());
400 let s = err.to_string();
401 assert!(s.contains("email"));
402 assert!(s.contains("2 new columns"));
403 }
404
405 #[test]
406 fn circuit_open_is_not_retriable_and_displays() {
407 let err = FaucetError::CircuitOpen {
408 failures: 5,
409 cooldown: std::time::Duration::from_secs(60),
410 };
411 assert!(!err.is_retriable());
412 let s = err.to_string();
413 assert!(s.contains("5"));
414 assert!(s.contains("Circuit open"));
415 }
416}