Skip to main content

opendal_service_alluxio/
core.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;
19
20use bytes::Buf;
21use http::Request;
22use http::Response;
23use http::StatusCode;
24use serde::Deserialize;
25use serde::Serialize;
26
27use opendal_core::raw::*;
28use opendal_core::*;
29
30/// Alluxio core
31#[derive(Clone)]
32pub struct AlluxioCore {
33    pub info: ServiceInfo,
34    pub capability: Capability,
35    /// root of this backend.
36    pub root: String,
37    /// endpoint of alluxio
38    pub endpoint: String,
39}
40
41impl Debug for AlluxioCore {
42    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
43        f.debug_struct("AlluxioCore")
44            .field("root", &self.root)
45            .field("endpoint", &self.endpoint)
46            .finish_non_exhaustive()
47    }
48}
49
50impl AlluxioCore {
51    pub async fn create_dir(&self, ctx: &OperationContext, path: &str) -> Result<()> {
52        let path = build_rooted_abs_path(&self.root, path);
53
54        let r = CreateDirRequest {
55            recursive: Some(true),
56            allow_exists: Some(true),
57        };
58
59        let body = serde_json::to_vec(&r).map_err(new_json_serialize_error)?;
60        let body = bytes::Bytes::from(body);
61
62        let mut req = Request::post(format!(
63            "{}/api/v1/paths/{}/create-directory",
64            self.endpoint,
65            percent_encode_path(&path)
66        ));
67
68        req = req.header("Content-Type", "application/json");
69
70        let req = req
71            .extension(Operation::CreateDir)
72            .extension(ServiceOperation("CreateDirectory"));
73
74        let req = req
75            .body(Buffer::from(body))
76            .map_err(new_request_build_error)?;
77
78        let resp = ctx.http_transport().send(req).await?;
79
80        let status = resp.status();
81        match status {
82            StatusCode::OK => Ok(()),
83            _ => Err(parse_error(
84                ErrorContext::new(ServiceOperation("CreateDirectory")),
85                resp,
86            )),
87        }
88    }
89
90    pub async fn create_file(&self, ctx: &OperationContext, path: &str) -> Result<u64> {
91        let path = build_rooted_abs_path(&self.root, path);
92
93        let r = CreateFileRequest {
94            recursive: Some(true),
95        };
96
97        let body = serde_json::to_vec(&r).map_err(new_json_serialize_error)?;
98        let body = bytes::Bytes::from(body);
99        let mut req = Request::post(format!(
100            "{}/api/v1/paths/{}/create-file",
101            self.endpoint,
102            percent_encode_path(&path)
103        ));
104
105        req = req.header("Content-Type", "application/json");
106
107        let req = req
108            .extension(Operation::Write)
109            .extension(ServiceOperation("CreateFile"));
110
111        let req = req
112            .body(Buffer::from(body))
113            .map_err(new_request_build_error)?;
114
115        let resp = ctx.http_transport().send(req).await?;
116        let status = resp.status();
117
118        match status {
119            StatusCode::OK => {
120                let body = resp.into_body();
121                let steam_id: u64 =
122                    serde_json::from_reader(body.reader()).map_err(new_json_serialize_error)?;
123                Ok(steam_id)
124            }
125            _ => Err(parse_error(
126                ErrorContext::new(ServiceOperation("CreateFile")),
127                resp,
128            )),
129        }
130    }
131
132    pub(super) async fn open_file(&self, ctx: &OperationContext, path: &str) -> Result<u64> {
133        let path = build_rooted_abs_path(&self.root, path);
134
135        let req = Request::post(format!(
136            "{}/api/v1/paths/{}/open-file",
137            self.endpoint,
138            percent_encode_path(&path)
139        ));
140
141        let req = req
142            .extension(Operation::Read)
143            .extension(ServiceOperation("OpenFile"));
144
145        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
146        let resp = ctx.http_transport().send(req).await?;
147
148        let status = resp.status();
149
150        match status {
151            StatusCode::OK => {
152                let body = resp.into_body();
153                let steam_id: u64 =
154                    serde_json::from_reader(body.reader()).map_err(new_json_serialize_error)?;
155                Ok(steam_id)
156            }
157            _ => Err(parse_error(
158                ErrorContext::new(ServiceOperation("OpenFile")),
159                resp,
160            )),
161        }
162    }
163
164    pub(super) async fn delete(&self, ctx: &OperationContext, path: &str) -> Result<()> {
165        let path = build_rooted_abs_path(&self.root, path);
166
167        let req = Request::post(format!(
168            "{}/api/v1/paths/{}/delete",
169            self.endpoint,
170            percent_encode_path(&path)
171        ));
172
173        let req = req
174            .extension(Operation::Delete)
175            .extension(ServiceOperation("Delete"));
176
177        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
178        let resp = ctx.http_transport().send(req).await?;
179
180        let status = resp.status();
181
182        match status {
183            StatusCode::OK => Ok(()),
184            _ => {
185                let err = parse_error(ErrorContext::new(ServiceOperation("Delete")), resp);
186                if err.kind() == ErrorKind::NotFound {
187                    return Ok(());
188                }
189                Err(err)
190            }
191        }
192    }
193
194    pub(super) async fn rename(&self, ctx: &OperationContext, path: &str, dst: &str) -> Result<()> {
195        let path = build_rooted_abs_path(&self.root, path);
196        let dst = build_rooted_abs_path(&self.root, dst);
197
198        let req = Request::post(format!(
199            "{}/api/v1/paths/{}/rename?dst={}",
200            self.endpoint,
201            percent_encode_path(&path),
202            percent_encode_path(&dst)
203        ));
204
205        let req = req
206            .extension(Operation::Rename)
207            .extension(ServiceOperation("Rename"));
208
209        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
210
211        let resp = ctx.http_transport().send(req).await?;
212
213        let status = resp.status();
214
215        match status {
216            StatusCode::OK => Ok(()),
217            _ => Err(parse_error(
218                ErrorContext::new(ServiceOperation("Rename")),
219                resp,
220            )),
221        }
222    }
223
224    pub(super) async fn get_status(&self, ctx: &OperationContext, path: &str) -> Result<FileInfo> {
225        let path = build_rooted_abs_path(&self.root, path);
226
227        let req = Request::post(format!(
228            "{}/api/v1/paths/{}/get-status",
229            self.endpoint,
230            percent_encode_path(&path)
231        ));
232
233        let req = req
234            .extension(Operation::Stat)
235            .extension(ServiceOperation("GetStatus"));
236
237        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
238
239        let resp = ctx.http_transport().send(req).await?;
240
241        let status = resp.status();
242
243        match status {
244            StatusCode::OK => {
245                let body = resp.into_body();
246                let file_info: FileInfo =
247                    serde_json::from_reader(body.reader()).map_err(new_json_serialize_error)?;
248                Ok(file_info)
249            }
250            _ => Err(parse_error(
251                ErrorContext::new(ServiceOperation("GetStatus")),
252                resp,
253            )),
254        }
255    }
256
257    pub(super) async fn list_status(
258        &self,
259        ctx: &OperationContext,
260        path: &str,
261    ) -> Result<Vec<FileInfo>> {
262        let path = build_rooted_abs_path(&self.root, path);
263
264        let req = Request::post(format!(
265            "{}/api/v1/paths/{}/list-status",
266            self.endpoint,
267            percent_encode_path(&path)
268        ));
269
270        let req = req
271            .extension(Operation::List)
272            .extension(ServiceOperation("ListStatus"));
273
274        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
275
276        let resp = ctx.http_transport().send(req).await?;
277
278        let status = resp.status();
279
280        match status {
281            StatusCode::OK => {
282                let body = resp.into_body();
283                let file_infos: Vec<FileInfo> =
284                    serde_json::from_reader(body.reader()).map_err(new_json_deserialize_error)?;
285                Ok(file_infos)
286            }
287            _ => Err(parse_error(
288                ErrorContext::new(ServiceOperation("ListStatus")),
289                resp,
290            )),
291        }
292    }
293
294    pub async fn read(
295        &self,
296        ctx: &OperationContext,
297        stream_id: u64,
298        range: BytesRange,
299    ) -> Result<Response<HttpBody>> {
300        if !range.is_full() {
301            return Err(Error::new(
302                ErrorKind::Unsupported,
303                "alluxio stream read doesn't support range",
304            )
305            .with_context("range", format!("{range:?}")));
306        }
307
308        let req = Request::post(format!(
309            "{}/api/v1/streams/{}/read",
310            self.endpoint, stream_id,
311        ));
312
313        let req = req
314            .extension(Operation::Read)
315            .extension(ServiceOperation("ReadStream"));
316
317        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
318
319        ctx.http_transport().fetch(req).await
320    }
321
322    pub(super) async fn write(
323        &self,
324        ctx: &OperationContext,
325        stream_id: u64,
326        body: Buffer,
327    ) -> Result<usize> {
328        let req = Request::post(format!(
329            "{}/api/v1/streams/{}/write",
330            self.endpoint, stream_id
331        ));
332
333        let req = req
334            .extension(Operation::Write)
335            .extension(ServiceOperation("WriteStream"));
336
337        let req = req.body(body).map_err(new_request_build_error)?;
338
339        let resp = ctx.http_transport().send(req).await?;
340
341        let status = resp.status();
342
343        match status {
344            StatusCode::OK => {
345                let body = resp.into_body();
346                let size: usize =
347                    serde_json::from_reader(body.reader()).map_err(new_json_serialize_error)?;
348                Ok(size)
349            }
350            _ => Err(parse_error(
351                ErrorContext::new(ServiceOperation("WriteStream")),
352                resp,
353            )),
354        }
355    }
356
357    pub(super) async fn close(&self, ctx: &OperationContext, stream_id: u64) -> Result<()> {
358        let req = Request::post(format!(
359            "{}/api/v1/streams/{}/close",
360            self.endpoint, stream_id
361        ));
362
363        let req = req
364            .extension(Operation::Write)
365            .extension(ServiceOperation("CloseStream"));
366
367        let req = req.body(Buffer::new()).map_err(new_request_build_error)?;
368
369        let resp = ctx.http_transport().send(req).await?;
370
371        let status = resp.status();
372
373        match status {
374            StatusCode::OK => Ok(()),
375            _ => Err(parse_error(
376                ErrorContext::new(ServiceOperation("CloseStream")),
377                resp,
378            )),
379        }
380    }
381}
382
383#[cfg(test)]
384mod tests {
385    use http::StatusCode;
386
387    use super::*;
388
389    #[tokio::test]
390    async fn test_read_rejects_range() {
391        let core = AlluxioCore {
392            info: ServiceInfo::new("alluxio", "", ""),
393            capability: Capability::default(),
394            root: "/".to_string(),
395            endpoint: "http://127.0.0.1:1".to_string(),
396        };
397
398        let ctx = OperationContext::new();
399        let err = match core.read(&ctx, 1, BytesRange::from(0_u64..1)).await {
400            Ok(_) => panic!("range read should be rejected"),
401            Err(err) => err,
402        };
403
404        assert_eq!(err.kind(), ErrorKind::Unsupported);
405    }
406
407    /// Error response example is from https://docs.aws.amazon.com/AmazonS3/latest/API/ErrorResponses.html
408    #[test]
409    fn test_parse_error() {
410        let err_res = vec![
411            (
412                r#"{"statusCode":"ALREADY_EXISTS","message":"The resource you requested already exist"}"#,
413                ErrorKind::AlreadyExists,
414            ),
415            (
416                r#"{"statusCode":"NOT_FOUND","message":"The resource you requested does not exist"}"#,
417                ErrorKind::NotFound,
418            ),
419            (
420                r#"{"statusCode":"INTERNAL_SERVER_ERROR","message":"Internal server error"}"#,
421                ErrorKind::Unexpected,
422            ),
423        ];
424
425        for res in err_res {
426            let bs = bytes::Bytes::from(res.0);
427            let body = Buffer::from(bs);
428            let resp = Response::builder()
429                .status(StatusCode::INTERNAL_SERVER_ERROR)
430                .body(body)
431                .unwrap();
432
433            let err = parse_error(ErrorContext::new(ServiceOperation("Test")), resp);
434
435            assert_eq!(err.kind(), res.1);
436        }
437    }
438}
439
440#[derive(Debug, Serialize)]
441struct CreateFileRequest {
442    #[serde(skip_serializing_if = "Option::is_none")]
443    recursive: Option<bool>,
444}
445
446#[derive(Debug, Serialize)]
447#[serde(rename_all = "camelCase")]
448struct CreateDirRequest {
449    #[serde(skip_serializing_if = "Option::is_none")]
450    recursive: Option<bool>,
451    #[serde(skip_serializing_if = "Option::is_none")]
452    allow_exists: Option<bool>,
453}
454
455/// Metadata of alluxio object
456#[derive(Debug, Deserialize)]
457#[serde(rename_all = "camelCase")]
458pub(super) struct FileInfo {
459    /// The path of the object
460    pub path: String,
461    /// The last modification time of the object
462    pub last_modification_time_ms: i64,
463    /// Whether the object is a folder
464    pub folder: bool,
465    /// The length of the object in bytes
466    pub length: u64,
467}
468
469impl TryFrom<FileInfo> for Metadata {
470    type Error = Error;
471
472    fn try_from(file_info: FileInfo) -> Result<Metadata> {
473        let mut metadata = if file_info.folder {
474            MetadataBuilder::dir()
475        } else {
476            MetadataBuilder::file(file_info.length)
477        };
478        metadata.last_modified(Timestamp::from_millisecond(
479            file_info.last_modification_time_ms,
480        )?);
481        Ok(metadata.build())
482    }
483}
484
485/// the error response of alluxio
486#[derive(Default, Debug, Deserialize)]
487#[serde(rename_all = "camelCase")]
488#[allow(dead_code)]
489struct AlluxioError {
490    status_code: String,
491    message: String,
492}
493
494/// Context needed to classify an error from this service.
495#[derive(Clone, Copy, Debug)]
496pub(crate) struct ErrorContext {
497    service_operation: ServiceOperation,
498}
499
500impl ErrorContext {
501    pub(crate) const fn new(service_operation: ServiceOperation) -> Self {
502        Self { service_operation }
503    }
504}
505
506/// Parse an error response using its service request context.
507pub(crate) fn parse_error(ctx: ErrorContext, resp: Response<Buffer>) -> Error {
508    let (parts, body) = resp.into_parts();
509    let bs = body.to_bytes();
510
511    let mut kind = match parts.status.as_u16() {
512        500 => ErrorKind::Unexpected,
513        _ => ErrorKind::Unexpected,
514    };
515
516    let (message, alluxio_err) = serde_json::from_reader::<_, AlluxioError>(bs.clone().reader())
517        .map(|alluxio_err| (format!("{alluxio_err:?}"), Some(alluxio_err)))
518        .unwrap_or_else(|_| (String::from_utf8_lossy(&bs).into_owned(), None));
519
520    if let Some(alluxio_err) = alluxio_err {
521        kind = match alluxio_err.status_code.as_str() {
522            "ALREADY_EXISTS" => ErrorKind::AlreadyExists,
523            "NOT_FOUND" => ErrorKind::NotFound,
524            _ => ErrorKind::Unexpected,
525        }
526    }
527
528    let mut err = Error::new(kind, message);
529
530    err = err.with_context("service_operation", ctx.service_operation.0);
531    err = with_error_response_context(err, parts);
532
533    err
534}