Skip to main content

mz_persist/
hedge.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! A [Blob] decorator that hedges slow `get` requests.
11//!
12//! Established connections to the blob store occasionally die in ways that
13//! surface only after multiple seconds (a TCP reset after a hang, or a black
14//! hole), well before any client timeout fires. A `get` riding such a
15//! connection stalls everything downstream of it, while other connections on
16//! the same process serve the same store normally. The mitigation, endorsed by
17//! the major object stores for idempotent reads, is a hedged request: if the
18//! first `get` has not completed within a short delay, race a second one on a
19//! connection the first cannot have poisoned, and take whichever succeeds
20//! first.
21//!
22//! Only `get` is hedged. All other [Blob] methods are forwarded to the
23//! primary handle untouched: writes, deletes, and restores have side
24//! effects, and lists are not latency-critical enough to justify racing a
25//! streaming interface. Extending hedging to any of them is forbidden.
26//!
27//! The hedge handle must not share a connection pool (or DNS state) with the
28//! primary, otherwise the hedge can be handed a connection dying in the same
29//! event that stalled the primary, exactly when a hedge matters most. See
30//! [crate::cfg::open_hedge_sibling] for how that isolation is constructed
31//! per backend.
32//!
33//! Hedging operates within a single `retry_external` attempt, before any
34//! failure surfaces. The retrying in `retry_external`, which is what
35//! recovers this failure class when hedging is off (at the cost of the full
36//! hang), stays untouched as the backstop. The governing principle for
37//! every race below: the primary's outcome is authoritative, and the hedge
38//! is opportunistic, invisible unless it wins. Nothing here assumes callers
39//! retry: every branch of the race degrades to the outcome of the un-hedged
40//! get, delayed by at most one hedge delay, so a caller that treats a get
41//! error as fatal sees the same error it would have seen without hedging,
42//! at most that one delay later.
43//!
44//! NOTE: enabling hedging largely suppresses the old fingerprints of the
45//! dead-connection class (client timeout counters, the SDK's
46//! connection-poisoning log lines), because the hung request is cancelled
47//! before they trigger. The `hedges_won` counter is the replacement signal.
48
49use 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    // Bounds the extra in-flight bytes (which the fetch memory semaphore
86    // cannot see) to about two batch parts. The warmer holds this many
87    // sockets open, so every admitted hedge can be served warm.
88    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
102// NOTE: the warmer, and any retries of a failed sibling open, only run while
103// hedging is enabled, so `enabled` stops the sibling's traffic too. Setting
104// this knob to 0 additionally stops the warmer while keeping hedging on,
105// which is why it must be changeable at runtime.
106pub(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
114/// The cost of one hedge in bucket tokens. Micro-token granularity keeps
115/// small `budget_ratio` values (down to 1e-6) from rounding to "never
116/// refill".
117const HEDGE_COST_MICRO_TOKENS: u64 = 1_000_000;
118
119/// Token-bucket capacity: 32 hedges.
120///
121/// The bucket's shape, not an operational lever: the tuning lever is
122/// `persist_blob_hedged_get_budget_ratio` and the kill switch is
123/// `persist_blob_hedged_get_enabled`. NOTE: because the bucket starts full,
124/// `budget_ratio = 0` still permits ~32 banked hedges before draining. It is
125/// not an instant stop, `enabled` is.
126const BUDGET_BURST_MICRO_TOKENS: u64 = 32 * HEDGE_COST_MICRO_TOKENS;
127
128/// Backoff between sibling open attempts. The ceiling bounds a long outage to
129/// one attempt per minute per process (a credential resolution, a health-check
130/// get and a warning) and still re-arms within a minute of recovery.
131const SIBLING_OPEN_RETRY_INITIAL_BACKOFF: Duration = Duration::from_secs(1);
132const SIBLING_OPEN_RETRY_MAX_BACKOFF: Duration = Duration::from_secs(60);
133
134/// Why a hedge was not fired for a get that exceeded the delay.
135enum HedgeRefused {
136    Concurrency,
137    Budget,
138}
139
140/// Bounds hedge amplification with two independent guards: the concurrency
141/// cap bounds memory held by raced gets, the token bucket bounds long-run
142/// request-rate/egress amplification (e.g. a store-wide brownout making
143/// every get slow, or large gets that legitimately exceed the delay, must
144/// not settle into hedging every request).
145#[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    /// Attempts to acquire both guards. The returned guard releases the
160    /// concurrency slot on drop. Spent tokens come back only via
161    /// [HedgeBudget::replenish].
162    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        // Constructed before the token take so its drop releases the slot
173        // on the budget-refusal path.
174        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    /// Adds `ratio` tokens, called once per completed get (hedged or not),
189    /// so under sustained slowness hedging settles at `ratio` of traffic.
190    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/// The sibling handle a [HedgedBlob] runs hedge requests on, produced by
212/// [crate::cfg::open_hedge_sibling].
213#[derive(Debug)]
214pub enum HedgeSibling {
215    /// A handle onto the same durable store with fully separate connection
216    /// state, kept warm by the wrapper.
217    Isolated(Arc<dyn Blob>),
218    /// The backend has no connection state to isolate (or a second open
219    /// would observe a different store): hedge on the primary instance
220    /// itself, with nothing to warm.
221    SharedWithPrimary,
222    /// Opening the sibling failed. Hedging stays unavailable until a later
223    /// attempt succeeds (see [HedgedBlob::new_arming]).
224    Unavailable,
225}
226
227/// Opens the sibling handle for a store. [HedgedBlob::new_arming] calls it
228/// until it returns something other than [HedgeSibling::Unavailable], and
229/// relies on it to log the cause of each failed attempt. See
230/// [crate::cfg::open_hedge_sibling].
231pub type HedgeSiblingOpener =
232    Box<dyn Fn() -> Pin<Box<dyn Future<Output = HedgeSibling> + Send>> + Send>;
233
234/// A [Blob] decorator that hedges slow `get` requests, per the module docs.
235#[derive(Debug)]
236pub struct HedgedBlob {
237    primary: Arc<dyn Blob>,
238    /// The handle hedge requests run on. Empty = hedging unavailable,
239    /// visible as `hedges_skipped{reason="unavailable"}`. Shared with the
240    /// sibling task, which installs the handle once the sibling is open.
241    hedge: Arc<OnceLock<Arc<dyn Blob>>>,
242    cfg: Arc<ConfigSet>,
243    metrics: BlobHedgeMetrics,
244    budget: HedgeBudget,
245    /// Opens the sibling, then keeps it warm. Aborted on drop. Only the
246    /// test constructor `new` leaves it out, when there is nothing to warm.
247    _sibling_task: Option<AbortOnDropHandle<()>>,
248}
249
250/// How often an idle sibling loop re-reads the dyncfgs, so a flip takes
251/// effect without a restart: the warm interval, or its default while
252/// warming is set to 0.
253fn 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
260/// Keeps the sibling's connection pool warm with periodic concurrent
261/// liveness gets, while hedging is enabled. A cold hedge can stall up to
262/// the connect timeout, which during correlated connection events is
263/// exactly when it must not. While hedging is disabled the warmer idles
264/// and the sibling sees no traffic at all, so a freshly enabled flag can
265/// find a cold pool for up to one warm interval plus a handshake. Hedges
266/// in that window are no better than no hedge, and never worse except for
267/// the grace window's delay on a late primary error.
268async 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            // Nothing to keep warm.
273            tokio::time::sleep(recheck_interval(&cfg)).await;
274            continue;
275        }
276        // Ping first, sleep second, so the pool is warm as soon as the
277        // sibling is armed. As many concurrent pings as hedges can run at once
278        // (HTTP/1.1 allows one in-flight request per connection, so N
279        // concurrent pings force N warm sockets).
280        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        // Bound the cycle: an unbounded hung ping would block warming
285        // past hyper's pool idle eviction, going cold exactly during the
286        // correlated events warming exists for. The timeout also drops
287        // the hung request, which closes its dying socket.
288        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                // A failing or hung warm path means hedges cannot be
294                // trusted to be fast. Surface it, and do not update the
295                // gauge: a fast failure must not report as a fast
296                // healthy path.
297                metrics.warm_errors.inc();
298            }
299        }
300        tokio::time::sleep(interval).await;
301    }
302}
303
304/// Installs the sibling into `slot` once `opener` yields one, then keeps it
305/// warm. See [HedgedBlob::new_arming] for the retry contract.
306async 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            // The opener logs the cause of each failure.
324            HedgeSibling::Unavailable => {
325                tokio::time::sleep(backoff).await;
326                backoff = (backoff * 2).min(SIBLING_OPEN_RETRY_MAX_BACKOFF);
327                // Retries are sibling traffic, so they wait for `enabled`
328                // like the warmer. The first attempt above is not gated, so
329                // that a process is usually armed before hedging is enabled
330                // instead of one recheck after it.
331                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    /// Returns a new [HedgedBlob] from an already-open sibling, for tests that
356    /// need arming to be synchronous.
357    ///
358    /// Must be called from within a tokio runtime: it spawns the sibling
359    /// warming task.
360    #[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    /// Returns a new [HedgedBlob] whose sibling is opened in the background.
396    ///
397    /// `opener` is called from a spawned task, so this constructor never waits
398    /// on the blob store. While it returns [HedgeSibling::Unavailable] the task
399    /// retries with backoff, but only while hedging is enabled, so a disabled
400    /// process makes exactly one attempt until the flag flips. Until an attempt
401    /// succeeds gets are not hedged, visible as
402    /// `hedges_skipped{reason="unavailable"}` and `hedge_armed` at 0.
403    ///
404    /// Dropping the returned [HedgedBlob] aborts the task, including an open
405    /// in progress, so `opener` must be safe to cancel.
406    ///
407    /// Must be called from within a tokio runtime.
408    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    /// Returns the sibling handle and a budget guard, or `None` (having
437    /// already recorded why) if this get must not hedge.
438    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        // NOTE: `Timeout` polls the wrapped future before checking the
472        // deadline, so a primary that is ready exactly at the delay
473        // boundary wins here without firing a hedge. A deadline-first
474        // combinator would not be incorrect, just wasteful: it would fire
475        // a redundant hedge (spending budget and skewing metrics) whenever
476        // the primary completes right at the boundary.
477        if let Ok(res) = tokio::time::timeout(delay, primary.as_mut()).await {
478            // A fast error is returned verbatim: hedging targets hangs,
479            // not failures.
480            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        // The losing future is dropped, which cancels the request in flight
488        // (there is no task boundary between here and the backend). An error
489        // on one leg does not end the race: the slow leg is expected to be a
490        // hung request, and a fast-failing hedge must not convert a get that
491        // was about to succeed into an error.
492        // NOTE: `select` polls its first argument first, so a primary that
493        // is ready simultaneously with the hedge is never miscredited as a
494        // hedge win, which matters because hedges_won is the detection
495        // signal that replaces the suppressed timeout counters (see the
496        // module doc). tokio::select! does NOT have this property unless
497        // marked `biased`.
498        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                // The primary failed after the hedge fired. Give the hedge
506                // a bounded grace window before returning the error: the
507                // window is only there to let an already-healthy hedge win,
508                // which takes about one round trip; reusing the hedge delay
509                // as its length caps the added latency of every branch at
510                // one delay (see the module doc). Beyond that, returning the
511                // primary's error is the better move: it is the un-hedged
512                // outcome, and the callers that wrap gets in retry_external
513                // recover this failure class promptly on a fresh
514                // connection, while an unbounded wait would gamble that
515                // recovery on the hedge leg's health, holding the get and
516                // its hedge slot for up to the blob client's per-attempt
517                // timeout when both legs are unhealthy. The guard stays
518                // held across the window on purpose: the hedge is still in
519                // flight, so the slot still bounds real work (contrast the
520                // hedge-error branch below).
521                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                        // Do not attach the hedge error as context (the
530                        // warning above records it): the error surface
531                        // must not depend on whether a hedge fired (see
532                        // the module doc), and hedge text mentioning
533                        // timeouts could make a string-matching consumer
534                        // like ExternalError::is_timeout misclassify a
535                        // non-timeout error.
536                        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                // The hedge leg is gone, so the concurrency slot no longer
548                // bounds any in-flight memory. Release it rather than
549                // pinning it for the primary's remaining hang, which could
550                // starve other gets of their hedges during exactly the
551                // events hedging exists for.
552                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    /// A test [Blob] whose `get` sleeps a fixed delay and then returns a
606    /// fixed outcome, counting calls.
607    #[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        // The hedge's value won, and it won at exactly the hedge delay, not
698        // at the primary's 3600s: the primary was cancelled while pending.
699        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        // First success wins, not first completion: the hedge fails fast at
723        // the 2s mark but the primary's later success is returned.
724        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        // Both legs fail with the primary failing first: the hedge's error
742        // within the grace window is counted, the primary's error returned.
743        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        // The primary fails after the hedge fired, and the hedge is slow:
755        // the get returns the primary's error after a bounded extra wait
756        // instead of holding on the hedge indefinitely.
757        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        // Primary error at 3s plus the delay-sized grace window.
764        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        // Dropping a hedged get mid-race must release the concurrency slot,
770        // else abandoned gets would permanently disable hedging.
771        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        // Both legs fail with the hedge failing first: the get falls back to
786        // awaiting the primary and returns the primary's error.
787        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        // No refill, so the bucket only ever drains.
846        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        // Turn refill up to one token per completed get. The next get still
856        // finds an empty bucket (refill lands at completion), the one after
857        // hedges again.
858        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        // The warmer pings immediately at start, then every 20s, holding as
887        // many sockets as hedges can run at once.
888        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        // Disabled: the warmer idles, the sibling sees no traffic.
904        tokio::time::sleep(SECS(120)).await;
905        tokio::task::yield_now().await;
906        assert_eq!(hedge.gets.load(Ordering::SeqCst), 0);
907        // Enabling at runtime starts warming within one warm interval.
908        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    /// An opener that reports the sibling unavailable `failures` times before
934    /// returning `hedge`, counting attempts.
935    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        // Unarmed: a slow get is not hedged, and says so.
973        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        // Backoff 1s, 2s, 4s: the fourth attempt succeeds at 7s, and the get
977        // above already spent 3s of it.
978        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        // Armed: the same slow get is hedged and the hedge wins.
983        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        // The arming task continued into the warm loop: one ping per socket
987        // at install time, plus the winning hedge.
988        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        // Nothing to warm: the primary sees no pings.
1001        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        // Disabled: the first attempt runs and fails, then nothing.
1019        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        // Enabling at runtime resumes the retries within one recheck.
1024        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    /// A test [Blob] that delays gets so the hedge (delay 0) fires and wins
1034    /// on every get in the conformance run below. Non-get methods pass
1035    /// through undelayed.
1036    #[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    /// Runs the full [Blob] conformance suite with a hedge racing on every
1068    /// single get: the primary's gets are artificially delayed while the
1069    /// hedge reads the same underlying store undelayed, so the hedge fires
1070    /// and wins throughout.
1071    #[mz_ore::test(tokio::test)]
1072    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait is not yet implemented
1073    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(&registry);
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        // Held by the opener and by its one in-flight open, which never ends.
1129        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}