mixermt commented on code in PR #3111: URL: https://github.com/apache/iceberg-rust/pull/3111#discussion_r4216882220
########## 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: Checked this end to end. hdfs-native on its own does resolve a `HADOOP_CONF_DIR` nameservice, but opendal never hands it the authority: since apache/opendal#7248 (0.56) it builds the client as `hdfs://nameservice` and writes each `name_node` entry as `dfs.namenode.rpc-address.nameservice.nnN`. So `hdfs://ns1` passed through, with `ns1` defined in hdfs-site.xml, becomes an rpc-address without a port, and so does a portless host: neither is ever dialed (a listener on 8020 saw nothing), and hdfs-native retries `invalid socket address` through its default 15 failovers with backoff before giving up. hdfs-native has no default port either; given `hdfs://127.0.0.1` directly it returns `No NameNode hosts found`. So pass-through would turn today's immediate error into a slow, opaque one. Kept strict for now. The `HDFS_NAME_NODE` doc no longer says "as in Hadoop" and states the divergence, the error now reads "`ns1` has no port and is not a declared nameservice; add the port or set `hdfs.name-node.ns1`", and the README notes the limitation. The real fix is upstream, letting opendal take a logical nameservice as `name_node` and pass it as the client URL; once that lands, undeclared names can pass through as you suggest. ########## 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: Fixed, and the leak was wider than these two sites: the `hdfs.name-node`, `hdfs.name-node.<nameservice>`, `hdfs.host` and `fs.defaultFS` errors echoed their values too. Every `Invalid hdfs path` error now goes through one helper that masks everything between the scheme and the last `@` (`hdfs://***@nn:bad/x`), so the path keeps its context and a password with a raw `/` or `#` is covered as well; the config errors use the same masking. A test covers each of these errors and asserts none contains the user or the password. `OpenDalResolvingStorage`'s scheme lookup echoes a bad path the same way for every scheme, which predates this PR; worth a separate fix. ########## 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: Done: both sites recover the guard with `PoisonError::into_inner`, `get`/`insert` return plain values, and the error helper is gone. Recovery is safe since writes only ever replace whole entries. A test poisons the lock and checks the cache still builds and serves operators. ########## 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: Fixed the sampling. One correction on the premise: the first stat never reaches the 10s timeout. It fails in about 1.5s after up to 4 dials, and over 30 runs no dial arrived after the sample. A reused stale operator also panics inside hdfs-native before the assertion, as the original bug did. Still, the order was wrong in principle, so the test now drops the first runtime, opens a marker connection and waits for the listener to close it. The listener accepts in order, so every earlier dial is counted before `after_first` is sampled, with no sleeps. The cache is crate-private, so the identity check went into the runtime-shutdown unit test: it asserts the rebuilt operator is a different instance from the stale one (`Operator::into_parts` and `Arc::ptr_eq`). A mutation that refreshes the sentinel but keeps the old operator fails only that assertion. -- 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]
