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::future::Future;
23use std::sync::Arc;
24use std::time::Duration;
25
26use opendal_core::raw::*;
27use opendal_core::*;
28
29#[derive(Clone, Debug)]
122pub struct TimeoutLayer {
123 timeout: Duration,
124 io_timeout: Duration,
125}
126
127impl Default for TimeoutLayer {
128 fn default() -> Self {
129 Self {
130 timeout: Duration::from_secs(60),
131 io_timeout: Duration::from_secs(10),
132 }
133 }
134}
135
136impl TimeoutLayer {
137 pub fn new() -> Self {
139 Self::default()
140 }
141
142 pub fn with_timeout(mut self, timeout: Duration) -> Self {
146 self.timeout = timeout;
147 self
148 }
149
150 pub fn with_io_timeout(mut self, timeout: Duration) -> Self {
154 self.io_timeout = timeout;
155 self
156 }
157}
158
159impl Layer for TimeoutLayer {
160 fn apply_service(&self, inner: Servicer) -> Servicer {
161 Arc::new(self.layer(inner))
162 }
163
164 fn apply_context(&self, _srv: Servicer, inner: OperationContext) -> OperationContext {
165 let executor = Executor::with(TimeoutExecutor::new(
167 inner.executor().clone().into_inner(),
168 self.io_timeout,
169 ));
170 inner.with_executor(executor)
171 }
172}
173
174impl TimeoutLayer {
175 fn layer(&self, inner: Servicer) -> TimeoutService {
176 TimeoutService {
177 inner,
178 timeout: self.timeout,
179 io_timeout: self.io_timeout,
180 }
181 }
182}
183
184#[doc(hidden)]
185#[derive(Debug)]
186pub struct TimeoutService {
187 inner: Servicer,
188 timeout: Duration,
189 io_timeout: Duration,
190}
191
192impl TimeoutService {
193 async fn timeout<F: Future<Output = Result<T>>, T>(&self, op: Operation, fut: F) -> Result<T> {
194 tokio::time::timeout(self.timeout, fut).await.map_err(|_| {
195 Error::new(ErrorKind::Unexpected, "operation timeout reached")
196 .with_operation(op)
197 .with_context("timeout", self.timeout.as_secs_f64().to_string())
198 .set_temporary()
199 })?
200 }
201}
202
203impl Service for TimeoutService {
204 type Reader = TimeoutWrapper<oio::Reader>;
205 type Writer = TimeoutWrapper<oio::Writer>;
206 type Lister = TimeoutWrapper<oio::Lister>;
207 type Deleter = TimeoutWrapper<oio::Deleter>;
208 type Copier = TimeoutWrapper<oio::Copier>;
209 type Composer = TimeoutWrapper<oio::Composer>;
210
211 fn info(&self) -> ServiceInfo {
212 self.inner.info()
213 }
214
215 fn capability(&self) -> Capability {
216 self.inner.capability()
217 }
218
219 async fn create_dir(
220 &self,
221 ctx: &OperationContext,
222 path: &str,
223 args: OpCreateDir,
224 ) -> Result<RpCreateDir> {
225 self.timeout(Operation::CreateDir, self.inner.create_dir(ctx, path, args))
226 .await
227 }
228
229 fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
230 self.inner
231 .read(ctx, path, args)
232 .map(|r| TimeoutWrapper::new(r, self.io_timeout))
233 }
234
235 fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
236 self.inner
237 .write(ctx, path, args)
238 .map(|r| TimeoutWrapper::new(r, self.io_timeout))
239 }
240
241 fn copy(
242 &self,
243 ctx: &OperationContext,
244 from: &str,
245 to: &str,
246 args: OpCopy,
247 ) -> Result<Self::Copier> {
248 self.inner
249 .copy(ctx, from, to, args)
250 .map(|c| TimeoutWrapper::new(c, self.io_timeout))
251 }
252
253 fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
254 self.inner
255 .compose(ctx, to, args)
256 .map(|c| TimeoutWrapper::new(c, self.io_timeout))
257 }
258
259 async fn rename(
260 &self,
261 ctx: &OperationContext,
262 from: &str,
263 to: &str,
264 args: OpRename,
265 ) -> Result<RpRename> {
266 self.timeout(Operation::Rename, self.inner.rename(ctx, from, to, args))
267 .await
268 }
269
270 async fn restore(
271 &self,
272 ctx: &OperationContext,
273 path: &str,
274 args: OpRestore,
275 ) -> Result<RpRestore> {
276 self.timeout(Operation::Restore, self.inner.restore(ctx, path, args))
277 .await
278 }
279
280 async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
281 self.timeout(Operation::Stat, self.inner.stat(ctx, path, args))
282 .await
283 }
284
285 fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
286 self.inner
287 .delete(ctx)
288 .map(|r| TimeoutWrapper::new(r, self.io_timeout))
289 }
290
291 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
292 self.inner
293 .list(ctx, path, args)
294 .map(|r| TimeoutWrapper::new(r, self.io_timeout))
295 }
296
297 async fn presign(
298 &self,
299 ctx: &OperationContext,
300 path: &str,
301 args: OpPresign,
302 ) -> Result<RpPresign> {
303 self.timeout(Operation::Presign, self.inner.presign(ctx, path, args))
304 .await
305 }
306}
307
308struct TimeoutExecutor {
309 exec: Arc<dyn Execute>,
310 timeout: Duration,
311}
312
313impl TimeoutExecutor {
314 fn new(exec: Arc<dyn Execute>, timeout: Duration) -> Self {
315 Self { exec, timeout }
316 }
317}
318
319impl Execute for TimeoutExecutor {
320 fn execute(&self, f: BoxedStaticFuture<()>) {
321 self.exec.execute(f)
322 }
323
324 fn timeout(&self) -> Option<BoxedStaticFuture<()>> {
325 Some(Box::pin(tokio::time::sleep(self.timeout)))
326 }
327}
328
329#[doc(hidden)]
330pub struct TimeoutWrapper<R> {
331 inner: R,
332
333 timeout: Duration,
334}
335
336impl<R> TimeoutWrapper<R> {
337 fn new(inner: R, timeout: Duration) -> Self {
338 Self { inner, timeout }
339 }
340
341 #[inline]
342 async fn io_timeout<F: Future<Output = Result<T>>, T>(
343 timeout: Duration,
344 op: &'static str,
345 fut: F,
346 ) -> Result<T> {
347 tokio::time::timeout(timeout, fut).await.map_err(|_| {
348 Error::new(ErrorKind::Unexpected, "io operation timeout reached")
349 .with_operation(op)
350 .with_context("timeout", timeout.as_secs_f64().to_string())
351 .set_temporary()
352 })?
353 }
354}
355
356impl<R: oio::ReadStream> oio::ReadStream for TimeoutWrapper<R> {
357 async fn read(&mut self) -> Result<Buffer> {
358 let fut = self.inner.read();
359 Self::io_timeout(self.timeout, Operation::Read.into_static(), fut).await
360 }
361}
362
363impl<R: oio::Read> oio::Read for TimeoutWrapper<R> {
364 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
365 let (rp, stream) = Self::io_timeout(
366 self.timeout,
367 Operation::Read.into_static(),
368 self.inner.open(range),
369 )
370 .await?;
371 Ok((
372 rp,
373 Box::new(TimeoutWrapper::new(stream, self.timeout)) as Box<dyn oio::ReadStreamDyn>,
374 ))
375 }
376
377 async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
378 Self::io_timeout(
379 self.timeout,
380 Operation::Read.into_static(),
381 self.inner.read(range),
382 )
383 .await
384 }
385}
386
387impl<R: oio::Write> oio::Write for TimeoutWrapper<R> {
388 async fn write(&mut self, bs: Buffer) -> Result<()> {
389 let fut = self.inner.write(bs);
390 Self::io_timeout(self.timeout, Operation::Write.into_static(), fut).await
391 }
392
393 async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
394 let fut = self.inner.copy_from(path, args, range);
395 Self::io_timeout(self.timeout, Operation::Write.into_static(), fut).await
396 }
397
398 async fn close(&mut self) -> Result<Metadata> {
399 let fut = self.inner.close();
400 Self::io_timeout(self.timeout, Operation::Write.into_static(), fut).await
401 }
402
403 async fn abort(&mut self) -> Result<()> {
404 let fut = self.inner.abort();
405 Self::io_timeout(self.timeout, Operation::Write.into_static(), fut).await
406 }
407}
408
409impl<R: oio::List> oio::List for TimeoutWrapper<R> {
410 async fn next(&mut self) -> Result<Option<oio::Entry>> {
411 let fut = self.inner.next();
412 Self::io_timeout(self.timeout, Operation::List.into_static(), fut).await
413 }
414}
415
416impl<R: oio::Delete> oio::Delete for TimeoutWrapper<R> {
417 async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
418 let fut = self.inner.delete(path, args);
419 Self::io_timeout(self.timeout, Operation::Delete.into_static(), fut).await
420 }
421
422 async fn close(&mut self) -> Result<()> {
423 let fut = self.inner.close();
424 Self::io_timeout(self.timeout, Operation::Delete.into_static(), fut).await
425 }
426}
427
428impl<C: oio::Copy> oio::Copy for TimeoutWrapper<C> {
429 async fn next(&mut self) -> Result<Option<usize>> {
430 let fut = self.inner.next();
431 Self::io_timeout(self.timeout, Operation::Copy.into_static(), fut).await
432 }
433
434 async fn close(&mut self) -> Result<Metadata> {
435 let fut = self.inner.close();
436 Self::io_timeout(self.timeout, Operation::Copy.into_static(), fut).await
437 }
438
439 async fn abort(&mut self) -> Result<()> {
440 let fut = self.inner.abort();
441 Self::io_timeout(self.timeout, Operation::Copy.into_static(), fut).await
442 }
443}
444
445impl<C: oio::Compose> oio::Compose for TimeoutWrapper<C> {
446 async fn compose(&mut self, path: &str, args: OpRead) -> Result<()> {
447 let fut = self.inner.compose(path, args);
448 Self::io_timeout(self.timeout, Operation::Compose.into_static(), fut).await
449 }
450
451 async fn close(&mut self) -> Result<Metadata> {
452 let fut = self.inner.close();
453 Self::io_timeout(self.timeout, Operation::Compose.into_static(), fut).await
454 }
455}
456
457#[cfg(test)]
458mod tests {
459 use std::future::pending;
460
461 use futures::StreamExt;
462 use tokio::time::timeout;
463
464 use super::*;
465
466 #[derive(Debug, Clone, Default)]
467 struct MockService;
468
469 impl Service for MockService {
470 type Reader = MockReader;
471 type Writer = ();
472 type Lister = MockLister;
473 type Deleter = MockDeleter;
474 type Copier = MockCopier;
475 type Composer = ();
476
477 fn info(&self) -> ServiceInfo {
478 ServiceInfo::with_scheme("mock")
479 }
480
481 fn capability(&self) -> Capability {
482 Capability {
483 read: true,
484 delete: true,
485 list: true,
486 copy: true,
487 ..Default::default()
488 }
489 }
490
491 async fn create_dir(
492 &self,
493 _: &OperationContext,
494 _: &str,
495 _: OpCreateDir,
496 ) -> Result<RpCreateDir> {
497 Err(Error::new(
498 ErrorKind::Unsupported,
499 "operation is not supported",
500 ))
501 }
502
503 async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
504 Err(Error::new(
505 ErrorKind::Unsupported,
506 "operation is not supported",
507 ))
508 }
509
510 fn read(&self, _ctx: &OperationContext, _: &str, _: OpRead) -> Result<Self::Reader> {
512 Ok(MockReader)
513 }
514
515 fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
516 Err(Error::new(
517 ErrorKind::Unsupported,
518 "operation is not supported",
519 ))
520 }
521
522 fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
524 Ok(MockDeleter)
525 }
526
527 fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
528 Ok(MockLister)
529 }
530
531 fn copy(&self, _: &OperationContext, _: &str, _: &str, _: OpCopy) -> Result<Self::Copier> {
532 Ok(MockCopier)
533 }
534
535 async fn rename(
536 &self,
537 _: &OperationContext,
538 _: &str,
539 _: &str,
540 _: OpRename,
541 ) -> Result<RpRename> {
542 Err(Error::new(
543 ErrorKind::Unsupported,
544 "operation is not supported",
545 ))
546 }
547
548 async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
549 Err(Error::new(
550 ErrorKind::Unsupported,
551 "operation is not supported",
552 ))
553 }
554 }
555
556 #[derive(Debug, Clone, Default)]
557 struct MockReader;
558
559 impl oio::Read for MockReader {
560 async fn open(&self, _: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
561 pending().await
562 }
563
564 async fn read(&self, _: BytesRange) -> Result<(RpRead, Buffer)> {
565 pending().await
566 }
567 }
568
569 #[derive(Debug, Clone, Default)]
570 struct MockLister;
571
572 impl oio::List for MockLister {
573 async fn next(&mut self) -> Result<Option<oio::Entry>> {
574 pending().await
575 }
576 }
577
578 #[derive(Debug, Clone, Default)]
579 struct MockDeleter;
580
581 impl oio::Delete for MockDeleter {
582 async fn delete(&mut self, _: &str, _: OpDelete) -> Result<()> {
583 pending().await
584 }
585
586 async fn close(&mut self) -> Result<()> {
587 Ok(())
588 }
589 }
590
591 #[derive(Debug, Clone, Default)]
592 struct MockCopier;
593
594 impl oio::Copy for MockCopier {
595 async fn next(&mut self) -> Result<Option<usize>> {
596 pending().await
597 }
598
599 async fn close(&mut self) -> Result<Metadata> {
600 pending().await
601 }
602
603 async fn abort(&mut self) -> Result<()> {
604 pending().await
605 }
606 }
607
608 #[tokio::test]
609 async fn test_delete_timeout() {
610 let srv = MockService;
611 let op = Operator::from_parts(OperationContext::default(), Arc::new(srv))
612 .layer(TimeoutLayer::default().with_io_timeout(Duration::from_secs(1)));
613
614 let fut = async {
615 let res = op.delete("test").await;
616 assert!(res.is_err());
617 let err = res.unwrap_err();
618 assert_eq!(err.kind(), ErrorKind::Unexpected);
619 assert!(err.to_string().contains("timeout"))
620 };
621
622 timeout(Duration::from_secs(2), fut)
623 .await
624 .expect("this test should not exceed 2 seconds")
625 }
626
627 #[tokio::test]
628 async fn test_io_timeout() {
629 let srv = MockService;
630 let op = Operator::from_parts(OperationContext::default(), Arc::new(srv))
631 .layer(TimeoutLayer::default().with_io_timeout(Duration::from_secs(1)));
632
633 let reader = op.reader("test").await.unwrap();
634
635 let res = reader.read(0..4).await;
636 assert!(res.is_err());
637 let err = res.unwrap_err();
638 assert_eq!(err.kind(), ErrorKind::Unexpected);
639 assert!(err.to_string().contains("timeout"))
640 }
641
642 #[tokio::test]
643 async fn test_list_timeout() {
644 let srv = MockService;
645 let op = Operator::from_parts(OperationContext::default(), Arc::new(srv)).layer(
646 TimeoutLayer::default()
647 .with_timeout(Duration::from_secs(1))
648 .with_io_timeout(Duration::from_secs(1)),
649 );
650
651 let mut lister = op.lister("test").await.unwrap();
652
653 let res = lister.next().await.unwrap();
654 assert!(res.is_err());
655 let err = res.unwrap_err();
656 assert_eq!(err.kind(), ErrorKind::Unexpected);
657 assert!(err.to_string().contains("timeout"))
658 }
659
660 #[tokio::test]
661 async fn test_delete_io_timeout() {
662 use oio::Delete;
663
664 let mut deleter = TimeoutWrapper::new(MockDeleter, Duration::from_secs(1));
665
666 let res = deleter.delete("test", OpDelete::default()).await;
667 assert!(res.is_err());
668 let err = res.unwrap_err();
669 assert_eq!(err.kind(), ErrorKind::Unexpected);
670 assert!(err.to_string().contains("timeout"));
671 }
672
673 #[tokio::test]
674 async fn test_copy_io_timeout() {
675 use oio::Copy;
676
677 let service = TimeoutLayer::default()
678 .with_io_timeout(Duration::from_millis(100))
679 .apply_service(Arc::new(MockService));
680 let ctx = OperationContext::new();
681 let mut copier = service.copy(&ctx, "f", "t", OpCopy::default()).unwrap();
682
683 let err = copier.next().await.unwrap_err();
684 assert!(err.to_string().contains("timeout"));
685 }
686
687 #[tokio::test]
688 async fn test_list_timeout_raw() {
689 use oio::List;
690
691 let timeout_layer = TimeoutLayer::default()
692 .with_timeout(Duration::from_secs(1))
693 .with_io_timeout(Duration::from_secs(1));
694 let service = timeout_layer.apply_service(Arc::new(MockService));
695 let ctx = OperationContext::new();
696
697 let mut lister = service.list(&ctx, "test", OpList::default()).unwrap();
698
699 let res = lister.next().await;
700 assert!(res.is_err());
701 let err = res.unwrap_err();
702 assert_eq!(err.kind(), ErrorKind::Unexpected);
703 assert!(err.to_string().contains("timeout"));
704 }
705}