Skip to main content

faucet_core/
error.rs

1//! Error types for faucet-stream.
2
3use std::time::Duration;
4use thiserror::Error;
5
6/// All possible errors returned by faucet-stream.
7///
8/// `#[non_exhaustive]`: this enum demonstrably grows (`QualityFailure`,
9/// `CircuitOpen`, `SchemaDrift`, `ContractViolation` all arrived post-1.0) and
10/// it is the one type every third-party connector author touches. Marking it
11/// keeps the next variant a **minor** release instead of a forced major for
12/// `faucet-core` and everything downstream. Match with a `_` arm.
13#[derive(Debug, Error)]
14#[non_exhaustive]
15pub enum FaucetError {
16    #[error("HTTP error: {0}")]
17    Http(#[from] reqwest::Error),
18
19    /// An HTTP response with a non-success status code.
20    ///
21    /// Contains the status code, URL, and (truncated) response body for
22    /// debugging.  Whether this error is retriable depends on the status code
23    /// — see [`FaucetError::is_retriable`].
24    #[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    /// The server responded with HTTP 429 Too Many Requests.
41    /// The inner value is the duration to wait before retrying,
42    /// parsed from the `Retry-After` response header (default: 60 s).
43    #[error("Rate limited: retry after {0:?}")]
44    RateLimited(Duration),
45
46    /// A URL could not be constructed or parsed.
47    #[error("URL error: {0}")]
48    Url(String),
49
50    /// A record transform could not be compiled or applied (e.g. invalid regex).
51    #[error("Transform error: {0}")]
52    Transform(String),
53
54    /// A configuration or validation error (e.g. invalid endpoint, missing descriptor).
55    #[error("Config error: {0}")]
56    Config(String),
57
58    /// A source operation failed (e.g. database query error, file read error).
59    #[error("Source error: {0}")]
60    Source(String),
61
62    /// A sink operation failed (e.g. BigQuery insert error).
63    #[error("Sink error: {0}")]
64    Sink(String),
65
66    /// A data-quality check failed under an `abort` policy.
67    #[error("Quality check '{check}' failed: {message}")]
68    QualityFailure { check: String, message: String },
69
70    /// An incoming page's shape diverged from the destination schema under an
71    /// `on_drift: fail` (or `on_incompatible: fail`) policy.
72    #[error("Schema drift on columns {columns:?}: {message}")]
73    SchemaDrift {
74        columns: Vec<String>,
75        message: String,
76    },
77
78    /// A data-flow policy rule was violated (#702): at run time a labelled
79    /// value headed for a sink the rule does not allow, or a static check
80    /// refused the run before it started.
81    #[error("Policy `{rule}` violated on column `{column}`: {message}")]
82    PolicyViolation {
83        rule: String,
84        column: String,
85        message: String,
86    },
87
88    /// A run budget ceiling was crossed (#703): `budget` names the ceiling
89    /// (`max_records` / `max_bytes` / `max_duration_secs`), `limit` its value
90    /// and `actual` what the run reached or would have reached. A records /
91    /// bytes crossing refuses the whole page **before** it is written, so
92    /// the bookmark never advances past it; a duration crossing cancels the
93    /// run cooperatively at the next page boundary.
94    #[error("Budget `{budget}` exceeded: limit {limit}, actual {actual}")]
95    BudgetExceeded {
96        budget: String,
97        limit: u64,
98        actual: u64,
99    },
100
101    /// A run's learned column profile drifted from its baseline under a
102    /// `profiling.on_drift: fail` policy (#708). Raised after the run's data
103    /// is written: the failure marks the run, it does not undo it.
104    #[error("Profile drift on columns {columns:?}: {message}")]
105    ProfileDrift {
106        columns: Vec<String>,
107        message: String,
108    },
109
110    /// A record breached the pipeline's data contract under an
111    /// `on_breach: fail` policy. `version` is the contract version the
112    /// record was validated against.
113    #[error("Contract v{version} violated: {message}")]
114    ContractViolation { version: String, message: String },
115
116    /// A state-store operation failed (read/write/delete of a replication
117    /// bookmark, checkpoint, or other persisted pipeline state).
118    #[error("State error: {0}")]
119    State(String),
120
121    /// A stored state value this release or this source cannot read (#736):
122    /// written by a newer faucet (a newer state format or bookmark schema) or
123    /// by a different source. Raised before anything is read from the source,
124    /// instead of guessing at the shape — which would silently re-sync or skip.
125    #[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        /// The state key.
132        key: String,
133        /// What is stored.
134        found: String,
135        /// What this release and source read.
136        expected: String,
137    },
138
139    /// The resilience circuit breaker opened after repeated failures; the run
140    /// is aborted fast. `cooldown` is advisory for the orchestration layer
141    /// (e.g. `faucet schedule` delays re-entry by this duration).
142    #[error("Circuit open after {failures} consecutive failures; cooldown {cooldown:?}")]
143    CircuitOpen { failures: u32, cooldown: Duration },
144
145    /// A custom error from a third-party connector.
146    ///
147    /// Use this to wrap your own error types without losing the error chain:
148    /// ```rust
149    /// use faucet_core::FaucetError;
150    /// let err = FaucetError::Custom(Box::new(std::io::Error::new(
151    ///     std::io::ErrorKind::Other,
152    ///     "my connector failed",
153    /// )));
154    /// ```
155    #[error("Connector error: {0}")]
156    Custom(#[from] Box<dyn std::error::Error + Send + Sync>),
157}
158
159impl FaucetError {
160    /// A storage client's failure, typed by its HTTP status when it has one:
161    /// a 429 or 5xx becomes [`HttpStatus`](Self::HttpStatus) (so the
162    /// resilience policy retries it) with `message` as the body; anything
163    /// else, including a failure with no response, is a
164    /// [`Sink`](Self::Sink) error carrying `message`.
165    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    /// Whether this error is transient and the request should be retried.
177    ///
178    /// Retriable errors:
179    /// - Network / connection errors (`Http` from reqwest)
180    /// - Server errors (5xx status codes)
181    /// - Rate limiting (429 — handled separately with `Retry-After`)
182    ///
183    /// Non-retriable errors:
184    /// - Client errors (4xx except 429)
185    /// - JSON parse / JSONPath / auth / transform errors
186    pub fn is_retriable(&self) -> bool {
187        match self {
188            // reqwest errors: connection timeouts, DNS failures, etc. are retriable.
189            FaucetError::Http(e) => {
190                // If it's a status error that leaked through, check the code.
191                if let Some(status) = e.status() {
192                    status.is_server_error()
193                } else {
194                    // Connection errors, timeouts, etc.
195                    true
196                }
197            }
198            // 5xx are retriable; 429 (Too Many Requests) is too — sources that
199            // surface a rate-limit as a plain HttpStatus rather than the
200            // dedicated RateLimited variant (XML, GraphQL) would otherwise abort
201            // on the first 429 (audit #146 H3).
202            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        // H3 (audit #146): a 429 surfaced as a plain HttpStatus (XML/GraphQL
266        // sources) must be retriable, not aborted on the first hit.
267        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}