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 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 let disabled = MetadataColumnsSpec {
334 enabled: false,
335 ..Default::default()
336 };
337 assert!(CompiledMetadata::compile(&disabled).unwrap().is_none());
338 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 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 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 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 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}