Skip to main content

faucet_core/
metadata.rs

1//! Optional `_faucet_*` run/lineage metadata columns (#510).
2//!
3//! An opt-in, connector-agnostic **sink decorator** that stamps a small set of
4//! metadata columns onto every row — the faucet-native analogue of Singer's
5//! `_sdc_*` columns, for freshness/lineage and debugging. Because it's a
6//! decorator (like [`CleanupTracker`](crate::cleanup::CleanupTracker)) it works
7//! for every sink without per-connector code.
8//!
9//! ```yaml
10//! metadata_columns:
11//!   prefix: "_faucet"                       # column-name prefix (default)
12//!   columns: [extracted_at, loaded_at, run_id, source]
13//! ```
14
15use crate::error::FaucetError;
16use crate::traits::{RowOutcome, Sink};
17use chrono::Utc;
18use schemars::JsonSchema;
19use serde::{Deserialize, Serialize};
20use serde_json::{Map, Value};
21use std::collections::BTreeMap;
22use std::sync::atomic::{AtomicU64, Ordering};
23
24/// A metadata column to stamp onto every row.
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
26#[serde(rename_all = "snake_case")]
27pub enum MetadataColumn {
28    /// When the record entered the pipeline (ingest time). In the sink-decorator
29    /// model this is captured just before the write, so it's ≈ `loaded_at`.
30    ExtractedAt,
31    /// Sink write time.
32    LoadedAt,
33    /// The pipeline run id.
34    RunId,
35    /// The source connector kind / stream id.
36    Source,
37    /// A monotonic per-run row ordinal.
38    Sequence,
39}
40
41impl MetadataColumn {
42    /// The column-name suffix (appended to the prefix as `{prefix}_{suffix}`).
43    pub fn suffix(self) -> &'static str {
44        match self {
45            MetadataColumn::ExtractedAt => "extracted_at",
46            MetadataColumn::LoadedAt => "loaded_at",
47            MetadataColumn::RunId => "run_id",
48            MetadataColumn::Source => "source",
49            MetadataColumn::Sequence => "sequence",
50        }
51    }
52}
53
54fn default_prefix() -> String {
55    "_faucet".to_owned()
56}
57fn default_true() -> bool {
58    true
59}
60
61/// The `metadata_columns:` config block.
62#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
63#[serde(deny_unknown_fields)]
64pub struct MetadataColumnsSpec {
65    /// Master switch (default `true` when the block is present).
66    #[serde(default = "default_true")]
67    pub enabled: bool,
68    /// Column-name prefix (default `_faucet`).
69    #[serde(default = "default_prefix")]
70    pub prefix: String,
71    /// Columns to stamp. Empty = the default set
72    /// (`extracted_at`, `loaded_at`, `run_id`, `source`).
73    #[serde(default)]
74    pub columns: Vec<MetadataColumn>,
75}
76
77impl Default for MetadataColumnsSpec {
78    fn default() -> Self {
79        Self {
80            enabled: true,
81            prefix: default_prefix(),
82            columns: Vec::new(),
83        }
84    }
85}
86
87const DEFAULT_COLUMNS: &[MetadataColumn] = &[
88    MetadataColumn::ExtractedAt,
89    MetadataColumn::LoadedAt,
90    MetadataColumn::RunId,
91    MetadataColumn::Source,
92];
93
94/// A validated metadata-columns policy.
95#[derive(Debug, Clone)]
96pub struct CompiledMetadata {
97    prefix: String,
98    columns: Vec<MetadataColumn>,
99}
100
101impl CompiledMetadata {
102    /// Compile a spec: resolve the default column set and validate the prefix.
103    /// Returns `Ok(None)` when disabled (so the caller skips the decorator).
104    pub fn compile(spec: &MetadataColumnsSpec) -> Result<Option<Self>, FaucetError> {
105        if !spec.enabled {
106            return Ok(None);
107        }
108        if spec.prefix.trim().is_empty() {
109            return Err(FaucetError::Config(
110                "metadata_columns: `prefix` must not be empty".into(),
111            ));
112        }
113        let columns = if spec.columns.is_empty() {
114            DEFAULT_COLUMNS.to_vec()
115        } else {
116            spec.columns.clone()
117        };
118        Ok(Some(Self {
119            prefix: spec.prefix.clone(),
120            columns,
121        }))
122    }
123
124    fn column_name(&self, col: MetadataColumn) -> String {
125        format!("{}_{}", self.prefix, col.suffix())
126    }
127}
128
129/// Per-run values the decorator stamps.
130#[derive(Debug, Clone)]
131pub struct MetadataContext {
132    /// The pipeline run id.
133    pub run_id: String,
134    /// The source connector kind / stream id.
135    pub source: String,
136}
137
138/// A [`Sink`] decorator that stamps `_faucet_*` metadata columns onto every row
139/// before delegating the write.
140pub struct MetadataSink {
141    inner: Box<dyn Sink>,
142    meta: CompiledMetadata,
143    ctx: MetadataContext,
144    seq: AtomicU64,
145}
146
147impl std::fmt::Debug for MetadataSink {
148    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
149        f.debug_struct("MetadataSink")
150            .field("prefix", &self.meta.prefix)
151            .field("columns", &self.meta.columns)
152            .finish()
153    }
154}
155
156impl MetadataSink {
157    /// Wrap a sink so every written row carries the configured metadata columns.
158    pub fn new(inner: Box<dyn Sink>, meta: CompiledMetadata, ctx: MetadataContext) -> Self {
159        Self {
160            inner,
161            meta,
162            ctx,
163            seq: AtomicU64::new(0),
164        }
165    }
166
167    /// Return a copy of `records` with the metadata columns injected. Non-object
168    /// records pass through unchanged (nowhere to stamp).
169    fn stamp(&self, records: &[Value]) -> Vec<Value> {
170        let now = Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true);
171        records
172            .iter()
173            .map(|rec| match rec {
174                Value::Object(map) => {
175                    let mut map = map.clone();
176                    self.inject(&mut map, &now);
177                    Value::Object(map)
178                }
179                other => other.clone(),
180            })
181            .collect()
182    }
183
184    fn inject(&self, map: &mut Map<String, Value>, now: &str) {
185        for &col in &self.meta.columns {
186            let name = self.meta.column_name(col);
187            let value = match col {
188                MetadataColumn::ExtractedAt | MetadataColumn::LoadedAt => {
189                    Value::String(now.to_owned())
190                }
191                MetadataColumn::RunId => Value::String(self.ctx.run_id.clone()),
192                MetadataColumn::Source => Value::String(self.ctx.source.clone()),
193                MetadataColumn::Sequence => Value::from(self.seq.fetch_add(1, Ordering::Relaxed)),
194            };
195            map.insert(name, value);
196        }
197    }
198}
199
200#[async_trait::async_trait]
201impl Sink for MetadataSink {
202    async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
203        self.inner.write_batch(&self.stamp(records)).await
204    }
205
206    async fn write_batch_partial(&self, records: &[Value]) -> Result<Vec<RowOutcome>, FaucetError> {
207        self.inner.write_batch_partial(&self.stamp(records)).await
208    }
209
210    async fn write_batch_idempotent(
211        &self,
212        records: &[Value],
213        scope: &str,
214        token: &str,
215    ) -> Result<usize, FaucetError> {
216        self.inner
217            .write_batch_idempotent(&self.stamp(records), scope, token)
218            .await
219    }
220
221    async fn flush(&self) -> Result<(), FaucetError> {
222        self.inner.flush().await
223    }
224
225    // ── Pure forwarding ──────────────────────────────────────────────────────
226    async fn local_outputs(&self) -> Vec<crate::local_outputs::LocalOutput> {
227        self.inner.local_outputs().await
228    }
229    async fn check(
230        &self,
231        ctx: &crate::check::CheckContext,
232    ) -> Result<crate::check::CheckReport, FaucetError> {
233        self.inner.check(ctx).await
234    }
235    fn supports_cleanup(&self) -> bool {
236        self.inner.supports_cleanup()
237    }
238    async fn cleanup_scope(
239        &self,
240        scope: &BTreeMap<String, Value>,
241        seen: &crate::cleanup::SeenKeys,
242    ) -> Result<u64, FaucetError> {
243        self.inner.cleanup_scope(scope, seen).await
244    }
245    fn supports_idempotent_writes(&self) -> bool {
246        self.inner.supports_idempotent_writes()
247    }
248    async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
249        self.inner.last_committed_token(scope).await
250    }
251    fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
252        self.inner.supported_write_modes()
253    }
254    fn dedups_by_key(&self) -> bool {
255        self.inner.dedups_by_key()
256    }
257    fn batch_atomicity(&self) -> crate::dlq::BatchAtomicity {
258        self.inner.batch_atomicity()
259    }
260    fn sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
261        self.inner.sink_guarantee()
262    }
263    async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
264        self.inner.current_schema().await
265    }
266    fn supports_schema_evolution(&self) -> bool {
267        self.inner.supports_schema_evolution()
268    }
269    async fn evolve_schema(
270        &self,
271        evolution: &crate::drift::SchemaEvolution,
272    ) -> Result<(), FaucetError> {
273        self.inner.evolve_schema(evolution).await
274    }
275    fn config_schema(&self) -> Value {
276        self.inner.config_schema()
277    }
278    fn connector_name(&self) -> &'static str {
279        self.inner.connector_name()
280    }
281    fn dataset_uri(&self) -> String {
282        self.inner.dataset_uri()
283    }
284    fn is_overwrite(&self) -> bool {
285        self.inner.is_overwrite()
286    }
287    async fn begin_overwrite(&self) -> Result<(), FaucetError> {
288        self.inner.begin_overwrite().await
289    }
290    async fn commit_overwrite(&self) -> Result<(), FaucetError> {
291        self.inner.commit_overwrite().await
292    }
293    async fn abort_overwrite(&self) -> Result<(), FaucetError> {
294        self.inner.abort_overwrite().await
295    }
296    async fn complete_run(&self) -> Result<(), FaucetError> {
297        self.inner.complete_run().await
298    }
299}
300
301#[cfg(test)]
302mod tests {
303    use super::*;
304    use serde_json::json;
305    use std::sync::Mutex;
306
307    #[derive(Debug, Default)]
308    struct CapturingSink {
309        rows: Mutex<Vec<Value>>,
310    }
311    #[async_trait::async_trait]
312    impl Sink for CapturingSink {
313        async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
314            self.rows.lock().unwrap().extend_from_slice(records);
315            Ok(records.len())
316        }
317        fn config_schema(&self) -> Value {
318            json!({})
319        }
320    }
321
322    fn spec(cols: &[MetadataColumn]) -> MetadataColumnsSpec {
323        MetadataColumnsSpec {
324            enabled: true,
325            prefix: "_faucet".into(),
326            columns: cols.to_vec(),
327        }
328    }
329
330    fn ctx() -> MetadataContext {
331        MetadataContext {
332            run_id: "run-42".into(),
333            source: "rest".into(),
334        }
335    }
336
337    #[test]
338    fn compile_resolves_defaults_and_respects_disable() {
339        let c = CompiledMetadata::compile(&spec(&[])).unwrap().unwrap();
340        assert_eq!(c.columns, DEFAULT_COLUMNS.to_vec());
341        // disabled → None
342        let disabled = MetadataColumnsSpec {
343            enabled: false,
344            ..Default::default()
345        };
346        assert!(CompiledMetadata::compile(&disabled).unwrap().is_none());
347        // empty prefix → error
348        let bad = MetadataColumnsSpec {
349            prefix: " ".into(),
350            ..Default::default()
351        };
352        assert!(CompiledMetadata::compile(&bad).is_err());
353    }
354
355    #[tokio::test]
356    async fn stamps_all_column_kinds_with_prefix() {
357        let meta = CompiledMetadata::compile(&spec(&[
358            MetadataColumn::ExtractedAt,
359            MetadataColumn::LoadedAt,
360            MetadataColumn::RunId,
361            MetadataColumn::Source,
362            MetadataColumn::Sequence,
363        ]))
364        .unwrap()
365        .unwrap();
366        let inner = Box::new(CapturingSink::default());
367        let sink = MetadataSink::new(inner, meta, ctx());
368        assert_eq!(
369            sink.batch_atomicity(),
370            crate::dlq::BatchAtomicity::BestEffort
371        );
372        let n = sink
373            .write_batch(&[json!({"id": 1}), json!({"id": 2})])
374            .await
375            .unwrap();
376        assert_eq!(n, 2);
377        // Re-wrap to read the captured rows would need the inner; instead assert
378        // via a fresh capturing sink shared by Arc — simpler: re-stamp directly.
379        let stamped = sink.stamp(&[json!({"id": 1}), json!({"id": 2})]);
380        let r0 = &stamped[0];
381        assert_eq!(r0["id"], 1);
382        assert_eq!(r0["_faucet_run_id"], "run-42");
383        assert_eq!(r0["_faucet_source"], "rest");
384        assert!(r0["_faucet_extracted_at"].is_string());
385        assert!(r0["_faucet_loaded_at"].is_string());
386        // Sequence is monotonic across records.
387        assert_eq!(
388            stamped[0]["_faucet_sequence"].as_u64().unwrap() + 1,
389            stamped[1]["_faucet_sequence"].as_u64().unwrap()
390        );
391    }
392
393    #[test]
394    fn non_object_records_pass_through() {
395        let meta = CompiledMetadata::compile(&spec(&[MetadataColumn::RunId]))
396            .unwrap()
397            .unwrap();
398        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
399        let out = sink.stamp(&[json!("scalar"), json!(42)]);
400        assert_eq!(out, vec![json!("scalar"), json!(42)]);
401    }
402
403    #[tokio::test]
404    async fn delegates_every_sink_method_to_inner() {
405        let meta = CompiledMetadata::compile(&spec(&[MetadataColumn::RunId]))
406            .unwrap()
407            .unwrap();
408        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
409
410        // Write paths stamp + forward.
411        assert_eq!(sink.write_batch(&[json!({"a": 1})]).await.unwrap(), 1);
412        assert_eq!(
413            sink.write_batch_partial(&[json!({"a": 1})])
414                .await
415                .unwrap()
416                .len(),
417            1
418        );
419        let _ = sink
420            .write_batch_idempotent(&[json!({"a": 1})], "scope", "tok")
421            .await;
422        sink.flush().await.unwrap();
423
424        // Pure forwards (inner uses trait defaults).
425        let _ = sink.check(&crate::check::CheckContext::default()).await;
426        assert!(!sink.supports_cleanup());
427        let _ = sink
428            .cleanup_scope(&BTreeMap::new(), &crate::cleanup::SeenKeys::new())
429            .await;
430        assert!(!sink.supports_idempotent_writes());
431        assert!(sink.last_committed_token("s").await.unwrap().is_none());
432        let _ = sink.supported_write_modes();
433        let _ = sink.dedups_by_key();
434        let _ = sink.sink_guarantee();
435        let _ = sink.current_schema().await.unwrap();
436        let _ = sink.supports_schema_evolution();
437        let _ = sink
438            .evolve_schema(&crate::drift::SchemaEvolution::default())
439            .await;
440        let _ = sink.config_schema();
441        let _ = sink.connector_name();
442        let _ = sink.dataset_uri();
443        assert!(!sink.is_overwrite());
444        let _ = sink.begin_overwrite().await;
445        let _ = sink.commit_overwrite().await;
446        sink.abort_overwrite().await.unwrap();
447        sink.complete_run().await.unwrap();
448
449        assert!(format!("{sink:?}").contains("MetadataSink"));
450    }
451
452    #[test]
453    fn custom_prefix_and_subset() {
454        let meta = CompiledMetadata::compile(&MetadataColumnsSpec {
455            enabled: true,
456            prefix: "_dt".into(),
457            columns: vec![MetadataColumn::RunId],
458        })
459        .unwrap()
460        .unwrap();
461        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
462        let out = sink.stamp(&[json!({"a": 1})]);
463        assert_eq!(out[0]["_dt_run_id"], "run-42");
464        assert!(out[0].get("_dt_loaded_at").is_none());
465    }
466}