1use 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#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
21#[serde(rename_all = "snake_case")]
22pub enum OnBatchError {
23 #[default]
25 Propagate,
26 DlqAll,
32}
33
34#[derive(Clone)]
36pub struct DlqConfig {
37 pub sink: Arc<dyn Sink>,
39 pub on_batch_error: OnBatchError,
41 pub max_failures_per_page: Option<usize>,
48 pub max_failures_total: Option<usize>,
54 pub include_original_payload: bool,
56}
57
58impl DlqConfig {
59 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#[derive(Debug, Clone, Default, PartialEq, Eq)]
87pub struct DlqStats {
88 pub records_dlq: usize,
90 pub pages_with_failures: usize,
92}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
97pub enum DlqReason {
98 Partial,
101 DlqAll,
104 Quality,
106 SchemaDrift,
109 Contract,
112}
113
114impl DlqReason {
115 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 pub const ALL: [DlqReason; 5] = [
130 DlqReason::Partial,
131 DlqReason::DlqAll,
132 DlqReason::Quality,
133 DlqReason::SchemaDrift,
134 DlqReason::Contract,
135 ];
136
137 pub fn from_serde_str(s: &str) -> Option<DlqReason> {
140 DlqReason::ALL.into_iter().find(|r| r.as_str() == s)
141 }
142}
143
144pub 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 let message = crate::redact::redact(&error.to_string());
168 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#[derive(Debug, Clone, PartialEq)]
192pub struct UnwrappedEnvelope {
193 pub payload: Value,
195 pub reason: Option<String>,
198 pub error_kind: Option<String>,
201 pub error_message: Option<String>,
203 pub record_index: Option<u64>,
205 pub pipeline: Option<String>,
207 pub row: Option<String>,
209 pub sink: Option<String>,
211 pub ts_ms: Option<i64>,
213}
214
215#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
219pub enum EnvelopeError {
220 #[error("DLQ envelope is not a JSON object")]
222 NotObject,
223 #[error("DLQ envelope has no `payload` field")]
225 MissingPayload,
226}
227
228pub 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema)]
267#[serde(rename_all = "snake_case")]
268#[non_exhaustive]
269pub enum BatchAtomicity {
270 Atomic,
273 PerRow,
276 #[default]
279 BestEffort,
280}
281
282impl BatchAtomicity {
283 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
299pub fn dlq_all_is_safe(atomicity: BatchAtomicity, dedups_by_key: bool) -> bool {
303 dedups_by_key || !matches!(atomicity, BatchAtomicity::BestEffort)
304}
305
306pub 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
317pub 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
337#[non_exhaustive]
338pub enum BatchOutcome {
339 Committed,
341 DlqPartial,
344 DlqAll,
347 Failed,
349}
350
351impl BatchOutcome {
352 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
366#[serde(default)]
367pub struct BatchOutcomes {
368 pub attempted: u64,
370 pub committed: u64,
372 pub dlq_partial: u64,
374 pub dlq_all: u64,
376 pub failed: u64,
378}
379
380impl BatchOutcomes {
381 pub fn is_empty(&self) -> bool {
383 self.attempted == 0
384 }
385
386 pub fn unclean(&self) -> u64 {
388 self.dlq_partial + self.dlq_all + self.failed
389 }
390}
391
392#[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 pub fn new() -> Self {
406 Self::default()
407 }
408
409 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 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#[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 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 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}