1use std::future::Future;
50use std::pin::Pin;
51use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
52use std::sync::{Arc, OnceLock};
53use std::time::{Duration, Instant};
54
55use async_trait::async_trait;
56use bytes::Bytes;
57use futures_util::future::{Either, select};
58use mz_dyncfg::{Config, ConfigSet, ParameterScope};
59use mz_ore::bytes::SegmentedBytes;
60use mz_ore::cast::CastLossy;
61use mz_ore::task::AbortOnDropHandle;
62use tracing::{debug, info, warn};
63
64use crate::location::{BLOB_GET_LIVENESS_KEY, Blob, BlobMetadata, ExternalError};
65use crate::metrics::BlobHedgeMetrics;
66
67pub(crate) const BLOB_HEDGED_GET_ENABLED: Config<bool> = Config::new(
68 "persist_blob_hedged_get_enabled",
69 true,
70 "Whether to hedge slow blob gets with a second request on a separate \
71 connection pool (Materialize).",
72 ParameterScope::Environment,
73);
74
75pub(crate) const BLOB_HEDGED_GET_DELAY: Config<Duration> = Config::new(
76 "persist_blob_hedged_get_delay",
77 Duration::from_secs(2),
78 "How long a blob get may be in flight before a hedge request is fired \
79 (Materialize).",
80 ParameterScope::Environment,
81);
82
83pub(crate) const BLOB_HEDGED_GET_MAX_CONCURRENT: Config<usize> = Config::new(
84 "persist_blob_hedged_get_max_concurrent",
85 2,
89 "Maximum concurrent hedge requests per blob handle, bounding the memory \
90 held by raced gets and the number of warm sockets (Materialize).",
91 ParameterScope::Environment,
92);
93
94pub(crate) const BLOB_HEDGED_GET_BUDGET_RATIO: Config<f64> = Config::new(
95 "persist_blob_hedged_get_budget_ratio",
96 0.01,
97 "Long-run bound on hedge requests as a fraction of blob gets \
98 (Materialize).",
99 ParameterScope::Environment,
100);
101
102pub(crate) const BLOB_HEDGED_GET_WARM_INTERVAL: Config<Duration> = Config::new(
107 "persist_blob_hedged_get_warm_interval",
108 Duration::from_secs(20),
109 "How often to issue liveness gets that keep the hedge connection pool \
110 warm, 0 disables warming without disabling hedging (Materialize).",
111 ParameterScope::Environment,
112);
113
114const HEDGE_COST_MICRO_TOKENS: u64 = 1_000_000;
118
119const BUDGET_BURST_MICRO_TOKENS: u64 = 32 * HEDGE_COST_MICRO_TOKENS;
127
128const SIBLING_OPEN_RETRY_INITIAL_BACKOFF: Duration = Duration::from_secs(1);
132const SIBLING_OPEN_RETRY_MAX_BACKOFF: Duration = Duration::from_secs(60);
133
134enum HedgeRefused {
136 Concurrency,
137 Budget,
138}
139
140#[derive(Debug)]
146struct HedgeBudget {
147 concurrent: AtomicUsize,
148 micro_tokens: AtomicU64,
149}
150
151impl HedgeBudget {
152 fn new() -> Self {
153 HedgeBudget {
154 concurrent: AtomicUsize::new(0),
155 micro_tokens: AtomicU64::new(BUDGET_BURST_MICRO_TOKENS),
156 }
157 }
158
159 fn try_acquire(&self, max_concurrent: usize) -> Result<HedgeGuard<'_>, HedgeRefused> {
163 let got_slot = self
164 .concurrent
165 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |c| {
166 (c < max_concurrent).then_some(c + 1)
167 })
168 .is_ok();
169 if !got_slot {
170 return Err(HedgeRefused::Concurrency);
171 }
172 let guard = HedgeGuard(self);
175 let took_token = self
176 .micro_tokens
177 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |t| {
178 t.checked_sub(HEDGE_COST_MICRO_TOKENS)
179 })
180 .is_ok();
181 if took_token {
182 Ok(guard)
183 } else {
184 Err(HedgeRefused::Budget)
185 }
186 }
187
188 fn replenish(&self, ratio: f64) {
191 let add = u64::cast_lossy(ratio.clamp(0.0, 1.0) * f64::cast_lossy(HEDGE_COST_MICRO_TOKENS));
192 if add == 0 {
193 return;
194 }
195 let _ = self
196 .micro_tokens
197 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |t| {
198 Some((t + add).min(BUDGET_BURST_MICRO_TOKENS))
199 });
200 }
201}
202
203struct HedgeGuard<'a>(&'a HedgeBudget);
204
205impl Drop for HedgeGuard<'_> {
206 fn drop(&mut self) {
207 self.0.concurrent.fetch_sub(1, Ordering::Relaxed);
208 }
209}
210
211#[derive(Debug)]
214pub enum HedgeSibling {
215 Isolated(Arc<dyn Blob>),
218 SharedWithPrimary,
222 Unavailable,
225}
226
227pub type HedgeSiblingOpener =
232 Box<dyn Fn() -> Pin<Box<dyn Future<Output = HedgeSibling> + Send>> + Send>;
233
234#[derive(Debug)]
236pub struct HedgedBlob {
237 primary: Arc<dyn Blob>,
238 hedge: Arc<OnceLock<Arc<dyn Blob>>>,
242 cfg: Arc<ConfigSet>,
243 metrics: BlobHedgeMetrics,
244 budget: HedgeBudget,
245 _sibling_task: Option<AbortOnDropHandle<()>>,
248}
249
250fn recheck_interval(cfg: &ConfigSet) -> Duration {
254 match BLOB_HEDGED_GET_WARM_INTERVAL.get(cfg) {
255 Duration::ZERO => *BLOB_HEDGED_GET_WARM_INTERVAL.default(),
256 interval => interval,
257 }
258}
259
260async fn warm(hedge: Arc<dyn Blob>, cfg: Arc<ConfigSet>, metrics: BlobHedgeMetrics) {
269 loop {
270 let interval = BLOB_HEDGED_GET_WARM_INTERVAL.get(&cfg);
271 if !BLOB_HEDGED_GET_ENABLED.get(&cfg) || interval == Duration::ZERO {
272 tokio::time::sleep(recheck_interval(&cfg)).await;
274 continue;
275 }
276 let start = Instant::now();
281 let sockets = BLOB_HEDGED_GET_MAX_CONCURRENT.get(&cfg);
282 let pings = (0..sockets).map(|_| hedge.get(BLOB_GET_LIVENESS_KEY));
283 let pings = futures_util::future::join_all(pings);
284 match tokio::time::timeout(interval, pings).await {
289 Ok(results) if results.iter().all(|r| r.is_ok()) => {
290 metrics.rtt_latency.set(start.elapsed().as_secs_f64());
291 }
292 Ok(_) | Err(_) => {
293 metrics.warm_errors.inc();
298 }
299 }
300 tokio::time::sleep(interval).await;
301 }
302}
303
304async fn arm_then_warm(
307 primary: Arc<dyn Blob>,
308 opener: HedgeSiblingOpener,
309 slot: Arc<OnceLock<Arc<dyn Blob>>>,
310 cfg: Arc<ConfigSet>,
311 metrics: BlobHedgeMetrics,
312) {
313 let mut backoff = SIBLING_OPEN_RETRY_INITIAL_BACKOFF;
314 let mut attempts = 0u64;
315 let hedge = loop {
316 attempts += 1;
317 match opener().await {
318 HedgeSibling::Isolated(hedge) => break hedge,
319 HedgeSibling::SharedWithPrimary => {
320 install(&slot, &metrics, primary);
321 return;
322 }
323 HedgeSibling::Unavailable => {
325 tokio::time::sleep(backoff).await;
326 backoff = (backoff * 2).min(SIBLING_OPEN_RETRY_MAX_BACKOFF);
327 while !BLOB_HEDGED_GET_ENABLED.get(&cfg) {
332 tokio::time::sleep(recheck_interval(&cfg)).await;
333 }
334 }
335 }
336 };
337 if attempts > 1 {
338 info!(
339 attempts,
340 "hedged blob gets armed after retrying the sibling open"
341 );
342 }
343 install(&slot, &metrics, Arc::clone(&hedge));
344 warm(hedge, cfg, metrics).await
345}
346
347fn install(slot: &OnceLock<Arc<dyn Blob>>, metrics: &BlobHedgeMetrics, hedge: Arc<dyn Blob>) {
348 if slot.set(hedge).is_err() {
349 mz_ore::soft_panic_or_log!("hedge sibling installed twice");
350 }
351 metrics.armed.set(1);
352}
353
354impl HedgedBlob {
355 #[cfg(test)]
361 pub(crate) fn new(
362 primary: Arc<dyn Blob>,
363 sibling: HedgeSibling,
364 cfg: Arc<ConfigSet>,
365 metrics: BlobHedgeMetrics,
366 ) -> HedgedBlob {
367 let slot = Arc::new(OnceLock::new());
368 let task = match sibling {
369 HedgeSibling::Isolated(h) => {
370 install(&slot, &metrics, Arc::clone(&h));
371 Some(
372 mz_ore::task::spawn(
373 || "persist::blob_hedge_warmer",
374 warm(h, Arc::clone(&cfg), metrics.clone()),
375 )
376 .abort_on_drop(),
377 )
378 }
379 HedgeSibling::SharedWithPrimary => {
380 install(&slot, &metrics, Arc::clone(&primary));
381 None
382 }
383 HedgeSibling::Unavailable => None,
384 };
385 HedgedBlob {
386 primary,
387 hedge: slot,
388 cfg,
389 metrics,
390 budget: HedgeBudget::new(),
391 _sibling_task: task,
392 }
393 }
394
395 pub fn new_arming(
409 primary: Arc<dyn Blob>,
410 opener: HedgeSiblingOpener,
411 cfg: Arc<ConfigSet>,
412 metrics: BlobHedgeMetrics,
413 ) -> HedgedBlob {
414 let slot: Arc<OnceLock<Arc<dyn Blob>>> = Arc::new(OnceLock::new());
415 let task = mz_ore::task::spawn(
416 || "persist::blob_hedge_sibling",
417 arm_then_warm(
418 Arc::clone(&primary),
419 opener,
420 Arc::clone(&slot),
421 Arc::clone(&cfg),
422 metrics.clone(),
423 ),
424 )
425 .abort_on_drop();
426 HedgedBlob {
427 primary,
428 hedge: slot,
429 cfg,
430 metrics,
431 budget: HedgeBudget::new(),
432 _sibling_task: Some(task),
433 }
434 }
435
436 fn admit(&self) -> Option<(&Arc<dyn Blob>, HedgeGuard<'_>)> {
439 let Some(hedge_blob) = self.hedge.get() else {
440 self.metrics.skipped_unavailable.inc();
441 return None;
442 };
443 match self
444 .budget
445 .try_acquire(BLOB_HEDGED_GET_MAX_CONCURRENT.get(&self.cfg))
446 {
447 Ok(guard) => Some((hedge_blob, guard)),
448 Err(HedgeRefused::Concurrency) => {
449 self.metrics.skipped_concurrency.inc();
450 None
451 }
452 Err(HedgeRefused::Budget) => {
453 self.metrics.skipped_budget.inc();
454 None
455 }
456 }
457 }
458
459 fn record_win(&self, key: &str, start: Instant) {
460 self.metrics.won.inc();
461 self.metrics
462 .won_seconds
463 .observe(start.elapsed().as_secs_f64());
464 debug!(%key, elapsed = ?start.elapsed(), "blob get won by hedge request");
465 }
466
467 async fn get_hedged(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
468 let start = Instant::now();
469 let delay = BLOB_HEDGED_GET_DELAY.get(&self.cfg);
470 let mut primary = std::pin::pin!(self.primary.get(key));
471 if let Ok(res) = tokio::time::timeout(delay, primary.as_mut()).await {
478 return res;
481 }
482 let Some((hedge_blob, guard)) = self.admit() else {
483 return primary.await;
484 };
485 self.metrics.fired.inc();
486 let mut hedge = std::pin::pin!(hedge_blob.get(key));
487 match select(primary.as_mut(), hedge.as_mut()).await {
499 Either::Left((Ok(res), _hedge)) => Ok(res),
500 Either::Right((Ok(res), _primary)) => {
501 self.record_win(key, start);
502 Ok(res)
503 }
504 Either::Left((Err(primary_err), hedge)) => {
505 match tokio::time::timeout(delay, hedge).await {
522 Ok(Ok(res)) => {
523 self.record_win(key, start);
524 Ok(res)
525 }
526 Ok(Err(hedge_err)) => {
527 self.metrics.errors.inc();
528 warn!(%key, %hedge_err, "hedged blob get: both requests failed");
529 Err(primary_err)
537 }
538 Err(_elapsed) => {
539 warn!(%key, "hedged blob get: primary failed, hedge still pending");
540 Err(primary_err)
541 }
542 }
543 }
544 Either::Right((Err(hedge_err), primary)) => {
545 self.metrics.errors.inc();
546 warn!(%key, %hedge_err, "hedge request failed, awaiting primary");
547 drop(guard);
553 primary.await
554 }
555 }
556 }
557}
558
559#[async_trait]
560impl Blob for HedgedBlob {
561 async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
562 if !BLOB_HEDGED_GET_ENABLED.get(&self.cfg) {
563 return self.primary.get(key).await;
564 }
565 let res = self.get_hedged(key).await;
566 self.budget
567 .replenish(BLOB_HEDGED_GET_BUDGET_RATIO.get(&self.cfg));
568 res
569 }
570
571 async fn list_keys_and_metadata(
572 &self,
573 key_prefix: &str,
574 f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
575 ) -> Result<(), ExternalError> {
576 self.primary.list_keys_and_metadata(key_prefix, f).await
577 }
578
579 async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError> {
580 self.primary.set(key, value).await
581 }
582
583 async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError> {
584 self.primary.delete(key).await
585 }
586
587 async fn restore(&self, key: &str) -> Result<(), ExternalError> {
588 self.primary.restore(key).await
589 }
590}
591
592#[cfg(test)]
593mod tests {
594 use std::sync::atomic::{AtomicUsize, Ordering};
595
596 use anyhow::anyhow;
597 use mz_dyncfg::ConfigUpdates;
598 use mz_ore::metrics::MetricsRegistry;
599
600 use crate::location::tests::blob_impl_test;
601 use crate::mem::MemMultiRegistry;
602
603 use super::*;
604
605 #[derive(Debug)]
608 struct TestBlob {
609 delay: Duration,
610 outcome: Result<Option<&'static str>, &'static str>,
611 gets: AtomicUsize,
612 }
613
614 impl TestBlob {
615 fn new(
616 delay: Duration,
617 outcome: Result<Option<&'static str>, &'static str>,
618 ) -> Arc<TestBlob> {
619 Arc::new(TestBlob {
620 delay,
621 outcome,
622 gets: AtomicUsize::new(0),
623 })
624 }
625 }
626
627 #[async_trait]
628 impl Blob for TestBlob {
629 async fn get(&self, _key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
630 self.gets.fetch_add(1, Ordering::SeqCst);
631 tokio::time::sleep(self.delay).await;
632 match self.outcome {
633 Ok(x) => Ok(x.map(|x| SegmentedBytes::from(Bytes::from(x)))),
634 Err(msg) => Err(ExternalError::from(anyhow!(msg))),
635 }
636 }
637
638 async fn list_keys_and_metadata(
639 &self,
640 _key_prefix: &str,
641 _f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
642 ) -> Result<(), ExternalError> {
643 unreachable!("test blob only supports get")
644 }
645
646 async fn set(&self, _key: &str, _value: Bytes) -> Result<(), ExternalError> {
647 unreachable!("test blob only supports get")
648 }
649
650 async fn delete(&self, _key: &str) -> Result<Option<usize>, ExternalError> {
651 unreachable!("test blob only supports get")
652 }
653
654 async fn restore(&self, _key: &str) -> Result<(), ExternalError> {
655 unreachable!("test blob only supports get")
656 }
657 }
658
659 fn test_cfg(customize: impl FnOnce(&mut ConfigUpdates)) -> Arc<ConfigSet> {
660 let cfg = crate::cfg::all_dyn_configs(ConfigSet::default());
661 let mut updates = ConfigUpdates::default();
662 updates.add(&BLOB_HEDGED_GET_ENABLED, true);
663 customize(&mut updates);
664 updates.apply(&cfg);
665 Arc::new(cfg)
666 }
667
668 fn metrics() -> BlobHedgeMetrics {
669 BlobHedgeMetrics::new(&MetricsRegistry::new())
670 }
671
672 fn hedged(primary: &Arc<TestBlob>, hedge: &Arc<TestBlob>, cfg: Arc<ConfigSet>) -> HedgedBlob {
673 let primary: Arc<dyn Blob> = Arc::<TestBlob>::clone(primary);
674 let hedge: Arc<dyn Blob> = Arc::<TestBlob>::clone(hedge);
675 HedgedBlob::new(primary, HedgeSibling::Isolated(hedge), cfg, metrics())
676 }
677
678 const SECS: fn(u64) -> Duration = Duration::from_secs;
679
680 #[mz_ore::test(tokio::test(start_paused = true))]
681 async fn fast_primary_no_hedge() {
682 let primary = TestBlob::new(SECS(0), Ok(Some("x")));
683 let hedge = TestBlob::new(SECS(0), Ok(Some("x")));
684 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
685 assert!(blob.get("k").await.unwrap().is_some());
686 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
687 assert_eq!(blob.metrics.fired.get(), 0);
688 }
689
690 #[mz_ore::test(tokio::test(start_paused = true))]
691 async fn hedge_wins_and_cancels_primary() {
692 let primary = TestBlob::new(SECS(3600), Ok(Some("slow")));
693 let hedge = TestBlob::new(SECS(0), Ok(Some("fast")));
694 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
695 let start = tokio::time::Instant::now();
696 let res = blob.get("k").await.unwrap().expect("some");
697 assert_eq!(res.into_contiguous(), b"fast".to_vec());
700 assert_eq!(start.elapsed(), SECS(2));
701 assert_eq!(blob.metrics.fired.get(), 1);
702 assert_eq!(blob.metrics.won.get(), 1);
703 assert_eq!(blob.metrics.won_seconds.get_sample_count(), 1);
704 }
705
706 #[mz_ore::test(tokio::test(start_paused = true))]
707 async fn primary_wins_after_hedge_fired() {
708 let primary = TestBlob::new(SECS(3), Ok(Some("primary")));
709 let hedge = TestBlob::new(SECS(3600), Ok(Some("hedge")));
710 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
711 let res = blob.get("k").await.unwrap().expect("some");
712 assert_eq!(res.into_contiguous(), b"primary".to_vec());
713 assert_eq!(blob.metrics.fired.get(), 1);
714 assert_eq!(blob.metrics.won.get(), 0);
715 }
716
717 #[mz_ore::test(tokio::test(start_paused = true))]
718 async fn hedge_error_does_not_fail_get() {
719 let primary = TestBlob::new(SECS(5), Ok(Some("primary")));
720 let hedge = TestBlob::new(SECS(0), Err("hedge boom"));
721 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
722 let res = blob.get("k").await.unwrap().expect("some");
725 assert_eq!(res.into_contiguous(), b"primary".to_vec());
726 assert_eq!(blob.metrics.errors.get(), 1);
727 }
728
729 #[mz_ore::test(tokio::test(start_paused = true))]
730 async fn primary_error_then_hedge_success() {
731 let primary = TestBlob::new(SECS(3), Err("primary boom"));
732 let hedge = TestBlob::new(SECS(2), Ok(Some("hedge")));
733 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
734 let res = blob.get("k").await.unwrap().expect("some");
735 assert_eq!(res.into_contiguous(), b"hedge".to_vec());
736 assert_eq!(blob.metrics.won.get(), 1);
737 }
738
739 #[mz_ore::test(tokio::test(start_paused = true))]
740 async fn primary_error_then_hedge_error() {
741 let primary = TestBlob::new(SECS(3), Err("primary boom"));
744 let hedge = TestBlob::new(Duration::from_millis(1500), Err("hedge boom"));
745 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
746 let err = blob.get("k").await.unwrap_err();
747 assert!(err.to_string().contains("primary boom"), "{}", err);
748 assert!(!err.to_string().contains("hedge boom"), "{}", err);
749 assert_eq!(blob.metrics.errors.get(), 1);
750 }
751
752 #[mz_ore::test(tokio::test(start_paused = true))]
753 async fn primary_error_hedge_timeout() {
754 let primary = TestBlob::new(SECS(3), Err("primary boom"));
758 let hedge = TestBlob::new(SECS(3600), Ok(Some("hedge")));
759 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
760 let start = tokio::time::Instant::now();
761 let err = blob.get("k").await.unwrap_err();
762 assert!(err.to_string().contains("primary boom"), "{}", err);
763 assert_eq!(start.elapsed(), SECS(5));
765 }
766
767 #[mz_ore::test(tokio::test(start_paused = true))]
768 async fn dropped_get_releases_concurrency_slot() {
769 let primary = TestBlob::new(SECS(3600), Ok(Some("slow")));
772 let hedge = TestBlob::new(SECS(3600), Ok(Some("slow")));
773 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_MAX_CONCURRENT, 1));
774 let blob = hedged(&primary, &hedge, cfg);
775 for expected_fired in [1, 2] {
776 let res = tokio::time::timeout(SECS(10), blob.get("k")).await;
777 assert!(res.is_err(), "get should still be pending at timeout");
778 assert_eq!(blob.metrics.fired.get(), expected_fired);
779 }
780 assert_eq!(blob.metrics.skipped_concurrency.get(), 0);
781 }
782
783 #[mz_ore::test(tokio::test(start_paused = true))]
784 async fn hedge_error_then_primary_error() {
785 let primary = TestBlob::new(SECS(3), Err("primary boom"));
788 let hedge = TestBlob::new(SECS(0), Err("hedge boom"));
789 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
790 let err = blob.get("k").await.unwrap_err();
791 assert!(err.to_string().contains("primary boom"), "{}", err);
792 assert!(!err.to_string().contains("hedge boom"), "{}", err);
793 assert!(!err.is_timeout());
794 }
795
796 #[mz_ore::test(tokio::test(start_paused = true))]
797 async fn fast_primary_error_passthrough() {
798 let primary = TestBlob::new(SECS(0), Err("fast fail"));
799 let hedge = TestBlob::new(SECS(0), Ok(Some("hedge")));
800 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
801 assert!(blob.get("k").await.is_err());
802 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
803 assert_eq!(blob.metrics.fired.get(), 0);
804 }
805
806 #[mz_ore::test(tokio::test(start_paused = true))]
807 async fn ok_none_wins() {
808 let primary = TestBlob::new(SECS(3600), Ok(None));
809 let hedge = TestBlob::new(SECS(0), Ok(None));
810 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
811 assert!(blob.get("k").await.unwrap().is_none());
812 assert_eq!(blob.metrics.won.get(), 1);
813 }
814
815 #[mz_ore::test(tokio::test(start_paused = true))]
816 async fn disabled_passthrough() {
817 let primary = TestBlob::new(SECS(0), Ok(Some("x")));
818 let hedge = TestBlob::new(SECS(0), Ok(Some("x")));
819 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_ENABLED, false));
820 let blob = hedged(&primary, &hedge, cfg);
821 assert!(blob.get("k").await.unwrap().is_some());
822 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
823 assert_eq!(blob.metrics.fired.get(), 0);
824 }
825
826 #[mz_ore::test(tokio::test(start_paused = true))]
827 async fn unavailable_sibling() {
828 let primary = TestBlob::new(SECS(3), Ok(Some("x")));
829 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
830 let blob = HedgedBlob::new(
831 primary_blob,
832 HedgeSibling::Unavailable,
833 test_cfg(|_| {}),
834 metrics(),
835 );
836 assert_eq!(blob.metrics.armed.get(), 0);
837 assert!(blob.get("k").await.unwrap().is_some());
838 assert_eq!(blob.metrics.skipped_unavailable.get(), 1);
839 }
840
841 #[mz_ore::test(tokio::test(start_paused = true))]
842 async fn budget_exhausts_and_refills() {
843 let primary = TestBlob::new(SECS(10), Ok(Some("slow")));
844 let hedge = TestBlob::new(SECS(0), Ok(Some("fast")));
845 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 0.0));
847 let blob = hedged(&primary, &hedge, Arc::clone(&cfg));
848 for _ in 0..32 {
849 assert!(blob.get("k").await.unwrap().is_some());
850 }
851 assert_eq!(blob.metrics.fired.get(), 32);
852 assert!(blob.get("k").await.unwrap().is_some());
853 assert_eq!(blob.metrics.fired.get(), 32);
854 assert_eq!(blob.metrics.skipped_budget.get(), 1);
855 let mut updates = ConfigUpdates::default();
859 updates.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 1.0);
860 updates.apply(&cfg);
861 assert!(blob.get("k").await.unwrap().is_some());
862 assert_eq!(blob.metrics.skipped_budget.get(), 2);
863 assert!(blob.get("k").await.unwrap().is_some());
864 assert_eq!(blob.metrics.fired.get(), 33);
865 }
866
867 #[mz_ore::test(tokio::test(start_paused = true))]
868 async fn concurrency_cap() {
869 let primary = TestBlob::new(SECS(10), Ok(Some("slow")));
870 let hedge = TestBlob::new(SECS(5), Ok(Some("fast")));
871 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_MAX_CONCURRENT, 1));
872 let blob = hedged(&primary, &hedge, cfg);
873 let (a, b) = tokio::join!(blob.get("k1"), blob.get("k2"));
874 assert!(a.is_ok() && b.is_ok());
875 assert_eq!(blob.metrics.fired.get(), 1);
876 assert_eq!(blob.metrics.skipped_concurrency.get(), 1);
877 }
878
879 #[mz_ore::test(tokio::test(start_paused = true))]
880 async fn warmer_pings_isolated_sibling() {
881 let primary = TestBlob::new(SECS(0), Ok(None));
882 let hedge = TestBlob::new(SECS(0), Ok(None));
883 let cfg = test_cfg(|_| {});
884 let sockets = BLOB_HEDGED_GET_MAX_CONCURRENT.get(&cfg);
885 let blob = hedged(&primary, &hedge, cfg);
886 tokio::time::sleep(SECS(1)).await;
889 tokio::task::yield_now().await;
890 assert_eq!(hedge.gets.load(Ordering::SeqCst), sockets);
891 tokio::time::sleep(SECS(20)).await;
892 tokio::task::yield_now().await;
893 assert_eq!(hedge.gets.load(Ordering::SeqCst), 2 * sockets);
894 drop(blob);
895 }
896
897 #[mz_ore::test(tokio::test(start_paused = true))]
898 async fn warmer_gated_on_enabled() {
899 let primary = TestBlob::new(SECS(0), Ok(None));
900 let hedge = TestBlob::new(SECS(0), Ok(None));
901 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_ENABLED, false));
902 let blob = hedged(&primary, &hedge, Arc::clone(&cfg));
903 tokio::time::sleep(SECS(120)).await;
905 tokio::task::yield_now().await;
906 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
907 let mut updates = ConfigUpdates::default();
909 updates.add(&BLOB_HEDGED_GET_ENABLED, true);
910 updates.apply(&cfg);
911 tokio::time::sleep(*BLOB_HEDGED_GET_WARM_INTERVAL.default() + SECS(1)).await;
912 tokio::task::yield_now().await;
913 assert!(hedge.gets.load(Ordering::SeqCst) > 0);
914 drop(blob);
915 }
916
917 #[mz_ore::test(tokio::test(start_paused = true))]
918 async fn shared_sibling_gets_no_warmer() {
919 let primary = TestBlob::new(SECS(0), Ok(None));
920 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
921 let blob = HedgedBlob::new(
922 primary_blob,
923 HedgeSibling::SharedWithPrimary,
924 test_cfg(|_| {}),
925 metrics(),
926 );
927 assert!(blob._sibling_task.is_none());
928 assert_eq!(blob.metrics.armed.get(), 1);
929 tokio::time::sleep(SECS(60)).await;
930 assert_eq!(primary.gets.load(Ordering::SeqCst), 0);
931 }
932
933 fn flaky_opener(
936 hedge: &Arc<TestBlob>,
937 failures: usize,
938 attempts: &Arc<AtomicUsize>,
939 ) -> HedgeSiblingOpener {
940 let hedge: Arc<dyn Blob> = Arc::<TestBlob>::clone(hedge);
941 let attempts = Arc::clone(attempts);
942 Box::new(move || {
943 let hedge = Arc::clone(&hedge);
944 let attempts = Arc::clone(&attempts);
945 Box::pin(async move {
946 if attempts.fetch_add(1, Ordering::SeqCst) < failures {
947 HedgeSibling::Unavailable
948 } else {
949 HedgeSibling::Isolated(hedge)
950 }
951 })
952 })
953 }
954
955 #[mz_ore::test(tokio::test(start_paused = true))]
956 async fn arming_retries_until_sibling_opens() {
957 let primary = TestBlob::new(SECS(3), Ok(Some("primary")));
958 let hedge = TestBlob::new(SECS(0), Ok(Some("hedge")));
959 let attempts = Arc::new(AtomicUsize::new(0));
960 let cfg = test_cfg(|_| {});
961 let sockets = BLOB_HEDGED_GET_MAX_CONCURRENT.get(&cfg);
962 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
963 let blob = HedgedBlob::new_arming(
964 primary_blob,
965 flaky_opener(&hedge, 3, &attempts),
966 cfg,
967 metrics(),
968 );
969 tokio::task::yield_now().await;
970 assert_eq!(attempts.load(Ordering::SeqCst), 1);
971 assert_eq!(blob.metrics.armed.get(), 0);
972 let res = blob.get("k").await.unwrap().expect("some");
974 assert_eq!(res.into_contiguous(), b"primary".to_vec());
975 assert_eq!(blob.metrics.skipped_unavailable.get(), 1);
976 tokio::time::sleep(SECS(5)).await;
979 tokio::task::yield_now().await;
980 assert_eq!(attempts.load(Ordering::SeqCst), 4);
981 assert_eq!(blob.metrics.armed.get(), 1);
982 let res = blob.get("k").await.unwrap().expect("some");
984 assert_eq!(res.into_contiguous(), b"hedge".to_vec());
985 assert_eq!(blob.metrics.won.get(), 1);
986 assert_eq!(hedge.gets.load(Ordering::SeqCst), sockets + 1);
989 }
990
991 #[mz_ore::test(tokio::test(start_paused = true))]
992 async fn arming_shared_with_primary() {
993 let primary = TestBlob::new(SECS(0), Ok(None));
994 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
995 let opener: HedgeSiblingOpener =
996 Box::new(|| Box::pin(async { HedgeSibling::SharedWithPrimary }));
997 let blob = HedgedBlob::new_arming(primary_blob, opener, test_cfg(|_| {}), metrics());
998 tokio::task::yield_now().await;
999 assert_eq!(blob.metrics.armed.get(), 1);
1000 tokio::time::sleep(SECS(60)).await;
1002 assert_eq!(primary.gets.load(Ordering::SeqCst), 0);
1003 }
1004
1005 #[mz_ore::test(tokio::test(start_paused = true))]
1006 async fn arming_retries_only_while_enabled() {
1007 let primary = TestBlob::new(SECS(0), Ok(None));
1008 let hedge = TestBlob::new(SECS(0), Ok(None));
1009 let attempts = Arc::new(AtomicUsize::new(0));
1010 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_ENABLED, false));
1011 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
1012 let blob = HedgedBlob::new_arming(
1013 primary_blob,
1014 flaky_opener(&hedge, 1, &attempts),
1015 Arc::clone(&cfg),
1016 metrics(),
1017 );
1018 tokio::time::sleep(SECS(600)).await;
1020 tokio::task::yield_now().await;
1021 assert_eq!(attempts.load(Ordering::SeqCst), 1);
1022 assert_eq!(blob.metrics.armed.get(), 0);
1023 let mut updates = ConfigUpdates::default();
1025 updates.add(&BLOB_HEDGED_GET_ENABLED, true);
1026 updates.apply(&cfg);
1027 tokio::time::sleep(*BLOB_HEDGED_GET_WARM_INTERVAL.default() + SECS(1)).await;
1028 tokio::task::yield_now().await;
1029 assert_eq!(attempts.load(Ordering::SeqCst), 2);
1030 assert_eq!(blob.metrics.armed.get(), 1);
1031 }
1032
1033 #[derive(Debug)]
1037 struct SlowGetBlob(Arc<dyn Blob>);
1038
1039 #[async_trait]
1040 impl Blob for SlowGetBlob {
1041 async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
1042 tokio::time::sleep(Duration::from_millis(2)).await;
1043 self.0.get(key).await
1044 }
1045
1046 async fn list_keys_and_metadata(
1047 &self,
1048 key_prefix: &str,
1049 f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
1050 ) -> Result<(), ExternalError> {
1051 self.0.list_keys_and_metadata(key_prefix, f).await
1052 }
1053
1054 async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError> {
1055 self.0.set(key, value).await
1056 }
1057
1058 async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError> {
1059 self.0.delete(key).await
1060 }
1061
1062 async fn restore(&self, key: &str) -> Result<(), ExternalError> {
1063 self.0.restore(key).await
1064 }
1065 }
1066
1067 #[mz_ore::test(tokio::test)]
1072 #[cfg_attr(miri, ignore)] async fn hedged_blob_conformance() {
1074 let registry = Arc::new(tokio::sync::Mutex::new(MemMultiRegistry::new(false)));
1075 let cfg = test_cfg(|u| {
1076 u.add(&BLOB_HEDGED_GET_DELAY, Duration::ZERO);
1077 u.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 1.0);
1078 });
1079 let metrics = metrics();
1080 let metrics_check = metrics.clone();
1081 blob_impl_test(move |path| {
1082 let path = path.to_owned();
1083 let registry = Arc::clone(®istry);
1084 let cfg = Arc::clone(&cfg);
1085 let metrics = metrics.clone();
1086 async move {
1087 let store: Arc<dyn Blob> = Arc::new(registry.lock().await.blob(&path));
1088 let primary: Arc<dyn Blob> = Arc::new(SlowGetBlob(Arc::clone(&store)));
1089 Ok(HedgedBlob::new(
1090 primary,
1091 HedgeSibling::Isolated(store),
1092 cfg,
1093 metrics,
1094 ))
1095 }
1096 })
1097 .await
1098 .expect("conformance");
1099 assert!(metrics_check.fired.get() > 0, "no hedge ever fired");
1100 assert!(metrics_check.won.get() > 0, "no hedge ever won");
1101 }
1102
1103 #[mz_ore::test(tokio::test(start_paused = true))]
1104 async fn dropping_blob_stops_arming() {
1105 let primary: Arc<dyn Blob> = TestBlob::new(SECS(0), Ok(None));
1106 let hedge = TestBlob::new(SECS(0), Ok(None));
1107 let attempts = Arc::new(AtomicUsize::new(0));
1108 let blob = HedgedBlob::new_arming(
1109 primary,
1110 flaky_opener(&hedge, usize::MAX, &attempts),
1111 test_cfg(|_| {}),
1112 metrics(),
1113 );
1114 tokio::time::sleep(SECS(10)).await;
1115 tokio::task::yield_now().await;
1116 let seen = attempts.load(Ordering::SeqCst);
1117 assert!(seen >= 3, "{seen}");
1118 drop(blob);
1119 tokio::task::yield_now().await;
1120 tokio::time::sleep(SECS(3600)).await;
1121 tokio::task::yield_now().await;
1122 assert_eq!(attempts.load(Ordering::SeqCst), seen);
1123 }
1124
1125 #[mz_ore::test(tokio::test(start_paused = true))]
1126 async fn dropping_blob_cancels_open() {
1127 let primary: Arc<dyn Blob> = TestBlob::new(SECS(0), Ok(None));
1128 let sentinel = Arc::new(());
1130 let opener: HedgeSiblingOpener = {
1131 let sentinel = Arc::clone(&sentinel);
1132 Box::new(move || {
1133 let held = Arc::clone(&sentinel);
1134 Box::pin(async move {
1135 let _held = held;
1136 std::future::pending::<HedgeSibling>().await
1137 })
1138 })
1139 };
1140 let blob = HedgedBlob::new_arming(primary, opener, test_cfg(|_| {}), metrics());
1141 tokio::task::yield_now().await;
1142 assert_eq!(Arc::strong_count(&sentinel), 3);
1143 drop(blob);
1144 tokio::task::yield_now().await;
1145 tokio::task::yield_now().await;
1146 assert_eq!(Arc::strong_count(&sentinel), 1);
1147 }
1148
1149 #[mz_ore::test(tokio::test(start_paused = true))]
1150 async fn arming_backoff_doubles_up_to_the_ceiling() {
1151 let primary: Arc<dyn Blob> = TestBlob::new(SECS(0), Ok(None));
1152 let attempt_times = Arc::new(std::sync::Mutex::new(Vec::new()));
1153 let opener: HedgeSiblingOpener = {
1154 let attempt_times = Arc::clone(&attempt_times);
1155 Box::new(move || {
1156 attempt_times
1157 .lock()
1158 .expect("lock")
1159 .push(tokio::time::Instant::now());
1160 Box::pin(async { HedgeSibling::Unavailable })
1161 })
1162 };
1163 let _blob = HedgedBlob::new_arming(primary, opener, test_cfg(|_| {}), metrics());
1164 tokio::time::sleep(SECS(300)).await;
1165 let attempt_times = attempt_times.lock().expect("lock").clone();
1166 let gaps: Vec<u64> = attempt_times
1167 .windows(2)
1168 .map(|w| (w[1] - w[0]).as_secs())
1169 .take(8)
1170 .collect();
1171 assert_eq!(gaps, vec![1, 2, 4, 8, 16, 32, 60, 60]);
1172 }
1173}