1use std::fmt::{self, Debug};
16use std::sync::Arc;
17
18use hdfs_native::{Client, ClientBuilder};
19use opendal::raw::oio::{Delete as _, Read as _, ReadStream as _, Write as _};
20use opendal::raw::{
21 Layer, OpCompose, OpCopy, OpCreateDir, OpDelete, OpList, OpPresign, OpRead, OpRename, OpStat,
22 OpWrite, RpCreateDir, RpPresign, RpRename, RpStat, Service, ServiceInfo, Servicer, oio,
23};
24use opendal::{Buffer, Capability, ErrorKind, Metadata, OperationContext, Result};
25use uuid::Uuid;
26
27#[derive(Debug, Clone)]
29pub struct HdfsCompatibilityLayer {
30 renamer: AtomicRenamer,
31}
32
33impl HdfsCompatibilityLayer {
34 pub fn new(
36 name_node: &str,
37 root: &str,
38 options: &std::collections::HashMap<String, String>,
39 ) -> Result<Self> {
40 let mut config = std::collections::HashMap::new();
41 let namenodes = name_node
42 .split(',')
43 .filter_map(|value| {
44 let value = value
45 .trim()
46 .trim_start_matches("hdfs://")
47 .trim_end_matches('/');
48 (!value.is_empty()).then_some(value)
49 })
50 .collect::<Vec<_>>();
51 for (index, namenode) in namenodes.iter().enumerate() {
52 config.insert(
53 format!("dfs.namenode.rpc-address.nameservice.nn{index}"),
54 (*namenode).to_string(),
55 );
56 }
57 config.insert(
58 "dfs.ha.namenodes.nameservice".to_string(),
59 (0..namenodes.len())
60 .map(|index| format!("nn{index}"))
61 .collect::<Vec<_>>()
62 .join(","),
63 );
64 config.extend(options.clone());
65 let client = ClientBuilder::new()
66 .with_url("hdfs://nameservice")
67 .with_config(config)
68 .build()
69 .map_err(hdfs_error)?;
70 Ok(Self {
71 renamer: AtomicRenamer::Native {
72 client,
73 root: opendal::raw::normalize_root(root),
74 },
75 })
76 }
77
78 #[cfg(any(test, feature = "testing"))]
80 pub fn new_for_test() -> Self {
81 Self {
82 renamer: AtomicRenamer::Raw,
83 }
84 }
85}
86
87#[derive(Debug, Clone)]
88enum AtomicRenamer {
89 Native {
90 client: Client,
91 root: String,
92 },
93 #[cfg(any(test, feature = "testing"))]
94 Raw,
95}
96
97impl AtomicRenamer {
98 async fn rename(
99 &self,
100 _inner: &Servicer,
101 _ctx: &OperationContext,
102 from: &str,
103 to: &str,
104 ) -> Result<()> {
105 match self {
106 Self::Native { client, root } => {
107 client
110 .rename(
111 &opendal::raw::build_rooted_abs_path(root, from),
112 &opendal::raw::build_rooted_abs_path(root, to),
113 true,
114 )
115 .await
116 .map_err(hdfs_error)
117 }
118 #[cfg(any(test, feature = "testing"))]
119 Self::Raw => _inner
120 .rename(_ctx, from, to, OpRename::new())
121 .await
122 .map(|_| ()),
123 }
124 }
125}
126
127fn hdfs_error(error: hdfs_native::HdfsError) -> opendal::Error {
128 opendal::Error::new(ErrorKind::Unexpected, "native HDFS operation failed").set_source(error)
129}
130
131impl Layer for HdfsCompatibilityLayer {
132 fn apply_service(&self, inner: Servicer) -> Servicer {
133 Arc::new(HdfsCompatibilityService {
134 inner,
135 renamer: self.renamer.clone(),
136 })
137 }
138}
139
140struct HdfsCompatibilityService {
141 inner: Servicer,
142 renamer: AtomicRenamer,
143}
144
145impl Debug for HdfsCompatibilityService {
146 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
147 f.debug_struct("HdfsCompatibilityService")
148 .field("inner", &self.inner)
149 .finish()
150 }
151}
152
153struct HdfsWriter(HdfsWriterInner);
155
156enum HdfsWriterInner {
157 Direct(oio::Writer),
158 Atomic {
159 inner: Servicer,
160 context: OperationContext,
161 renamer: AtomicRenamer,
162 writer: Option<oio::Writer>,
163 temporary_path: String,
164 target_path: String,
165 },
166}
167
168impl oio::Write for HdfsWriter {
169 async fn write(&mut self, buffer: Buffer) -> Result<()> {
170 match &mut self.0 {
171 HdfsWriterInner::Direct(writer) => writer.write(buffer).await,
172 HdfsWriterInner::Atomic { writer, .. } => {
173 writer
174 .as_mut()
175 .ok_or_else(writer_unavailable)?
176 .write(buffer)
177 .await
178 }
179 }
180 }
181
182 async fn close(&mut self) -> Result<Metadata> {
183 match &mut self.0 {
184 HdfsWriterInner::Direct(writer) => writer.close().await,
185 HdfsWriterInner::Atomic {
186 inner,
187 context,
188 renamer,
189 writer,
190 temporary_path,
191 target_path,
192 } => {
193 let mut writer = writer.take().ok_or_else(writer_unavailable)?;
194 let metadata = match writer.close().await {
195 Ok(metadata) => metadata,
196 Err(error) => {
197 drop(writer);
198 let _ = delete_path(inner, context, temporary_path).await;
199 return Err(error);
200 }
201 };
202 drop(writer);
203
204 if let Err(error) = renamer
205 .rename(inner, context, temporary_path, target_path)
206 .await
207 {
208 let _ = delete_path(inner, context, temporary_path).await;
209 return Err(error);
210 }
211
212 Ok(metadata)
213 }
214 }
215 }
216
217 async fn abort(&mut self) -> Result<()> {
218 match &mut self.0 {
219 HdfsWriterInner::Direct(writer) => writer.abort().await,
220 HdfsWriterInner::Atomic {
221 inner,
222 context,
223 writer,
224 temporary_path,
225 ..
226 } => {
227 let abort_result = if let Some(mut writer) = writer.take() {
228 let result = writer.abort().await;
229 drop(writer);
230 result
231 } else {
232 Ok(())
233 };
234 let cleanup_result = delete_path(inner, context, temporary_path).await;
235
236 match abort_result {
237 Err(error) if error.kind() != ErrorKind::Unsupported => Err(error),
238 _ => cleanup_result,
239 }
240 }
241 }
242 }
243}
244
245fn writer_unavailable() -> opendal::Error {
246 opendal::Error::new(
247 ErrorKind::Unexpected,
248 "HDFS writer is unavailable after close or abort",
249 )
250}
251
252impl Service for HdfsCompatibilityService {
253 type Reader = oio::Reader;
254 type Writer = HdfsWriter;
255 type Lister = oio::Lister;
256 type Deleter = oio::Deleter;
257 type Copier = oio::OneShotCopier;
258 type Composer = oio::Composer;
259
260 fn info(&self) -> ServiceInfo {
261 self.inner.info()
262 }
263
264 fn capability(&self) -> Capability {
265 let mut capability = self.inner.capability();
266 capability.copy = true;
267 capability
268 }
269
270 async fn create_dir(
271 &self,
272 ctx: &OperationContext,
273 path: &str,
274 args: OpCreateDir,
275 ) -> Result<RpCreateDir> {
276 self.inner.create_dir(ctx, path, args).await
277 }
278
279 async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
280 self.inner.stat(ctx, path, args).await
281 }
282
283 fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
284 self.inner.read(ctx, path, args)
285 }
286
287 fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
288 if args.append() {
289 return self
290 .inner
291 .write(ctx, path, args)
292 .map(|writer| HdfsWriter(HdfsWriterInner::Direct(writer)));
293 }
294
295 let temporary_path = temporary_path(path);
296 let writer = self.inner.write(ctx, &temporary_path, args)?;
297 Ok(HdfsWriter(HdfsWriterInner::Atomic {
298 inner: Arc::clone(&self.inner),
299 context: ctx.clone(),
300 renamer: self.renamer.clone(),
301 writer: Some(writer),
302 temporary_path,
303 target_path: path.to_string(),
304 }))
305 }
306
307 fn copy(
308 &self,
309 ctx: &OperationContext,
310 from: &str,
311 to: &str,
312 args: OpCopy,
313 ) -> Result<Self::Copier> {
314 if args.if_not_exists() || args.if_match().is_some() {
315 return Err(opendal::Error::new(
316 ErrorKind::Unsupported,
317 "conditional copy is not supported by the HDFS fallback",
318 ));
319 }
320
321 let inner = Arc::clone(&self.inner);
322 let context = ctx.clone();
323 let renamer = self.renamer.clone();
324 let from = from.to_string();
325 let to = to.to_string();
326 Ok(oio::OneShotCopier::new(async move {
327 copy_via_read_write(inner, &context, renamer, &from, &to).await
328 }))
329 }
330
331 fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
332 self.inner.delete(ctx)
333 }
334
335 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
336 self.inner.list(ctx, path, args)
337 }
338
339 fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
340 self.inner.compose(ctx, to, args)
341 }
342
343 async fn rename(
344 &self,
345 ctx: &OperationContext,
346 from: &str,
347 to: &str,
348 args: OpRename,
349 ) -> Result<RpRename> {
350 self.inner.rename(ctx, from, to, args).await
351 }
352
353 async fn presign(
354 &self,
355 ctx: &OperationContext,
356 path: &str,
357 args: OpPresign,
358 ) -> Result<RpPresign> {
359 self.inner.presign(ctx, path, args).await
360 }
361}
362
363async fn copy_via_read_write(
364 inner: Servicer,
365 context: &OperationContext,
366 renamer: AtomicRenamer,
367 source_path: &str,
368 target_path: &str,
369) -> Result<Metadata> {
370 let reader = inner.read(context, source_path, OpRead::new())?;
371 let (_, mut reader) = reader.open(opendal::BytesRange::from(..)).await?;
372 let temporary_path = temporary_path(target_path);
373 let mut writer = inner.write(context, &temporary_path, OpWrite::new())?;
374
375 loop {
376 let buffer = match reader.read().await {
377 Ok(buffer) => buffer,
378 Err(error) => {
379 abort_and_delete(&inner, context, writer, &temporary_path).await;
380 return Err(error);
381 }
382 };
383 if buffer.is_empty() {
384 break;
385 }
386 if let Err(error) = writer.write(buffer).await {
387 abort_and_delete(&inner, context, writer, &temporary_path).await;
388 return Err(error);
389 }
390 }
391
392 let metadata = match writer.close().await {
393 Ok(metadata) => metadata,
394 Err(error) => {
395 drop(writer);
396 let _ = delete_path(&inner, context, &temporary_path).await;
397 return Err(error);
398 }
399 };
400 drop(writer);
401
402 if let Err(error) = renamer
403 .rename(&inner, context, &temporary_path, target_path)
404 .await
405 {
406 let _ = delete_path(&inner, context, &temporary_path).await;
407 return Err(error);
408 }
409
410 Ok(metadata)
411}
412
413async fn abort_and_delete(
414 inner: &Servicer,
415 context: &OperationContext,
416 mut writer: oio::Writer,
417 path: &str,
418) {
419 let _ = writer.abort().await;
420 drop(writer);
421 let _ = delete_path(inner, context, path).await;
422}
423
424async fn delete_path(inner: &Servicer, context: &OperationContext, path: &str) -> Result<()> {
425 let mut deleter = inner.delete(context)?;
426 deleter.delete(path, OpDelete::new()).await?;
427 deleter.close().await
428}
429
430fn temporary_path(path: &str) -> String {
432 let suffix = format!(".greptime-{}.tmp", Uuid::new_v4());
433 match path.rsplit_once('/') {
434 Some((parent, name)) => format!("{parent}/.{name}{suffix}"),
435 None => format!(".{path}{suffix}"),
436 }
437}
438
439#[cfg(test)]
440mod tests {
441 use opendal::services::Fs;
442 use opendal::{Operator, Writer};
443 use tempfile::TempDir;
444
445 use super::*;
446
447 fn test_store() -> (TempDir, Operator) {
448 let directory = tempfile::tempdir().unwrap();
449 let store = Operator::new(Fs::default().root(directory.path().to_str().unwrap()))
450 .unwrap()
451 .layer(HdfsCompatibilityLayer::new_for_test());
452 (directory, store)
453 }
454
455 #[tokio::test]
456 async fn test_service_operations() {
457 let (_directory, store) = test_store();
458 store.create_dir("data/").await.unwrap();
459 store.write("data/source", "contents").await.unwrap();
460 assert_eq!(8, store.stat("data/source").await.unwrap().content_length());
461 assert_eq!(
462 b"onte",
463 store
464 .read_with("data/source")
465 .range(1..5)
466 .await
467 .unwrap()
468 .to_bytes()
469 .as_ref()
470 );
471 store.rename("data/source", "data/target").await.unwrap();
472 let entries = store.list("data/").await.unwrap();
473 assert!(entries.iter().any(|entry| entry.path() == "data/target"));
474 assert!(!store.exists("data/source").await.unwrap());
475 store.delete("data/target").await.unwrap();
476 assert!(!store.exists("data/target").await.unwrap());
477 }
478
479 #[tokio::test]
480 async fn test_atomic_write_keeps_old_data_after_abort() {
481 let (_directory, store) = test_store();
482 store.write("manifest.json", "old").await.unwrap();
483
484 let mut writer: Writer = store.writer("manifest.json").await.unwrap();
485 writer.write("new").await.unwrap();
486 writer.abort().await.unwrap();
487
488 assert_eq!(
489 b"old",
490 store
491 .read("manifest.json")
492 .await
493 .unwrap()
494 .to_bytes()
495 .as_ref()
496 );
497 assert!(
498 store
499 .list("")
500 .await
501 .unwrap()
502 .iter()
503 .all(|entry| !entry.path().contains(".greptime-"))
504 );
505 }
506
507 #[tokio::test]
508 async fn test_atomic_write_replaces_on_close() {
509 let (_directory, store) = test_store();
510 store.write("manifest.json", "old").await.unwrap();
511 store.write("manifest.json", "new").await.unwrap();
512
513 assert_eq!(
514 b"new",
515 store
516 .read("manifest.json")
517 .await
518 .unwrap()
519 .to_bytes()
520 .as_ref()
521 );
522 assert!(
523 store
524 .list("")
525 .await
526 .unwrap()
527 .iter()
528 .all(|entry| !entry.path().contains(".greptime-"))
529 );
530 }
531
532 #[tokio::test]
533 async fn test_copy_fallback_streams_to_target() {
534 let (_directory, store) = test_store();
535 store.write("source.parquet", "contents").await.unwrap();
536 store
537 .copy("source.parquet", "target.parquet")
538 .await
539 .unwrap();
540
541 assert_eq!(
542 b"contents",
543 store
544 .read("target.parquet")
545 .await
546 .unwrap()
547 .to_bytes()
548 .as_ref()
549 );
550 }
551}