Skip to main content

opendal_layer_timeout/
lib.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18#![doc = include_str!("../README.md")]
19#![cfg_attr(docsrs, feature(doc_cfg))]
20#![cfg_attr(docsrs, doc(auto_cfg))]
21#![deny(missing_docs)]
22use std::future::Future;
23use std::sync::Arc;
24use std::time::Duration;
25
26use opendal_core::raw::*;
27use opendal_core::*;
28
29/// `TimeoutLayer` adds deadlines to operations so slow or stalled work cannot
30/// hang indefinitely.
31///
32/// For example, a dead connection could stall a database query. `TimeoutLayer`
33/// interrupts the operation and returns an error that users can retry or report.
34///
35/// # Notes
36///
37/// `TimeoutLayer` applies two timeout budgets:
38///
39/// - `timeout` bounds control operations such as `stat`, `create_dir`, `rename`,
40///   and `presign`.
41/// - `io_timeout` bounds operations that open IO bodies, such as `read`, `write`,
42///   and `list`, and every method call on returned readers, writers, listers,
43///   deleters, and copiers.
44///
45/// # Default
46///
47/// - timeout: 60 seconds
48/// - io_timeout: 10 seconds
49///
50/// # Cancellation Safety
51///
52/// `TimeoutLayer` enforces deadlines by dropping the in-flight future when a
53/// timeout is reached. This can break lower layers that rely on a future being
54/// resolved to restore internal state.
55///
56/// For example, while using `TimeoutLayer` with `RetryLayer` at the same time,
57/// please make sure timeout layer is added before retry layer.
58///
59/// ```no_run
60/// # use std::time::Duration;
61/// #
62/// # use opendal_core::services;
63/// # use opendal_core::Operator;
64/// # use opendal_core::Result;
65/// # use opendal_layer_retry::RetryLayer;
66/// # use opendal_layer_timeout::TimeoutLayer;
67/// #
68/// # fn main() -> Result<()> {
69/// let op = Operator::new(services::Memory::default())?
70///     // This is fine: each retry attempt is timed out.
71///     .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)))
72///     .layer(RetryLayer::default())
73///     // This is wrong: timeout can drop RetryLayer's future before it restores body state.
74///     .layer(TimeoutLayer::default().with_io_timeout(Duration::from_nanos(1)));
75/// # Ok(())
76/// # }
77/// ```
78///
79/// # Examples
80///
81/// The following example creates a timeout layer with a 10-second timeout for
82/// control operations and a 3-second timeout for IO operations.
83///
84/// ```no_run
85/// # use std::time::Duration;
86/// #
87/// # use opendal_core::services;
88/// # use opendal_core::Operator;
89/// # use opendal_core::Result;
90/// # use opendal_layer_timeout::TimeoutLayer;
91/// #
92/// # fn main() -> Result<()> {
93/// let _ = Operator::new(services::Memory::default())?
94///     .layer(
95///         TimeoutLayer::default()
96///             .with_timeout(Duration::from_secs(10))
97///             .with_io_timeout(Duration::from_secs(3)),
98///     );
99/// # Ok(())
100/// # }
101/// ```
102///
103/// # Implementation Notes
104///
105/// `TimeoutLayer` uses [`tokio::time::timeout`] to bound service calls and IO
106/// body methods. It also supplies an executor timeout so concurrent block write
107/// and copy tasks can fail instead of waiting forever.
108///
109/// This introduces a small amount of overhead for IO operations, but it is needed
110/// to implement timeouts correctly. OpenDAL used to implement this as a
111/// zero-cost deadline check that only stored an [`Instant`] and compared it with
112/// the current time. However, that approach does not work for all cases.
113///
114/// For example, a user's TCP connection could enter the
115/// [Busy ESTAB](https://blog.cloudflare.com/when-tcp-sockets-refuse-to-die)
116/// state. In this state, the connection emits no IO events, so the runtime never
117/// polls the future again. The future hangs until Linux closes the connection
118/// after it reaches the
119/// [net.ipv4.tcp_retries2](https://man7.org/linux/man-pages/man7/tcp.7.html)
120/// limit.
121#[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    /// Create a new [`TimeoutLayer`] with default settings.
138    pub fn new() -> Self {
139        Self::default()
140    }
141
142    /// Set the timeout for control operations.
143    ///
144    /// This timeout is for all non-io operations like `stat`, `delete`.
145    pub fn with_timeout(mut self, timeout: Duration) -> Self {
146        self.timeout = timeout;
147        self
148    }
149
150    /// Set the timeout for IO operations and body methods.
151    ///
152    /// This timeout is for all io operations like `read`, `Reader::read` and `Writer::write`.
153    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        // Concurrent block IO paths read this timeout from the operation context's executor.
166        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        /// Return a reader whose operations never complete.
511        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        /// Return a deleter whose operations never complete.
523        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}