#![doc = include_str!("../README.md")]
#![cfg_attr(docsrs, feature(doc_cfg))]
#![cfg_attr(docsrs, doc(auto_cfg))]
#![deny(missing_docs)]
use std::fmt::Debug;
use std::sync::Arc;
use backon::BlockingRetryable;
use backon::ExponentialBuilder;
use backon::Retryable;
use opendal_core::raw::*;
use opendal_core::*;
pub struct RetryLayer<I: RetryInterceptor = DefaultRetryInterceptor> {
builder: ExponentialBuilder,
notify: Arc<I>,
}
impl<I: RetryInterceptor> Debug for RetryLayer<I> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RetryLayer")
.field("builder", &self.builder)
.finish_non_exhaustive()
}
}
impl<I: RetryInterceptor> Clone for RetryLayer<I> {
fn clone(&self) -> Self {
Self {
builder: self.builder,
notify: self.notify.clone(),
}
}
}
impl Default for RetryLayer {
fn default() -> Self {
Self {
builder: ExponentialBuilder::default(),
notify: Arc::new(DefaultRetryInterceptor),
}
}
}
impl RetryLayer {
pub fn new() -> RetryLayer {
Self::default()
}
}
impl<I: RetryInterceptor> RetryLayer<I> {
pub fn with_notify<NI: RetryInterceptor>(self, notify: NI) -> RetryLayer<NI> {
RetryLayer {
builder: self.builder,
notify: Arc::new(notify),
}
}
pub fn with_jitter(mut self) -> Self {
self.builder = self.builder.with_jitter();
self
}
pub fn with_factor(mut self, factor: f32) -> Self {
self.builder = self.builder.with_factor(factor);
self
}
pub fn with_min_delay(mut self, min_delay: Duration) -> Self {
self.builder = self.builder.with_min_delay(min_delay);
self
}
pub fn with_max_delay(mut self, max_delay: Duration) -> Self {
self.builder = self.builder.with_max_delay(max_delay);
self
}
pub fn with_max_times(mut self, max_times: usize) -> Self {
self.builder = self.builder.with_max_times(max_times);
self
}
}
impl<I: RetryInterceptor> Layer for RetryLayer<I> {
fn apply_service(&self, inner: Servicer) -> Servicer {
Arc::new(self.layer(inner))
}
}
impl<I: RetryInterceptor> RetryLayer<I> {
fn layer(&self, inner: Servicer) -> RetryService<I> {
RetryService {
inner,
notify: self.notify.clone(),
builder: self.builder,
}
}
}
#[non_exhaustive]
#[derive(Debug)]
pub struct RetryEvent<'a> {
pub op: Operation,
pub err: &'a Error,
pub retry_after: Duration,
pub attempt: u32,
}
pub trait RetryInterceptor: Send + Sync + 'static {
fn intercept(&self, event: RetryEvent<'_>);
}
impl<F> RetryInterceptor for F
where
F: for<'a> Fn(RetryEvent<'a>) + Send + Sync + 'static,
{
fn intercept(&self, event: RetryEvent<'_>) {
self(event);
}
}
pub struct DefaultRetryInterceptor;
impl RetryInterceptor for DefaultRetryInterceptor {
fn intercept(&self, event: RetryEvent<'_>) {
log::warn!(
target: "opendal::layers::retry",
"will retry {:?} (attempt {}) after {}s because: {:?}",
event.op, event.attempt, event.retry_after.as_secs_f64(), event.err
);
}
}
#[doc(hidden)]
pub struct RetryService<I: RetryInterceptor> {
inner: Servicer,
notify: Arc<I>,
builder: ExponentialBuilder,
}
impl<I: RetryInterceptor> Debug for RetryService<I> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RetryService")
.field("inner", &self.inner)
.finish_non_exhaustive()
}
}
impl<I: RetryInterceptor> Service for RetryService<I> {
type Reader = RetryReader<oio::Reader, I>;
type Writer = RetryWrapper<oio::Writer, I>;
type Lister = RetryWrapper<oio::Lister, I>;
type Deleter = RetryWrapper<oio::Deleter, I>;
type Copier = RetryWrapper<oio::Copier, I>;
type Composer = oio::Composer;
fn info(&self) -> ServiceInfo {
self.inner.info()
}
fn capability(&self) -> Capability {
self.inner.capability()
}
fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
self.inner.compose(ctx, to, args)
}
async fn create_dir(
&self,
ctx: &OperationContext,
path: &str,
args: OpCreateDir,
) -> Result<RpCreateDir> {
let mut attempt: u32 = 0;
{ || self.inner.create_dir(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::CreateDir,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|err| err.set_persistent())
}
fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
let mut attempt: u32 = 0;
let reader = { || self.inner.read(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Read,
err,
retry_after: dur,
attempt,
})
})
.call()
.map_err(|err| err.set_persistent())?;
Ok(RetryReader::new(reader, self.notify.clone(), self.builder))
}
fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
let mut attempt: u32 = 0;
let writer = { || self.inner.write(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Write,
err,
retry_after: dur,
attempt,
})
})
.call()
.map_err(|err| err.set_persistent())?;
Ok(RetryWrapper::new(writer, self.notify.clone(), self.builder))
}
async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
let mut attempt: u32 = 0;
{ || self.inner.stat(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Stat,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|err| err.set_persistent())
}
fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
let mut attempt: u32 = 0;
let deleter = { || self.inner.delete(ctx) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Delete,
err,
retry_after: dur,
attempt,
})
})
.call()
.map_err(|err| err.set_persistent())?;
Ok(RetryWrapper::new(
deleter,
self.notify.clone(),
self.builder,
))
}
fn copy(
&self,
ctx: &OperationContext,
from: &str,
to: &str,
args: OpCopy,
) -> Result<Self::Copier> {
let mut attempt: u32 = 0;
let copier = { || self.inner.copy(ctx, from, to, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Copy,
err,
retry_after: dur,
attempt,
})
})
.call()
.map_err(|err| err.set_persistent())?;
Ok(RetryWrapper::new(copier, self.notify.clone(), self.builder))
}
async fn rename(
&self,
ctx: &OperationContext,
from: &str,
to: &str,
args: OpRename,
) -> Result<RpRename> {
let mut attempt: u32 = 0;
{ || self.inner.rename(ctx, from, to, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Rename,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|err| err.set_persistent())
}
async fn restore(
&self,
ctx: &OperationContext,
path: &str,
args: OpRestore,
) -> Result<RpRestore> {
let mut attempt: u32 = 0;
{ || self.inner.restore(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Restore,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|err| err.set_persistent())
}
fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
let mut attempt: u32 = 0;
let lister = { || self.inner.list(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::List,
err,
retry_after: dur,
attempt,
})
})
.call()
.map_err(|err| err.set_persistent())?;
Ok(RetryWrapper::new(lister, self.notify.clone(), self.builder))
}
async fn presign(
&self,
ctx: &OperationContext,
path: &str,
args: OpPresign,
) -> Result<RpPresign> {
let mut attempt: u32 = 0;
{ || self.inner.presign(ctx, path, args.clone()) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Presign,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|err| err.set_persistent())
}
}
#[doc(hidden)]
pub struct RetryReader<R, I> {
inner: Arc<R>,
notify: Arc<I>,
builder: ExponentialBuilder,
}
impl<R, I> RetryReader<R, I> {
fn new(inner: R, notify: Arc<I>, builder: ExponentialBuilder) -> Self {
Self {
inner: Arc::new(inner),
notify,
builder,
}
}
}
impl<R: oio::Read + 'static, I: RetryInterceptor> oio::Read for RetryReader<R, I> {
async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
use backon::Retryable;
let mut attempt: u32 = 0;
let (rp, stream) = { || self.inner.open(range) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Read,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|e| e.set_persistent())?;
Ok((
rp,
Box::new(RetryReadStream::new(
self.inner.clone(),
stream,
range,
self.notify.clone(),
self.builder,
)) as Box<dyn oio::ReadStreamDyn>,
))
}
async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
use backon::Retryable;
let mut attempt: u32 = 0;
{ || self.inner.read(range) }
.retry(self.builder)
.when(|e| e.is_temporary())
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Read,
err,
retry_after: dur,
attempt,
})
})
.await
.map_err(|e| e.set_persistent())
}
}
#[doc(hidden)]
pub struct RetryReadStream<R, I> {
reader: Arc<R>,
stream: Option<Box<dyn oio::ReadStreamDyn>>,
range: BytesRange,
read: u64,
notify: Arc<I>,
builder: ExponentialBuilder,
}
impl<R, I> RetryReadStream<R, I> {
fn new(
reader: Arc<R>,
stream: Box<dyn oio::ReadStreamDyn>,
range: BytesRange,
notify: Arc<I>,
builder: ExponentialBuilder,
) -> Self {
Self {
reader,
stream: Some(stream),
range,
read: 0,
notify,
builder,
}
}
}
impl<R: oio::Read, I: RetryInterceptor> oio::ReadStream for RetryReadStream<R, I> {
async fn read(&mut self) -> Result<Buffer> {
use backon::RetryableWithContext;
let reader = self.reader.clone();
let stream = self.stream.take();
let range = self.range;
let read = self.read;
let mut attempt: u32 = 0;
let ((stream, range, read), res) = {
|(stream, mut range, mut read): (
Option<Box<dyn oio::ReadStreamDyn>>,
BytesRange,
u64,
)| {
let reader = reader.clone();
async move {
let mut stream = match stream {
Some(stream) => stream,
None => {
range.advance(read);
read = 0;
match reader.open(range).await {
Ok((_, stream)) => stream,
Err(err) => return ((None, range, read), Err(err)),
}
}
};
let res = match stream.read().await {
Ok(buf) => {
if !buf.is_empty() {
read += buf.len() as u64;
}
(Some(stream), Ok(buf))
}
Err(err) => (None, Err(err)),
};
((res.0, range, read), res.1)
}
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context((stream, range, read))
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Read,
err,
retry_after: dur,
attempt,
})
})
.await;
self.stream = stream;
self.range = range;
self.read = read;
res.map_err(|err| err.set_persistent())
}
}
#[doc(hidden)]
pub struct RetryWrapper<R, I> {
inner: Option<R>,
notify: Arc<I>,
builder: ExponentialBuilder,
}
impl<R, I> RetryWrapper<R, I> {
fn new(inner: R, notify: Arc<I>, backoff: ExponentialBuilder) -> Self {
Self {
inner: Some(inner),
notify,
builder: backoff,
}
}
fn take_inner(&mut self) -> Result<R> {
self.inner.take().ok_or_else(|| {
Error::new(
ErrorKind::Unexpected,
"retry layer is in bad state, please make sure future not dropped before ready",
)
})
}
}
impl<R: oio::ReadStream, I: RetryInterceptor> oio::ReadStream for RetryWrapper<R, I> {
async fn read(&mut self) -> Result<Buffer> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut r: R| async move {
let res = r.read().await;
(r, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Read,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
}
impl<R: oio::Write, I: RetryInterceptor> oio::Write for RetryWrapper<R, I> {
async fn write(&mut self, bs: Buffer) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let ((inner, _), res) = {
|(mut r, bs): (R, Buffer)| async move {
let res = r.write(bs.clone()).await;
((r, bs), res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context((inner, bs))
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Write,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let path = path.to_string();
let mut attempt: u32 = 0;
let ((inner, _, _, _), res) = {
|(mut r, path, args, range): (R, String, OpRead, BytesRange)| async move {
let res = r.copy_from(&path, args.clone(), range).await;
((r, path, args, range), res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context((inner, path, args, range))
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Write,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
async fn abort(&mut self) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut r: R| async move {
let res = r.abort().await;
(r, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Write,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
async fn close(&mut self) -> Result<Metadata> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut r: R| async move {
let res = r.close().await;
(r, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Write,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
}
impl<P: oio::List, I: RetryInterceptor> oio::List for RetryWrapper<P, I> {
async fn next(&mut self) -> Result<Option<oio::Entry>> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut p: P| async move {
let res = p.next().await;
(p, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::List,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
}
impl<P: oio::Delete, I: RetryInterceptor> oio::Delete for RetryWrapper<P, I> {
async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let path = path.to_string();
let args_cloned = args.clone();
let mut attempt: u32 = 0;
let (inner, res) = {
|mut p: P| {
let path = path.clone();
let args = args_cloned.clone();
async move {
let res = p.delete(&path, args).await;
(p, res)
}
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Delete,
err,
retry_after: dur,
attempt,
});
})
.await;
self.inner = Some(inner);
res.map_err(|e| e.set_persistent())
}
async fn close(&mut self) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut p: P| async move {
let res = p.close().await;
(p, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Delete,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
}
impl<C: oio::Copy, I: RetryInterceptor> oio::Copy for RetryWrapper<C, I> {
async fn next(&mut self) -> Result<Option<usize>> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut c: C| async move {
let res = c.next().await;
(c, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Copy,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
async fn close(&mut self) -> Result<Metadata> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut c: C| async move {
let res = c.close().await;
(c, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Copy,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
async fn abort(&mut self) -> Result<()> {
use backon::RetryableWithContext;
let inner = self.take_inner()?;
let mut attempt: u32 = 0;
let (inner, res) = {
|mut c: C| async move {
let res = c.abort().await;
(c, res)
}
}
.retry(self.builder)
.when(|e| e.is_temporary())
.context(inner)
.notify(|err, dur| {
attempt += 1;
self.notify.intercept(RetryEvent {
op: Operation::Copy,
err,
retry_after: dur,
attempt,
})
})
.await;
self.inner = Some(inner);
res.map_err(|err| err.set_persistent())
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use bytes::Bytes;
use futures::TryStreamExt;
use futures::stream;
use logforth::append::Testing;
use logforth::filter::rustlog::RustLogFilterBuilder;
use logforth::layout::TextLayout;
use opendal_layer_logging::LoggingLayer;
use super::*;
#[derive(Default, Clone)]
struct MockBuilder {
attempt: Arc<Mutex<usize>>,
}
impl Builder for MockBuilder {
type Config = ();
fn build(self) -> Result<impl Service> {
Ok(MockService {
attempt: self.attempt,
})
}
}
#[derive(Debug, Clone, Default)]
struct MockService {
attempt: Arc<Mutex<usize>>,
}
pub struct MockReader {
backend: MockService,
}
impl MockReader {
fn new(backend: MockService, _: &str, _: OpRead) -> Self {
Self { backend }
}
}
impl oio::StreamRead for MockReader {
async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
let backend = &self.backend;
let rp = RpRead::new({
let metadata = MetadataBuilder::file(13);
metadata.build()
});
let stream = MockReadStream {
buf: Bytes::from("Hello, World!").into(),
range,
attempt: backend.attempt.clone(),
};
Ok((rp, Box::new(stream) as Box<dyn oio::ReadStreamDyn>))
}
}
impl Service for MockService {
type Reader = oio::StreamReader<MockReader>;
type Writer = MockWriter;
type Lister = MockLister;
type Deleter = MockDeleter;
type Copier = MockCopier;
type Composer = ();
fn info(&self) -> ServiceInfo {
ServiceInfo::with_scheme("mock")
}
fn capability(&self) -> Capability {
Capability {
read: true,
write: true,
write_can_multi: true,
delete: true,
delete_max_size: Some(10),
stat: true,
list: true,
list_with_recursive: true,
copy: true,
..Default::default()
}
}
async fn create_dir(
&self,
_: &OperationContext,
_: &str,
_: OpCreateDir,
) -> Result<RpCreateDir> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
Ok(RpStat::new({
let metadata = MetadataBuilder::file(13);
metadata.build()
}))
}
fn read(&self, _: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
Ok(oio::StreamReader::new(MockReader::new(
self.clone(),
path,
args,
)))
}
fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
Ok(MockDeleter {
size: 0,
attempt: self.attempt.clone(),
})
}
fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
Ok(MockWriter {})
}
fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
let lister = MockLister::default();
Ok(lister)
}
fn copy(&self, _: &OperationContext, _: &str, _: &str, _: OpCopy) -> Result<Self::Copier> {
Ok(MockCopier {
attempt: self.attempt.clone(),
})
}
async fn rename(
&self,
_: &OperationContext,
_: &str,
_: &str,
_: OpRename,
) -> Result<RpRename> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
}
#[derive(Debug, Clone, Default)]
struct MockReadStream {
buf: Buffer,
range: BytesRange,
attempt: Arc<Mutex<usize>>,
}
impl oio::ReadStream for MockReadStream {
async fn read(&mut self) -> Result<Buffer> {
let mut attempt = self.attempt.lock().unwrap();
*attempt += 1;
match *attempt {
1 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from reader")
.set_temporary(),
),
2 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from reader")
.set_temporary(),
),
3 => Ok(self.buf.slice(self.range.to_range_as_usize())),
4 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from reader")
.set_temporary(),
),
5 => Ok(self.buf.slice(self.range.to_range_as_usize())),
_ => unreachable!(),
}
}
}
#[derive(Debug, Clone, Default)]
struct MockWriter {}
impl oio::Write for MockWriter {
async fn write(&mut self, _: Buffer) -> Result<()> {
Ok(())
}
async fn close(&mut self) -> Result<Metadata> {
Err(Error::new(ErrorKind::Unexpected, "always close failed").set_temporary())
}
async fn abort(&mut self) -> Result<()> {
Ok(())
}
}
#[derive(Debug, Clone, Default)]
struct MockLister {
attempt: usize,
}
impl oio::List for MockLister {
async fn next(&mut self) -> Result<Option<oio::Entry>> {
self.attempt += 1;
match self.attempt {
1 => Err(Error::new(
ErrorKind::RateLimited,
"retryable rate limited error from lister",
)
.set_temporary()),
2 => Ok(Some(oio::Entry::new(
"hello",
MetadataBuilder::file(0).build(),
))),
3 => Ok(Some(oio::Entry::new(
"world",
MetadataBuilder::file(0).build(),
))),
4 => Err(
Error::new(ErrorKind::Unexpected, "retryable internal server error")
.set_temporary(),
),
5 => Ok(Some(oio::Entry::new(
"2023/",
MetadataBuilder::dir().build(),
))),
6 => Ok(Some(oio::Entry::new(
"0208/",
MetadataBuilder::dir().build(),
))),
7 => Ok(None),
_ => {
unreachable!()
}
}
}
}
#[derive(Debug, Clone, Default)]
struct MockDeleter {
size: usize,
attempt: Arc<Mutex<usize>>,
}
impl oio::Delete for MockDeleter {
async fn delete(&mut self, _: &str, _: OpDelete) -> Result<()> {
self.size += 1;
Ok(())
}
async fn close(&mut self) -> Result<()> {
let mut attempt = self.attempt.lock().unwrap();
*attempt += 1;
match *attempt {
1 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
.set_temporary(),
),
2..=4 => {
self.size = self.size.saturating_sub(1);
Err(
Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
.set_temporary(),
)
}
5 => {
self.size = self.size.saturating_sub(1);
if self.size == 0 {
Ok(())
} else {
Err(
Error::new(ErrorKind::Unexpected, "retryable_error from deleter")
.set_temporary(),
)
}
}
_ => unreachable!(),
}
}
}
#[derive(Debug, Clone, Default)]
struct MockCopier {
attempt: Arc<Mutex<usize>>,
}
impl oio::Copy for MockCopier {
async fn next(&mut self) -> Result<Option<usize>> {
let mut attempt = self.attempt.lock().unwrap();
*attempt += 1;
match *attempt {
1 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from copier")
.set_temporary(),
),
2 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from copier")
.set_temporary(),
),
3 => Ok(Some(8)),
4 => Err(
Error::new(ErrorKind::Unexpected, "retryable_error from copier")
.set_temporary(),
),
5 => Ok(Some(5)),
6 => Ok(None),
_ => unreachable!(),
}
}
async fn close(&mut self) -> Result<Metadata> {
Ok(MetadataBuilder::unknown().build())
}
async fn abort(&mut self) -> Result<()> {
Ok(())
}
}
fn setup() {
let _ = logforth::starter_log::builder()
.dispatch(|d| {
d.filter(RustLogFilterBuilder::from_default_env().build())
.append(Testing::default().with_layout(TextLayout::default()))
})
.try_apply();
}
#[tokio::test]
async fn test_retry_read() -> Result<()> {
setup();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?
.layer(LoggingLayer::default())
.layer(RetryLayer::default());
let r = op.reader("retryable_error").await?;
let mut content = Vec::new();
let size = r
.read_into(&mut content, ..)
.await
.expect("read must succeed");
assert_eq!(size, 13);
assert_eq!(content, "Hello, World!".as_bytes());
assert_eq!(*builder.attempt.lock().unwrap(), 5);
Ok(())
}
#[tokio::test]
async fn test_retry_write_fail_on_close() -> Result<()> {
setup();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?
.layer(
RetryLayer::default()
.with_min_delay(Duration::from_millis(1))
.with_max_delay(Duration::from_millis(1))
.with_jitter(),
)
.layer(LoggingLayer::default());
let mut w = op.writer("test_write").await?;
w.write("aaa").await?;
w.write("bbb").await?;
match w.close().await {
Ok(_) => (),
Err(_) => {
w.abort().await?;
}
};
Ok(())
}
#[tokio::test]
async fn test_retry_list() -> Result<()> {
setup();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?.layer(RetryLayer::default());
let expected = vec!["hello", "world", "2023/", "0208/"];
let mut lister = op
.lister("retryable_error/")
.await
.expect("service must support list");
let mut actual = Vec::new();
while let Some(obj) = lister.try_next().await.expect("must success") {
actual.push(obj.name().to_owned());
}
assert_eq!(actual, expected);
Ok(())
}
#[tokio::test]
async fn test_retry_event_attempt_and_op() -> Result<()> {
setup();
#[derive(Default, Clone)]
struct Recorder {
events: Arc<Mutex<Vec<(Operation, u32)>>>,
}
impl RetryInterceptor for Recorder {
fn intercept(&self, event: RetryEvent<'_>) {
self.events.lock().unwrap().push((event.op, event.attempt));
}
}
let recorder = Recorder::default();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?.layer(
RetryLayer::default()
.with_min_delay(Duration::from_millis(1))
.with_max_delay(Duration::from_millis(1))
.with_notify(recorder.clone()),
);
let r = op.reader("retryable_error").await?;
let mut content = Vec::new();
let _ = r.read_into(&mut content, ..).await?;
let events = recorder.events.lock().unwrap().clone();
assert_eq!(
events,
vec![
(Operation::Read, 1),
(Operation::Read, 2),
(Operation::Read, 1),
],
);
Ok(())
}
#[tokio::test]
async fn test_retry_read_stream_error_enters_retry_budget() -> Result<()> {
setup();
#[derive(Default, Clone)]
struct Recorder {
events: Arc<Mutex<Vec<(Operation, u32)>>>,
}
impl RetryInterceptor for Recorder {
fn intercept(&self, event: RetryEvent<'_>) {
self.events.lock().unwrap().push((event.op, event.attempt));
}
}
let recorder = Recorder::default();
let backend = MockService::default();
let reader = RetryReader::new(
oio::StreamReader::new(MockReader::new(
backend.clone(),
"retryable_error",
OpRead::default(),
)),
Arc::new(recorder.clone()),
ExponentialBuilder::default()
.with_min_delay(Duration::from_millis(1))
.with_max_delay(Duration::from_millis(1)),
);
let (_, mut stream) = oio::Read::open(&reader, BytesRange::default()).await?;
let buf = oio::ReadStream::read_all(&mut stream).await?;
assert_eq!(buf.to_bytes(), Bytes::from_static(b"Hello, World!"));
assert_eq!(*backend.attempt.lock().unwrap(), 5);
let events = recorder.events.lock().unwrap().clone();
assert_eq!(
events,
vec![
(Operation::Read, 1),
(Operation::Read, 2),
(Operation::Read, 1),
],
);
Ok(())
}
#[tokio::test]
async fn test_retry_batch() -> Result<()> {
setup();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?.layer(
RetryLayer::default()
.with_min_delay(Duration::from_secs_f32(0.1))
.with_max_times(5),
);
let paths = vec!["hello", "world", "test", "batch"];
op.delete_stream(stream::iter(paths)).await?;
assert_eq!(*builder.attempt.lock().unwrap(), 5);
Ok(())
}
#[tokio::test]
async fn test_retry_copy() -> Result<()> {
setup();
#[derive(Default, Clone)]
struct Recorder {
events: Arc<Mutex<Vec<(Operation, u32)>>>,
}
impl RetryInterceptor for Recorder {
fn intercept(&self, event: RetryEvent<'_>) {
self.events.lock().unwrap().push((event.op, event.attempt));
}
}
let recorder = Recorder::default();
let builder = MockBuilder::default();
let op = Operator::new(builder.clone())?.layer(
RetryLayer::default()
.with_min_delay(Duration::from_millis(1))
.with_max_delay(Duration::from_millis(1))
.with_notify(recorder.clone()),
);
op.copy("from", "to").await.expect("copy must succeed");
assert_eq!(*builder.attempt.lock().unwrap(), 6);
let events = recorder.events.lock().unwrap().clone();
assert_eq!(
events,
vec![
(Operation::Copy, 1),
(Operation::Copy, 2),
(Operation::Copy, 1),
],
);
Ok(())
}
}