1use std::sync::Arc;
50use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
51use std::time::{Duration, Instant};
52
53use async_trait::async_trait;
54use bytes::Bytes;
55use futures_util::future::{Either, select};
56use mz_dyncfg::{Config, ConfigSet, ParameterScope};
57use mz_ore::bytes::SegmentedBytes;
58use mz_ore::cast::CastLossy;
59use mz_ore::task::AbortOnDropHandle;
60use tracing::{debug, warn};
61
62use crate::location::{BLOB_GET_LIVENESS_KEY, Blob, BlobMetadata, ExternalError};
63use crate::metrics::BlobHedgeMetrics;
64
65pub(crate) const BLOB_HEDGED_GET_ENABLED: Config<bool> = Config::new(
66 "persist_blob_hedged_get_enabled",
67 false,
68 "Whether to hedge slow blob gets with a second request on a separate \
69 connection pool (Materialize).",
70 ParameterScope::Environment,
71);
72
73pub(crate) const BLOB_HEDGED_GET_DELAY: Config<Duration> = Config::new(
74 "persist_blob_hedged_get_delay",
75 Duration::from_secs(2),
76 "How long a blob get may be in flight before a hedge request is fired \
77 (Materialize).",
78 ParameterScope::Environment,
79);
80
81pub(crate) const BLOB_HEDGED_GET_MAX_CONCURRENT: Config<usize> = Config::new(
82 "persist_blob_hedged_get_max_concurrent",
83 2,
87 "Maximum concurrent hedge requests per blob handle, bounding the memory \
88 held by raced gets and the number of warm sockets (Materialize).",
89 ParameterScope::Environment,
90);
91
92pub(crate) const BLOB_HEDGED_GET_BUDGET_RATIO: Config<f64> = Config::new(
93 "persist_blob_hedged_get_budget_ratio",
94 0.01,
95 "Long-run bound on hedge requests as a fraction of blob gets \
96 (Materialize).",
97 ParameterScope::Environment,
98);
99
100pub(crate) const BLOB_HEDGED_GET_WARM_INTERVAL: Config<Duration> = Config::new(
104 "persist_blob_hedged_get_warm_interval",
105 Duration::from_secs(20),
106 "How often to issue liveness gets that keep the hedge connection pool \
107 warm, 0 disables warming without disabling hedging (Materialize).",
108 ParameterScope::Environment,
109);
110
111const HEDGE_COST_MICRO_TOKENS: u64 = 1_000_000;
115
116const BUDGET_BURST_MICRO_TOKENS: u64 = 32 * HEDGE_COST_MICRO_TOKENS;
124
125enum HedgeRefused {
127 Concurrency,
128 Budget,
129}
130
131#[derive(Debug)]
137struct HedgeBudget {
138 concurrent: AtomicUsize,
139 micro_tokens: AtomicU64,
140}
141
142impl HedgeBudget {
143 fn new() -> Self {
144 HedgeBudget {
145 concurrent: AtomicUsize::new(0),
146 micro_tokens: AtomicU64::new(BUDGET_BURST_MICRO_TOKENS),
147 }
148 }
149
150 fn try_acquire(&self, max_concurrent: usize) -> Result<HedgeGuard<'_>, HedgeRefused> {
154 let got_slot = self
155 .concurrent
156 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |c| {
157 (c < max_concurrent).then_some(c + 1)
158 })
159 .is_ok();
160 if !got_slot {
161 return Err(HedgeRefused::Concurrency);
162 }
163 let guard = HedgeGuard(self);
166 let took_token = self
167 .micro_tokens
168 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |t| {
169 t.checked_sub(HEDGE_COST_MICRO_TOKENS)
170 })
171 .is_ok();
172 if took_token {
173 Ok(guard)
174 } else {
175 Err(HedgeRefused::Budget)
176 }
177 }
178
179 fn replenish(&self, ratio: f64) {
182 let add = u64::cast_lossy(ratio.clamp(0.0, 1.0) * f64::cast_lossy(HEDGE_COST_MICRO_TOKENS));
183 if add == 0 {
184 return;
185 }
186 let _ = self
187 .micro_tokens
188 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |t| {
189 Some((t + add).min(BUDGET_BURST_MICRO_TOKENS))
190 });
191 }
192}
193
194struct HedgeGuard<'a>(&'a HedgeBudget);
195
196impl Drop for HedgeGuard<'_> {
197 fn drop(&mut self) {
198 self.0.concurrent.fetch_sub(1, Ordering::Relaxed);
199 }
200}
201
202#[derive(Debug)]
205pub enum HedgeSibling {
206 Isolated(Arc<dyn Blob>),
209 SharedWithPrimary,
213 Unavailable,
216}
217
218#[derive(Debug)]
220pub struct HedgedBlob {
221 primary: Arc<dyn Blob>,
222 hedge: Option<Arc<dyn Blob>>,
225 cfg: Arc<ConfigSet>,
226 metrics: BlobHedgeMetrics,
227 budget: HedgeBudget,
228 _warmer: Option<AbortOnDropHandle<()>>,
229}
230
231fn spawn_warmer(
239 hedge: Arc<dyn Blob>,
240 cfg: Arc<ConfigSet>,
241 metrics: BlobHedgeMetrics,
242) -> AbortOnDropHandle<()> {
243 mz_ore::task::spawn(|| "persist::blob_hedge_warmer", async move {
244 loop {
245 let interval = BLOB_HEDGED_GET_WARM_INTERVAL.get(&cfg);
246 if !BLOB_HEDGED_GET_ENABLED.get(&cfg) || interval == Duration::ZERO {
247 let recheck = if interval == Duration::ZERO {
251 *BLOB_HEDGED_GET_WARM_INTERVAL.default()
252 } else {
253 interval
254 };
255 tokio::time::sleep(recheck).await;
256 continue;
257 }
258 let start = Instant::now();
263 let sockets = BLOB_HEDGED_GET_MAX_CONCURRENT.get(&cfg);
264 let pings = (0..sockets).map(|_| hedge.get(BLOB_GET_LIVENESS_KEY));
265 let pings = futures_util::future::join_all(pings);
266 match tokio::time::timeout(interval, pings).await {
271 Ok(results) if results.iter().all(|r| r.is_ok()) => {
272 metrics.rtt_latency.set(start.elapsed().as_secs_f64());
273 }
274 Ok(_) | Err(_) => {
275 metrics.warm_errors.inc();
280 }
281 }
282 tokio::time::sleep(interval).await;
283 }
284 })
285 .abort_on_drop()
286}
287
288impl HedgedBlob {
289 pub fn new(
294 primary: Arc<dyn Blob>,
295 sibling: HedgeSibling,
296 cfg: Arc<ConfigSet>,
297 metrics: BlobHedgeMetrics,
298 ) -> HedgedBlob {
299 let (hedge, warmer) = match sibling {
300 HedgeSibling::Isolated(h) => {
301 let warmer = spawn_warmer(Arc::clone(&h), Arc::clone(&cfg), metrics.clone());
302 (Some(h), Some(warmer))
303 }
304 HedgeSibling::SharedWithPrimary => (Some(Arc::clone(&primary)), None),
305 HedgeSibling::Unavailable => (None, None),
306 };
307 metrics.armed.set(i64::from(hedge.is_some()));
308 HedgedBlob {
309 primary,
310 hedge,
311 cfg,
312 metrics,
313 budget: HedgeBudget::new(),
314 _warmer: warmer,
315 }
316 }
317
318 fn admit(&self) -> Option<(&Arc<dyn Blob>, HedgeGuard<'_>)> {
321 let Some(hedge_blob) = &self.hedge else {
322 self.metrics.skipped_unavailable.inc();
323 return None;
324 };
325 match self
326 .budget
327 .try_acquire(BLOB_HEDGED_GET_MAX_CONCURRENT.get(&self.cfg))
328 {
329 Ok(guard) => Some((hedge_blob, guard)),
330 Err(HedgeRefused::Concurrency) => {
331 self.metrics.skipped_concurrency.inc();
332 None
333 }
334 Err(HedgeRefused::Budget) => {
335 self.metrics.skipped_budget.inc();
336 None
337 }
338 }
339 }
340
341 fn record_win(&self, key: &str, start: Instant) {
342 self.metrics.won.inc();
343 self.metrics
344 .won_seconds
345 .observe(start.elapsed().as_secs_f64());
346 debug!(%key, elapsed = ?start.elapsed(), "blob get won by hedge request");
347 }
348
349 async fn get_hedged(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
350 let start = Instant::now();
351 let delay = BLOB_HEDGED_GET_DELAY.get(&self.cfg);
352 let mut primary = std::pin::pin!(self.primary.get(key));
353 if let Ok(res) = tokio::time::timeout(delay, primary.as_mut()).await {
360 return res;
363 }
364 let Some((hedge_blob, guard)) = self.admit() else {
365 return primary.await;
366 };
367 self.metrics.fired.inc();
368 let mut hedge = std::pin::pin!(hedge_blob.get(key));
369 match select(primary.as_mut(), hedge.as_mut()).await {
381 Either::Left((Ok(res), _hedge)) => Ok(res),
382 Either::Right((Ok(res), _primary)) => {
383 self.record_win(key, start);
384 Ok(res)
385 }
386 Either::Left((Err(primary_err), hedge)) => {
387 match tokio::time::timeout(delay, hedge).await {
404 Ok(Ok(res)) => {
405 self.record_win(key, start);
406 Ok(res)
407 }
408 Ok(Err(hedge_err)) => {
409 self.metrics.errors.inc();
410 warn!(%key, %hedge_err, "hedged blob get: both requests failed");
411 Err(primary_err)
419 }
420 Err(_elapsed) => {
421 warn!(%key, "hedged blob get: primary failed, hedge still pending");
422 Err(primary_err)
423 }
424 }
425 }
426 Either::Right((Err(hedge_err), primary)) => {
427 self.metrics.errors.inc();
428 warn!(%key, %hedge_err, "hedge request failed, awaiting primary");
429 drop(guard);
435 primary.await
436 }
437 }
438 }
439}
440
441#[async_trait]
442impl Blob for HedgedBlob {
443 async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
444 if !BLOB_HEDGED_GET_ENABLED.get(&self.cfg) {
445 return self.primary.get(key).await;
446 }
447 let res = self.get_hedged(key).await;
448 self.budget
449 .replenish(BLOB_HEDGED_GET_BUDGET_RATIO.get(&self.cfg));
450 res
451 }
452
453 async fn list_keys_and_metadata(
454 &self,
455 key_prefix: &str,
456 f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
457 ) -> Result<(), ExternalError> {
458 self.primary.list_keys_and_metadata(key_prefix, f).await
459 }
460
461 async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError> {
462 self.primary.set(key, value).await
463 }
464
465 async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError> {
466 self.primary.delete(key).await
467 }
468
469 async fn restore(&self, key: &str) -> Result<(), ExternalError> {
470 self.primary.restore(key).await
471 }
472}
473
474#[cfg(test)]
475mod tests {
476 use std::sync::atomic::{AtomicUsize, Ordering};
477
478 use anyhow::anyhow;
479 use mz_dyncfg::ConfigUpdates;
480 use mz_ore::metrics::MetricsRegistry;
481
482 use crate::location::tests::blob_impl_test;
483 use crate::mem::MemMultiRegistry;
484
485 use super::*;
486
487 #[derive(Debug)]
490 struct TestBlob {
491 delay: Duration,
492 outcome: Result<Option<&'static str>, &'static str>,
493 gets: AtomicUsize,
494 }
495
496 impl TestBlob {
497 fn new(
498 delay: Duration,
499 outcome: Result<Option<&'static str>, &'static str>,
500 ) -> Arc<TestBlob> {
501 Arc::new(TestBlob {
502 delay,
503 outcome,
504 gets: AtomicUsize::new(0),
505 })
506 }
507 }
508
509 #[async_trait]
510 impl Blob for TestBlob {
511 async fn get(&self, _key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
512 self.gets.fetch_add(1, Ordering::SeqCst);
513 tokio::time::sleep(self.delay).await;
514 match self.outcome {
515 Ok(x) => Ok(x.map(|x| SegmentedBytes::from(Bytes::from(x)))),
516 Err(msg) => Err(ExternalError::from(anyhow!(msg))),
517 }
518 }
519
520 async fn list_keys_and_metadata(
521 &self,
522 _key_prefix: &str,
523 _f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
524 ) -> Result<(), ExternalError> {
525 unreachable!("test blob only supports get")
526 }
527
528 async fn set(&self, _key: &str, _value: Bytes) -> Result<(), ExternalError> {
529 unreachable!("test blob only supports get")
530 }
531
532 async fn delete(&self, _key: &str) -> Result<Option<usize>, ExternalError> {
533 unreachable!("test blob only supports get")
534 }
535
536 async fn restore(&self, _key: &str) -> Result<(), ExternalError> {
537 unreachable!("test blob only supports get")
538 }
539 }
540
541 fn test_cfg(customize: impl FnOnce(&mut ConfigUpdates)) -> Arc<ConfigSet> {
542 let cfg = crate::cfg::all_dyn_configs(ConfigSet::default());
543 let mut updates = ConfigUpdates::default();
544 updates.add(&BLOB_HEDGED_GET_ENABLED, true);
545 customize(&mut updates);
546 updates.apply(&cfg);
547 Arc::new(cfg)
548 }
549
550 fn metrics() -> BlobHedgeMetrics {
551 BlobHedgeMetrics::new(&MetricsRegistry::new())
552 }
553
554 fn hedged(primary: &Arc<TestBlob>, hedge: &Arc<TestBlob>, cfg: Arc<ConfigSet>) -> HedgedBlob {
555 let primary: Arc<dyn Blob> = Arc::<TestBlob>::clone(primary);
556 let hedge: Arc<dyn Blob> = Arc::<TestBlob>::clone(hedge);
557 HedgedBlob::new(primary, HedgeSibling::Isolated(hedge), cfg, metrics())
558 }
559
560 const SECS: fn(u64) -> Duration = Duration::from_secs;
561
562 #[mz_ore::test(tokio::test(start_paused = true))]
563 async fn fast_primary_no_hedge() {
564 let primary = TestBlob::new(SECS(0), Ok(Some("x")));
565 let hedge = TestBlob::new(SECS(0), Ok(Some("x")));
566 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
567 assert!(blob.get("k").await.unwrap().is_some());
568 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
569 assert_eq!(blob.metrics.fired.get(), 0);
570 }
571
572 #[mz_ore::test(tokio::test(start_paused = true))]
573 async fn hedge_wins_and_cancels_primary() {
574 let primary = TestBlob::new(SECS(3600), Ok(Some("slow")));
575 let hedge = TestBlob::new(SECS(0), Ok(Some("fast")));
576 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
577 let start = tokio::time::Instant::now();
578 let res = blob.get("k").await.unwrap().expect("some");
579 assert_eq!(res.into_contiguous(), b"fast".to_vec());
582 assert_eq!(start.elapsed(), SECS(2));
583 assert_eq!(blob.metrics.fired.get(), 1);
584 assert_eq!(blob.metrics.won.get(), 1);
585 assert_eq!(blob.metrics.won_seconds.get_sample_count(), 1);
586 }
587
588 #[mz_ore::test(tokio::test(start_paused = true))]
589 async fn primary_wins_after_hedge_fired() {
590 let primary = TestBlob::new(SECS(3), Ok(Some("primary")));
591 let hedge = TestBlob::new(SECS(3600), Ok(Some("hedge")));
592 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
593 let res = blob.get("k").await.unwrap().expect("some");
594 assert_eq!(res.into_contiguous(), b"primary".to_vec());
595 assert_eq!(blob.metrics.fired.get(), 1);
596 assert_eq!(blob.metrics.won.get(), 0);
597 }
598
599 #[mz_ore::test(tokio::test(start_paused = true))]
600 async fn hedge_error_does_not_fail_get() {
601 let primary = TestBlob::new(SECS(5), Ok(Some("primary")));
602 let hedge = TestBlob::new(SECS(0), Err("hedge boom"));
603 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
604 let res = blob.get("k").await.unwrap().expect("some");
607 assert_eq!(res.into_contiguous(), b"primary".to_vec());
608 assert_eq!(blob.metrics.errors.get(), 1);
609 }
610
611 #[mz_ore::test(tokio::test(start_paused = true))]
612 async fn primary_error_then_hedge_success() {
613 let primary = TestBlob::new(SECS(3), Err("primary boom"));
614 let hedge = TestBlob::new(SECS(2), Ok(Some("hedge")));
615 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
616 let res = blob.get("k").await.unwrap().expect("some");
617 assert_eq!(res.into_contiguous(), b"hedge".to_vec());
618 assert_eq!(blob.metrics.won.get(), 1);
619 }
620
621 #[mz_ore::test(tokio::test(start_paused = true))]
622 async fn primary_error_then_hedge_error() {
623 let primary = TestBlob::new(SECS(3), Err("primary boom"));
626 let hedge = TestBlob::new(Duration::from_millis(1500), Err("hedge boom"));
627 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
628 let err = blob.get("k").await.unwrap_err();
629 assert!(err.to_string().contains("primary boom"), "{}", err);
630 assert!(!err.to_string().contains("hedge boom"), "{}", err);
631 assert_eq!(blob.metrics.errors.get(), 1);
632 }
633
634 #[mz_ore::test(tokio::test(start_paused = true))]
635 async fn primary_error_hedge_timeout() {
636 let primary = TestBlob::new(SECS(3), Err("primary boom"));
640 let hedge = TestBlob::new(SECS(3600), Ok(Some("hedge")));
641 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
642 let start = tokio::time::Instant::now();
643 let err = blob.get("k").await.unwrap_err();
644 assert!(err.to_string().contains("primary boom"), "{}", err);
645 assert_eq!(start.elapsed(), SECS(5));
647 }
648
649 #[mz_ore::test(tokio::test(start_paused = true))]
650 async fn dropped_get_releases_concurrency_slot() {
651 let primary = TestBlob::new(SECS(3600), Ok(Some("slow")));
654 let hedge = TestBlob::new(SECS(3600), Ok(Some("slow")));
655 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_MAX_CONCURRENT, 1));
656 let blob = hedged(&primary, &hedge, cfg);
657 for expected_fired in [1, 2] {
658 let res = tokio::time::timeout(SECS(10), blob.get("k")).await;
659 assert!(res.is_err(), "get should still be pending at timeout");
660 assert_eq!(blob.metrics.fired.get(), expected_fired);
661 }
662 assert_eq!(blob.metrics.skipped_concurrency.get(), 0);
663 }
664
665 #[mz_ore::test(tokio::test(start_paused = true))]
666 async fn hedge_error_then_primary_error() {
667 let primary = TestBlob::new(SECS(3), Err("primary boom"));
670 let hedge = TestBlob::new(SECS(0), Err("hedge boom"));
671 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
672 let err = blob.get("k").await.unwrap_err();
673 assert!(err.to_string().contains("primary boom"), "{}", err);
674 assert!(!err.to_string().contains("hedge boom"), "{}", err);
675 assert!(!err.is_timeout());
676 }
677
678 #[mz_ore::test(tokio::test(start_paused = true))]
679 async fn fast_primary_error_passthrough() {
680 let primary = TestBlob::new(SECS(0), Err("fast fail"));
681 let hedge = TestBlob::new(SECS(0), Ok(Some("hedge")));
682 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
683 assert!(blob.get("k").await.is_err());
684 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
685 assert_eq!(blob.metrics.fired.get(), 0);
686 }
687
688 #[mz_ore::test(tokio::test(start_paused = true))]
689 async fn ok_none_wins() {
690 let primary = TestBlob::new(SECS(3600), Ok(None));
691 let hedge = TestBlob::new(SECS(0), Ok(None));
692 let blob = hedged(&primary, &hedge, test_cfg(|_| {}));
693 assert!(blob.get("k").await.unwrap().is_none());
694 assert_eq!(blob.metrics.won.get(), 1);
695 }
696
697 #[mz_ore::test(tokio::test(start_paused = true))]
698 async fn disabled_passthrough() {
699 let primary = TestBlob::new(SECS(0), Ok(Some("x")));
700 let hedge = TestBlob::new(SECS(0), Ok(Some("x")));
701 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_ENABLED, false));
702 let blob = hedged(&primary, &hedge, cfg);
703 assert!(blob.get("k").await.unwrap().is_some());
704 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
705 assert_eq!(blob.metrics.fired.get(), 0);
706 }
707
708 #[mz_ore::test(tokio::test(start_paused = true))]
709 async fn unavailable_sibling() {
710 let primary = TestBlob::new(SECS(3), Ok(Some("x")));
711 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
712 let blob = HedgedBlob::new(
713 primary_blob,
714 HedgeSibling::Unavailable,
715 test_cfg(|_| {}),
716 metrics(),
717 );
718 assert_eq!(blob.metrics.armed.get(), 0);
719 assert!(blob.get("k").await.unwrap().is_some());
720 assert_eq!(blob.metrics.skipped_unavailable.get(), 1);
721 }
722
723 #[mz_ore::test(tokio::test(start_paused = true))]
724 async fn budget_exhausts_and_refills() {
725 let primary = TestBlob::new(SECS(10), Ok(Some("slow")));
726 let hedge = TestBlob::new(SECS(0), Ok(Some("fast")));
727 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 0.0));
729 let blob = hedged(&primary, &hedge, Arc::clone(&cfg));
730 for _ in 0..32 {
731 assert!(blob.get("k").await.unwrap().is_some());
732 }
733 assert_eq!(blob.metrics.fired.get(), 32);
734 assert!(blob.get("k").await.unwrap().is_some());
735 assert_eq!(blob.metrics.fired.get(), 32);
736 assert_eq!(blob.metrics.skipped_budget.get(), 1);
737 let mut updates = ConfigUpdates::default();
741 updates.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 1.0);
742 updates.apply(&cfg);
743 assert!(blob.get("k").await.unwrap().is_some());
744 assert_eq!(blob.metrics.skipped_budget.get(), 2);
745 assert!(blob.get("k").await.unwrap().is_some());
746 assert_eq!(blob.metrics.fired.get(), 33);
747 }
748
749 #[mz_ore::test(tokio::test(start_paused = true))]
750 async fn concurrency_cap() {
751 let primary = TestBlob::new(SECS(10), Ok(Some("slow")));
752 let hedge = TestBlob::new(SECS(5), Ok(Some("fast")));
753 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_MAX_CONCURRENT, 1));
754 let blob = hedged(&primary, &hedge, cfg);
755 let (a, b) = tokio::join!(blob.get("k1"), blob.get("k2"));
756 assert!(a.is_ok() && b.is_ok());
757 assert_eq!(blob.metrics.fired.get(), 1);
758 assert_eq!(blob.metrics.skipped_concurrency.get(), 1);
759 }
760
761 #[mz_ore::test(tokio::test(start_paused = true))]
762 async fn warmer_pings_isolated_sibling() {
763 let primary = TestBlob::new(SECS(0), Ok(None));
764 let hedge = TestBlob::new(SECS(0), Ok(None));
765 let cfg = test_cfg(|_| {});
766 let sockets = BLOB_HEDGED_GET_MAX_CONCURRENT.get(&cfg);
767 let blob = hedged(&primary, &hedge, cfg);
768 tokio::time::sleep(SECS(1)).await;
771 tokio::task::yield_now().await;
772 assert_eq!(hedge.gets.load(Ordering::SeqCst), sockets);
773 tokio::time::sleep(SECS(20)).await;
774 tokio::task::yield_now().await;
775 assert_eq!(hedge.gets.load(Ordering::SeqCst), 2 * sockets);
776 drop(blob);
777 }
778
779 #[mz_ore::test(tokio::test(start_paused = true))]
780 async fn warmer_gated_on_enabled() {
781 let primary = TestBlob::new(SECS(0), Ok(None));
782 let hedge = TestBlob::new(SECS(0), Ok(None));
783 let cfg = test_cfg(|u| u.add(&BLOB_HEDGED_GET_ENABLED, false));
784 let blob = hedged(&primary, &hedge, Arc::clone(&cfg));
785 tokio::time::sleep(SECS(120)).await;
787 tokio::task::yield_now().await;
788 assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
789 let mut updates = ConfigUpdates::default();
791 updates.add(&BLOB_HEDGED_GET_ENABLED, true);
792 updates.apply(&cfg);
793 tokio::time::sleep(*BLOB_HEDGED_GET_WARM_INTERVAL.default() + SECS(1)).await;
794 tokio::task::yield_now().await;
795 assert!(hedge.gets.load(Ordering::SeqCst) > 0);
796 drop(blob);
797 }
798
799 #[mz_ore::test(tokio::test(start_paused = true))]
800 async fn shared_sibling_gets_no_warmer() {
801 let primary = TestBlob::new(SECS(0), Ok(None));
802 let primary_blob: Arc<dyn Blob> = Arc::<TestBlob>::clone(&primary);
803 let blob = HedgedBlob::new(
804 primary_blob,
805 HedgeSibling::SharedWithPrimary,
806 test_cfg(|_| {}),
807 metrics(),
808 );
809 assert!(blob._warmer.is_none());
810 assert_eq!(blob.metrics.armed.get(), 1);
811 tokio::time::sleep(SECS(60)).await;
812 assert_eq!(primary.gets.load(Ordering::SeqCst), 0);
813 }
814
815 #[derive(Debug)]
819 struct SlowGetBlob(Arc<dyn Blob>);
820
821 #[async_trait]
822 impl Blob for SlowGetBlob {
823 async fn get(&self, key: &str) -> Result<Option<SegmentedBytes>, ExternalError> {
824 tokio::time::sleep(Duration::from_millis(2)).await;
825 self.0.get(key).await
826 }
827
828 async fn list_keys_and_metadata(
829 &self,
830 key_prefix: &str,
831 f: &mut (dyn FnMut(BlobMetadata) + Send + Sync),
832 ) -> Result<(), ExternalError> {
833 self.0.list_keys_and_metadata(key_prefix, f).await
834 }
835
836 async fn set(&self, key: &str, value: Bytes) -> Result<(), ExternalError> {
837 self.0.set(key, value).await
838 }
839
840 async fn delete(&self, key: &str) -> Result<Option<usize>, ExternalError> {
841 self.0.delete(key).await
842 }
843
844 async fn restore(&self, key: &str) -> Result<(), ExternalError> {
845 self.0.restore(key).await
846 }
847 }
848
849 #[mz_ore::test(tokio::test)]
854 #[cfg_attr(miri, ignore)] async fn hedged_blob_conformance() {
856 let registry = Arc::new(tokio::sync::Mutex::new(MemMultiRegistry::new(false)));
857 let cfg = test_cfg(|u| {
858 u.add(&BLOB_HEDGED_GET_DELAY, Duration::ZERO);
859 u.add(&BLOB_HEDGED_GET_BUDGET_RATIO, 1.0);
860 });
861 let metrics = metrics();
862 let metrics_check = metrics.clone();
863 blob_impl_test(move |path| {
864 let path = path.to_owned();
865 let registry = Arc::clone(®istry);
866 let cfg = Arc::clone(&cfg);
867 let metrics = metrics.clone();
868 async move {
869 let store: Arc<dyn Blob> = Arc::new(registry.lock().await.blob(&path));
870 let primary: Arc<dyn Blob> = Arc::new(SlowGetBlob(Arc::clone(&store)));
871 Ok(HedgedBlob::new(
872 primary,
873 HedgeSibling::Isolated(store),
874 cfg,
875 metrics,
876 ))
877 }
878 })
879 .await
880 .expect("conformance");
881 assert!(metrics_check.fired.get() > 0, "no hedge ever fired");
882 assert!(metrics_check.won.get() > 0, "no hedge ever won");
883 }
884}