diff --git a/Cargo.lock b/Cargo.lock index c6d279229..a9fac1098 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9503,6 +9503,7 @@ dependencies = [ "clap", "const-str", "datafusion", + "faster-hex", "flatbuffers", "flate2", "futures", @@ -9527,6 +9528,7 @@ dependencies = [ "metrics", "metrics-util", "mime_guess", + "moka", "opentelemetry", "opentelemetry_sdk", "p256 0.14.0", @@ -9619,6 +9621,7 @@ dependencies = [ "urlencoding", "uuid", "x509-parser", + "xxhash-rust", "zeroize", "zip", "zstd 0.14.0", diff --git a/crates/ecstore/src/bucket/metadata_sys.rs b/crates/ecstore/src/bucket/metadata_sys.rs index 5e30841ad..7b36a1647 100644 --- a/crates/ecstore/src/bucket/metadata_sys.rs +++ b/crates/ecstore/src/bucket/metadata_sys.rs @@ -1043,11 +1043,8 @@ pub async fn get_on_demand_migration_config(bucket: &str) -> Result Result, OffsetDateTime)>> { - let sys = bucket_metadata_sys_of(ctx)?; +pub async fn get_on_demand_migration_config_in(api: &ECStore, bucket: &str) -> Result, OffsetDateTime)>> { + let sys = bucket_metadata_sys_of(&api.ctx)?; let lock = sys.read().await; lock.get_on_demand_migration_config(bucket).await } diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 63d2b5fdd..764e48391 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -353,6 +353,11 @@ async fn resume_rebalance_after_init(store: Arc, rx: CancellationToken) } impl ECStore { + /// Shutdown token owned by this store instance. + pub fn background_cancel_token(&self) -> Option { + self.ctx.background_cancel_token() + } + /// Validate topology and process storage-class overrides before any disk is opened. pub fn validate_startup_storage_class(endpoint_pools: &EndpointServerPools) -> Result<()> { let drive_counts = startup_pool_drive_counts(endpoint_pools); diff --git a/crates/madmin/src/on_demand_migration.rs b/crates/madmin/src/on_demand_migration.rs index 0fce8af6c..e92899717 100644 --- a/crates/madmin/src/on_demand_migration.rs +++ b/crates/madmin/src/on_demand_migration.rs @@ -17,7 +17,7 @@ //! Wire types for `PUT`/`GET`/`DELETE /v3/on-demand-migration/{bucket}`, //! `GET .../status`, `POST .../backfill?op=start|cancel` and //! `GET .../backfill` (ODM-12), mirroring the server's config model -//! (`crates/ecstore/src/bucket/on_demand_migration/config.rs`) and handler +//! (`rustfs/src/on_demand_migration/config.rs`) and handler //! responses (`rustfs/src/admin/handlers/on_demand_migration.rs`). The SDK //! owns its own copies, madmin-go style; the fixtures under //! `fixtures/on_demand_migration/` are the contract both sides pin diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 03b6b378e..52a7e6598 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -276,6 +276,7 @@ rcgen = { workspace = true } # Async Runtime and Networking async-trait = { workspace = true } axum.workspace = true +faster-hex.workspace = true futures.workspace = true futures-lite.workspace = true futures-util.workspace = true @@ -291,6 +292,8 @@ tokio-rustls = { workspace = true, default-features = false, features = ["loggin aws-smithy-runtime-api = { workspace = true, features = ["http-1x"] } aws-smithy-types = { workspace = true } google-cloud-auth = { workspace = true, optional = true } +moka = { workspace = true, features = ["sync"] } +xxhash-rust = { workspace = true, features = ["xxh3"] } aws-sdk-s3 = { workspace = true, default-features = false, features = ["sigv4a", "default-https-client", "rt-tokio"] } tokio-stream.workspace = true tokio-util = { workspace = true, features = ["io", "compat", "time"] } diff --git a/rustfs/src/app/bucket_list_through.rs b/rustfs/src/app/bucket_list_through.rs index c385fdb9c..586f5b0ae 100644 --- a/rustfs/src/app/bucket_list_through.rs +++ b/rustfs/src/app/bucket_list_through.rs @@ -443,7 +443,9 @@ mod tests { use super::*; use crate::app::bucket_usecase::DefaultBucketUsecase; use crate::app::gating_test_env::{run_large_stack_test, shared_gating_ecstore}; - use crate::app::storage_api::bucket_usecase::s3::{ListObjectsV2Input, ListObjectsV2Output, S3Request, S3Response}; + use crate::app::storage_api::bucket_usecase::s3::{ + ListObjectsInput, ListObjectsV2Input, ListObjectsV2Output, S3Request, S3Response, XmlSerialize, XmlSerializer, + }; use crate::app::storage_api::test::StoragePutObjReader; use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions}; use crate::app::storage_api::test::contract::object::ObjectIO as _; @@ -451,7 +453,6 @@ mod tests { FilterConfig, MAX_LIST_NO_PROGRESS_PAGES, OnDemandMigrationConfig, PathStyle, PolicyConfig, Provider, SourceConfig, SourceCredentials, TlsConfig, }; - use s3s::dto::ListObjectsInput; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt}; @@ -864,7 +865,7 @@ mod tests { assert_eq!(output.next_marker.as_deref(), (index == 0).then_some(expected_key)); let mut xml = Vec::new(); - s3s::xml::Serialize::serialize(&output, &mut s3s::xml::Serializer::new(&mut xml)) + XmlSerialize::serialize(&output, &mut XmlSerializer::new(&mut xml)) .expect("serialize the real v1 response"); assert!(!xml.contains(&0), "XML 1.0 forbids NUL in NextMarker"); let mut reader = quick_xml::Reader::from_reader(xml.as_slice()); diff --git a/rustfs/src/app/storage_api.rs b/rustfs/src/app/storage_api.rs index 6c810e213..da7af7c86 100644 --- a/rustfs/src/app/storage_api.rs +++ b/rustfs/src/app/storage_api.rs @@ -27,12 +27,16 @@ pub(crate) fn EndpointServerPools( /// S3 wire types for app-layer modules, funneled here so new files stay off /// the direct s3s surface (s3s footprint ratchet, `scripts/check_s3s_footprint.sh`). pub(crate) mod s3 { + #[cfg(test)] + pub(crate) use s3s::dto::ListObjectsInput; #[cfg(test)] pub(crate) use s3s::dto::{ BucketVersioningStatus, DeleteMarkerReplication, DeleteMarkerReplicationStatus, Destination, ListObjectsV2Input, ListObjectsV2Output, ReplicationConfiguration, ReplicationRule, ReplicationRuleFilter, ReplicationRuleStatus, ServerSideEncryptionByDefault, ServerSideEncryptionConfiguration, ServerSideEncryptionRule, Tag, VersioningConfiguration, }; + #[cfg(test)] + pub(crate) use s3s::xml::{Serialize as XmlSerialize, Serializer as XmlSerializer}; pub(crate) use s3s::{S3Error, S3ErrorCode, S3Result}; #[cfg(test)] pub(crate) use s3s::{S3Request, S3Response}; diff --git a/rustfs/src/on_demand_migration/backfill.rs b/rustfs/src/on_demand_migration/backfill.rs index 8cc384ec8..f2d14b61d 100644 --- a/rustfs/src/on_demand_migration/backfill.rs +++ b/rustfs/src/on_demand_migration/backfill.rs @@ -471,7 +471,7 @@ impl BackfillContext for BucketBackfillContext { async fn config_updated_at(&self) -> Result, StorageError> { Ok( - super::config::decode_stored_config(get_on_demand_migration_config_in(&self.api.ctx, self.state.bucket()).await?)? + super::config::decode_stored_config(get_on_demand_migration_config_in(&self.api, self.state.bucket()).await?)? .map(|(_, updated_at)| updated_at), ) } @@ -1472,7 +1472,7 @@ impl Job { /// Spawns [`run_backfill_recovery_loop`] on the store's shutdown token; /// `false` (nothing spawned) when the store has no background token. pub fn spawn_backfill_recovery_loop(runner: Arc) -> bool { - let Some(cancel) = runner.api.ctx.background_cancel_token() else { + let Some(cancel) = runner.api.background_cancel_token() else { return false; }; tokio::spawn(run_backfill_recovery_loop(runner, cancel));