Skip to main content

faucet_core/observability/
decorator.rs

1//! Pipeline-internal decorators that emit spans + metrics around every
2//! source / sink trait call. See the design spec for the full vocabulary.
3
4use 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
22/// Guard an inner connector's `connector_name()` so an empty string maps to
23/// the `"unknown"` fallback. Used both for the `connector` metric label and the
24/// `connector_name()` passthrough so the two never disagree.
25fn guarded_connector_name(raw: &'static str) -> &'static str {
26    if raw.is_empty() { "unknown" } else { raw }
27}
28
29/// Build the base `pipeline` / `row` / `connector` label vec once. The two
30/// `pipeline` / `row` heap allocations and the vec construction happen a single
31/// time at decorator construction; per-call sites `clone()` this instead of
32/// rebuilding from the `Arc<str>` labels on every page / write / flush.
33fn 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
41/// Wraps a `&dyn Source` (or any `&S: Source`) and emits spans + metrics
42/// around every call. Constructed by `Pipeline::run` and never exposed to
43/// end users; the wrapped source remains the user-facing object.
44pub struct InstrumentedSource<'a, S: Source + ?Sized> {
45    inner: &'a S,
46    labels: Labels,
47    connector: SharedString,
48    /// Precomputed `pipeline` / `row` / `connector` labels, cloned per call.
49    base_labels: Vec<Label>,
50    page_index: Arc<AtomicUsize>,
51    /// The run's usage meter (#704); `None` = count nothing beyond metrics.
52    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    /// Tally records and estimated bytes into a run's usage meter (#704).
75    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    /// Returns `metric_labels()` with an additional `kind` label appended.
85    /// Used by `InstrumentedSink::write_batch` (Task 9) and any future
86    /// instrumentation paths where `self` is in scope.
87    #[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        // Return the guarded name so an inner connector that returns "" maps to
99        // the "unknown" fallback — keeping this passthrough consistent with the
100        // `connector` metric label rather than leaking an empty string.
101        guarded_connector_name(self.inner.connector_name())
102    }
103
104    /// Forward the round-trip recorder to the wrapped connector (#638) — the
105    /// decorator sits between the pipeline and the connector, so without this
106    /// the hook would never reach the code that performs the I/O.
107    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        // Library-call path; the pipeline drives through stream_pages.
166        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    // Columnar fast path (feature `arrow`): forward transparently to the inner
177    // source. The columnar streaming loop in `pipeline.rs` emits the source
178    // metrics itself, so no instrumentation is layered here (RFC 0002 / #375).
179    #[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    // Native byte-passthrough capability (#633) forwards to the inner source; the
195    // native streaming loop in `pipeline.rs` emits the source metrics itself.
196    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            // In-flight gauge tracks open streams. Decrement on drop so
226            // cancellation leaves the gauge consistent.
227            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                // Armed across the poll so a cancelled / panicking page-fetch
248                // still records the time spent. Disarmed on the terminal empty
249                // poll (`Ok(None)`) so end-of-stream doesn't record a spurious
250                // ~0 sample into the page-duration histogram.
251                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                        // Close the timing window BEFORE yielding: in an
276                        // `async_stream` the timer local persists across the
277                        // yield, so dropping it at scope-exit would fold the
278                        // downstream sink/consumer latency into the source's
279                        // page-duration histogram (audit #321 M10).
280                        _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
308/// Map a `FaucetError` variant to its stable `kind` label value. The match
309/// must be exhaustive; update when new variants are added.
310pub(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
336/// Wraps a `&dyn Sink` (or any `&S: Sink`) and emits spans + metrics around
337/// `write_batch` and `flush`. Constructed by `Pipeline::run`.
338pub struct InstrumentedSink<'a, S: Sink + ?Sized> {
339    inner: &'a S,
340    labels: Labels,
341    connector: SharedString,
342    /// Precomputed `pipeline` / `row` / `connector` labels, cloned per call.
343    base_labels: Vec<Label>,
344    /// The run's usage meter (#704); `None` = count nothing beyond metrics.
345    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    /// Tally accepted records and estimated bytes into a run's usage meter
367    /// (#704).
368    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    /// Count `accepted` records of `records` as written; when the sink
384    /// accepted a prefix, only that prefix's estimated size is attributed.
385    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        // Return the guarded name so an inner connector that returns "" maps to
403        // the "unknown" fallback — keeping this passthrough consistent with the
404        // `connector` metric label rather than leaking an empty string.
405        guarded_connector_name(self.inner.connector_name())
406    }
407
408    /// Forward the round-trip recorder to the wrapped connector (#638) — the
409    /// decorator sits between the pipeline and the connector, so without this
410    /// the hook would never reach the code that performs the I/O.
411    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    // Identity + provenance passthroughs. Instrumentation must be invisible to
419    // anything asking the sink *what* it is or *what it wrote* — a decorator that
420    // falls back to the trait defaults reports `"<name>://unknown"` and an empty
421    // output list, which for `local_outputs` means the retention GC (#587) never
422    // learns about files this sink created and can never reclaim them. Silent, and
423    // only observable as disk filling up.
424    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    // Columnar fast path (feature `arrow`): forward transparently to the inner
433    // sink; the columnar loop in `pipeline.rs` emits the sink metrics (RFC 0002).
434    #[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    // Native byte-passthrough load (#633): forward to the inner sink; the native
448    // loop in `pipeline.rs` emits the sink metrics.
449    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        // RAII guard ensures the gauge is decremented even if write_batch
475        // panics or the future is cancelled.
476        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        // RAII guard ensures the gauge is decremented even if write_batch_partial
535        // panics or the future is cancelled.
536        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    // ── Non-instrumented passthroughs ────────────────────────────────────────
630    // These carry no per-call metric/span of their own, but they MUST delegate
631    // to the inner sink — the `Sink` trait gives each a default that disables
632    // the corresponding feature (schema-drift, upsert, exactly-once). Because
633    // the pipeline drives the *wrapped* sink, failing to forward them silently
634    // makes those features inert through the entire CLI/observability path.
635
636    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    /// A change-stream double with its own multi-table hooks (#731).
730    #[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    // Process-global recorder shared across all observability tests in this
795    // crate. Task 5 established the same pattern.
796    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            // First test installs; the OnceLock guarantees we never install
804            // twice. If something else (e.g. the timer test) already installed
805            // a recorder, `set_global_recorder` will Err — but in that case
806            // *our* snapshotter is disconnected from the live recorder. The
807            // workaround is for all observability tests to share one source of
808            // truth — this file. If a future test elsewhere installs a
809            // recorder first, restructure so all tests share this OnceLock.
810            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    // Inner connector that returns an empty name. The instrumented wrapper must
848    // map this to the `"unknown"` fallback so the `connector_name()` passthrough
849    // never disagrees with the `connector` metric label.
850    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        // `InstrumentedSource::new` debug_asserts on an empty inner name, so
868        // build the wrapper directly with the fallback name to exercise the
869        // passthrough without tripping the assertion in debug builds.
870        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    // Source with a unique connector name so the page-duration histogram for
915    // this run can be isolated in the shared global recorder.
916    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        // 5 records at batch_size 2 → pages [2, 2, 1] = 3 yielded pages. The
934        // terminal empty poll must NOT add a 4th (spurious ~0) sample.
935        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        // Process did not abort — implicit by reaching this line.
986    }
987
988    // ── error_kind: exhaustive variant → label mapping ───────────────────────
989
990    #[test]
991    fn error_kind_covers_all_variants() {
992        use std::time::Duration;
993        // Build one of every non-`Http` FaucetError variant and assert its
994        // stable label. (`Http` wraps a `reqwest::Error`, which has no public
995        // constructor; it is exercised through the live request paths in the
996        // connector crates' tests.)
997        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    // ── Source passthrough methods ───────────────────────────────────────────
1047
1048    // A source that overrides every passthrough so the instrumented wrapper's
1049    // delegating methods (state_key / apply_start_bookmark / fetch_with_context
1050    // / fetch_with_context_incremental) are exercised.
1051    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        // state_key passthrough
1088        assert_eq!(wrapped.state_key(), Some("passthrough_key".to_string()));
1089
1090        // fetch_with_context passthrough
1091        let ctx = HashMap::new();
1092        assert_eq!(
1093            wrapped.fetch_with_context(&ctx).await.unwrap(),
1094            vec![json!({"fwc": 1})]
1095        );
1096
1097        // fetch_with_context_incremental passthrough
1098        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        // apply_start_bookmark passthrough
1103        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        // capability passthroughs: defaults for this inner source…
1111        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    /// A source advertising exactly-once — the decorator must not mask it
1121    /// (the pipeline's mechanism selection reads these through the wrapper).
1122    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        // `InstrumentedSink::new` debug_asserts on an empty inner name, so build
1205        // the wrapper directly with the fallback name to exercise the
1206        // passthrough without tripping the assertion in debug builds.
1207        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    /// Regression (#194): the pipeline drives the *wrapped* sink, so
1222    /// `InstrumentedSink` MUST forward the capability methods to the inner sink.
1223    /// Before this was fixed, the trait defaults (`current_schema -> None`,
1224    /// `supports_schema_evolution -> false`, `supports_idempotent_writes ->
1225    /// false`) silently disabled schema-drift, evolution, and exactly-once
1226    /// detection through the entire observability/CLI path even when the real
1227    /// sink supported them.
1228    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        // Schema-drift (#194): the wrapper must surface the inner schema, not the
1268        // `None` default — otherwise drift detection is inert through the pipeline.
1269        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        // Pre-existing capabilities the wrapper must also forward.
1279        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        // Typed delivery capabilities (#292): the pipeline's mechanism
1295        // selection reads these through the wrapper.
1296        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        // Sink that returns 2 Ok + 1 Err.
1362        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        // faucet_sink_records_total should reflect 2 (Ok count), not 3.
1394        // Filter to this test's own labels (connector="mixed") — prior tests in
1395        // the same `mod sink_tests` (e.g. records_writes_and_records_counters
1396        // for connector="mock-sink") leave entries in the shared global
1397        // recorder, and the HashMap-iteration order of `Snapshot::into_vec()`
1398        // is non-deterministic, so a naïve `find_map` returns an arbitrary
1399        // entry.
1400        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    // ── flush error path ─────────────────────────────────────────────────────
1424
1425    #[tokio::test]
1426    #[allow(clippy::await_holding_lock)]
1427    async fn flush_error_increments_errors_total_and_propagates() {
1428        // A sink whose flush() returns Err must surface the error and emit
1429        // faucet_sink_errors_total with the matching kind label.
1430        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    // ── panic isolation on every sink call ───────────────────────────────────
1471
1472    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    /// `InstrumentedSink` must forward the identity + provenance methods.
1540    ///
1541    /// It wraps every sink whenever observability is active — i.e. the default
1542    /// build — and it forwarded neither of these, which review caught. The
1543    /// failure mode is silent: `local_outputs()` falling back to the trait
1544    /// default hides every file the inner sink created from the retention GC
1545    /// (#587), so those files are never reclaimed and nothing logs or errors.
1546    /// `dataset_uri()` falling back records `jsonl://unknown` in lineage and the
1547    /// catalog.
1548    ///
1549    /// The sibling decorators are covered in
1550    /// `tests/local_output_forwarding.rs`; this one lives here because the
1551    /// module is private.
1552    #[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}