Merge branch 'main' of https://github.com/rustfs/s3-rustfs into feature/ilm

# Conflicts:
#	rustfs/src/storage/ecfs.rs
This commit is contained in:
likewu
2025-06-23 16:44:37 +08:00
10 changed files with 146 additions and 139 deletions
+8 -8
View File
@@ -12,9 +12,9 @@ mod service;
mod storage;
use crate::auth::IAMAuth;
use crate::console::{init_console_cfg, CONSOLE_CONFIG};
use crate::console::{CONSOLE_CONFIG, init_console_cfg};
// Ensure the correct path for parse_license is imported
use crate::server::{wait_for_shutdown, ServiceState, ServiceStateManager, ShutdownSignal, SHUTDOWN_TIMEOUT};
use crate::server::{SHUTDOWN_TIMEOUT, ServiceState, ServiceStateManager, ShutdownSignal, wait_for_shutdown};
use bytes::Bytes;
use chrono::Datelike;
use clap::Parser;
@@ -22,6 +22,7 @@ use common::{
// error::{Error, Result},
globals::set_global_addr,
};
use ecstore::StorageAPI;
use ecstore::bucket::metadata_sys::init_bucket_metadata_sys;
use ecstore::cmd::bucket_replication::init_bucket_replication_pool;
use ecstore::config as ecconfig;
@@ -29,12 +30,11 @@ use ecstore::config::GLOBAL_ConfigSys;
use ecstore::heal::background_heal_ops::init_auto_heal;
use ecstore::rpc::make_server;
use ecstore::store_api::BucketOptions;
use ecstore::StorageAPI;
use ecstore::{
endpoints::EndpointServerPools,
heal::data_scanner::init_data_scanner,
set_global_endpoints,
store::{init_local_disks, ECStore},
store::{ECStore, init_local_disks},
update_erasure_type,
};
use ecstore::{global::set_global_rustfs_port, notification_sys::new_global_notification_sys};
@@ -49,7 +49,7 @@ use iam::init_iam_sys;
use license::init_license;
use protos::proto_gen::node_service::node_service_server::NodeServiceServer;
use rustfs_config::{DEFAULT_ACCESS_KEY, DEFAULT_SECRET_KEY, RUSTFS_TLS_CERT, RUSTFS_TLS_KEY};
use rustfs_obs::{init_obs, set_global_guard, SystemObserver};
use rustfs_obs::{SystemObserver, init_obs, set_global_guard};
use rustfs_utils::net::parse_and_resolve_address;
use rustls::ServerConfig;
use s3s::{host::MultiDomain, service::S3ServiceBuilder};
@@ -60,12 +60,12 @@ use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use tokio::net::TcpListener;
use tokio::signal::unix::{signal, SignalKind};
use tokio::signal::unix::{SignalKind, signal};
use tokio_rustls::TlsAcceptor;
use tonic::{metadata::MetadataValue, Request, Status};
use tonic::{Request, Status, metadata::MetadataValue};
use tower_http::cors::CorsLayer;
use tower_http::trace::TraceLayer;
use tracing::{debug, error, info, instrument, warn, Span};
use tracing::{Span, debug, error, info, instrument, warn};
const MI_B: usize = 1024 * 1024;
+61 -30
View File
@@ -74,6 +74,7 @@ use query::instance::make_rustfsms;
use rustfs_filemeta::headers::RESERVED_METADATA_PREFIX_LOWER;
use rustfs_filemeta::headers::{AMZ_DECODED_CONTENT_LENGTH, AMZ_OBJECT_TAGGING};
use rustfs_rio::CompressReader;
use rustfs_rio::EtagReader;
use rustfs_rio::HashReader;
use rustfs_rio::Reader;
use rustfs_rio::WarpReader;
@@ -341,7 +342,7 @@ impl S3 for FS {
..Default::default()
};
let dst_opts = copy_dst_opts(&bucket, &key, version_id, &req.headers, None)
let dst_opts = copy_dst_opts(&bucket, &key, version_id, &req.headers, HashMap::new())
.await
.map_err(ApiError::from)?;
@@ -368,13 +369,50 @@ impl S3 for FS {
src_info.metadata_only = true;
}
let reader = Box::new(WarpReader::new(gr.stream));
let hrd = HashReader::new(reader, gr.object_info.size, gr.object_info.size, None, false).map_err(ApiError::from)?;
let mut reader: Box<dyn Reader> = Box::new(WarpReader::new(gr.stream));
let actual_size = src_info.get_actual_size().map_err(ApiError::from)?;
let mut length = actual_size;
let mut compress_metadata = HashMap::new();
if is_compressible(&req.headers, &key) && actual_size > MIN_COMPRESSIBLE_SIZE as i64 {
compress_metadata.insert(
format!("{}compression", RESERVED_METADATA_PREFIX_LOWER),
CompressionAlgorithm::default().to_string(),
);
compress_metadata.insert(format!("{}actual-size", RESERVED_METADATA_PREFIX_LOWER,), actual_size.to_string());
let hrd = EtagReader::new(reader, None);
// let hrd = HashReader::new(reader, length, actual_size, None, false).map_err(ApiError::from)?;
reader = Box::new(CompressReader::new(hrd, CompressionAlgorithm::default()));
length = -1;
} else {
src_info
.user_defined
.remove(&format!("{}compression", RESERVED_METADATA_PREFIX_LOWER));
src_info
.user_defined
.remove(&format!("{}actual-size", RESERVED_METADATA_PREFIX_LOWER));
src_info
.user_defined
.remove(&format!("{}compression-size", RESERVED_METADATA_PREFIX_LOWER));
}
let hrd = HashReader::new(reader, length, actual_size, None, false).map_err(ApiError::from)?;
src_info.put_object_reader = Some(PutObjReader::new(hrd));
// check quota
// TODO: src metadada
for (k, v) in compress_metadata {
src_info.user_defined.insert(k, v);
}
// TODO: src tags
let oi = store
@@ -429,7 +467,7 @@ impl S3 for FS {
let metadata = extract_metadata(&req.headers);
let opts: ObjectOptions = del_opts(&bucket, &key, version_id, &req.headers, Some(metadata))
let opts: ObjectOptions = del_opts(&bucket, &key, version_id, &req.headers, metadata)
.await
.map_err(ApiError::from)?;
@@ -499,7 +537,7 @@ impl S3 for FS {
let metadata = extract_metadata(&req.headers);
let opts: ObjectOptions = del_opts(&bucket, "", None, &req.headers, Some(metadata))
let opts: ObjectOptions = del_opts(&bucket, "", None, &req.headers, metadata)
.await
.map_err(ApiError::from)?;
@@ -736,7 +774,7 @@ impl S3 for FS {
content_type,
last_modified,
e_tag: info.etag,
metadata,
metadata: Some(metadata),
version_id: info.version_id.map(|v| v.to_string()),
// metadata: object_metadata,
..Default::default()
@@ -856,7 +894,7 @@ impl S3 for FS {
let mut obj = Object {
key: Some(v.name.to_owned()),
last_modified: v.mod_time.map(Timestamp::from),
size: Some(v.size),
size: Some(v.get_actual_size().unwrap_or_default()),
e_tag: v.etag.clone(),
..Default::default()
};
@@ -1055,7 +1093,7 @@ impl S3 for FS {
let mt = metadata.clone();
let mt2 = metadata.clone();
let mut opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, Some(mt))
let mut opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, mt)
.await
.map_err(ApiError::from)?;
@@ -1065,14 +1103,12 @@ impl S3 for FS {
let dsc = must_replicate(&bucket, &key, &repoptions).await;
// warn!("dsc {}", &dsc.replicate_any().clone());
if dsc.replicate_any() {
if let Some(metadata) = opts.user_defined.as_mut() {
let k = format!("{}{}", RESERVED_METADATA_PREFIX_LOWER, "replication-timestamp");
let now: DateTime<Utc> = Utc::now();
let formatted_time = now.to_rfc3339();
metadata.insert(k, formatted_time);
opts.user_defined.insert(k, formatted_time);
let k = format!("{}{}", RESERVED_METADATA_PREFIX_LOWER, "replication-status");
metadata.insert(k, dsc.pending_status());
}
opts.user_defined.insert(k, dsc.pending_status());
}
let obj_info = store
@@ -1133,7 +1169,7 @@ impl S3 for FS {
);
}
let opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, Some(metadata))
let opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, metadata)
.await
.map_err(ApiError::from)?;
@@ -2314,11 +2350,10 @@ impl S3 for FS {
s3_error!(InternalError, "{}", e.to_string())
})?;
let legal_hold = if let Some(ud) = object_info.user_defined {
ud.get("x-amz-object-lock-legal-hold").map(|v| v.as_str().to_string())
} else {
None
};
let legal_hold = object_info
.user_defined
.get("x-amz-object-lock-legal-hold")
.map(|v| v.as_str().to_string());
let status = if let Some(v) = legal_hold {
v
@@ -2415,20 +2450,16 @@ impl S3 for FS {
s3_error!(InternalError, "{}", e.to_string())
})?;
let mode = if let Some(ref ud) = object_info.user_defined {
ud.get("x-amz-object-lock-mode")
.map(|v| ObjectLockRetentionMode::from(v.as_str().to_string()))
} else {
None
};
let mode = object_info
.user_defined
.get("x-amz-object-lock-mode")
.map(|v| ObjectLockRetentionMode::from(v.as_str().to_string()));
let retain_until_date = if let Some(ref ud) = object_info.user_defined {
ud.get("x-amz-object-lock-retain-until-date")
let retain_until_date = object_info
.user_defined
.get("x-amz-object-lock-retain-until-date")
.and_then(|v| OffsetDateTime::parse(v.as_str(), &Rfc3339).ok())
.map(Timestamp::from)
} else {
None
};
.map(Timestamp::from);
Ok(S3Response::new(GetObjectRetentionOutput {
retention: Some(ObjectLockRetention { mode, retain_until_date }),
+29 -32
View File
@@ -14,7 +14,7 @@ pub async fn del_opts(
object: &str,
vid: Option<String>,
headers: &HeaderMap<HeaderValue>,
metadata: Option<HashMap<String, String>>,
metadata: HashMap<String, String>,
) -> Result<ObjectOptions> {
let versioned = BucketVersioningSys::prefix_enabled(bucket, object).await;
let version_suspended = BucketVersioningSys::suspended(bucket).await;
@@ -33,7 +33,7 @@ pub async fn del_opts(
}
}
let mut opts = put_opts_from_headers(headers, metadata)
let mut opts = put_opts_from_headers(headers, metadata.clone())
.map_err(|err| StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), err.to_string()))?;
opts.version_id = {
@@ -72,7 +72,7 @@ pub async fn get_opts(
}
}
let mut opts = get_default_opts(headers, None, false)
let mut opts = get_default_opts(headers, HashMap::new(), false)
.map_err(|err| StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), err.to_string()))?;
opts.version_id = {
@@ -97,7 +97,7 @@ pub async fn put_opts(
object: &str,
vid: Option<String>,
headers: &HeaderMap<HeaderValue>,
metadata: Option<HashMap<String, String>>,
metadata: HashMap<String, String>,
) -> Result<ObjectOptions> {
let versioned = BucketVersioningSys::prefix_enabled(bucket, object).await;
let version_suspended = BucketVersioningSys::prefix_suspended(bucket, object).await;
@@ -136,30 +136,27 @@ pub async fn copy_dst_opts(
object: &str,
vid: Option<String>,
headers: &HeaderMap<HeaderValue>,
metadata: Option<HashMap<String, String>>,
metadata: HashMap<String, String>,
) -> Result<ObjectOptions> {
put_opts(bucket, object, vid, headers, metadata).await
}
pub fn copy_src_opts(_bucket: &str, _object: &str, headers: &HeaderMap<HeaderValue>) -> Result<ObjectOptions> {
get_default_opts(headers, None, false)
get_default_opts(headers, HashMap::new(), false)
}
pub fn put_opts_from_headers(
headers: &HeaderMap<HeaderValue>,
metadata: Option<HashMap<String, String>>,
) -> Result<ObjectOptions> {
pub fn put_opts_from_headers(headers: &HeaderMap<HeaderValue>, metadata: HashMap<String, String>) -> Result<ObjectOptions> {
get_default_opts(headers, metadata, false)
}
/// Creates default options for getting an object from a bucket.
pub fn get_default_opts(
_headers: &HeaderMap<HeaderValue>,
metadata: Option<HashMap<String, String>>,
metadata: HashMap<String, String>,
_copy_source: bool,
) -> Result<ObjectOptions> {
Ok(ObjectOptions {
user_defined: metadata.clone(),
user_defined: metadata,
..Default::default()
})
}
@@ -244,13 +241,13 @@ mod tests {
#[tokio::test]
async fn test_del_opts_basic() {
let headers = create_test_headers();
let metadata = Some(create_test_metadata());
let metadata = create_test_metadata();
let result = del_opts("test-bucket", "test-object", None, &headers, metadata).await;
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_some());
assert!(!opts.user_defined.is_empty());
assert_eq!(opts.version_id, None);
}
@@ -258,7 +255,7 @@ mod tests {
async fn test_del_opts_with_directory_object() {
let headers = create_test_headers();
let result = del_opts("test-bucket", "test-dir/", None, &headers, None).await;
let result = del_opts("test-bucket", "test-dir/", None, &headers, HashMap::new()).await;
assert!(result.is_ok());
let opts = result.unwrap();
@@ -270,7 +267,7 @@ mod tests {
let headers = create_test_headers();
let valid_uuid = Uuid::new_v4().to_string();
let result = del_opts("test-bucket", "test-object", Some(valid_uuid.clone()), &headers, None).await;
let result = del_opts("test-bucket", "test-object", Some(valid_uuid.clone()), &headers, HashMap::new()).await;
// This test may fail if versioning is not enabled for the bucket
// In a real test environment, you would mock BucketVersioningSys
@@ -289,7 +286,7 @@ mod tests {
let headers = create_test_headers();
let invalid_uuid = "invalid-uuid".to_string();
let result = del_opts("test-bucket", "test-object", Some(invalid_uuid), &headers, None).await;
let result = del_opts("test-bucket", "test-object", Some(invalid_uuid), &headers, HashMap::new()).await;
assert!(result.is_err());
if let Err(err) = result {
@@ -361,13 +358,13 @@ mod tests {
#[tokio::test]
async fn test_put_opts_basic() {
let headers = create_test_headers();
let metadata = Some(create_test_metadata());
let metadata = create_test_metadata();
let result = put_opts("test-bucket", "test-object", None, &headers, metadata).await;
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_some());
assert!(!opts.user_defined.is_empty());
assert_eq!(opts.version_id, None);
}
@@ -375,7 +372,7 @@ mod tests {
async fn test_put_opts_with_directory_object() {
let headers = create_test_headers();
let result = put_opts("test-bucket", "test-dir/", None, &headers, None).await;
let result = put_opts("test-bucket", "test-dir/", None, &headers, HashMap::new()).await;
assert!(result.is_ok());
let opts = result.unwrap();
@@ -387,7 +384,7 @@ mod tests {
let headers = create_test_headers();
let invalid_uuid = "invalid-uuid".to_string();
let result = put_opts("test-bucket", "test-object", Some(invalid_uuid), &headers, None).await;
let result = put_opts("test-bucket", "test-object", Some(invalid_uuid), &headers, HashMap::new()).await;
assert!(result.is_err());
if let Err(err) = result {
@@ -405,13 +402,13 @@ mod tests {
#[tokio::test]
async fn test_copy_dst_opts() {
let headers = create_test_headers();
let metadata = Some(create_test_metadata());
let metadata = create_test_metadata();
let result = copy_dst_opts("test-bucket", "test-object", None, &headers, metadata).await;
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_some());
assert!(!opts.user_defined.is_empty());
}
#[test]
@@ -422,20 +419,20 @@ mod tests {
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_none());
assert!(opts.user_defined.is_empty());
}
#[test]
fn test_put_opts_from_headers() {
let headers = create_test_headers();
let metadata = Some(create_test_metadata());
let metadata = create_test_metadata();
let result = put_opts_from_headers(&headers, metadata);
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_some());
let user_defined = opts.user_defined.unwrap();
assert!(!opts.user_defined.is_empty());
let user_defined = opts.user_defined;
assert_eq!(user_defined.get("key1"), Some(&"value1".to_string()));
assert_eq!(user_defined.get("key2"), Some(&"value2".to_string()));
}
@@ -443,14 +440,14 @@ mod tests {
#[test]
fn test_get_default_opts_with_metadata() {
let headers = create_test_headers();
let metadata = Some(create_test_metadata());
let metadata = create_test_metadata();
let result = get_default_opts(&headers, metadata, false);
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_some());
let user_defined = opts.user_defined.unwrap();
assert!(!opts.user_defined.is_empty());
let user_defined = opts.user_defined;
assert_eq!(user_defined.get("key1"), Some(&"value1".to_string()));
assert_eq!(user_defined.get("key2"), Some(&"value2".to_string()));
}
@@ -459,11 +456,11 @@ mod tests {
fn test_get_default_opts_without_metadata() {
let headers = create_test_headers();
let result = get_default_opts(&headers, None, false);
let result = get_default_opts(&headers, HashMap::new(), false);
assert!(result.is_ok());
let opts = result.unwrap();
assert!(opts.user_defined.is_none());
assert!(opts.user_defined.is_empty());
}
#[test]