Skip to main content

opendal_core/raw/oio/write/
api.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
18use std::future::Future;
19use std::ops::DerefMut;
20
21use crate::raw::*;
22use crate::*;
23
24/// Writer is a type-erased [`Write`].
25pub type Writer = Box<dyn WriteDyn>;
26
27/// Write is the async sink used by services and layers.
28pub trait Write: Unpin + Send + Sync {
29    /// Write the entire buffer into the writer.
30    ///
31    /// `Ok(())` means all bytes from `bs` have been accepted. Implementations
32    /// must return an error instead of treating a partial write as success.
33    fn write(&mut self, bs: Buffer) -> impl Future<Output = Result<()>> + MaybeSend;
34
35    /// Copy one absolute bounded source range into this writer.
36    ///
37    /// Callers invoke this operation only when the composed service declares
38    /// [`Capability::write_can_copy_from`]. Every error is an execution failure
39    /// and must not trigger streaming fallback.
40    fn copy_from(
41        &mut self,
42        _path: &str,
43        _args: OpRead,
44        _range: BytesRange,
45    ) -> impl Future<Output = Result<()>> + MaybeSend {
46        async {
47            Err(Error::new(
48                ErrorKind::Unsupported,
49                "writer doesn't support native copy",
50            ))
51        }
52    }
53
54    /// Close the writer and make sure all data has been flushed.
55    fn close(&mut self) -> impl Future<Output = Result<Metadata>> + MaybeSend;
56
57    /// Abort the pending writer.
58    fn abort(&mut self) -> impl Future<Output = Result<()>> + MaybeSend;
59}
60
61impl Write for () {
62    async fn write(&mut self, _: Buffer) -> Result<()> {
63        unimplemented!("write is required to be implemented for oio::Write")
64    }
65
66    async fn close(&mut self) -> Result<Metadata> {
67        Err(Error::new(
68            ErrorKind::Unsupported,
69            "output writer doesn't support close",
70        ))
71    }
72
73    async fn abort(&mut self) -> Result<()> {
74        Err(Error::new(
75            ErrorKind::Unsupported,
76            "output writer doesn't support abort",
77        ))
78    }
79}
80
81/// WriteDyn is the object-safe version of [`Write`] used by [`Writer`].
82pub trait WriteDyn: Unpin + Send + Sync {
83    /// The dyn version of [`Write::write`].
84    fn write_dyn(&mut self, bs: Buffer) -> BoxedFuture<'_, Result<()>>;
85
86    /// The dyn version of [`Write::copy_from`].
87    fn copy_from_dyn<'a>(
88        &'a mut self,
89        path: &'a str,
90        args: OpRead,
91        range: BytesRange,
92    ) -> BoxedFuture<'a, Result<()>>;
93
94    /// The dyn version of [`Write::close`].
95    fn close_dyn(&mut self) -> BoxedFuture<'_, Result<Metadata>>;
96
97    /// The dyn version of [`Write::abort`].
98    fn abort_dyn(&mut self) -> BoxedFuture<'_, Result<()>>;
99}
100
101impl<T: Write + ?Sized> WriteDyn for T {
102    fn write_dyn(&mut self, bs: Buffer) -> BoxedFuture<'_, Result<()>> {
103        Box::pin(self.write(bs))
104    }
105
106    fn copy_from_dyn<'a>(
107        &'a mut self,
108        path: &'a str,
109        args: OpRead,
110        range: BytesRange,
111    ) -> BoxedFuture<'a, Result<()>> {
112        Box::pin(self.copy_from(path, args, range))
113    }
114
115    fn close_dyn(&mut self) -> BoxedFuture<'_, Result<Metadata>> {
116        Box::pin(self.close())
117    }
118
119    fn abort_dyn(&mut self) -> BoxedFuture<'_, Result<()>> {
120        Box::pin(self.abort())
121    }
122}
123
124impl<T: WriteDyn + ?Sized> Write for Box<T> {
125    async fn write(&mut self, bs: Buffer) -> Result<()> {
126        self.deref_mut().write_dyn(bs).await
127    }
128
129    async fn copy_from(&mut self, path: &str, args: OpRead, range: BytesRange) -> Result<()> {
130        self.deref_mut().copy_from_dyn(path, args, range).await
131    }
132
133    async fn close(&mut self) -> Result<Metadata> {
134        self.deref_mut().close_dyn().await
135    }
136
137    async fn abort(&mut self) -> Result<()> {
138        self.deref_mut().abort_dyn().await
139    }
140}