Skip to main content

iceberg_storage_opendal/
lib.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
18//! OpenDAL-based storage implementation for Apache Iceberg.
19//!
20//! This crate provides [`OpenDalStorage`] and [`OpenDalStorageFactory`],
21//! which implement the [`Storage`] and
22//! [`StorageFactory`] traits from the `iceberg` crate
23//! using [OpenDAL](https://opendal.apache.org/) as the backend.
24
25mod 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
106/// Timeout in milliseconds for each IO call on a reader, writer, lister or deleter, applied per
107/// retry attempt via OpenDAL's `TimeoutLayer::with_io_timeout`. Defaults to 10 seconds.
108///
109/// OpenDAL's separate `timeout`, which bounds whole control operations such as `stat`, is not
110/// affected and stays at its 60-second default.
111pub const OPENDAL_IO_TIMEOUT_MS: &str = "opendal.io-timeout-ms";
112
113/// Matches OpenDAL's `TimeoutLayer` default.
114const DEFAULT_IO_TIMEOUT_MS: NonZeroU64 = NonZeroU64::new(10_000).unwrap();
115
116/// Backend-independent client settings shared by every [`OpenDalStorage`] variant.
117#[derive(Clone, Debug, Properties, Serialize, Deserialize)]
118#[serde(default)]
119pub struct OpenDalClientConfig {
120    /// IO timeout in milliseconds.
121    #[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/// OpenDAL-based storage factory.
155///
156/// Maps scheme to the corresponding OpenDalStorage storage variant.
157/// Use this factory with `FileIOBuilder::new(factory)` to create FileIO instances.
158///
159/// # Serialization
160///
161/// The receiving binary must enable the feature corresponding to the serialized backend variant.
162/// For example, deserializing `OpenDalStorageFactory::S3` requires the `opendal-s3` feature.
163///
164/// Serialization fails when the `OpenDalStorageFactory::S3` variant contains a custom AWS
165/// credential loader because the loader holds process-local state that cannot be reconstructed in
166/// another process. Construct the factory without a custom loader before serializing it.
167#[derive(Clone, Debug, Serialize, Deserialize)]
168pub enum OpenDalStorageFactory {
169    /// Memory storage factory.
170    #[cfg(feature = "opendal-memory")]
171    Memory,
172    /// Local filesystem storage factory.
173    #[cfg(feature = "opendal-fs")]
174    Fs,
175    /// S3 storage factory.
176    #[cfg(feature = "opendal-s3")]
177    S3 {
178        /// Custom AWS credential loader.
179        #[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    /// GCS storage factory.
187    #[cfg(feature = "opendal-gcs")]
188    Gcs,
189    /// OSS storage factory.
190    #[cfg(feature = "opendal-oss")]
191    Oss,
192    /// Azure Data Lake Storage factory.
193    #[cfg(feature = "opendal-azdls")]
194    Azdls,
195    /// HuggingFace Hub storage factory.
196    #[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/// Default memory operator for serde deserialization.
272#[cfg(feature = "opendal-memory")]
273fn default_memory_operator() -> Operator {
274    memory_config_build().expect("Failed to create default memory operator")
275}
276
277/// OpenDAL-based storage implementation.
278///
279/// The serialized representation is not a stable format and may change between crate versions.
280#[derive(Clone, Debug, Serialize, Deserialize)]
281pub enum OpenDalStorage {
282    /// Memory storage variant.
283    #[cfg(feature = "opendal-memory")]
284    Memory {
285        /// Pre-built memory operator.
286        #[serde(skip, default = "self::default_memory_operator")]
287        operator: Operator,
288        /// Backend-independent client settings.
289        #[serde(default)]
290        client_config: OpenDalClientConfig,
291    },
292    /// Local filesystem storage variant.
293    #[cfg(feature = "opendal-fs")]
294    LocalFs {
295        /// Backend-independent client settings.
296        #[serde(default)]
297        client_config: OpenDalClientConfig,
298    },
299    /// S3 storage variant.
300    ///
301    /// Accepts any S3-family URL (`s3://`, `s3a://`, `s3n://`); the scheme is
302    /// derived from the path at call time.
303    #[cfg(feature = "opendal-s3")]
304    S3 {
305        /// S3 configuration.
306        config: Arc<S3Config>,
307        /// Custom AWS credential loader.
308        #[serde(skip)]
309        customized_credential_load: Option<CustomAwsCredentialLoader>,
310        /// Backend-independent client settings.
311        #[serde(default)]
312        client_config: OpenDalClientConfig,
313    },
314    /// GCS storage variant.
315    #[cfg(feature = "opendal-gcs")]
316    Gcs {
317        /// GCS configuration.
318        config: Arc<GcsConfig>,
319        /// Backend-independent client settings.
320        #[serde(default)]
321        client_config: OpenDalClientConfig,
322    },
323    /// OSS storage variant.
324    #[cfg(feature = "opendal-oss")]
325    Oss {
326        /// OSS configuration.
327        config: Arc<OssConfig>,
328        /// Backend-independent client settings.
329        #[serde(default)]
330        client_config: OpenDalClientConfig,
331    },
332    /// Azure Data Lake Storage variant.
333    ///
334    /// Accepts paths of the form
335    /// `abfs[s]://<filesystem>@<account>.dfs.<endpoint-suffix>/<path>` or
336    /// `wasb[s]://<container>@<account>.blob.<endpoint-suffix>/<path>`.
337    /// The scheme is derived from the path at call time.
338    #[cfg(feature = "opendal-azdls")]
339    Azdls {
340        /// Azure DLS configuration.
341        config: Arc<AzdlsConfig>,
342        /// Backend-independent client settings.
343        #[serde(default)]
344        client_config: OpenDalClientConfig,
345    },
346    /// HuggingFace Hub storage variant.
347    ///
348    /// Accepts paths of the form
349    /// `hf://<repo_type>/<owner>/<repo>[@<revision>]/<path_in_repo>`,
350    /// where `<repo_type>` must be one of `models`, `datasets`, `spaces`, or `buckets`.
351    #[cfg(feature = "opendal-hf")]
352    Hf {
353        /// HuggingFace Hub configuration (token + endpoint).
354        config: Arc<HfConfig>,
355        /// Backend-independent client settings.
356        #[serde(default)]
357        client_config: OpenDalClientConfig,
358    },
359}
360
361impl OpenDalStorage {
362    /// Creates operator from path.
363    ///
364    /// # Arguments
365    ///
366    /// * path: It should be *absolute* path starting with scheme string used to construct [`FileIO`](iceberg::io::FileIO).
367    ///
368    /// # Returns
369    ///
370    /// The return value consists of two parts:
371    ///
372    /// * An [`opendal::Operator`] instance used to operate on file.
373    /// * Relative path to the root uri of [`opendal::Operator`].
374    #[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                // Use the URL scheme in the path for prefix matching. This enables
408                // use of S3-compatible storage backends using custom schemes (e.g., `minio://`, `r2://`).
409                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        // Apply observability/resilience layers. TimeoutLayer must be
472        // inside RetryLayer so each retry attempt is independently
473        // bounded — without a per-attempt timeout, a future parked on a
474        // silently dropped TCP connection never produces an `Err` and
475        // RetryLayer cannot retry, leaving the caller hung indefinitely.
476        // See: https://opendal.apache.org/docs/rust/opendal/layers/struct.TimeoutLayer.html
477        //
478        // Transient errors are common for object stores; we retry temporary
479        // failures with exponential backoff. The retry behavior also
480        // benefits non-object-store backends.
481        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    /// Returns a cache key used by `delete_stream` to group paths by storage operator.
517    ///
518    /// For most backends the URL host (bucket name) is sufficient. For HF the host
519    /// encodes the repo type, not the repo identity, so a more specific key is used.
520    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    /// Extracts the relative path from an absolute path without building an operator.
532    ///
533    /// This is a lightweight alternative to [`create_operator`](Self::create_operator) for cases
534    /// where only the relative path is needed (e.g. bulk deletes where the operator is already
535    /// available).
536    #[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
739// Newtype wrappers for opendal types to satisfy orphan rules.
740// We can't implement iceberg's FileRead/FileWrite traits directly on opendal's
741// Reader/Writer since neither trait nor type is defined in this crate.
742
743/// Wrapper around `opendal::Reader` that implements `FileRead`.
744pub(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
756/// Wrapper around `opendal::Writer` that implements `FileWrite`.
757pub(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        // Object stores may omit the size (reported as 0); validate only a reported nonzero size.
788        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        // `TimeoutLayer` has no getters, so compare through `Debug`.
845        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            // `build` returns `dyn Storage`, so read the config back from its serialized form.
935            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        // A zero concurrency limit stalls every IO call; paused time skips the waits.
953        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        // Note: the memory service does report a content length, so this only pins the happy
1011        // path. The counter in `OpenDalWriter` is what covers services that don't, such as S3.
1012        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        // Without the scheme prefix, falls back to stripping the leading slash
1053        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        // All S3-family schemes are accepted by the same storage instance.
1088        // Custom schemes for S3-compatible stores (e.g., `minio://`) are also
1089        // accepted because the path's scheme is used as-is for prefix matching.
1090        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}