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
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 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 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 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 assert_eq!(*builder.attempt.lock().unwrap(), 5);
1393 Ok(())
1394 }
1395
1396 #[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 .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 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 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}