Skip to main content

faucet_core/
transforming_source.rs

1//! Wrap any [`Source`] with a fixed list of [`TransformStage`]s applied to
2//! every emitted record. The canonical way for library callers to attach
3//! stages (transforms wrapped via [`TransformStage::Map`], plus `Filter` /
4//! `Explode` / `Custom`); the CLI uses this same type internally.
5
6use crate::error::FaucetError;
7use crate::observability::{Labels, instrumented_apply_stages};
8use crate::pipeline::StreamPage;
9use crate::stage::{CompiledStage, TransformStage, compile_stage};
10use crate::traits::Source;
11use async_trait::async_trait;
12use futures::StreamExt;
13use futures_core::Stream;
14use serde_json::Value;
15use std::collections::HashMap;
16use std::pin::Pin;
17
18/// Source decorator that applies a fixed list of compiled stages to every
19/// record. Emits `faucet_transform_*` metrics per page via
20/// [`instrumented_apply_stages`].
21///
22/// # Example
23///
24/// ```no_run
25/// use faucet_core::{RecordTransform, Source, TransformingSource};
26/// use faucet_core::observability::Labels;
27/// use faucet_core::stage::TransformStage;
28/// use faucet_core::transform::KeyCaseMode;
29///
30/// # fn build_inner() -> Box<dyn Source> { unimplemented!() }
31/// let inner: Box<dyn Source> = build_inner();
32/// let wrapped = TransformingSource::new(
33///     inner,
34///     vec![TransformStage::Map(RecordTransform::KeysCase { mode: KeyCaseMode::Snake, on_collision: Default::default() })],
35///     Labels::for_named("rest"),
36/// ).unwrap();
37/// ```
38pub struct TransformingSource {
39    inner: Box<dyn Source>,
40    stages: Vec<CompiledStage>,
41    labels: Labels,
42    /// Optional Arrow `RecordBatch → RecordBatch` form for each stage, parallel
43    /// to `stages` (#375). `Some` only for columnar-capable stages (today the
44    /// SQL transform, supplied via [`new_with_batches`](Self::new_with_batches));
45    /// `None` for `Value`-only stages. When every entry is `Some` and the inner
46    /// source is columnar, the whole chain runs on the columnar fast path.
47    #[cfg(feature = "arrow")]
48    batch_fns: Vec<Option<crate::stage::PageFnBatchBox>>,
49}
50
51impl TransformingSource {
52    /// Compile `stages` and wrap `inner`. Returns
53    /// [`FaucetError::Transform`] if any stage's compilation fails (e.g.
54    /// invalid regex in `RenameKeys`). The chain stays on the `Value` path
55    /// (no columnar batch forms).
56    pub fn new(
57        inner: Box<dyn Source>,
58        stages: Vec<TransformStage>,
59        labels: Labels,
60    ) -> Result<Self, FaucetError> {
61        let compiled = stages
62            .iter()
63            .map(compile_stage)
64            .collect::<Result<Vec<_>, _>>()?;
65        // Auto-derive the Arrow batch form for each stage (#636): a `Map` over a
66        // vectorizable record transform (select/drop/rename_field/set/redact)
67        // gets a columnar kernel, so a chain of only those over a columnar
68        // source+sink keeps the fast path. Any stage without one (a
69        // value-inspecting transform, filter/explode/cdc-unwrap, a custom or
70        // page fn) leaves a `None`, which makes `supports_columnar` false and
71        // holds the whole chain on the `Value` path.
72        #[cfg(feature = "arrow")]
73        let batch_fns: Vec<Option<crate::stage::PageFnBatchBox>> = stages
74            .iter()
75            .map(|s| match s {
76                TransformStage::Map(t) => crate::columnar_transform::batch_form(t),
77                _ => None,
78            })
79            .collect();
80        Ok(Self {
81            inner,
82            stages: compiled,
83            labels,
84            #[cfg(feature = "arrow")]
85            batch_fns,
86        })
87    }
88
89    /// Like [`new`](Self::new), but each stage may carry an Arrow `RecordBatch`
90    /// form (`batch_fns[i]` parallels `stages[i]`), so the chain can run on the
91    /// columnar fast path (#375) when the inner source and sink are Arrow-native
92    /// and **every** stage supplies one. Used by the CLI for `sql` transforms.
93    #[cfg(feature = "arrow")]
94    pub fn new_with_batches(
95        inner: Box<dyn Source>,
96        stages: Vec<TransformStage>,
97        batch_fns: Vec<Option<crate::stage::PageFnBatchBox>>,
98        labels: Labels,
99    ) -> Result<Self, FaucetError> {
100        if batch_fns.len() != stages.len() {
101            return Err(FaucetError::Transform(format!(
102                "TransformingSource::new_with_batches: {} batch fns for {} stages",
103                batch_fns.len(),
104                stages.len()
105            )));
106        }
107        let compiled = stages
108            .iter()
109            .map(compile_stage)
110            .collect::<Result<Vec<_>, _>>()?;
111        Ok(Self {
112            inner,
113            stages: compiled,
114            labels,
115            batch_fns,
116        })
117    }
118}
119
120#[async_trait]
121impl Source for TransformingSource {
122    async fn fetch_with_context(
123        &self,
124        ctx: &HashMap<String, Value>,
125    ) -> Result<Vec<Value>, FaucetError> {
126        let records = self.inner.fetch_with_context(ctx).await?;
127        instrumented_apply_stages(records, &self.stages, &self.labels)
128    }
129
130    async fn fetch_with_context_incremental(
131        &self,
132        ctx: &HashMap<String, Value>,
133    ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
134        let (records, bookmark) = self.inner.fetch_with_context_incremental(ctx).await?;
135        let transformed = instrumented_apply_stages(records, &self.stages, &self.labels)?;
136        Ok((transformed, bookmark))
137    }
138
139    fn stream_pages<'a>(
140        &'a self,
141        ctx: &'a HashMap<String, Value>,
142        batch_size: usize,
143    ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
144        Box::pin(async_stream::try_stream! {
145            let mut pages = self.inner.stream_pages(ctx, batch_size);
146            while let Some(page) = pages.next().await {
147                let page = page?;
148                // The inner source already sized this page per its own config
149                // `batch_size` (the authoritative knob — the pipeline-supplied
150                // hint is only informational). Re-chunking the transformed
151                // output *below* that inner page size would silently defeat an
152                // explicit source `batch_size` (e.g. a 200k-row page shrunk to
153                // the 1k default hint → 200 tiny sink writes / load jobs). So
154                // never chunk smaller than the inner page; only bound *growth*
155                // from a 1→N stage (explode) at that inner size.
156                let page_len = page.records.len();
157                let out = instrumented_apply_stages(
158                    page.records, &self.stages, &self.labels,
159                )?;
160                if out.is_empty() {
161                    yield StreamPage { records: vec![], bookmark: page.bookmark };
162                    continue;
163                }
164                if batch_size == 0 {
165                    yield StreamPage { records: out, bookmark: page.bookmark };
166                    continue;
167                }
168                let effective = std::cmp::max(batch_size, page_len);
169                let total = out.len();
170                let mut start = 0usize;
171                while start < total {
172                    let end = std::cmp::min(start + effective, total);
173                    let is_last = end == total;
174                    let chunk: Vec<Value> = out[start..end].to_vec();
175                    yield StreamPage {
176                        records: chunk,
177                        bookmark: if is_last { page.bookmark.clone() } else { None },
178                    };
179                    start = end;
180                }
181            }
182        })
183    }
184
185    /// Columnar only when the inner source is columnar **and** every stage has
186    /// an Arrow batch form (today: the SQL transform). Any `Value`-only stage
187    /// (`Map` / `Filter` / `Explode` / `CdcUnwrap` / `Custom` / plain `PageFn`)
188    /// keeps the whole chain on the `Value` path (#375).
189    #[cfg(feature = "arrow")]
190    fn supports_columnar(&self) -> bool {
191        self.inner.supports_columnar()
192            && !self.batch_fns.is_empty()
193            && self.batch_fns.iter().all(Option::is_some)
194    }
195
196    /// Stream the inner source's Arrow batches with every stage's batch form
197    /// applied in declared order — so `parquet → sql → parquet` runs Arrow
198    /// end-to-end. Only reached when [`supports_columnar`](Self::supports_columnar)
199    /// is `true`, i.e. every stage has a batch form.
200    #[cfg(feature = "arrow")]
201    fn stream_batches<'a>(
202        &'a self,
203        ctx: &'a HashMap<String, Value>,
204        batch_size: usize,
205    ) -> Pin<Box<dyn Stream<Item = Result<crate::columnar::ColumnarPage, FaucetError>> + Send + 'a>>
206    {
207        Box::pin(async_stream::try_stream! {
208            use metrics::{Label, SharedString, counter};
209            let metric_labels = vec![
210                Label::new("pipeline", SharedString::from(self.labels.pipeline.to_string())),
211                Label::new("row", SharedString::from(self.labels.row.to_string())),
212            ];
213            let mut pages = self.inner.stream_batches(ctx, batch_size);
214            while let Some(page) = pages.next().await {
215                let crate::columnar::ColumnarPage { mut batch, bookmark } = page?;
216                // Metric parity with the Value path (#636): the same
217                // `faucet_transform_records_{in,out}_total` counters, per page.
218                // A pure-projection kernel is 1→1, so in == out, but emitting
219                // both keeps a dashboard identical across the two paths.
220                let n_in = batch.num_rows();
221                for bf in self.batch_fns.iter().flatten() {
222                    batch = bf(batch).inspect_err(|_| {
223                        counter!("faucet_transform_errors_total", metric_labels.clone())
224                            .increment(1);
225                    })?;
226                }
227                counter!("faucet_transform_records_in_total", metric_labels.clone())
228                    .increment(n_in as u64);
229                counter!("faucet_transform_records_out_total", metric_labels.clone())
230                    .increment(batch.num_rows() as u64);
231                yield crate::columnar::ColumnarPage { batch, bookmark };
232            }
233        })
234    }
235
236    fn state_key(&self) -> Option<String> {
237        self.inner.state_key()
238    }
239
240    async fn apply_start_bookmark(&self, bookmark: Value) -> Result<(), FaucetError> {
241        self.inner.apply_start_bookmark(bookmark).await
242    }
243
244    fn supports_exactly_once(&self) -> bool {
245        self.inner.supports_exactly_once()
246    }
247
248    fn replay_guarantee(&self) -> crate::idempotency::ReplayGuarantee {
249        self.inner.replay_guarantee()
250    }
251
252    async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
253        self.inner.capture_resume_position().await
254    }
255    async fn lag(&self) -> Result<Option<crate::lag::SourceLag>, FaucetError> {
256        self.inner.lag().await
257    }
258
259    fn state_schema(&self) -> u32 {
260        self.inner.state_schema()
261    }
262
263    fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError> {
264        self.inner.migrate_state(from, data)
265    }
266
267    fn record_table(&self, record: &Value) -> Option<String> {
268        self.inner.record_table(record)
269    }
270
271    fn position_le(&self, a: &Value, b: &Value) -> Option<bool> {
272        self.inner.position_le(a, b)
273    }
274
275    fn position_min(&self, positions: &[Value]) -> Option<Value> {
276        self.inner.position_min(positions)
277    }
278
279    fn connector_name(&self) -> &'static str {
280        self.inner.connector_name()
281    }
282
283    fn dataset_uri(&self) -> String {
284        // Forward the wrapped connector's identity — without this, lineage and
285        // the Data Movement Catalog would see the default
286        // `<connector>://unknown` whenever transforms are attached.
287        self.inner.dataset_uri()
288    }
289
290    fn set_roundtrip_recorder(
291        &self,
292        recorder: std::sync::Arc<crate::observability::RoundtripRecorder>,
293    ) {
294        self.inner.set_roundtrip_recorder(recorder);
295    }
296
297    fn set_run_clock(&self, now: chrono::DateTime<chrono::Utc>) {
298        self.inner.set_run_clock(now);
299    }
300}
301
302#[cfg(test)]
303mod tests {
304    use super::*;
305
306    /// A change-stream double with its own multi-table hooks (#731).
307    #[tokio::test]
308    async fn routed_double_fetches_nothing() {
309        use crate::Source as _;
310        assert!(
311            RoutedSource
312                .fetch_with_context(&Default::default())
313                .await
314                .unwrap()
315                .is_empty()
316        );
317    }
318
319    struct RoutedSource;
320
321    #[async_trait::async_trait]
322    impl crate::Source for RoutedSource {
323        async fn fetch_with_context(
324            &self,
325            _ctx: &std::collections::HashMap<String, serde_json::Value>,
326        ) -> Result<Vec<serde_json::Value>, crate::FaucetError> {
327            Ok(Vec::new())
328        }
329
330        fn record_table(&self, record: &serde_json::Value) -> Option<String> {
331            record.get("t")?.as_str().map(str::to_string)
332        }
333
334        fn position_le(&self, a: &serde_json::Value, b: &serde_json::Value) -> Option<bool> {
335            Some(a.as_u64()? <= b.as_u64()?)
336        }
337
338        fn position_min(&self, _positions: &[serde_json::Value]) -> Option<serde_json::Value> {
339            Some(serde_json::json!("inner-min"))
340        }
341    }
342
343    fn assert_forwards_multi_table_hooks(s: &dyn crate::Source) {
344        assert_eq!(
345            s.record_table(&serde_json::json!({"t": "public.a"}))
346                .as_deref(),
347            Some("public.a")
348        );
349        assert_eq!(
350            s.position_le(&serde_json::json!(1), &serde_json::json!(2)),
351            Some(true)
352        );
353        assert_eq!(
354            s.position_le(&serde_json::json!(3), &serde_json::json!(2)),
355            Some(false)
356        );
357        assert_eq!(
358            s.position_min(&[serde_json::json!(1), serde_json::json!(2)]),
359            Some(serde_json::json!("inner-min"))
360        );
361    }
362
363    #[test]
364    fn multi_table_hooks_are_forwarded_to_the_inner_source() {
365        let wrapped =
366            TransformingSource::new(Box::new(RoutedSource), vec![], Labels::for_named("test"))
367                .unwrap();
368        assert_forwards_multi_table_hooks(&wrapped);
369    }
370    use crate::stage::TransformStage;
371    use crate::transform::{KeyCaseMode, RecordTransform};
372    use serde_json::json;
373    use std::sync::Arc;
374    use std::sync::atomic::{AtomicBool, Ordering};
375
376    struct MockSource(Vec<Value>);
377
378    #[async_trait]
379    impl Source for MockSource {
380        async fn fetch_with_context(
381            &self,
382            _ctx: &HashMap<String, Value>,
383        ) -> Result<Vec<Value>, FaucetError> {
384            Ok(self.0.clone())
385        }
386    }
387
388    #[tokio::test]
389    async fn fetch_with_context_transforms_records() {
390        let inner: Box<dyn Source> = Box::new(MockSource(vec![json!({"FooBar": 1})]));
391        let wrapped = TransformingSource::new(
392            inner,
393            vec![TransformStage::Map(RecordTransform::KeysCase {
394                mode: KeyCaseMode::Snake,
395                on_collision: Default::default(),
396            })],
397            Labels::for_named("test"),
398        )
399        .expect("compile succeeds");
400        let out = wrapped.fetch_with_context(&HashMap::new()).await.unwrap();
401        assert_eq!(out, vec![json!({"foo_bar": 1})]);
402    }
403
404    struct VersionedSource;
405
406    #[async_trait]
407    impl Source for VersionedSource {
408        async fn fetch_with_context(
409            &self,
410            _ctx: &HashMap<String, Value>,
411        ) -> Result<Vec<Value>, FaucetError> {
412            Ok(Vec::new())
413        }
414        fn state_schema(&self) -> u32 {
415            2
416        }
417        fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError> {
418            Ok(json!({"from": from, "data": data}))
419        }
420    }
421
422    #[test]
423    fn state_versioning_is_forwarded_to_the_inner_source() {
424        let wrapped =
425            TransformingSource::new(Box::new(VersionedSource), vec![], Labels::for_named("test"))
426                .expect("compile succeeds");
427        assert_eq!(wrapped.state_schema(), 2);
428        assert_eq!(
429            wrapped.migrate_state(1, json!("x")).unwrap(),
430            json!({"from": 1, "data": "x"})
431        );
432    }
433
434    struct IncrementalSource {
435        records: Vec<Value>,
436        bookmark: Value,
437    }
438
439    #[async_trait]
440    impl Source for IncrementalSource {
441        async fn fetch_with_context(
442            &self,
443            _ctx: &HashMap<String, Value>,
444        ) -> Result<Vec<Value>, FaucetError> {
445            Ok(self.records.clone())
446        }
447
448        async fn fetch_with_context_incremental(
449            &self,
450            _ctx: &HashMap<String, Value>,
451        ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
452            Ok((self.records.clone(), Some(self.bookmark.clone())))
453        }
454    }
455
456    #[tokio::test]
457    async fn fetch_with_context_incremental_transforms_and_preserves_bookmark() {
458        let inner: Box<dyn Source> = Box::new(IncrementalSource {
459            records: vec![json!({"FooBar": 1})],
460            bookmark: json!("2026-05-28T00:00:00Z"),
461        });
462        let wrapped = TransformingSource::new(
463            inner,
464            vec![TransformStage::Map(RecordTransform::KeysCase {
465                mode: KeyCaseMode::Snake,
466                on_collision: Default::default(),
467            })],
468            Labels::for_named("test"),
469        )
470        .unwrap();
471        let (records, bookmark) = wrapped
472            .fetch_with_context_incremental(&HashMap::new())
473            .await
474            .unwrap();
475        assert_eq!(records, vec![json!({"foo_bar": 1})]);
476        assert_eq!(bookmark, Some(json!("2026-05-28T00:00:00Z")));
477    }
478
479    /// Emits records as N predetermined pages with the bookmark only on the last.
480    /// Overrides `stream_pages` directly so the test catches whether the wrapper
481    /// delegates to the native streaming path (correct) or falls back to the
482    /// chunk-the-buffer default (wrong — the bug we're fixing).
483    struct PagedSource {
484        pages: Vec<Vec<Value>>,
485        final_bookmark: Value,
486    }
487
488    #[async_trait]
489    impl Source for PagedSource {
490        async fn fetch_with_context(
491            &self,
492            _ctx: &HashMap<String, Value>,
493        ) -> Result<Vec<Value>, FaucetError> {
494            Ok(self.pages.iter().flatten().cloned().collect())
495        }
496
497        fn stream_pages<'a>(
498            &'a self,
499            _ctx: &'a HashMap<String, Value>,
500            _batch_size: usize,
501        ) -> Pin<Box<dyn futures_core::Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>>
502        {
503            let pages = self.pages.clone();
504            let bookmark = self.final_bookmark.clone();
505            Box::pin(async_stream::try_stream! {
506                let n = pages.len();
507                for (i, records) in pages.into_iter().enumerate() {
508                    let bm = if i + 1 == n { Some(bookmark.clone()) } else { None };
509                    yield StreamPage { records, bookmark: bm };
510                }
511            })
512        }
513    }
514
515    #[tokio::test]
516    async fn stream_pages_transforms_each_page_and_preserves_bookmarks() {
517        let inner: Box<dyn Source> = Box::new(PagedSource {
518            pages: vec![
519                vec![json!({"FooBar": 1})],
520                vec![json!({"FooBar": 2})],
521                vec![json!({"FooBar": 3})],
522            ],
523            final_bookmark: json!("v1"),
524        });
525        let wrapped = TransformingSource::new(
526            inner,
527            vec![TransformStage::Map(RecordTransform::KeysCase {
528                mode: KeyCaseMode::Snake,
529                on_collision: Default::default(),
530            })],
531            Labels::for_named("test"),
532        )
533        .unwrap();
534
535        let ctx = HashMap::new();
536        let mut stream = wrapped.stream_pages(&ctx, 1000);
537        let mut collected: Vec<StreamPage> = Vec::new();
538        while let Some(page) = stream.next().await {
539            collected.push(page.unwrap());
540        }
541
542        assert_eq!(collected.len(), 3);
543        assert_eq!(collected[0].records, vec![json!({"foo_bar": 1})]);
544        assert!(collected[0].bookmark.is_none());
545        assert_eq!(collected[1].records, vec![json!({"foo_bar": 2})]);
546        assert!(collected[1].bookmark.is_none());
547        assert_eq!(collected[2].records, vec![json!({"foo_bar": 3})]);
548        assert_eq!(collected[2].bookmark, Some(json!("v1")));
549    }
550
551    /// Regression: a 1→1 transform must NOT re-chunk the inner page below the
552    /// size the inner source already chose. A source that yields one large page
553    /// (its config `batch_size` honored) followed by a `keys_case` stage must
554    /// still emit that page as ONE `StreamPage`, not many hint-sized sub-pages —
555    /// otherwise an explicit source `batch_size` is silently defeated whenever a
556    /// transform is present (the cause of 60 tiny BigQuery load jobs instead of
557    /// one).
558    #[tokio::test]
559    async fn stream_pages_does_not_rechunk_large_page_below_inner_size() {
560        let big: Vec<Value> = (0..2500).map(|i| json!({"FooBar": i})).collect();
561        let inner: Box<dyn Source> = Box::new(PagedSource {
562            pages: vec![big],
563            final_bookmark: json!("v1"),
564        });
565        let wrapped = TransformingSource::new(
566            inner,
567            vec![TransformStage::Map(RecordTransform::KeysCase {
568                mode: KeyCaseMode::Snake,
569                on_collision: Default::default(),
570            })],
571            Labels::for_named("t"),
572        )
573        .unwrap();
574        let ctx = HashMap::new();
575        // Pipeline hint is the 1000-row default; it must NOT shrink the page.
576        let mut stream = wrapped.stream_pages(&ctx, 1000);
577        let mut pages: Vec<StreamPage> = Vec::new();
578        while let Some(p) = stream.next().await {
579            pages.push(p.unwrap());
580        }
581        assert_eq!(pages.len(), 1, "one inner page must stay one page");
582        assert_eq!(pages[0].records.len(), 2500);
583        assert_eq!(pages[0].records[0], json!({"foo_bar": 0}));
584        assert_eq!(pages[0].bookmark, Some(json!("v1")));
585    }
586
587    #[tokio::test]
588    async fn stream_pages_passes_through_empty_records_page_with_bookmark() {
589        struct EmptyWithBookmark;
590        #[async_trait]
591        impl Source for EmptyWithBookmark {
592            async fn fetch_with_context(
593                &self,
594                _ctx: &HashMap<String, Value>,
595            ) -> Result<Vec<Value>, FaucetError> {
596                Ok(Vec::new())
597            }
598            fn stream_pages<'a>(
599                &'a self,
600                _ctx: &'a HashMap<String, Value>,
601                _batch_size: usize,
602            ) -> Pin<
603                Box<dyn futures_core::Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>,
604            > {
605                Box::pin(async_stream::try_stream! {
606                    yield StreamPage { records: Vec::new(), bookmark: Some(json!("v1")) };
607                })
608            }
609        }
610        let wrapped = TransformingSource::new(
611            Box::new(EmptyWithBookmark),
612            vec![TransformStage::Map(RecordTransform::KeysCase {
613                mode: KeyCaseMode::Snake,
614                on_collision: Default::default(),
615            })],
616            Labels::for_named("test"),
617        )
618        .unwrap();
619        let ctx = HashMap::new();
620        let mut stream = wrapped.stream_pages(&ctx, 1000);
621        let page = stream.next().await.unwrap().unwrap();
622        assert!(page.records.is_empty());
623        assert_eq!(page.bookmark, Some(json!("v1")));
624        assert!(stream.next().await.is_none());
625    }
626
627    struct InstrumentedSource {
628        started: Arc<AtomicBool>,
629    }
630
631    #[async_trait]
632    impl Source for InstrumentedSource {
633        async fn fetch_with_context(
634            &self,
635            _ctx: &HashMap<String, Value>,
636        ) -> Result<Vec<Value>, FaucetError> {
637            Ok(vec![])
638        }
639        fn connector_name(&self) -> &'static str {
640            "instrumented"
641        }
642        fn state_key(&self) -> Option<String> {
643            Some("instrumented::key".to_string())
644        }
645        async fn apply_start_bookmark(&self, _bookmark: Value) -> Result<(), FaucetError> {
646            self.started.store(true, Ordering::Relaxed);
647            Ok(())
648        }
649        fn supports_exactly_once(&self) -> bool {
650            true
651        }
652        async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
653            Ok(Some(json!("captured")))
654        }
655    }
656
657    struct RecorderProbe(Arc<AtomicBool>);
658
659    #[async_trait]
660    impl Source for RecorderProbe {
661        async fn fetch_with_context(
662            &self,
663            _ctx: &HashMap<String, Value>,
664        ) -> Result<Vec<Value>, FaucetError> {
665            Ok(vec![])
666        }
667        fn set_roundtrip_recorder(&self, _recorder: Arc<crate::observability::RoundtripRecorder>) {
668            self.0.store(true, Ordering::Relaxed);
669        }
670        fn set_run_clock(&self, _now: chrono::DateTime<chrono::Utc>) {
671            self.0.store(true, Ordering::Relaxed);
672        }
673    }
674
675    #[test]
676    fn run_clock_reaches_the_wrapped_source() {
677        let got = Arc::new(AtomicBool::new(false));
678        let wrapped = TransformingSource::new(
679            Box::new(RecorderProbe(got.clone())),
680            vec![],
681            Labels::for_named("test"),
682        )
683        .unwrap();
684        wrapped.set_run_clock(chrono::Utc::now());
685        assert!(got.load(Ordering::Relaxed));
686    }
687
688    #[test]
689    fn roundtrip_recorder_reaches_the_wrapped_source() {
690        let got = Arc::new(AtomicBool::new(false));
691        let wrapped = TransformingSource::new(
692            Box::new(RecorderProbe(got.clone())),
693            vec![],
694            Labels::for_named("test"),
695        )
696        .unwrap();
697        wrapped.set_roundtrip_recorder(Arc::new(crate::observability::RoundtripRecorder::new(
698            crate::observability::RoundtripSide::Source,
699            "p",
700            "r",
701            "probe",
702        )));
703        assert!(got.load(Ordering::Relaxed));
704    }
705
706    #[tokio::test]
707    async fn connector_name_state_key_and_start_bookmark_delegate_to_inner() {
708        let started = Arc::new(AtomicBool::new(false));
709        let inner = InstrumentedSource {
710            started: started.clone(),
711        };
712        let wrapped = TransformingSource::new(
713            Box::new(inner),
714            vec![TransformStage::Map(RecordTransform::KeysCase {
715                mode: KeyCaseMode::Snake,
716                on_collision: Default::default(),
717            })],
718            Labels::for_named("test"),
719        )
720        .unwrap();
721        assert_eq!(wrapped.connector_name(), "instrumented");
722        assert_eq!(wrapped.state_key(), Some("instrumented::key".to_string()));
723        wrapped.apply_start_bookmark(json!("bm")).await.unwrap();
724        assert!(started.load(Ordering::Relaxed));
725        // Exactly-once capabilities must survive the transform wrap — the
726        // pipeline's mechanism selection reads them through this layer.
727        assert!(wrapped.supports_exactly_once());
728        assert_eq!(
729            wrapped.replay_guarantee(),
730            crate::idempotency::ReplayGuarantee::Deterministic
731        );
732        assert_eq!(
733            wrapped.capture_resume_position().await.unwrap(),
734            Some(json!("captured"))
735        );
736        assert_eq!(wrapped.lag().await.unwrap(), None);
737    }
738
739    #[tokio::test]
740    async fn new_fails_fast_on_invalid_regex() {
741        let inner: Box<dyn Source> = Box::new(MockSource(vec![]));
742        let result = TransformingSource::new(
743            inner,
744            vec![TransformStage::Map(RecordTransform::RenameKeys {
745                pattern: "[invalid".to_string(),
746                replacement: "x".to_string(),
747            })],
748            Labels::for_named("test"),
749        );
750        let err = match result {
751            Ok(_) => panic!("invalid regex must fail at new()"),
752            Err(e) => e,
753        };
754        assert!(matches!(err, FaucetError::Transform(_)));
755    }
756
757    #[tokio::test]
758    async fn custom_closure_transform_runs() {
759        let inner: Box<dyn Source> = Box::new(MockSource(vec![json!({"x": 1})]));
760        let wrapped = TransformingSource::new(
761            inner,
762            vec![TransformStage::Map(RecordTransform::custom(
763                |mut record| {
764                    if let Some(obj) = record.as_object_mut() {
765                        obj.insert("added".to_string(), json!(true));
766                    }
767                    record
768                },
769            ))],
770            Labels::for_named("test"),
771        )
772        .unwrap();
773        let out = wrapped.fetch_with_context(&HashMap::new()).await.unwrap();
774        assert_eq!(out, vec![json!({"x": 1, "added": true})]);
775    }
776
777    #[tokio::test]
778    async fn usable_as_boxed_dyn_source() {
779        let inner: Box<dyn Source> = Box::new(MockSource(vec![json!({"FooBar": 1})]));
780        let wrapped: Box<dyn Source> = Box::new(
781            TransformingSource::new(
782                inner,
783                vec![TransformStage::Map(RecordTransform::KeysCase {
784                    mode: KeyCaseMode::Snake,
785                    on_collision: Default::default(),
786                })],
787                Labels::for_named("test"),
788            )
789            .unwrap(),
790        );
791        let out = wrapped.fetch_with_context(&HashMap::new()).await.unwrap();
792        assert_eq!(out, vec![json!({"foo_bar": 1})]);
793    }
794
795    /// A source that emits a single page with the given records and bookmark.
796    ///
797    /// Only the bookmark-forwarding tests construct it, and those are gated on
798    /// a transform feature, so a build without one compiles it unused.
799    #[allow(dead_code)]
800    struct OnePageSource {
801        records: Vec<Value>,
802        bookmark: Option<Value>,
803    }
804
805    #[async_trait]
806    impl Source for OnePageSource {
807        async fn fetch_with_context(
808            &self,
809            _ctx: &HashMap<String, Value>,
810        ) -> Result<Vec<Value>, FaucetError> {
811            Ok(self.records.clone())
812        }
813        async fn fetch_with_context_incremental(
814            &self,
815            _ctx: &HashMap<String, Value>,
816        ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
817            Ok((self.records.clone(), self.bookmark.clone()))
818        }
819        fn stream_pages<'a>(
820            &'a self,
821            _ctx: &'a HashMap<String, Value>,
822            _batch_size: usize,
823        ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
824            let page = StreamPage {
825                records: self.records.clone(),
826                bookmark: self.bookmark.clone(),
827            };
828            Box::pin(async_stream::stream! { yield Ok(page); })
829        }
830    }
831
832    #[cfg(feature = "transform-explode")]
833    fn explode_stage() -> TransformStage {
834        TransformStage::Explode(crate::stage::ExplodeSpec {
835            path: "items".to_owned(),
836            prefix: None,
837            separator: "_".to_owned(),
838            on_missing: crate::stage::OnMissing::Drop,
839        })
840    }
841
842    /// Build N records each with a 10-element `items` array, so explode 10×s them.
843    #[cfg(feature = "transform-explode")]
844    fn explode_10x_records(n: usize) -> Vec<Value> {
845        (0..n)
846            .map(|i| {
847                json!({
848                    "id": i,
849                    "items": (0..10).map(|j| json!({"k": j})).collect::<Vec<_>>(),
850                })
851            })
852            .collect()
853    }
854
855    #[cfg(feature = "transform-explode")]
856    #[tokio::test]
857    async fn stream_pages_rechunks_explosion_with_bookmark_on_last() {
858        let inner: Box<dyn Source> = Box::new(OnePageSource {
859            records: explode_10x_records(100), // 100 → 1000 after explode
860            bookmark: Some(json!("bm")),
861        });
862        let wrapped =
863            TransformingSource::new(inner, vec![explode_stage()], Labels::for_named("t")).unwrap();
864        let ctx = HashMap::new();
865        let mut stream = wrapped.stream_pages(&ctx, 200);
866        let mut sub_pages: Vec<StreamPage> = Vec::new();
867        while let Some(p) = stream.next().await {
868            sub_pages.push(p.unwrap());
869        }
870        assert_eq!(sub_pages.len(), 5, "1000 records / 200 batch = 5 sub-pages");
871        for (i, p) in sub_pages.iter().enumerate() {
872            assert_eq!(p.records.len(), 200, "sub-page {i} should be size 200");
873            if i < 4 {
874                assert!(
875                    p.bookmark.is_none(),
876                    "non-final sub-page {i} carries no bookmark"
877                );
878            } else {
879                assert_eq!(p.bookmark, Some(json!("bm")), "final sub-page has bookmark");
880            }
881        }
882    }
883
884    #[cfg(feature = "transform-explode")]
885    #[tokio::test]
886    async fn stream_pages_batch_size_zero_emits_one_page() {
887        let inner: Box<dyn Source> = Box::new(OnePageSource {
888            records: explode_10x_records(10), // 10 → 100 after explode
889            bookmark: Some(json!("bm")),
890        });
891        let wrapped =
892            TransformingSource::new(inner, vec![explode_stage()], Labels::for_named("t")).unwrap();
893        let ctx = HashMap::new();
894        let mut stream = wrapped.stream_pages(&ctx, 0);
895        let mut sub_pages: Vec<StreamPage> = Vec::new();
896        while let Some(p) = stream.next().await {
897            sub_pages.push(p.unwrap());
898        }
899        assert_eq!(sub_pages.len(), 1, "batch_size=0 means one sub-page");
900        assert_eq!(sub_pages[0].records.len(), 100);
901        assert_eq!(sub_pages[0].bookmark, Some(json!("bm")));
902    }
903
904    #[cfg(feature = "transform-filter")]
905    #[tokio::test]
906    async fn stream_pages_filter_drops_all_still_yields_bookmark() {
907        let inner: Box<dyn Source> = Box::new(OnePageSource {
908            records: vec![json!({"deleted": true}), json!({"deleted": true})],
909            bookmark: Some(json!("bm")),
910        });
911        let drop_all = TransformStage::Filter(crate::stage::FilterSpec {
912            path: "deleted".to_owned(),
913            op: crate::stage::FilterOp::Ne,
914            value: Some(json!(true)),
915        });
916        let wrapped =
917            TransformingSource::new(inner, vec![drop_all], Labels::for_named("t")).unwrap();
918        let ctx = HashMap::new();
919        let mut stream = wrapped.stream_pages(&ctx, 100);
920        let mut sub_pages: Vec<StreamPage> = Vec::new();
921        while let Some(p) = stream.next().await {
922            sub_pages.push(p.unwrap());
923        }
924        assert_eq!(sub_pages.len(), 1);
925        assert!(sub_pages[0].records.is_empty());
926        assert_eq!(sub_pages[0].bookmark, Some(json!("bm")));
927    }
928}
929
930#[cfg(all(test, feature = "arrow"))]
931mod columnar_tests {
932    use super::*;
933    use crate::columnar::{ColumnarPage, record_batch_to_values, values_to_record_batch_inferred};
934    use crate::stage::TransformStage;
935    use serde_json::json;
936    use std::sync::Arc;
937
938    /// A source that emits one Arrow batch (columnar-capable).
939    struct ColumnarMock(Vec<Value>);
940    #[async_trait]
941    impl Source for ColumnarMock {
942        async fn fetch_with_context(
943            &self,
944            _ctx: &HashMap<String, Value>,
945        ) -> Result<Vec<Value>, FaucetError> {
946            Ok(self.0.clone())
947        }
948        fn supports_columnar(&self) -> bool {
949            true
950        }
951        fn stream_batches<'a>(
952            &'a self,
953            _ctx: &'a HashMap<String, Value>,
954            _bs: usize,
955        ) -> Pin<Box<dyn Stream<Item = Result<ColumnarPage, FaucetError>> + Send + 'a>> {
956            let batch = values_to_record_batch_inferred(&self.0).unwrap();
957            Box::pin(async_stream::stream! {
958                yield Ok(ColumnarPage { batch, bookmark: Some(json!("bm")) });
959            })
960        }
961    }
962
963    /// A source with no columnar support (default `supports_columnar` = false).
964    struct RowOnlyMock(Vec<Value>);
965    #[async_trait]
966    impl Source for RowOnlyMock {
967        async fn fetch_with_context(
968            &self,
969            _ctx: &HashMap<String, Value>,
970        ) -> Result<Vec<Value>, FaucetError> {
971            Ok(self.0.clone())
972        }
973    }
974
975    /// An identity page stage + its identity batch form — enough to exercise
976    /// the columnar wiring.
977    fn identity_stage() -> (TransformStage, Option<crate::stage::PageFnBatchBox>) {
978        let rows: crate::stage::PageFnBox = Arc::new(Ok);
979        let batch: crate::stage::PageFnBatchBox = Arc::new(Ok);
980        (TransformStage::PageFn(rows), Some(batch))
981    }
982
983    #[tokio::test]
984    async fn columnar_inner_plus_columnar_stage_is_supported_and_streams() {
985        let inner: Box<dyn Source> =
986            Box::new(ColumnarMock(vec![json!({"id": 1}), json!({"id": 2})]));
987        let (stage, batch) = identity_stage();
988        let wrapped = TransformingSource::new_with_batches(
989            inner,
990            vec![stage],
991            vec![batch],
992            Labels::for_named("t"),
993        )
994        .unwrap();
995        assert!(wrapped.supports_columnar());
996        let ctx = HashMap::new();
997        let mut s = wrapped.stream_batches(&ctx, 0);
998        let page = s.next().await.unwrap().unwrap();
999        let rows = record_batch_to_values(&page.batch).unwrap();
1000        assert_eq!(rows.len(), 2);
1001        assert_eq!(page.bookmark, Some(json!("bm")));
1002    }
1003
1004    #[tokio::test]
1005    async fn value_only_stage_disables_columnar() {
1006        // A `Map` stage has no Arrow batch form (batch_fn None) → the whole
1007        // chain drops off the columnar path even though the inner is columnar.
1008        let inner: Box<dyn Source> = Box::new(ColumnarMock(vec![json!({"FooBar": 1})]));
1009        let wrapped = TransformingSource::new(
1010            inner,
1011            vec![TransformStage::Map(
1012                crate::transform::RecordTransform::KeysCase {
1013                    mode: crate::transform::KeyCaseMode::Snake,
1014                    on_collision: crate::transform::KeyCollision::Error,
1015                },
1016            )],
1017            Labels::for_named("t"),
1018        )
1019        .unwrap();
1020        assert!(!wrapped.supports_columnar());
1021    }
1022
1023    /// #636: a chain of *built-in vectorizable* transforms (no hand-supplied
1024    /// batch fns) over a columnar source keeps the fast path automatically —
1025    /// `new` now derives their Arrow kernels. This is the whole point: a
1026    /// `parquet → select/drop/set → parquet` run no longer drops to `Value`.
1027    #[cfg(all(feature = "transform-drop", feature = "transform-set"))]
1028    #[tokio::test]
1029    async fn built_in_vectorizable_transforms_keep_the_columnar_path() {
1030        use crate::transform::RecordTransform;
1031        let inner: Box<dyn Source> = Box::new(ColumnarMock(vec![
1032            json!({"id": 1, "name": "ada", "secret": "x"}),
1033            json!({"id": 2, "name": "grace", "secret": "y"}),
1034        ]));
1035        let mut set_vals = serde_json::Map::new();
1036        set_vals.insert("stage".into(), json!("prod"));
1037        let wrapped = TransformingSource::new(
1038            inner,
1039            vec![
1040                TransformStage::Map(RecordTransform::Drop {
1041                    fields: vec!["secret".into()],
1042                }),
1043                TransformStage::Map(RecordTransform::Set { values: set_vals }),
1044            ],
1045            Labels::for_named("t"),
1046        )
1047        .unwrap();
1048        assert!(
1049            wrapped.supports_columnar(),
1050            "a chain of vectorizable built-ins must stay columnar"
1051        );
1052        let ctx = HashMap::new();
1053        let mut s = wrapped.stream_batches(&ctx, 0);
1054        let page = s.next().await.unwrap().unwrap();
1055        let rows = record_batch_to_values(&page.batch).unwrap();
1056        assert_eq!(rows.len(), 2);
1057        // `drop` removed `secret`, `set` added `stage` — the transforms really
1058        // ran on the columnar path, not just passed through.
1059        assert!(rows[0].get("secret").is_none(), "drop ran: {:?}", rows[0]);
1060        assert_eq!(rows[0]["stage"], json!("prod"), "set ran: {:?}", rows[0]);
1061    }
1062
1063    /// One non-vectorizable transform anywhere in the chain holds the whole
1064    /// chain on the `Value` path — the vectorizable ones must not silently run
1065    /// columnar while the opaque one is skipped.
1066    #[cfg(all(feature = "transform-select", feature = "transform-flatten"))]
1067    #[tokio::test]
1068    async fn a_mixed_chain_with_one_opaque_transform_stays_on_value() {
1069        use crate::transform::RecordTransform;
1070        let inner: Box<dyn Source> = Box::new(ColumnarMock(vec![json!({"id": 1})]));
1071        let wrapped = TransformingSource::new(
1072            inner,
1073            vec![
1074                TransformStage::Map(RecordTransform::Select {
1075                    fields: vec!["id".into()],
1076                }),
1077                // `flatten` has no kernel → None → whole chain off columnar.
1078                TransformStage::Map(RecordTransform::Flatten {
1079                    separator: "_".into(),
1080                }),
1081            ],
1082            Labels::for_named("t"),
1083        )
1084        .unwrap();
1085        assert!(!wrapped.supports_columnar());
1086    }
1087
1088    #[tokio::test]
1089    async fn columnar_stage_over_row_only_inner_is_disabled() {
1090        let inner: Box<dyn Source> = Box::new(RowOnlyMock(vec![json!({"id": 1})]));
1091        let (stage, batch) = identity_stage();
1092        let wrapped = TransformingSource::new_with_batches(
1093            inner,
1094            vec![stage],
1095            vec![batch],
1096            Labels::for_named("t"),
1097        )
1098        .unwrap();
1099        assert!(!wrapped.supports_columnar(), "inner is not columnar");
1100    }
1101}