1use 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#[derive(Clone)]
32pub struct AlluxioCore {
33 pub info: ServiceInfo,
34 pub capability: Capability,
35 pub root: String,
37 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 #[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#[derive(Debug, Deserialize)]
457#[serde(rename_all = "camelCase")]
458pub(super) struct FileInfo {
459 pub path: String,
461 pub last_modification_time_ms: i64,
463 pub folder: bool,
465 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#[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#[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
506pub(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}