1use std::fmt::Debug;
19use std::sync::Arc;
20
21use http::StatusCode;
22use log::debug;
23use opendal_core::raw::*;
24use opendal_core::*;
25
26use super::UPYUN_SCHEME;
27use super::config::UpyunConfig;
28use super::core::parse_error;
29use super::core::*;
30use super::deleter::UpyunDeleter;
31use super::lister::UpyunLister;
32use super::reader::*;
33use super::writer::UpyunWriter;
34use super::writer::UpyunWriters;
35
36#[doc = include_str!("docs.md")]
38#[derive(Default)]
39pub struct UpyunBuilder {
40 pub(super) config: UpyunConfig,
41}
42
43impl Debug for UpyunBuilder {
44 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
45 f.debug_struct("UpyunBuilder")
46 .field("config", &self.config)
47 .finish_non_exhaustive()
48 }
49}
50
51impl UpyunBuilder {
52 pub fn root(mut self, root: &str) -> Self {
56 self.config.root = if root.is_empty() {
57 None
58 } else {
59 Some(root.to_string())
60 };
61
62 self
63 }
64
65 pub fn bucket(mut self, bucket: &str) -> Self {
69 self.config.bucket = bucket.to_string();
70
71 self
72 }
73
74 pub fn operator(mut self, operator: &str) -> Self {
78 self.config.operator = if operator.is_empty() {
79 None
80 } else {
81 Some(operator.to_string())
82 };
83
84 self
85 }
86
87 pub fn password(mut self, password: &str) -> Self {
91 self.config.password = if password.is_empty() {
92 None
93 } else {
94 Some(password.to_string())
95 };
96
97 self
98 }
99}
100
101impl Builder for UpyunBuilder {
102 type Config = UpyunConfig;
103
104 fn build(self) -> Result<impl Service> {
106 debug!("backend build started: {:?}", self);
107
108 let root = normalize_root(&self.config.root.clone().unwrap_or_default());
109 debug!("backend use root {}", root);
110
111 if self.config.bucket.is_empty() {
113 return Err(Error::new(ErrorKind::ConfigInvalid, "bucket is empty")
114 .with_operation("Builder::build")
115 .with_context("service", UPYUN_SCHEME));
116 }
117
118 debug!("backend use bucket {}", self.config.bucket);
119
120 let operator = match &self.config.operator {
121 Some(operator) => Ok(operator.clone()),
122 None => Err(Error::new(ErrorKind::ConfigInvalid, "operator is empty")
123 .with_operation("Builder::build")
124 .with_context("service", UPYUN_SCHEME)),
125 }?;
126
127 let password = match &self.config.password {
128 Some(password) => Ok(password.clone()),
129 None => Err(Error::new(ErrorKind::ConfigInvalid, "password is empty")
130 .with_operation("Builder::build")
131 .with_context("service", UPYUN_SCHEME)),
132 }?;
133
134 let signer = UpyunSigner {
135 operator: operator.clone(),
136 password: password.clone(),
137 };
138
139 Ok(UpyunBackend {
140 core: Arc::new(UpyunCore {
141 info: ServiceInfo::new(UPYUN_SCHEME, &root, ""),
142 capability: Capability {
143 stat: true,
144
145 create_dir: true,
146
147 read: true,
148 read_with_suffix: true,
149
150 write: true,
151 write_can_empty: true,
152 write_can_multi: true,
153 write_with_cache_control: true,
154 write_with_content_type: true,
155
156 write_multi_min_size: Some(1024 * 1024),
158 write_multi_max_size: Some(50 * 1024 * 1024),
159
160 delete: true,
161 rename: true,
162 copy: true,
163
164 list: true,
165 list_with_limit: true,
166
167 shared: true,
168
169 ..Default::default()
170 },
171 root,
172 operator,
173 bucket: self.config.bucket.clone(),
174 signer,
175 }),
176 })
177 }
178}
179
180#[derive(Debug, Clone)]
182pub struct UpyunBackend {
183 pub(crate) core: Arc<UpyunCore>,
184}
185
186impl Service for UpyunBackend {
187 type Reader = oio::StreamReader<UpyunReader>;
188 type Writer = UpyunWriters;
189 type Lister = oio::PageLister<UpyunLister>;
190 type Deleter = oio::OneShotDeleter<UpyunDeleter>;
191 type Copier = oio::OneShotCopier;
192 type Composer = ();
193
194 fn info(&self) -> ServiceInfo {
195 self.core.info.clone()
196 }
197
198 fn capability(&self) -> Capability {
199 self.core.capability
200 }
201
202 async fn create_dir(
203 &self,
204 ctx: &OperationContext,
205 path: &str,
206 _: OpCreateDir,
207 ) -> Result<RpCreateDir> {
208 let resp = self.core.create_dir(ctx, path).await?;
209
210 let status = resp.status();
211
212 match status {
213 StatusCode::OK => Ok(RpCreateDir::default()),
214 _ => Err(parse_error(
215 ErrorContext::new(ServiceOperation("CreateFolder")),
216 resp,
217 )),
218 }
219 }
220
221 async fn stat(&self, ctx: &OperationContext, path: &str, _args: OpStat) -> Result<RpStat> {
222 let resp = self.core.info(ctx, path).await?;
223
224 let status = resp.status();
225
226 match status {
227 StatusCode::OK => parse_info(resp.headers()).map(RpStat::new),
228 _ => Err(parse_error(
229 ErrorContext::new(ServiceOperation("GetFileInfo")),
230 resp,
231 )),
232 }
233 }
234 fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
235 let output: oio::StreamReader<UpyunReader> = {
236 Ok(oio::StreamReader::new(UpyunReader::new(
237 self.clone(),
238 ctx.clone(),
239 path,
240 args,
241 )))
242 }?;
243
244 Ok(output)
245 }
246
247 fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
248 let output: UpyunWriters = {
249 let concurrent = args.concurrent();
250 let writer = UpyunWriter::new(self.core.clone(), ctx.clone(), args, path.to_string());
251
252 let w = oio::MultipartWriter::new(ctx.executor().clone(), writer, concurrent);
253
254 Ok(w)
255 }?;
256
257 Ok(output)
258 }
259
260 fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
261 let output: oio::OneShotDeleter<UpyunDeleter> = {
262 Ok(oio::OneShotDeleter::new(UpyunDeleter::new(
263 self.core.clone(),
264 ctx.clone(),
265 )))
266 }?;
267
268 Ok(output)
269 }
270
271 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
272 let output: oio::PageLister<UpyunLister> = {
273 let l = UpyunLister::new(self.core.clone(), ctx.clone(), path, args.limit());
274 Ok(oio::PageLister::new(l))
275 }?;
276
277 Ok(output)
278 }
279
280 fn copy(
281 &self,
282 ctx: &OperationContext,
283 from: &str,
284 to: &str,
285 args: OpCopy,
286 ) -> Result<Self::Copier> {
287 let backend = self.clone();
288 let core = self.core.clone();
289 let ctx = ctx.clone();
290 let from = from.to_string();
291 let to = to.to_string();
292 let source_content_length_hint = args.source_content_length_hint();
293
294 Ok(oio::OneShotCopier::new(async move {
295 let source_size = match source_content_length_hint {
296 Some(size) => size,
297 None => backend
298 .stat(&ctx, &from, OpStat::default())
299 .await?
300 .into_metadata()
301 .content_length(),
302 };
303
304 let resp = core.copy(&ctx, &from, &to).await?;
305 let status = resp.status();
306
307 match status {
308 StatusCode::OK => Ok(MetadataBuilder::file(source_size).build()),
309 _ => Err(parse_error(
310 ErrorContext::new(ServiceOperation("CopyFile")),
311 resp,
312 )),
313 }
314 }))
315 }
316
317 async fn rename(
318 &self,
319 ctx: &OperationContext,
320 from: &str,
321 to: &str,
322 _args: OpRename,
323 ) -> Result<RpRename> {
324 let resp = self.core.move_object(ctx, from, to).await?;
325
326 let status = resp.status();
327
328 match status {
329 StatusCode::OK => Ok(RpRename::default()),
330 _ => Err(parse_error(
331 ErrorContext::new(ServiceOperation("MoveFile")),
332 resp,
333 )),
334 }
335 }
336
337 async fn presign(
338 &self,
339 _ctx: &OperationContext,
340 _path: &str,
341 _args: OpPresign,
342 ) -> Result<RpPresign> {
343 Err(Error::new(
344 ErrorKind::Unsupported,
345 "operation is not supported",
346 ))
347 }
348}