1use 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
18pub struct TransformingSource {
39 inner: Box<dyn Source>,
40 stages: Vec<CompiledStage>,
41 labels: Labels,
42 #[cfg(feature = "arrow")]
48 batch_fns: Vec<Option<crate::stage::PageFnBatchBox>>,
49}
50
51impl TransformingSource {
52 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 #[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 #[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 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 #[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 #[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 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 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 #[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 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 #[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 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 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 #[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 #[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), 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), 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 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 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 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 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 #[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 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 #[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 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}