1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum RoundtripSide {
33 Source,
35 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#[derive(Debug, Default)]
79pub struct ThrottleTally {
80 throttled: AtomicU64,
81 wait_nanos: AtomicU64,
82}
83
84impl ThrottleTally {
85 pub fn throttled(&self) -> u64 {
87 self.throttled.load(Ordering::Relaxed)
88 }
89
90 pub fn wait(&self) -> Duration {
92 Duration::from_nanos(self.wait_nanos.load(Ordering::Relaxed))
93 }
94}
95
96#[derive(Debug, Clone)]
104pub struct RoundtripRecorder {
105 side: RoundtripSide,
106 base: Vec<Label>,
108 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 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 pub fn throttle_tally(&self) -> Arc<ThrottleTally> {
148 Arc::clone(&self.throttle)
149 }
150
151 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 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 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 pub fn with_meter(mut self, meter: Arc<UsageMeter>) -> Self {
196 self.meter = Some(meter);
197 self
198 }
199
200 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 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 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 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 #[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#[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 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
304pub 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
328pub 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
343pub 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 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 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#[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 pub fn install(&self, recorder: Arc<RoundtripRecorder>) {
749 let _ = self.0.set(recorder);
750 }
751
752 pub fn recorder(&self) -> Option<Arc<RoundtripRecorder>> {
755 self.0.get().cloned()
756 }
757
758 pub fn record(&self, op: &'static str) {
760 if let Some(r) = self.0.get() {
761 r.record(op);
762 }
763 }
764
765 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 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 pub fn throttled(&self) {
781 if let Some(r) = self.0.get() {
782 r.throttled();
783 }
784 }
785
786 pub fn retry(&self, class: RetryClass) {
788 if let Some(r) = self.0.get() {
789 r.retry(class);
790 }
791 }
792
793 pub fn throttle_wait_timer(&self) -> ThrottleWait {
795 ThrottleWait::start(self.recorder())
796 }
797
798 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}