Skip to main content

opendal_service_vercel_blob/
backend.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::fmt::Debug;
19use std::sync::Arc;
20
21use bytes::Buf;
22use http::StatusCode;
23use log::debug;
24
25use super::VERCEL_BLOB_SCHEME;
26use super::config::VercelBlobConfig;
27use super::core::Blob;
28use super::core::VercelBlobCore;
29use super::core::parse_blob;
30use super::core::parse_error;
31use super::deleter::VercelBlobDeleter;
32use super::lister::VercelBlobLister;
33use super::reader::*;
34use super::writer::VercelBlobWriter;
35use super::writer::VercelBlobWriters;
36use opendal_core::raw::*;
37use opendal_core::*;
38
39/// [VercelBlob](https://vercel.com/docs/storage/vercel-blob) services support.
40#[doc = include_str!("docs.md")]
41#[derive(Default)]
42pub struct VercelBlobBuilder {
43    pub(super) config: VercelBlobConfig,
44}
45
46impl Debug for VercelBlobBuilder {
47    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48        f.debug_struct("VercelBlobBuilder")
49            .field("config", &self.config)
50            .finish_non_exhaustive()
51    }
52}
53
54impl VercelBlobBuilder {
55    /// Set root of this backend.
56    ///
57    /// All operations will happen under this root.
58    pub fn root(mut self, root: &str) -> Self {
59        self.config.root = if root.is_empty() {
60            None
61        } else {
62            Some(root.to_string())
63        };
64
65        self
66    }
67
68    /// Vercel Blob token.
69    ///
70    /// Get from Vercel environment variable `BLOB_READ_WRITE_TOKEN`.
71    /// It is required.
72    pub fn token(mut self, token: &str) -> Self {
73        if !token.is_empty() {
74            self.config.token = Some(token.to_string());
75        }
76        self
77    }
78}
79
80impl Builder for VercelBlobBuilder {
81    type Config = VercelBlobConfig;
82
83    /// Builds the backend and returns the result of VercelBlobBackend.
84    fn build(self) -> Result<impl Service> {
85        debug!("backend build started: {:?}", self);
86
87        let root = normalize_root(&self.config.root.clone().unwrap_or_default());
88        debug!("backend use root {}", root);
89
90        // Handle token.
91        let Some(token) = self.config.token.clone() else {
92            return Err(Error::new(ErrorKind::ConfigInvalid, "token is empty")
93                .with_operation("Builder::build")
94                .with_context("service", VERCEL_BLOB_SCHEME));
95        };
96
97        Ok(VercelBlobBackend {
98            core: Arc::new(VercelBlobCore {
99                info: ServiceInfo::new(VERCEL_BLOB_SCHEME, &root, ""),
100                capability: Capability {
101                    stat: true,
102
103                    read: true,
104                    read_with_suffix: true,
105
106                    write: true,
107                    write_can_empty: true,
108                    write_can_multi: true,
109                    write_multi_min_size: Some(5 * 1024 * 1024),
110
111                    copy: true,
112
113                    list: true,
114                    list_with_limit: true,
115
116                    delete: true,
117
118                    shared: true,
119
120                    ..Default::default()
121                },
122                root,
123                token,
124            }),
125        })
126    }
127}
128
129/// Backend for VercelBlob services.
130#[derive(Debug, Clone)]
131pub struct VercelBlobBackend {
132    pub(crate) core: Arc<VercelBlobCore>,
133}
134
135impl Service for VercelBlobBackend {
136    type Reader = oio::StreamReader<VercelBlobReader>;
137    type Writer = VercelBlobWriters;
138    type Lister = oio::PageLister<VercelBlobLister>;
139    type Deleter = oio::OneShotDeleter<VercelBlobDeleter>;
140    type Copier = oio::OneShotCopier;
141
142    fn info(&self) -> ServiceInfo {
143        self.core.info.clone()
144    }
145
146    fn capability(&self) -> Capability {
147        self.core.capability
148    }
149
150    async fn create_dir(
151        &self,
152        _ctx: &OperationContext,
153        _path: &str,
154        _args: OpCreateDir,
155    ) -> Result<RpCreateDir> {
156        Err(Error::new(
157            ErrorKind::Unsupported,
158            "operation is not supported",
159        ))
160    }
161
162    async fn stat(&self, ctx: &OperationContext, path: &str, _args: OpStat) -> Result<RpStat> {
163        let resp = self.core.head(ctx, path).await?;
164
165        let status = resp.status();
166
167        match status {
168            StatusCode::OK => {
169                let bs = resp.into_body();
170
171                let resp: Blob =
172                    serde_json::from_reader(bs.reader()).map_err(new_json_deserialize_error)?;
173
174                parse_blob(&resp).map(RpStat::new)
175            }
176            _ => Err(parse_error(resp)),
177        }
178    }
179    fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
180        let output: oio::StreamReader<VercelBlobReader> = {
181            Ok(oio::StreamReader::new(VercelBlobReader::new(
182                self.clone(),
183                ctx.clone(),
184                path,
185                args,
186            )))
187        }?;
188
189        Ok(output)
190    }
191
192    fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
193        let output: VercelBlobWriters = {
194            let concurrent = args.concurrent();
195            let writer =
196                VercelBlobWriter::new(self.core.clone(), ctx.clone(), args, path.to_string());
197
198            let w = oio::MultipartWriter::new(ctx.executor().clone(), writer, concurrent);
199
200            Ok(w)
201        }?;
202
203        Ok(output)
204    }
205
206    fn copy(
207        &self,
208        ctx: &OperationContext,
209        from: &str,
210        to: &str,
211        _args: OpCopy,
212        _opts: OpCopier,
213    ) -> Result<Self::Copier> {
214        let core = self.core.clone();
215        let ctx = ctx.clone();
216        let from = from.to_string();
217        let to = to.to_string();
218
219        Ok(oio::OneShotCopier::new(async move {
220            let resp = core.copy(&ctx, &from, &to).await?;
221            let status = resp.status();
222
223            match status {
224                StatusCode::OK => Ok(Metadata::default()),
225                _ => Err(parse_error(resp)),
226            }
227        }))
228    }
229
230    fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
231        let output: oio::PageLister<VercelBlobLister> = {
232            let l = VercelBlobLister::new(self.core.clone(), ctx.clone(), path, args.limit());
233            Ok(oio::PageLister::new(l))
234        }?;
235
236        Ok(output)
237    }
238
239    fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
240        let output: oio::OneShotDeleter<VercelBlobDeleter> = {
241            Ok(oio::OneShotDeleter::new(VercelBlobDeleter::new(
242                self.core.clone(),
243                ctx.clone(),
244            )))
245        }?;
246
247        Ok(output)
248    }
249
250    async fn rename(
251        &self,
252        _ctx: &OperationContext,
253        _from: &str,
254        _to: &str,
255        _args: OpRename,
256    ) -> Result<RpRename> {
257        Err(Error::new(
258            ErrorKind::Unsupported,
259            "operation is not supported",
260        ))
261    }
262
263    async fn presign(
264        &self,
265        _ctx: &OperationContext,
266        _path: &str,
267        _args: OpPresign,
268    ) -> Result<RpPresign> {
269        Err(Error::new(
270            ErrorKind::Unsupported,
271            "operation is not supported",
272        ))
273    }
274}