laskoviymishka commented on code in PR #3111:
URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4204790164


##########
crates/storage/opendal/tests/hdfs_native_runtime_test.rs:
##########
@@ -0,0 +1,78 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! A cached HDFS operator is bound to the tokio runtime that built it; once
+//! that runtime is gone it must be rebuilt rather than reused. Needs no
+//! HDFS: the NameNode is a local listener that only counts dials.
+
+#[cfg(feature = "opendal-hdfs-native")]
+mod tests {
+    use std::net::TcpListener;
+    use std::sync::Arc;
+    use std::sync::atomic::{AtomicUsize, Ordering};
+    use std::time::Duration;
+
+    use iceberg::io::{FileIO, FileIOBuilder, HDFS_NAME_NODE};
+    use iceberg_storage_opendal::OpenDalStorageFactory;
+
+    /// Accepts and immediately closes connections, counting them.
+    fn fake_name_node() -> (u16, Arc<AtomicUsize>) {
+        let listener = TcpListener::bind("127.0.0.1:0").unwrap();
+        let port = listener.local_addr().unwrap().port();
+        let dials = Arc::new(AtomicUsize::new(0));
+        let counter = dials.clone();
+        std::thread::spawn(move || {
+            for stream in listener.incoming() {
+                let Ok(_stream) = stream else { break };
+                counter.fetch_add(1, Ordering::SeqCst);
+            }
+        });
+        (port, dials)
+    }
+
+    fn stat(runtime: &tokio::runtime::Runtime, file_io: &FileIO) {
+        runtime.block_on(async {
+            let input = file_io.new_input("hdfs:///f").unwrap();
+            // The fake NameNode never answers, so this fails; only the dial 
matters.
+            let _ = tokio::time::timeout(Duration::from_secs(10), 
input.metadata()).await;
+        });
+    }
+
+    #[test]
+    fn test_hdfs_operator_is_rebuilt_after_its_runtime_is_dropped() {
+        let (port, dials) = fake_name_node();
+        let file_io = 
FileIOBuilder::new(Arc::new(OpenDalStorageFactory::HdfsNative))
+            .with_prop(HDFS_NAME_NODE, format!("hdfs://127.0.0.1:{port}"))
+            .with_prop("hadoop.dfs.client.failover.max.attempts", "1")
+            .build();
+
+        let first = tokio::runtime::Runtime::new().unwrap();
+        stat(&first, &file_io);
+        let after_first = dials.load(Ordering::SeqCst);
+        assert!(after_first >= 1, "the NameNode was never dialed");
+        drop(first);
+
+        // Reusing the FileIO from another runtime used to panic inside
+        // hdfs-native, whose client was bound to the dropped runtime.
+        let second = tokio::runtime::Runtime::new().unwrap();
+        stat(&second, &file_io);
+        assert!(
+            dials.load(Ordering::SeqCst) > after_first,

Review Comment:
   Nit — this can pass without a dial from the second runtime. `after_first` is 
sampled right after the 10s timeout, and the counter is bumped asynchronously 
by the listener thread, so a straggling retry from the first runtime satisfies 
`> after_first` on its own. I'd let the counter settle before sampling 
`after_first`, and ideally assert the cache entry's identity actually changed 
(a test-only generation counter, or the `Operator` pointer) so we're proving 
the rebuild rather than just a dial.



##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1264 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no 
JNI).
+
+use std::collections::HashMap;
+use std::collections::hash_map::Entry;
+use std::sync::{Arc, RwLock, Weak};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, 
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use serde::{Deserialize, Serialize};
+use tokio::runtime::Handle;
+use tokio::task::JoinHandle;
+use url::Url;
+
+use crate::OpenDalClientConfig;
+use crate::utils::from_opendal_error;
+
+/// Hadoop's default filesystem, which serves authority-less paths.
+const FS_DEFAULT_FS: &str = "fs.defaultFS";
+const HDFS_DEFAULT_PORT: u16 = 8020;
+/// PyIceberg keys with no equivalent in opendal's config.
+const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"];
+/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>`
+/// expands to them, and they are honored as well when passed through 
`hadoop.`.
+const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes";
+const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address";
+
+/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path
+/// authorities take, so every source shares cache keys. `hdfs-native` dials
+/// a socket address and has no default port, so anything else is `None`:
+/// another scheme, a logical name, port 0, userinfo, a path.
+fn hdfs_native_name_node(entry: &str) -> Option<String> {
+    let rest = entry.trim().trim_end_matches('/');
+    let rest = rest.strip_prefix("hdfs://").unwrap_or(rest);
+    if rest.is_empty() || rest.contains("://") {
+        return None;
+    }
+    let url = Url::parse(&format!("hdfs://{rest}")).ok()?;
+    let (host, port) = (url.host_str()?, url.port()?);
+    let plain = host.is_empty()
+        || port == 0
+        || !url.username().is_empty()
+        || url.password().is_some()
+        || !matches!(url.path(), "" | "/")
+        || url.query().is_some()
+        || url.fragment().is_some();
+    (!plain).then(|| format!("hdfs://{host}:{port}"))
+}
+
+/// Parses a comma-separated NameNode list, each entry normalized; an empty
+/// list is `Ok` and empty.
+fn hdfs_native_name_node_list(property: &str, value: &str) -> 
Result<Vec<String>> {
+    value
+        .split(',')
+        .map(str::trim)
+        .filter(|entry| !entry.is_empty())
+        .map(|entry| {
+            hdfs_native_name_node(entry).ok_or_else(|| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!(
+                        "Invalid `{property}` entry: {entry}, expected 
host:port (hdfs:// optional)"
+                    ),
+                )
+            })
+        })
+        .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) -> 
Result<HdfsNativeConfig> {
+    hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+    mut m: HashMap<String, String>,
+    mut warn: impl FnMut(String),
+) -> Result<HdfsNativeConfig> {
+    let mut cfg = HdfsNativeConfig::default();
+
+    // Entries are trimmed one by one: opendal splits the list on `,` as is,
+    // so a space after a comma would break failover to that NameNode. An
+    // empty result is dropped because `Operator::from_config` bypasses the
+    // builder's empty-string guard and `Some("")` would shadow the
+    // path-authority fallback below.
+    if let Some(name_node) = m.remove(HDFS_NAME_NODE) {
+        let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?;
+        if !entries.is_empty() {
+            cfg.name_node = Some(entries.join(","));
+        }
+    }
+
+    // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of
+    // an HA nameservice, which the resolver reads back from the options.
+    let nameservice_prefix = format!("{HDFS_NAME_NODE}.");
+    let declared_keys: Vec<String> = m
+        .keys()
+        .filter(|key| key.starts_with(&nameservice_prefix))
+        .cloned()
+        .collect();
+    let mut declared = Vec::new();
+    for key in declared_keys {
+        let value = m.remove(&key).unwrap_or_default();
+        let nameservice = key[nameservice_prefix.len()..].to_string();
+        if nameservice.is_empty() || 
nameservice.chars().any(char::is_whitespace) {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!(
+                    "Invalid property `{key}`: a nameservice name must follow 
`{nameservice_prefix}`"
+                ),
+            ));
+        }
+        let entries = hdfs_native_name_node_list(&key, &value)?;
+        if entries.is_empty() {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!("Invalid property `{key}`: no NameNodes"),
+            ));
+        }
+        declared.push((nameservice, entries));
+    }
+
+    // A config carried over from PyIceberg would otherwise change identity
+    // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache.
+    for key in HDFS_UNSUPPORTED_KEYS {
+        if m.remove(key).is_some() {
+            warn(format!(
+                "`{key}` is not supported by the hdfs-native backend and is 
ignored"
+            ));
+        }
+    }
+    if m.contains_key(HDFS_HADOOP_CONF_PREFIX) {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!(
+                "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key 
must follow the prefix"
+            ),
+        ));
+    }
+
+    let host = m
+        .remove(HDFS_HOST)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty());
+    let port = m
+        .remove(HDFS_PORT)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty())
+        .map(|port| {
+            port.parse::<u16>().map_err(|e| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+                )
+            })
+        })
+        .transpose()?;
+
+    let mut options: HashMap<String, String> = m
+        .into_iter()
+        .filter_map(|(key, value)| {
+            key.strip_prefix(HDFS_HADOOP_CONF_PREFIX)
+                .map(|stripped| (stripped.to_string(), value))
+        })
+        .collect();
+    // Explicit `hadoop.` keys win over the sugar.
+    for (nameservice, entries) in declared {
+        let ids: Vec<String> = (0..entries.len()).map(|i| 
format!("nn{i}")).collect();
+        options
+            .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}"))
+            .or_insert_with(|| ids.join(","));
+        for (id, entry) in ids.iter().zip(&entries) {
+            options
+                .entry(format!(
+                    "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}"
+                ))
+                .or_insert_with(|| 
entry.trim_start_matches("hdfs://").to_string());
+        }
+    }
+    // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for
+    // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an
+    // explicit `hadoop.fs.defaultFS` wins.
+    match host {
+        Some(host) => {
+            // An IPv6 literal needs brackets in a URI authority.
+            let bracketed = if host.contains(':') && !host.starts_with('[') {
+                format!("[{host}]")
+            } else {
+                host.clone()
+            };
+            let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+            // Validated like every other NameNode spelling, so a host that
+            // carries a scheme or port fails here and not at the first I/O.
+            let default_fs = 
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+                .ok_or_else(|| {
+                    Error::new(
+                        ErrorKind::DataInvalid,
+                        format!(
+                            "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: 
{host}:{port}, expected a host name or IP and a port"
+                        ),
+                    )
+                })?;
+            options
+                .entry(FS_DEFAULT_FS.to_string())
+                .or_insert(default_fs);
+        }
+        None if port.is_some() => {
+            warn(format!(
+                "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is 
ignored"
+            ));
+        }
+        None => {}
+    }
+    if !options.is_empty() {
+        cfg.options = Some(options);
+    }
+
+    Ok(cfg)
+}
+
+/// Parse an HDFS path into `Some("hdfs://<authority>")` (`None` when
+/// authority-less) and the relative path (no leading `/`, opendal style).
+pub(crate) fn hdfs_native_parse_path(path: &str) -> Result<(Option<String>, 
&str)> {
+    let url = Url::parse(path).map_err(|e| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}: {e}"),
+        )
+    })?;
+    // Non-special schemes parse even without `//` (e.g. `hdfs:x` is a valid
+    // non-hierarchical URL), so require the literal prefix before slicing.
+    let (Some(after_scheme), "hdfs") = (path.strip_prefix("hdfs://"), 
url.scheme()) else {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, expected scheme `hdfs://`"),
+        ));
+    };
+    // Userinfo has nowhere to go and port 0 cannot be dialed; silently
+    // dropping the one or reading the other as a logical name would mislead.
+    if !url.username().is_empty() || url.password().is_some() {
+        // Not echoing the path: it may carry a password.
+        let host = url.host_str().unwrap_or_default();
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path for host `{host}`: userinfo is not 
supported"),
+        ));
+    }
+    if url.port() == Some(0) {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, port 0 cannot be dialed"),
+        ));
+    }
+
+    let name_node = url.host_str().filter(|h| !h.is_empty()).map(|host| {
+        url.port()
+            .map(|port| format!("hdfs://{host}:{port}"))
+            .unwrap_or_else(|| format!("hdfs://{host}"))
+    });
+
+    // `url.path()` borrows from `url` and can't be returned with the input's
+    // lifetime. Slice the path component out of the original input instead;
+    // it starts after the first `/` following the `hdfs://` prefix. Opendal
+    // paths must not start with `/` (`Deleter::delete` rejects them).
+    let rel = match after_scheme.find('/') {
+        Some(i) => after_scheme[i..].trim_start_matches('/'),
+        None => "",
+    };
+
+    Ok((name_node, rel))
+}
+
+/// Resolves the effective NameNode for a path, plus the relative path. As in
+/// Hadoop, an authority with a port is used as is; a logical nameservice
+/// authority (no port) resolves through its declaration, and an
+/// authority-less path through `hdfs.name-node`, else `fs.defaultFS`. The
+/// operator cache, `delete_stream` batching and `relativize_path` all go
+/// through this, so they cannot drift apart.
+pub(crate) fn hdfs_native_effective_name_node<'a>(
+    config: &HdfsNativeConfig,
+    path: &'a str,
+) -> Result<(String, &'a str)> {
+    let (authority, relative_path) = hdfs_native_parse_path(path)?;
+    let invalid = |reason: String| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, {reason}"),
+        )
+    };
+    let name_node = match authority {
+        Some(authority) if hdfs_native_name_node(&authority).is_some() => 
authority,
+        Some(logical) => {
+            let nameservice = logical.trim_start_matches("hdfs://");
+            hdfs_native_nameservice(config, nameservice)?.ok_or_else(|| {
+                invalid(format!(
+                    "logical nameservice `{nameservice}` is not declared; set 
`{HDFS_NAME_NODE}.{nameservice}`"
+                ))
+            })?
+        }
+        None => match (&config.name_node, hdfs_native_default_fs(config)) {
+            (Some(name_node), _) => name_node.clone(),
+            (None, Some(default_fs)) => 
hdfs_native_name_node(default_fs).ok_or_else(|| {
+                invalid(format!(
+                    "`{FS_DEFAULT_FS}` {default_fs} is not an HDFS host:port, 
a logical nameservice requires `{HDFS_NAME_NODE}`"
+                ))
+            })?,
+            (None, None) => {
+                return Err(invalid(format!(
+                    "authority-less paths require `{HDFS_NAME_NODE}` or 
`{HDFS_HOST}`"
+                )));
+            }
+        },
+    };
+    Ok((name_node, relative_path))
+}
+
+/// `fs.defaultFS` from the forwarded options, as written.
+fn hdfs_native_default_fs(config: &HdfsNativeConfig) -> Option<&str> {
+    config
+        .options
+        .as_ref()?
+        .get(FS_DEFAULT_FS)
+        .map(|s| s.trim())
+        .filter(|s| !s.is_empty())
+}
+
+/// The NameNodes that Hadoop's keys in the forwarded options declare for a
+/// nameservice, if any. A declaration with a missing or malformed address is
+/// an error rather than a silent fallback.
+fn hdfs_native_nameservice(config: &HdfsNativeConfig, nameservice: &str) -> 
Result<Option<String>> {
+    let Some(options) = config.options.as_ref() else {
+        return Ok(None);
+    };
+    let Some(ids) = 
options.get(&format!("{HA_NAMENODES_PREFIX}.{nameservice}")) else {
+        return Ok(None);
+    };
+    let name_nodes = ids
+        .split(',')
+        .map(str::trim)
+        .filter(|id| !id.is_empty())
+        .map(|id| {
+            let key = 
format!("{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}");
+            options
+                .get(&key)
+                .and_then(|value| hdfs_native_name_node(value))
+                .ok_or_else(|| {
+                    Error::new(
+                        ErrorKind::DataInvalid,
+                        format!(
+                            "Nameservice `{nameservice}` declares NameNode 
`{id}` but `{key}` is missing or not host:port"
+                        ),
+                    )
+                })
+        })
+        .collect::<Result<Vec<_>>>()?;
+    if name_nodes.is_empty() {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Nameservice `{nameservice}` declares no NameNodes"),
+        ));
+    }
+    Ok(Some(name_nodes.join(",")))
+}
+
+/// State of [`OpenDalStorage::HdfsNative`](crate::OpenDalStorage::HdfsNative):
+/// the parsed configuration and the per-NameNode operator cache. Only the
+/// storage factories build it.
+#[derive(Clone, Debug, Serialize, Deserialize)]
+pub struct HdfsNativeStorage {
+    pub(crate) config: Arc<HdfsNativeConfig>,
+    #[serde(skip, default)]
+    pub(crate) operators: HdfsNativeOperatorCache,
+    #[serde(default)]
+    pub(crate) client_config: OpenDalClientConfig,
+}
+
+impl HdfsNativeStorage {
+    pub(crate) fn new(config: HdfsNativeConfig, client_config: 
OpenDalClientConfig) -> Self {
+        Self {
+            config: Arc::new(config),
+            operators: HdfsNativeOperatorCache::default(),
+            client_config,
+        }
+    }
+}
+
+/// Operators cached per effective NameNode: each holds an `hdfs-native`
+/// client with live RPC connections, whose tasks run on the tokio runtime
+/// current when it was built. An entry is rebuilt once that runtime is
+/// gone, as `hdfs-native` panics when
+/// it spawns onto a dead one. The cache lives as long as the storage that
+/// owns it (clones share it).
+#[derive(Clone, Debug, Default)]
+pub(crate) struct HdfsNativeOperatorCache(Arc<RwLock<HashMap<String, 
CachedOperator>>>);
+
+#[derive(Debug)]
+struct CachedOperator {
+    operator: Operator,
+    sentinel: RuntimeSentinel,
+}
+
+/// A task parked on the building runtime that owns the token, so the token
+/// outlives it only while that runtime is alive. Aborted on drop so entries
+/// do not leave parked tasks behind.
+#[derive(Debug)]
+struct RuntimeSentinel {
+    alive: Weak<()>,
+    task: JoinHandle<()>,
+}
+
+impl RuntimeSentinel {
+    fn spawn(handle: &Handle) -> Self {
+        let token = Arc::new(());
+        let alive = Arc::downgrade(&token);
+        let task = handle.spawn(async move {
+            let _token = token;
+            std::future::pending::<()>().await
+        });
+        Self { alive, task }
+    }
+}
+
+impl Drop for RuntimeSentinel {
+    fn drop(&mut self) {
+        self.task.abort();
+    }
+}
+
+impl CachedOperator {
+    fn new(operator: Operator, handle: &Handle) -> Self {
+        Self {
+            operator,
+            sentinel: RuntimeSentinel::spawn(handle),
+        }
+    }
+
+    fn runtime_alive(&self) -> bool {
+        self.sentinel.alive.strong_count() > 0
+    }
+}
+
+impl HdfsNativeOperatorCache {
+    pub(crate) fn get(&self, name_node: &str) -> Result<Option<Operator>> {
+        Ok(self
+            .0
+            .read()
+            .map_err(poisoned)?

Review Comment:
   Nit, and separate from the lock thread from last round (that one's fine — no 
guard held across an `.await`): a poisoned lock here turns every subsequent 
HDFS call into a permanent `Unexpected` error. The map is only ever mutated as 
whole-entry operations, so a poison just means some other thread panicked — it 
shouldn't take down all HDFS I/O with it. I'd recover the guard with 
`unwrap_or_else(PoisonError::into_inner)` at both the read and write sites 
instead of mapping it to an error.



##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1264 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no 
JNI).
+
+use std::collections::HashMap;
+use std::collections::hash_map::Entry;
+use std::sync::{Arc, RwLock, Weak};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, 
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use serde::{Deserialize, Serialize};
+use tokio::runtime::Handle;
+use tokio::task::JoinHandle;
+use url::Url;
+
+use crate::OpenDalClientConfig;
+use crate::utils::from_opendal_error;
+
+/// Hadoop's default filesystem, which serves authority-less paths.
+const FS_DEFAULT_FS: &str = "fs.defaultFS";
+const HDFS_DEFAULT_PORT: u16 = 8020;
+/// PyIceberg keys with no equivalent in opendal's config.
+const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"];
+/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>`
+/// expands to them, and they are honored as well when passed through 
`hadoop.`.
+const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes";
+const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address";
+
+/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path
+/// authorities take, so every source shares cache keys. `hdfs-native` dials
+/// a socket address and has no default port, so anything else is `None`:
+/// another scheme, a logical name, port 0, userinfo, a path.
+fn hdfs_native_name_node(entry: &str) -> Option<String> {
+    let rest = entry.trim().trim_end_matches('/');
+    let rest = rest.strip_prefix("hdfs://").unwrap_or(rest);
+    if rest.is_empty() || rest.contains("://") {
+        return None;
+    }
+    let url = Url::parse(&format!("hdfs://{rest}")).ok()?;
+    let (host, port) = (url.host_str()?, url.port()?);
+    let plain = host.is_empty()
+        || port == 0
+        || !url.username().is_empty()
+        || url.password().is_some()
+        || !matches!(url.path(), "" | "/")
+        || url.query().is_some()
+        || url.fragment().is_some();
+    (!plain).then(|| format!("hdfs://{host}:{port}"))
+}
+
+/// Parses a comma-separated NameNode list, each entry normalized; an empty
+/// list is `Ok` and empty.
+fn hdfs_native_name_node_list(property: &str, value: &str) -> 
Result<Vec<String>> {
+    value
+        .split(',')
+        .map(str::trim)
+        .filter(|entry| !entry.is_empty())
+        .map(|entry| {
+            hdfs_native_name_node(entry).ok_or_else(|| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!(
+                        "Invalid `{property}` entry: {entry}, expected 
host:port (hdfs:// optional)"
+                    ),
+                )
+            })
+        })
+        .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) -> 
Result<HdfsNativeConfig> {
+    hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+    mut m: HashMap<String, String>,
+    mut warn: impl FnMut(String),
+) -> Result<HdfsNativeConfig> {
+    let mut cfg = HdfsNativeConfig::default();
+
+    // Entries are trimmed one by one: opendal splits the list on `,` as is,
+    // so a space after a comma would break failover to that NameNode. An
+    // empty result is dropped because `Operator::from_config` bypasses the
+    // builder's empty-string guard and `Some("")` would shadow the
+    // path-authority fallback below.
+    if let Some(name_node) = m.remove(HDFS_NAME_NODE) {
+        let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?;
+        if !entries.is_empty() {
+            cfg.name_node = Some(entries.join(","));
+        }
+    }
+
+    // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of
+    // an HA nameservice, which the resolver reads back from the options.
+    let nameservice_prefix = format!("{HDFS_NAME_NODE}.");
+    let declared_keys: Vec<String> = m
+        .keys()
+        .filter(|key| key.starts_with(&nameservice_prefix))
+        .cloned()
+        .collect();
+    let mut declared = Vec::new();
+    for key in declared_keys {
+        let value = m.remove(&key).unwrap_or_default();
+        let nameservice = key[nameservice_prefix.len()..].to_string();
+        if nameservice.is_empty() || 
nameservice.chars().any(char::is_whitespace) {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!(
+                    "Invalid property `{key}`: a nameservice name must follow 
`{nameservice_prefix}`"
+                ),
+            ));
+        }
+        let entries = hdfs_native_name_node_list(&key, &value)?;
+        if entries.is_empty() {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!("Invalid property `{key}`: no NameNodes"),
+            ));
+        }
+        declared.push((nameservice, entries));
+    }
+
+    // A config carried over from PyIceberg would otherwise change identity
+    // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache.
+    for key in HDFS_UNSUPPORTED_KEYS {
+        if m.remove(key).is_some() {
+            warn(format!(
+                "`{key}` is not supported by the hdfs-native backend and is 
ignored"
+            ));
+        }
+    }
+    if m.contains_key(HDFS_HADOOP_CONF_PREFIX) {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!(
+                "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key 
must follow the prefix"
+            ),
+        ));
+    }
+
+    let host = m
+        .remove(HDFS_HOST)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty());
+    let port = m
+        .remove(HDFS_PORT)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty())
+        .map(|port| {
+            port.parse::<u16>().map_err(|e| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+                )
+            })
+        })
+        .transpose()?;
+
+    let mut options: HashMap<String, String> = m
+        .into_iter()
+        .filter_map(|(key, value)| {
+            key.strip_prefix(HDFS_HADOOP_CONF_PREFIX)
+                .map(|stripped| (stripped.to_string(), value))
+        })
+        .collect();
+    // Explicit `hadoop.` keys win over the sugar.
+    for (nameservice, entries) in declared {
+        let ids: Vec<String> = (0..entries.len()).map(|i| 
format!("nn{i}")).collect();
+        options
+            .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}"))
+            .or_insert_with(|| ids.join(","));
+        for (id, entry) in ids.iter().zip(&entries) {
+            options
+                .entry(format!(
+                    "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}"
+                ))
+                .or_insert_with(|| 
entry.trim_start_matches("hdfs://").to_string());
+        }
+    }
+    // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for
+    // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an
+    // explicit `hadoop.fs.defaultFS` wins.
+    match host {
+        Some(host) => {
+            // An IPv6 literal needs brackets in a URI authority.
+            let bracketed = if host.contains(':') && !host.starts_with('[') {
+                format!("[{host}]")
+            } else {
+                host.clone()
+            };
+            let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+            // Validated like every other NameNode spelling, so a host that
+            // carries a scheme or port fails here and not at the first I/O.
+            let default_fs = 
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+                .ok_or_else(|| {
+                    Error::new(
+                        ErrorKind::DataInvalid,
+                        format!(
+                            "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: 
{host}:{port}, expected a host name or IP and a port"
+                        ),
+                    )
+                })?;
+            options
+                .entry(FS_DEFAULT_FS.to_string())
+                .or_insert(default_fs);
+        }
+        None if port.is_some() => {
+            warn(format!(
+                "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is 
ignored"
+            ));
+        }
+        None => {}
+    }
+    if !options.is_empty() {
+        cfg.options = Some(options);
+    }
+
+    Ok(cfg)
+}
+
+/// Parse an HDFS path into `Some("hdfs://<authority>")` (`None` when
+/// authority-less) and the relative path (no leading `/`, opendal style).
+pub(crate) fn hdfs_native_parse_path(path: &str) -> Result<(Option<String>, 
&str)> {
+    let url = Url::parse(path).map_err(|e| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}: {e}"),

Review Comment:
   Nit — the comment a few lines down says we don't echo the path because it 
may carry a password, but this `Url::parse` failure (and the scheme-mismatch 
`return` right below it) both format the full `{path}`. 
`hdfs://user:pw@nn:bad/x` fails `parse` before the userinfo check runs, so the 
password lands in the error. I'd redact userinfo before echoing, or echo only 
the scheme plus the parser reason.



##########
crates/storage/opendal/src/hdfs_native.rs:
##########
@@ -0,0 +1,1264 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! HDFS storage backend via OpenDAL's `services-hdfs-native` (pure Rust, no 
JNI).
+
+use std::collections::HashMap;
+use std::collections::hash_map::Entry;
+use std::sync::{Arc, RwLock, Weak};
+
+use iceberg::io::{HDFS_HADOOP_CONF_PREFIX, HDFS_HOST, HDFS_NAME_NODE, 
HDFS_PORT};
+use iceberg::{Error, ErrorKind, Result};
+use opendal::Operator;
+use opendal::services::HdfsNativeConfig;
+use serde::{Deserialize, Serialize};
+use tokio::runtime::Handle;
+use tokio::task::JoinHandle;
+use url::Url;
+
+use crate::OpenDalClientConfig;
+use crate::utils::from_opendal_error;
+
+/// Hadoop's default filesystem, which serves authority-less paths.
+const FS_DEFAULT_FS: &str = "fs.defaultFS";
+const HDFS_DEFAULT_PORT: u16 = 8020;
+/// PyIceberg keys with no equivalent in opendal's config.
+const HDFS_UNSUPPORTED_KEYS: [&str; 2] = ["hdfs.user", "hdfs.kerberos_ticket"];
+/// Hadoop's keys declaring an HA nameservice: `hdfs.name-node.<nameservice>`
+/// expands to them, and they are honored as well when passed through 
`hadoop.`.
+const HA_NAMENODES_PREFIX: &str = "dfs.ha.namenodes";
+const HA_NAMENODE_RPC_ADDRESS_PREFIX: &str = "dfs.namenode.rpc-address";
+
+/// Normalizes one NameNode spelling to `hdfs://host:port`, the form path
+/// authorities take, so every source shares cache keys. `hdfs-native` dials
+/// a socket address and has no default port, so anything else is `None`:
+/// another scheme, a logical name, port 0, userinfo, a path.
+fn hdfs_native_name_node(entry: &str) -> Option<String> {
+    let rest = entry.trim().trim_end_matches('/');
+    let rest = rest.strip_prefix("hdfs://").unwrap_or(rest);
+    if rest.is_empty() || rest.contains("://") {
+        return None;
+    }
+    let url = Url::parse(&format!("hdfs://{rest}")).ok()?;
+    let (host, port) = (url.host_str()?, url.port()?);
+    let plain = host.is_empty()
+        || port == 0
+        || !url.username().is_empty()
+        || url.password().is_some()
+        || !matches!(url.path(), "" | "/")
+        || url.query().is_some()
+        || url.fragment().is_some();
+    (!plain).then(|| format!("hdfs://{host}:{port}"))
+}
+
+/// Parses a comma-separated NameNode list, each entry normalized; an empty
+/// list is `Ok` and empty.
+fn hdfs_native_name_node_list(property: &str, value: &str) -> 
Result<Vec<String>> {
+    value
+        .split(',')
+        .map(str::trim)
+        .filter(|entry| !entry.is_empty())
+        .map(|entry| {
+            hdfs_native_name_node(entry).ok_or_else(|| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!(
+                        "Invalid `{property}` entry: {entry}, expected 
host:port (hdfs:// optional)"
+                    ),
+                )
+            })
+        })
+        .collect()
+}
+
+/// Parse iceberg properties to [`HdfsNativeConfig`]; dropped properties are
+/// logged as warnings.
+pub(crate) fn hdfs_native_config_parse(m: HashMap<String, String>) -> 
Result<HdfsNativeConfig> {
+    hdfs_native_config_parse_with(m, |warning| tracing::warn!("{warning}"))
+}
+
+/// The parser proper, with the warning sink injected so tests can see what
+/// was dropped without a tracing subscriber.
+fn hdfs_native_config_parse_with(
+    mut m: HashMap<String, String>,
+    mut warn: impl FnMut(String),
+) -> Result<HdfsNativeConfig> {
+    let mut cfg = HdfsNativeConfig::default();
+
+    // Entries are trimmed one by one: opendal splits the list on `,` as is,
+    // so a space after a comma would break failover to that NameNode. An
+    // empty result is dropped because `Operator::from_config` bypasses the
+    // builder's empty-string guard and `Some("")` would shadow the
+    // path-authority fallback below.
+    if let Some(name_node) = m.remove(HDFS_NAME_NODE) {
+        let entries = hdfs_native_name_node_list(HDFS_NAME_NODE, &name_node)?;
+        if !entries.is_empty() {
+            cfg.name_node = Some(entries.join(","));
+        }
+    }
+
+    // `hdfs.name-node.<nameservice>` is sugar for Hadoop's own declaration of
+    // an HA nameservice, which the resolver reads back from the options.
+    let nameservice_prefix = format!("{HDFS_NAME_NODE}.");
+    let declared_keys: Vec<String> = m
+        .keys()
+        .filter(|key| key.starts_with(&nameservice_prefix))
+        .cloned()
+        .collect();
+    let mut declared = Vec::new();
+    for key in declared_keys {
+        let value = m.remove(&key).unwrap_or_default();
+        let nameservice = key[nameservice_prefix.len()..].to_string();
+        if nameservice.is_empty() || 
nameservice.chars().any(char::is_whitespace) {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!(
+                    "Invalid property `{key}`: a nameservice name must follow 
`{nameservice_prefix}`"
+                ),
+            ));
+        }
+        let entries = hdfs_native_name_node_list(&key, &value)?;
+        if entries.is_empty() {
+            return Err(Error::new(
+                ErrorKind::DataInvalid,
+                format!("Invalid property `{key}`: no NameNodes"),
+            ));
+        }
+        declared.push((nameservice, entries));
+    }
+
+    // A config carried over from PyIceberg would otherwise change identity
+    // silently; the client reads `HADOOP_USER_NAME` and the Kerberos cache.
+    for key in HDFS_UNSUPPORTED_KEYS {
+        if m.remove(key).is_some() {
+            warn(format!(
+                "`{key}` is not supported by the hdfs-native backend and is 
ignored"
+            ));
+        }
+    }
+    if m.contains_key(HDFS_HADOOP_CONF_PREFIX) {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!(
+                "Invalid property `{HDFS_HADOOP_CONF_PREFIX}`: a Hadoop key 
must follow the prefix"
+            ),
+        ));
+    }
+
+    let host = m
+        .remove(HDFS_HOST)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty());
+    let port = m
+        .remove(HDFS_PORT)
+        .map(|s| s.trim().to_string())
+        .filter(|s| !s.is_empty())
+        .map(|port| {
+            port.parse::<u16>().map_err(|e| {
+                Error::new(
+                    ErrorKind::DataInvalid,
+                    format!("Invalid `{HDFS_PORT}`: {port}: {e}"),
+                )
+            })
+        })
+        .transpose()?;
+
+    let mut options: HashMap<String, String> = m
+        .into_iter()
+        .filter_map(|(key, value)| {
+            key.strip_prefix(HDFS_HADOOP_CONF_PREFIX)
+                .map(|stripped| (stripped.to_string(), value))
+        })
+        .collect();
+    // Explicit `hadoop.` keys win over the sugar.
+    for (nameservice, entries) in declared {
+        let ids: Vec<String> = (0..entries.len()).map(|i| 
format!("nn{i}")).collect();
+        options
+            .entry(format!("{HA_NAMENODES_PREFIX}.{nameservice}"))
+            .or_insert_with(|| ids.join(","));
+        for (id, entry) in ids.iter().zip(&entries) {
+            options
+                .entry(format!(
+                    "{HA_NAMENODE_RPC_ADDRESS_PREFIX}.{nameservice}.{id}"
+                ))
+                .or_insert_with(|| 
entry.trim_start_matches("hdfs://").to_string());
+        }
+    }
+    // PyIceberg's `hdfs.host`/`hdfs.port` name the filesystem for
+    // authority-less paths, which is what Hadoop's `fs.defaultFS` means; an
+    // explicit `hadoop.fs.defaultFS` wins.
+    match host {
+        Some(host) => {
+            // An IPv6 literal needs brackets in a URI authority.
+            let bracketed = if host.contains(':') && !host.starts_with('[') {
+                format!("[{host}]")
+            } else {
+                host.clone()
+            };
+            let port = port.unwrap_or(HDFS_DEFAULT_PORT);
+            // Validated like every other NameNode spelling, so a host that
+            // carries a scheme or port fails here and not at the first I/O.
+            let default_fs = 
hdfs_native_name_node(&format!("hdfs://{bracketed}:{port}"))
+                .ok_or_else(|| {
+                    Error::new(
+                        ErrorKind::DataInvalid,
+                        format!(
+                            "Invalid `{HDFS_HOST}`/`{HDFS_PORT}`: 
{host}:{port}, expected a host name or IP and a port"
+                        ),
+                    )
+                })?;
+            options
+                .entry(FS_DEFAULT_FS.to_string())
+                .or_insert(default_fs);
+        }
+        None if port.is_some() => {
+            warn(format!(
+                "`{HDFS_PORT}` has no effect without `{HDFS_HOST}` and is 
ignored"
+            ));
+        }
+        None => {}
+    }
+    if !options.is_empty() {
+        cfg.options = Some(options);
+    }
+
+    Ok(cfg)
+}
+
+/// Parse an HDFS path into `Some("hdfs://<authority>")` (`None` when
+/// authority-less) and the relative path (no leading `/`, opendal style).
+pub(crate) fn hdfs_native_parse_path(path: &str) -> Result<(Option<String>, 
&str)> {
+    let url = Url::parse(path).map_err(|e| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}: {e}"),
+        )
+    })?;
+    // Non-special schemes parse even without `//` (e.g. `hdfs:x` is a valid
+    // non-hierarchical URL), so require the literal prefix before slicing.
+    let (Some(after_scheme), "hdfs") = (path.strip_prefix("hdfs://"), 
url.scheme()) else {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, expected scheme `hdfs://`"),
+        ));
+    };
+    // Userinfo has nowhere to go and port 0 cannot be dialed; silently
+    // dropping the one or reading the other as a logical name would mislead.
+    if !url.username().is_empty() || url.password().is_some() {
+        // Not echoing the path: it may carry a password.
+        let host = url.host_str().unwrap_or_default();
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path for host `{host}`: userinfo is not 
supported"),
+        ));
+    }
+    if url.port() == Some(0) {
+        return Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, port 0 cannot be dialed"),
+        ));
+    }
+
+    let name_node = url.host_str().filter(|h| !h.is_empty()).map(|host| {
+        url.port()
+            .map(|port| format!("hdfs://{host}:{port}"))
+            .unwrap_or_else(|| format!("hdfs://{host}"))
+    });
+
+    // `url.path()` borrows from `url` and can't be returned with the input's
+    // lifetime. Slice the path component out of the original input instead;
+    // it starts after the first `/` following the `hdfs://` prefix. Opendal
+    // paths must not start with `/` (`Deleter::delete` rejects them).
+    let rel = match after_scheme.find('/') {
+        Some(i) => after_scheme[i..].trim_start_matches('/'),
+        None => "",
+    };
+
+    Ok((name_node, rel))
+}
+
+/// Resolves the effective NameNode for a path, plus the relative path. As in
+/// Hadoop, an authority with a port is used as is; a logical nameservice
+/// authority (no port) resolves through its declaration, and an
+/// authority-less path through `hdfs.name-node`, else `fs.defaultFS`. The
+/// operator cache, `delete_stream` batching and `relativize_path` all go
+/// through this, so they cannot drift apart.
+pub(crate) fn hdfs_native_effective_name_node<'a>(
+    config: &HdfsNativeConfig,
+    path: &'a str,
+) -> Result<(String, &'a str)> {
+    let (authority, relative_path) = hdfs_native_parse_path(path)?;
+    let invalid = |reason: String| {
+        Error::new(
+            ErrorKind::DataInvalid,
+            format!("Invalid hdfs path: {path}, {reason}"),
+        )
+    };
+    let name_node = match authority {
+        Some(authority) if hdfs_native_name_node(&authority).is_some() => 
authority,
+        Some(logical) => {

Review Comment:
   Followup, not blocking — this arm treats any portless authority as a logical 
nameservice that must be declared through a property, and errors otherwise. The 
common HA deployment writes Iceberg locations as 
`hdfs://nameservice1/warehouse/...`, where `nameservice1` is defined only in 
the cluster's hdfs-site.xml under `HADOOP_CONF_DIR` — those tables fail here 
unless the user re-declares every nameservice as `hdfs.name-node.<ns>`.
   
   hdfs-native already reads `HADOOP_CONF_DIR` and resolves nameservices 
itself, so I'd rather not pre-empt it with an error: passing 
`hdfs://<authority>` straight through as the `name_node` and letting the client 
resolve it (failing with its own error if it can't) would cover this, and 
`hdfs://host/path` with no port, which Hadoop reads as `host:8020`.
   
   Fine to take as a follow-on. If the strict behavior is deliberate for now, 
the "as in Hadoop" line in the `HDFS_NAME_NODE` doc should go — it diverges 
here — and the limitation is worth a README note.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to