Skip to main content

opendal_layer_retry/
lib.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18#![doc = include_str!("../README.md")]
19#![cfg_attr(docsrs, feature(doc_cfg))]
20#![cfg_attr(docsrs, doc(auto_cfg))]
21#![deny(missing_docs)]
22use std::fmt::Debug;
23use std::sync::Arc;
24
25use backon::BlockingRetryable;
26use backon::ExponentialBuilder;
27use backon::Retryable;
28use opendal_core::raw::*;
29use opendal_core::*;
30
31/// `RetryLayer` retries temporarily failed operations.
32///
33/// # Notes
34///
35/// This layer retries an operation when [`Error::is_temporary`] returns `true`.
36/// If the operation still fails, the layer marks the error as `Persistent` to
37/// indicate that retries did not resolve it.
38///
39/// # Stateful operation bodies
40///
41/// While retrying stateful operation bodies, please make sure either:
42///
43/// - All futures generated by the body methods are resolved to `Ready`.
44/// - Or, no methods on the same body are called after retry returns a final error.
45///
46/// Otherwise, `RetryLayer` can lose the inner body state and fail subsequent calls with an
47/// `Unexpected` error.
48///
49/// For example, while composing `RetryLayer` with `TimeoutLayer`. The order of layer is sensitive.
50///
51/// ```no_run
52/// # use std::time::Duration;
53/// #
54/// # use opendal_core::services;
55/// # use opendal_core::Operator;
56/// # use opendal_core::Result;
57/// # use opendal_layer_retry::RetryLayer;
58/// # use opendal_layer_timeout::TimeoutLayer;
59/// #
60/// # fn main() -> Result<()> {
61/// let op = Operator::new(services::Memory::default())?
62///     // This is fine, since timeout happens during retry.
63///     .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)))
64///     .layer(RetryLayer::default())
65///     // This is wrong. Timeout layer can drop the retry future before it restores body state.
66///     .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)));
67/// # Ok(())
68/// # }
69/// ```
70///
71/// # Examples
72///
73/// ```no_run
74/// # use opendal_core::services;
75/// # use opendal_core::Operator;
76/// # use opendal_core::Result;
77/// # use opendal_layer_retry::RetryLayer;
78/// #
79/// # fn main() -> Result<()> {
80/// let _ = Operator::new(services::Memory::default())?
81///     .layer(RetryLayer::default());
82/// # Ok(())
83/// # }
84/// ```
85///
86/// ## Customize retry interceptor
87///
88/// RetryLayer accepts [`RetryInterceptor`] to allow users to customize
89/// their own retry interceptor logic.
90///
91/// ```no_run
92/// # use std::time::Duration;
93/// #
94/// # use opendal_core::services;
95/// # use opendal_core::Error;
96/// # use opendal_core::Operator;
97/// # use opendal_core::Result;
98/// # use opendal_layer_retry::RetryEvent;
99/// # use opendal_layer_retry::RetryInterceptor;
100/// # use opendal_layer_retry::RetryLayer;
101/// #
102/// struct MyRetryInterceptor;
103///
104/// impl RetryInterceptor for MyRetryInterceptor {
105///     fn intercept(&self, event: RetryEvent<'_>) {
106///         // do something with event.op, event.err, event.retry_after, event.attempt
107///     }
108/// }
109///
110/// # fn main() -> Result<()> {
111/// let _ = Operator::new(services::Memory::default())?
112///     .layer(RetryLayer::default().with_notify(MyRetryInterceptor));
113/// # Ok(())
114/// # }
115/// ```
116pub struct RetryLayer<I: RetryInterceptor = DefaultRetryInterceptor> {
117    builder: ExponentialBuilder,
118    notify: Arc<I>,
119}
120
121impl<I: RetryInterceptor> Debug for RetryLayer<I> {
122    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
123        f.debug_struct("RetryLayer")
124            .field("builder", &self.builder)
125            .finish_non_exhaustive()
126    }
127}
128
129impl<I: RetryInterceptor> Clone for RetryLayer<I> {
130    fn clone(&self) -> Self {
131        Self {
132            builder: self.builder,
133            notify: self.notify.clone(),
134        }
135    }
136}
137
138impl Default for RetryLayer {
139    fn default() -> Self {
140        Self {
141            builder: ExponentialBuilder::default(),
142            notify: Arc::new(DefaultRetryInterceptor),
143        }
144    }
145}
146
147impl RetryLayer {
148    /// Create a new [`RetryLayer`].
149    pub fn new() -> RetryLayer {
150        Self::default()
151    }
152}
153
154impl<I: RetryInterceptor> RetryLayer<I> {
155    /// Set the retry interceptor as new notify.
156    ///
157    /// ```no_run
158    /// use opendal_core::services;
159    /// use opendal_core::Operator;
160    /// use opendal_layer_retry::RetryLayer;
161    ///
162    /// fn notify(_event: opendal_layer_retry::RetryEvent<'_>) {}
163    ///
164    /// let _ = Operator::new(services::Memory::default())
165    ///     .expect("must init")
166    ///     .layer(RetryLayer::default().with_notify(notify));
167    /// ```
168    pub fn with_notify<NI: RetryInterceptor>(self, notify: NI) -> RetryLayer<NI> {
169        RetryLayer {
170            builder: self.builder,
171            notify: Arc::new(notify),
172        }
173    }
174
175    /// Set jitter of current backoff.
176    ///
177    /// If jitter is enabled, ExponentialBackoff will add a random jitter in `[0, min_delay)
178    /// to current delay.
179    pub fn with_jitter(mut self) -> Self {
180        self.builder = self.builder.with_jitter();
181        self
182    }
183
184    /// Set factor of current backoff.
185    ///
186    /// # Panics
187    ///
188    /// This function will panic if input factor smaller than `1.0`.
189    pub fn with_factor(mut self, factor: f32) -> Self {
190        self.builder = self.builder.with_factor(factor);
191        self
192    }
193
194    /// Set min_delay of current backoff.
195    pub fn with_min_delay(mut self, min_delay: Duration) -> Self {
196        self.builder = self.builder.with_min_delay(min_delay);
197        self
198    }
199
200    /// Set max_delay of current backoff.
201    ///
202    /// Delay will not increase if current delay is larger than max_delay.
203    pub fn with_max_delay(mut self, max_delay: Duration) -> Self {
204        self.builder = self.builder.with_max_delay(max_delay);
205        self
206    }
207
208    /// Set max_times of current backoff.
209    ///
210    /// Backoff will return `None` if max times is reaching.
211    pub fn with_max_times(mut self, max_times: usize) -> Self {
212        self.builder = self.builder.with_max_times(max_times);
213        self
214    }
215}
216
217impl<I: RetryInterceptor> Layer for RetryLayer<I> {
218    fn apply_service(&self, inner: Servicer) -> Servicer {
219        Arc::new(self.layer(inner))
220    }
221}
222
223impl<I: RetryInterceptor> RetryLayer<I> {
224    fn layer(&self, inner: Servicer) -> RetryService<I> {
225        RetryService {
226            inner,
227            notify: self.notify.clone(),
228            builder: self.builder,
229        }
230    }
231}
232
233/// Context passed to [`RetryInterceptor`] before each retry sleep.
234#[non_exhaustive]
235#[derive(Debug)]
236pub struct RetryEvent<'a> {
237    /// The operation being retried.
238    pub op: Operation,
239    /// The error that triggered the retry.
240    pub err: &'a Error,
241    /// The duration to wait before the next retry attempt.
242    pub retry_after: Duration,
243    /// 1-based retry attempt number.
244    pub attempt: u32,
245}
246
247/// RetryInterceptor observes retry attempts before the retry sleep.
248pub trait RetryInterceptor: Send + Sync + 'static {
249    /// Called before each retry sleep.
250    ///
251    /// # Notes
252    ///
253    /// The intercept must be quick and non-blocking. No heavy IO is
254    /// allowed. Otherwise, the retry will be blocked.
255    fn intercept(&self, event: RetryEvent<'_>);
256}
257
258impl<F> RetryInterceptor for F
259where
260    F: for<'a> Fn(RetryEvent<'a>) + Send + Sync + 'static,
261{
262    fn intercept(&self, event: RetryEvent<'_>) {
263        self(event);
264    }
265}
266
267/// The DefaultRetryInterceptor logs each retry error at warn level.
268pub struct DefaultRetryInterceptor;
269
270impl RetryInterceptor for DefaultRetryInterceptor {
271    fn intercept(&self, event: RetryEvent<'_>) {
272        log::warn!(
273            target: "opendal::layers::retry",
274            "will retry {:?} (attempt {}) after {}s because: {:?}",
275            event.op, event.attempt, event.retry_after.as_secs_f64(), event.err
276        );
277    }
278}
279
280#[doc(hidden)]
281pub struct RetryService<I: RetryInterceptor> {
282    inner: Servicer,
283    notify: Arc<I>,
284    builder: ExponentialBuilder,
285}
286
287impl<I: RetryInterceptor> Debug for RetryService<I> {
288    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
289        f.debug_struct("RetryService")
290            .field("inner", &self.inner)
291            .finish_non_exhaustive()
292    }
293}
294
295impl<I: RetryInterceptor> Service for RetryService<I> {
296    type Reader = RetryReader<oio::Reader, I>;
297    type Writer = RetryWrapper<oio::Writer, I>;
298    type Lister = RetryWrapper<oio::Lister, I>;
299    type Deleter = RetryWrapper<oio::Deleter, I>;
300    type Copier = RetryWrapper<oio::Copier, I>;
301    type Composer = oio::Composer;
302
303    fn info(&self) -> ServiceInfo {
304        self.inner.info()
305    }
306
307    fn capability(&self) -> Capability {
308        self.inner.capability()
309    }
310
311    fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
312        self.inner.compose(ctx, to, args)
313    }
314
315    async fn create_dir(
316        &self,
317        ctx: &OperationContext,
318        path: &str,
319        args: OpCreateDir,
320    ) -> Result<RpCreateDir> {
321        let mut attempt: u32 = 0;
322        { || self.inner.create_dir(ctx, path, args.clone()) }
323            .retry(self.builder)
324            .when(|e| e.is_temporary())
325            .notify(|err, dur| {
326                attempt += 1;
327                self.notify.intercept(RetryEvent {
328                    op: Operation::CreateDir,
329                    err,
330                    retry_after: dur,
331                    attempt,
332                })
333            })
334            .await
335            .map_err(|err| err.set_persistent())
336    }
337
338    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
339        let mut attempt: u32 = 0;
340        let reader = { || self.inner.read(ctx, path, args.clone()) }
341            .retry(self.builder)
342            .when(|e| e.is_temporary())
343            .notify(|err, dur| {
344                attempt += 1;
345                self.notify.intercept(RetryEvent {
346                    op: Operation::Read,
347                    err,
348                    retry_after: dur,
349                    attempt,
350                })
351            })
352            .call()
353            .map_err(|err| err.set_persistent())?;
354
355        Ok(RetryReader::new(reader, self.notify.clone(), self.builder))
356    }
357
358    fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
359        let mut attempt: u32 = 0;
360        let writer = { || self.inner.write(ctx, path, args.clone()) }
361            .retry(self.builder)
362            .when(|e| e.is_temporary())
363            .notify(|err, dur| {
364                attempt += 1;
365                self.notify.intercept(RetryEvent {
366                    op: Operation::Write,
367                    err,
368                    retry_after: dur,
369                    attempt,
370                })
371            })
372            .call()
373            .map_err(|err| err.set_persistent())?;
374
375        Ok(RetryWrapper::new(writer, self.notify.clone(), self.builder))
376    }
377
378    async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
379        let mut attempt: u32 = 0;
380        { || self.inner.stat(ctx, path, args.clone()) }
381            .retry(self.builder)
382            .when(|e| e.is_temporary())
383            .notify(|err, dur| {
384                attempt += 1;
385                self.notify.intercept(RetryEvent {
386                    op: Operation::Stat,
387                    err,
388                    retry_after: dur,
389                    attempt,
390                })
391            })
392            .await
393            .map_err(|err| err.set_persistent())
394    }
395
396    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
397        let mut attempt: u32 = 0;
398        let deleter = { || self.inner.delete(ctx) }
399            .retry(self.builder)
400            .when(|e| e.is_temporary())
401            .notify(|err, dur| {
402                attempt += 1;
403                self.notify.intercept(RetryEvent {
404                    op: Operation::Delete,
405                    err,
406                    retry_after: dur,
407                    attempt,
408                })
409            })
410            .call()
411            .map_err(|err| err.set_persistent())?;
412
413        Ok(RetryWrapper::new(
414            deleter,
415            self.notify.clone(),
416            self.builder,
417        ))
418    }
419
420    fn copy(
421        &self,
422        ctx: &OperationContext,
423        from: &str,
424        to: &str,
425        args: OpCopy,
426    ) -> Result<Self::Copier> {
427        let mut attempt: u32 = 0;
428        let copier = { || self.inner.copy(ctx, from, to, args.clone()) }
429            .retry(self.builder)
430            .when(|e| e.is_temporary())
431            .notify(|err, dur| {
432                attempt += 1;
433                self.notify.intercept(RetryEvent {
434                    op: Operation::Copy,
435                    err,
436                    retry_after: dur,
437                    attempt,
438                })
439            })
440            .call()
441            .map_err(|err| err.set_persistent())?;
442
443        Ok(RetryWrapper::new(copier, self.notify.clone(), self.builder))
444    }
445
446    async fn rename(
447        &self,
448        ctx: &OperationContext,
449        from: &str,
450        to: &str,
451        args: OpRename,
452    ) -> Result<RpRename> {
453        let mut attempt: u32 = 0;
454        { || self.inner.rename(ctx, from, to, args.clone()) }
455            .retry(self.builder)
456            .when(|e| e.is_temporary())
457            .notify(|err, dur| {
458                attempt += 1;
459                self.notify.intercept(RetryEvent {
460                    op: Operation::Rename,
461                    err,
462                    retry_after: dur,
463                    attempt,
464                })
465            })
466            .await
467            .map_err(|err| err.set_persistent())
468    }
469
470    async fn restore(
471        &self,
472        ctx: &OperationContext,
473        path: &str,
474        args: OpRestore,
475    ) -> Result<RpRestore> {
476        let mut attempt: u32 = 0;
477        { || self.inner.restore(ctx, path, args.clone()) }
478            .retry(self.builder)
479            .when(|e| e.is_temporary())
480            .notify(|err, dur| {
481                attempt += 1;
482                self.notify.intercept(RetryEvent {
483                    op: Operation::Restore,
484                    err,
485                    retry_after: dur,
486                    attempt,
487                })
488            })
489            .await
490            .map_err(|err| err.set_persistent())
491    }
492
493    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
494        let mut attempt: u32 = 0;
495        let lister = { || self.inner.list(ctx, path, args.clone()) }
496            .retry(self.builder)
497            .when(|e| e.is_temporary())
498            .notify(|err, dur| {
499                attempt += 1;
500                self.notify.intercept(RetryEvent {
501                    op: Operation::List,
502                    err,
503                    retry_after: dur,
504                    attempt,
505                })
506            })
507            .call()
508            .map_err(|err| err.set_persistent())?;
509
510        Ok(RetryWrapper::new(lister, self.notify.clone(), self.builder))
511    }
512
513    async fn presign(
514        &self,
515        ctx: &OperationContext,
516        path: &str,
517        args: OpPresign,
518    ) -> Result<RpPresign> {
519        let mut attempt: u32 = 0;
520        { || self.inner.presign(ctx, path, args.clone()) }
521            .retry(self.builder)
522            .when(|e| e.is_temporary())
523            .notify(|err, dur| {
524                attempt += 1;
525                self.notify.intercept(RetryEvent {
526                    op: Operation::Presign,
527                    err,
528                    retry_after: dur,
529                    attempt,
530                })
531            })
532            .await
533            .map_err(|err| err.set_persistent())
534    }
535}
536
537#[doc(hidden)]
538pub struct RetryReader<R, I> {
539    inner: Arc<R>,
540    notify: Arc<I>,
541    builder: ExponentialBuilder,
542}
543
544impl<R, I> RetryReader<R, I> {
545    fn new(inner: R, notify: Arc<I>, builder: ExponentialBuilder) -> Self {
546        Self {
547            inner: Arc::new(inner),
548            notify,
549            builder,
550        }
551    }
552}
553
554impl<R: oio::Read + 'static, I: RetryInterceptor> oio::Read for RetryReader<R, I> {
555    async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
556        use backon::Retryable;
557
558        let mut attempt: u32 = 0;
559        let (rp, stream) = { || self.inner.open(range) }
560            .retry(self.builder)
561            .when(|e| e.is_temporary())
562            .notify(|err, dur| {
563                attempt += 1;
564                self.notify.intercept(RetryEvent {
565                    op: Operation::Read,
566                    err,
567                    retry_after: dur,
568                    attempt,
569                })
570            })
571            .await
572            .map_err(|e| e.set_persistent())?;
573
574        Ok((
575            rp,
576            Box::new(RetryReadStream::new(
577                self.inner.clone(),
578                stream,
579                range,
580                self.notify.clone(),
581                self.builder,
582            )) as Box<dyn oio::ReadStreamDyn>,
583        ))
584    }
585
586    async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
587        use backon::Retryable;
588
589        let mut attempt: u32 = 0;
590        { || self.inner.read(range) }
591            .retry(self.builder)
592            .when(|e| e.is_temporary())
593            .notify(|err, dur| {
594                attempt += 1;
595                self.notify.intercept(RetryEvent {
596                    op: Operation::Read,
597                    err,
598                    retry_after: dur,
599                    attempt,
600                })
601            })
602            .await
603            .map_err(|e| e.set_persistent())
604    }
605}
606
607#[doc(hidden)]
608pub struct RetryReadStream<R, I> {
609    reader: Arc<R>,
610    stream: Option<Box<dyn oio::ReadStreamDyn>>,
611    range: BytesRange,
612    read: u64,
613    notify: Arc<I>,
614    builder: ExponentialBuilder,
615}
616
617impl<R, I> RetryReadStream<R, I> {
618    fn new(
619        reader: Arc<R>,
620        stream: Box<dyn oio::ReadStreamDyn>,
621        range: BytesRange,
622        notify: Arc<I>,
623        builder: ExponentialBuilder,
624    ) -> Self {
625        Self {
626            reader,
627            stream: Some(stream),
628            range,
629            read: 0,
630            notify,
631            builder,
632        }
633    }
634}
635
636impl<R: oio::Read, I: RetryInterceptor> oio::ReadStream for RetryReadStream<R, I> {
637    async fn read(&mut self) -> Result<Buffer> {
638        use backon::RetryableWithContext;
639
640        let reader = self.reader.clone();
641        let stream = self.stream.take();
642        let range = self.range;
643        let read = self.read;
644        let mut attempt: u32 = 0;
645
646        let ((stream, range, read), res) = {
647            |(stream, mut range, mut read): (
648                Option<Box<dyn oio::ReadStreamDyn>>,
649                BytesRange,
650                u64,
651            )| {
652                let reader = reader.clone();
653                async move {
654                    let mut stream = match stream {
655                        Some(stream) => stream,
656                        None => {
657                            range.advance(read);
658                            read = 0;
659
660                            match reader.open(range).await {
661                                Ok((_, stream)) => stream,
662                                Err(err) => return ((None, range, read), Err(err)),
663                            }
664                        }
665                    };
666
667                    let res = match stream.read().await {
668                        Ok(buf) => {
669                            if !buf.is_empty() {
670                                read += buf.len() as u64;
671                            }
672                            (Some(stream), Ok(buf))
673                        }
674                        Err(err) => (None, Err(err)),
675                    };
676
677                    ((res.0, range, read), res.1)
678                }
679            }
680        }
681        .retry(self.builder)
682        .when(|e| e.is_temporary())
683        .context((stream, range, read))
684        .notify(|err, dur| {
685            attempt += 1;
686            self.notify.intercept(RetryEvent {
687                op: Operation::Read,
688                err,
689                retry_after: dur,
690                attempt,
691            })
692        })
693        .await;
694
695        self.stream = stream;
696        self.range = range;
697        self.read = read;
698
699        res.map_err(|err| err.set_persistent())
700    }
701}
702
703#[doc(hidden)]
704pub struct RetryWrapper<R, I> {
705    inner: Option<R>,
706    notify: Arc<I>,
707
708    builder: ExponentialBuilder,
709}
710
711impl<R, I> RetryWrapper<R, I> {
712    fn new(inner: R, notify: Arc<I>, backoff: ExponentialBuilder) -> Self {
713        Self {
714            inner: Some(inner),
715            notify,
716            builder: backoff,
717        }
718    }
719
720    fn take_inner(&mut self) -> Result<R> {
721        self.inner.take().ok_or_else(|| {
722            Error::new(
723                ErrorKind::Unexpected,
724                "retry layer is in bad state, please make sure future not dropped before ready",
725            )
726        })
727    }
728}
729
730impl<R: oio::ReadStream, I: RetryInterceptor> oio::ReadStream for RetryWrapper<R, I> {
731    async fn read(&mut self) -> Result<Buffer> {
732        use backon::RetryableWithContext;
733
734        let inner = self.take_inner()?;
735        let mut attempt: u32 = 0;
736
737        let (inner, res) = {
738            |mut r: R| async move {
739                let res = r.read().await;
740
741                (r, res)
742            }
743        }
744        .retry(self.builder)
745        .when(|e| e.is_temporary())
746        .context(inner)
747        .notify(|err, dur| {
748            attempt += 1;
749            self.notify.intercept(RetryEvent {
750                op: Operation::Read,
751                err,
752                retry_after: dur,
753                attempt,
754            })
755        })
756        .await;
757
758        self.inner = Some(inner);
759        res.map_err(|err| err.set_persistent())
760    }
761}
762
763impl<R: oio::Write, I: RetryInterceptor> oio::Write for RetryWrapper<R, I> {
764    async fn write(&mut self, bs: Buffer) -> Result<()> {
765        use backon::RetryableWithContext;
766
767        let inner = self.take_inner()?;
768        let mut attempt: u32 = 0;
769
770        let ((inner, _), res) = {
771            |(mut r, bs): (R, Buffer)| async move {
772                let res = r.write(bs.clone()).await;
773
774                ((r, bs), res)
775            }
776        }
777        .retry(self.builder)
778        .when(|e| e.is_temporary())
779        .context((inner, bs))
780        .notify(|err, dur| {
781            attempt += 1;
782            self.notify.intercept(RetryEvent {
783                op: Operation::Write,
784                err,
785                retry_after: dur,
786                attempt,
787            })
788        })
789        .await;
790
791        self.inner = Some(inner);
792        res.map_err(|err| err.set_persistent())
793    }
794
795    async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
796        use backon::RetryableWithContext;
797
798        let inner = self.take_inner()?;
799        let path = path.to_string();
800        let mut attempt: u32 = 0;
801
802        let ((inner, _, _, _), res) = {
803            |(mut r, path, args, range): (R, String, OpRead, BytesRange)| async move {
804                let res = r.copy_from(&path, args.clone(), range).await;
805
806                ((r, path, args, range), res)
807            }
808        }
809        .retry(self.builder)
810        .when(|e| e.is_temporary())
811        .context((inner, path, args, range))
812        .notify(|err, dur| {
813            attempt += 1;
814            self.notify.intercept(RetryEvent {
815                op: Operation::Write,
816                err,
817                retry_after: dur,
818                attempt,
819            })
820        })
821        .await;
822
823        self.inner = Some(inner);
824        res.map_err(|err| err.set_persistent())
825    }
826
827    async fn abort(&mut self) -> Result<()> {
828        use backon::RetryableWithContext;
829
830        let inner = self.take_inner()?;
831        let mut attempt: u32 = 0;
832
833        let (inner, res) = {
834            |mut r: R| async move {
835                let res = r.abort().await;
836
837                (r, res)
838            }
839        }
840        .retry(self.builder)
841        .when(|e| e.is_temporary())
842        .context(inner)
843        .notify(|err, dur| {
844            attempt += 1;
845            self.notify.intercept(RetryEvent {
846                op: Operation::Write,
847                err,
848                retry_after: dur,
849                attempt,
850            })
851        })
852        .await;
853
854        self.inner = Some(inner);
855        res.map_err(|err| err.set_persistent())
856    }
857
858    async fn close(&mut self) -> Result<Metadata> {
859        use backon::RetryableWithContext;
860
861        let inner = self.take_inner()?;
862        let mut attempt: u32 = 0;
863
864        let (inner, res) = {
865            |mut r: R| async move {
866                let res = r.close().await;
867
868                (r, res)
869            }
870        }
871        .retry(self.builder)
872        .when(|e| e.is_temporary())
873        .context(inner)
874        .notify(|err, dur| {
875            attempt += 1;
876            self.notify.intercept(RetryEvent {
877                op: Operation::Write,
878                err,
879                retry_after: dur,
880                attempt,
881            })
882        })
883        .await;
884
885        self.inner = Some(inner);
886        res.map_err(|err| err.set_persistent())
887    }
888}
889
890impl<P: oio::List, I: RetryInterceptor> oio::List for RetryWrapper<P, I> {
891    async fn next(&mut self) -> Result<Option<oio::Entry>> {
892        use backon::RetryableWithContext;
893
894        let inner = self.take_inner()?;
895        let mut attempt: u32 = 0;
896
897        let (inner, res) = {
898            |mut p: P| async move {
899                let res = p.next().await;
900
901                (p, res)
902            }
903        }
904        .retry(self.builder)
905        .when(|e| e.is_temporary())
906        .context(inner)
907        .notify(|err, dur| {
908            attempt += 1;
909            self.notify.intercept(RetryEvent {
910                op: Operation::List,
911                err,
912                retry_after: dur,
913                attempt,
914            })
915        })
916        .await;
917
918        self.inner = Some(inner);
919        res.map_err(|err| err.set_persistent())
920    }
921}
922
923impl<P: oio::Delete, I: RetryInterceptor> oio::Delete for RetryWrapper<P, I> {
924    async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
925        use backon::RetryableWithContext;
926
927        let inner = self.take_inner()?;
928        let path = path.to_string();
929        let args_cloned = args.clone();
930        let mut attempt: u32 = 0;
931
932        let (inner, res) = {
933            |mut p: P| {
934                let path = path.clone();
935                let args = args_cloned.clone();
936                async move {
937                    let res = p.delete(&path, args).await;
938                    (p, res)
939                }
940            }
941        }
942        .retry(self.builder)
943        .when(|e| e.is_temporary())
944        .context(inner)
945        .notify(|err, dur| {
946            attempt += 1;
947            self.notify.intercept(RetryEvent {
948                op: Operation::Delete,
949                err,
950                retry_after: dur,
951                attempt,
952            });
953        })
954        .await;
955
956        self.inner = Some(inner);
957        res.map_err(|e| e.set_persistent())
958    }
959
960    async fn close(&mut self) -> Result<()> {
961        use backon::RetryableWithContext;
962
963        let inner = self.take_inner()?;
964        let mut attempt: u32 = 0;
965
966        let (inner, res) = {
967            |mut p: P| async move {
968                let res = p.close().await;
969
970                (p, res)
971            }
972        }
973        .retry(self.builder)
974        .when(|e| e.is_temporary())
975        .context(inner)
976        .notify(|err, dur| {
977            attempt += 1;
978            self.notify.intercept(RetryEvent {
979                op: Operation::Delete,
980                err,
981                retry_after: dur,
982                attempt,
983            })
984        })
985        .await;
986
987        self.inner = Some(inner);
988        res.map_err(|err| err.set_persistent())
989    }
990}
991
992impl<C: oio::Copy, I: RetryInterceptor> oio::Copy for RetryWrapper<C, I> {
993    async fn next(&mut self) -> Result<Option<usize>> {
994        use backon::RetryableWithContext;
995
996        let inner = self.take_inner()?;
997        let mut attempt: u32 = 0;
998
999        let (inner, res) = {
1000            |mut c: C| async move {
1001                let res = c.next().await;
1002
1003                (c, res)
1004            }
1005        }
1006        .retry(self.builder)
1007        .when(|e| e.is_temporary())
1008        .context(inner)
1009        .notify(|err, dur| {
1010            attempt += 1;
1011            self.notify.intercept(RetryEvent {
1012                op: Operation::Copy,
1013                err,
1014                retry_after: dur,
1015                attempt,
1016            })
1017        })
1018        .await;
1019
1020        self.inner = Some(inner);
1021        res.map_err(|err| err.set_persistent())
1022    }
1023
1024    async fn close(&mut self) -> Result<Metadata> {
1025        use backon::RetryableWithContext;
1026
1027        let inner = self.take_inner()?;
1028        let mut attempt: u32 = 0;
1029
1030        let (inner, res) = {
1031            |mut c: C| async move {
1032                let res = c.close().await;
1033
1034                (c, res)
1035            }
1036        }
1037        .retry(self.builder)
1038        .when(|e| e.is_temporary())
1039        .context(inner)
1040        .notify(|err, dur| {
1041            attempt += 1;
1042            self.notify.intercept(RetryEvent {
1043                op: Operation::Copy,
1044                err,
1045                retry_after: dur,
1046                attempt,
1047            })
1048        })
1049        .await;
1050
1051        self.inner = Some(inner);
1052        res.map_err(|err| err.set_persistent())
1053    }
1054
1055    async fn abort(&mut self) -> Result<()> {
1056        use backon::RetryableWithContext;
1057
1058        let inner = self.take_inner()?;
1059        let mut attempt: u32 = 0;
1060
1061        let (inner, res) = {
1062            |mut c: C| async move {
1063                let res = c.abort().await;
1064
1065                (c, res)
1066            }
1067        }
1068        .retry(self.builder)
1069        .when(|e| e.is_temporary())
1070        .context(inner)
1071        .notify(|err, dur| {
1072            attempt += 1;
1073            self.notify.intercept(RetryEvent {
1074                op: Operation::Copy,
1075                err,
1076                retry_after: dur,
1077                attempt,
1078            })
1079        })
1080        .await;
1081
1082        self.inner = Some(inner);
1083        res.map_err(|err| err.set_persistent())
1084    }
1085}
1086
1087#[cfg(test)]
1088mod tests {
1089    use std::sync::Mutex;
1090
1091    use bytes::Bytes;
1092    use futures::TryStreamExt;
1093    use futures::stream;
1094    use logforth::append::Testing;
1095    use logforth::filter::rustlog::RustLogFilterBuilder;
1096    use logforth::layout::TextLayout;
1097    use opendal_layer_logging::LoggingLayer;
1098
1099    use super::*;
1100
1101    #[derive(Default, Clone)]
1102    struct MockBuilder {
1103        attempt: Arc<Mutex<usize>>,
1104    }
1105
1106    impl Builder for MockBuilder {
1107        type Config = ();
1108
1109        fn build(self) -> Result<impl Service> {
1110            Ok(MockService {
1111                attempt: self.attempt,
1112            })
1113        }
1114    }
1115
1116    #[derive(Debug, Clone, Default)]
1117    struct MockService {
1118        attempt: Arc<Mutex<usize>>,
1119    }
1120
1121    /// Reader returned by this backend.
1122    pub struct MockReader {
1123        backend: MockService,
1124    }
1125
1126    impl MockReader {
1127        fn new(backend: MockService, _: &str, _: OpRead) -> Self {
1128            Self { backend }
1129        }
1130    }
1131
1132    impl oio::StreamRead for MockReader {
1133        async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
1134            let backend = &self.backend;
1135            let rp = RpRead::new({
1136                let metadata = MetadataBuilder::file(13);
1137                metadata.build()
1138            });
1139            let stream = MockReadStream {
1140                buf: Bytes::from("Hello, World!").into(),
1141                range,
1142                attempt: backend.attempt.clone(),
1143            };
1144
1145            Ok((rp, Box::new(stream) as Box<dyn oio::ReadStreamDyn>))
1146        }
1147    }
1148
1149    impl Service for MockService {
1150        type Reader = oio::StreamReader<MockReader>;
1151        type Writer = MockWriter;
1152        type Lister = MockLister;
1153        type Deleter = MockDeleter;
1154        type Copier = MockCopier;
1155        type Composer = ();
1156
1157        fn info(&self) -> ServiceInfo {
1158            ServiceInfo::with_scheme("mock")
1159        }
1160
1161        fn capability(&self) -> Capability {
1162            Capability {
1163                read: true,
1164                write: true,
1165                write_can_multi: true,
1166                delete: true,
1167                delete_max_size: Some(10),
1168                stat: true,
1169                list: true,
1170                list_with_recursive: true,
1171                copy: true,
1172                ..Default::default()
1173            }
1174        }
1175
1176        async fn create_dir(
1177            &self,
1178            _: &OperationContext,
1179            _: &str,
1180            _: OpCreateDir,
1181        ) -> Result<RpCreateDir> {
1182            Err(Error::new(
1183                ErrorKind::Unsupported,
1184                "operation is not supported",
1185            ))
1186        }
1187
1188        async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
1189            Ok(RpStat::new({
1190                let metadata = MetadataBuilder::file(13);
1191                metadata.build()
1192            }))
1193        }
1194
1195        fn read(&self, _: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
1196            Ok(oio::StreamReader::new(MockReader::new(
1197                self.clone(),
1198                path,
1199                args,
1200            )))
1201        }
1202
1203        fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
1204            Ok(MockDeleter {
1205                size: 0,
1206                attempt: self.attempt.clone(),
1207            })
1208        }
1209
1210        fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
1211            Ok(MockWriter {})
1212        }
1213
1214        fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
1215            let lister = MockLister::default();
1216            Ok(lister)
1217        }
1218
1219        fn copy(&self, _: &OperationContext, _: &str, _: &str, _: OpCopy) -> Result<Self::Copier> {
1220            Ok(MockCopier {
1221                attempt: self.attempt.clone(),
1222            })
1223        }
1224
1225        async fn rename(
1226            &self,
1227            _: &OperationContext,
1228            _: &str,
1229            _: &str,
1230            _: OpRename,
1231        ) -> Result<RpRename> {
1232            Err(Error::new(
1233                ErrorKind::Unsupported,
1234                "operation is not supported",
1235            ))
1236        }
1237
1238        async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
1239            Err(Error::new(
1240                ErrorKind::Unsupported,
1241                "operation is not supported",
1242            ))
1243        }
1244    }
1245
1246    #[derive(Debug, Clone, Default)]
1247    struct MockReadStream {
1248        buf: Buffer,
1249        range: BytesRange,
1250        attempt: Arc<Mutex<usize>>,
1251    }
1252
1253    impl oio::ReadStream for MockReadStream {
1254        async fn read(&mut self) -> Result<Buffer> {
1255            let mut attempt = self.attempt.lock().unwrap();
1256            *attempt += 1;
1257
1258            match *attempt {
1259                1 => Err(
1260                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1261                        .set_temporary(),
1262                ),
1263                2 => Err(
1264                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1265                        .set_temporary(),
1266                ),
1267                // Should read out all data.
1268                3 => Ok(self.buf.slice(self.range.to_range_as_usize())),
1269                4 => Err(
1270                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1271                        .set_temporary(),
1272                ),
1273                // Should be empty.
1274                5 => Ok(self.buf.slice(self.range.to_range_as_usize())),
1275                _ => unreachable!(),
1276            }
1277        }
1278    }
1279
1280    #[derive(Debug, Clone, Default)]
1281    struct MockWriter {}
1282
1283    impl oio::Write for MockWriter {
1284        async fn write(&mut self, _: Buffer) -> Result<()> {
1285            Ok(())
1286        }
1287
1288        async fn close(&mut self) -> Result<Metadata> {
1289            Err(Error::new(ErrorKind::Unexpected, "always close failed").set_temporary())
1290        }
1291
1292        async fn abort(&mut self) -> Result<()> {
1293            Ok(())
1294        }
1295    }
1296
1297    #[derive(Debug, Clone, Default)]
1298    struct MockLister {
1299        attempt: usize,
1300    }
1301
1302    impl oio::List for MockLister {
1303        async fn next(&mut self) -> Result<Option<oio::Entry>> {
1304            self.attempt += 1;
1305            match self.attempt {
1306                1 => Err(Error::new(
1307                    ErrorKind::RateLimited,
1308                    "retryable rate limited error from lister",
1309                )
1310                .set_temporary()),
1311                2 => Ok(Some(oio::Entry::new(
1312                    "hello",
1313                    MetadataBuilder::file(0).build(),
1314                ))),
1315                3 => Ok(Some(oio::Entry::new(
1316                    "world",
1317                    MetadataBuilder::file(0).build(),
1318                ))),
1319                4 => Err(
1320                    Error::new(ErrorKind::Unexpected, "retryable internal server error")
1321                        .set_temporary(),
1322                ),
1323                5 => Ok(Some(oio::Entry::new(
1324                    "2023/",
1325                    MetadataBuilder::dir().build(),
1326                ))),
1327                6 => Ok(Some(oio::Entry::new(
1328                    "0208/",
1329                    MetadataBuilder::dir().build(),
1330                ))),
1331                7 => Ok(None),
1332                _ => {
1333                    unreachable!()
1334                }
1335            }
1336        }
1337    }
1338
1339    #[derive(Debug, Clone, Default)]
1340    struct MockDeleter {
1341        size: usize,
1342        attempt: Arc<Mutex<usize>>,
1343    }
1344
1345    impl oio::Delete for MockDeleter {
1346        async fn delete(&mut self, _: &str, _: OpDelete) -> Result<()> {
1347            self.size += 1;
1348            Ok(())
1349        }
1350
1351        async fn close(&mut self) -> Result<()> {
1352            let mut attempt = self.attempt.lock().unwrap();
1353            *attempt += 1;
1354
1355            match *attempt {
1356                1 => Err(
1357                    Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1358                        .set_temporary(),
1359                ),
1360                2..=4 => {
1361                    self.size = self.size.saturating_sub(1);
1362                    Err(
1363                        Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1364                            .set_temporary(),
1365                    )
1366                }
1367                5 => {
1368                    self.size = self.size.saturating_sub(1);
1369                    if self.size == 0 {
1370                        Ok(())
1371                    } else {
1372                        Err(
1373                            Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1374                                .set_temporary(),
1375                        )
1376                    }
1377                }
1378                _ => unreachable!(),
1379            }
1380        }
1381    }
1382
1383    #[derive(Debug, Clone, Default)]
1384    struct MockCopier {
1385        attempt: Arc<Mutex<usize>>,
1386    }
1387
1388    impl oio::Copy for MockCopier {
1389        async fn next(&mut self) -> Result<Option<usize>> {
1390            let mut attempt = self.attempt.lock().unwrap();
1391            *attempt += 1;
1392
1393            match *attempt {
1394                1 => Err(
1395                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1396                        .set_temporary(),
1397                ),
1398                2 => Err(
1399                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1400                        .set_temporary(),
1401                ),
1402                3 => Ok(Some(8)),
1403                4 => Err(
1404                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1405                        .set_temporary(),
1406                ),
1407                5 => Ok(Some(5)),
1408                6 => Ok(None),
1409                _ => unreachable!(),
1410            }
1411        }
1412
1413        async fn close(&mut self) -> Result<Metadata> {
1414            Ok(MetadataBuilder::unknown().build())
1415        }
1416
1417        async fn abort(&mut self) -> Result<()> {
1418            Ok(())
1419        }
1420    }
1421
1422    fn setup() {
1423        let _ = logforth::starter_log::builder()
1424            .dispatch(|d| {
1425                d.filter(RustLogFilterBuilder::from_default_env().build())
1426                    .append(Testing::default().with_layout(TextLayout::default()))
1427            })
1428            .try_apply();
1429    }
1430
1431    #[tokio::test]
1432    async fn test_retry_read() -> Result<()> {
1433        setup();
1434
1435        let builder = MockBuilder::default();
1436        let op = Operator::new(builder.clone())?
1437            .layer(LoggingLayer::default())
1438            .layer(RetryLayer::default());
1439
1440        let r = op.reader("retryable_error").await?;
1441        let mut content = Vec::new();
1442        let size = r
1443            .read_into(&mut content, ..)
1444            .await
1445            .expect("read must succeed");
1446        assert_eq!(size, 13);
1447        assert_eq!(content, "Hello, World!".as_bytes());
1448        // The error is retryable, we should request it 3 times.
1449        assert_eq!(*builder.attempt.lock().unwrap(), 5);
1450        Ok(())
1451    }
1452
1453    /// This test is used to reproduce the panic issue while composing retry layer with timeout layer.
1454    #[tokio::test]
1455    async fn test_retry_write_fail_on_close() -> Result<()> {
1456        setup();
1457
1458        let builder = MockBuilder::default();
1459        let op = Operator::new(builder.clone())?
1460            .layer(
1461                RetryLayer::default()
1462                    .with_min_delay(Duration::from_millis(1))
1463                    .with_max_delay(Duration::from_millis(1))
1464                    .with_jitter(),
1465            )
1466            // Uncomment this to reproduce timeout layer panic.
1467            // .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)))
1468            .layer(LoggingLayer::default());
1469
1470        let mut w = op.writer("test_write").await?;
1471        w.write("aaa").await?;
1472        w.write("bbb").await?;
1473        match w.close().await {
1474            Ok(_) => (),
1475            Err(_) => {
1476                w.abort().await?;
1477            }
1478        };
1479        Ok(())
1480    }
1481
1482    #[tokio::test]
1483    async fn test_retry_list() -> Result<()> {
1484        setup();
1485
1486        let builder = MockBuilder::default();
1487        let op = Operator::new(builder.clone())?.layer(RetryLayer::default());
1488
1489        let expected = vec!["hello", "world", "2023/", "0208/"];
1490
1491        let mut lister = op
1492            .lister("retryable_error/")
1493            .await
1494            .expect("service must support list");
1495        let mut actual = Vec::new();
1496        while let Some(obj) = lister.try_next().await.expect("must success") {
1497            actual.push(obj.name().to_owned());
1498        }
1499
1500        assert_eq!(actual, expected);
1501        Ok(())
1502    }
1503
1504    #[tokio::test]
1505    async fn test_retry_event_attempt_and_op() -> Result<()> {
1506        setup();
1507
1508        #[derive(Default, Clone)]
1509        struct Recorder {
1510            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1511        }
1512
1513        impl RetryInterceptor for Recorder {
1514            fn intercept(&self, event: RetryEvent<'_>) {
1515                self.events.lock().unwrap().push((event.op, event.attempt));
1516            }
1517        }
1518
1519        let recorder = Recorder::default();
1520        let builder = MockBuilder::default();
1521        let op = Operator::new(builder.clone())?.layer(
1522            RetryLayer::default()
1523                .with_min_delay(Duration::from_millis(1))
1524                .with_max_delay(Duration::from_millis(1))
1525                .with_notify(recorder.clone()),
1526        );
1527
1528        let r = op.reader("retryable_error").await?;
1529        let mut content = Vec::new();
1530        let _ = r.read_into(&mut content, ..).await?;
1531
1532        let events = recorder.events.lock().unwrap().clone();
1533        assert_eq!(
1534            events,
1535            vec![
1536                (Operation::Read, 1),
1537                (Operation::Read, 2),
1538                (Operation::Read, 1),
1539            ],
1540        );
1541        Ok(())
1542    }
1543
1544    #[tokio::test]
1545    async fn test_retry_read_stream_error_enters_retry_budget() -> Result<()> {
1546        setup();
1547
1548        #[derive(Default, Clone)]
1549        struct Recorder {
1550            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1551        }
1552
1553        impl RetryInterceptor for Recorder {
1554            fn intercept(&self, event: RetryEvent<'_>) {
1555                self.events.lock().unwrap().push((event.op, event.attempt));
1556            }
1557        }
1558
1559        let recorder = Recorder::default();
1560        let backend = MockService::default();
1561        let reader = RetryReader::new(
1562            oio::StreamReader::new(MockReader::new(
1563                backend.clone(),
1564                "retryable_error",
1565                OpRead::default(),
1566            )),
1567            Arc::new(recorder.clone()),
1568            ExponentialBuilder::default()
1569                .with_min_delay(Duration::from_millis(1))
1570                .with_max_delay(Duration::from_millis(1)),
1571        );
1572
1573        let (_, mut stream) = oio::Read::open(&reader, BytesRange::default()).await?;
1574        let buf = oio::ReadStream::read_all(&mut stream).await?;
1575
1576        assert_eq!(buf.to_bytes(), Bytes::from_static(b"Hello, World!"));
1577        assert_eq!(*backend.attempt.lock().unwrap(), 5);
1578
1579        let events = recorder.events.lock().unwrap().clone();
1580        assert_eq!(
1581            events,
1582            vec![
1583                (Operation::Read, 1),
1584                (Operation::Read, 2),
1585                (Operation::Read, 1),
1586            ],
1587        );
1588        Ok(())
1589    }
1590
1591    #[tokio::test]
1592    async fn test_retry_batch() -> Result<()> {
1593        setup();
1594
1595        let builder = MockBuilder::default();
1596        // set to a lower delay to make it run faster
1597        let op = Operator::new(builder.clone())?.layer(
1598            RetryLayer::default()
1599                .with_min_delay(Duration::from_secs_f32(0.1))
1600                .with_max_times(5),
1601        );
1602
1603        let paths = vec!["hello", "world", "test", "batch"];
1604        op.delete_stream(stream::iter(paths)).await?;
1605        assert_eq!(*builder.attempt.lock().unwrap(), 5);
1606        Ok(())
1607    }
1608
1609    #[tokio::test]
1610    async fn test_retry_copy() -> Result<()> {
1611        setup();
1612
1613        #[derive(Default, Clone)]
1614        struct Recorder {
1615            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1616        }
1617
1618        impl RetryInterceptor for Recorder {
1619            fn intercept(&self, event: RetryEvent<'_>) {
1620                self.events.lock().unwrap().push((event.op, event.attempt));
1621            }
1622        }
1623
1624        let recorder = Recorder::default();
1625        let builder = MockBuilder::default();
1626        let op = Operator::new(builder.clone())?.layer(
1627            RetryLayer::default()
1628                .with_min_delay(Duration::from_millis(1))
1629                .with_max_delay(Duration::from_millis(1))
1630                .with_notify(recorder.clone()),
1631        );
1632
1633        op.copy("from", "to").await.expect("copy must succeed");
1634
1635        // The MockCopier returns errors on attempts 1, 2, 4 and progress
1636        // on attempts 3, 5; finishing on attempt 6. The retry layer must
1637        // retry `Copier::next` to drive the operation to completion.
1638        assert_eq!(*builder.attempt.lock().unwrap(), 6);
1639
1640        let events = recorder.events.lock().unwrap().clone();
1641        assert_eq!(
1642            events,
1643            vec![
1644                (Operation::Copy, 1),
1645                (Operation::Copy, 2),
1646                (Operation::Copy, 1),
1647            ],
1648        );
1649        Ok(())
1650    }
1651}