Skip to main content

faucet_core/
dlq.rs

1//! Dead-letter queue (DLQ) wiring shared by the pipeline runner.
2//!
3//! The types defined here are config-shaped: they describe *what* the
4//! pipeline should do with row-level failures, not *how* the routing is
5//! executed. The execution lives in [`run_stream`](crate::run_stream).
6//!
7//! See `docs/superpowers/specs/2026-05-24-dlq-design.md`.
8
9use crate::FaucetError;
10use crate::traits::Sink;
11use schemars::JsonSchema;
12use serde::{Deserialize, Serialize};
13use serde_json::{Value, json};
14use std::fmt;
15use std::sync::Arc;
16use std::time::{SystemTime, UNIX_EPOCH};
17
18/// Policy applied when a sink reports an outer failure (the whole batch
19/// failed, no per-row info).
20#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
21#[serde(rename_all = "snake_case")]
22pub enum OnBatchError {
23    /// Surface the underlying [`FaucetError`] and fail the pipeline (default).
24    #[default]
25    Propagate,
26    /// Treat every row in the failed page as a DLQ candidate. Unsafe with
27    /// best-effort APIs that haven't overridden
28    /// [`Sink::write_batch_partial`] — already-committed rows would land in
29    /// the DLQ as duplicates. Use with atomic sinks (single-statement
30    /// INSERT, file writes) where the failure mode is "nothing landed".
31    DlqAll,
32}
33
34/// Pipeline-level DLQ wiring.
35#[derive(Clone)]
36pub struct DlqConfig {
37    /// Sink that receives DLQ envelopes.
38    pub sink: Arc<dyn Sink>,
39    /// What to do when the main sink fails wholesale.
40    pub on_batch_error: OnBatchError,
41    /// Per-page failure budget. `None` = unlimited.
42    ///
43    /// This budget is **shared across both sink-side row failures and
44    /// quality-check quarantines**: a record routed to the DLQ by a
45    /// `quarantine` quality check counts against it just as a sink-side
46    /// row failure does.
47    pub max_failures_per_page: Option<usize>,
48    /// Cumulative failure budget across the run. `None` = unlimited.
49    ///
50    /// This budget is **shared across both sink-side row failures and
51    /// quality-check quarantines**: records quarantined by the quality pass
52    /// accumulate in this counter alongside sink-side failures.
53    pub max_failures_total: Option<usize>,
54    /// Always `true` in v1. Reserved for a future "headers-only" mode.
55    pub include_original_payload: bool,
56}
57
58impl DlqConfig {
59    /// Convenience constructor: `propagate` policy, no budgets, payload
60    /// included.
61    pub fn new(sink: Arc<dyn Sink>) -> Self {
62        Self {
63            sink,
64            on_batch_error: OnBatchError::Propagate,
65            max_failures_per_page: None,
66            max_failures_total: None,
67            include_original_payload: true,
68        }
69    }
70}
71
72impl fmt::Debug for DlqConfig {
73    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
74        f.debug_struct("DlqConfig")
75            .field("sink", &self.sink.connector_name())
76            .field("on_batch_error", &self.on_batch_error)
77            .field("max_failures_per_page", &self.max_failures_per_page)
78            .field("max_failures_total", &self.max_failures_total)
79            .field("include_original_payload", &self.include_original_payload)
80            .finish()
81    }
82}
83
84/// Counters returned alongside [`PipelineResult`](crate::PipelineResult)
85/// when a DLQ is wired.
86#[derive(Debug, Clone, Default, PartialEq, Eq)]
87pub struct DlqStats {
88    /// Total rows routed to the DLQ across the run.
89    pub records_dlq: usize,
90    /// Pages that produced at least one DLQ record.
91    pub pages_with_failures: usize,
92}
93
94/// Reason a page produced DLQ traffic. Used as a metric label and span
95/// attribute; closed-set enum so cardinality stays bounded.
96#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub enum DlqReason {
98    /// At least one per-row outcome was `Err`, surfaced by an overriding
99    /// [`Sink::write_batch_partial`].
100    Partial,
101    /// The whole batch failed and the configured policy was
102    /// [`OnBatchError::DlqAll`].
103    DlqAll,
104    /// A record was quarantined (or batch-quarantined) by a data-quality check.
105    Quality,
106    /// A record was routed to the DLQ by an `on_drift`/`on_incompatible`
107    /// quarantine policy.
108    SchemaDrift,
109    /// A record was routed to the DLQ by a data-contract `on_breach:
110    /// quarantine` policy.
111    Contract,
112}
113
114impl DlqReason {
115    /// Returns the stable Prometheus label value for this reason.
116    /// Closed-set values: `"partial"`, `"dlq_all"`, or `"quality"`.
117    pub fn as_str(self) -> &'static str {
118        match self {
119            DlqReason::Partial => "partial",
120            DlqReason::DlqAll => "dlq_all",
121            DlqReason::Quality => "quality",
122            DlqReason::SchemaDrift => "schema_drift",
123            DlqReason::Contract => "contract",
124        }
125    }
126
127    /// Every closed-set reason value, for validating a user-supplied
128    /// `--reason` filter against the exact serde strings.
129    pub const ALL: [DlqReason; 5] = [
130        DlqReason::Partial,
131        DlqReason::DlqAll,
132        DlqReason::Quality,
133        DlqReason::SchemaDrift,
134        DlqReason::Contract,
135    ];
136
137    /// Parse a reason from its stable serde string (the inverse of
138    /// [`as_str`](Self::as_str)). Returns `None` for an unknown value.
139    pub fn from_serde_str(s: &str) -> Option<DlqReason> {
140        DlqReason::ALL.into_iter().find(|r| r.as_str() == s)
141    }
142}
143
144/// Build a single DLQ envelope.
145///
146/// The schema is fixed; see the design spec for the rationale. `payload`
147/// is included verbatim — no truncation, no transformation. `reason`
148/// records *which stage* quarantined the row (as the closed-set
149/// [`DlqReason`] serde value) so tools like `faucet dlq inspect` /
150/// `faucet dlq replay` can group and filter without re-deriving it from
151/// the free-form error message. It is written as a top-level `reason`
152/// field alongside the structured `error`.
153pub fn build_envelope(
154    payload: &Value,
155    error: &FaucetError,
156    reason: DlqReason,
157    sink_name: &str,
158    pipeline_name: &str,
159    row: &str,
160    record_index: usize,
161) -> Value {
162    let kind = crate::observability::decorator::error_kind(error);
163    // The envelope is written to a file / object store, so it leaves the process:
164    // scrub any resolved secret the error text picked up (a `reqwest` error
165    // embeds the request URL, which may carry an API key in a query parameter).
166    // No-op unless the host installed a redactor (#456 H5).
167    let message = crate::redact::redact(&error.to_string());
168    // `as_millis()` returns u128. Convert via TryFrom so we saturate at
169    // i64::MAX instead of silently wrapping to a negative number. The
170    // saturation ceiling (year ~292,000,000) is impossible in practice,
171    // so this only ever fires on a corrupt clock. `unwrap_or(0)` covers
172    // the (also impossible on modern systems) clock-before-epoch case.
173    let ts_ms = SystemTime::now()
174        .duration_since(UNIX_EPOCH)
175        .map(|d| i64::try_from(d.as_millis()).unwrap_or(i64::MAX))
176        .unwrap_or(0);
177    json!({
178        "error": { "kind": kind, "message": message },
179        "reason": reason.as_str(),
180        "payload": payload,
181        "ts_ms": ts_ms,
182        "sink": sink_name,
183        "pipeline": pipeline_name,
184        "row": row,
185        "record_index": record_index,
186    })
187}
188
189/// A DLQ envelope parsed back into its original payload plus the metadata
190/// needed to inspect and replay it. Produced by [`unwrap_envelope`].
191#[derive(Debug, Clone, PartialEq)]
192pub struct UnwrappedEnvelope {
193    /// The original record that was quarantined — replayed verbatim.
194    pub payload: Value,
195    /// The stage that quarantined the row (`build_envelope`'s `reason`
196    /// field). `None` for envelopes written before the field existed.
197    pub reason: Option<String>,
198    /// The [`FaucetError`] variant name (`error.kind`), e.g. `"Sink"`,
199    /// `"QualityFailure"`. `None` if the envelope omits it.
200    pub error_kind: Option<String>,
201    /// Human-readable failure message (`error.message`), if present.
202    pub error_message: Option<String>,
203    /// Position of the record within its original page.
204    pub record_index: Option<u64>,
205    /// Pipeline name that produced the envelope, if present.
206    pub pipeline: Option<String>,
207    /// Matrix row id that produced the envelope, if present.
208    pub row: Option<String>,
209    /// Sink name the record was destined for, if present.
210    pub sink: Option<String>,
211    /// Epoch-millis timestamp the envelope was written, if present.
212    pub ts_ms: Option<i64>,
213}
214
215/// Error returned by [`unwrap_envelope`] when a value is not a usable DLQ
216/// envelope. Only the *payload* is mandatory — every other field is
217/// optional so envelopes written by older versions still replay.
218#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
219pub enum EnvelopeError {
220    /// The value was not a JSON object.
221    #[error("DLQ envelope is not a JSON object")]
222    NotObject,
223    /// The mandatory `payload` field was absent — nothing to replay.
224    #[error("DLQ envelope has no `payload` field")]
225    MissingPayload,
226}
227
228/// Parse a DLQ envelope produced by [`build_envelope`] back into its
229/// original payload plus metadata.
230///
231/// Only `payload` is required; all other fields are optional so envelopes
232/// written before a field existed still round-trip (forward-compatible
233/// read). Callers reading a DLQ location back (e.g. `faucet dlq inspect`)
234/// should treat an [`EnvelopeError`] as "skip + count", never as fatal —
235/// a DLQ file may legitimately contain arbitrary lines.
236pub fn unwrap_envelope(value: &Value) -> Result<UnwrappedEnvelope, EnvelopeError> {
237    let obj = value.as_object().ok_or(EnvelopeError::NotObject)?;
238    let payload = obj.get("payload").ok_or(EnvelopeError::MissingPayload)?;
239    let error = obj.get("error").and_then(|e| e.as_object());
240    let str_field = |k: &str| obj.get(k).and_then(|v| v.as_str()).map(str::to_owned);
241    Ok(UnwrappedEnvelope {
242        payload: payload.clone(),
243        reason: str_field("reason"),
244        error_kind: error
245            .and_then(|e| e.get("kind"))
246            .and_then(|v| v.as_str())
247            .map(str::to_owned),
248        error_message: error
249            .and_then(|e| e.get("message"))
250            .and_then(|v| v.as_str())
251            .map(str::to_owned),
252        record_index: obj.get("record_index").and_then(Value::as_u64),
253        pipeline: str_field("pipeline"),
254        row: str_field("row"),
255        sink: str_field("sink"),
256        ts_ms: obj.get("ts_ms").and_then(Value::as_i64),
257    })
258}
259
260/// What a sink promises about one failed batch write (#737): whether rows of a
261/// batch whose write failed may already have landed.
262///
263/// It decides whether [`OnBatchError::DlqAll`] is safe: routing a failed batch
264/// to the DLQ only avoids duplicates when nothing of it committed, because a
265/// DLQ replay writes every routed row again.
266#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
267#[serde(rename_all = "snake_case")]
268#[non_exhaustive]
269pub enum BatchAtomicity {
270    /// All or nothing: a failed write lands no row (one statement, a
271    /// transaction around every chunk, one object upload, one table commit).
272    Atomic,
273    /// Per-row outcomes through [`Sink::write_batch_partial`]; an outer `Err`
274    /// from it means no row of that call committed.
275    PerRow,
276    /// A failed write may have committed some rows and reports no per-row
277    /// detail. The default, because it is the only safe assumption.
278    #[default]
279    BestEffort,
280}
281
282impl BatchAtomicity {
283    /// Stable label (`atomic` / `per_row` / `best_effort`).
284    pub fn as_str(self) -> &'static str {
285        match self {
286            BatchAtomicity::Atomic => "atomic",
287            BatchAtomicity::PerRow => "per_row",
288            BatchAtomicity::BestEffort => "best_effort",
289        }
290    }
291}
292
293impl fmt::Display for BatchAtomicity {
294    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
295        f.write_str(self.as_str())
296    }
297}
298
299/// Whether [`OnBatchError::DlqAll`] can route a failed batch without
300/// duplicating rows downstream: the sink lands nothing on failure, or writes
301/// by key so a replayed row overwrites itself.
302pub fn dlq_all_is_safe(atomicity: BatchAtomicity, dedups_by_key: bool) -> bool {
303    dedups_by_key || !matches!(atomicity, BatchAtomicity::BestEffort)
304}
305
306/// The refusal for `dlq_all` on a sink that may commit part of a failed batch.
307pub fn dlq_all_refusal(sink: &str, atomicity: BatchAtomicity) -> FaucetError {
308    FaucetError::Config(format!(
309        "dlq: on_batch_error 'dlq_all' is unsafe with sink '{sink}' (batch atomicity \
310         '{atomicity}'): a failed batch may already have written some rows, and replaying \
311         the DLQ would write them again. Use on_batch_error 'propagate', configure the sink \
312         with write_mode 'upsert' and a key so a replay overwrites instead of duplicating, or \
313         set allow_duplicates_on_dlq_all: true to accept the duplicates"
314    ))
315}
316
317/// Refuse [`OnBatchError::DlqAll`] against a sink that may commit part of a
318/// failed batch, unless the caller explicitly accepts duplicates.
319pub fn check_dlq_all_policy(
320    sink: &dyn Sink,
321    on_batch_error: OnBatchError,
322    allow_duplicates: bool,
323) -> Result<(), FaucetError> {
324    if on_batch_error != OnBatchError::DlqAll || allow_duplicates {
325        return Ok(());
326    }
327    let atomicity = sink.batch_atomicity();
328    if dlq_all_is_safe(atomicity, sink.dedups_by_key()) {
329        Ok(())
330    } else {
331        Err(dlq_all_refusal(sink.connector_name(), atomicity))
332    }
333}
334
335/// How one sink write (a page, or an adaptive sub-batch of one) ended (#737).
336#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
337#[non_exhaustive]
338pub enum BatchOutcome {
339    /// Every row committed.
340    Committed,
341    /// Some rows committed; the sink reported the others per row and they went
342    /// to the DLQ.
343    DlqPartial,
344    /// The whole write failed and `on_batch_error: dlq_all` sent every row to
345    /// the DLQ.
346    DlqAll,
347    /// The write failed and the error propagated (or was retried).
348    Failed,
349}
350
351impl BatchOutcome {
352    /// Stable label (`committed` / `dlq_partial` / `dlq_all` / `failed`).
353    pub fn as_str(self) -> &'static str {
354        match self {
355            BatchOutcome::Committed => "committed",
356            BatchOutcome::DlqPartial => "dlq_partial",
357            BatchOutcome::DlqAll => "dlq_all",
358            BatchOutcome::Failed => "failed",
359        }
360    }
361}
362
363/// Per-run batch outcome counts, as reported on a run (#737). `attempted` is
364/// the sum of the other four.
365#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
366#[serde(default)]
367pub struct BatchOutcomes {
368    /// Sink writes attempted.
369    pub attempted: u64,
370    /// Writes where every row committed.
371    pub committed: u64,
372    /// Writes where some rows failed per row and went to the DLQ.
373    pub dlq_partial: u64,
374    /// Failed writes routed whole to the DLQ by `on_batch_error: dlq_all`.
375    pub dlq_all: u64,
376    /// Failed writes whose error propagated.
377    pub failed: u64,
378}
379
380impl BatchOutcomes {
381    /// Whether no write was attempted.
382    pub fn is_empty(&self) -> bool {
383        self.attempted == 0
384    }
385
386    /// Writes that did not fully commit.
387    pub fn unclean(&self) -> u64 {
388        self.dlq_partial + self.dlq_all + self.failed
389    }
390}
391
392/// Shared counters a pipeline run fills in as it writes (#737). Attach with
393/// [`Pipeline::with_batch_outcomes`](crate::Pipeline::with_batch_outcomes) and
394/// read [`snapshot`](Self::snapshot) afterwards — on failure too.
395#[derive(Debug, Default)]
396pub struct BatchOutcomeCounters {
397    committed: std::sync::atomic::AtomicU64,
398    dlq_partial: std::sync::atomic::AtomicU64,
399    dlq_all: std::sync::atomic::AtomicU64,
400    failed: std::sync::atomic::AtomicU64,
401}
402
403impl BatchOutcomeCounters {
404    /// Empty counters.
405    pub fn new() -> Self {
406        Self::default()
407    }
408
409    /// Count one write.
410    pub fn record(&self, outcome: BatchOutcome) {
411        use std::sync::atomic::Ordering::Relaxed;
412        let cell = match outcome {
413            BatchOutcome::Committed => &self.committed,
414            BatchOutcome::DlqPartial => &self.dlq_partial,
415            BatchOutcome::DlqAll => &self.dlq_all,
416            BatchOutcome::Failed => &self.failed,
417        };
418        cell.fetch_add(1, Relaxed);
419    }
420
421    /// The counts so far.
422    pub fn snapshot(&self) -> BatchOutcomes {
423        use std::sync::atomic::Ordering::Relaxed;
424        let committed = self.committed.load(Relaxed);
425        let dlq_partial = self.dlq_partial.load(Relaxed);
426        let dlq_all = self.dlq_all.load(Relaxed);
427        let failed = self.failed.load(Relaxed);
428        BatchOutcomes {
429            attempted: committed + dlq_partial + dlq_all + failed,
430            committed,
431            dlq_partial,
432            dlq_all,
433            failed,
434        }
435    }
436}
437
438/// Where one pipeline run reports its batch outcomes: the
439/// `faucet_batch_outcomes_total` counter and, when attached, the caller's
440/// [`BatchOutcomeCounters`].
441#[derive(Debug, Clone)]
442pub(crate) struct BatchOutcomeSink {
443    labels: Vec<metrics::Label>,
444    counters: Option<Arc<BatchOutcomeCounters>>,
445}
446
447impl BatchOutcomeSink {
448    pub(crate) fn new(
449        pipeline: &str,
450        row: &str,
451        sink: &str,
452        counters: Option<Arc<BatchOutcomeCounters>>,
453    ) -> Self {
454        use metrics::{Label, SharedString};
455        Self {
456            labels: vec![
457                Label::new("pipeline", SharedString::from(pipeline.to_string())),
458                Label::new("row", SharedString::from(row.to_string())),
459                Label::new("sink", SharedString::from(sink.to_string())),
460            ],
461            counters,
462        }
463    }
464
465    pub(crate) fn record(&self, outcome: BatchOutcome) {
466        let mut labels = self.labels.clone();
467        labels.push(metrics::Label::new(
468            "outcome",
469            metrics::SharedString::const_str(outcome.as_str()),
470        ));
471        metrics::counter!("faucet_batch_outcomes_total", labels).increment(1);
472        if let Some(c) = &self.counters {
473            c.record(outcome);
474        }
475    }
476
477    /// Record the outcome of a plain (non-partial) write and pass it through.
478    pub(crate) fn observe<T>(&self, result: Result<T, FaucetError>) -> Result<T, FaucetError> {
479        self.record(if result.is_ok() {
480            BatchOutcome::Committed
481        } else {
482            BatchOutcome::Failed
483        });
484        result
485    }
486}
487
488#[cfg(test)]
489mod tests {
490    use super::*;
491
492    struct AtomSink {
493        atomicity: BatchAtomicity,
494        keyed: bool,
495    }
496
497    #[async_trait::async_trait]
498    impl Sink for AtomSink {
499        async fn write_batch(&self, r: &[Value]) -> Result<usize, FaucetError> {
500            Ok(r.len())
501        }
502        fn batch_atomicity(&self) -> BatchAtomicity {
503            self.atomicity
504        }
505        fn dedups_by_key(&self) -> bool {
506            self.keyed
507        }
508        fn connector_name(&self) -> &'static str {
509            "atom"
510        }
511    }
512
513    #[test]
514    fn batch_atomicity_labels_and_default() {
515        assert_eq!(BatchAtomicity::default(), BatchAtomicity::BestEffort);
516        assert_eq!(BatchAtomicity::Atomic.as_str(), "atomic");
517        assert_eq!(BatchAtomicity::PerRow.to_string(), "per_row");
518        assert_eq!(BatchAtomicity::BestEffort.as_str(), "best_effort");
519        assert_eq!(
520            serde_json::to_value(BatchAtomicity::PerRow).unwrap(),
521            json!("per_row")
522        );
523    }
524
525    #[test]
526    fn dlq_all_is_safe_unless_best_effort_and_unkeyed() {
527        assert!(dlq_all_is_safe(BatchAtomicity::Atomic, false));
528        assert!(dlq_all_is_safe(BatchAtomicity::PerRow, false));
529        assert!(dlq_all_is_safe(BatchAtomicity::BestEffort, true));
530        assert!(!dlq_all_is_safe(BatchAtomicity::BestEffort, false));
531    }
532
533    #[test]
534    fn check_dlq_all_policy_decision_paths() {
535        let best = AtomSink {
536            atomicity: BatchAtomicity::BestEffort,
537            keyed: false,
538        };
539        assert!(check_dlq_all_policy(&best, OnBatchError::Propagate, false).is_ok());
540        assert!(check_dlq_all_policy(&best, OnBatchError::DlqAll, true).is_ok());
541        let err = check_dlq_all_policy(&best, OnBatchError::DlqAll, false).unwrap_err();
542        let msg = err.to_string();
543        assert!(
544            msg.contains("'atom'") && msg.contains("best_effort"),
545            "{msg}"
546        );
547        assert!(msg.contains("allow_duplicates_on_dlq_all"), "{msg}");
548        let keyed = AtomSink {
549            atomicity: BatchAtomicity::BestEffort,
550            keyed: true,
551        };
552        assert!(check_dlq_all_policy(&keyed, OnBatchError::DlqAll, false).is_ok());
553        let atomic = AtomSink {
554            atomicity: BatchAtomicity::Atomic,
555            keyed: false,
556        };
557        assert!(check_dlq_all_policy(&atomic, OnBatchError::DlqAll, false).is_ok());
558    }
559
560    #[tokio::test]
561    async fn atom_sink_writes_every_record() {
562        let sink = AtomSink {
563            atomicity: BatchAtomicity::Atomic,
564            keyed: false,
565        };
566        assert_eq!(sink.write_batch(&[json!(1), json!(2)]).await.unwrap(), 2);
567    }
568
569    #[test]
570    fn batch_outcome_counters_snapshot() {
571        let c = BatchOutcomeCounters::new();
572        assert!(c.snapshot().is_empty());
573        for o in [
574            BatchOutcome::Committed,
575            BatchOutcome::Committed,
576            BatchOutcome::DlqPartial,
577            BatchOutcome::DlqAll,
578            BatchOutcome::Failed,
579        ] {
580            c.record(o);
581        }
582        let s = c.snapshot();
583        assert_eq!(
584            s,
585            BatchOutcomes {
586                attempted: 5,
587                committed: 2,
588                dlq_partial: 1,
589                dlq_all: 1,
590                failed: 1,
591            }
592        );
593        assert_eq!(s.unclean(), 3);
594        assert!(!s.is_empty());
595        assert_eq!(BatchOutcome::DlqPartial.as_str(), "dlq_partial");
596        assert_eq!(BatchOutcome::DlqAll.as_str(), "dlq_all");
597        assert_eq!(BatchOutcome::Failed.as_str(), "failed");
598        assert_eq!(BatchOutcome::Committed.as_str(), "committed");
599    }
600
601    #[test]
602    fn batch_outcome_sink_observes_and_counts() {
603        let counters = Arc::new(BatchOutcomeCounters::new());
604        let sink = BatchOutcomeSink::new("p", "r", "s", Some(Arc::clone(&counters)));
605        assert_eq!(sink.observe(Ok::<usize, FaucetError>(3)).unwrap(), 3);
606        assert!(
607            sink.observe(Err::<usize, _>(FaucetError::Sink("x".into())))
608                .is_err()
609        );
610        sink.record(BatchOutcome::DlqAll);
611        let snap = counters.snapshot();
612        assert_eq!((snap.committed, snap.failed, snap.dlq_all), (1, 1, 1));
613        BatchOutcomeSink::new("p", "r", "s", None).record(BatchOutcome::Committed);
614    }
615
616    #[test]
617    fn envelope_has_all_required_fields() {
618        let payload = json!({"user_id": 7, "name": "Alice"});
619        let err = FaucetError::Sink("row rejected: bad timestamp".into());
620        let env = build_envelope(
621            &payload,
622            &err,
623            DlqReason::Partial,
624            "bigquery",
625            "users_etl",
626            "us",
627            3,
628        );
629
630        assert_eq!(env["error"]["kind"], "Sink");
631        assert_eq!(env["reason"], "partial");
632        assert!(
633            env["error"]["message"]
634                .as_str()
635                .unwrap()
636                .contains("row rejected")
637        );
638        assert_eq!(env["payload"], payload);
639        assert!(env["ts_ms"].as_i64().unwrap() > 0);
640        assert_eq!(env["sink"], "bigquery");
641        assert_eq!(env["pipeline"], "users_etl");
642        assert_eq!(env["row"], "us");
643        assert_eq!(env["record_index"], 3);
644    }
645
646    #[test]
647    fn envelope_preserves_payload_byte_for_byte() {
648        let payload = json!({
649            "nested": { "a": [1, 2, 3], "b": null, "c": true },
650            "unicode": "café — résumé"
651        });
652        let env = build_envelope(
653            &payload,
654            &FaucetError::Sink("x".into()),
655            DlqReason::Quality,
656            "s",
657            "p",
658            "",
659            0,
660        );
661        assert_eq!(env["payload"], payload);
662    }
663
664    #[test]
665    fn envelope_empty_row_serializes_as_empty_string() {
666        let env = build_envelope(
667            &json!({}),
668            &FaucetError::Sink("x".into()),
669            DlqReason::DlqAll,
670            "s",
671            "",
672            "",
673            0,
674        );
675        assert_eq!(env["row"], "");
676        assert_eq!(env["pipeline"], "");
677    }
678
679    #[test]
680    fn dlq_reason_from_serde_str_round_trips() {
681        for r in DlqReason::ALL {
682            assert_eq!(DlqReason::from_serde_str(r.as_str()), Some(r));
683        }
684        assert_eq!(DlqReason::from_serde_str("nope"), None);
685        assert_eq!(DlqReason::from_serde_str("sink_error"), None);
686    }
687
688    #[test]
689    fn unwrap_envelope_round_trips_build_envelope() {
690        let payload = json!({"id": 42, "name": "Zoe"});
691        let err = FaucetError::QualityFailure {
692            check: "not_null(email)".into(),
693            message: "email is null".into(),
694        };
695        let env = build_envelope(&payload, &err, DlqReason::Quality, "pg", "etl", "eu", 5);
696        let u = unwrap_envelope(&env).expect("valid envelope");
697        assert_eq!(u.payload, payload);
698        assert_eq!(u.reason.as_deref(), Some("quality"));
699        assert_eq!(u.error_kind.as_deref(), Some("QualityFailure"));
700        assert!(u.error_message.unwrap().contains("email is null"));
701        assert_eq!(u.record_index, Some(5));
702        assert_eq!(u.pipeline.as_deref(), Some("etl"));
703        assert_eq!(u.row.as_deref(), Some("eu"));
704        assert_eq!(u.sink.as_deref(), Some("pg"));
705        assert!(u.ts_ms.unwrap() > 0);
706    }
707
708    #[test]
709    fn unwrap_envelope_tolerates_legacy_envelope_without_reason() {
710        // An envelope written before `reason`/`error` existed still yields its
711        // payload; the missing metadata comes back as `None`, never a panic.
712        let legacy = json!({ "payload": { "x": 1 } });
713        let u = unwrap_envelope(&legacy).expect("payload present");
714        assert_eq!(u.payload, json!({ "x": 1 }));
715        assert_eq!(u.reason, None);
716        assert_eq!(u.error_kind, None);
717        assert_eq!(u.record_index, None);
718    }
719
720    #[test]
721    fn unwrap_envelope_errors_on_non_object_and_missing_payload() {
722        assert_eq!(
723            unwrap_envelope(&json!("just a string")),
724            Err(EnvelopeError::NotObject)
725        );
726        assert_eq!(
727            unwrap_envelope(&json!([1, 2, 3])),
728            Err(EnvelopeError::NotObject)
729        );
730        assert_eq!(
731            unwrap_envelope(&json!({ "error": { "kind": "Sink" } })),
732            Err(EnvelopeError::MissingPayload)
733        );
734    }
735
736    #[test]
737    fn on_batch_error_defaults_to_propagate() {
738        assert_eq!(OnBatchError::default(), OnBatchError::Propagate);
739    }
740
741    #[test]
742    fn on_batch_error_serializes_snake_case() {
743        let prop = serde_json::to_string(&OnBatchError::Propagate).unwrap();
744        let all = serde_json::to_string(&OnBatchError::DlqAll).unwrap();
745        assert_eq!(prop, "\"propagate\"");
746        assert_eq!(all, "\"dlq_all\"");
747    }
748
749    #[test]
750    fn on_batch_error_deserializes_snake_case() {
751        let prop: OnBatchError = serde_json::from_str("\"propagate\"").unwrap();
752        let all: OnBatchError = serde_json::from_str("\"dlq_all\"").unwrap();
753        assert_eq!(prop, OnBatchError::Propagate);
754        assert_eq!(all, OnBatchError::DlqAll);
755    }
756
757    #[test]
758    fn dlq_reason_strings() {
759        assert_eq!(DlqReason::Partial.as_str(), "partial");
760        assert_eq!(DlqReason::DlqAll.as_str(), "dlq_all");
761    }
762
763    #[test]
764    fn dlq_reason_quality_string() {
765        assert_eq!(DlqReason::Quality.as_str(), "quality");
766    }
767
768    #[test]
769    fn dlq_reason_schema_drift_string() {
770        assert_eq!(DlqReason::SchemaDrift.as_str(), "schema_drift");
771    }
772
773    #[test]
774    fn dlq_reason_contract_string() {
775        assert_eq!(DlqReason::Contract.as_str(), "contract");
776    }
777}