1#![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
31pub 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 pub fn new() -> RetryLayer {
150 Self::default()
151 }
152}
153
154impl<I: RetryInterceptor> RetryLayer<I> {
155 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 pub fn with_jitter(mut self) -> Self {
180 self.builder = self.builder.with_jitter();
181 self
182 }
183
184 pub fn with_factor(mut self, factor: f32) -> Self {
190 self.builder = self.builder.with_factor(factor);
191 self
192 }
193
194 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 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 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#[non_exhaustive]
235#[derive(Debug)]
236pub struct RetryEvent<'a> {
237 pub op: Operation,
239 pub err: &'a Error,
241 pub retry_after: Duration,
243 pub attempt: u32,
245}
246
247pub trait RetryInterceptor: Send + Sync + 'static {
249 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
267pub 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 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 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 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 assert_eq!(*builder.attempt.lock().unwrap(), 5);
1450 Ok(())
1451 }
1452
1453 #[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 .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 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 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}