1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
26#[serde(rename_all = "snake_case")]
27pub enum MetadataColumn {
28 ExtractedAt,
31 LoadedAt,
33 RunId,
35 Source,
37 Sequence,
39}
40
41impl MetadataColumn {
42 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#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)]
63#[serde(deny_unknown_fields)]
64pub struct MetadataColumnsSpec {
65 #[serde(default = "default_true")]
67 pub enabled: bool,
68 #[serde(default = "default_prefix")]
70 pub prefix: String,
71 #[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#[derive(Debug, Clone)]
96pub struct CompiledMetadata {
97 prefix: String,
98 columns: Vec<MetadataColumn>,
99}
100
101impl CompiledMetadata {
102 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#[derive(Debug, Clone)]
131pub struct MetadataContext {
132 pub run_id: String,
134 pub source: String,
136}
137
138pub 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 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 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 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 let disabled = MetadataColumnsSpec {
343 enabled: false,
344 ..Default::default()
345 };
346 assert!(CompiledMetadata::compile(&disabled).unwrap().is_none());
347 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 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 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 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 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}