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    /// Whether this error is transient and the request should be retried.
161    ///
162    /// Retriable errors:
163    /// - Network / connection errors (`Http` from reqwest)
164    /// - Server errors (5xx status codes)
165    /// - Rate limiting (429 — handled separately with `Retry-After`)
166    ///
167    /// Non-retriable errors:
168    /// - Client errors (4xx except 429)
169    /// - JSON parse / JSONPath / auth / transform errors
170    pub fn is_retriable(&self) -> bool {
171        match self {
172            // reqwest errors: connection timeouts, DNS failures, etc. are retriable.
173            FaucetError::Http(e) => {
174                // If it's a status error that leaked through, check the code.
175                if let Some(status) = e.status() {
176                    status.is_server_error()
177                } else {
178                    // Connection errors, timeouts, etc.
179                    true
180                }
181            }
182            // 5xx are retriable; 429 (Too Many Requests) is too — sources that
183            // surface a rate-limit as a plain HttpStatus rather than the
184            // dedicated RateLimited variant (XML, GraphQL) would otherwise abort
185            // on the first 429 (audit #146 H3).
186            FaucetError::HttpStatus { status, .. } => *status >= 500 || *status == 429,
187            FaucetError::RateLimited(_) => true,
188            _ => false,
189        }
190    }
191}
192
193#[cfg(test)]
194mod tests {
195    use super::*;
196
197    #[test]
198    fn http_status_5xx_is_retriable() {
199        let err = FaucetError::HttpStatus {
200            status: 500,
201            url: "https://example.com".into(),
202            body: "Internal Server Error".into(),
203        };
204        assert!(err.is_retriable());
205
206        let err = FaucetError::HttpStatus {
207            status: 503,
208            url: "https://example.com".into(),
209            body: "".into(),
210        };
211        assert!(err.is_retriable());
212    }
213
214    #[test]
215    fn http_status_4xx_is_not_retriable() {
216        let err = FaucetError::HttpStatus {
217            status: 400,
218            url: "https://example.com".into(),
219            body: "Bad Request".into(),
220        };
221        assert!(!err.is_retriable());
222
223        let err = FaucetError::HttpStatus {
224            status: 404,
225            url: "https://example.com".into(),
226            body: "".into(),
227        };
228        assert!(!err.is_retriable());
229    }
230
231    #[test]
232    fn http_status_429_is_retriable() {
233        // H3 (audit #146): a 429 surfaced as a plain HttpStatus (XML/GraphQL
234        // sources) must be retriable, not aborted on the first hit.
235        let err = FaucetError::HttpStatus {
236            status: 429,
237            url: "https://example.com".into(),
238            body: "Too Many Requests".into(),
239        };
240        assert!(err.is_retriable());
241    }
242
243    #[test]
244    fn rate_limited_is_retriable() {
245        let err = FaucetError::RateLimited(Duration::from_secs(30));
246        assert!(err.is_retriable());
247    }
248
249    #[test]
250    fn json_error_is_not_retriable() {
251        let serde_err = serde_json::from_str::<serde_json::Value>("not json").unwrap_err();
252        let err = FaucetError::Json(serde_err);
253        assert!(!err.is_retriable());
254    }
255
256    #[test]
257    fn jsonpath_error_is_not_retriable() {
258        let err = FaucetError::JsonPath("bad path".into());
259        assert!(!err.is_retriable());
260    }
261
262    #[test]
263    fn auth_error_is_not_retriable() {
264        let err = FaucetError::Auth("invalid token".into());
265        assert!(!err.is_retriable());
266    }
267
268    #[test]
269    fn url_error_is_not_retriable() {
270        let err = FaucetError::Url("bad url".into());
271        assert!(!err.is_retriable());
272    }
273
274    #[test]
275    fn transform_error_is_not_retriable() {
276        let err = FaucetError::Transform("bad regex".into());
277        assert!(!err.is_retriable());
278    }
279
280    #[test]
281    fn http_status_display_includes_url_and_body() {
282        let err = FaucetError::HttpStatus {
283            status: 422,
284            url: "https://api.example.com/test".into(),
285            body: "Unprocessable Entity".into(),
286        };
287        let msg = err.to_string();
288        assert!(msg.contains("422"));
289        assert!(msg.contains("https://api.example.com/test"));
290        assert!(msg.contains("Unprocessable Entity"));
291    }
292
293    #[test]
294    fn config_error_is_not_retriable() {
295        let err = FaucetError::Config("bad endpoint".into());
296        assert!(!err.is_retriable());
297    }
298
299    #[test]
300    fn config_error_display() {
301        let err = FaucetError::Config("missing descriptor".into());
302        assert_eq!(err.to_string(), "Config error: missing descriptor");
303    }
304
305    #[test]
306    fn source_error_is_not_retriable() {
307        let err = FaucetError::Source("query failed".into());
308        assert!(!err.is_retriable());
309    }
310
311    #[test]
312    fn source_error_display() {
313        let err = FaucetError::Source("connection refused".into());
314        assert_eq!(err.to_string(), "Source error: connection refused");
315    }
316
317    #[test]
318    fn custom_error_is_not_retriable() {
319        let err = FaucetError::Custom(Box::new(std::io::Error::other("custom failure")));
320        assert!(!err.is_retriable());
321    }
322
323    #[test]
324    fn custom_error_display() {
325        let err = FaucetError::Custom(Box::new(std::io::Error::other("custom failure")));
326        assert_eq!(err.to_string(), "Connector error: custom failure");
327    }
328
329    #[test]
330    fn custom_error_from_boxed() {
331        let io_err = std::io::Error::other("file missing");
332        let boxed: Box<dyn std::error::Error + Send + Sync> = Box::new(io_err);
333        let err: FaucetError = boxed.into();
334        assert!(matches!(err, FaucetError::Custom(_)));
335    }
336
337    #[test]
338    fn sink_error_is_not_retriable() {
339        let err = FaucetError::Sink("BigQuery insert failed".into());
340        assert!(!err.is_retriable());
341    }
342
343    #[test]
344    fn sink_error_display() {
345        let err = FaucetError::Sink("connection refused".into());
346        assert_eq!(err.to_string(), "Sink error: connection refused");
347    }
348
349    #[test]
350    fn quality_failure_is_not_retriable_and_displays() {
351        let err = FaucetError::QualityFailure {
352            check: "not_null".into(),
353            message: "field 'user_id' was null".into(),
354        };
355        assert!(!err.is_retriable());
356        let s = err.to_string();
357        assert!(s.contains("not_null"));
358        assert!(s.contains("user_id"));
359    }
360
361    #[test]
362    fn schema_drift_is_not_retriable_and_displays() {
363        let err = FaucetError::SchemaDrift {
364            columns: vec!["email".into(), "score".into()],
365            message: "2 new columns".into(),
366        };
367        assert!(!err.is_retriable());
368        let s = err.to_string();
369        assert!(s.contains("email"));
370        assert!(s.contains("2 new columns"));
371    }
372
373    #[test]
374    fn circuit_open_is_not_retriable_and_displays() {
375        let err = FaucetError::CircuitOpen {
376            failures: 5,
377            cooldown: std::time::Duration::from_secs(60),
378        };
379        assert!(!err.is_retriable());
380        let s = err.to_string();
381        assert!(s.contains("5"));
382        assert!(s.contains("Circuit open"));
383    }
384}