JingsongLi commented on code in PR #1028: URL: https://github.com/apache/paimon-rust/pull/1028#discussion_r4177532097
########## crates/paimon/src/io/storage_oss_cpp.rs: ########## @@ -0,0 +1,780 @@ +// 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. + +use std::collections::HashMap; +use std::ffi::{c_char, c_void, CStr, CString}; +use std::fmt::{Debug, Formatter}; +use std::sync::atomic::{AtomicU32, Ordering}; +use std::sync::Arc; + +use libloading::Library; +use opendal::raw::*; +use opendal::{Buffer, Builder, BytesRange, Capability, EntryMode, Error, ErrorKind, Metadata}; +use opendal::{OperationContext, Operator, Result}; +use tokio::sync::{OnceCell, Semaphore}; + +const PREFIX: &str = "fs.oss.cpp."; +const DEFAULT_CONCURRENCY: usize = 8; + +#[derive(Clone)] +pub struct OssCppStorageConfig { + library: String, + strings: Vec<CString>, + concurrency: usize, + gate: Arc<Semaphore>, + connect_timeout: i64, + request_timeout: i64, + retry_attempts: i64, + path_style: u8, +} + +impl Default for OssCppStorageConfig { + fn default() -> Self { + Self { + library: String::new(), + strings: Vec::new(), + concurrency: 0, + gate: Arc::new(Semaphore::new(0)), + connect_timeout: 0, + request_timeout: 0, + retry_attempts: 0, + path_style: 0, + } + } +} + +impl Debug for OssCppStorageConfig { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.debug_struct("OssCppStorageConfig") + .field("library", &self.library) + .field("concurrency", &self.concurrency) + .finish_non_exhaustive() + } +} + +pub(crate) fn oss_cpp_config_parse( + props: HashMap<String, String>, +) -> crate::Result<OssCppStorageConfig> { + let parse = || -> Result<OssCppStorageConfig> { + let required = |key: &str| -> Result<String> { + props + .get(key) + .filter(|s| !s.is_empty()) + .cloned() + .ok_or_else(|| Error::new(ErrorKind::ConfigInvalid, format!("Missing {key}"))) + }; + let number = |name: &str, default: i64| -> Result<i64> { + let key = format!("{PREFIX}{name}"); + match props.get(&key) { + None => Ok(default), + Some(value) => value.parse::<i64>().ok().filter(|n| *n > 0).ok_or_else(|| { + Error::new(ErrorKind::ConfigInvalid, format!("{key} must be positive")) + }), + } + }; + let endpoint = required("fs.oss.endpoint")?; + let endpoint = if endpoint.contains("://") { + endpoint + } else { + format!("https://{endpoint}") + }; + let url = url::Url::parse(&endpoint) + .map_err(|_| Error::new(ErrorKind::ConfigInvalid, "Invalid OSS endpoint"))?; + if !matches!(url.scheme(), "http" | "https") + || url.host_str().is_none() + || !url.username().is_empty() + || url.password().is_some() + || url.query().is_some() + || url.fragment().is_some() + || url.path() != "/" + { + return Err(Error::new( + ErrorKind::ConfigInvalid, + "OSS endpoint must be an HTTP(S) origin", + )); + } + let strings = [ + endpoint, + required("fs.oss.region")?, + required("fs.oss.accessKeyId")?, + required("fs.oss.accessKeySecret")?, + props + .get("fs.oss.securityToken") + .cloned() + .unwrap_or_default(), + // The native SDK appends this to its own User-Agent. + format!("{} oss-cpp", super::user_agent::oss_user_agent(&props)), + ] + .into_iter() + .map(|s| cstring(&s)) + .collect::<Result<Vec<_>>>()?; + let concurrency = number("max.concurrent.requests", DEFAULT_CONCURRENCY as i64)? as usize; + if concurrency > Semaphore::MAX_PERMITS { + return Err(Error::new( + ErrorKind::ConfigInvalid, + "OSS C++ concurrency exceeds semaphore limit", + )); + } + let path_style = match props.get("fs.oss.cpp.path-style").map(String::as_str) { + None | Some("false") => 0, + Some("true") => 1, + _ => { + return Err(Error::new( + ErrorKind::ConfigInvalid, + "Invalid OSS C++ path-style", + )) + } + }; + Ok(OssCppStorageConfig { + library: required("fs.oss.cpp.library.path")?, + strings, + concurrency, + gate: Arc::new(Semaphore::new(concurrency)), + connect_timeout: number("connect.timeout-ms", 10_000)?, + request_timeout: number("request.timeout-ms", 30_000)?, + retry_attempts: number("retry.max-attempts", 3)?, + path_style, + }) + }; + parse().map_err(|e| crate::Error::ConfigInvalid { + message: e.to_string(), + }) +} + +pub(crate) fn oss_cpp_config_build( + config: &OssCppStorageConfig, + bucket: &str, +) -> crate::Result<Operator> { + Operator::new(CppBuilder { + config: config.clone(), + bucket: bucket.to_string(), + }) + .map_err(|e| crate::Error::IoUnexpected { + message: "Cannot build OSS C++ operator".to_string(), + source: Box::new(e), + }) +} + +#[derive(Default)] +struct CppBuilder { + config: OssCppStorageConfig, + bucket: String, +} + +impl Builder for CppBuilder { + type Config = (); + fn build(self) -> Result<impl Service> { + if self.config.strings.len() != 6 || self.config.concurrency == 0 { + return Err(Error::new( + ErrorKind::ConfigInvalid, + "Missing OSS C++ configuration", + )); + } + Ok(CppService { + info: ServiceInfo::new("oss-cpp", "/", &self.bucket), + client: Arc::new(LazyClient { + gate: self.config.gate.clone(), + config: self.config, + bucket: self.bucket, + client: OnceCell::new(), + }), + }) + } +} + +struct CppService { + info: ServiceInfo, + client: Arc<LazyClient>, +} + +impl Debug for CppService { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + f.debug_struct("OssCppService") + .field("info", &self.info) + .finish() + } +} + +impl Service for CppService { + type Reader = CppReader; + type Writer = (); + type Lister = CppLister; + type Deleter = (); + type Copier = (); + + fn info(&self) -> ServiceInfo { + self.info.clone() + } + fn capability(&self) -> Capability { + Capability { + stat: true, + read: true, + read_with_suffix: true, + list: true, + list_with_recursive: true, + shared: true, + ..Default::default() + } + } + async fn stat(&self, _: &OperationContext, path: &str, _: OpStat) -> Result<RpStat> { + let path = path.to_string(); + let metadata = self.client.call(move |c| c.stat(&path)).await?; + Ok(RpStat::new(metadata)) + } + fn read(&self, _: &OperationContext, path: &str, _: OpRead) -> Result<Self::Reader> { + Ok(CppReader { + client: self.client.clone(), + path: path.to_string(), + }) + } + fn list(&self, _: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> { + Ok(CppLister { + client: self.client.clone(), + path: path.to_string(), + recursive: args.recursive(), + max_keys: args.limit().unwrap_or(1000).clamp(1, 1000) as u32, + token: String::new(), + done: false, + entries: Vec::new().into_iter(), + }) + } + async fn create_dir( + &self, + _: &OperationContext, + _: &str, + _: OpCreateDir, + ) -> Result<RpCreateDir> { + unsupported() + } + fn write(&self, _: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> { + unsupported() + } + fn delete(&self, _: &OperationContext) -> Result<Self::Deleter> { + unsupported() + } + fn copy( + &self, + _: &OperationContext, + _: &str, + _: &str, + _: OpCopy, + _: OpCopier, + ) -> Result<Self::Copier> { + unsupported() + } + async fn rename( + &self, + _: &OperationContext, + _: &str, + _: &str, + _: OpRename, + ) -> Result<RpRename> { + unsupported() + } + async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> { + unsupported() + } +} + +fn unsupported<T>() -> Result<T> { + Err(Error::new( + ErrorKind::Unsupported, + "OSS C++ backend is read-only", + )) +} + +struct CppReader { + client: Arc<LazyClient>, + path: String, +} + +impl oio::Read for CppReader { + async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> { + let (_, buffer) = self.read(range).await?; + Ok((RpRead::default(), Box::new(buffer))) + } + async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> { + let path = self.path.clone(); + let data = self.client.call(move |c| c.read(&path, range)).await?; + Ok((RpRead::default(), Buffer::from(data))) + } +} + +struct CppLister { + client: Arc<LazyClient>, + path: String, + recursive: bool, + max_keys: u32, + token: String, + done: bool, + entries: std::vec::IntoIter<oio::Entry>, +} + +impl oio::List for CppLister { + async fn next(&mut self) -> Result<Option<oio::Entry>> { + loop { + if let Some(entry) = self.entries.next() { + return Ok(Some(entry)); + } + if self.done { + return Ok(None); + } + let path = self.path.clone(); + let token = self.token.clone(); + let recursive = self.recursive; + let max_keys = self.max_keys; + let (entries, next) = self + .client + .call(move |c| c.list(&path, &token, recursive, max_keys)) + .await?; + self.done = next.is_empty(); + self.token = next; + self.entries = entries.into_iter(); + } + } +} + +struct LazyClient { + config: OssCppStorageConfig, + bucket: String, + client: OnceCell<Arc<Client>>, + gate: Arc<Semaphore>, +} + +// Native clients/connection pools cannot be inherited across fork. Check before +// touching OnceCell or Semaphore: either may have been locked by a parent thread. +fn ensure_process() -> Result<()> { + static PID: AtomicU32 = AtomicU32::new(0); + let current = std::process::id(); + match PID.compare_exchange(0, current, Ordering::AcqRel, Ordering::Acquire) { + Ok(_) => Ok(()), + Err(pid) if pid == current => Ok(()), + Err(_) => Err(Error::new( + ErrorKind::Unsupported, + "OSS C++ SDK cannot be used after fork; use spawn workers", Review Comment: [P2] Connect this rejection to the public fork-safety error category. After a parent has initialized this backend, an inherited client is correctly rejected by ensure_process, but this new Unsupported message is not recognized by Error::from_opendal_with_context/is_process_fork_unsupported: both only recognize JINDO_FORK_ERROR. Consequently Python to_py_err and DataFusion error conversion produce ordinary ValueError rather than the documented ForkSafetyError ("native state inherited across process fork"). Callers catching that exception to recreate spawned workers miss this rejection. A public Error::from(opendal::Error::new(Unsupported, this exact message)) classification probe fails on head and the main-integrated tree. Add the OSS marker to the shared recognition, or use a backend-neutral fork marker, and cover propagation through FileIO and the binding. -- 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]
