1mod utils;
26
27use std::collections::HashMap;
28use std::collections::hash_map::Entry;
29use std::num::NonZeroU64;
30use std::sync::Arc;
31use std::time::Duration;
32
33use async_trait::async_trait;
34use bytes::Bytes;
35use cfg_if::cfg_if;
36use futures::StreamExt;
37use futures::stream::BoxStream;
38use iceberg::io::{
39 FileMetadata, FileRead, FileWrite, InputFile, OutputFile, Storage, StorageConfig,
40 StorageFactory,
41};
42use iceberg::{Error, ErrorKind, Result};
43use iceberg_property_macro::Properties;
44use opendal::Operator;
45use opendal::layers::{RetryLayer, TimeoutLayer};
46use serde::{Deserialize, Serialize};
47use utils::from_opendal_error;
48
49cfg_if! {
50 if #[cfg(feature = "opendal-azdls")] {
51 mod azdls;
52 use azdls::*;
53 use opendal::services::AzdlsConfig;
54 }
55}
56
57cfg_if! {
58 if #[cfg(feature = "opendal-hf")] {
59 mod hf;
60 use hf::*;
61 use opendal::services::HfConfig;
62 }
63}
64
65cfg_if! {
66 if #[cfg(feature = "opendal-fs")] {
67 mod fs;
68 use fs::*;
69 }
70}
71
72cfg_if! {
73 if #[cfg(feature = "opendal-gcs")] {
74 mod gcs;
75 use gcs::*;
76 use opendal::services::GcsConfig;
77 }
78}
79
80cfg_if! {
81 if #[cfg(feature = "opendal-memory")] {
82 mod memory;
83 use memory::*;
84 }
85}
86
87cfg_if! {
88 if #[cfg(feature = "opendal-oss")] {
89 mod oss;
90 use opendal::services::OssConfig;
91 use oss::*;
92 }
93}
94
95cfg_if! {
96 if #[cfg(feature = "opendal-s3")] {
97 mod s3;
98 use opendal::services::S3Config;
99 pub use s3::*;
100 }
101}
102
103mod resolving;
104pub use resolving::{OpenDalResolvingStorage, OpenDalResolvingStorageFactory};
105
106pub const OPENDAL_IO_TIMEOUT_MS: &str = "opendal.io-timeout-ms";
112
113const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 = NonZeroU64::new(10_000).unwrap();
115
116#[derive(Clone, Debug, Properties, Serialize, Deserialize)]
118#[serde(default)]
119pub struct OpenDalClientConfig {
120 #[property(
122 key = OPENDAL_IO_TIMEOUT_MS,
123 default = DEFAULT_IO_TIMEOUT_MS,
124 parse_with = parse_io_timeout_ms
125 )]
126 io_timeout_ms: NonZeroU64,
127}
128
129impl Default for OpenDalClientConfig {
130 fn default() -> Self {
131 Self {
132 io_timeout_ms: DEFAULT_IO_TIMEOUT_MS,
133 }
134 }
135}
136
137impl OpenDalClientConfig {
138 pub(crate) fn io_timeout(&self) -> Duration {
139 Duration::from_millis(self.io_timeout_ms.get())
140 }
141}
142
143fn parse_io_timeout_ms(value: &str) -> Result<NonZeroU64> {
144 value.parse().map_err(|error| {
145 Error::new(
146 ErrorKind::DataInvalid,
147 "Expected a positive integer number of milliseconds",
148 )
149 .with_context("value", format!("{value:?}"))
150 .with_source(error)
151 })
152}
153
154#[derive(Clone, Debug, Serialize, Deserialize)]
168pub enum OpenDalStorageFactory {
169 #[cfg(feature = "opendal-memory")]
171 Memory,
172 #[cfg(feature = "opendal-fs")]
174 Fs,
175 #[cfg(feature = "opendal-s3")]
177 S3 {
178 #[serde(
180 skip_deserializing,
181 skip_serializing_if = "Option::is_none",
182 serialize_with = "serialize_custom_credential_loader"
183 )]
184 customized_credential_load: Option<CustomAwsCredentialLoader>,
185 },
186 #[cfg(feature = "opendal-gcs")]
188 Gcs,
189 #[cfg(feature = "opendal-oss")]
191 Oss,
192 #[cfg(feature = "opendal-azdls")]
194 Azdls,
195 #[cfg(feature = "opendal-hf")]
197 Hf,
198}
199
200#[cfg(feature = "opendal-s3")]
201pub(crate) fn serialize_custom_credential_loader<S>(
202 _loader: &Option<CustomAwsCredentialLoader>,
203 _serializer: S,
204) -> std::result::Result<S::Ok, S::Error>
205where
206 S: serde::Serializer,
207{
208 Err(serde::ser::Error::custom(
209 "custom AWS credential loaders cannot be serialized",
210 ))
211}
212
213#[typetag::serde(name = "OpenDalStorageFactory")]
214impl StorageFactory for OpenDalStorageFactory {
215 #[allow(unused_variables)]
216 fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>> {
217 let client_config = OpenDalClientConfig::from_properties(config.props())?;
218 match self {
219 #[cfg(feature = "opendal-memory")]
220 OpenDalStorageFactory::Memory => Ok(Arc::new(OpenDalStorage::Memory {
221 operator: memory_config_build()?,
222 client_config,
223 })),
224 #[cfg(feature = "opendal-fs")]
225 OpenDalStorageFactory::Fs => Ok(Arc::new(OpenDalStorage::LocalFs { client_config })),
226 #[cfg(feature = "opendal-s3")]
227 OpenDalStorageFactory::S3 {
228 customized_credential_load,
229 } => Ok(Arc::new(OpenDalStorage::S3 {
230 config: s3_config_parse(config.props().clone())?.into(),
231 customized_credential_load: customized_credential_load.clone(),
232 client_config,
233 })),
234 #[cfg(feature = "opendal-gcs")]
235 OpenDalStorageFactory::Gcs => Ok(Arc::new(OpenDalStorage::Gcs {
236 config: gcs_config_parse(config.props().clone())?.into(),
237 client_config,
238 })),
239 #[cfg(feature = "opendal-oss")]
240 OpenDalStorageFactory::Oss => Ok(Arc::new(OpenDalStorage::Oss {
241 config: oss_config_parse(config.props().clone())?.into(),
242 client_config,
243 })),
244 #[cfg(feature = "opendal-azdls")]
245 OpenDalStorageFactory::Azdls => Ok(Arc::new(OpenDalStorage::Azdls {
246 config: azdls_config_parse(config.props().clone())?.into(),
247 client_config,
248 })),
249 #[cfg(feature = "opendal-hf")]
250 OpenDalStorageFactory::Hf => Ok(Arc::new(OpenDalStorage::Hf {
251 config: hf_config_parse(config.props().clone())?.into(),
252 client_config,
253 })),
254 #[cfg(all(
255 not(feature = "opendal-memory"),
256 not(feature = "opendal-fs"),
257 not(feature = "opendal-s3"),
258 not(feature = "opendal-gcs"),
259 not(feature = "opendal-oss"),
260 not(feature = "opendal-azdls"),
261 not(feature = "opendal-hf"),
262 ))]
263 _ => Err(Error::new(
264 ErrorKind::FeatureUnsupported,
265 "No storage service has been enabled",
266 )),
267 }
268 }
269}
270
271#[cfg(feature = "opendal-memory")]
273fn default_memory_operator() -> Operator {
274 memory_config_build().expect("Failed to create default memory operator")
275}
276
277#[derive(Clone, Debug, Serialize, Deserialize)]
281pub enum OpenDalStorage {
282 #[cfg(feature = "opendal-memory")]
284 Memory {
285 #[serde(skip, default = "self::default_memory_operator")]
287 operator: Operator,
288 #[serde(default)]
290 client_config: OpenDalClientConfig,
291 },
292 #[cfg(feature = "opendal-fs")]
294 LocalFs {
295 #[serde(default)]
297 client_config: OpenDalClientConfig,
298 },
299 #[cfg(feature = "opendal-s3")]
304 S3 {
305 config: Arc<S3Config>,
307 #[serde(skip)]
309 customized_credential_load: Option<CustomAwsCredentialLoader>,
310 #[serde(default)]
312 client_config: OpenDalClientConfig,
313 },
314 #[cfg(feature = "opendal-gcs")]
316 Gcs {
317 config: Arc<GcsConfig>,
319 #[serde(default)]
321 client_config: OpenDalClientConfig,
322 },
323 #[cfg(feature = "opendal-oss")]
325 Oss {
326 config: Arc<OssConfig>,
328 #[serde(default)]
330 client_config: OpenDalClientConfig,
331 },
332 #[cfg(feature = "opendal-azdls")]
339 Azdls {
340 config: Arc<AzdlsConfig>,
342 #[serde(default)]
344 client_config: OpenDalClientConfig,
345 },
346 #[cfg(feature = "opendal-hf")]
352 Hf {
353 config: Arc<HfConfig>,
355 #[serde(default)]
357 client_config: OpenDalClientConfig,
358 },
359}
360
361impl OpenDalStorage {
362 #[allow(unreachable_code, unused_variables)]
375 pub(crate) fn create_operator<'a>(
376 &self,
377 path: &'a impl AsRef<str>,
378 ) -> Result<(Operator, &'a str)> {
379 let path = path.as_ref();
380 let (operator, relative_path): (Operator, &str) = match self {
381 #[cfg(feature = "opendal-memory")]
382 OpenDalStorage::Memory { operator: op, .. } => {
383 if let Some(stripped) = path.strip_prefix("memory:/") {
384 (op.clone(), stripped)
385 } else {
386 (op.clone(), &path[1..])
387 }
388 }
389 #[cfg(feature = "opendal-fs")]
390 OpenDalStorage::LocalFs { .. } => {
391 let op = fs_config_build()?;
392 if let Some(stripped) = path.strip_prefix("file:/") {
393 (op, stripped)
394 } else {
395 (op, &path[1..])
396 }
397 }
398 #[cfg(feature = "opendal-s3")]
399 OpenDalStorage::S3 {
400 config,
401 customized_credential_load,
402 ..
403 } => {
404 let op = s3_config_build(config, customized_credential_load, path)?;
405 let op_info = op.info();
406
407 let url = url::Url::parse(path).map_err(|e| {
410 Error::new(
411 ErrorKind::DataInvalid,
412 format!("Invalid s3 url: {path}: {e}"),
413 )
414 })?;
415 let prefix = format!("{}://{}/", url.scheme(), op_info.name());
416 if path.starts_with(&prefix) {
417 (op, &path[prefix.len()..])
418 } else {
419 return Err(Error::new(
420 ErrorKind::DataInvalid,
421 format!("Invalid s3 url: {path}, should start with {prefix}"),
422 ));
423 }
424 }
425 #[cfg(feature = "opendal-gcs")]
426 OpenDalStorage::Gcs { config, .. } => {
427 let operator = gcs_config_build(config, path)?;
428 let prefix = format!("gs://{}/", operator.info().name());
429 if path.starts_with(&prefix) {
430 (operator, &path[prefix.len()..])
431 } else {
432 return Err(Error::new(
433 ErrorKind::DataInvalid,
434 format!("Invalid gcs url: {path}, should start with {prefix}"),
435 ));
436 }
437 }
438 #[cfg(feature = "opendal-oss")]
439 OpenDalStorage::Oss { config, .. } => {
440 let op = oss_config_build(config, path)?;
441 let prefix = format!("oss://{}/", op.info().name());
442 if path.starts_with(&prefix) {
443 (op, &path[prefix.len()..])
444 } else {
445 return Err(Error::new(
446 ErrorKind::DataInvalid,
447 format!("Invalid oss url: {path}, should start with {prefix}"),
448 ));
449 }
450 }
451 #[cfg(feature = "opendal-azdls")]
452 OpenDalStorage::Azdls { config, .. } => azdls_create_operator(path, config)?,
453 #[cfg(feature = "opendal-hf")]
454 OpenDalStorage::Hf { config, .. } => hf_config_build(config, path)?,
455 #[cfg(all(
456 not(feature = "opendal-s3"),
457 not(feature = "opendal-fs"),
458 not(feature = "opendal-gcs"),
459 not(feature = "opendal-oss"),
460 not(feature = "opendal-azdls"),
461 not(feature = "opendal-hf"),
462 ))]
463 _ => {
464 return Err(Error::new(
465 ErrorKind::FeatureUnsupported,
466 "No storage service has been enabled",
467 ));
468 }
469 };
470
471 let operator = operator
482 .layer(TimeoutLayer::new().with_io_timeout(self.client_config().io_timeout()))
483 .layer(RetryLayer::new());
484 Ok((operator, relative_path))
485 }
486
487 pub(crate) fn client_config(&self) -> &OpenDalClientConfig {
488 match self {
489 #[cfg(feature = "opendal-memory")]
490 OpenDalStorage::Memory { client_config, .. } => client_config,
491 #[cfg(feature = "opendal-fs")]
492 OpenDalStorage::LocalFs { client_config } => client_config,
493 #[cfg(feature = "opendal-s3")]
494 OpenDalStorage::S3 { client_config, .. } => client_config,
495 #[cfg(feature = "opendal-gcs")]
496 OpenDalStorage::Gcs { client_config, .. } => client_config,
497 #[cfg(feature = "opendal-oss")]
498 OpenDalStorage::Oss { client_config, .. } => client_config,
499 #[cfg(feature = "opendal-azdls")]
500 OpenDalStorage::Azdls { client_config, .. } => client_config,
501 #[cfg(feature = "opendal-hf")]
502 OpenDalStorage::Hf { client_config, .. } => client_config,
503 #[cfg(all(
504 not(feature = "opendal-memory"),
505 not(feature = "opendal-s3"),
506 not(feature = "opendal-fs"),
507 not(feature = "opendal-gcs"),
508 not(feature = "opendal-oss"),
509 not(feature = "opendal-azdls"),
510 not(feature = "opendal-hf"),
511 ))]
512 _ => unreachable!(),
513 }
514 }
515
516 fn batch_key_for_path(&self, path: &str) -> String {
521 match self {
522 #[cfg(feature = "opendal-hf")]
523 OpenDalStorage::Hf { .. } => hf_batch_key(path),
524 _ => url::Url::parse(path)
525 .ok()
526 .and_then(|u| u.host_str().map(|s| s.to_string()))
527 .unwrap_or_default(),
528 }
529 }
530
531 #[allow(unreachable_code, unused_variables)]
537 pub(crate) fn relativize_path<'a>(&self, path: &'a str) -> Result<&'a str> {
538 match self {
539 #[cfg(feature = "opendal-memory")]
540 OpenDalStorage::Memory { .. } => {
541 Ok(path.strip_prefix("memory:/").unwrap_or(&path[1..]))
542 }
543 #[cfg(feature = "opendal-fs")]
544 OpenDalStorage::LocalFs { .. } => Ok(path.strip_prefix("file:/").unwrap_or(&path[1..])),
545 #[cfg(feature = "opendal-s3")]
546 OpenDalStorage::S3 { .. } => {
547 let url = url::Url::parse(path)?;
548 let bucket = url.host_str().ok_or_else(|| {
549 Error::new(
550 ErrorKind::DataInvalid,
551 format!("Invalid s3 url: {path}, missing bucket"),
552 )
553 })?;
554 let prefix = format!("{}://{}/", url.scheme(), bucket);
555 if path.starts_with(&prefix) {
556 Ok(&path[prefix.len()..])
557 } else {
558 Err(Error::new(
559 ErrorKind::DataInvalid,
560 format!("Invalid s3 url: {path}, should start with {prefix}"),
561 ))
562 }
563 }
564 #[cfg(feature = "opendal-gcs")]
565 OpenDalStorage::Gcs { .. } => {
566 let url = url::Url::parse(path)?;
567 let bucket = url.host_str().ok_or_else(|| {
568 Error::new(
569 ErrorKind::DataInvalid,
570 format!("Invalid gcs url: {path}, missing bucket"),
571 )
572 })?;
573 let prefix = format!("gs://{}/", bucket);
574 if path.starts_with(&prefix) {
575 Ok(&path[prefix.len()..])
576 } else {
577 Err(Error::new(
578 ErrorKind::DataInvalid,
579 format!("Invalid gcs url: {path}, should start with {prefix}"),
580 ))
581 }
582 }
583 #[cfg(feature = "opendal-oss")]
584 OpenDalStorage::Oss { .. } => {
585 let url = url::Url::parse(path)?;
586 let bucket = url.host_str().ok_or_else(|| {
587 Error::new(
588 ErrorKind::DataInvalid,
589 format!("Invalid oss url: {path}, missing bucket"),
590 )
591 })?;
592 let prefix = format!("oss://{}/", bucket);
593 if path.starts_with(&prefix) {
594 Ok(&path[prefix.len()..])
595 } else {
596 Err(Error::new(
597 ErrorKind::DataInvalid,
598 format!("Invalid oss url: {path}, should start with {prefix}"),
599 ))
600 }
601 }
602 #[cfg(feature = "opendal-azdls")]
603 OpenDalStorage::Azdls { config, .. } => {
604 let azure_path = path.parse::<AzureStoragePath>()?;
605 match_path_with_config(&azure_path, config)?;
606 Ok(azure_path.relative_path(path))
607 }
608 #[cfg(feature = "opendal-hf")]
609 OpenDalStorage::Hf { .. } => {
610 let parsed = HfUri::parse(path).ok_or_else(|| {
611 Error::new(ErrorKind::DataInvalid, format!("Invalid hf url: {path}"))
612 })?;
613 Ok(&path[path.len() - parsed.path.len()..])
614 }
615 #[cfg(all(
616 not(feature = "opendal-s3"),
617 not(feature = "opendal-fs"),
618 not(feature = "opendal-gcs"),
619 not(feature = "opendal-oss"),
620 not(feature = "opendal-azdls"),
621 not(feature = "opendal-hf"),
622 ))]
623 _ => Err(Error::new(
624 ErrorKind::FeatureUnsupported,
625 "No storage service has been enabled",
626 )),
627 }
628 }
629}
630
631#[typetag::serde(name = "OpenDalStorage")]
632#[async_trait]
633impl Storage for OpenDalStorage {
634 async fn exists(&self, path: &str) -> Result<bool> {
635 let (op, relative_path) = self.create_operator(&path)?;
636 Ok(op.exists(relative_path).await.map_err(from_opendal_error)?)
637 }
638
639 async fn metadata(&self, path: &str) -> Result<FileMetadata> {
640 let (op, relative_path) = self.create_operator(&path)?;
641 let meta = op.stat(relative_path).await.map_err(from_opendal_error)?;
642 Ok(FileMetadata {
643 size: meta.content_length(),
644 })
645 }
646
647 async fn read(&self, path: &str) -> Result<Bytes> {
648 let (op, relative_path) = self.create_operator(&path)?;
649 Ok(op
650 .read(relative_path)
651 .await
652 .map_err(from_opendal_error)?
653 .to_bytes())
654 }
655
656 async fn reader(&self, path: &str) -> Result<Box<dyn FileRead>> {
657 let (op, relative_path) = self.create_operator(&path)?;
658 Ok(Box::new(OpenDalReader(
659 op.reader(relative_path).await.map_err(from_opendal_error)?,
660 )))
661 }
662
663 async fn write(&self, path: &str, bs: Bytes) -> Result<()> {
664 let (op, relative_path) = self.create_operator(&path)?;
665 op.write(relative_path, bs)
666 .await
667 .map_err(from_opendal_error)?;
668 Ok(())
669 }
670
671 async fn writer(&self, path: &str) -> Result<Box<dyn FileWrite>> {
672 let (op, relative_path) = self.create_operator(&path)?;
673 Ok(Box::new(OpenDalWriter::new(
674 op.writer(relative_path).await.map_err(from_opendal_error)?,
675 )))
676 }
677
678 async fn delete(&self, path: &str) -> Result<()> {
679 let (op, relative_path) = self.create_operator(&path)?;
680 Ok(op.delete(relative_path).await.map_err(from_opendal_error)?)
681 }
682
683 async fn delete_prefix(&self, path: &str) -> Result<()> {
684 let (op, relative_path) = self.create_operator(&path)?;
685 let path = if relative_path.ends_with('/') {
686 relative_path.to_string()
687 } else {
688 format!("{relative_path}/")
689 };
690 Ok(op
691 .delete_with(&path)
692 .recursive(true)
693 .await
694 .map_err(from_opendal_error)?)
695 }
696
697 async fn delete_stream(&self, mut paths: BoxStream<'static, String>) -> Result<()> {
698 let mut deleters: HashMap<String, opendal::Deleter> = HashMap::new();
699
700 while let Some(path) = paths.next().await {
701 let bucket = self.batch_key_for_path(&path);
702
703 let (relative_path, deleter) = match deleters.entry(bucket) {
704 Entry::Occupied(entry) => {
705 (self.relativize_path(&path)?.to_string(), entry.into_mut())
706 }
707 Entry::Vacant(entry) => {
708 let (op, rel) = self.create_operator(&path)?;
709 let rel = rel.to_string();
710 let deleter = op.deleter().await.map_err(from_opendal_error)?;
711 (rel, entry.insert(deleter))
712 }
713 };
714
715 deleter
716 .delete(relative_path)
717 .await
718 .map_err(from_opendal_error)?;
719 }
720
721 for (_, mut deleter) in deleters {
722 deleter.close().await.map_err(from_opendal_error)?;
723 }
724
725 Ok(())
726 }
727
728 #[allow(unreachable_code, unused_variables)]
729 fn new_input(&self, path: &str) -> Result<InputFile> {
730 Ok(InputFile::new(Arc::new(self.clone()), path.to_string()))
731 }
732
733 #[allow(unreachable_code, unused_variables)]
734 fn new_output(&self, path: &str) -> Result<OutputFile> {
735 Ok(OutputFile::new(Arc::new(self.clone()), path.to_string()))
736 }
737}
738
739pub(crate) struct OpenDalReader(pub(crate) opendal::Reader);
745
746#[async_trait]
747impl FileRead for OpenDalReader {
748 async fn read(&self, range: std::ops::Range<u64>) -> Result<Bytes> {
749 Ok(opendal::Reader::read(&self.0, range)
750 .await
751 .map_err(from_opendal_error)?
752 .to_bytes())
753 }
754}
755
756pub(crate) struct OpenDalWriter {
758 inner: opendal::Writer,
759 bytes_written: u64,
760}
761
762impl OpenDalWriter {
763 pub(crate) fn new(inner: opendal::Writer) -> Self {
764 Self {
765 inner,
766 bytes_written: 0,
767 }
768 }
769}
770
771#[async_trait]
772impl FileWrite for OpenDalWriter {
773 async fn write(&mut self, bs: Bytes) -> Result<()> {
774 let len = bs.len() as u64;
775 opendal::Writer::write(&mut self.inner, bs)
776 .await
777 .map_err(from_opendal_error)?;
778 self.bytes_written += len;
779 Ok(())
780 }
781
782 async fn close(&mut self) -> Result<FileMetadata> {
783 let metadata = opendal::Writer::close(&mut self.inner)
784 .await
785 .map_err(from_opendal_error)?;
786
787 let reported_size = metadata.content_length();
789 if reported_size != 0 && reported_size != self.bytes_written {
790 return Err(Error::new(
791 ErrorKind::Unexpected,
792 format!(
793 "Wrote {} bytes but storage reports {reported_size}",
794 self.bytes_written
795 ),
796 ));
797 }
798
799 Ok(FileMetadata {
800 size: self.bytes_written,
801 })
802 }
803}
804
805#[cfg(test)]
806mod tests {
807 use super::*;
808
809 fn client_config(value: &str) -> Result<OpenDalClientConfig> {
810 OpenDalClientConfig::from_properties(&HashMap::from([(
811 OPENDAL_IO_TIMEOUT_MS.to_string(),
812 value.to_string(),
813 )]))
814 }
815
816 #[test]
817 fn test_io_timeout_parsing() {
818 let unset = OpenDalClientConfig::from_properties(&HashMap::new()).unwrap();
819 assert_eq!(
820 unset.io_timeout(),
821 Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get())
822 );
823
824 let max = u64::MAX.to_string();
825 for (valid, ms) in [("45000", 45_000), ("1", 1), (max.as_str(), u64::MAX)] {
826 assert_eq!(
827 client_config(valid).unwrap().io_timeout(),
828 Duration::from_millis(ms),
829 "{valid}"
830 );
831 }
832
833 for invalid in ["0", "-1", "12.5", "abc", "", "18446744073709551616"] {
834 let err = client_config(invalid).unwrap_err().to_string();
835 assert!(err.contains(OPENDAL_IO_TIMEOUT_MS), "{invalid}");
836 assert!(err.contains(&format!("value: {invalid:?}")), "{err}");
837 let reason = invalid.parse::<NonZeroU64>().unwrap_err().to_string();
838 assert!(err.contains(&reason), "{err}");
839 }
840 }
841
842 #[test]
843 fn test_default_timeouts_match_opendal() {
844 let opendal_default = format!("{:?}", TimeoutLayer::new());
846 let layer = |timeout, io_timeout| {
847 format!(
848 "{:?}",
849 TimeoutLayer::new()
850 .with_timeout(timeout)
851 .with_io_timeout(io_timeout)
852 )
853 };
854 let control = Duration::from_secs(60);
855 let io = Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get());
856 assert_eq!(opendal_default, layer(control, io));
857 assert_ne!(opendal_default, layer(control + Duration::from_secs(1), io));
858 assert_ne!(
859 opendal_default,
860 layer(control, io + Duration::from_millis(1))
861 );
862 }
863
864 #[cfg(feature = "opendal-s3")]
865 #[test]
866 fn test_client_config_serde_round_trip() {
867 let storage = OpenDalStorage::S3 {
868 config: Arc::new(S3Config::default()),
869 customized_credential_load: None,
870 client_config: client_config("45000").unwrap(),
871 };
872
873 let mut value = serde_json::to_value(&storage).unwrap();
874 let restored: OpenDalStorage = serde_json::from_value(value.clone()).unwrap();
875 assert_eq!(
876 restored.client_config().io_timeout(),
877 Duration::from_secs(45)
878 );
879
880 let mut zero = value.clone();
881 zero["S3"]["client_config"]["io_timeout_ms"] = 0.into();
882 assert!(serde_json::from_value::<OpenDalStorage>(zero).is_err());
883
884 value["S3"].as_object_mut().unwrap().remove("client_config");
885 let restored: OpenDalStorage = serde_json::from_value(value).unwrap();
886 assert_eq!(
887 restored.client_config().io_timeout(),
888 Duration::from_millis(DEFAULT_IO_TIMEOUT_MS.get())
889 );
890 }
891
892 #[cfg(all(feature = "opendal-fs", feature = "opendal-memory"))]
893 #[test]
894 fn test_old_unit_variant_forms_are_rejected() {
895 for old in [r#""LocalFs""#, r#""Memory""#, r#"{"Memory":null}"#] {
896 assert!(
897 serde_json::from_str::<OpenDalStorage>(old).is_err(),
898 "{old}"
899 );
900 }
901 }
902
903 #[cfg(feature = "opendal-memory")]
904 #[test]
905 fn test_factory_rejects_invalid_io_timeout() {
906 let config = StorageConfig::new().with_prop(OPENDAL_IO_TIMEOUT_MS, "nope");
907
908 let err = OpenDalStorageFactory::Memory.build(&config).unwrap_err();
909 assert!(err.to_string().contains(OPENDAL_IO_TIMEOUT_MS));
910 }
911
912 #[test]
913 fn test_factory_propagates_io_timeout() {
914 let config = StorageConfig::new().with_prop(OPENDAL_IO_TIMEOUT_MS, "45000");
915 let factories = [
916 #[cfg(feature = "opendal-memory")]
917 OpenDalStorageFactory::Memory,
918 #[cfg(feature = "opendal-fs")]
919 OpenDalStorageFactory::Fs,
920 #[cfg(feature = "opendal-s3")]
921 OpenDalStorageFactory::S3 {
922 customized_credential_load: None,
923 },
924 #[cfg(feature = "opendal-gcs")]
925 OpenDalStorageFactory::Gcs,
926 #[cfg(feature = "opendal-oss")]
927 OpenDalStorageFactory::Oss,
928 #[cfg(feature = "opendal-azdls")]
929 OpenDalStorageFactory::Azdls,
930 #[cfg(feature = "opendal-hf")]
931 OpenDalStorageFactory::Hf,
932 ];
933 for factory in factories {
934 let storage = factory.build(&config).unwrap();
936 let json = serde_json::to_value(&*storage).unwrap();
937 let client_config = json
938 .as_object()
939 .unwrap()
940 .values()
941 .find_map(|variant| variant.get("client_config"))
942 .unwrap_or_else(|| panic!("{json}"));
943 assert_eq!(client_config["io_timeout_ms"], 45_000, "{factory:?}");
944 }
945 }
946
947 #[cfg(feature = "opendal-memory")]
948 #[tokio::test(start_paused = true)]
949 async fn test_io_timeout_reaches_timeout_layer() {
950 use opendal::layers::ConcurrentLimitLayer;
951
952 let storage = OpenDalStorage::Memory {
954 operator: default_memory_operator().layer(ConcurrentLimitLayer::new(0)),
955 client_config: client_config("45000").unwrap(),
956 };
957
958 let err = storage
959 .read("memory:/stalled")
960 .await
961 .unwrap_err()
962 .to_string();
963 assert!(err.contains("io operation timeout reached"), "{err}");
964 assert!(err.contains("{ timeout: 45 }"), "{err}");
965 }
966
967 #[cfg(feature = "opendal-s3")]
968 #[derive(Debug)]
969 struct EmptyCredentialLoader;
970
971 #[cfg(feature = "opendal-s3")]
972 impl ProvideCredential for EmptyCredentialLoader {
973 type Credential = AwsCredential;
974
975 async fn provide_credential(
976 &self,
977 _ctx: &reqsign_core::Context,
978 ) -> reqsign_core::Result<Option<AwsCredential>> {
979 Ok(None)
980 }
981 }
982
983 #[cfg(feature = "opendal-s3")]
984 #[test]
985 fn test_s3_factory_custom_credential_loader_serialization_fails() {
986 let file_io = iceberg::io::FileIOBuilder::new(Arc::new(OpenDalStorageFactory::S3 {
987 customized_credential_load: Some(CustomAwsCredentialLoader::new(EmptyCredentialLoader)),
988 }))
989 .build();
990
991 let err = file_io.serialize_all().unwrap_err();
992 assert!(
993 err.to_string()
994 .contains("custom AWS credential loaders cannot be serialized")
995 );
996 }
997
998 #[cfg(feature = "opendal-memory")]
999 #[test]
1000 fn test_default_memory_operator() {
1001 let op = default_memory_operator();
1002 assert_eq!(op.info().scheme().to_string(), "memory");
1003 }
1004
1005 #[cfg(feature = "opendal-memory")]
1006 #[tokio::test]
1007 async fn test_writer_close_returns_stored_size() {
1008 use iceberg::encryption::{EncryptedOutputFile, StandardKeyMetadata};
1009
1010 let storage = Arc::new(OpenDalStorage::Memory {
1013 operator: default_memory_operator(),
1014 client_config: OpenDalClientConfig::default(),
1015 });
1016 let path = "memory:///stored-size";
1017 for plaintext in [
1018 Bytes::new(),
1019 Bytes::from_static(b"test data"),
1020 Bytes::from(vec![7; 3 * 1024]),
1021 ] {
1022 let mut writer = storage.writer(path).await.unwrap();
1023 for chunk in plaintext.chunks(1024) {
1024 writer.write(Bytes::copy_from_slice(chunk)).await.unwrap();
1025 }
1026 let metadata = writer.close().await.unwrap();
1027 assert_eq!(metadata.size, plaintext.len() as u64);
1028 assert_eq!(metadata.size, storage.metadata(path).await.unwrap().size);
1029
1030 let output = EncryptedOutputFile::new(
1031 OutputFile::new(storage.clone(), path.to_string()),
1032 StandardKeyMetadata::try_new(b"0123456789abcdef").unwrap(),
1033 );
1034 let metadata = output.write(plaintext.clone()).await.unwrap();
1035 assert!(metadata.size > plaintext.len() as u64);
1036 assert_eq!(metadata.size, storage.metadata(path).await.unwrap().size);
1037 }
1038 }
1039
1040 #[cfg(feature = "opendal-memory")]
1041 #[test]
1042 fn test_relativize_path_memory() {
1043 let storage = OpenDalStorage::Memory {
1044 operator: default_memory_operator(),
1045 client_config: OpenDalClientConfig::default(),
1046 };
1047
1048 assert_eq!(
1049 storage.relativize_path("memory:/path/to/file").unwrap(),
1050 "path/to/file"
1051 );
1052 assert_eq!(
1054 storage.relativize_path("/path/to/file").unwrap(),
1055 "path/to/file"
1056 );
1057 }
1058
1059 #[cfg(feature = "opendal-fs")]
1060 #[test]
1061 fn test_relativize_path_fs() {
1062 let storage = OpenDalStorage::LocalFs {
1063 client_config: OpenDalClientConfig::default(),
1064 };
1065
1066 assert_eq!(
1067 storage
1068 .relativize_path("file:/tmp/data/file.parquet")
1069 .unwrap(),
1070 "tmp/data/file.parquet"
1071 );
1072 assert_eq!(
1073 storage.relativize_path("/tmp/data/file.parquet").unwrap(),
1074 "tmp/data/file.parquet"
1075 );
1076 }
1077
1078 #[cfg(feature = "opendal-s3")]
1079 #[test]
1080 fn test_relativize_path_s3() {
1081 let storage = OpenDalStorage::S3 {
1082 config: Arc::new(S3Config::default()),
1083 customized_credential_load: None,
1084 client_config: OpenDalClientConfig::default(),
1085 };
1086
1087 for scheme in ["s3", "s3a", "s3n", "minio"] {
1091 assert_eq!(
1092 storage
1093 .relativize_path(&format!("{scheme}://my-bucket/path/to/file.parquet"))
1094 .unwrap(),
1095 "path/to/file.parquet"
1096 );
1097 }
1098 }
1099
1100 #[cfg(feature = "opendal-gcs")]
1101 #[test]
1102 fn test_relativize_path_gcs() {
1103 let storage = OpenDalStorage::Gcs {
1104 config: Arc::new(GcsConfig::default()),
1105 client_config: OpenDalClientConfig::default(),
1106 };
1107
1108 assert_eq!(
1109 storage
1110 .relativize_path("gs://my-bucket/path/to/file.parquet")
1111 .unwrap(),
1112 "path/to/file.parquet"
1113 );
1114 }
1115
1116 #[cfg(feature = "opendal-gcs")]
1117 #[test]
1118 fn test_relativize_path_gcs_invalid_scheme() {
1119 let storage = OpenDalStorage::Gcs {
1120 config: Arc::new(GcsConfig::default()),
1121 client_config: OpenDalClientConfig::default(),
1122 };
1123
1124 assert!(
1125 storage
1126 .relativize_path("s3://my-bucket/path/to/file.parquet")
1127 .is_err()
1128 );
1129 }
1130
1131 #[cfg(feature = "opendal-oss")]
1132 #[test]
1133 fn test_relativize_path_oss() {
1134 let storage = OpenDalStorage::Oss {
1135 config: Arc::new(OssConfig::default()),
1136 client_config: OpenDalClientConfig::default(),
1137 };
1138
1139 assert_eq!(
1140 storage
1141 .relativize_path("oss://my-bucket/path/to/file.parquet")
1142 .unwrap(),
1143 "path/to/file.parquet"
1144 );
1145 }
1146
1147 #[cfg(feature = "opendal-oss")]
1148 #[test]
1149 fn test_relativize_path_oss_invalid_scheme() {
1150 let storage = OpenDalStorage::Oss {
1151 config: Arc::new(OssConfig::default()),
1152 client_config: OpenDalClientConfig::default(),
1153 };
1154
1155 assert!(
1156 storage
1157 .relativize_path("s3://my-bucket/path/to/file.parquet")
1158 .is_err()
1159 );
1160 }
1161
1162 #[cfg(feature = "opendal-azdls")]
1163 #[test]
1164 fn test_relativize_path_azdls() {
1165 let storage = OpenDalStorage::Azdls {
1166 config: Arc::new(AzdlsConfig {
1167 account_name: Some("myaccount".to_string()),
1168 endpoint: Some("https://myaccount.dfs.core.windows.net".to_string()),
1169 ..Default::default()
1170 }),
1171 client_config: OpenDalClientConfig::default(),
1172 };
1173
1174 assert_eq!(
1175 storage
1176 .relativize_path("abfss://myfs@myaccount.dfs.core.windows.net/path/to/file.parquet")
1177 .unwrap(),
1178 "path/to/file.parquet"
1179 );
1180 }
1181}