Skip to main content

faucet_core/observability/
roundtrip.rs

1//! Protocol-agnostic upstream round-trip counting (#638).
2//!
3//! faucet counts records, pages, and errors, but until now had no metric for
4//! **how many times a connector actually talked to its backend**. That number
5//! is what drives API-quota consumption, egress cost, database load, and poll
6//! overhead — and `faucet_source_pages_total` only proxies it for paged HTTP
7//! sources, missing job submits, poll loops, and every non-HTTP connector.
8//!
9//! The counter is deliberately not HTTP-shaped: each connector decides what one
10//! round trip means for it and names it with a **closed** `op` label (a SQL
11//! source counts `query`, an object store `list`/`get`, Kafka `poll`, a REST
12//! async job `submit`/`poll`/`fetch`/`page`).
13//!
14//! ## Why a recorder handle rather than a free function
15//!
16//! A connector's I/O sites sit far below the pipeline, which is the only place
17//! that knows the `pipeline` / `row` / `connector` labels every other metric
18//! carries. Handing the connector a pre-labelled handle keeps this metric
19//! consistent with the rest, and — unlike a tokio task-local — the `Arc`
20//! survives `tokio::spawn`, which the S3 and Parquet fan-out paths rely on.
21
22use crate::resilience::RetryClass;
23use crate::usage::{CostSignal, UsageMeter, UsageSide};
24use metrics::{Label, SharedString, counter, histogram};
25use std::sync::atomic::{AtomicU64, Ordering};
26use std::sync::{Arc, OnceLock};
27use std::time::{Duration, Instant};
28
29/// Which side of the pipeline a round trip belongs to. Selects the metric
30/// name, so the counter reads the same way as every other source/sink pair.
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum RoundtripSide {
33    /// `faucet_source_roundtrips_total` / `faucet_source_roundtrip_duration_seconds`.
34    Source,
35    /// `faucet_sink_roundtrips_total` / `faucet_sink_roundtrip_duration_seconds`.
36    Sink,
37}
38
39impl RoundtripSide {
40    const fn counter_name(self) -> &'static str {
41        match self {
42            Self::Source => "faucet_source_roundtrips_total",
43            Self::Sink => "faucet_sink_roundtrips_total",
44        }
45    }
46
47    const fn histogram_name(self) -> &'static str {
48        match self {
49            Self::Source => "faucet_source_roundtrip_duration_seconds",
50            Self::Sink => "faucet_sink_roundtrip_duration_seconds",
51        }
52    }
53
54    const fn throttled_name(self) -> &'static str {
55        match self {
56            Self::Source => "faucet_source_throttled_total",
57            Self::Sink => "faucet_sink_throttled_total",
58        }
59    }
60
61    const fn throttle_wait_name(self) -> &'static str {
62        match self {
63            Self::Source => "faucet_source_throttle_wait_seconds",
64            Self::Sink => "faucet_sink_throttle_wait_seconds",
65        }
66    }
67
68    const fn retries_name(self) -> &'static str {
69        match self {
70            Self::Source => "faucet_source_retries_total",
71            Self::Sink => "faucet_sink_retries_total",
72        }
73    }
74}
75
76/// Running totals of the rate limiting one recorder observed (#734), shared by
77/// every clone so the pipeline can compare the wait against the run's length.
78#[derive(Debug, Default)]
79pub struct ThrottleTally {
80    throttled: AtomicU64,
81    wait_nanos: AtomicU64,
82}
83
84impl ThrottleTally {
85    /// Rate-limit responses received.
86    pub fn throttled(&self) -> u64 {
87        self.throttled.load(Ordering::Relaxed)
88    }
89
90    /// Time actually slept because of them.
91    pub fn wait(&self) -> Duration {
92        Duration::from_nanos(self.wait_nanos.load(Ordering::Relaxed))
93    }
94}
95
96/// A pre-labelled handle a connector uses to count its own backend calls.
97///
98/// Cheap to clone (the label vec is built once at construction and cloned per
99/// call, exactly like the decorators' `base_labels`). Held behind an
100/// `Arc`/`OnceLock` by connectors that opt in; a connector that never calls
101/// [`record`](Self::record) emits nothing at all, so instrumentation can land
102/// connector by connector without any behaviour change in between.
103#[derive(Debug, Clone)]
104pub struct RoundtripRecorder {
105    side: RoundtripSide,
106    /// `pipeline` / `row` / `connector`, resolved once by the pipeline.
107    base: Vec<Label>,
108    /// The run's usage meter (#704), when one is attached: every round trip
109    /// and cost signal is tallied there as well as emitted as a metric.
110    meter: Option<Arc<UsageMeter>>,
111    connector: SharedString,
112    throttle: Arc<ThrottleTally>,
113}
114
115impl RoundtripSide {
116    fn usage_side(self) -> UsageSide {
117        match self {
118            Self::Source => UsageSide::Source,
119            Self::Sink => UsageSide::Sink,
120        }
121    }
122}
123
124impl RoundtripRecorder {
125    /// Build a recorder for one connector instance.
126    pub fn new(
127        side: RoundtripSide,
128        pipeline: impl Into<SharedString>,
129        row: impl Into<SharedString>,
130        connector: impl Into<SharedString>,
131    ) -> Self {
132        let connector: SharedString = connector.into();
133        Self {
134            side,
135            base: vec![
136                Label::new("pipeline", pipeline.into()),
137                Label::new("row", row.into()),
138                Label::new("connector", connector.clone()),
139            ],
140            meter: None,
141            connector,
142            throttle: Arc::new(ThrottleTally::default()),
143        }
144    }
145
146    /// The rate-limit totals this recorder has observed so far.
147    pub fn throttle_tally(&self) -> Arc<ThrottleTally> {
148        Arc::clone(&self.throttle)
149    }
150
151    /// Count one rate-limit response (HTTP 429, a `RateLimited` error, a
152    /// backend's throttling code) — every one received, whether or not it is
153    /// retried. Emits `faucet_source_throttled_total`.
154    pub fn throttled(&self) {
155        counter!(self.side.throttled_name(), self.base.clone()).increment(1);
156        self.throttle.throttled.fetch_add(1, Ordering::Relaxed);
157        if let Some(m) = self.metered_source() {
158            m.add_throttled();
159        }
160    }
161
162    /// Record time actually slept because of a rate limit. Pass the measured
163    /// sleep, never the server's `Retry-After` value; [`ThrottleWait`] measures
164    /// it for you, including a sleep cut short by cancellation.
165    pub fn throttle_wait(&self, slept: Duration) {
166        histogram!(self.side.throttle_wait_name(), self.base.clone()).record(slept.as_secs_f64());
167        self.throttle.wait_nanos.fetch_add(
168            u64::try_from(slept.as_nanos()).unwrap_or(u64::MAX),
169            Ordering::Relaxed,
170        );
171        if let Some(m) = self.metered_source() {
172            m.add_throttle_wait(slept);
173        }
174    }
175
176    /// Count one retry the connector is about to make, by class. Emits
177    /// `faucet_source_retries_total{class}`.
178    pub fn retry(&self, class: RetryClass) {
179        let mut labels = self.base.clone();
180        labels.push(Label::new("class", SharedString::const_str(class.as_str())));
181        counter!(self.side.retries_name(), labels).increment(1);
182        if let Some(m) = self.metered_source() {
183            m.add_source_retry(class.as_str());
184        }
185    }
186
187    fn metered_source(&self) -> Option<&Arc<UsageMeter>> {
188        match self.side {
189            RoundtripSide::Source => self.meter.as_ref(),
190            RoundtripSide::Sink => None,
191        }
192    }
193
194    /// Also tally into a run's usage meter (#704).
195    pub fn with_meter(mut self, meter: Arc<UsageMeter>) -> Self {
196        self.meter = Some(meter);
197        self
198    }
199
200    /// Report a backend-measured usage figure (#704) — BigQuery's bytes
201    /// billed for a job, the payload size of a streaming insert, a
202    /// warehouse's credits. Emitted as
203    /// `faucet_cost_signals_total{pipeline,row,connector,kind,unit}` (the
204    /// quantity rounded to a whole unit) and, when a meter is attached, kept
205    /// verbatim for the run's usage record. `kind` and `unit` are a closed
206    /// set per connector, documented in its README.
207    pub fn signal(&self, kind: &'static str, unit: &'static str, quantity: f64) {
208        let mut labels = self.base.clone();
209        labels.push(Label::new("kind", SharedString::const_str(kind)));
210        labels.push(Label::new("unit", SharedString::const_str(unit)));
211        counter!("faucet_cost_signals_total", labels).increment(quantity.max(0.0).round() as u64);
212        if let Some(m) = &self.meter {
213            m.add_signal(CostSignal {
214                kind: kind.to_string(),
215                unit: unit.to_string(),
216                quantity,
217                side: self.side.usage_side(),
218                connector: self.connector.to_string(),
219            });
220        }
221    }
222
223    /// Count one round trip to the backend.
224    ///
225    /// `op` must come from the connector's own **closed** set — never a URL,
226    /// query, status code, or host, which would make the series unbounded.
227    /// A retried call is a real round trip and must be counted again.
228    pub fn record(&self, op: &'static str) {
229        counter!(self.side.counter_name(), self.labels_for(op)).increment(1);
230        if let Some(m) = &self.meter {
231            m.add_roundtrip(self.side.usage_side(), op);
232        }
233    }
234
235    /// Count records whose incremental replication key was missing or `null`
236    /// (#747). Emits `faucet_source_replication_key_missing_total`.
237    pub fn replication_key_missing(&self, n: u64) {
238        if n > 0 {
239            counter!(
240                "faucet_source_replication_key_missing_total",
241                self.base.clone()
242            )
243            .increment(n);
244        }
245    }
246
247    /// Count one round trip and record how long it took.
248    pub fn record_timed(&self, op: &'static str, elapsed: Duration) {
249        let labels = self.labels_for(op);
250        counter!(self.side.counter_name(), labels.clone()).increment(1);
251        histogram!(self.side.histogram_name(), labels).record(elapsed.as_secs_f64());
252        if let Some(m) = &self.meter {
253            m.add_roundtrip(self.side.usage_side(), op);
254        }
255    }
256
257    fn labels_for(&self, op: &'static str) -> Vec<Label> {
258        let mut labels = self.base.clone();
259        labels.push(Label::new("op", SharedString::const_str(op)));
260        labels
261    }
262
263    /// The label set this recorder emits for `op` — exposed so a connector's
264    /// tests can assert what they would produce without a metrics recorder
265    /// installed.
266    #[doc(hidden)]
267    pub fn labels_for_test(&self, op: &'static str) -> Vec<(String, String)> {
268        self.labels_for(op)
269            .into_iter()
270            .map(|l| (l.key().to_string(), l.value().to_string()))
271            .collect()
272    }
273}
274
275/// Measures one rate-limit sleep (#734): hold it across the sleep and the
276/// elapsed time is recorded when it drops — after the sleep completes, when a
277/// cancellation branch wins a `select!`, or when the future is dropped mid-sleep
278/// — so the recorded wait is always what was actually slept.
279#[derive(Debug)]
280#[must_use = "the wait is recorded when the guard drops"]
281pub struct ThrottleWait {
282    recorder: Option<Arc<RoundtripRecorder>>,
283    start: Instant,
284}
285
286impl ThrottleWait {
287    /// Start timing a sleep on behalf of `recorder` (a no-op guard for `None`).
288    pub fn start(recorder: Option<Arc<RoundtripRecorder>>) -> Self {
289        Self {
290            recorder,
291            start: Instant::now(),
292        }
293    }
294}
295
296impl Drop for ThrottleWait {
297    fn drop(&mut self) {
298        if let Some(r) = &self.recorder {
299            r.throttle_wait(self.start.elapsed());
300        }
301    }
302}
303
304/// Sleep `wait` on behalf of a rate limit, recording the time actually slept
305/// (#734). Returns `false` when `cancel` fired first; the partial wait is still
306/// recorded.
307pub async fn throttle_sleep(
308    recorder: Option<Arc<RoundtripRecorder>>,
309    wait: Duration,
310    cancel: Option<&tokio_util::sync::CancellationToken>,
311) -> bool {
312    let _timer = ThrottleWait::start(recorder);
313    match cancel {
314        Some(token) => {
315            tokio::select! {
316                biased;
317                _ = token.cancelled() => false,
318                _ = tokio::time::sleep(wait) => true,
319            }
320        }
321        None => {
322            tokio::time::sleep(wait).await;
323            true
324        }
325    }
326}
327
328/// The one-line warning a run logs when rate-limit waits took more than a
329/// tenth of it (#734), or `None` when they did not.
330pub fn throttle_warning(throttled: u64, wait: Duration, run: Duration) -> Option<String> {
331    if wait.is_zero() || run.is_zero() || wait.as_secs_f64() * 10.0 <= run.as_secs_f64() {
332        return None;
333    }
334    let pct = (wait.as_secs_f64() / run.as_secs_f64() * 100.0).min(100.0);
335    Some(format!(
336        "source spent {:.1}s of a {:.1}s run ({pct:.0}%) waiting on rate limits \
337         ({throttled} throttled responses); lower concurrency, stagger schedules or raise the quota",
338        wait.as_secs_f64(),
339        run.as_secs_f64(),
340    ))
341}
342
343/// Register descriptions for both sides' counters and histograms. Called once
344/// by `install_observability`.
345pub fn describe_roundtrip_metrics() {
346    metrics::describe_counter!(
347        "faucet_cost_signals_total",
348        "Backend-reported usage a connector measured during a run (BigQuery bytes billed, streamed payload bytes, …), by kind and unit"
349    );
350    metrics::describe_counter!(
351        "faucet_source_roundtrips_total",
352        "Calls a source made to its upstream backend, by connector-defined op"
353    );
354    metrics::describe_counter!(
355        "faucet_sink_roundtrips_total",
356        "Calls a sink made to its upstream backend, by connector-defined op"
357    );
358    metrics::describe_histogram!(
359        "faucet_source_roundtrip_duration_seconds",
360        metrics::Unit::Seconds,
361        "Duration of one source round trip to its upstream backend"
362    );
363    metrics::describe_histogram!(
364        "faucet_sink_roundtrip_duration_seconds",
365        metrics::Unit::Seconds,
366        "Duration of one sink round trip to its upstream backend"
367    );
368    metrics::describe_counter!(
369        "faucet_source_throttled_total",
370        "Rate-limit responses (HTTP 429 and equivalents) a source received"
371    );
372    metrics::describe_histogram!(
373        "faucet_source_throttle_wait_seconds",
374        metrics::Unit::Seconds,
375        "Time a source actually slept because of one rate-limit response"
376    );
377    metrics::describe_counter!(
378        "faucet_source_replication_key_missing_total",
379        "Records an incremental source received without its replication key (kept, dropped or failed per on_missing_key)"
380    );
381    metrics::describe_counter!(
382        "faucet_source_retries_total",
383        "Retries a source made against its upstream backend, by retry class"
384    );
385}
386
387#[cfg(test)]
388mod tests {
389    use super::*;
390
391    #[test]
392    fn metered_recorders_feed_the_usage_meter_through_a_slot() {
393        let meter = Arc::new(crate::usage::UsageMeter::new());
394        for side in [RoundtripSide::Source, RoundtripSide::Sink] {
395            let slot = RecorderSlot::new();
396            slot.record("get");
397            slot.record_timed("get", Duration::from_millis(1));
398            slot.signal("bytes_billed", "bytes", 10.0);
399            assert!(slot.recorder().is_none());
400            slot.install(Arc::new(
401                RoundtripRecorder::new(side, "p", "r", "c").with_meter(meter.clone()),
402            ));
403            slot.record("get");
404            slot.record_timed("put", Duration::from_millis(2));
405            slot.signal("bytes_billed", "bytes", 1024.4);
406            assert!(slot.recorder().is_some());
407        }
408        let snap = meter.snapshot();
409        assert_eq!(snap.source_roundtrips["get"], 1);
410        assert_eq!(snap.source_roundtrips["put"], 1);
411        assert_eq!(snap.sink_roundtrips["get"], 1);
412        assert_eq!(snap.signals.len(), 2);
413        assert_eq!(snap.signals[0].kind, "bytes_billed");
414        assert_eq!(snap.signals[0].connector, "c");
415        assert_eq!(snap.signals[1].side, crate::usage::UsageSide::Sink);
416    }
417
418    #[test]
419    fn side_selects_the_metric_names() {
420        assert_eq!(
421            RoundtripSide::Source.counter_name(),
422            "faucet_source_roundtrips_total"
423        );
424        assert_eq!(
425            RoundtripSide::Sink.counter_name(),
426            "faucet_sink_roundtrips_total"
427        );
428        assert_eq!(
429            RoundtripSide::Source.histogram_name(),
430            "faucet_source_roundtrip_duration_seconds"
431        );
432        assert_eq!(
433            RoundtripSide::Sink.histogram_name(),
434            "faucet_sink_roundtrip_duration_seconds"
435        );
436    }
437
438    #[test]
439    fn labels_carry_the_universal_trio_plus_op() {
440        let r = RoundtripRecorder::new(RoundtripSide::Source, "p", "rowA", "rest");
441        let labels = r.labels_for_test("poll");
442        assert_eq!(
443            labels,
444            vec![
445                ("pipeline".to_string(), "p".to_string()),
446                ("row".to_string(), "rowA".to_string()),
447                ("connector".to_string(), "rest".to_string()),
448                ("op".to_string(), "poll".to_string()),
449            ],
450            "the trio must match every other metric, or this one can't be joined to them"
451        );
452    }
453
454    #[test]
455    fn op_is_the_only_thing_that_varies_between_calls() {
456        // Guards against a future refactor that rebuilds the base labels
457        // per call and lets them drift.
458        let r = RoundtripRecorder::new(RoundtripSide::Sink, "p", "", "s3");
459        let a = r.labels_for_test("put");
460        let b = r.labels_for_test("list");
461        assert_eq!(a[..3], b[..3]);
462        assert_eq!(a[3].1, "put");
463        assert_eq!(b[3].1, "list");
464    }
465
466    #[test]
467    fn record_emits_the_counter_under_an_installed_recorder() {
468        use crate::observability::decorator::source_tests::{LOCK, snapshotter};
469        use metrics_util::debugging::DebugValue;
470
471        let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
472        let snap = snapshotter();
473        let r = RoundtripRecorder::new(RoundtripSide::Source, "pipe", "rowA", "rest");
474        r.record("submit");
475        r.record("poll");
476        r.record("poll");
477
478        let counts: Vec<(String, u64)> = snap
479            .snapshot()
480            .into_vec()
481            .into_iter()
482            .filter(|(k, _, _, _)| k.key().name() == "faucet_source_roundtrips_total")
483            .filter_map(|(k, _, _, v)| {
484                let op = k
485                    .key()
486                    .labels()
487                    .find(|l| l.key() == "op")?
488                    .value()
489                    .to_string();
490                match v {
491                    DebugValue::Counter(c) => Some((op, c)),
492                    _ => None,
493                }
494            })
495            .collect();
496
497        let poll = counts.iter().find(|(op, _)| op == "poll").map(|(_, c)| *c);
498        let submit = counts
499            .iter()
500            .find(|(op, _)| op == "submit")
501            .map(|(_, c)| *c);
502        assert_eq!(submit, Some(1), "one submit: {counts:?}");
503        assert_eq!(
504            poll,
505            Some(2),
506            "each poll counts — the poll loop's overhead is the whole signal: {counts:?}"
507        );
508    }
509
510    #[test]
511    fn replication_key_missing_counts_through_the_slot() {
512        use crate::observability::decorator::source_tests::{LOCK, snapshotter};
513        use metrics_util::debugging::DebugValue;
514
515        let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
516        let snap = snapshotter();
517        let slot = RecorderSlot::new();
518        slot.replication_key_missing(4);
519        slot.install(Arc::new(RoundtripRecorder::new(
520            RoundtripSide::Source,
521            "pipe",
522            "rowK",
523            "rest",
524        )));
525        slot.replication_key_missing(0);
526        slot.replication_key_missing(3);
527        let total: u64 = snap
528            .snapshot()
529            .into_vec()
530            .into_iter()
531            .filter(|(k, _, _, _)| {
532                k.key().name() == "faucet_source_replication_key_missing_total"
533                    && k.key().labels().any(|l| l.value() == "rowK")
534            })
535            .map(|(_, _, _, v)| match v {
536                DebugValue::Counter(c) => c,
537                _ => 0,
538            })
539            .sum();
540        assert_eq!(total, 3);
541    }
542
543    #[test]
544    fn recording_without_an_installed_recorder_is_a_no_op() {
545        // The metrics facade discards into a no-op recorder when none is
546        // installed, so a connector can always call these — there is no
547        // "is observability on?" check to get wrong.
548        let r = RoundtripRecorder::new(RoundtripSide::Source, "p", "r", "postgres");
549        r.record("query");
550        r.record_timed("query", Duration::from_millis(5));
551    }
552
553    #[test]
554    fn throttling_feeds_the_tally_the_meter_and_the_metrics() {
555        use crate::observability::decorator::source_tests::{LOCK, snapshotter};
556        use metrics_util::debugging::DebugValue;
557
558        let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
559        let snap = snapshotter();
560        let meter = Arc::new(UsageMeter::new());
561        let r = RoundtripRecorder::new(RoundtripSide::Source, "pipe", "rowT", "rest")
562            .with_meter(meter.clone());
563        let tally = r.throttle_tally();
564        r.throttled();
565        r.throttled();
566        r.throttle_wait(Duration::from_millis(300));
567        r.retry(RetryClass::RateLimited);
568        r.retry(RetryClass::Http5xx);
569        assert_eq!(tally.throttled(), 2);
570        assert_eq!(tally.wait(), Duration::from_millis(300));
571        let usage = meter.snapshot();
572        assert_eq!(usage.throttled, 2);
573        assert!((usage.throttle_wait_secs - 0.3).abs() < 1e-9);
574        assert_eq!(usage.source_retries["rate_limited"], 1);
575        assert_eq!(usage.source_retries["http_5xx"], 1);
576
577        let entries: Vec<_> = snap
578            .snapshot()
579            .into_vec()
580            .into_iter()
581            .filter(|(k, _, _, _)| k.key().labels().any(|l| l.value() == "rowT"))
582            .collect();
583        let counter = |name: &str, class: Option<&str>| {
584            entries.iter().find_map(|(k, _, _, v)| {
585                let class_ok = class.is_none_or(|c| {
586                    k.key()
587                        .labels()
588                        .any(|l| l.key() == "class" && l.value() == c)
589                });
590                match v {
591                    DebugValue::Counter(n) if k.key().name() == name && class_ok => Some(*n),
592                    _ => None,
593                }
594            })
595        };
596        assert_eq!(counter("faucet_source_throttled_total", None), Some(2));
597        assert_eq!(
598            counter("faucet_source_retries_total", Some("rate_limited")),
599            Some(1)
600        );
601        assert_eq!(
602            counter("faucet_source_retries_total", Some("http_5xx")),
603            Some(1)
604        );
605        assert!(entries.iter().any(|(k, _, _, v)| {
606            k.key().name() == "faucet_source_throttle_wait_seconds"
607                && matches!(v, DebugValue::Histogram(h) if h.len() == 1)
608        }));
609    }
610
611    #[test]
612    fn sink_side_throttling_emits_metrics_but_stays_off_the_usage_record() {
613        let meter = Arc::new(UsageMeter::new());
614        let r =
615            RoundtripRecorder::new(RoundtripSide::Sink, "p", "r", "http").with_meter(meter.clone());
616        r.throttled();
617        r.throttle_wait(Duration::from_millis(5));
618        r.retry(RetryClass::Timeout);
619        assert_eq!(r.throttle_tally().throttled(), 1);
620        let usage = meter.snapshot();
621        assert_eq!(usage.throttled, 0);
622        assert!(usage.source_retries.is_empty());
623        assert_eq!(
624            RoundtripSide::Sink.throttled_name(),
625            "faucet_sink_throttled_total"
626        );
627        assert_eq!(
628            RoundtripSide::Sink.throttle_wait_name(),
629            "faucet_sink_throttle_wait_seconds"
630        );
631        assert_eq!(
632            RoundtripSide::Sink.retries_name(),
633            "faucet_sink_retries_total"
634        );
635    }
636
637    #[tokio::test]
638    async fn throttle_sleep_records_the_time_actually_slept() {
639        let r = Arc::new(RoundtripRecorder::new(
640            RoundtripSide::Source,
641            "p",
642            "r",
643            "rest",
644        ));
645        let tally = r.throttle_tally();
646        assert!(throttle_sleep(Some(r.clone()), Duration::from_millis(40), None).await);
647        let full = tally.wait();
648        assert!(full >= Duration::from_millis(40), "{full:?}");
649
650        let token = tokio_util::sync::CancellationToken::new();
651        let t = token.clone();
652        tokio::spawn(async move {
653            tokio::time::sleep(Duration::from_millis(30)).await;
654            t.cancel();
655        });
656        let completed =
657            throttle_sleep(Some(r.clone()), Duration::from_secs(30), Some(&token)).await;
658        assert!(!completed, "cancellation wins");
659        let partial = tally.wait() - full;
660        assert!(
661            partial >= Duration::from_millis(25) && partial < Duration::from_secs(5),
662            "the partial wait is recorded, not the requested 30 s: {partial:?}"
663        );
664
665        let token = tokio_util::sync::CancellationToken::new();
666        assert!(throttle_sleep(None, Duration::from_millis(1), Some(&token)).await);
667    }
668
669    #[tokio::test]
670    async fn a_dropped_sleep_still_records_its_partial_wait() {
671        let slot = RecorderSlot::new();
672        slot.throttled();
673        slot.retry(RetryClass::Connection);
674        drop(slot.throttle_wait_timer());
675        let r = Arc::new(RoundtripRecorder::new(
676            RoundtripSide::Source,
677            "p",
678            "r",
679            "rest",
680        ));
681        slot.install(r.clone());
682        slot.throttled();
683        slot.retry(RetryClass::Connection);
684        let tally = r.throttle_tally();
685        let sleeping = async {
686            let _t = slot.throttle_wait_timer();
687            tokio::time::sleep(Duration::from_secs(30)).await;
688        };
689        let _ = tokio::time::timeout(Duration::from_millis(30), sleeping).await;
690        assert_eq!(tally.throttled(), 1);
691        assert!(
692            tally.wait() >= Duration::from_millis(25),
693            "{:?}",
694            tally.wait()
695        );
696        assert!(tally.wait() < Duration::from_secs(5));
697    }
698
699    #[test]
700    fn warns_only_when_waiting_exceeds_a_tenth_of_the_run() {
701        assert_eq!(
702            throttle_warning(0, Duration::ZERO, Duration::from_secs(10)),
703            None
704        );
705        assert_eq!(
706            throttle_warning(3, Duration::from_secs(1), Duration::from_secs(10)),
707            None,
708            "exactly 10% does not warn"
709        );
710        assert_eq!(
711            throttle_warning(3, Duration::from_secs(1), Duration::ZERO),
712            None
713        );
714        let msg = throttle_warning(312, Duration::from_secs(2460), Duration::from_secs(3600))
715            .expect("68% warns");
716        assert!(msg.contains("2460.0s of a 3600.0s run (68%)"), "{msg}");
717        assert!(msg.contains("312 throttled responses"), "{msg}");
718        let capped = throttle_warning(1, Duration::from_secs(20), Duration::from_secs(10)).unwrap();
719        assert!(capped.contains("(100%)"), "{capped}");
720    }
721
722    #[test]
723    fn clone_shares_the_prebuilt_labels() {
724        let r = RoundtripRecorder::new(RoundtripSide::Source, "p", "r", "kafka");
725        let c = r.clone();
726        assert_eq!(r.labels_for_test("poll"), c.labels_for_test("poll"));
727    }
728}
729
730/// A connector's slot for the recorder the pipeline installs (#638 / #704).
731///
732/// Connectors keep one of these in their struct and forward
733/// [`set_roundtrip_recorder`](crate::Source::set_roundtrip_recorder) to
734/// [`install`](Self::install); every call site then does
735/// `self.roundtrips.record("get")` without checking whether a pipeline
736/// installed anything. First install wins — the pipeline installs exactly
737/// once per run, and a re-used connector instance keeps the labels it is
738/// already counting under.
739#[derive(Debug, Default)]
740pub struct RecorderSlot(OnceLock<Arc<RoundtripRecorder>>);
741
742impl RecorderSlot {
743    pub const fn new() -> Self {
744        Self(OnceLock::new())
745    }
746
747    /// Install the pipeline's recorder (no-op when one is already installed).
748    pub fn install(&self, recorder: Arc<RoundtripRecorder>) {
749        let _ = self.0.set(recorder);
750    }
751
752    /// The installed recorder, for handing to a helper that performs I/O on
753    /// the connector's behalf.
754    pub fn recorder(&self) -> Option<Arc<RoundtripRecorder>> {
755        self.0.get().cloned()
756    }
757
758    /// Count one round trip when a recorder is installed.
759    pub fn record(&self, op: &'static str) {
760        if let Some(r) = self.0.get() {
761            r.record(op);
762        }
763    }
764
765    /// Count one timed round trip when a recorder is installed.
766    pub fn record_timed(&self, op: &'static str, elapsed: Duration) {
767        if let Some(r) = self.0.get() {
768            r.record_timed(op, elapsed);
769        }
770    }
771
772    /// Count records missing their replication key when a recorder is installed.
773    pub fn replication_key_missing(&self, n: u64) {
774        if let Some(r) = self.0.get() {
775            r.replication_key_missing(n);
776        }
777    }
778
779    /// Count one rate-limit response when a recorder is installed.
780    pub fn throttled(&self) {
781        if let Some(r) = self.0.get() {
782            r.throttled();
783        }
784    }
785
786    /// Count one retry of `class` when a recorder is installed.
787    pub fn retry(&self, class: RetryClass) {
788        if let Some(r) = self.0.get() {
789            r.retry(class);
790        }
791    }
792
793    /// Start timing a rate-limit sleep ([`ThrottleWait`]).
794    pub fn throttle_wait_timer(&self) -> ThrottleWait {
795        ThrottleWait::start(self.recorder())
796    }
797
798    /// Report a backend-measured usage figure when a recorder is installed.
799    pub fn signal(&self, kind: &'static str, unit: &'static str, quantity: f64) {
800        if let Some(r) = self.0.get() {
801            r.signal(kind, unit, quantity);
802        }
803    }
804}