1use crate::error::FaucetError;
5use crate::observability::labels::Labels;
6use crate::observability::timer::DurationGuard;
7use crate::pipeline::StreamPage;
8use crate::traits::{Sink, Source};
9use crate::usage::{UsageMeter, estimate_page_bytes};
10use async_trait::async_trait;
11use futures::FutureExt;
12use futures_core::Stream;
13use metrics::{Label, SharedString, counter, gauge};
14use serde_json::Value;
15use std::collections::HashMap;
16use std::panic::AssertUnwindSafe;
17use std::pin::Pin;
18use std::sync::Arc;
19use std::sync::atomic::{AtomicUsize, Ordering};
20use tracing::{Instrument, info_span};
21
22fn guarded_connector_name(raw: &'static str) -> &'static str {
26 if raw.is_empty() { "unknown" } else { raw }
27}
28
29fn base_metric_labels(labels: &Labels, connector: &SharedString) -> Vec<Label> {
34 vec![
35 Label::new("pipeline", SharedString::from(labels.pipeline.to_string())),
36 Label::new("row", SharedString::from(labels.row.to_string())),
37 Label::new("connector", connector.clone()),
38 ]
39}
40
41pub struct InstrumentedSource<'a, S: Source + ?Sized> {
45 inner: &'a S,
46 labels: Labels,
47 connector: SharedString,
48 base_labels: Vec<Label>,
50 page_index: Arc<AtomicUsize>,
51 meter: Option<Arc<UsageMeter>>,
53}
54
55impl<'a, S: Source + ?Sized> InstrumentedSource<'a, S> {
56 pub fn new(inner: &'a S, labels: Labels) -> Self {
57 let raw = inner.connector_name();
58 debug_assert!(
59 !raw.is_empty(),
60 "connector_name() must return a non-empty string"
61 );
62 let connector: SharedString = SharedString::const_str(guarded_connector_name(raw));
63 let base_labels = base_metric_labels(&labels, &connector);
64 Self {
65 inner,
66 labels,
67 connector,
68 base_labels,
69 page_index: Arc::new(AtomicUsize::new(0)),
70 meter: None,
71 }
72 }
73
74 pub fn with_meter(mut self, meter: Arc<UsageMeter>) -> Self {
76 self.meter = Some(meter);
77 self
78 }
79
80 fn metric_labels(&self) -> Vec<Label> {
81 self.base_labels.clone()
82 }
83
84 #[allow(dead_code)]
88 fn error_labels(&self, kind: &'static str) -> Vec<Label> {
89 let mut l = self.metric_labels();
90 l.push(Label::new("kind", SharedString::const_str(kind)));
91 l
92 }
93}
94
95#[async_trait]
96impl<'a, S: Source + ?Sized> Source for InstrumentedSource<'a, S> {
97 fn connector_name(&self) -> &'static str {
98 guarded_connector_name(self.inner.connector_name())
102 }
103
104 fn set_roundtrip_recorder(
108 &self,
109 recorder: std::sync::Arc<crate::observability::RoundtripRecorder>,
110 ) {
111 self.inner.set_roundtrip_recorder(recorder);
112 }
113
114 fn set_run_clock(&self, now: chrono::DateTime<chrono::Utc>) {
115 self.inner.set_run_clock(now);
116 }
117
118 fn state_key(&self) -> Option<String> {
119 self.inner.state_key()
120 }
121
122 async fn apply_start_bookmark(&self, bookmark: Value) -> Result<(), FaucetError> {
123 self.inner.apply_start_bookmark(bookmark).await
124 }
125
126 fn supports_exactly_once(&self) -> bool {
127 self.inner.supports_exactly_once()
128 }
129
130 fn replay_guarantee(&self) -> crate::idempotency::ReplayGuarantee {
131 self.inner.replay_guarantee()
132 }
133
134 async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
135 self.inner.capture_resume_position().await
136 }
137 async fn lag(&self) -> Result<Option<crate::lag::SourceLag>, FaucetError> {
138 self.inner.lag().await
139 }
140
141 fn state_schema(&self) -> u32 {
142 self.inner.state_schema()
143 }
144
145 fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError> {
146 self.inner.migrate_state(from, data)
147 }
148
149 fn record_table(&self, record: &Value) -> Option<String> {
150 self.inner.record_table(record)
151 }
152
153 fn position_le(&self, a: &Value, b: &Value) -> Option<bool> {
154 self.inner.position_le(a, b)
155 }
156
157 fn position_min(&self, positions: &[Value]) -> Option<Value> {
158 self.inner.position_min(positions)
159 }
160
161 async fn fetch_with_context(
162 &self,
163 context: &HashMap<String, Value>,
164 ) -> Result<Vec<Value>, FaucetError> {
165 self.inner.fetch_with_context(context).await
167 }
168
169 async fn fetch_with_context_incremental(
170 &self,
171 context: &HashMap<String, Value>,
172 ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
173 self.inner.fetch_with_context_incremental(context).await
174 }
175
176 #[cfg(feature = "arrow")]
180 fn supports_columnar(&self) -> bool {
181 self.inner.supports_columnar()
182 }
183
184 #[cfg(feature = "arrow")]
185 fn stream_batches<'b>(
186 &'b self,
187 context: &'b HashMap<String, Value>,
188 batch_size: usize,
189 ) -> Pin<Box<dyn Stream<Item = Result<crate::columnar::ColumnarPage, FaucetError>> + Send + 'b>>
190 {
191 self.inner.stream_batches(context, batch_size)
192 }
193
194 fn native_output_formats(&self) -> &'static [crate::native::NativeFormat] {
197 self.inner.native_output_formats()
198 }
199
200 fn stream_native<'b>(
201 &'b self,
202 context: &'b HashMap<String, Value>,
203 format: crate::native::NativeFormat,
204 batch_size: usize,
205 ) -> Pin<Box<dyn Stream<Item = Result<crate::native::NativeBatch, FaucetError>> + Send + 'b>>
206 {
207 self.inner.stream_native(context, format, batch_size)
208 }
209
210 fn stream_pages<'b>(
211 &'b self,
212 context: &'b HashMap<String, Value>,
213 batch_size: usize,
214 ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'b>> {
215 let inner_stream = self.inner.stream_pages(context, batch_size);
216 let labels = self.labels.clone();
217 let connector = self.connector.clone();
218 let page_index = Arc::clone(&self.page_index);
219 let meter = self.meter.clone();
220 let metric_labels = self.metric_labels();
221 let pipeline = self.labels.pipeline.clone();
222 let row = self.labels.row.clone();
223
224 Box::pin(async_stream::try_stream! {
225 struct InFlightGuard(Vec<Label>);
228 impl Drop for InFlightGuard {
229 fn drop(&mut self) {
230 gauge!("faucet_source_in_flight", self.0.clone()).decrement(1.0);
231 }
232 }
233 gauge!("faucet_source_in_flight", metric_labels.clone()).increment(1.0);
234 let _in_flight = InFlightGuard(metric_labels.clone());
235
236 let mut inner = inner_stream;
237 loop {
238 let idx = page_index.fetch_add(1, Ordering::Relaxed);
239 let span = info_span!(
240 "faucet.source.page",
241 pipeline = %pipeline,
242 row = %row,
243 run_id = %labels.run_id,
244 connector = %connector,
245 page_index = idx,
246 );
247 let mut _timer = DurationGuard::new(
252 "faucet_source_page_duration_seconds",
253 metric_labels.clone(),
254 );
255
256 let next = AssertUnwindSafe(async {
257 use futures::StreamExt;
258 inner.next().await
259 })
260 .catch_unwind()
261 .instrument(span)
262 .await;
263
264 match next {
265 Ok(Some(Ok(page))) => {
266 counter!("faucet_source_pages_total", metric_labels.clone()).increment(1);
267 counter!("faucet_source_records_total", metric_labels.clone())
268 .increment(page.records.len() as u64);
269 if let Some(m) = &meter {
270 let bytes = estimate_page_bytes(&page.records);
271 counter!("faucet_source_bytes_total", metric_labels.clone())
272 .increment(bytes);
273 m.add_read(page.records.len() as u64, bytes);
274 }
275 _timer.record_now();
281 yield page;
282 }
283 Ok(Some(Err(e))) => {
284 let mut l = metric_labels.clone();
285 l.push(Label::new("kind", SharedString::const_str(error_kind(&e))));
286 counter!("faucet_source_errors_total", l).increment(1);
287 Err(e)?;
288 }
289 Ok(None) => {
290 _timer.disarm();
291 break;
292 }
293 Err(panic) => {
294 let mut l = metric_labels.clone();
295 l.push(Label::new("kind", SharedString::const_str("Panic")));
296 counter!("faucet_source_errors_total", l).increment(1);
297 let msg = panic.downcast_ref::<&'static str>().map(|s| (*s).to_string())
298 .or_else(|| panic.downcast_ref::<String>().cloned())
299 .unwrap_or_else(|| "<non-string panic payload>".to_string());
300 Err(FaucetError::Custom(format!("panic in source: {msg}").into()))?;
301 }
302 }
303 }
304 })
305 }
306}
307
308pub(crate) fn error_kind(e: &FaucetError) -> &'static str {
311 match e {
312 FaucetError::Http(_) => "Http",
313 FaucetError::HttpStatus { .. } => "HttpStatus",
314 FaucetError::Json(_) => "Json",
315 FaucetError::JsonPath(_) => "JsonPath",
316 FaucetError::Auth(_) => "Auth",
317 FaucetError::RateLimited { .. } => "RateLimited",
318 FaucetError::Url(_) => "Url",
319 FaucetError::Transform(_) => "Transform",
320 FaucetError::Config(_) => "Config",
321 FaucetError::Source(_) => "Source",
322 FaucetError::Sink(_) => "Sink",
323 FaucetError::QualityFailure { .. } => "QualityFailure",
324 FaucetError::SchemaDrift { .. } => "SchemaDrift",
325 FaucetError::ProfileDrift { .. } => "ProfileDrift",
326 FaucetError::PolicyViolation { .. } => "PolicyViolation",
327 FaucetError::BudgetExceeded { .. } => "BudgetExceeded",
328 FaucetError::ContractViolation { .. } => "ContractViolation",
329 FaucetError::State(_) => "State",
330 FaucetError::StateIncompatible { .. } => "StateIncompatible",
331 FaucetError::CircuitOpen { .. } => "CircuitOpen",
332 FaucetError::Custom(_) => "Custom",
333 }
334}
335
336pub struct InstrumentedSink<'a, S: Sink + ?Sized> {
339 inner: &'a S,
340 labels: Labels,
341 connector: SharedString,
342 base_labels: Vec<Label>,
344 meter: Option<Arc<UsageMeter>>,
346}
347
348impl<'a, S: Sink + ?Sized> InstrumentedSink<'a, S> {
349 pub fn new(inner: &'a S, labels: Labels) -> Self {
350 let raw = inner.connector_name();
351 debug_assert!(
352 !raw.is_empty(),
353 "connector_name() must return a non-empty string"
354 );
355 let connector: SharedString = SharedString::const_str(guarded_connector_name(raw));
356 let base_labels = base_metric_labels(&labels, &connector);
357 Self {
358 inner,
359 labels,
360 connector,
361 base_labels,
362 meter: None,
363 }
364 }
365
366 pub fn with_meter(mut self, meter: Arc<UsageMeter>) -> Self {
369 self.meter = Some(meter);
370 self
371 }
372
373 fn metric_labels(&self) -> Vec<Label> {
374 self.base_labels.clone()
375 }
376
377 fn error_labels(&self, kind: &'static str) -> Vec<Label> {
378 let mut l = self.metric_labels();
379 l.push(Label::new("kind", SharedString::const_str(kind)));
380 l
381 }
382
383 fn meter_written(&self, records: &[Value], accepted: usize) {
386 let Some(m) = &self.meter else {
387 return;
388 };
389 let bytes = if accepted >= records.len() {
390 estimate_page_bytes(records)
391 } else {
392 estimate_page_bytes(&records[..accepted])
393 };
394 counter!("faucet_sink_bytes_total", self.metric_labels()).increment(bytes);
395 m.add_written(accepted as u64, bytes);
396 }
397}
398
399#[async_trait]
400impl<'a, S: Sink + ?Sized> Sink for InstrumentedSink<'a, S> {
401 fn connector_name(&self) -> &'static str {
402 guarded_connector_name(self.inner.connector_name())
406 }
407
408 fn set_roundtrip_recorder(
412 &self,
413 recorder: std::sync::Arc<crate::observability::RoundtripRecorder>,
414 ) {
415 self.inner.set_roundtrip_recorder(recorder);
416 }
417
418 fn dataset_uri(&self) -> String {
425 self.inner.dataset_uri()
426 }
427
428 async fn local_outputs(&self) -> Vec<crate::local_outputs::LocalOutput> {
429 self.inner.local_outputs().await
430 }
431
432 #[cfg(feature = "arrow")]
435 fn supports_columnar(&self) -> bool {
436 self.inner.supports_columnar()
437 }
438
439 #[cfg(feature = "arrow")]
440 async fn write_batch_columnar(
441 &self,
442 batch: &arrow::array::RecordBatch,
443 ) -> Result<usize, FaucetError> {
444 self.inner.write_batch_columnar(batch).await
445 }
446
447 fn native_load_capabilities(&self) -> Vec<crate::native::NativeLoadCapability> {
450 self.inner.native_load_capabilities()
451 }
452
453 async fn load_native(
454 &self,
455 batch: crate::native::NativeBatch,
456 scope: &str,
457 ctx: crate::native::NativeLoadContext,
458 ) -> Result<usize, FaucetError> {
459 self.inner.load_native(batch, scope, ctx).await
460 }
461
462 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
463 let span = info_span!(
464 "faucet.sink.write",
465 pipeline = %self.labels.pipeline,
466 row = %self.labels.row,
467 run_id = %self.labels.run_id,
468 connector = %self.connector,
469 records = records.len(),
470 );
471 let metric_labels = self.metric_labels();
472 gauge!("faucet_sink_in_flight", metric_labels.clone()).increment(1.0);
473
474 struct InFlightGuard(Vec<Label>);
477 impl Drop for InFlightGuard {
478 fn drop(&mut self) {
479 gauge!("faucet_sink_in_flight", self.0.clone()).decrement(1.0);
480 }
481 }
482 let _in_flight = InFlightGuard(metric_labels.clone());
483
484 let _timer =
485 DurationGuard::new("faucet_sink_write_duration_seconds", metric_labels.clone());
486
487 let result = AssertUnwindSafe(self.inner.write_batch(records))
488 .catch_unwind()
489 .instrument(span)
490 .await;
491
492 match result {
493 Ok(Ok(n)) => {
494 counter!("faucet_sink_writes_total", metric_labels.clone()).increment(1);
495 counter!("faucet_sink_records_total", metric_labels.clone()).increment(n as u64);
496 self.meter_written(records, n);
497 Ok(n)
498 }
499 Ok(Err(e)) => {
500 counter!(
501 "faucet_sink_errors_total",
502 self.error_labels(error_kind(&e))
503 )
504 .increment(1);
505 Err(e)
506 }
507 Err(panic) => {
508 counter!("faucet_sink_errors_total", self.error_labels("Panic")).increment(1);
509 let msg = panic
510 .downcast_ref::<&'static str>()
511 .map(|s| (*s).to_string())
512 .or_else(|| panic.downcast_ref::<String>().cloned())
513 .unwrap_or_else(|| "<non-string panic payload>".to_string());
514 Err(FaucetError::Custom(format!("panic in sink: {msg}").into()))
515 }
516 }
517 }
518
519 async fn write_batch_partial(
520 &self,
521 records: &[Value],
522 ) -> Result<Vec<crate::traits::RowOutcome>, FaucetError> {
523 let span = info_span!(
524 "faucet.sink.write_partial",
525 pipeline = %self.labels.pipeline,
526 row = %self.labels.row,
527 run_id = %self.labels.run_id,
528 connector = %self.connector,
529 records = records.len(),
530 );
531 let metric_labels = self.metric_labels();
532 gauge!("faucet_sink_in_flight", metric_labels.clone()).increment(1.0);
533
534 struct InFlightGuard(Vec<Label>);
537 impl Drop for InFlightGuard {
538 fn drop(&mut self) {
539 gauge!("faucet_sink_in_flight", self.0.clone()).decrement(1.0);
540 }
541 }
542 let _in_flight = InFlightGuard(metric_labels.clone());
543
544 let _timer =
545 DurationGuard::new("faucet_sink_write_duration_seconds", metric_labels.clone());
546
547 let result = AssertUnwindSafe(self.inner.write_batch_partial(records))
548 .catch_unwind()
549 .instrument(span)
550 .await;
551
552 match result {
553 Ok(Ok(outcomes)) => {
554 let success_count = outcomes.iter().filter(|o| o.is_ok()).count();
555 counter!("faucet_sink_writes_total", metric_labels.clone()).increment(1);
556 counter!("faucet_sink_records_total", metric_labels.clone())
557 .increment(success_count as u64);
558 if let Some(m) = &self.meter {
559 let bytes: u64 = outcomes
560 .iter()
561 .zip(records.iter())
562 .filter(|(o, _)| o.is_ok())
563 .map(|(_, r)| crate::usage::estimate_json_bytes(r))
564 .sum();
565 counter!("faucet_sink_bytes_total", metric_labels.clone()).increment(bytes);
566 m.add_written(success_count as u64, bytes);
567 }
568 Ok(outcomes)
569 }
570 Ok(Err(e)) => {
571 counter!(
572 "faucet_sink_errors_total",
573 self.error_labels(error_kind(&e))
574 )
575 .increment(1);
576 Err(e)
577 }
578 Err(panic) => {
579 counter!("faucet_sink_errors_total", self.error_labels("Panic")).increment(1);
580 let msg = panic
581 .downcast_ref::<&'static str>()
582 .map(|s| (*s).to_string())
583 .or_else(|| panic.downcast_ref::<String>().cloned())
584 .unwrap_or_else(|| "<non-string panic payload>".to_string());
585 Err(FaucetError::Custom(format!("panic in sink: {msg}").into()))
586 }
587 }
588 }
589
590 async fn flush(&self) -> Result<(), FaucetError> {
591 let span = info_span!(
592 "faucet.sink.flush",
593 pipeline = %self.labels.pipeline,
594 row = %self.labels.row,
595 run_id = %self.labels.run_id,
596 connector = %self.connector,
597 );
598 let metric_labels = self.metric_labels();
599 let _timer =
600 DurationGuard::new("faucet_sink_flush_duration_seconds", metric_labels.clone());
601
602 let result = AssertUnwindSafe(self.inner.flush())
603 .catch_unwind()
604 .instrument(span)
605 .await;
606
607 match result {
608 Ok(Ok(())) => Ok(()),
609 Ok(Err(e)) => {
610 counter!(
611 "faucet_sink_errors_total",
612 self.error_labels(error_kind(&e))
613 )
614 .increment(1);
615 Err(e)
616 }
617 Err(panic) => {
618 counter!("faucet_sink_errors_total", self.error_labels("Panic")).increment(1);
619 let msg = panic
620 .downcast_ref::<&'static str>()
621 .map(|s| (*s).to_string())
622 .or_else(|| panic.downcast_ref::<String>().cloned())
623 .unwrap_or_else(|| "<non-string panic payload>".to_string());
624 Err(FaucetError::Custom(format!("panic in flush: {msg}").into()))
625 }
626 }
627 }
628
629 async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
637 self.inner.current_schema().await
638 }
639
640 fn supports_schema_evolution(&self) -> bool {
641 self.inner.supports_schema_evolution()
642 }
643
644 async fn evolve_schema(
645 &self,
646 evolution: &crate::drift::SchemaEvolution,
647 ) -> Result<(), FaucetError> {
648 self.inner.evolve_schema(evolution).await
649 }
650
651 fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
652 self.inner.supported_write_modes()
653 }
654
655 fn supports_cleanup(&self) -> bool {
656 self.inner.supports_cleanup()
657 }
658
659 async fn cleanup_scope(
660 &self,
661 scope: &std::collections::BTreeMap<String, Value>,
662 seen: &crate::cleanup::SeenKeys,
663 ) -> Result<u64, FaucetError> {
664 self.inner.cleanup_scope(scope, seen).await
665 }
666
667 fn supports_idempotent_writes(&self) -> bool {
668 self.inner.supports_idempotent_writes()
669 }
670
671 fn sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
672 self.inner.sink_guarantee()
673 }
674
675 fn dedups_by_key(&self) -> bool {
676 self.inner.dedups_by_key()
677 }
678 fn batch_atomicity(&self) -> crate::dlq::BatchAtomicity {
679 self.inner.batch_atomicity()
680 }
681
682 async fn write_batch_idempotent(
683 &self,
684 records: &[Value],
685 scope: &str,
686 token: &str,
687 ) -> Result<usize, FaucetError> {
688 let n = self
689 .inner
690 .write_batch_idempotent(records, scope, token)
691 .await?;
692 self.meter_written(records, n);
693 Ok(n)
694 }
695
696 async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
697 self.inner.last_committed_token(scope).await
698 }
699
700 fn is_overwrite(&self) -> bool {
701 self.inner.is_overwrite()
702 }
703
704 async fn begin_overwrite(&self) -> Result<(), FaucetError> {
705 self.inner.begin_overwrite().await
706 }
707
708 async fn commit_overwrite(&self) -> Result<(), FaucetError> {
709 self.inner.commit_overwrite().await
710 }
711
712 async fn abort_overwrite(&self) -> Result<(), FaucetError> {
713 self.inner.abort_overwrite().await
714 }
715 async fn complete_run(&self) -> Result<(), FaucetError> {
716 self.inner.complete_run().await
717 }
718}
719
720#[cfg(test)]
721pub(crate) mod source_tests {
722 use super::*;
723 use async_trait::async_trait;
724 use futures::StreamExt;
725 use metrics_util::debugging::{DebugValue, DebuggingRecorder, Snapshotter};
726 use serde_json::json;
727 use std::sync::{Mutex, OnceLock};
728
729 #[tokio::test]
731 async fn routed_double_fetches_nothing() {
732 use crate::Source as _;
733 assert!(
734 RoutedSource
735 .fetch_with_context(&Default::default())
736 .await
737 .unwrap()
738 .is_empty()
739 );
740 }
741
742 struct RoutedSource;
743
744 #[async_trait::async_trait]
745 impl crate::Source for RoutedSource {
746 async fn fetch_with_context(
747 &self,
748 _ctx: &std::collections::HashMap<String, serde_json::Value>,
749 ) -> Result<Vec<serde_json::Value>, crate::FaucetError> {
750 Ok(Vec::new())
751 }
752
753 fn record_table(&self, record: &serde_json::Value) -> Option<String> {
754 record.get("t")?.as_str().map(str::to_string)
755 }
756
757 fn position_le(&self, a: &serde_json::Value, b: &serde_json::Value) -> Option<bool> {
758 Some(a.as_u64()? <= b.as_u64()?)
759 }
760
761 fn position_min(&self, _positions: &[serde_json::Value]) -> Option<serde_json::Value> {
762 Some(serde_json::json!("inner-min"))
763 }
764 }
765
766 fn assert_forwards_multi_table_hooks(s: &dyn crate::Source) {
767 assert_eq!(
768 s.record_table(&serde_json::json!({"t": "public.a"}))
769 .as_deref(),
770 Some("public.a")
771 );
772 assert_eq!(
773 s.position_le(&serde_json::json!(1), &serde_json::json!(2)),
774 Some(true)
775 );
776 assert_eq!(
777 s.position_le(&serde_json::json!(3), &serde_json::json!(2)),
778 Some(false)
779 );
780 assert_eq!(
781 s.position_min(&[serde_json::json!(1), serde_json::json!(2)]),
782 Some(serde_json::json!("inner-min"))
783 );
784 }
785
786 #[test]
787 fn multi_table_hooks_are_forwarded_to_the_inner_source() {
788 let inner = RoutedSource;
789 let wrapped = InstrumentedSource::new(&inner, labels());
790 assert_forwards_multi_table_hooks(&wrapped);
791 wrapped.set_run_clock(chrono::Utc::now());
792 }
793
794 pub(crate) static LOCK: Mutex<()> = Mutex::new(());
797 static SNAPSHOTTER: OnceLock<Snapshotter> = OnceLock::new();
798
799 pub(crate) fn snapshotter() -> &'static Snapshotter {
800 SNAPSHOTTER.get_or_init(|| {
801 let recorder = DebuggingRecorder::new();
802 let snap = recorder.snapshotter();
803 let _ = metrics::set_global_recorder(recorder);
811 snap
812 })
813 }
814
815 pub(in crate::observability) fn labels() -> Labels {
816 Labels::new("p", "r", "rid")
817 }
818
819 struct MockSource(Vec<Value>);
820 #[async_trait]
821 impl Source for MockSource {
822 async fn fetch_with_context(
823 &self,
824 _: &HashMap<String, Value>,
825 ) -> Result<Vec<Value>, FaucetError> {
826 Ok(self.0.clone())
827 }
828 fn connector_name(&self) -> &'static str {
829 "mock"
830 }
831 }
832
833 struct PanickingSource;
834 #[async_trait]
835 impl Source for PanickingSource {
836 async fn fetch_with_context(
837 &self,
838 _: &HashMap<String, Value>,
839 ) -> Result<Vec<Value>, FaucetError> {
840 panic!("kaboom")
841 }
842 fn connector_name(&self) -> &'static str {
843 "panic-test"
844 }
845 }
846
847 struct EmptyNameSource;
851 #[async_trait]
852 impl Source for EmptyNameSource {
853 async fn fetch_with_context(
854 &self,
855 _: &HashMap<String, Value>,
856 ) -> Result<Vec<Value>, FaucetError> {
857 Ok(vec![])
858 }
859 fn connector_name(&self) -> &'static str {
860 ""
861 }
862 }
863
864 #[test]
865 fn empty_inner_connector_name_falls_back_to_unknown() {
866 let inner = EmptyNameSource;
867 let wrapped = InstrumentedSource {
871 inner: &inner,
872 labels: labels(),
873 connector: SharedString::const_str("unknown"),
874 base_labels: Vec::new(),
875 page_index: Arc::new(AtomicUsize::new(0)),
876 meter: None,
877 };
878 assert_eq!(
879 Source::connector_name(&wrapped),
880 "unknown",
881 "instrumented source must not leak an empty connector name"
882 );
883 }
884
885 #[tokio::test]
886 #[allow(clippy::await_holding_lock)]
887 async fn records_records_counter_per_page() {
888 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
889 let snap = snapshotter();
890 let inner = MockSource((0..5).map(|i| json!({"i": i})).collect());
891 let wrapped = InstrumentedSource::new(&inner, labels());
892 let ctx = HashMap::new();
893 let mut s = wrapped.stream_pages(&ctx, 2);
894 while s.next().await.is_some() {}
895 let snapshot = snap.snapshot();
896 let records: u64 = snapshot
897 .into_vec()
898 .into_iter()
899 .filter_map(|(key, _u, _d, v)| {
900 if key.key().name() == "faucet_source_records_total"
901 && let DebugValue::Counter(c) = v
902 {
903 return Some(c);
904 }
905 None
906 })
907 .sum();
908 assert!(
909 records >= 5,
910 "expected at least 5 records counted, got {records}"
911 );
912 }
913
914 struct PageCountSource(Vec<Value>);
917 #[async_trait]
918 impl Source for PageCountSource {
919 async fn fetch_with_context(
920 &self,
921 _: &HashMap<String, Value>,
922 ) -> Result<Vec<Value>, FaucetError> {
923 Ok(self.0.clone())
924 }
925 fn connector_name(&self) -> &'static str {
926 "page-count-probe"
927 }
928 }
929
930 #[tokio::test]
931 #[allow(clippy::await_holding_lock)]
932 async fn page_duration_records_one_sample_per_yielded_page() {
933 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
936 let snap = snapshotter();
937 let inner = PageCountSource((0..5).map(|i| json!({"i": i})).collect());
938 let wrapped = InstrumentedSource::new(&inner, labels());
939 let ctx = HashMap::new();
940 let mut s = wrapped.stream_pages(&ctx, 2);
941 let mut pages = 0usize;
942 while s.next().await.is_some() {
943 pages += 1;
944 }
945 assert_eq!(pages, 3, "expected 3 yielded pages");
946
947 let snapshot = snap.snapshot();
948 let samples: usize = snapshot
949 .into_vec()
950 .into_iter()
951 .filter_map(|(key, _u, _d, v)| {
952 if key.key().name() == "faucet_source_page_duration_seconds"
953 && key
954 .key()
955 .labels()
956 .any(|l| l.key() == "connector" && l.value() == "page-count-probe")
957 && let DebugValue::Histogram(h) = v
958 {
959 return Some(h.len());
960 }
961 None
962 })
963 .sum();
964 assert_eq!(
965 samples, pages,
966 "page-duration histogram must have exactly one sample per yielded \
967 page ({pages}), not page+1 (no spurious terminal sample)"
968 );
969 }
970
971 #[tokio::test]
972 #[allow(clippy::await_holding_lock)]
973 async fn maps_panic_to_custom_error_with_kind_panic() {
974 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
975 let _snap = snapshotter();
976 let inner = PanickingSource;
977 let wrapped = InstrumentedSource::new(&inner, labels());
978 let ctx = HashMap::new();
979 let mut s = wrapped.stream_pages(&ctx, 10);
980 let first = s
981 .next()
982 .await
983 .expect("stream yields at least one item before terminating");
984 assert!(matches!(first, Err(FaucetError::Custom(_))));
985 }
987
988 #[test]
991 fn error_kind_covers_all_variants() {
992 use std::time::Duration;
993 let cases: Vec<(FaucetError, &str)> = vec![
998 (
999 FaucetError::HttpStatus {
1000 status: 500,
1001 url: "u".into(),
1002 body: "b".into(),
1003 },
1004 "HttpStatus",
1005 ),
1006 (
1007 FaucetError::Json(serde_json::from_str::<Value>("nope").unwrap_err()),
1008 "Json",
1009 ),
1010 (FaucetError::JsonPath("bad".into()), "JsonPath"),
1011 (FaucetError::Auth("a".into()), "Auth"),
1012 (
1013 FaucetError::RateLimited(Duration::from_secs(1)),
1014 "RateLimited",
1015 ),
1016 (FaucetError::Url("bad url".into()), "Url"),
1017 (FaucetError::Transform("t".into()), "Transform"),
1018 (FaucetError::Config("c".into()), "Config"),
1019 (FaucetError::Source("s".into()), "Source"),
1020 (FaucetError::Sink("s".into()), "Sink"),
1021 (
1022 FaucetError::QualityFailure {
1023 check: "chk".into(),
1024 message: "m".into(),
1025 },
1026 "QualityFailure",
1027 ),
1028 (FaucetError::State("st".into()), "State"),
1029 (
1030 FaucetError::CircuitOpen {
1031 failures: 3,
1032 cooldown: Duration::from_secs(60),
1033 },
1034 "CircuitOpen",
1035 ),
1036 (
1037 FaucetError::Custom(Box::new(std::io::Error::other("boom"))),
1038 "Custom",
1039 ),
1040 ];
1041 for (err, expected) in cases {
1042 assert_eq!(error_kind(&err), expected, "mismatch for {err:?}");
1043 }
1044 }
1045
1046 struct PassthroughSource {
1052 seen_bookmark: Mutex<Option<Value>>,
1053 }
1054 #[async_trait]
1055 impl Source for PassthroughSource {
1056 async fn fetch_with_context(
1057 &self,
1058 _: &HashMap<String, Value>,
1059 ) -> Result<Vec<Value>, FaucetError> {
1060 Ok(vec![json!({"fwc": 1})])
1061 }
1062 async fn fetch_with_context_incremental(
1063 &self,
1064 _: &HashMap<String, Value>,
1065 ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
1066 Ok((vec![json!({"inc": 1})], Some(json!("bm"))))
1067 }
1068 fn state_key(&self) -> Option<String> {
1069 Some("passthrough_key".into())
1070 }
1071 async fn apply_start_bookmark(&self, bookmark: Value) -> Result<(), FaucetError> {
1072 *self.seen_bookmark.lock().unwrap() = Some(bookmark);
1073 Ok(())
1074 }
1075 fn connector_name(&self) -> &'static str {
1076 "passthrough"
1077 }
1078 }
1079
1080 #[tokio::test]
1081 async fn source_passthroughs_delegate_to_inner() {
1082 let inner = PassthroughSource {
1083 seen_bookmark: Mutex::new(None),
1084 };
1085 let wrapped = InstrumentedSource::new(&inner, labels());
1086
1087 assert_eq!(wrapped.state_key(), Some("passthrough_key".to_string()));
1089
1090 let ctx = HashMap::new();
1092 assert_eq!(
1093 wrapped.fetch_with_context(&ctx).await.unwrap(),
1094 vec![json!({"fwc": 1})]
1095 );
1096
1097 let (recs, bm) = wrapped.fetch_with_context_incremental(&ctx).await.unwrap();
1099 assert_eq!(recs, vec![json!({"inc": 1})]);
1100 assert_eq!(bm, Some(json!("bm")));
1101
1102 wrapped.apply_start_bookmark(json!("resume")).await.unwrap();
1104 assert_eq!(
1105 *inner.seen_bookmark.lock().unwrap(),
1106 Some(json!("resume")),
1107 "apply_start_bookmark must reach the inner source"
1108 );
1109
1110 assert!(!wrapped.supports_exactly_once());
1112 assert_eq!(
1113 wrapped.replay_guarantee(),
1114 crate::idempotency::ReplayGuarantee::NonDeterministic
1115 );
1116 assert_eq!(wrapped.capture_resume_position().await.unwrap(), None);
1117 assert_eq!(wrapped.lag().await.unwrap(), None);
1118 }
1119
1120 struct ExactlyOnceSource;
1123 #[async_trait]
1124 impl Source for ExactlyOnceSource {
1125 async fn fetch_with_context(
1126 &self,
1127 _context: &HashMap<String, Value>,
1128 ) -> Result<Vec<Value>, FaucetError> {
1129 Ok(vec![])
1130 }
1131 fn supports_exactly_once(&self) -> bool {
1132 true
1133 }
1134 async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
1135 Ok(Some(json!("pos")))
1136 }
1137 fn connector_name(&self) -> &'static str {
1138 "eo-source"
1139 }
1140 }
1141
1142 #[tokio::test]
1143 async fn source_capability_passthroughs_delegate_to_inner() {
1144 let inner = ExactlyOnceSource;
1145 let wrapped = InstrumentedSource::new(&inner, labels());
1146 assert!(wrapped.supports_exactly_once());
1147 assert_eq!(
1148 wrapped.replay_guarantee(),
1149 crate::idempotency::ReplayGuarantee::Deterministic,
1150 "typed capability derives through the wrapper"
1151 );
1152 assert_eq!(
1153 wrapped.capture_resume_position().await.unwrap(),
1154 Some(json!("pos"))
1155 );
1156 }
1157}
1158
1159#[cfg(test)]
1160mod sink_tests {
1161 use super::source_tests::{LOCK, labels, snapshotter};
1162 use super::*;
1163 use async_trait::async_trait;
1164 use metrics_util::debugging::DebugValue;
1165 use serde_json::json;
1166
1167 struct MockSink(std::sync::Mutex<Vec<Value>>);
1168 #[async_trait]
1169 impl Sink for MockSink {
1170 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
1171 self.0.lock().unwrap().extend(records.iter().cloned());
1172 Ok(records.len())
1173 }
1174 fn connector_name(&self) -> &'static str {
1175 "mock-sink"
1176 }
1177 }
1178
1179 struct FailingSink;
1180 #[async_trait]
1181 impl Sink for FailingSink {
1182 async fn write_batch(&self, _: &[Value]) -> Result<usize, FaucetError> {
1183 Err(FaucetError::Sink("nope".into()))
1184 }
1185 fn connector_name(&self) -> &'static str {
1186 "failing-sink"
1187 }
1188 }
1189
1190 struct EmptyNameSink;
1191 #[async_trait]
1192 impl Sink for EmptyNameSink {
1193 async fn write_batch(&self, _: &[Value]) -> Result<usize, FaucetError> {
1194 Ok(0)
1195 }
1196 fn connector_name(&self) -> &'static str {
1197 ""
1198 }
1199 }
1200
1201 #[test]
1202 fn empty_inner_connector_name_falls_back_to_unknown() {
1203 let inner = EmptyNameSink;
1204 let wrapped = InstrumentedSink {
1208 inner: &inner,
1209 labels: labels(),
1210 connector: SharedString::const_str("unknown"),
1211 base_labels: Vec::new(),
1212 meter: None,
1213 };
1214 assert_eq!(
1215 Sink::connector_name(&wrapped),
1216 "unknown",
1217 "instrumented sink must not leak an empty connector name"
1218 );
1219 }
1220
1221 struct CapableSink;
1229 #[async_trait]
1230 impl Sink for CapableSink {
1231 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
1232 Ok(records.len())
1233 }
1234 fn connector_name(&self) -> &'static str {
1235 "capable-sink"
1236 }
1237 async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
1238 Ok(Some(
1239 json!({"type": "object", "properties": {"id": {"type": "integer"}}}),
1240 ))
1241 }
1242 fn supports_schema_evolution(&self) -> bool {
1243 true
1244 }
1245 fn supports_idempotent_writes(&self) -> bool {
1246 true
1247 }
1248 fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
1249 &[
1250 crate::write_mode::WriteMode::Append,
1251 crate::write_mode::WriteMode::Upsert,
1252 ]
1253 }
1254 async fn last_committed_token(&self, _scope: &str) -> Result<Option<String>, FaucetError> {
1255 Ok(Some("tok-1".into()))
1256 }
1257 fn dedups_by_key(&self) -> bool {
1258 true
1259 }
1260 }
1261
1262 #[tokio::test]
1263 async fn instrumented_sink_forwards_capability_methods_to_inner() {
1264 let inner = CapableSink;
1265 let wrapped = InstrumentedSink::new(&inner, labels());
1266
1267 assert_eq!(
1270 wrapped.current_schema().await.unwrap(),
1271 Some(json!({"type": "object", "properties": {"id": {"type": "integer"}}})),
1272 "current_schema must delegate to the inner sink"
1273 );
1274 assert!(
1275 wrapped.supports_schema_evolution(),
1276 "supports_schema_evolution must delegate"
1277 );
1278 assert!(
1280 wrapped.supports_idempotent_writes(),
1281 "supports_idempotent_writes must delegate (exactly-once)"
1282 );
1283 assert!(
1284 wrapped
1285 .supported_write_modes()
1286 .contains(&crate::write_mode::WriteMode::Upsert),
1287 "supported_write_modes must delegate"
1288 );
1289 assert_eq!(
1290 wrapped.last_committed_token("scope").await.unwrap(),
1291 Some("tok-1".to_string()),
1292 "last_committed_token must delegate"
1293 );
1294 assert_eq!(
1297 wrapped.sink_guarantee(),
1298 crate::idempotency::SinkGuarantee::AtomicWatermark,
1299 "sink_guarantee must delegate"
1300 );
1301 assert!(wrapped.dedups_by_key(), "dedups_by_key must delegate");
1302 assert_eq!(
1303 wrapped.batch_atomicity(),
1304 crate::dlq::BatchAtomicity::BestEffort
1305 );
1306 }
1307
1308 #[tokio::test]
1309 #[allow(clippy::await_holding_lock)]
1310 async fn records_writes_and_records_counters() {
1311 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1312 let snap = snapshotter();
1313 let inner = MockSink(std::sync::Mutex::new(Vec::new()));
1314 let wrapped = InstrumentedSink::new(&inner, labels());
1315 wrapped
1316 .write_batch(&[json!({"a": 1}), json!({"a": 2})])
1317 .await
1318 .unwrap();
1319 let snapshot = snap.snapshot();
1320 let writes: u64 = snapshot
1321 .into_vec()
1322 .into_iter()
1323 .filter_map(|(key, _u, _d, v)| {
1324 if key.key().name() == "faucet_sink_writes_total"
1325 && let DebugValue::Counter(c) = v
1326 {
1327 return Some(c);
1328 }
1329 None
1330 })
1331 .sum();
1332 assert!(writes >= 1, "expected at least one write counted");
1333 }
1334
1335 #[tokio::test]
1336 #[allow(clippy::await_holding_lock)]
1337 async fn error_increments_errors_total_with_kind() {
1338 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1339 let snap = snapshotter();
1340 let inner = FailingSink;
1341 let wrapped = InstrumentedSink::new(&inner, labels());
1342 let _ = wrapped.write_batch(&[json!({})]).await;
1343 let snapshot = snap.snapshot();
1344 let found = snapshot.into_vec().into_iter().any(|(key, _u, _d, v)| {
1345 key.key().name() == "faucet_sink_errors_total"
1346 && key
1347 .key()
1348 .labels()
1349 .any(|l| l.key() == "kind" && l.value() == "Sink")
1350 && matches!(v, DebugValue::Counter(c) if c >= 1)
1351 });
1352 assert!(found, "expected sink_errors_total with kind=Sink");
1353 }
1354
1355 #[tokio::test]
1356 #[allow(clippy::await_holding_lock)]
1357 async fn instrumented_sink_write_batch_partial_counts_successful_outcomes() {
1358 use crate::traits::RowOutcome;
1359 use metrics_util::debugging::DebugValue;
1360
1361 struct MixedSink;
1363 #[async_trait]
1364 impl Sink for MixedSink {
1365 async fn write_batch(&self, _r: &[Value]) -> Result<usize, FaucetError> {
1366 unreachable!()
1367 }
1368 async fn write_batch_partial(
1369 &self,
1370 _r: &[Value],
1371 ) -> Result<Vec<RowOutcome>, FaucetError> {
1372 Ok(vec![
1373 Ok(()),
1374 Err(FaucetError::Sink("bad row".into())),
1375 Ok(()),
1376 ])
1377 }
1378 fn connector_name(&self) -> &'static str {
1379 "mixed"
1380 }
1381 }
1382
1383 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1384 let snap = snapshotter();
1385
1386 let inner = MixedSink;
1387 let wrapped = InstrumentedSink::new(&inner, labels());
1388 let _ = wrapped
1389 .write_batch_partial(&[json!({}), json!({}), json!({})])
1390 .await
1391 .unwrap();
1392
1393 let snapshot = snap.snapshot();
1401 let records: u64 = snapshot
1402 .into_vec()
1403 .into_iter()
1404 .filter_map(|(k, _u, _d, v): (metrics_util::CompositeKey, _, _, _)| {
1405 if k.key().name() == "faucet_sink_records_total"
1406 && k.key()
1407 .labels()
1408 .any(|l| l.key() == "connector" && l.value() == "mixed")
1409 && let DebugValue::Counter(c) = v
1410 {
1411 Some(c)
1412 } else {
1413 None
1414 }
1415 })
1416 .sum();
1417 assert!(
1418 records >= 2,
1419 "expected faucet_sink_records_total{{connector=mixed}} >= 2, got {records}"
1420 );
1421 }
1422
1423 #[tokio::test]
1426 #[allow(clippy::await_holding_lock)]
1427 async fn flush_error_increments_errors_total_and_propagates() {
1428 struct FlushFailSink;
1431 #[async_trait]
1432 impl Sink for FlushFailSink {
1433 async fn write_batch(&self, r: &[Value]) -> Result<usize, FaucetError> {
1434 Ok(r.len())
1435 }
1436 async fn flush(&self) -> Result<(), FaucetError> {
1437 Err(FaucetError::Sink("flush boom".into()))
1438 }
1439 fn connector_name(&self) -> &'static str {
1440 "flush-fail-sink"
1441 }
1442 }
1443
1444 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1445 let snap = snapshotter();
1446 let inner = FlushFailSink;
1447 let wrapped = InstrumentedSink::new(&inner, labels());
1448 let err = wrapped.flush().await.unwrap_err();
1449 assert!(matches!(&err, FaucetError::Sink(m) if m.contains("flush boom")));
1450
1451 let snapshot = snap.snapshot();
1452 let found = snapshot.into_vec().into_iter().any(|(key, _u, _d, v)| {
1453 key.key().name() == "faucet_sink_errors_total"
1454 && key
1455 .key()
1456 .labels()
1457 .any(|l| l.key() == "connector" && l.value() == "flush-fail-sink")
1458 && key
1459 .key()
1460 .labels()
1461 .any(|l| l.key() == "kind" && l.value() == "Sink")
1462 && matches!(v, DebugValue::Counter(c) if c >= 1)
1463 });
1464 assert!(
1465 found,
1466 "expected sink_errors_total{{connector=flush-fail-sink,kind=Sink}}"
1467 );
1468 }
1469
1470 struct PanickingSink;
1473 #[async_trait]
1474 impl Sink for PanickingSink {
1475 async fn write_batch(&self, _: &[Value]) -> Result<usize, FaucetError> {
1476 panic!("write kaboom")
1477 }
1478 async fn write_batch_partial(
1479 &self,
1480 _: &[Value],
1481 ) -> Result<Vec<crate::traits::RowOutcome>, FaucetError> {
1482 panic!("partial kaboom")
1483 }
1484 async fn flush(&self) -> Result<(), FaucetError> {
1485 panic!("flush kaboom")
1486 }
1487 fn connector_name(&self) -> &'static str {
1488 "panic-sink"
1489 }
1490 }
1491
1492 #[tokio::test]
1493 #[allow(clippy::await_holding_lock)]
1494 async fn write_batch_panic_maps_to_custom_error() {
1495 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1496 let _snap = snapshotter();
1497 let inner = PanickingSink;
1498 let wrapped = InstrumentedSink::new(&inner, labels());
1499 let err = wrapped.write_batch(&[json!({})]).await.unwrap_err();
1500 match err {
1501 FaucetError::Custom(b) => {
1502 assert!(b.to_string().contains("panic in sink: write kaboom"))
1503 }
1504 other => panic!("expected Custom panic error, got {other:?}"),
1505 }
1506 }
1507
1508 #[tokio::test]
1509 #[allow(clippy::await_holding_lock)]
1510 async fn write_batch_partial_panic_maps_to_custom_error() {
1511 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1512 let _snap = snapshotter();
1513 let inner = PanickingSink;
1514 let wrapped = InstrumentedSink::new(&inner, labels());
1515 let err = wrapped.write_batch_partial(&[json!({})]).await.unwrap_err();
1516 match err {
1517 FaucetError::Custom(b) => {
1518 assert!(b.to_string().contains("panic in sink: partial kaboom"))
1519 }
1520 other => panic!("expected Custom panic error, got {other:?}"),
1521 }
1522 }
1523
1524 #[tokio::test]
1525 #[allow(clippy::await_holding_lock)]
1526 async fn flush_panic_maps_to_custom_error() {
1527 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
1528 let _snap = snapshotter();
1529 let inner = PanickingSink;
1530 let wrapped = InstrumentedSink::new(&inner, labels());
1531 let err = wrapped.flush().await.unwrap_err();
1532 match err {
1533 FaucetError::Custom(b) => {
1534 assert!(b.to_string().contains("panic in flush: flush kaboom"))
1535 }
1536 other => panic!("expected Custom panic error, got {other:?}"),
1537 }
1538 }
1539 #[tokio::test]
1553 async fn instrumented_sink_forwards_identity_and_local_outputs() {
1554 use crate::local_outputs::{LocalOutput, LocalOutputLog};
1555
1556 struct FileSink {
1557 outputs: LocalOutputLog,
1558 }
1559
1560 #[async_trait]
1561 impl Sink for FileSink {
1562 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
1563 Ok(records.len())
1564 }
1565 fn connector_name(&self) -> &'static str {
1566 "jsonl"
1567 }
1568 fn dataset_uri(&self) -> String {
1569 "file:///tmp/out.jsonl".to_string()
1570 }
1571 async fn local_outputs(&self) -> Vec<LocalOutput> {
1572 self.outputs.snapshot()
1573 }
1574 }
1575
1576 let outputs = LocalOutputLog::new();
1577 outputs.record_open("/tmp/out.jsonl", false);
1578 let inner = FileSink { outputs };
1579 let sink = InstrumentedSink::new(&inner, Labels::new("p", "r", "run-1"));
1580
1581 let reported = sink.local_outputs().await;
1582 assert_eq!(
1583 reported.len(),
1584 1,
1585 "InstrumentedSink must not hide the inner sink's files from the GC"
1586 );
1587 assert_eq!(reported[0].path, std::path::PathBuf::from("/tmp/out.jsonl"));
1588 assert!(
1589 !reported[0].pre_existing,
1590 "classification must survive verbatim"
1591 );
1592 assert_eq!(sink.dataset_uri(), "file:///tmp/out.jsonl");
1593 }
1594}