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 check(
227        &self,
228        ctx: &crate::check::CheckContext,
229    ) -> Result<crate::check::CheckReport, FaucetError> {
230        self.inner.check(ctx).await
231    }
232    fn supports_cleanup(&self) -> bool {
233        self.inner.supports_cleanup()
234    }
235    async fn cleanup_scope(
236        &self,
237        scope: &BTreeMap<String, Value>,
238        seen: &crate::cleanup::SeenKeys,
239    ) -> Result<u64, FaucetError> {
240        self.inner.cleanup_scope(scope, seen).await
241    }
242    fn supports_idempotent_writes(&self) -> bool {
243        self.inner.supports_idempotent_writes()
244    }
245    async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
246        self.inner.last_committed_token(scope).await
247    }
248    fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
249        self.inner.supported_write_modes()
250    }
251    fn dedups_by_key(&self) -> bool {
252        self.inner.dedups_by_key()
253    }
254    fn sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
255        self.inner.sink_guarantee()
256    }
257    async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
258        self.inner.current_schema().await
259    }
260    fn supports_schema_evolution(&self) -> bool {
261        self.inner.supports_schema_evolution()
262    }
263    async fn evolve_schema(
264        &self,
265        evolution: &crate::drift::SchemaEvolution,
266    ) -> Result<(), FaucetError> {
267        self.inner.evolve_schema(evolution).await
268    }
269    fn config_schema(&self) -> Value {
270        self.inner.config_schema()
271    }
272    fn connector_name(&self) -> &'static str {
273        self.inner.connector_name()
274    }
275    fn dataset_uri(&self) -> String {
276        self.inner.dataset_uri()
277    }
278    fn is_overwrite(&self) -> bool {
279        self.inner.is_overwrite()
280    }
281    async fn begin_overwrite(&self) -> Result<(), FaucetError> {
282        self.inner.begin_overwrite().await
283    }
284    async fn commit_overwrite(&self) -> Result<(), FaucetError> {
285        self.inner.commit_overwrite().await
286    }
287    async fn abort_overwrite(&self) -> Result<(), FaucetError> {
288        self.inner.abort_overwrite().await
289    }
290}
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295    use serde_json::json;
296    use std::sync::Mutex;
297
298    #[derive(Debug, Default)]
299    struct CapturingSink {
300        rows: Mutex<Vec<Value>>,
301    }
302    #[async_trait::async_trait]
303    impl Sink for CapturingSink {
304        async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
305            self.rows.lock().unwrap().extend_from_slice(records);
306            Ok(records.len())
307        }
308        fn config_schema(&self) -> Value {
309            json!({})
310        }
311    }
312
313    fn spec(cols: &[MetadataColumn]) -> MetadataColumnsSpec {
314        MetadataColumnsSpec {
315            enabled: true,
316            prefix: "_faucet".into(),
317            columns: cols.to_vec(),
318        }
319    }
320
321    fn ctx() -> MetadataContext {
322        MetadataContext {
323            run_id: "run-42".into(),
324            source: "rest".into(),
325        }
326    }
327
328    #[test]
329    fn compile_resolves_defaults_and_respects_disable() {
330        let c = CompiledMetadata::compile(&spec(&[])).unwrap().unwrap();
331        assert_eq!(c.columns, DEFAULT_COLUMNS.to_vec());
332        // disabled → None
333        let disabled = MetadataColumnsSpec {
334            enabled: false,
335            ..Default::default()
336        };
337        assert!(CompiledMetadata::compile(&disabled).unwrap().is_none());
338        // empty prefix → error
339        let bad = MetadataColumnsSpec {
340            prefix: " ".into(),
341            ..Default::default()
342        };
343        assert!(CompiledMetadata::compile(&bad).is_err());
344    }
345
346    #[tokio::test]
347    async fn stamps_all_column_kinds_with_prefix() {
348        let meta = CompiledMetadata::compile(&spec(&[
349            MetadataColumn::ExtractedAt,
350            MetadataColumn::LoadedAt,
351            MetadataColumn::RunId,
352            MetadataColumn::Source,
353            MetadataColumn::Sequence,
354        ]))
355        .unwrap()
356        .unwrap();
357        let inner = Box::new(CapturingSink::default());
358        let sink = MetadataSink::new(inner, meta, ctx());
359        let n = sink
360            .write_batch(&[json!({"id": 1}), json!({"id": 2})])
361            .await
362            .unwrap();
363        assert_eq!(n, 2);
364        // Re-wrap to read the captured rows would need the inner; instead assert
365        // via a fresh capturing sink shared by Arc — simpler: re-stamp directly.
366        let stamped = sink.stamp(&[json!({"id": 1}), json!({"id": 2})]);
367        let r0 = &stamped[0];
368        assert_eq!(r0["id"], 1);
369        assert_eq!(r0["_faucet_run_id"], "run-42");
370        assert_eq!(r0["_faucet_source"], "rest");
371        assert!(r0["_faucet_extracted_at"].is_string());
372        assert!(r0["_faucet_loaded_at"].is_string());
373        // Sequence is monotonic across records.
374        assert_eq!(
375            stamped[0]["_faucet_sequence"].as_u64().unwrap() + 1,
376            stamped[1]["_faucet_sequence"].as_u64().unwrap()
377        );
378    }
379
380    #[test]
381    fn non_object_records_pass_through() {
382        let meta = CompiledMetadata::compile(&spec(&[MetadataColumn::RunId]))
383            .unwrap()
384            .unwrap();
385        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
386        let out = sink.stamp(&[json!("scalar"), json!(42)]);
387        assert_eq!(out, vec![json!("scalar"), json!(42)]);
388    }
389
390    #[tokio::test]
391    async fn delegates_every_sink_method_to_inner() {
392        let meta = CompiledMetadata::compile(&spec(&[MetadataColumn::RunId]))
393            .unwrap()
394            .unwrap();
395        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
396
397        // Write paths stamp + forward.
398        assert_eq!(sink.write_batch(&[json!({"a": 1})]).await.unwrap(), 1);
399        assert_eq!(
400            sink.write_batch_partial(&[json!({"a": 1})])
401                .await
402                .unwrap()
403                .len(),
404            1
405        );
406        let _ = sink
407            .write_batch_idempotent(&[json!({"a": 1})], "scope", "tok")
408            .await;
409        sink.flush().await.unwrap();
410
411        // Pure forwards (inner uses trait defaults).
412        let _ = sink.check(&crate::check::CheckContext::default()).await;
413        assert!(!sink.supports_cleanup());
414        let _ = sink
415            .cleanup_scope(&BTreeMap::new(), &crate::cleanup::SeenKeys::new())
416            .await;
417        assert!(!sink.supports_idempotent_writes());
418        assert!(sink.last_committed_token("s").await.unwrap().is_none());
419        let _ = sink.supported_write_modes();
420        let _ = sink.dedups_by_key();
421        let _ = sink.sink_guarantee();
422        let _ = sink.current_schema().await.unwrap();
423        let _ = sink.supports_schema_evolution();
424        let _ = sink
425            .evolve_schema(&crate::drift::SchemaEvolution::default())
426            .await;
427        let _ = sink.config_schema();
428        let _ = sink.connector_name();
429        let _ = sink.dataset_uri();
430        assert!(!sink.is_overwrite());
431        let _ = sink.begin_overwrite().await;
432        let _ = sink.commit_overwrite().await;
433        sink.abort_overwrite().await.unwrap();
434
435        assert!(format!("{sink:?}").contains("MetadataSink"));
436    }
437
438    #[test]
439    fn custom_prefix_and_subset() {
440        let meta = CompiledMetadata::compile(&MetadataColumnsSpec {
441            enabled: true,
442            prefix: "_dt".into(),
443            columns: vec![MetadataColumn::RunId],
444        })
445        .unwrap()
446        .unwrap();
447        let sink = MetadataSink::new(Box::new(CapturingSink::default()), meta, ctx());
448        let out = sink.stamp(&[json!({"a": 1})]);
449        assert_eq!(out[0]["_dt_run_id"], "run-42");
450        assert!(out[0].get("_dt_loaded_at").is_none());
451    }
452}