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
302    fn info(&self) -> ServiceInfo {
303        self.inner.info()
304    }
305
306    fn capability(&self) -> Capability {
307        self.inner.capability()
308    }
309
310    async fn create_dir(
311        &self,
312        ctx: &OperationContext,
313        path: &str,
314        args: OpCreateDir,
315    ) -> Result<RpCreateDir> {
316        let mut attempt: u32 = 0;
317        { || self.inner.create_dir(ctx, path, args.clone()) }
318            .retry(self.builder)
319            .when(|e| e.is_temporary())
320            .notify(|err, dur| {
321                attempt += 1;
322                self.notify.intercept(RetryEvent {
323                    op: Operation::CreateDir,
324                    err,
325                    retry_after: dur,
326                    attempt,
327                })
328            })
329            .await
330            .map_err(|err| err.set_persistent())
331    }
332
333    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
334        let mut attempt: u32 = 0;
335        let reader = { || self.inner.read(ctx, path, args.clone()) }
336            .retry(self.builder)
337            .when(|e| e.is_temporary())
338            .notify(|err, dur| {
339                attempt += 1;
340                self.notify.intercept(RetryEvent {
341                    op: Operation::Read,
342                    err,
343                    retry_after: dur,
344                    attempt,
345                })
346            })
347            .call()
348            .map_err(|err| err.set_persistent())?;
349
350        Ok(RetryReader::new(reader, self.notify.clone(), self.builder))
351    }
352
353    fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
354        let mut attempt: u32 = 0;
355        let writer = { || self.inner.write(ctx, path, args.clone()) }
356            .retry(self.builder)
357            .when(|e| e.is_temporary())
358            .notify(|err, dur| {
359                attempt += 1;
360                self.notify.intercept(RetryEvent {
361                    op: Operation::Write,
362                    err,
363                    retry_after: dur,
364                    attempt,
365                })
366            })
367            .call()
368            .map_err(|err| err.set_persistent())?;
369
370        Ok(RetryWrapper::new(writer, self.notify.clone(), self.builder))
371    }
372
373    async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
374        let mut attempt: u32 = 0;
375        { || self.inner.stat(ctx, path, args.clone()) }
376            .retry(self.builder)
377            .when(|e| e.is_temporary())
378            .notify(|err, dur| {
379                attempt += 1;
380                self.notify.intercept(RetryEvent {
381                    op: Operation::Stat,
382                    err,
383                    retry_after: dur,
384                    attempt,
385                })
386            })
387            .await
388            .map_err(|err| err.set_persistent())
389    }
390
391    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
392        let mut attempt: u32 = 0;
393        let deleter = { || self.inner.delete(ctx) }
394            .retry(self.builder)
395            .when(|e| e.is_temporary())
396            .notify(|err, dur| {
397                attempt += 1;
398                self.notify.intercept(RetryEvent {
399                    op: Operation::Delete,
400                    err,
401                    retry_after: dur,
402                    attempt,
403                })
404            })
405            .call()
406            .map_err(|err| err.set_persistent())?;
407
408        Ok(RetryWrapper::new(
409            deleter,
410            self.notify.clone(),
411            self.builder,
412        ))
413    }
414
415    fn copy(
416        &self,
417        ctx: &OperationContext,
418        from: &str,
419        to: &str,
420        args: OpCopy,
421        opts: OpCopier,
422    ) -> Result<Self::Copier> {
423        let mut attempt: u32 = 0;
424        let copier = { || self.inner.copy(ctx, from, to, args.clone(), opts.clone()) }
425            .retry(self.builder)
426            .when(|e| e.is_temporary())
427            .notify(|err, dur| {
428                attempt += 1;
429                self.notify.intercept(RetryEvent {
430                    op: Operation::Copy,
431                    err,
432                    retry_after: dur,
433                    attempt,
434                })
435            })
436            .call()
437            .map_err(|err| err.set_persistent())?;
438
439        Ok(RetryWrapper::new(copier, self.notify.clone(), self.builder))
440    }
441
442    async fn rename(
443        &self,
444        ctx: &OperationContext,
445        from: &str,
446        to: &str,
447        args: OpRename,
448    ) -> Result<RpRename> {
449        let mut attempt: u32 = 0;
450        { || self.inner.rename(ctx, from, to, args.clone()) }
451            .retry(self.builder)
452            .when(|e| e.is_temporary())
453            .notify(|err, dur| {
454                attempt += 1;
455                self.notify.intercept(RetryEvent {
456                    op: Operation::Rename,
457                    err,
458                    retry_after: dur,
459                    attempt,
460                })
461            })
462            .await
463            .map_err(|err| err.set_persistent())
464    }
465
466    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
467        let mut attempt: u32 = 0;
468        let lister = { || self.inner.list(ctx, path, args.clone()) }
469            .retry(self.builder)
470            .when(|e| e.is_temporary())
471            .notify(|err, dur| {
472                attempt += 1;
473                self.notify.intercept(RetryEvent {
474                    op: Operation::List,
475                    err,
476                    retry_after: dur,
477                    attempt,
478                })
479            })
480            .call()
481            .map_err(|err| err.set_persistent())?;
482
483        Ok(RetryWrapper::new(lister, self.notify.clone(), self.builder))
484    }
485
486    async fn presign(
487        &self,
488        ctx: &OperationContext,
489        path: &str,
490        args: OpPresign,
491    ) -> Result<RpPresign> {
492        let mut attempt: u32 = 0;
493        { || self.inner.presign(ctx, path, args.clone()) }
494            .retry(self.builder)
495            .when(|e| e.is_temporary())
496            .notify(|err, dur| {
497                attempt += 1;
498                self.notify.intercept(RetryEvent {
499                    op: Operation::Presign,
500                    err,
501                    retry_after: dur,
502                    attempt,
503                })
504            })
505            .await
506            .map_err(|err| err.set_persistent())
507    }
508}
509
510#[doc(hidden)]
511pub struct RetryReader<R, I> {
512    inner: Arc<R>,
513    notify: Arc<I>,
514    builder: ExponentialBuilder,
515}
516
517impl<R, I> RetryReader<R, I> {
518    fn new(inner: R, notify: Arc<I>, builder: ExponentialBuilder) -> Self {
519        Self {
520            inner: Arc::new(inner),
521            notify,
522            builder,
523        }
524    }
525}
526
527impl<R: oio::Read + 'static, I: RetryInterceptor> oio::Read for RetryReader<R, I> {
528    async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
529        use backon::Retryable;
530
531        let mut attempt: u32 = 0;
532        let (rp, stream) = { || self.inner.open(range) }
533            .retry(self.builder)
534            .when(|e| e.is_temporary())
535            .notify(|err, dur| {
536                attempt += 1;
537                self.notify.intercept(RetryEvent {
538                    op: Operation::Read,
539                    err,
540                    retry_after: dur,
541                    attempt,
542                })
543            })
544            .await
545            .map_err(|e| e.set_persistent())?;
546
547        Ok((
548            rp,
549            Box::new(RetryReadStream::new(
550                self.inner.clone(),
551                stream,
552                range,
553                self.notify.clone(),
554                self.builder,
555            )) as Box<dyn oio::ReadStreamDyn>,
556        ))
557    }
558
559    async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
560        use backon::Retryable;
561
562        let mut attempt: u32 = 0;
563        { || self.inner.read(range) }
564            .retry(self.builder)
565            .when(|e| e.is_temporary())
566            .notify(|err, dur| {
567                attempt += 1;
568                self.notify.intercept(RetryEvent {
569                    op: Operation::Read,
570                    err,
571                    retry_after: dur,
572                    attempt,
573                })
574            })
575            .await
576            .map_err(|e| e.set_persistent())
577    }
578}
579
580#[doc(hidden)]
581pub struct RetryReadStream<R, I> {
582    reader: Arc<R>,
583    stream: Option<Box<dyn oio::ReadStreamDyn>>,
584    range: BytesRange,
585    read: u64,
586    notify: Arc<I>,
587    builder: ExponentialBuilder,
588}
589
590impl<R, I> RetryReadStream<R, I> {
591    fn new(
592        reader: Arc<R>,
593        stream: Box<dyn oio::ReadStreamDyn>,
594        range: BytesRange,
595        notify: Arc<I>,
596        builder: ExponentialBuilder,
597    ) -> Self {
598        Self {
599            reader,
600            stream: Some(stream),
601            range,
602            read: 0,
603            notify,
604            builder,
605        }
606    }
607}
608
609impl<R: oio::Read, I: RetryInterceptor> oio::ReadStream for RetryReadStream<R, I> {
610    async fn read(&mut self) -> Result<Buffer> {
611        use backon::RetryableWithContext;
612
613        let reader = self.reader.clone();
614        let stream = self.stream.take();
615        let range = self.range;
616        let read = self.read;
617        let mut attempt: u32 = 0;
618
619        let ((stream, range, read), res) = {
620            |(stream, mut range, mut read): (
621                Option<Box<dyn oio::ReadStreamDyn>>,
622                BytesRange,
623                u64,
624            )| {
625                let reader = reader.clone();
626                async move {
627                    let mut stream = match stream {
628                        Some(stream) => stream,
629                        None => {
630                            range.advance(read);
631                            read = 0;
632
633                            match reader.open(range).await {
634                                Ok((_, stream)) => stream,
635                                Err(err) => return ((None, range, read), Err(err)),
636                            }
637                        }
638                    };
639
640                    let res = match stream.read().await {
641                        Ok(buf) => {
642                            if !buf.is_empty() {
643                                read += buf.len() as u64;
644                            }
645                            (Some(stream), Ok(buf))
646                        }
647                        Err(err) => (None, Err(err)),
648                    };
649
650                    ((res.0, range, read), res.1)
651                }
652            }
653        }
654        .retry(self.builder)
655        .when(|e| e.is_temporary())
656        .context((stream, range, read))
657        .notify(|err, dur| {
658            attempt += 1;
659            self.notify.intercept(RetryEvent {
660                op: Operation::Read,
661                err,
662                retry_after: dur,
663                attempt,
664            })
665        })
666        .await;
667
668        self.stream = stream;
669        self.range = range;
670        self.read = read;
671
672        res.map_err(|err| err.set_persistent())
673    }
674}
675
676#[doc(hidden)]
677pub struct RetryWrapper<R, I> {
678    inner: Option<R>,
679    notify: Arc<I>,
680
681    builder: ExponentialBuilder,
682}
683
684impl<R, I> RetryWrapper<R, I> {
685    fn new(inner: R, notify: Arc<I>, backoff: ExponentialBuilder) -> Self {
686        Self {
687            inner: Some(inner),
688            notify,
689            builder: backoff,
690        }
691    }
692
693    fn take_inner(&mut self) -> Result<R> {
694        self.inner.take().ok_or_else(|| {
695            Error::new(
696                ErrorKind::Unexpected,
697                "retry layer is in bad state, please make sure future not dropped before ready",
698            )
699        })
700    }
701}
702
703impl<R: oio::ReadStream, I: RetryInterceptor> oio::ReadStream for RetryWrapper<R, I> {
704    async fn read(&mut self) -> Result<Buffer> {
705        use backon::RetryableWithContext;
706
707        let inner = self.take_inner()?;
708        let mut attempt: u32 = 0;
709
710        let (inner, res) = {
711            |mut r: R| async move {
712                let res = r.read().await;
713
714                (r, res)
715            }
716        }
717        .retry(self.builder)
718        .when(|e| e.is_temporary())
719        .context(inner)
720        .notify(|err, dur| {
721            attempt += 1;
722            self.notify.intercept(RetryEvent {
723                op: Operation::Read,
724                err,
725                retry_after: dur,
726                attempt,
727            })
728        })
729        .await;
730
731        self.inner = Some(inner);
732        res.map_err(|err| err.set_persistent())
733    }
734}
735
736impl<R: oio::Write, I: RetryInterceptor> oio::Write for RetryWrapper<R, I> {
737    async fn write(&mut self, bs: Buffer) -> Result<()> {
738        use backon::RetryableWithContext;
739
740        let inner = self.take_inner()?;
741        let mut attempt: u32 = 0;
742
743        let ((inner, _), res) = {
744            |(mut r, bs): (R, Buffer)| async move {
745                let res = r.write(bs.clone()).await;
746
747                ((r, bs), res)
748            }
749        }
750        .retry(self.builder)
751        .when(|e| e.is_temporary())
752        .context((inner, bs))
753        .notify(|err, dur| {
754            attempt += 1;
755            self.notify.intercept(RetryEvent {
756                op: Operation::Write,
757                err,
758                retry_after: dur,
759                attempt,
760            })
761        })
762        .await;
763
764        self.inner = Some(inner);
765        res.map_err(|err| err.set_persistent())
766    }
767
768    async fn abort(&mut self) -> Result<()> {
769        use backon::RetryableWithContext;
770
771        let inner = self.take_inner()?;
772        let mut attempt: u32 = 0;
773
774        let (inner, res) = {
775            |mut r: R| async move {
776                let res = r.abort().await;
777
778                (r, res)
779            }
780        }
781        .retry(self.builder)
782        .when(|e| e.is_temporary())
783        .context(inner)
784        .notify(|err, dur| {
785            attempt += 1;
786            self.notify.intercept(RetryEvent {
787                op: Operation::Write,
788                err,
789                retry_after: dur,
790                attempt,
791            })
792        })
793        .await;
794
795        self.inner = Some(inner);
796        res.map_err(|err| err.set_persistent())
797    }
798
799    async fn close(&mut self) -> Result<Metadata> {
800        use backon::RetryableWithContext;
801
802        let inner = self.take_inner()?;
803        let mut attempt: u32 = 0;
804
805        let (inner, res) = {
806            |mut r: R| async move {
807                let res = r.close().await;
808
809                (r, res)
810            }
811        }
812        .retry(self.builder)
813        .when(|e| e.is_temporary())
814        .context(inner)
815        .notify(|err, dur| {
816            attempt += 1;
817            self.notify.intercept(RetryEvent {
818                op: Operation::Write,
819                err,
820                retry_after: dur,
821                attempt,
822            })
823        })
824        .await;
825
826        self.inner = Some(inner);
827        res.map_err(|err| err.set_persistent())
828    }
829}
830
831impl<P: oio::List, I: RetryInterceptor> oio::List for RetryWrapper<P, I> {
832    async fn next(&mut self) -> Result<Option<oio::Entry>> {
833        use backon::RetryableWithContext;
834
835        let inner = self.take_inner()?;
836        let mut attempt: u32 = 0;
837
838        let (inner, res) = {
839            |mut p: P| async move {
840                let res = p.next().await;
841
842                (p, res)
843            }
844        }
845        .retry(self.builder)
846        .when(|e| e.is_temporary())
847        .context(inner)
848        .notify(|err, dur| {
849            attempt += 1;
850            self.notify.intercept(RetryEvent {
851                op: Operation::List,
852                err,
853                retry_after: dur,
854                attempt,
855            })
856        })
857        .await;
858
859        self.inner = Some(inner);
860        res.map_err(|err| err.set_persistent())
861    }
862}
863
864impl<P: oio::Delete, I: RetryInterceptor> oio::Delete for RetryWrapper<P, I> {
865    async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
866        use backon::RetryableWithContext;
867
868        let inner = self.take_inner()?;
869        let path = path.to_string();
870        let args_cloned = args.clone();
871        let mut attempt: u32 = 0;
872
873        let (inner, res) = {
874            |mut p: P| {
875                let path = path.clone();
876                let args = args_cloned.clone();
877                async move {
878                    let res = p.delete(&path, args).await;
879                    (p, res)
880                }
881            }
882        }
883        .retry(self.builder)
884        .when(|e| e.is_temporary())
885        .context(inner)
886        .notify(|err, dur| {
887            attempt += 1;
888            self.notify.intercept(RetryEvent {
889                op: Operation::Delete,
890                err,
891                retry_after: dur,
892                attempt,
893            });
894        })
895        .await;
896
897        self.inner = Some(inner);
898        res.map_err(|e| e.set_persistent())
899    }
900
901    async fn close(&mut self) -> Result<()> {
902        use backon::RetryableWithContext;
903
904        let inner = self.take_inner()?;
905        let mut attempt: u32 = 0;
906
907        let (inner, res) = {
908            |mut p: P| async move {
909                let res = p.close().await;
910
911                (p, res)
912            }
913        }
914        .retry(self.builder)
915        .when(|e| e.is_temporary())
916        .context(inner)
917        .notify(|err, dur| {
918            attempt += 1;
919            self.notify.intercept(RetryEvent {
920                op: Operation::Delete,
921                err,
922                retry_after: dur,
923                attempt,
924            })
925        })
926        .await;
927
928        self.inner = Some(inner);
929        res.map_err(|err| err.set_persistent())
930    }
931}
932
933impl<C: oio::Copy, I: RetryInterceptor> oio::Copy for RetryWrapper<C, I> {
934    async fn next(&mut self) -> Result<Option<usize>> {
935        use backon::RetryableWithContext;
936
937        let inner = self.take_inner()?;
938        let mut attempt: u32 = 0;
939
940        let (inner, res) = {
941            |mut c: C| async move {
942                let res = c.next().await;
943
944                (c, res)
945            }
946        }
947        .retry(self.builder)
948        .when(|e| e.is_temporary())
949        .context(inner)
950        .notify(|err, dur| {
951            attempt += 1;
952            self.notify.intercept(RetryEvent {
953                op: Operation::Copy,
954                err,
955                retry_after: dur,
956                attempt,
957            })
958        })
959        .await;
960
961        self.inner = Some(inner);
962        res.map_err(|err| err.set_persistent())
963    }
964
965    async fn close(&mut self) -> Result<Metadata> {
966        use backon::RetryableWithContext;
967
968        let inner = self.take_inner()?;
969        let mut attempt: u32 = 0;
970
971        let (inner, res) = {
972            |mut c: C| async move {
973                let res = c.close().await;
974
975                (c, res)
976            }
977        }
978        .retry(self.builder)
979        .when(|e| e.is_temporary())
980        .context(inner)
981        .notify(|err, dur| {
982            attempt += 1;
983            self.notify.intercept(RetryEvent {
984                op: Operation::Copy,
985                err,
986                retry_after: dur,
987                attempt,
988            })
989        })
990        .await;
991
992        self.inner = Some(inner);
993        res.map_err(|err| err.set_persistent())
994    }
995
996    async fn abort(&mut self) -> Result<()> {
997        use backon::RetryableWithContext;
998
999        let inner = self.take_inner()?;
1000        let mut attempt: u32 = 0;
1001
1002        let (inner, res) = {
1003            |mut c: C| async move {
1004                let res = c.abort().await;
1005
1006                (c, res)
1007            }
1008        }
1009        .retry(self.builder)
1010        .when(|e| e.is_temporary())
1011        .context(inner)
1012        .notify(|err, dur| {
1013            attempt += 1;
1014            self.notify.intercept(RetryEvent {
1015                op: Operation::Copy,
1016                err,
1017                retry_after: dur,
1018                attempt,
1019            })
1020        })
1021        .await;
1022
1023        self.inner = Some(inner);
1024        res.map_err(|err| err.set_persistent())
1025    }
1026}
1027
1028#[cfg(test)]
1029mod tests {
1030    use std::sync::Mutex;
1031
1032    use bytes::Bytes;
1033    use futures::TryStreamExt;
1034    use futures::stream;
1035    use logforth::append::Testing;
1036    use logforth::filter::rustlog::RustLogFilterBuilder;
1037    use logforth::layout::TextLayout;
1038    use opendal_layer_logging::LoggingLayer;
1039
1040    use super::*;
1041
1042    #[derive(Default, Clone)]
1043    struct MockBuilder {
1044        attempt: Arc<Mutex<usize>>,
1045    }
1046
1047    impl Builder for MockBuilder {
1048        type Config = ();
1049
1050        fn build(self) -> Result<impl Service> {
1051            Ok(MockService {
1052                attempt: self.attempt,
1053            })
1054        }
1055    }
1056
1057    #[derive(Debug, Clone, Default)]
1058    struct MockService {
1059        attempt: Arc<Mutex<usize>>,
1060    }
1061
1062    /// Reader returned by this backend.
1063    pub struct MockReader {
1064        backend: MockService,
1065    }
1066
1067    impl MockReader {
1068        fn new(backend: MockService, _: &str, _: OpRead) -> Self {
1069            Self { backend }
1070        }
1071    }
1072
1073    impl oio::StreamRead for MockReader {
1074        async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
1075            let backend = &self.backend;
1076            let rp = RpRead::new(Metadata::new(EntryMode::FILE).with_content_length(0));
1077            let stream = MockReadStream {
1078                buf: Bytes::from("Hello, World!").into(),
1079                range,
1080                attempt: backend.attempt.clone(),
1081            };
1082
1083            Ok((rp, Box::new(stream) as Box<dyn oio::ReadStreamDyn>))
1084        }
1085    }
1086
1087    impl Service for MockService {
1088        type Reader = oio::StreamReader<MockReader>;
1089        type Writer = MockWriter;
1090        type Lister = MockLister;
1091        type Deleter = MockDeleter;
1092        type Copier = MockCopier;
1093
1094        fn info(&self) -> ServiceInfo {
1095            ServiceInfo::with_scheme("mock")
1096        }
1097
1098        fn capability(&self) -> Capability {
1099            Capability {
1100                read: true,
1101                write: true,
1102                write_can_multi: true,
1103                delete: true,
1104                delete_max_size: Some(10),
1105                stat: true,
1106                list: true,
1107                list_with_recursive: true,
1108                copy: true,
1109                ..Default::default()
1110            }
1111        }
1112
1113        async fn create_dir(
1114            &self,
1115            _: &OperationContext,
1116            _: &str,
1117            _: OpCreateDir,
1118        ) -> Result<RpCreateDir> {
1119            Err(Error::new(
1120                ErrorKind::Unsupported,
1121                "operation is not supported",
1122            ))
1123        }
1124
1125        async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
1126            Ok(RpStat::new(
1127                Metadata::new(EntryMode::FILE).with_content_length(13),
1128            ))
1129        }
1130
1131        fn read(&self, _: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
1132            Ok(oio::StreamReader::new(MockReader::new(
1133                self.clone(),
1134                path,
1135                args,
1136            )))
1137        }
1138
1139        fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
1140            Ok(MockDeleter {
1141                size: 0,
1142                attempt: self.attempt.clone(),
1143            })
1144        }
1145
1146        fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
1147            Ok(MockWriter {})
1148        }
1149
1150        fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
1151            let lister = MockLister::default();
1152            Ok(lister)
1153        }
1154
1155        fn copy(
1156            &self,
1157            _: &OperationContext,
1158            _: &str,
1159            _: &str,
1160            _: OpCopy,
1161            _: OpCopier,
1162        ) -> Result<Self::Copier> {
1163            Ok(MockCopier {
1164                attempt: self.attempt.clone(),
1165            })
1166        }
1167
1168        async fn rename(
1169            &self,
1170            _: &OperationContext,
1171            _: &str,
1172            _: &str,
1173            _: OpRename,
1174        ) -> Result<RpRename> {
1175            Err(Error::new(
1176                ErrorKind::Unsupported,
1177                "operation is not supported",
1178            ))
1179        }
1180
1181        async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
1182            Err(Error::new(
1183                ErrorKind::Unsupported,
1184                "operation is not supported",
1185            ))
1186        }
1187    }
1188
1189    #[derive(Debug, Clone, Default)]
1190    struct MockReadStream {
1191        buf: Buffer,
1192        range: BytesRange,
1193        attempt: Arc<Mutex<usize>>,
1194    }
1195
1196    impl oio::ReadStream for MockReadStream {
1197        async fn read(&mut self) -> Result<Buffer> {
1198            let mut attempt = self.attempt.lock().unwrap();
1199            *attempt += 1;
1200
1201            match *attempt {
1202                1 => Err(
1203                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1204                        .set_temporary(),
1205                ),
1206                2 => Err(
1207                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1208                        .set_temporary(),
1209                ),
1210                // Should read out all data.
1211                3 => Ok(self.buf.slice(self.range.to_range_as_usize())),
1212                4 => Err(
1213                    Error::new(ErrorKind::Unexpected, "retryable_error from reader")
1214                        .set_temporary(),
1215                ),
1216                // Should be empty.
1217                5 => Ok(self.buf.slice(self.range.to_range_as_usize())),
1218                _ => unreachable!(),
1219            }
1220        }
1221    }
1222
1223    #[derive(Debug, Clone, Default)]
1224    struct MockWriter {}
1225
1226    impl oio::Write for MockWriter {
1227        async fn write(&mut self, _: Buffer) -> Result<()> {
1228            Ok(())
1229        }
1230
1231        async fn close(&mut self) -> Result<Metadata> {
1232            Err(Error::new(ErrorKind::Unexpected, "always close failed").set_temporary())
1233        }
1234
1235        async fn abort(&mut self) -> Result<()> {
1236            Ok(())
1237        }
1238    }
1239
1240    #[derive(Debug, Clone, Default)]
1241    struct MockLister {
1242        attempt: usize,
1243    }
1244
1245    impl oio::List for MockLister {
1246        async fn next(&mut self) -> Result<Option<oio::Entry>> {
1247            self.attempt += 1;
1248            match self.attempt {
1249                1 => Err(Error::new(
1250                    ErrorKind::RateLimited,
1251                    "retryable rate limited error from lister",
1252                )
1253                .set_temporary()),
1254                2 => Ok(Some(oio::Entry::new(
1255                    "hello",
1256                    Metadata::new(EntryMode::FILE),
1257                ))),
1258                3 => Ok(Some(oio::Entry::new(
1259                    "world",
1260                    Metadata::new(EntryMode::FILE),
1261                ))),
1262                4 => Err(
1263                    Error::new(ErrorKind::Unexpected, "retryable internal server error")
1264                        .set_temporary(),
1265                ),
1266                5 => Ok(Some(oio::Entry::new(
1267                    "2023/",
1268                    Metadata::new(EntryMode::DIR),
1269                ))),
1270                6 => Ok(Some(oio::Entry::new(
1271                    "0208/",
1272                    Metadata::new(EntryMode::DIR),
1273                ))),
1274                7 => Ok(None),
1275                _ => {
1276                    unreachable!()
1277                }
1278            }
1279        }
1280    }
1281
1282    #[derive(Debug, Clone, Default)]
1283    struct MockDeleter {
1284        size: usize,
1285        attempt: Arc<Mutex<usize>>,
1286    }
1287
1288    impl oio::Delete for MockDeleter {
1289        async fn delete(&mut self, _: &str, _: OpDelete) -> Result<()> {
1290            self.size += 1;
1291            Ok(())
1292        }
1293
1294        async fn close(&mut self) -> Result<()> {
1295            let mut attempt = self.attempt.lock().unwrap();
1296            *attempt += 1;
1297
1298            match *attempt {
1299                1 => Err(
1300                    Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1301                        .set_temporary(),
1302                ),
1303                2..=4 => {
1304                    self.size = self.size.saturating_sub(1);
1305                    Err(
1306                        Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1307                            .set_temporary(),
1308                    )
1309                }
1310                5 => {
1311                    self.size = self.size.saturating_sub(1);
1312                    if self.size == 0 {
1313                        Ok(())
1314                    } else {
1315                        Err(
1316                            Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
1317                                .set_temporary(),
1318                        )
1319                    }
1320                }
1321                _ => unreachable!(),
1322            }
1323        }
1324    }
1325
1326    #[derive(Debug, Clone, Default)]
1327    struct MockCopier {
1328        attempt: Arc<Mutex<usize>>,
1329    }
1330
1331    impl oio::Copy for MockCopier {
1332        async fn next(&mut self) -> Result<Option<usize>> {
1333            let mut attempt = self.attempt.lock().unwrap();
1334            *attempt += 1;
1335
1336            match *attempt {
1337                1 => Err(
1338                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1339                        .set_temporary(),
1340                ),
1341                2 => Err(
1342                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1343                        .set_temporary(),
1344                ),
1345                3 => Ok(Some(8)),
1346                4 => Err(
1347                    Error::new(ErrorKind::Unexpected, "retryable_error from copier")
1348                        .set_temporary(),
1349                ),
1350                5 => Ok(Some(5)),
1351                6 => Ok(None),
1352                _ => unreachable!(),
1353            }
1354        }
1355
1356        async fn close(&mut self) -> Result<Metadata> {
1357            Ok(Metadata::default())
1358        }
1359
1360        async fn abort(&mut self) -> Result<()> {
1361            Ok(())
1362        }
1363    }
1364
1365    fn setup() {
1366        let _ = logforth::starter_log::builder()
1367            .dispatch(|d| {
1368                d.filter(RustLogFilterBuilder::from_default_env().build())
1369                    .append(Testing::default().with_layout(TextLayout::default()))
1370            })
1371            .try_apply();
1372    }
1373
1374    #[tokio::test]
1375    async fn test_retry_read() -> Result<()> {
1376        setup();
1377
1378        let builder = MockBuilder::default();
1379        let op = Operator::new(builder.clone())?
1380            .layer(LoggingLayer::default())
1381            .layer(RetryLayer::default());
1382
1383        let r = op.reader("retryable_error").await?;
1384        let mut content = Vec::new();
1385        let size = r
1386            .read_into(&mut content, ..)
1387            .await
1388            .expect("read must succeed");
1389        assert_eq!(size, 13);
1390        assert_eq!(content, "Hello, World!".as_bytes());
1391        // The error is retryable, we should request it 3 times.
1392        assert_eq!(*builder.attempt.lock().unwrap(), 5);
1393        Ok(())
1394    }
1395
1396    /// This test is used to reproduce the panic issue while composing retry layer with timeout layer.
1397    #[tokio::test]
1398    async fn test_retry_write_fail_on_close() -> Result<()> {
1399        setup();
1400
1401        let builder = MockBuilder::default();
1402        let op = Operator::new(builder.clone())?
1403            .layer(
1404                RetryLayer::default()
1405                    .with_min_delay(Duration::from_millis(1))
1406                    .with_max_delay(Duration::from_millis(1))
1407                    .with_jitter(),
1408            )
1409            // Uncomment this to reproduce timeout layer panic.
1410            // .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)))
1411            .layer(LoggingLayer::default());
1412
1413        let mut w = op.writer("test_write").await?;
1414        w.write("aaa").await?;
1415        w.write("bbb").await?;
1416        match w.close().await {
1417            Ok(_) => (),
1418            Err(_) => {
1419                w.abort().await?;
1420            }
1421        };
1422        Ok(())
1423    }
1424
1425    #[tokio::test]
1426    async fn test_retry_list() -> Result<()> {
1427        setup();
1428
1429        let builder = MockBuilder::default();
1430        let op = Operator::new(builder.clone())?.layer(RetryLayer::default());
1431
1432        let expected = vec!["hello", "world", "2023/", "0208/"];
1433
1434        let mut lister = op
1435            .lister("retryable_error/")
1436            .await
1437            .expect("service must support list");
1438        let mut actual = Vec::new();
1439        while let Some(obj) = lister.try_next().await.expect("must success") {
1440            actual.push(obj.name().to_owned());
1441        }
1442
1443        assert_eq!(actual, expected);
1444        Ok(())
1445    }
1446
1447    #[tokio::test]
1448    async fn test_retry_event_attempt_and_op() -> Result<()> {
1449        setup();
1450
1451        #[derive(Default, Clone)]
1452        struct Recorder {
1453            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1454        }
1455
1456        impl RetryInterceptor for Recorder {
1457            fn intercept(&self, event: RetryEvent<'_>) {
1458                self.events.lock().unwrap().push((event.op, event.attempt));
1459            }
1460        }
1461
1462        let recorder = Recorder::default();
1463        let builder = MockBuilder::default();
1464        let op = Operator::new(builder.clone())?.layer(
1465            RetryLayer::default()
1466                .with_min_delay(Duration::from_millis(1))
1467                .with_max_delay(Duration::from_millis(1))
1468                .with_notify(recorder.clone()),
1469        );
1470
1471        let r = op.reader("retryable_error").await?;
1472        let mut content = Vec::new();
1473        let _ = r.read_into(&mut content, ..).await?;
1474
1475        let events = recorder.events.lock().unwrap().clone();
1476        assert_eq!(
1477            events,
1478            vec![
1479                (Operation::Read, 1),
1480                (Operation::Read, 2),
1481                (Operation::Read, 1),
1482            ],
1483        );
1484        Ok(())
1485    }
1486
1487    #[tokio::test]
1488    async fn test_retry_read_stream_error_enters_retry_budget() -> Result<()> {
1489        setup();
1490
1491        #[derive(Default, Clone)]
1492        struct Recorder {
1493            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1494        }
1495
1496        impl RetryInterceptor for Recorder {
1497            fn intercept(&self, event: RetryEvent<'_>) {
1498                self.events.lock().unwrap().push((event.op, event.attempt));
1499            }
1500        }
1501
1502        let recorder = Recorder::default();
1503        let backend = MockService::default();
1504        let reader = RetryReader::new(
1505            oio::StreamReader::new(MockReader::new(
1506                backend.clone(),
1507                "retryable_error",
1508                OpRead::default(),
1509            )),
1510            Arc::new(recorder.clone()),
1511            ExponentialBuilder::default()
1512                .with_min_delay(Duration::from_millis(1))
1513                .with_max_delay(Duration::from_millis(1)),
1514        );
1515
1516        let (_, mut stream) = oio::Read::open(&reader, BytesRange::default()).await?;
1517        let buf = oio::ReadStream::read_all(&mut stream).await?;
1518
1519        assert_eq!(buf.to_bytes(), Bytes::from_static(b"Hello, World!"));
1520        assert_eq!(*backend.attempt.lock().unwrap(), 5);
1521
1522        let events = recorder.events.lock().unwrap().clone();
1523        assert_eq!(
1524            events,
1525            vec![
1526                (Operation::Read, 1),
1527                (Operation::Read, 2),
1528                (Operation::Read, 1),
1529            ],
1530        );
1531        Ok(())
1532    }
1533
1534    #[tokio::test]
1535    async fn test_retry_batch() -> Result<()> {
1536        setup();
1537
1538        let builder = MockBuilder::default();
1539        // set to a lower delay to make it run faster
1540        let op = Operator::new(builder.clone())?.layer(
1541            RetryLayer::default()
1542                .with_min_delay(Duration::from_secs_f32(0.1))
1543                .with_max_times(5),
1544        );
1545
1546        let paths = vec!["hello", "world", "test", "batch"];
1547        op.delete_stream(stream::iter(paths)).await?;
1548        assert_eq!(*builder.attempt.lock().unwrap(), 5);
1549        Ok(())
1550    }
1551
1552    #[tokio::test]
1553    async fn test_retry_copy() -> Result<()> {
1554        setup();
1555
1556        #[derive(Default, Clone)]
1557        struct Recorder {
1558            events: Arc<Mutex<Vec<(Operation, u32)>>>,
1559        }
1560
1561        impl RetryInterceptor for Recorder {
1562            fn intercept(&self, event: RetryEvent<'_>) {
1563                self.events.lock().unwrap().push((event.op, event.attempt));
1564            }
1565        }
1566
1567        let recorder = Recorder::default();
1568        let builder = MockBuilder::default();
1569        let op = Operator::new(builder.clone())?.layer(
1570            RetryLayer::default()
1571                .with_min_delay(Duration::from_millis(1))
1572                .with_max_delay(Duration::from_millis(1))
1573                .with_notify(recorder.clone()),
1574        );
1575
1576        op.copy("from", "to").await.expect("copy must succeed");
1577
1578        // The MockCopier returns errors on attempts 1, 2, 4 and progress
1579        // on attempts 3, 5; finishing on attempt 6. The retry layer must
1580        // retry `Copier::next` to drive the operation to completion.
1581        assert_eq!(*builder.attempt.lock().unwrap(), 6);
1582
1583        let events = recorder.events.lock().unwrap().clone();
1584        assert_eq!(
1585            events,
1586            vec![
1587                (Operation::Copy, 1),
1588                (Operation::Copy, 2),
1589                (Operation::Copy, 1),
1590            ],
1591        );
1592        Ok(())
1593    }
1594}